diff --git a/crates/asap-physical-operators/src/physical_planner/mod.rs b/crates/asap-physical-operators/src/physical_planner/mod.rs index 67ae02afe..ce358a5f5 100644 --- a/crates/asap-physical-operators/src/physical_planner/mod.rs +++ b/crates/asap-physical-operators/src/physical_planner/mod.rs @@ -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() { diff --git a/crates/integration-tests/tests/automatic_window_composition.rs b/crates/integration-tests/tests/automatic_window_composition.rs index d7fad8f1a..514fbbe5e 100644 --- a/crates/integration-tests/tests/automatic_window_composition.rs +++ b/crates/integration-tests/tests/automatic_window_composition.rs @@ -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(), @@ -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 { @@ -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, })) @@ -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(), @@ -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( @@ -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() ); diff --git a/docs/develop_docs/planner-layering-status.md b/docs/develop_docs/planner-layering-status.md index 02b61ef3c..c99db97a8 100644 --- a/docs/develop_docs/planner-layering-status.md +++ b/docs/develop_docs/planner-layering-status.md @@ -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.