From 341590a09a93a18daff04d56b8c849114ca35d50 Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 10 Sep 2026 20:39:46 -0600 Subject: [PATCH 1/7] fix(sds): reject typed groups on string-label sources --- crates/asap_types/src/grouping_projection.rs | 27 ++++++++++++++++++-- crates/asap_types/src/precompute_plan.rs | 8 ++++++ crates/asap_types/src/sds.rs | 8 ++++++ 3 files changed, 41 insertions(+), 2 deletions(-) diff --git a/crates/asap_types/src/grouping_projection.rs b/crates/asap_types/src/grouping_projection.rs index 0b34cb89..561b29f8 100644 --- a/crates/asap_types/src/grouping_projection.rs +++ b/crates/asap_types/src/grouping_projection.rs @@ -174,8 +174,31 @@ mod identity_tests { assert_ne!(legacy.id, numeric.id); assert_ne!(legacy.id, nullable.id); assert_ne!(numeric.id, nullable.id); - numeric.validate().unwrap(); - nullable.validate().unwrap(); + assert!(numeric + .validate() + .unwrap_err() + .to_string() + .contains("string labels")); + assert!(nullable + .validate() + .unwrap_err() + .to_string() + .contains("string labels")); + for grouping in [numeric.group_by_keys, nullable.group_by_keys] { + let table = crate::sds::DataDescriptor::new_typed( + crate::sds::DataSourceIdentity::Table { + table_ref: "samples".into(), + }, + crate::sds::ValueProjectionIdentity::Column { + name: "value".into(), + }, + "", + Vec::::new(), + "table.samples.v1", + ) + .with_grouping_projection(grouping); + table.validate().unwrap(); + } } } diff --git a/crates/asap_types/src/precompute_plan.rs b/crates/asap_types/src/precompute_plan.rs index bf8df3f1..257d99dd 100644 --- a/crates/asap_types/src/precompute_plan.rs +++ b/crates/asap_types/src/precompute_plan.rs @@ -417,6 +417,14 @@ impl PrecomputePlan { } let mut materializations = BTreeSet::new(); for materialization in &self.materializations { + if materialization.table_name.is_none() + && !materialization.grouping_labels.is_legacy_labels() + { + return Err(PrecomputePlanError::CatalogContract( + "time-series grouping requires non-null string labels".into(), + )); + } + materialization .grouping_labels .validate() diff --git a/crates/asap_types/src/sds.rs b/crates/asap_types/src/sds.rs index 45434db2..23c3cc48 100644 --- a/crates/asap_types/src/sds.rs +++ b/crates/asap_types/src/sds.rs @@ -931,6 +931,14 @@ impl DataDescriptor { pub fn validate(&self) -> Result<(), SdsError> { self.value_projection.validate().map_err(SdsError)?; self.group_by_keys.validate().map_err(SdsError)?; + if matches!(self.source, DataSourceIdentity::TimeSeries { .. }) + && !self.group_by_keys.is_legacy_labels() + { + return Err(SdsError( + "time-series grouping requires non-null string labels".into(), + )); + } + if matches!(self.source, DataSourceIdentity::Table { .. }) { for column in self.group_by_keys.columns() { crate::table_population::validate_column_name(&column.name).map_err(SdsError)?; From 600faf87fed1a202d4b244e5e95145979eb7a98e Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 10 Sep 2026 21:14:19 -0600 Subject: [PATCH 2/7] Execute typed ClickHouse grouping through shared summary plans --- Cargo.lock | 11 +- control_plane/Cargo.toml | 8 +- control_plane/src/clickhouse.rs | 27 ++- crates/asap_types/Cargo.toml | 3 +- crates/asap_types/src/grouping_projection.rs | 46 ++++ crates/asap_types/src/policy_fingerprint.rs | 6 + .../asap_types/src/precompute_plan/catalog.rs | 9 + crates/asap_types/src/summary_catalog.rs | 6 +- data_plane/Cargo.toml | 6 +- .../accelerator.rs | 9 + .../clickhouse_result_adapter.rs | 226 +++++++++++++++++- .../relational_adapter.rs | 192 ++++++++++++++- .../asap_clickhouse_query_engine/server.rs | 1 + .../sketch_db/backfill/clickhouse_reader.rs | 59 ++++- .../tests/clickhouse_differential_e2e.rs | 73 +++++- 15 files changed, 642 insertions(+), 40 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 7030d4fb..16fee0c0 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=e17a53715b445ab490f9c5f07d25fda730d252ec#e17a53715b445ab490f9c5f07d25fda730d252ec" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=7261dab97f8e43070e6b9e33c85493b12b9a351c#7261dab97f8e43070e6b9e33c85493b12b9a351c" 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=e17a53715b445ab490f9c5f07d25fda730d252ec#e17a53715b445ab490f9c5f07d25fda730d252ec" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=7261dab97f8e43070e6b9e33c85493b12b9a351c#7261dab97f8e43070e6b9e33c85493b12b9a351c" 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=e17a53715b445ab490f9c5f07d25fda730d252ec#e17a53715b445ab490f9c5f07d25fda730d252ec" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=7261dab97f8e43070e6b9e33c85493b12b9a351c#7261dab97f8e43070e6b9e33c85493b12b9a351c" 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=e17a53715b445ab490f9c5f07d25fda730d252ec#e17a53715b445ab490f9c5f07d25fda730d252ec" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=7261dab97f8e43070e6b9e33c85493b12b9a351c#7261dab97f8e43070e6b9e33c85493b12b9a351c" [[package]] name = "asap-types" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=e17a53715b445ab490f9c5f07d25fda730d252ec#e17a53715b445ab490f9c5f07d25fda730d252ec" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=7261dab97f8e43070e6b9e33c85493b12b9a351c#7261dab97f8e43070e6b9e33c85493b12b9a351c" dependencies = [ "serde", "serde_json", @@ -460,6 +460,7 @@ version = "0.1.0" dependencies = [ "anyhow", "asap-types", + "base64 0.21.7", "clap 4.6.1", "serde", "serde_json", diff --git a/control_plane/Cargo.toml b/control_plane/Cargo.toml index 994965da..537c5e96 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 = "e17a53715b445ab490f9c5f07d25fda730d252ec" } -asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "e17a53715b445ab490f9c5f07d25fda730d252ec" } +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" } # 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 = "e17a53715b445ab490f9c5f07d25fda730d252ec" } -asap-frontend-sql = { git = "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/ProjectASAP/ASAPPlanner", rev = "e17a53715b445ab490f9c5f07d25fda730d252ec" } +asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "7261dab97f8e43070e6b9e33c85493b12b9a351c" } +asap-frontend-sql = { git = "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/ProjectASAP/ASAPPlanner", rev = "7261dab97f8e43070e6b9e33c85493b12b9a351c" } [dev-dependencies] tokio = { version = "1", features = ["full", "test-util"] } diff --git a/control_plane/src/clickhouse.rs b/control_plane/src/clickhouse.rs index 95da9ec3..6317ee6b 100644 --- a/control_plane/src/clickhouse.rs +++ b/control_plane/src/clickhouse.rs @@ -220,14 +220,34 @@ fn materialize_selected_sql( use planner_types::{post_asap::SummaryExpr, pre_asap::Reduction}; let SummaryExpr::SummaryAgg { reduction: Reduction::Reduce(keys), + child, .. } = &node.expr else { return Err("SQL materialization requires a supported reduction".into()); }; - if keys.is_without() || !keys.keys().is_empty() { - return Err("SQL grouped source projection requires a typed grouping reader".into()); + if keys.is_without() { + return Err("SQL grouping exclusion requires a resolved projection".into()); + } + let SummaryExpr::KeepPreAsap(source) = &child.expr else { + return Err("SQL grouping requires a typed source subtree".into()); + }; + let source_schema = source.output_schema().map_err(|error| error.to_string())?; + let mut columns = Vec::new(); + for key in keys.keys() { + let mut column = source_schema + .columns + .get(*key) + .cloned() + .ok_or("SQL grouping column is absent from source schema")?; + if column.nullable { + return Err("nullable table grouping requires an explicit null-key encoding".into()); + } + column.table = None; + columns.push(column); } + let grouping = asap_types::GroupingProjection::new(columns); + grouping.validate()?; let (table, value, window, population, timestamp) = clickhouse_materialization_leaf_contract(node, query.start_ms, query.end_ms)?; let window_secs = window.ok_or("SQL materialization requires a bounded window")?; @@ -237,7 +257,7 @@ fn materialize_selected_sql( family: crate::physical::compiler::physical_materialization_family(family), window_secs, spatial_filter: String::new(), - grouping: Vec::new(), + grouping: grouping.names(), item_label: None, heap_update_mode: None, aggregation_input: AggregationInput::Raw, @@ -248,6 +268,7 @@ fn materialize_selected_sql( ) .map_err(|error| error.to_string())?; config.table_name = Some(table); + config.grouping_labels = grouping; config.value_projection = Some(value); config.table_timestamp_column = Some(timestamp); config.table_population = Some(population); diff --git a/crates/asap_types/Cargo.toml b/crates/asap_types/Cargo.toml index a044eb6b..21d80432 100644 --- a/crates/asap_types/Cargo.toml +++ b/crates/asap_types/Cargo.toml @@ -4,6 +4,7 @@ version.workspace = true edition.workspace = true [dependencies] +base64 = "0.21" tracing.workspace = true serde.workspace = true serde_json.workspace = true @@ -32,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 = "e17a53715b445ab490f9c5f07d25fda730d252ec" } +planner-types = { package = "asap-types", git = "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/ProjectASAP/ASAPPlanner", rev = "7261dab97f8e43070e6b9e33c85493b12b9a351c" } diff --git a/crates/asap_types/src/grouping_projection.rs b/crates/asap_types/src/grouping_projection.rs index 561b29f8..2ba03aaa 100644 --- a/crates/asap_types/src/grouping_projection.rs +++ b/crates/asap_types/src/grouping_projection.rs @@ -24,6 +24,13 @@ impl GroupingProjection { } Ok(()) } + pub fn validate_table_columns(&self) -> Result<(), String> { + self.validate()?; + for column in self.columns() { + crate::table_population::validate_column_name(&column.name)?; + } + Ok(()) + } pub fn push(&mut self, name: String) { self.0.push(Column::new(name, DataType::Utf8, false)); self.0.sort_by(|a, b| a.name.cmp(&b.name)); @@ -223,3 +230,42 @@ pub(crate) fn serialize_config_grouping( grouping.serialize(serializer) } } + +/// Table group values travel through the string-key index as lossless JSON bytes. +/// Map values use ordered key/value pairs, preserving duplicate keys. +pub const TABLE_GROUP_OBSERVATION_SEMANTICS: &str = "asap.table-column-groups.base64-json.v1"; + +pub fn encode_table_group_value(value: &serde_json::Value) -> Result { + use base64::Engine; + let bytes = serde_json::to_vec(value).map_err(|error| error.to_string())?; + Ok(base64::engine::general_purpose::STANDARD.encode(bytes)) +} + +pub fn decode_table_group_value(encoded: &str) -> Result { + use base64::Engine; + let bytes = base64::engine::general_purpose::STANDARD + .decode(encoded) + .map_err(|error| error.to_string())?; + serde_json::from_slice(&bytes).map_err(|error| error.to_string()) +} + +#[cfg(test)] +mod codec_tests { + use super::*; + /// The grouping codec preserves null, escaping, duplicate map keys and exact integers. + #[test] + fn typed_group_codec_is_lossless_at_numeric_and_string_boundaries() { + for value in [ + serde_json::Value::Null, + serde_json::json!("a,\"b\\c\n"), + serde_json::json!(i64::MAX), + serde_json::json!(i64::MIN), + serde_json::json!([["k", i64::MAX], ["k", null]]), + ] { + let encoded = encode_table_group_value(&value).unwrap(); + assert!(!encoded.contains(['\"', ',', '\\'])); + assert_eq!(decode_table_group_value(&encoded).unwrap(), value); + } + assert!(decode_table_group_value("not base64").is_err()); + } +} diff --git a/crates/asap_types/src/policy_fingerprint.rs b/crates/asap_types/src/policy_fingerprint.rs index 19daa348..f9e75165 100644 --- a/crates/asap_types/src/policy_fingerprint.rs +++ b/crates/asap_types/src/policy_fingerprint.rs @@ -127,6 +127,12 @@ impl PolicyFingerprint { } buf.push(0); + if cfg.table_name.is_some() && !cfg.grouping_labels.is_empty() { + buf.extend_from_slice( + crate::grouping_projection::TABLE_GROUP_OBSERVATION_SEMANTICS.as_bytes(), + ); + buf.push(0); + } if !cfg.grouping_labels.is_legacy_labels() { buf.extend_from_slice(b"typed-grouping:"); buf.extend_from_slice( diff --git a/crates/asap_types/src/precompute_plan/catalog.rs b/crates/asap_types/src/precompute_plan/catalog.rs index b6845a2d..bd1caaff 100644 --- a/crates/asap_types/src/precompute_plan/catalog.rs +++ b/crates/asap_types/src/precompute_plan/catalog.rs @@ -78,6 +78,15 @@ impl PrecomputePlan { table_ref: table_ref.clone(), }, ); + if config.table_name.is_some() + && !config.grouping_labels.is_empty() + && data.observation_semantics + != crate::grouping_projection::TABLE_GROUP_OBSERVATION_SEMANTICS + { + return Err(invalid( + "table grouping codec differs from installed contract", + )); + } let expected_projection = config.effective_value_projection(); if data.partitioning != config.partitioning || data.timestamp_column != config.table_timestamp_column diff --git a/crates/asap_types/src/summary_catalog.rs b/crates/asap_types/src/summary_catalog.rs index b7864c2c..5f31ee43 100644 --- a/crates/asap_types/src/summary_catalog.rs +++ b/crates/asap_types/src/summary_catalog.rs @@ -119,7 +119,11 @@ impl SummaryCatalog { .population_filter_canonical() .map_err(SummaryCatalogError::Descriptor)?, config.grouping_labels.names(), - "asap.timestamped-observations.v2", + if config.table_name.is_some() && !config.grouping_labels.is_empty() { + crate::grouping_projection::TABLE_GROUP_OBSERVATION_SEMANTICS + } else { + "asap.timestamped-observations.v2" + }, ) .with_grouping_projection(config.grouping_labels.clone()) .with_partitioning(config.partitioning) diff --git a/data_plane/Cargo.toml b/data_plane/Cargo.toml index 248cef48..14c29ec9 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 = "e17a53715b445ab490f9c5f07d25fda730d252ec" } -asap-frontend-promql = { git = "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/ProjectASAP/ASAPPlanner", rev = "e17a53715b445ab490f9c5f07d25fda730d252ec" } +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" } # 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 = "e17a53715b445ab490f9c5f07d25fda730d252ec" } +asap-aware-mapping = { git = "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/ProjectASAP/ASAPPlanner", rev = "7261dab97f8e43070e6b9e33c85493b12b9a351c" } 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/accelerator.rs b/data_plane/src/query_engines/asap_clickhouse_query_engine/accelerator.rs index 4cc33dd6..d8c162a6 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 @@ -102,6 +102,15 @@ impl CatalogClickHouseAccelerator { parameters.insert(format!("param_{name}"), end_ms.to_string()); } parameters.insert("default_format".into(), "JSONCompact".into()); + parameters.insert( + "output_format_json_map_as_array_of_tuples".into(), + "1".into(), + ); + parameters.insert( + "output_format_json_named_tuples_as_objects".into(), + "0".into(), + ); + parameters.insert("output_format_json_quote_64bit_integers".into(), "0".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/clickhouse_result_adapter.rs b/data_plane/src/query_engines/asap_clickhouse_query_engine/clickhouse_result_adapter.rs index bcca5f52..32883fc8 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 @@ -1,12 +1,16 @@ use super::fallback::ClickHouseRawResponse; use arrow::{ - array::{Float64Array, StringArray, TimestampMillisecondArray}, + array::{Array, Float64Array, StringArray, TimestampMillisecondArray}, datatypes::{DataType, Field, Schema}, json::LineDelimitedWriter, record_batch::RecordBatch, util::display::array_value_to_string, }; use axum::response::{IntoResponse, Response}; +use serde::{ + ser::{SerializeMap, SerializeSeq, SerializeStruct}, + Serialize, Serializer, +}; use std::{collections::BTreeMap, sync::Arc}; #[derive(Debug, Clone, Copy, PartialEq, Eq)] @@ -89,6 +93,12 @@ impl ClickHouseQueryResult { if column > 0 { output.push(b'\t'); } + if matches!(batch.column(column).data_type(), DataType::Map(..)) { + output.extend_from_slice( + map_literal(batch.column(column).as_ref(), row)?.as_bytes(), + ); + continue; + } let value = array_value_to_string(batch.column(column).as_ref(), row) .map_err(|error| { ClickHouseResultError::Arrow(error.to_string()) @@ -98,6 +108,17 @@ impl ClickHouseQueryResult { output.push(b'\n'); } ClickHouseFormat::JsonEachRow => { + if batch + .schema() + .fields() + .iter() + .any(|field| matches!(field.data_type(), DataType::Map(..))) + { + serde_json::to_writer(&mut output, &JsonArrowRow { batch, row })?; + output.push(b'\n'); + continue; + } + let row = batch.slice(row, 1); let mut writer = LineDelimitedWriter::new(&mut output); writer @@ -134,6 +155,27 @@ impl ClickHouseQueryResult { .collect::>() }) .unwrap_or_default(); + if self.batches.iter().any(|batch| { + batch + .schema() + .fields() + .iter() + .any(|field| matches!(field.data_type(), DataType::Map(..))) + }) { + let count: usize = self.batches.iter().map(RecordBatch::num_rows).sum(); + let mut output = Vec::new(); + let mut serializer = serde_json::Serializer::new(&mut output); + let mut document = serializer.serialize_struct("ClickHouseResult", 4)?; + document.serialize_field("meta", &meta)?; + document.serialize_field("data", &JsonArrowRows(&self.batches))?; + document.serialize_field("rows", &count)?; + document.serialize_field( + "statistics", + &serde_json::json!({"elapsed":0.0,"rows_read":count,"bytes_read":0}), + )?; + SerializeStruct::end(document)?; + return Ok(output); + } let mut rows = Vec::new(); for batch in &self.batches { let mut encoded = Vec::new(); @@ -160,7 +202,7 @@ impl ClickHouseQueryResult { } } -fn clickhouse_type(data_type: &arrow::datatypes::DataType, nullable: bool) -> String { +pub(super) fn clickhouse_type(data_type: &arrow::datatypes::DataType, nullable: bool) -> String { use arrow::datatypes::DataType; let base = match data_type { DataType::Boolean => "Bool".into(), @@ -176,6 +218,14 @@ fn clickhouse_type(data_type: &arrow::datatypes::DataType, nullable: bool) -> St DataType::Float64 => "Float64".into(), DataType::Utf8 | DataType::LargeUtf8 => "String".into(), DataType::Timestamp(_, _) => "DateTime64(3)".into(), + DataType::Map(entries, _) => match entries.data_type() { + DataType::Struct(fields) if fields.len() == 2 => format!( + "Map({}, {})", + clickhouse_type(fields[0].data_type(), false), + clickhouse_type(fields[1].data_type(), fields[1].is_nullable()) + ), + other => other.to_string(), + }, other => other.to_string(), }; if nullable { @@ -243,3 +293,175 @@ mod tests { assert_eq!(document["rows"], 1); } } + +struct JsonArrowRows<'a>(&'a [RecordBatch]); +impl Serialize for JsonArrowRows<'_> { + fn serialize(&self, serializer: S) -> Result { + let mut rows = + serializer.serialize_seq(Some(self.0.iter().map(RecordBatch::num_rows).sum()))?; + for batch in self.0 { + for row in 0..batch.num_rows() { + rows.serialize_element(&JsonArrowRow { batch, row })?; + } + } + rows.end() + } +} +struct JsonArrowRow<'a> { + batch: &'a RecordBatch, + row: usize, +} +impl Serialize for JsonArrowRow<'_> { + fn serialize(&self, serializer: S) -> Result { + let mut object = serializer.serialize_map(Some(self.batch.num_columns()))?; + for (field, column) in self + .batch + .schema() + .fields() + .iter() + .zip(self.batch.columns()) + { + object.serialize_entry( + field.name(), + &JsonArrowValue { + array: column.as_ref(), + row: self.row, + }, + )?; + } + object.end() + } +} +struct JsonArrowValue<'a> { + array: &'a dyn Array, + row: usize, +} +impl Serialize for JsonArrowValue<'_> { + fn serialize(&self, serializer: S) -> Result { + use arrow::array::*; + use serde::ser::Error; + if self.array.is_null(self.row) { + return serializer.serialize_none(); + } + macro_rules! scalar { + ($ty:ty) => { + self.array + .as_any() + .downcast_ref::<$ty>() + .ok_or_else(|| S::Error::custom("Arrow scalar type mismatch"))? + .value(self.row) + .serialize(serializer) + }; + } + match self.array.data_type() { + DataType::Map(..) => { + let map = self + .array + .as_any() + .downcast_ref::() + .ok_or_else(|| S::Error::custom("Arrow map type mismatch"))?; + let entries = map.value(self.row); + let mut object = serializer.serialize_map(Some(entries.len()))?; + for row in 0..entries.len() { + object.serialize_entry( + &JsonArrowValue { + array: entries.column(0).as_ref(), + row, + }, + &JsonArrowValue { + array: entries.column(1).as_ref(), + row, + }, + )?; + } + object.end() + } + DataType::Boolean => scalar!(BooleanArray), + DataType::Int8 => scalar!(Int8Array), + DataType::Int16 => scalar!(Int16Array), + DataType::Int32 => scalar!(Int32Array), + DataType::Int64 => scalar!(Int64Array), + DataType::UInt8 => scalar!(UInt8Array), + DataType::UInt16 => scalar!(UInt16Array), + DataType::UInt32 => scalar!(UInt32Array), + DataType::UInt64 => scalar!(UInt64Array), + DataType::Float32 => scalar!(Float32Array), + DataType::Float64 => scalar!(Float64Array), + DataType::Utf8 => scalar!(StringArray), + DataType::LargeUtf8 => scalar!(LargeStringArray), + DataType::Timestamp(..) => array_value_to_string(self.array, self.row) + .map_err(S::Error::custom)? + .serialize(serializer), + dtype => Err(S::Error::custom(format!( + "unsupported nested Arrow value {dtype}" + ))), + } + } +} + +fn map_literal(array: &dyn Array, row: usize) -> Result { + use arrow::array::{Float32Array, LargeStringArray, MapArray}; + if array.is_null(row) { + return Ok("NULL".into()); + } + match array.data_type() { + DataType::Map(..) => { + let map = array + .as_any() + .downcast_ref::() + .ok_or_else(|| ClickHouseResultError::Arrow("invalid map array".into()))?; + let entries = map.value(row); + let mut values = Vec::with_capacity(entries.len()); + for row in 0..entries.len() { + values.push(format!( + "{}:{}", + map_literal(entries.column(0).as_ref(), row)?, + map_literal(entries.column(1).as_ref(), row)? + )); + } + Ok(format!("{{{}}}", values.join(","))) + } + DataType::Utf8 | DataType::LargeUtf8 => { + let value = if let Some(array) = array.as_any().downcast_ref::() { + array.value(row) + } else if let Some(array) = array.as_any().downcast_ref::() { + array.value(row) + } else { + return Err(ClickHouseResultError::Arrow("invalid string array".into())); + }; + let mut escaped = String::from("'"); + for ch in value.chars() { + match ch { + '\\' => escaped.push_str("\\\\"), + '\'' => escaped.push_str("\\'"), + '\n' => escaped.push_str("\\n"), + '\r' => escaped.push_str("\\r"), + '\t' => escaped.push_str("\\t"), + '\0' => escaped.push_str("\\0"), + '\u{0008}' => escaped.push_str("\\b"), + '\u{000c}' => escaped.push_str("\\f"), + ch => escaped.push(ch), + } + } + escaped.push('\''); + Ok(escaped) + } + DataType::Float64 => Ok(array + .as_any() + .downcast_ref::() + .unwrap() + .value(row) + .to_string()), + DataType::Float32 => Ok(array + .as_any() + .downcast_ref::() + .unwrap() + .value(row) + .to_string()), + DataType::Timestamp(..) => Err(ClickHouseResultError::Arrow( + "timestamp map TSV encoding is unsupported".into(), + )), + _ => array_value_to_string(array, row) + .map_err(|error| ClickHouseResultError::Arrow(error.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 d33799a7..56cc44f8 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,7 +4,8 @@ use std::{cmp::Ordering, collections::BTreeMap, sync::Arc}; use arrow::{ array::{ - ArrayRef, BooleanArray, Float64Array, Int64Array, StringArray, TimestampMillisecondArray, + ArrayRef, BooleanArray, Float64Array, Int64Array, MapArray, StringArray, StructArray, + TimestampMillisecondArray, }, datatypes::{DataType as ArrowDataType, Field, Schema}, record_batch::RecordBatch, @@ -45,6 +46,7 @@ enum Cell { Utf8(String), Bool(bool), Timestamp(i64), + Map(Vec<(Cell, Cell)>), } fn json_cell( @@ -69,6 +71,26 @@ fn json_cell( .map(|value| Cell::Utf8(value.into())) .ok_or_else(invalid), DataType::Bool => value.as_bool().map(Cell::Bool).ok_or_else(invalid), + DataType::Map { + key, + value: value_type, + value_nullable, + } => { + let (key_type, item_type) = map_type_parts(clickhouse_type).ok_or_else(invalid)?; + let entries = value.as_array().ok_or_else(invalid)?; + let mut result = Vec::with_capacity(entries.len()); + for entry in entries { + let pair = entry + .as_array() + .filter(|pair| pair.len() == 2) + .ok_or_else(invalid)?; + result.push(( + json_cell(&pair[0], key, false, key_type)?, + json_cell(&pair[1], value_type, *value_nullable, item_type)?, + )); + } + Ok(Cell::Map(result)) + } DataType::Timestamp => parse_clickhouse_timestamp(value, clickhouse_type) .map(Cell::Timestamp) .ok_or_else(invalid), @@ -251,6 +273,24 @@ 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 +} + fn clickhouse_type_matches(actual: Option<&str>, expected: &DataType, nullable: bool) -> bool { let Some(mut actual) = actual else { return false; @@ -271,6 +311,14 @@ fn clickhouse_type_matches(actual: Option<&str>, expected: &DataType, nullable: DataType::Float64 => actual == "Float64", DataType::Utf8 => actual == "String", DataType::Bool => actual == "Bool", + DataType::Map { + key, + value, + value_nullable, + } => map_type_parts(actual).is_some_and(|(key_type, value_type)| { + clickhouse_type_matches(Some(key_type), key, false) + && clickhouse_type_matches(Some(value_type), value, *value_nullable) + }), DataType::Timestamp => { actual == "Int64" || actual == "DateTime" @@ -490,7 +538,17 @@ fn row_from_value( .iter() .map(|(name, dtype, nullable)| { if let Some(value) = group.get(name) { - return Ok(Cell::Utf8(value.clone())); + let encoded = asap_types::grouping_projection::decode_table_group_value(value) + .map_err(|error| { + ClickHouseRelationalError::Invalid(format!( + "invalid typed group {name}: {error}" + )) + })?; + let column_type = super::clickhouse_result_adapter::clickhouse_type( + &arrow_type(dtype), + *nullable, + ); + return json_cell(&encoded, dtype, *nullable, &column_type); } if *dtype == DataType::Timestamp { return Ok(Cell::Timestamp(timestamp)); @@ -649,6 +707,24 @@ fn cell_cmp(left: &Cell, right: &Cell) -> Option { (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)), + (Cell::Map(left), Cell::Map(right)) => { + for ((left_key, left_value), (right_key, right_value)) in left.iter().zip(right) { + let order = cell_cmp(left_key, right_key)?; + if order != Ordering::Equal { + return Some(order); + } + let order = match (left_value, right_value) { + (Cell::Null, Cell::Null) => Ordering::Equal, + (Cell::Null, _) => Ordering::Greater, + (_, Cell::Null) => Ordering::Less, + _ => cell_cmp(left_value, right_value)?, + }; + if order != Ordering::Equal { + return Some(order); + } + } + Some(left.len().cmp(&right.len())) + } _ => None, } } @@ -667,6 +743,24 @@ fn arrow_type(dtype: &DataType) -> ArrowDataType { DataType::Float64 => ArrowDataType::Float64, DataType::Utf8 => ArrowDataType::Utf8, DataType::Bool => ArrowDataType::Boolean, + DataType::Map { + key, + value, + value_nullable, + } => ArrowDataType::Map( + Arc::new(Field::new( + "entries", + ArrowDataType::Struct( + vec![ + Field::new("key", arrow_type(key), false), + Field::new("value", arrow_type(value), *value_nullable), + ] + .into(), + ), + false, + )), + false, + ), DataType::Timestamp => { ArrowDataType::Timestamp(arrow::datatypes::TimeUnit::Millisecond, None) } @@ -696,6 +790,57 @@ fn build_array( DataType::Float64 => Arc::new(Float64Array::from(values!(Float64))) as ArrayRef, DataType::Utf8 => Arc::new(StringArray::from(values!(Utf8))) as ArrayRef, DataType::Bool => Arc::new(BooleanArray::from(values!(Bool))) as ArrayRef, + DataType::Map { key, value, .. } => { + let mut offsets = vec![0_i32]; + let mut valid = Vec::with_capacity(rows.len()); + let mut entries = Vec::new(); + for row in rows { + match row.get(column) { + Some(Cell::Map(pairs)) => { + valid.push(true); + entries.extend( + pairs + .iter() + .map(|(key, value)| vec![key.clone(), value.clone()]), + ); + } + Some(Cell::Null) => valid.push(false), + _ => { + return Err(ClickHouseRelationalError::Invalid( + "incompatible map value".into(), + )) + } + } + offsets.push(i32::try_from(entries.len()).map_err(|_| { + ClickHouseRelationalError::Invalid("map offset exceeds Arrow limit".into()) + })?); + } + let ArrowDataType::Map(field, ordered) = arrow_type(dtype) else { + unreachable!() + }; + let ArrowDataType::Struct(fields) = field.data_type() else { + unreachable!() + }; + let values = StructArray::try_new( + fields.clone(), + vec![ + build_array(&entries, 0, key)?, + build_array(&entries, 1, value)?, + ], + None, + ) + .map_err(|error| ClickHouseRelationalError::Arrow(error.to_string()))?; + Arc::new( + MapArray::try_new( + field, + arrow::buffer::OffsetBuffer::new(offsets.into()), + values, + Some(arrow::buffer::NullBuffer::from(valid)), + ordered, + ) + .map_err(|error| ClickHouseRelationalError::Arrow(error.to_string()))?, + ) as ArrayRef + } DataType::Timestamp => { Arc::new(TimestampMillisecondArray::from(values!(Timestamp))) as ArrayRef } @@ -785,6 +930,49 @@ mod tests { }) } + /// Map entries remain ordered pairs, including duplicate keys and null values. + #[test] + fn map_transport_retains_duplicate_keys_and_null_values() { + use super::super::clickhouse_result_adapter::{ClickHouseFormat, ClickHouseQueryResult}; + let dtype = DataType::Map { + key: Box::new(DataType::Utf8), + value: Box::new(DataType::Utf8), + value_nullable: true, + }; + let cell = json_cell( + &serde_json::json!([["job", "a"], ["job", "b"], ["zone", null]]), + &dtype, + false, + "Map(String, Nullable(String))", + ) + .unwrap(); + let Cell::Map(entries) = &cell else { + panic!("expected map") + }; + assert_eq!(entries.len(), 3); + let array = build_array(&[vec![cell]], 0, &dtype).unwrap(); + let batch = RecordBatch::try_new( + Arc::new(Schema::new(vec![Field::new( + "labels", + arrow_type(&dtype), + false, + )])), + vec![array], + ) + .unwrap(); + let result = ClickHouseQueryResult { + batches: vec![batch], + }; + let body = String::from_utf8(result.encode(ClickHouseFormat::Json).unwrap()).unwrap(); + assert!(body.contains("Map(String, Nullable(String))"), "{body}"); + assert!(body.contains("\"job\":\"a\",\"job\":\"b\""), "{body}"); + assert!(body.contains("\"zone\":null"), "{body}"); + assert_eq!( + String::from_utf8(result.encode(ClickHouseFormat::TabSeparated).unwrap()).unwrap(), + "{'job':'a','job':'b','zone':NULL}\n" + ); + } + #[test] fn executes_filter_project_arithmetic_sort_and_limit_chain() { let input_schema = schema(&[("ts", DataType::Timestamp), ("sum", DataType::Float64)]); diff --git a/data_plane/src/query_engines/asap_clickhouse_query_engine/server.rs b/data_plane/src/query_engines/asap_clickhouse_query_engine/server.rs index 73a10b2e..1560560c 100644 --- a/data_plane/src/query_engines/asap_clickhouse_query_engine/server.rs +++ b/data_plane/src/query_engines/asap_clickhouse_query_engine/server.rs @@ -375,6 +375,7 @@ async fn execute_or_fallback(state: &ServerState, request: &ClickHouseQueryReque tracing::info!( failure_stage = reason.stage(), failure_reason = reason.reason_code(), + failure_detail = ?reason, "ClickHouse acceleration routed to exact fallback" ); let stage = reason.stage(); diff --git a/data_plane/src/storage_engines/sketch_db/backfill/clickhouse_reader.rs b/data_plane/src/storage_engines/sketch_db/backfill/clickhouse_reader.rs index 596c2a7f..563a94b3 100644 --- a/data_plane/src/storage_engines/sketch_db/backfill/clickhouse_reader.rs +++ b/data_plane/src/storage_engines/sketch_db/backfill/clickhouse_reader.rs @@ -51,6 +51,7 @@ pub struct ClickHouseReader { population: Option, value_projection: Option, output_metric: Option, + grouping_projection: Option, } #[derive(Deserialize)] @@ -58,6 +59,7 @@ pub struct ClickHouseReader { enum ClickHouseLabels { Series(String), Map(std::collections::BTreeMap), + Columns(Vec<(String, String)>), } #[derive(Deserialize)] @@ -76,11 +78,27 @@ impl ClickHouseReader { population: None, value_projection: None, output_metric: None, + grouping_projection: None, }) } fn sql(&self) -> String { let c = &self.config; + let labels = self.grouping_projection.as_ref().map_or_else( + || c.labels_column.clone(), + |grouping| { + let entries = grouping + .columns() + .iter() + .map(|column| { + let value = format!("base64Encode(toJSONString({}))", column.name); + format!("'{}', {value}", column.name) + }) + .collect::>() + .join(", "); + format!("map({entries})") + }, + ); let population = self.population.as_ref().map_or_else( || format!("{} = {{metric:String}}", c.metric_column), |population| { @@ -132,7 +150,7 @@ impl ClickHouseReader { FROM {database}.{table} WHERE {population} \ AND {timestamp} >= {{start_ms:Int64}} AND {timestamp} < {{end_ms:Int64}} \ ORDER BY labels, timestamp_ms FORMAT JSONEachRow", - labels = c.labels_column, + labels = labels, timestamp = c.timestamp_ms_column, value = value, database = c.database, @@ -153,8 +171,13 @@ pub fn clickhouse_reader_factory(config: ClickHouseReaderConfig) -> ReaderFactor } .into()); } - if !materialization.grouping_labels.is_legacy_labels() { - return Err("typed table grouping requires a typed grouping reader".into()); + materialization.grouping_labels.validate_table_columns()?; + for column in materialization.grouping_labels.columns() { + if column.nullable { + return Err( + "nullable table grouping requires an explicit null-key encoding".into(), + ); + } } let mut source_config = config.clone(); source_config.database = database.clone(); @@ -182,6 +205,7 @@ pub fn clickhouse_reader_factory(config: ClickHouseReaderConfig) -> ReaderFactor reader.population = Some(materialization.table_population.clone().unwrap_or_default()); reader.value_projection = Some(materialization.effective_value_projection().clone()); reader.output_metric = Some(materialization.metric.clone()); + reader.grouping_projection = Some(materialization.grouping_labels.clone()); Ok(Arc::new(reader) as Arc) } source => fallback(source, materialization), @@ -208,6 +232,13 @@ impl RawSampleReader for ClickHouseReader { ("param_start_ms", start_ms.as_str()), ("param_end_ms", end_ms.as_str()), ]); + if self.grouping_projection.is_some() { + request = request.query(&[ + ("output_format_json_map_as_array_of_tuples", "1"), + ("output_format_json_named_tuples_as_objects", "0"), + ("output_format_json_quote_64bit_integers", "0"), + ]); + } if let Some(asap_types::sds::ValueProjectionIdentity::Constant { value }) = &self.value_projection { @@ -261,8 +292,24 @@ impl RawSampleReader for ClickHouseReader { serde_json::from_str(line).map_err(|error| RawSampleReaderError::Decode { reason: error.to_string(), })?; + let labels = match row.labels { + ClickHouseLabels::Columns(columns) => { + let count = columns.len(); + let labels = columns + .into_iter() + .collect::>(); + if labels.len() != count { + return Err(RawSampleReaderError::Decode { + reason: "duplicate source grouping columns".into(), + }); + } + ClickHouseLabels::Map(labels) + } + labels => labels, + }; let sample = RawSample { - labels: match row.labels { + labels: match labels { + ClickHouseLabels::Columns(_) => unreachable!("columns normalized above"), ClickHouseLabels::Series(series) => { if self.population.is_some() { let metric = self.output_metric.as_deref().unwrap_or(&filter.metric); @@ -412,9 +459,7 @@ mod tests { }, &typed, ); - assert!( - matches!(rejected, Err(error) if error.to_string().contains("typed grouping reader")) - ); + assert!(matches!(rejected, Err(error) if error.to_string().contains("null-key encoding"))); assert!(factory( &BackfillSource::ClickHouse { diff --git a/data_plane/tests/clickhouse_differential_e2e.rs b/data_plane/tests/clickhouse_differential_e2e.rs index fa41ada2..6f2b870a 100644 --- a/data_plane/tests/clickhouse_differential_e2e.rs +++ b/data_plane/tests/clickhouse_differential_e2e.rs @@ -152,7 +152,7 @@ async fn exact_proxy_matches_clickhouse_for_sql_and_grafana_smoke_queries() { #[tokio::test] async fn compiled_publication_executes_mixed_dag_in_data_plane_process() { - for aggregate in ["sum(value)", "count(*)"] { + for aggregate in ["sum(value)", "count(*)", "max(value)"] { run_mixed_aggregate(aggregate).await; } } @@ -167,6 +167,13 @@ async fn run_mixed_aggregate(aggregate: &str) { let client = reqwest::Client::new(); eprintln!("checking mixed aggregate {aggregate}"); let sql = format!("SELECT sums.timestamp, sums.total / divisors.divisor AS ratio FROM (SELECT 2000 AS timestamp, {aggregate} AS total FROM telemetry WHERE metric = 'requests' AND timestamp_ms >= 0 AND timestamp_ms < 2000) AS sums INNER JOIN divisors ON sums.timestamp = divisors.timestamp"); + let grouped = aggregate == "max(value)"; + let sql = if grouped { + "SELECT labels, max(value) AS value FROM telemetry WHERE metric = 'requests' AND timestamp_ms >= 0 AND timestamp_ms < 2000 GROUP BY labels ORDER BY labels".to_string() + } else { + sql + }; + let format = if grouped { "JSON" } else { "TabSeparated" }; let value_type = if aggregate == "count(*)" { "Nullable(Float64)" } else { @@ -199,13 +206,31 @@ async fn run_mixed_aggregate(aggregate: &str) { } let mut exact_request = client .post(&clickhouse_url) - .body(format!("{sql} FORMAT TabSeparated")); + .body(format!("{sql} FORMAT {format}")); if let Some(user) = &user { exact_request = exact_request.basic_auth(user, password.as_ref()); } let exact = exact_request.send().await.unwrap().bytes().await.unwrap(); let mut workload = mixed_workload(&sql); + if grouped { + use planner_types::pre_asap::{Column, DataType}; + workload + .tables + .get_mut("telemetry") + .unwrap() + .columns + .push(Column::new( + "labels", + DataType::Map { + key: Box::new(DataType::Utf8), + value: Box::new(DataType::Utf8), + value_nullable: false, + }, + false, + )); + } + workload.tables.get_mut("telemetry").unwrap().columns[1].nullable = aggregate == "count(*)"; if aggregate == "count(*)" { let mut nullable_count = mixed_workload(&sql.replace("count(*)", "count(value)")); @@ -242,10 +267,13 @@ async fn run_mixed_aggregate(aggregate: &str) { ); } let entry = publication.query_plan.entries.values().next().unwrap(); - assert!(entry.nodes.values().any(|node| matches!( - node, - control_plane::query_plan::QueryPlanNode::ExternalExact { .. } - ))); + assert_eq!( + entry.nodes.values().any(|node| matches!( + node, + control_plane::query_plan::QueryPlanNode::ExternalExact { .. } + )), + !grouped + ); assert!(entry.nodes.values().any(|node| matches!( node, control_plane::query_plan::QueryPlanNode::ReadMaterialization { .. } @@ -282,7 +310,7 @@ async fn run_mixed_aggregate(aggregate: &str) { .arg("0") .arg("--output-dir") .arg(output.path()) - .stdout(Stdio::null()) + .stdout(Stdio::inherit()) .stderr(Stdio::inherit()) .env("RUST_LOG", "data_plane=info"); if let Some(user) = &user { @@ -376,7 +404,7 @@ async fn run_mixed_aggregate(aggregate: &str) { let mut mutated_request = client .post(&clickhouse_url) - .body(format!("{sql} FORMAT TabSeparated")); + .body(format!("{sql} FORMAT {format}")); if let Some(user) = &user { mutated_request = mutated_request.basic_auth(user, password.as_ref()); } @@ -388,7 +416,7 @@ async fn run_mixed_aggregate(aggregate: &str) { let mut mixed = client .get(format!("http://127.0.0.1:{sql_port}/")) - .query(&[("query", sql.as_str())]); + .query(&[("query", sql.as_str()), ("default_format", format)]); if let Some(user) = &user { mixed = mixed.header("x-clickhouse-user", user); } @@ -411,9 +439,30 @@ async fn run_mixed_aggregate(aggregate: &str) { .to_owned(); let actual = response.bytes().await.unwrap(); assert!(status.is_success(), "mixed listener returned {status}"); - assert_eq!(execution, "hybrid", "failure reason: {failure_reason}"); assert_eq!( - actual, exact, - "compiled mixed result must equal pre-mutation exact baseline" + execution, + if grouped { "warm" } else { "hybrid" }, + "failure reason: {failure_reason}" ); + if grouped { + let actual: serde_json::Value = serde_json::from_slice(&actual).unwrap(); + let exact: serde_json::Value = serde_json::from_slice(&exact).unwrap(); + assert_eq!(actual["meta"], exact["meta"]); + let actual_rows = actual["data"].as_array().unwrap(); + let exact_rows = exact["data"].as_array().unwrap(); + assert_eq!(actual_rows.len(), exact_rows.len()); + for (actual, exact) in actual_rows.iter().zip(exact_rows) { + assert_eq!(actual["labels"], exact["labels"]); + assert_eq!( + actual["value"].as_f64().unwrap().to_bits(), + exact["value"].as_f64().unwrap().to_bits() + ); + } + assert_eq!(actual["data"].as_array().unwrap().len(), 2); + } else { + assert_eq!( + actual, exact, + "compiled mixed result must equal pre-mutation exact baseline" + ); + } } From 56805e0cfe8b56f15d865196a909d4951c8978ff Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 10 Sep 2026 21:28:13 -0600 Subject: [PATCH 3/7] Resolve exact integer expressions in SQL source boundaries --- control_plane/src/clickhouse.rs | 52 ++++++++++++++++++- .../tests/clickhouse_differential_e2e.rs | 2 +- 2 files changed, 52 insertions(+), 2 deletions(-) diff --git a/control_plane/src/clickhouse.rs b/control_plane/src/clickhouse.rs index 6317ee6b..f6c3c3e3 100644 --- a/control_plane/src/clickhouse.rs +++ b/control_plane/src/clickhouse.rs @@ -461,6 +461,26 @@ fn bind_selected_node( }) } +/// Evaluate only exact integer constant arithmetic at the installation boundary. +/// No SQL text rewrite, floating coercion, or runtime-column evaluation is allowed. +fn constant_int64(expr: &QueryExpr) -> Option { + use planner_types::pre_asap::{ArithmeticOpKind, ScalarValue}; + match expr { + QueryExpr::Literal(ScalarValue::Int64(value)) => Some(*value), + QueryExpr::Arithmetic { op, left, right } => { + let left = constant_int64(left)?; + let right = constant_int64(right)?; + match op { + ArithmeticOpKind::Add => left.checked_add(right), + ArithmeticOpKind::Sub => left.checked_sub(right), + ArithmeticOpKind::Mul => left.checked_mul(right), + _ => None, + } + } + _ => None, + } +} + fn clickhouse_materialization_leaf_contract( node: &planner_types::post_asap::SummaryNode, evaluation_start_ms: u64, @@ -594,7 +614,10 @@ fn clickhouse_materialization_leaf_contract( .ok_or_else(|| { "SQL materialization predicate references an unknown column".to_string() })?; - match (name, op, right.as_ref()) { + let folded = + constant_int64(right).map(|value| QueryExpr::Literal(ScalarValue::Int64(value))); + let right = folded.as_ref().unwrap_or(right.as_ref()); + match (name, op, right) { ( name, CompareOpKind::Gt | CompareOpKind::Ge, @@ -927,6 +950,33 @@ mod tests { .contains("ambiguous")); } + #[test] + fn constant_integer_boundaries_reject_overflow_and_dynamic_values() { + use planner_types::pre_asap::{ArithmeticOpKind, ScalarValue}; + let literal = |value| QueryExpr::Literal(ScalarValue::Int64(value)); + let subtract = |left, right| QueryExpr::Arithmetic { + op: ArithmeticOpKind::Sub, + left: std::rc::Rc::new(left), + right: std::rc::Rc::new(right), + }; + assert_eq!( + constant_int64(&subtract(literal(1_788_891_296_000), literal(43_200_000))), + Some(1_788_848_096_000) + ); + assert_eq!( + constant_int64(&subtract(literal(i64::MIN), literal(1))), + None + ); + assert_eq!( + constant_int64(&subtract(QueryExpr::Column(0), literal(1))), + None + ); + assert_eq!( + constant_int64(&QueryExpr::Literal(ScalarValue::Float64(1.0))), + None + ); + } + #[test] fn sql_evaluation_rejects_empty_reversed_and_unrepresentable_ranges() { for (start_ms, end_ms) in [(2, 1), (1, 1), (0, u64::MAX)] { diff --git a/data_plane/tests/clickhouse_differential_e2e.rs b/data_plane/tests/clickhouse_differential_e2e.rs index 6f2b870a..871bb0b1 100644 --- a/data_plane/tests/clickhouse_differential_e2e.rs +++ b/data_plane/tests/clickhouse_differential_e2e.rs @@ -169,7 +169,7 @@ async fn run_mixed_aggregate(aggregate: &str) { let sql = format!("SELECT sums.timestamp, sums.total / divisors.divisor AS ratio FROM (SELECT 2000 AS timestamp, {aggregate} AS total FROM telemetry WHERE metric = 'requests' AND timestamp_ms >= 0 AND timestamp_ms < 2000) AS sums INNER JOIN divisors ON sums.timestamp = divisors.timestamp"); let grouped = aggregate == "max(value)"; let sql = if grouped { - "SELECT labels, max(value) AS value FROM telemetry WHERE metric = 'requests' AND timestamp_ms >= 0 AND timestamp_ms < 2000 GROUP BY labels ORDER BY labels".to_string() + "SELECT labels, max(value) AS value FROM telemetry WHERE metric = 'requests' AND timestamp_ms > 1999 - 2000 AND timestamp_ms <= 1999 GROUP BY labels ORDER BY labels".to_string() } else { sql }; From 91ed5b7c7eea03ba6697d7356f38a83835591568 Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 10 Sep 2026 21:32:44 -0600 Subject: [PATCH 4/7] Check nullable grouping rejection after typed transport support --- .../sketch_db/backfill/clickhouse_reader.rs | 14 ++++++++++++++ 1 file changed, 14 insertions(+) diff --git a/data_plane/src/storage_engines/sketch_db/backfill/clickhouse_reader.rs b/data_plane/src/storage_engines/sketch_db/backfill/clickhouse_reader.rs index 563a94b3..d8852029 100644 --- a/data_plane/src/storage_engines/sketch_db/backfill/clickhouse_reader.rs +++ b/data_plane/src/storage_engines/sketch_db/backfill/clickhouse_reader.rs @@ -452,6 +452,20 @@ mod tests { planner_types::pre_asap::DataType::Int64, false, )]); + assert!(factory( + &BackfillSource::ClickHouse { + database: "metrics".into(), + table: "another_table".into() + }, + &typed, + ) + .is_ok()); + typed.grouping_labels = + asap_types::GroupingProjection::new(vec![planner_types::pre_asap::Column::new( + "tenant", + planner_types::pre_asap::DataType::Int64, + true, + )]); let rejected = factory( &BackfillSource::ClickHouse { database: "metrics".into(), From c825e8041fc3f9bcca45139dc76d0065b71c1be9 Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 10 Sep 2026 21:46:02 -0600 Subject: [PATCH 5/7] Enforce canonical group keys and supported SQL output policies --- control_plane/src/clickhouse.rs | 5 +- crates/asap_types/src/grouping_projection.rs | 70 +++++++++++++++- .../asap_types/src/precompute_plan/catalog.rs | 6 ++ .../accelerator.rs | 27 ++++++ .../clickhouse_result_adapter.rs | 84 +++++++++++++++++-- .../sketch_db/backfill/clickhouse_reader.rs | 23 +++-- 6 files changed, 192 insertions(+), 23 deletions(-) diff --git a/control_plane/src/clickhouse.rs b/control_plane/src/clickhouse.rs index f6c3c3e3..19fe715c 100644 --- a/control_plane/src/clickhouse.rs +++ b/control_plane/src/clickhouse.rs @@ -240,14 +240,11 @@ fn materialize_selected_sql( .get(*key) .cloned() .ok_or("SQL grouping column is absent from source schema")?; - if column.nullable { - return Err("nullable table grouping requires an explicit null-key encoding".into()); - } column.table = None; columns.push(column); } let grouping = asap_types::GroupingProjection::new(columns); - grouping.validate()?; + grouping.validate_table_group_codec()?; let (table, value, window, population, timestamp) = clickhouse_materialization_leaf_contract(node, query.start_ms, query.end_ms)?; let window_secs = window.ok_or("SQL materialization requires a bounded window")?; diff --git a/crates/asap_types/src/grouping_projection.rs b/crates/asap_types/src/grouping_projection.rs index 2ba03aaa..0e421cce 100644 --- a/crates/asap_types/src/grouping_projection.rs +++ b/crates/asap_types/src/grouping_projection.rs @@ -24,10 +24,21 @@ impl GroupingProjection { } Ok(()) } - pub fn validate_table_columns(&self) -> Result<(), String> { + pub fn validate_table_group_codec(&self) -> Result<(), String> { self.validate()?; for column in self.columns() { crate::table_population::validate_column_name(&column.name)?; + if column.nullable { + return Err( + "nullable table grouping requires an explicit null-key encoding".into(), + ); + } + if !table_group_codec_type(&column.dtype) { + return Err( + "table group codec requires integer, string, boolean, or supported Map values" + .into(), + ); + } } Ok(()) } @@ -67,6 +78,21 @@ impl GroupingProjection { serde_json::from_value(value.clone()) } } +// Float grouping requires SQL equality canonicalization (+0/-0 and NaNs), +// and timestamps need explicit unit/timezone semantics. Neither is claimed by v1. +fn table_group_codec_type(dtype: &DataType) -> bool { + match dtype { + DataType::Int64 | DataType::Utf8 | DataType::Bool => true, + DataType::Map { key, value, .. } => { + matches!( + key.as_ref(), + DataType::Int64 | DataType::Utf8 | DataType::Bool + ) && table_group_codec_type(value) + } + _ => false, + } +} + impl From for GroupingProjection { fn from(mut names: KeyByLabelNames) -> Self { names.labels.sort(); @@ -252,6 +278,48 @@ pub fn decode_table_group_value(encoded: &str) -> Result Result { + if let Some(setting) = request + .parameters + .keys() + .find(|key| key.starts_with("output_format_")) + { + return Err(format!("unsupported output setting {setting}")); + } match request .format() .unwrap_or("TabSeparated") @@ -277,6 +284,26 @@ impl ClickHouseAccelerator for CatalogClickHouseAccelerator { #[cfg(test)] mod tests { use super::*; + #[test] + fn client_output_settings_cannot_silently_change_warm_format() { + let mut request = ClickHouseQueryRequest { + method: axum::http::Method::GET, + sql: "SELECT labels FROM samples".into(), + body: Bytes::new(), + parameters: BTreeMap::from([("default_format".into(), "JSON".into())]), + headers: HeaderMap::new(), + }; + assert!(requested_format(&request).is_ok()); + for setting in [ + "output_format_json_map_as_array_of_tuples", + "output_format_json_quote_64bit_integers", + ] { + request.parameters.insert(setting.into(), "1".into()); + assert!(requested_format(&request).is_err()); + request.parameters.remove(setting); + } + } + use crate::{ precompute_engine::operators::SumAccumulator, storage_engines::sketch_db::index::{AggKind, Capability, SketchInstanceMetadata}, 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 32883fc8..12cec2f6 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 @@ -84,6 +84,31 @@ pub fn from_series_rows( impl ClickHouseQueryResult { pub fn encode(&self, format: ClickHouseFormat) -> Result, ClickHouseResultError> { + if self.batches.iter().any(|batch| { + batch + .columns() + .iter() + .any(|column| contains_map_timestamp(column.data_type(), false)) + }) { + return Err(ClickHouseResultError::Arrow( + "nested Map timestamp output requires an explicit formatting contract".into(), + )); + } + // ClickHouse's 64-bit JSON quoting default is deployment-configurable. + // Until the client output policy is explicit, never guess it for warm output. + if matches!( + format, + ClickHouseFormat::Json | ClickHouseFormat::JsonEachRow + ) && self.batches.iter().any(|batch| { + batch + .columns() + .iter() + .any(|column| contains_json_integer64(column.data_type())) + }) { + return Err(ClickHouseResultError::Arrow( + "64-bit JSON integer output requires an explicit quoting contract".into(), + )); + } let mut output = Vec::new(); for batch in &self.batches { for row in 0..batch.num_rows() { @@ -235,6 +260,28 @@ pub(super) fn clickhouse_type(data_type: &arrow::datatypes::DataType, nullable: } } +fn contains_map_timestamp(dtype: &DataType, in_map: bool) -> bool { + match dtype { + DataType::Timestamp(..) => in_map, + DataType::Map(entries, _) => contains_map_timestamp(entries.data_type(), true), + DataType::Struct(fields) => fields + .iter() + .any(|field| contains_map_timestamp(field.data_type(), in_map)), + _ => false, + } +} + +fn contains_json_integer64(dtype: &DataType) -> bool { + match dtype { + DataType::Int64 | DataType::UInt64 => true, + DataType::Map(entries, _) => contains_json_integer64(entries.data_type()), + DataType::Struct(fields) => fields + .iter() + .any(|field| contains_json_integer64(field.data_type())), + _ => false, + } +} + fn escape_tsv(value: &str) -> String { value .replace('\\', "\\\\") @@ -261,6 +308,32 @@ mod tests { }; use std::sync::Arc; + #[test] + fn map_timestamp_transport_is_not_assumed_to_match_native_formatting() { + let entries = DataType::Struct( + vec![ + Field::new("key", DataType::Utf8, false), + Field::new( + "value", + DataType::Timestamp(arrow::datatypes::TimeUnit::Millisecond, None), + false, + ), + ] + .into(), + ); + let dtype = DataType::Map(Arc::new(Field::new("entries", entries, false)), false); + let batch = RecordBatch::try_new( + Arc::new(Schema::new(vec![Field::new("m", dtype.clone(), false)])), + vec![arrow::array::new_empty_array(&dtype)], + ) + .unwrap(); + let result = ClickHouseQueryResult { + batches: vec![batch], + }; + assert!(result.encode(ClickHouseFormat::Json).is_err()); + assert!(result.encode(ClickHouseFormat::TabSeparated).is_err()); + } + #[test] fn encodes_table_without_using_promql_query_result() { let schema = Arc::new(Schema::new(vec![ @@ -282,15 +355,8 @@ mod tests { result.encode(ClickHouseFormat::TabSeparated).unwrap(), b"a\\tb\t7\n" ); - assert_eq!( - result.encode(ClickHouseFormat::JsonEachRow).unwrap(), - b"{\"zone\":\"a\\tb\",\"count\":7}\n" - ); - let document: serde_json::Value = - serde_json::from_slice(&result.encode(ClickHouseFormat::Json).unwrap()).unwrap(); - assert_eq!(document["data"][0]["count"], 7); - assert_eq!(document["meta"][1]["type"], "Int64"); - assert_eq!(document["rows"], 1); + assert!(result.encode(ClickHouseFormat::JsonEachRow).is_err()); + assert!(result.encode(ClickHouseFormat::Json).is_err()); } } diff --git a/data_plane/src/storage_engines/sketch_db/backfill/clickhouse_reader.rs b/data_plane/src/storage_engines/sketch_db/backfill/clickhouse_reader.rs index d8852029..fb169e89 100644 --- a/data_plane/src/storage_engines/sketch_db/backfill/clickhouse_reader.rs +++ b/data_plane/src/storage_engines/sketch_db/backfill/clickhouse_reader.rs @@ -171,14 +171,9 @@ pub fn clickhouse_reader_factory(config: ClickHouseReaderConfig) -> ReaderFactor } .into()); } - materialization.grouping_labels.validate_table_columns()?; - for column in materialization.grouping_labels.columns() { - if column.nullable { - return Err( - "nullable table grouping requires an explicit null-key encoding".into(), - ); - } - } + materialization + .grouping_labels + .validate_table_group_codec()?; let mut source_config = config.clone(); source_config.database = database.clone(); source_config.table = table.clone(); @@ -297,7 +292,17 @@ impl RawSampleReader for ClickHouseReader { let count = columns.len(); let labels = columns .into_iter() - .collect::>(); + .map(|(name, value)| { + asap_types::grouping_projection::decode_table_group_value(&value) + .and_then(|value| { + asap_types::grouping_projection::encode_table_group_value( + &value, + ) + }) + .map(|value| (name, value)) + .map_err(|reason| RawSampleReaderError::Decode { reason }) + }) + .collect::, _>>()?; if labels.len() != count { return Err(RawSampleReaderError::Decode { reason: "duplicate source grouping columns".into(), From 55051d9e44bcb8c2aed722c02b8367b8d7d65bb8 Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 10 Sep 2026 20:54:45 -0600 Subject: [PATCH 6/7] style: format typed grouping consumers --- data_plane/src/drivers/ingest/otel.rs | 3 ++- .../drivers/ingest/prometheus_remote_write.rs | 6 +++-- .../src/precompute_engine/ingest_handler.rs | 22 +++++++++---------- .../src/precompute_engine/output_sink.rs | 6 +++-- .../sketch_db/backfill/processor.rs | 3 ++- .../storage_engines/sketch_db/index/mod.rs | 3 +-- 6 files changed, 23 insertions(+), 20 deletions(-) diff --git a/data_plane/src/drivers/ingest/otel.rs b/data_plane/src/drivers/ingest/otel.rs index c0627865..daf43318 100644 --- a/data_plane/src/drivers/ingest/otel.rs +++ b/data_plane/src/drivers/ingest/otel.rs @@ -606,7 +606,8 @@ fn resolve_bucket_sid_for_agg_config( point_labels: &HashMap, ) -> (u64, asap_types::PolicyFingerprint) { let grouping_pairs: Vec<(&str, &str)> = config - .grouping_labels.iter() + .grouping_labels + .iter() .map(|name| { let v = point_labels.get(name).map(|s| s.as_str()).unwrap_or(""); (name.as_str(), v) diff --git a/data_plane/src/drivers/ingest/prometheus_remote_write.rs b/data_plane/src/drivers/ingest/prometheus_remote_write.rs index 985079ac..5593e66a 100644 --- a/data_plane/src/drivers/ingest/prometheus_remote_write.rs +++ b/data_plane/src/drivers/ingest/prometheus_remote_write.rs @@ -363,7 +363,8 @@ impl PrometheusRemoteWriteReceiver { let mut labels = group_key.as_population_labels(); if labels.is_empty() { labels = config - .grouping_labels.iter() + .grouping_labels + .iter() .cloned() .zip(group_key.values().labels) .collect(); @@ -667,7 +668,8 @@ fn route_messages( Vec::new() } else { config - .grouping_labels.iter() + .grouping_labels + .iter() .map(|name| { ( name.as_str(), diff --git a/data_plane/src/precompute_engine/ingest_handler.rs b/data_plane/src/precompute_engine/ingest_handler.rs index 00958f31..e50323b3 100644 --- a/data_plane/src/precompute_engine/ingest_handler.rs +++ b/data_plane/src/precompute_engine/ingest_handler.rs @@ -278,14 +278,14 @@ impl IngestState { labels: &std::collections::HashMap, config: &AggregationConfig, ) -> Arc { - crate::precompute_engine::group_key::intern_pairs( - config.grouping_labels.iter().map(|name| { + crate::precompute_engine::group_key::intern_pairs(config.grouping_labels.iter().map( + |name| { ( name.as_str(), labels.get(name).map(String::as_str).unwrap_or(""), ) - }), - ) + }, + )) } } @@ -296,14 +296,12 @@ fn extract_group_key( config: &AggregationConfig, ) -> Arc { let labels = parse_labels_from_series_key(series_key); - crate::precompute_engine::group_key::intern_pairs(config.grouping_labels.iter().map( - |name| { - ( - name.as_str(), - labels.get(name.as_str()).copied().unwrap_or(""), - ) - }, - )) + crate::precompute_engine::group_key::intern_pairs(config.grouping_labels.iter().map(|name| { + ( + name.as_str(), + labels.get(name.as_str()).copied().unwrap_or(""), + ) + })) } #[cfg(test)] diff --git a/data_plane/src/precompute_engine/output_sink.rs b/data_plane/src/precompute_engine/output_sink.rs index 8bbef531..b7412abf 100644 --- a/data_plane/src/precompute_engine/output_sink.rs +++ b/data_plane/src/precompute_engine/output_sink.rs @@ -174,7 +174,8 @@ impl SketchStoreSink { if let Some(revision) = &output.input_revision { let group_values = output.population_labels.clone().unwrap_or_else(|| { agg_cfg - .grouping_labels.iter() + .grouping_labels + .iter() .cloned() .zip(output.key.clone().unwrap_or_default().labels) .collect() @@ -370,7 +371,8 @@ mod tests { parameters: HashMap::new(), grouping_labels: KeyByLabelNames::new( grouping_keys.iter().map(|s| s.to_string()).collect(), - ).into(), + ) + .into(), aggregated_labels: KeyByLabelNames::empty(), rollup_labels: KeyByLabelNames::empty(), original_yaml: String::new(), diff --git a/data_plane/src/storage_engines/sketch_db/backfill/processor.rs b/data_plane/src/storage_engines/sketch_db/backfill/processor.rs index 70b6d7c4..21fc95ff 100644 --- a/data_plane/src/storage_engines/sketch_db/backfill/processor.rs +++ b/data_plane/src/storage_engines/sketch_db/backfill/processor.rs @@ -127,7 +127,8 @@ fn resolve_backfill_bucket_sid( ) -> u64 { let labels = parse_labels_from_series_key(series_key); let grouping_pairs: Vec<(&str, &str)> = config - .grouping_labels.iter() + .grouping_labels + .iter() .map(|name| (name.as_str(), *labels.get(name.as_str()).unwrap_or(&""))) .collect(); let attrs_fp = canonical_attrs_fingerprint(&grouping_pairs); diff --git a/data_plane/src/storage_engines/sketch_db/index/mod.rs b/data_plane/src/storage_engines/sketch_db/index/mod.rs index 930ca6a9..4a0959a0 100644 --- a/data_plane/src/storage_engines/sketch_db/index/mod.rs +++ b/data_plane/src/storage_engines/sketch_db/index/mod.rs @@ -2666,8 +2666,7 @@ impl SketchStore { let target_metric = agg_cfg.metric.as_str(); let target_agg_type = agg_cfg.aggregation_type; let target_params = canonical_parameters(&agg_cfg.parameters); - let target_group_keys: BTreeSet = - agg_cfg.grouping_labels.iter().cloned().collect(); + let target_group_keys: BTreeSet = agg_cfg.grouping_labels.iter().cloned().collect(); // Collect the matching sids under a short read lock; then call // `remove_instance` per sid (which takes its own write lock). From eff81f480a29bc7733f1bfe1a10835cc1e25dd3c Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 10 Sep 2026 21:53:30 -0600 Subject: [PATCH 7/7] Format combined process regression modules --- data_plane/tests/asapquery_compatibility_process_e2e.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/data_plane/tests/asapquery_compatibility_process_e2e.rs b/data_plane/tests/asapquery_compatibility_process_e2e.rs index 669d9151..5b5e1e8a 100644 --- a/data_plane/tests/asapquery_compatibility_process_e2e.rs +++ b/data_plane/tests/asapquery_compatibility_process_e2e.rs @@ -21,10 +21,10 @@ use tokio::sync::Mutex; #[path = "support/erp_planning_process.rs"] mod erp_planning_process; -#[path = "support/durable_summary_process.rs"] -mod durable_summary_process; #[path = "support/distinct_planning_process.rs"] mod distinct_planning_process; +#[path = "support/durable_summary_process.rs"] +mod durable_summary_process; struct ChildGuard(Child);