diff --git a/control_plane/src/main.rs b/control_plane/src/main.rs index 51775f6b5..fbbe1b3c0 100644 --- a/control_plane/src/main.rs +++ b/control_plane/src/main.rs @@ -661,6 +661,7 @@ fn compile_physical_plan_request( planner_selection_trace: planner_selection_trace.into(), query_workload: Some(query_workload), data_workload: Some(request.data_workload), + source_ingestion_rates: Default::default(), canonical_roots, queries, allow_mixed_summary_and_exact_execution: request.target diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index 573bf6764..deec6bf6c 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -216,6 +216,8 @@ pub struct PhysicalCompilationRequest { pub query_workload: Option, /// Source evidence supplied independently from query demand. pub data_workload: Option, + /// Per-metric samples/second used to price raw selector folds. + pub source_ingestion_rates: BTreeMap>, /// Workload-lowered roots retained across physical candidate enumeration. pub canonical_roots: Vec>, pub queries: Vec, @@ -438,6 +440,10 @@ pub struct BackendLocalPhysicalInputs { #[serde(default, skip_serializing_if = "std::ops::Not::not")] pub require_backend_local_execution: bool, pub lifecycle_costs: LifecycleUnitCosts, + /// Samples/second by metric, before label filtering. Missing or stale + /// evidence falls back to the conservative workload-wide ingestion rate. + #[serde(default, skip_serializing_if = "BTreeMap::is_empty")] + pub source_ingestion_rates: BTreeMap>, pub evidence_observed_at_unix_ms: u64, pub evidence_valid_for_ms: u64, pub horizon_seconds: f64, @@ -855,6 +861,18 @@ impl BackendLocalPlanningInput { "compatibility workload requires backend_local_remote_write target".into(), )); } + for (metric, evidence) in &self.physical_inputs.source_ingestion_rates { + if metric.trim().is_empty() + || evidence + .value + .is_some_and(|rate| !rate.0.is_finite() || rate.0 < 0.0) + { + return Err(CompileError::Snapshot( + "source_ingestion_rates requires named metrics and finite nonnegative rates" + .into(), + )); + } + } let workload = self.query_workload; let mut data_workload = self.data_workload.clone(); let scoped_snapshot_id = self @@ -1185,6 +1203,7 @@ impl BackendLocalPlanningInput { require_backend_local_execution: self.physical_inputs.require_backend_local_execution, query_workload: Some(workload), data_workload: Some(data_workload), + source_ingestion_rates: self.physical_inputs.source_ingestion_rates, canonical_roots, queries, topk_membership_evidence_by_query_id: topk_evidence_by_id, @@ -6043,6 +6062,7 @@ pub(crate) mod tests { require_backend_local_execution: false, query_workload: None, data_workload: None, + source_ingestion_rates: BTreeMap::new(), queries: vec![QueryCompilationInput { query_id: query_id.into(), query_string: promql.into(), @@ -8944,6 +8964,7 @@ pub(crate) mod tests { physical_inputs: BackendLocalPhysicalInputs { require_backend_local_execution: false, lifecycle_costs: template.summary_lifecycle_inputs.costs, + source_ingestion_rates: BTreeMap::new(), evidence_observed_at_unix_ms: 9_500, evidence_valid_for_ms: 60_000, horizon_seconds: 300.0, diff --git a/control_plane/src/physical/compiler/placement.rs b/control_plane/src/physical/compiler/placement.rs index c2221fff3..91e6bd1d9 100644 --- a/control_plane/src/physical/compiler/placement.rs +++ b/control_plane/src/physical/compiler/placement.rs @@ -163,8 +163,7 @@ impl CostModel for LifecycleCosts { /// The window implementation compilation installs for `state` when `query` /// retains it, possibly `beside_raw` inputs, with the number of states that -/// layout keeps in the store for the state's own window. A derived state -/// reading a longer window over it retains more, so this is a lower bound there. +/// layout keeps for both query readouts and downstream maintenance windows. fn installed_window( query: &QueryCompilationInput, state: &SelectedMaterialization, @@ -196,8 +195,22 @@ fn installed_window( id.as_ref() == Some(&candidate.realization_id) && &candidate.framework == framework })? .clone(); + let maintenance_lookback = query_states + .iter() + .filter(|consumer| { + immutable_materialization_sources(&consumer.node) + .is_some_and(|sources| sources.iter().any(|source| Rc::ptr_eq(source, &state.node))) + }) + .map(|consumer| { + consumer + .window_secs + .map(|seconds| seconds.saturating_mul(1_000)) + .unwrap_or(query.query_lookback_ms) + }) + .max() + .unwrap_or(0); let retained = retained_state_count( - branch.query_lookback_ms, + branch.query_lookback_ms.max(maintenance_lookback), retention_margin_ms, window.slide_secs.saturating_mul(1_000), &window.layout, @@ -205,6 +218,66 @@ fn installed_window( Some((window, retained)) } +/// Identity of a raw, non-cohort installation after window selection. This +/// uses the deployment fingerprint so logical lookbacks sharing panes share +/// maintenance cost too. Derived/native cohorts are priced by their own layouts. +fn raw_installation_id( + request: &PhysicalCompilationRequest, + query_index: usize, + state: &SelectedMaterialization, + states: &[SelectedMaterialization], + window: &WindowRealizationCandidate, + environment: &PhysicalDeploymentContext, +) -> Option { + let asap_types::WindowMaterializationLayout::Pane { pane_secs } = window.layout else { + return None; + }; + if !window.derived + || pane_secs == 0 + || leaf_selector(&state.node).is_none() + || !matches!( + state.family, + SummaryFamilyType::ExactAggregate(planner_types::post_asap::ExactKind::Sum, _) + ) + || windows::cohort_nodes(states).contains(&(Rc::as_ptr(&state.node) as usize)) + { + return None; + } + let query = &request.queries[query_index]; + // With identical pane widths, reuse keeps all read costs and strictly + // saves one producer iff its non-read implementation cost is positive. + let mut maintenance = query.summary_lifecycle_inputs.clone(); + maintenance.costs.read = 0.0; + let cost = derived_window_cost( + &window.cost, + &maintenance, + window.window_secs, + window.slide_secs, + &window.layout, + request.query_retention_margin_ms, + ) + .weighted_cost; + if !cost.is_finite() || cost <= 0.0 { + return None; + } + let aggregation = + physical_aggregation(query, state, query.query_id.clone(), environment.target); + let (mut config, computation) = scoped_materialization(&aggregation, &state.node).ok()?; + config.window_size = pane_secs; + config.slide_interval = pane_secs; + config.window_type = asap_types::WindowKind::Tumbling; + config.window_layout = window.layout.clone(); + config.pane_origin_ms = shared_pane_origin_ms( + request.query_workload.as_ref(), + [query_index], + pane_secs.saturating_mul(1_000), + ) + .ok()?; + config.pane_origin_ms?; + config.allocate_stored_output_id(&computation); + Some(config.policy_fingerprint()) +} + /// A state's index, raw bindability, retained and rebuilt costs, and /// retained store states. type Decision = (usize, bool, Option, Option, Option); @@ -248,6 +321,24 @@ fn rebuild_is_cheaper(retained: Option, rebuilt: Option) -> bool { matches!((retained, rebuilt), (Some(retained), Some(rebuilt)) if rebuilt.0 < retained.0) } +/// Apply the same source observation to retained maintenance and raw folds. +fn source_data( + data: &DataWorkload, + rates: &BTreeMap>, + metric: &str, + now: u64, +) -> DataWorkload { + let mut scoped = data.clone(); + if let Some(evidence) = rates.get(metric).filter(|evidence| { + evidence + .value_at(now) + .is_some_and(|rate| rate.0.is_finite() && rate.0 >= 0.0) + }) { + scoped.ingestion_rate = evidence.clone(); + } + scoped +} + /// 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. @@ -315,7 +406,22 @@ pub(super) fn place( .collect(); let horizon = Some(Horizon(first.summary_lifecycle_inputs.horizon_seconds)); let now = environment.observed_at_unix_ms; - let update_rate = data.ingestion_rate.value_at(now).map(|rate| rate.0); + let update_rate = |scan: &QueryTimeOperator| { + let metric_rate = match scan { + QueryTimeOperator::Scan { + metric: Some(metric), + .. + } => request + .source_ingestion_rates + .get(metric) + .and_then(|evidence| evidence.value_at(now)), + _ => None, + }; + metric_rate + .or_else(|| data.ingestion_rate.value_at(now)) + .filter(|rate| rate.0.is_finite() && rate.0 >= 0.0) + .map(|rate| rate.0) + }; let model = |costs: LifecycleUnitCosts, retained_states, evaluation_interval_ms| LifecycleCosts { costs, @@ -329,6 +435,10 @@ pub(super) fn place( consumers: &[usize], lifecycle: SummaryMaintenanceLifecycle, model: &LifecycleCosts| { + let scoped = selected_input_contract(state) + .ok() + .map(|(metric, _, _)| source_data(data, &request.source_ingestion_rates, &metric, now)); + let data = scoped.as_ref().unwrap_or(data); let candidates = enumerate_summary_maintenance_lifecycles( Rc::clone(state), WorkloadDemand::new_with_data(workload, data, consumers), @@ -437,20 +547,19 @@ pub(super) fn place( .iter() .map(|&query| { let raw = raw_programs[query].as_ref()?; - let scanned_seconds = raw + let scanned_samples = raw .scans .iter() .map(|(_, scan)| match scan { QueryTimeOperator::Scan { range_ms: Some(range_ms), .. - } => Some(*range_ms as f64 / 1_000.0), + } => Some(update_rate(scan)? * *range_ms as f64 / 1_000.0), _ => None, }) .sum::>()?; let lifecycle = &queries[query].summary_lifecycle_inputs; - let fold = - update_rate? * scanned_seconds * lifecycle.costs.maintenance_per_update; + let fold = scanned_samples * lifecycle.costs.maintenance_per_update; let costs = LifecycleUnitCosts { build: lifecycle.costs.build + fold / query_states[query].len() as f64, ..lifecycle.costs.clone() @@ -474,6 +583,73 @@ pub(super) fn place( .sum::>(); decisions.push((state_index, bindable, retained, rebuilt, retained_states)); } + // Distinct logical lookbacks can install the same raw panes. Within one + // query their placement moves together, so charge shared maintenance and + // the largest retention once; each logical read still pays its read cost. + let mut physical_groups = BTreeMap::<(usize, asap_types::PolicyFingerprint), Vec>::new(); + for (index, (state, consumers)) in states.iter().enumerate() { + let [query] = consumers.as_slice() else { + continue; + }; + let selected = &selected_states[*query]; + let Some(state) = selected + .iter() + .find(|selected| Rc::ptr_eq(&selected.node, state)) + else { + continue; + }; + let Some((window, _)) = installed_window( + &queries[*query], + state, + selected, + environment, + request.query_retention_margin_ms, + false, + ) else { + continue; + }; + if let Some(identity) = + raw_installation_id(request, *query, state, selected, &window, environment) + { + physical_groups + .entry((*query, identity)) + .or_default() + .push(index); + } + } + for ((query, _), group) in physical_groups + .into_iter() + .filter(|(_, group)| group.len() > 1) + { + if group + .iter() + .any(|&index| decisions[index].2.is_none() || decisions[index].4.is_none()) + { + continue; + } + let owner = *group + .iter() + .max_by_key(|&&index| decisions[index].4) + .unwrap(); + let lifecycle = &queries[query].summary_lifecycle_inputs; + for index in group.into_iter().filter(|&index| index != owner) { + let costs = LifecycleUnitCosts { + build: 0.0, + maintenance_per_update: 0.0, + read: lifecycle.costs.read, + retention_per_second: 0.0, + retirement: 0.0, + store_per_byte_second: 0.0, + }; + decisions[index].2 = price( + &states[index].0, + &[query], + SummaryMaintenanceLifecycle::ContinuouslyMaintained, + &model(costs, Some(0), lifecycle.evaluation_interval_ms), + ); + decisions[index].4 = Some(0); + } + } // First choose a group baseline: states linked through a query move // together when rebuilding all consumers costs less than retaining them. // The bounded-lag search below may replace it with an admissible mix. @@ -528,19 +704,17 @@ pub(super) fn place( // Rebuilt beside retained state, a program folds only the samples of // the rebuilt states' own selectors. let alone = |member: usize| { - let ( - _, - QueryTimeOperator::Scan { - range_ms: Some(range_ms), - .. - }, - ) = leaf_selector(&states[member].0)? + let (_, scan) = leaf_selector(&states[member].0)?; + let QueryTimeOperator::Scan { + range_ms: Some(range_ms), + .. + } = &scan else { return None; }; let costs = LifecycleUnitCosts { build: lifecycle.costs.build - + update_rate? * range_ms as f64 / 1_000.0 + + update_rate(&scan)? * *range_ms as f64 / 1_000.0 * lifecycle.costs.maintenance_per_update, ..lifecycle.costs.clone() }; @@ -1221,6 +1395,16 @@ pub(super) fn time_native_candidate( trace: Vec::new(), }) }; + let scoped_data = match &query.legacy_query_source { + Source::TimeSeries { metric } => source_data( + data, + &inputs.source_ingestion_rates, + metric, + environment.observed_at_unix_ms, + ), + _ => data.clone(), + }; + let data = &scoped_data; let lifecycle = &query.summary_lifecycle_inputs; // Enumerate lifecycle choices first; each retained native alternative is // repriced below using its hypothetical installed window layout. @@ -1353,6 +1537,42 @@ fn find(root: &Rc, summary: &SummaryNode) -> Option mod tests { use super::*; + // A derived consumer's longer window retains and prices the source panes + // that maintenance needs, even when the source has a shorter own readout. + #[test] + fn derived_consumer_extends_priced_source_retention() { + let mut wire: Value = serde_json::from_str(include_str!( + "../../../../docs/examples/asapquery-planning-snapshot.json" + )) + .unwrap(); + wire["query_workload"]["repeating_queries"][0]["query"] = + "quantile(0.9, sum_over_time(m[1m]))".into(); + let snapshot: BackendLocalPlanningInput = serde_json::from_value(wire).unwrap(); + let (request, environment) = snapshot.into_physical_compilation_request().unwrap(); + let query = &request.queries[0]; + let mut states = + collect_selected_materializations(&query.selected_plan_root, true).unwrap(); + let derived = states + .iter_mut() + .find(|state| immutable_materialization_sources(&state.node).is_some()) + .expect("derived consumer"); + let source = immutable_materialization_sources(&derived.node) + .unwrap() + .remove(0); + derived.window_secs = Some(600); + let source = states + .iter() + .find(|state| Rc::ptr_eq(&state.node, &source)) + .unwrap(); + let (window, retained) = + installed_window(query, source, &states, &environment, 0, false).unwrap(); + assert_eq!( + retained, + retained_state_count(600_000, 0, window.slide_secs * 1_000, &window.layout) + ); + assert_eq!(retained, 11); + } + fn option(beside_raw: Option, alone: Option) -> MixedOption { MixedOption { beside_raw: beside_raw.map(Cost), diff --git a/control_plane/tests/lifecycle_placement.rs b/control_plane/tests/lifecycle_placement.rs index e8c9b965f..0c34680af 100644 --- a/control_plane/tests/lifecycle_placement.rs +++ b/control_plane/tests/lifecycle_placement.rs @@ -387,3 +387,124 @@ fn inadmissible_mix_keeps_the_group_decision() { } } } + +// Explicit per-metric sample-rate evidence prices each selector's raw fold; +// absent evidence retains the conservative workload-wide rate. +#[test] +fn raw_selector_folds_use_their_own_source_rates() { + let mut wire = serde_json::to_value(fixture(1.0, false)).unwrap(); + wire["query_workload"]["repeating_queries"][0]["query"] = MIXED_QUERY.into(); + let baseline = selected_plan(serde_json::from_value(wire.clone()).unwrap()); + let retained_cost = |plan: &CompiledPhysicalPlan| { + plan.planner_selection_trace + .iter() + .filter(|event| event["stage"] == "deployment.lifecycle_placement") + .map(|event| cost(event, "continuously_maintained_cost")) + .sum::() + }; + let mut a = wire["data_workload"]["ingestion_rate"].clone(); + a["value"] = 1.0.into(); + let mut b = a.clone(); + b["value"] = 2.0.into(); + wire["implementation"]["source_ingestion_rates"] = serde_json::json!({"a": a, "b": b}); + let input: BackendLocalPlanningInput = serde_json::from_value(wire).unwrap(); + let costs = input.physical_inputs.lifecycle_costs.clone(); + let evaluations = input.physical_inputs.horizon_seconds / 10.0; + let plan = selected_plan(input); + let priced: f64 = plan + .planner_selection_trace + .iter() + .filter(|event| event["stage"] == "deployment.lifecycle_placement") + .map(|event| cost(event, "ephemeral_cost")) + .sum(); + let expected = evaluations + * (2.0 * (costs.build + costs.read + costs.retirement) + + (60.0 + 2.0 * 600.0) * costs.maintenance_per_update); + assert!((priced - expected).abs() < 1e-9, "{priced} != {expected}"); + let expected_saving = evaluations * 10.0 * (200.0 - 3.0) * costs.maintenance_per_update; + assert!((retained_cost(&baseline) - retained_cost(&plan) - expected_saving).abs() < 1e-9); +} + +// Native lifecycle candidates price the complete windows their retained +// realization installs, including the configured retention margin. +#[test] +fn native_lifecycle_prices_complete_window_retention() { + for (margin, expected) in [(0, 1), (25_000, 4)] { + let mut wire = serde_json::to_value(fixture(1e-12, false)).unwrap(); + wire["query_workload"]["repeating_queries"][0]["query"] = "sum(rate(m[1m]))".into(); + wire["implementation"]["query_staleness_margin_ms"] = margin.into(); + let (request, _) = serde_json::from_value::(wire) + .unwrap() + .into_physical_compilation_request() + .unwrap(); + let native: Vec<_> = request + .planner_selection_trace + .iter() + .filter(|event| event["stage"] == "deployment.native_candidate_placement") + .collect(); + assert!(!native.is_empty()); + for event in native { + assert_eq!(event["retained_states"].as_u64(), Some(expected), "{event}"); + } + } +} + +// Stale per-metric evidence uses the workload rate, and invalid observations +// are rejected instead of allowing a negative raw execution cost. +#[test] +fn source_rate_evidence_requires_fresh_nonnegative_values() { + let mut input = fixture(1.0, false); + let baseline = decision(&selected_plan(input.clone())); + let mut rate = input.data_workload.ingestion_rate.clone(); + rate.value = Some(planner_types::workload::Rate(1.0)); + rate.observed_at_ms = Some(input.environment.observed_at_unix_ms.saturating_sub(1)); + rate.valid_for_ms = Some(0); + input + .physical_inputs + .source_ingestion_rates + .insert("m".into(), rate.clone()); + let stale = decision(&selected_plan(input.clone())); + assert_eq!( + cost(&stale, "ephemeral_cost"), + cost(&baseline, "ephemeral_cost") + ); + rate.value = Some(planner_types::workload::Rate(-1.0)); + input + .physical_inputs + .source_ingestion_rates + .insert("m".into(), rate); + let error = input.into_physical_compilation_request().unwrap_err(); + assert!(error.to_string().contains("finite nonnegative rates")); +} + +// Different logical lookbacks sharing one installed raw producer charge its +// maintenance and maximum retained pane population once, while keeping both reads. +#[test] +fn shared_physical_panes_are_priced_once_across_logical_lookbacks() { + let plan = two_state_plan( + "sum(sum_over_time(m[1m])) / sum(sum_over_time(m[10m]))", + 1e-12, + ); + assert_eq!(plan.precompute_plan.materializations.len(), 1); + let installed = plan.precompute_plan.materializations[0] + .num_aggregates_to_retain + .unwrap(); + assert_eq!(installed, 61); + let priced: u64 = plan + .planner_selection_trace + .iter() + .filter(|event| event["stage"] == "deployment.lifecycle_placement") + .map(|event| event["retained_states"].as_u64().unwrap()) + .sum(); + assert_eq!(priced, installed); + let retained_cost = |plan: &CompiledPhysicalPlan| { + plan.planner_selection_trace + .iter() + .filter(|event| event["stage"] == "deployment.lifecycle_placement") + .map(|event| cost(event, "continuously_maintained_cost")) + .sum::() + }; + let longest = two_state_plan("sum(sum_over_time(m[10m]))", 1e-12); + // Same producer and 61 panes, plus the shorter logical read every 10s. + assert!((retained_cost(&plan) - retained_cost(&longest) - 3.0).abs() < 1e-9); +} diff --git a/docs/examples/workload-cost-evidence.md b/docs/examples/workload-cost-evidence.md index 2e0b8f296..aa928b4ba 100644 --- a/docs/examples/workload-cost-evidence.md +++ b/docs/examples/workload-cost-evidence.md @@ -98,13 +98,27 @@ 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 +under `require_backend_local_execution`). A query may retain admissible native +branches beside raw inputs when their bounded-lag mixed assignment costs less; +otherwise its states retain their group placement. With no retained state it +runs natively over raw series. The manifest of the 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. +Retention prices use the chosen installed layout, including downstream +maintenance lookbacks and the configured retention margin. Raw additive +states within one query that share the same generated panes charge their +producer and longest retention once, while preserving each logical read cost. + +`implementation.source_ingestion_rates` optionally maps metric names to +samples-per-second evidence, using the same `Evidence` format as +`data_workload.ingestion_rate`. Fresh observations price both retained +maintenance and raw selector folds. Rates must be finite and nonnegative. +Missing or expired metric evidence uses the workload-wide rate as a +conservative bound; label filters do not imply an invented selectivity. + An exact alternative without an accessible native backend is unavailable even if its numeric quote would be cheap.