From 353f4971eaffd21b9653d1fed8e5a745503f5f5d Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 10 Sep 2026 07:55:25 -0600 Subject: [PATCH] feat(clickhouse): compile bounded SQL summaries in shared plans --- control_plane/src/clickhouse.rs | 213 ++++++++++++++++++++++++++------ control_plane/src/main.rs | 15 +-- 2 files changed, 179 insertions(+), 49 deletions(-) diff --git a/control_plane/src/clickhouse.rs b/control_plane/src/clickhouse.rs index 4f547b5d..fbb85a28 100644 --- a/control_plane/src/clickhouse.rs +++ b/control_plane/src/clickhouse.rs @@ -15,8 +15,9 @@ use crate::query_plan::{ MaterializationBinding, PhysicalGrouping, QueryLanguage, QueryPlan, QueryPlanEntry, }; use asap_types::summary_catalog::SummaryCatalog; -use serde::{Deserialize, Serialize}; +use serde::Deserialize; use std::collections::HashMap; +use std::rc::Rc; #[derive(Debug, thiserror::Error)] pub enum ClickHousePlanningError { @@ -43,8 +44,20 @@ pub async fn plan_clickhouse_sql( // SQL keeps relational parents such as Project and Filter above a // summary-capable Aggregate. Use ASAPPlanner's recursive selector here; // the PromQL deployment lowering retains its existing conservative rules. - let cost_model = ControlPlaneCostModel::new(accuracy); - let selected = crate::planner_selection::select_summary(&canonical, &cost_model)?; + let cost_model = ControlPlaneCostModel::new(accuracy.clone()); + let selected = crate::planner_selection::select_workload( + vec![(0, Rc::new(canonical.clone()))], + accuracy, + &cost_model, + )? + .into_iter() + .next() + .map(|(_, node)| node) + .ok_or_else(|| { + crate::planner_selection::SelectionError::Workload( + "SQL workload search returned no root".into(), + ) + })?; let physical = PhysicalExpr::committed(selected); Ok(ClickHousePlannedQuery { canonical_sql: canonical_sql_identity(&canonical), @@ -90,20 +103,9 @@ pub struct ClickHouseSqlWorkloadEntry { pub cumulative: bool, } -/// Physical-plan components produced for the normal atomic install path. -#[derive(Debug, Serialize)] -pub struct ClickHouseCompiledBundle { - pub sds: SummaryCatalog, - pub tables: HashMap, - pub accuracy: AccuracyTarget, - pub query_plan: QueryPlan, - pub precompute_plan: PrecomputePlan, - pub transmission_plan: TransmissionPlan, -} - pub async fn compile_clickhouse_workload( request: &ClickHouseSqlWorkload, -) -> Result { +) -> Result { request .precompute_plan .validate_against_catalog(&request.sds) @@ -185,10 +187,11 @@ pub async fn compile_clickhouse_workload( ))); } } - Ok(ClickHouseCompiledBundle { - sds: request.sds.clone(), - tables: request.tables.clone(), - accuracy: request.accuracy.clone(), + let publication = crate::physical::publication::PhysicalPlanPublication { + summary_catalog: request.sds.clone(), + precompute_plan: request.precompute_plan.clone(), + collector_plans: Vec::new(), + transmission_plan: request.transmission_plan.clone(), query_plan: QueryPlan { plan_id: request.sds.plan_id, plan_version: request.sds.plan_version, @@ -198,9 +201,11 @@ pub async fn compile_clickhouse_workload( }), entries, }, - precompute_plan: request.precompute_plan.clone(), - transmission_plan: request.transmission_plan.clone(), - }) + }; + publication + .validate() + .map_err(ClickHousePlanningError::Lower)?; + Ok(publication) } fn bind_selected_node( @@ -209,13 +214,14 @@ fn bind_selected_node( query: &ClickHouseSqlWorkloadEntry, request: &ClickHouseSqlWorkload, ) -> Result { - let (metric, source_window, spatial_filter) = - crate::physical::compiler::materialization_leaf_contract(node) + let (table_ref, value_column, source_window, spatial_filter) = + clickhouse_materialization_leaf_contract(node) .map_err(crate::query_plan::QueryPlanError::Invalid)?; let expected = crate::physical::compiler::physical_materialization_family(family); let selected = select_materialization( &request.precompute_plan.materializations, - &metric, + &table_ref, + &value_column, &spatial_filter, &expected, source_window.unwrap_or((query.end_ms.saturating_sub(query.start_ms)) / 1000), @@ -230,15 +236,144 @@ fn bind_selected_node( }) } +fn clickhouse_materialization_leaf_contract( + node: &planner_types::post_asap::SummaryNode, +) -> Result<(String, String, Option, String), String> { + use planner_types::{ + post_asap::SummaryExpr, + pre_asap::{CompareOpKind, QueryExpr, ScalarValue, Source}, + }; + let SummaryExpr::SummaryAgg { child, input, .. } = &node.expr else { + return Err("SQL materialization leaf is not a summary aggregate".into()); + }; + let SummaryExpr::KeepPreAsap(expr) = &child.expr else { + return Err("SQL materialization leaf has no tabular source".into()); + }; + fn table_scan( + expr: &QueryExpr, + ) -> Option<( + &str, + &[planner_types::pre_asap::Predicate], + &planner_types::pre_asap::Schema, + )> { + match expr { + QueryExpr::Project { child, .. } => table_scan(child), + QueryExpr::Scan { + source: Source::Table { table_ref }, + predicates, + schema, + } => Some((table_ref, predicates, schema)), + _ => None, + } + } + let (source, explicit_window) = match expr.as_ref() { + QueryExpr::TimeRange { child, range } + if range.as_millis() > 0 && range.as_millis() % 1_000 == 0 => + { + (child.as_ref(), Some(range.as_secs())) + } + source => (source, None), + }; + let Some((table_ref, predicates, schema)) = table_scan(source) else { + return Err("SQL materialization leaf has no table scan".into()); + }; + let planner_types::post_asap::SummaryInputExpr::Column(value_input) = &input.weight else { + return Err("SQL summary requires a column-valued update".into()); + }; + let value_column = match value_input { + planner_types::pre_asap::ColumnRef::Named(name) + | planner_types::pre_asap::ColumnRef::Qualified { name, .. } => name.clone(), + _ => return Err("SQL summary requires a named value column".into()), + }; + + fn comparisons<'a>(expr: &'a QueryExpr, out: &mut Vec<&'a QueryExpr>) { + if let QueryExpr::BoolAnd(children) = expr { + for child in children { + comparisons(child, out); + } + } else { + out.push(expr); + } + } + let mut lower_ms = None; + let mut upper_ms = None; + let mut leaves = Vec::new(); + for predicate in predicates { + comparisons(&predicate.0, &mut leaves); + } + for leaf in leaves { + let QueryExpr::Compare { left, op, right } = leaf else { + return Err("SQL materialization predicate is not a comparison".into()); + }; + let QueryExpr::Column(column) = left.as_ref() else { + return Err("SQL materialization predicate must reference a column".into()); + }; + let name = schema + .columns + .get(*column) + .map(|column| column.name.as_str()) + .ok_or_else(|| { + "SQL materialization predicate references an unknown column".to_string() + })?; + match (name, op, right.as_ref()) { + ( + name, + CompareOpKind::Gt | CompareOpKind::Ge, + QueryExpr::Literal(ScalarValue::Int64(value)), + ) if schema + .time_index + .is_some_and(|index| schema.columns[index].name == name) => + { + lower_ms = Some(*value); + } + ( + name, + CompareOpKind::Lt | CompareOpKind::Le, + QueryExpr::Literal(ScalarValue::Int64(value)), + ) if schema + .time_index + .is_some_and(|index| schema.columns[index].name == name) => + { + upper_ms = Some(*value); + } + _ => { + return Err(format!( + "SQL population predicate on {name} needs a canonical catalog filter" + )) + } + } + } + let inferred_window = match (lower_ms, upper_ms) { + (Some(lower), Some(upper)) if upper > lower && (upper - lower) % 1_000 == 0 => { + Some((upper - lower) as u64 / 1_000) + } + (None, None) => None, + _ => { + return Err("SQL timestamp range must provide compatible lower and upper bounds".into()) + } + }; + let window_secs = explicit_window.or(inferred_window).ok_or_else(|| { + "SQL table summary requires a positive whole-second timestamp range".to_string() + })?; + Ok(( + table_ref.to_owned(), + value_column, + Some(window_secs), + String::new(), + )) +} + fn select_materialization<'a>( materializations: &'a [asap_types::PrecomputeMaterialization], - metric: &str, + table_ref: &str, + value_column: &str, spatial_filter: &str, expected: &planner_types::post_asap::SummaryFamilyType, semantic_window_seconds: u64, ) -> Result<&'a asap_types::PrecomputeMaterialization, crate::query_plan::QueryPlanError> { let mut matches = materializations.iter().filter(|candidate| { - candidate.metric == metric + candidate.table_name.as_deref() == Some(table_ref) + && candidate.value_column.as_deref() == Some(value_column) && candidate.spatial_filter_normalized == spatial_filter && candidate .accumulator_spec() @@ -250,12 +385,12 @@ fn select_materialization<'a>( }); let selected = matches.next().ok_or_else(|| { crate::query_plan::QueryPlanError::Invalid(format!( - "no precompute materialization matches {metric}/{expected:?}" + "no precompute materialization matches {table_ref}.{value_column}/{expected:?}" )) })?; if matches.next().is_some() { return Err(crate::query_plan::QueryPlanError::Invalid(format!( - "ambiguous precompute materializations match {metric}/{expected:?}" + "ambiguous precompute materializations match {table_ref}.{value_column}/{expected:?}" ))); } Ok(selected) @@ -268,7 +403,7 @@ mod tests { fn materialization( agg: AggregationType, - metric: &str, + value_column: &str, window: u64, slide: u64, parameter: (&str, serde_json::Value), @@ -285,10 +420,10 @@ mod tests { slide, WindowKind::Tumbling, String::new(), - metric.into(), - None, - None, + format!("telemetry.{value_column}"), None, + Some("telemetry".into()), + Some(value_column.into()), ); value.pane_origin_ms = Some(0); value @@ -349,29 +484,31 @@ mod tests { let sum_family = sum_60.accumulator_spec().unwrap().family; let count_family = count_60.accumulator_spec().unwrap().family; assert_eq!( - select_materialization(&configs, "requests", "", &sum_family, 60) + select_materialization(&configs, "telemetry", "requests", "", &sum_family, 60) .unwrap() .policy_fingerprint(), sum_60.policy_fingerprint() ); let dd_family = dd_2.accumulator_spec().unwrap().family; assert_eq!( - select_materialization(&configs, "requests", "", &dd_family, 60) + select_materialization(&configs, "telemetry", "requests", "", &dd_family, 60) .unwrap() .policy_fingerprint(), dd_2.policy_fingerprint() ); assert_eq!( - select_materialization(&configs, "requests", "", &count_family, 60) + select_materialization(&configs, "telemetry", "requests", "", &count_family, 60) .unwrap() .policy_fingerprint(), count_60.policy_fingerprint() ); - assert!(select_materialization(&configs, "missing", "", &sum_family, 60).is_err()); + assert!( + select_materialization(&configs, "telemetry", "missing", "", &sum_family, 60).is_err() + ); let mut ambiguous = configs.clone(); ambiguous.push(sum_60); assert!( - select_materialization(&ambiguous, "requests", "", &sum_family, 60) + select_materialization(&ambiguous, "telemetry", "requests", "", &sum_family, 60) .unwrap_err() .to_string() .contains("ambiguous") diff --git a/control_plane/src/main.rs b/control_plane/src/main.rs index 8c282bb1..65b5a9d2 100644 --- a/control_plane/src/main.rs +++ b/control_plane/src/main.rs @@ -809,12 +809,12 @@ async fn handle_compile_and_publish_clickhouse_plan( State(state): State, Json(request): Json, ) -> impl IntoResponse { - let bundle = match clickhouse::compile_clickhouse_workload(&request).await { - Ok(bundle) => bundle, + let publication = match clickhouse::compile_clickhouse_workload(&request).await { + Ok(publication) => publication, Err(error) => return (StatusCode::BAD_REQUEST, error.to_string()).into_response(), }; - let plan_id = bundle.sds.plan_id; - let plan_version = bundle.sds.plan_version; + let plan_id = publication.summary_catalog.plan_id; + let plan_version = publication.summary_catalog.plan_version; let Some(client) = state.backend_client.as_ref() else { return ( StatusCode::SERVICE_UNAVAILABLE, @@ -822,13 +822,6 @@ async fn handle_compile_and_publish_clickhouse_plan( ) .into_response(); }; - let publication = physical::publication::PhysicalPlanPublication { - summary_catalog: bundle.sds, - precompute_plan: bundle.precompute_plan, - collector_plans: Vec::new(), - transmission_plan: bundle.transmission_plan, - query_plan: bundle.query_plan, - }; if let Err(error) = client .post_catalog_plan_typed(&publication, None, &[]) .await