diff --git a/crates/asap-aware-mapping/src/replacement.rs b/crates/asap-aware-mapping/src/replacement.rs index fb697251..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( @@ -5249,6 +5144,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 +5154,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 +6853,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/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/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/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")); - } -} 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")); +} 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.