From b184638fff11caaf959f84de93b00f5c043ba123 Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 10 Sep 2026 17:34:26 -0600 Subject: [PATCH] fix: reject unevaluated maintenance summary updates --- .../precompute_engine/maintenance_runtime.rs | 36 +++++++++++++++++-- docs/developer_docs/maintenance-replay.md | 8 +++++ 2 files changed, 42 insertions(+), 2 deletions(-) diff --git a/data_plane/src/precompute_engine/maintenance_runtime.rs b/data_plane/src/precompute_engine/maintenance_runtime.rs index a74340b9..f894fb2e 100644 --- a/data_plane/src/precompute_engine/maintenance_runtime.rs +++ b/data_plane/src/precompute_engine/maintenance_runtime.rs @@ -42,8 +42,10 @@ impl PrecomputeOperatorRegistry for OperatorAdapter<'_> { return Ok(Arc::clone(&self.source)); } match &node.payload { - ExecutableOperatorPayload::SummaryAgg { .. } - | ExecutableOperatorPayload::SummaryMerge => merge_inputs(inputs), + ExecutableOperatorPayload::SummaryMerge => merge_inputs(inputs), + ExecutableOperatorPayload::SummaryAgg { .. } => Err( + "maintenance SummaryAgg requires a typed update evaluator; merging input state does not execute its update expression".into(), + ), payload => Err(format!( "maintenance operator {:?} has no summary-state implementation", payload.operator() @@ -477,6 +479,36 @@ mod tests { Arc::new(accumulator) } + #[test] + fn summary_aggregation_does_not_silently_reuse_input_family() { + use planner_types::post_asap::{ + ExactKind, ExactParams, GroupingStrategy, SummaryFamilyType, SummaryUpdate, + }; + use planner_types::pre_asap::{ColumnRef, Reduction}; + + let binding = BackendExecutableBinding { + nodes: BTreeMap::new(), + query_sink: PostAsapNodeId(2), + query_plan_sink: control_plane::query_plan::QueryNodeId(2), + precompute_sinks: vec![PostAsapNodeId(1)], + }; + let adapter = OperatorAdapter { + binding: &binding, + source_definition: definition(1), + source: sum(7.0), + }; + let mut aggregate = node(1); + aggregate.operator = ExecutableOperator::SummaryAgg; + aggregate.payload = ExecutableOperatorPayload::SummaryAgg { + family: SummaryFamilyType::ExactAggregate(ExactKind::Count, ExactParams::Count), + input: SummaryUpdate::column(ColumnRef::SampleValue), + reduction: Reduction::by(vec![]), + grouping: GroupingStrategy::default(), + }; + let error = adapter.execute(&aggregate, &[Arc::new(sum(7.0))]); + assert!(matches!(error, Err(reason) if reason.contains("typed update evaluator"))); + } + // Receipts follow the declared event-time horizon and never retain accepted // summary payloads; expired retries fail instead of becoming duplicate writes. #[test] diff --git a/docs/developer_docs/maintenance-replay.md b/docs/developer_docs/maintenance-replay.md index 24ed3f5a..59e24170 100644 --- a/docs/developer_docs/maintenance-replay.md +++ b/docs/developer_docs/maintenance-replay.md @@ -1,5 +1,13 @@ # Maintenance replay and retention +The production maintenance adapter accepts an already computed source summary +and executes homogeneous `SummaryMerge` dependencies. An additional `SummaryAgg` +requires evaluation of its declared update expression and target family; it is +rejected until that evaluator exists. Copying or merging its input would silently +ignore those semantics (for example, count over a sum state is not that sum). +The scheduler's ability to traverse a DAG does not imply every operator is +implemented by the production adapter. + The maintenance sink retains in-process publication receipts for the active physical plan. Receipt identity includes the plan generation, target summary definition, output window, and a SHA-256 digest of the source definition,