From f5485130dafdad75f041cb38c6fd80a87298dce7 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 04:38:29 +0000 Subject: [PATCH 1/5] chore: repin Planner to 9a13ae4 for PromQL Fallback physical coverage Planner #486/#487 compile retained PromQL Fallback subtrees to native operators over raw-series inputs, which query-time lifecycle placement needs. Co-Authored-By: Claude Opus 5.5 --- Cargo.lock | 12 ++++++------ Cargo.toml | 10 +++++----- 2 files changed, 11 insertions(+), 11 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index a4600e45b..f75e5382b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -364,7 +364,7 @@ dependencies = [ [[package]] name = "asap-aware-mapping" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=c98281a59df59740f006ca0b5b8da1ea9bb7fc41#c98281a59df59740f006ca0b5b8da1ea9bb7fc41" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=9a13ae43eecc8487fe6a25a98f6ec09399d2073e#9a13ae43eecc8487fe6a25a98f6ec09399d2073e" dependencies = [ "asap-types", "asap_sketchlib 0.3.0 (git+https://github.com/ProjectASAP/asap_sketchlib)", @@ -376,7 +376,7 @@ dependencies = [ [[package]] name = "asap-frontend-promql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=c98281a59df59740f006ca0b5b8da1ea9bb7fc41#c98281a59df59740f006ca0b5b8da1ea9bb7fc41" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=9a13ae43eecc8487fe6a25a98f6ec09399d2073e#9a13ae43eecc8487fe6a25a98f6ec09399d2073e" dependencies = [ "asap-types", "promql-parser 0.10.0 (git+https://github.com/ProjectASAP/promql-parser?rev=9fede7eecca923c9882fe256484d00d37f8706cb)", @@ -385,7 +385,7 @@ dependencies = [ [[package]] name = "asap-frontend-sql" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=c98281a59df59740f006ca0b5b8da1ea9bb7fc41#c98281a59df59740f006ca0b5b8da1ea9bb7fc41" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=9a13ae43eecc8487fe6a25a98f6ec09399d2073e#9a13ae43eecc8487fe6a25a98f6ec09399d2073e" dependencies = [ "asap-sql-function-catalog", "asap-types", @@ -396,7 +396,7 @@ dependencies = [ [[package]] name = "asap-physical-operators" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=c98281a59df59740f006ca0b5b8da1ea9bb7fc41#c98281a59df59740f006ca0b5b8da1ea9bb7fc41" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=9a13ae43eecc8487fe6a25a98f6ec09399d2073e#9a13ae43eecc8487fe6a25a98f6ec09399d2073e" dependencies = [ "asap-types", "asap_sketchlib 0.3.0 (git+https://github.com/ProjectASAP/asap_sketchlib?rev=5f03ccbd798ed5fec62bdd839bcb331123cab369)", @@ -410,12 +410,12 @@ dependencies = [ [[package]] name = "asap-sql-function-catalog" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=c98281a59df59740f006ca0b5b8da1ea9bb7fc41#c98281a59df59740f006ca0b5b8da1ea9bb7fc41" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=9a13ae43eecc8487fe6a25a98f6ec09399d2073e#9a13ae43eecc8487fe6a25a98f6ec09399d2073e" [[package]] name = "asap-types" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=c98281a59df59740f006ca0b5b8da1ea9bb7fc41#c98281a59df59740f006ca0b5b8da1ea9bb7fc41" +source = "git+https://github.com/ProjectASAP/ASAPPlanner?rev=9a13ae43eecc8487fe6a25a98f6ec09399d2073e#9a13ae43eecc8487fe6a25a98f6ec09399d2073e" dependencies = [ "serde", "serde_json", diff --git a/Cargo.toml b/Cargo.toml index 121c5539a..a416f3a41 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -16,10 +16,10 @@ version = "0.1.0" [workspace.dependencies] # Keep Planner frontends, selection, and IR on the same immutable revision. # Alias upstream asap-types because this workspace also defines asap_types. -planner-types = { package = "asap-types", git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "c98281a59df59740f006ca0b5b8da1ea9bb7fc41" } -asap-aware-mapping = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "c98281a59df59740f006ca0b5b8da1ea9bb7fc41" } -asap-frontend-promql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "c98281a59df59740f006ca0b5b8da1ea9bb7fc41" } -asap-frontend-sql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "c98281a59df59740f006ca0b5b8da1ea9bb7fc41" } +planner-types = { package = "asap-types", git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "9a13ae43eecc8487fe6a25a98f6ec09399d2073e" } +asap-aware-mapping = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "9a13ae43eecc8487fe6a25a98f6ec09399d2073e" } +asap-frontend-promql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "9a13ae43eecc8487fe6a25a98f6ec09399d2073e" } +asap-frontend-sql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "9a13ae43eecc8487fe6a25a98f6ec09399d2073e" } # Shared external deps (used by 2+ crates) serde = { version = "1.0", features = ["derive"] } @@ -39,7 +39,7 @@ arc-swap = "1.7" reqwest = { version = "0.12", default-features = false, features = ["json", "rustls-tls"] } # Internal crates -asap-physical-operators = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "c98281a59df59740f006ca0b5b8da1ea9bb7fc41" } +asap-physical-operators = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "9a13ae43eecc8487fe6a25a98f6ec09399d2073e" } asap_sketch_codec = { path = "crates/asap_sketch_codec" } asap_summary_state = { path = "crates/asap_summary_state" } asap_types = { path = "crates/asap_types" } From 7b2ef9b3959a1854235d837fcaa78276d27db3e4 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 04:59:45 +0000 Subject: [PATCH 2/5] feat: place summary state by backend-priced lifecycle, not masks The backend decided precompute-versus-Prometheus placement by enumerating materialization-key masks and compiling every mask as its own candidate. Placement is now a summary-maintenance lifecycle chosen per unique state: Planner enumerates each state's alternatives and the backend prices them with its own unit costs, charging shared state once with all consumers' demand. ContinuouslyMaintained state is precomputed. Ephemeral state is rebuilt at query time and is offered only when raw series are readable from Prometheus; a query that keeps no state then runs natively over range-selector Scans, compiled once by Planner's PromQL Fallback lowering. Retention adds estimated state bytes x retained panes x the new lifecycle_costs.store_per_byte_second (default 0, preserving placements). Window-layout tests that exercise retained state pin it by disallowing the query-time raw source. Co-Authored-By: Claude Opus 5.5 --- .../examples/calibration_candidates.rs | 6 +- control_plane/src/main.rs | 1 - control_plane/src/physical/compiler.rs | 129 ++++-- .../src/physical/compiler/placement.rs | 397 ++++++++++++++++++ .../src/physical/post_asap/cost_model.rs | 2 +- control_plane/src/physical/workload_cost.rs | 196 ++------- .../materialization_candidates.rs | 104 ----- .../src/physical/workload_cost/status.rs | 42 -- control_plane/tests/lifecycle_placement.rs | 131 ++++++ .../asap_query_engine/post_asap_readout.rs | 3 + docs/examples/workload-cost-evidence.md | 15 +- 11 files changed, 692 insertions(+), 334 deletions(-) create mode 100644 control_plane/src/physical/compiler/placement.rs delete mode 100644 control_plane/src/physical/workload_cost/materialization_candidates.rs create mode 100644 control_plane/tests/lifecycle_placement.rs diff --git a/control_plane/examples/calibration_candidates.rs b/control_plane/examples/calibration_candidates.rs index 928dc1475..330f4822e 100644 --- a/control_plane/examples/calibration_candidates.rs +++ b/control_plane/examples/calibration_candidates.rs @@ -126,7 +126,6 @@ fn main() -> Result<(), Box> { .enumerate() { let queries = candidate.queries.clone(); - let enabled_materialization_keys = candidate.enabled_materialization_keys.clone(); let planner_selected_queries = planner_forest(&queries); let compiled = if metricsql { DeploymentPlanCompiler.compile_metricsql(candidate, environment.clone()) @@ -137,7 +136,7 @@ fn main() -> Result<(), Box> { Ok(plan) => plan, Err(error) => { results.push( - json!({"candidate_index": index, "materialization_policy": enabled_materialization_keys, "planner_selected_queries": planner_selected_queries, "unavailable_reason": error.to_string()}), + json!({"candidate_index": index, "planner_selected_queries": planner_selected_queries, "unavailable_reason": error.to_string()}), ); continue; } @@ -146,14 +145,13 @@ fn main() -> Result<(), Box> { Ok(manifest) => manifest, Err(error) => { results.push( - json!({"candidate_index": index, "materialization_policy": enabled_materialization_keys, "planner_selected_queries": planner_selected_queries, "unavailable_reason": error.to_string()}), + json!({"candidate_index": index, "planner_selected_queries": planner_selected_queries, "unavailable_reason": error.to_string()}), ); continue; } }; results.push(json!({ "candidate_index": index, - "materialization_policy": enabled_materialization_keys, "planner_selected_queries": planner_selected_queries, "manifest": manifest, "lifecycle_estimates": plan.lifecycle_estimates, diff --git a/control_plane/src/main.rs b/control_plane/src/main.rs index 326044635..51775f6b5 100644 --- a/control_plane/src/main.rs +++ b/control_plane/src/main.rs @@ -666,7 +666,6 @@ fn compile_physical_plan_request( allow_mixed_summary_and_exact_execution: request.target == physical::compiler::PhysicalDeploymentTarget::BackendLocalRemoteWrite, require_backend_local_execution: false, - enabled_materialization_keys: None, topk_membership_evidence_by_query_id: request.evidence, exact_composition_costs: request.exact_composition_costs, erp: request.erp, diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index 8078fad2c..d8fff2597 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -40,6 +40,7 @@ use crate::query_plan::{ use crate::types::AccuracyTarget; use planner_types::pre_asap::Source; +mod placement; mod rate_placement; mod windows; pub(super) use windows::gcd; @@ -184,6 +185,15 @@ pub struct LifecycleUnitCosts { pub read: f64, pub retention_per_second: f64, pub retirement: f64, + /// Summary-store price per retained byte-second. Retaining a state costs + /// its estimated bytes times its retained panes times this price, on top + /// of `retention_per_second`. + #[serde(default, skip_serializing_if = "f64_is_zero")] + pub store_per_byte_second: f64, +} + +fn f64_is_zero(value: &f64) -> bool { + *value == 0.0 } #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] @@ -208,9 +218,6 @@ pub struct PhysicalCompilationRequest { pub allow_mixed_summary_and_exact_execution: bool, /// Deployment feasibility: external exact dependencies cannot be bound. pub require_backend_local_execution: bool, - /// Enabled optional candidate keys: None enables all eligible keys; an - /// explicitly empty set enables none. These are not catalog definition IDs. - pub enabled_materialization_keys: Option>, /// Original dashboard demand, in the same order as queries. None is legacy input. pub query_workload: Option, /// Source evidence supplied independently from query demand. @@ -1149,7 +1156,6 @@ impl BackendLocalPlanningInput { require_backend_local_execution: self .physical_inputs .require_backend_local_execution, - enabled_materialization_keys: None, query_workload: Some(workload), data_workload: Some(data_workload), canonical_roots, @@ -1471,6 +1477,7 @@ impl DeploymentPlanCompiler { for (id, root) in planner_types::post_asap::share_common_summary_subtrees(roots) { request.queries[id].selected_plan_root = root; } + let placement = placement::place(&request, &environment, frontend); let population_operators = super::maintained_population::operators(&request)?; let mut compiled_materializations = Vec::with_capacity(request.queries.len()); let mut collector_materializations = Vec::with_capacity(request.queries.len()); @@ -1555,25 +1562,8 @@ impl DeploymentPlanCompiler { .ok() .flatten() .is_some()) - && request.enabled_materialization_keys.as_ref().is_none_or(|policy| { - let key = crate::query_plan::query_time::selected_counter_materialization( - &query.query_string, - &state.node, - ) - .ok() - .flatten() - .or_else(|| { - crate::query_plan::query_time::selected_range_max_materialization( - &query.query_string, - &state.node, - ) - .ok() - .flatten() - }); - // Masks enumerate counter/max choices only. Other selected - // summaries remain required by this physical candidate. - key.is_none_or(|key| policy.contains(&key)) - }) + // An Ephemeral state is rebuilt at query time, not maintained. + && !placement.is_ephemeral(query_index, &state.node) }) .collect::>(); // An exact native fallback has no maintained state and must not @@ -2064,7 +2054,6 @@ impl DeploymentPlanCompiler { stable_workload_plan_id(&plan_materializations, &request.queries).hash(&mut hash); "typed-local-residual-v3-counter-index".hash(&mut hash); super::maintained_population::supported(&request).hash(&mut hash); - request.enabled_materialization_keys.hash(&mut hash); for query in &request.queries { format!("{:?}", query.selected_plan_root).hash(&mut hash); } @@ -2224,7 +2213,9 @@ impl DeploymentPlanCompiler { } else { false }; - let native_rate = if population_operators[query_index].is_none() { + let native_rate = if population_operators[query_index].is_none() + && placement.raw_program(query_index).is_none() + { query .retained_physical()? .map(|physical| { @@ -2257,7 +2248,9 @@ impl DeploymentPlanCompiler { } else { None }; - let mut entry = if let Some((source, program)) = native_rate { + let mut entry = if let Some(raw) = placement.raw_program(query_index) { + Ok(raw_query_time_entry(query, canonical.clone(), raw)?) + } else if let Some((source, program)) = native_rate { let native_state_binding = if let SummaryExpr::SummaryAgg { family: SummaryFamilyType::Sketch(..) @@ -2438,7 +2431,9 @@ impl DeploymentPlanCompiler { }, ) }?; - if request.allow_mixed_summary_and_exact_execution { + if request.allow_mixed_summary_and_exact_execution + && placement.raw_program(query_index).is_none() + { // Any Planner-selected leaf without a physical summary binding // is an exact subtree boundary. Deployed plans never retain a // backend-local range index leaf. @@ -2701,7 +2696,8 @@ impl DeploymentPlanCompiler { crate::query_plan::QueryPlanNode::ExactFallback { .. } | crate::query_plan::QueryPlanNode::Logical { operator: crate::query_plan::query_time::QueryTimeOperator::ExactSubquery { .. } - | crate::query_plan::query_time::QueryTimeOperator::CandidateExactSubquery { .. }, .. })) { + | crate::query_plan::query_time::QueryTimeOperator::CandidateExactSubquery { .. } + | crate::query_plan::query_time::QueryTimeOperator::Scan { .. }, .. })) { return Err(CompileError::Query { query_id: entry.query_id.clone(), reason: "external execution is unavailable in this deployment".into(), @@ -2724,11 +2720,74 @@ impl DeploymentPlanCompiler { storage_routing, lifecycle_estimates: lifecycle_estimates.into_values().collect(), cost_comparison: None, - planner_selection_trace: request.planner_selection_trace, + planner_selection_trace: if placement.trace.is_empty() { + request.planner_selection_trace + } else { + let mut trace = (*request.planner_selection_trace).clone(); + trace.extend(placement.trace); + trace.into() + }, }) } } +/// Native query-time execution of a query that keeps no state: each raw input +/// of the retained program is read by its range-selector `Scan`. +fn raw_query_time_entry( + query: &QueryCompilationInput, + canonical: String, + raw: &placement::RawQueryTimeProgram, +) -> Result { + let mut nodes = BTreeMap::new(); + let mut inputs = Vec::new(); + let mut source_nodes = Vec::new(); + for (ordinal, (slot, scan)) in raw.scans.iter().enumerate() { + let id = crate::query_plan::QueryNodeId(ordinal as u64); + nodes.insert( + id, + crate::query_plan::QueryPlanNode::Logical { + operator: scan.clone(), + inputs: vec![], + }, + ); + inputs.push(id); + source_nodes.push(*slot); + } + let root = crate::query_plan::QueryNodeId(raw.scans.len() as u64); + nodes.insert( + root, + crate::query_plan::QueryPlanNode::Physical { + inputs, + source_nodes, + max_bytes: 64 * 1024 * 1024, + }, + ); + let entry = QueryPlanEntry { + physical_dag: Some( + serde_json::from_slice( + &raw.program + .encode() + .map_err(|e| CompileError::Snapshot(e.to_string()))?, + ) + .map_err(|e| CompileError::Snapshot(e.to_string()))?, + ), + language: crate::query_plan::QueryLanguage::PromQl, + query_id: query.query_id.clone(), + canonical_query: canonical, + fixed_evaluation: None, + root, + nodes, + instant: InstantExecution { + lookback_ms: query.query_lookback_ms, + full_history: false, + cumulative_readout: false, + }, + fallback: FallbackPolicy::ExactBackend, + }; + entry.recover_vector_physical_dag()?; + Ok(entry) +} + fn summary_agg_metric(node: &SummaryNode) -> Option { fn walk(node: &SummaryNode, metrics: &mut BTreeSet) { match &node.expr { @@ -5644,6 +5703,7 @@ pub(crate) mod tests { read: 0.1, retention_per_second: 0.001, retirement: 1.0, + store_per_byte_second: 0.0, }, }; let post_asap = select_post_asap(&parsed, accuracy.clone(), &lifecycle, evidence.as_ref())?; @@ -5657,7 +5717,6 @@ pub(crate) mod tests { canonical_roots: Vec::new(), allow_mixed_summary_and_exact_execution: false, require_backend_local_execution: false, - enabled_materialization_keys: None, query_workload: None, data_workload: None, queries: vec![QueryCompilationInput { @@ -7089,6 +7148,7 @@ pub(crate) mod tests { read: 0.0, retention_per_second: 1.0, retirement: 0.0, + store_per_byte_second: 0.0, }; let full = derived_window_cost( &template, @@ -7686,6 +7746,9 @@ pub(crate) mod tests { continue; } let mut snapshot = planning_snapshot(); + // Window layouts are properties of retained state; with no query-time + // raw source, sparse or costly maintenance cannot move it to query time. + snapshot.physical_inputs.require_backend_local_execution = true; let entry = &mut snapshot.query_workload.repeating_queries.as_mut().unwrap()[0]; entry.query = Query("sum_over_time(a[1m])".into()); entry.requirements.accuracy = @@ -7731,6 +7794,9 @@ pub(crate) mod tests { fn different_cadences_share_common_panes_only_when_cheaper() { for read_cost in [0.0, 1_000_000.0] { let mut snapshot = planning_snapshot(); + // Window layouts are properties of retained state; with no query-time + // raw source, sparse or costly maintenance cannot move it to query time. + snapshot.physical_inputs.require_backend_local_execution = true; snapshot.physical_inputs.lifecycle_costs.read = read_cost; snapshot .physical_inputs @@ -7780,6 +7846,9 @@ pub(crate) mod tests { fn sharing_keeps_profitable_subsets_with_other_phases_or_finer_cadences() { for (third_phase, third_interval) in [(5_000, 20_000), (0, 1_000)] { let mut snapshot = planning_snapshot(); + // Window layouts are properties of retained state; with no query-time + // raw source, sparse or costly maintenance cannot move it to query time. + snapshot.physical_inputs.require_backend_local_execution = true; snapshot .physical_inputs .lifecycle_costs diff --git a/control_plane/src/physical/compiler/placement.rs b/control_plane/src/physical/compiler/placement.rs new file mode 100644 index 000000000..e582a6d83 --- /dev/null +++ b/control_plane/src/physical/compiler/placement.rs @@ -0,0 +1,397 @@ +//! Precompute-or-query-time placement, decided only as a summary-maintenance +//! lifecycle per unique summary state. +//! +//! Planner enumerates each state's lifecycle alternatives; this backend prices +//! them with its own unit costs and picks the cheapest for the whole workload. +//! A `ContinuouslyMaintained` state is precomputed at ingestion. An `Ephemeral` +//! state is rebuilt for each query from raw data read from Prometheus at query +//! time, so it is offered only when that raw source is bindable. +use super::*; +use asap_aware_mapping::{enumerate_summary_maintenance_lifecycles, CostModel}; +use asap_physical_operators::physical_planner::{ + compile, promql_fallback, CompiledPhysicalDag, InputContract, +}; +use planner_types::post_asap::PostAsapOperatorPayload; +use planner_types::pre_asap::{AggIntent, CompareOpKind, ScalarValue}; + +use crate::query_plan::query_time::{LabelMatch, LabelMatcher, QueryTimeOperator}; + +/// A query whose every state is `Ephemeral` keeps nothing between queries. +/// Planner's PromQL Fallback lowering compiles it once over raw-series inputs, +/// each read at query time by the matching range-selector `Scan`. +pub(super) struct RawQueryTimeProgram { + pub(super) program: CompiledPhysicalDag, + pub(super) scans: Vec<(u64, QueryTimeOperator)>, +} + +#[derive(Default)] +pub(super) struct Placement { + ephemeral: Vec>>, + raw: Vec>, + pub(super) trace: Vec, +} + +impl Placement { + pub(super) fn is_ephemeral(&self, query: usize, state: &Rc) -> bool { + self.ephemeral + .get(query) + .is_some_and(|states| states.iter().any(|s| Rc::ptr_eq(s, state))) + } + + pub(super) fn raw_program(&self, query: usize) -> Option<&RawQueryTimeProgram> { + self.raw.get(query).and_then(Option::as_ref) + } +} + +/// Lifecycle prices of one state. Retention charges the state's estimated +/// bytes for every retained pane at the summary-store price. +struct StateCosts { + inputs: SummaryMaintenanceLifecycleCostInputs, + capabilities: SummaryMaintenanceCapabilities, +} + +impl CostModel for StateCosts { + fn rank_candidates( + &self, + _intent: &AggIntent, + candidates: &[SketchAlgorithm], + ) -> Vec { + candidates.to_vec() + } + + fn summary_maintenance_lifecycle_cost_inputs( + &self, + _summary: &SummaryNode, + ) -> SummaryMaintenanceLifecycleCostInputs { + self.inputs.clone() + } + + fn summary_maintenance_capabilities( + &self, + _summary: &SummaryNode, + ) -> SummaryMaintenanceCapabilities { + self.capabilities + } +} + +fn state_costs( + state: &SummaryNode, + costs: &LifecycleUnitCosts, + evaluation_interval_ms: u32, + delete: bool, +) -> StateCosts { + let panes = selected_input_contract(state) + .ok() + .and_then(|(_, window, _)| window) + .map_or(1.0, |seconds| { + (seconds.saturating_mul(1_000) as f64 / f64::from(evaluation_interval_ms.max(1))) + .ceil() + .max(1.0) + }); + let store = match &state.expr { + SummaryExpr::SummaryAgg { family, .. } if costs.store_per_byte_second > 0.0 => { + crate::physical::post_asap::cost_model::analytical_state_bytes(family) + .map(|bytes| bytes * panes * costs.store_per_byte_second) + } + _ => Some(0.0), + }; + StateCosts { + inputs: SummaryMaintenanceLifecycleCostInputs { + build_cost: Some(Cost(costs.build)), + maintenance_cost_per_update: Some(Cost(costs.maintenance_per_update)), + summary_read_cost: Some(Cost(costs.read)), + retention_cost_rate: store.map(|store| CostRate(costs.retention_per_second + store)), + retirement_cost: Some(Cost(costs.retirement)), + }, + capabilities: SummaryMaintenanceCapabilities { + incremental_update: true, + merge: true, + delete, + }, + } +} + +/// Choose one lifecycle per unique state reachable from the (already shared) +/// selected roots. Shared state is priced once with the demand of all its +/// consumers. Without complete workload evidence every state is retained. +pub(super) fn place( + request: &PhysicalCompilationRequest, + environment: &PhysicalDeploymentContext, + frontend: QueryFrontend, +) -> Placement { + let queries = &request.queries; + let mut placement = Placement { + ephemeral: vec![Vec::new(); queries.len()], + raw: (0..queries.len()).map(|_| None).collect(), + trace: Vec::new(), + }; + let (Some(workload), Some(data), Some(first)) = ( + &request.query_workload, + &request.data_workload, + queries.first(), + ) else { + return placement; + }; + // Query-time raw reads and exact subtrees both need the Prometheus source. + let raw_bindable = request.allow_mixed_summary_and_exact_execution + && !request.require_backend_local_execution + && frontend == QueryFrontend::PromQl + && environment.target == PhysicalDeploymentTarget::BackendLocalRemoteWrite; + let mut states: Vec<(Rc, Vec)> = Vec::new(); + let mut query_states = vec![Vec::new(); queries.len()]; + for (index, query) in queries.iter().enumerate() { + if super::super::maintained_population::supported_node(&query.selected_plan_root) { + continue; + } + let Ok(selected) = collect_selected_materializations(&query.selected_plan_root, true) + else { + continue; + }; + for state in selected { + if query_states[index] + .iter() + .any(|known: &Rc| Rc::ptr_eq(known, &state.node)) + { + continue; + } + query_states[index].push(Rc::clone(&state.node)); + match states + .iter_mut() + .find(|(known, _)| Rc::ptr_eq(known, &state.node)) + { + Some((_, consumers)) => consumers.push(index), + None => states.push((state.node, vec![index])), + } + } + } + let raw_programs: Vec> = (0..queries.len()) + .map(|index| { + (raw_bindable && !query_states[index].is_empty()) + .then(|| request.canonical_roots.get(index)) + .flatten() + .and_then(|root| raw_query_time_program(root).ok()) + }) + .collect(); + // A leaf the query-time lowering can externalize reads Prometheus directly. + let externalizable = |state: &SummaryNode, query: usize| { + raw_bindable + && (crate::query_plan::query_time::selected_counter_materialization( + &queries[query].query_string, + state, + ) + .ok() + .flatten() + .is_some() + || crate::query_plan::query_time::selected_range_max_materialization( + &queries[query].query_string, + state, + ) + .ok() + .flatten() + .is_some()) + }; + let horizon = first.summary_lifecycle_inputs.horizon_seconds; + let mut ephemeral = vec![false; states.len()]; + let mut decisions = Vec::new(); + for (state_index, (state, consumers)) in states.iter().enumerate() { + let bindable = consumers + .iter() + .all(|&query| raw_programs[query].is_some() || externalizable(state, query)); + let lead = &queries[consumers[0]].summary_lifecycle_inputs; + let interval = consumers + .iter() + .map(|&query| { + queries[query] + .summary_lifecycle_inputs + .evaluation_interval_ms + }) + .min() + .unwrap_or(lead.evaluation_interval_ms); + let model = state_costs( + state, + &lead.costs, + interval, + environment.target == PhysicalDeploymentTarget::BackendLocalRemoteWrite, + ); + let Ok(candidates) = enumerate_summary_maintenance_lifecycles( + Rc::clone(state), + WorkloadDemand::new_with_data(workload, data, consumers), + environment.observed_at_unix_ms, + Some(Horizon(horizon)), + SummaryMaintenanceLifecycleCapabilities { + supports_ephemeral: bindable, + supports_prepared: false, + supports_shared: false, + supports_continuously_maintained: true, + }, + &model, + ) else { + continue; + }; + let Some(deployment) = candidates + .deployments() + .iter() + .find(|deployment| Rc::ptr_eq(&deployment.summary, state)) + else { + continue; + }; + let cost = |lifecycle: &SummaryMaintenanceLifecycle| { + deployment + .alternatives + .iter() + .find(|alternative| { + alternative.rejection.is_none() + && &alternative.summary_maintenance_lifecycle == lifecycle + }) + .and_then(|alternative| alternative.total_cost) + }; + let retained = cost(&SummaryMaintenanceLifecycle::ContinuouslyMaintained); + let rebuilt = cost(&SummaryMaintenanceLifecycle::Ephemeral); + // An unpriced retained state stays retained unless rebuilding is priced: + // unknown cost never makes a lifecycle win. + ephemeral[state_index] = match (retained, rebuilt) { + (Some(retained), Some(rebuilt)) => rebuilt.0 < retained.0, + (None, Some(_)) => true, + _ => false, + }; + decisions.push((state_index, bindable, retained, rebuilt)); + } + // A query rebuilds either all of its states or none: raw query-time inputs + // and exact subtrees share no snapshot with installed state. Retaining is + // always realizable, so a query that keeps any state keeps all of them. + let index_of = |state: &Rc| states.iter().position(|(s, _)| Rc::ptr_eq(s, state)); + loop { + let mut changed = false; + for (query, owned) in query_states.iter().enumerate() { + let realizable = owned.iter().all(|state| { + index_of(state).is_some_and(|i| ephemeral[i]) + && (raw_programs[query].is_some() || externalizable(state, query)) + }); + if realizable { + continue; + } + for state in owned { + if let Some(i) = index_of(state).filter(|&i| ephemeral[i]) { + ephemeral[i] = false; + changed = true; + } + } + } + if !changed { + break; + } + } + for (state_index, bindable, retained, rebuilt) in decisions { + let (state, consumers) = &states[state_index]; + placement.trace.push(json!({ + "stage": "deployment.lifecycle_placement", + "query_ids": consumers.iter().map(|&q| &queries[q].query_id).collect::>(), + "logical_root_id": crate::planner_selection::explained_root_id(state, &queries[consumers[0]].accuracy_target), + "ephemeral_bindable": bindable, + "continuously_maintained_cost": retained.map(|cost| cost.0), + "ephemeral_cost": rebuilt.map(|cost| cost.0), + "selected": if ephemeral[state_index] { "ephemeral" } else { "continuously_maintained" }, + })); + } + for (query, owned) in query_states.iter().enumerate() { + let chosen: Vec<_> = owned + .iter() + .filter(|state| index_of(state).is_some_and(|i| ephemeral[i])) + .cloned() + .collect(); + if !owned.is_empty() && chosen.len() == owned.len() { + placement.raw[query] = raw_programs[query].as_ref().map(|raw| RawQueryTimeProgram { + program: raw.program.clone(), + scans: raw.scans.clone(), + }); + } + placement.ephemeral[query] = chosen; + } + placement +} + +/// Compile the whole query over raw-series inputs and name the Prometheus +/// range selector that supplies each input. +fn raw_query_time_program(root: &QueryExpr) -> Result { + let typed = asap_physical_operators::physical_planner::promql_rows::with_series_identity(root) + .map_err(|error| error.to_string())?; + let keep = crate::planner_selection::keep_pre_asap(&typed).map_err(|e| e.to_string())?; + let dag = planner_types::post_asap::compile_post_asap_dag(&keep).map_err(|e| e.to_string())?; + let mut inputs = BTreeMap::new(); + let mut scans = Vec::new(); + for node in &dag.nodes { + let PostAsapOperatorPayload::Fallback { expression } = &node.payload else { + continue; + }; + let selectors = promql_fallback::raw_series(expression).map_err(|e| e.to_string())?; + for (ordinal, (selector, schema)) in selectors.into_iter().enumerate() { + let slot = promql_fallback::raw_series_input(u64::from(node.id.0), ordinal); + scans.push((slot, range_selector_scan(&selector)?)); + inputs.insert(slot, InputContract::bounded(schema)); + } + } + if scans.is_empty() { + return Err("query reads no raw series".into()); + } + let program = compile(&dag, inputs, &[u64::from(dag.root.0)]).map_err(|e| e.to_string())?; + Ok(RawQueryTimeProgram { program, scans }) +} + +/// `TimeRange { range, [TimeShift { offset }], Scan }` as a Prometheus range +/// selector. `@` modifiers and non-label predicates have no such selector. +fn range_selector_scan(selector: &QueryExpr) -> Result { + let QueryExpr::TimeRange { range, child } = selector else { + return Err("raw input is not a range selector".into()); + }; + let (offset_ms, scan) = match child.as_ref() { + QueryExpr::TimeShift { shift, child } if shift.at.is_none() => { + (shift.offset_ms, child.as_ref()) + } + QueryExpr::TimeShift { .. } => return Err("raw selector uses @".into()), + scan => (0, scan), + }; + let QueryExpr::Scan { + source: Source::TimeSeries { metric }, + predicates, + schema, + } = scan + else { + return Err("raw input does not read a named time series".into()); + }; + let matchers = predicates + .iter() + .map(|predicate| { + let QueryExpr::Compare { left, op, right } = predicate.0.as_ref() else { + return Err("raw selector predicate is not a label comparison".to_string()); + }; + let (QueryExpr::Column(column), QueryExpr::Literal(ScalarValue::Utf8(value))) = + (left.as_ref(), right.as_ref()) + else { + return Err("raw selector predicate must compare a label with a string".into()); + }; + let name = schema + .columns + .get(*column) + .map(|field| field.name.clone()) + .ok_or("raw selector predicate names an unknown label")?; + let operation = match op { + CompareOpKind::Eq => LabelMatch::Equal, + CompareOpKind::Ne => LabelMatch::NotEqual, + CompareOpKind::Regex => LabelMatch::Regex, + CompareOpKind::NotRegex => LabelMatch::NotRegex, + _ => return Err("raw selector predicate is not a PromQL matcher".into()), + }; + Ok(LabelMatcher { + name, + value: value.clone(), + operation, + }) + }) + .collect::>()?; + Ok(QueryTimeOperator::Scan { + metric: (!metric.is_empty()).then(|| metric.clone()), + matchers, + range_ms: Some(u64::try_from(range.as_millis()).map_err(|e| e.to_string())?), + offset_ms, + }) +} diff --git a/control_plane/src/physical/post_asap/cost_model.rs b/control_plane/src/physical/post_asap/cost_model.rs index 98ee99864..47034249b 100644 --- a/control_plane/src/physical/post_asap/cost_model.rs +++ b/control_plane/src/physical/post_asap/cost_model.rs @@ -42,7 +42,7 @@ pub struct CandidateCostEstimate { pub erp_record_ids: Vec, } -fn analytical_state_bytes(family: &SummaryFamilyType) -> Option { +pub(crate) fn analytical_state_bytes(family: &SummaryFamilyType) -> Option { use asap_types::AggregationType as A; use planner_types::post_asap::ExactKind; let (aggregation, params) = match family { diff --git a/control_plane/src/physical/workload_cost.rs b/control_plane/src/physical/workload_cost.rs index 88b0a607b..ed946d6d9 100644 --- a/control_plane/src/physical/workload_cost.rs +++ b/control_plane/src/physical/workload_cost.rs @@ -4,9 +4,8 @@ //! Planner owns semantic legality. A manifest describes the exact physical //! demand to price; provider quotes and candidate evaluations are separate. -mod materialization_candidates; mod status; -pub use status::{CandidateEvaluationStatus, CandidateSearchScope}; +pub use status::CandidateEvaluationStatus; #[cfg(test)] use super::compiler::DeploymentPlanCompiler; @@ -97,24 +96,11 @@ pub struct CandidatePlanEvaluation { pub unavailable_reason: Option, } -#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] -pub struct MaterializationSearchCoverage { - #[serde(rename = "eligible_leaves")] - pub eligible_materialization_count: usize, - #[serde(rename = "enumerated_local_masks")] - pub enumerated_candidate_key_sets: usize, - pub exhaustive: bool, - #[serde(rename = "scope")] - pub search_scope: CandidateSearchScope, -} - #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] pub struct CandidatePlanSelectionReport { #[serde(default)] #[serde(rename = "logical_selection")] pub planner_selection_trace: std::sync::Arc>, - #[serde(default)] - pub materialization_search_coverage: Option, pub data_snapshot_id: String, pub model_version: String, pub selected_plan_id: u64, @@ -294,16 +280,20 @@ pub fn manifest( let source = json!({"source": planner_types::pre_asap::Source::TimeSeries { metric: population.metric.clone() }, "location": "backend", "ingest": plan.precompute_plan.ingest}); add(format!("source:{source}"), source, "horizon", 1.0); } - if matches!( - node, - crate::query_plan::QueryPlanNode::Logical { - operator: crate::query_plan::query_time::QueryTimeOperator::Scan { .. }, - .. - } - ) { - return Err(invalid( - "generic backend raw scans are outside the ASAP/Prometheus execution contract", - )); + if let crate::query_plan::QueryPlanNode::Logical { + operator: crate::query_plan::query_time::QueryTimeOperator::Scan { metric, .. }, + .. + } = node + { + // Only a native program reads raw series, from Prometheus at + // query time; its per-evaluation read is the query node below. + let (Some(metric), Some(_)) = (metric, &entry.physical_dag) else { + return Err(invalid( + "generic backend raw scans are outside the ASAP/Prometheus execution contract", + )); + }; + let source = json!({"source": planner_types::pre_asap::Source::TimeSeries { metric: metric.clone() }, "location": "exact_backend"}); + add(format!("source:{source}"), source, "horizon", 1.0); } if let crate::query_plan::QueryPlanNode::Logical { operator: @@ -545,7 +535,6 @@ fn candidate_description<'a>( &( &logical_root_ids, candidate.allow_mixed_summary_and_exact_execution, - &candidate.enabled_materialization_keys, ), ) }), @@ -706,23 +695,6 @@ fn select_candidates( "candidate inventory must contain 1..=4096 candidates", )); } - let candidate_key_sets: BTreeSet<_> = candidates - .iter() - .filter(|c| c.allow_mixed_summary_and_exact_execution) - .filter_map(|c| c.enabled_materialization_keys.clone()) - .collect(); - let eligible_keys: BTreeSet<_> = candidate_key_sets - .iter() - .flat_map(|p| p.iter().cloned()) - .collect(); - let materialization_search_coverage = - (!candidate_key_sets.is_empty()).then(|| MaterializationSearchCoverage { - eligible_materialization_count: eligible_keys.len(), - enumerated_candidate_key_sets: candidate_key_sets.len(), - exhaustive: eligible_keys.len() < usize::BITS as usize - && candidate_key_sets.len() == (1usize << eligible_keys.len()), - search_scope: CandidateSearchScope::PlannerAuthorizedMaterializations, - }); let planner_selection_trace = candidates[0].planner_selection_trace.clone(); let mut comparison_workload = None; let mut candidate_evaluations = Vec::new(); @@ -793,7 +765,6 @@ fn select_candidates( candidate_evaluations[best_index].status = CandidateEvaluationStatus::Selected; plan.cost_comparison = Some(CandidatePlanSelectionReport { planner_selection_trace, - materialization_search_coverage, data_snapshot_id: evidence.data_snapshot_id.clone(), model_version: evidence.model_version.clone(), selected_plan_id: plan.envelope.plan_id, @@ -823,8 +794,6 @@ pub fn enumerate_exact_and_materialized_candidates( if !candidates.iter().any(|existing| { existing.allow_mixed_summary_and_exact_execution == candidate.allow_mixed_summary_and_exact_execution - && existing.enabled_materialization_keys - == candidate.enabled_materialization_keys && existing .queries .iter() @@ -893,7 +862,6 @@ fn enumerate_frontier_candidates( } if maintained_roots.iter().all(Option::is_some) { candidate.allow_mixed_summary_and_exact_execution = false; - candidate.enabled_materialization_keys = None; } Some(candidate) }) @@ -902,8 +870,6 @@ fn enumerate_frontier_candidates( if !candidates.iter().any(|existing| { existing.allow_mixed_summary_and_exact_execution == candidate.allow_mixed_summary_and_exact_execution - && existing.enabled_materialization_keys - == candidate.enabled_materialization_keys && existing .queries .iter() @@ -934,12 +900,14 @@ fn enumerate_frontier_candidates( Ok(candidates) } +/// The Planner-selected candidate and the whole-workload native exact +/// alternative. Placement within a candidate is a lifecycle decision made +/// during compilation, so no per-state variants are enumerated here. fn materialization_candidates( request: PhysicalCompilationRequest, ) -> Result, CompileError> { let mut exact = request.clone(); exact.allow_mixed_summary_and_exact_execution = false; - exact.enabled_materialization_keys = None; for (index, query) in exact.queries.iter_mut().enumerate() { let parsed = if let Some(root) = request.canonical_roots.get(index) { root.as_ref().clone() @@ -962,42 +930,7 @@ fn materialization_candidates( { Ok(vec![request]) } else { - if !request.allow_mixed_summary_and_exact_execution - || request.enabled_materialization_keys.is_some() - { - return Ok(vec![request, exact]); - } - let mut keys = BTreeSet::new(); - for query in &request.queries { - match crate::query_plan::query_time::eligible_materialization_keys( - &query.query_string, - &query.selected_plan_root, - ) { - Ok(found) => keys.extend(found), - // A failed local projection must not make the native candidate - // disappear. Compile/select_lowest_cost_candidate retains its concrete unavailability. - Err(_) => return Ok(vec![request, exact]), - } - } - if keys.is_empty() { - return Ok(vec![request, exact]); - } - let inventory = materialization_candidates::enumerate(keys); - debug_assert_eq!( - inventory.exhaustive, - inventory.eligible_materialization_count <= 4 - ); - let mut candidate_requests: Vec<_> = inventory - .candidate_key_sets - .into_iter() - .map(|enabled_keys| { - let mut candidate = request.clone(); - candidate.enabled_materialization_keys = Some(enabled_keys); - candidate - }) - .collect(); - candidate_requests.push(exact); - Ok(candidate_requests) + Ok(vec![request, exact]) } } @@ -1425,59 +1358,37 @@ mod tests { ); } + /// One query keeps all of its states or rebuilds all of them at query + /// time, so a binary never mixes installed state with exact subtrees. #[test] - fn materialization_masks_reject_mixed_snapshots_before_costing() { - let mut snapshot = fixture(); - let q = &mut snapshot.query_workload.repeating_queries.as_mut().unwrap()[0]; - q.query = - planner_types::workload::Query("max_over_time(a[1m]) + max_over_time(b[1m])".into()); - q.requirements.accuracy = planner_types::workload::AccuracyRequirement::Explicit( - crate::types::AccuracyTarget::Exact, - ); - let (request, environment) = snapshot.into_physical_compilation_request().unwrap(); - let candidates = enumerate_frontier_candidates(request).unwrap(); - assert_eq!( - candidates.len(), - 5, - "four proposed candidate key sets plus native" - ); - let mut identities = BTreeSet::new(); - let mut rejected = 0; - for candidate in &candidates[..4] { - let enabled = candidate - .enabled_materialization_keys - .as_ref() - .unwrap() - .len(); - let result = - DeploymentPlanCompiler.compile_promql(candidate.clone(), environment.clone()); - if enabled == 1 { - let Err(error) = result else { - panic!("mixed snapshots must fail binding"); - }; - assert!( - error.to_string().contains("common snapshot proof"), - "{error}" - ); - rejected += 1; - continue; - } - let plan = result.unwrap(); - assert!(identities.insert(plan.envelope.plan_id)); - let cost = manifest(&plan, &candidate.queries).unwrap(); - assert_eq!( - cost.components - .keys() - .filter(|k| k.starts_with("state:backend:")) - .count(), - enabled * 4 + fn placement_never_mixes_installed_and_query_time_operands() { + for (store, materialized) in [(0.0, 2), (1.0, 0)] { + let mut snapshot = fixture(); + snapshot + .physical_inputs + .lifecycle_costs + .store_per_byte_second = store; + let q = &mut snapshot.query_workload.repeating_queries.as_mut().unwrap()[0]; + q.query = planner_types::workload::Query( + "max_over_time(a[1m]) + max_over_time(b[1m])".into(), ); + q.requirements.accuracy = planner_types::workload::AccuracyRequirement::Explicit( + crate::types::AccuracyTarget::Exact, + ); + let (request, environment) = snapshot.into_physical_compilation_request().unwrap(); + let candidates = enumerate_frontier_candidates(request).unwrap(); + assert_eq!(candidates.len(), 2, "selected candidate plus native"); + let plan = DeploymentPlanCompiler + .compile_promql(candidates[0].clone(), environment) + .unwrap(); + assert_eq!(plan.precompute_plan.materializations.len(), materialized); + let cost = manifest(&plan, &candidates[0].queries).unwrap(); assert_eq!( cost.components .keys() - .filter(|k| k.starts_with("raw-state:")) + .filter(|k| k.starts_with("state:backend:")) .count(), - 0 + materialized * 4 ); assert_eq!( cost.components @@ -1486,25 +1397,10 @@ mod tests { && v.implementation.get("location").and_then(Value::as_str) == Some("exact_backend")) .count(), - 2 - enabled + 2 - materialized ); - assert!(!plan - .query_plan - .entries - .values() - .any(|entry| entry.nodes.values().any(|node| matches!( - node, - crate::query_plan::QueryPlanNode::ExactFallback { .. } - )))); + assert!(!candidates[1].allow_mixed_summary_and_exact_execution); } - assert_eq!(rejected, 2); - assert_eq!(identities.len(), 2); - assert!( - !candidates - .last() - .unwrap() - .allow_mixed_summary_and_exact_execution - ); } fn quoted() -> ( diff --git a/control_plane/src/physical/workload_cost/materialization_candidates.rs b/control_plane/src/physical/workload_cost/materialization_candidates.rs deleted file mode 100644 index 2344192f2..000000000 --- a/control_plane/src/physical/workload_cost/materialization_candidates.rs +++ /dev/null @@ -1,104 +0,0 @@ -//! Bounded physical implementation search; semantic leaf legality is established -//! by the Planner-witness lowering before these stable keys are supplied. -use std::collections::BTreeSet; - -#[derive(Debug)] -pub(super) struct MaterializationCandidateSets { - pub candidate_key_sets: Vec>, - pub exhaustive: bool, - pub eligible_materialization_count: usize, -} - -/// Enumerate every enabled_keys for up to four leaves. Larger forests retain all-materialized, -/// all-exact, then singleton/complement pairs in stable key order. The caller must -/// disclose bounded coverage; no unenumerated optimum is claimed. Reserve one -/// of the selector's 64 candidate slots for native execution. -pub(super) fn enumerate(keys: BTreeSet) -> MaterializationCandidateSets { - let eligible_materialization_count = keys.len(); - let ordered: Vec<_> = keys.iter().cloned().collect(); - let exhaustive = eligible_materialization_count <= 4; - let mut candidate_key_sets = vec![keys.clone()]; - if exhaustive { - for bits in 0..(1usize << eligible_materialization_count) { - let enabled_keys = ordered - .iter() - .enumerate() - .filter(|(i, _)| bits & (1 << i) != 0) - .map(|(_, key)| key.clone()) - .collect(); - if !candidate_key_sets.contains(&enabled_keys) { - candidate_key_sets.push(enabled_keys); - } - } - } else { - candidate_key_sets.push(BTreeSet::new()); - for key in ordered { - for enabled_keys in [ - BTreeSet::from([key.clone()]), - keys.difference(&BTreeSet::from([key])).cloned().collect(), - ] { - if candidate_key_sets.len() >= 63 { - break; - } - if !candidate_key_sets.contains(&enabled_keys) { - candidate_key_sets.push(enabled_keys); - } - } - } - } - MaterializationCandidateSets { - candidate_key_sets, - exhaustive, - eligible_materialization_count, - } -} - -#[cfg(test)] -mod tests { - use super::*; - #[test] - fn small_forests_cover_every_mixed_path_once() { - let result = enumerate(BTreeSet::from(["a".into(), "b".into()])); - assert!(result.exhaustive); - assert_eq!(result.eligible_materialization_count, 2); - assert_eq!(result.candidate_key_sets.len(), 4); - assert_eq!( - result - .candidate_key_sets - .iter() - .collect::>() - .len(), - 4 - ); - assert!(result - .candidate_key_sets - .contains(&BTreeSet::from(["a".into()]))); - assert!(result - .candidate_key_sets - .contains(&BTreeSet::from(["b".into()]))); - } - #[test] - fn large_inventory_reserves_native_slot_and_discloses_truncation() { - let keys = (0..100).map(|i| format!("{i:03}")).collect(); - let result = enumerate(keys); - assert!(!result.exhaustive); - assert_eq!(result.candidate_key_sets.len(), 63); - assert_eq!(result.candidate_key_sets[0].len(), 100); - assert!(result.candidate_key_sets[1].is_empty()); - assert_eq!( - result - .candidate_key_sets - .iter() - .collect::>() - .len(), - 63 - ); - } - #[test] - fn no_materialization_has_one_exact_implementation() { - assert_eq!( - enumerate(BTreeSet::new()).candidate_key_sets, - vec![BTreeSet::new()] - ); - } -} diff --git a/control_plane/src/physical/workload_cost/status.rs b/control_plane/src/physical/workload_cost/status.rs index e73051e8a..4927884a3 100644 --- a/control_plane/src/physical/workload_cost/status.rs +++ b/control_plane/src/physical/workload_cost/status.rs @@ -55,38 +55,6 @@ impl From for String { } } -// Preserve the existing report text, including historical terminology, on the -// wire. The typed variant describes the actual scope for new Rust consumers. -const MATERIALIZATION_SEARCH_SCOPE: &str = "Backend materialization versus Prometheus exact-subquery masks over Planner-authorized leaves; native alternative separate; bounded inventory does not claim an unenumerated optimum"; - -#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] -#[serde(from = "String", into = "String")] -pub enum CandidateSearchScope { - PlannerAuthorizedMaterializations, - Other(String), -} - -impl From for CandidateSearchScope { - fn from(value: String) -> Self { - if value == MATERIALIZATION_SEARCH_SCOPE { - Self::PlannerAuthorizedMaterializations - } else { - Self::Other(value) - } - } -} - -impl From for String { - fn from(value: CandidateSearchScope) -> Self { - match value { - CandidateSearchScope::PlannerAuthorizedMaterializations => { - MATERIALIZATION_SEARCH_SCOPE.to_owned() - } - CandidateSearchScope::Other(value) => value, - } - } -} - #[cfg(test)] mod tests { use super::*; @@ -111,14 +79,4 @@ mod tests { assert_eq!(serde_json::to_value(status).unwrap(), value); } } - - /// Old search descriptions retain their exact serialized representation. - #[test] - fn scope_wire_values_round_trip() { - for value in [MATERIALIZATION_SEARCH_SCOPE, "future_scope"] { - let scope: CandidateSearchScope = - serde_json::from_value(serde_json::json!(value)).unwrap(); - assert_eq!(serde_json::to_value(scope).unwrap(), value); - } - } } diff --git a/control_plane/tests/lifecycle_placement.rs b/control_plane/tests/lifecycle_placement.rs new file mode 100644 index 000000000..e89b2751f --- /dev/null +++ b/control_plane/tests/lifecycle_placement.rs @@ -0,0 +1,131 @@ +//! Precompute-or-query-time placement follows the backend's lifecycle costs. +use control_plane::physical::compiler::{ + BackendLocalPlanningInput, CompiledPhysicalPlan, DeploymentPlanCompiler, +}; +use control_plane::physical::workload_cost::enumerate_exact_and_materialized_candidates; +use control_plane::query_plan::{query_time::QueryTimeOperator, QueryPlanNode}; +use serde_json::Value; + +fn fixture(store_per_byte_second: f64, require_local: bool) -> BackendLocalPlanningInput { + let mut wire: Value = serde_json::from_str(include_str!( + "../../docs/examples/asapquery-planning-snapshot.json" + )) + .unwrap(); + wire["implementation"]["lifecycle_costs"]["store_per_byte_second"] = + store_per_byte_second.into(); + wire["implementation"]["require_backend_local_execution"] = require_local.into(); + serde_json::from_value(wire).unwrap() +} + +/// The Planner-selected candidate, compiled; the native exact alternative is last. +fn selected_plan(input: BackendLocalPlanningInput) -> CompiledPhysicalPlan { + let (request, environment) = input.into_physical_compilation_request().unwrap(); + let candidate = enumerate_exact_and_materialized_candidates(request) + .unwrap() + .into_iter() + .next() + .unwrap(); + assert!(candidate.allow_mixed_summary_and_exact_execution); + DeploymentPlanCompiler + .compile_promql(candidate, environment) + .unwrap() +} + +fn raw_scans(plan: &CompiledPhysicalPlan) -> Vec<&QueryTimeOperator> { + plan.query_plan + .entries + .values() + .flat_map(|entry| entry.nodes.values()) + .filter_map(|node| match node { + QueryPlanNode::Logical { + operator: operator @ QueryTimeOperator::Scan { .. }, + .. + } => Some(operator), + _ => None, + }) + .collect() +} + +fn placements(plan: &CompiledPhysicalPlan) -> Vec { + plan.planner_selection_trace + .iter() + .filter(|entry| entry["stage"] == "deployment.lifecycle_placement") + .map(|entry| entry["selected"].as_str().unwrap().to_owned()) + .collect() +} + +// A cheap summary store keeps the quantile sketch continuously maintained. +#[test] +fn cheap_summary_store_precomputes_the_state() { + let plan = selected_plan(fixture(0.0, false)); + assert_eq!(placements(&plan), ["continuously_maintained"]); + assert_eq!(plan.precompute_plan.materializations.len(), 1); + assert!(raw_scans(&plan).is_empty()); +} + +// An expensive summary store rebuilds the state per query from raw Prometheus series. +#[test] +fn expensive_summary_store_moves_the_state_to_query_time() { + let plan = selected_plan(fixture(1.0, false)); + assert_eq!(placements(&plan), ["ephemeral"]); + assert!(plan.precompute_plan.materializations.is_empty()); + let entry = plan.query_plan.entries.values().next().unwrap(); + entry.recover_vector_physical_dag().unwrap(); + assert_eq!( + raw_scans(&plan), + [&QueryTimeOperator::Scan { + metric: Some("m".into()), + matchers: vec![], + range_ms: Some(60_000), + offset_ms: 0, + }] + ); +} + +// Without a query-time raw source, the state stays maintained whatever the store costs. +#[test] +fn ephemeral_requires_a_bindable_raw_source() { + let plan = selected_plan(fixture(1.0, true)); + assert_eq!(placements(&plan), ["continuously_maintained"]); + assert_eq!(plan.precompute_plan.materializations.len(), 1); + assert!(raw_scans(&plan).is_empty()); +} + +fn decisions(queries: &[&str]) -> Vec { + let mut wire = serde_json::to_value(fixture(0.0, false)).unwrap(); + let template = wire["query_workload"]["repeating_queries"][0].clone(); + wire["query_workload"]["repeating_queries"] = queries + .iter() + .map(|query| { + let mut entry = template.clone(); + entry["query"] = (*query).into(); + entry["requirements"]["accuracy"] = serde_json::json!({"explicit": "Exact"}); + entry + }) + .collect(); + let plan = selected_plan(serde_json::from_value(wire).unwrap()); + plan.planner_selection_trace + .iter() + .filter(|entry| entry["stage"] == "deployment.lifecycle_placement") + .cloned() + .collect() +} + +// A state shared by two queries is one lifecycle decision: maintenance and +// retention are charged once, while both queries' reads are counted. +#[test] +fn shared_state_is_priced_once_with_all_reads() { + let [alone] = decisions(&["sum_over_time(m[1m])"]).try_into().unwrap(); + let [shared] = decisions(&["sum_over_time(m[1m])", "sum_over_time(m[1m]) * 2"]) + .try_into() + .unwrap(); + assert_eq!(shared["query_ids"].as_array().unwrap().len(), 2); + let cost = |decision: &Value, field: &str| decision[field].as_f64().unwrap(); + // Rebuilding scales with reads; maintaining adds only the extra read cost. + let extra_reads = cost(&shared, "ephemeral_cost") / cost(&alone, "ephemeral_cost"); + assert_eq!(extra_reads, 2.0); + assert!( + cost(&shared, "continuously_maintained_cost") + < 2.0 * cost(&alone, "continuously_maintained_cost") + ); +} diff --git a/data_plane/src/query_engines/asap_query_engine/post_asap_readout.rs b/data_plane/src/query_engines/asap_query_engine/post_asap_readout.rs index 6c959d0cb..16ef57150 100644 --- a/data_plane/src/query_engines/asap_query_engine/post_asap_readout.rs +++ b/data_plane/src/query_engines/asap_query_engine/post_asap_readout.rs @@ -1856,6 +1856,9 @@ mod tests { entry["demand"]["fixed_interval_at"] = serde_json::json!({ "interval": evaluation_secs * 1_000, "evaluation_phase": phase_ms }); + // Window schedules belong to retained state; without a + // query-time raw source sparse reads cannot move it to query time. + snapshot["implementation"]["require_backend_local_execution"] = true.into(); let snapshot: BackendLocalPlanningInput = serde_json::from_value(snapshot).unwrap(); let (mut request, env) = snapshot.into_physical_compilation_request().unwrap(); diff --git a/docs/examples/workload-cost-evidence.md b/docs/examples/workload-cost-evidence.md index fc5430d55..4647c7572 100644 --- a/docs/examples/workload-cost-evidence.md +++ b/docs/examples/workload-cost-evidence.md @@ -84,8 +84,19 @@ measurements. ## Supported migration scope -The default inventory compares the Planner-selected continuously maintained -workload with its whole-workload exact fallback. `workload_cost::select` also +The default inventory compares the Planner-selected workload with its +whole-workload exact fallback. 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 +maintained state costs build, per-update maintenance over the ingestion rate, +reads, retention and retirement; retention adds the state's estimated bytes +times its retained panes times `store_per_byte_second` (default 0). An +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 claim exhaustive search over every lifecycle, engine or Planner algorithm. An exact alternative without an accessible native backend is unavailable even From a163a57cd1f176848f699f9d785d4d3a8440ae99 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 04:59:46 +0000 Subject: [PATCH 3/5] test: summary-store price flips placement end to end A cheap store serves the query from precomputed state without raw reads; an expensive store makes the backend read the range selector from Prometheus at query time and compute the quantile itself. Co-Authored-By: Claude Opus 5.5 --- .../asapquery_compatibility_process_e2e.rs | 3 + .../support/lifecycle_placement_process.rs | 182 ++++++++++++++++++ 2 files changed, 185 insertions(+) create mode 100644 data_plane/tests/support/lifecycle_placement_process.rs diff --git a/data_plane/tests/asapquery_compatibility_process_e2e.rs b/data_plane/tests/asapquery_compatibility_process_e2e.rs index fa69c4ccf..7bb99e832 100644 --- a/data_plane/tests/asapquery_compatibility_process_e2e.rs +++ b/data_plane/tests/asapquery_compatibility_process_e2e.rs @@ -38,6 +38,9 @@ mod current_series_process; #[path = "support/issue_701_702_process.rs"] mod issue_701_702_process; +#[path = "support/lifecycle_placement_process.rs"] +mod lifecycle_placement_process; + // Test-only quotes preserve the fixture's local candidate without a production bypass. fn quote_snapshot_for_test( snapshot: control_plane::physical::compiler::BackendLocalPlanningInput, diff --git a/data_plane/tests/support/lifecycle_placement_process.rs b/data_plane/tests/support/lifecycle_placement_process.rs new file mode 100644 index 000000000..a65763dc8 --- /dev/null +++ b/data_plane/tests/support/lifecycle_placement_process.rs @@ -0,0 +1,182 @@ +//! The summary-store price, not an enumerated placement, decides whether a +//! state is precomputed at ingestion or rebuilt from raw series at query time. +use super::*; +use control_plane::physical::compiler::BackendLocalPlanningInput; + +const QUERY: &str = "quantile_over_time(0.99, m[1m])"; +const RAW_SELECTOR: &str = r#"{__name__="m"}[60000ms]"#; + +struct Deployment { + backend: String, + requests: Arc>>>, + _child: ChildGuard, + _output: tempfile::TempDir, + prometheus: tokio::task::JoinHandle<()>, +} + +/// Start the production binary from the fixture priced with `store` per +/// retained byte-second. Its Prometheus serves `samples` of `m` for raw reads. +async fn deploy(store: f64, samples: Vec<(i64, f64)>) -> Deployment { + let mut fixture: Value = serde_json::from_str(include_str!( + "../../../docs/examples/asapquery-planning-snapshot.json" + )) + .unwrap(); + fixture["implementation"]["scrape_interval_ms"] = 1000.into(); + fixture["data_workload"]["data_ingestion_interval"]["value"] = 1000.into(); + fixture["implementation"]["lifecycle_costs"]["store_per_byte_second"] = store.into(); + let snapshot: BackendLocalPlanningInput = serde_json::from_value(fixture).unwrap(); + let priced = quote_snapshot_for_test(snapshot); + + let requests = Arc::new(Mutex::new(Vec::>::new())); + let recorded = requests.clone(); + let matrix = serde_json::json!({"status": "success", "data": {"resultType": "matrix", + "result": [{"metric": {"__name__": "m", "instance": "a"}, + "values": samples.iter().map(|(ms, value)| serde_json::json!([*ms as f64 / 1000.0, value.to_string()])).collect::>()}]}}); + let empty = + serde_json::json!({"status": "success", "data": {"resultType": "vector", "result": []}}); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let prometheus_url = format!("http://{}", listener.local_addr().unwrap()); + let prometheus = tokio::spawn(async move { + axum::serve( + listener, + Router::new() + .route("/-/healthy", get(|| async { "healthy" })) + .route( + "/api/v1/query", + get(move |Query(params): Query>| { + let recorded = recorded.clone(); + let raw = params.get("query").map(String::as_str) == Some(RAW_SELECTOR); + let response = if raw { matrix.clone() } else { empty.clone() }; + async move { + recorded.lock().await.push(params); + Json(response) + } + }), + ), + ) + .await + .unwrap(); + }); + let output = tempfile::tempdir().unwrap(); + let path = output.path().join("planning.json"); + std::fs::write(&path, serde_json::to_vec(&priced).unwrap()).unwrap(); + let port = unused_port(); + let mut child = ChildGuard( + Command::new(env!("CARGO_BIN_EXE_data_plane")) + .args(["--profile", "asapquery", "--planning-snapshot"]) + .arg(path) + .args([ + "--prometheus-server", + &prometheus_url, + "--forward-unsupported-queries", + "--precompute-allowed-lateness-ms", + "0", + "--precompute-flush-interval-ms", + "25", + "--http-port", + &port.to_string(), + "--output-dir", + ]) + .arg(output.path()) + .stdout(Stdio::null()) + .stderr(Stdio::inherit()) + .spawn() + .unwrap(), + ); + let backend = format!("http://127.0.0.1:{port}"); + wait_until_ready( + &reqwest::Client::new(), + &format!("{backend}/api/v1/health"), + &mut child.0, + ) + .await; + Deployment { + backend, + requests, + _child: child, + _output: output, + prometheus, + } +} + +fn origin_ms() -> i64 { + let now = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap() + .as_millis() as i64; + now - now.rem_euclid(600_000) - 1_200_000 +} + +// A cheap summary store precomputes the sketch: the query is answered from +// maintained state and never reads raw series from Prometheus. +#[tokio::test] +async fn cheap_summary_store_answers_from_precomputed_state() { + let origin = origin_ms(); + let deployment = deploy(0.0, vec![]).await; + let client = reqwest::Client::new(); + let samples: Vec<_> = (1..=121) + .map(|i| (origin + i * 1000, (1 + i % 7) as f64)) + .collect(); + let wire = WriteRequest { + timeseries: vec![series_with_labels("m", &[("instance", "a")], &samples)], + }; + assert_eq!(remote_write(&client, &deployment.backend, &wire).await, 204); + let output = deployment._output.path().join("query_engine.log"); + let at = (origin + 120_000) as f64 / 1000.0; + let response = wait_for_warm_instant(&client, &deployment.backend, QUERY, at, &output).await; + assert!(first_value(&response, "value").is_some(), "{response}"); + assert!( + !deployment + .requests + .lock() + .await + .iter() + .any(|request| request.get("query").map(String::as_str) == Some(RAW_SELECTOR)), + "precomputed state must not read raw series" + ); + deployment.prometheus.abort(); +} + +// An expensive summary store rebuilds the state per query: the backend reads the +// range selector's raw series from Prometheus and computes the answer itself. +#[tokio::test] +async fn expensive_summary_store_rebuilds_state_from_raw_series() { + let at_ms = origin_ms() + 120_000; + let values = [5.0, 1.0, 4.0, 2.0, 3.0]; + let samples = values + .iter() + .enumerate() + .map(|(i, value)| (at_ms - 50_000 + i as i64 * 10_000, *value)) + .collect(); + let deployment = deploy(1.0, samples).await; + let response: Value = reqwest::Client::new() + .get(format!("{}/api/v1/query", deployment.backend)) + .query(&[ + ("query", QUERY.to_string()), + ("time", (at_ms as f64 / 1000.0).to_string()), + ]) + .send() + .await + .unwrap() + .json() + .await + .unwrap(); + // PromQL quantile_over_time(0.99) of 1..=5: rank 3.96 between 4 and 5. + assert_eq!(first_value(&response, "value"), Some(4.96), "{response}"); + let requests = deployment.requests.lock().await; + assert!( + requests.iter().any(|request| { + request.get("query").map(String::as_str) == Some(RAW_SELECTOR) + && request.get("time").map(String::as_str) + == Some(format!("{:.3}", at_ms as f64 / 1000.0).as_str()) + }), + "query time must read raw series: {requests:?}" + ); + assert!( + !requests + .iter() + .any(|request| request.get("query").map(String::as_str) == Some(QUERY)), + "the backend, not Prometheus, answers the query: {requests:?}" + ); + deployment.prometheus.abort(); +} From db891bfb49676be5c71ff58cfa6164b6df991175 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 05:04:57 +0000 Subject: [PATCH 4/5] refactor: price lifecycle inputs per summary in one backend cost model One CostModel prices every summary a Planner enumeration lists, so a root with several states can be enumerated and bound in one call. Co-Authored-By: Claude Opus 5.5 --- .../src/physical/compiler/placement.rs | 138 +++++++++--------- 1 file changed, 70 insertions(+), 68 deletions(-) diff --git a/control_plane/src/physical/compiler/placement.rs b/control_plane/src/physical/compiler/placement.rs index e582a6d83..bf5fdd3fd 100644 --- a/control_plane/src/physical/compiler/placement.rs +++ b/control_plane/src/physical/compiler/placement.rs @@ -43,14 +43,16 @@ impl Placement { } } -/// Lifecycle prices of one state. Retention charges the state's estimated -/// bytes for every retained pane at the summary-store price. -struct StateCosts { - inputs: SummaryMaintenanceLifecycleCostInputs, - capabilities: SummaryMaintenanceCapabilities, +/// Backend lifecycle prices for any summary state. Retention charges the +/// state's estimated bytes for every retained pane at the summary-store price; +/// an unknown size under a positive price leaves retention unpriced. +struct LifecycleCosts<'a> { + costs: &'a LifecycleUnitCosts, + evaluation_interval_ms: u32, + delete: bool, } -impl CostModel for StateCosts { +impl CostModel for LifecycleCosts<'_> { fn rank_candidates( &self, _intent: &AggIntent, @@ -61,53 +63,67 @@ impl CostModel for StateCosts { fn summary_maintenance_lifecycle_cost_inputs( &self, - _summary: &SummaryNode, + summary: &SummaryNode, ) -> SummaryMaintenanceLifecycleCostInputs { - self.inputs.clone() - } - - fn summary_maintenance_capabilities( - &self, - _summary: &SummaryNode, - ) -> SummaryMaintenanceCapabilities { - self.capabilities - } -} - -fn state_costs( - state: &SummaryNode, - costs: &LifecycleUnitCosts, - evaluation_interval_ms: u32, - delete: bool, -) -> StateCosts { - let panes = selected_input_contract(state) - .ok() - .and_then(|(_, window, _)| window) - .map_or(1.0, |seconds| { - (seconds.saturating_mul(1_000) as f64 / f64::from(evaluation_interval_ms.max(1))) + let costs = self.costs; + let panes = selected_input_contract(summary) + .ok() + .and_then(|(_, window, _)| window) + .map_or(1.0, |seconds| { + (seconds.saturating_mul(1_000) as f64 + / f64::from(self.evaluation_interval_ms.max(1))) .ceil() .max(1.0) - }); - let store = match &state.expr { - SummaryExpr::SummaryAgg { family, .. } if costs.store_per_byte_second > 0.0 => { - crate::physical::post_asap::cost_model::analytical_state_bytes(family) - .map(|bytes| bytes * panes * costs.store_per_byte_second) - } - _ => Some(0.0), - }; - StateCosts { - inputs: SummaryMaintenanceLifecycleCostInputs { + }); + let store = match &summary.expr { + SummaryExpr::SummaryAgg { family, .. } if costs.store_per_byte_second > 0.0 => { + crate::physical::post_asap::cost_model::analytical_state_bytes(family) + .map(|bytes| bytes * panes * costs.store_per_byte_second) + } + _ => Some(0.0), + }; + SummaryMaintenanceLifecycleCostInputs { build_cost: Some(Cost(costs.build)), maintenance_cost_per_update: Some(Cost(costs.maintenance_per_update)), summary_read_cost: Some(Cost(costs.read)), retention_cost_rate: store.map(|store| CostRate(costs.retention_per_second + store)), retirement_cost: Some(Cost(costs.retirement)), - }, - capabilities: SummaryMaintenanceCapabilities { + } + } + + fn summary_maintenance_capabilities( + &self, + _summary: &SummaryNode, + ) -> SummaryMaintenanceCapabilities { + SummaryMaintenanceCapabilities { incremental_update: true, merge: true, - delete, - }, + delete: self.delete, + } + } +} + +/// Selectable total cost of `lifecycle` for `summary`, if Planner listed it. +fn alternative_cost( + deployment: &asap_aware_mapping::SummaryMaintenanceDeployment, + lifecycle: &SummaryMaintenanceLifecycle, +) -> Option { + deployment + .alternatives + .iter() + .find(|alternative| { + alternative.rejection.is_none() + && &alternative.summary_maintenance_lifecycle == lifecycle + }) + .and_then(|alternative| alternative.total_cost) +} + +/// Unknown cost never makes a lifecycle win. +fn rebuild_is_cheaper(retained: Option, rebuilt: Option) -> bool { + match (retained, rebuilt) { + (Some(retained), Some(rebuilt)) => rebuilt.0 < retained.0, + (None, Some(_)) => true, + _ => false, } } @@ -207,12 +223,11 @@ pub(super) fn place( }) .min() .unwrap_or(lead.evaluation_interval_ms); - let model = state_costs( - state, - &lead.costs, - interval, - environment.target == PhysicalDeploymentTarget::BackendLocalRemoteWrite, - ); + let model = LifecycleCosts { + costs: &lead.costs, + evaluation_interval_ms: interval, + delete: environment.target == PhysicalDeploymentTarget::BackendLocalRemoteWrite, + }; let Ok(candidates) = enumerate_summary_maintenance_lifecycles( Rc::clone(state), WorkloadDemand::new_with_data(workload, data, consumers), @@ -235,25 +250,12 @@ pub(super) fn place( else { continue; }; - let cost = |lifecycle: &SummaryMaintenanceLifecycle| { - deployment - .alternatives - .iter() - .find(|alternative| { - alternative.rejection.is_none() - && &alternative.summary_maintenance_lifecycle == lifecycle - }) - .and_then(|alternative| alternative.total_cost) - }; - let retained = cost(&SummaryMaintenanceLifecycle::ContinuouslyMaintained); - let rebuilt = cost(&SummaryMaintenanceLifecycle::Ephemeral); - // An unpriced retained state stays retained unless rebuilding is priced: - // unknown cost never makes a lifecycle win. - ephemeral[state_index] = match (retained, rebuilt) { - (Some(retained), Some(rebuilt)) => rebuilt.0 < retained.0, - (None, Some(_)) => true, - _ => false, - }; + let retained = alternative_cost( + deployment, + &SummaryMaintenanceLifecycle::ContinuouslyMaintained, + ); + let rebuilt = alternative_cost(deployment, &SummaryMaintenanceLifecycle::Ephemeral); + ephemeral[state_index] = rebuild_is_cheaper(retained, rebuilt); decisions.push((state_index, bindable, retained, rebuilt)); } // A query rebuilds either all of its states or none: raw query-time inputs From bed9d911c02fd6bad51eb4fcc458a90a163e3313 Mon Sep 17 00:00:00 2001 From: zzylol Date: Wed, 30 Sep 2026 05:45:02 +0000 Subject: [PATCH 5/5] fix: keep retention when a lifecycle is unpriced; price partitions An unpriced retained state was moved to query time whenever rebuilding was priced; it now stays retained, the placement used without lifecycle evidence. Retention scales by input cardinality for per-series and grouped state, and raw selectors without a metric name are not offered as query-time sources, since their Scan cannot be priced. Co-Authored-By: Claude Opus 5.5 --- .../src/physical/compiler/placement.rs | 56 +++++++++++++++---- 1 file changed, 45 insertions(+), 11 deletions(-) diff --git a/control_plane/src/physical/compiler/placement.rs b/control_plane/src/physical/compiler/placement.rs index bf5fdd3fd..4e46eede1 100644 --- a/control_plane/src/physical/compiler/placement.rs +++ b/control_plane/src/physical/compiler/placement.rs @@ -44,11 +44,14 @@ impl Placement { } /// Backend lifecycle prices for any summary state. Retention charges the -/// state's estimated bytes for every retained pane at the summary-store price; -/// an unknown size under a positive price leaves retention unpriced. +/// state's estimated bytes for every retained pane and partition at the +/// summary-store price; an unknown size under a positive price leaves +/// retention unpriced. Panes are estimated from the evaluation interval +/// because the window layout is chosen only for retained state. struct LifecycleCosts<'a> { costs: &'a LifecycleUnitCosts, evaluation_interval_ms: u32, + input_cardinality: Option, delete: bool, } @@ -76,9 +79,22 @@ impl CostModel for LifecycleCosts<'_> { .max(1.0) }); let store = match &summary.expr { - SummaryExpr::SummaryAgg { family, .. } if costs.store_per_byte_second > 0.0 => { + SummaryExpr::SummaryAgg { + family, reduction, .. + } if costs.store_per_byte_second > 0.0 => { + // Per-series and grouped state keeps one instance per input series at most. + let partitioned = + matches!(reduction, planner_types::pre_asap::Reduction::PerEntity) + || reduction + .group_keys() + .is_some_and(|keys| !keys.keys().is_empty()); + let partitions = if partitioned { + self.input_cardinality.unwrap_or(1).max(1) as f64 + } else { + 1.0 + }; crate::physical::post_asap::cost_model::analytical_state_bytes(family) - .map(|bytes| bytes * panes * costs.store_per_byte_second) + .map(|bytes| bytes * panes * partitions * costs.store_per_byte_second) } _ => Some(0.0), }; @@ -118,13 +134,10 @@ fn alternative_cost( .and_then(|alternative| alternative.total_cost) } -/// Unknown cost never makes a lifecycle win. +/// Rebuilding must be priced cheaper than retaining. An unpriced alternative +/// never displaces retention, the placement used without lifecycle evidence. fn rebuild_is_cheaper(retained: Option, rebuilt: Option) -> bool { - match (retained, rebuilt) { - (Some(retained), Some(rebuilt)) => rebuilt.0 < retained.0, - (None, Some(_)) => true, - _ => false, - } + matches!((retained, rebuilt), (Some(retained), Some(rebuilt)) if rebuilt.0 < retained.0) } /// Choose one lifecycle per unique state reachable from the (already shared) @@ -226,6 +239,10 @@ pub(super) fn place( let model = LifecycleCosts { costs: &lead.costs, evaluation_interval_ms: interval, + input_cardinality: data + .input_cardinality + .value_at(environment.observed_at_unix_ms) + .copied(), delete: environment.target == PhysicalDeploymentTarget::BackendLocalRemoteWrite, }; let Ok(candidates) = enumerate_summary_maintenance_lifecycles( @@ -360,6 +377,9 @@ fn range_selector_scan(selector: &QueryExpr) -> Result Result>()?; Ok(QueryTimeOperator::Scan { - metric: (!metric.is_empty()).then(|| metric.clone()), + metric: Some(metric.clone()), matchers, range_ms: Some(u64::try_from(range.as_millis()).map_err(|e| e.to_string())?), offset_ms, }) } + +#[cfg(test)] +mod tests { + use super::*; + + // Only a priced, strictly cheaper rebuild moves a state to query time. + #[test] + fn unpriced_alternatives_never_displace_retention() { + assert!(rebuild_is_cheaper(Some(Cost(2.0)), Some(Cost(1.0)))); + assert!(!rebuild_is_cheaper(Some(Cost(1.0)), Some(Cost(1.0)))); + assert!(!rebuild_is_cheaper(None, Some(Cost(1.0)))); + assert!(!rebuild_is_cheaper(Some(Cost(1.0)), None)); + } +}