diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index d8fff2597..575d8609f 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -41,7 +41,6 @@ use crate::types::AccuracyTarget; use planner_types::pre_asap::Source; mod placement; -mod rate_placement; mod windows; pub(super) use windows::gcd; pub use windows::{prepare_window_implementations, WindowCostModel}; @@ -95,24 +94,19 @@ pub struct QueryCompilationInput { } impl QueryCompilationInput { - /// Physical planning precedes deployment feasibility and cost selection. + /// Retain Planner's native realization of the selected root under the + /// timing its lifecycle choice wrote into it, before deployment pricing. pub(crate) fn retain_physical_candidate(&mut self) -> Result<(), CompileError> { - use asap_physical_operators::physical_planner::{promql_rows, PhysicalCandidate}; - let candidate = - rate_placement::compile_fixed_window_rate_aggregation(&self.selected_plan_root) - .or_else(|_| { - promql_rows::compile_current_series_readout(&self.selected_plan_root) - .or_else(|_| { - promql_rows::compile_rate_ranking(&self.selected_plan_root) - .map(|(_, dag)| dag) - }) - .map(|query| PhysicalCandidate { - precompute: None, - query, - materialized_outputs: BTreeMap::new(), - }) - }) - .ok(); + let candidate = placement::root_fixed_window_candidate(&self.selected_plan_root) + .ok() + .or_else(|| placement::query_time_candidate(&self.selected_plan_root)); + self.retain(candidate) + } + + fn retain( + &mut self, + candidate: Option, + ) -> Result<(), CompileError> { self.physical_candidate = candidate .map(|candidate| candidate.encode()) .transpose() @@ -1012,8 +1006,7 @@ impl BackendLocalPlanningInput { exact_costs_by_id.insert(format!("compat-query-{index}"), rows.clone()); } } - let mut candidate_roots = Vec::new(); - let mut planner_selection_trace = logical_roots_and_candidates( + let mut planner_selection_trace = select_logical_roots_with_scoped_evidence_and_trace( &mut queries, canonical_roots.clone(), &topk_evidence_by_id, @@ -1021,131 +1014,122 @@ impl BackendLocalPlanningInput { &exact_costs_by_id, self.physical_inputs.erp.as_ref(), self.environment.observed_at_unix_ms, - Some(&mut candidate_roots), )?; - // Resolve physical row identity before asking Planner for snapshot heap - // candidates. Keep canonical query semantics and exact candidates intact. - for (index, (query, root)) in queries.iter().zip(&canonical_roots).enumerate() { + for query in &mut queries { + prepare_window_implementations( + query, + &self.physical_inputs.window_cost_model, + self.environment.target, + self.physical_inputs.query_retention_margin_ms, + )?; + query.retain_physical_candidate()?; + } + // Native physical realizations need the complete series identity in + // their rows. Planner's PlanSpace proposes them for the identity-typed + // root; each is a logical alternative whose readout-built states are + // placed by lifecycle, then substituted into the preferred workload. + let mut planner_candidate_forests = Vec::new(); + for (index, root) in canonical_roots.iter().enumerate() { let Ok(typed) = asap_physical_operators::physical_planner::promql_rows::with_series_identity(root) else { continue; }; + let query = &queries[index]; let evidence = QueryEvidence { topk: topk_evidence_by_id.get(&query.query_id), scoped: scoped_evidence_by_id.get(&query.query_id), now_ms: self.environment.observed_at_unix_ms, }; - let strategy = - asap_aware_mapping::SketchAlgorithmStrategy::new_with_planning_inputs_and_evidence( - &asap_aware_mapping::cost_model::DefaultCostModel, - &asap_aware_mapping::accuracy::DefaultAccuracyModel, - &asap_aware_mapping::accuracy::EqualSplitAllocator, + // The identity-typed root reaches row-level realizations; the + // canonical root's search proposes whole-root realizations (such as + // a current-series heap) that add the identity themselves. + let mut candidates: Vec<(bool, Rc)> = Vec::new(); + for (typed_search, search_root) in [(true, Rc::new(typed)), (false, Rc::clone(root))] { + let inventory = crate::planner_selection::enumerate_workload_candidates( + vec![(index, search_root)], + query.accuracy_target.clone(), + &ControlPlaneCostModel::new(query.accuracy_target.clone()), &evidence, - ); - let typed = Rc::new(typed); - let mut proposed = - strategy.current_series_topk_candidates(&typed, &query.accuracy_target); - let direct = asap_aware_mapping::ReplacementStrategy::propose( - &strategy, - &asap_aware_mapping::TargetSubDAG::new(&typed), - ); - proposed - .candidates - .extend(rate_placement::fixed_window_rate_candidates( - &direct.candidates, - &typed, - )); - proposed - .candidates - .extend(rate_placement::query_time_rate_aggregation_candidates( - &direct.candidates, - &typed, - )); - proposed.candidates.extend(direct.candidates); - proposed.rejected.extend(direct.rejected); - for candidate in proposed.candidates { - let asap_aware_mapping::Replacement::Summary(root) = candidate.replacement else { - continue; - }; - let root = asap_aware_mapping::replacement::finalize_query_candidate(root, &typed) - .map_err(|error| CompileError::Snapshot(error.to_string()))?; - let compiled = asap_physical_operators::physical_planner::promql_rows::compile_current_series_readout(&root) - .or_else(|_| asap_physical_operators::physical_planner::promql_rows::compile_rate_ranking(&root).map(|(_, program)| program)); - if let Ok(physical) = rate_placement::compile_fixed_window_rate_aggregation(&root) { + &asap_aware_mapping::DefaultAccuracyModel, + ) + .map_err(|error| CompileError::Snapshot(error.to_string()))?; + for reason in inventory + .rejected_assemblies + .into_iter() + .filter(|_| typed_search) + { planner_selection_trace.push(serde_json::json!({ - "stage":"planner.physical_candidate", "query_id":query.query_id, - "logical_root_id":crate::planner_selection::explained_root_id(&root, &query.accuracy_target), - "rationale":candidate.rationale, "physical_candidate":serde_json::from_slice::(&physical.encode().map_err(|e| CompileError::Snapshot(e.to_string()))?).map_err(|e| CompileError::Snapshot(e.to_string()))?, - "guarantee":root.guarantee, + "stage": "planner.physical_candidate", "query_id": query.query_id, + "status": "rejected", "reason": reason, })); - candidate_roots.push(vec![(index, root)]); - continue; } - match compiled { - Ok(program) => { + for (_, candidate) in inventory.candidates.into_iter().flatten() { + if !candidates.iter().any(|(_, known)| **known == *candidate) { + candidates.push((typed_search, candidate)); + } + } + } + for (typed_search, candidate) in candidates { + let Some(timed) = placement::time_native_candidate( + &candidate, + query, + &workload, + &data_workload, + index, + &self.environment, + ) else { + if typed_search { planner_selection_trace.push(serde_json::json!({ "stage": "planner.physical_candidate", "query_id": query.query_id, - "logical_root_id": crate::planner_selection::explained_root_id(&root, &query.accuracy_target), - "rationale": candidate.rationale, "physical_dag": serde_json::from_slice::(&program.encode().map_err(|error| CompileError::Snapshot(error.to_string()))?).map_err(|error| CompileError::Snapshot(error.to_string()))?, - "guarantee": root.guarantee, + "logical_root_id": crate::planner_selection::explained_root_id(&candidate, &query.accuracy_target), + "status": "unsupported", + "reason": "no native physical realization under any lifecycle choice", })); - candidate_roots.push(vec![(index, root)]); } - Err(error) => planner_selection_trace.push(serde_json::json!({ - "stage": "planner.physical_candidate", "query_id": query.query_id, - "status": "unsupported", "reason": error.to_string(), - })), - } - } - for rejected in proposed.rejected { - planner_selection_trace.push(serde_json::json!({ - "stage": "planner.physical_candidate", "query_id": query.query_id, - "status": "rejected", "reason": rejected.error.to_string(), - })); - } - } - for query in &mut queries { - prepare_window_implementations( - query, - &self.physical_inputs.window_cost_model, - self.environment.target, - self.physical_inputs.query_retention_margin_ms, - )?; - } - let mut planner_candidate_forests = Vec::new(); - for roots in candidate_roots { - let mut forest = queries.clone(); - let changed: BTreeSet<_> = roots.iter().map(|(index, _)| *index).collect(); - for (index, root) in roots { - forest[index].selected_plan_root = root; - } - // Unchanged roots already have prepared windows in `queries`. - let preparation = changed.into_iter().try_for_each(|index| { - prepare_window_implementations( + continue; + }; + let mut forest = queries.clone(); + forest[index].selected_plan_root = Rc::clone(&timed.root); + // Unchanged roots already have prepared windows in `queries`. + if let Err(error) = prepare_window_implementations( &mut forest[index], &self.physical_inputs.window_cost_model, self.environment.target, self.physical_inputs.query_retention_margin_ms, - ) - }); - if let Err(error) = preparation { - planner_selection_trace.push(serde_json::json!({ - "stage": "deployment.window_feasibility", - "status": "rejected", - "logical_root_ids": forest.iter().map(|query| crate::planner_selection::explained_root_id( - &query.selected_plan_root, &query.accuracy_target)).collect::>(), - "reason": error.to_string(), - })); - continue; + ) { + planner_selection_trace.push(serde_json::json!({ + "stage": "deployment.window_feasibility", + "status": "rejected", + "query_id": query.query_id, + "logical_root_id": crate::planner_selection::explained_root_id(&timed.root, &query.accuracy_target), + "reason": error.to_string(), + })); + continue; + } + let encode = |bytes: Result, asap_physical_operators::Error>| { + bytes + .map_err(|e| CompileError::Snapshot(e.to_string())) + .and_then(|bytes| { + serde_json::from_slice::(&bytes) + .map_err(|e| CompileError::Snapshot(e.to_string())) + }) + }; + let mut event = serde_json::json!({ + "stage": "planner.physical_candidate", "query_id": query.query_id, + "logical_root_id": crate::planner_selection::explained_root_id(&timed.root, &query.accuracy_target), + "guarantee": timed.root.guarantee, + }); + if timed.physical.precompute.is_some() { + event["physical_candidate"] = encode(timed.physical.encode())?; + } else { + event["physical_dag"] = encode(timed.physical.query.encode())?; + } + planner_selection_trace.push(event); + planner_selection_trace.extend(timed.trace); + forest[index].retain(Some(timed.physical))?; + planner_candidate_forests.push(forest); } - planner_candidate_forests.push(forest); - } - for query in queries - .iter_mut() - .chain(planner_candidate_forests.iter_mut().flatten()) - { - query.retain_physical_candidate()?; } // Composable lowering residualizes unsafe leaves individually; retain Planner siblings. Ok(( @@ -2935,7 +2919,6 @@ pub fn select_logical_roots_with_scoped_evidence_and_trace( exact_costs, erp, now_ms, - None, ) } @@ -2947,7 +2930,6 @@ fn logical_roots_and_candidates( exact_costs: &HashMap>, erp: Option<&super::erp::ErpPlanningInput>, now_ms: u64, - mut candidates_out: Option<&mut Vec)>>>, ) -> Result, CompileError> { let mut traces = Vec::new(); if roots.len() != queries.len() { @@ -3023,27 +3005,6 @@ fn logical_roots_and_candidates( AccuracyTarget::Exact => 0.0, }, }; - let mut candidate_assembly_rejections = Vec::new(); - if let Some(output) = candidates_out.as_deref_mut() { - let inventory = crate::planner_selection::enumerate_workload_candidates( - roots.clone(), - accuracy.clone(), - &model, - &QueryEvidence { - topk: certificate, - scoped: scoped_certificate, - now_ms, - }, - &accuracy_model, - ) - .map_err(|error| CompileError::Snapshot(error.to_string()))?; - candidate_assembly_rejections = inventory.rejected_assemblies; - let candidates = inventory.candidates; - // Retain every root's candidates without multiplying independent - // cohorts. Deployment evaluates each root substitution in the - // preferred workload context; this is not exhaustive joint search. - output.extend(candidates); - } let (selected, mut trace) = crate::planner_selection::select_workload_with_accuracy_model_and_trace( roots, @@ -3057,14 +3018,6 @@ fn logical_roots_and_candidates( &accuracy_model, ) .map_err(|error| CompileError::Snapshot(error.to_string()))?; - trace["candidate_assembly_rejections"] = serde_json::json!(candidate_assembly_rejections); - if candidates_out.is_some() { - trace["computation_search_scope"] = serde_json::json!({ - "inventory": "all_root_candidates", - "deployment_evaluation": "single_root_substitutions_in_preferred_workload", - "joint_workload_search_exhaustive": false, - }); - } if let Some(evidence) = scoped_certificate { trace["accuracy_evidence_scope"] = serde_json::json!({ "query_id": scope, @@ -4718,7 +4671,7 @@ fn collect_selected_materializations( physical_source.as_ref().unwrap_or(node), None, composable, - rate_placement::compile_fixed_window_rate_aggregation(node).is_ok(), + placement::root_fixed_window_candidate(node).is_ok(), None, &mut selected, )?; diff --git a/control_plane/src/physical/compiler/placement.rs b/control_plane/src/physical/compiler/placement.rs index 4e46eede1..68034522c 100644 --- a/control_plane/src/physical/compiler/placement.rs +++ b/control_plane/src/physical/compiler/placement.rs @@ -9,7 +9,7 @@ use super::*; use asap_aware_mapping::{enumerate_summary_maintenance_lifecycles, CostModel}; use asap_physical_operators::physical_planner::{ - compile, promql_fallback, CompiledPhysicalDag, InputContract, + compile, promql_fallback, promql_rows, CompiledPhysicalDag, InputContract, PhysicalCandidate, }; use planner_types::post_asap::PostAsapOperatorPayload; use planner_types::pre_asap::{AggIntent, CompareOpKind, ScalarValue}; @@ -418,6 +418,210 @@ fn range_selector_scan(selector: &QueryExpr) -> Result Result { + promql_rows::compile_fixed_window_rate_aggregation(dag) +} + +/// [`fixed_window_candidate`] under the timing written into `root`. +pub(super) fn root_fixed_window_candidate( + root: &Rc, +) -> Result { + let dag = planner_types::post_asap::compile_post_asap_dag(root) + .map_err(|e| asap_physical_operators::Error::Invalid(e.to_string()))?; + fixed_window_candidate(&dag) +} + +/// Planner's query-time realization above a maintained population or exact +/// per-series Rate readouts; the input is bound by the backend. +pub(super) fn query_time_candidate(root: &Rc) -> Option { + promql_rows::compile_current_series_readout(root) + .or_else(|_| promql_rows::compile_rate_ranking(root).map(|(_, dag)| dag)) + .ok() + .map(|query| PhysicalCandidate { + precompute: None, + query, + materialized_outputs: BTreeMap::new(), + }) +} + +/// A summary built from another retained state's finalized readouts, such as +/// a heap or grouped Sum over per-series Rate. +fn over_readouts(summary: &SummaryNode) -> bool { + matches!(&summary.expr, SummaryExpr::SummaryAgg { child, .. } + if matches!(&child.expr, SummaryExpr::ValueOperation { + operation: planner_types::post_asap::ValueOperation::FinalizeExactAccumulator, + child, + .. + } if matches!(child.expr, SummaryExpr::SummaryAgg { .. }))) +} + +/// Write `timing` onto the readouts feeding `target`, rebuilding its path. +fn retime( + node: &Rc, + target: &Rc, + timing: planner_types::post_asap::ExecutionTiming, +) -> Option> { + let mut next = node.as_ref().clone(); + if Rc::ptr_eq(node, target) { + let SummaryExpr::SummaryAgg { child, .. } = &mut next.expr else { + return None; + }; + let mut readout = child.as_ref().clone(); + let SummaryExpr::ValueOperation { timing: old, .. } = &mut readout.expr else { + return None; + }; + *old = timing; + *child = Rc::new(readout); + return Some(Rc::new(next)); + } + match &mut next.expr { + SummaryExpr::ValueOperation { child, .. } | SummaryExpr::SummaryAgg { child, .. } => { + *child = retime(child, target, timing)? + } + SummaryExpr::SummaryEstimate { summary_input, .. } => { + *summary_input = retime(summary_input, target, timing)? + } + _ => return None, + } + Some(Rc::new(next)) +} + +/// A Planner candidate with a native physical realization after its +/// readout-built states are placed by lifecycle. +pub(super) struct TimedCandidate { + pub(super) root: Rc, + pub(super) physical: PhysicalCandidate, + pub(super) trace: Vec, +} + +/// Place each state of `root` built from retained readouts: maintained, it +/// consumes complete per-series states at ingestion; ephemeral, it is rebuilt +/// for each query from the retained state's readouts. Planner's lifecycle plan +/// times the DAG and its native compiler reads that timing. Other states stay +/// retained here; compilation places them. `None` means no native realization. +pub(super) fn time_native_candidate( + root: &Rc, + query: &QueryCompilationInput, + workload: &QueryWorkload, + data: &DataWorkload, + index: usize, + environment: &PhysicalDeploymentContext, +) -> Option { + use planner_types::post_asap::ExecutionTiming; + let untimed = || { + query_time_candidate(root).map(|physical| TimedCandidate { + root: Rc::clone(root), + physical, + trace: Vec::new(), + }) + }; + let lifecycle = &query.summary_lifecycle_inputs; + let model = LifecycleCosts { + costs: &lifecycle.costs, + evaluation_interval_ms: lifecycle.evaluation_interval_ms, + input_cardinality: data + .input_cardinality + .value_at(environment.observed_at_unix_ms) + .copied(), + delete: environment.target == PhysicalDeploymentTarget::BackendLocalRemoteWrite, + }; + let enumerate = || { + enumerate_summary_maintenance_lifecycles( + Rc::clone(root), + WorkloadDemand::new_with_data(workload, data, std::slice::from_ref(&index)), + environment.observed_at_unix_ms, + Some(Horizon(lifecycle.horizon_seconds)), + SummaryMaintenanceLifecycleCapabilities { + supports_ephemeral: true, + supports_prepared: false, + supports_shared: false, + supports_continuously_maintained: true, + }, + &model, + ) + .ok() + }; + let candidates = enumerate()?; + let mut choices = Vec::new(); + let mut retimed = Rc::clone(root); + let mut trace = Vec::new(); + let mut maintained_over_readouts = false; + for deployment in candidates.deployments() { + let retained = alternative_cost( + deployment, + &SummaryMaintenanceLifecycle::ContinuouslyMaintained, + ); + let lifecycle = if over_readouts(&deployment.summary) { + let rebuilt = alternative_cost(deployment, &SummaryMaintenanceLifecycle::Ephemeral); + let ephemeral = rebuild_is_cheaper(retained, rebuilt); + maintained_over_readouts |= !ephemeral; + let timing = if ephemeral { + ExecutionTiming::QueryTime + } else { + ExecutionTiming::IngestionTime + }; + retimed = retime(&retimed, &find(&retimed, &deployment.summary)?, timing)?; + // Recorded per candidate forest; compilation records the final + // placement of the states it installs as `lifecycle_placement`. + trace.push(json!({ + "stage": "deployment.native_candidate_placement", + "query_ids": [&query.query_id], + "candidate_root_id": crate::planner_selection::explained_root_id(root, &query.accuracy_target), + "logical_root_id": crate::planner_selection::explained_root_id(&deployment.summary, &query.accuracy_target), + "ephemeral_bindable": true, + "continuously_maintained_cost": retained.map(|cost| cost.0), + "ephemeral_cost": rebuilt.map(|cost| cost.0), + "selected": if ephemeral { "ephemeral" } else { "continuously_maintained" }, + })); + if ephemeral { + SummaryMaintenanceLifecycle::Ephemeral + } else { + SummaryMaintenanceLifecycle::ContinuouslyMaintained + } + } else { + SummaryMaintenanceLifecycle::ContinuouslyMaintained + }; + choices.push((deployment.post_asap_node_id, lifecycle)); + } + if trace.is_empty() { + return untimed(); + } + let timed = candidates + .select(&choices) + .ok()? + .execution_timed_dag() + .ok()?; + let physical = if maintained_over_readouts { + fixed_window_candidate(&timed).ok()? + } else { + query_time_candidate(&retimed)? + }; + Some(TimedCandidate { + root: retimed, + physical, + trace, + }) +} + +/// The node of `root` structurally equal to `summary`, found after earlier +/// rewrites replaced the original `Rc`s. +fn find(root: &Rc, summary: &SummaryNode) -> Option> { + if root.as_ref() == summary { + return Some(Rc::clone(root)); + } + match &root.expr { + SummaryExpr::ValueOperation { child, .. } | SummaryExpr::SummaryAgg { child, .. } => { + find(child, summary) + } + SummaryExpr::SummaryEstimate { summary_input, .. } => find(summary_input, summary), + _ => None, + } +} + #[cfg(test)] mod tests { use super::*; diff --git a/control_plane/src/physical/compiler/rate_placement.rs b/control_plane/src/physical/compiler/rate_placement.rs deleted file mode 100644 index 413a99d04..000000000 --- a/control_plane/src/physical/compiler/rate_placement.rs +++ /dev/null @@ -1,225 +0,0 @@ -//! Rate-placement variants of Planner heap and grouped Sum candidates. -//! -//! Planner no longer lists these: it treats placement as a lifecycle choice. -//! Until the compiler selects placement through lifecycle timing, it derives -//! the same two variants that Planner used to offer, with timing written into -//! the candidate root, and compiles them as before. -use asap_aware_mapping::{Replacement, ReplacementSubDAG}; -use asap_physical_operators::physical_planner::{promql_rows, PhysicalCandidate}; -use asap_physical_operators::Error; -use planner_types::post_asap::{ - compile_post_asap_dag, ExactKind, ExecutionTiming, PostAsapOperatorPayload, SketchAlgorithm, - SummaryExpr, SummaryFamilyType, SummaryNode, ValueOperation, -}; -use planner_types::pre_asap::{QueryExpr, Reduction}; -use std::rc::Rc; - -/// Compile a candidate whose Rate finalization timing is part of its root. -pub(crate) fn compile_fixed_window_rate_aggregation( - selected: &Rc, -) -> Result { - let dag = compile_post_asap_dag(selected).map_err(|e| Error::Invalid(e.to_string()))?; - promql_rows::compile_fixed_window_rate_aggregation(&dag) -} - -/// Fixed-window maintenance finalizes each series' counter state and builds a -/// fresh heap or grouped Sum for that evaluation window. Deployment must -/// provide a complete, synchronized population and bind the matching window; -/// this never adds one window's rates to another. -pub(crate) fn fixed_window_rate_candidates( - direct: &[ReplacementSubDAG], - root: &Rc, -) -> Vec { - 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 candidates = direct.to_vec(); - candidates.retain_mut(|candidate| { - let Replacement::Summary(node) = &candidate.replacement else { - return false; - }; - let Ok(dag) = compile_post_asap_dag(node) else { - return false; - }; - if !dag.nodes.iter().any(|node| match &node.payload { - PostAsapOperatorPayload::SummaryAgg { - family: SummaryFamilyType::Sketch(kind, _), - .. - } => matches!( - kind.algorithm(), - SketchAlgorithm::CmsWithHeap | SketchAlgorithm::CountSketchWithHeap - ), - PostAsapOperatorPayload::SummaryAgg { - family: SummaryFamilyType::ExactAggregate(ExactKind::Sum, _), - .. - } => true, - _ => false, - }) { - return false; - } - let Some(placed) = place(node) else { - return false; - }; - if compile_post_asap_dag(&placed).is_err() { - return false; - } - let Ok(placed) = asap_aware_mapping::replacement::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 - }); - candidates -} - -/// Retain grouped Sum after a per-series Rate readout as a query-time -/// candidate alongside its complete-window maintenance placement. -pub(crate) fn query_time_rate_aggregation_candidates( - direct: &[ReplacementSubDAG], - root: &Rc, -) -> Vec { - 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 candidates = fixed_window_rate_candidates(direct, root); - 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 - }); - candidates -} - -#[cfg(test)] -mod tests { - use super::*; - use asap_aware_mapping::{ReplacementStrategy, SketchAlgorithmStrategy, TargetSubDAG}; - use planner_types::types::AccuracyTarget; - - fn candidates(query: &str) -> (Vec, Vec) { - let root = - crate::query_parser::parse_query_expr_canonical(query, AccuracyTarget::Exact).unwrap(); - let typed = Rc::new(promql_rows::with_series_identity(&root).unwrap()); - let strategy = - SketchAlgorithmStrategy::new(&asap_aware_mapping::cost_model::DefaultCostModel); - let direct = strategy.propose(&TargetSubDAG::new(&typed)).candidates; - ( - fixed_window_rate_candidates(&direct, &typed), - query_time_rate_aggregation_candidates(&direct, &typed), - ) - } - - fn rate_timing(node: &Rc) -> Option { - match &node.expr { - SummaryExpr::ValueOperation { - child, - operation: ValueOperation::FinalizeExactAccumulator, - timing, - } if matches!( - &child.expr, - SummaryExpr::SummaryAgg { - family: SummaryFamilyType::ExactAggregate(ExactKind::Rate, _), - .. - } - ) => - { - Some(*timing) - } - SummaryExpr::ValueOperation { child, .. } | SummaryExpr::SummaryAgg { child, .. } => { - rate_timing(child) - } - SummaryExpr::SummaryEstimate { summary_input, .. } => rate_timing(summary_input), - _ => None, - } - } - - // Grouped Sum over Rate yields a fixed-window precompute placement and a - // query-time placement of the same computation. - #[test] - fn grouped_rate_sum_has_both_placements() { - let (fixed, query_time) = candidates("sum by (job) (rate(requests_total[1m]))"); - let [fixed] = fixed.as_slice() else { - panic!("expected one fixed-window candidate, got {}", fixed.len()); - }; - let Replacement::Summary(fixed) = &fixed.replacement else { - panic!("expected a summary replacement"); - }; - assert_eq!(rate_timing(fixed), Some(ExecutionTiming::IngestionTime)); - let physical = compile_fixed_window_rate_aggregation(fixed).unwrap(); - assert!(physical.precompute.is_some()); - - let [query_time] = query_time.as_slice() else { - panic!( - "expected one query-time candidate, got {}", - query_time.len() - ); - }; - let Replacement::Summary(query_time) = &query_time.replacement else { - panic!("expected a summary replacement"); - }; - assert_eq!(rate_timing(query_time), Some(ExecutionTiming::QueryTime)); - assert!(compile_fixed_window_rate_aggregation(query_time).is_err()); - } - - // Candidates without a heap or grouped Sum have no Rate placement variant. - #[test] - fn per_series_rate_has_no_placement_variant() { - let (fixed, query_time) = candidates("rate(requests_total[1m])"); - assert!(fixed.is_empty()); - assert!(query_time.is_empty()); - } -} diff --git a/control_plane/src/physical/compiler/windows.rs b/control_plane/src/physical/compiler/windows.rs index 02d81bee1..fa0fb311a 100644 --- a/control_plane/src/physical/compiler/windows.rs +++ b/control_plane/src/physical/compiler/windows.rs @@ -158,8 +158,7 @@ pub fn prepare_window_implementations( })?; let cohorts = cohort_nodes(&states); let native_cohort = - super::rate_placement::compile_fixed_window_rate_aggregation(&query.selected_plan_root) - .is_ok(); + super::placement::root_fixed_window_candidate(&query.selected_plan_root).is_ok(); let requirements = states .iter() .map(|state| { diff --git a/control_plane/src/physical/workload_cost.rs b/control_plane/src/physical/workload_cost.rs index ed946d6d9..32968be90 100644 --- a/control_plane/src/physical/workload_cost.rs +++ b/control_plane/src/physical/workload_cost.rs @@ -939,10 +939,10 @@ mod tests { use super::super::compiler::BackendLocalPlanningInput; use super::*; - /// Deployment pricing must see candidate sketch families, not only an - /// upstream winner chosen before runtime resource evidence is applied. + /// PlanSpace decides the computation: every quantile family it can + /// realize is ranked in the selection trace, and one is committed. #[test] - fn planner_quantile_inventory_reaches_backend_before_family_selection() { + fn planner_ranks_quantile_families_before_committing_one() { let mut input = fixture(); let queries = input.query_workload.repeating_queries.as_mut().unwrap(); queries.truncate(1); @@ -951,20 +951,11 @@ mod tests { crate::types::AccuracyTarget::Epsilon(0.05), ); let (request, _) = input.into_physical_compilation_request().unwrap(); - let candidates = enumerate_exact_and_materialized_candidates(request).unwrap(); - let roots = candidates - .iter() - .flat_map(|c| &c.queries) - .map(|q| format!("{:?}", q.selected_plan_root)) - .collect::>(); - assert!( - roots.iter().any(|r| r.contains("Kll")), - "KLL disappeared before backend pricing" - ); - assert!( - roots.iter().any(|r| r.contains("DDSketch")), - "DDSketch disappeared before backend pricing" - ); + let trace = serde_json::to_string(&request.planner_selection_trace).unwrap(); + assert!(trace.contains("Kll sketch"), "KLL was not ranked"); + assert!(trace.contains("DDSketch sketch"), "DDSketch was not ranked"); + let committed = format!("{:?}", request.queries[0].selected_plan_root); + assert!(committed.contains("Kll") != committed.contains("DDSketch")); } /// A deployment without external execution rejects native candidates before pricing. @@ -1218,12 +1209,7 @@ mod tests { let evidence = input.workload_cost_evidence.clone().unwrap(); let (request, env) = input.into_physical_compilation_request().unwrap(); let candidates = enumerate_exact_and_materialized_candidates(request).unwrap(); - assert!(candidates.len() > 2); - assert!(candidates[1..].iter().any(|candidate| candidate - .queries - .iter() - .zip(&candidates[0].queries) - .any(|(a, b)| Rc::ptr_eq(&a.selected_plan_root, &b.selected_plan_root)))); + assert!(candidates.len() >= 2); let reference = candidates .iter() diff --git a/control_plane/tests/native_rate_topk.rs b/control_plane/tests/native_rate_topk.rs index 88ed8725e..31f22628f 100644 --- a/control_plane/tests/native_rate_topk.rs +++ b/control_plane/tests/native_rate_topk.rs @@ -7,6 +7,13 @@ use control_plane::physical::{ use serde_json::{json, Value}; fn fixture(certified: bool) -> BackendLocalPlanningInput { + placed_fixture(certified, 0.0) +} + +/// A positive summary-store price makes rebuilding a heap per query from the +/// retained counter readouts cheaper than maintaining it. Local-only execution +/// keeps the counter state itself precomputed rather than read raw. +fn placed_fixture(certified: bool, store_per_byte_second: f64) -> BackendLocalPlanningInput { let mut wire: Value = serde_json::from_str(include_str!( "../../docs/examples/asapquery-compatibility-demo-snapshot.json" )) @@ -18,6 +25,10 @@ fn fixture(certified: bool) -> BackendLocalPlanningInput { json!({"explicit":{"EpsilonDelta":{"epsilon":0.1,"delta":0.1}}}); wire["query_workload"]["repeating_queries"] = json!([entry]); wire["implementation"]["topk_evidence"] = json!({}); + wire["implementation"]["lifecycle_costs"]["store_per_byte_second"] = + store_per_byte_second.into(); + wire["implementation"]["require_backend_local_execution"] = + (store_per_byte_second > 0.0).into(); wire["implementation"]["data_snapshot_id"] = "snapshot-topk-test".into(); if certified { wire["implementation"]["accuracy_evidence"] = json!({query: { @@ -36,7 +47,9 @@ fn fixture(certified: bool) -> BackendLocalPlanningInput { // external whole-query fallback or accumulate counter samples as heap weights. #[test] fn rate_heap_candidates_bind_durable_counter_windows() { - let (request, environment) = fixture(true).into_physical_compilation_request().unwrap(); + let (request, environment) = placed_fixture(true, 1.0) + .into_physical_compilation_request() + .unwrap(); let mut families = std::collections::BTreeSet::new(); for candidate in enumerate_exact_and_materialized_candidates(request).unwrap() { let Ok(plan) = DeploymentPlanCompiler.compile_promql(candidate, environment.clone()) else { @@ -122,7 +135,7 @@ fn rate_heap_costs_can_select_each_compiled_candidate() { workload_cost::{manifest, WorkloadCostEvidence, WorkloadQuote}, }; for preferred in ["exact", "CmsWithHeap", "CountSketchWithHeap"] { - let mut input = fixture(true); + let mut input = placed_fixture(true, 1.0); let (request, environment) = input.clone().into_physical_compilation_request().unwrap(); let quotes = enumerate_exact_and_materialized_candidates(request) .unwrap() @@ -233,49 +246,57 @@ fn fixed_window_rate_heap_candidates_install_both_physical_graphs() { ); } -// Rate must precede grouped Sum in both placements; deployment chooses ownership. +// Rate precedes grouped Sum in both placements; the summary-store price +// chooses between maintaining the Sum and rebuilding it per query. #[test] -fn grouped_rate_has_native_query_and_maintenance_candidates() { - let mut wire = serde_json::to_value(fixture(false)).unwrap(); - let entry = &mut wire["query_workload"]["repeating_queries"][0]; - entry["query"] = "sum by(job)(rate(requests_total[1m]))".into(); - entry["requirements"]["accuracy"] = json!({"explicit":"Exact"}); - entry["demand"]["fixed_interval_at"]["interval"] = 5_000.into(); - entry["demand"]["fixed_interval_at"]["evaluation_phase"] = 0.into(); - wire["implementation"]["accuracy_evidence"] = json!({}); - let input: BackendLocalPlanningInput = serde_json::from_value(wire).unwrap(); - let (request, environment) = input.into_physical_compilation_request().unwrap(); - let mut placements = std::collections::BTreeSet::new(); - let mut errors = Vec::new(); - for candidate in enumerate_exact_and_materialized_candidates(request).unwrap() { - let plan = match DeploymentPlanCompiler.compile_promql(candidate, environment.clone()) { - Ok(plan) => plan, - Err(error) => { - errors.push(error.to_string()); +fn grouped_rate_placement_follows_summary_store_cost() { + for (store, maintained) in [(0.0, true), (1.0, false)] { + let mut wire = serde_json::to_value(placed_fixture(false, store)).unwrap(); + let entry = &mut wire["query_workload"]["repeating_queries"][0]; + entry["query"] = "sum by(job)(rate(requests_total[1m]))".into(); + entry["requirements"]["accuracy"] = json!({"explicit":"Exact"}); + entry["demand"]["fixed_interval_at"]["interval"] = 5_000.into(); + entry["demand"]["fixed_interval_at"]["evaluation_phase"] = 0.into(); + wire["implementation"]["accuracy_evidence"] = json!({}); + let input: BackendLocalPlanningInput = serde_json::from_value(wire).unwrap(); + let (request, environment) = input.into_physical_compilation_request().unwrap(); + let mut placements = std::collections::BTreeSet::new(); + let mut errors = Vec::new(); + for candidate in enumerate_exact_and_materialized_candidates(request).unwrap() { + let plan = match DeploymentPlanCompiler.compile_promql(candidate, environment.clone()) { + Ok(plan) => plan, + Err(error) => { + errors.push(error.to_string()); + continue; + } + }; + let entry = plan.query_plan.entries.values().next().unwrap(); + let Some(program) = &entry.physical_dag else { + continue; + }; + if entry.physical_vector_binding().is_none() { continue; } - }; - let entry = plan.query_plan.entries.values().next().unwrap(); - let Some(program) = &entry.physical_dag else { - continue; - }; - if entry.physical_vector_binding().is_none() { - continue; + entry.recover_vector_physical_dag().unwrap(); + let stored = plan + .precompute_plan + .executable_dags + .values() + .any(|dag| !dag.native_programs.is_empty()); + let rebuilt = program.to_string().contains("SummaryBuild"); + assert!(!(stored && rebuilt)); + // Planner also offers a relational Sum over the readouts, which has + // no Sum state to place. + if stored || rebuilt { + placements.insert(stored); + } } - entry.recover_vector_physical_dag().unwrap(); - let stored = plan - .precompute_plan - .executable_dags - .values() - .any(|dag| !dag.native_programs.is_empty()); - assert_eq!(program.to_string().contains("SummaryBuild"), !stored); - placements.insert(stored); + assert_eq!( + placements, + std::collections::BTreeSet::from([maintained]), + "{errors:#?}" + ); } - assert_eq!( - placements, - std::collections::BTreeSet::from([false, true]), - "{errors:#?}" - ); } // Deployment must install the exact retained Planner graphs, including roots, diff --git a/docs/examples/workload-cost-evidence.md b/docs/examples/workload-cost-evidence.md index 4647c7572..2e0b8f296 100644 --- a/docs/examples/workload-cost-evidence.md +++ b/docs/examples/workload-cost-evidence.md @@ -85,7 +85,11 @@ measurements. ## Supported migration scope The default inventory compares the Planner-selected workload with its -whole-workload exact fallback. Within a candidate, whether each summary state +whole-workload exact fallback, maintained-population variants, and one +substitution per native physical realization that Planner's PlanSpace proposes +for a query (for example a heap over per-series Rate readouts). Planner's +global selection, not a quote per alternative, decides every other summary +choice such as the sketch family. Within a candidate, whether each summary state is precomputed or rebuilt at query time is not a separate candidate: the compiler chooses a summary-maintenance lifecycle per unique state from `implementation.lifecycle_costs`, pricing shared state once. Continuously @@ -96,8 +100,10 @@ ephemeral state costs build, read and retirement per read, and is offered only when the deployment can read raw series from Prometheus at query time (not under `require_backend_local_execution`). A query rebuilds all of its states or none; with no state left it runs natively over raw series. The manifest of the -resulting placement is quoted like any other. `workload_cost::select` also -accepts additional Planner-authorized, already-bindable forests. This does not +resulting placement is quoted like any other. A state built from another +retained state's readouts (such as that heap) is placed the same way: retained, +it is maintained over complete per-series states each window; ephemeral, it is +rebuilt per query from the readouts. This does not claim exhaustive search over every lifecycle, engine or Planner algorithm. An exact alternative without an accessible native backend is unavailable even if its numeric quote would be cheap.