From cf943d463f2ee431ee68dc3aae6e7d6fb08dd8db Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 11 Sep 2026 09:06:42 -0600 Subject: [PATCH 01/13] feat(clickhouse): execute typed array element access --- Cargo.lock | 10 +- control_plane/Cargo.toml | 8 +- .../src/query_plan/clickhouse_exact.rs | 24 ++- crates/asap_types/Cargo.toml | 2 +- data_plane/Cargo.toml | 6 +- .../relational_adapter.rs | 194 +++++++++++++++++- 6 files changed, 221 insertions(+), 23 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 5b68bc67..c2824200 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -373,7 +373,7 @@ dependencies = [ [[package]] name = "asap-aware-mapping" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=0402384e589df6e087d6d2d22b463ddc2eea0774#0402384e589df6e087d6d2d22b463ddc2eea0774" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=5e9396033a5a5d5f350bfa675d12770c972ee991#5e9396033a5a5d5f350bfa675d12770c972ee991" dependencies = [ "asap-types", "serde", @@ -384,7 +384,7 @@ dependencies = [ [[package]] name = "asap-frontend-promql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=0402384e589df6e087d6d2d22b463ddc2eea0774#0402384e589df6e087d6d2d22b463ddc2eea0774" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=5e9396033a5a5d5f350bfa675d12770c972ee991#5e9396033a5a5d5f350bfa675d12770c972ee991" dependencies = [ "asap-types", "promql-parser", @@ -393,7 +393,7 @@ dependencies = [ [[package]] name = "asap-frontend-sql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=0402384e589df6e087d6d2d22b463ddc2eea0774#0402384e589df6e087d6d2d22b463ddc2eea0774" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=5e9396033a5a5d5f350bfa675d12770c972ee991#5e9396033a5a5d5f350bfa675d12770c972ee991" dependencies = [ "asap-sql-function-catalog", "asap-types", @@ -415,12 +415,12 @@ dependencies = [ [[package]] name = "asap-sql-function-catalog" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=0402384e589df6e087d6d2d22b463ddc2eea0774#0402384e589df6e087d6d2d22b463ddc2eea0774" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=5e9396033a5a5d5f350bfa675d12770c972ee991#5e9396033a5a5d5f350bfa675d12770c972ee991" [[package]] name = "asap-types" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=0402384e589df6e087d6d2d22b463ddc2eea0774#0402384e589df6e087d6d2d22b463ddc2eea0774" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=5e9396033a5a5d5f350bfa675d12770c972ee991#5e9396033a5a5d5f350bfa675d12770c972ee991" dependencies = [ "serde", "serde_json", diff --git a/control_plane/Cargo.toml b/control_plane/Cargo.toml index 7cb29606..26453434 100644 --- a/control_plane/Cargo.toml +++ b/control_plane/Cargo.toml @@ -76,8 +76,8 @@ asap_types.workspace = true # scaffolding, unaware that `data_plane`'s `summary_executor.rs` in *this* # repo is a real one. Vendored locally instead of chased upstream -- see # `data_plane/src/query_engines/asap_query_engine/summary_exec.rs`. -planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "0402384e589df6e087d6d2d22b463ddc2eea0774" } -asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "0402384e589df6e087d6d2d22b463ddc2eea0774" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "5e9396033a5a5d5f350bfa675d12770c972ee991" } +asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "5e9396033a5a5d5f350bfa675d12770c972ee991" } # L1 adoption (design-target-architecture.md Part B): the PromQL front # end itself, replacing control_plane's own query_parser/promql.rs. @@ -85,8 +85,8 @@ asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = # `planner-types`/`asap-aware-mapping` above -- these three MUST move # together (two revs of the same upstream repo's types in one workspace # resolve to distinct Rust types that won't unify). -asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "0402384e589df6e087d6d2d22b463ddc2eea0774" } -asap-frontend-sql = { git = "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/ProjectASAP/ASAPPlanner", rev = "0402384e589df6e087d6d2d22b463ddc2eea0774" } +asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "5e9396033a5a5d5f350bfa675d12770c972ee991" } +asap-frontend-sql = { git = "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/ProjectASAP/ASAPPlanner", rev = "5e9396033a5a5d5f350bfa675d12770c972ee991" } [dev-dependencies] tokio = { version = "1", features = ["full", "test-util"] } diff --git a/control_plane/src/query_plan/clickhouse_exact.rs b/control_plane/src/query_plan/clickhouse_exact.rs index 9c845542..654cc80c 100644 --- a/control_plane/src/query_plan/clickhouse_exact.rs +++ b/control_plane/src/query_plan/clickhouse_exact.rs @@ -77,7 +77,7 @@ fn scalar(expr: &QueryExpr, schema: &Schema) -> Result { let function = match name.as_str() { "map" => "map", "mapconcat" => "mapConcat", - "asap_map_access" => "arrayElement", + "asap_map_access" | "asap_element_access" => "arrayElement", _ => return Err(format!("unsupported exact scalar function {name}")), }; expr.scalar_type(schema).map_err(|e| e.to_string())?; @@ -301,6 +301,28 @@ mod tests { assert!(!sql.contains("{from:")); assert!(!sql.contains("{to:")); } + #[test] + fn typed_list_access_renders_native_element_lookup() { + let schema = Schema::new(vec![Column::new( + "samples", + DataType::List { + element: Box::new(Column::new("item", DataType::Float64, false)), + }, + false, + )]); + let expr = QueryExpr::FunctionCall { + name: "asap_element_access".into(), + args: vec![ + QueryExpr::Column(0), + QueryExpr::Literal(ScalarValue::Int64(-1)), + ], + }; + assert_eq!( + scalar(&expr, &schema).unwrap(), + "arrayElement(`samples`, -1)" + ); + } + #[test] fn unsupported_scalar_is_not_forwarded_as_arbitrary_native_code() { let schema = Schema::new(vec![]); diff --git a/crates/asap_types/Cargo.toml b/crates/asap_types/Cargo.toml index a83ea156..63c84655 100644 --- a/crates/asap_types/Cargo.toml +++ b/crates/asap_types/Cargo.toml @@ -33,4 +33,4 @@ sha2 = "0.10" # exactly (`control_plane/Cargo.toml`) -- two different revs of the same # git dependency in one workspace resolve to two distinct Rust types that # won't unify. -planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "0402384e589df6e087d6d2d22b463ddc2eea0774" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "5e9396033a5a5d5f350bfa675d12770c972ee991" } diff --git a/data_plane/Cargo.toml b/data_plane/Cargo.toml index 60ae2ef0..cc531049 100644 --- a/data_plane/Cargo.toml +++ b/data_plane/Cargo.toml @@ -39,8 +39,8 @@ sha2 = "0.10" # reduction: Reduction, .. }`) are `pre_asap` types, in the same crate now # (not a separate `asap-ir` import). Query serving consumes the compiled # QueryPlan; these types are used at physical-plan compilation boundaries. -planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "0402384e589df6e087d6d2d22b463ddc2eea0774" } -asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "0402384e589df6e087d6d2d22b463ddc2eea0774" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "5e9396033a5a5d5f350bfa675d12770c972ee991" } +asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "5e9396033a5a5d5f350bfa675d12770c972ee991" } # Shared external (workspace) serde.workspace = true @@ -133,7 +133,7 @@ fs2 = "0.4" # none of them. [dev-dependencies] -asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "0402384e589df6e087d6d2d22b463ddc2eea0774" } +asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "5e9396033a5a5d5f350bfa675d12770c972ee991" } tempfile = "3.20.0" criterion = { version = "0.5", features = ["html_reports"] } tokio-tungstenite = "0.21" diff --git a/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs b/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs index 15945be2..04196130 100644 --- a/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs +++ b/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs @@ -41,6 +41,7 @@ enum Cell { Bool(bool), Timestamp(i64), Map(Vec<(Cell, Cell)>), + List(Arc<[Cell]>), } fn json_cell( @@ -60,9 +61,23 @@ fn json_cell( match dtype { DataType::Null if value.is_null() => Ok(Cell::Null), DataType::Null => Err(invalid()), - DataType::List { .. } | DataType::Struct { .. } => Err( - ClickHouseRelationalError::Unsupported("collection value transport".into()), - ), + DataType::List { element } => { + let item_type = clickhouse_type + .strip_prefix("Array(") + .and_then(|inner| inner.strip_suffix(')')) + .ok_or_else(invalid)?; + let items = value.as_array().ok_or_else(invalid)?; + Ok(Cell::List( + items + .iter() + .map(|item| json_cell(item, &element.dtype, element.nullable, item_type)) + .collect::, _>>()? + .into(), + )) + } + DataType::Struct { .. } => Err(ClickHouseRelationalError::Unsupported( + "struct value transport".into(), + )), DataType::Int64 => value.as_i64().map(Cell::Int64).ok_or_else(invalid), DataType::Float64 => value.as_f64().map(Cell::Float64).ok_or_else(invalid), DataType::Utf8 => value @@ -307,7 +322,16 @@ fn clickhouse_type_matches(actual: Option<&str>, expected: &DataType, nullable: } match expected { DataType::Null => actual == "Nothing", - DataType::List { .. } | DataType::Struct { .. } => false, + DataType::List { element } => { + !nullable + && actual + .strip_prefix("Array(") + .and_then(|inner| inner.strip_suffix(')')) + .is_some_and(|inner| { + clickhouse_type_matches(Some(inner), &element.dtype, element.nullable) + }) + } + DataType::Struct { .. } => false, DataType::Int64 => actual == "Int64", DataType::Float64 => actual == "Float64", DataType::Utf8 => actual == "String", @@ -428,11 +452,17 @@ impl ClickHouseRelationalAdapter { } for row in &input.rows { for key in keys { - if contains_nan(&eval(&key.expr, row, &schema)?) { + let value = eval(&key.expr, row, &schema)?; + if contains_nan(&value) { return Err(ClickHouseRelationalError::Unsupported( "NaN sort key".into(), )); } + if !matches!(value, Cell::Null) && cell_cmp(&value, &value).is_none() { + return Err(ClickHouseRelationalError::Unsupported( + "unsupported sort key value type".into(), + )); + } } } input @@ -548,7 +578,50 @@ fn eval( } QueryExpr::FunctionCall { name, args } => { use planner_types::pre_asap::scalar_signature::MapScalarFunction; - let function = MapScalarFunction::from_name(name).ok_or_else(|| { + if name.eq_ignore_ascii_case("asap_element_access") { + let (output_type, _) = expr + .scalar_type(schema) + .map_err(|error| ClickHouseRelationalError::Invalid(error.to_string()))?; + if let DataType::List { element } = args[0] + .scalar_type(schema) + .map_err(|error| ClickHouseRelationalError::Invalid(error.to_string()))? + .0 + { + let Cell::List(values) = eval(&args[0], row, schema)? else { + return Err(ClickHouseRelationalError::Invalid( + "array access input".into(), + )); + }; + let index = match eval(&args[1], row, schema)? { + Cell::Null => return Ok(Cell::Null), + Cell::Int64(index) => index, + _ => { + return Err(ClickHouseRelationalError::Invalid( + "array access index".into(), + )) + } + }; + let offset = if index > 0 { + usize::try_from(index - 1).ok() + } else if index < 0 { + usize::try_from(index.unsigned_abs()) + .ok() + .and_then(|distance| values.len().checked_sub(distance)) + } else { + None + }; + return match offset.and_then(|offset| values.get(offset)) { + Some(value) => Ok(value.clone()), + None => default_collection_element(&output_type, element.nullable), + }; + } + } + let function = (if name.eq_ignore_ascii_case("asap_element_access") { + Some(MapScalarFunction::Access) + } else { + MapScalarFunction::from_name(name) + }) + .ok_or_else(|| { ClickHouseRelationalError::Unsupported(format!("scalar function {name}")) })?; expr.scalar_type(schema) @@ -619,7 +692,7 @@ fn eval( else { unreachable!() }; - default_map_value(&value, value_nullable) + default_collection_element(&value, value_nullable) } } } @@ -629,7 +702,7 @@ fn eval( } } -fn default_map_value(dtype: &DataType, nullable: bool) -> Result { +fn default_collection_element(dtype: &DataType, nullable: bool) -> Result { if nullable { return Ok(Cell::Null); } @@ -640,9 +713,10 @@ fn default_map_value(dtype: &DataType, nullable: bool) -> Result Cell::Utf8(String::new()), DataType::Bool => Cell::Bool(false), DataType::Map { .. } => Cell::Map(Vec::new()), + DataType::List { .. } => Cell::List(Arc::from([])), _ => { return Err(ClickHouseRelationalError::Unsupported( - "map missing-key default type".into(), + "collection missing-element default type".into(), )) } }) @@ -768,6 +842,7 @@ fn compare_sort_keys( fn contains_nan(value: &Cell) -> bool { match value { Cell::Float64(value) => value.is_nan(), + Cell::List(values) => values.iter().any(contains_nan), Cell::Map(entries) => entries .iter() .any(|(key, value)| contains_nan(key) || contains_nan(value)), @@ -974,6 +1049,107 @@ mod tests { }; use std::rc::Rc; + #[test] + fn decodes_declared_array_elements_without_losing_nullability() { + use planner_types::pre_asap::Column; + let dtype = DataType::List { + element: Box::new(Column { + name: "item".into(), + dtype: DataType::Int64, + nullable: true, + table: None, + }), + }; + assert!(clickhouse_type_matches( + Some("Array(Nullable(Int64))"), + &dtype, + false + )); + assert!(!clickhouse_type_matches( + Some("Array(Int64)"), + &dtype, + false + )); + assert!(!clickhouse_type_matches( + Some("Nullable(Array(Nullable(Int64)))"), + &dtype, + true + )); + let value = json_cell( + &serde_json::json!([9007199254740993_i64, null, -7]), + &dtype, + false, + "Array(Nullable(Int64))", + ) + .unwrap(); + let Cell::List(items) = &value else { + panic!("expected list") + }; + assert_eq!( + items.as_ref(), + &[Cell::Int64(9007199254740993), Cell::Null, Cell::Int64(-7)] + ); + let Cell::List(copy) = value.clone() else { + unreachable!() + }; + assert!(Arc::ptr_eq(items, ©)); + assert!(json_cell( + &serde_json::json!(["wrong"]), + &dtype, + false, + "Array(Nullable(Int64))" + ) + .is_err()); + } + + #[test] + fn array_access_uses_signed_indices_and_element_defaults() { + use planner_types::pre_asap::{Column, Schema}; + let function = |name: &str, args| QueryExpr::FunctionCall { + name: name.into(), + args, + }; + let dtype = DataType::List { + element: Box::new(Column::new("item", DataType::Int64, false)), + }; + let schema = Schema::new(vec![ + Column::new("items", dtype, false), + Column::new("index", DataType::Int64, true), + ]); + let access = function( + "asap_element_access", + vec![QueryExpr::Column(0), QueryExpr::Column(1)], + ); + let items = Cell::List(vec![Cell::Int64(10), Cell::Int64(20)].into()); + for (index, expected) in [ + (1, 10), + (2, 20), + (-1, 20), + (-2, 10), + (0, 0), + (3, 0), + (i64::MIN, 0), + (i64::MAX, 0), + ] { + assert_eq!( + eval(&access, &[items.clone(), Cell::Int64(index)], &schema).unwrap(), + Cell::Int64(expected) + ); + } + assert_eq!( + eval(&access, &[items, Cell::Null], &schema).unwrap(), + Cell::Null + ); + let zero = function( + "asap_element_access", + vec![ + QueryExpr::Column(0), + QueryExpr::Literal(ScalarValue::Int64(0)), + ], + ); + assert!(eval(&zero, &[Cell::List(Arc::from([])), Cell::Int64(0)], &schema).is_err()); + } + fn schema(fields: &[(&str, DataType)]) -> SummarySchema { SummarySchema { fields: fields From 32562405af8c62a7101f8a880618f3cf8ce686b6 Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 11 Sep 2026 09:10:56 -0600 Subject: [PATCH 02/13] style(clickhouse): format collection default helper --- .../asap_clickhouse_query_engine/relational_adapter.rs | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs b/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs index 04196130..072c230a 100644 --- a/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs +++ b/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs @@ -702,7 +702,10 @@ fn eval( } } -fn default_collection_element(dtype: &DataType, nullable: bool) -> Result { +fn default_collection_element( + dtype: &DataType, + nullable: bool, +) -> Result { if nullable { return Ok(Cell::Null); } From f6c59bfd798cf860c2e219be0e4140a18fdb3c7c Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 11 Sep 2026 09:11:15 -0600 Subject: [PATCH 03/13] feat(clickhouse): read typed tuple fields inside query DAGs --- .../src/query_plan/clickhouse_exact.rs | 1 + .../relational_adapter.rs | 130 +++++++++++++++++- .../relational_adapter/collection.rs | 120 ++++++++++++++++ 3 files changed, 245 insertions(+), 6 deletions(-) create mode 100644 data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter/collection.rs diff --git a/control_plane/src/query_plan/clickhouse_exact.rs b/control_plane/src/query_plan/clickhouse_exact.rs index 654cc80c..786943ce 100644 --- a/control_plane/src/query_plan/clickhouse_exact.rs +++ b/control_plane/src/query_plan/clickhouse_exact.rs @@ -78,6 +78,7 @@ fn scalar(expr: &QueryExpr, schema: &Schema) -> Result { "map" => "map", "mapconcat" => "mapConcat", "asap_map_access" | "asap_element_access" => "arrayElement", + "asap_struct_field" => "tupleElement", _ => return Err(format!("unsupported exact scalar function {name}")), }; expr.scalar_type(schema).map_err(|e| e.to_string())?; diff --git a/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs b/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs index 04196130..04de311d 100644 --- a/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs +++ b/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs @@ -1,6 +1,7 @@ //! ClickHouse row semantics for planner-owned relational wrappers. mod aggregate; +mod collection; use std::{cmp::Ordering, collections::BTreeMap, sync::Arc}; @@ -42,6 +43,7 @@ enum Cell { Timestamp(i64), Map(Vec<(Cell, Cell)>), List(Arc<[Cell]>), + Struct(Arc<[Cell]>), } fn json_cell( @@ -75,9 +77,25 @@ fn json_cell( .into(), )) } - DataType::Struct { .. } => Err(ClickHouseRelationalError::Unsupported( - "struct value transport".into(), - )), + DataType::Struct { fields } => { + let types = + collection::tuple_field_types(clickhouse_type, fields).ok_or_else(invalid)?; + let items = value + .as_array() + .filter(|items| items.len() == fields.len()) + .ok_or_else(invalid)?; + Ok(Cell::Struct( + items + .iter() + .zip(fields) + .zip(types) + .map(|((item, field), native)| { + json_cell(item, &field.dtype, field.nullable, native) + }) + .collect::, _>>()? + .into(), + )) + } DataType::Int64 => value.as_i64().map(Cell::Int64).ok_or_else(invalid), DataType::Float64 => value.as_f64().map(Cell::Float64).ok_or_else(invalid), DataType::Utf8 => value @@ -331,7 +349,9 @@ fn clickhouse_type_matches(actual: Option<&str>, expected: &DataType, nullable: clickhouse_type_matches(Some(inner), &element.dtype, element.nullable) }) } - DataType::Struct { .. } => false, + DataType::Struct { fields } => { + !nullable && collection::tuple_field_types(actual, fields).is_some() + } DataType::Int64 => actual == "Int64", DataType::Float64 => actual == "Float64", DataType::Utf8 => actual == "String", @@ -578,6 +598,37 @@ fn eval( } QueryExpr::FunctionCall { name, args } => { use planner_types::pre_asap::scalar_signature::MapScalarFunction; + if name.eq_ignore_ascii_case("asap_struct_field") { + expr.scalar_type(schema) + .map_err(|error| ClickHouseRelationalError::Invalid(error.to_string()))?; + let DataType::Struct { fields } = args[0] + .scalar_type(schema) + .map_err(|error| ClickHouseRelationalError::Invalid(error.to_string()))? + .0 + else { + unreachable!() + }; + let offset = match &args[1] { + QueryExpr::Literal(ScalarValue::Int64(index)) => { + usize::try_from(index - 1).ok() + } + QueryExpr::Literal(ScalarValue::Utf8(name)) => { + fields.iter().position(|field| &field.name == name) + } + _ => None, + } + .ok_or_else(|| { + ClickHouseRelationalError::Invalid("struct field selector".into()) + })?; + let Cell::Struct(values) = eval(&args[0], row, schema)? else { + return Err(ClickHouseRelationalError::Invalid( + "struct field input".into(), + )); + }; + return values.get(offset).cloned().ok_or_else(|| { + ClickHouseRelationalError::Invalid("struct field value".into()) + }); + } if name.eq_ignore_ascii_case("asap_element_access") { let (output_type, _) = expr .scalar_type(schema) @@ -702,7 +753,10 @@ fn eval( } } -fn default_collection_element(dtype: &DataType, nullable: bool) -> Result { +fn default_collection_element( + dtype: &DataType, + nullable: bool, +) -> Result { if nullable { return Ok(Cell::Null); } @@ -714,6 +768,13 @@ fn default_collection_element(dtype: &DataType, nullable: bool) -> Result Cell::Bool(false), DataType::Map { .. } => Cell::Map(Vec::new()), DataType::List { .. } => Cell::List(Arc::from([])), + DataType::Struct { fields } => Cell::Struct( + fields + .iter() + .map(|field| default_collection_element(&field.dtype, field.nullable)) + .collect::, _>>()? + .into(), + ), _ => { return Err(ClickHouseRelationalError::Unsupported( "collection missing-element default type".into(), @@ -842,7 +903,7 @@ fn compare_sort_keys( fn contains_nan(value: &Cell) -> bool { match value { Cell::Float64(value) => value.is_nan(), - Cell::List(values) => values.iter().any(contains_nan), + Cell::List(values) | Cell::Struct(values) => values.iter().any(contains_nan), Cell::Map(entries) => entries .iter() .any(|(key, value)| contains_nan(key) || contains_nan(value)), @@ -1150,6 +1211,63 @@ mod tests { assert!(eval(&zero, &[Cell::List(Arc::from([])), Cell::Int64(0)], &schema).is_err()); } + #[test] + fn nested_array_tuple_access_preserves_fields_and_defaults() { + use planner_types::pre_asap::{Column, Schema}; + let tuple = DataType::Struct { + fields: vec![ + Column::new("ts", DataType::Int64, false), + Column::new("value", DataType::Float64, true), + ], + }; + let dtype = DataType::List { + element: Box::new(Column::new("item", tuple, false)), + }; + let native = "Array(Tuple(ts Int64, value Nullable(Float64)))"; + assert!(clickhouse_type_matches(Some(native), &dtype, false)); + let samples = json_cell( + &serde_json::json!([[9007199254740993_i64, 2.5], [7, null]]), + &dtype, + false, + native, + ) + .unwrap(); + let schema = Schema::new(vec![Column::new("samples", dtype, false)]); + let field = |index, name: &str| QueryExpr::FunctionCall { + name: "asap_struct_field".into(), + args: vec![ + QueryExpr::FunctionCall { + name: "asap_element_access".into(), + args: vec![ + QueryExpr::Column(0), + QueryExpr::Literal(ScalarValue::Int64(index)), + ], + }, + QueryExpr::Literal(ScalarValue::Utf8(name.into())), + ], + }; + assert_eq!( + eval(&field(1, "ts"), std::slice::from_ref(&samples), &schema).unwrap(), + Cell::Int64(9007199254740993) + ); + assert_eq!( + eval(&field(1, "value"), std::slice::from_ref(&samples), &schema).unwrap(), + Cell::Float64(2.5) + ); + assert_eq!( + eval(&field(-1, "value"), std::slice::from_ref(&samples), &schema).unwrap(), + Cell::Null + ); + assert_eq!( + eval(&field(99, "ts"), std::slice::from_ref(&samples), &schema).unwrap(), + Cell::Int64(0) + ); + assert_eq!( + eval(&field(99, "value"), std::slice::from_ref(&samples), &schema).unwrap(), + Cell::Null + ); + } + fn schema(fields: &[(&str, DataType)]) -> SummarySchema { SummarySchema { fields: fields diff --git a/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter/collection.rs b/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter/collection.rs new file mode 100644 index 00000000..e676f9b7 --- /dev/null +++ b/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter/collection.rs @@ -0,0 +1,120 @@ +//! Native collection metadata is checked against the existing shared schema. + +use planner_types::pre_asap::Column; + +/// Split native type arguments without splitting nested types or quoted names. +fn arguments(input: &str) -> Option> { + let mut result = Vec::new(); + let mut start = 0; + let mut depth = 0_usize; + let mut quote = None; + let mut chars = input.char_indices().peekable(); + while let Some((position, ch)) = chars.next() { + if let Some(delimiter) = quote { + if ch == '\\' { + chars.next()?; + } else if ch == delimiter { + if chars.peek().is_some_and(|(_, next)| *next == delimiter) { + chars.next(); + } else { + quote = None; + } + } + continue; + } + match ch { + '\'' | '`' | '"' => quote = Some(ch), + '(' => depth = depth.checked_add(1)?, + ')' => depth = depth.checked_sub(1)?, + ',' if depth == 0 => { + result.push(input[start..position].trim()); + start = position + 1; + } + _ => {} + } + } + if depth != 0 || quote.is_some() { + return None; + } + if !input.is_empty() { + result.push(input[start..].trim()); + } + (!result.iter().any(|argument| argument.is_empty())).then_some(result) +} + +/// Anonymous native Tuple fields have explicit one-based names in the shared +/// Struct schema. Named fields must match their native names exactly. Other +/// Arrow names are never interpreted as an anonymous Tuple. +pub(super) fn tuple_field_types<'a>(native: &'a str, fields: &[Column]) -> Option> { + let inner = native.strip_prefix("Tuple(")?.strip_suffix(')')?; + let args = arguments(inner)?; + if args.len() != fields.len() { + return None; + } + let mut types = Vec::with_capacity(fields.len()); + for (index, (argument, field)) in args.into_iter().zip(fields).enumerate() { + if field.table.is_some() { + return None; + } + if field.name == (index + 1).to_string() + && super::clickhouse_type_matches(Some(argument), &field.dtype, field.nullable) + { + types.push(argument); + continue; + } + // Initial named transport accepts ordinary identifiers. Quoted names + // remain unsupported until a native identifier-decoding contract exists. + let (name, native_type) = argument.split_once(char::is_whitespace)?; + if name != field.name + || name.is_empty() + || !name + .chars() + .all(|ch| ch.is_ascii_alphanumeric() || ch == '_') + || !super::clickhouse_type_matches( + Some(native_type.trim()), + &field.dtype, + field.nullable, + ) + { + return None; + } + types.push(native_type.trim()); + } + Some(types) +} + +#[cfg(test)] +mod tests { + use super::*; + use planner_types::pre_asap::DataType; + + #[test] + fn tuple_metadata_preserves_names_order_and_nested_types() { + let fields = vec![ + Column::new("ts", DataType::Int64, false), + Column::new( + "samples", + DataType::List { + element: Box::new(Column::new("item", DataType::Float64, true)), + }, + false, + ), + ]; + assert_eq!( + tuple_field_types("Tuple(ts Int64, samples Array(Nullable(Float64)))", &fields), + Some(vec!["Int64", "Array(Nullable(Float64))"]) + ); + assert!( + tuple_field_types("Tuple(samples Int64, ts Array(Nullable(Float64)))", &fields) + .is_none() + ); + assert!(tuple_field_types("Tuple(Int64, Array(Nullable(Float64)))", &fields).is_none()); + let anonymous = vec![ + Column::new("1", DataType::Int64, false), + Column::new("2", DataType::Utf8, false), + ]; + assert!(tuple_field_types("Tuple(Int64, String)", &anonymous).is_some()); + assert!(arguments("Map(String, Tuple(Int64, String)), DateTime64(3, 'UTC')").is_some()); + assert!(arguments("Map(String, Tuple(Int64, String)").is_none()); + } +} From cce8a1e7710e1e47e8744b5090d251f168b343e6 Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 11 Sep 2026 09:13:09 -0600 Subject: [PATCH 04/13] refactor(clickhouse): share nested type argument parsing --- .../relational_adapter.rs | 19 +++++-------------- .../relational_adapter/collection.rs | 2 +- 2 files changed, 6 insertions(+), 15 deletions(-) diff --git a/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs b/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs index 04de311d..0bfe338a 100644 --- a/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs +++ b/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs @@ -307,20 +307,11 @@ impl ClickHouseRelation { fn map_type_parts(actual: &str) -> Option<(&str, &str)> { let inner = actual.trim().strip_prefix("Map(")?.strip_suffix(')')?; - let mut depth = 0_i32; - let mut quoted = false; - for (index, ch) in inner.char_indices() { - match ch { - '\'' => quoted = !quoted, - '(' if !quoted => depth += 1, - ')' if !quoted => depth -= 1, - ',' if !quoted && depth == 0 => { - return Some((inner[..index].trim(), inner[index + 1..].trim())) - } - _ => {} - } - } - None + let args = collection::arguments(inner)?; + let [key, value] = args.as_slice() else { + return None; + }; + Some((*key, *value)) } fn clickhouse_type_matches(actual: Option<&str>, expected: &DataType, nullable: bool) -> bool { diff --git a/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter/collection.rs b/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter/collection.rs index e676f9b7..2bb43e7b 100644 --- a/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter/collection.rs +++ b/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter/collection.rs @@ -3,7 +3,7 @@ use planner_types::pre_asap::Column; /// Split native type arguments without splitting nested types or quoted names. -fn arguments(input: &str) -> Option> { +pub(super) fn arguments(input: &str) -> Option> { let mut result = Vec::new(); let mut start = 0; let mut depth = 0_usize; From 10909feabf05a88b60726402a0d701e4926eec82 Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 11 Sep 2026 09:13:53 -0600 Subject: [PATCH 05/13] fix(clickhouse): distinguish nonfinite exact values from nulls --- .../accelerator.rs | 3 +++ .../relational_adapter.rs | 27 +++++++++++++++++++ 2 files changed, 30 insertions(+) diff --git a/data_plane/src/query_engines/asap_clickhouse_query_engine/accelerator.rs b/data_plane/src/query_engines/asap_clickhouse_query_engine/accelerator.rs index 907b7fea..8c893375 100644 --- a/data_plane/src/query_engines/asap_clickhouse_query_engine/accelerator.rs +++ b/data_plane/src/query_engines/asap_clickhouse_query_engine/accelerator.rs @@ -111,6 +111,9 @@ impl CatalogClickHouseAccelerator { "0".into(), ); parameters.insert("output_format_json_quote_64bit_integers".into(), "0".into()); + // Preserve the distinction between NULL and unsupported NaN/Inf. + // The typed decoder rejects quoted non-finite values and falls back. + parameters.insert("output_format_json_quote_denormals".into(), "1".into()); if let Some(database) = request_context.database() { parameters.insert("database".into(), database.into()); } diff --git a/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs b/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs index 072c230a..ea399ba2 100644 --- a/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs +++ b/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs @@ -1105,6 +1105,33 @@ mod tests { .is_err()); } + #[test] + fn nonfinite_external_array_values_cannot_become_nulls() { + use planner_types::pre_asap::Column; + let dtype = DataType::List { + element: Box::new(Column::new("item", DataType::Float64, true)), + }; + for value in ["inf", "-inf", "nan"] { + assert!(json_cell( + &serde_json::json!([value]), + &dtype, + false, + "Array(Nullable(Float64))" + ) + .is_err()); + } + assert_eq!( + json_cell( + &serde_json::json!([null]), + &dtype, + false, + "Array(Nullable(Float64))" + ) + .unwrap(), + Cell::List(vec![Cell::Null].into()) + ); + } + #[test] fn array_access_uses_signed_indices_and_element_defaults() { use planner_types::pre_asap::{Column, Schema}; From f8261340035a85f1ec8f8abd16be20c9935192e8 Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 11 Sep 2026 09:14:14 -0600 Subject: [PATCH 06/13] test(clickhouse): verify exact denormal transport policy --- .../asap_clickhouse_query_engine/accelerator.rs | 4 ++++ .../asap_clickhouse_query_engine/relational_adapter.rs | 7 +++++++ 2 files changed, 11 insertions(+) diff --git a/data_plane/src/query_engines/asap_clickhouse_query_engine/accelerator.rs b/data_plane/src/query_engines/asap_clickhouse_query_engine/accelerator.rs index 8c893375..8e34ceaf 100644 --- a/data_plane/src/query_engines/asap_clickhouse_query_engine/accelerator.rs +++ b/data_plane/src/query_engines/asap_clickhouse_query_engine/accelerator.rs @@ -331,6 +331,10 @@ mod tests { request: &ClickHouseQueryRequest, ) -> Result { + assert_eq!( + request.parameters.get("output_format_json_quote_denormals"), + Some(&"1".into()) + ); assert_eq!(request.parameters.get("param_from"), Some(&"0".into())); assert_eq!(request.parameters.get("param_to"), Some(&"2000".into())); Ok(ClickHouseRawResponse { diff --git a/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs b/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs index ea399ba2..cdecb879 100644 --- a/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs +++ b/data_plane/src/query_engines/asap_clickhouse_query_engine/relational_adapter.rs @@ -1112,6 +1112,13 @@ mod tests { element: Box::new(Column::new("item", DataType::Float64, true)), }; for value in ["inf", "-inf", "nan"] { + assert!(json_cell( + &serde_json::json!(value), + &DataType::Float64, + true, + "Nullable(Float64)" + ) + .is_err()); assert!(json_cell( &serde_json::json!([value]), &dtype, From 05b57db76045298522ab7470e9e982e4506d3965 Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 11 Sep 2026 09:15:57 -0600 Subject: [PATCH 07/13] test(clickhouse): verify native tuple field rendering --- .../src/query_plan/clickhouse_exact.rs | 25 +++++++++++++++++++ 1 file changed, 25 insertions(+) diff --git a/control_plane/src/query_plan/clickhouse_exact.rs b/control_plane/src/query_plan/clickhouse_exact.rs index 3f8161b3..fa227841 100644 --- a/control_plane/src/query_plan/clickhouse_exact.rs +++ b/control_plane/src/query_plan/clickhouse_exact.rs @@ -324,6 +324,31 @@ mod tests { ); } + #[test] + fn typed_struct_field_renders_native_lookup() { + let schema = Schema::new(vec![Column::new( + "sample", + DataType::Struct { + fields: vec![ + Column::new("ts", DataType::Int64, false), + Column::new("value", DataType::Float64, true), + ], + }, + false, + )]); + let expr = QueryExpr::FunctionCall { + name: "asap_struct_field".into(), + args: vec![ + QueryExpr::Column(0), + QueryExpr::Literal(ScalarValue::Utf8("value".into())), + ], + }; + assert_eq!( + scalar(&expr, &schema).unwrap(), + "tupleElement(`sample`, 'value')" + ); + } + #[test] fn unsupported_scalar_is_not_forwarded_as_arbitrary_native_code() { let schema = Schema::new(vec![]); From 31bdeb8f65353291c4f9492f76f2d587ca3d85ee Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 11 Sep 2026 09:37:36 -0600 Subject: [PATCH 08/13] fix(clickhouse): preserve SQL nulls in result formats --- .../accelerator.rs | 9 ++-- .../clickhouse_result_adapter.rs | 46 +++++++++++++++++-- 2 files changed, 47 insertions(+), 8 deletions(-) diff --git a/data_plane/src/query_engines/asap_clickhouse_query_engine/accelerator.rs b/data_plane/src/query_engines/asap_clickhouse_query_engine/accelerator.rs index 8e34ceaf..dc774056 100644 --- a/data_plane/src/query_engines/asap_clickhouse_query_engine/accelerator.rs +++ b/data_plane/src/query_engines/asap_clickhouse_query_engine/accelerator.rs @@ -147,11 +147,9 @@ impl CatalogClickHouseAccelerator { } fn requested_format(request: &ClickHouseQueryRequest) -> Result { - if let Some(setting) = request - .parameters - .keys() - .find(|key| key.starts_with("output_format_")) - { + if let Some(setting) = request.parameters.keys().find(|key| { + key.starts_with("output_format_") || key.as_str() == "format_tsv_null_representation" + }) { return Err(format!("unsupported output setting {setting}")); } match request @@ -302,6 +300,7 @@ mod tests { for setting in [ "output_format_json_map_as_array_of_tuples", "output_format_json_quote_64bit_integers", + "format_tsv_null_representation", ] { request.parameters.insert(setting.into(), "1".into()); assert!(requested_format(&request).is_err()); diff --git a/data_plane/src/query_engines/asap_clickhouse_query_engine/clickhouse_result_adapter.rs b/data_plane/src/query_engines/asap_clickhouse_query_engine/clickhouse_result_adapter.rs index a829c548..61d0558b 100644 --- a/data_plane/src/query_engines/asap_clickhouse_query_engine/clickhouse_result_adapter.rs +++ b/data_plane/src/query_engines/asap_clickhouse_query_engine/clickhouse_result_adapter.rs @@ -2,7 +2,7 @@ use super::fallback::ClickHouseRawResponse; use arrow::{ array::{Array, Float64Array, StringArray, TimestampMillisecondArray}, datatypes::{DataType, Field, Schema}, - json::LineDelimitedWriter, + json::{LineDelimitedWriter, WriterBuilder}, record_batch::RecordBatch, util::display::array_value_to_string, }; @@ -118,6 +118,10 @@ impl ClickHouseQueryResult { if column > 0 { output.push(b'\t'); } + if batch.column(column).is_null(row) { + output.extend_from_slice(b"\\N"); + continue; + } if matches!(batch.column(column).data_type(), DataType::Map(..)) { output.extend_from_slice( map_literal(batch.column(column).as_ref(), row)?.as_bytes(), @@ -145,7 +149,9 @@ impl ClickHouseQueryResult { } let row = batch.slice(row, 1); - let mut writer = LineDelimitedWriter::new(&mut output); + let mut writer: LineDelimitedWriter<_> = WriterBuilder::new() + .with_explicit_nulls(true) + .build(&mut output); writer .write_batches(&[&row]) .map_err(|error| ClickHouseResultError::Arrow(error.to_string()))?; @@ -204,7 +210,9 @@ impl ClickHouseQueryResult { let mut rows = Vec::new(); for batch in &self.batches { let mut encoded = Vec::new(); - let mut writer = LineDelimitedWriter::new(&mut encoded); + let mut writer: LineDelimitedWriter<_> = WriterBuilder::new() + .with_explicit_nulls(true) + .build(&mut encoded); writer .write_batches(&[batch]) .map_err(|error| ClickHouseResultError::Arrow(error.to_string()))?; @@ -309,6 +317,38 @@ mod tests { }; use std::sync::Arc; + #[test] + fn nullable_fields_remain_explicit_in_json_and_tsv() { + let batch = RecordBatch::try_new( + Arc::new(Schema::new(vec![ + Field::new("value", DataType::Float64, true), + Field::new("label", DataType::Utf8, true), + ])), + vec![ + Arc::new(Float64Array::from(vec![Some(1.25), None])), + Arc::new(StringArray::from(vec![Some(""), None])), + ], + ) + .unwrap(); + let result = ClickHouseQueryResult { + batches: vec![batch], + }; + let json: serde_json::Value = + serde_json::from_slice(&result.encode(ClickHouseFormat::Json).unwrap()).unwrap(); + assert_eq!( + json["data"], + serde_json::json!([{ "value":1.25,"label":"" }, { "value":null,"label":null }]) + ); + let lines = result.encode(ClickHouseFormat::JsonEachRow).unwrap(); + let null_row: serde_json::Value = + serde_json::from_slice(lines.split(|byte| *byte == b'\n').nth(1).unwrap()).unwrap(); + assert_eq!(null_row, serde_json::json!({"value":null,"label":null})); + assert_eq!( + result.encode(ClickHouseFormat::TabSeparated).unwrap(), + b"1.25\t\n\\N\t\\N\n" + ); + } + #[test] fn empty_map_bottom_type_uses_clickhouse_nothing() { let entries = DataType::Struct( From 8e0ddd6f87c959775e91337f67933c1fbeb2fd56 Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 11 Sep 2026 09:37:36 -0600 Subject: [PATCH 09/13] fix(clickhouse): preserve SQL nulls in result formats --- .../accelerator.rs | 9 ++-- .../clickhouse_result_adapter.rs | 46 +++++++++++++++++-- 2 files changed, 47 insertions(+), 8 deletions(-) diff --git a/data_plane/src/query_engines/asap_clickhouse_query_engine/accelerator.rs b/data_plane/src/query_engines/asap_clickhouse_query_engine/accelerator.rs index 8e34ceaf..dc774056 100644 --- a/data_plane/src/query_engines/asap_clickhouse_query_engine/accelerator.rs +++ b/data_plane/src/query_engines/asap_clickhouse_query_engine/accelerator.rs @@ -147,11 +147,9 @@ impl CatalogClickHouseAccelerator { } fn requested_format(request: &ClickHouseQueryRequest) -> Result { - if let Some(setting) = request - .parameters - .keys() - .find(|key| key.starts_with("output_format_")) - { + if let Some(setting) = request.parameters.keys().find(|key| { + key.starts_with("output_format_") || key.as_str() == "format_tsv_null_representation" + }) { return Err(format!("unsupported output setting {setting}")); } match request @@ -302,6 +300,7 @@ mod tests { for setting in [ "output_format_json_map_as_array_of_tuples", "output_format_json_quote_64bit_integers", + "format_tsv_null_representation", ] { request.parameters.insert(setting.into(), "1".into()); assert!(requested_format(&request).is_err()); diff --git a/data_plane/src/query_engines/asap_clickhouse_query_engine/clickhouse_result_adapter.rs b/data_plane/src/query_engines/asap_clickhouse_query_engine/clickhouse_result_adapter.rs index a829c548..61d0558b 100644 --- a/data_plane/src/query_engines/asap_clickhouse_query_engine/clickhouse_result_adapter.rs +++ b/data_plane/src/query_engines/asap_clickhouse_query_engine/clickhouse_result_adapter.rs @@ -2,7 +2,7 @@ use super::fallback::ClickHouseRawResponse; use arrow::{ array::{Array, Float64Array, StringArray, TimestampMillisecondArray}, datatypes::{DataType, Field, Schema}, - json::LineDelimitedWriter, + json::{LineDelimitedWriter, WriterBuilder}, record_batch::RecordBatch, util::display::array_value_to_string, }; @@ -118,6 +118,10 @@ impl ClickHouseQueryResult { if column > 0 { output.push(b'\t'); } + if batch.column(column).is_null(row) { + output.extend_from_slice(b"\\N"); + continue; + } if matches!(batch.column(column).data_type(), DataType::Map(..)) { output.extend_from_slice( map_literal(batch.column(column).as_ref(), row)?.as_bytes(), @@ -145,7 +149,9 @@ impl ClickHouseQueryResult { } let row = batch.slice(row, 1); - let mut writer = LineDelimitedWriter::new(&mut output); + let mut writer: LineDelimitedWriter<_> = WriterBuilder::new() + .with_explicit_nulls(true) + .build(&mut output); writer .write_batches(&[&row]) .map_err(|error| ClickHouseResultError::Arrow(error.to_string()))?; @@ -204,7 +210,9 @@ impl ClickHouseQueryResult { let mut rows = Vec::new(); for batch in &self.batches { let mut encoded = Vec::new(); - let mut writer = LineDelimitedWriter::new(&mut encoded); + let mut writer: LineDelimitedWriter<_> = WriterBuilder::new() + .with_explicit_nulls(true) + .build(&mut encoded); writer .write_batches(&[batch]) .map_err(|error| ClickHouseResultError::Arrow(error.to_string()))?; @@ -309,6 +317,38 @@ mod tests { }; use std::sync::Arc; + #[test] + fn nullable_fields_remain_explicit_in_json_and_tsv() { + let batch = RecordBatch::try_new( + Arc::new(Schema::new(vec![ + Field::new("value", DataType::Float64, true), + Field::new("label", DataType::Utf8, true), + ])), + vec![ + Arc::new(Float64Array::from(vec![Some(1.25), None])), + Arc::new(StringArray::from(vec![Some(""), None])), + ], + ) + .unwrap(); + let result = ClickHouseQueryResult { + batches: vec![batch], + }; + let json: serde_json::Value = + serde_json::from_slice(&result.encode(ClickHouseFormat::Json).unwrap()).unwrap(); + assert_eq!( + json["data"], + serde_json::json!([{ "value":1.25,"label":"" }, { "value":null,"label":null }]) + ); + let lines = result.encode(ClickHouseFormat::JsonEachRow).unwrap(); + let null_row: serde_json::Value = + serde_json::from_slice(lines.split(|byte| *byte == b'\n').nth(1).unwrap()).unwrap(); + assert_eq!(null_row, serde_json::json!({"value":null,"label":null})); + assert_eq!( + result.encode(ClickHouseFormat::TabSeparated).unwrap(), + b"1.25\t\n\\N\t\\N\n" + ); + } + #[test] fn empty_map_bottom_type_uses_clickhouse_nothing() { let entries = DataType::Struct( From f6aa7feb082ff4326e992f725d34a9fbb1e1b4c2 Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 11 Sep 2026 09:39:04 -0600 Subject: [PATCH 10/13] test(clickhouse): distinguish literal null marker strings --- .../clickhouse_result_adapter.rs | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/data_plane/src/query_engines/asap_clickhouse_query_engine/clickhouse_result_adapter.rs b/data_plane/src/query_engines/asap_clickhouse_query_engine/clickhouse_result_adapter.rs index 61d0558b..cdc46d26 100644 --- a/data_plane/src/query_engines/asap_clickhouse_query_engine/clickhouse_result_adapter.rs +++ b/data_plane/src/query_engines/asap_clickhouse_query_engine/clickhouse_result_adapter.rs @@ -325,8 +325,8 @@ mod tests { Field::new("label", DataType::Utf8, true), ])), vec![ - Arc::new(Float64Array::from(vec![Some(1.25), None])), - Arc::new(StringArray::from(vec![Some(""), None])), + Arc::new(Float64Array::from(vec![Some(1.25), None, Some(2.5)])), + Arc::new(StringArray::from(vec![Some(""), None, Some("\\N")])), ], ) .unwrap(); @@ -337,7 +337,7 @@ mod tests { serde_json::from_slice(&result.encode(ClickHouseFormat::Json).unwrap()).unwrap(); assert_eq!( json["data"], - serde_json::json!([{ "value":1.25,"label":"" }, { "value":null,"label":null }]) + serde_json::json!([{ "value":1.25,"label":"" }, { "value":null,"label":null }, {"value":2.5,"label":"\\N"}]) ); let lines = result.encode(ClickHouseFormat::JsonEachRow).unwrap(); let null_row: serde_json::Value = @@ -345,7 +345,7 @@ mod tests { assert_eq!(null_row, serde_json::json!({"value":null,"label":null})); assert_eq!( result.encode(ClickHouseFormat::TabSeparated).unwrap(), - b"1.25\t\n\\N\t\\N\n" + b"1.25\t\n\\N\t\\N\n2.5\t\\\\N\n" ); } From 2754be8ed9e6e44b731ca2cd8d4f854b9f7b48c9 Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 11 Sep 2026 09:39:04 -0600 Subject: [PATCH 11/13] test(clickhouse): distinguish literal null marker strings --- .../clickhouse_result_adapter.rs | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/data_plane/src/query_engines/asap_clickhouse_query_engine/clickhouse_result_adapter.rs b/data_plane/src/query_engines/asap_clickhouse_query_engine/clickhouse_result_adapter.rs index 61d0558b..cdc46d26 100644 --- a/data_plane/src/query_engines/asap_clickhouse_query_engine/clickhouse_result_adapter.rs +++ b/data_plane/src/query_engines/asap_clickhouse_query_engine/clickhouse_result_adapter.rs @@ -325,8 +325,8 @@ mod tests { Field::new("label", DataType::Utf8, true), ])), vec![ - Arc::new(Float64Array::from(vec![Some(1.25), None])), - Arc::new(StringArray::from(vec![Some(""), None])), + Arc::new(Float64Array::from(vec![Some(1.25), None, Some(2.5)])), + Arc::new(StringArray::from(vec![Some(""), None, Some("\\N")])), ], ) .unwrap(); @@ -337,7 +337,7 @@ mod tests { serde_json::from_slice(&result.encode(ClickHouseFormat::Json).unwrap()).unwrap(); assert_eq!( json["data"], - serde_json::json!([{ "value":1.25,"label":"" }, { "value":null,"label":null }]) + serde_json::json!([{ "value":1.25,"label":"" }, { "value":null,"label":null }, {"value":2.5,"label":"\\N"}]) ); let lines = result.encode(ClickHouseFormat::JsonEachRow).unwrap(); let null_row: serde_json::Value = @@ -345,7 +345,7 @@ mod tests { assert_eq!(null_row, serde_json::json!({"value":null,"label":null})); assert_eq!( result.encode(ClickHouseFormat::TabSeparated).unwrap(), - b"1.25\t\n\\N\t\\N\n" + b"1.25\t\n\\N\t\\N\n2.5\t\\\\N\n" ); } From fe0ba21bb5d888e8bd4fdc720cfe7a78966021e4 Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 11 Sep 2026 09:43:51 -0600 Subject: [PATCH 12/13] test(clickhouse): execute planner-selected array SQL in a process --- Cargo.lock | 10 +- control_plane/Cargo.toml | 8 +- crates/asap_types/Cargo.toml | 2 +- data_plane/Cargo.toml | 6 +- .../tests/clickhouse_differential_e2e.rs | 278 +++++++++++++++--- 5 files changed, 246 insertions(+), 58 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index c0dfaf2b..13a4be1f 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -373,7 +373,7 @@ dependencies = [ [[package]] name = "asap-aware-mapping" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=5e9396033a5a5d5f350bfa675d12770c972ee991#5e9396033a5a5d5f350bfa675d12770c972ee991" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=cf1f9acf0e09ce31b8319e674fb2508185a31148#cf1f9acf0e09ce31b8319e674fb2508185a31148" dependencies = [ "asap-types", "serde", @@ -384,7 +384,7 @@ dependencies = [ [[package]] name = "asap-frontend-promql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=5e9396033a5a5d5f350bfa675d12770c972ee991#5e9396033a5a5d5f350bfa675d12770c972ee991" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=cf1f9acf0e09ce31b8319e674fb2508185a31148#cf1f9acf0e09ce31b8319e674fb2508185a31148" dependencies = [ "asap-types", "promql-parser", @@ -393,7 +393,7 @@ dependencies = [ [[package]] name = "asap-frontend-sql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=5e9396033a5a5d5f350bfa675d12770c972ee991#5e9396033a5a5d5f350bfa675d12770c972ee991" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=cf1f9acf0e09ce31b8319e674fb2508185a31148#cf1f9acf0e09ce31b8319e674fb2508185a31148" dependencies = [ "asap-sql-function-catalog", "asap-types", @@ -415,12 +415,12 @@ dependencies = [ [[package]] name = "asap-sql-function-catalog" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=5e9396033a5a5d5f350bfa675d12770c972ee991#5e9396033a5a5d5f350bfa675d12770c972ee991" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=cf1f9acf0e09ce31b8319e674fb2508185a31148#cf1f9acf0e09ce31b8319e674fb2508185a31148" [[package]] name = "asap-types" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=5e9396033a5a5d5f350bfa675d12770c972ee991#5e9396033a5a5d5f350bfa675d12770c972ee991" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=cf1f9acf0e09ce31b8319e674fb2508185a31148#cf1f9acf0e09ce31b8319e674fb2508185a31148" dependencies = [ "serde", "serde_json", diff --git a/control_plane/Cargo.toml b/control_plane/Cargo.toml index 26453434..fb664da3 100644 --- a/control_plane/Cargo.toml +++ b/control_plane/Cargo.toml @@ -76,8 +76,8 @@ asap_types.workspace = true # scaffolding, unaware that `data_plane`'s `summary_executor.rs` in *this* # repo is a real one. Vendored locally instead of chased upstream -- see # `data_plane/src/query_engines/asap_query_engine/summary_exec.rs`. -planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "5e9396033a5a5d5f350bfa675d12770c972ee991" } -asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "5e9396033a5a5d5f350bfa675d12770c972ee991" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "cf1f9acf0e09ce31b8319e674fb2508185a31148" } +asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "cf1f9acf0e09ce31b8319e674fb2508185a31148" } # L1 adoption (design-target-architecture.md Part B): the PromQL front # end itself, replacing control_plane's own query_parser/promql.rs. @@ -85,8 +85,8 @@ asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = # `planner-types`/`asap-aware-mapping` above -- these three MUST move # together (two revs of the same upstream repo's types in one workspace # resolve to distinct Rust types that won't unify). -asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "5e9396033a5a5d5f350bfa675d12770c972ee991" } -asap-frontend-sql = { git = "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/ProjectASAP/ASAPPlanner", rev = "5e9396033a5a5d5f350bfa675d12770c972ee991" } +asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "cf1f9acf0e09ce31b8319e674fb2508185a31148" } +asap-frontend-sql = { git = "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/ProjectASAP/ASAPPlanner", rev = "cf1f9acf0e09ce31b8319e674fb2508185a31148" } [dev-dependencies] tokio = { version = "1", features = ["full", "test-util"] } diff --git a/crates/asap_types/Cargo.toml b/crates/asap_types/Cargo.toml index d6480e60..16ed0aba 100644 --- a/crates/asap_types/Cargo.toml +++ b/crates/asap_types/Cargo.toml @@ -34,4 +34,4 @@ sha2 = "0.10" # exactly (`control_plane/Cargo.toml`) -- two different revs of the same # git dependency in one workspace resolve to two distinct Rust types that # won't unify. -planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "5e9396033a5a5d5f350bfa675d12770c972ee991" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "cf1f9acf0e09ce31b8319e674fb2508185a31148" } diff --git a/data_plane/Cargo.toml b/data_plane/Cargo.toml index cc531049..536f24fe 100644 --- a/data_plane/Cargo.toml +++ b/data_plane/Cargo.toml @@ -39,8 +39,8 @@ sha2 = "0.10" # reduction: Reduction, .. }`) are `pre_asap` types, in the same crate now # (not a separate `asap-ir` import). Query serving consumes the compiled # QueryPlan; these types are used at physical-plan compilation boundaries. -planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "5e9396033a5a5d5f350bfa675d12770c972ee991" } -asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "5e9396033a5a5d5f350bfa675d12770c972ee991" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "cf1f9acf0e09ce31b8319e674fb2508185a31148" } +asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "cf1f9acf0e09ce31b8319e674fb2508185a31148" } # Shared external (workspace) serde.workspace = true @@ -133,7 +133,7 @@ fs2 = "0.4" # none of them. [dev-dependencies] -asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "5e9396033a5a5d5f350bfa675d12770c972ee991" } +asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "cf1f9acf0e09ce31b8319e674fb2508185a31148" } tempfile = "3.20.0" criterion = { version = "0.5", features = ["html_reports"] } tokio-tungstenite = "0.21" diff --git a/data_plane/tests/clickhouse_differential_e2e.rs b/data_plane/tests/clickhouse_differential_e2e.rs index 75a07fd4..8023fd22 100644 --- a/data_plane/tests/clickhouse_differential_e2e.rs +++ b/data_plane/tests/clickhouse_differential_e2e.rs @@ -51,6 +51,67 @@ async fn wait_http(client: &reqwest::Client, url: &str, child: &mut Child) { panic!("data plane did not become ready at {url}"); } +async fn spawn_backend( + clickhouse_url: &str, + user: Option<&str>, + password: Option<&str>, + output: &std::path::Path, + bootstrap: &std::path::Path, + api_port: u16, + sql_port: u16, +) -> ChildGuard { + let client = reqwest::Client::new(); + let mut command = Command::new(env!("CARGO_BIN_EXE_data_plane")); + command + .arg("--streaming-config") + .arg(bootstrap) + .arg("--http-port") + .arg(api_port.to_string()) + .arg("--clickhouse-http-port") + .arg(sql_port.to_string()) + .arg("--clickhouse-url") + .arg(&clickhouse_url) + .arg("--clickhouse-backfill-table") + .arg("deployment_default_not_the_job_table") + .arg("--clickhouse-backfill-database") + .arg("default") + .arg("--clickhouse-backfill-value-column") + .arg("wrong_value") + .arg("--enable-backfill-worker") + .arg("--precompute-allowed-lateness-ms") + .arg("0") + .arg("--precompute-flush-interval-ms") + .arg("50") + .arg("--persistence-delete-older-than-secs") + .arg("0") + .arg("--output-dir") + .arg(output) + .stdout(Stdio::inherit()) + .stderr(Stdio::inherit()) + .env("RUST_LOG", "data_plane=info"); + if let Some(user) = user { + command.arg("--clickhouse-user").arg(user); + } + if let Some(password) = password { + command.arg("--clickhouse-password").arg(password); + } + let mut process = ChildGuard(command.spawn().unwrap()); + wait_http( + &client, + &format!("http://127.0.0.1:{api_port}/api/v1/health"), + &mut process.0, + ) + .await; + wait_http( + &client, + &format!("http://127.0.0.1:{sql_port}/ping"), + &mut process.0, + ) + .await; + + process +} + fn mixed_workload(sql: &str) -> control_plane::clickhouse::ClickHouseSqlAutomaticWorkload { use control_plane::physical::compiler::{PlanEnvelope, BACKEND_COMPAT, PLANNER_REVISION}; use planner_types::pre_asap::{Column, DataType, Schema}; @@ -285,51 +346,14 @@ async fn run_mixed_aggregate(aggregate: &str) { let output = tempfile::tempdir().unwrap(); let mut bootstrap = tempfile::NamedTempFile::new().unwrap(); writeln!(bootstrap, "aggregations: []").unwrap(); - let mut command = Command::new(env!("CARGO_BIN_EXE_data_plane")); - command - .arg("--streaming-config") - .arg(bootstrap.path()) - .arg("--http-port") - .arg(api_port.to_string()) - .arg("--clickhouse-http-port") - .arg(sql_port.to_string()) - .arg("--clickhouse-url") - .arg(&clickhouse_url) - .arg("--clickhouse-backfill-table") - .arg("deployment_default_not_the_job_table") - .arg("--clickhouse-backfill-database") - .arg("default") - .arg("--clickhouse-backfill-value-column") - .arg("wrong_value") - .arg("--enable-backfill-worker") - .arg("--precompute-allowed-lateness-ms") - .arg("0") - .arg("--precompute-flush-interval-ms") - .arg("50") - .arg("--persistence-delete-older-than-secs") - .arg("0") - .arg("--output-dir") - .arg(output.path()) - .stdout(Stdio::inherit()) - .stderr(Stdio::inherit()) - .env("RUST_LOG", "data_plane=info"); - if let Some(user) = &user { - command.arg("--clickhouse-user").arg(user); - } - if let Some(password) = &password { - command.arg("--clickhouse-password").arg(password); - } - let mut process = ChildGuard(command.spawn().unwrap()); - wait_http( - &client, - &format!("http://127.0.0.1:{api_port}/api/v1/health"), - &mut process.0, - ) - .await; - wait_http( - &client, - &format!("http://127.0.0.1:{sql_port}/ping"), - &mut process.0, + let _process = spawn_backend( + &clickhouse_url, + user.as_deref(), + password.as_deref(), + output.path(), + bootstrap.path(), + api_port, + sql_port, ) .await; @@ -466,3 +490,167 @@ async fn run_mixed_aggregate(aggregate: &str) { ); } } + +#[tokio::test] +async fn array_sql_executes_local_element_after_typed_exact_leaf() { + use asap_types::query_plan::QueryPlanNode; + use planner_types::{ + post_asap::ValueOperation, + pre_asap::{Column, DataType, QueryExpr, Schema}, + }; + let Ok(clickhouse_url) = std::env::var("CLICKHOUSE_URL") else { + eprintln!("skipping collection process E2E because CLICKHOUSE_URL is unset"); + return; + }; + let user = std::env::var("CLICKHOUSE_USER").ok(); + let password = std::env::var("CLICKHOUSE_PASSWORD").ok(); + let client = reqwest::Client::new(); + let table = format!("asap_collection_elements_{}", std::process::id()); + for sql in [ + format!("CREATE TABLE default.{table}(timestamp Int64, samples Array(Float64), nullable_samples Array(Nullable(Float64)), position Nullable(Int64)) ENGINE=Memory"), + format!("INSERT INTO default.{table} VALUES (100,[10,20],[10,20],-1),(200,[3.5],[3.5],9),(300,[],[],1),(400,[7.5],[NULL],1),(500,[1],[1],NULL),(2000,[999],[999],1)"), + ] { + let mut request = client.post(&clickhouse_url).body(sql); + if let Some(user) = &user { request = request.basic_auth(user, password.as_ref()); } + let response = request.send().await.unwrap(); + let status = response.status(); + let body = response.text().await.unwrap(); + assert!(status.is_success(), "native collection fixture: {body}"); + } + let sql = format!("SELECT arrayElement(samples, position) AS result, arrayElement(nullable_samples, position) AS nullable_result FROM {table} WHERE timestamp >= 0 AND timestamp < 2000 ORDER BY result NULLS FIRST"); + let mut workload = mixed_workload(&sql); + workload.tables = HashMap::from([( + table.clone(), + Schema::with_time_index( + vec![ + Column::new("timestamp", DataType::Int64, false), + Column::new( + "samples", + DataType::List { + element: Box::new(Column::new("item", DataType::Float64, false)), + }, + false, + ), + Column::new( + "nullable_samples", + DataType::List { + element: Box::new(Column::new("item", DataType::Float64, true)), + }, + false, + ), + Column::new("position", DataType::Int64, true), + ], + 0, + vec![], + ), + )]); + let (publication, trace) = + control_plane::clickhouse::compile_automatic_clickhouse_workload(&workload) + .await + .unwrap(); + assert!(publication.precompute_plan.materializations.is_empty()); + let entry = publication.query_plan.entries.values().next().unwrap(); + assert!(entry + .nodes + .values() + .any(|node| matches!(node, QueryPlanNode::ExternalExact { .. }))); + assert!(entry.nodes.values().any(|node| { + let QueryPlanNode::Relational { operation, .. } = node else { return false; }; + let ValueOperation::Project { cols, .. } = serde_json::from_value(operation.clone()).unwrap() else { return false; }; + cols.iter().any(|column| matches!(&column.expr, QueryExpr::FunctionCall { name, .. } if name == "asap_element_access")) + }), "Planner-selected DAG must preserve local element evaluation"); + eprintln!( + "collection Planner selection: {}", + serde_json::to_string(&trace).unwrap() + ); + let install = publication.install_request(None, Vec::new()).unwrap(); + let api_port = unused_port(); + let sql_port = unused_port(); + let output = tempfile::tempdir().unwrap(); + let mut bootstrap = tempfile::NamedTempFile::new().unwrap(); + writeln!(bootstrap, "aggregations: []").unwrap(); + let _process = spawn_backend( + &clickhouse_url, + user.as_deref(), + password.as_deref(), + output.path(), + bootstrap.path(), + api_port, + sql_port, + ) + .await; + let response = client + .post(format!("http://127.0.0.1:{api_port}/api/v1/physical-plan")) + .json(&install) + .send() + .await + .unwrap(); + let status = response.status(); + let body = response.text().await.unwrap(); + assert!(status.is_success(), "stage collection publication: {body}"); + let response = client + .post(format!( + "http://127.0.0.1:{api_port}/api/v1/physical-plan/activate" + )) + .json(&serde_json::json!({"plan_id":72,"plan_version":1})) + .send() + .await + .unwrap(); + assert!(response.status().is_success()); + let mut native = client + .post(&clickhouse_url) + .body(format!("{sql} FORMAT JSON")); + if let Some(user) = &user { + native = native.basic_auth(user, password.as_ref()); + } + let expected: serde_json::Value = native.send().await.unwrap().json().await.unwrap(); + let mut actual_request = client + .get(format!("http://127.0.0.1:{sql_port}/")) + .query(&[("query", sql.as_str()), ("default_format", "JSON")]); + if let Some(user) = &user { + actual_request = actual_request.basic_auth(user, password.as_ref()); + } + let response = actual_request.send().await.unwrap(); + assert!(response.status().is_success()); + assert_eq!( + response.headers().get("x-asap-execution").unwrap(), + "exact_fallback" + ); + assert_eq!( + response.headers().get("x-asap-execution-detail").unwrap(), + "external_dag" + ); + let actual: serde_json::Value = response.json().await.unwrap(); + for field in ["meta", "rows"] { + assert_eq!(actual[field], expected[field], "{field}"); + } + let actual_rows = actual["data"].as_array().unwrap(); + let expected_rows = expected["data"].as_array().unwrap(); + assert_eq!(actual_rows.len(), expected_rows.len()); + for (actual, expected) in actual_rows.iter().zip(expected_rows) { + for field in ["result", "nullable_result"] { + if expected[field].is_null() { + assert!(actual[field].is_null()); + } else { + // Both declared columns are Float64; JSON 20 and 20.0 encode + // the same value despite different serde_json Number variants. + assert_eq!( + actual[field].as_f64().unwrap().to_bits(), + expected[field].as_f64().unwrap().to_bits(), + "{field}" + ); + } + } + } + assert_eq!( + actual["data"], + serde_json::json!([{ "result":null, "nullable_result":null },{ "result":0.0, "nullable_result":null },{ "result":0.0, "nullable_result":null },{ "result":7.5, "nullable_result":null },{ "result":20.0, "nullable_result":20.0 }]) + ); + let mut cleanup = client + .post(&clickhouse_url) + .body(format!("DROP TABLE default.{table}")); + if let Some(user) = &user { + cleanup = cleanup.basic_auth(user, password.as_ref()); + } + assert!(cleanup.send().await.unwrap().status().is_success()); +} From f711679831c35f09838f016f8545343847085c8d Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 11 Sep 2026 09:46:04 -0600 Subject: [PATCH 13/13] test(clickhouse): include nested tuple fields in SQL process replay --- Cargo.lock | 10 ++--- control_plane/Cargo.toml | 8 ++-- crates/asap_types/Cargo.toml | 2 +- data_plane/Cargo.toml | 6 +-- .../tests/clickhouse_differential_e2e.rs | 42 +++++++++++++------ 5 files changed, 43 insertions(+), 25 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 13a4be1f..75f85534 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -373,7 +373,7 @@ dependencies = [ [[package]] name = "asap-aware-mapping" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=cf1f9acf0e09ce31b8319e674fb2508185a31148#cf1f9acf0e09ce31b8319e674fb2508185a31148" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=23270ba33009953b6c1e293ba95c6a430ab8d222#23270ba33009953b6c1e293ba95c6a430ab8d222" dependencies = [ "asap-types", "serde", @@ -384,7 +384,7 @@ dependencies = [ [[package]] name = "asap-frontend-promql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=cf1f9acf0e09ce31b8319e674fb2508185a31148#cf1f9acf0e09ce31b8319e674fb2508185a31148" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=23270ba33009953b6c1e293ba95c6a430ab8d222#23270ba33009953b6c1e293ba95c6a430ab8d222" dependencies = [ "asap-types", "promql-parser", @@ -393,7 +393,7 @@ dependencies = [ [[package]] name = "asap-frontend-sql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=cf1f9acf0e09ce31b8319e674fb2508185a31148#cf1f9acf0e09ce31b8319e674fb2508185a31148" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=23270ba33009953b6c1e293ba95c6a430ab8d222#23270ba33009953b6c1e293ba95c6a430ab8d222" dependencies = [ "asap-sql-function-catalog", "asap-types", @@ -415,12 +415,12 @@ dependencies = [ [[package]] name = "asap-sql-function-catalog" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=cf1f9acf0e09ce31b8319e674fb2508185a31148#cf1f9acf0e09ce31b8319e674fb2508185a31148" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=23270ba33009953b6c1e293ba95c6a430ab8d222#23270ba33009953b6c1e293ba95c6a430ab8d222" [[package]] name = "asap-types" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=cf1f9acf0e09ce31b8319e674fb2508185a31148#cf1f9acf0e09ce31b8319e674fb2508185a31148" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=23270ba33009953b6c1e293ba95c6a430ab8d222#23270ba33009953b6c1e293ba95c6a430ab8d222" dependencies = [ "serde", "serde_json", diff --git a/control_plane/Cargo.toml b/control_plane/Cargo.toml index fb664da3..22ae1147 100644 --- a/control_plane/Cargo.toml +++ b/control_plane/Cargo.toml @@ -76,8 +76,8 @@ asap_types.workspace = true # scaffolding, unaware that `data_plane`'s `summary_executor.rs` in *this* # repo is a real one. Vendored locally instead of chased upstream -- see # `data_plane/src/query_engines/asap_query_engine/summary_exec.rs`. -planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "cf1f9acf0e09ce31b8319e674fb2508185a31148" } -asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "cf1f9acf0e09ce31b8319e674fb2508185a31148" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "23270ba33009953b6c1e293ba95c6a430ab8d222" } +asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "23270ba33009953b6c1e293ba95c6a430ab8d222" } # L1 adoption (design-target-architecture.md Part B): the PromQL front # end itself, replacing control_plane's own query_parser/promql.rs. @@ -85,8 +85,8 @@ asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = # `planner-types`/`asap-aware-mapping` above -- these three MUST move # together (two revs of the same upstream repo's types in one workspace # resolve to distinct Rust types that won't unify). -asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "cf1f9acf0e09ce31b8319e674fb2508185a31148" } -asap-frontend-sql = { git = "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/ProjectASAP/ASAPPlanner", rev = "cf1f9acf0e09ce31b8319e674fb2508185a31148" } +asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "23270ba33009953b6c1e293ba95c6a430ab8d222" } +asap-frontend-sql = { git = "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/ProjectASAP/ASAPPlanner", rev = "23270ba33009953b6c1e293ba95c6a430ab8d222" } [dev-dependencies] tokio = { version = "1", features = ["full", "test-util"] } diff --git a/crates/asap_types/Cargo.toml b/crates/asap_types/Cargo.toml index 16ed0aba..88ce9cef 100644 --- a/crates/asap_types/Cargo.toml +++ b/crates/asap_types/Cargo.toml @@ -34,4 +34,4 @@ sha2 = "0.10" # exactly (`control_plane/Cargo.toml`) -- two different revs of the same # git dependency in one workspace resolve to two distinct Rust types that # won't unify. -planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "cf1f9acf0e09ce31b8319e674fb2508185a31148" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "23270ba33009953b6c1e293ba95c6a430ab8d222" } diff --git a/data_plane/Cargo.toml b/data_plane/Cargo.toml index 536f24fe..39906783 100644 --- a/data_plane/Cargo.toml +++ b/data_plane/Cargo.toml @@ -39,8 +39,8 @@ sha2 = "0.10" # reduction: Reduction, .. }`) are `pre_asap` types, in the same crate now # (not a separate `asap-ir` import). Query serving consumes the compiled # QueryPlan; these types are used at physical-plan compilation boundaries. -planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "cf1f9acf0e09ce31b8319e674fb2508185a31148" } -asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "cf1f9acf0e09ce31b8319e674fb2508185a31148" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "23270ba33009953b6c1e293ba95c6a430ab8d222" } +asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "23270ba33009953b6c1e293ba95c6a430ab8d222" } # Shared external (workspace) serde.workspace = true @@ -133,7 +133,7 @@ fs2 = "0.4" # none of them. [dev-dependencies] -asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "cf1f9acf0e09ce31b8319e674fb2508185a31148" } +asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "23270ba33009953b6c1e293ba95c6a430ab8d222" } tempfile = "3.20.0" criterion = { version = "0.5", features = ["html_reports"] } tokio-tungstenite = "0.21" diff --git a/data_plane/tests/clickhouse_differential_e2e.rs b/data_plane/tests/clickhouse_differential_e2e.rs index 8023fd22..a470580e 100644 --- a/data_plane/tests/clickhouse_differential_e2e.rs +++ b/data_plane/tests/clickhouse_differential_e2e.rs @@ -492,7 +492,7 @@ async fn run_mixed_aggregate(aggregate: &str) { } #[tokio::test] -async fn array_sql_executes_local_element_after_typed_exact_leaf() { +async fn collection_sql_executes_local_elements_after_typed_exact_leaf() { use asap_types::query_plan::QueryPlanNode; use planner_types::{ post_asap::ValueOperation, @@ -507,8 +507,8 @@ async fn array_sql_executes_local_element_after_typed_exact_leaf() { let client = reqwest::Client::new(); let table = format!("asap_collection_elements_{}", std::process::id()); for sql in [ - format!("CREATE TABLE default.{table}(timestamp Int64, samples Array(Float64), nullable_samples Array(Nullable(Float64)), position Nullable(Int64)) ENGINE=Memory"), - format!("INSERT INTO default.{table} VALUES (100,[10,20],[10,20],-1),(200,[3.5],[3.5],9),(300,[],[],1),(400,[7.5],[NULL],1),(500,[1],[1],NULL),(2000,[999],[999],1)"), + format!("CREATE TABLE default.{table}(timestamp Int64, samples Array(Float64), nullable_samples Array(Nullable(Float64)), tuples Array(Tuple(ts Int64, value Nullable(Float64))), position Nullable(Int64)) ENGINE=Memory"), + format!("INSERT INTO default.{table} VALUES (100,[10,20],[10,20],[(1,5.5)],-1),(200,[3.5],[3.5],[(1,NULL)],9),(300,[],[],[],1),(400,[7.5],[NULL],[(1,3.5)],1),(500,[1],[1],[(1,9.5)],NULL),(2000,[999],[999],[(1,999)],1)"), ] { let mut request = client.post(&clickhouse_url).body(sql); if let Some(user) = &user { request = request.basic_auth(user, password.as_ref()); } @@ -517,7 +517,7 @@ async fn array_sql_executes_local_element_after_typed_exact_leaf() { let body = response.text().await.unwrap(); assert!(status.is_success(), "native collection fixture: {body}"); } - let sql = format!("SELECT arrayElement(samples, position) AS result, arrayElement(nullable_samples, position) AS nullable_result FROM {table} WHERE timestamp >= 0 AND timestamp < 2000 ORDER BY result NULLS FIRST"); + let sql = format!("SELECT arrayElement(samples, position) AS result, arrayElement(nullable_samples, position) AS nullable_result, tupleElement(arrayElement(tuples, 1), 'value') AS tuple_result FROM {table} WHERE timestamp >= 0 AND timestamp < 2000 ORDER BY result NULLS FIRST"); let mut workload = mixed_workload(&sql); workload.tables = HashMap::from([( table.clone(), @@ -538,6 +538,22 @@ async fn array_sql_executes_local_element_after_typed_exact_leaf() { }, false, ), + Column::new( + "tuples", + DataType::List { + element: Box::new(Column::new( + "item", + DataType::Struct { + fields: vec![ + Column::new("ts", DataType::Int64, false), + Column::new("value", DataType::Float64, true), + ], + }, + false, + )), + }, + false, + ), Column::new("position", DataType::Int64, true), ], 0, @@ -554,11 +570,13 @@ async fn array_sql_executes_local_element_after_typed_exact_leaf() { .nodes .values() .any(|node| matches!(node, QueryPlanNode::ExternalExact { .. }))); - assert!(entry.nodes.values().any(|node| { - let QueryPlanNode::Relational { operation, .. } = node else { return false; }; - let ValueOperation::Project { cols, .. } = serde_json::from_value(operation.clone()).unwrap() else { return false; }; - cols.iter().any(|column| matches!(&column.expr, QueryExpr::FunctionCall { name, .. } if name == "asap_element_access")) - }), "Planner-selected DAG must preserve local element evaluation"); + for function in ["asap_element_access", "asap_struct_field"] { + assert!(entry.nodes.values().any(|node| { + let QueryPlanNode::Relational { operation, .. } = node else { return false; }; + let ValueOperation::Project { cols, .. } = serde_json::from_value(operation.clone()).unwrap() else { return false; }; + cols.iter().any(|column| matches!(&column.expr, QueryExpr::FunctionCall { name, .. } if name == function)) + }), "Planner-selected DAG must preserve local {function} evaluation"); + } eprintln!( "collection Planner selection: {}", serde_json::to_string(&trace).unwrap() @@ -628,11 +646,11 @@ async fn array_sql_executes_local_element_after_typed_exact_leaf() { let expected_rows = expected["data"].as_array().unwrap(); assert_eq!(actual_rows.len(), expected_rows.len()); for (actual, expected) in actual_rows.iter().zip(expected_rows) { - for field in ["result", "nullable_result"] { + for field in ["result", "nullable_result", "tuple_result"] { if expected[field].is_null() { assert!(actual[field].is_null()); } else { - // Both declared columns are Float64; JSON 20 and 20.0 encode + // The declared result columns are Float64; JSON 20 and 20.0 encode // the same value despite different serde_json Number variants. assert_eq!( actual[field].as_f64().unwrap().to_bits(), @@ -644,7 +662,7 @@ async fn array_sql_executes_local_element_after_typed_exact_leaf() { } assert_eq!( actual["data"], - serde_json::json!([{ "result":null, "nullable_result":null },{ "result":0.0, "nullable_result":null },{ "result":0.0, "nullable_result":null },{ "result":7.5, "nullable_result":null },{ "result":20.0, "nullable_result":20.0 }]) + serde_json::json!([{ "result":null, "nullable_result":null, "tuple_result":9.5 },{ "result":0.0, "nullable_result":null, "tuple_result":null },{ "result":0.0, "nullable_result":null, "tuple_result":null },{ "result":7.5, "nullable_result":null, "tuple_result":3.5 },{ "result":20.0, "nullable_result":20.0, "tuple_result":5.5 }]) ); let mut cleanup = client .post(&clickhouse_url)