Skip to content
Draft
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
Original file line number Diff line number Diff line change
Expand Up @@ -2635,9 +2635,9 @@ mod tests {
}

#[test]
fn state_only_needs_no_evaluation_and_summary_merge_child_fails_closed() {
fn state_only_needs_no_evaluation_and_summary_subtract_child_fails_closed() {
// A state-only root is costable without evaluation evidence; a
// SummaryAgg over a reserved SummaryMerge is rejected at planning.
// SummaryAgg over a reserved SummarySubtract is rejected at planning.
let workload = streaming_workload();
let target = streaming_sum_query();
let estimated = summary_with_operations(false, false, false);
Expand Down Expand Up @@ -2666,7 +2666,7 @@ mod tests {
let nested = std::rc::Rc::new(
OperatorNode::with_schema(
asap_types::ir::Operator::ASAP(ASAPOp::SummaryAgg {
child: evaluation_state(&summary_with_operations(true, false, false)),
child: evaluation_state(&summary_with_operations(false, true, false)),
family: FieldDataType::ExactAggregate(ExactKind::Count, ExactParams::Count),
input: SummaryUpdate {
item: None,
Expand Down Expand Up @@ -2694,7 +2694,7 @@ mod tests {
assert!(matches!(
streaming_planning_error(nested, &nested_model),
asap_types::post_asap::ExecutionDataStateError::UnimplementedOperator {
operator: "SummaryMerge"
operator: "SummarySubtract"
}
));
}
Expand Down
145 changes: 145 additions & 0 deletions crates/integration-tests/tests/planner_layering_merge.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,145 @@
//! The window-composition example merges KLL panes through the unified IR and native runtime.
mod physical_common;
use asap_physical_operators::{
physical_planner::{compile, cut_candidate, InputContract},
runtime::Scope,
values::{Batch, Value},
};
use asap_types::{
ir::operator_properties::{Reduction, Source},
ir::{ASAPOp, NonASAPOp, Operator, OperatorNode},
post_asap::{
GroupingStrategy, SketchAlgorithm, SketchKind, SketchParams, SketchStatistic, SummaryUpdate,
},
pre_asap::{ColumnRef, DataType, Field, FieldDataType, Schema},
};
use std::{collections::BTreeMap, sync::Arc};

/// Build five one-minute KLLs, merge them, then read p99 through the real wire compiler.
#[test]
fn five_minute_quantile_merges_five_one_minute_states() {
let family = FieldDataType::Sketch(
SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k: 200 }),
Default::default(),
);
let states = (0..5)
.map(|pane| {
let scan = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Scan {
source: Source::Table {
table_ref: format!("pane_{pane}"),
},
predicates: vec![],
schema: Schema::new(vec![Field::plain("value", DataType::Float64, false)]),
}))
.unwrap();
OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryAgg {
child: scan,
family: family.clone(),
input: SummaryUpdate::column(ColumnRef::SampleValue),
reduction: Reduction::by(vec![]),
grouping: GroupingStrategy::default(),
filter: None,
}))
.unwrap()
})
.collect();
let merged =
OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryMerge { children: states }))
.unwrap();
let root = OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryEstimate {
summary_input: merged,
query: SketchStatistic::Quantile { q: 0.99 },
}))
.unwrap();
root.validate_structure().unwrap();
let wire = physical_common::compile_post_asap_dag(&root).unwrap();
wire.validate().unwrap();
let mut inputs = BTreeMap::new();
let mut batches = BTreeMap::new();
for node in &wire.nodes {
if let asap_types::ir::export::PostAsapOperatorPayload::Relational {
operator:
asap_types::ir::export::NonASAPOpKind::Scan {
source: Source::Table { table_ref },
..
},
} = &node.payload
{
let pane: usize = table_ref.strip_prefix("pane_").unwrap().parse().unwrap();
let schema = Arc::new(node.output_schema.clone());
let id = u64::from(node.id.0);
inputs.insert(id, InputContract::bounded(schema.clone()));
batches.insert(
id,
Batch::try_new(
schema,
(0..20)
.map(|i| vec![Value::Float64((pane * 20 + i) as f64)])
.collect(),
)
.unwrap(),
);
}
}
let plan = compile(&wire, inputs, &[u64::from(wire.root.0)]).unwrap();
let outputs = physical_common::execute(
&plan,
batches.clone(),
Scope::Query {
evaluation_time_ms: 300_000,
revision: 1,
},
);
// Pattern B1 stores five built panes; its query reads the same merge DAG
// as Pattern B2, which builds every pane from raw rows on each read.
let frontier: Vec<_> = wire
.nodes
.iter()
.filter_map(|node| {
matches!(
node.payload,
asap_types::ir::export::PostAsapOperatorPayload::SummaryAgg { .. }
)
.then_some(u64::from(node.id.0))
})
.collect();
let materialized = cut_candidate(&plan, &frontier).unwrap();
let precompute = materialized.precompute.as_ref().unwrap();
let stored = physical_common::execute(
precompute,
batches,
Scope::Ingestion {
window_start_ms: 0,
window_end_ms: 300_000,
revision: 1,
},
);
let retained = precompute
.roots()
.iter()
.copied()
.zip(stored)
.map(|(id, mut batches)| {
assert_eq!(batches.len(), 1);
(id, batches.remove(0))
})
.collect();
let retained_outputs = physical_common::execute(
&materialized.query,
retained,
Scope::Query {
evaluation_time_ms: 300_000,
revision: 1,
},
);
assert_eq!(retained_outputs[0][0].rows().len(), 1);
assert!(
matches!(retained_outputs[0][0].rows()[0].as_slice(), [Value::Float64(value)] if *value == 98.0)
);
assert_eq!(outputs[0][0].rows().len(), 1);
assert!(
matches!(outputs[0][0].rows()[0].as_slice(), [Value::Float64(value)] if *value == 98.0),
"{:?}",
outputs[0][0].rows()
);
}
63 changes: 49 additions & 14 deletions crates/types/src/ir/asap.rs
Original file line number Diff line number Diff line change
Expand Up @@ -39,9 +39,7 @@ pub enum ASAPOp {
},
/// Read an exact accumulator's state as its finalized value: the
/// maintenance-to-read boundary before query-time operators.
FinalizeExactAccumulator {
child: Rc<OperatorNode>,
},
FinalizeExactAccumulator { child: Rc<OperatorNode> },
/// Maintain the full declared population, including membership changes.
MaintainPopulation {
child: Rc<OperatorNode>,
Expand All @@ -52,10 +50,9 @@ pub enum ASAPOp {
child: Rc<OperatorNode>,
evaluation: PopulationStatistic,
},
// ── Reserved: migrated but unimplemented (§1.3 of the proposal) ──
SummaryMerge {
children: Vec<Rc<OperatorNode>>,
},
/// Merge compatible partial states for the same grouping and family.
SummaryMerge { children: Vec<Rc<OperatorNode>> },
// ── Reserved: migrated but unimplemented ──
SummarySubtract {
left: Rc<OperatorNode>,
right: Rc<OperatorNode>,
Expand Down Expand Up @@ -186,11 +183,7 @@ impl ASAPOp {
use ASAPOp::*;
matches!(
self,
SummaryMerge { .. }
| SummarySubtract { .. }
| SummaryDelete { .. }
| SummaryJoin { .. }
| Extension { .. }
SummarySubtract { .. } | SummaryDelete { .. } | SummaryJoin { .. } | Extension { .. }
)
}

Expand All @@ -202,6 +195,14 @@ impl ASAPOp {
pub fn produced_state(&self) -> Option<&FieldDataType> {
match self {
ASAPOp::SummaryAgg { family, .. } | ASAPOp::SummaryJoin { family, .. } => Some(family),
ASAPOp::SummaryMerge { children } => children.first().and_then(|child| {
child
.schema
.fields
.iter()
.find(|field| !field.is_plain())
.map(|field| &field.dtype)
}),
_ => None,
}
}
Expand Down Expand Up @@ -405,8 +406,11 @@ impl ASAPOp {
.output_schema()?
}
}
SummaryMerge { .. }
| SummarySubtract { .. }
SummaryMerge { children } => {
self.validate_inputs()?;
children[0].schema.clone()
}
SummarySubtract { .. }
| SummaryDelete { .. }
| SummaryJoin { .. }
| Extension { .. } => return Err(Self::unimplemented()),
Expand Down Expand Up @@ -444,6 +448,37 @@ impl ASAPOp {
}
};
match self {
SummaryMerge { children } => {
let Some(first) = children.first() else {
return Err(SchemaDerivationError::InvalidScalarSignature(
"summary merge requires at least one state input".into(),
));
};
// Matching state parameters and grouping positions are necessary;
// matching names alone cannot prove two states compatible.
if first
.schema
.fields
.iter()
.filter(|field| !field.is_plain())
.count()
!= 1
{
return Err(SchemaDerivationError::InvalidScalarSignature(
"summary merge requires exactly one state column".into(),
));
}
for child in children {
needs_state(child, "SummaryMerge")?;
if child.schema != first.schema {
return Err(SchemaDerivationError::InvalidScalarSignature(
"summary merge inputs must have identical state and grouping schemas"
.into(),
));
}
}
Ok(())
}
SummaryEstimate {
summary_input,
query,
Expand Down
17 changes: 15 additions & 2 deletions crates/types/src/ir/timing.rs
Original file line number Diff line number Diff line change
Expand Up @@ -396,8 +396,21 @@ fn validate_asap(
}
Ok(())
}
ASAPOp::SummaryMerge { .. }
| ASAPOp::SummarySubtract { .. }
ASAPOp::SummaryMerge { children } => {
for child in children {
let state = state_of(child);
if state.primitive != DataPrimitive::SummaryState
|| (timing == ExecutionTiming::IngestionTime && state.timing != timing)
{
return Err(ExecutionDataStateError::IllegalChildDataState {
edge: "SummaryMerge.children",
child: state,
});
}
}
Ok(())
}
ASAPOp::SummarySubtract { .. }
| ASAPOp::SummaryDelete { .. }
| ASAPOp::SummaryJoin { .. }
| ASAPOp::Extension { .. } => Err(ExecutionDataStateError::UnimplementedOperator {
Expand Down
4 changes: 2 additions & 2 deletions crates/types/src/post_asap/execution_data_state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -134,8 +134,8 @@ pub enum ExecutionDataStateError {
/// declared data_state.
#[error("exact operator consumes non-plain column {column:?} ({dtype})")]
NonPlainOperand { column: String, dtype: String },
/// A reserved ASAP operator (`SummaryMerge`, `SummarySubtract`,
/// `SummaryDelete`, `SummaryJoin`, `Extension`) in an executable plan.
/// A reserved ASAP operator (`SummarySubtract`, `SummaryDelete`,
/// `SummaryJoin`, `Extension`) in an executable plan.
#[error("{operator} is a reserved operator with no execution contract yet")]
UnimplementedOperator { operator: &'static str },
/// A node reached by export without a timing: the lifecycle timing pass
Expand Down
72 changes: 72 additions & 0 deletions crates/types/tests/summary_merge.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,72 @@
//! Window composition merges compatible summary states without consuming raw rows.
use asap_types::{
ir::operator_properties::{Reduction, Source},
ir::{ASAPOp, NonASAPOp, Operator, OperatorNode},
post_asap::{
ExecutionTiming, GroupingStrategy, SketchAlgorithm, SketchKind, SketchParams, SummaryUpdate,
},
pre_asap::{ColumnRef, DataType, Field, FieldDataType, Schema},
};
use std::rc::Rc;
fn state(k: u32) -> Rc<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();
OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryAgg {
child: scan,
family: FieldDataType::Sketch(
SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k }),
Default::default(),
),
input: SummaryUpdate::column(ColumnRef::SampleValue),
reduction: Reduction::by(vec![]),
grouping: GroupingStrategy::default(),
filter: None,
}))
.unwrap()
}
/// Two KLL panes compose into one typed state at either execution phase.
#[test]
fn compatible_panes_merge_and_export() {
let root = OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryMerge {
children: vec![state(200), state(200)],
}))
.unwrap();
root.validate_structure().unwrap();
assert_eq!(root.schema.fields.len(), 1);
let timed = asap_types::ir::apply_lifecycle_timings(
&root,
&Default::default(),
&mut Default::default(),
)
.unwrap();
asap_types::ir::export::compile_post_asap_dag(&timed)
.unwrap()
.validate()
.unwrap();
for timing in [ExecutionTiming::IngestionTime, ExecutionTiming::QueryTime] {
asap_types::ir::timing::validate_default(&root, timing).unwrap();
assert_eq!(
asap_types::ir::timing::planned_data_state(&root, timing).primitive,
asap_types::post_asap::DataPrimitive::SummaryState
);
}
}
/// An empty merge, raw rows and differently sized state cannot masquerade as compatible panes.
#[test]
fn incompatible_merge_inputs_fail() {
for children in [
vec![],
vec![state(200), state(300)],
vec![state(200).children()[0].clone()],
] {
assert!(
OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryMerge { children })).is_err()
);
}
}
Loading
Loading