From 0ca3062dce3cf81a873f26ece55b163a1eae5cc7 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 03:06:21 +0000 Subject: [PATCH 1/3] fix: keep a selected grouped Sum over realized Rate readouts DAG assembly replaced any selected outer Sum over an inner aggregate with a query-time exact Sum, even when the selected summary realizes the inner Rate itself. The grouped Sum state therefore never reached the inventory, and no lifecycle choice could move grouped Sum into precompute. Assembly now keeps such a selected summary; the query-time residual still applies when the outer summary would hide its inner aggregate in KeepPreAsap. Default selection for sum by(job)(rate(...)) now yields Rate -> grouped Sum state; the physical frontier test reads its query root accordingly. Conflicts with earlier stack changes resolved to the integration tree: - crates/integration-tests/tests/summary_maintenance_lifecycle_e2e.rs: c98281a Merge remote-tracking branch 'origin/feat/compile-once-cuts' into integration/planner-for-backend Co-Authored-By: Claude Opus 5.5 --- crates/asap-aware-mapping/src/replacement.rs | 76 +++++++++- .../tests/precompute_candidates.rs | 7 +- .../summary_maintenance_lifecycle_e2e.rs | 130 ++++++++++++++++++ 3 files changed, 210 insertions(+), 3 deletions(-) diff --git a/crates/asap-aware-mapping/src/replacement.rs b/crates/asap-aware-mapping/src/replacement.rs index fb697251..59fc05be 100644 --- a/crates/asap-aware-mapping/src/replacement.rs +++ b/crates/asap-aware-mapping/src/replacement.rs @@ -5249,6 +5249,9 @@ impl<'a> GlobalSelection<'a> { if let Some(node) = self.assembled_nodes.borrow().get(&ptr) { return Ok(Rc::clone(node)); } + // A selected summary that realizes its inner aggregate, instead of + // hiding it in `KeepPreAsap`, is kept; lifecycle assignment decides + // whether it runs in precompute or at query time. let selected_composed_summary = self .groups .get(&ptr) @@ -5256,7 +5259,7 @@ impl<'a> GlobalSelection<'a> { .is_some_and(|candidate| matches!(&candidate.replacement, Replacement::Summary(node) if matches!(&node.expr, SummaryExpr::SummaryAgg { child, .. } - if matches!(&child.expr, SummaryExpr::KeepPreAsap(raw) if !contains_aggregate(raw))))); + if !matches!(&child.expr, SummaryExpr::KeepPreAsap(raw) if contains_aggregate(raw))))); let node = if query_time_nested_sum(target) && !selected_composed_summary { self.assemble_residual(target)? } else { @@ -6955,6 +6958,77 @@ mod tests { use asap_types::types::AccuracyTarget; use std::collections::HashMap; + // Candidate shape without execution timing: what is computed, not where. + fn timing_free_shape(node: &Rc) -> serde_json::Value { + fn strip(value: &mut serde_json::Value) { + match value { + serde_json::Value::Object(fields) => { + fields.remove("timing"); + fields.values_mut().for_each(strip); + } + serde_json::Value::Array(values) => values.iter_mut().for_each(strip), + _ => {} + } + } + let mut shape = + serde_json::to_value(asap_types::post_asap::compile_post_asap_dag(node).unwrap()) + .unwrap(); + strip(&mut shape); + shape + } + + // Rate inventories never offer two candidates that differ only in timing. + #[test] + fn rate_candidate_inventories_have_no_timing_only_duplicates() { + for (query, accuracy) in [ + ("sum by(job)(rate(m[1m]))", AccuracyTarget::Exact), + ("topk by(job)(2, rate(m[1m]))", AccuracyTarget::Epsilon(0.1)), + ] { + let root = Rc::new(lower_promql(query, accuracy)); + let inventory = search_workload(vec![(0usize, root)]) + .enumerate_candidate_dags(4096) + .unwrap(); + let shapes = inventory + .candidates + .iter() + .map(|forest| timing_free_shape(&forest[0].1)) + .collect::>(); + for (i, shape) in shapes.iter().enumerate() { + assert!(!shapes[..i].contains(shape), "{query}: duplicate {i}"); + } + } + } + + // Grouped Sum over Rate readouts stays a summary state in the inventory, + // so lifecycle assignment can place it in precompute or at query time. + #[test] + fn grouped_rate_sum_inventory_keeps_sum_state_for_lifecycle_placement() { + let root = Rc::new(lower_promql( + "sum by(job)(rate(m[1m]))", + AccuracyTarget::Exact, + )); + let inventory = search_workload(vec![(0usize, root)]) + .enumerate_candidate_dags(4096) + .unwrap(); + let is_exact = |node: &SummaryNode, kind: ExactKind| { + matches!(&node.expr, SummaryExpr::SummaryAgg { + family: SummaryFamilyType::ExactAggregate(k, _), .. + } if *k == kind) + }; + assert!(inventory.candidates.iter().any(|forest| { + let SummaryExpr::ValueOperation { child: sum, .. } = &forest[0].1.expr else { + return false; + }; + let SummaryExpr::SummaryAgg { child: rate, .. } = &sum.expr else { + return false; + }; + is_exact(sum, ExactKind::Sum) + && matches!(&rate.expr, SummaryExpr::ValueOperation { + child, operation: ValueOperation::FinalizeExactAccumulator, .. + } if is_exact(child, ExactKind::Rate)) + })); + } + // Every exposed query result has a readout; internal accumulator frontiers stay states. #[test] fn query_candidate_roots_do_not_leak_exact_accumulator_state() { diff --git a/crates/asap-physical-operators/tests/precompute_candidates.rs b/crates/asap-physical-operators/tests/precompute_candidates.rs index e1bf9080..f8050118 100644 --- a/crates/asap-physical-operators/tests/precompute_candidates.rs +++ b/crates/asap-physical-operators/tests/precompute_candidates.rs @@ -56,7 +56,7 @@ fn grouped_rate() -> PostAsapDag { let space = grouped_rate_space(); let selected = space .global_selection(&DefaultCostModel) - .assemble_selected_dag(&space.roots[0].1) + .assemble_selected_query(&space.roots[0].1) .unwrap() .unwrap(); compile_post_asap_dag(&selected).unwrap() @@ -108,7 +108,10 @@ fn grouped_rate_can_be_materialized_before_or_after_grouped_sum() { PostAsapOperatorPayload::Value { operation: ValueOperation::FinalizeExactAccumulator } - ) + ) && dag + .edges + .iter() + .any(|edge| edge.producer == state.id && edge.consumer == node.id) }) .unwrap(); let input_schema = Arc::new(state.output_schema.clone()); diff --git a/crates/integration-tests/tests/summary_maintenance_lifecycle_e2e.rs b/crates/integration-tests/tests/summary_maintenance_lifecycle_e2e.rs index f58e7ac7..e1a1f57b 100644 --- a/crates/integration-tests/tests/summary_maintenance_lifecycle_e2e.rs +++ b/crates/integration-tests/tests/summary_maintenance_lifecycle_e2e.rs @@ -686,3 +686,133 @@ fn lifecycle_timing_cuts_one_compilation() { } } } + +/// Grouped Rate→Sum is one inventory candidate: retaining the Sum state puts +/// Rate and Sum in precompute, while an `Ephemeral` Sum over a retained Rate +/// state leaves Sum in the query DAG. +#[test] +fn grouped_rate_sum_placement_is_a_lifecycle_choice() { + use asap_aware_mapping::enumerate_summary_maintenance_lifecycles; + use asap_physical_operators::physical_planner::{compile_candidate, InputContract}; + use asap_types::post_asap::{ + ExactKind, PostAsapOperatorPayload, SummaryExpr, SummaryFamilyType, + }; + use std::{collections::BTreeMap, sync::Arc}; + + let workload = quantile_workload("sum by(job)(rate(m[1m]))"); + let root = Rc::new( + asap_physical_operators::physical_planner::promql_rows::with_series_identity( + &lower_promql_workload(&workload, 0).unwrap().remove(0), + ) + .unwrap(), + ); + let is_exact = |node: &SummaryNode, kind: ExactKind| { + matches!(&node.expr, SummaryExpr::SummaryAgg { + family: SummaryFamilyType::ExactAggregate(k, _), .. + } if *k == kind) + }; + let inventory = asap_aware_mapping::search_workload(vec![("q", root)]) + .enumerate_candidate_dags(4096) + .unwrap(); + let candidates = inventory + .candidates + .into_iter() + .map(|mut forest| forest.remove(0).1) + .filter(|candidate| { + matches!(&candidate.expr, SummaryExpr::ValueOperation { child, .. } + if is_exact(child, ExactKind::Sum)) + }) + .collect::>(); + let [candidate] = candidates.as_slice() else { + panic!("one grouped Sum candidate, got {}", candidates.len()); + }; + let mut placements = Vec::new(); + for sum_lifecycle in [ + SummaryMaintenanceLifecycle::ContinuouslyMaintained, + SummaryMaintenanceLifecycle::Ephemeral, + ] { + let lifecycles = enumerate_summary_maintenance_lifecycles( + Rc::clone(candidate), + 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 = lifecycles + .deployments() + .iter() + .map(|deployment| { + let lifecycle = if is_exact(&deployment.summary, ExactKind::Sum) { + sum_lifecycle.clone() + } else { + SummaryMaintenanceLifecycle::ContinuouslyMaintained + }; + (deployment.post_asap_node_id, lifecycle) + }) + .collect::>(); + assert_eq!(choices.len(), 2, "Rate and Sum states"); + let dag = lifecycles + .select(&choices) + .unwrap() + .execution_timed_dag() + .unwrap(); + let raw = dag + .nodes + .iter() + .find(|node| matches!(node.payload, PostAsapOperatorPayload::Fallback { .. })) + .unwrap(); + let frontier = + asap_physical_operators::physical_planner::frontier_from_timing(&dag).unwrap(); + let [boundary] = frontier.as_slice() else { + panic!("one precompute output, got {frontier:?}"); + }; + let boundary = dag + .nodes + .iter() + .find(|node| u64::from(node.id.0) == *boundary) + .unwrap(); + let physical = compile_candidate( + &dag, + BTreeMap::from([( + u64::from(raw.id.0), + InputContract::bounded(Arc::new(raw.output_schema.clone())), + )]), + &[u64::from(dag.root.0)], + &frontier, + ) + .unwrap(); + let json = |value| String::from_utf8(serde_json::to_vec(value).unwrap()).unwrap(); + placements.push(( + boundary.payload.clone(), + json(physical.precompute.as_ref().unwrap()), + json(&physical.query), + )); + } + let builds = |json: &str, kind: &str| { + json.contains(&format!( + r#"{{"SummaryBuild":{{"family":{{"ExactAggregate":["{kind}","{kind}"]}}"# + )) + }; + let [(retained, retained_pre, retained_query), (ephemeral, ephemeral_pre, ephemeral_query)] = + placements.as_slice() + else { + unreachable!() + }; + let state = |payload: &PostAsapOperatorPayload, kind: ExactKind| { + matches!(payload, PostAsapOperatorPayload::SummaryAgg { + family: SummaryFamilyType::ExactAggregate(k, _), .. + } if *k == kind) + }; + assert!(state(retained, ExactKind::Sum)); + assert!(builds(retained_pre, "Rate") && builds(retained_pre, "Sum")); + assert!(!retained_query.contains("SummaryBuild")); + assert!(state(ephemeral, ExactKind::Rate)); + assert!(builds(ephemeral_pre, "Rate") && !builds(ephemeral_pre, "Sum")); + assert!(builds(ephemeral_query, "Sum")); +} From ff18c033503145f6bc174f5cf9bf44afd41fef52 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 03:06:22 +0000 Subject: [PATCH 2/3] refactor!: remove timing-only Rate placement candidates SketchAlgorithmStrategy::fixed_window_rate_candidates and query_time_rate_aggregation_candidates returned the same logical DAG as the ordinary heap or grouped Sum candidate with Rate finalization flipped between ingestion and query time. Placement now comes only from a chosen lifecycle via SummaryMaintenanceLifecyclePlan::execution_timed_dag. compile_fixed_window_rate_aggregation takes that lifecycle-timed PostAsapDag instead of a SummaryNode with baked-in timing. The fixed-window heap test binds continuously maintained lifecycles; the grouped Sum placement pair is covered by the lifecycle end-to-end test. Co-Authored-By: Claude Opus 5.5 --- crates/asap-aware-mapping/src/replacement.rs | 117 +----------- .../src/physical_planner/promql_rows.rs | 15 +- .../tests/weighted_topk_binding.rs | 178 +++++++++++------- 3 files changed, 118 insertions(+), 192 deletions(-) diff --git a/crates/asap-aware-mapping/src/replacement.rs b/crates/asap-aware-mapping/src/replacement.rs index 59fc05be..8c826b37 100644 --- a/crates/asap-aware-mapping/src/replacement.rs +++ b/crates/asap-aware-mapping/src/replacement.rs @@ -1356,114 +1356,6 @@ impl<'a> SketchAlgorithmStrategy<'a> { self.propose_with(&ranked, None) } - /// Fixed-window maintenance can finalize each series' counter state and - /// build a fresh heap or grouped Sum for that evaluation window. Deployment must provide - /// a complete, synchronized population and bind the matching window; this - /// candidate never incrementally adds one window's rates to another. - pub fn fixed_window_rate_candidates(&self, root: &Rc) -> Proposals { - fn place(node: &Rc) -> Option> { - let mut next = node.as_ref().clone(); - match &mut next.expr { - SummaryExpr::ValueOperation { - child, - operation: ValueOperation::FinalizeExactAccumulator, - timing, - } if matches!(&child.expr, SummaryExpr::SummaryAgg { - family: SummaryFamilyType::ExactAggregate(ExactKind::Rate, _), - reduction: Reduction::PerEntity, child: source, .. - } if matches!(&source.expr, SummaryExpr::KeepPreAsap(source) if matches!(source.as_ref(), QueryExpr::TimeRange { .. }))) => - { - *timing = ExecutionTiming::IngestionTime; - } - SummaryExpr::ValueOperation { child, .. } - | SummaryExpr::SummaryAgg { child, .. } => *child = place(child)?, - SummaryExpr::SummaryEstimate { summary_input, .. } => { - *summary_input = place(summary_input)? - } - _ => return None, - } - Some(Rc::new(next)) - } - let mut proposals = self.propose_with(root, None); - proposals.candidates.retain_mut(|candidate| { - let Replacement::Summary(node) = &candidate.replacement else { - return false; - }; - let Ok(dag) = asap_types::post_asap::compile_post_asap_dag(node) else { - return false; - }; - if !dag.nodes.iter().any(|node| match &node.payload { - asap_types::post_asap::PostAsapOperatorPayload::SummaryAgg { - family: SummaryFamilyType::Sketch(kind, _), - .. - } => matches!( - kind.algorithm(), - SketchAlgorithm::CmsWithHeap | SketchAlgorithm::CountSketchWithHeap - ), - asap_types::post_asap::PostAsapOperatorPayload::SummaryAgg { - family: SummaryFamilyType::ExactAggregate(ExactKind::Sum, _), - .. - } => true, - _ => false, - }) { - return false; - } - let Some(placed) = place(node) else { - return false; - }; - if asap_types::post_asap::compile_post_asap_dag(&placed).is_err() { - return false; - } - let Ok(placed) = finalize_query_candidate(placed, root) else { - return false; - }; - candidate.replacement = Replacement::Summary(placed); - candidate - .rationale - .push_str("; fixed-window precompute over complete per-series counter states"); - true - }); - proposals - } - - /// Retain grouped Sum after a per-series Rate readout as a query-time - /// candidate alongside its complete-window maintenance placement. - pub fn query_time_rate_aggregation_candidates(&self, root: &Rc) -> Proposals { - fn query_time(node: &Rc) -> Rc { - let mut next = node.as_ref().clone(); - match &mut next.expr { - SummaryExpr::ValueOperation { - child, - operation: ValueOperation::FinalizeExactAccumulator, - timing, - } if matches!( - &child.expr, - SummaryExpr::SummaryAgg { - family: SummaryFamilyType::ExactAggregate(ExactKind::Rate, _), - .. - } - ) => - { - *timing = ExecutionTiming::QueryTime; - } - SummaryExpr::ValueOperation { child, .. } - | SummaryExpr::SummaryAgg { child, .. } => *child = query_time(child), - _ => {} - } - Rc::new(next) - } - let mut proposals = self.fixed_window_rate_candidates(root); - proposals.candidates.retain_mut(|candidate| { - let Replacement::Summary(node) = &candidate.replacement else { return false }; - if !matches!(&node.expr, SummaryExpr::ValueOperation { child, operation: ValueOperation::FinalizeExactAccumulator, .. } - if matches!(&child.expr, SummaryExpr::SummaryAgg { family: SummaryFamilyType::ExactAggregate(ExactKind::Sum, _), .. })) { return false; } - candidate.replacement = Replacement::Summary(query_time(node)); - candidate.rationale = "query-time grouped Sum over complete per-series Rate readouts".into(); - true - }); - proposals - } - pub(crate) fn from_planning_inputs(planning_inputs: CandidatePlanningInputs<'a>) -> Self { Self { planning_inputs } } @@ -2764,7 +2656,9 @@ fn realize_physical_summary_input( /// Emit `SummaryAgg` (recursively binding the child), plus the /// `SummaryEstimate` readout when `estimate` is set. // Retain the exact expression and schema while placing its value production -// on the update path. Read-time consumers keep their original shared nodes. +// on the update path. This is the initial layout for values feeding a summary; +// lifecycle timing is authoritative. Read-time consumers keep their original +// shared nodes. fn maintenance_exact_values(node: Rc) -> Option> { let expr = match &node.expr { // These guards can fall back at read time, but cannot recover a parent @@ -2993,8 +2887,9 @@ fn construct_summary_agg( }; Rc::clone(child) } else if snapshot_weighted { - // A fresh query-time summary consumes this evaluation's finalized rates. - // Moving rate snapshots must never accumulate across evaluations. + // Each evaluation's finalized rates feed a fresh summary; rate snapshots + // must never accumulate across evaluations. Query time is only the + // initial layout; a retained summary's lifecycle moves it to ingestion. finalize_query_candidate(bound_child, &input.child)? } else { let child = finalize_exact_accumulator_at( diff --git a/crates/asap-physical-operators/src/physical_planner/promql_rows.rs b/crates/asap-physical-operators/src/physical_planner/promql_rows.rs index 8b887580..d2c13328 100644 --- a/crates/asap-physical-operators/src/physical_planner/promql_rows.rs +++ b/crates/asap-physical-operators/src/physical_planner/promql_rows.rs @@ -249,16 +249,13 @@ pub fn compile_rate_ranking( Ok((source, program)) } -/// The selected logical placement requires fresh aggregate state per closed window. -/// Compile both physical graphs before deployment chooses storage or scheduling. -/// The input is the complete collection of per-series exact counter states. +/// Compile a lifecycle-timed DAG whose heap or grouped Sum over per-series +/// Rate readouts runs at ingestion time: fresh aggregate state per closed +/// window. The input is the complete collection of per-series counter states. pub fn compile_fixed_window_rate_aggregation( - selected: &Rc, + dag: &planner_types::post_asap::PostAsapDag, ) -> Result { - use planner_types::post_asap::{ - compile_post_asap_dag, ExactKind, ExecutionTiming, SketchAlgorithm, - }; - let dag = compile_post_asap_dag(selected).map_err(|e| invalid(e.to_string()))?; + use planner_types::post_asap::{ExactKind, ExecutionTiming, SketchAlgorithm}; let sources = dag .nodes .iter() @@ -310,7 +307,7 @@ pub fn compile_fixed_window_rate_aggregation( )); } compile_candidate( - &dag, + dag, BTreeMap::from([( u64::from(source.id.0), InputContract::bounded(Arc::new(source.output_schema.clone())), diff --git a/crates/asap-physical-operators/tests/weighted_topk_binding.rs b/crates/asap-physical-operators/tests/weighted_topk_binding.rs index 664ae799..ef086b1b 100644 --- a/crates/asap-physical-operators/tests/weighted_topk_binding.rs +++ b/crates/asap-physical-operators/tests/weighted_topk_binding.rs @@ -756,10 +756,102 @@ fn spatial_topk_exposes_signed_heap_candidate_over_complete_snapshot() { } } -// Placement changes execution ownership only. Every fixed-window candidate -// contains Rate finalization before a fresh heap, with query readout downstream. +/// Deployment-side lifecycle choice: every summary state of `candidate` is +/// continuously maintained, and the chosen lifecycles set execution timing. +fn continuously_maintained_dag(candidate: &Rc) -> PostAsapDag { + use asap_aware_mapping::{ + cost_model::{Cost, CostModel}, + enumerate_summary_maintenance_lifecycles, CostRate, Horizon, + SummaryMaintenanceCapabilities, SummaryMaintenanceLifecycleCapabilities, + SummaryMaintenanceLifecycleCostInputs, WorkloadDemand, + }; + use planner_types::workload::{ + DataArrival, Rate, RepeatedDemand, RepeatingEntry, RepetitionInterval, + }; + struct Costed; + impl CostModel for Costed { + fn rank_candidates( + &self, + _: &planner_types::pre_asap::agg_intent::AggIntent, + candidates: &[SketchAlgorithm], + ) -> Vec { + candidates.to_vec() + } + fn summary_maintenance_lifecycle_cost_inputs( + &self, + _: &SummaryNode, + ) -> SummaryMaintenanceLifecycleCostInputs { + SummaryMaintenanceLifecycleCostInputs { + build_cost: Some(Cost(10.)), + maintenance_cost_per_update: Some(Cost(1.)), + summary_read_cost: Some(Cost(1.)), + retention_cost_rate: Some(CostRate(0.1)), + retirement_cost: Some(Cost(1.)), + } + } + fn summary_maintenance_capabilities( + &self, + _: &SummaryNode, + ) -> SummaryMaintenanceCapabilities { + SummaryMaintenanceCapabilities { + incremental_update: true, + merge: true, + delete: true, + } + } + } + const NOW_MS: u64 = 1_000_000; + let queries = QueryWorkload { + language: QueryLanguage::PromQL, + query_batch: None, + repeating_queries: Some(vec![RepeatingEntry { + query: Query("topk by(job)(2, rate(m[1m]))".into()), + demand: RepeatedDemand::FixedInterval(RepetitionInterval(60_000)), + requirements: QueryRequirements::default(), + predictability: Predictability::Predictable { known_at: None }, + time_selection: TimeSelection::default(), + }]), + }; + let data = DataWorkload { + arrival: DataArrival::ContinuouslyIngesting, + ingestion_rate: WorkloadEvidence { + value: Some(Rate(1.)), + source: planner_types::workload::EvidenceSource::Observed, + observed_at_ms: Some(NOW_MS), + valid_for_ms: Some(60_000), + }, + ..Default::default() + }; + let lifecycles = enumerate_summary_maintenance_lifecycles( + Rc::clone(candidate), + WorkloadDemand::new_with_data(&queries, &data, &[0]), + NOW_MS, + Some(Horizon(100.)), + SummaryMaintenanceLifecycleCapabilities::ALL, + &Costed, + ) + .unwrap(); + let choices = lifecycles + .deployments() + .iter() + .map(|deployment| { + ( + deployment.post_asap_node_id, + SummaryMaintenanceLifecycle::ContinuouslyMaintained, + ) + }) + .collect::>(); + lifecycles + .select(&choices) + .unwrap() + .execution_timed_dag() + .unwrap() +} + +// A maintained heap over finalized per-series Rate is the fixed-window +// placement: lifecycle timing, not a separate candidate, puts it in precompute. #[test] -fn planner_exposes_fixed_window_rate_heap_precompute_candidates() { +fn maintained_rate_heap_lifecycle_compiles_fixed_window_precompute() { use asap_physical_operators::physical_planner::{ compile_candidate, promql_rows::with_series_identity, }; @@ -775,13 +867,17 @@ fn planner_exposes_fixed_window_rate_heap_precompute_candidates() { &EqualSplitAllocator, &Evidence, ); - let candidates = strategy.fixed_window_rate_candidates(&root).candidates; + let candidates = strategy + .replacements(&TargetSubDAG::new(&root)) + .into_iter() + .filter_map(|candidate| match candidate.replacement { + Replacement::Summary(root) if candidate.rationale.contains("WithHeap") => Some(root), + _ => None, + }) + .collect::>(); assert_eq!(candidates.len(), 2); - for candidate in candidates { - let Replacement::Summary(root) = candidate.replacement else { - panic!() - }; - let dag = compile_post_asap_dag(&root).unwrap(); + for root in candidates { + let dag = continuously_maintained_dag(&root); let state = dag .nodes .iter() @@ -819,16 +915,11 @@ fn planner_exposes_fixed_window_rate_heap_precompute_candidates() { &[u64::from(heap.id.0)], ) .unwrap(); - let exported = asap_physical_operators::physical_planner::promql_rows::compile_fixed_window_rate_aggregation(&root).unwrap(); + let exported = asap_physical_operators::physical_planner::promql_rows::compile_fixed_window_rate_aggregation(&dag).unwrap(); assert_eq!( serde_json::to_vec(&exported).unwrap(), serde_json::to_vec(&physical).unwrap() ); - assert!( - asap_physical_operators::physical_planner::promql_rows::compile_rate_ranking(&root) - .is_err(), - "query binding must not move the selected precompute frontier" - ); // Execute the selected split across a state serialization boundary. // Each run builds fresh weights from that window's counters. let execute = |plan: &asap_physical_operators::physical_planner::CompiledPhysicalDag, @@ -959,60 +1050,3 @@ fn planner_exposes_fixed_window_rate_heap_precompute_candidates() { ); } } - -// Grouped Rate has a legal stored Sum candidate as well as query-time reduction. -#[test] -fn grouped_rate_exposes_precomputed_sum_with_query_readout() { - let root = Rc::new( - asap_physical_operators::physical_planner::promql_rows::with_series_identity( - &lower_promql("sum by(job)(rate(m[1m]))", AccuracyTarget::Exact).unwrap(), - ) - .unwrap(), - ); - let strategy = SketchAlgorithmStrategy::new_with_planning_inputs_and_evidence( - &DefaultCostModel, - &DefaultAccuracyModel, - &EqualSplitAllocator, - &Evidence, - ); - let direct = strategy.query_time_rate_aggregation_candidates(&root); - assert!( - direct.candidates.iter().any(|candidate| { - let Replacement::Summary(root) = &candidate.replacement else { - return false; - }; - let Ok((_, program)) = - asap_physical_operators::physical_planner::promql_rows::compile_rate_ranking(root) - else { - return false; - }; - let output = program.output_contract(program.roots()[0]).unwrap(); - output - .schema - .fields - .iter() - .all(|field| matches!(field.dtype, SummaryFamilyType::Plain(_))) - }), - "query-time grouped Rate must finalize Sum inside the physical graph" - ); - let candidates = strategy.fixed_window_rate_candidates(&root).candidates; - assert!( - !candidates.is_empty(), - "Planner must expose Rate -> grouped Sum at ingestion" - ); - for candidate in candidates { - let Replacement::Summary(root) = candidate.replacement else { - panic!() - }; - let physical = asap_physical_operators::physical_planner::promql_rows::compile_fixed_window_rate_aggregation(&root).unwrap(); - let precompute = - String::from_utf8(serde_json::to_vec(&physical.precompute.unwrap()).unwrap()).unwrap(); - assert!( - precompute.contains("SummaryBuild") - && precompute.contains("Rate") - && precompute.contains("Sum") - ); - let query = String::from_utf8(serde_json::to_vec(&physical.query).unwrap()).unwrap(); - assert!(query.contains("Readout") && !query.contains("SummaryBuild")); - } -} From eb5885ca7cb3a6a25ae01f5d4652c2cde5e547c9 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 03:06:22 +0000 Subject: [PATCH 3/3] docs: describe grouped Rate->Sum placement as a lifecycle choice Co-Authored-By: Claude Opus 5.5 --- .../architecture/input-output-workflow.md | 4 +- .../physical-planning-and-deployment.md | 38 ++++++++++--------- 2 files changed, 24 insertions(+), 18 deletions(-) diff --git a/docs/design_docs/architecture/input-output-workflow.md b/docs/design_docs/architecture/input-output-workflow.md index fb9a7f42..32c0b093 100644 --- a/docs/design_docs/architecture/input-output-workflow.md +++ b/docs/design_docs/architecture/input-output-workflow.md @@ -34,7 +34,9 @@ fields and [frontend dependencies](#frontend-specific-dependencies). DAG assembly](#selection-and-dag-assembly), and [summary-maintenance lifecycle](#summary-maintenance-lifecycle-aware-helper) APIs operate on this `PlanSpace`. These are alternative uses of the candidate space, not mandatory sequential -stages. `PlanSpace` itself has no selected summary-maintenance lifecycle. +stages. `PlanSpace` itself has no selected summary-maintenance lifecycle, and +its candidates do not choose precompute versus query-time placement: a chosen +lifecycle assignment sets each node's execution timing. The candidate DAGs are logical planning artifacts. ASAPPlanner does **not** produce a deployed executable plan; downstream systems bind physical operators, diff --git a/docs/design_docs/physical-planning-and-deployment.md b/docs/design_docs/physical-planning-and-deployment.md index 6bddfa6b..b53ee835 100644 --- a/docs/design_docs/physical-planning-and-deployment.md +++ b/docs/design_docs/physical-planning-and-deployment.md @@ -89,10 +89,12 @@ separate unsupported compilation, deployment infeasibility, missing evidence, and a feasible candidate that loses on cost. Absence is not a cost comparison. For `sum by(job)(rate(m[1m]))`, Rate remains per series before grouped Sum. -When lifecycle requirements permit it, a candidate may finalize Rate and Sum -within a bounded precompute run and persist the grouped value. Another may leave -those operators in the query DAG. Storing a value requires its exact evaluation -window, revision, readiness and serving cadence to match the query contract. +`PlanSpace` offers one such candidate, with a per-series Rate state and a grouped +Sum state. Its lifecycle assignment places it: a retained Sum state finalizes +Rate and builds Sum within a bounded precompute run; an `Ephemeral` Sum over a +retained Rate state leaves the Rate readout and Sum in the query DAG. Storing a +value requires its exact evaluation window, revision, readiness and serving +cadence to match the query contract. For instant-vector TopK, CMS/CountSketch with a candidate heap requires explicit series identity and a supported latest-value input protocol. Appending historical @@ -385,23 +387,25 @@ Materialization frontiers are Planner decisions. A candidate records both the precompute Physical DAG and the query Physical DAG, with typed outputs connecting them. The deployment compiler binds those outputs; it does not move operators. Lifecycle timing gives the frontier: ingestion-time nodes read by query-time -nodes. Moving further bounded consumers into precompute, as in Candidate B -below, is not yet expressed as a lifecycle choice. - -For `sum by(job)(rate(m[1m]))`, legal physical candidates can include: +nodes. For `sum by(job)(rate(m[1m]))`, the two lifecycle choices of the single +logical candidate give: ```text -Candidate A: - precompute: compatible per-series counter states → per-series Rate - materialized output: per-series rate values for window/evaluation/revision - query: stored per-series rate values → grouped Sum - -Candidate B: - precompute: compatible per-series counter states → per-series Rate → grouped Sum - materialized output: grouped values for window/evaluation/revision - query: stored grouped values → result +Candidate A (Rate state retained, Sum Ephemeral): + precompute: counter samples → per-series Rate state + materialized output: per-series Rate states for window/evaluation/revision + query: stored Rate states → Rate readout → grouped Sum → result + +Candidate B (Rate and Sum states retained): + precompute: counter samples → per-series Rate → grouped Sum state + materialized output: grouped Sum states for window/evaluation/revision + query: stored grouped Sum states → Sum readout → result ``` +Explicit frontiers passed to `compile_candidates` can also persist per-series +rate values; deriving that frontier from timing inside the physical planner is +not yet implemented. + Both preserve reset-aware Rate before Sum. Summing raw counters before Rate is not equivalent. The counter-state build may be another precompute DAG; typed state inputs do not imply that a deployment can construct or bind those states.