From 53911070590d0950064cce85bc767cea1fe4fd07 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Sat, 3 Oct 2026 03:49:38 +0000 Subject: [PATCH] feat(ir): enable compatible summary state merges --- .../src/summary_maintenance_cost/model.rs | 8 +- .../tests/planner_layering_merge.rs | 145 ++++++++++++++++++ crates/types/src/ir/asap.rs | 63 ++++++-- crates/types/src/ir/timing.rs | 17 +- .../src/post_asap/execution_data_state.rs | 4 +- crates/types/tests/summary_merge.rs | 72 +++++++++ .../operator-design-acceptance.md | 12 ++ 7 files changed, 299 insertions(+), 22 deletions(-) create mode 100644 crates/integration-tests/tests/planner_layering_merge.rs create mode 100644 crates/types/tests/summary_merge.rs 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 03d2bd329..c3c98fcda 100644 --- a/crates/asap-aware-mapping/src/summary_maintenance_cost/model.rs +++ b/crates/asap-aware-mapping/src/summary_maintenance_cost/model.rs @@ -2635,9 +2635,9 @@ mod tests { } #[test] - fn state_only_needs_no_evaluation_and_summary_merge_child_fails_closed() { + fn state_only_needs_no_evaluation_and_summary_subtract_child_fails_closed() { // A state-only root is costable without evaluation evidence; a - // SummaryAgg over a reserved SummaryMerge is rejected at planning. + // SummaryAgg over a reserved SummarySubtract is rejected at planning. let workload = streaming_workload(); let target = streaming_sum_query(); let estimated = summary_with_operations(false, false, false); @@ -2666,7 +2666,7 @@ mod tests { let nested = std::rc::Rc::new( OperatorNode::with_schema( asap_types::ir::Operator::ASAP(ASAPOp::SummaryAgg { - child: evaluation_state(&summary_with_operations(true, false, false)), + child: evaluation_state(&summary_with_operations(false, true, false)), family: FieldDataType::ExactAggregate(ExactKind::Count, ExactParams::Count), input: SummaryUpdate { item: None, @@ -2694,7 +2694,7 @@ mod tests { assert!(matches!( streaming_planning_error(nested, &nested_model), asap_types::post_asap::ExecutionDataStateError::UnimplementedOperator { - operator: "SummaryMerge" + operator: "SummarySubtract" } )); } diff --git a/crates/integration-tests/tests/planner_layering_merge.rs b/crates/integration-tests/tests/planner_layering_merge.rs new file mode 100644 index 000000000..9512582db --- /dev/null +++ b/crates/integration-tests/tests/planner_layering_merge.rs @@ -0,0 +1,145 @@ +//! The window-composition example merges KLL panes through the unified IR and native runtime. +mod physical_common; +use asap_physical_operators::{ + physical_planner::{compile, cut_candidate, InputContract}, + runtime::Scope, + values::{Batch, Value}, +}; +use asap_types::{ + ir::operator_properties::{Reduction, Source}, + ir::{ASAPOp, NonASAPOp, Operator, OperatorNode}, + post_asap::{ + GroupingStrategy, SketchAlgorithm, SketchKind, SketchParams, SketchStatistic, SummaryUpdate, + }, + pre_asap::{ColumnRef, DataType, Field, FieldDataType, Schema}, +}; +use std::{collections::BTreeMap, sync::Arc}; + +/// Build five one-minute KLLs, merge them, then read p99 through the real wire compiler. +#[test] +fn five_minute_quantile_merges_five_one_minute_states() { + let family = FieldDataType::Sketch( + SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k: 200 }), + Default::default(), + ); + let states = (0..5) + .map(|pane| { + let scan = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Scan { + source: Source::Table { + table_ref: format!("pane_{pane}"), + }, + predicates: vec![], + schema: Schema::new(vec![Field::plain("value", DataType::Float64, false)]), + })) + .unwrap(); + OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryAgg { + child: scan, + family: family.clone(), + input: SummaryUpdate::column(ColumnRef::SampleValue), + reduction: Reduction::by(vec![]), + grouping: GroupingStrategy::default(), + filter: None, + })) + .unwrap() + }) + .collect(); + let merged = + OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryMerge { children: states })) + .unwrap(); + let root = OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryEstimate { + summary_input: merged, + query: SketchStatistic::Quantile { q: 0.99 }, + })) + .unwrap(); + root.validate_structure().unwrap(); + let wire = physical_common::compile_post_asap_dag(&root).unwrap(); + wire.validate().unwrap(); + let mut inputs = BTreeMap::new(); + let mut batches = BTreeMap::new(); + for node in &wire.nodes { + if let asap_types::ir::export::PostAsapOperatorPayload::Relational { + operator: + asap_types::ir::export::NonASAPOpKind::Scan { + source: Source::Table { table_ref }, + .. + }, + } = &node.payload + { + let pane: usize = table_ref.strip_prefix("pane_").unwrap().parse().unwrap(); + let schema = Arc::new(node.output_schema.clone()); + let id = u64::from(node.id.0); + inputs.insert(id, InputContract::bounded(schema.clone())); + batches.insert( + id, + Batch::try_new( + schema, + (0..20) + .map(|i| vec![Value::Float64((pane * 20 + i) as f64)]) + .collect(), + ) + .unwrap(), + ); + } + } + let plan = compile(&wire, inputs, &[u64::from(wire.root.0)]).unwrap(); + let outputs = physical_common::execute( + &plan, + batches.clone(), + Scope::Query { + evaluation_time_ms: 300_000, + revision: 1, + }, + ); + // Pattern B1 stores five built panes; its query reads the same merge DAG + // as Pattern B2, which builds every pane from raw rows on each read. + let frontier: Vec<_> = wire + .nodes + .iter() + .filter_map(|node| { + matches!( + node.payload, + asap_types::ir::export::PostAsapOperatorPayload::SummaryAgg { .. } + ) + .then_some(u64::from(node.id.0)) + }) + .collect(); + let materialized = cut_candidate(&plan, &frontier).unwrap(); + let precompute = materialized.precompute.as_ref().unwrap(); + let stored = physical_common::execute( + precompute, + batches, + Scope::Ingestion { + window_start_ms: 0, + window_end_ms: 300_000, + revision: 1, + }, + ); + let retained = precompute + .roots() + .iter() + .copied() + .zip(stored) + .map(|(id, mut batches)| { + assert_eq!(batches.len(), 1); + (id, batches.remove(0)) + }) + .collect(); + let retained_outputs = physical_common::execute( + &materialized.query, + retained, + Scope::Query { + evaluation_time_ms: 300_000, + revision: 1, + }, + ); + assert_eq!(retained_outputs[0][0].rows().len(), 1); + assert!( + matches!(retained_outputs[0][0].rows()[0].as_slice(), [Value::Float64(value)] if *value == 98.0) + ); + assert_eq!(outputs[0][0].rows().len(), 1); + assert!( + matches!(outputs[0][0].rows()[0].as_slice(), [Value::Float64(value)] if *value == 98.0), + "{:?}", + outputs[0][0].rows() + ); +} diff --git a/crates/types/src/ir/asap.rs b/crates/types/src/ir/asap.rs index c6d9867b3..7cf4c48d8 100644 --- a/crates/types/src/ir/asap.rs +++ b/crates/types/src/ir/asap.rs @@ -39,9 +39,7 @@ pub enum ASAPOp { }, /// Read an exact accumulator's state as its finalized value: the /// maintenance-to-read boundary before query-time operators. - FinalizeExactAccumulator { - child: Rc, - }, + FinalizeExactAccumulator { child: Rc }, /// Maintain the full declared population, including membership changes. MaintainPopulation { child: Rc, @@ -52,10 +50,9 @@ pub enum ASAPOp { child: Rc, evaluation: PopulationStatistic, }, - // ── Reserved: migrated but unimplemented (§1.3 of the proposal) ── - SummaryMerge { - children: Vec>, - }, + /// Merge compatible partial states for the same grouping and family. + SummaryMerge { children: Vec> }, + // ── Reserved: migrated but unimplemented ── SummarySubtract { left: Rc, right: Rc, @@ -186,11 +183,7 @@ impl ASAPOp { use ASAPOp::*; matches!( self, - SummaryMerge { .. } - | SummarySubtract { .. } - | SummaryDelete { .. } - | SummaryJoin { .. } - | Extension { .. } + SummarySubtract { .. } | SummaryDelete { .. } | SummaryJoin { .. } | Extension { .. } ) } @@ -202,6 +195,14 @@ impl ASAPOp { pub fn produced_state(&self) -> Option<&FieldDataType> { match self { ASAPOp::SummaryAgg { family, .. } | ASAPOp::SummaryJoin { family, .. } => Some(family), + ASAPOp::SummaryMerge { children } => children.first().and_then(|child| { + child + .schema + .fields + .iter() + .find(|field| !field.is_plain()) + .map(|field| &field.dtype) + }), _ => None, } } @@ -405,8 +406,11 @@ impl ASAPOp { .output_schema()? } } - SummaryMerge { .. } - | SummarySubtract { .. } + SummaryMerge { children } => { + self.validate_inputs()?; + children[0].schema.clone() + } + SummarySubtract { .. } | SummaryDelete { .. } | SummaryJoin { .. } | Extension { .. } => return Err(Self::unimplemented()), @@ -444,6 +448,37 @@ impl ASAPOp { } }; match self { + SummaryMerge { children } => { + let Some(first) = children.first() else { + return Err(SchemaDerivationError::InvalidScalarSignature( + "summary merge requires at least one state input".into(), + )); + }; + // Matching state parameters and grouping positions are necessary; + // matching names alone cannot prove two states compatible. + if first + .schema + .fields + .iter() + .filter(|field| !field.is_plain()) + .count() + != 1 + { + return Err(SchemaDerivationError::InvalidScalarSignature( + "summary merge requires exactly one state column".into(), + )); + } + for child in children { + needs_state(child, "SummaryMerge")?; + if child.schema != first.schema { + return Err(SchemaDerivationError::InvalidScalarSignature( + "summary merge inputs must have identical state and grouping schemas" + .into(), + )); + } + } + Ok(()) + } SummaryEstimate { summary_input, query, diff --git a/crates/types/src/ir/timing.rs b/crates/types/src/ir/timing.rs index eab02c99a..cb0bffc34 100644 --- a/crates/types/src/ir/timing.rs +++ b/crates/types/src/ir/timing.rs @@ -396,8 +396,21 @@ fn validate_asap( } Ok(()) } - ASAPOp::SummaryMerge { .. } - | ASAPOp::SummarySubtract { .. } + ASAPOp::SummaryMerge { children } => { + for child in children { + let state = state_of(child); + if state.primitive != DataPrimitive::SummaryState + || (timing == ExecutionTiming::IngestionTime && state.timing != timing) + { + return Err(ExecutionDataStateError::IllegalChildDataState { + edge: "SummaryMerge.children", + child: state, + }); + } + } + Ok(()) + } + ASAPOp::SummarySubtract { .. } | ASAPOp::SummaryDelete { .. } | ASAPOp::SummaryJoin { .. } | ASAPOp::Extension { .. } => Err(ExecutionDataStateError::UnimplementedOperator { diff --git a/crates/types/src/post_asap/execution_data_state.rs b/crates/types/src/post_asap/execution_data_state.rs index de764ea2c..53b9857d9 100644 --- a/crates/types/src/post_asap/execution_data_state.rs +++ b/crates/types/src/post_asap/execution_data_state.rs @@ -134,8 +134,8 @@ pub enum ExecutionDataStateError { /// declared data_state. #[error("exact operator consumes non-plain column {column:?} ({dtype})")] NonPlainOperand { column: String, dtype: String }, - /// A reserved ASAP operator (`SummaryMerge`, `SummarySubtract`, - /// `SummaryDelete`, `SummaryJoin`, `Extension`) in an executable plan. + /// A reserved ASAP operator (`SummarySubtract`, `SummaryDelete`, + /// `SummaryJoin`, `Extension`) in an executable plan. #[error("{operator} is a reserved operator with no execution contract yet")] UnimplementedOperator { operator: &'static str }, /// A node reached by export without a timing: the lifecycle timing pass diff --git a/crates/types/tests/summary_merge.rs b/crates/types/tests/summary_merge.rs new file mode 100644 index 000000000..f4d438d1a --- /dev/null +++ b/crates/types/tests/summary_merge.rs @@ -0,0 +1,72 @@ +//! Window composition merges compatible summary states without consuming raw rows. +use asap_types::{ + ir::operator_properties::{Reduction, Source}, + ir::{ASAPOp, NonASAPOp, Operator, OperatorNode}, + post_asap::{ + ExecutionTiming, GroupingStrategy, SketchAlgorithm, SketchKind, SketchParams, SummaryUpdate, + }, + pre_asap::{ColumnRef, DataType, Field, FieldDataType, Schema}, +}; +use std::rc::Rc; +fn state(k: u32) -> Rc { + let scan = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Scan { + source: Source::Table { + table_ref: "latencies".into(), + }, + predicates: vec![], + schema: Schema::new(vec![Field::plain("value", DataType::Float64, false)]), + })) + .unwrap(); + OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryAgg { + child: scan, + family: FieldDataType::Sketch( + SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k }), + Default::default(), + ), + input: SummaryUpdate::column(ColumnRef::SampleValue), + reduction: Reduction::by(vec![]), + grouping: GroupingStrategy::default(), + filter: None, + })) + .unwrap() +} +/// Two KLL panes compose into one typed state at either execution phase. +#[test] +fn compatible_panes_merge_and_export() { + let root = OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryMerge { + children: vec![state(200), state(200)], + })) + .unwrap(); + root.validate_structure().unwrap(); + assert_eq!(root.schema.fields.len(), 1); + let timed = asap_types::ir::apply_lifecycle_timings( + &root, + &Default::default(), + &mut Default::default(), + ) + .unwrap(); + asap_types::ir::export::compile_post_asap_dag(&timed) + .unwrap() + .validate() + .unwrap(); + for timing in [ExecutionTiming::IngestionTime, ExecutionTiming::QueryTime] { + asap_types::ir::timing::validate_default(&root, timing).unwrap(); + assert_eq!( + asap_types::ir::timing::planned_data_state(&root, timing).primitive, + asap_types::post_asap::DataPrimitive::SummaryState + ); + } +} +/// An empty merge, raw rows and differently sized state cannot masquerade as compatible panes. +#[test] +fn incompatible_merge_inputs_fail() { + for children in [ + vec![], + vec![state(200), state(300)], + vec![state(200).children()[0].clone()], + ] { + assert!( + OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryMerge { children })).is_err() + ); + } +} diff --git a/docs/develop_docs/operator-design-acceptance.md b/docs/develop_docs/operator-design-acceptance.md index 390d49733..daee31c74 100644 --- a/docs/develop_docs/operator-design-acceptance.md +++ b/docs/develop_docs/operator-design-acceptance.md @@ -119,3 +119,15 @@ skipped because Node.js is unavailable in this environment). The vendored Metric verifying its existing 21 library and 3 doctest failures. Rust 1.98 changes one compiler-diagnostic fingerprint; no baseline hashes or vendored sources were changed to accommodate that older toolchain. Formatting and clippy pass on 1.99. + +## Planner-layering follow-up: mergeable state + +`ASAPOp::SummaryMerge` now derives and checks its shared schema, preserves the +state family and parameters, and validates dependencies at either execution +phase. Empty merges, raw inputs and differing state/grouping schemas fail. +`types/tests/summary_merge.rs` checks these contracts and executable export. +`integration-tests/tests/planner_layering_merge.rs` builds five one-minute KLL +states, merges them through unified export/native compilation, and checks p99 +for both rebuilding raw inputs and reading materialized panes (Examples 3B/4B). +This proves the merge building block; automatic window candidate generation and +Exponential Histogram construction remain separate work.