From 55f4af83410d214db2181b86902d9cc1bcac8232 Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 10 Sep 2026 18:03:52 -0600 Subject: [PATCH] fix: stop maintenance execution at materialized inputs --- .../precompute_engine/maintenance_runtime.rs | 36 ++++++----- .../src/precompute_engine/subdag_scheduler.rs | 59 +++++++++++++++++++ 2 files changed, 81 insertions(+), 14 deletions(-) diff --git a/data_plane/src/precompute_engine/maintenance_runtime.rs b/data_plane/src/precompute_engine/maintenance_runtime.rs index 6cd345858..3db532e44 100644 --- a/data_plane/src/precompute_engine/maintenance_runtime.rs +++ b/data_plane/src/precompute_engine/maintenance_runtime.rs @@ -28,19 +28,20 @@ struct OperatorAdapter<'a> { impl PrecomputeOperatorRegistry for OperatorAdapter<'_> { type Error = String; + fn materialized_input(&self, node: &ExecutableDagNode) -> Result, String> { + Ok(matches!( + self.binding.node(node.id), + Some(BackendNodeBinding::Materialization { summary_definition }) + if *summary_definition == self.source_definition + ) + .then(|| Arc::clone(&self.source))) + } + fn execute( &self, node: &ExecutableDagNode, inputs: &[Arc], ) -> Result { - let is_source = matches!( - self.binding.node(node.id), - Some(BackendNodeBinding::Materialization { summary_definition }) - if *summary_definition == self.source_definition - ); - if is_source && inputs.is_empty() { - return Ok(Arc::clone(&self.source)); - } match &node.payload { ExecutableOperatorPayload::SummaryMerge => merge_inputs(inputs), ExecutableOperatorPayload::SummaryAgg { .. } => Err( @@ -326,7 +327,7 @@ impl MaintenanceDagSink { if !depends_on_any(&dag, *sink_node, &source_nodes) { continue; } - let reachable = dependencies(&dag, *sink_node); + let reachable = dependencies_until(&dag, *sink_node, &source_nodes); let foreign_source = reachable.iter().any(|node| { let has_input = dag.edges.iter().any(|edge| edge.consumer == *node); !has_input @@ -344,6 +345,7 @@ impl MaintenanceDagSink { } if dag.edges.iter().any(|edge| { reachable.contains(&edge.consumer) + && !source_nodes.contains(&edge.consumer) && !matches!( edge.grouping, planner_types::post_asap::GroupingEdgeCompatibility::Identical @@ -430,11 +432,19 @@ fn depends_on_any( fn dependencies( dag: &planner_types::post_asap::ExecutableDag, sink: PostAsapNodeId, +) -> BTreeSet { + dependencies_until(dag, sink, &BTreeSet::new()) +} + +fn dependencies_until( + dag: &planner_types::post_asap::ExecutableDag, + sink: PostAsapNodeId, + frontier: &BTreeSet, ) -> BTreeSet { let mut pending = vec![sink]; let mut seen = BTreeSet::new(); while let Some(node) = pending.pop() { - if seen.insert(node) { + if seen.insert(node) && !frontier.contains(&node) { pending.extend( dag.edges .iter() @@ -696,9 +706,7 @@ mod tests { use crate::storage_engines::types::{ ActivePhysicalPlan, HotReloadActivePhysicalPlan, StreamingConfig, }; - use control_plane::physical::executable_binding::{ - InstalledPostAsapDag, PostAsapDagDocument, - }; + use control_plane::physical::executable_binding::{InstalledPostAsapDag, OwnedPostAsapDag}; use std::sync::atomic::{AtomicUsize, Ordering}; #[derive(Default)] @@ -774,7 +782,7 @@ mod tests { bundle.precompute_plan.executable_dags = BTreeMap::from([( "retry".into(), InstalledPostAsapDag { - document: PostAsapDagDocument::from_executable("retry".into(), &dag).unwrap(), + document: OwnedPostAsapDag::from_executable("retry".into(), &dag).unwrap(), binding, }, )]); diff --git a/data_plane/src/precompute_engine/subdag_scheduler.rs b/data_plane/src/precompute_engine/subdag_scheduler.rs index 7d21bcf67..6f632f44c 100644 --- a/data_plane/src/precompute_engine/subdag_scheduler.rs +++ b/data_plane/src/precompute_engine/subdag_scheduler.rs @@ -21,6 +21,11 @@ pub struct MaterializationCommitKey { pub trait PrecomputeOperatorRegistry { type Error; + /// A materialized input is an execution frontier: its absorbed semantic + /// dependencies have already run and must not be evaluated again. + fn materialized_input(&self, _node: &ExecutableDagNode) -> Result, Self::Error> { + Ok(None) + } fn execute(&self, node: &ExecutableDagNode, inputs: &[Arc]) -> Result; } @@ -120,6 +125,14 @@ where "query-time node {id} in precompute dependency path" ))); } + if let Some(value) = registry + .materialized_input(node) + .map_err(ScheduleError::Operator)? + { + values.insert(id, Arc::new(value)); + active.remove(&id); + return Ok(()); + } let child_ids = inputs.get(&id).cloned().unwrap_or_default(); for child in &child_ids { visit(*child, nodes, inputs, active, values, registry)?; @@ -301,6 +314,52 @@ mod tests { assert_eq!(registry.0.lock().unwrap().values().sum::(), 4); } + #[test] + fn supplied_materialization_cuts_absorbed_dependencies_and_is_shared() { + struct FrontierRegistry(Registry); + impl PrecomputeOperatorRegistry for FrontierRegistry { + type Error = String; + fn materialized_input(&self, node: &ExecutableDagNode) -> Result, String> { + Ok((node.id == PostAsapNodeId(1)).then_some(10)) + } + fn execute( + &self, + node: &ExecutableDagNode, + inputs: &[Arc], + ) -> Result { + assert_ne!( + node.id, + PostAsapNodeId(0), + "absorbed source subtree must not execute" + ); + self.0.execute(node, inputs) + } + } + let mut raw = node(0); + raw.output_state = ExecutionDataState::READ_ROWS; + let mut query = node(4); + query.output_state = ExecutionDataState::READ_ROWS; + let dag = ExecutableDag { + nodes: vec![raw, node(1), node(2), node(3), query], + edges: vec![edge(0, 1), edge(1, 2), edge(1, 3), edge(2, 3), edge(3, 4)], + root: PostAsapNodeId(4), + }; + let mut bindings = binding(); + bindings + .nodes + .insert(PostAsapNodeId(0), BackendNodeBinding::QueryInput); + let registry = FrontierRegistry(Registry::default()); + let sink = Sink::default(); + let result = + execute_precompute_sink(&dag, &bindings, PostAsapNodeId(3), key(3), ®istry, &sink) + .unwrap(); + assert_eq!(*result, 25); + assert_eq!( + *registry.0 .0.lock().unwrap(), + BTreeMap::from([(2, 1), (3, 1)]) + ); + } + #[test] fn rejects_query_node_in_precompute_path_and_mismatched_lineage_key() { let mut query_child = node(0);