From bca6641bb4d5af702bfcaab6a913cdd5c8aab9a1 Mon Sep 17 00:00:00 2001 From: zzylol Date: Tue, 29 Sep 2026 22:58:56 +0000 Subject: [PATCH 1/8] feat: derive physical candidates from one compilation `compile_candidate` re-lowered the whole Post-ASAP DAG for every materialization frontier, and `enumerate_frontiers` compiled it once more. Number helper operators from their Planner node (`u64::MAX - node_id`, at most one helper per node) so every boundary choice is a subgraph of one lowering. `cut_candidate` partitions a `compile` result for one frontier and `enumerate_compiled_frontiers` enumerates over it; both produce candidates byte-identical to per-frontier recompilation. Co-Authored-By: Claude Opus 5.5 --- .../src/physical_planner/candidates.rs | 152 ++++++++++++--- .../src/physical_planner/compiled.rs | 56 +++++- .../src/physical_planner/mod.rs | 20 +- .../tests/precompute_candidates.rs | 183 +++++++++++++++++- 4 files changed, 366 insertions(+), 45 deletions(-) diff --git a/crates/asap-physical-operators/src/physical_planner/candidates.rs b/crates/asap-physical-operators/src/physical_planner/candidates.rs index a8f476f3..ebab0d31 100644 --- a/crates/asap-physical-operators/src/physical_planner/candidates.rs +++ b/crates/asap-physical-operators/src/physical_planner/candidates.rs @@ -25,24 +25,42 @@ pub fn compile_candidate( inputs: BTreeMap, roots: &[NodeId], frontier: &[NodeId], +) -> Result { + cut_candidate(&compile(dag, inputs, roots)?, frontier) +} + +/// Derive one frontier's candidate from a complete [`compile`] result by +/// partitioning its operators; nothing is lowered again. A deployment compiles +/// each query DAG once and derives every placement choice from that result. +/// The candidate is identical to [`compile_candidate`] for the same frontier. +pub fn cut_candidate( + compiled: &CompiledPhysicalDag, + frontier: &[NodeId], ) -> Result { if frontier.is_empty() { return Ok(PhysicalCandidate { precompute: None, - query: compile(dag, inputs, roots)?, + query: compiled.clone(), materialized_outputs: BTreeMap::new(), }); } let frontier_set: BTreeSet<_> = frontier.iter().copied().collect(); - if frontier_set.len() != frontier.len() || frontier.iter().any(|id| inputs.contains_key(id)) { + // `compile` retains only reachable nodes and numbers its helper operators + // above the u32 Planner ID range; only Planner outputs are boundaries. + if frontier_set.len() != frontier.len() + || frontier + .iter() + .any(|&id| !compiled.is_operator(id) || u32::try_from(id).is_err()) + { return Err(invalid("frontier must contain distinct computed outputs")); } - let full = compile(dag, inputs.clone(), roots)?; - let precompute = compile(dag, inputs.clone(), frontier)?; + let inputs: BTreeMap<_, _> = compiled + .input_contracts() + .map(|(id, contract)| (id, contract.clone())) + .collect(); + let precompute = compiled.cut(&inputs, frontier)?; let mut materialized_outputs = BTreeMap::new(); for &id in frontier { - // Also proves that the frontier is reachable from the requested roots. - full.output_contract(id)?; let mut output = precompute.output_contract(id)?; if output.properties.boundedness != Boundedness::Bounded { return Err(invalid("materialized output requires bounded execution")); @@ -54,7 +72,7 @@ pub fn compile_candidate( } let mut query_inputs = inputs; query_inputs.extend(materialized_outputs.clone()); - let query = compile(dag, query_inputs, roots)?; + let query = compiled.cut(&query_inputs, compiled.roots())?; let used: BTreeSet<_> = query.input_contracts().map(|(id, _)| id).collect(); if !frontier.iter().all(|id| used.contains(id)) { return Err(invalid( @@ -78,43 +96,40 @@ pub fn enumerate_frontiers( inputs: &BTreeMap, roots: &[NodeId], max_candidates: usize, +) -> Result>, Error> { + enumerate_compiled_frontiers(&compile(dag, inputs.clone(), roots)?, max_candidates) +} + +/// [`enumerate_frontiers`] over an existing [`compile`] result, so enumeration +/// and [`cut_candidate`] share one lowering. +pub fn enumerate_compiled_frontiers( + compiled: &CompiledPhysicalDag, + max_candidates: usize, ) -> Result>, Error> { if max_candidates == 0 { return Err(invalid( "frontier search requires a positive candidate budget", )); } - let compiled = compile(dag, inputs.clone(), roots)?; let mut ancestors = BTreeMap::>::new(); let mut eligible = Vec::new(); - for node in &dag.nodes { - let id = u64::from(node.id.0); - if inputs.contains_key(&id) { - continue; - } - let Ok(contract) = compiled.output_contract(id) else { - continue; - }; - if contract.properties.boundedness != Boundedness::Bounded { + for (id, properties) in compiled.output_properties()? { + if !compiled.is_operator(id) + || u32::try_from(id).is_err() + || properties.boundedness != Boundedness::Bounded + { continue; } let mut seen = BTreeSet::new(); let mut pending = vec![id]; while let Some(current) = pending.pop() { - if !seen.insert(current) || inputs.contains_key(¤t) { - continue; + if seen.insert(current) { + pending.extend(compiled.dependencies(current)); } - pending.extend( - dag.edges - .iter() - .filter(|edge| u64::from(edge.consumer.0) == current) - .map(|edge| u64::from(edge.producer.0)), - ); } ancestors.insert(id, seen); eligible.push(id); } - eligible.sort_unstable(); let mut frontiers = vec![vec![]]; for id in eligible { let additions = frontiers @@ -142,16 +157,20 @@ pub fn enumerate_frontiers( /// Lower every maintenance candidate before feasibility/cost evaluation. Keep /// individual failures visible; do not substitute another computation on error. +/// The DAG is lowered once; each frontier is a [`cut_candidate`] of it. pub fn compile_candidates( dag: &PostAsapDag, inputs: BTreeMap, roots: &[NodeId], frontiers: &[Vec], ) -> Vec> { - frontiers - .iter() - .map(|frontier| compile_candidate(dag, inputs.clone(), roots, frontier)) - .collect() + match compile(dag, inputs, roots) { + Ok(compiled) => frontiers + .iter() + .map(|frontier| cut_candidate(&compiled, frontier)) + .collect(), + Err(error) => frontiers.iter().map(|_| Err(error.clone())).collect(), + } } /// Complete workload cost supplied by scoped optimizer/deployment evidence. @@ -273,3 +292,76 @@ impl PhysicalCandidate { Ok(()) } } + +#[cfg(test)] +mod tests { + use super::*; + use planner_types::workload::*; + + fn grouped_rate() -> (PostAsapDag, BTreeMap, NodeId) { + let workload = PlanningWorkload { + query_workload: QueryWorkload { + language: QueryLanguage::PromQL, + query_batch: Some(vec![BatchEntry { + query: Query("sum by(job)(rate(m[1m]))".into()), + requirements: QueryRequirements { + accuracy: AccuracyRequirement::Explicit( + planner_types::types::AccuracyTarget::Exact, + ), + ..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 = std::rc::Rc::new(promql_rows::with_series_identity(&root).unwrap()); + let space = asap_aware_mapping::search_workload(vec![("q", root)]); + let selected = space + .global_selection(&asap_aware_mapping::cost_model::DefaultCostModel) + .assemble_selected_dag(&space.roots[0].1) + .unwrap() + .unwrap(); + let dag = planner_types::post_asap::compile_post_asap_dag(&selected).unwrap(); + let state = dag + .nodes + .iter() + .find(|node| matches!(node.payload, Payload::SummaryAgg { .. })) + .unwrap(); + let inputs = BTreeMap::from([( + u64::from(state.id.0), + InputContract::bounded(Arc::new(state.output_schema.clone())), + )]); + (dag.clone(), inputs, u64::from(dag.root.0)) + } + + /// Enumerating and cutting every frontier lowers each Planner node once. + #[test] + fn candidates_for_all_frontiers_share_one_lowering() { + let (dag, inputs, root) = grouped_rate(); + let lowered = || crate::physical_planner::LOWERED_NODES.with(|count| count.get()); + let before = lowered(); + let compiled = compile(&dag, inputs, &[root]).unwrap(); + let once = lowered() - before; + let frontiers = enumerate_compiled_frontiers(&compiled, 4096).unwrap(); + assert!(frontiers.len() >= 3, "{frontiers:?}"); + for frontier in &frontiers { + cut_candidate(&compiled, frontier).unwrap(); + } + assert!(once > 0); + assert_eq!(lowered() - before, once); + } +} diff --git a/crates/asap-physical-operators/src/physical_planner/compiled.rs b/crates/asap-physical-operators/src/physical_planner/compiled.rs index 70d6a9ff..af0bbe39 100644 --- a/crates/asap-physical-operators/src/physical_planner/compiled.rs +++ b/crates/asap-physical-operators/src/physical_planner/compiled.rs @@ -212,13 +212,8 @@ impl CompiledPhysicalDag { } /// Derive a reachable output contract without opening deployment readers. pub fn output_contract(&self, id: NodeId) -> Result { - let sources = self - .input_contracts() - .map(|(id, contract)| (id, Box::new(contract.clone()) as Source<'_>)) - .collect(); - let graph = self.instantiate(sources)?; - let properties = graph.properties(&self.roots)?; - let properties = *properties + let properties = *self + .output_properties()? .get(&id) .ok_or_else(|| invalid("output is not reachable"))?; let schema = match self @@ -231,6 +226,53 @@ impl CompiledPhysicalDag { }; Ok(InputContract { schema, properties }) } + /// Properties of every reachable node, derived in one contract-only pass. + pub(super) fn output_properties(&self) -> Result, Error> { + let sources = self + .input_contracts() + .map(|(id, contract)| (id, Box::new(contract.clone()) as Source<'_>)) + .collect(); + self.instantiate(sources)?.properties(&self.roots) + } + /// Direct physical dependencies; empty for inputs and unknown IDs. + pub(super) fn dependencies(&self, id: NodeId) -> &[NodeId] { + match self.nodes.get(&id) { + Some(Node::Operator { inputs, .. }) => inputs, + _ => &[], + } + } + pub(super) fn is_operator(&self, id: NodeId) -> bool { + matches!(self.nodes.get(&id), Some(Node::Operator { .. })) + } + /// Keep the already-lowered operators reachable from `roots`, replacing + /// each node in `boundaries` by a typed input. Nothing is lowered again. + pub(super) fn cut( + &self, + boundaries: &BTreeMap, + roots: &[NodeId], + ) -> Result { + let mut result = Self::new(roots.to_vec()); + let mut pending = roots.to_vec(); + while let Some(id) = pending.pop() { + if result.nodes.contains_key(&id) { + continue; + } + let node = match boundaries.get(&id) { + Some(contract) => Node::Input(contract.clone()), + None => self + .nodes + .get(&id) + .cloned() + .ok_or_else(|| invalid(format!("missing physical node {id}")))?, + }; + if let Node::Operator { inputs, .. } = &node { + pending.extend(inputs); + } + result.nodes.insert(id, node); + } + result.validate()?; + Ok(result) + } /// Validate using contract-only sources. No deployment reader is available. pub fn validate(&self) -> Result<(), Error> { let sources = self diff --git a/crates/asap-physical-operators/src/physical_planner/mod.rs b/crates/asap-physical-operators/src/physical_planner/mod.rs index 142a2c25..213dfa8b 100644 --- a/crates/asap-physical-operators/src/physical_planner/mod.rs +++ b/crates/asap-physical-operators/src/physical_planner/mod.rs @@ -36,8 +36,8 @@ pub mod promql_values; mod candidates; pub use candidates::{ - compile_candidate, compile_candidates, enumerate_frontiers, select_candidate, CandidateCost, - CandidateSelection, PhysicalCandidate, + compile_candidate, compile_candidates, cut_candidate, enumerate_compiled_frontiers, + enumerate_frontiers, select_candidate, CandidateCost, CandidateSelection, PhysicalCandidate, }; mod temporal_panes; @@ -109,6 +109,12 @@ pub fn bind_with_data_sources<'a>( bind(dag, sources, roots) } +#[cfg(test)] +thread_local! { + /// Planner nodes lowered by this thread, for compile-once tests. + static LOWERED_NODES: std::cell::Cell = const { std::cell::Cell::new(0) }; +} + fn compile_internal( dag: &PostAsapDag, mut sources: BTreeMap, @@ -165,9 +171,12 @@ fn compile_internal( } } let mut graph = CompiledPhysicalDag::new(roots.to_vec()); - let mut auxiliary = u64::MAX; for id in ordered { let node = nodes[&id]; + // At most one helper operator per node, numbered above the u32 Planner + // ID range by its node alone, so every boundary choice yields a subgraph + // of the same lowering and candidate cuts need not renumber operators. + let auxiliary = u64::MAX - id; let output = Arc::new(node.output_schema.clone()); crate::values::validate_schema(&output)?; if let Some(source) = sources.remove(&id) { @@ -176,6 +185,8 @@ fn compile_internal( } graph.add_input(id, source)?; } else { + #[cfg(test)] + LOWERED_NODES.with(|count| count.set(count.get() + 1)); let mut inputs = dependencies.get(&id).cloned().unwrap_or_default(); let mut schemas = inputs .iter() @@ -191,7 +202,6 @@ fn compile_internal( Operator::union(schemas[0].clone(), schemas.len())?, )?; inputs = vec![auxiliary]; - auxiliary -= 1; schemas.truncate(1); } if let Payload::Value { @@ -286,7 +296,6 @@ fn compile_internal( vec![auxiliary], Operator::limit(input, *k as u64, 0, groups)?.with_output_schema(output)?, )?; - auxiliary -= 1; continue; } // A closed row must include either all source labels or the explicit @@ -346,7 +355,6 @@ fn compile_internal( vec![auxiliary], Operator::scope_timestamp(compact, output)?, )?; - auxiliary -= 1; continue; } let mut operator = compile_node(node, &schemas) diff --git a/crates/asap-physical-operators/tests/precompute_candidates.rs b/crates/asap-physical-operators/tests/precompute_candidates.rs index c00b388f..f402e3dc 100644 --- a/crates/asap-physical-operators/tests/precompute_candidates.rs +++ b/crates/asap-physical-operators/tests/precompute_candidates.rs @@ -4,8 +4,9 @@ use asap_physical_operators::{ factory::create_planner_accumulator, operators::Operator, physical_planner::{ - compile_candidates, select_candidate, CandidateCost, CompiledPhysicalDag, InputContract, - Source, + compile, compile_candidate, compile_candidates, cut_candidate, + enumerate_compiled_frontiers, enumerate_frontiers, select_candidate, CandidateCost, + CompiledPhysicalDag, InputContract, PhysicalCandidate, Source, }, runtime::{Limits, RunContext, Scope}, values::{Batch, Value}, @@ -546,3 +547,181 @@ fn enumerated_grouped_rate_candidates_execute_numeric_query_outputs() { "must execute both stored and query-time grouped Rate candidates: {executed}" ); } + +/// The per-frontier lowering used before compile-once cuts: each boundary +/// choice lowers the precompute and query DAGs from the logical DAG again. +fn recompiled_candidate( + dag: &PostAsapDag, + inputs: &BTreeMap, + roots: &[u64], + frontier: &[u64], +) -> Result { + use asap_physical_operators::plan::Emission; + if frontier.is_empty() { + return Ok(PhysicalCandidate { + precompute: None, + query: compile(dag, inputs.clone(), roots)?, + materialized_outputs: BTreeMap::new(), + }); + } + let precompute = compile(dag, inputs.clone(), frontier)?; + let mut materialized_outputs = BTreeMap::new(); + for &id in frontier { + let mut output = precompute.output_contract(id)?; + output.properties.emission = Emission::Unknown; + materialized_outputs.insert(id, output); + } + let mut query_inputs = inputs.clone(); + query_inputs.extend(materialized_outputs.clone()); + Ok(PhysicalCandidate { + precompute: Some(precompute), + query: compile(dag, query_inputs, roots)?, + materialized_outputs, + }) +} + +fn assert_cuts_match_recompilation( + dag: &PostAsapDag, + inputs: BTreeMap, + roots: &[u64], + min_frontiers: usize, +) { + let compiled = compile(dag, inputs.clone(), roots).unwrap(); + let frontiers = enumerate_compiled_frontiers(&compiled, 4096).unwrap(); + assert_eq!( + frontiers, + enumerate_frontiers(dag, &inputs, roots, 4096).unwrap() + ); + assert!(frontiers.len() >= min_frontiers, "{frontiers:?}"); + for frontier in &frontiers { + let cut = cut_candidate(&compiled, frontier).unwrap(); + let expected = recompiled_candidate(dag, &inputs, roots, frontier).unwrap(); + assert_eq!( + cut.encode().unwrap(), + expected.encode().unwrap(), + "{frontier:?}" + ); + } +} + +/// Every enumerated grouped Rate→Sum frontier (query-only, stored Rate, +/// stored Sum) cuts to exactly the candidate that per-frontier lowering builds. +#[test] +fn grouped_rate_cuts_equal_per_frontier_compilation() { + let dag = grouped_rate(); + let state = dag + .nodes + .iter() + .find(|node| matches!(node.payload, PostAsapOperatorPayload::SummaryAgg { .. })) + .unwrap(); + let inputs = BTreeMap::from([( + u64::from(state.id.0), + InputContract::bounded(Arc::new(state.output_schema.clone())), + )]); + assert_cuts_match_recompilation(&dag, inputs, &[u64::from(dag.root.0)], 3); +} + +/// Cuts of a DAG whose nodes lower to helper operators (current-series +/// population read by Sort→Limit) keep the same operator IDs as recompilation. +#[test] +fn population_topk_cuts_equal_per_frontier_compilation() { + let workload = PlanningWorkload { + query_workload: QueryWorkload { + language: QueryLanguage::PromQL, + query_batch: Some(vec![BatchEntry { + query: Query("topk by(job)(1, m)".into()), + requirements: QueryRequirements { + accuracy: AccuracyRequirement::Explicit(AccuracyTarget::Exact), + ..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(60_000)), + ..Default::default() + }, + ..Default::default() + }), + }; + let original = asap_frontend_promql::lower_promql_workload(&workload, 0) + .unwrap() + .remove(0); + let root = Rc::new( + asap_physical_operators::physical_planner::promql_rows::with_series_identity(&original) + .unwrap(), + ); + let selected = asap_aware_mapping::maintained_population::MaintainedPopulationStrategy::new( + std::slice::from_ref(&root), + ) + .candidate(&root) + .unwrap(); + let dag = compile_post_asap_dag(&selected).unwrap(); + let raw = dag + .nodes + .iter() + .find(|node| matches!(node.payload, PostAsapOperatorPayload::Fallback { .. })) + .unwrap(); + let inputs = BTreeMap::from([( + u64::from(raw.id.0), + InputContract::bounded(Arc::new(raw.output_schema.clone())), + )]); + let roots = [u64::from(dag.root.0)]; + let compiled = compile(&dag, inputs.clone(), &roots).unwrap(); + // The root reads its population through a Sort helper numbered by the root. + let helper = u64::MAX - roots[0]; + assert_eq!(compiled.operator_name(helper), Some("Sort")); + assert!( + cut_candidate(&compiled, &[helper]).is_err(), + "helper operators are not Planner boundaries" + ); + assert_cuts_match_recompilation(&dag, inputs, &roots, 2); +} + +/// Cuts reject frontiers that recompilation rejects: duplicates, inputs, +/// unknown IDs, and an output shadowed by its descendant. +#[test] +fn cut_candidate_rejects_invalid_frontiers() { + let dag = grouped_rate(); + let state = dag + .nodes + .iter() + .find(|node| matches!(node.payload, PostAsapOperatorPayload::SummaryAgg { .. })) + .unwrap(); + let readout = dag + .nodes + .iter() + .find(|node| { + matches!( + node.payload, + PostAsapOperatorPayload::Value { + operation: ValueOperation::FinalizeExactAccumulator + } + ) + }) + .unwrap(); + let (state_id, rate_id, root) = ( + u64::from(state.id.0), + u64::from(readout.id.0), + u64::from(dag.root.0), + ); + let inputs = BTreeMap::from([( + state_id, + InputContract::bounded(Arc::new(state.output_schema.clone())), + )]); + let compiled = compile(&dag, inputs.clone(), &[root]).unwrap(); + for frontier in [ + vec![rate_id, rate_id], + vec![state_id], + vec![999], + vec![root, rate_id], + ] { + assert!(cut_candidate(&compiled, &frontier).is_err(), "{frontier:?}"); + assert!(compile_candidate(&dag, inputs.clone(), &[root], &frontier).is_err()); + } +} From 753c704cd5bb86f59d6bf7120c025e5d1b819891 Mon Sep 17 00:00:00 2001 From: zzylol Date: Tue, 29 Sep 2026 22:58:56 +0000 Subject: [PATCH 2/8] docs: describe compile-once candidate cuts Co-Authored-By: Claude Opus 5.5 --- .../physical-planning-and-deployment.md | 11 +++++++++ docs/develop_docs/library-api.md | 24 +++++++++++++++++++ 2 files changed, 35 insertions(+) diff --git a/docs/design_docs/physical-planning-and-deployment.md b/docs/design_docs/physical-planning-and-deployment.md index 1d2b11e8..44c9c812 100644 --- a/docs/design_docs/physical-planning-and-deployment.md +++ b/docs/design_docs/physical-planning-and-deployment.md @@ -408,6 +408,17 @@ feasibility is rejected before pricing. The optimizer supplies candidate frontiers and cost evidence, including updates, retention, recurrence and sharing. `enumerate_frontiers` constructs bounded, reachable antichain frontiers above explicit input boundaries, including query-only and fully precomputed results. It fails explicitly when the candidate budget is exceeded. Maintenance selection must still reject frontiers that violate window, freshness, or reuse requirements; deployment feasibility is checked before pricing. +Lowering a node does not depend on the chosen frontier, so Planner lowers each +query DAG once. `compile(dag, inputs, roots)` yields the complete +`CompiledPhysicalDag`; `enumerate_compiled_frontiers(&compiled, max)` and +`cut_candidate(&compiled, frontier)` then derive each placement by partitioning +its operators: nodes above the frontier and below the roots form the query DAG, +and the frontier's ancestors form the precompute DAG. Helper operators are +numbered by their Planner node, so a cut is byte-identical to lowering that +frontier directly. `compile_candidate(s)` and `enumerate_frontiers` are wrappers +over this path. A deployment compiles each query DAG once, not once per +placement choice. Temporal pane candidates remain a separate lowering. + Physical compilation opens no readers. Bounded precompute outputs become typed query inputs. Their source, filters, grouping, build window, evaluation time, readiness and revision contracts must accompany the selected lifecycle and be checked during diff --git a/docs/develop_docs/library-api.md b/docs/develop_docs/library-api.md index 61be5e65..3d9652a7 100644 --- a/docs/develop_docs/library-api.md +++ b/docs/develop_docs/library-api.md @@ -631,6 +631,30 @@ an alternative with `MissingCostEvidence` is accepted only when the cost model's complete-candidate hook covers lifecycle costs. Window frameworks and totals come from that hook, as in Planner selection. +A lifecycle choice then fixes each physical placement: for example, a +continuously maintained state places its producer in precompute, while an +ephemeral one keeps it in the query. Compile each query's `PostAsapDag` once +and derive every placement from that result: + +```rust +use asap_physical_operators::physical_planner::{ + compile, cut_candidate, enumerate_compiled_frontiers, +}; + +let compiled = compile(&dag, inputs, &roots)?; // each node lowered once +for frontier in enumerate_compiled_frontiers(&compiled, 4096)? { + // Precompute/query DAGs split at `frontier`; no logical lowering. + let candidate = cut_candidate(&compiled, &frontier)?; + // Check feasibility and price `candidate`; bind the selected one as is. +} +``` + +`cut_candidate` returns exactly the `PhysicalCandidate` that +`compile_candidate(&dag, inputs, &roots, &frontier)` returns, and rejects the +same invalid frontiers. The mapping from lifecycle choice to frontier stays with +the caller. Temporal pane candidates are a different lowering and still use +`compile_temporal_pane_candidate`. + ## Optional whole-plan selection and DAG assembly ### What does global selection mean? From ad858fae0b20eb91dae26e724853bb572b4a63ab Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 02:14:40 +0000 Subject: [PATCH 3/8] test: compare cut candidates by their serialized form The #462 split no longer exposes PhysicalCandidate::encode; its serde form gives the same byte-for-byte comparison. Co-Authored-By: Claude Opus 5.5 --- crates/asap-physical-operators/tests/precompute_candidates.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/crates/asap-physical-operators/tests/precompute_candidates.rs b/crates/asap-physical-operators/tests/precompute_candidates.rs index f402e3dc..f0d63488 100644 --- a/crates/asap-physical-operators/tests/precompute_candidates.rs +++ b/crates/asap-physical-operators/tests/precompute_candidates.rs @@ -597,8 +597,8 @@ fn assert_cuts_match_recompilation( let cut = cut_candidate(&compiled, frontier).unwrap(); let expected = recompiled_candidate(dag, &inputs, roots, frontier).unwrap(); assert_eq!( - cut.encode().unwrap(), - expected.encode().unwrap(), + serde_json::to_vec(&cut).unwrap(), + serde_json::to_vec(&expected).unwrap(), "{frontier:?}" ); } From e6903724b709fc78e250dc1012fd07e674d1c8ac Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 02:47:53 +0000 Subject: [PATCH 4/8] feat: derive the physical frontier from lifecycle timing The lifecycle layer assigns each node's timing; physical compilation now reads it. `frontier_from_timing` returns the ingestion-time nodes read by query-time nodes (or an ingestion-time root) and rejects a query-time node feeding an ingestion-time one, so each lifecycle assignment is a `cut_candidate` of one `compile` result. `enumerate_compiled_frontiers` is private: placement comes from timing, and its only caller is `enumerate_frontiers`. Co-Authored-By: Claude Opus 5.5 --- .../src/physical_planner/candidates.rs | 122 +++++++++++++++++- .../src/physical_planner/mod.rs | 4 +- .../tests/precompute_candidates.rs | 12 +- 3 files changed, 125 insertions(+), 13 deletions(-) diff --git a/crates/asap-physical-operators/src/physical_planner/candidates.rs b/crates/asap-physical-operators/src/physical_planner/candidates.rs index ebab0d31..f5658a9a 100644 --- a/crates/asap-physical-operators/src/physical_planner/candidates.rs +++ b/crates/asap-physical-operators/src/physical_planner/candidates.rs @@ -86,6 +86,41 @@ pub fn cut_candidate( }) } +/// Materialization frontier implied by lifecycle-assigned timing: ingestion-time +/// nodes read by a query-time node, plus the root when it is ingestion-timed. +/// `cut_candidate` of one [`compile`] result with this frontier realizes the +/// assignment, so different assignments are different cuts of one lowering. +/// That holds while timing-dependent lowering (an ingestion-time `Binary` +/// aligns by value column) has the same timing at compile time as here. +/// A query-time node feeding an ingestion-time node has no valid placement. +pub fn frontier_from_timing(dag: &PostAsapDag) -> Result, Error> { + use planner_types::post_asap::ExecutionTiming::IngestionTime; + let timing = dag + .nodes + .iter() + .map(|node| (node.id, node.output_state.timing)) + .collect::>(); + let mut frontier = BTreeSet::new(); + if timing.get(&dag.root) == Some(&IngestionTime) { + frontier.insert(u64::from(dag.root.0)); + } + for edge in &dag.edges { + let (Some(&producer), Some(&consumer)) = + (timing.get(&edge.producer), timing.get(&edge.consumer)) + else { + return Err(invalid("timed DAG edge names an unknown node")); + }; + match (producer == IngestionTime, consumer == IngestionTime) { + (true, false) => { + frontier.insert(u64::from(edge.producer.0)); + } + (false, true) => return Err(invalid("query-time node feeds an ingestion-time node")), + _ => {} + } + } + Ok(frontier.into_iter().collect()) +} + /// Enumerate bounded, reachable materialization frontiers above explicit inputs. /// Each frontier is an antichain: storing an output and its ancestor together /// would leave the ancestor unused by query execution. Lifecycle eligibility @@ -100,9 +135,7 @@ pub fn enumerate_frontiers( enumerate_compiled_frontiers(&compile(dag, inputs.clone(), roots)?, max_candidates) } -/// [`enumerate_frontiers`] over an existing [`compile`] result, so enumeration -/// and [`cut_candidate`] share one lowering. -pub fn enumerate_compiled_frontiers( +fn enumerate_compiled_frontiers( compiled: &CompiledPhysicalDag, max_candidates: usize, ) -> Result>, Error> { @@ -364,4 +397,87 @@ mod tests { assert!(once > 0); assert_eq!(lowered() - before, once); } + + fn with_timing( + dag: &PostAsapDag, + timing: impl Fn(&PostAsapDagNode) -> planner_types::post_asap::ExecutionTiming, + ) -> PostAsapDag { + let mut timed = dag.clone(); + for node in &mut timed.nodes { + node.output_state.timing = timing(node); + } + for edge in &mut timed.edges { + let producer = timed.nodes.iter().find(|node| node.id == edge.producer); + edge.data_state = producer.unwrap().output_state; + } + timed + } + + fn raw_input(dag: &PostAsapDag) -> BTreeMap { + let raw = dag + .nodes + .iter() + .find(|node| matches!(node.payload, Payload::Fallback { .. })) + .unwrap(); + BTreeMap::from([( + u64::from(raw.id.0), + InputContract::bounded(Arc::new(raw.output_schema.clone())), + )]) + } + + /// Cutting one compilation by a retained-state timing and by the all + /// query-time timing (what ContinuouslyMaintained and Ephemeral assign) + /// lowers each Planner node once and matches `compile_candidate`. + #[test] + fn timing_cuts_share_one_lowering() { + use planner_types::post_asap::ExecutionTiming::QueryTime; + let (retained, _, root) = grouped_rate(); + let ephemeral = with_timing(&retained, |_| QueryTime); + let inputs = raw_input(&retained); + let lowered = || crate::physical_planner::LOWERED_NODES.with(|count| count.get()); + let before = lowered(); + let compiled = compile(&ephemeral, inputs.clone(), &[root]).unwrap(); + let once = lowered() - before; + let cuts = [&retained, &ephemeral].map(|timed| { + let frontier = frontier_from_timing(timed).unwrap(); + let cut = cut_candidate(&compiled, &frontier).unwrap(); + (timed, frontier, cut) + }); + assert!(once > 0); + assert_eq!(lowered() - before, once); + assert_eq!(cuts[0].1.len(), 1); + assert!(cuts[1].1.is_empty()); + for (timed, frontier, cut) in cuts { + let expected = compile_candidate(timed, inputs.clone(), &[root], &frontier).unwrap(); + assert_eq!( + serde_json::to_vec(&cut).unwrap(), + serde_json::to_vec(&expected).unwrap() + ); + } + } + + /// The frontier is the ingestion-time nodes read at query time; an + /// ingestion-time root is itself the frontier. + #[test] + fn frontier_from_timing_includes_ingestion_root() { + use planner_types::post_asap::ExecutionTiming::IngestionTime; + let (dag, _, root) = grouped_rate(); + let timed = with_timing(&dag, |_| IngestionTime); + assert_eq!(frontier_from_timing(&timed).unwrap(), [root]); + } + + /// A query-time node feeding an ingestion-time node is rejected. + #[test] + fn frontier_from_timing_rejects_query_time_input_to_ingestion() { + use planner_types::post_asap::ExecutionTiming::{IngestionTime, QueryTime}; + let (dag, _, _) = grouped_rate(); + let timed = with_timing(&dag, |node| { + if node.id == dag.root { + IngestionTime + } else { + QueryTime + } + }); + assert!(frontier_from_timing(&timed).is_err()); + } } diff --git a/crates/asap-physical-operators/src/physical_planner/mod.rs b/crates/asap-physical-operators/src/physical_planner/mod.rs index 213dfa8b..6a6ce940 100644 --- a/crates/asap-physical-operators/src/physical_planner/mod.rs +++ b/crates/asap-physical-operators/src/physical_planner/mod.rs @@ -36,8 +36,8 @@ pub mod promql_values; mod candidates; pub use candidates::{ - compile_candidate, compile_candidates, cut_candidate, enumerate_compiled_frontiers, - enumerate_frontiers, select_candidate, CandidateCost, CandidateSelection, PhysicalCandidate, + compile_candidate, compile_candidates, cut_candidate, enumerate_frontiers, + frontier_from_timing, select_candidate, CandidateCost, CandidateSelection, PhysicalCandidate, }; mod temporal_panes; diff --git a/crates/asap-physical-operators/tests/precompute_candidates.rs b/crates/asap-physical-operators/tests/precompute_candidates.rs index f0d63488..65345c76 100644 --- a/crates/asap-physical-operators/tests/precompute_candidates.rs +++ b/crates/asap-physical-operators/tests/precompute_candidates.rs @@ -4,9 +4,9 @@ use asap_physical_operators::{ factory::create_planner_accumulator, operators::Operator, physical_planner::{ - compile, compile_candidate, compile_candidates, cut_candidate, - enumerate_compiled_frontiers, enumerate_frontiers, select_candidate, CandidateCost, - CompiledPhysicalDag, InputContract, PhysicalCandidate, Source, + compile, compile_candidate, compile_candidates, cut_candidate, enumerate_frontiers, + select_candidate, CandidateCost, CompiledPhysicalDag, InputContract, PhysicalCandidate, + Source, }, runtime::{Limits, RunContext, Scope}, values::{Batch, Value}, @@ -587,11 +587,7 @@ fn assert_cuts_match_recompilation( min_frontiers: usize, ) { let compiled = compile(dag, inputs.clone(), roots).unwrap(); - let frontiers = enumerate_compiled_frontiers(&compiled, 4096).unwrap(); - assert_eq!( - frontiers, - enumerate_frontiers(dag, &inputs, roots, 4096).unwrap() - ); + let frontiers = enumerate_frontiers(dag, &inputs, roots, 4096).unwrap(); assert!(frontiers.len() >= min_frontiers, "{frontiers:?}"); for frontier in &frontiers { let cut = cut_candidate(&compiled, frontier).unwrap(); From 68a069e3ffab36c5cbb23cb478606196431cd57c Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 02:47:53 +0000 Subject: [PATCH 5/8] test: cut one compilation by chosen lifecycle timings For the KLL quantile and grouped Rate->Sum fixtures, ContinuouslyMaintained and Ephemeral timed DAGs cut one compilation into exactly the candidates `compile_candidate` builds. The hand-written timing frontier in the chosen lifecycle test now uses `frontier_from_timing`. Co-Authored-By: Claude Opus 5.5 --- .../summary_maintenance_lifecycle_e2e.rs | 186 +++++++++++------- 1 file changed, 119 insertions(+), 67 deletions(-) diff --git a/crates/integration-tests/tests/summary_maintenance_lifecycle_e2e.rs b/crates/integration-tests/tests/summary_maintenance_lifecycle_e2e.rs index 281ea4e8..959304a9 100644 --- a/crates/integration-tests/tests/summary_maintenance_lifecycle_e2e.rs +++ b/crates/integration-tests/tests/summary_maintenance_lifecycle_e2e.rs @@ -212,6 +212,15 @@ fn selected_plan_with_horizon( .into_iter() .next() .expect("one normalized workload entry"); + selected_plan_for_lowered(workload, lowered, model, horizon) +} + +fn selected_plan_for_lowered( + workload: &PlanningWorkload, + lowered: asap_types::pre_asap::QueryExpr, + model: &dyn CostModel, + horizon: Horizon, +) -> asap_aware_mapping::SummaryMaintenanceLifecyclePlan { let root = Rc::new(lowered); let strategies = asap_aware_mapping::default_strategies_with(model); let space = search_workload_with(vec![("dashboard", Rc::clone(&root))], &strategies); @@ -877,32 +886,70 @@ fn quantile_workload(query: &str) -> PlanningWorkload { workload } -/// Precompute outputs implied by timing: ingestion-time nodes read by a -/// query-time node, or the root when it is itself ingestion-timed. -fn ingestion_frontier(dag: &asap_types::post_asap::PostAsapDag) -> Vec { - use asap_types::post_asap::ExecutionTiming::IngestionTime; - let timing = |id| { - dag.nodes - .iter() - .find(|node| node.id == id) - .unwrap() - .output_state - .timing - }; - let mut frontier: Vec<_> = dag +/// Timed DAG for `query` after binding every summary state to `lifecycle`. +/// Grouped queries carry a physical series identity, as per-entity state needs. +fn lifecycle_timed_dag( + query: &str, + lifecycle: &SummaryMaintenanceLifecycle, +) -> (asap_types::post_asap::PostAsapDag, Vec) { + use asap_aware_mapping::enumerate_summary_maintenance_lifecycles; + let workload = quantile_workload(query); + let mut lowered = lower_promql_workload(&workload, 0).unwrap().remove(0); + if query.contains(" by(") { + lowered = + asap_physical_operators::physical_planner::promql_rows::with_series_identity(&lowered) + .unwrap(); + } + let root = + selected_plan_for_lowered(&workload, lowered, &FullyCostedRuntime, Horizon(100.)).root; + let candidates = enumerate_summary_maintenance_lifecycles( + root, + WorkloadDemand::new_with_data( + &workload.query_workload, + workload.data_workload.as_ref().unwrap(), + &[1], + ), + NOW_MS, + Some(Horizon(100.)), + SummaryMaintenanceLifecycleCapabilities::ALL, + &FullyCostedRuntime, + ) + .unwrap(); + let choices: Vec<_> = candidates + .deployments() + .iter() + .map(|deployment| (deployment.post_asap_node_id, lifecycle.clone())) + .collect(); + let mut states: Vec<_> = choices.iter().map(|(id, _)| u64::from(id.0)).collect(); + states.sort_unstable(); + let dag = candidates + .select(&choices) + .unwrap() + .execution_timed_dag() + .unwrap(); + (dag, states) +} + +/// Compile inputs for a timed DAG: its raw source, available at either phase. +fn raw_inputs( + dag: &asap_types::post_asap::PostAsapDag, +) -> std::collections::BTreeMap { + let raw = dag .nodes .iter() - .filter(|node| { - node.output_state.timing == IngestionTime - && (node.id == dag.root - || dag.edges.iter().any(|edge| { - edge.producer == node.id && timing(edge.consumer) != IngestionTime - })) + .find(|node| { + matches!( + node.payload, + asap_types::post_asap::PostAsapOperatorPayload::Fallback { .. } + ) }) - .map(|node| u64::from(node.id.0)) - .collect(); - frontier.sort_unstable(); - frontier + .unwrap(); + std::collections::BTreeMap::from([( + u64::from(raw.id.0), + asap_physical_operators::physical_planner::InputContract::bounded(std::sync::Arc::new( + raw.output_schema.clone(), + )), + )]) } /// For existing PromQL fixtures, Planner's own retained lifecycle selection @@ -935,62 +982,29 @@ fn planner_lifecycle_selection_reproduces_strategy_timing() { /// precompute empty and reads the raw source at query time; both answer alike. #[test] fn chosen_lifecycle_timing_decides_precompute_contents() { - use asap_aware_mapping::enumerate_summary_maintenance_lifecycles; use asap_physical_operators::{ - physical_planner::{compile_candidate, InputContract}, + physical_planner::{compile_candidate, frontier_from_timing}, runtime::Scope, values::{Batch, Value}, }; - use asap_types::{ - post_asap::{PostAsapOperatorPayload, SummaryFamilyType}, - pre_asap::DataType, - }; - use std::{collections::BTreeMap, sync::Arc}; + use asap_types::{post_asap::SummaryFamilyType, pre_asap::DataType}; + use std::collections::BTreeMap; - let workload = quantile_workload("quantile(0.99, latency)"); - let root = selected_plan(&workload).root; let mut answers = Vec::new(); for lifecycle in [ SummaryMaintenanceLifecycle::ContinuouslyMaintained, SummaryMaintenanceLifecycle::Ephemeral, ] { - let candidates = enumerate_summary_maintenance_lifecycles( - Rc::clone(&root), - WorkloadDemand::new_with_data( - &workload.query_workload, - workload.data_workload.as_ref().unwrap(), - &[1], - ), - NOW_MS, - Some(Horizon(100.)), - SummaryMaintenanceLifecycleCapabilities::ALL, - &FullyCostedRuntime, - ) - .unwrap(); - let [deployment] = candidates.deployments() else { + let (dag, states) = lifecycle_timed_dag("quantile(0.99, latency)", &lifecycle); + let [state] = states[..] else { panic!("one summary state"); }; - let id = deployment.post_asap_node_id; - let state = u64::from(id.0); - let dag = candidates - .select(&[(id, lifecycle.clone())]) - .unwrap() - .execution_timed_dag() - .unwrap(); - let raw = dag - .nodes - .iter() - .find(|node| matches!(node.payload, PostAsapOperatorPayload::Fallback { .. })) - .unwrap(); - let (raw_id, schema) = (u64::from(raw.id.0), Arc::new(raw.output_schema.clone())); - let frontier = ingestion_frontier(&dag); - let candidate = compile_candidate( - &dag, - BTreeMap::from([(raw_id, InputContract::bounded(schema.clone()))]), - &[u64::from(dag.root.0)], - &frontier, - ) - .unwrap(); + let inputs = raw_inputs(&dag); + let (&raw_id, contract) = inputs.iter().next().unwrap(); + let schema = contract.schema.clone(); + let frontier = frontier_from_timing(&dag).unwrap(); + let candidate = + compile_candidate(&dag, inputs, &[u64::from(dag.root.0)], &frontier).unwrap(); let rows = (1..=100) .map(|value| { schema @@ -1059,3 +1073,41 @@ fn chosen_lifecycle_timing_decides_precompute_contents() { assert_eq!(answers[0], answers[1]); assert_eq!(answers[0].len(), 1); } + +/// One compilation, cut by each lifecycle assignment's timing, yields exactly +/// the candidate `compile_candidate` builds for that timed DAG: the retained +/// state is the frontier under ContinuouslyMaintained, and nothing under +/// Ephemeral. Covers the KLL quantile fixture and grouped Rate→Sum. +#[test] +fn lifecycle_timing_cuts_one_compilation() { + use asap_physical_operators::physical_planner::{ + compile, compile_candidate, cut_candidate, frontier_from_timing, + }; + for query in ["quantile(0.99, latency)", "sum by(job)(rate(m[1m]))"] { + let ephemeral = SummaryMaintenanceLifecycle::Ephemeral; + let (compiled_dag, _) = lifecycle_timed_dag(query, &ephemeral); + let inputs = raw_inputs(&compiled_dag); + let roots = [u64::from(compiled_dag.root.0)]; + let compiled = compile(&compiled_dag, inputs.clone(), &roots).unwrap(); + for lifecycle in [ + SummaryMaintenanceLifecycle::ContinuouslyMaintained, + ephemeral, + ] { + let (dag, states) = lifecycle_timed_dag(query, &lifecycle); + let frontier = frontier_from_timing(&dag).unwrap(); + let expected_frontier = if lifecycle == SummaryMaintenanceLifecycle::Ephemeral { + vec![] + } else { + states + }; + assert_eq!(frontier, expected_frontier, "{query} {lifecycle:?}"); + let cut = cut_candidate(&compiled, &frontier).unwrap(); + let expected = compile_candidate(&dag, inputs.clone(), &roots, &frontier).unwrap(); + assert_eq!( + serde_json::to_vec(&cut).unwrap(), + serde_json::to_vec(&expected).unwrap(), + "{query} {lifecycle:?}" + ); + } + } +} From abb0a8595b48031742f3cb32f8610ad16824ba73 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 02:47:54 +0000 Subject: [PATCH 6/8] docs: describe timing-derived candidate cuts Co-Authored-By: Claude Opus 5.5 --- .../physical-planning-and-deployment.md | 26 ++++++++++++------- docs/develop_docs/library-api.md | 24 ++++++++++------- 2 files changed, 30 insertions(+), 20 deletions(-) diff --git a/docs/design_docs/physical-planning-and-deployment.md b/docs/design_docs/physical-planning-and-deployment.md index 44c9c812..a37ba8eb 100644 --- a/docs/design_docs/physical-planning-and-deployment.md +++ b/docs/design_docs/physical-planning-and-deployment.md @@ -408,16 +408,21 @@ feasibility is rejected before pricing. The optimizer supplies candidate frontiers and cost evidence, including updates, retention, recurrence and sharing. `enumerate_frontiers` constructs bounded, reachable antichain frontiers above explicit input boundaries, including query-only and fully precomputed results. It fails explicitly when the candidate budget is exceeded. Maintenance selection must still reject frontiers that violate window, freshness, or reuse requirements; deployment feasibility is checked before pricing. -Lowering a node does not depend on the chosen frontier, so Planner lowers each -query DAG once. `compile(dag, inputs, roots)` yields the complete -`CompiledPhysicalDag`; `enumerate_compiled_frontiers(&compiled, max)` and -`cut_candidate(&compiled, frontier)` then derive each placement by partitioning -its operators: nodes above the frontier and below the roots form the query DAG, -and the frontier's ancestors form the precompute DAG. Helper operators are -numbered by their Planner node, so a cut is byte-identical to lowering that -frontier directly. `compile_candidate(s)` and `enumerate_frontiers` are wrappers -over this path. A deployment compiles each query DAG once, not once per -placement choice. Temporal pane candidates remain a separate lowering. +The lifecycle layer decides timing; physical compilation reads it. Lowering a +node does not depend on the frontier, so each query DAG is lowered once and +different lifecycle assignments are different cuts of that lowering. +`compile(dag, inputs, roots)` yields the complete `CompiledPhysicalDag`. +`frontier_from_timing(&timed_dag)` reads an assignment's timed DAG (from +`execution_timed_dag`) and returns its frontier: ingestion-time nodes read by +query-time nodes, or an ingestion-time root; a query-time node feeding an +ingestion-time node is rejected. `cut_candidate(&compiled, &frontier)` then +partitions the lowered operators: the frontier's ancestors form the precompute +DAG and the rest form the query DAG. Helper operators are numbered by their +Planner node (`u64::MAX - node_id`), so a cut is byte-identical to +`compile_candidate` for that frontier. One exception: an ingestion-time +`Binary` lowers differently, so its timing must match at compile time. +`compile_candidate(s)` and `enumerate_frontiers` wrap the same path. Temporal +pane candidates remain a separate lowering. Physical compilation opens no readers. Bounded precompute outputs become typed query inputs. Their source, filters, grouping, build window, evaluation time, readiness and @@ -538,6 +543,7 @@ pane compilation, alongside independent operator/runtime fixtures: | `summary_maintenance_lifecycle_e2e::selected_temporal_lifecycle_compiles_panes_and_executes` | PromQL p50/p99 workloads → selected continuous lifecycle and Sliding framework → automatically generated precompute/query DAGs → real codec round-trip → adjacent aligned windows; checks filters, entity identity, sample counts, missing/duplicate panes and phase rejection before opening readers | | `summary_maintenance_lifecycle_e2e::continuous_lifecycle_compiles_and_executes_spatial_kll` | PromQL workload → selected continuous lifecycle → logical DAG → compiled precompute/query candidate → results in independent revisions; an unbounded candidate fails before pricing, and a bounded request candidate summarizes the same input samples | | `summary_maintenance_lifecycle_e2e::chosen_lifecycle_timing_decides_precompute_contents` | PromQL workload → enumerated lifecycles → explicit choice → timed DAG → compiled candidate; ContinuouslyMaintained stores the state in precompute, Ephemeral leaves precompute empty and reads the raw source at query time; both return the same p99 | +| `summary_maintenance_lifecycle_e2e::lifecycle_timing_cuts_one_compilation` | KLL quantile and grouped Rate→Sum: one compilation cut by the ContinuouslyMaintained and Ephemeral timed DAGs equals `compile_candidate` for each; the frontier is the retained state or empty | | `summary_maintenance_lifecycle_e2e::planner_lifecycle_selection_reproduces_strategy_timing` | For PromQL fixtures, the timed DAG from Planner's retained selection equals the DAG realization strategies produce today | | `kll_pane_execution::five_panes_roundtrip_and_shared_merge_runs_once` | Explicit one-minute precompute DAGs → real MessagePack state bytes → five required query inputs → shared native merge → p50/p99; counts every sample once, checks adjacent aligned windows and instruments one merge start per run | | `kll_pane_execution::restored_panes_reject_corruption_parameters_schema_and_missing_binding` | Corrupt bytes, parameter relabelling, incompatible schemas and absent bindings fail explicitly | diff --git a/docs/develop_docs/library-api.md b/docs/develop_docs/library-api.md index 3d9652a7..d046b18a 100644 --- a/docs/develop_docs/library-api.md +++ b/docs/develop_docs/library-api.md @@ -631,28 +631,32 @@ an alternative with `MissingCostEvidence` is accepted only when the cost model's complete-candidate hook covers lifecycle costs. Window frameworks and totals come from that hook, as in Planner selection. -A lifecycle choice then fixes each physical placement: for example, a -continuously maintained state places its producer in precompute, while an -ephemeral one keeps it in the query. Compile each query's `PostAsapDag` once -and derive every placement from that result: +A lifecycle choice then fixes each physical placement through timing: a +continuously maintained state and its inputs run at ingestion time, while an +ephemeral one stays at query time. Compile each query's `PostAsapDag` once and +cut every chosen assignment from that result: ```rust use asap_physical_operators::physical_planner::{ - compile, cut_candidate, enumerate_compiled_frontiers, + compile, cut_candidate, frontier_from_timing, }; let compiled = compile(&dag, inputs, &roots)?; // each node lowered once -for frontier in enumerate_compiled_frontiers(&compiled, 4096)? { +for plan in lifecycle_plans { + let frontier = frontier_from_timing(&plan.execution_timed_dag()?)?; // Precompute/query DAGs split at `frontier`; no logical lowering. let candidate = cut_candidate(&compiled, &frontier)?; // Check feasibility and price `candidate`; bind the selected one as is. } ``` -`cut_candidate` returns exactly the `PhysicalCandidate` that -`compile_candidate(&dag, inputs, &roots, &frontier)` returns, and rejects the -same invalid frontiers. The mapping from lifecycle choice to frontier stays with -the caller. Temporal pane candidates are a different lowering and still use +The frontier is the set of ingestion-time nodes read by query-time nodes (or an +ingestion-time root). `frontier_from_timing` rejects a query-time node feeding +an ingestion-time node. `cut_candidate` returns exactly what +`compile_candidate(&dag, inputs, &roots, &frontier)` returns and rejects the +same invalid frontiers. If the DAG has an ingestion-time `Binary`, compile with +the same timing for that node, because it lowers differently. Temporal pane +candidates are a different lowering and still use `compile_temporal_pane_candidate`. ## Optional whole-plan selection and DAG assembly From d41f201a3c645e483ddeef32a7c1bc415c500be7 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 03:14:09 +0000 Subject: [PATCH 7/8] refactor: allow several helper operators per Planner node Number helpers as u64::MAX - (node << 16) - index so a node that lowers to an operator chain keeps deterministic, traversal-independent helper IDs. Co-Authored-By: Claude Opus 5.5 --- .../src/physical_planner/mod.rs | 14 ++++++++++---- .../tests/precompute_candidates.rs | 2 +- .../physical-planning-and-deployment.md | 2 +- 3 files changed, 12 insertions(+), 6 deletions(-) diff --git a/crates/asap-physical-operators/src/physical_planner/mod.rs b/crates/asap-physical-operators/src/physical_planner/mod.rs index 6a6ce940..9fdc6c84 100644 --- a/crates/asap-physical-operators/src/physical_planner/mod.rs +++ b/crates/asap-physical-operators/src/physical_planner/mod.rs @@ -115,6 +115,15 @@ thread_local! { static LOWERED_NODES: std::cell::Cell = const { std::cell::Cell::new(0) }; } +/// Helper operators are numbered from their Planner node alone, above the u32 +/// Planner ID range, so every boundary choice yields a subgraph of the same +/// lowering and candidate cuts need not renumber operators. A node lowering to +/// several helpers takes consecutive indices below its base. +fn helper_id(node: NodeId, index: u64) -> NodeId { + debug_assert!(node <= u64::from(u32::MAX) && index < 1 << 16); + u64::MAX - (node << 16) - index +} + fn compile_internal( dag: &PostAsapDag, mut sources: BTreeMap, @@ -173,10 +182,7 @@ fn compile_internal( let mut graph = CompiledPhysicalDag::new(roots.to_vec()); for id in ordered { let node = nodes[&id]; - // At most one helper operator per node, numbered above the u32 Planner - // ID range by its node alone, so every boundary choice yields a subgraph - // of the same lowering and candidate cuts need not renumber operators. - let auxiliary = u64::MAX - id; + let auxiliary = helper_id(id, 0); let output = Arc::new(node.output_schema.clone()); crate::values::validate_schema(&output)?; if let Some(source) = sources.remove(&id) { diff --git a/crates/asap-physical-operators/tests/precompute_candidates.rs b/crates/asap-physical-operators/tests/precompute_candidates.rs index 65345c76..e1bf9080 100644 --- a/crates/asap-physical-operators/tests/precompute_candidates.rs +++ b/crates/asap-physical-operators/tests/precompute_candidates.rs @@ -670,7 +670,7 @@ fn population_topk_cuts_equal_per_frontier_compilation() { let roots = [u64::from(dag.root.0)]; let compiled = compile(&dag, inputs.clone(), &roots).unwrap(); // The root reads its population through a Sort helper numbered by the root. - let helper = u64::MAX - roots[0]; + let helper = u64::MAX - (roots[0] << 16); assert_eq!(compiled.operator_name(helper), Some("Sort")); assert!( cut_candidate(&compiled, &[helper]).is_err(), diff --git a/docs/design_docs/physical-planning-and-deployment.md b/docs/design_docs/physical-planning-and-deployment.md index a37ba8eb..aea9cd6c 100644 --- a/docs/design_docs/physical-planning-and-deployment.md +++ b/docs/design_docs/physical-planning-and-deployment.md @@ -418,7 +418,7 @@ query-time nodes, or an ingestion-time root; a query-time node feeding an ingestion-time node is rejected. `cut_candidate(&compiled, &frontier)` then partitions the lowered operators: the frontier's ancestors form the precompute DAG and the rest form the query DAG. Helper operators are numbered by their -Planner node (`u64::MAX - node_id`), so a cut is byte-identical to +Planner node (`u64::MAX - (node_id << 16) - index`), so a cut is byte-identical to `compile_candidate` for that frontier. One exception: an ingestion-time `Binary` lowers differently, so its timing must match at compile time. `compile_candidate(s)` and `enumerate_frontiers` wrap the same path. Temporal From e45c6213b21c7a314c9ad47174caaa9c7527fcd4 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 03:19:33 +0000 Subject: [PATCH 8/8] test: expect retained frontier states read at query time or at the root With several retained states, only those read by a query-time node or forming the root are cut points. Co-Authored-By: Claude Opus 5.5 --- .../tests/summary_maintenance_lifecycle_e2e.rs | 18 ++++++++++++++++++ 1 file changed, 18 insertions(+) diff --git a/crates/integration-tests/tests/summary_maintenance_lifecycle_e2e.rs b/crates/integration-tests/tests/summary_maintenance_lifecycle_e2e.rs index 959304a9..0d778d0f 100644 --- a/crates/integration-tests/tests/summary_maintenance_lifecycle_e2e.rs +++ b/crates/integration-tests/tests/summary_maintenance_lifecycle_e2e.rs @@ -1095,10 +1095,28 @@ fn lifecycle_timing_cuts_one_compilation() { ] { let (dag, states) = lifecycle_timed_dag(query, &lifecycle); let frontier = frontier_from_timing(&dag).unwrap(); + // Retained states read by a query-time consumer, or the root itself. + let query_time = |id: u64| { + dag.nodes.iter().any(|node| { + u64::from(node.id.0) == id + && node.output_state.timing + == asap_types::post_asap::ExecutionTiming::QueryTime + }) + }; let expected_frontier = if lifecycle == SummaryMaintenanceLifecycle::Ephemeral { vec![] } else { states + .iter() + .copied() + .filter(|state| { + *state == u64::from(dag.root.0) + || dag.edges.iter().any(|edge| { + u64::from(edge.producer.0) == *state + && query_time(u64::from(edge.consumer.0)) + }) + }) + .collect() }; assert_eq!(frontier, expected_frontier, "{query} {lifecycle:?}"); let cut = cut_candidate(&compiled, &frontier).unwrap();