diff --git a/control_plane/src/clickhouse.rs b/control_plane/src/clickhouse.rs index 66e8aa186..6486063e6 100644 --- a/control_plane/src/clickhouse.rs +++ b/control_plane/src/clickhouse.rs @@ -118,6 +118,7 @@ pub async fn compile_clickhouse_workload( tables: request.tables.clone(), }; let mut entries = std::collections::BTreeMap::new(); + let mut installed_dags = std::collections::BTreeMap::new(); for query in &request.queries { let planned = plan_clickhouse_sql(&query.sql, &catalog, request.accuracy.clone()).await?; let PhysicalExpr::Committed(crate::physical::post_asap::PostAsapPlan::Summary(root)) = @@ -127,7 +128,11 @@ pub async fn compile_clickhouse_workload( "SQL did not produce a summary DAG".into(), )); }; - let executable = QueryPlanEntry::compile_bound_relational( + let semantic = planner_types::post_asap::compile_executable_dag_with_node_ids(&root) + .map_err(|error| ClickHousePlanningError::Lower(error.to_string()))?; + let mut materialization_nodes = std::collections::BTreeMap::new(); + let mut query_nodes = std::collections::BTreeMap::new(); + let executable = QueryPlanEntry::compile_bound_relational_mapped( query.sql.clone(), planned.canonical_sql.clone(), &root, @@ -142,9 +147,34 @@ pub async fn compile_clickhouse_workload( cumulative_readout: query.cumulative, }, FallbackPolicy::ExactBackend, - |node, family| bind_selected_node(node, family, query, request), + |node, family| { + let binding = bind_selected_node(node, family, query, request)?; + let id = semantic.node_ids.node_id(node).ok_or_else(|| { + crate::query_plan::QueryPlanError::Invalid( + "selected SQL node is absent from semantic DAG".into(), + ) + })?; + materialization_nodes.insert(id, binding.materialization); + Ok(binding) + }, + |node, query_node| { + if let Some(id) = semantic.node_ids.node_id(node) { + query_nodes.insert(id, query_node); + } + }, ) .map_err(|error| ClickHousePlanningError::Lower(error.to_string()))?; + let installed = crate::physical::executable_binding::install_selected_dag( + query.sql.clone(), + &semantic.dag, + executable.root, + |id| materialization_nodes.get(&id).copied(), + |id| query_nodes.get(&id).copied(), + ) + .map_err(ClickHousePlanningError::Lower)?; + crate::physical::executable_binding::validate_query_plan(&installed, &executable) + .map_err(ClickHousePlanningError::Lower)?; + installed_dags.insert(query.sql.clone(), installed); if executable .nodes .values() @@ -154,32 +184,6 @@ pub async fn compile_clickhouse_workload( "compiled SQL contains an unsupported operator; publication refused".into(), )); } - let bindings = executable.materialization_bindings(); - let identities = bindings - .iter() - .map(|binding| { - request - .sds - .materializations - .get(&binding.materialization) - .ok_or_else(|| { - ClickHousePlanningError::Lower( - "compiled SQL binding is absent from SDS".into(), - ) - }) - }) - .collect::, _>>()?; - // Descriptor references are already represented by each DAG's - // MaterializationBinding and validated through SummaryCatalog. - let _descriptor_ids = identities - .iter() - .map(|identity| { - ( - &identity.summary_descriptor_id, - &identity.data_descriptor_id, - ) - }) - .collect::>(); let identity = QueryPlan::catalog_key(QueryLanguage::ClickHouseSql, &planned.canonical_sql); if entries.insert(identity.clone(), executable).is_some() { return Err(ClickHousePlanningError::Lower(format!( @@ -187,9 +191,11 @@ pub async fn compile_clickhouse_workload( ))); } } + let mut precompute_plan = request.precompute_plan.clone(); + precompute_plan.executable_dags = installed_dags; let publication = crate::physical::publication::PhysicalPlanPublication { summary_catalog: request.sds.clone(), - precompute_plan: request.precompute_plan.clone(), + precompute_plan, collector_plans: Vec::new(), transmission_plan: request.transmission_plan.clone(), query_plan: QueryPlan { @@ -632,6 +638,21 @@ mod tests { }], }; let publication = compile_clickhouse_workload(&request).await.unwrap(); + let installed = publication + .precompute_plan + .executable_dags + .get(&request.queries[0].sql) + .unwrap(); + installed.validate().unwrap(); + assert_eq!(installed.binding.precompute_sinks.len(), 1); + assert_eq!( + installed.binding.nodes.len(), + installed.document.nodes.len() + ); + assert!(installed.binding.nodes.values().any(|binding| matches!( + binding, + crate::physical::executable_binding::BackendNodeBinding::Query { .. } + ))); let entry = publication.query_plan.entries.values().next().unwrap(); assert!(entry .nodes diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index 83b60d836..5de2565fc 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -2264,52 +2264,25 @@ impl PhysicalCompiler { for (query_index, compiled) in executable_dags.iter().enumerate() { let Some(compiled) = compiled else { continue }; let query_id = request.queries[query_index].query_id.clone(); - let mut placements = BTreeMap::new(); - let mut precompute_sinks = Vec::new(); - for node in &compiled.dag.nodes { - let placement = - if let Some(definition) = node_bindings.get(&(query_index, node.id)).copied() { - precompute_sinks.push(node.id); - crate::physical::executable_binding::BackendNodeBinding::Materialization { - summary_definition: definition.into(), - } - } else if node.output_state.timing - == planner_types::post_asap::ExecutionTiming::MaintenanceTime - { - super::executable_binding::BackendNodeBinding::MaintenanceInput - } else { - match query_node_bindings.get(&(query_index, node.id)).copied() { - Some(query_node) => { - super::executable_binding::BackendNodeBinding::Query { query_node } - } - None => super::executable_binding::BackendNodeBinding::QueryInput, - } - }; - placements.insert(node.id, placement); - } - precompute_sinks.sort(); - let installed = crate::physical::executable_binding::InstalledPostAsapDag { - document: super::executable_binding::OwnedPostAsapDag::from_executable( - query_id.clone(), - &compiled.dag, - ) - .map_err(|reason| CompileError::Query { - query_id: query_id.clone(), - reason, - })?, - binding: super::executable_binding::BackendExecutableBinding { - nodes: placements, - query_sink: compiled.dag.root, - query_plan_sink: query_plan - .entries - .values() - .find(|entry| entry.query_id == query_id) - .expect("compiled query entry exists") - .root, - precompute_sinks, + let query_plan_sink = query_plan + .entries + .values() + .find(|entry| entry.query_id == query_id) + .expect("compiled query entry exists") + .root; + let installed = super::executable_binding::install_selected_dag( + query_id.clone(), + &compiled.dag, + query_plan_sink, + |id| { + node_bindings + .get(&(query_index, id)) + .copied() + .map(Into::into) }, - }; - installed.validate().map_err(|reason| CompileError::Query { + |id| query_node_bindings.get(&(query_index, id)).copied(), + ) + .map_err(|reason| CompileError::Query { query_id: query_id.clone(), reason, })?; diff --git a/control_plane/src/physical/executable_binding.rs b/control_plane/src/physical/executable_binding.rs index 6aff751f5..28f82a57c 100644 --- a/control_plane/src/physical/executable_binding.rs +++ b/control_plane/src/physical/executable_binding.rs @@ -2,6 +2,47 @@ pub use asap_types::executable_plan::*; +/// Assign backend phases to a selected semantic DAG without changing its nodes. +pub fn install_selected_dag( + query_id: String, + dag: &planner_types::post_asap::ExecutableDag, + query_plan_sink: QueryNodeId, + materialization: impl Fn( + planner_types::post_asap::PostAsapNodeId, + ) -> Option, + query_node: impl Fn(planner_types::post_asap::PostAsapNodeId) -> Option, +) -> Result { + let mut nodes = std::collections::BTreeMap::new(); + let mut precompute_sinks = Vec::new(); + for node in &dag.nodes { + let binding = if let Some(summary_definition) = materialization(node.id) { + precompute_sinks.push(node.id); + BackendNodeBinding::Materialization { summary_definition } + } else if node.output_state.timing + == planner_types::post_asap::ExecutionTiming::MaintenanceTime + { + BackendNodeBinding::MaintenanceInput + } else { + query_node(node.id).map_or(BackendNodeBinding::QueryInput, |query_node| { + BackendNodeBinding::Query { query_node } + }) + }; + nodes.insert(node.id, binding); + } + precompute_sinks.sort(); + let installed = InstalledPostAsapDag { + document: OwnedPostAsapDag::from_executable(query_id, dag)?, + binding: BackendExecutableBinding { + nodes, + query_sink: dag.root, + query_plan_sink, + precompute_sinks, + }, + }; + installed.validate()?; + Ok(installed) +} + pub fn validate_query_plan( installed: &InstalledPostAsapDag, query: &crate::query_plan::QueryPlanEntry, diff --git a/control_plane/src/query_plan.rs b/control_plane/src/query_plan.rs index ff9d0e78b..1fbf7827f 100644 --- a/control_plane/src/query_plan.rs +++ b/control_plane/src/query_plan.rs @@ -402,6 +402,34 @@ impl QueryPlanEntry { } pub fn compile_bound_relational( + query_id: String, + canonical_query: String, + root: &Rc, + fixed_evaluation: FixedEvaluationRange, + instant: InstantExecution, + fallback: FallbackPolicy, + bind: F, + ) -> Result + where + F: FnMut( + &Rc, + &SummaryFamilyType, + ) -> Result, + { + Self::compile_bound_relational_mapped( + query_id, + canonical_query, + root, + fixed_evaluation, + instant, + fallback, + bind, + |_, _| {}, + ) + } + + /// Preserve Planner-to-runtime node identities for installed SQL DAGs. + pub fn compile_bound_relational_mapped( query_id: String, canonical_query: String, root: &Rc, @@ -409,12 +437,14 @@ impl QueryPlanEntry { instant: InstantExecution, fallback: FallbackPolicy, mut bind: F, + mut lowered: G, ) -> Result where F: FnMut( &Rc, &SummaryFamilyType, ) -> Result, + G: FnMut(&Rc, QueryNodeId), { let mut compiler = DagCompiler { next_id: 0, @@ -423,7 +453,7 @@ impl QueryPlanEntry { bind: &mut bind, logical_source: None, preserve_relational: true, - lowered: None, + lowered: Some(&mut lowered), }; let root = compiler.lower(root)?; Ok(Self { diff --git a/docs/developer_docs/query-engine/clickhouse-sql-support.md b/docs/developer_docs/query-engine/clickhouse-sql-support.md index 3668aca7b..e2d62f534 100644 --- a/docs/developer_docs/query-engine/clickhouse-sql-support.md +++ b/docs/developer_docs/query-engine/clickhouse-sql-support.md @@ -41,6 +41,13 @@ pane duration, and `pane_origin_ms` through the authoritative SummaryCatalog. Any invalid SQL entry rejects the complete candidate snapshot before activation; the active generation remains unchanged. +SQL compilation also retains the selected Planner semantic DAG in +`PrecomputePlan.executable_dags`. The compiler records materialization and query +node bindings during lowering and assigns phases with the same placement builder +as PromQL. Planner node IDs remain distinct from SummaryDefinitionId and +QueryNodeId. This preserves the actual selected DAG across publication instead +of reconstructing it from materialization configs later. + The query listener snapshots `HotReloadActivePhysicalPlan` once per request. It uses the SQL parsing context and the matching `QueryPlanEntry` from that same snapshot, reads SummaryStore state, executes relational operators, and