From 1149b4e40e2caeddbd107eb5e1c7280f2e6dfc21 Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 11 Sep 2026 12:41:54 -0600 Subject: [PATCH 1/3] feat: represent exact maintenance binary timing explicitly --- crates/asap-aware-mapping/src/replacement.rs | 52 ++++++++++++++ .../src/summary_maintenance_cost/model.rs | 1 + .../tests/promql_to_post_asap.rs | 38 ++++++++++ crates/types/src/post_asap/cse.rs | 5 +- crates/types/src/post_asap/executable_dag.rs | 7 +- .../src/post_asap/execution_data_state.rs | 71 ++++++++++++++++--- crates/types/src/post_asap/expr.rs | 2 + 7 files changed, 166 insertions(+), 10 deletions(-) diff --git a/crates/asap-aware-mapping/src/replacement.rs b/crates/asap-aware-mapping/src/replacement.rs index bc073e0b..6ef02331 100644 --- a/crates/asap-aware-mapping/src/replacement.rs +++ b/crates/asap-aware-mapping/src/replacement.rs @@ -1796,6 +1796,7 @@ fn realize_binary( Ok(Some(Rc::new(SummaryNode { expr: SummaryExpr::BinaryOp { + timing: ExecutionTiming::ReadTime, lhs: lhs_node, rhs: rhs_node, operator: asap_types::post_asap::BinaryOperator { @@ -2179,6 +2180,56 @@ fn realize_physical_summary_input( /// Emit `SummaryAgg` (recursively binding the child), plus the /// `SummaryEstimate` readout when `estimate` is set. #[allow(clippy::too_many_arguments)] +// Retain the exact expression and schema while placing its value production +// on the update path. Read-time consumers keep their original shared nodes. +fn maintenance_exact_values(node: Rc) -> Rc { + let expr = match &node.expr { + SummaryExpr::BinaryOp { + lhs, rhs, operator, .. + } if operator.vector_match.is_none() + && matches!( + operator.kind, + asap_types::pre_asap::BinaryOpKind::Arithmetic(_) + ) + && node + .guarantee + .as_ref() + .is_some_and(ResultGuarantee::is_exact) => + { + SummaryExpr::BinaryOp { + lhs: maintenance_exact_values(lhs.clone()), + rhs: maintenance_exact_values(rhs.clone()), + operator: operator.clone(), + timing: ExecutionTiming::MaintenanceTime, + } + } + SummaryExpr::ValueOperation { + child, + operation: ValueOperation::FinalizeExactAccumulator, + .. + } if matches!( + child.expr, + SummaryExpr::SummaryAgg { + family: SummaryFamilyType::ExactAggregate(..), + .. + } + ) => + { + SummaryExpr::ValueOperation { + child: child.clone(), + operation: ValueOperation::FinalizeExactAccumulator, + timing: ExecutionTiming::MaintenanceTime, + } + } + _ => return node, + }; + Rc::new(SummaryNode { + expr, + schema: node.schema.clone(), + guarantee: node.guarantee.clone(), + }) +} + fn construct_summary_agg( node: &QueryExpr, reduction: &Reduction, @@ -2241,6 +2292,7 @@ fn construct_summary_agg( // an exact scalar accumulator currently stores its value directly. let bound_child = finalize_exact_accumulator_at(bound_child, &input.child, ExecutionTiming::MaintenanceTime)?; + let bound_child = maintenance_exact_values(bound_child); // ── Guarantee (issue #172) ────────────────────────────────────────── // Derived *before* the node exists, so an illegal composition is never diff --git a/crates/asap-aware-mapping/src/summary_maintenance_cost/model.rs b/crates/asap-aware-mapping/src/summary_maintenance_cost/model.rs index 408f0c3f..00ef21d8 100644 --- a/crates/asap-aware-mapping/src/summary_maintenance_cost/model.rs +++ b/crates/asap-aware-mapping/src/summary_maintenance_cost/model.rs @@ -3000,6 +3000,7 @@ mod tests { let operand = summary_with_operations(false, false, false); Rc::new(SummaryNode { expr: SummaryExpr::BinaryOp { + timing: asap_types::post_asap::ExecutionTiming::ReadTime, lhs: Rc::clone(&operand), rhs: operand, operator: asap_types::post_asap::BinaryOperator { diff --git a/crates/integration-tests/tests/promql_to_post_asap.rs b/crates/integration-tests/tests/promql_to_post_asap.rs index 6892c16a..0668a4d3 100644 --- a/crates/integration-tests/tests/promql_to_post_asap.rs +++ b/crates/integration-tests/tests/promql_to_post_asap.rs @@ -943,3 +943,41 @@ fn nested_summary_explicitly_finalizes_exact_child_at_maintenance_time() { .any(|field| matches!(field.dtype, SummaryFamilyType::Plain(DataType::Float64)))); compile_executable_dag(&plan).expect("explicit boundary is a valid executable DAG"); } + +#[test] +fn exact_binary_maintenance_has_explicit_timing_and_legacy_wire_default() { + use asap_types::post_asap::{ExecutableOperatorPayload, ExecutionTiming}; + for (query, expected) in [ + ( + "quantile(0.9, sum_over_time(m[1m]) + sum_over_time(n[1m]))", + ExecutionTiming::MaintenanceTime, + ), + ( + "sum_over_time(m[1m]) + sum_over_time(n[1m])", + ExecutionTiming::ReadTime, + ), + ] { + let input = lower_promql(query, AccuracyTarget::Epsilon(0.05)).unwrap(); + let search = search_workload(vec![("q", Rc::new(input))]); + let choice = search.global_selection(&DefaultCostModel); + let plan = choice.materialize(&search.roots[0].1).unwrap().unwrap(); + let dag = compile_executable_dag(&plan).unwrap(); + let payload = dag + .nodes + .iter() + .find_map(|node| { + matches!(node.payload, ExecutableOperatorPayload::Binary { .. }) + .then_some(&node.payload) + }) + .unwrap(); + assert!( + matches!(payload, ExecutableOperatorPayload::Binary { timing, .. } if *timing == expected) + ); + let wire = serde_json::to_value(payload).unwrap(); + if expected == ExecutionTiming::ReadTime { + assert!(wire.get("timing").is_none()); + } + let restored: ExecutableOperatorPayload = serde_json::from_value(wire).unwrap(); + assert_eq!(&restored, payload); + } +} diff --git a/crates/types/src/post_asap/cse.rs b/crates/types/src/post_asap/cse.rs index fb4de31b..92e29913 100644 --- a/crates/types/src/post_asap/cse.rs +++ b/crates/types/src/post_asap/cse.rs @@ -31,13 +31,15 @@ fn same_node(left: &SummaryNode, right: &SummaryNode) -> bool { lhs: al, rhs: ar, operator: ao, + timing: at, }, BinaryOp { lhs: bl, rhs: br, operator: bo, + timing: bt, }, - ) => Rc::ptr_eq(al, bl) && Rc::ptr_eq(ar, br) && ao == bo, + ) => Rc::ptr_eq(al, bl) && Rc::ptr_eq(ar, br) && ao == bo && at == bt, ( CandidateTopK { candidates: ac, @@ -411,6 +413,7 @@ mod tests { for _ in 0..24 { current = Rc::new(SummaryNode { expr: SummaryExpr::BinaryOp { + timing: super::super::ExecutionTiming::ReadTime, lhs: Rc::clone(¤t), rhs: current, operator: super::super::BinaryOperator { diff --git a/crates/types/src/post_asap/executable_dag.rs b/crates/types/src/post_asap/executable_dag.rs index 3f357b75..f13a415f 100644 --- a/crates/types/src/post_asap/executable_dag.rs +++ b/crates/types/src/post_asap/executable_dag.rs @@ -70,6 +70,8 @@ pub enum ExecutableOperatorPayload { expression: QueryExpr, }, Binary { + #[serde(default, skip_serializing_if = "ExecutionTiming::is_read_time")] + timing: ExecutionTiming, operator: BinaryOperator, }, CandidateTopK { @@ -444,7 +446,10 @@ pub fn compile_executable_dag_with_node_ids( SummaryExpr::KeepPreAsap(expression) => ExecutableOperatorPayload::Fallback { expression: (**expression).clone(), }, - SummaryExpr::BinaryOp { operator, .. } => ExecutableOperatorPayload::Binary { + SummaryExpr::BinaryOp { + operator, timing, .. + } => ExecutableOperatorPayload::Binary { + timing: *timing, operator: operator.clone(), }, SummaryExpr::CandidateTopK { diff --git a/crates/types/src/post_asap/execution_data_state.rs b/crates/types/src/post_asap/execution_data_state.rs index 8f5bfb14..3befce97 100644 --- a/crates/types/src/post_asap/execution_data_state.rs +++ b/crates/types/src/post_asap/execution_data_state.rs @@ -58,7 +58,15 @@ pub enum ExecutionTiming { ReadTime, } +impl Default for ExecutionTiming { + fn default() -> Self { + Self::ReadTime + } +} impl ExecutionTiming { + pub fn is_read_time(&self) -> bool { + *self == Self::ReadTime + } pub fn as_str(self) -> &'static str { match self { Self::MaintenanceTime => "maintenance_time", @@ -188,6 +196,8 @@ pub enum ExecutionDataStateError { /// nothing maintains state above it, so its output is never read. #[error("A maintenance-time value operation cannot be a plan root: its update-path output feeds nothing")] MaintenanceRowsAtRoot, + #[error("unsupported maintenance binary schema or operator")] + InvalidMaintenanceBinary, /// An `ExactOperation` whose input columns are not all `Plain` at its /// declared data_state. #[error("exact operator consumes non-plain column {column:?} ({dtype})")] @@ -223,9 +233,13 @@ impl ExecutionDataStateAssignment { pub fn produced_data_state(expr: &SummaryExpr) -> Option { Some(match expr { SummaryExpr::KeepPreAsap(_) => return None, - SummaryExpr::BinaryOp { .. } - | SummaryExpr::CandidateTopK { .. } - | SummaryExpr::RelationalJoin { .. } => ExecutionDataState::READ_ROWS, + SummaryExpr::BinaryOp { timing, .. } => ExecutionDataState { + timing: *timing, + primitive: DataPrimitive::Raw, + }, + SummaryExpr::CandidateTopK { .. } | SummaryExpr::RelationalJoin { .. } => { + ExecutionDataState::READ_ROWS + } SummaryExpr::SummaryAgg { .. } | SummaryExpr::SummaryJoin { .. } | SummaryExpr::SummarySubtract { .. } @@ -312,11 +326,45 @@ fn visit( match &node.expr { SummaryExpr::KeepPreAsap(_) => Ok(()), - SummaryExpr::BinaryOp { lhs, rhs, .. } => { + SummaryExpr::BinaryOp { + lhs, + rhs, + timing, + operator, + } => { + if *timing == ExecutionTiming::MaintenanceTime { + use crate::pre_asap::{BinaryOpKind, DataType}; + if operator.vector_match.is_some() + || !matches!(operator.kind, BinaryOpKind::Arithmetic(_)) + || lhs.schema != rhs.schema + || lhs.schema != node.schema + || !node.schema.fields.iter().all(|field| { + !field.nullable + && matches!( + field.dtype, + SummaryFamilyType::Plain(DataType::Float64 | DataType::Timestamp) + ) + }) + || node + .schema + .fields + .iter() + .filter(|field| { + matches!(field.dtype, SummaryFamilyType::Plain(DataType::Float64)) + }) + .count() + != 1 + { + return Err(ExecutionDataStateError::InvalidMaintenanceBinary); + } + } + let expected = ExecutionDataState { + timing: *timing, + primitive: DataPrimitive::Raw, + }; for input in [lhs, rhs] { - let state = - produced_data_state(&input.expr).unwrap_or(ExecutionDataState::READ_ROWS); - if state != ExecutionDataState::READ_ROWS { + let state = produced_data_state(&input.expr).unwrap_or(expected); + if state != expected { return Err(ExecutionDataStateError::IllegalChildDataState { edge: "BinaryOp operand", child: state, @@ -453,9 +501,16 @@ pub fn assigned_child_data_state(parent: &SummaryExpr, child: &SummaryNode) -> E SummaryExpr::ValueOperation { timing: ExecutionTiming::ReadTime, .. + } + | SummaryExpr::BinaryOp { + timing: ExecutionTiming::ReadTime, + .. } => ExecutionDataState::READ_ROWS, SummaryExpr::KeepPreAsap(_) - | SummaryExpr::BinaryOp { .. } + | SummaryExpr::BinaryOp { + timing: ExecutionTiming::MaintenanceTime, + .. + } | SummaryExpr::CandidateTopK { .. } | SummaryExpr::RelationalJoin { .. } | SummaryExpr::SummaryAgg { .. } diff --git a/crates/types/src/post_asap/expr.rs b/crates/types/src/post_asap/expr.rs index 2e82150d..f0110717 100644 --- a/crates/types/src/post_asap/expr.rs +++ b/crates/types/src/post_asap/expr.rs @@ -1,3 +1,4 @@ +use super::ExecutionTiming; use std::rc::Rc; use super::guarantee::ResultGuarantee; @@ -120,6 +121,7 @@ pub enum SummaryExpr { /// This keeps realizable summary/readout leaves visible instead of /// hiding the complete expression inside `KeepPreAsap`. BinaryOp { + timing: ExecutionTiming, lhs: Rc, rhs: Rc, operator: BinaryOperator, From 24add29e04172c800df101679af820f136dade3d Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 11 Sep 2026 13:03:36 -0600 Subject: [PATCH 2/3] style: derive the legacy read-time default --- crates/types/src/post_asap/execution_data_state.rs | 10 ++++------ 1 file changed, 4 insertions(+), 6 deletions(-) diff --git a/crates/types/src/post_asap/execution_data_state.rs b/crates/types/src/post_asap/execution_data_state.rs index 3befce97..fba174b0 100644 --- a/crates/types/src/post_asap/execution_data_state.rs +++ b/crates/types/src/post_asap/execution_data_state.rs @@ -52,17 +52,15 @@ use crate::pre_asap::query_expr::{aggregate_output_schema, QueryExprError}; use crate::pre_asap::schema::{Column, Schema}; /// When a post-ASAP value is produced. -#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)] +#[derive( + Default, Debug, Clone, Copy, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize, +)] pub enum ExecutionTiming { MaintenanceTime, + #[default] ReadTime, } -impl Default for ExecutionTiming { - fn default() -> Self { - Self::ReadTime - } -} impl ExecutionTiming { pub fn is_read_time(&self) -> bool { *self == Self::ReadTime From 73cf53188fba9ebdb6bb7d53d1626361deca658d Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 11 Sep 2026 13:07:54 -0600 Subject: [PATCH 3/3] style: retain constructor lint annotation --- crates/asap-aware-mapping/src/replacement.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/crates/asap-aware-mapping/src/replacement.rs b/crates/asap-aware-mapping/src/replacement.rs index 6ef02331..fccc188d 100644 --- a/crates/asap-aware-mapping/src/replacement.rs +++ b/crates/asap-aware-mapping/src/replacement.rs @@ -2179,7 +2179,6 @@ fn realize_physical_summary_input( /// Emit `SummaryAgg` (recursively binding the child), plus the /// `SummaryEstimate` readout when `estimate` is set. -#[allow(clippy::too_many_arguments)] // Retain the exact expression and schema while placing its value production // on the update path. Read-time consumers keep their original shared nodes. fn maintenance_exact_values(node: Rc) -> Rc { @@ -2230,6 +2229,7 @@ fn maintenance_exact_values(node: Rc) -> Rc { }) } +#[allow(clippy::too_many_arguments)] fn construct_summary_agg( node: &QueryExpr, reduction: &Reduction,