From 4b4abd556b7f67d36b45dac3a714cc7fc5f81337 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Sun, 4 Oct 2026 16:33:24 +0000 Subject: [PATCH 1/4] fix(executor): build per-entity summaries over a TimeShift pane A tumbling pane (#580) reads TimeRange(w) over TimeShift(i*w) over Scan. The per-entity build path required the Scan directly under the TimeRange and rejected panes with "per-entity summary requires a resolved source". The deployment supplies the shifted raw rows, so look through one TimeShift to find the source schema. Ported from #566. Co-Authored-By: Claude Opus 5.5 --- crates/executor/src/physical_planner/mod.rs | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/crates/executor/src/physical_planner/mod.rs b/crates/executor/src/physical_planner/mod.rs index d307a6a5..b102585e 100644 --- a/crates/executor/src/physical_planner/mod.rs +++ b/crates/executor/src/physical_planner/mod.rs @@ -404,7 +404,13 @@ fn compile_internal( "per-entity summary requires a resolved raw time range", )); }; - let Some(NonASAPOp::Scan { schema, .. }) = child.non_asap() else { + // A tumbling pane reads its window through a TimeShift (#580); + // the deployment supplies the shifted raw rows. + let source = match child.non_asap() { + Some(NonASAPOp::TimeShift { child, .. }) => child, + _ => child, + }; + let Some(NonASAPOp::Scan { schema, .. }) = source.non_asap() else { return Err(invalid("per-entity summary requires a resolved source")); }; if !schema.closed || update.item.is_some() { From dc8ea53a579b810b18b07c323ef3fefd29305fb1 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Sun, 4 Oct 2026 16:33:25 +0000 Subject: [PATCH 2/4] fix(executor): keep the evaluation timestamp on temporal pane merges A per-series SummaryMerge grouped away the time column, so its native output disagreed with the Planner schema ("native output type differs from Planner output"). Merge by series identity, then attach the run scope's timestamp with scope_timestamp, as per-series SummaryAgg does. Pane build timestamps never reach the merged state. Ported from #569. Regression tests merge 5 one-minute panes for an exact Sum and a KLL quantile, rebuilt at query time and retained with distinct build timestamps, and compare them with one 5-minute build. Co-Authored-By: Claude Opus 5.5 --- crates/executor/src/physical_planner/mod.rs | 16 + crates/executor/tests/tumbling_pane_merge.rs | 389 +++++++++++++++++++ 2 files changed, 405 insertions(+) create mode 100644 crates/executor/tests/tumbling_pane_merge.rs diff --git a/crates/executor/src/physical_planner/mod.rs b/crates/executor/src/physical_planner/mod.rs index b102585e..f652712a 100644 --- a/crates/executor/src/physical_planner/mod.rs +++ b/crates/executor/src/physical_planner/mod.rs @@ -555,6 +555,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/executor/tests/tumbling_pane_merge.rs b/crates/executor/tests/tumbling_pane_merge.rs new file mode 100644 index 00000000..d5d9000f --- /dev/null +++ b/crates/executor/tests/tumbling_pane_merge.rs @@ -0,0 +1,389 @@ +//! Tumbling panes merged by `SummaryMerge` answer the same per-series query as +//! one build over the whole window (#509 Example 3/4, Pattern B; #580). +use asap_executor::{ + operators::Operator as PhysicalOperator, + physical_planner::{ + compile, promql_rows::encode_series_identity, CompiledPhysicalDAG, InputContract, Source, + }, + runtime::{Limits, RunContext, Scope}, + values::{Batch, Value}, +}; +use asap_logical_optimizer::search_workload; +use asap_plan_selection::candidate_selection::global_selection; +use asap_plan_selection::cost::cost_model::DefaultCostModel; +use futures::{executor::block_on, StreamExt}; +use planner_types::ir::export::PhysicalASAPDAG; +use planner_types::ir::operator::operator_properties::TimeShift; +use planner_types::ir::physical_export::compile_physical_asap_dag_with_node_ids; +use planner_types::ir::properties::summary_coverage::{CoverageRegion, SummaryCoverage}; +use planner_types::ir::{properties::ExecutionTiming, ASAPOp, NonASAPOp, Operator, OperatorNode}; +use planner_types::types::AccuracyTarget; +use planner_types::workload::*; +use std::{collections::BTreeMap, rc::Rc, sync::Arc}; + +const EVALUATION_MS: i64 = 300_000; +const PANE_MS: i64 = 60_000; +const PANES: i64 = 5; + +/// The selected single-build plan for a 5-minute per-series PromQL query. +fn single_build(query: &str, accuracy: AccuracyTarget) -> Rc { + let workload = PlanningWorkload { + query_workload: QueryWorkload { + language: QueryLanguage::PromQL, + query_batch: Some(vec![BatchEntry { + query: Query(query.into()), + requirements: QueryRequirements { + accuracy: AccuracyRequirement::Explicit(accuracy), + ..Default::default() + }, + predictability: Predictability::Unknown, + invocations: 1, + execute_at: None, + time_selection: TimeSelection::default(), + }]), + repeating_queries: None, + }, + data_workload: Some(DataWorkload { + data_ingestion_interval: Evidence { + value: Some(DurationMs(1000)), + ..Default::default() + }, + ..Default::default() + }), + }; + let root = asap_frontend_promql::lower_promql_workload(&workload, 0) + .unwrap() + .remove(0); + let root = asap_executor::physical_planner::promql_rows::with_series_identity(&root).unwrap(); + let space = search_workload(vec![("q", root)]); + global_selection(&space, &DefaultCostModel) + .assemble_selected_query(&space.roots[0].1) + .unwrap() + .unwrap() +} + +/// Rewrite the root's per-entity `SummaryAgg` over `TimeRange(5m)` into the +/// shape #580 emits: pane `i` is `TimeRange(1m)` over `TimeShift(i·1m)` over +/// the same `Scan`, and a `SummaryMerge` combines the panes. Coverage regions +/// are absolute, ending at the evaluation time. +fn tumbling_panes(root: &Rc) -> Rc { + let [agg] = root.operator.children()[..] else { + panic!("expected one summary input"); + }; + let Some(NonASAPOp::TimeRange { kind, child, .. }) = agg.operator.children()[0].non_asap() + else { + panic!("expected a raw time range under the summary"); + }; + let source = agg.coverage.as_ref().unwrap().source.clone(); + let panes = (0..PANES) + .map(|i| { + let shifted = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::TimeShift { + shift: TimeShift { + offset_ms: i * PANE_MS, + at: None, + }, + child: child.clone(), + })) + .unwrap(); + let range = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::TimeRange { + range: std::time::Duration::from_millis(PANE_MS as u64), + kind: *kind, + child: shifted, + })) + .unwrap(); + let end = EVALUATION_MS - i * PANE_MS; + Rc::new( + agg.map_children(|_| range.clone()) + .unwrap() + .with_guarantee(agg.guarantee.clone()) + .with_coverage(SummaryCoverage { + source: source.clone(), + regions: vec![CoverageRegion { + time_ms: Some(end - PANE_MS..end), + population: BTreeMap::new(), + }], + }) + .unwrap(), + ) + }) + .collect(); + let merged = + OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryMerge { children: panes })).unwrap(); + let merged = Rc::new((*merged).clone().with_guarantee(agg.guarantee.clone())); + let root = Rc::new( + root.map_children(|_| merged.clone()) + .unwrap() + .with_guarantee(root.guarantee.clone()), + ); + root.validate_structure().unwrap(); + root +} + +/// Two series, four samples per minute in [0, 5m): row `(series, minute, j)`. +fn samples() -> Vec<(&'static str, i64, f64)> { + let mut rows = vec![]; + for (s, series) in ["a", "b"].into_iter().enumerate() { + for minute in 0..PANES { + for j in 0..4 { + let ts = minute * PANE_MS + j * 15_000; + rows.push((series, ts, (s * 100) as f64 + (minute * 4 + j) as f64)); + } + } + } + rows +} + +struct Exported { + dag: PhysicalASAPDAG, + /// Each raw `TimeRange` input with the pane it reads (0 = most recent). + ranges: Vec<(u64, i64)>, + /// Each `SummaryAgg` node with its pane. + builds: Vec<(u64, i64)>, +} + +/// Place every node at query time. `apply_materialization_timings` does not +/// yet accept `SummaryMerge` (it reports `UnimplementedOperator`), so this test +/// assigns the timing a planner would choose for Pattern B2 directly. +fn at_query_time(node: &Rc) -> Rc { + let mut timed = (**node).clone(); + timed.operator = node.operator.map_children(at_query_time); + timed.timing = Some(ExecutionTiming::QueryTime); + Rc::new(timed) +} + +fn export(root: &Rc) -> Exported { + let root = at_query_time(root); + let compiled = compile_physical_asap_dag_with_node_ids(&root).unwrap(); + let pane = |node: &OperatorNode| -> i64 { + match node.operator.children()[0].non_asap() { + Some(NonASAPOp::TimeShift { shift, .. }) => shift.offset_ms / PANE_MS, + _ => 0, + } + }; + let (mut ranges, mut builds) = (vec![], vec![]); + for node in &compiled.dag.nodes { + let logical = compiled.node_ids.operator_node(node.id).unwrap(); + match &logical.operator { + Operator::NonASAP(NonASAPOp::TimeRange { .. }) => { + ranges.push((u64::from(node.id.0), pane(logical))) + } + Operator::ASAP(ASAPOp::SummaryAgg { child, .. }) => { + builds.push((u64::from(node.id.0), pane(child))) + } + _ => {} + } + } + Exported { + dag: compiled.dag, + ranges, + builds, + } +} + +fn schema_of(dag: &PhysicalASAPDAG, id: u64) -> Arc { + Arc::new( + dag.nodes + .iter() + .find(|node| u64::from(node.id.0) == id) + .unwrap() + .output_schema + .clone(), + ) +} + +/// Raw rows for each `TimeRange` input. A deployment supplies pane `i` with +/// the samples in `[T - (i+1)·w, T - i·w)`; a single build reads all samples. +fn raw_inputs(exported: &Exported) -> BTreeMap { + exported + .ranges + .iter() + .map(|&(id, pane)| { + let schema = schema_of(&exported.dag, id); + let end = EVALUATION_MS - pane * PANE_MS; + let start = if exported.ranges.len() == 1 { + 0 + } else { + end - PANE_MS + }; + let rows = samples() + .into_iter() + .filter(|(_, ts, _)| (start..end).contains(ts)) + .map(|(series, ts, value)| { + schema + .fields + .iter() + .map(|field| match field.name.as_str() { + "ts" => Value::Timestamp(ts), + "value" => Value::Float64(value), + _ => Value::Utf8( + encode_series_identity(&BTreeMap::from([ + ("__name__".to_string(), "m".to_string()), + ("series".to_string(), series.to_string()), + ])) + .unwrap() + .into(), + ), + }) + .collect() + }) + .collect(); + (id, Batch::try_new(schema, rows).unwrap()) + }) + .collect() +} + +fn run(plan: &CompiledPhysicalDAG, inputs: BTreeMap, scope: Scope) -> Vec> { + let sources = inputs + .into_iter() + .map(|(id, batch)| { + let source = PhysicalOperator::source(batch.schema().clone(), vec![batch]).unwrap(); + (id, Box::new(source) as Source<'_>) + }) + .collect(); + let dag = plan.instantiate(sources).unwrap(); + let mut rows = block_on(async { + let context = RunContext::new(scope, Limits::default()).unwrap(); + let mut output = dag.execute(plan.roots(), context).unwrap().remove(0); + let mut rows = vec![]; + while let Some(batch) = output.next().await { + rows.extend(batch.unwrap().rows().iter().cloned()); + } + rows + }); + rows.sort_by_key(|row| format!("{row:?}")); + rows +} + +/// `Value` has no equality; its debug form identifies every scalar exactly. +fn render(rows: &[Vec]) -> Vec { + rows.iter().map(|row| format!("{row:?}")).collect() +} + +fn has_timestamp(row: &[Value], ms: i64) -> bool { + row.iter() + .any(|v| matches!(v, Value::Timestamp(t) if *t == ms)) +} + +fn query_scope() -> Scope { + Scope::Query { + evaluation_time_ms: EVALUATION_MS, + revision: 1, + } +} + +/// Compile and run `root` with every raw time range as a deployment input. +fn execute(root: &Rc) -> Vec> { + let exported = export(root); + let inputs = raw_inputs(&exported); + let contracts = inputs + .iter() + .map(|(&id, batch)| (id, InputContract::bounded(batch.schema().clone()))) + .collect(); + let plan = compile( + &exported.dag, + contracts, + &[u64::from(exported.dag.roots[0].0)], + ) + .unwrap(); + run(&plan, inputs, query_scope()) +} + +/// Build panes in an ingestion run, retain them with distinct build +/// timestamps, then merge them in the query run (Pattern B1). +fn execute_retained(root: &Rc) -> Vec> { + let exported = export(root); + let inputs = raw_inputs(&exported); + let contracts: BTreeMap<_, _> = inputs + .iter() + .map(|(&id, batch)| (id, InputContract::bounded(batch.schema().clone()))) + .collect(); + let mut retained = BTreeMap::new(); + for &(build, pane) in &exported.builds { + let range = exported.ranges.iter().find(|r| r.1 == pane).unwrap().0; + let plan = compile( + &exported.dag, + BTreeMap::from([(range, contracts[&range].clone())]), + &[build], + ) + .unwrap(); + let end = EVALUATION_MS - pane * PANE_MS; + let rows = run( + &plan, + BTreeMap::from([(range, inputs[&range].clone())]), + Scope::Ingestion { + window_start_ms: end - PANE_MS, + window_end_ms: end, + revision: 1, + }, + ); + assert!(rows.iter().all(|row| has_timestamp(row, end))); + retained.insert( + build, + Batch::try_new(schema_of(&exported.dag, build), rows).unwrap(), + ); + } + let plan = compile( + &exported.dag, + retained + .iter() + .map(|(&id, batch)| (id, InputContract::bounded(batch.schema().clone()))) + .collect(), + &[u64::from(exported.dag.roots[0].0)], + ) + .unwrap(); + run(&plan, retained, query_scope()) +} + +fn values(rows: &[Vec]) -> Vec { + rows.iter() + .flat_map(|row| { + row.iter().filter_map(|v| match v { + Value::Float64(v) => Some(*v), + _ => None, + }) + }) + .collect() +} + +/// Returns the per-series values, which every execution path agrees on. +fn assert_pane_merge_matches_single_build(query: &str, accuracy: AccuracyTarget) -> Vec { + let single = single_build(query, accuracy); + let panes = tumbling_panes(&single); + let expected = execute(&single); + assert_eq!(expected.len(), 2, "{expected:?}"); + // The single build already carries the evaluation timestamp. + assert!(expected.iter().all(|row| has_timestamp(row, EVALUATION_MS))); + let result = values(&expected); + let expected = render(&expected); + assert_eq!( + render(&execute(&panes)), + expected, + "panes rebuilt at query time" + ); + assert_eq!( + render(&execute_retained(&panes)), + expected, + "retained panes" + ); + result +} + +/// Exact per-series Sum over 5 merged 1-minute panes equals one 5-minute build, +/// timestamped at the evaluation time. +#[test] +fn exact_sum_over_tumbling_panes_matches_single_build() { + let sums = + assert_pane_merge_matches_single_build("sum_over_time(m[5m])", AccuracyTarget::Exact); + assert_eq!(sums, vec![190., 2190.]); +} + +/// A KLL quantile over 5 merged 1-minute panes equals one 5-minute build. The +/// 20 samples per series are below k, so both states are exact. +#[test] +fn kll_quantile_over_tumbling_panes_matches_single_build() { + let p99 = assert_pane_merge_matches_single_build( + "quantile_over_time(0.99, m[5m])", + AccuracyTarget::Epsilon(0.05), + ); + assert_eq!(p99, vec![19., 119.]); +} From aec031f99ede3f7e961b4009108a9056bc7c6c7e Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Sun, 4 Oct 2026 18:07:13 +0000 Subject: [PATCH 3/4] fix: integrate with #592 Pane coverage time is a CoverageTime since #592; the tumbling pane-merge test keeps absolute time ranges. Co-Authored-By: Claude Opus 5.5 --- crates/executor/tests/tumbling_pane_merge.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/crates/executor/tests/tumbling_pane_merge.rs b/crates/executor/tests/tumbling_pane_merge.rs index d5d9000f..357d7057 100644 --- a/crates/executor/tests/tumbling_pane_merge.rs +++ b/crates/executor/tests/tumbling_pane_merge.rs @@ -99,7 +99,7 @@ fn tumbling_panes(root: &Rc) -> Rc { .with_coverage(SummaryCoverage { source: source.clone(), regions: vec![CoverageRegion { - time_ms: Some(end - PANE_MS..end), + time_ms: Some((end - PANE_MS..end).into()), population: BTreeMap::new(), }], }) From 4dce4aa4aa5438810d4de58116a16cbc7b1031e6 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Thu, 8 Oct 2026 14:33:44 +0000 Subject: [PATCH 4/4] test(executor): tumbling panes derive their coverage from TimeRange over TimeShift Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01W7qG9aFyPij5uWsyAJCxDW --- crates/executor/src/physical_planner/mod.rs | 4 ++- crates/executor/tests/tumbling_pane_merge.rs | 38 ++++++-------------- 2 files changed, 14 insertions(+), 28 deletions(-) diff --git a/crates/executor/src/physical_planner/mod.rs b/crates/executor/src/physical_planner/mod.rs index f652712a..906fb737 100644 --- a/crates/executor/src/physical_planner/mod.rs +++ b/crates/executor/src/physical_planner/mod.rs @@ -555,7 +555,9 @@ fn compile_internal( continue; } } - if matches!(node.payload, Payload::SummaryMerge) && output.time_index.is_some() { + if matches!(node.payload, Payload::ASAP(ASAPOp::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 diff --git a/crates/executor/tests/tumbling_pane_merge.rs b/crates/executor/tests/tumbling_pane_merge.rs index 357d7057..24fd5052 100644 --- a/crates/executor/tests/tumbling_pane_merge.rs +++ b/crates/executor/tests/tumbling_pane_merge.rs @@ -12,10 +12,9 @@ use asap_logical_optimizer::search_workload; use asap_plan_selection::candidate_selection::global_selection; use asap_plan_selection::cost::cost_model::DefaultCostModel; use futures::{executor::block_on, StreamExt}; -use planner_types::ir::export::PhysicalASAPDAG; use planner_types::ir::operator::operator_properties::TimeShift; use planner_types::ir::physical_export::compile_physical_asap_dag_with_node_ids; -use planner_types::ir::properties::summary_coverage::{CoverageRegion, SummaryCoverage}; +use planner_types::ir::physical_export::PhysicalASAPDAG; use planner_types::ir::{properties::ExecutionTiming, ASAPOp, NonASAPOp, Operator, OperatorNode}; use planner_types::types::AccuracyTarget; use planner_types::workload::*; @@ -64,8 +63,8 @@ fn single_build(query: &str, accuracy: AccuracyTarget) -> Rc { /// Rewrite the root's per-entity `SummaryAgg` over `TimeRange(5m)` into the /// shape #580 emits: pane `i` is `TimeRange(1m)` over `TimeShift(i·1m)` over -/// the same `Scan`, and a `SummaryMerge` combines the panes. Coverage regions -/// are absolute, ending at the evaluation time. +/// the same `Scan`, and a `SummaryMerge` combines the panes. Each pane's +/// coverage is derived from its `TimeRange` over `TimeShift`. fn tumbling_panes(root: &Rc) -> Rc { let [agg] = root.operator.children()[..] else { panic!("expected one summary input"); @@ -74,7 +73,6 @@ fn tumbling_panes(root: &Rc) -> Rc { else { panic!("expected a raw time range under the summary"); }; - let source = agg.coverage.as_ref().unwrap().source.clone(); let panes = (0..PANES) .map(|i| { let shifted = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::TimeShift { @@ -91,19 +89,10 @@ fn tumbling_panes(root: &Rc) -> Rc { child: shifted, })) .unwrap(); - let end = EVALUATION_MS - i * PANE_MS; Rc::new( - agg.map_children(|_| range.clone()) + agg.with_new_children(|_| range.clone()) .unwrap() - .with_guarantee(agg.guarantee.clone()) - .with_coverage(SummaryCoverage { - source: source.clone(), - regions: vec![CoverageRegion { - time_ms: Some((end - PANE_MS..end).into()), - population: BTreeMap::new(), - }], - }) - .unwrap(), + .with_guarantee(agg.guarantee.clone()), ) }) .collect(); @@ -111,7 +100,7 @@ fn tumbling_panes(root: &Rc) -> Rc { OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryMerge { children: panes })).unwrap(); let merged = Rc::new((*merged).clone().with_guarantee(agg.guarantee.clone())); let root = Rc::new( - root.map_children(|_| merged.clone()) + root.with_new_children(|_| merged.clone()) .unwrap() .with_guarantee(root.guarantee.clone()), ); @@ -165,10 +154,10 @@ fn export(root: &Rc) -> Exported { let logical = compiled.node_ids.operator_node(node.id).unwrap(); match &logical.operator { Operator::NonASAP(NonASAPOp::TimeRange { .. }) => { - ranges.push((u64::from(node.id.0), pane(logical))) + ranges.push((node.id as u64, pane(logical))) } Operator::ASAP(ASAPOp::SummaryAgg { child, .. }) => { - builds.push((u64::from(node.id.0), pane(child))) + builds.push((node.id as u64, pane(child))) } _ => {} } @@ -184,7 +173,7 @@ fn schema_of(dag: &PhysicalASAPDAG, id: u64) -> Arc) -> Vec> { .iter() .map(|(&id, batch)| (id, InputContract::bounded(batch.schema().clone()))) .collect(); - let plan = compile( - &exported.dag, - contracts, - &[u64::from(exported.dag.roots[0].0)], - ) - .unwrap(); + let plan = compile(&exported.dag, contracts, &[exported.dag.roots[0] as u64]).unwrap(); run(&plan, inputs, query_scope()) } @@ -328,7 +312,7 @@ fn execute_retained(root: &Rc) -> Vec> { .iter() .map(|(&id, batch)| (id, InputContract::bounded(batch.schema().clone()))) .collect(), - &[u64::from(exported.dag.roots[0].0)], + &[exported.dag.roots[0] as u64], ) .unwrap(); run(&plan, retained, query_scope())