Skip to main content

query/datafusion/
json_expr_planner.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::sync::{Arc, LazyLock};
16
17use arrow_schema::Field;
18use arrow_schema::extension::ExtensionType;
19use common_function::scalars::json::json_get::JsonGetWithType;
20use common_function::scalars::udf::create_udf;
21use datafusion_common::arrow::datatypes::DataType;
22use datafusion_common::{
23    Column, DFSchema, DataFusionError, Result, ScalarValue, TableReference, plan_datafusion_err,
24};
25use datafusion_expr::expr::{BinaryExpr, ScalarFunction};
26use datafusion_expr::planner::{
27    ExprPlanner, PlannerResult, RawAggregateExpr, RawBinaryExpr, RawFieldAccessExpr, RawScalarExpr,
28    RawWindowExpr,
29};
30use datafusion_expr::type_coercion::functions::{UDFCoercionExt, fields_with_udf};
31use datafusion_expr::{
32    Expr, ExprSchemable, GetFieldAccess, Operator, ScalarUDF, WindowFunctionDefinition,
33};
34use datatypes::extension::json::{
35    Json2ExtensionType, is_json2_extension_type, parse_legacy_json2_settings,
36};
37use datatypes::types::json_type::JsonNativeType;
38use sqlparser::ast::BinaryOperator;
39
40/// Rewrites JSON-aware SQL expressions into DataFusion expressions.
41///
42/// This planner handles the following cases:
43/// - Rewrites compound identifiers on JSON extension columns into `json_get` function.
44///   For example, `select a.b.c` => `select json_get(a, '$.b.c')`.
45/// - Extends a JSON path with list indexes and fields following an index.
46///   For example, `select a.b[0].c` => `select json_get(a, '$.b[0].c')`.
47/// - Pushes an "expected type" argument into the `json_get` function when it participates in a
48///   binary operator. So that `json_get` knows the wanted data type when dealing with variant
49///   JSON values.
50///   For example, `select json_get(a, "b.c") + 1` => `select json_get(a, "b.c", NULL::Int64) + 1`.
51/// - Infers the expected type from scalar, aggregate, and window function signatures.
52///   For example, `select abs(a.b.c)` => `select abs(json_get(a, '$.b.c', NULL::Float64))`.
53#[derive(Debug)]
54pub(crate) struct JsonExprPlanner;
55
56impl ExprPlanner for JsonExprPlanner {
57    fn plan_binary_op(
58        &self,
59        expr: RawBinaryExpr,
60        schema: &DFSchema,
61    ) -> Result<PlannerResult<RawBinaryExpr>> {
62        let RawBinaryExpr {
63            op,
64            mut left,
65            mut right,
66        } = expr;
67
68        if !is_untyped_json_get(&left) && !is_untyped_json_get(&right) {
69            return Ok(PlannerResult::Original(RawBinaryExpr { op, left, right }));
70        }
71
72        let Some(expr_op) = parse_sql_op(&op) else {
73            return Ok(PlannerResult::Original(RawBinaryExpr { op, left, right }));
74        };
75
76        let left_type = left.get_type(schema)?;
77        let right_type = right.get_type(schema)?;
78        let left_changed = push_json_get_type_arg(&mut left, &right_type)?;
79        let right_changed = push_json_get_type_arg(&mut right, &left_type)?;
80        if left_changed || right_changed {
81            Ok(PlannerResult::Planned(Expr::BinaryExpr(BinaryExpr::new(
82                Box::new(left),
83                expr_op,
84                Box::new(right),
85            ))))
86        } else {
87            Ok(PlannerResult::Original(RawBinaryExpr { op, left, right }))
88        }
89    }
90
91    /// Extends the path of an untyped `json_get` with one field access.
92    ///
93    /// For `j.o.l[1].inner.l[2]`, `plan_compound_identifier` first produces
94    /// `json_get(j, '$.o.l')`. DataFusion then calls this method successively
95    /// with a list index, two named fields, and another list index, producing
96    /// the final path `$.o.l[1].inner.l[2]`.
97    fn plan_field_access(
98        &self,
99        mut expr: RawFieldAccessExpr,
100        _schema: &DFSchema,
101    ) -> Result<PlannerResult<RawFieldAccessExpr>> {
102        // See `normalize_field_access_after_subscript` for the reason why we construct the
103        // "suffix" like this.
104        let suffix = match &expr.field_access {
105            GetFieldAccess::ListIndex { key } => {
106                // DataFusion parses ordinary integer literals within the i64 range as Int64.
107                let Expr::Literal(ScalarValue::Int64(Some(index)), _) = key.as_ref() else {
108                    return Ok(PlannerResult::Original(expr));
109                };
110                format!("[{index}]")
111            }
112            GetFieldAccess::NamedStructField { name } => {
113                let Some(name) = name.try_as_str().flatten() else {
114                    return Ok(PlannerResult::Original(expr));
115                };
116                format!(".{}", json_path_field(name)?)
117            }
118            GetFieldAccess::ListRange { .. } => return Ok(PlannerResult::Original(expr)),
119        };
120        let Some(json_get) = extract_untyped_json_get(&mut expr.expr) else {
121            return Ok(PlannerResult::Original(expr));
122        };
123        let Some(Expr::Literal(ScalarValue::Utf8(Some(path)), _)) = json_get.args.get_mut(1) else {
124            return Ok(PlannerResult::Original(expr));
125        };
126
127        path.push_str(&suffix);
128        Ok(PlannerResult::Planned(expr.expr))
129    }
130
131    fn plan_compound_identifier(
132        &self,
133        field: &Field,
134        qualifier: Option<&TableReference>,
135        nested_names: &[String],
136    ) -> Result<PlannerResult<Vec<Expr>>> {
137        if !is_json2_extension_type(field) {
138            return Ok(PlannerResult::Original(Vec::new()));
139        }
140
141        static JSON_GET_UDF: LazyLock<Arc<ScalarUDF>> =
142            LazyLock::new(|| Arc::new(create_udf(Arc::new(JsonGetWithType::default()))));
143
144        let json_get = JSON_GET_UDF.clone();
145        let mut path = "$".to_string();
146        for name in nested_names {
147            path.push('.');
148            path.push_str(&json_path_field(name)?);
149        }
150
151        let mut args = vec![
152            Expr::Column(Column::from((qualifier, field))),
153            Expr::Literal(ScalarValue::Utf8(Some(path)), None),
154        ];
155        if let Some(json_type) = json_type_hint(field, nested_names)? {
156            args.push(Expr::Literal(
157                ScalarValue::try_new_null(&json_type.as_arrow_type())?,
158                None,
159            ));
160        }
161
162        Ok(PlannerResult::Planned(Expr::ScalarFunction(
163            ScalarFunction::new_udf(json_get, args),
164        )))
165    }
166
167    /// Rewrites JSON2 arguments without taking over the final function planning.
168    ///
169    /// `Original` carries the possibly modified raw expression to subsequent planners and then
170    /// DataFusion's default function construction. Returning `Planned` would short-circuit both.
171    fn plan_scalar(&self, mut expr: RawScalarExpr) -> Result<PlannerResult<RawScalarExpr>> {
172        push_function_arg_types(expr.func.as_ref(), &mut expr.args)?;
173        Ok(PlannerResult::Original(expr))
174    }
175
176    /// Rewrites JSON2 arguments while preserving subsequent aggregate planning.
177    fn plan_aggregate(
178        &self,
179        mut expr: RawAggregateExpr,
180    ) -> Result<PlannerResult<RawAggregateExpr>> {
181        push_function_arg_types(expr.func.as_ref(), &mut expr.args)?;
182        Ok(PlannerResult::Original(expr))
183    }
184
185    /// Rewrites JSON2 arguments while preserving subsequent window planning.
186    fn plan_window(&self, mut expr: RawWindowExpr) -> Result<PlannerResult<RawWindowExpr>> {
187        match &expr.func_def {
188            WindowFunctionDefinition::AggregateUDF(func) => {
189                push_function_arg_types(func.as_ref(), &mut expr.args)?;
190            }
191            WindowFunctionDefinition::WindowUDF(func) => {
192                push_function_arg_types(func.as_ref(), &mut expr.args)?;
193            }
194        }
195        Ok(PlannerResult::Original(expr))
196    }
197}
198
199/// Returns the configured native type for an exact JSON2 object path.
200fn json_type_hint(field: &Field, path: &[String]) -> Result<Option<JsonNativeType>> {
201    let settings = if field.extension_type_name() == Some(Json2ExtensionType::NAME) {
202        let extension = field
203            .try_extension_type::<Json2ExtensionType>()
204            .map_err(|e| plan_datafusion_err!("invalid JSON2 extension metadata: {e}"))?;
205        Some(extension.metadata().json_settings().clone())
206    } else {
207        parse_legacy_json2_settings(field.metadata())
208            .map_err(|e| plan_datafusion_err!("invalid JSON2 extension metadata: {e}"))?
209    };
210
211    Ok(settings.and_then(|settings| {
212        settings
213            .type_hints()
214            .iter()
215            .find(|hint| hint.path == path)
216            .map(|hint| JsonNativeType::from(&hint.data_type))
217    }))
218}
219
220/// Quotes field names containing JSONPath punctuation, preserving literal keys.
221fn json_path_field(name: &str) -> Result<String> {
222    if !name.is_empty() && name.chars().all(|c| c.is_alphanumeric() || c == '_') {
223        Ok(name.to_string())
224    } else {
225        // Reuse serde_json's string escaping.
226        serde_json::to_string(name).map_err(|e| DataFusionError::External(Box::new(e)))
227    }
228}
229
230enum JsonGetTypeResolution {
231    Fallback,
232    Typed(Vec<(usize, DataType)>),
233}
234
235/// Infers static output types for untyped `json_get` arguments from a function signature.
236///
237/// DataFusion requires every expression to have one Arrow data type during planning. A JSON path
238/// may contain heterogeneous values across rows, but it cannot expose those values as different
239/// Arrow types in one result column. Preserving their runtime types would require a single
240/// Variant-like data type and Variant-aware functions instead. Maybe we can wait for
241/// https://github.com/apache/datafusion/issues/16116
242///
243/// This helper uses the function's coercion rules to select a supported output type, then appends
244/// a typed NULL argument to each relevant `json_get`. The typed argument makes `json_get` project
245/// compatible JSON values to that type and return NULL for incompatible values. Functions that
246/// accept json_get's default `Utf8View` output keep the two-argument form so later rewrites can
247/// still push down an outer cast.
248fn push_function_arg_types<F>(func: &F, args: &mut [Expr]) -> Result<()>
249where
250    F: UDFCoercionExt,
251{
252    if !args.iter().any(is_untyped_json_get) {
253        return Ok(());
254    }
255
256    let fields = args.iter().map(function_arg_field).collect::<Vec<_>>();
257    match infer_json_get_types(func, args, &fields) {
258        JsonGetTypeResolution::Fallback => {
259            let Some(data_type) = fallback_json_get_type(func, args, &fields) else {
260                return Ok(());
261            };
262            for arg in args.iter_mut() {
263                if is_untyped_json_get(arg) {
264                    let _ = push_json_get_type_arg(arg, &data_type)?;
265                }
266            }
267        }
268        JsonGetTypeResolution::Typed(types) => {
269            for (index, data_type) in types {
270                let _ = push_json_get_type_arg(&mut args[index], &data_type)?;
271            }
272        }
273    }
274    Ok(())
275}
276
277fn infer_json_get_types<F>(func: &F, args: &[Expr], fields: &[Arc<Field>]) -> JsonGetTypeResolution
278where
279    F: UDFCoercionExt,
280{
281    // Only untyped json_get arguments use Null placeholders; preserve every other known argument
282    // type. fields_with_udf performs contextual coercion rather than reverse inference from a
283    // signature alone. Numeric signatures may preserve all-Null inputs, while Comparable
284    // signatures may default them to Utf8. For example, retaining the Float64 peer in
285    // coalesce(json_get(...), 1.0) lets DataFusion resolve json_get to Float64 instead of Utf8.
286    //
287    // This is a best-effort probe: a failure does not mean the actual function call is invalid, so
288    // try concrete JSON types before leaving final validation to DataFusion's default planner.
289    let Ok(coerced) = fields_with_udf(fields, func) else {
290        return JsonGetTypeResolution::Fallback;
291    };
292
293    let mut inferred_types = Vec::with_capacity(coerced.len());
294    for (index, (arg, field)) in args.iter().zip(coerced).enumerate() {
295        if !is_untyped_json_get(arg) || field.data_type().is_null() {
296            continue;
297        }
298        let Some(data_type) = json_get_output_type(field.data_type()) else {
299            return JsonGetTypeResolution::Fallback;
300        };
301        inferred_types.push((index, data_type));
302    }
303    if inferred_types.is_empty() {
304        JsonGetTypeResolution::Fallback
305    } else {
306        JsonGetTypeResolution::Typed(inferred_types)
307    }
308}
309
310fn fallback_json_get_type<F>(func: &F, args: &[Expr], fields: &[Arc<Field>]) -> Option<DataType>
311where
312    F: UDFCoercionExt,
313{
314    // Prefer json_get's default Utf8View type. If the function rejects strings but accepts numeric
315    // values, prefer Float64 so both integers and fractions remain usable.
316    let mut candidate_fields = fields.to_vec();
317    for data_type in [
318        DataType::Utf8View,
319        DataType::Float64,
320        DataType::Int64,
321        DataType::Boolean,
322    ] {
323        for (index, arg) in args.iter().enumerate() {
324            if is_untyped_json_get(arg) {
325                candidate_fields[index] = Arc::new(
326                    fields[index]
327                        .as_ref()
328                        .clone()
329                        .with_data_type(data_type.clone()),
330                );
331            }
332        }
333        if fields_with_udf(&candidate_fields, func).is_ok() {
334            return Some(data_type);
335        }
336    }
337    None
338}
339
340fn function_arg_field(expr: &Expr) -> Arc<Field> {
341    let data_type = if is_untyped_json_get(expr) {
342        DataType::Null
343    } else if let Some(data_type) = extract_json_get_type(expr) {
344        data_type
345    } else {
346        // Treat unresolved expressions as untyped NULL. This lets signatures such as `power`
347        // infer a JSON type, while functions such as `coalesce` can leave it untyped for default
348        // planning. This is only best-effort: overloaded or user-defined functions may select a
349        // different signature for NULL than for the expression's actual type.
350        // TODO(LFC): Use the input schema once DataFusion passes it to ExprPlanner::plan_*().
351        expr.get_type(&DFSchema::empty()).unwrap_or(DataType::Null)
352    };
353    Arc::new(Field::new("", data_type, true))
354}
355
356fn json_get_output_type(data_type: &DataType) -> Option<DataType> {
357    let output_type = match data_type {
358        DataType::Boolean => DataType::Boolean,
359        data_type if data_type.is_integer() => DataType::Int64,
360        data_type if data_type.is_floating() => DataType::Float64,
361        DataType::Decimal128(_, _) | DataType::Decimal256(_, _) => DataType::Float64,
362        data_type if data_type.is_string() => DataType::Utf8View,
363        _ => return None,
364    };
365    Some(output_type)
366}
367
368macro_rules! is_untyped_json_get_func {
369    ($func:expr) => {
370        $func
371            .func
372            .name()
373            .eq_ignore_ascii_case(JsonGetWithType::NAME)
374            && $func.args.len() == 2
375    };
376}
377
378macro_rules! is_typed_json_get_func {
379    ($func:expr) => {
380        $func
381            .func
382            .name()
383            .eq_ignore_ascii_case(JsonGetWithType::NAME)
384            && $func.args.len() == 3
385    };
386}
387
388fn extract_untyped_json_get(expr: &mut Expr) -> Option<&mut ScalarFunction> {
389    match expr {
390        Expr::ScalarFunction(f) if is_untyped_json_get_func!(f) => Some(f),
391        _ => None,
392    }
393}
394
395fn extract_json_get_type(expr: &Expr) -> Option<DataType> {
396    match expr {
397        Expr::ScalarFunction(f) if is_typed_json_get_func!(f) => f
398            .args
399            .get(2)
400            .and_then(|x| x.as_literal())
401            .map(|x| x.data_type()),
402        _ => None,
403    }
404}
405
406fn is_untyped_json_get(expr: &Expr) -> bool {
407    matches!(
408        expr,
409        Expr::ScalarFunction(f) if is_untyped_json_get_func!(f)
410    )
411}
412
413fn push_json_get_type_arg(expr: &mut Expr, data_type: &DataType) -> Result<bool> {
414    let Some(json_get) = extract_untyped_json_get(expr) else {
415        return Ok(false);
416    };
417
418    // The two-argument form already returns Utf8View. Keep it so JsonGetRewriter can still absorb
419    // a cast added by subsequent function coercion.
420    if data_type.is_string() {
421        return Ok(false);
422    }
423    let with_type = ScalarValue::try_new_null(data_type).map(|x| Expr::Literal(x, None))?;
424    json_get.args.push(with_type);
425    Ok(true)
426}
427
428fn parse_sql_op(op: &BinaryOperator) -> Option<Operator> {
429    match *op {
430        BinaryOperator::Plus => Some(Operator::Plus),
431        BinaryOperator::Minus => Some(Operator::Minus),
432        BinaryOperator::Multiply => Some(Operator::Multiply),
433        BinaryOperator::Divide => Some(Operator::Divide),
434        BinaryOperator::Modulo => Some(Operator::Modulo),
435        BinaryOperator::Gt => Some(Operator::Gt),
436        BinaryOperator::GtEq => Some(Operator::GtEq),
437        BinaryOperator::Lt => Some(Operator::Lt),
438        BinaryOperator::LtEq => Some(Operator::LtEq),
439        BinaryOperator::Eq => Some(Operator::Eq),
440        BinaryOperator::NotEq => Some(Operator::NotEq),
441        BinaryOperator::And => Some(Operator::And),
442        BinaryOperator::Or => Some(Operator::Or),
443        BinaryOperator::BitwiseAnd => Some(Operator::BitwiseAnd),
444        BinaryOperator::BitwiseOr => Some(Operator::BitwiseOr),
445        BinaryOperator::BitwiseXor => Some(Operator::BitwiseXor),
446        _ => None,
447    }
448}
449
450#[cfg(test)]
451mod tests {
452    use arrow_schema::Fields;
453    use datafusion::functions_aggregate::count::count_udaf;
454    use datafusion::functions_aggregate::sum::sum_udaf;
455    use datafusion_expr::WindowFrame;
456    use datafusion_functions::core::coalesce;
457    use datafusion_functions::math::{abs, power};
458    use datatypes::extension::json::{Json2ExtensionType, JsonMetadata};
459    use datatypes::json::{JsonSettings, JsonTypeHint};
460    use datatypes::prelude::ConcreteDataType;
461
462    use super::*;
463
464    fn json_get_expr(base: Expr, path: &str) -> Expr {
465        let json_get = Arc::new(create_udf(Arc::new(JsonGetWithType::default())));
466        Expr::ScalarFunction(ScalarFunction::new_udf(
467            json_get,
468            vec![
469                base,
470                Expr::Literal(ScalarValue::Utf8(Some(path.to_string())), None),
471            ],
472        ))
473    }
474
475    #[test]
476    fn test_plan_binary_op() -> Result<()> {
477        let planner = JsonExprPlanner;
478        let schema = DFSchema::from_unqualified_fields(
479            Fields::from(vec![Field::new("value", DataType::Int64, true)]),
480            Default::default(),
481        )?;
482
483        let planned = planner.plan_binary_op(
484            RawBinaryExpr {
485                op: BinaryOperator::Eq,
486                left: json_get_expr(
487                    Expr::Literal(ScalarValue::Binary(Some(b"{\"a\": 1}".to_vec())), None),
488                    "a",
489                ),
490                right: Expr::Column(Column::new_unqualified("value")),
491            },
492            &schema,
493        )?;
494
495        match planned {
496            PlannerResult::Planned(Expr::BinaryExpr(expr)) => {
497                assert_eq!(expr.op, Operator::Eq);
498
499                match expr.left.as_ref() {
500                    Expr::ScalarFunction(func) => {
501                        assert_eq!(func.func.name(), JsonGetWithType::NAME);
502                        assert_eq!(func.args.len(), 3);
503                        assert_eq!(func.args[2], Expr::Literal(ScalarValue::Int64(None), None));
504                    }
505                    other => panic!("expected json_get on left side, got {other:?}"),
506                }
507
508                assert_eq!(
509                    expr.right.as_ref(),
510                    &Expr::Column(Column::new_unqualified("value"))
511                );
512            }
513            other => panic!("expected planned binary expression, got {other:?}"),
514        }
515
516        let original = planner.plan_binary_op(
517            RawBinaryExpr {
518                op: BinaryOperator::StringConcat,
519                left: Expr::Column(Column::new_unqualified("value")),
520                right: Expr::Literal(ScalarValue::Utf8(Some("x".to_string())), None),
521            },
522            &schema,
523        )?;
524
525        match original {
526            PlannerResult::Original(expr) => {
527                assert!(matches!(expr.op, BinaryOperator::StringConcat));
528                assert_eq!(expr.left, Expr::Column(Column::new_unqualified("value")));
529                assert_eq!(
530                    expr.right,
531                    Expr::Literal(ScalarValue::Utf8(Some("x".to_string())), None)
532                );
533            }
534            other => panic!(
535                "expected original expression for unsupported operator, got {:?}",
536                other,
537            ),
538        }
539
540        Ok(())
541    }
542
543    #[test]
544    fn test_plan_list_index() -> Result<()> {
545        let planner = JsonExprPlanner;
546        let planned = planner.plan_field_access(
547            RawFieldAccessExpr {
548                field_access: GetFieldAccess::ListIndex {
549                    key: Box::new(Expr::Literal(ScalarValue::Int64(Some(0)), None)),
550                },
551                expr: json_get_expr(Expr::Column(Column::new_unqualified("j")), "list"),
552            },
553            &DFSchema::empty(),
554        )?;
555        let PlannerResult::Planned(Expr::ScalarFunction(func)) = planned else {
556            unreachable!()
557        };
558        assert_eq!(func.func.name(), JsonGetWithType::NAME);
559        assert_eq!(func.args.len(), 2);
560        assert_eq!(
561            func.args[1],
562            Expr::Literal(ScalarValue::Utf8(Some("list[0]".to_string())), None)
563        );
564        Ok(())
565    }
566
567    #[test]
568    fn test_plan_field_after_list_index() -> Result<()> {
569        let planner = JsonExprPlanner;
570        let planned = planner.plan_field_access(
571            RawFieldAccessExpr {
572                field_access: GetFieldAccess::NamedStructField {
573                    name: ScalarValue::Utf8(Some("a.b".to_string())),
574                },
575                expr: json_get_expr(Expr::Column(Column::new_unqualified("j")), "list[0]"),
576            },
577            &DFSchema::empty(),
578        )?;
579        let PlannerResult::Planned(Expr::ScalarFunction(func)) = planned else {
580            unreachable!()
581        };
582        assert_eq!(
583            func.args[1],
584            Expr::Literal(
585                ScalarValue::Utf8(Some(r#"list[0]."a.b""#.to_string())),
586                None
587            )
588        );
589        Ok(())
590    }
591
592    #[test]
593    fn test_plan_compound_identifier() -> Result<()> {
594        let planner = JsonExprPlanner;
595        let qualifier = TableReference::bare("events");
596        let nested_names = vec!["payload".to_string(), "cpu".to_string()];
597
598        let planned = planner.plan_compound_identifier(
599            &Field::new("labels", DataType::Struct(Fields::empty()), true)
600                .with_extension_type(Json2ExtensionType::default()),
601            Some(&qualifier),
602            &nested_names,
603        )?;
604
605        match planned {
606            PlannerResult::Planned(Expr::ScalarFunction(func)) => {
607                assert_eq!(func.func.name(), JsonGetWithType::NAME);
608                assert_eq!(func.args.len(), 2);
609                assert_eq!(
610                    func.args[0],
611                    Expr::Column(Column::new(Some(qualifier.clone()), "labels"))
612                );
613                assert_eq!(
614                    func.args[1],
615                    Expr::Literal(ScalarValue::Utf8(Some("$.payload.cpu".to_string())), None)
616                );
617            }
618            other => panic!("expected json_get scalar function, got {other:?}"),
619        }
620
621        for (key, path) in [
622            ("http.status_code", r#"$."http.status_code""#),
623            ("a\"b", r#"$."a\"b""#),
624            ("a\\b", r#"$."a\\b""#),
625        ] {
626            let planned = planner.plan_compound_identifier(
627                &Field::new("labels", DataType::Struct(Fields::empty()), true)
628                    .with_extension_type(Json2ExtensionType::default()),
629                Some(&qualifier),
630                &[key.to_string()],
631            )?;
632            let PlannerResult::Planned(Expr::ScalarFunction(func)) = planned else {
633                unreachable!()
634            };
635            assert_eq!(
636                func.args[1],
637                Expr::Literal(ScalarValue::Utf8(Some(path.to_string())), None)
638            );
639        }
640
641        let PlannerResult::Planned(Expr::ScalarFunction(func)) = planner.plan_compound_identifier(
642            &Field::new("labels", DataType::Struct(Fields::empty()), true)
643                .with_extension_type(Json2ExtensionType::default()),
644            Some(&qualifier),
645            &["resource".to_string(), "http.status_code".to_string()],
646        )?
647        else {
648            unreachable!()
649        };
650        assert_eq!(
651            func.args[1],
652            Expr::Literal(
653                ScalarValue::Utf8(Some(r#"$.resource."http.status_code""#.to_string())),
654                None
655            )
656        );
657
658        let original = planner.plan_compound_identifier(
659            &Field::new("plain", DataType::Utf8, true),
660            Some(&qualifier),
661            &nested_names,
662        )?;
663
664        match original {
665            PlannerResult::Original(exprs) => assert!(exprs.is_empty()),
666            other => panic!(
667                "expected original empty result for non-json field, got {:?}",
668                other,
669            ),
670        }
671
672        Ok(())
673    }
674
675    #[test]
676    fn test_plan_compound_identifier_applies_json2_type_hint() -> Result<()> {
677        let settings = JsonSettings::try_new(
678            vec![JsonTypeHint {
679                path: vec!["payload".to_string(), "cpu".to_string()],
680                data_type: ConcreteDataType::int64_datatype(),
681                inverted_index: false,
682            }],
683            None,
684        )
685        .unwrap();
686        let field = Field::new("labels", DataType::Struct(Fields::empty()), true)
687            .with_extension_type(Json2ExtensionType::new(Arc::new(JsonMetadata::new(
688                settings,
689            ))));
690
691        let PlannerResult::Planned(Expr::ScalarFunction(func)) = JsonExprPlanner
692            .plan_compound_identifier(
693                &field,
694                Some(&TableReference::bare("events")),
695                &["payload".to_string(), "cpu".to_string()],
696            )?
697        else {
698            unreachable!()
699        };
700
701        assert_eq!(
702            Some(DataType::Int64),
703            extract_json_get_type(&Expr::ScalarFunction(func))
704        );
705        Ok(())
706    }
707
708    #[test]
709    fn test_plan_functions() -> Result<()> {
710        let planner = JsonExprPlanner;
711        let json_get = || json_get_expr(Expr::Column(Column::new_unqualified("j")), "a.b");
712
713        let PlannerResult::Original(scalar) = planner.plan_scalar(RawScalarExpr {
714            func: abs(),
715            args: vec![json_get()],
716        })?
717        else {
718            unreachable!();
719        };
720        assert_eq!(
721            Some(DataType::Float64),
722            extract_json_get_type(&scalar.args[0])
723        );
724
725        let PlannerResult::Original(scalar) = planner.plan_scalar(RawScalarExpr {
726            func: power(),
727            args: vec![
728                json_get(),
729                Expr::Column(Column::new_unqualified("exponent")),
730            ],
731        })?
732        else {
733            unreachable!();
734        };
735        assert_eq!(
736            Some(DataType::Float64),
737            extract_json_get_type(&scalar.args[0])
738        );
739
740        let PlannerResult::Original(aggregate) = planner.plan_aggregate(RawAggregateExpr {
741            func: sum_udaf(),
742            args: vec![json_get()],
743            distinct: false,
744            filter: None,
745            order_by: vec![],
746            null_treatment: None,
747        })?
748        else {
749            unreachable!();
750        };
751        assert_eq!(
752            Some(DataType::Float64),
753            extract_json_get_type(&aggregate.args[0])
754        );
755
756        let PlannerResult::Original(count) = planner.plan_aggregate(RawAggregateExpr {
757            func: count_udaf(),
758            args: vec![json_get()],
759            distinct: false,
760            filter: None,
761            order_by: vec![],
762            null_treatment: None,
763        })?
764        else {
765            unreachable!();
766        };
767        assert_eq!(None, extract_json_get_type(&count.args[0]));
768
769        let PlannerResult::Original(window) = planner.plan_window(RawWindowExpr {
770            func_def: WindowFunctionDefinition::AggregateUDF(sum_udaf()),
771            args: vec![json_get()],
772            partition_by: vec![],
773            order_by: vec![],
774            window_frame: WindowFrame::new(None),
775            filter: None,
776            null_treatment: None,
777            distinct: false,
778        })?
779        else {
780            unreachable!();
781        };
782        assert_eq!(
783            Some(DataType::Float64),
784            extract_json_get_type(&window.args[0])
785        );
786        Ok(())
787    }
788
789    #[test]
790    fn test_plan_function_with_mixed_json_get_types() -> Result<()> {
791        let planner = JsonExprPlanner;
792        let json_get = || json_get_expr(Expr::Column(Column::new_unqualified("j")), "a.b");
793        let mut typed = json_get();
794        push_json_get_type_arg(&mut typed, &DataType::Float64)?;
795
796        let PlannerResult::Original(scalar) = planner.plan_scalar(RawScalarExpr {
797            func: coalesce(),
798            args: vec![json_get(), typed],
799        })?
800        else {
801            unreachable!();
802        };
803        assert_eq!(
804            Some(DataType::Float64),
805            extract_json_get_type(&scalar.args[0])
806        );
807        assert_eq!(
808            Some(DataType::Float64),
809            extract_json_get_type(&scalar.args[1])
810        );
811        Ok(())
812    }
813}