From 87192d75240de8003f47ea8ff52584aaf609b818 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Sat, 3 Oct 2026 14:35:57 +0000 Subject: [PATCH 01/13] feat(ir): define compatible logical summary merges --- crates/types/src/ir/asap.rs | 63 ++++++++++++++----- crates/types/tests/summary_merge_structure.rs | 53 ++++++++++++++++ 2 files changed, 102 insertions(+), 14 deletions(-) create mode 100644 crates/types/tests/summary_merge_structure.rs diff --git a/crates/types/src/ir/asap.rs b/crates/types/src/ir/asap.rs index 90b239463..ecd7262f2 100644 --- a/crates/types/src/ir/asap.rs +++ b/crates/types/src/ir/asap.rs @@ -41,9 +41,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: C, - }, + FinalizeExactAccumulator { child: C }, /// Maintain the full declared population, including membership changes. MaintainPopulation { child: C, @@ -54,10 +52,9 @@ pub enum ASAPOp> { child: C, 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: C, right: C, @@ -192,11 +189,7 @@ impl ASAPOp { use ASAPOp::*; matches!( self, - SummaryMerge { .. } - | SummarySubtract { .. } - | SummaryDelete { .. } - | SummaryJoin { .. } - | Extension { .. } + SummarySubtract { .. } | SummaryDelete { .. } | SummaryJoin { .. } | Extension { .. } ) } @@ -208,6 +201,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, } } @@ -411,8 +412,11 @@ impl ASAPOp { .output_schema()? } } - SummaryMerge { .. } - | SummarySubtract { .. } + SummaryMerge { children } => { + self.validate_inputs()?; + children[0].schema.clone() + } + SummarySubtract { .. } | SummaryDelete { .. } | SummaryJoin { .. } | Extension { .. } => return Err(Self::unimplemented()), @@ -450,6 +454,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/tests/summary_merge_structure.rs b/crates/types/tests/summary_merge_structure.rs new file mode 100644 index 000000000..999ee95f7 --- /dev/null +++ b/crates/types/tests/summary_merge_structure.rs @@ -0,0 +1,53 @@ +//! 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::{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 logical state without timing assignment. +#[test] +fn compatible_panes_merge_structurally() { + 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); +} +/// 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() + ); + } +} From 4fa2db9d1bb4f015d4aebb83095abdbd4d46e88c Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Sat, 3 Oct 2026 16:03:55 +0000 Subject: [PATCH 02/13] feat(ir): require disjoint coverage for summary merges --- crates/types/src/ir/asap.rs | 24 +++++++ crates/types/src/ir/node.rs | 13 +++- crates/types/tests/summary_merge_structure.rs | 62 ++++++++++++++++++- 3 files changed, 95 insertions(+), 4 deletions(-) diff --git a/crates/types/src/ir/asap.rs b/crates/types/src/ir/asap.rs index ecd7262f2..8bce8f369 100644 --- a/crates/types/src/ir/asap.rs +++ b/crates/types/src/ir/asap.rs @@ -213,6 +213,29 @@ impl ASAPOp { } } + /// Derive joint observation coverage; unknown or overlapping inputs fail closed. + pub fn merged_extent( + &self, + ) -> Result { + let ASAPOp::SummaryMerge { children } = self else { + return Err(SchemaDerivationError::InvalidScalarSignature( + "coverage merge requires SummaryMerge".into(), + )); + }; + let inputs = children + .iter() + .map(|child| { + child.observation_extent.clone().ok_or_else(|| { + SchemaDerivationError::InvalidScalarSignature( + "summary merge requires known observation coverage".into(), + ) + }) + }) + .collect::, _>>()?; + super::observation_extent::ObservationExtent::merge_disjoint(&inputs) + .map_err(|error| SchemaDerivationError::InvalidScalarSignature(error.to_string())) + } + /// Output schema derived from the operator and its children. Summary /// planning may retain a more specific schema (evaluation column naming) /// through [`OperatorNode::with_schema`]; all structural metadata must @@ -483,6 +506,7 @@ impl ASAPOp { )); } } + self.merged_extent()?; Ok(()) } SummaryEstimate { diff --git a/crates/types/src/ir/node.rs b/crates/types/src/ir/node.rs index 47ffa0ea3..42ab57604 100644 --- a/crates/types/src/ir/node.rs +++ b/crates/types/src/ir/node.rs @@ -110,7 +110,11 @@ impl OperatorNode { /// ASAP operator, ...). pub fn new(operator: Operator) -> Result { let schema = operator.output_schema()?; - Ok(Self::with_schema(operator, schema)) + let mut node = Self::with_schema(operator, schema); + if let Some(op @ ASAPOp::SummaryMerge { .. }) = node.asap() { + node.observation_extent = Some(op.merged_extent()?); + } + Ok(node) } /// Build a node with caller-supplied output names and qualifiers. For @@ -322,6 +326,13 @@ impl OperatorNode { None if node.requires_coverage() => return Err(CoverageError::Missing.into()), None => {} } + if let Some(op @ ASAPOp::SummaryMerge { .. }) = node.asap() { + if node.observation_extent.as_ref() != Some(&op.merged_extent()?) { + return Err(SchemaDerivationError::InvalidScalarSignature( + "retained merge coverage disagrees with input union".into(), + )); + } + } node.operator.validate_inputs()?; if node.result_kind != node.operator.output_kind() { return Err(SchemaDerivationError::InvalidScalarSignature( diff --git a/crates/types/tests/summary_merge_structure.rs b/crates/types/tests/summary_merge_structure.rs index 999ee95f7..5c8b0c132 100644 --- a/crates/types/tests/summary_merge_structure.rs +++ b/crates/types/tests/summary_merge_structure.rs @@ -15,7 +15,7 @@ fn state(k: u32) -> Rc { schema: Schema::new(vec![Field::plain("value", DataType::Float64, false)]), })) .unwrap(); - OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryAgg { + let summary = OperatorNode::new(Operator::ASAP(ASAPOp::SummaryAgg { child: scan, family: FieldDataType::Sketch( SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k }), @@ -26,17 +26,38 @@ fn state(k: u32) -> Rc { grouping: GroupingStrategy::default(), filter: None, })) - .unwrap() + .unwrap(); + std::rc::Rc::new( + summary + .with_observation_extent(asap_types::ir::observation_extent::ObservationExtent { + source: "latencies:timestamp".into(), + revision: "fixture-1".into(), + input: SummaryUpdate::column(ColumnRef::SampleValue), + grouping: Reduction::by(vec![]), + multiplicity: + asap_types::ir::observation_extent::ObservationMultiplicity::OncePerObservation, + regions: vec![asap_types::ir::observation_extent::ExtentRegion { + start_ms: 0, + end_ms: 1, + population: Default::default(), + }], + }) + .unwrap(), + ) } /// Two KLL panes compose into one typed logical state without timing assignment. #[test] fn compatible_panes_merge_structurally() { let root = OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryMerge { - children: vec![state(200), state(200)], + children: vec![state(200), shifted_state(200, 1, 2)], })) .unwrap(); root.validate_structure().unwrap(); assert_eq!(root.schema.fields.len(), 1); + assert_eq!( + root.observation_extent.as_ref().unwrap().regions[0].end_ms, + 2 + ); } /// An empty merge, raw rows and differently sized state cannot masquerade as compatible panes. #[test] @@ -51,3 +72,38 @@ fn incompatible_merge_inputs_fail() { ); } } + +fn shifted_state(k: u32, start: i64, end: i64) -> Rc { + let mut node = (*state(k)).clone(); + let region = &mut node.observation_extent.as_mut().unwrap().regions[0]; + region.start_ms = start; + region.end_ms = end; + Rc::new(node) +} +/// Schema equality cannot authorize overlapping or unknown observation coverage. +#[test] +fn unsafe_coverage_merge_is_rejected() { + let mut unknown = (*state(200)).clone(); + unknown.observation_extent = None; + for children in [ + vec![state(200), state(200)], + vec![state(200), Rc::new(unknown)], + ] { + assert!( + OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryMerge { children })).is_err() + ); + } +} + +/// Gapped time coverage remains disconnected, and forged output metadata is rejected. +#[test] +fn merge_derives_coverage_and_validates_retained_metadata() { + let root = OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryMerge { + children: vec![state(200), shifted_state(200, 2, 3)], + })) + .unwrap(); + assert_eq!(root.observation_extent.as_ref().unwrap().regions.len(), 2); + let mut forged = (*root).clone(); + forged.observation_extent.as_mut().unwrap().regions[0].end_ms = 2; + assert!(Rc::new(forged).validate_structure().is_err()); +} From 4596a7a2535e9e80bcd07beed9303ecc275343f2 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Sat, 3 Oct 2026 16:57:54 +0000 Subject: [PATCH 03/13] refactor(ir): use SummaryCoverage names in summary merges Co-Authored-By: Claude Opus 5.5 --- crates/types/src/ir/asap.rs | 20 +++++--------- crates/types/src/ir/node.rs | 13 ++++----- crates/types/src/ir/summary_coverage.rs | 4 +++ crates/types/tests/summary_merge_structure.rs | 27 ++++++++----------- 4 files changed, 28 insertions(+), 36 deletions(-) diff --git a/crates/types/src/ir/asap.rs b/crates/types/src/ir/asap.rs index 8bce8f369..658fb4ed4 100644 --- a/crates/types/src/ir/asap.rs +++ b/crates/types/src/ir/asap.rs @@ -6,6 +6,7 @@ use std::rc::Rc; use serde::{Deserialize, Serialize}; use super::node::{OperatorNode, OperatorResultKind}; +use super::summary_coverage::{CoverageError, SummaryCoverage}; use crate::ir::operator_properties::Reduction; use crate::ir::SchemaDerivationError; use crate::post_asap::maintained_population::{MaintainedPopulation, PopulationStatistic}; @@ -213,10 +214,8 @@ impl ASAPOp { } } - /// Derive joint observation coverage; unknown or overlapping inputs fail closed. - pub fn merged_extent( - &self, - ) -> Result { + /// Derive the merged node's coverage; unknown or overlapping inputs fail closed. + pub fn merged_coverage(&self) -> Result { let ASAPOp::SummaryMerge { children } = self else { return Err(SchemaDerivationError::InvalidScalarSignature( "coverage merge requires SummaryMerge".into(), @@ -224,16 +223,9 @@ impl ASAPOp { }; let inputs = children .iter() - .map(|child| { - child.observation_extent.clone().ok_or_else(|| { - SchemaDerivationError::InvalidScalarSignature( - "summary merge requires known observation coverage".into(), - ) - }) - }) + .map(|child| child.coverage.clone().ok_or(CoverageError::UnknownInput)) .collect::, _>>()?; - super::observation_extent::ObservationExtent::merge_disjoint(&inputs) - .map_err(|error| SchemaDerivationError::InvalidScalarSignature(error.to_string())) + Ok(SummaryCoverage::merge_disjoint(&inputs)?) } /// Output schema derived from the operator and its children. Summary @@ -506,7 +498,7 @@ impl ASAPOp { )); } } - self.merged_extent()?; + self.merged_coverage()?; Ok(()) } SummaryEstimate { diff --git a/crates/types/src/ir/node.rs b/crates/types/src/ir/node.rs index 42ab57604..c33a87f91 100644 --- a/crates/types/src/ir/node.rs +++ b/crates/types/src/ir/node.rs @@ -112,7 +112,7 @@ impl OperatorNode { let schema = operator.output_schema()?; let mut node = Self::with_schema(operator, schema); if let Some(op @ ASAPOp::SummaryMerge { .. }) = node.asap() { - node.observation_extent = Some(op.merged_extent()?); + node.coverage = Some(op.merged_coverage()?); } Ok(node) } @@ -165,7 +165,10 @@ impl OperatorNode { /// Summary nodes whose state can be composed must declare coverage. pub fn requires_coverage(&self) -> bool { - matches!(self.asap(), Some(ASAPOp::SummaryAgg { .. })) + matches!( + self.asap(), + Some(ASAPOp::SummaryAgg { .. } | ASAPOp::SummaryMerge { .. }) + ) } pub fn non_asap(&self) -> Option<&NonASAPOp> { @@ -327,10 +330,8 @@ impl OperatorNode { None => {} } if let Some(op @ ASAPOp::SummaryMerge { .. }) = node.asap() { - if node.observation_extent.as_ref() != Some(&op.merged_extent()?) { - return Err(SchemaDerivationError::InvalidScalarSignature( - "retained merge coverage disagrees with input union".into(), - )); + if node.coverage.as_ref() != Some(&op.merged_coverage()?) { + return Err(CoverageError::MergeOutputMismatch.into()); } } node.operator.validate_inputs()?; diff --git a/crates/types/src/ir/summary_coverage.rs b/crates/types/src/ir/summary_coverage.rs index 3f3798a88..120fb1212 100644 --- a/crates/types/src/ir/summary_coverage.rs +++ b/crates/types/src/ir/summary_coverage.rs @@ -48,6 +48,10 @@ pub enum CoverageError { NotState, #[error("summary node requires coverage")] Missing, + #[error("summary merge requires known coverage on every input")] + UnknownInput, + #[error("retained merge coverage disagrees with input union")] + MergeOutputMismatch, } impl SummaryCoverage { diff --git a/crates/types/tests/summary_merge_structure.rs b/crates/types/tests/summary_merge_structure.rs index 5c8b0c132..6489159af 100644 --- a/crates/types/tests/summary_merge_structure.rs +++ b/crates/types/tests/summary_merge_structure.rs @@ -29,16 +29,12 @@ fn state(k: u32) -> Rc { .unwrap(); std::rc::Rc::new( summary - .with_observation_extent(asap_types::ir::observation_extent::ObservationExtent { + .with_coverage(asap_types::ir::summary_coverage::SummaryCoverage { source: "latencies:timestamp".into(), - revision: "fixture-1".into(), input: SummaryUpdate::column(ColumnRef::SampleValue), - grouping: Reduction::by(vec![]), - multiplicity: - asap_types::ir::observation_extent::ObservationMultiplicity::OncePerObservation, - regions: vec![asap_types::ir::observation_extent::ExtentRegion { - start_ms: 0, - end_ms: 1, + reduction: Reduction::by(vec![]), + regions: vec![asap_types::ir::summary_coverage::CoverageRegion { + time_ms: Some(0..1), population: Default::default(), }], }) @@ -55,8 +51,8 @@ fn compatible_panes_merge_structurally() { root.validate_structure().unwrap(); assert_eq!(root.schema.fields.len(), 1); assert_eq!( - root.observation_extent.as_ref().unwrap().regions[0].end_ms, - 2 + root.coverage.as_ref().unwrap().regions[0].time_ms, + Some(0..2) ); } /// An empty merge, raw rows and differently sized state cannot masquerade as compatible panes. @@ -75,16 +71,15 @@ fn incompatible_merge_inputs_fail() { fn shifted_state(k: u32, start: i64, end: i64) -> Rc { let mut node = (*state(k)).clone(); - let region = &mut node.observation_extent.as_mut().unwrap().regions[0]; - region.start_ms = start; - region.end_ms = end; + let region = &mut node.coverage.as_mut().unwrap().regions[0]; + region.time_ms = Some(start..end); Rc::new(node) } /// Schema equality cannot authorize overlapping or unknown observation coverage. #[test] fn unsafe_coverage_merge_is_rejected() { let mut unknown = (*state(200)).clone(); - unknown.observation_extent = None; + unknown.coverage = None; for children in [ vec![state(200), state(200)], vec![state(200), Rc::new(unknown)], @@ -102,8 +97,8 @@ fn merge_derives_coverage_and_validates_retained_metadata() { children: vec![state(200), shifted_state(200, 2, 3)], })) .unwrap(); - assert_eq!(root.observation_extent.as_ref().unwrap().regions.len(), 2); + assert_eq!(root.coverage.as_ref().unwrap().regions.len(), 2); let mut forged = (*root).clone(); - forged.observation_extent.as_mut().unwrap().regions[0].end_ms = 2; + forged.coverage.as_mut().unwrap().regions[0].time_ms = Some(0..2); assert!(Rc::new(forged).validate_structure().is_err()); } From fe35b52bae3841f6047f43e3df027b01bac88821 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Sat, 3 Oct 2026 17:33:53 +0000 Subject: [PATCH 04/13] feat(ir): check merge update/reduction on producers; test documented coverage examples SummaryCoverage no longer repeats input/reduction, so SummaryMerge compares them through OperatorNode::summary_update. summary_coverage_examples.rs builds each example in docs/develop_docs/summary-coverage.md as a SummaryAgg -> SummaryMerge plan. Co-Authored-By: Claude Opus 5.5 --- crates/types/src/ir/asap.rs | 8 + crates/types/src/ir/node.rs | 14 + .../types/tests/summary_coverage_examples.rs | 264 ++++++++++++++++++ crates/types/tests/summary_merge_structure.rs | 6 +- 4 files changed, 289 insertions(+), 3 deletions(-) create mode 100644 crates/types/tests/summary_coverage_examples.rs diff --git a/crates/types/src/ir/asap.rs b/crates/types/src/ir/asap.rs index 658fb4ed4..80a5fd9b8 100644 --- a/crates/types/src/ir/asap.rs +++ b/crates/types/src/ir/asap.rs @@ -498,6 +498,14 @@ impl ASAPOp { )); } } + // Coverage records only time and population; what each state + // summarizes and how it is grouped come from the producers. + let update = first.summary_update(); + if update.is_none() || children.iter().any(|c| c.summary_update() != update) { + return Err(SchemaDerivationError::InvalidScalarSignature( + "summary merge inputs must share update expression and reduction".into(), + )); + } self.merged_coverage()?; Ok(()) } diff --git a/crates/types/src/ir/node.rs b/crates/types/src/ir/node.rs index c33a87f91..3cdbe8ecf 100644 --- a/crates/types/src/ir/node.rs +++ b/crates/types/src/ir/node.rs @@ -10,10 +10,12 @@ use serde::{Deserialize, Serialize}; use super::asap::ASAPOp; use super::non_asap::NonASAPOp; +use super::operator_properties::Reduction; use super::summary_coverage::{CoverageError, SummaryCoverage}; use crate::ir::SchemaDerivationError; use crate::post_asap::execution_data_state::ExecutionTiming; use crate::post_asap::guarantee::ResultGuarantee; +use crate::post_asap::SummaryUpdate; use crate::pre_asap::schema::Schema; /// The output category of an operator, derived from the operation and its @@ -163,6 +165,18 @@ impl OperatorNode { Ok(self) } + /// What a summary state is updated with and how it is grouped: the + /// `SummaryAgg` fields, or those shared by a `SummaryMerge`'s inputs. + pub fn summary_update(&self) -> Option<(&SummaryUpdate, &Reduction)> { + match self.asap()? { + ASAPOp::SummaryAgg { + input, reduction, .. + } => Some((input, reduction)), + ASAPOp::SummaryMerge { children } => children.first()?.summary_update(), + _ => None, + } + } + /// Summary nodes whose state can be composed must declare coverage. pub fn requires_coverage(&self) -> bool { matches!( diff --git a/crates/types/tests/summary_coverage_examples.rs b/crates/types/tests/summary_coverage_examples.rs new file mode 100644 index 000000000..e2fad07ab --- /dev/null +++ b/crates/types/tests/summary_coverage_examples.rs @@ -0,0 +1,264 @@ +//! The examples in docs/develop_docs/summary-coverage.md, built as real +//! SummaryAgg -> SummaryMerge plans. Every input has the same schema +//! `(job: Utf8, state: KLL{k=200})`; only coverage differs. +use asap_types::ir::summary_coverage::{CoverageError, CoverageRegion, SummaryCoverage}; +use asap_types::ir::{ + ASAPOp, ExprSemantics, NonASAPOp, Operator, OperatorNode, Predicate, ScalarExpr, + SchemaDerivationError, +}; +use asap_types::post_asap::{ + GroupingStrategy, SketchAlgorithm, SketchKind, SketchParams, SummaryUpdate, +}; +use asap_types::pre_asap::{ + ColumnRef, CompareOpKind, DataType, Field, FieldDataType, Reduction, ScalarValue, Schema, + Source, +}; +use std::rc::Rc; + +const MIN: i64 = 60_000; + +fn requests() -> Source { + Source::Table { + table_ref: "requests".into(), + } +} + +fn scan() -> Rc { + OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Scan { + source: requests(), + predicates: vec![], + schema: Schema::new(vec![ + Field::plain("job", DataType::Utf8, false), + Field::plain("region", DataType::Utf8, false), + Field::plain("tier", DataType::Utf8, false), + Field::plain("latency", DataType::Float64, false), + Field::plain("size", DataType::Float64, false), + ]), + })) + .unwrap() +} + +/// `region = value`, evaluated on the scan schema. +fn region_is(value: &str) -> Predicate { + Predicate(ScalarExpr::Compare { + left: Box::new(ScalarExpr::Column(1)), + op: CompareOpKind::Eq, + right: Box::new(ScalarExpr::Literal(ScalarValue::Utf8(value.into()))), + semantics: ExprSemantics::Sql, + }) +} + +fn region(time_ms: Option>, population: &[(&str, &str)]) -> CoverageRegion { + CoverageRegion { + time_ms, + population: population + .iter() + .map(|(k, v)| (k.to_string(), v.to_string())) + .collect(), + } +} + +/// p99-ready KLL over `column`, grouped by job, with declared coverage. +fn kll_over( + column: &str, + filter: Option, + coverage: SummaryCoverage, +) -> Rc { + let node = OperatorNode::new(Operator::ASAP(ASAPOp::SummaryAgg { + child: scan(), + family: FieldDataType::Sketch( + SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k: 200 }), + GroupingStrategy::default(), + ), + input: SummaryUpdate::column(ColumnRef::Named(column.into())), + reduction: Reduction::by(vec![0]), + grouping: GroupingStrategy::default(), + filter, + })) + .unwrap(); + Rc::new(node.with_coverage(coverage).unwrap()) +} + +fn kll(time_ms: Option>, population: &[(&str, &str)]) -> Rc { + kll_over( + "latency", + None, + SummaryCoverage { + source: requests(), + regions: vec![region(time_ms, population)], + }, + ) +} + +fn merge(children: Vec>) -> Result, SchemaDerivationError> { + let schema = children[0].schema.clone(); + let merged = OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryMerge { children }))?; + merged.validate_structure()?; + // Schema never changes; only coverage does. + assert_eq!(merged.schema, schema); + Ok(merged) +} + +fn regions(node: &OperatorNode) -> Vec { + node.coverage.as_ref().unwrap().regions.clone() +} + +fn rejected(result: Result, SchemaDerivationError>, expected: CoverageError) { + match result { + Err(SchemaDerivationError::Coverage(actual)) => assert_eq!(actual, expected), + other => panic!("expected {expected:?}, got {other:?}"), + } +} + +/// Example 1: adjacent panes coalesce, gaps stay, overlapping windows are rejected. +#[test] +fn example_1_time() { + let adjacent = merge(vec![kll(Some(0..MIN), &[]), kll(Some(MIN..2 * MIN), &[])]).unwrap(); + assert_eq!(regions(&adjacent), vec![region(Some(0..2 * MIN), &[])]); + + let gapped = merge(vec![ + kll(Some(0..MIN), &[]), + kll(Some(2 * MIN..3 * MIN), &[]), + ]) + .unwrap(); + assert_eq!( + regions(&gapped), + vec![ + region(Some(0..MIN), &[]), + region(Some(2 * MIN..3 * MIN), &[]) + ] + ); + + rejected( + merge(vec![ + kll(Some(0..2 * MIN), &[]), + kll(Some(MIN..3 * MIN), &[]), + ]), + CoverageError::PossibleOverlap, + ); +} + +/// Example 2: disjoint label values merge; different labels or equal values are rejected. +#[test] +fn example_2_population() { + let t = Some(0..MIN); + let us_eu = merge(vec![ + kll(t.clone(), &[("region", "us")]), + kll(t.clone(), &[("region", "eu")]), + ]) + .unwrap(); + assert_eq!( + regions(&us_eu), + vec![ + region(t.clone(), &[("region", "eu")]), + region(t.clone(), &[("region", "us")]), + ] + ); + for other in [("tier", "premium"), ("region", "us")] { + rejected( + merge(vec![ + kll(t.clone(), &[("region", "us")]), + kll(t.clone(), &[other]), + ]), + CoverageError::PossibleOverlap, + ); + } +} + +/// Example 3: time and population stay paired; never widened to {us,eu} × [0,2). +#[test] +fn example_3_joint_regions() { + let joint = merge(vec![ + kll(Some(0..MIN), &[("region", "us")]), + kll(Some(MIN..2 * MIN), &[("region", "eu")]), + ]) + .unwrap(); + assert_eq!( + regions(&joint), + vec![ + region(Some(MIN..2 * MIN), &[("region", "eu")]), + region(Some(0..MIN), &[("region", "us")]), + ] + ); +} + +/// A table without a time column declares no time bounds. +#[test] +fn tabular_source_without_time_bounds() { + let by_region = merge(vec![ + kll(None, &[("region", "us")]), + kll(None, &[("region", "eu")]), + ]) + .unwrap(); + assert_eq!(regions(&by_region).len(), 2); + rejected( + merge(vec![ + kll(None, &[("region", "us")]), + kll(Some(0..MIN), &[("region", "us")]), + ]), + CoverageError::PossibleOverlap, + ); +} + +/// Inputs must read the same source and share the producer's update and reduction. +#[test] +fn incompatible_inputs() { + let other_source = kll_over( + "latency", + None, + SummaryCoverage { + source: Source::Table { + table_ref: "other".into(), + }, + regions: vec![region(Some(MIN..2 * MIN), &[])], + }, + ); + rejected( + merge(vec![kll(Some(0..MIN), &[]), other_source]), + CoverageError::SourceMismatch, + ); + + // Same schema (both Float64 columns), different update expression. + let size = kll_over( + "size", + None, + SummaryCoverage { + source: requests(), + regions: vec![region(Some(MIN..2 * MIN), &[])], + }, + ); + assert!(matches!( + merge(vec![kll(Some(0..MIN), &[]), size]), + Err(SchemaDerivationError::InvalidScalarSignature(message)) + if message.contains("update expression and reduction") + )); +} + +/// Summary nodes must carry coverage, and merges reject inputs without it. +#[test] +fn coverage_is_required() { + let mut missing = (*kll(Some(MIN..2 * MIN), &[])).clone(); + missing.coverage = None; + let missing = Rc::new(missing); + assert!(matches!( + missing.validate_structure(), + Err(SchemaDerivationError::Coverage(CoverageError::Missing)) + )); + rejected( + merge(vec![kll(Some(0..MIN), &[]), missing]), + CoverageError::UnknownInput, + ); +} + +/// Trusted declarations: population is not checked against the filter (#570). +/// Both states hold US data, yet the wrong declaration lets them merge. +#[test] +fn wrong_population_declaration_is_accepted_until_570() { + let declared = |value: &str| SummaryCoverage { + source: requests(), + regions: vec![region(Some(0..MIN), &[("region", value)])], + }; + let a = kll_over("latency", Some(region_is("us")), declared("eu")); + let b = kll_over("latency", Some(region_is("us")), declared("us")); + assert!(merge(vec![a, b]).is_ok()); +} diff --git a/crates/types/tests/summary_merge_structure.rs b/crates/types/tests/summary_merge_structure.rs index 6489159af..383b2bfbd 100644 --- a/crates/types/tests/summary_merge_structure.rs +++ b/crates/types/tests/summary_merge_structure.rs @@ -30,9 +30,9 @@ fn state(k: u32) -> Rc { std::rc::Rc::new( summary .with_coverage(asap_types::ir::summary_coverage::SummaryCoverage { - source: "latencies:timestamp".into(), - input: SummaryUpdate::column(ColumnRef::SampleValue), - reduction: Reduction::by(vec![]), + source: Source::Table { + table_ref: "latencies".into(), + }, regions: vec![asap_types::ir::summary_coverage::CoverageRegion { time_ms: Some(0..1), population: Default::default(), From 32c4b24f5a9c011a2828a7ca8c9db0c9e1948464 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Sat, 3 Oct 2026 19:32:26 +0000 Subject: [PATCH 05/13] test(ir): point coverage examples at the design document Co-Authored-By: Claude Opus 5.5 --- crates/types/tests/summary_coverage_examples.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/crates/types/tests/summary_coverage_examples.rs b/crates/types/tests/summary_coverage_examples.rs index e2fad07ab..00785319e 100644 --- a/crates/types/tests/summary_coverage_examples.rs +++ b/crates/types/tests/summary_coverage_examples.rs @@ -1,4 +1,4 @@ -//! The examples in docs/develop_docs/summary-coverage.md, built as real +//! The examples in docs/design_docs/proposals/summary-coverage.md, built as real //! SummaryAgg -> SummaryMerge plans. Every input has the same schema //! `(job: Utf8, state: KLL{k=200})`; only coverage differs. use asap_types::ir::summary_coverage::{CoverageError, CoverageRegion, SummaryCoverage}; From 4c94aec10bd605e7c2da070b472a247973484ed2 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Sat, 3 Oct 2026 19:45:47 +0000 Subject: [PATCH 06/13] test(ir): point coverage examples at the ASAP primitive schema design doc Co-Authored-By: Claude Opus 5.5 --- crates/types/tests/summary_coverage_examples.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/crates/types/tests/summary_coverage_examples.rs b/crates/types/tests/summary_coverage_examples.rs index 00785319e..bfaa28483 100644 --- a/crates/types/tests/summary_coverage_examples.rs +++ b/crates/types/tests/summary_coverage_examples.rs @@ -1,4 +1,4 @@ -//! The examples in docs/design_docs/proposals/summary-coverage.md, built as real +//! The examples in docs/design_docs/proposals/asap-primitive-schema.md, built as real //! SummaryAgg -> SummaryMerge plans. Every input has the same schema //! `(job: Utf8, state: KLL{k=200})`; only coverage differs. use asap_types::ir::summary_coverage::{CoverageError, CoverageRegion, SummaryCoverage}; From 548339f61c8c5d579b56141057bf8c2a5925822d Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Tue, 6 Oct 2026 19:30:55 +0000 Subject: [PATCH 07/13] feat(ir): reject summary merges of heap-based sketches A merged top-k heap can miss an item that is heavy in only one input, so CmsWithHeap, CountSketchWithHeap and UnivMon states do not merge. Also correct the merge_disjoint doc: SummaryMerge checks schema, update, reduction and heap families, not accuracy. Co-Authored-By: Claude Opus 5.5 --- crates/types/src/ir/asap.rs | 27 +++++++++++-- crates/types/src/ir/summary_coverage.rs | 6 ++- crates/types/tests/summary_merge_structure.rs | 38 ++++++++++++++++--- 3 files changed, 61 insertions(+), 10 deletions(-) diff --git a/crates/types/src/ir/asap.rs b/crates/types/src/ir/asap.rs index 80a5fd9b8..5723ac72d 100644 --- a/crates/types/src/ir/asap.rs +++ b/crates/types/src/ir/asap.rs @@ -10,7 +10,7 @@ use super::summary_coverage::{CoverageError, SummaryCoverage}; use crate::ir::operator_properties::Reduction; use crate::ir::SchemaDerivationError; use crate::post_asap::maintained_population::{MaintainedPopulation, PopulationStatistic}; -use crate::post_asap::sketch::{GroupingStrategy, SketchStatistic, SummaryUpdate}; +use crate::post_asap::sketch::{GroupingStrategy, SketchAlgorithm, SketchStatistic, SummaryUpdate}; use crate::pre_asap::schema::{ColumnId, DataType, Field, FieldDataType, Schema}; /// Why an ASAP operator cannot be used yet. @@ -42,7 +42,9 @@ pub enum ASAPOp> { }, /// Read an exact accumulator's state as its finalized value: the /// maintenance-to-read boundary before query-time operators. - FinalizeExactAccumulator { child: C }, + FinalizeExactAccumulator { + child: C, + }, /// Maintain the full declared population, including membership changes. MaintainPopulation { child: C, @@ -54,7 +56,9 @@ pub enum ASAPOp> { evaluation: PopulationStatistic, }, /// Merge compatible partial states for the same grouping and family. - SummaryMerge { children: Vec }, + SummaryMerge { + children: Vec, + }, // ── Reserved: migrated but unimplemented ── SummarySubtract { left: C, @@ -498,6 +502,23 @@ impl ASAPOp { )); } } + // A heap of top-k candidates does not merge exactly: an item heavy + // in only one input can be missing from the merged heap. + if first.schema.fields.iter().any(|field| { + matches!( + &field.dtype, + FieldDataType::Sketch(kind, _) if matches!( + kind.algorithm(), + SketchAlgorithm::CmsWithHeap + | SketchAlgorithm::CountSketchWithHeap + | SketchAlgorithm::UnivMon + ) + ) + }) { + return Err(SchemaDerivationError::InvalidScalarSignature( + "summary merge does not support heap-based sketches".into(), + )); + } // Coverage records only time and population; what each state // summarizes and how it is grouped come from the producers. let update = first.summary_update(); diff --git a/crates/types/src/ir/summary_coverage.rs b/crates/types/src/ir/summary_coverage.rs index 120fb1212..032426340 100644 --- a/crates/types/src/ir/summary_coverage.rs +++ b/crates/types/src/ir/summary_coverage.rs @@ -75,8 +75,10 @@ impl SummaryCoverage { /// Every observation in a region is assumed to contribute once to the state. /// Compose once-per-observation summaries only when their joint regions are - /// provably disjoint. Update/reduction compatibility, family merge capability - /// and accuracy are checked by `SummaryMerge`, not here. + /// provably disjoint. `SummaryMerge` checks the rest: equal schemas (so the + /// same family and parameters), equal update and reduction, and no + /// heap-based family. Accuracy of the merged state is not assessed; its + /// `guarantee` is `None`. pub fn merge_disjoint(inputs: &[Self]) -> Result { let first = inputs.first().ok_or(CoverageError::EmptyMerge)?; let mut merged = first.clone(); diff --git a/crates/types/tests/summary_merge_structure.rs b/crates/types/tests/summary_merge_structure.rs index 383b2bfbd..eb92054a3 100644 --- a/crates/types/tests/summary_merge_structure.rs +++ b/crates/types/tests/summary_merge_structure.rs @@ -7,6 +7,13 @@ use asap_types::{ }; use std::rc::Rc; fn state(k: u32) -> Rc { + state_with(FieldDataType::Sketch( + SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k }), + Default::default(), + )) +} + +fn state_with(family: FieldDataType) -> Rc { let scan = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Scan { source: Source::Table { table_ref: "latencies".into(), @@ -17,10 +24,7 @@ fn state(k: u32) -> Rc { .unwrap(); let summary = OperatorNode::new(Operator::ASAP(ASAPOp::SummaryAgg { child: scan, - family: FieldDataType::Sketch( - SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k }), - Default::default(), - ), + family, input: SummaryUpdate::column(ColumnRef::SampleValue), reduction: Reduction::by(vec![]), grouping: GroupingStrategy::default(), @@ -70,11 +74,35 @@ fn incompatible_merge_inputs_fail() { } fn shifted_state(k: u32, start: i64, end: i64) -> Rc { - let mut node = (*state(k)).clone(); + shifted(&state(k), start, end) +} + +fn shifted(state: &Rc, start: i64, end: i64) -> Rc { + let mut node = (**state).clone(); let region = &mut node.coverage.as_mut().unwrap().regions[0]; region.time_ms = Some(start..end); Rc::new(node) } + +/// A heap of top-k candidates does not merge exactly, even over disjoint panes. +#[test] +fn heap_based_sketches_do_not_merge() { + let heap = state_with(FieldDataType::Sketch( + SketchKind::new( + SketchAlgorithm::CmsWithHeap, + SketchParams::CmsWithHeap { + width: 1024, + depth: 4, + heap_size: 10, + }, + ), + Default::default(), + )); + let result = OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryMerge { + children: vec![Rc::clone(&heap), shifted(&heap, 1, 2)], + })); + assert!(result.is_err()); +} /// Schema equality cannot authorize overlapping or unknown observation coverage. #[test] fn unsafe_coverage_merge_is_rejected() { From d951cbb0262d6702156a341097d6a9b75ba754bd Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Tue, 6 Oct 2026 21:25:07 +0000 Subject: [PATCH 08/13] refactor(ir): compute merge coverage with SummaryCoverage::of_merge `ASAPOp::merged_coverage` only applied to SummaryMerge and returned an error for every other operator. Replace it with `SummaryCoverage::of_merge`, which takes the merge's inputs; the callers already have them. Also explain what `OperatorNode::summary_update` returns and why merging compares it. Co-Authored-By: Claude Opus 5.5 --- crates/types/src/ir/asap.rs | 18 ++---------------- crates/types/src/ir/node.rs | 16 ++++++++++------ crates/types/src/ir/summary_coverage.rs | 12 ++++++++++++ 3 files changed, 24 insertions(+), 22 deletions(-) diff --git a/crates/types/src/ir/asap.rs b/crates/types/src/ir/asap.rs index 5723ac72d..17a20de0c 100644 --- a/crates/types/src/ir/asap.rs +++ b/crates/types/src/ir/asap.rs @@ -6,7 +6,7 @@ use std::rc::Rc; use serde::{Deserialize, Serialize}; use super::node::{OperatorNode, OperatorResultKind}; -use super::summary_coverage::{CoverageError, SummaryCoverage}; +use super::summary_coverage::SummaryCoverage; use crate::ir::operator_properties::Reduction; use crate::ir::SchemaDerivationError; use crate::post_asap::maintained_population::{MaintainedPopulation, PopulationStatistic}; @@ -218,20 +218,6 @@ impl ASAPOp { } } - /// Derive the merged node's coverage; unknown or overlapping inputs fail closed. - pub fn merged_coverage(&self) -> Result { - let ASAPOp::SummaryMerge { children } = self else { - return Err(SchemaDerivationError::InvalidScalarSignature( - "coverage merge requires SummaryMerge".into(), - )); - }; - let inputs = children - .iter() - .map(|child| child.coverage.clone().ok_or(CoverageError::UnknownInput)) - .collect::, _>>()?; - Ok(SummaryCoverage::merge_disjoint(&inputs)?) - } - /// Output schema derived from the operator and its children. Summary /// planning may retain a more specific schema (evaluation column naming) /// through [`OperatorNode::with_schema`]; all structural metadata must @@ -527,7 +513,7 @@ impl ASAPOp { "summary merge inputs must share update expression and reduction".into(), )); } - self.merged_coverage()?; + SummaryCoverage::of_merge(children)?; Ok(()) } SummaryEstimate { diff --git a/crates/types/src/ir/node.rs b/crates/types/src/ir/node.rs index 3cdbe8ecf..d2c345d14 100644 --- a/crates/types/src/ir/node.rs +++ b/crates/types/src/ir/node.rs @@ -113,8 +113,8 @@ impl OperatorNode { pub fn new(operator: Operator) -> Result { let schema = operator.output_schema()?; let mut node = Self::with_schema(operator, schema); - if let Some(op @ ASAPOp::SummaryMerge { .. }) = node.asap() { - node.coverage = Some(op.merged_coverage()?); + if let Some(ASAPOp::SummaryMerge { children }) = node.asap() { + node.coverage = Some(SummaryCoverage::of_merge(children)?); } Ok(node) } @@ -165,8 +165,12 @@ impl OperatorNode { Ok(self) } - /// What a summary state is updated with and how it is grouped: the - /// `SummaryAgg` fields, or those shared by a `SummaryMerge`'s inputs. + /// What a summary state was built from: the update expression (`input`) + /// and grouping (`reduction`) of the `SummaryAgg` that produced it. For a + /// `SummaryMerge` these are its first input's, which merge validation + /// requires every input to share. `None` for any other node. Merging + /// compares it so that all inputs summarize the same expression with the + /// same grouping. pub fn summary_update(&self) -> Option<(&SummaryUpdate, &Reduction)> { match self.asap()? { ASAPOp::SummaryAgg { @@ -343,8 +347,8 @@ impl OperatorNode { None if node.requires_coverage() => return Err(CoverageError::Missing.into()), None => {} } - if let Some(op @ ASAPOp::SummaryMerge { .. }) = node.asap() { - if node.coverage.as_ref() != Some(&op.merged_coverage()?) { + if let Some(ASAPOp::SummaryMerge { children }) = node.asap() { + if node.coverage.as_ref() != Some(&SummaryCoverage::of_merge(children)?) { return Err(CoverageError::MergeOutputMismatch.into()); } } diff --git a/crates/types/src/ir/summary_coverage.rs b/crates/types/src/ir/summary_coverage.rs index 032426340..0583a7fbb 100644 --- a/crates/types/src/ir/summary_coverage.rs +++ b/crates/types/src/ir/summary_coverage.rs @@ -73,6 +73,18 @@ impl SummaryCoverage { Ok(()) } + /// Coverage of a `SummaryMerge` over `inputs`: the disjoint union of their + /// coverage. Fails when an input has none or two inputs may overlap. + pub fn of_merge( + inputs: &[std::rc::Rc], + ) -> Result { + let coverage = inputs + .iter() + .map(|input| input.coverage.clone().ok_or(CoverageError::UnknownInput)) + .collect::, _>>()?; + Self::merge_disjoint(&coverage) + } + /// Every observation in a region is assumed to contribute once to the state. /// Compose once-per-observation summaries only when their joint regions are /// provably disjoint. `SummaryMerge` checks the rest: equal schemas (so the From 8caceedb54fdf124081f38406facd7274dcfcc96 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Wed, 7 Oct 2026 13:24:38 +0000 Subject: [PATCH 09/13] refactor(ir): leave summary merge coverage to #646; allow heap-based sketches SummaryMerge no longer derives or checks coverage: of_merge, UnknownInput and MergeOutputMismatch are removed and summary_coverage.rs matches main. Coverage for all summary nodes will be derived by one function in #646. The heap-based sketch restriction is also dropped. Co-Authored-By: Claude Opus 5.5 --- crates/types/src/ir/asap.rs | 25 +- crates/types/src/ir/node.rs | 16 +- crates/types/src/ir/summary_coverage.rs | 22 +- .../types/tests/summary_coverage_examples.rs | 264 ------------------ crates/types/tests/summary_merge_structure.rs | 72 ++--- 5 files changed, 31 insertions(+), 368 deletions(-) delete mode 100644 crates/types/tests/summary_coverage_examples.rs diff --git a/crates/types/src/ir/asap.rs b/crates/types/src/ir/asap.rs index 17a20de0c..3e56c191f 100644 --- a/crates/types/src/ir/asap.rs +++ b/crates/types/src/ir/asap.rs @@ -6,11 +6,10 @@ use std::rc::Rc; use serde::{Deserialize, Serialize}; use super::node::{OperatorNode, OperatorResultKind}; -use super::summary_coverage::SummaryCoverage; use crate::ir::operator_properties::Reduction; use crate::ir::SchemaDerivationError; use crate::post_asap::maintained_population::{MaintainedPopulation, PopulationStatistic}; -use crate::post_asap::sketch::{GroupingStrategy, SketchAlgorithm, SketchStatistic, SummaryUpdate}; +use crate::post_asap::sketch::{GroupingStrategy, SketchStatistic, SummaryUpdate}; use crate::pre_asap::schema::{ColumnId, DataType, Field, FieldDataType, Schema}; /// Why an ASAP operator cannot be used yet. @@ -488,32 +487,14 @@ impl ASAPOp { )); } } - // A heap of top-k candidates does not merge exactly: an item heavy - // in only one input can be missing from the merged heap. - if first.schema.fields.iter().any(|field| { - matches!( - &field.dtype, - FieldDataType::Sketch(kind, _) if matches!( - kind.algorithm(), - SketchAlgorithm::CmsWithHeap - | SketchAlgorithm::CountSketchWithHeap - | SketchAlgorithm::UnivMon - ) - ) - }) { - return Err(SchemaDerivationError::InvalidScalarSignature( - "summary merge does not support heap-based sketches".into(), - )); - } - // Coverage records only time and population; what each state - // summarizes and how it is grouped come from the producers. + // Equal schemas do not prove the states summarize the same + // expression with the same grouping; the producers do. let update = first.summary_update(); if update.is_none() || children.iter().any(|c| c.summary_update() != update) { return Err(SchemaDerivationError::InvalidScalarSignature( "summary merge inputs must share update expression and reduction".into(), )); } - SummaryCoverage::of_merge(children)?; Ok(()) } SummaryEstimate { diff --git a/crates/types/src/ir/node.rs b/crates/types/src/ir/node.rs index d2c345d14..ffd74b0eb 100644 --- a/crates/types/src/ir/node.rs +++ b/crates/types/src/ir/node.rs @@ -112,11 +112,7 @@ impl OperatorNode { /// ASAP operator, ...). pub fn new(operator: Operator) -> Result { let schema = operator.output_schema()?; - let mut node = Self::with_schema(operator, schema); - if let Some(ASAPOp::SummaryMerge { children }) = node.asap() { - node.coverage = Some(SummaryCoverage::of_merge(children)?); - } - Ok(node) + Ok(Self::with_schema(operator, schema)) } /// Build a node with caller-supplied output names and qualifiers. For @@ -183,10 +179,7 @@ impl OperatorNode { /// Summary nodes whose state can be composed must declare coverage. pub fn requires_coverage(&self) -> bool { - matches!( - self.asap(), - Some(ASAPOp::SummaryAgg { .. } | ASAPOp::SummaryMerge { .. }) - ) + matches!(self.asap(), Some(ASAPOp::SummaryAgg { .. })) } pub fn non_asap(&self) -> Option<&NonASAPOp> { @@ -347,11 +340,6 @@ impl OperatorNode { None if node.requires_coverage() => return Err(CoverageError::Missing.into()), None => {} } - if let Some(ASAPOp::SummaryMerge { children }) = node.asap() { - if node.coverage.as_ref() != Some(&SummaryCoverage::of_merge(children)?) { - return Err(CoverageError::MergeOutputMismatch.into()); - } - } node.operator.validate_inputs()?; if node.result_kind != node.operator.output_kind() { return Err(SchemaDerivationError::InvalidScalarSignature( diff --git a/crates/types/src/ir/summary_coverage.rs b/crates/types/src/ir/summary_coverage.rs index 0583a7fbb..3f3798a88 100644 --- a/crates/types/src/ir/summary_coverage.rs +++ b/crates/types/src/ir/summary_coverage.rs @@ -48,10 +48,6 @@ pub enum CoverageError { NotState, #[error("summary node requires coverage")] Missing, - #[error("summary merge requires known coverage on every input")] - UnknownInput, - #[error("retained merge coverage disagrees with input union")] - MergeOutputMismatch, } impl SummaryCoverage { @@ -73,24 +69,10 @@ impl SummaryCoverage { Ok(()) } - /// Coverage of a `SummaryMerge` over `inputs`: the disjoint union of their - /// coverage. Fails when an input has none or two inputs may overlap. - pub fn of_merge( - inputs: &[std::rc::Rc], - ) -> Result { - let coverage = inputs - .iter() - .map(|input| input.coverage.clone().ok_or(CoverageError::UnknownInput)) - .collect::, _>>()?; - Self::merge_disjoint(&coverage) - } - /// Every observation in a region is assumed to contribute once to the state. /// Compose once-per-observation summaries only when their joint regions are - /// provably disjoint. `SummaryMerge` checks the rest: equal schemas (so the - /// same family and parameters), equal update and reduction, and no - /// heap-based family. Accuracy of the merged state is not assessed; its - /// `guarantee` is `None`. + /// provably disjoint. Update/reduction compatibility, family merge capability + /// and accuracy are checked by `SummaryMerge`, not here. pub fn merge_disjoint(inputs: &[Self]) -> Result { let first = inputs.first().ok_or(CoverageError::EmptyMerge)?; let mut merged = first.clone(); diff --git a/crates/types/tests/summary_coverage_examples.rs b/crates/types/tests/summary_coverage_examples.rs deleted file mode 100644 index bfaa28483..000000000 --- a/crates/types/tests/summary_coverage_examples.rs +++ /dev/null @@ -1,264 +0,0 @@ -//! The examples in docs/design_docs/proposals/asap-primitive-schema.md, built as real -//! SummaryAgg -> SummaryMerge plans. Every input has the same schema -//! `(job: Utf8, state: KLL{k=200})`; only coverage differs. -use asap_types::ir::summary_coverage::{CoverageError, CoverageRegion, SummaryCoverage}; -use asap_types::ir::{ - ASAPOp, ExprSemantics, NonASAPOp, Operator, OperatorNode, Predicate, ScalarExpr, - SchemaDerivationError, -}; -use asap_types::post_asap::{ - GroupingStrategy, SketchAlgorithm, SketchKind, SketchParams, SummaryUpdate, -}; -use asap_types::pre_asap::{ - ColumnRef, CompareOpKind, DataType, Field, FieldDataType, Reduction, ScalarValue, Schema, - Source, -}; -use std::rc::Rc; - -const MIN: i64 = 60_000; - -fn requests() -> Source { - Source::Table { - table_ref: "requests".into(), - } -} - -fn scan() -> Rc { - OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Scan { - source: requests(), - predicates: vec![], - schema: Schema::new(vec![ - Field::plain("job", DataType::Utf8, false), - Field::plain("region", DataType::Utf8, false), - Field::plain("tier", DataType::Utf8, false), - Field::plain("latency", DataType::Float64, false), - Field::plain("size", DataType::Float64, false), - ]), - })) - .unwrap() -} - -/// `region = value`, evaluated on the scan schema. -fn region_is(value: &str) -> Predicate { - Predicate(ScalarExpr::Compare { - left: Box::new(ScalarExpr::Column(1)), - op: CompareOpKind::Eq, - right: Box::new(ScalarExpr::Literal(ScalarValue::Utf8(value.into()))), - semantics: ExprSemantics::Sql, - }) -} - -fn region(time_ms: Option>, population: &[(&str, &str)]) -> CoverageRegion { - CoverageRegion { - time_ms, - population: population - .iter() - .map(|(k, v)| (k.to_string(), v.to_string())) - .collect(), - } -} - -/// p99-ready KLL over `column`, grouped by job, with declared coverage. -fn kll_over( - column: &str, - filter: Option, - coverage: SummaryCoverage, -) -> Rc { - let node = OperatorNode::new(Operator::ASAP(ASAPOp::SummaryAgg { - child: scan(), - family: FieldDataType::Sketch( - SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k: 200 }), - GroupingStrategy::default(), - ), - input: SummaryUpdate::column(ColumnRef::Named(column.into())), - reduction: Reduction::by(vec![0]), - grouping: GroupingStrategy::default(), - filter, - })) - .unwrap(); - Rc::new(node.with_coverage(coverage).unwrap()) -} - -fn kll(time_ms: Option>, population: &[(&str, &str)]) -> Rc { - kll_over( - "latency", - None, - SummaryCoverage { - source: requests(), - regions: vec![region(time_ms, population)], - }, - ) -} - -fn merge(children: Vec>) -> Result, SchemaDerivationError> { - let schema = children[0].schema.clone(); - let merged = OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryMerge { children }))?; - merged.validate_structure()?; - // Schema never changes; only coverage does. - assert_eq!(merged.schema, schema); - Ok(merged) -} - -fn regions(node: &OperatorNode) -> Vec { - node.coverage.as_ref().unwrap().regions.clone() -} - -fn rejected(result: Result, SchemaDerivationError>, expected: CoverageError) { - match result { - Err(SchemaDerivationError::Coverage(actual)) => assert_eq!(actual, expected), - other => panic!("expected {expected:?}, got {other:?}"), - } -} - -/// Example 1: adjacent panes coalesce, gaps stay, overlapping windows are rejected. -#[test] -fn example_1_time() { - let adjacent = merge(vec![kll(Some(0..MIN), &[]), kll(Some(MIN..2 * MIN), &[])]).unwrap(); - assert_eq!(regions(&adjacent), vec![region(Some(0..2 * MIN), &[])]); - - let gapped = merge(vec![ - kll(Some(0..MIN), &[]), - kll(Some(2 * MIN..3 * MIN), &[]), - ]) - .unwrap(); - assert_eq!( - regions(&gapped), - vec![ - region(Some(0..MIN), &[]), - region(Some(2 * MIN..3 * MIN), &[]) - ] - ); - - rejected( - merge(vec![ - kll(Some(0..2 * MIN), &[]), - kll(Some(MIN..3 * MIN), &[]), - ]), - CoverageError::PossibleOverlap, - ); -} - -/// Example 2: disjoint label values merge; different labels or equal values are rejected. -#[test] -fn example_2_population() { - let t = Some(0..MIN); - let us_eu = merge(vec![ - kll(t.clone(), &[("region", "us")]), - kll(t.clone(), &[("region", "eu")]), - ]) - .unwrap(); - assert_eq!( - regions(&us_eu), - vec![ - region(t.clone(), &[("region", "eu")]), - region(t.clone(), &[("region", "us")]), - ] - ); - for other in [("tier", "premium"), ("region", "us")] { - rejected( - merge(vec![ - kll(t.clone(), &[("region", "us")]), - kll(t.clone(), &[other]), - ]), - CoverageError::PossibleOverlap, - ); - } -} - -/// Example 3: time and population stay paired; never widened to {us,eu} × [0,2). -#[test] -fn example_3_joint_regions() { - let joint = merge(vec![ - kll(Some(0..MIN), &[("region", "us")]), - kll(Some(MIN..2 * MIN), &[("region", "eu")]), - ]) - .unwrap(); - assert_eq!( - regions(&joint), - vec![ - region(Some(MIN..2 * MIN), &[("region", "eu")]), - region(Some(0..MIN), &[("region", "us")]), - ] - ); -} - -/// A table without a time column declares no time bounds. -#[test] -fn tabular_source_without_time_bounds() { - let by_region = merge(vec![ - kll(None, &[("region", "us")]), - kll(None, &[("region", "eu")]), - ]) - .unwrap(); - assert_eq!(regions(&by_region).len(), 2); - rejected( - merge(vec![ - kll(None, &[("region", "us")]), - kll(Some(0..MIN), &[("region", "us")]), - ]), - CoverageError::PossibleOverlap, - ); -} - -/// Inputs must read the same source and share the producer's update and reduction. -#[test] -fn incompatible_inputs() { - let other_source = kll_over( - "latency", - None, - SummaryCoverage { - source: Source::Table { - table_ref: "other".into(), - }, - regions: vec![region(Some(MIN..2 * MIN), &[])], - }, - ); - rejected( - merge(vec![kll(Some(0..MIN), &[]), other_source]), - CoverageError::SourceMismatch, - ); - - // Same schema (both Float64 columns), different update expression. - let size = kll_over( - "size", - None, - SummaryCoverage { - source: requests(), - regions: vec![region(Some(MIN..2 * MIN), &[])], - }, - ); - assert!(matches!( - merge(vec![kll(Some(0..MIN), &[]), size]), - Err(SchemaDerivationError::InvalidScalarSignature(message)) - if message.contains("update expression and reduction") - )); -} - -/// Summary nodes must carry coverage, and merges reject inputs without it. -#[test] -fn coverage_is_required() { - let mut missing = (*kll(Some(MIN..2 * MIN), &[])).clone(); - missing.coverage = None; - let missing = Rc::new(missing); - assert!(matches!( - missing.validate_structure(), - Err(SchemaDerivationError::Coverage(CoverageError::Missing)) - )); - rejected( - merge(vec![kll(Some(0..MIN), &[]), missing]), - CoverageError::UnknownInput, - ); -} - -/// Trusted declarations: population is not checked against the filter (#570). -/// Both states hold US data, yet the wrong declaration lets them merge. -#[test] -fn wrong_population_declaration_is_accepted_until_570() { - let declared = |value: &str| SummaryCoverage { - source: requests(), - regions: vec![region(Some(0..MIN), &[("region", value)])], - }; - let a = kll_over("latency", Some(region_is("us")), declared("eu")); - let b = kll_over("latency", Some(region_is("us")), declared("us")); - assert!(merge(vec![a, b]).is_ok()); -} diff --git a/crates/types/tests/summary_merge_structure.rs b/crates/types/tests/summary_merge_structure.rs index eb92054a3..f76807f0c 100644 --- a/crates/types/tests/summary_merge_structure.rs +++ b/crates/types/tests/summary_merge_structure.rs @@ -1,7 +1,7 @@ //! Window composition merges compatible summary states without consuming raw rows. use asap_types::{ ir::operator_properties::{Reduction, Source}, - ir::{ASAPOp, NonASAPOp, Operator, OperatorNode}, + ir::{ASAPOp, NonASAPOp, Operator, OperatorNode, SchemaDerivationError}, post_asap::{GroupingStrategy, SketchAlgorithm, SketchKind, SketchParams, SummaryUpdate}, pre_asap::{ColumnRef, DataType, Field, FieldDataType, Schema}, }; @@ -14,6 +14,10 @@ fn state(k: u32) -> Rc { } fn state_with(family: FieldDataType) -> Rc { + state_over(family, SummaryUpdate::column(ColumnRef::SampleValue)) +} + +fn state_over(family: FieldDataType, input: SummaryUpdate) -> Rc { let scan = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Scan { source: Source::Table { table_ref: "latencies".into(), @@ -25,7 +29,7 @@ fn state_with(family: FieldDataType) -> Rc { let summary = OperatorNode::new(Operator::ASAP(ASAPOp::SummaryAgg { child: scan, family, - input: SummaryUpdate::column(ColumnRef::SampleValue), + input, reduction: Reduction::by(vec![]), grouping: GroupingStrategy::default(), filter: None, @@ -54,10 +58,6 @@ fn compatible_panes_merge_structurally() { .unwrap(); root.validate_structure().unwrap(); assert_eq!(root.schema.fields.len(), 1); - assert_eq!( - root.coverage.as_ref().unwrap().regions[0].time_ms, - Some(0..2) - ); } /// An empty merge, raw rows and differently sized state cannot masquerade as compatible panes. #[test] @@ -84,49 +84,25 @@ fn shifted(state: &Rc, start: i64, end: i64) -> Rc { Rc::new(node) } -/// A heap of top-k candidates does not merge exactly, even over disjoint panes. +/// Equal schemas do not prove both states summarize the same expression. #[test] -fn heap_based_sketches_do_not_merge() { - let heap = state_with(FieldDataType::Sketch( - SketchKind::new( - SketchAlgorithm::CmsWithHeap, - SketchParams::CmsWithHeap { - width: 1024, - depth: 4, - heap_size: 10, - }, - ), - Default::default(), - )); +fn different_update_expressions_do_not_merge() { + let kll = || { + FieldDataType::Sketch( + SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k: 200 }), + Default::default(), + ) + }; + let named = state_over( + kll(), + SummaryUpdate::column(ColumnRef::Named("value".into())), + ); let result = OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryMerge { - children: vec![Rc::clone(&heap), shifted(&heap, 1, 2)], + children: vec![state(200), shifted(&named, 1, 2)], })); - assert!(result.is_err()); -} -/// Schema equality cannot authorize overlapping or unknown observation coverage. -#[test] -fn unsafe_coverage_merge_is_rejected() { - let mut unknown = (*state(200)).clone(); - unknown.coverage = None; - for children in [ - vec![state(200), state(200)], - vec![state(200), Rc::new(unknown)], - ] { - assert!( - OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryMerge { children })).is_err() - ); - } -} - -/// Gapped time coverage remains disconnected, and forged output metadata is rejected. -#[test] -fn merge_derives_coverage_and_validates_retained_metadata() { - let root = OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryMerge { - children: vec![state(200), shifted_state(200, 2, 3)], - })) - .unwrap(); - assert_eq!(root.coverage.as_ref().unwrap().regions.len(), 2); - let mut forged = (*root).clone(); - forged.coverage.as_mut().unwrap().regions[0].time_ms = Some(0..2); - assert!(Rc::new(forged).validate_structure().is_err()); + assert!(matches!( + result, + Err(SchemaDerivationError::InvalidScalarSignature(message)) + if message.contains("update expression and reduction") + )); } From cf3fc265cc0c041beec43096320d27f41d938805 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Wed, 7 Oct 2026 14:02:02 +0000 Subject: [PATCH 10/13] refactor(ir): rename summary_update to summary_input_data It reads the producing SummaryAgg's update expression and reduction; it does not update the summary. Co-Authored-By: Claude Opus 5.5 --- crates/types/src/ir/asap.rs | 8 ++++++-- crates/types/src/ir/node.rs | 16 +++++++++------- 2 files changed, 15 insertions(+), 9 deletions(-) diff --git a/crates/types/src/ir/asap.rs b/crates/types/src/ir/asap.rs index 3e56c191f..d9f5ca27f 100644 --- a/crates/types/src/ir/asap.rs +++ b/crates/types/src/ir/asap.rs @@ -489,8 +489,12 @@ impl ASAPOp { } // Equal schemas do not prove the states summarize the same // expression with the same grouping; the producers do. - let update = first.summary_update(); - if update.is_none() || children.iter().any(|c| c.summary_update() != update) { + let input_data = first.summary_input_data(); + if input_data.is_none() + || children + .iter() + .any(|c| c.summary_input_data() != input_data) + { return Err(SchemaDerivationError::InvalidScalarSignature( "summary merge inputs must share update expression and reduction".into(), )); diff --git a/crates/types/src/ir/node.rs b/crates/types/src/ir/node.rs index ffd74b0eb..a3e453cdb 100644 --- a/crates/types/src/ir/node.rs +++ b/crates/types/src/ir/node.rs @@ -161,18 +161,20 @@ impl OperatorNode { Ok(self) } - /// What a summary state was built from: the update expression (`input`) - /// and grouping (`reduction`) of the `SummaryAgg` that produced it. For a + /// What a summary state summarizes from each row: the update expression + /// fed into the state (`input`, e.g. the `latency` column) and the + /// grouping (`reduction`, e.g. by `job`) of the `SummaryAgg` that produced + /// it. Which rows were included is `coverage`, not this. For a /// `SummaryMerge` these are its first input's, which merge validation - /// requires every input to share. `None` for any other node. Merging - /// compares it so that all inputs summarize the same expression with the - /// same grouping. - pub fn summary_update(&self) -> Option<(&SummaryUpdate, &Reduction)> { + /// requires every input to share. `None` for any other node. Equal schemas + /// cannot tell a KLL over `latency` from one over `size`; merging compares + /// this instead. + pub fn summary_input_data(&self) -> Option<(&SummaryUpdate, &Reduction)> { match self.asap()? { ASAPOp::SummaryAgg { input, reduction, .. } => Some((input, reduction)), - ASAPOp::SummaryMerge { children } => children.first()?.summary_update(), + ASAPOp::SummaryMerge { children } => children.first()?.summary_input_data(), _ => None, } } From c8d6564cbe2bbfe6adb66fe25b514eebd591b804 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Wed, 7 Oct 2026 14:19:53 +0000 Subject: [PATCH 11/13] feat(ir): record summarized columns in SummaryCoverage SummaryCoverage now records which columns a state summarizes as well as which rows: `input` (the SummaryAgg update expression) and `group_by` (its reduction). with_coverage rejects a SummaryAgg declaration whose columns differ from the node's own (ColumnMismatch), and SummaryMerge requires coverage on every input (UnknownInput) with identical columns. merge_disjoint checks columns too. summary_input_data is removed. A nested SummaryMerge carries no coverage until #646, so it is rejected. Co-Authored-By: Claude Opus 5.5 --- crates/types/src/ir/asap.rs | 19 ++--- crates/types/src/ir/node.rs | 31 +++----- crates/types/src/ir/summary_coverage.rs | 34 +++++++-- crates/types/tests/schema_rebuilding.rs | 2 + crates/types/tests/structure_contract.rs | 4 + crates/types/tests/summary_coverage.rs | 13 ++++ crates/types/tests/summary_merge_structure.rs | 74 +++++++++++++++---- 7 files changed, 125 insertions(+), 52 deletions(-) diff --git a/crates/types/src/ir/asap.rs b/crates/types/src/ir/asap.rs index d9f5ca27f..9c9a7d549 100644 --- a/crates/types/src/ir/asap.rs +++ b/crates/types/src/ir/asap.rs @@ -6,6 +6,7 @@ use std::rc::Rc; use serde::{Deserialize, Serialize}; use super::node::{OperatorNode, OperatorResultKind}; +use super::summary_coverage::CoverageError; use crate::ir::operator_properties::Reduction; use crate::ir::SchemaDerivationError; use crate::post_asap::maintained_population::{MaintainedPopulation, PopulationStatistic}; @@ -487,17 +488,13 @@ impl ASAPOp { )); } } - // Equal schemas do not prove the states summarize the same - // expression with the same grouping; the producers do. - let input_data = first.summary_input_data(); - if input_data.is_none() - || children - .iter() - .any(|c| c.summary_input_data() != input_data) - { - return Err(SchemaDerivationError::InvalidScalarSignature( - "summary merge inputs must share update expression and reduction".into(), - )); + // Equal schemas cannot tell a KLL over `latency` from one over + // `size`; the columns in each input's coverage can. A nested + // `SummaryMerge` carries no coverage yet (#646), so it is rejected. + let columns = first.coverage.as_ref().ok_or(CoverageError::UnknownInput)?; + for child in children { + let coverage = child.coverage.as_ref().ok_or(CoverageError::UnknownInput)?; + coverage.check_columns(columns)?; } Ok(()) } diff --git a/crates/types/src/ir/node.rs b/crates/types/src/ir/node.rs index a3e453cdb..4333ee7fb 100644 --- a/crates/types/src/ir/node.rs +++ b/crates/types/src/ir/node.rs @@ -10,12 +10,10 @@ use serde::{Deserialize, Serialize}; use super::asap::ASAPOp; use super::non_asap::NonASAPOp; -use super::operator_properties::Reduction; use super::summary_coverage::{CoverageError, SummaryCoverage}; use crate::ir::SchemaDerivationError; use crate::post_asap::execution_data_state::ExecutionTiming; use crate::post_asap::guarantee::ResultGuarantee; -use crate::post_asap::SummaryUpdate; use crate::pre_asap::schema::Schema; /// The output category of an operator, derived from the operation and its @@ -148,7 +146,8 @@ impl OperatorNode { } /// Attach caller-established coverage. Required on summary nodes; see - /// [`Self::requires_coverage`]. + /// [`Self::requires_coverage`]. On a `SummaryAgg` the coverage columns + /// must equal the node's `input` and `reduction`. pub fn with_coverage( mut self, coverage: SummaryCoverage, @@ -157,28 +156,18 @@ impl OperatorNode { if self.result_kind != OperatorResultKind::State { return Err(CoverageError::NotState.into()); } + if let Some(ASAPOp::SummaryAgg { + input, reduction, .. + }) = self.asap() + { + if &coverage.input != input || &coverage.group_by != reduction { + return Err(CoverageError::ColumnMismatch.into()); + } + } self.coverage = Some(coverage); Ok(self) } - /// What a summary state summarizes from each row: the update expression - /// fed into the state (`input`, e.g. the `latency` column) and the - /// grouping (`reduction`, e.g. by `job`) of the `SummaryAgg` that produced - /// it. Which rows were included is `coverage`, not this. For a - /// `SummaryMerge` these are its first input's, which merge validation - /// requires every input to share. `None` for any other node. Equal schemas - /// cannot tell a KLL over `latency` from one over `size`; merging compares - /// this instead. - pub fn summary_input_data(&self) -> Option<(&SummaryUpdate, &Reduction)> { - match self.asap()? { - ASAPOp::SummaryAgg { - input, reduction, .. - } => Some((input, reduction)), - ASAPOp::SummaryMerge { children } => children.first()?.summary_input_data(), - _ => None, - } - } - /// Summary nodes whose state can be composed must declare coverage. pub fn requires_coverage(&self) -> bool { matches!(self.asap(), Some(ASAPOp::SummaryAgg { .. })) diff --git a/crates/types/src/ir/summary_coverage.rs b/crates/types/src/ir/summary_coverage.rs index 3f3798a88..de9af8880 100644 --- a/crates/types/src/ir/summary_coverage.rs +++ b/crates/types/src/ir/summary_coverage.rs @@ -1,6 +1,10 @@ -//! Joint time/population coverage for summary composition, independent of schema. +//! What a summary state covers, independent of schema: which rows (source and +//! joint time/population regions) and which columns (the update expression fed +//! into the state and the grouping). //! Equality predicates are a deliberately narrow proof vocabulary. Unsupported //! predicates cannot be declared disjoint merely by giving them different names. +use crate::ir::operator_properties::Reduction; +use crate::post_asap::SummaryUpdate; use crate::pre_asap::Source; use serde::{Deserialize, Serialize}; use std::collections::BTreeMap; @@ -13,8 +17,15 @@ pub struct SummaryCoverage { /// Observation data source, as named by `Scan`: a table or a time series. /// Region time bounds refer to its time column. pub source: Source, - /// Union of joint regions; never the Cartesian product of independent bounds. + /// Rows: union of joint regions; never the Cartesian product of + /// independent bounds. pub regions: Vec, + /// Columns: the update expression fed into the state from each row, equal + /// to the producing `SummaryAgg`'s `input` (e.g. the `latency` column). + pub input: SummaryUpdate, + /// Columns: the grouping, equal to the producing `SummaryAgg`'s + /// `reduction` (e.g. by `job`). Compared by child-schema column index. + pub group_by: Reduction, } #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] @@ -48,6 +59,10 @@ pub enum CoverageError { NotState, #[error("summary node requires coverage")] Missing, + #[error("summary coverage columns differ from the summary's input or grouping")] + ColumnMismatch, + #[error("summary merge requires coverage on every input")] + UnknownInput, } impl SummaryCoverage { @@ -70,9 +85,17 @@ impl SummaryCoverage { } /// Every observation in a region is assumed to contribute once to the state. - /// Compose once-per-observation summaries only when their joint regions are - /// provably disjoint. Update/reduction compatibility, family merge capability - /// and accuracy are checked by `SummaryMerge`, not here. + /// Compose once-per-observation summaries only when they read the same + /// columns and their joint regions are provably disjoint. Family merge + /// capability and accuracy are not checked here. + /// Rejects `other` unless it reads the same columns as `self`. + pub fn check_columns(&self, other: &Self) -> Result<(), CoverageError> { + if self.input != other.input || self.group_by != other.group_by { + return Err(CoverageError::ColumnMismatch); + } + Ok(()) + } + pub fn merge_disjoint(inputs: &[Self]) -> Result { let first = inputs.first().ok_or(CoverageError::EmptyMerge)?; let mut merged = first.clone(); @@ -82,6 +105,7 @@ impl SummaryCoverage { if input.source != first.source { return Err(CoverageError::SourceMismatch); } + input.check_columns(first)?; merged.regions.extend(input.regions.iter().cloned()); } merged.validate()?; diff --git a/crates/types/tests/schema_rebuilding.rs b/crates/types/tests/schema_rebuilding.rs index b0fc4c666..18924338f 100644 --- a/crates/types/tests/schema_rebuilding.rs +++ b/crates/types/tests/schema_rebuilding.rs @@ -17,6 +17,8 @@ fn coverage() -> SummaryCoverage { time_ms: None, population: Default::default(), }], + input: SummaryUpdate::column(ColumnRef::Named("value".into())), + group_by: Reduction::by(vec![0]), } } /// Rewrites clear coverage; a rewriter must declare it again for summary nodes. diff --git a/crates/types/tests/structure_contract.rs b/crates/types/tests/structure_contract.rs index d7c144f4e..04a03969d 100644 --- a/crates/types/tests/structure_contract.rs +++ b/crates/types/tests/structure_contract.rs @@ -24,6 +24,10 @@ fn whole_table() -> asap_types::ir::summary_coverage::SummaryCoverage { time_ms: None, population: Default::default(), }], + input: asap_types::post_asap::SummaryUpdate::column( + asap_types::pre_asap::ColumnRef::Named("x".into()), + ), + group_by: asap_types::pre_asap::Reduction::by(vec![]), } } /// Resolved filters cannot hide invalid scalar types or out-of-scope columns. diff --git a/crates/types/tests/summary_coverage.rs b/crates/types/tests/summary_coverage.rs index e36c3d7a4..8433ab6a7 100644 --- a/crates/types/tests/summary_coverage.rs +++ b/crates/types/tests/summary_coverage.rs @@ -20,6 +20,8 @@ fn coverage(start: i64, end: i64, population: &[(&str, &str)]) -> SummaryCoverag .map(|(k, v)| (k.to_string(), v.to_string())) .collect(), }], + input: SummaryUpdate::column(ColumnRef::Named("latency".into())), + group_by: Reduction::by(vec![]), } } /// Adjacent panes coalesce; gaps remain disconnected rather than becoming a hull. @@ -148,3 +150,14 @@ fn regions_without_time_bounds() { Err(CoverageError::PossibleOverlap) ); } + +/// Disjoint rows do not make coverage of different columns composable. +#[test] +fn different_columns_do_not_compose() { + let mut by_job = coverage(1, 2, &[]); + by_job.group_by = Reduction::by(vec![0]); + assert_eq!( + SummaryCoverage::merge_disjoint(&[coverage(0, 1, &[]), by_job]), + Err(CoverageError::ColumnMismatch) + ); +} diff --git a/crates/types/tests/summary_merge_structure.rs b/crates/types/tests/summary_merge_structure.rs index f76807f0c..3edda31a2 100644 --- a/crates/types/tests/summary_merge_structure.rs +++ b/crates/types/tests/summary_merge_structure.rs @@ -1,6 +1,7 @@ //! Window composition merges compatible summary states without consuming raw rows. use asap_types::{ ir::operator_properties::{Reduction, Source}, + ir::summary_coverage::CoverageError, ir::{ASAPOp, NonASAPOp, Operator, OperatorNode, SchemaDerivationError}, post_asap::{GroupingStrategy, SketchAlgorithm, SketchKind, SketchParams, SummaryUpdate}, pre_asap::{ColumnRef, DataType, Field, FieldDataType, Schema}, @@ -29,7 +30,7 @@ fn state_over(family: FieldDataType, input: SummaryUpdate) -> Rc { let summary = OperatorNode::new(Operator::ASAP(ASAPOp::SummaryAgg { child: scan, family, - input, + input: input.clone(), reduction: Reduction::by(vec![]), grouping: GroupingStrategy::default(), filter: None, @@ -45,6 +46,8 @@ fn state_over(family: FieldDataType, input: SummaryUpdate) -> Rc { time_ms: Some(0..1), population: Default::default(), }], + input, + group_by: Reduction::by(vec![]), }) .unwrap(), ) @@ -84,25 +87,66 @@ fn shifted(state: &Rc, start: i64, end: i64) -> Rc { Rc::new(node) } -/// Equal schemas do not prove both states summarize the same expression. +fn kll() -> FieldDataType { + FieldDataType::Sketch( + SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k: 200 }), + Default::default(), + ) +} + +fn merge(children: Vec>) -> Result, SchemaDerivationError> { + OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryMerge { children })) +} + +/// Equal schemas cannot tell a KLL over one column from a KLL over another; +/// the columns in their coverage can. #[test] -fn different_update_expressions_do_not_merge() { - let kll = || { - FieldDataType::Sketch( - SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k: 200 }), - Default::default(), - ) - }; +fn different_coverage_columns_do_not_merge() { let named = state_over( kll(), SummaryUpdate::column(ColumnRef::Named("value".into())), ); - let result = OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryMerge { - children: vec![state(200), shifted(&named, 1, 2)], - })); assert!(matches!( - result, - Err(SchemaDerivationError::InvalidScalarSignature(message)) - if message.contains("update expression and reduction") + merge(vec![state(200), shifted(&named, 1, 2)]), + Err(SchemaDerivationError::Coverage( + CoverageError::ColumnMismatch + )) )); } + +/// A `SummaryAgg` cannot declare columns other than its own input and grouping. +#[test] +fn declared_columns_must_match_the_summary() { + let mut coverage = state(200).coverage.clone().unwrap(); + coverage.input = SummaryUpdate::column(ColumnRef::Named("value".into())); + let fresh = (*state(200)).clone(); + assert!(matches!( + fresh.clone().with_coverage(coverage.clone()), + Err(SchemaDerivationError::Coverage( + CoverageError::ColumnMismatch + )) + )); + let mut forged = fresh; + forged.coverage = Some(coverage); + assert!(matches!( + Rc::new(forged).validate_structure(), + Err(SchemaDerivationError::Coverage( + CoverageError::ColumnMismatch + )) + )); +} + +/// Every merge input needs coverage. A nested merge has none until #646, so +/// it is rejected for now. +#[test] +fn merge_inputs_without_coverage_are_rejected() { + let mut missing = (*state(200)).clone(); + missing.coverage = None; + let nested = merge(vec![state(200)]).unwrap(); + for input in [Rc::new(missing), nested] { + assert!(matches!( + merge(vec![state(200), input]), + Err(SchemaDerivationError::Coverage(CoverageError::UnknownInput)) + )); + } +} From 7ce0660f92abdc6f17bc9352e6a4f3db82b8c1c7 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Wed, 7 Oct 2026 14:37:03 +0000 Subject: [PATCH 12/13] docs(ir): attach the merge_disjoint doc comment to merge_disjoint Co-Authored-By: Claude Opus 5.5 --- crates/types/src/ir/summary_coverage.rs | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/crates/types/src/ir/summary_coverage.rs b/crates/types/src/ir/summary_coverage.rs index de9af8880..a06bf4f07 100644 --- a/crates/types/src/ir/summary_coverage.rs +++ b/crates/types/src/ir/summary_coverage.rs @@ -84,10 +84,6 @@ impl SummaryCoverage { Ok(()) } - /// Every observation in a region is assumed to contribute once to the state. - /// Compose once-per-observation summaries only when they read the same - /// columns and their joint regions are provably disjoint. Family merge - /// capability and accuracy are not checked here. /// Rejects `other` unless it reads the same columns as `self`. pub fn check_columns(&self, other: &Self) -> Result<(), CoverageError> { if self.input != other.input || self.group_by != other.group_by { @@ -96,6 +92,10 @@ impl SummaryCoverage { Ok(()) } + /// Every observation in a region is assumed to contribute once to the state. + /// Compose once-per-observation summaries only when they read the same + /// columns and their joint regions are provably disjoint. Family merge + /// capability and accuracy are not checked here. pub fn merge_disjoint(inputs: &[Self]) -> Result { let first = inputs.first().ok_or(CoverageError::EmptyMerge)?; let mut merged = first.clone(); From 4a0eab70d0c582c3df8292b80274e4a689a21e21 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Wed, 7 Oct 2026 22:49:35 +0000 Subject: [PATCH 13/13] refactor(ir): keep SummaryMerge structural; leave coverage to #646 Remove the coverage columns (input, group_by), check_columns and the merge's coverage checks. SummaryMerge now only checks structure: at least one input, every input is State with one state column and an identical schema. Whether a structurally valid merge is semantically valid (same computation, disjoint selections) is decided by summary coverage in #646, following the design in #573. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01W7qG9aFyPij5uWsyAJCxDW --- crates/types/src/ir/asap.rs | 10 +-- crates/types/src/ir/node.rs | 11 +-- crates/types/src/ir/summary_coverage.rs | 34 ++------- crates/types/tests/schema_rebuilding.rs | 2 - crates/types/tests/structure_contract.rs | 4 -- crates/types/tests/summary_coverage.rs | 13 ---- crates/types/tests/summary_merge_structure.rs | 71 ++----------------- 7 files changed, 14 insertions(+), 131 deletions(-) diff --git a/crates/types/src/ir/asap.rs b/crates/types/src/ir/asap.rs index 9c9a7d549..16945d925 100644 --- a/crates/types/src/ir/asap.rs +++ b/crates/types/src/ir/asap.rs @@ -6,7 +6,6 @@ use std::rc::Rc; use serde::{Deserialize, Serialize}; use super::node::{OperatorNode, OperatorResultKind}; -use super::summary_coverage::CoverageError; use crate::ir::operator_properties::Reduction; use crate::ir::SchemaDerivationError; use crate::post_asap::maintained_population::{MaintainedPopulation, PopulationStatistic}; @@ -489,13 +488,8 @@ impl ASAPOp { } } // Equal schemas cannot tell a KLL over `latency` from one over - // `size`; the columns in each input's coverage can. A nested - // `SummaryMerge` carries no coverage yet (#646), so it is rejected. - let columns = first.coverage.as_ref().ok_or(CoverageError::UnknownInput)?; - for child in children { - let coverage = child.coverage.as_ref().ok_or(CoverageError::UnknownInput)?; - coverage.check_columns(columns)?; - } + // `size`, nor prove the inputs disjoint; summary coverage (#646) + // decides whether a structurally valid merge is semantically valid. Ok(()) } SummaryEstimate { diff --git a/crates/types/src/ir/node.rs b/crates/types/src/ir/node.rs index 4333ee7fb..47ffa0ea3 100644 --- a/crates/types/src/ir/node.rs +++ b/crates/types/src/ir/node.rs @@ -146,8 +146,7 @@ impl OperatorNode { } /// Attach caller-established coverage. Required on summary nodes; see - /// [`Self::requires_coverage`]. On a `SummaryAgg` the coverage columns - /// must equal the node's `input` and `reduction`. + /// [`Self::requires_coverage`]. pub fn with_coverage( mut self, coverage: SummaryCoverage, @@ -156,14 +155,6 @@ impl OperatorNode { if self.result_kind != OperatorResultKind::State { return Err(CoverageError::NotState.into()); } - if let Some(ASAPOp::SummaryAgg { - input, reduction, .. - }) = self.asap() - { - if &coverage.input != input || &coverage.group_by != reduction { - return Err(CoverageError::ColumnMismatch.into()); - } - } self.coverage = Some(coverage); Ok(self) } diff --git a/crates/types/src/ir/summary_coverage.rs b/crates/types/src/ir/summary_coverage.rs index a06bf4f07..3f3798a88 100644 --- a/crates/types/src/ir/summary_coverage.rs +++ b/crates/types/src/ir/summary_coverage.rs @@ -1,10 +1,6 @@ -//! What a summary state covers, independent of schema: which rows (source and -//! joint time/population regions) and which columns (the update expression fed -//! into the state and the grouping). +//! Joint time/population coverage for summary composition, independent of schema. //! Equality predicates are a deliberately narrow proof vocabulary. Unsupported //! predicates cannot be declared disjoint merely by giving them different names. -use crate::ir::operator_properties::Reduction; -use crate::post_asap::SummaryUpdate; use crate::pre_asap::Source; use serde::{Deserialize, Serialize}; use std::collections::BTreeMap; @@ -17,15 +13,8 @@ pub struct SummaryCoverage { /// Observation data source, as named by `Scan`: a table or a time series. /// Region time bounds refer to its time column. pub source: Source, - /// Rows: union of joint regions; never the Cartesian product of - /// independent bounds. + /// Union of joint regions; never the Cartesian product of independent bounds. pub regions: Vec, - /// Columns: the update expression fed into the state from each row, equal - /// to the producing `SummaryAgg`'s `input` (e.g. the `latency` column). - pub input: SummaryUpdate, - /// Columns: the grouping, equal to the producing `SummaryAgg`'s - /// `reduction` (e.g. by `job`). Compared by child-schema column index. - pub group_by: Reduction, } #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] @@ -59,10 +48,6 @@ pub enum CoverageError { NotState, #[error("summary node requires coverage")] Missing, - #[error("summary coverage columns differ from the summary's input or grouping")] - ColumnMismatch, - #[error("summary merge requires coverage on every input")] - UnknownInput, } impl SummaryCoverage { @@ -84,18 +69,10 @@ impl SummaryCoverage { Ok(()) } - /// Rejects `other` unless it reads the same columns as `self`. - pub fn check_columns(&self, other: &Self) -> Result<(), CoverageError> { - if self.input != other.input || self.group_by != other.group_by { - return Err(CoverageError::ColumnMismatch); - } - Ok(()) - } - /// Every observation in a region is assumed to contribute once to the state. - /// Compose once-per-observation summaries only when they read the same - /// columns and their joint regions are provably disjoint. Family merge - /// capability and accuracy are not checked here. + /// Compose once-per-observation summaries only when their joint regions are + /// provably disjoint. Update/reduction compatibility, family merge capability + /// and accuracy are checked by `SummaryMerge`, not here. pub fn merge_disjoint(inputs: &[Self]) -> Result { let first = inputs.first().ok_or(CoverageError::EmptyMerge)?; let mut merged = first.clone(); @@ -105,7 +82,6 @@ impl SummaryCoverage { if input.source != first.source { return Err(CoverageError::SourceMismatch); } - input.check_columns(first)?; merged.regions.extend(input.regions.iter().cloned()); } merged.validate()?; diff --git a/crates/types/tests/schema_rebuilding.rs b/crates/types/tests/schema_rebuilding.rs index 18924338f..b0fc4c666 100644 --- a/crates/types/tests/schema_rebuilding.rs +++ b/crates/types/tests/schema_rebuilding.rs @@ -17,8 +17,6 @@ fn coverage() -> SummaryCoverage { time_ms: None, population: Default::default(), }], - input: SummaryUpdate::column(ColumnRef::Named("value".into())), - group_by: Reduction::by(vec![0]), } } /// Rewrites clear coverage; a rewriter must declare it again for summary nodes. diff --git a/crates/types/tests/structure_contract.rs b/crates/types/tests/structure_contract.rs index 04a03969d..d7c144f4e 100644 --- a/crates/types/tests/structure_contract.rs +++ b/crates/types/tests/structure_contract.rs @@ -24,10 +24,6 @@ fn whole_table() -> asap_types::ir::summary_coverage::SummaryCoverage { time_ms: None, population: Default::default(), }], - input: asap_types::post_asap::SummaryUpdate::column( - asap_types::pre_asap::ColumnRef::Named("x".into()), - ), - group_by: asap_types::pre_asap::Reduction::by(vec![]), } } /// Resolved filters cannot hide invalid scalar types or out-of-scope columns. diff --git a/crates/types/tests/summary_coverage.rs b/crates/types/tests/summary_coverage.rs index 8433ab6a7..e36c3d7a4 100644 --- a/crates/types/tests/summary_coverage.rs +++ b/crates/types/tests/summary_coverage.rs @@ -20,8 +20,6 @@ fn coverage(start: i64, end: i64, population: &[(&str, &str)]) -> SummaryCoverag .map(|(k, v)| (k.to_string(), v.to_string())) .collect(), }], - input: SummaryUpdate::column(ColumnRef::Named("latency".into())), - group_by: Reduction::by(vec![]), } } /// Adjacent panes coalesce; gaps remain disconnected rather than becoming a hull. @@ -150,14 +148,3 @@ fn regions_without_time_bounds() { Err(CoverageError::PossibleOverlap) ); } - -/// Disjoint rows do not make coverage of different columns composable. -#[test] -fn different_columns_do_not_compose() { - let mut by_job = coverage(1, 2, &[]); - by_job.group_by = Reduction::by(vec![0]); - assert_eq!( - SummaryCoverage::merge_disjoint(&[coverage(0, 1, &[]), by_job]), - Err(CoverageError::ColumnMismatch) - ); -} diff --git a/crates/types/tests/summary_merge_structure.rs b/crates/types/tests/summary_merge_structure.rs index 3edda31a2..7b7993ff4 100644 --- a/crates/types/tests/summary_merge_structure.rs +++ b/crates/types/tests/summary_merge_structure.rs @@ -1,7 +1,6 @@ //! Window composition merges compatible summary states without consuming raw rows. use asap_types::{ ir::operator_properties::{Reduction, Source}, - ir::summary_coverage::CoverageError, ir::{ASAPOp, NonASAPOp, Operator, OperatorNode, SchemaDerivationError}, post_asap::{GroupingStrategy, SketchAlgorithm, SketchKind, SketchParams, SummaryUpdate}, pre_asap::{ColumnRef, DataType, Field, FieldDataType, Schema}, @@ -15,10 +14,6 @@ fn state(k: u32) -> Rc { } fn state_with(family: FieldDataType) -> Rc { - state_over(family, SummaryUpdate::column(ColumnRef::SampleValue)) -} - -fn state_over(family: FieldDataType, input: SummaryUpdate) -> Rc { let scan = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Scan { source: Source::Table { table_ref: "latencies".into(), @@ -30,7 +25,7 @@ fn state_over(family: FieldDataType, input: SummaryUpdate) -> Rc { let summary = OperatorNode::new(Operator::ASAP(ASAPOp::SummaryAgg { child: scan, family, - input: input.clone(), + input: SummaryUpdate::column(ColumnRef::SampleValue), reduction: Reduction::by(vec![]), grouping: GroupingStrategy::default(), filter: None, @@ -46,8 +41,6 @@ fn state_over(family: FieldDataType, input: SummaryUpdate) -> Rc { time_ms: Some(0..1), population: Default::default(), }], - input, - group_by: Reduction::by(vec![]), }) .unwrap(), ) @@ -87,66 +80,14 @@ fn shifted(state: &Rc, start: i64, end: i64) -> Rc { Rc::new(node) } -fn kll() -> FieldDataType { - FieldDataType::Sketch( - SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k: 200 }), - Default::default(), - ) -} - fn merge(children: Vec>) -> Result, SchemaDerivationError> { OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryMerge { children })) } -/// Equal schemas cannot tell a KLL over one column from a KLL over another; -/// the columns in their coverage can. +/// A merge is a state, so merges nest structurally. #[test] -fn different_coverage_columns_do_not_merge() { - let named = state_over( - kll(), - SummaryUpdate::column(ColumnRef::Named("value".into())), - ); - assert!(matches!( - merge(vec![state(200), shifted(&named, 1, 2)]), - Err(SchemaDerivationError::Coverage( - CoverageError::ColumnMismatch - )) - )); -} - -/// A `SummaryAgg` cannot declare columns other than its own input and grouping. -#[test] -fn declared_columns_must_match_the_summary() { - let mut coverage = state(200).coverage.clone().unwrap(); - coverage.input = SummaryUpdate::column(ColumnRef::Named("value".into())); - let fresh = (*state(200)).clone(); - assert!(matches!( - fresh.clone().with_coverage(coverage.clone()), - Err(SchemaDerivationError::Coverage( - CoverageError::ColumnMismatch - )) - )); - let mut forged = fresh; - forged.coverage = Some(coverage); - assert!(matches!( - Rc::new(forged).validate_structure(), - Err(SchemaDerivationError::Coverage( - CoverageError::ColumnMismatch - )) - )); -} - -/// Every merge input needs coverage. A nested merge has none until #646, so -/// it is rejected for now. -#[test] -fn merge_inputs_without_coverage_are_rejected() { - let mut missing = (*state(200)).clone(); - missing.coverage = None; - let nested = merge(vec![state(200)]).unwrap(); - for input in [Rc::new(missing), nested] { - assert!(matches!( - merge(vec![state(200), input]), - Err(SchemaDerivationError::Coverage(CoverageError::UnknownInput)) - )); - } +fn merges_nest() { + let inner = merge(vec![state(200)]).unwrap(); + let outer = merge(vec![inner, shifted_state(200, 1, 2)]).unwrap(); + outer.validate_structure().unwrap(); }