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
24 changes: 23 additions & 1 deletion crates/executor/src/physical_planner/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -398,7 +398,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() {
Expand Down Expand Up @@ -545,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
Loading