From dfc891a01a59327a76ae8f0706d29b1ca5fd9415 Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 11 Sep 2026 07:14:40 -0600 Subject: [PATCH 1/6] feat(clickhouse): evaluate typed map scalar expressions --- Cargo.lock | 10 +- control_plane/Cargo.toml | 8 +- crates/asap_types/Cargo.toml | 2 +- data_plane/Cargo.toml | 6 +- .../relational_adapter.rs | 317 +++++++++++++++++- 5 files changed, 315 insertions(+), 28 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 16fee0c0..7195756f 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=7261dab97f8e43070e6b9e33c85493b12b9a351c#7261dab97f8e43070e6b9e33c85493b12b9a351c" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=3026d772678bf71d83a6ff93e4adf43f418babec#3026d772678bf71d83a6ff93e4adf43f418babec" 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=7261dab97f8e43070e6b9e33c85493b12b9a351c#7261dab97f8e43070e6b9e33c85493b12b9a351c" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=3026d772678bf71d83a6ff93e4adf43f418babec#3026d772678bf71d83a6ff93e4adf43f418babec" 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=7261dab97f8e43070e6b9e33c85493b12b9a351c#7261dab97f8e43070e6b9e33c85493b12b9a351c" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=3026d772678bf71d83a6ff93e4adf43f418babec#3026d772678bf71d83a6ff93e4adf43f418babec" 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=7261dab97f8e43070e6b9e33c85493b12b9a351c#7261dab97f8e43070e6b9e33c85493b12b9a351c" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=3026d772678bf71d83a6ff93e4adf43f418babec#3026d772678bf71d83a6ff93e4adf43f418babec" [[package]] name = "asap-types" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=7261dab97f8e43070e6b9e33c85493b12b9a351c#7261dab97f8e43070e6b9e33c85493b12b9a351c" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=3026d772678bf71d83a6ff93e4adf43f418babec#3026d772678bf71d83a6ff93e4adf43f418babec" dependencies = [ "serde", "serde_json", diff --git a/control_plane/Cargo.toml b/control_plane/Cargo.toml index 537c5e96..44bb74ca 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 = "7261dab97f8e43070e6b9e33c85493b12b9a351c" } -asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "7261dab97f8e43070e6b9e33c85493b12b9a351c" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "3026d772678bf71d83a6ff93e4adf43f418babec" } +asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "3026d772678bf71d83a6ff93e4adf43f418babec" } # 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 = "7261dab97f8e43070e6b9e33c85493b12b9a351c" } -asap-frontend-sql = { git = "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/ProjectASAP/ASAPPlanner", rev = "7261dab97f8e43070e6b9e33c85493b12b9a351c" } +asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "3026d772678bf71d83a6ff93e4adf43f418babec" } +asap-frontend-sql = { git = "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/ProjectASAP/ASAPPlanner", rev = "3026d772678bf71d83a6ff93e4adf43f418babec" } [dev-dependencies] tokio = { version = "1", features = ["full", "test-util"] } diff --git a/crates/asap_types/Cargo.toml b/crates/asap_types/Cargo.toml index 21d80432..ac657d36 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 = "7261dab97f8e43070e6b9e33c85493b12b9a351c" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "3026d772678bf71d83a6ff93e4adf43f418babec" } diff --git a/data_plane/Cargo.toml b/data_plane/Cargo.toml index 14c29ec9..c5537c08 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 = "7261dab97f8e43070e6b9e33c85493b12b9a351c" } -asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "7261dab97f8e43070e6b9e33c85493b12b9a351c" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "3026d772678bf71d83a6ff93e4adf43f418babec" } +asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "3026d772678bf71d83a6ff93e4adf43f418babec" } # 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 = "7261dab97f8e43070e6b9e33c85493b12b9a351c" } +asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "3026d772678bf71d83a6ff93e4adf43f418babec" } 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 56cc44f8..e90db384 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 @@ -4,8 +4,8 @@ use std::{cmp::Ordering, collections::BTreeMap, sync::Arc}; use arrow::{ array::{ - ArrayRef, BooleanArray, Float64Array, Int64Array, MapArray, StringArray, StructArray, - TimestampMillisecondArray, + ArrayRef, BooleanArray, Float64Array, Int64Array, MapArray, NullArray, StringArray, + StructArray, TimestampMillisecondArray, }, datatypes::{DataType as ArrowDataType, Field, Schema}, record_batch::RecordBatch, @@ -64,6 +64,11 @@ 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::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,6 +312,8 @@ fn clickhouse_type_matches(actual: Option<&str>, expected: &DataType, nullable: return false; } match expected { + DataType::Null => actual == "Nothing", + DataType::List { .. } | DataType::Struct { .. } => false, DataType::Int64 => actual == "Int64", DataType::Float64 => actual == "Float64", DataType::Utf8 => actual == "String", @@ -346,13 +353,16 @@ impl ClickHouseRelationalAdapter { } _ => None, }; + let mut fields = left.fields.clone(); + fields.extend(right.fields.clone()); + let schema = scalar_schema(&fields); let mut rows = Vec::new(); for left_row in &left.rows { for right_row in &right.rows { let mut joined = Vec::with_capacity(left_row.len() + right_row.len()); joined.extend(left_row.iter().cloned()); joined.extend(right_row.iter().cloned()); - if matches!(eval(&pred.0, &joined)?, Cell::Bool(true)) { + if matches!(eval(&pred.0, &joined, &schema)?, Cell::Bool(true)) { rows.push(joined); } } @@ -369,10 +379,11 @@ impl ClickHouseRelationalAdapter { pred: &planner_types::pre_asap::Predicate, mut input: ClickHouseRelation, ) -> Result { + let schema = scalar_schema(&input.fields); input.rows = input .rows .into_iter() - .filter_map(|row| match eval(&pred.0, &row) { + .filter_map(|row| match eval(&pred.0, &row, &schema) { Ok(Cell::Bool(true)) => Some(Ok(row)), Ok(_) => None, Err(error) => Some(Err(error)), @@ -387,13 +398,14 @@ impl ClickHouseRelationalAdapter { output_schema: &SummarySchema, mut input: ClickHouseRelation, ) -> Result { + let schema = scalar_schema(&input.fields); match operation { ValueOperation::Project { cols, .. } => { let mut rows = Vec::with_capacity(input.rows.len()); for row in &input.rows { rows.push( cols.iter() - .map(|item| eval(&item.expr, row)) + .map(|item| eval(&item.expr, row, &schema)) .collect::, _>>()?, ); } @@ -408,12 +420,12 @@ impl ClickHouseRelationalAdapter { } for row in &input.rows { for key in keys { - eval(&key.expr, row)?; + eval(&key.expr, row, &schema)?; } } input .rows - .sort_by(|left, right| compare_sort_keys(left, right, keys)); + .sort_by(|left, right| compare_sort_keys(left, right, keys, &schema)); } ValueOperation::Limit { n, offset } => { input.rows = input.rows.into_iter().skip(*offset).take(*n).collect(); @@ -571,7 +583,22 @@ fn row_from_value( .collect() } -fn eval(expr: &QueryExpr, row: &[Cell]) -> Result { +fn scalar_schema(fields: &[(String, DataType, bool)]) -> planner_types::pre_asap::Schema { + planner_types::pre_asap::Schema::new( + fields + .iter() + .map(|(name, dtype, nullable)| { + planner_types::pre_asap::Column::new(name.clone(), dtype.clone(), *nullable) + }) + .collect(), + ) +} + +fn eval( + expr: &QueryExpr, + row: &[Cell], + schema: &planner_types::pre_asap::Schema, +) -> Result { match expr { QueryExpr::Column(index) => { row.get(*index) @@ -589,12 +616,89 @@ fn eval(expr: &QueryExpr, row: &[Cell]) -> Result Cell::Null, }), QueryExpr::Compare { left, op, right } => { - let left = eval(left, row)?; - let right = eval(right, row)?; + let left = eval(left, row, schema)?; + let right = eval(right, row, schema)?; compare(op, left, right) } QueryExpr::Arithmetic { op, left, right } => { - arithmetic(op, eval(left, row)?, eval(right, row)?) + arithmetic(op, eval(left, row, schema)?, eval(right, row, schema)?) + } + QueryExpr::FunctionCall { name, args } => { + use planner_types::pre_asap::scalar_signature::MapScalarFunction; + let function = MapScalarFunction::from_name(name).ok_or_else(|| { + ClickHouseRelationalError::Unsupported(format!("scalar function {name}")) + })?; + expr.scalar_type(schema) + .map_err(|error| ClickHouseRelationalError::Invalid(error.to_string()))?; + let values = args + .iter() + .map(|arg| eval(arg, row, schema)) + .collect::, _>>()?; + match function { + MapScalarFunction::Construct => { + let mut values = values.into_iter(); + let mut entries = Vec::new(); + while let Some(key) = values.next() { + if !matches!(key, Cell::Int64(_) | Cell::Utf8(_) | Cell::Bool(_)) { + return Err(ClickHouseRelationalError::Unsupported( + "map key value type".into(), + )); + } + entries.push(( + key, + values.next().ok_or_else(|| { + ClickHouseRelationalError::Invalid("odd map argument count".into()) + })?, + )); + } + Ok(Cell::Map(entries)) + } + MapScalarFunction::Concat => { + let mut entries = Vec::new(); + for value in values { + let Cell::Map(mut next) = value else { + return Err(ClickHouseRelationalError::Invalid( + "map concat argument".into(), + )); + }; + entries.append(&mut next); + } + Ok(Cell::Map(entries)) + } + MapScalarFunction::Access => { + let [Cell::Map(entries), key] = values.as_slice() else { + return Err(ClickHouseRelationalError::Invalid( + "map access arguments".into(), + )); + }; + if matches!(key, Cell::Null) { + return Ok(Cell::Null); + } + if !matches!(key, Cell::Int64(_) | Cell::Utf8(_) | Cell::Bool(_)) { + return Err(ClickHouseRelationalError::Unsupported( + "map lookup key type".into(), + )); + } + if let Some((_, value)) = entries.iter().find(|(candidate, _)| candidate == key) + { + return Ok(value.clone()); + } + let ( + DataType::Map { + value, + value_nullable, + .. + }, + _, + ) = args[0] + .scalar_type(schema) + .map_err(|error| ClickHouseRelationalError::Invalid(error.to_string()))? + else { + unreachable!() + }; + default_map_value(&value, value_nullable) + } + } } other => Err(ClickHouseRelationalError::Unsupported(format!( "scalar expression {other:?}" @@ -602,6 +706,25 @@ fn eval(expr: &QueryExpr, row: &[Cell]) -> Result Result { + if nullable { + return Ok(Cell::Null); + } + Ok(match dtype { + DataType::Null => Cell::Null, + DataType::Int64 => Cell::Int64(0), + DataType::Float64 => Cell::Float64(0.0), + DataType::Utf8 => Cell::Utf8(String::new()), + DataType::Bool => Cell::Bool(false), + DataType::Map { .. } => Cell::Map(Vec::new()), + _ => { + return Err(ClickHouseRelationalError::Unsupported( + "map missing-key default type".into(), + )) + } + }) +} + fn compare(op: &CompareOpKind, left: Cell, right: Cell) -> Result { if matches!(left, Cell::Null) || matches!(right, Cell::Null) { return Ok(Cell::Null); @@ -633,6 +756,22 @@ fn arithmetic( if matches!(left, Cell::Null) || matches!(right, Cell::Null) { return Ok(Cell::Null); } + if let (Cell::Int64(left), Cell::Int64(right)) = (&left, &right) { + let integer = match op { + ArithmeticOpKind::Add => Some(left.checked_add(*right)), + ArithmeticOpKind::Sub => Some(left.checked_sub(*right)), + ArithmeticOpKind::Mul => Some(left.checked_mul(*right)), + ArithmeticOpKind::Mod => Some(left.checked_rem(*right)), + _ => None, + }; + if let Some(value) = integer { + return value.map(Cell::Int64).ok_or_else(|| { + ClickHouseRelationalError::Invalid( + "integer arithmetic overflow or zero divisor".into(), + ) + }); + } + } let (left, right) = match (left, right) { (Cell::Int64(left), Cell::Int64(right)) => (left as f64, right as f64), (Cell::Int64(left), Cell::Float64(right)) => (left as f64, right), @@ -660,12 +799,17 @@ fn arithmetic( Ok(Cell::Float64(value)) } -fn compare_sort_keys(left: &[Cell], right: &[Cell], keys: &[SortKey]) -> Ordering { +fn compare_sort_keys( + left: &[Cell], + right: &[Cell], + keys: &[SortKey], + schema: &planner_types::pre_asap::Schema, +) -> Ordering { for key in keys { - let Ok(left) = eval(&key.expr, left) else { + let Ok(left) = eval(&key.expr, left, schema) else { return Ordering::Equal; }; - let Ok(right) = eval(&key.expr, right) else { + let Ok(right) = eval(&key.expr, right, schema) else { return Ordering::Equal; }; let (ordering, order_depends_on_direction) = match (&left, &right) { @@ -739,6 +883,19 @@ fn intersect_coverage(current: Option<(u64, u64)>, next: Option<(u64, u64)>) -> fn arrow_type(dtype: &DataType) -> ArrowDataType { match dtype { + DataType::Null => ArrowDataType::Null, + DataType::List { element } => ArrowDataType::List(Arc::new(Field::new( + &element.name, + arrow_type(&element.dtype), + element.nullable, + ))), + DataType::Struct { fields } => ArrowDataType::Struct( + fields + .iter() + .map(|field| Field::new(&field.name, arrow_type(&field.dtype), field.nullable)) + .collect::>() + .into(), + ), DataType::Int64 => ArrowDataType::Int64, DataType::Float64 => ArrowDataType::Float64, DataType::Utf8 => ArrowDataType::Utf8, @@ -786,6 +943,22 @@ fn build_array( }}; } Ok(match dtype { + DataType::Null => { + if rows + .iter() + .any(|row| !matches!(row.get(column), Some(Cell::Null))) + { + return Err(ClickHouseRelationalError::Invalid( + "non-null value in bottom-typed column".into(), + )); + } + Arc::new(NullArray::new(rows.len())) as ArrayRef + } + DataType::List { .. } | DataType::Struct { .. } => { + return Err(ClickHouseRelationalError::Unsupported( + "collection value transport".into(), + )) + } DataType::Int64 => Arc::new(Int64Array::from(values!(Int64))) as ArrayRef, DataType::Float64 => Arc::new(Float64Array::from(values!(Float64))) as ArrayRef, DataType::Utf8 => Arc::new(StringArray::from(values!(Utf8))) as ArrayRef, @@ -1054,7 +1227,12 @@ mod tests { #[test] fn unsupported_scalar_expression_fails_closed() { let row = vec![Cell::Float64(1.0)]; - let error = eval(&QueryExpr::BoolAnd(vec![]), &row).unwrap_err(); + let error = eval( + &QueryExpr::BoolAnd(vec![]), + &row, + &planner_types::pre_asap::Schema::new(vec![]), + ) + .unwrap_err(); assert!(matches!(error, ClickHouseRelationalError::Unsupported(_))); } @@ -1124,3 +1302,112 @@ mod tests { ); } } + +#[cfg(test)] +mod scalar_contract_tests { + use super::*; + use planner_types::pre_asap::{Column, Schema}; + + fn function(name: &str, args: Vec) -> QueryExpr { + QueryExpr::FunctionCall { + name: name.into(), + args, + } + } + fn text(value: &str) -> QueryExpr { + QueryExpr::Literal(ScalarValue::Utf8(value.into())) + } + + #[test] + fn map_access_uses_declared_default_and_first_duplicate() { + let dtype = DataType::Map { + key: Box::new(DataType::Utf8), + value: Box::new(DataType::Int64), + value_nullable: false, + }; + let schema = Schema::new(vec![Column::new("m", dtype, false)]); + let access = function("asap_map_access", vec![QueryExpr::Column(0), text("a")]); + assert_eq!( + eval(&access, &[Cell::Map(vec![])], &schema).unwrap(), + Cell::Int64(0) + ); + assert_eq!( + eval( + &access, + &[Cell::Map(vec![ + (Cell::Utf8("a".into()), Cell::Int64(7)), + (Cell::Utf8("a".into()), Cell::Int64(9)) + ])], + &schema + ) + .unwrap(), + Cell::Int64(7) + ); + let nullable = Schema::new(vec![Column::new( + "m", + DataType::Map { + key: Box::new(DataType::Utf8), + value: Box::new(DataType::Int64), + value_nullable: true, + }, + false, + )]); + assert_eq!( + eval(&access, &[Cell::Map(vec![])], &nullable).unwrap(), + Cell::Null + ); + } + + #[test] + fn map_concat_preserves_duplicates_and_empty_map() { + let map = |value| { + function( + "map", + vec![text("a"), QueryExpr::Literal(ScalarValue::Int64(value))], + ) + }; + let concat = function("mapConcat", vec![function("map", vec![]), map(7), map(9)]); + let schema = Schema::new(vec![]); + assert_eq!( + eval(&concat, &[], &schema).unwrap(), + Cell::Map(vec![ + (Cell::Utf8("a".into()), Cell::Int64(7)), + (Cell::Utf8("a".into()), Cell::Int64(9)) + ]) + ); + let mixed = function( + "map", + vec![ + text("a"), + QueryExpr::Literal(ScalarValue::Int64(1)), + text("b"), + QueryExpr::Literal(ScalarValue::Float64(2.5)), + ], + ); + assert!(eval(&mixed, &[], &schema).is_err()); + } + + #[test] + fn integer_modulo_never_rounds_through_float() { + assert_eq!( + arithmetic( + &ArithmeticOpKind::Mod, + Cell::Int64(9_007_199_254_740_993), + Cell::Int64(2) + ) + .unwrap(), + Cell::Int64(1) + ); + assert_eq!( + arithmetic(&ArithmeticOpKind::Mod, Cell::Int64(-7), Cell::Int64(3)).unwrap(), + Cell::Int64(-1) + ); + assert!(arithmetic(&ArithmeticOpKind::Mod, Cell::Int64(7), Cell::Int64(0)).is_err()); + assert!(arithmetic( + &ArithmeticOpKind::Mod, + Cell::Int64(i64::MIN), + Cell::Int64(-1) + ) + .is_err()); + } +} From f66462bf852b4aeef48340aef196ed9365e0f246 Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 11 Sep 2026 07:23:54 -0600 Subject: [PATCH 2/6] test(clickhouse): retain null lookup-key semantics --- .../asap_clickhouse_query_engine/relational_adapter.rs | 8 ++++++++ 1 file changed, 8 insertions(+) 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 e90db384..069275ad 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 @@ -1356,6 +1356,14 @@ mod scalar_contract_tests { eval(&access, &[Cell::Map(vec![])], &nullable).unwrap(), Cell::Null ); + let null_key = function( + "asap_map_access", + vec![QueryExpr::Column(0), QueryExpr::Literal(ScalarValue::Null)], + ); + assert_eq!( + eval(&null_key, &[Cell::Map(vec![])], &schema).unwrap(), + Cell::Null + ); } #[test] From ed1d17b7304553b5b9c2b80ac412fa295d26dcb2 Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 11 Sep 2026 07:34:16 -0600 Subject: [PATCH 3/6] fix(clickhouse): preserve mixed numeric comparison precision --- Cargo.lock | 10 ++-- control_plane/Cargo.toml | 8 +-- crates/asap_types/Cargo.toml | 2 +- data_plane/Cargo.toml | 6 +- .../relational_adapter.rs | 55 ++++++++++++++++++- 5 files changed, 65 insertions(+), 16 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 7195756f..76616c37 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=3026d772678bf71d83a6ff93e4adf43f418babec#3026d772678bf71d83a6ff93e4adf43f418babec" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=b2b05628dd9a58db555ba309bf7201e10513f4d2#b2b05628dd9a58db555ba309bf7201e10513f4d2" 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=3026d772678bf71d83a6ff93e4adf43f418babec#3026d772678bf71d83a6ff93e4adf43f418babec" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=b2b05628dd9a58db555ba309bf7201e10513f4d2#b2b05628dd9a58db555ba309bf7201e10513f4d2" 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=3026d772678bf71d83a6ff93e4adf43f418babec#3026d772678bf71d83a6ff93e4adf43f418babec" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=b2b05628dd9a58db555ba309bf7201e10513f4d2#b2b05628dd9a58db555ba309bf7201e10513f4d2" 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=3026d772678bf71d83a6ff93e4adf43f418babec#3026d772678bf71d83a6ff93e4adf43f418babec" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=b2b05628dd9a58db555ba309bf7201e10513f4d2#b2b05628dd9a58db555ba309bf7201e10513f4d2" [[package]] name = "asap-types" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=3026d772678bf71d83a6ff93e4adf43f418babec#3026d772678bf71d83a6ff93e4adf43f418babec" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=b2b05628dd9a58db555ba309bf7201e10513f4d2#b2b05628dd9a58db555ba309bf7201e10513f4d2" dependencies = [ "serde", "serde_json", diff --git a/control_plane/Cargo.toml b/control_plane/Cargo.toml index 44bb74ca..2aec4683 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 = "3026d772678bf71d83a6ff93e4adf43f418babec" } -asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "3026d772678bf71d83a6ff93e4adf43f418babec" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "b2b05628dd9a58db555ba309bf7201e10513f4d2" } +asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "b2b05628dd9a58db555ba309bf7201e10513f4d2" } # 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 = "3026d772678bf71d83a6ff93e4adf43f418babec" } -asap-frontend-sql = { git = "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/ProjectASAP/ASAPPlanner", rev = "3026d772678bf71d83a6ff93e4adf43f418babec" } +asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "b2b05628dd9a58db555ba309bf7201e10513f4d2" } +asap-frontend-sql = { git = "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/ProjectASAP/ASAPPlanner", rev = "b2b05628dd9a58db555ba309bf7201e10513f4d2" } [dev-dependencies] tokio = { version = "1", features = ["full", "test-util"] } diff --git a/crates/asap_types/Cargo.toml b/crates/asap_types/Cargo.toml index ac657d36..a50afa72 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 = "3026d772678bf71d83a6ff93e4adf43f418babec" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "b2b05628dd9a58db555ba309bf7201e10513f4d2" } diff --git a/data_plane/Cargo.toml b/data_plane/Cargo.toml index c5537c08..6bf88325 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 = "3026d772678bf71d83a6ff93e4adf43f418babec" } -asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "3026d772678bf71d83a6ff93e4adf43f418babec" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "b2b05628dd9a58db555ba309bf7201e10513f4d2" } +asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "b2b05628dd9a58db555ba309bf7201e10513f4d2" } # 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 = "3026d772678bf71d83a6ff93e4adf43f418babec" } +asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "b2b05628dd9a58db555ba309bf7201e10513f4d2" } 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 069275ad..d9ef0a8c 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 @@ -420,7 +420,12 @@ impl ClickHouseRelationalAdapter { } for row in &input.rows { for key in keys { - eval(&key.expr, row, &schema)?; + if matches!(eval(&key.expr, row, &schema)?, Cell::Float64(value) if value.is_nan()) + { + return Err(ClickHouseRelationalError::Unsupported( + "NaN sort key".into(), + )); + } } } input @@ -842,12 +847,32 @@ fn compare_sort_keys( Ordering::Equal } +fn integer_float_cmp(integer: i64, float: f64) -> Option { + if float.is_nan() { + return None; + } + // These bounds are powers of two, exactly representable as Float64. + if float >= 9_223_372_036_854_775_808.0 { + return Some(Ordering::Less); + } + if float < -9_223_372_036_854_775_808.0 { + return Some(Ordering::Greater); + } + let integral = float as i64; + match integer.cmp(&integral) { + Ordering::Equal => 0.0_f64.partial_cmp(&float.fract()), + other => Some(other), + } +} + fn cell_cmp(left: &Cell, right: &Cell) -> Option { match (left, right) { (Cell::Int64(left), Cell::Int64(right)) => Some(left.cmp(right)), (Cell::Float64(left), Cell::Float64(right)) => left.partial_cmp(right), - (Cell::Int64(left), Cell::Float64(right)) => (*left as f64).partial_cmp(right), - (Cell::Float64(left), Cell::Int64(right)) => left.partial_cmp(&(*right as f64)), + (Cell::Int64(left), Cell::Float64(right)) => integer_float_cmp(*left, *right), + (Cell::Float64(left), Cell::Int64(right)) => { + integer_float_cmp(*right, *left).map(Ordering::reverse) + } (Cell::Utf8(left), Cell::Utf8(right)) => Some(left.cmp(right)), (Cell::Bool(left), Cell::Bool(right)) => Some(left.cmp(right)), (Cell::Timestamp(left), Cell::Timestamp(right)) => Some(left.cmp(right)), @@ -1395,6 +1420,30 @@ mod scalar_contract_tests { assert!(eval(&mixed, &[], &schema).is_err()); } + #[test] + fn mixed_comparison_preserves_integer_precision_and_boundaries() { + assert_eq!( + integer_float_cmp(9_007_199_254_740_993, 9_007_199_254_740_992.0), + Some(Ordering::Greater) + ); + assert_eq!( + integer_float_cmp(i64::MAX, 9_223_372_036_854_775_808.0), + Some(Ordering::Less) + ); + assert_eq!( + integer_float_cmp(i64::MIN, -9_223_372_036_854_775_808.0), + Some(Ordering::Equal) + ); + assert_eq!(integer_float_cmp(-1, -1.5), Some(Ordering::Greater)); + assert_eq!(integer_float_cmp(1, 1.5), Some(Ordering::Less)); + assert_eq!(integer_float_cmp(0, f64::INFINITY), Some(Ordering::Less)); + assert_eq!( + integer_float_cmp(0, f64::NEG_INFINITY), + Some(Ordering::Greater) + ); + assert_eq!(integer_float_cmp(0, f64::NAN), None); + } + #[test] fn integer_modulo_never_rounds_through_float() { assert_eq!( From 487a7c532269193a5020dd8b2d3e73d4e267663b Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 11 Sep 2026 07:36:20 -0600 Subject: [PATCH 4/6] fix(clickhouse): reject unordered nested Map sort keys --- .../relational_adapter.rs | 49 ++++++++++++++++++- 1 file changed, 47 insertions(+), 2 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 d9ef0a8c..d14e5719 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 @@ -420,8 +420,7 @@ impl ClickHouseRelationalAdapter { } for row in &input.rows { for key in keys { - if matches!(eval(&key.expr, row, &schema)?, Cell::Float64(value) if value.is_nan()) - { + if contains_nan(&eval(&key.expr, row, &schema)?) { return Err(ClickHouseRelationalError::Unsupported( "NaN sort key".into(), )); @@ -847,6 +846,16 @@ fn compare_sort_keys( Ordering::Equal } +fn contains_nan(value: &Cell) -> bool { + match value { + Cell::Float64(value) => value.is_nan(), + Cell::Map(entries) => entries + .iter() + .any(|(key, value)| contains_nan(key) || contains_nan(value)), + _ => false, + } +} + fn integer_float_cmp(integer: i64, float: f64) -> Option { if float.is_nan() { return None; @@ -1420,6 +1429,42 @@ mod scalar_contract_tests { assert!(eval(&mixed, &[], &schema).is_err()); } + #[test] + fn sorting_nested_nan_fails_before_comparator_can_treat_it_as_equal() { + let dtype = DataType::Map { + key: Box::new(DataType::Utf8), + value: Box::new(DataType::Float64), + value_nullable: false, + }; + let input = ClickHouseRelation { + rows: vec![vec![Cell::Map(vec![( + Cell::Utf8("a".into()), + Cell::Float64(f64::NAN), + )])]], + fields: vec![("m".into(), dtype.clone(), false)], + coverage: None, + }; + let schema = SummarySchema { + fields: vec![planner_types::post_asap::SummaryField { + name: "m".into(), + dtype: SummaryFamilyType::Plain(dtype), + nullable: false, + }], + time_index: None, + }; + let operation = ValueOperation::Sort { + keys: vec![SortKey { + expr: QueryExpr::Column(0), + ascending: true, + nulls_first: false, + }], + partition_by: planner_types::pre_asap::GroupKeys::none(), + }; + assert!(ClickHouseRelationalAdapter + .apply_operation(&operation, &schema, input) + .is_err()); + } + #[test] fn mixed_comparison_preserves_integer_precision_and_boundaries() { assert_eq!( From d97490b1636b4092795a0207114fd3760391d071 Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 11 Sep 2026 07:50:47 -0600 Subject: [PATCH 5/6] feat(clickhouse): render canonical exact subtrees --- Cargo.lock | 10 +- control_plane/Cargo.toml | 8 +- control_plane/src/query_plan.rs | 59 +-- .../src/query_plan/clickhouse_exact.rs | 372 ++++++++++++++++++ .../tests/fixtures/sql_exact_cuts/q07.sql | 1 + .../tests/fixtures/sql_exact_cuts/q09.sql | 1 + .../tests/fixtures/sql_exact_cuts/q12.sql | 1 + .../tests/fixtures/sql_exact_cuts/q27.sql | 1 + crates/asap_types/Cargo.toml | 2 +- data_plane/Cargo.toml | 6 +- 10 files changed, 409 insertions(+), 52 deletions(-) create mode 100644 control_plane/src/query_plan/clickhouse_exact.rs create mode 100644 control_plane/tests/fixtures/sql_exact_cuts/q07.sql create mode 100644 control_plane/tests/fixtures/sql_exact_cuts/q09.sql create mode 100644 control_plane/tests/fixtures/sql_exact_cuts/q12.sql create mode 100644 control_plane/tests/fixtures/sql_exact_cuts/q27.sql diff --git a/Cargo.lock b/Cargo.lock index 76616c37..5b68bc67 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=b2b05628dd9a58db555ba309bf7201e10513f4d2#b2b05628dd9a58db555ba309bf7201e10513f4d2" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=0402384e589df6e087d6d2d22b463ddc2eea0774#0402384e589df6e087d6d2d22b463ddc2eea0774" 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=b2b05628dd9a58db555ba309bf7201e10513f4d2#b2b05628dd9a58db555ba309bf7201e10513f4d2" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=0402384e589df6e087d6d2d22b463ddc2eea0774#0402384e589df6e087d6d2d22b463ddc2eea0774" 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=b2b05628dd9a58db555ba309bf7201e10513f4d2#b2b05628dd9a58db555ba309bf7201e10513f4d2" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=0402384e589df6e087d6d2d22b463ddc2eea0774#0402384e589df6e087d6d2d22b463ddc2eea0774" 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=b2b05628dd9a58db555ba309bf7201e10513f4d2#b2b05628dd9a58db555ba309bf7201e10513f4d2" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=0402384e589df6e087d6d2d22b463ddc2eea0774#0402384e589df6e087d6d2d22b463ddc2eea0774" [[package]] name = "asap-types" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=b2b05628dd9a58db555ba309bf7201e10513f4d2#b2b05628dd9a58db555ba309bf7201e10513f4d2" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=0402384e589df6e087d6d2d22b463ddc2eea0774#0402384e589df6e087d6d2d22b463ddc2eea0774" dependencies = [ "serde", "serde_json", diff --git a/control_plane/Cargo.toml b/control_plane/Cargo.toml index 2aec4683..7cb29606 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 = "b2b05628dd9a58db555ba309bf7201e10513f4d2" } -asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "b2b05628dd9a58db555ba309bf7201e10513f4d2" } +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" } # 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 = "b2b05628dd9a58db555ba309bf7201e10513f4d2" } -asap-frontend-sql = { git = "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/ProjectASAP/ASAPPlanner", rev = "b2b05628dd9a58db555ba309bf7201e10513f4d2" } +asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "0402384e589df6e087d6d2d22b463ddc2eea0774" } +asap-frontend-sql = { git = "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/ProjectASAP/ASAPPlanner", rev = "0402384e589df6e087d6d2d22b463ddc2eea0774" } [dev-dependencies] tokio = { version = "1", features = ["full", "test-util"] } diff --git a/control_plane/src/query_plan.rs b/control_plane/src/query_plan.rs index af3078f0..58ec03f2 100644 --- a/control_plane/src/query_plan.rs +++ b/control_plane/src/query_plan.rs @@ -5,6 +5,7 @@ //! node IDs. Serving executes this graph without reconstructing Planner IR or //! searching for compatible materializations. +mod clickhouse_exact; pub mod logical; use std::collections::{BTreeMap, BTreeSet}; @@ -1324,48 +1325,28 @@ where } } SummaryExpr::KeepPreAsap(expr) if self.preserve_relational => { - let planner_types::pre_asap::QueryExpr::Scan { - source: planner_types::pre_asap::Source::Table { table_ref }, - predicates, - schema, + let mut expression = + clickhouse_exact::render(expr).map_err(QueryPlanError::UnsupportedNode)?; + let mut bounded = false; + if let planner_types::pre_asap::QueryExpr::Scan { + predicates, schema, .. } = expr.as_ref() - else { - return Err(QueryPlanError::UnsupportedNode( - "SQL exact cut is not a direct table scan".into(), - )); - }; - if !predicates.is_empty() { - return Err(QueryPlanError::UnsupportedNode( - "SQL exact table cut contains unrendered predicates".into(), - )); - } - fn quoted(identifier: &str) -> String { - identifier - .split('.') - .map(|part| format!("`{}`", part.replace('`', "``"))) - .collect::>() - .join(".") + { + if predicates.is_empty() { + if let Some(column) = schema.time_index.and_then(|i| schema.columns.get(i)) + { + let name = format!("`{}`", column.name.replace('`', "``")); + expression.push_str(&format!( + " WHERE {name} >= {{from:UInt64}} AND {name} <= {{to:UInt64}}" + )); + bounded = true; + } + } } - let columns = schema - .columns - .iter() - .map(|column| quoted(&column.name)) - .collect::>() - .join(", "); - let time_filter = schema.time_index.and_then(|index| { - schema.columns.get(index).map(|source_column| { - let column = quoted(&source_column.name); - format!(" WHERE {column} >= {{from:UInt64}} AND {column} <= {{to:UInt64}}") - }) - }); QueryPlanNode::ExternalExact { request: ExternalExactRequest { language: QueryLanguage::ClickHouseSql, - expression: format!( - "SELECT {columns} FROM {}{}", - quoted(table_ref), - time_filter.unwrap_or_default() - ), + expression, output: ExternalExactOutput::Relation { schema: serde_json::to_value(&node.schema).map_err(|error| { QueryPlanError::Invalid(format!( @@ -1374,8 +1355,8 @@ where })?, }, parameters: BTreeMap::new(), - start_parameter: Some("from".into()), - end_parameter: Some("to".into()), + start_parameter: bounded.then(|| "from".into()), + end_parameter: bounded.then(|| "to".into()), input_contracts: Vec::new(), }, inputs: Vec::new(), diff --git a/control_plane/src/query_plan/clickhouse_exact.rs b/control_plane/src/query_plan/clickhouse_exact.rs new file mode 100644 index 00000000..9190b3c7 --- /dev/null +++ b/control_plane/src/query_plan/clickhouse_exact.rs @@ -0,0 +1,372 @@ +//! Render a supported canonical relational cut without changing its row population. +//! Unsupported operators remain admission errors, never guessed SQL semantics. +use planner_types::pre_asap::{ + AggIntent, ArithmeticOpKind, CompareOpKind, QueryExpr, Reduction, ScalarValue, Schema, Source, +}; + +fn quoted(name: &str) -> String { + format!("`{}`", name.replace('`', "``")) +} +fn column(index: usize, schema: &Schema) -> Result { + schema + .columns + .get(index) + .map(|c| quoted(&c.name)) + .ok_or_else(|| format!("unresolved exact column {index}")) +} +fn scalar(expr: &QueryExpr, schema: &Schema) -> Result { + Ok(match expr { + QueryExpr::Column(id) => column(*id, schema)?, + QueryExpr::Literal(value) => match value { + ScalarValue::Int64(v) => v.to_string(), + ScalarValue::Float64(v) if v.is_finite() => format!("toFloat64('{}')", v), + ScalarValue::Utf8(v) => format!("'{}'", v.replace('\\', "\\\\").replace('\'', "\\'")), + ScalarValue::Boolean(v) => if *v { "true" } else { "false" }.into(), + ScalarValue::Null => "NULL".into(), + _ => return Err("nonfinite exact literal".into()), + }, + QueryExpr::Arithmetic { op, left, right } => { + let op = match op { + ArithmeticOpKind::Add => "+", + ArithmeticOpKind::Sub => "-", + ArithmeticOpKind::Mul => "*", + ArithmeticOpKind::Div => "/", + ArithmeticOpKind::Mod => "%", + _ => return Err("unsupported exact arithmetic".into()), + }; + format!( + "({} {op} {})", + scalar(left, schema)?, + scalar(right, schema)? + ) + } + QueryExpr::Compare { op, left, right } => { + let op = match op { + CompareOpKind::Eq => "=", + CompareOpKind::Ne => "!=", + CompareOpKind::Lt => "<", + CompareOpKind::Le => "<=", + CompareOpKind::Gt => ">", + CompareOpKind::Ge => ">=", + _ => return Err("unsupported exact comparison".into()), + }; + format!( + "({} {op} {})", + scalar(left, schema)?, + scalar(right, schema)? + ) + } + QueryExpr::BoolAnd(args) | QueryExpr::BoolOr(args) => { + let and = matches!(expr, QueryExpr::BoolAnd(_)); + if args.is_empty() { + if and { "true" } else { "false" }.into() + } else { + format!( + "({})", + args.iter() + .map(|e| scalar(e, schema)) + .collect::, _>>()? + .join(if and { " AND " } else { " OR " }) + ) + } + } + QueryExpr::Not(arg) => format!("NOT ({})", scalar(arg, schema)?), + QueryExpr::IsNull(arg) => format!("({} IS NULL)", scalar(arg, schema)?), + QueryExpr::IsNotNull(arg) => format!("({} IS NOT NULL)", scalar(arg, schema)?), + QueryExpr::FunctionCall { name, args } => { + let function = match name.as_str() { + "map" => "map", + "mapconcat" => "mapConcat", + "asap_map_access" => "arrayElement", + _ => return Err(format!("unsupported exact scalar function {name}")), + }; + expr.scalar_type(schema).map_err(|e| e.to_string())?; + format!( + "{function}({})", + args.iter() + .map(|e| scalar(e, schema)) + .collect::, _>>()? + .join(", ") + ) + } + _ => return Err("unsupported exact scalar expression".into()), + }) +} +fn aggregate(intent: &AggIntent, schema: &Schema) -> Result { + if let Some((arg, order)) = intent.arg_selector_columns(schema)? { + let AggIntent::Extension { ext_kind, .. } = intent else { + unreachable!() + }; + let function = if ext_kind == "arg_max" { + "argMax" + } else { + "argMin" + }; + return Ok(format!( + "{function}({}, {})", + column(arg, schema)?, + column(order, schema)? + )); + } + let (function, col) = match intent { + AggIntent::Count { .. } => return Ok("count()".into()), + AggIntent::Sum { col } => ("sum", col), + AggIntent::Min { col } => ("min", col), + AggIntent::Max { col } => ("max", col), + AggIntent::Avg { col } => ("avg", col), + _ => return Err("unsupported exact aggregate contract".into()), + }; + Ok(format!( + "{function}({})", + column( + col.ok_or("SQL aggregate requires explicit input column")?, + schema + )? + )) +} +/// Composite cuts retain their own canonical predicates; caller bounds are not +/// injected into descendant scans (which may belong to independent windows). +pub(super) fn render(expr: &QueryExpr) -> Result { + let output = expr.output_schema().map_err(|e| e.to_string())?; + // SQL names must identify a unique positional field at each nested boundary. + let mut names = std::collections::BTreeSet::new(); + if output.columns.iter().any(|c| !names.insert(&c.name)) { + return Err("ambiguous exact output column names".into()); + } + match expr { + QueryExpr::Scan { + source: Source::Table { table_ref }, + predicates, + schema, + } => { + let table = table_ref + .split('.') + .map(quoted) + .collect::>() + .join("."); + let columns = schema + .columns + .iter() + .map(|c| quoted(&c.name)) + .collect::>() + .join(", "); + let filters = predicates + .iter() + .map(|p| scalar(&p.0, schema)) + .collect::, _>>()?; + Ok(format!( + "SELECT {columns} FROM {table}{}", + if filters.is_empty() { + String::new() + } else { + format!(" WHERE {}", filters.join(" AND ")) + } + )) + } + QueryExpr::Filter { pred, child } => { + let schema = child.output_schema().map_err(|e| e.to_string())?; + Ok(format!( + "SELECT * FROM ({}) WHERE {}", + render(child)?, + scalar(&pred.0, &schema)? + )) + } + QueryExpr::Project { cols, child, .. } => { + let schema = child.output_schema().map_err(|e| e.to_string())?; + if cols.len() != output.columns.len() { + return Err("exact projection width mismatch".into()); + } + let columns = cols + .iter() + .zip(&output.columns) + .map(|(item, col)| { + Ok(format!( + "{} AS {}", + scalar(&item.expr, &schema)?, + quoted(&col.name) + )) + }) + .collect::, String>>()?; + Ok(format!( + "SELECT {} FROM ({})", + columns.join(", "), + render(child)? + )) + } + QueryExpr::Aggregate { + reduction: Reduction::Reduce(keys), + measures, + having: None, + child, + .. + } if !keys.is_without() => { + let schema = child.output_schema().map_err(|e| e.to_string())?; + let groups = keys + .keys() + .iter() + .map(|id| column(*id, &schema)) + .collect::, _>>()?; + let mut values = groups.clone(); + values.extend( + measures + .iter() + .map(|m| aggregate(m, &schema)) + .collect::, _>>()?, + ); + if values.len() != output.columns.len() { + return Err("exact aggregate width mismatch".into()); + } + let values = values + .iter() + .zip(&output.columns) + .map(|(v, c)| format!("{v} AS {}", quoted(&c.name))) + .collect::>(); + Ok(format!( + "SELECT {} FROM ({}){}", + values.join(", "), + render(child)?, + if groups.is_empty() { + String::new() + } else { + format!(" GROUP BY {}", groups.join(", ")) + } + )) + } + QueryExpr::Sort { + keys, + partition_by, + child, + } if partition_by.keys().is_empty() && !partition_by.is_without() => { + let schema = child.output_schema().map_err(|e| e.to_string())?; + let keys = keys + .iter() + .map(|key| { + Ok(format!( + "{} {} NULLS {}", + scalar(&key.expr, &schema)?, + if key.ascending { "ASC" } else { "DESC" }, + if key.nulls_first { "FIRST" } else { "LAST" } + )) + }) + .collect::, String>>()?; + if keys.is_empty() { + return Err("empty exact sort".into()); + } + Ok(format!( + "SELECT * FROM ({}) ORDER BY {}", + render(child)?, + keys.join(", ") + )) + } + QueryExpr::Limit { n, offset, child } => Ok(format!( + "SELECT * FROM ({}) LIMIT {n} OFFSET {offset}", + render(child)? + )), + _ => Err("unsupported canonical ClickHouse exact subtree".into()), + } +} + +#[cfg(test)] +mod tests { + use super::*; + use planner_types::pre_asap::{Column, DataType, Predicate, ProjectItem}; + use std::rc::Rc; + #[test] + fn composite_cut_preserves_branch_time_and_positional_projection() { + let scan = Rc::new(QueryExpr::Scan { + source: Source::Table { + table_ref: "db.samples".into(), + }, + schema: Schema::new(vec![ + Column::new("ts", DataType::Int64, false), + Column::new("v", DataType::Float64, false), + ]), + predicates: vec![Predicate(Rc::new(QueryExpr::Compare { + left: Rc::new(QueryExpr::Column(0)), + op: CompareOpKind::Lt, + right: Rc::new(QueryExpr::Literal(ScalarValue::Int64(-100))), + }))], + }); + let project = QueryExpr::Project { + cols: vec![ProjectItem { + alias: Some("result".into()), + expr: QueryExpr::Column(1), + }], + qualifier: None, + child: scan, + }; + let sql = render(&project).unwrap(); + assert!(sql.contains("`ts` < -100")); + assert!(sql.starts_with("SELECT `v` AS `result`")); + assert!(!sql.contains("{from:")); + assert!(!sql.contains("{to:")); + } + #[test] + fn unsupported_scalar_is_not_forwarded_as_arbitrary_native_code() { + let schema = Schema::new(vec![]); + let expr = QueryExpr::FunctionCall { + name: "unreviewedFunction".into(), + args: vec![], + }; + assert!(scalar(&expr, &schema).is_err()); + } +} + +#[cfg(test)] +mod original_tests { + use super::*; + use asap_frontend_sql::{lower_sql_dialect, SqlCatalog}; + use planner_types::{ + pre_asap::{Column, DataType}, + types::AccuracyTarget, + workload::SqlDialect, + }; + #[tokio::test] + async fn original_exact_shapes_retain_native_aggregates_and_bounds() { + let catalog = SqlCatalog::new().with_table( + "raw_samples", + Schema::new(vec![ + Column::new("metric", DataType::Utf8, false), + Column::new("ts_ms", DataType::Int64, false), + Column::new("value", DataType::Float64, false), + Column::new( + "labels", + DataType::Map { + key: Box::new(DataType::Utf8), + value: Box::new(DataType::Utf8), + value_nullable: false, + }, + false, + ), + ]), + ); + for sql in [ + include_str!("../../tests/fixtures/sql_exact_cuts/q07.sql"), + include_str!("../../tests/fixtures/sql_exact_cuts/q09.sql"), + include_str!("../../tests/fixtures/sql_exact_cuts/q12.sql"), + include_str!("../../tests/fixtures/sql_exact_cuts/q27.sql"), + ] { + let canonical = lower_sql_dialect( + sql, + &catalog, + SqlDialect::ClickhouseSQL, + AccuracyTarget::Exact, + ) + .await + .unwrap(); + let rendered = render(&canonical).unwrap(); + assert!(rendered.contains("1788891296000")); + assert!(rendered.contains("`ts_ms`")); + assert!(!rendered.contains("{from:")); + } + } + #[test] + fn literal_quotes_and_backslashes_are_escaped_independently() { + let rendered = scalar( + &QueryExpr::Literal(ScalarValue::Utf8("a\\'b\n".into())), + &Schema::new(vec![]), + ) + .unwrap(); + assert_eq!(rendered, "'a\\\\\\'b\n'"); + } +} diff --git a/control_plane/tests/fixtures/sql_exact_cuts/q07.sql b/control_plane/tests/fixtures/sql_exact_cuts/q07.sql new file mode 100644 index 00000000..fd214356 --- /dev/null +++ b/control_plane/tests/fixtures/sql_exact_cuts/q07.sql @@ -0,0 +1 @@ +SELECT mapConcat(labels,map('__name__','user_service_cache_refresh_lag_seconds')) AS labels, argMax(value,ts_ms) AS value FROM raw_samples WHERE metric='user_service_cache_refresh_lag_seconds' AND ts_ms>1788891296000-300000 AND ts_ms<=1788891296000 GROUP BY labels ORDER BY labels diff --git a/control_plane/tests/fixtures/sql_exact_cuts/q09.sql b/control_plane/tests/fixtures/sql_exact_cuts/q09.sql new file mode 100644 index 00000000..d134dc55 --- /dev/null +++ b/control_plane/tests/fixtures/sql_exact_cuts/q09.sql @@ -0,0 +1 @@ +SELECT map() AS labels, sum(value) AS value FROM (SELECT labels,argMax(value,ts_ms) AS value FROM raw_samples WHERE metric='backend_process_resident_memory_bytes' AND ts_ms>1788891296000-300000 AND ts_ms<=1788891296000 GROUP BY labels) diff --git a/control_plane/tests/fixtures/sql_exact_cuts/q12.sql b/control_plane/tests/fixtures/sql_exact_cuts/q12.sql new file mode 100644 index 00000000..3d860f5e --- /dev/null +++ b/control_plane/tests/fixtures/sql_exact_cuts/q12.sql @@ -0,0 +1 @@ +SELECT map('job',job) AS labels,sum(value) AS value FROM (SELECT labels['job'] AS job,labels,argMax(value,ts_ms) AS value FROM raw_samples WHERE metric='backend_process_resident_memory_bytes' AND ts_ms>1788891296000-300000 AND ts_ms<=1788891296000 GROUP BY job,labels) GROUP BY job ORDER BY value DESC,job LIMIT 2 diff --git a/control_plane/tests/fixtures/sql_exact_cuts/q27.sql b/control_plane/tests/fixtures/sql_exact_cuts/q27.sql new file mode 100644 index 00000000..f3fbb3d0 --- /dev/null +++ b/control_plane/tests/fixtures/sql_exact_cuts/q27.sql @@ -0,0 +1 @@ +SELECT map('job',job) labels,avg(value) value FROM (SELECT labels['job'] job,ts_ms,sum(value) value FROM raw_samples WHERE metric='backend_process_resident_memory_bytes' AND ts_ms+4000>=1788891296000-21600000 AND ts_ms+4000<=1788891296000 AND modulo(ts_ms+4000,60000)=0 GROUP BY job,ts_ms) GROUP BY job ORDER BY value DESC,job LIMIT 3 diff --git a/crates/asap_types/Cargo.toml b/crates/asap_types/Cargo.toml index a50afa72..a83ea156 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 = "b2b05628dd9a58db555ba309bf7201e10513f4d2" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "0402384e589df6e087d6d2d22b463ddc2eea0774" } diff --git a/data_plane/Cargo.toml b/data_plane/Cargo.toml index 6bf88325..60ae2ef0 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 = "b2b05628dd9a58db555ba309bf7201e10513f4d2" } -asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "b2b05628dd9a58db555ba309bf7201e10513f4d2" } +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" } # 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 = "b2b05628dd9a58db555ba309bf7201e10513f4d2" } +asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "0402384e589df6e087d6d2d22b463ddc2eea0774" } tempfile = "3.20.0" criterion = { version = "0.5", features = ["html_reports"] } tokio-tungstenite = "0.21" From ef1caef738bf3efdf9f62fdd057188e3bb3b9dff Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 11 Sep 2026 08:07:53 -0600 Subject: [PATCH 6/6] fix(clickhouse): retain typed relational aggregate operations --- control_plane/src/query_plan.rs | 3 ++ .../src/query_plan/clickhouse_exact.rs | 51 +++++++++++++++++++ 2 files changed, 54 insertions(+) diff --git a/control_plane/src/query_plan.rs b/control_plane/src/query_plan.rs index 58ec03f2..17943b88 100644 --- a/control_plane/src/query_plan.rs +++ b/control_plane/src/query_plan.rs @@ -991,6 +991,9 @@ where | planner_types::post_asap::ValueOperation::Filter { .. } | planner_types::post_asap::ValueOperation::Sort { .. } | planner_types::post_asap::ValueOperation::Limit { .. } + | planner_types::post_asap::ValueOperation::Exact( + planner_types::post_asap::ExactOperation::Aggregate { .. } + ) ) => { QueryPlanNode::Relational { diff --git a/control_plane/src/query_plan/clickhouse_exact.rs b/control_plane/src/query_plan/clickhouse_exact.rs index 9190b3c7..9c845542 100644 --- a/control_plane/src/query_plan/clickhouse_exact.rs +++ b/control_plane/src/query_plan/clickhouse_exact.rs @@ -358,6 +358,57 @@ mod original_tests { assert!(rendered.contains("1788891296000")); assert!(rendered.contains("`ts_ms`")); assert!(!rendered.contains("{from:")); + if sql.contains("sum(value)") && sql.contains("argMax") { + use crate::physical::post_asap::{PhysicalExpr, PostAsapPlan}; + use crate::query_plan::{ + FallbackPolicy, FixedEvaluationRange, InstantExecution, QueryPlanEntry, + QueryPlanError, QueryPlanNode, + }; + let planned = + crate::clickhouse::plan_clickhouse_sql(sql, &catalog, AccuracyTarget::Exact) + .await + .unwrap(); + let PhysicalExpr::Committed(PostAsapPlan::Summary(root)) = planned.physical else { + panic!("missing selected SQL DAG") + }; + let entry = QueryPlanEntry::compile_bound_relational( + "test".into(), + planned.canonical_sql, + &root, + FixedEvaluationRange { + start_ms: 1788890996000, + end_ms: 1788891296000, + cumulative: false, + }, + InstantExecution { + lookback_ms: 300000, + full_history: false, + cumulative_readout: false, + }, + FallbackPolicy::ExactBackend, + |_, _| Err(QueryPlanError::Invalid("unexpected summary binding".into())), + ) + .unwrap(); + assert!( + !entry + .nodes + .values() + .any(|node| matches!(node, QueryPlanNode::Logical { .. })), + "SQL must not acquire PromQL operators" + ); + assert!(entry.nodes.values().any(|node| match node { + QueryPlanNode::Relational { operation, .. } => matches!( + serde_json::from_value::( + operation.clone() + ) + .unwrap(), + planner_types::post_asap::ValueOperation::Exact( + planner_types::post_asap::ExactOperation::Aggregate { .. } + ) + ), + _ => false, + })); + } } } #[test]