diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index 5b53924fc..fc5dc4365 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -1780,36 +1780,7 @@ impl DeploymentPlanCompiler { .retained_physical()? .is_some_and(|candidate| candidate.precompute.is_some()); for (ordinal, selected) in selected.into_iter().enumerate() { - let mut branch_query = query.clone(); - branch_query.query_lookback_ms = selected - .window_secs - .map(|seconds| seconds.saturating_mul(1_000)) - .unwrap_or(query.query_lookback_ms); - branch_query.group_by_labels = selected - .group_by - .clone() - .unwrap_or_else(|| query.group_by_labels.clone()); - branch_query - .window_realization_candidates - .retain(|candidate| { - candidate.window_secs.saturating_mul(1_000) - == branch_query.query_lookback_ms - && if cohort_nodes.contains(&(Rc::as_ptr(&selected.node) as usize)) { - if native_cohort { - windows::is_complete_window(candidate) - && candidate.slide_secs.saturating_mul(1000) - == u64::from( - query - .summary_lifecycle_inputs - .evaluation_interval_ms, - ) - } else { - windows::is_full_cohort(candidate) - } - } else { - !candidate.cohort_only - } - }); + let branch_query = state_query(query, &selected, &cohort_nodes, native_cohort); let query = &branch_query; let lifecycle_costs = SummaryMaintenanceLifecycleCostInputs { build_cost: Some(Cost(query.summary_lifecycle_inputs.costs.build)), @@ -3885,6 +3856,42 @@ struct PlannerPhysicalSelection { lifecycle_cost: f64, } +/// `query` as the consumer of `state` alone: its window, grouping, and the +/// window implementations that can install it. +fn state_query( + query: &QueryCompilationInput, + state: &SelectedMaterialization, + cohort_nodes: &BTreeSet, + native_cohort: bool, +) -> QueryCompilationInput { + let mut branch = query.clone(); + branch.query_lookback_ms = state + .window_secs + .map(|seconds| seconds.saturating_mul(1_000)) + .unwrap_or(query.query_lookback_ms); + branch.group_by_labels = state + .group_by + .clone() + .unwrap_or_else(|| query.group_by_labels.clone()); + let lookback_ms = branch.query_lookback_ms; + let cohort = cohort_nodes.contains(&(Rc::as_ptr(&state.node) as usize)); + branch.window_realization_candidates.retain(|candidate| { + candidate.window_secs.saturating_mul(1_000) == lookback_ms + && if cohort { + if native_cohort { + windows::is_complete_window(candidate) + && candidate.slide_secs.saturating_mul(1000) + == u64::from(query.summary_lifecycle_inputs.evaluation_interval_ms) + } else { + windows::is_full_cohort(candidate) + } + } else { + !candidate.cohort_only + } + }); + branch +} + fn select_lifecycle( query: &QueryCompilationInput, node: &SummaryNode, diff --git a/control_plane/src/physical/compiler/placement.rs b/control_plane/src/physical/compiler/placement.rs index abb1e1cb9..156786e3a 100644 --- a/control_plane/src/physical/compiler/placement.rs +++ b/control_plane/src/physical/compiler/placement.rs @@ -3,8 +3,9 @@ //! //! 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 +//! A `ContinuouslyMaintained` state is precomputed at ingestion into the window +//! layout the compiler installs for it. An `Ephemeral` state is never built: +//! its queries run an exact program over 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}; @@ -46,16 +47,18 @@ impl Placement { /// Backend lifecycle prices for any summary state. Retention charges the /// 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, +/// retention unpriced. +struct LifecycleCosts { + costs: LifecycleUnitCosts, + /// States the installed window layout retains. Without one, panes are + /// estimated as the state's window over `evaluation_interval_ms`. + retained_states: Option, evaluation_interval_ms: u32, input_cardinality: Option, delete: bool, } -impl CostModel for LifecycleCosts<'_> { +impl CostModel for LifecycleCosts { fn rank_candidates( &self, _intent: &AggIntent, @@ -68,16 +71,21 @@ impl CostModel for LifecycleCosts<'_> { &self, summary: &SummaryNode, ) -> SummaryMaintenanceLifecycleCostInputs { - 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 costs = &self.costs; + let panes = self.retained_states.map_or_else( + || { + 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) + }) + }, + |states| states as f64, + ); let store = match &summary.expr { SummaryExpr::SummaryAgg { family, reduction, .. @@ -119,6 +127,70 @@ impl CostModel for LifecycleCosts<'_> { } } +/// The window implementation compilation installs for `state` when `query` +/// retains it, 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. +fn installed_window( + query: &QueryCompilationInput, + state: &SelectedMaterialization, + query_states: &[SelectedMaterialization], + environment: &PhysicalDeploymentContext, + retention_margin_ms: u64, +) -> Option<(WindowRealizationCandidate, u64)> { + let native_cohort = query + .retained_physical() + .ok()? + .is_some_and(|candidate| candidate.precompute.is_some()); + let branch = state_query( + query, + state, + &super::windows::cohort_nodes(query_states), + native_cohort, + ); + let model = ControlPlaneCostModel::new(branch.accuracy_target.clone()) + .with_window_implementation_costs( + validate_window_implementations(&branch, environment).ok()?, + ); + let (id, framework, _) = model.cheapest_window_implementation()?; + let window = branch + .window_realization_candidates + .iter() + .find(|candidate| { + id.as_ref() == Some(&candidate.realization_id) && &candidate.framework == framework + })? + .clone(); + let retained = retained_state_count( + branch.query_lookback_ms, + retention_margin_ms, + window.slide_secs.saturating_mul(1_000), + &window.layout, + ); + Some((window, retained)) +} + +/// A state's index, raw bindability, retained and rebuilt costs, and +/// retained store states. +type Decision = (usize, bool, Option, Option, Option); + +/// One installed copy of a retained state and the consumers reading it. +struct Install<'a> { + /// Window, slide, layout, evaluation cadence, and phase modulo cadence and + /// window; `None` when no installable window or phase is known. + layout: Option<( + u64, + u64, + asap_types::WindowMaterializationLayout, + u64, + u64, + u64, + )>, + costs: &'a LifecycleUnitCosts, + retained_states: Option, + evaluation_interval_ms: u32, + consumers: Vec, +} + /// Selectable total cost of `lifecycle` for `summary`, if Planner listed it. fn alternative_cost( deployment: &asap_aware_mapping::SummaryMaintenanceDeployment, @@ -168,6 +240,8 @@ pub(super) fn place( && environment.target == PhysicalDeploymentTarget::BackendLocalRemoteWrite; let mut states: Vec<(Rc, Vec)> = Vec::new(); let mut query_states = vec![Vec::new(); queries.len()]; + let mut selected_states: Vec> = + (0..queries.len()).map(|_| Vec::new()).collect(); for (index, query) in queries.iter().enumerate() { if super::super::maintained_population::supported_node(&query.selected_plan_root) { continue; @@ -176,7 +250,7 @@ pub(super) fn place( else { continue; }; - for state in selected { + for state in &selected { if query_states[index] .iter() .any(|known: &Rc| Rc::ptr_eq(known, &state.node)) @@ -189,9 +263,10 @@ pub(super) fn place( .find(|(known, _)| Rc::ptr_eq(known, &state.node)) { Some((_, consumers)) => consumers.push(index), - None => states.push((state.node, vec![index])), + None => states.push((Rc::clone(&state.node), vec![index])), } } + selected_states[index] = selected; } let raw_programs: Vec> = (0..queries.len()) .map(|index| { @@ -201,91 +276,208 @@ pub(super) fn place( .and_then(|root| raw_query_time_program(root).ok()) }) .collect(); - 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()); - 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 = LifecycleCosts { - costs: &lead.costs, - evaluation_interval_ms: interval, - input_cardinality: data - .input_cardinality - .value_at(environment.observed_at_unix_ms) - .copied(), + 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 model = + |costs: LifecycleUnitCosts, retained_states, evaluation_interval_ms| LifecycleCosts { + costs, + retained_states, + evaluation_interval_ms, + input_cardinality: data.input_cardinality.value_at(now).copied(), delete: environment.target == PhysicalDeploymentTarget::BackendLocalRemoteWrite, }; - let Ok(candidates) = enumerate_summary_maintenance_lifecycles( + // Planner's alternative `lifecycle` for `state`, priced for `consumers`. + let price = |state: &Rc, + consumers: &[usize], + lifecycle: SummaryMaintenanceLifecycle, + model: &LifecycleCosts| { + let candidates = enumerate_summary_maintenance_lifecycles( Rc::clone(state), WorkloadDemand::new_with_data(workload, data, consumers), - environment.observed_at_unix_ms, - Some(Horizon(horizon)), + now, + horizon, SummaryMaintenanceLifecycleCapabilities { - supports_ephemeral: bindable, + supports_ephemeral: lifecycle == SummaryMaintenanceLifecycle::Ephemeral, supports_prepared: false, supports_shared: false, - supports_continuously_maintained: true, + supports_continuously_maintained: lifecycle + == SummaryMaintenanceLifecycle::ContinuouslyMaintained, }, - &model, - ) else { - continue; - }; - let Some(deployment) = candidates + model, + ) + .ok()?; + let deployment = candidates .deployments() .iter() - .find(|deployment| Rc::ptr_eq(&deployment.summary, state)) - else { - continue; + .find(|deployment| Rc::ptr_eq(&deployment.summary, state))?; + alternative_cost(deployment, &lifecycle) + }; + let phases: Vec> = workload + .entries() + .map(|entry| match entry.recurrence { + QueryRecurrence::Repeated(RepeatedDemand::FixedIntervalAt { + evaluation_phase, .. + }) => Some(evaluation_phase.0), + _ => None, + }) + .collect(); + let mut decisions: Vec = Vec::new(); + for (state_index, (state, consumers)) in states.iter().enumerate() { + let bindable = consumers.iter().all(|&query| raw_programs[query].is_some()); + // Consumers share one installed state only when compilation groups + // them: the same window layout and cadence, and the same evaluation + // phase within both cadence and window. Any other consumer installs its own. + let mut installs: Vec = Vec::new(); + for &query in consumers { + let lifecycle = &queries[query].summary_lifecycle_inputs; + let installed = selected_states[query] + .iter() + .find(|selected| Rc::ptr_eq(&selected.node, state)) + .and_then(|selected| { + installed_window( + &queries[query], + selected, + &selected_states[query], + environment, + request.query_retention_margin_ms, + ) + }); + let layout = installed + .as_ref() + .zip(phases[query]) + .map(|((window, _), phase)| { + let cadence_ms = u64::from(lifecycle.evaluation_interval_ms).max(1); + let window_ms = window.window_secs.saturating_mul(1_000).max(1); + ( + window.window_secs, + window.slide_secs, + window.layout.clone(), + cadence_ms, + phase % cadence_ms, + phase % window_ms, + ) + }); + match installs + .iter_mut() + .find(|install| layout.is_some() && install.layout == layout) + { + Some(install) => install.consumers.push(query), + None => installs.push(Install { + layout, + costs: &lifecycle.costs, + retained_states: installed.map(|(_, retained)| retained), + evaluation_interval_ms: lifecycle.evaluation_interval_ms, + consumers: vec![query], + }), + } + } + let retained = installs + .iter() + .map(|install| { + price( + state, + &install.consumers, + SummaryMaintenanceLifecycle::ContinuouslyMaintained, + &model( + install.costs.clone(), + install.retained_states, + install.evaluation_interval_ms, + ), + ) + .map(|cost| cost.0) + }) + .sum::>() + .map(Cost); + // Each consumer runs its raw program once per evaluation. Per state the + // program builds, finalizes and retires a transient accumulator, as + // Planner's per-read rebuild charges; it also folds every sample its + // Scans cover, at the per-update cost maintenance pays for the same + // source rate. The query's states share that fold. + let rebuilt = if bindable { + consumers + .iter() + .map(|&query| { + let raw = raw_programs[query].as_ref()?; + let scanned_seconds = raw + .scans + .iter() + .map(|(_, scan)| match scan { + QueryTimeOperator::Scan { + range_ms: Some(range_ms), + .. + } => Some(*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 costs = LifecycleUnitCosts { + build: lifecycle.costs.build + fold / query_states[query].len() as f64, + ..lifecycle.costs.clone() + }; + price( + state, + std::slice::from_ref(&query), + SummaryMaintenanceLifecycle::Ephemeral, + &model(costs, None, lifecycle.evaluation_interval_ms), + ) + .map(|cost| cost.0) + }) + .sum::>() + .map(Cost) + } else { + None }; - 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)); + let retained_states = installs + .iter() + .map(|install| install.retained_states) + .sum::>(); + decisions.push((state_index, bindable, retained, rebuilt, retained_states)); } // 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. + // and exact subtrees share no snapshot with installed state. States linked + // through a query therefore move together, and only when the raw programs + // of all their queries cost less than retaining 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() - }); - if realizable { - continue; - } - for state in owned { - if let Some(i) = index_of(state).filter(|&i| ephemeral[i]) { - ephemeral[i] = false; - changed = true; - } + let mut linked: Vec = (0..states.len()).collect(); + for owned in &query_states { + let members: Vec = owned.iter().filter_map(index_of).collect(); + if let Some(&first) = members.first() { + for member in members { + let (from, to) = (linked[member], linked[first]); + linked + .iter_mut() + .filter(|l| **l == from) + .for_each(|l| *l = to); } } - if !changed { - break; - } } - for (state_index, bindable, retained, rebuilt) in decisions { + let total = |group: usize, cost: fn(&Decision) -> Option| { + decisions + .iter() + .filter(|decision| linked[decision.0] == group) + .map(|decision| cost(decision).map(|cost| cost.0)) + .sum::>() + .map(Cost) + }; + let ephemeral: Vec = (0..states.len()) + .map(|state| { + rebuild_is_cheaper( + total(linked[state], |decision| decision.2), + total(linked[state], |decision| decision.3), + ) + }) + .collect(); + for (state_index, bindable, retained, rebuilt, retained_states) 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, + "retained_states": retained_states, "continuously_maintained_cost": retained.map(|cost| cost.0), "ephemeral_cost": rebuilt.map(|cost| cost.0), "selected": if ephemeral[state_index] { "ephemeral" } else { "continuously_maintained" }, @@ -499,8 +691,11 @@ pub(super) fn time_native_candidate( }) }; let lifecycle = &query.summary_lifecycle_inputs; + // Windows are prepared from the timed root, so the installed layout is not + // known while its timing is being chosen. let model = LifecycleCosts { - costs: &lifecycle.costs, + costs: lifecycle.costs.clone(), + retained_states: None, evaluation_interval_ms: lifecycle.evaluation_interval_ms, input_cardinality: data .input_cardinality diff --git a/control_plane/src/physical/post_asap/cost_model.rs b/control_plane/src/physical/post_asap/cost_model.rs index 47034249b..44102aa47 100644 --- a/control_plane/src/physical/post_asap/cost_model.rs +++ b/control_plane/src/physical/post_asap/cost_model.rs @@ -397,6 +397,28 @@ impl ControlPlaneCostModel { self } + /// The window implementation every complete candidate of this model + /// installs, whatever lifecycles its summary states take. + pub(crate) fn cheapest_window_implementation( + &self, + ) -> Option<&(Option, SummaryWindowFramework, Cost)> { + self.window_framework_costs + .iter() + // GOS/error propagation is introduced by the later adaptation + // slice. Until then, do not claim an approximate exponential + // histogram window is exact. + .filter(|(_, framework, _)| { + !matches!(framework, SummaryWindowFramework::ExponentialHistogram) + }) + .filter(|(_, _, cost)| cost.0.is_finite() && cost.0 >= 0.0) + .min_by(|left, right| { + left.2 + .0 + .total_cmp(&right.2 .0) + .then_with(|| left.0.cmp(&right.0)) + }) + } + pub fn with_summary_maintenance( mut self, lifecycle_costs: SummaryMaintenanceLifecycleCostInputs, @@ -626,21 +648,7 @@ impl CostModel for ControlPlaneCostModel { .iter() .map(|deployment| deployment.selected_cost.0) .sum(); - self.window_framework_costs - .iter() - // GOS/error propagation is introduced by the later adaptation - // slice. Until then, do not claim an approximate exponential - // histogram window is exact. - .filter(|(_, framework, _)| { - !matches!(framework, SummaryWindowFramework::ExponentialHistogram) - }) - .filter(|(_, _, cost)| cost.0.is_finite() && cost.0 >= 0.0) - .min_by(|left, right| { - left.2 - .0 - .total_cmp(&right.2 .0) - .then_with(|| left.0.cmp(&right.0)) - }) + self.cheapest_window_implementation() .map( |(id, framework, physical_cost)| CompleteSummaryCandidateEstimate { physical_plan_id: id.clone(), diff --git a/control_plane/tests/lifecycle_placement.rs b/control_plane/tests/lifecycle_placement.rs index 9f3c48ef8..fa8dd3991 100644 --- a/control_plane/tests/lifecycle_placement.rs +++ b/control_plane/tests/lifecycle_placement.rs @@ -1,6 +1,7 @@ //! Precompute-or-query-time placement follows the backend's lifecycle costs. use control_plane::physical::compiler::{ BackendLocalPlanningInput, CompiledPhysicalPlan, DeploymentPlanCompiler, + PhysicalCompilationRequest, }; use control_plane::physical::workload_cost::enumerate_exact_and_materialized_candidates; use control_plane::query_plan::{query_time::QueryTimeOperator, QueryPlanNode}; @@ -19,13 +20,22 @@ fn fixture(store_per_byte_second: f64, require_local: bool) -> BackendLocalPlann /// The Planner-selected candidate, compiled; the native exact alternative is last. fn selected_plan(input: BackendLocalPlanningInput) -> CompiledPhysicalPlan { + selected_plan_with(input, |_| {}) +} + +/// [`selected_plan`] after `edit` adjusts the per-query compilation inputs. +fn selected_plan_with( + input: BackendLocalPlanningInput, + edit: impl FnOnce(&mut PhysicalCompilationRequest), +) -> CompiledPhysicalPlan { let (request, environment) = input.into_physical_compilation_request().unwrap(); - let candidate = enumerate_exact_and_materialized_candidates(request) + let mut candidate = enumerate_exact_and_materialized_candidates(request) .unwrap() .into_iter() .next() .unwrap(); assert!(candidate.allow_mixed_summary_and_exact_execution); + edit(&mut candidate); DeploymentPlanCompiler .compile_promql(candidate, environment) .unwrap() @@ -92,18 +102,39 @@ fn ephemeral_requires_a_bindable_raw_source() { } fn decisions(queries: &[&str]) -> Vec { + let queries: Vec<_> = queries.iter().map(|query| (*query, 10_000, 0.1)).collect(); + priced_decisions(&queries) +} + +/// Lifecycle decisions for exact `(query, evaluation interval ms, read unit cost)`. +fn priced_decisions(queries: &[(&str, u32, f64)]) -> Vec { + let queries: Vec<_> = queries + .iter() + .map(|&(query, interval, read)| (query, interval, 0, read)) + .collect(); + phased_decisions(&queries) +} + +/// [`priced_decisions`] with each query's evaluation phase in ms. +fn phased_decisions(queries: &[(&str, u32, u64, f64)]) -> 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| { + .map(|(query, interval, phase, _)| { let mut entry = template.clone(); entry["query"] = (*query).into(); + entry["demand"]["fixed_interval_at"]["interval"] = (*interval).into(); + entry["demand"]["fixed_interval_at"]["evaluation_phase"] = (*phase).into(); entry["requirements"]["accuracy"] = serde_json::json!({"explicit": "Exact"}); entry }) .collect(); - let plan = selected_plan(serde_json::from_value(wire).unwrap()); + let plan = selected_plan_with(serde_json::from_value(wire).unwrap(), |request| { + for (query, (_, _, _, read)) in request.queries.iter_mut().zip(queries) { + query.summary_lifecycle_inputs.costs.read = *read; + } + }); plan.planner_selection_trace .iter() .filter(|entry| entry["stage"] == "deployment.lifecycle_placement") @@ -129,3 +160,125 @@ fn shared_state_is_priced_once_with_all_reads() { < 2.0 * cost(&alone, "continuously_maintained_cost") ); } + +fn decision(plan: &CompiledPhysicalPlan) -> Value { + let [decision] = plan + .planner_selection_trace + .iter() + .filter(|entry| entry["stage"] == "deployment.lifecycle_placement") + .cloned() + .collect::>() + .try_into() + .unwrap(); + decision +} + +fn cost(decision: &Value, field: &str) -> f64 { + decision[field].as_f64().unwrap() +} + +// Retention is priced for the panes the installed window layout retains, +// not for window / evaluation interval. +#[test] +fn retention_prices_the_installed_window_panes() { + let plan = selected_plan(fixture(1e-12, false)); + let decision = decision(&plan); + assert_eq!(decision["selected"], "continuously_maintained"); + let [installed] = plan.precompute_plan.materializations.as_slice() else { + panic!("one installed state"); + }; + assert_eq!( + decision["retained_states"].as_u64(), + installed.num_aggregates_to_retain + ); +} + +// Rebuilding runs the exact query-time program: every evaluation builds and +// retires a transient accumulator, folds each raw sample its Scan covers once, +// and reads the result once. +#[test] +fn ephemeral_cost_is_the_raw_program_over_its_scanned_samples() { + let input = fixture(1.0, false); + let costs = input.physical_inputs.lifecycle_costs.clone(); + let horizon = input.physical_inputs.horizon_seconds; + let rate = input.data_workload.ingestion_rate.value.unwrap().0; + let plan = selected_plan(input); + let [QueryTimeOperator::Scan { + range_ms: Some(range_ms), + .. + }] = raw_scans(&plan).as_slice() + else { + panic!("one raw range selector"); + }; + let evaluations = horizon / 10.0; + let raw_samples = rate * *range_ms as f64 / 1_000.0; + let expected = evaluations + * (costs.build + + raw_samples * costs.maintenance_per_update + + costs.read + + costs.retirement); + let decision = decision(&plan); + assert_eq!(decision["selected"], "ephemeral"); + assert!((cost(&decision, "ephemeral_cost") - expected).abs() < 1e-9); +} + +// The store price at which the fixture's quantile sketch moves to query time: +// retention grows linearly with the price and flips just past the point where +// it overtakes the raw program's cost. +#[test] +fn store_price_flips_placement_at_the_break_even_price() { + let free = decision(&selected_plan(fixture(0.0, false))); + let priced = decision(&selected_plan(fixture(1.0, false))); + let slope = + cost(&priced, "continuously_maintained_cost") - cost(&free, "continuously_maintained_cost"); + let threshold = + (cost(&free, "ephemeral_cost") - cost(&free, "continuously_maintained_cost")) / slope; + // Over the 300 s horizon at one read per 10 s: the raw program costs + // 30 * (10 + 100 samples/s * 60 s * 0.001 + 0.1 + 1) = 513; free retention costs + // 10 + 300 * 100 * 0.001 + 30 * 0.1 + 300 * 0.001 + 1 = 44.3; each unit of + // price adds 300 s * 7 retained panes * 65536 estimated sketch bytes. + let expected = (513.0 - 44.3) / (300.0 * 7.0 * 65_536.0); + assert!( + (threshold - expected).abs() < expected * 1e-9, + "{threshold}" + ); + for (scale, selected) in [(0.99, "continuously_maintained"), (1.01, "ephemeral")] { + let decision = decision(&selected_plan(fixture(threshold * scale, false))); + assert_eq!(decision["selected"], selected); + } +} + +// Consumers whose evaluation cadences install separate windows are priced as +// the states they install, each consumer's reads at its own read unit cost. +#[test] +fn shared_state_combines_every_consumers_demand_and_costs() { + let first = ("sum_over_time(m[1m])", 10_000, 0.1); + let second = ("sort(sum_over_time(m[1m]))", 20_000, 0.7); + let [alone_first] = priced_decisions(&[first]).try_into().unwrap(); + let [alone_second] = priced_decisions(&[second]).try_into().unwrap(); + let [shared] = priced_decisions(&[first, second]).try_into().unwrap(); + assert_eq!(shared["query_ids"].as_array().unwrap().len(), 2); + for field in ["continuously_maintained_cost", "ephemeral_cost"] { + let expected = cost(&alone_first, field) + cost(&alone_second, field); + assert!((cost(&shared, field) - expected).abs() < 1e-9, "{field}"); + } + assert_eq!( + shared["retained_states"].as_u64().unwrap(), + alone_first["retained_states"].as_u64().unwrap() + + alone_second["retained_states"].as_u64().unwrap() + ); +} + +// Consumers on one cadence but out of phase read different panes, so each +// installs and is priced as its own state. +#[test] +fn out_of_phase_consumers_are_priced_as_separate_installs() { + let first = ("sum_over_time(m[1m])", 10_000, 0, 0.1); + let second = ("sort(sum_over_time(m[1m]))", 10_000, 5_000, 0.1); + let [alone_first] = phased_decisions(&[first]).try_into().unwrap(); + let [alone_second] = phased_decisions(&[second]).try_into().unwrap(); + let [shared] = phased_decisions(&[first, second]).try_into().unwrap(); + let expected = cost(&alone_first, "continuously_maintained_cost") + + cost(&alone_second, "continuously_maintained_cost"); + assert!((cost(&shared, "continuously_maintained_cost") - expected).abs() < 1e-9); +}