diff --git a/crates/asap-physical-operators/src/physical_planner/candidates.rs b/crates/asap-physical-operators/src/physical_planner/candidates.rs index a8f476f3..f5658a9a 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( @@ -68,6 +86,41 @@ pub fn compile_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 @@ -78,43 +131,38 @@ pub fn enumerate_frontiers( inputs: &BTreeMap, roots: &[NodeId], max_candidates: usize, +) -> Result>, Error> { + enumerate_compiled_frontiers(&compile(dag, inputs.clone(), roots)?, max_candidates) +} + +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 +190,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 +325,159 @@ 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); + } + + 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/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..9fdc6c84 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_frontiers, + frontier_from_timing, select_candidate, CandidateCost, CandidateSelection, PhysicalCandidate, }; mod temporal_panes; @@ -109,6 +109,21 @@ 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) }; +} + +/// 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, @@ -165,9 +180,9 @@ fn compile_internal( } } let mut graph = CompiledPhysicalDag::new(roots.to_vec()); - let mut auxiliary = u64::MAX; for id in ordered { let node = nodes[&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) { @@ -176,6 +191,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 +208,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 +302,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 +361,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..e1bf9080 100644 --- a/crates/asap-physical-operators/tests/precompute_candidates.rs +++ b/crates/asap-physical-operators/tests/precompute_candidates.rs @@ -4,7 +4,8 @@ use asap_physical_operators::{ factory::create_planner_accumulator, operators::Operator, physical_planner::{ - compile_candidates, select_candidate, CandidateCost, CompiledPhysicalDag, InputContract, + compile, compile_candidate, compile_candidates, cut_candidate, enumerate_frontiers, + select_candidate, CandidateCost, CompiledPhysicalDag, InputContract, PhysicalCandidate, Source, }, runtime::{Limits, RunContext, Scope}, @@ -546,3 +547,177 @@ 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_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!( + serde_json::to_vec(&cut).unwrap(), + serde_json::to_vec(&expected).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] << 16); + 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()); + } +} diff --git a/crates/integration-tests/tests/summary_maintenance_lifecycle_e2e.rs b/crates/integration-tests/tests/summary_maintenance_lifecycle_e2e.rs index 281ea4e8..0d778d0f 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,59 @@ 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(); + // 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(); + 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:?}" + ); + } + } +} diff --git a/docs/design_docs/physical-planning-and-deployment.md b/docs/design_docs/physical-planning-and-deployment.md index 1d2b11e8..aea9cd6c 100644 --- a/docs/design_docs/physical-planning-and-deployment.md +++ b/docs/design_docs/physical-planning-and-deployment.md @@ -408,6 +408,22 @@ 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. +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 << 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 +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 @@ -527,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 61be5e65..d046b18a 100644 --- a/docs/develop_docs/library-api.md +++ b/docs/develop_docs/library-api.md @@ -631,6 +631,34 @@ 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 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, frontier_from_timing, +}; + +let compiled = compile(&dag, inputs, &roots)?; // each node lowered once +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. +} +``` + +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 ### What does global selection mean?