Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
58 changes: 50 additions & 8 deletions crates/types/src/ir/asap.rs
Original file line number Diff line number Diff line change
Expand Up @@ -54,10 +54,11 @@ pub enum ASAPOp<C = Rc<OperatorNode>> {
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<C>,
},
// ── Reserved: migrated but unimplemented ──
SummarySubtract {
left: C,
right: C,
Expand Down Expand Up @@ -192,11 +193,7 @@ impl ASAPOp {
use ASAPOp::*;
matches!(
self,
SummaryMerge { .. }
| SummarySubtract { .. }
| SummaryDelete { .. }
| SummaryJoin { .. }
| Extension { .. }
SummarySubtract { .. } | SummaryDelete { .. } | SummaryJoin { .. } | Extension { .. }
)
}

Expand All @@ -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,
}
}
Expand Down Expand Up @@ -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()),
Expand Down Expand Up @@ -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,
Expand Down
93 changes: 93 additions & 0 deletions crates/types/tests/summary_merge_structure.rs
Original file line number Diff line number Diff line change
@@ -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<OperatorNode> {
state_with(FieldDataType::Sketch(
SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k }),
Default::default(),
))
}

fn state_with(family: FieldDataType) -> Rc<OperatorNode> {
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<OperatorNode> {
shifted(&state(k), start, end)
}

fn shifted(state: &Rc<OperatorNode>, start: i64, end: i64) -> Rc<OperatorNode> {
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<Rc<OperatorNode>>) -> Result<Rc<OperatorNode>, 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();
}
Loading