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
16 changes: 16 additions & 0 deletions crates/asap-physical-operators/src/physical_planner/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -551,6 +551,22 @@ fn compile_internal(
continue;
}
}
if matches!(node.payload, Payload::SummaryMerge) && output.time_index.is_some() {
// Pane timestamps describe their individual builds. A merged
// per-series state represents this evaluation's entire window,
// so merge by series identity and attach the execution scope's
// timestamp after merging, as per-series SummaryAgg does.
let merged = bind_operation(node, &schemas)?;
let compact = merged.schema();
let merge_id = helper_id(id, 1);
physical_dag.add(merge_id, inputs, merged)?;
physical_dag.add(
id,
vec![merge_id],
Operator::scope_timestamp(compact, output)?,
)?;
continue;
}
let mut operator = compile_node(node, &schemas)
.map_err(|error| invalid(format!("node {id}: {error}")))?;
if operator.is_counter_evaluation() {
Expand Down
75 changes: 66 additions & 9 deletions crates/integration-tests/tests/automatic_window_composition.rs
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,16 @@ use std::{collections::BTreeMap, sync::Arc};
/// Generate panes from one five-minute state and find its producer/reader split automatically.
#[test]
fn automatic_panes_and_materialization_preserve_quantile() {
execute_generated_panes(false);
}

/// Per-series PromQL panes preserve temporal metadata through the native wire compiler.
#[test]
fn automatic_per_series_panes_preserve_quantile() {
execute_generated_panes(true);
}

fn execute_generated_panes(per_series: bool) {
let family = FieldDataType::Sketch(
SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k: 200 }),
Default::default(),
Expand All @@ -27,7 +37,17 @@ fn automatic_panes_and_materialization_preserve_quantile() {
metric: "events".into(),
},
predicates: vec![],
schema: Schema::new(vec![Field::plain("value", DataType::Float64, false)]),
schema: if per_series {
Schema::lifted(
vec![
Field::plain("ts", DataType::Timestamp, false),
Field::plain("value", DataType::Float64, false),
],
Some(0),
)
} else {
Schema::new(vec![Field::plain("value", DataType::Float64, false)])
},
}))
.unwrap();
let range = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::TimeRange {
Expand All @@ -40,7 +60,11 @@ fn automatic_panes_and_materialization_preserve_quantile() {
child: range,
family,
input: SummaryUpdate::column(ColumnRef::SampleValue),
reduction: Reduction::by(vec![]),
reduction: if per_series {
Reduction::PerEntity
} else {
Reduction::by(vec![])
},
grouping: GroupingStrategy::default(),
filter: None,
}))
Expand Down Expand Up @@ -92,7 +116,16 @@ fn automatic_panes_and_materialization_preserve_quantile() {
Batch::try_new(
schema,
(0..20)
.map(|i| vec![Value::Float64((pane * 20 + i) as f64)])
.map(|i| {
if per_series {
vec![
Value::Timestamp((pane * 20 + i) as i64 * 1000),
Value::Float64((pane * 20 + i) as f64),
]
} else {
vec![Value::Float64((pane * 20 + i) as f64)]
}
})
.collect(),
)
.unwrap(),
Expand Down Expand Up @@ -161,9 +194,21 @@ fn automatic_panes_and_materialization_preserve_quantile() {
.iter()
.copied()
.zip(stored)
.map(|(id, mut batches)| {
.enumerate()
.map(|(pane, (id, mut batches))| {
assert_eq!(batches.len(), 1);
(id, batches.remove(0))
let mut batch = batches.remove(0);
if per_series {
// Retained panes can have different build timestamps; the
// merged answer must carry the query's evaluation timestamp.
let mut rows = batch.rows().to_vec();
let time = batch.schema().time_index.unwrap();
for row in &mut rows {
row[time] = Value::Timestamp((pane as i64 + 1) * 60_000);
}
batch = Batch::try_new(batch.schema().clone(), rows).unwrap();
}
(id, batch)
})
.collect();
let retained_outputs = physical_common::execute(
Expand All @@ -174,13 +219,25 @@ fn automatic_panes_and_materialization_preserve_quantile() {
revision: 1,
},
);
if per_series {
assert!(matches!(
retained_outputs[0][0].rows()[0][0],
Value::Timestamp(300_000)
));
assert!(matches!(
outputs[0][0].rows()[0][0],
Value::Timestamp(300_000)
));
}
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!(retained_outputs[0][0].rows()[0]
.iter()
.any(|value| matches!(value, 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()[0]
.iter()
.any(|value| matches!(value, Value::Float64(value) if *value == 98.0)),
"{:?}",
outputs[0][0].rows()
);
Expand Down
4 changes: 3 additions & 1 deletion docs/develop_docs/planner-layering-status.md
Original file line number Diff line number Diff line change
Expand Up @@ -122,7 +122,9 @@ Unknown recurrence keeps the original; budget exhaustion returns an error.
producer/reader frontier, including rebuilding all panes at query time. It
retains individual failures and never returns a truncated inventory.
`integration-tests/tests/automatic_window_composition.rs` verifies generated
panes execute identically with and without retained outputs. The deployment
panes execute identically with and without retained outputs. Per-series native
merges ignore individual pane build timestamps and attach the query evaluation
timestamp; the regression covers differently timestamped retained panes. The deployment
still supplies each selector's exact raw window and binds retained outputs to
that window/revision. This does not introduce a rotating pane cache or claim
that cadence alone certifies compatibility with a catalog's pane origin.
Expand Down
Loading