diff --git a/crates/types/src/ir/asap.rs b/crates/types/src/ir/asap.rs index 90b23946..16945d92 100644 --- a/crates/types/src/ir/asap.rs +++ b/crates/types/src/ir/asap.rs @@ -54,10 +54,11 @@ pub enum ASAPOp> { child: C, evaluation: PopulationStatistic, }, - // ── Reserved: migrated but unimplemented (§1.3 of the proposal) ── + /// Merge compatible partial states for the same grouping and family. SummaryMerge { children: Vec, }, + // ── Reserved: migrated but unimplemented ── SummarySubtract { left: C, right: C, @@ -192,11 +193,7 @@ impl ASAPOp { use ASAPOp::*; matches!( self, - SummaryMerge { .. } - | SummarySubtract { .. } - | SummaryDelete { .. } - | SummaryJoin { .. } - | Extension { .. } + SummarySubtract { .. } | SummaryDelete { .. } | SummaryJoin { .. } | Extension { .. } ) } @@ -208,6 +205,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 +416,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 +458,40 @@ 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(), + )); + } + } + // Equal schemas cannot tell a KLL over `latency` from one over + // `size`, nor prove the inputs disjoint; summary coverage (#646) + // decides whether a structurally valid merge is semantically valid. + 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 00000000..7b7993ff --- /dev/null +++ b/crates/types/tests/summary_merge_structure.rs @@ -0,0 +1,93 @@ +//! Window composition merges compatible summary states without consuming raw rows. +use asap_types::{ + ir::operator_properties::{Reduction, Source}, + ir::{ASAPOp, NonASAPOp, Operator, OperatorNode, SchemaDerivationError}, + post_asap::{GroupingStrategy, SketchAlgorithm, SketchKind, SketchParams, SummaryUpdate}, + pre_asap::{ColumnRef, DataType, Field, FieldDataType, Schema}, +}; +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(), + }, + predicates: vec![], + schema: Schema::new(vec![Field::plain("value", DataType::Float64, false)]), + })) + .unwrap(); + let summary = OperatorNode::new(Operator::ASAP(ASAPOp::SummaryAgg { + child: scan, + family, + input: SummaryUpdate::column(ColumnRef::SampleValue), + reduction: Reduction::by(vec![]), + grouping: GroupingStrategy::default(), + filter: None, + })) + .unwrap(); + std::rc::Rc::new( + summary + .with_coverage(asap_types::ir::summary_coverage::SummaryCoverage { + source: Source::Table { + table_ref: "latencies".into(), + }, + regions: vec![asap_types::ir::summary_coverage::CoverageRegion { + time_ms: Some(0..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), shifted_state(200, 1, 2)], + })) + .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() + ); + } +} + +fn shifted_state(k: u32, start: i64, end: i64) -> Rc { + 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) +} + +fn merge(children: Vec>) -> Result, SchemaDerivationError> { + OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryMerge { children })) +} + +/// A merge is a state, so merges nest structurally. +#[test] +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(); +}