From 5e9396033a5a5d5f350bfa675d12770c972ee991 Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 11 Sep 2026 08:58:23 -0600 Subject: [PATCH 1/4] feat(types): derive generic collection element types --- crates/types/src/pre_asap/query_expr.rs | 5 +- crates/types/src/pre_asap/scalar_signature.rs | 123 ++++++++++++++++++ 2 files changed, 127 insertions(+), 1 deletion(-) diff --git a/crates/types/src/pre_asap/query_expr.rs b/crates/types/src/pre_asap/query_expr.rs index 848aed95..dbc14b48 100644 --- a/crates/types/src/pre_asap/query_expr.rs +++ b/crates/types/src/pre_asap/query_expr.rs @@ -1698,7 +1698,10 @@ fn infer_expr_type( (to.clone(), *try_cast || nullable) } QueryExpr::FunctionCall { name, args } => { - if name == "asap_struct_field" { + if name == "asap_element_access" { + super::scalar_signature::element_access_type(args, schema) + .map_err(QueryExprError::InvalidScalarSignature)? + } else if name == "asap_struct_field" { super::scalar_signature::struct_field_type(args, schema) .map_err(QueryExprError::InvalidScalarSignature)? } else if let Some(function) = diff --git a/crates/types/src/pre_asap/scalar_signature.rs b/crates/types/src/pre_asap/scalar_signature.rs index 97a69af8..7d61f624 100644 --- a/crates/types/src/pre_asap/scalar_signature.rs +++ b/crates/types/src/pre_asap/scalar_signature.rs @@ -379,3 +379,126 @@ mod struct_field_tests { .is_err()); } } + +/// Canonical element lookup over a declared Map or List. Map lookup retains its +/// existing key/default contract. List lookup is one-based, supports negative +/// indices, and returns the declared element default when a dynamic index is +/// out of range. Literal zero is conservatively rejected because native array +/// behavior depends on whether the input array is constant. Nullable containers +/// are unsupported; nullable indices produce nullable results. +pub fn element_access_type( + args: &[super::QueryExpr], + schema: &super::Schema, +) -> Result<(DataType, bool), String> { + use super::{QueryExpr, ScalarValue}; + let [input, index] = args else { + return Err("element access requires a collection and index".into()); + }; + let source = input.scalar_type(schema).map_err(|e| e.to_string())?; + let key = index.scalar_type(schema).map_err(|e| e.to_string())?; + match &source.0 { + DataType::Map { .. } => MapScalarFunction::Access.output_type(&[source, key]), + DataType::List { element } => { + if source.1 { + return Err("nullable List container access is unsupported".into()); + } + if !matches!(key.0, DataType::Int64 | DataType::Null) { + return Err("List index must have integer type".into()); + } + if matches!(index, QueryExpr::Literal(ScalarValue::Int64(0))) { + return Err( + "literal zero List index is unsupported without constant-array proof".into(), + ); + } + Ok(( + element.dtype.clone(), + element.nullable || key.1 || key.0 == DataType::Null, + )) + } + _ => Err("element access requires a Map or List".into()), + } +} + +#[cfg(test)] +mod element_access_tests { + use super::*; + use crate::pre_asap::{Column, QueryExpr, ScalarValue, Schema}; + fn access(index: QueryExpr) -> QueryExpr { + QueryExpr::FunctionCall { + name: "asap_element_access".into(), + args: vec![QueryExpr::Column(0), index], + } + } + #[test] + fn list_index_preserves_nested_element_metadata() { + let element = DataType::Struct { + fields: vec![ + Column::new("ts", DataType::Int64, false), + Column::new("value", DataType::Float64, true), + ], + }; + let schema = Schema::new(vec![ + Column::new( + "samples", + DataType::List { + element: Box::new(Column::new("item", element.clone(), false)), + }, + false, + ), + Column::new("i", DataType::Int64, true), + ]); + for index in [1, -1, 100] { + assert_eq!( + access(QueryExpr::Literal(ScalarValue::Int64(index))) + .scalar_type(&schema) + .unwrap(), + (element.clone(), false) + ); + } + assert_eq!( + access(QueryExpr::Column(1)).scalar_type(&schema).unwrap(), + (element.clone(), true) + ); + assert!(access(QueryExpr::Literal(ScalarValue::Int64(0))) + .scalar_type(&schema) + .is_err()); + assert!(access(QueryExpr::Literal(ScalarValue::Float64(1.0))) + .scalar_type(&schema) + .is_err()); + let nested = QueryExpr::FunctionCall { + name: "asap_struct_field".into(), + args: vec![ + access(QueryExpr::Literal(ScalarValue::Int64(1))), + QueryExpr::Literal(ScalarValue::Int64(2)), + ], + }; + assert_eq!( + nested.scalar_type(&schema).unwrap(), + (DataType::Float64, true) + ); + let roundtrip: QueryExpr = + serde_json::from_value(serde_json::to_value(&nested).unwrap()).unwrap(); + assert_eq!(roundtrip, nested); + } + #[test] + fn generic_map_lookup_reuses_legacy_signature() { + let schema = Schema::new(vec![Column::new( + "m", + DataType::Map { + key: Box::new(DataType::Utf8), + value: Box::new(DataType::Int64), + value_nullable: false, + }, + false, + )]); + let key = QueryExpr::Literal(ScalarValue::Utf8("k".into())); + let legacy = QueryExpr::FunctionCall { + name: "asap_map_access".into(), + args: vec![QueryExpr::Column(0), key.clone()], + }; + assert_eq!( + access(key).scalar_type(&schema).unwrap(), + legacy.scalar_type(&schema).unwrap() + ); + } +} From 4d4262357f4d4e21bb5e318d1a40d374079fed29 Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 11 Sep 2026 09:21:05 -0600 Subject: [PATCH 2/4] feat(sql): lower typed List element access --- crates/frontend-sql/src/sql/expr.rs | 2 +- crates/frontend-sql/src/sql/map_planning.rs | 52 +++++++++++++++++---- crates/frontend-sql/tests/sql_lowering.rs | 49 +++++++++++++++++++ 3 files changed, 93 insertions(+), 10 deletions(-) diff --git a/crates/frontend-sql/src/sql/expr.rs b/crates/frontend-sql/src/sql/expr.rs index 029842d5..cd57079f 100644 --- a/crates/frontend-sql/src/sql/expr.rs +++ b/crates/frontend-sql/src/sql/expr.rs @@ -203,7 +203,7 @@ pub(super) fn df_expr_to_unresolved(expr: &Expr) -> Result, _> = sf.args.iter().map(df_expr_to_unresolved).collect(); Ok(Unresolved::FunctionCall { name: if sf.func.name().eq_ignore_ascii_case("arrayelement") { - "asap_map_access".into() + "asap_element_access".into() } else { sf.func.name().to_string() }, diff --git a/crates/frontend-sql/src/sql/map_planning.rs b/crates/frontend-sql/src/sql/map_planning.rs index be99d4a8..7387d044 100644 --- a/crates/frontend-sql/src/sql/map_planning.rs +++ b/crates/frontend-sql/src/sql/map_planning.rs @@ -1,7 +1,8 @@ //! DataFusion planning adapters. Types come from the canonical signature rules; //! physical evaluation deliberately remains the query engine's responsibility. -use super::types::{arrow_to_dtype, dtype_to_arrow}; -use asap_types::pre_asap::scalar_signature::MapScalarFunction; +use super::types::{arrow_to_dtype, dtype_to_arrow, scalar_value_to_asap}; +use asap_types::pre_asap::scalar_signature::{element_access_type, MapScalarFunction}; +use asap_types::pre_asap::{Column, QueryExpr, Schema}; use datafusion::arrow::datatypes::DataType; use datafusion::common::{DataFusionError, ExprSchema, Result}; use datafusion::logical_expr::{ @@ -37,7 +38,12 @@ struct MapPlanningFunction { signature: Signature, } impl MapPlanningFunction { - fn output(&self, args: &[DataType], nullable: &[bool]) -> Result<(DataType, bool)> { + fn output( + &self, + args: &[DataType], + nullable: &[bool], + expressions: Option<&[Expr]>, + ) -> Result<(DataType, bool)> { let inputs = args .iter() .zip(nullable) @@ -47,10 +53,36 @@ impl MapPlanningFunction { .map_err(|e| DataFusionError::Plan(e.to_string())) }) .collect::>>()?; - let (dtype, nullable) = self - .function - .output_type(&inputs) - .map_err(DataFusionError::Plan)?; + let (dtype, nullable) = if self.name == "arrayelement" { + // DataFusion asks for argument-dependent types before canonical + // expression binding. Reuse the shared resolver over typed argument + // slots; final canonical binding also validates literal selectors. + let schema = Schema::new( + inputs + .into_iter() + .enumerate() + .map(|(index, (dtype, nullable))| { + Column::new(format!("argument_{index}"), dtype, nullable) + }) + .collect(), + ); + let args = (0..schema.columns.len()) + .map(|index| { + if let Some(Expr::Literal(value)) = expressions.and_then(|args| args.get(index)) + { + scalar_value_to_asap(value) + .map(QueryExpr::Literal) + .map_err(|error| DataFusionError::Plan(error.to_string())) + } else { + Ok(QueryExpr::Column(index)) + } + }) + .collect::>>()?; + element_access_type(&args, &schema) + } else { + self.function.output_type(&inputs) + } + .map_err(DataFusionError::Plan)?; Ok((dtype_to_arrow(&dtype), nullable)) } } @@ -71,6 +103,7 @@ impl ScalarUDFImpl for MapPlanningFunction { .iter() .map(|dtype| *dtype == DataType::Null) .collect::>(), + None, ) .map(|output| output.0) } @@ -84,7 +117,8 @@ impl ScalarUDFImpl for MapPlanningFunction { .iter() .map(|arg| arg.nullable(schema)) .collect::>>()?; - self.output(types, &nullable).map(|output| output.0) + self.output(types, &nullable, Some(args)) + .map(|output| output.0) } fn is_nullable(&self, args: &[Expr], schema: &dyn ExprSchema) -> bool { let types = args @@ -97,7 +131,7 @@ impl ScalarUDFImpl for MapPlanningFunction { .collect::>>(); match (types, nullable) { (Ok(types), Ok(nullable)) => self - .output(&types, &nullable) + .output(&types, &nullable, Some(args)) .map(|out| out.1) .unwrap_or(true), _ => true, diff --git a/crates/frontend-sql/tests/sql_lowering.rs b/crates/frontend-sql/tests/sql_lowering.rs index bb9d4318..556c0e5b 100644 --- a/crates/frontend-sql/tests/sql_lowering.rs +++ b/crates/frontend-sql/tests/sql_lowering.rs @@ -2344,3 +2344,52 @@ async fn arg_selector_result_schema_tracks_selected_argument() { assert_eq!(schema.columns[0].nullable, nullable); } } + +#[tokio::test] +async fn clickhouse_list_element_uses_canonical_typed_access() { + let catalog = SqlCatalog::new().with_table( + "t", + Schema::new(vec![ + Column::new( + "samples", + DataType::List { + element: Box::new(Column::new("item", DataType::Int64, false)), + }, + false, + ), + Column::new("index", DataType::Int64, true), + ]), + ); + for (sql, nullable) in [ + ("SELECT samples[1] AS selected FROM t", false), + ("SELECT arrayElement(samples, -1) AS selected FROM t", false), + ("SELECT samples[index] AS selected FROM t", true), + ] { + let query = lower_sql_dialect( + sql, + &catalog, + SqlDialect::ClickhouseSQL, + AccuracyTarget::Exact, + ) + .await + .unwrap(); + let output = query.output_schema().unwrap(); + assert_eq!(output.columns[0].dtype, DataType::Int64); + assert_eq!(output.columns[0].nullable, nullable); + let serialized = serde_json::to_string(&query).unwrap(); + assert!(serialized.contains("asap_element_access"), "{serialized}"); + } + for sql in ["SELECT samples[0] FROM t", "SELECT samples['bad'] FROM t"] { + assert!( + lower_sql_dialect( + sql, + &catalog, + SqlDialect::ClickhouseSQL, + AccuracyTarget::Exact + ) + .await + .is_err(), + "{sql}" + ); + } +} From 73c121b0488197f1622629a749ab3ac0d7a02584 Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 11 Sep 2026 09:21:59 -0600 Subject: [PATCH 3/4] refactor(sql): name shared collection planning adapter --- .../sql/{map_planning.rs => collection_planning.rs} | 12 ++++++------ crates/frontend-sql/src/sql/mod.rs | 4 ++-- 2 files changed, 8 insertions(+), 8 deletions(-) rename crates/frontend-sql/src/sql/{map_planning.rs => collection_planning.rs} (93%) diff --git a/crates/frontend-sql/src/sql/map_planning.rs b/crates/frontend-sql/src/sql/collection_planning.rs similarity index 93% rename from crates/frontend-sql/src/sql/map_planning.rs rename to crates/frontend-sql/src/sql/collection_planning.rs index 7387d044..d4582afd 100644 --- a/crates/frontend-sql/src/sql/map_planning.rs +++ b/crates/frontend-sql/src/sql/collection_planning.rs @@ -17,7 +17,7 @@ pub(super) fn register(context: &SessionContext) { ("mapconcat", MapScalarFunction::Concat), ("arrayelement", MapScalarFunction::Access), ] { - context.register_udf(ScalarUDF::from(MapPlanningFunction { + context.register_udf(ScalarUDF::from(CollectionPlanningFunction { name, function, signature: match function { @@ -32,12 +32,12 @@ pub(super) fn register(context: &SessionContext) { } } #[derive(Debug)] -struct MapPlanningFunction { +struct CollectionPlanningFunction { name: &'static str, function: MapScalarFunction, signature: Signature, } -impl MapPlanningFunction { +impl CollectionPlanningFunction { fn output( &self, args: &[DataType], @@ -86,7 +86,7 @@ impl MapPlanningFunction { Ok((dtype_to_arrow(&dtype), nullable)) } } -impl ScalarUDFImpl for MapPlanningFunction { +impl ScalarUDFImpl for CollectionPlanningFunction { fn as_any(&self) -> &dyn std::any::Any { self } @@ -138,7 +138,7 @@ impl ScalarUDFImpl for MapPlanningFunction { } } fn invoke_batch(&self, _args: &[ColumnarValue], _number_rows: usize) -> Result { - Err(DataFusionError::NotImplemented("map planning adapter cannot execute; use a capable query engine or external exact subtree".into())) + Err(DataFusionError::NotImplemented("collection planning adapter cannot execute; use a capable query engine or external exact subtree".into())) } } @@ -147,7 +147,7 @@ mod tests { use super::*; #[test] fn planning_adapter_explicitly_refuses_physical_execution() { - let adapter = MapPlanningFunction { + let adapter = CollectionPlanningFunction { name: "map", function: MapScalarFunction::Construct, signature: Signature::any(0, Volatility::Immutable), diff --git a/crates/frontend-sql/src/sql/mod.rs b/crates/frontend-sql/src/sql/mod.rs index 4cf9fad5..36831cc2 100644 --- a/crates/frontend-sql/src/sql/mod.rs +++ b/crates/frontend-sql/src/sql/mod.rs @@ -69,7 +69,7 @@ use crate::error::SqlError as LoweringError; mod clickhouse_ast; mod expr; -mod map_planning; +mod collection_planning; mod types; pub use types::SqlCatalog; @@ -201,7 +201,7 @@ impl<'a> SqlLowerer<'a> { let config = SessionConfig::new().set_str("datafusion.sql_parser.dialect", dialect_name); let ctx = SessionContext::new_with_config(config); if matches!(self.dialect, SqlDialect::ClickhouseSQL) { - map_planning::register(&ctx); + collection_planning::register(&ctx); } // A catalog key like "bgp.bgp_updates" schema-qualifies the table // (e.g. a ClickHouse database name). DataFusion requires the parent From cf1f9acf0e09ce31b8319e674fb2508185a31148 Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 11 Sep 2026 09:23:20 -0600 Subject: [PATCH 4/4] style(sql): order collection module declaration --- crates/frontend-sql/src/sql/mod.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/crates/frontend-sql/src/sql/mod.rs b/crates/frontend-sql/src/sql/mod.rs index 36831cc2..dc9116eb 100644 --- a/crates/frontend-sql/src/sql/mod.rs +++ b/crates/frontend-sql/src/sql/mod.rs @@ -68,8 +68,8 @@ use asap_types::workload::SqlDialect; use crate::error::SqlError as LoweringError; mod clickhouse_ast; -mod expr; mod collection_planning; +mod expr; mod types; pub use types::SqlCatalog;