diff --git a/crates/asap-aware-mapping/src/cost_model.rs b/crates/asap-aware-mapping/src/cost_model.rs index fae0243f..556afe89 100644 --- a/crates/asap-aware-mapping/src/cost_model.rs +++ b/crates/asap-aware-mapping/src/cost_model.rs @@ -56,7 +56,6 @@ use asap_types::pre_asap::agg_intent::AggIntent; use asap_types::pre_asap::expr_ir::ColumnRef; use asap_types::pre_asap::query_expr::QueryExpr; -use crate::lifecycle::{LifecycleCostInputs, SummaryLifecycleCapabilities}; use crate::recurrence::{ self, Horizon, RecurrenceCostExplanation, RecurrenceError, RecurrenceProfile, }; @@ -64,6 +63,9 @@ use crate::replacement::{ realize_child, Implementation, Replacement, ReplacementProvenance, ReplacementSubDAG, TargetSubDAG, }; +use crate::summary_maintenance_lifecycle::{ + SummaryMaintenanceCapabilities, SummaryMaintenanceLifecycleCostInputs, +}; /// A CSE-detected, legality-gated shared subtree with two or more consumers /// — the unit [`CostModel::cse_share_decision`] decides over. Built by @@ -510,17 +512,20 @@ pub trait CostModel { /// compare physical summary-state lifecycles. Unknown values stay /// unknown, preventing long-lived deployments from winning through /// optimistic zeroes. - fn summary_lifecycle_cost_inputs(&self, _summary: &SummaryNode) -> LifecycleCostInputs { - LifecycleCostInputs::default() + fn summary_maintenance_lifecycle_cost_inputs( + &self, + _summary: &SummaryNode, + ) -> SummaryMaintenanceLifecycleCostInputs { + SummaryMaintenanceLifecycleCostInputs::default() } /// Physical update/merge/delete support for one concrete summary. The /// conservative default advertises no long-lived maintenance capability. - fn summary_lifecycle_capabilities( + fn summary_maintenance_capabilities( &self, _summary: &SummaryNode, - ) -> SummaryLifecycleCapabilities { - SummaryLifecycleCapabilities::default() + ) -> SummaryMaintenanceCapabilities { + SummaryMaintenanceCapabilities::default() } /// Cost of evaluating `target` directly from its logical/raw inputs once. diff --git a/crates/asap-aware-mapping/src/lib.rs b/crates/asap-aware-mapping/src/lib.rs index be0f0122..ed1570e4 100644 --- a/crates/asap-aware-mapping/src/lib.rs +++ b/crates/asap-aware-mapping/src/lib.rs @@ -186,11 +186,11 @@ pub mod accuracy_reconciliation; pub mod cost_model; pub mod explanation; pub mod grouping; -pub mod lifecycle; pub mod recurrence; pub mod replacement; pub mod rewrite; pub mod rollup; +pub mod summary_maintenance_lifecycle; pub mod topk_reuse; pub use accuracy::{ @@ -204,11 +204,6 @@ pub use explanation::{ explain_replacements, explain_replacements_with, ExplanationKind, ReplacementExplanation, }; pub use grouping::{has_subpopulations, HydraGroupingStrategy}; -pub use lifecycle::{ - plan_summary_lifecycles, LifecycleAlternative, LifecycleCapabilities, LifecycleCostInputs, - LifecyclePlan, LifecyclePlanError, LifecycleRejection, StateDeployment, - SummaryLifecycleCapabilities, WorkloadDemand, -}; pub use recurrence::{ evaluation_rate_of, total_cost, update_rate_from_data_workload, CostRate, EvaluationRate, Horizon, RecurrenceCostExplanation, RecurrenceError, RecurrenceProfile, RootRecurrence, @@ -223,4 +218,14 @@ pub use replacement::{ MAX_SEARCH_ITERATIONS, }; pub use rewrite::AvgToSumOverCountStrategy; +pub use summary_maintenance_lifecycle::{ + global_selection_with_summary_maintenance_lifecycles, + materialize_with_summary_maintenance_lifecycles, plan_summary_maintenance_lifecycles, + MaterializeSummaryMaintenanceLifecycleError, SummaryMaintenanceCapabilities, + SummaryMaintenanceDeployment, SummaryMaintenanceLifecycleAlternative, + SummaryMaintenanceLifecycleCapabilities, SummaryMaintenanceLifecycleCostInputs, + SummaryMaintenanceLifecyclePlan, SummaryMaintenanceLifecyclePlanError, + SummaryMaintenanceLifecycleRejection, SummaryMaintenanceLifecycleSelectionError, + WorkloadDemand, +}; pub use topk_reuse::TopKLimitReuseStrategy; diff --git a/crates/asap-aware-mapping/src/lifecycle.rs b/crates/asap-aware-mapping/src/lifecycle.rs deleted file mode 100644 index 8a4a37f8..00000000 --- a/crates/asap-aware-mapping/src/lifecycle.rs +++ /dev/null @@ -1,1209 +0,0 @@ -//! Workload-aware physical lifecycle planning for summary state. -//! -//! This module answers how each unique `SummaryAgg` state is deployed for the -//! supplied query and data workloads. Unknown evidence stays unknown and -//! therefore cannot make a long-lived lifecycle win. - -use std::collections::HashSet; -use std::rc::Rc; - -use asap_types::post_asap::{EvaluationSchedule, OutputRepresentation, SummaryMaintenanceMode}; -use asap_types::post_asap::{StateLifecycle, SummaryExpr, SummaryNode}; -use asap_types::workload::{ - DataArrival, Predictability, QueryRecurrence, QueryWorkload, RepeatedDemand, TimestampMs, - WorkloadError, -}; - -use crate::cost_model::{Cost, CostModel}; -use crate::recurrence::{CostRate, EvaluationRate, Horizon, UpdateRate}; - -/// Runtime lifecycle shapes available to the planner. -#[derive(Debug, Clone, Copy, PartialEq, Eq)] -pub struct LifecycleCapabilities { - pub ephemeral: bool, - pub prepared: bool, - pub shared: bool, - pub continuously_maintained: bool, -} - -/// Capabilities of one concrete summary family/state representation. -#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] -pub struct SummaryLifecycleCapabilities { - pub incremental_update: bool, - pub merge: bool, - pub delete: bool, -} - -impl LifecycleCapabilities { - pub const ALL: Self = Self { - ephemeral: true, - prepared: true, - shared: true, - continuously_maintained: true, - }; -} - -impl Default for LifecycleCapabilities { - fn default() -> Self { - Self::ALL - } -} - -/// Primitive costs for one concrete summary state. Every field is optional: -/// missing statistics produce an uncosted alternative, never a zero. -#[derive(Debug, Clone, Default, PartialEq)] -pub struct LifecycleCostInputs { - pub build_cost: Option, - pub maintenance_cost_per_update: Option, - pub summary_read_cost: Option, - pub retention_cost_rate: Option, - pub retirement_cost: Option, -} - -#[derive(Debug, Clone, PartialEq, Eq)] -pub enum LifecycleRejection { - UnsupportedByRuntime, - RequiresPredictableOneTimeQuery, - RequiresMultipleReads, - RequiresHorizon, - RequiresContinuousData, - MissingOrStaleIngestionRate, - SummaryDoesNotSupportIncrementalUpdates, - SummaryDoesNotSupportDeletion, - MissingCostEvidence, -} - -#[derive(Debug, Clone, PartialEq)] -pub struct LifecycleAlternative { - pub lifecycle: StateLifecycle, - /// How this lifecycle obtains and refreshes its summary state. - pub maintenance_mode: SummaryMaintenanceMode, - pub total_cost: Option, - pub rejection: Option, - pub assumptions: Vec, -} - -impl LifecycleAlternative { - fn selectable(&self) -> bool { - self.rejection.is_none() && self.total_cost.is_some() - } -} - -/// One unique summary-state deployment. Shared `Rc` nodes are emitted once. -#[derive(Debug, Clone)] -pub struct StateDeployment { - pub summary_index: usize, - pub summary: Rc, - pub selected: Option, - pub selected_maintenance_mode: Option, - pub evaluation_schedule: Option, - pub output_representation: OutputRepresentation, - pub alternatives: Vec, -} - -#[derive(Debug, Clone)] -pub struct LifecyclePlan { - pub root: Rc, - pub deployments: Vec, - pub horizon: Option, - pub evaluation_rate: Option, - pub update_rate: Option, - pub expected_reads: Option, - pub selected_raw_recompute: bool, - pub summary_total_cost: Option, - pub raw_recompute_total_cost: Option, -} - -/// Explicit association between a materialized target and the normalized -/// workload entries whose demand consumes it. -#[derive(Debug, Clone, Copy)] -pub struct WorkloadDemand<'a> { - pub workload: &'a QueryWorkload, - pub entry_indices: &'a [usize], -} - -impl<'a> WorkloadDemand<'a> { - pub const fn new(workload: &'a QueryWorkload, entry_indices: &'a [usize]) -> Self { - Self { - workload, - entry_indices, - } - } -} - -#[derive(Debug, thiserror::Error)] -pub enum LifecyclePlanError { - #[error(transparent)] - InvalidWorkload(#[from] WorkloadError), - #[error("optimization horizon must be finite and strictly positive")] - InvalidHorizon, - #[error("workload entry index {index} is out of bounds for {entry_count} entries")] - InvalidWorkloadEntry { index: usize, entry_count: usize }, - #[error("a workload-demand binding must contain at least one entry")] - EmptyWorkloadDemand, - #[error("workload entry index {index} appears more than once in one demand binding")] - DuplicateWorkloadEntry { index: usize }, -} - -#[derive(Debug)] -struct WorkloadFacts { - reads: Option, - one_time_invocations: u64, - evaluation_rate: Option, - update_rate: Option, - arrival: DataArrival, - prepared_window: Option<(TimestampMs, TimestampMs)>, - prepared_eligible: bool, - requires_deletion: bool, -} - -/// Validate a materialized plan, enumerate lifecycle alternatives for each -/// unique summary state, and select the cheapest legal alternative whose cost -/// is fully known. -pub fn plan_summary_lifecycles( - root: Rc, - demand: WorkloadDemand<'_>, - now_ms: u64, - horizon: Option, - capabilities: LifecycleCapabilities, - cost_model: &dyn CostModel, -) -> Result { - demand.workload.validate()?; - if horizon.is_some_and(|h| !h.0.is_finite() || h.0 <= 0.0) { - return Err(LifecyclePlanError::InvalidHorizon); - } - let facts = workload_facts(demand.workload, demand.entry_indices, now_ms, horizon)?; - let mut summaries = Vec::new(); - collect_summary_aggs(&root, &mut HashSet::new(), &mut summaries); - let deployments: Vec = summaries - .into_iter() - .enumerate() - .map(|(summary_index, summary)| { - let alternatives = alternatives_for( - &facts, - horizon, - capabilities, - cost_model.summary_lifecycle_capabilities(&summary), - cost_model.summary_lifecycle_cost_inputs(&summary), - ); - let selected = alternatives - .iter() - .filter(|candidate| candidate.selectable()) - .min_by(|a, b| a.total_cost.unwrap().0.total_cmp(&b.total_cost.unwrap().0)) - .map(|candidate| (candidate.lifecycle.clone(), candidate.maintenance_mode)); - let evaluation_schedule = selected.as_ref().map(|(lifecycle, _)| match lifecycle { - StateLifecycle::Ephemeral => EvaluationSchedule::OneShot, - StateLifecycle::Prepared { .. } | StateLifecycle::Shared { .. } - if matches!( - facts.arrival, - DataArrival::ContinuouslyIngesting | DataArrival::Mixed - ) => - { - EvaluationSchedule::PerUpdate - } - StateLifecycle::Prepared { .. } => EvaluationSchedule::OneShot, - StateLifecycle::Shared { .. } => EvaluationSchedule::OnRead, - StateLifecycle::ContinuouslyMaintained => EvaluationSchedule::PerUpdate, - }); - StateDeployment { - summary_index, - summary, - selected: selected.as_ref().map(|(lifecycle, _)| lifecycle.clone()), - selected_maintenance_mode: selected.map(|(_, mode)| mode), - evaluation_schedule, - output_representation: OutputRepresentation::SummaryState, - alternatives, - } - }) - .collect(); - let summary_total_cost = deployments.iter().try_fold(Cost::ZERO, |sum, deployment| { - let selected = deployment.selected.as_ref()?; - let cost = deployment - .alternatives - .iter() - .find(|alternative| &alternative.lifecycle == selected)? - .total_cost?; - Some(Cost(sum.0 + cost.0)) - }); - Ok(LifecyclePlan { - root, - deployments, - horizon, - evaluation_rate: facts.evaluation_rate, - update_rate: facts.update_rate, - expected_reads: facts.reads, - selected_raw_recompute: false, - summary_total_cost, - raw_recompute_total_cost: None, - }) -} - -fn workload_facts( - workload: &QueryWorkload, - workload_entry_indices: &[usize], - now_ms: u64, - horizon: Option, -) -> Result { - let mut one_time_invocations = 0u64; - let mut recurring_reads = 0.0; - let mut recurring_known = true; - let mut evaluation_rate = 0.0; - let mut has_evaluation_rate = false; - let mut prepared_start: Option = None; - let mut prepared_end: Option = None; - let mut prepared_eligible = true; - let mut requires_deletion = false; - - let entries: Vec<_> = workload.entries().collect(); - if workload_entry_indices.is_empty() { - return Err(LifecyclePlanError::EmptyWorkloadDemand); - } - let mut seen_indices = HashSet::new(); - for &index in workload_entry_indices { - if !seen_indices.insert(index) { - return Err(LifecyclePlanError::DuplicateWorkloadEntry { index }); - } - let entry = entries - .get(index) - .ok_or(LifecyclePlanError::InvalidWorkloadEntry { - index, - entry_count: entries.len(), - })?; - requires_deletion |= entry.time_selection.lookback.is_some() - && entry.time_selection.as_of.is_none() - && matches!( - entry.time_selection.scope, - asap_types::workload::QueryTimeScope::RealTime - | asap_types::workload::QueryTimeScope::Mixed - ); - match &entry.recurrence { - QueryRecurrence::OneTime { - invocations, - execute_at, - } => { - one_time_invocations = one_time_invocations.saturating_add(*invocations); - let covered = if let ( - Predictability::Predictable { - known_at: Some(known), - }, - Some(execute), - ) = (&entry.predictability, execute_at) - { - if known < execute { - prepared_start = Some(prepared_start.map_or(*known, |old| old.min(*known))); - prepared_end = Some(prepared_end.map_or(*execute, |old| old.max(*execute))); - true - } else { - false - } - } else { - false - }; - prepared_eligible &= covered; - } - QueryRecurrence::Repeated(RepeatedDemand::FixedInterval(interval)) => { - prepared_eligible = false; - let rate = 1000.0 / f64::from(interval.0); - evaluation_rate += rate; - has_evaluation_rate = true; - if let Some(h) = horizon { - recurring_reads += h.0 * rate; - } else { - recurring_known = false; - } - } - QueryRecurrence::Repeated(RepeatedDemand::Scheduled(schedule)) => { - prepared_eligible = false; - if let Some(h) = horizon { - let end_ms = now_ms.saturating_add((h.0 * 1000.0) as u64); - let reads_in_horizon = schedule - .iter() - .filter(|at| at.0 >= now_ms && at.0 <= end_ms) - .count() as f64; - recurring_reads += reads_in_horizon; - evaluation_rate += reads_in_horizon / h.0; - has_evaluation_rate = true; - } else { - recurring_known = false; - } - } - QueryRecurrence::Repeated(RepeatedDemand::EstimatedRate(estimate)) => { - prepared_eligible = false; - if !estimate.is_fresh_at(now_ms) { - recurring_known = false; - continue; - } - let rate = match estimate.expected { - asap_types::workload::ExpectedDemand::AverageRate(rate) => Some(rate.0), - asap_types::workload::ExpectedDemand::InvocationCount(count) => { - let millis = estimate - .observation_window - .end - .0 - .saturating_sub(estimate.observation_window.start.0); - (millis > 0).then_some(count as f64 / (millis as f64 / 1000.0)) - } - }; - if let Some(rate) = rate { - evaluation_rate += rate; - has_evaluation_rate = true; - if let Some(h) = horizon { - recurring_reads += h.0 * rate; - } else { - recurring_known = false; - } - } else { - recurring_known = false; - } - } - QueryRecurrence::Unknown => { - prepared_eligible = false; - recurring_known = false; - } - } - } - - let data = workload.data_workload.as_ref(); - let arrival = data.map_or(DataArrival::Unknown, |data| data.arrival); - let update_rate = data - .and_then(|data| data.ingestion_rate.value_at(now_ms)) - .map(|rate| UpdateRate(rate.0)); - let reads = recurring_known.then_some(one_time_invocations as f64 + recurring_reads); - Ok(WorkloadFacts { - reads, - one_time_invocations, - evaluation_rate: has_evaluation_rate.then_some(EvaluationRate(evaluation_rate)), - update_rate, - arrival, - prepared_window: prepared_start.zip(prepared_end), - prepared_eligible, - requires_deletion, - }) -} - -fn alternatives_for( - facts: &WorkloadFacts, - horizon: Option, - capabilities: LifecycleCapabilities, - summary_capabilities: SummaryLifecycleCapabilities, - costs: LifecycleCostInputs, -) -> Vec { - let alternatives = vec![ - ephemeral(facts, capabilities, &costs), - prepared(facts, capabilities, summary_capabilities, &costs), - shared(facts, horizon, capabilities, summary_capabilities, &costs), - continuous(facts, horizon, capabilities, summary_capabilities, &costs), - ]; - alternatives -} - -fn ephemeral( - facts: &WorkloadFacts, - capabilities: LifecycleCapabilities, - costs: &LifecycleCostInputs, -) -> LifecycleAlternative { - let lifecycle = StateLifecycle::Ephemeral; - if !capabilities.ephemeral { - return rejected( - lifecycle, - SummaryMaintenanceMode::DirectBuild, - LifecycleRejection::UnsupportedByRuntime, - ); - } - let total_cost = zip_costs(&[ - costs.build_cost, - costs.summary_read_cost, - costs.retirement_cost, - ]) - .zip(facts.reads) - .map(|(per_read, reads)| Cost(per_read * reads)); - costed_or_unknown( - lifecycle, - SummaryMaintenanceMode::DirectBuild, - total_cost, - vec!["state is rebuilt per invocation".into()], - ) -} - -fn prepared( - facts: &WorkloadFacts, - capabilities: LifecycleCapabilities, - summary_capabilities: SummaryLifecycleCapabilities, - costs: &LifecycleCostInputs, -) -> LifecycleAlternative { - if !facts.prepared_eligible { - return rejected( - StateLifecycle::Prepared { - activate_at: TimestampMs(0), - retire_at: TimestampMs(0), - }, - retained_mode(facts), - LifecycleRejection::RequiresPredictableOneTimeQuery, - ); - } - let Some((activate_at, retire_at)) = facts.prepared_window else { - return rejected( - StateLifecycle::Prepared { - activate_at: TimestampMs(0), - retire_at: TimestampMs(0), - }, - retained_mode(facts), - LifecycleRejection::RequiresPredictableOneTimeQuery, - ); - }; - let lifecycle = StateLifecycle::Prepared { - activate_at, - retire_at, - }; - if !capabilities.prepared { - return rejected( - lifecycle, - retained_mode(facts), - LifecycleRejection::UnsupportedByRuntime, - ); - } - if let Some(rejection) = maintenance_capability_rejection(facts, summary_capabilities) { - return rejected(lifecycle, retained_mode(facts), rejection); - } - let seconds = retire_at.0.saturating_sub(activate_at.0) as f64 / 1000.0; - let maintenance = maintenance_cost(facts, costs, seconds); - let total_cost = match ( - costs.build_cost, - costs.summary_read_cost, - costs.retention_cost_rate, - costs.retirement_cost, - maintenance, - ) { - (Some(build), Some(read), Some(retention), Some(retire), Some(maintenance)) => Some(Cost( - build.0 - + read.0 * facts.one_time_invocations as f64 - + retention.0 * seconds - + retire.0 - + maintenance, - )), - _ => None, - }; - costed_or_unknown( - lifecycle, - retained_mode(facts), - total_cost, - vec!["activation and retirement come from the declared schedule".into()], - ) -} - -fn shared( - facts: &WorkloadFacts, - horizon: Option, - capabilities: LifecycleCapabilities, - summary_capabilities: SummaryLifecycleCapabilities, - costs: &LifecycleCostInputs, -) -> LifecycleAlternative { - let lifecycle = StateLifecycle::Shared { - retention: asap_types::workload::DurationMs(horizon.map_or(0, |h| (h.0 * 1000.0) as u64)), - }; - if !capabilities.shared { - return rejected( - lifecycle, - retained_mode(facts), - LifecycleRejection::UnsupportedByRuntime, - ); - } - if let Some(rejection) = maintenance_capability_rejection(facts, summary_capabilities) { - return rejected(lifecycle, retained_mode(facts), rejection); - } - if facts.reads.is_none_or(|reads| reads <= 1.0) { - return rejected( - lifecycle, - retained_mode(facts), - LifecycleRejection::RequiresMultipleReads, - ); - } - let Some(horizon) = horizon else { - return rejected( - lifecycle, - retained_mode(facts), - LifecycleRejection::RequiresHorizon, - ); - }; - let total_cost = retained_cost(facts, costs, horizon.0); - costed_or_unknown( - lifecycle, - retained_mode(facts), - total_cost, - vec!["one state is shared across reads".into()], - ) -} - -fn continuous( - facts: &WorkloadFacts, - horizon: Option, - capabilities: LifecycleCapabilities, - summary_capabilities: SummaryLifecycleCapabilities, - costs: &LifecycleCostInputs, -) -> LifecycleAlternative { - let lifecycle = StateLifecycle::ContinuouslyMaintained; - if !capabilities.continuously_maintained { - return rejected( - lifecycle, - SummaryMaintenanceMode::Incremental, - LifecycleRejection::UnsupportedByRuntime, - ); - } - if !matches!( - facts.arrival, - DataArrival::ContinuouslyIngesting | DataArrival::Mixed - ) { - return rejected( - lifecycle, - SummaryMaintenanceMode::Incremental, - LifecycleRejection::RequiresContinuousData, - ); - } - if facts.update_rate.is_none() { - return rejected( - lifecycle, - SummaryMaintenanceMode::Incremental, - LifecycleRejection::MissingOrStaleIngestionRate, - ); - } - if let Some(rejection) = maintenance_capability_rejection(facts, summary_capabilities) { - return rejected(lifecycle, SummaryMaintenanceMode::Incremental, rejection); - } - let Some(horizon) = horizon else { - return rejected( - lifecycle, - SummaryMaintenanceMode::Incremental, - LifecycleRejection::RequiresHorizon, - ); - }; - let total_cost = retained_cost(facts, costs, horizon.0); - costed_or_unknown( - lifecycle, - SummaryMaintenanceMode::Incremental, - total_cost, - vec!["updates are applied for the optimization horizon".into()], - ) -} - -fn maintenance_capability_rejection( - facts: &WorkloadFacts, - capabilities: SummaryLifecycleCapabilities, -) -> Option { - if matches!( - facts.arrival, - DataArrival::ContinuouslyIngesting | DataArrival::Mixed - ) && !capabilities.incremental_update - { - Some(LifecycleRejection::SummaryDoesNotSupportIncrementalUpdates) - } else if matches!( - facts.arrival, - DataArrival::ContinuouslyIngesting | DataArrival::Mixed - ) && facts.requires_deletion - && !capabilities.delete - { - Some(LifecycleRejection::SummaryDoesNotSupportDeletion) - } else { - None - } -} - -fn retained_cost(facts: &WorkloadFacts, costs: &LifecycleCostInputs, seconds: f64) -> Option { - let reads = facts.reads?; - let maintenance = maintenance_cost(facts, costs, seconds)?; - Some(Cost( - costs.build_cost?.0 - + maintenance - + reads * costs.summary_read_cost?.0 - + seconds * costs.retention_cost_rate?.0 - + costs.retirement_cost?.0, - )) -} - -fn maintenance_cost( - facts: &WorkloadFacts, - costs: &LifecycleCostInputs, - seconds: f64, -) -> Option { - match facts.arrival { - DataArrival::AtRest => Some(0.0), - DataArrival::ContinuouslyIngesting | DataArrival::Mixed => { - Some(seconds * facts.update_rate?.0 * costs.maintenance_cost_per_update?.0) - } - DataArrival::Unknown => None, - } -} - -fn zip_costs(costs: &[Option]) -> Option { - costs - .iter() - .try_fold(0.0, |sum, cost| Some(sum + cost.as_ref()?.0)) -} - -fn costed_or_unknown( - lifecycle: StateLifecycle, - maintenance_mode: SummaryMaintenanceMode, - total_cost: Option, - assumptions: Vec, -) -> LifecycleAlternative { - LifecycleAlternative { - lifecycle, - maintenance_mode, - total_cost, - rejection: total_cost - .is_none() - .then_some(LifecycleRejection::MissingCostEvidence), - assumptions, - } -} - -fn rejected( - lifecycle: StateLifecycle, - maintenance_mode: SummaryMaintenanceMode, - rejection: LifecycleRejection, -) -> LifecycleAlternative { - LifecycleAlternative { - lifecycle, - maintenance_mode, - total_cost: None, - rejection: Some(rejection), - assumptions: Vec::new(), - } -} - -fn retained_mode(facts: &WorkloadFacts) -> SummaryMaintenanceMode { - match facts.arrival { - DataArrival::ContinuouslyIngesting | DataArrival::Mixed => { - SummaryMaintenanceMode::Incremental - } - DataArrival::AtRest | DataArrival::Unknown => SummaryMaintenanceMode::DirectBuild, - } -} - -fn collect_summary_aggs( - node: &Rc, - seen: &mut HashSet<*const SummaryNode>, - output: &mut Vec>, -) { - if !seen.insert(Rc::as_ptr(node)) { - return; - } - match &node.expr { - SummaryExpr::SummaryAgg { child, .. } => { - output.push(Rc::clone(node)); - collect_summary_aggs(child, seen, output); - } - SummaryExpr::SummaryJoin { outer, inner, .. } - | SummaryExpr::SummarySubtract { - left: outer, - right: inner, - } => { - collect_summary_aggs(outer, seen, output); - collect_summary_aggs(inner, seen, output); - } - SummaryExpr::SummaryDelete { summary_input, .. } - | SummaryExpr::SummaryEstimate { summary_input, .. } => { - collect_summary_aggs(summary_input, seen, output) - } - SummaryExpr::SummaryMerge { children } => { - for child in children { - collect_summary_aggs(child, seen, output); - } - } - SummaryExpr::KeepPreAsap(_) => {} - } -} - -#[cfg(test)] -mod tests { - use super::*; - use asap_types::post_asap::{ - ExactKind, ExactParams, GroupingStrategy, ResultGuarantee, SummaryFamilyType, SummaryField, - SummarySchema, - }; - use asap_types::pre_asap::{Column, ColumnRef, DataType, QueryExpr, Reduction, Schema, Source}; - use asap_types::workload::{ - BatchEntry, DataWorkload, DurationMs, Evidence, EvidenceSource, Predictability, Query, - QueryLanguage, QueryRequirements, Rate, RepeatingEntry, RepetitionInterval, TimeSelection, - }; - - struct UnitCosts; - - impl CostModel for UnitCosts { - fn rank_candidates( - &self, - _intent: &asap_types::pre_asap::AggIntent, - candidates: &[asap_types::post_asap::SketchAlgorithm], - ) -> Vec { - candidates.to_vec() - } - - fn summary_lifecycle_cost_inputs(&self, _summary: &SummaryNode) -> LifecycleCostInputs { - LifecycleCostInputs { - build_cost: Some(Cost(10.0)), - maintenance_cost_per_update: Some(Cost(1.0)), - summary_read_cost: Some(Cost(1.0)), - retention_cost_rate: Some(CostRate(0.1)), - retirement_cost: Some(Cost(1.0)), - } - } - - fn summary_lifecycle_capabilities( - &self, - _summary: &SummaryNode, - ) -> SummaryLifecycleCapabilities { - SummaryLifecycleCapabilities { - incremental_update: true, - merge: true, - delete: true, - } - } - } - - struct NoDelete; - - impl CostModel for NoDelete { - fn rank_candidates( - &self, - _intent: &asap_types::pre_asap::AggIntent, - candidates: &[asap_types::post_asap::SketchAlgorithm], - ) -> Vec { - candidates.to_vec() - } - - fn summary_lifecycle_cost_inputs(&self, summary: &SummaryNode) -> LifecycleCostInputs { - UnitCosts.summary_lifecycle_cost_inputs(summary) - } - - fn summary_lifecycle_capabilities( - &self, - _summary: &SummaryNode, - ) -> SummaryLifecycleCapabilities { - SummaryLifecycleCapabilities { - incremental_update: true, - merge: true, - delete: false, - } - } - } - - fn query_root() -> Rc { - query_root_for("m") - } - - fn query_root_for(metric: &str) -> Rc { - Rc::new(QueryExpr::Scan { - source: Source::TimeSeries { - metric: metric.into(), - }, - predicates: vec![], - schema: Schema::with_time_index( - vec![ - Column::new("ts", DataType::Timestamp, false), - Column::new("value", DataType::Float64, false), - ], - 0, - vec![], - ), - }) - } - - fn summary() -> Rc { - let child = Rc::new(SummaryNode { - expr: SummaryExpr::KeepPreAsap(query_root()), - schema: SummarySchema { - fields: vec![], - time_index: None, - }, - guarantee: Some(ResultGuarantee::exact("raw")), - }); - let family = SummaryFamilyType::ExactAggregate(ExactKind::Sum, ExactParams::Sum); - Rc::new(SummaryNode { - expr: SummaryExpr::SummaryAgg { - child, - family: family.clone(), - col: ColumnRef::Named("value".into()), - reduction: Reduction::by(vec![]), - grouping: GroupingStrategy::default(), - }, - schema: SummarySchema { - fields: vec![SummaryField { - name: "state".into(), - dtype: family, - nullable: false, - }], - time_index: None, - }, - guarantee: Some(ResultGuarantee::exact("sum")), - }) - } - - fn batch(predictability: Predictability) -> BatchEntry { - BatchEntry { - query: Query("sum(m)".into()), - requirements: QueryRequirements::default(), - predictability, - invocations: 1, - execute_at: None, - time_selection: TimeSelection::default(), - } - } - - fn workload( - batches: Vec, - repeating: Vec, - data: DataWorkload, - ) -> QueryWorkload { - QueryWorkload { - language: QueryLanguage::PromQL, - query_batch: (!batches.is_empty()).then_some(batches), - repeating_queries: (!repeating.is_empty()).then_some(repeating), - data_workload: Some(data), - } - } - - fn at_rest() -> DataWorkload { - DataWorkload { - arrival: DataArrival::AtRest, - ..Default::default() - } - } - - fn continuous(observed_at_ms: u64, valid_for_ms: u64) -> DataWorkload { - DataWorkload { - arrival: DataArrival::ContinuouslyIngesting, - ingestion_rate: Evidence { - value: Some(Rate(1.0)), - source: EvidenceSource::Observed, - observed_at_ms: Some(observed_at_ms), - valid_for_ms: Some(valid_for_ms), - }, - ..Default::default() - } - } - - fn repeating() -> RepeatingEntry { - RepeatingEntry { - query: Query("sum(m)".into()), - demand: RepeatedDemand::FixedInterval(RepetitionInterval(1_000)), - requirements: QueryRequirements::default(), - predictability: Predictability::Predictable { known_at: None }, - time_selection: TimeSelection::default(), - } - } - - #[test] - fn unpredictable_one_time_at_rest_selects_ephemeral() { - let plan = plan_summary_lifecycles( - summary(), - WorkloadDemand::new( - &workload(vec![batch(Predictability::AdHoc)], vec![], at_rest()), - &[0], - ), - 1_000, - None, - LifecycleCapabilities::ALL, - &UnitCosts, - ) - .unwrap(); - assert_eq!(plan.deployments.len(), 1); - assert_eq!( - plan.deployments[0].selected, - Some(StateLifecycle::Ephemeral) - ); - assert_eq!( - plan.deployments[0].selected_maintenance_mode, - Some(SummaryMaintenanceMode::DirectBuild) - ); - assert_eq!( - plan.deployments[0].alternatives[0].total_cost, - Some(Cost(12.0)) - ); - } - - #[test] - fn predictable_scheduled_one_time_offers_prepared_state() { - let mut entry = batch(Predictability::Predictable { - known_at: Some(TimestampMs(1_000)), - }); - entry.execute_at = Some(TimestampMs(11_000)); - let plan = plan_summary_lifecycles( - summary(), - WorkloadDemand::new(&workload(vec![entry], vec![], at_rest()), &[0]), - 1_000, - None, - LifecycleCapabilities::ALL, - &UnitCosts, - ) - .unwrap(); - let prepared = &plan.deployments[0].alternatives[1]; - assert!(prepared.rejection.is_none()); - assert_eq!(prepared.total_cost, Some(Cost(13.0))); - } - - #[test] - fn repeated_at_rest_selects_shared_without_inventing_updates() { - let plan = plan_summary_lifecycles( - summary(), - WorkloadDemand::new(&workload(vec![], vec![repeating()], at_rest()), &[0]), - 1_000, - Some(Horizon(10.0)), - LifecycleCapabilities::ALL, - &UnitCosts, - ) - .unwrap(); - assert_eq!( - plan.deployments[0].selected, - Some(StateLifecycle::Shared { - retention: DurationMs(10_000) - }) - ); - assert_eq!( - plan.deployments[0].selected_maintenance_mode, - Some(SummaryMaintenanceMode::DirectBuild) - ); - assert_eq!( - plan.deployments[0].alternatives[3].rejection, - Some(LifecycleRejection::RequiresContinuousData) - ); - assert_eq!(plan.update_rate, None); - } - - #[test] - fn repeated_continuous_workload_can_select_continuous_maintenance() { - let capabilities = LifecycleCapabilities { - shared: false, - ..LifecycleCapabilities::ALL - }; - let plan = plan_summary_lifecycles( - summary(), - WorkloadDemand::new( - &workload(vec![], vec![repeating()], continuous(1_000, 60_000)), - &[0], - ), - 1_000, - Some(Horizon(10.0)), - capabilities, - &UnitCosts, - ) - .unwrap(); - assert_eq!( - plan.deployments[0].selected, - Some(StateLifecycle::ContinuouslyMaintained) - ); - assert_eq!( - plan.deployments[0].selected_maintenance_mode, - Some(SummaryMaintenanceMode::Incremental) - ); - assert_eq!(plan.evaluation_rate, Some(EvaluationRate(1.0))); - assert_eq!(plan.update_rate, Some(UpdateRate(1.0))); - } - - #[test] - fn stale_ingestion_evidence_cannot_enable_continuous_maintenance() { - let plan = plan_summary_lifecycles( - summary(), - WorkloadDemand::new( - &workload(vec![], vec![repeating()], continuous(1_000, 1_000)), - &[0], - ), - 3_000, - Some(Horizon(10.0)), - LifecycleCapabilities::ALL, - &UnitCosts, - ) - .unwrap(); - assert_eq!( - plan.deployments[0].alternatives[3].rejection, - Some(LifecycleRejection::MissingOrStaleIngestionRate) - ); - assert_eq!(plan.update_rate, None); - } - - #[test] - fn unknown_costs_do_not_make_a_long_lived_lifecycle_win() { - let plan = plan_summary_lifecycles( - summary(), - WorkloadDemand::new( - &workload(vec![], vec![repeating()], continuous(1_000, 60_000)), - &[0], - ), - 1_000, - Some(Horizon(10.0)), - LifecycleCapabilities::ALL, - &crate::cost_model::DefaultCostModel, - ) - .unwrap(); - assert_eq!(plan.deployments[0].selected, None); - assert!(plan.deployments[0] - .alternatives - .iter() - .all(|alternative| alternative.rejection.is_some())); - } - - #[test] - fn unrelated_workload_entries_do_not_create_reuse_for_a_target() { - let plan = plan_summary_lifecycles( - summary(), - WorkloadDemand::new( - &workload( - vec![batch(Predictability::AdHoc), batch(Predictability::AdHoc)], - vec![], - at_rest(), - ), - &[0], - ), - 1_000, - Some(Horizon(10.0)), - LifecycleCapabilities::ALL, - &UnitCosts, - ) - .unwrap(); - assert_eq!( - plan.deployments[0].selected, - Some(StateLifecycle::Ephemeral) - ); - assert_eq!( - plan.deployments[0].alternatives[2].rejection, - Some(LifecycleRejection::RequiresMultipleReads) - ); - } - - #[test] - fn scheduled_rate_counts_only_executions_inside_the_horizon() { - let mut entry = repeating(); - entry.demand = RepeatedDemand::Scheduled(vec![ - TimestampMs(999), - TimestampMs(5_000), - TimestampMs(20_000), - ]); - let plan = plan_summary_lifecycles( - summary(), - WorkloadDemand::new(&workload(vec![], vec![entry], at_rest()), &[0]), - 1_000, - Some(Horizon(10.0)), - LifecycleCapabilities::ALL, - &UnitCosts, - ) - .unwrap(); - assert_eq!(plan.evaluation_rate, Some(EvaluationRate(0.1))); - } - - #[test] - fn demand_binding_rejects_empty_and_duplicate_entries() { - let workload = workload(vec![batch(Predictability::AdHoc)], vec![], at_rest()); - assert!(matches!( - plan_summary_lifecycles( - summary(), - WorkloadDemand::new(&workload, &[]), - 1_000, - None, - LifecycleCapabilities::ALL, - &UnitCosts, - ), - Err(LifecyclePlanError::EmptyWorkloadDemand) - )); - assert!(matches!( - plan_summary_lifecycles( - summary(), - WorkloadDemand::new(&workload, &[0, 0]), - 1_000, - None, - LifecycleCapabilities::ALL, - &UnitCosts, - ), - Err(LifecyclePlanError::DuplicateWorkloadEntry { index: 0 }) - )); - } - - #[test] - fn prepared_requires_every_bound_consumer_to_be_scheduled_and_predictable() { - let mut predictable = batch(Predictability::Predictable { - known_at: Some(TimestampMs(1_000)), - }); - predictable.execute_at = Some(TimestampMs(2_000)); - let workload = workload( - vec![predictable, batch(Predictability::AdHoc)], - vec![], - at_rest(), - ); - let plan = plan_summary_lifecycles( - summary(), - WorkloadDemand::new(&workload, &[0, 1]), - 1_000, - Some(Horizon(10.0)), - LifecycleCapabilities::ALL, - &UnitCosts, - ) - .unwrap(); - assert_eq!( - plan.deployments[0].alternatives[1].rejection, - Some(LifecycleRejection::RequiresPredictableOneTimeQuery) - ); - } - - #[test] - fn moving_realtime_maintenance_requires_summary_deletion_support() { - let mut entry = repeating(); - entry.time_selection = TimeSelection { - scope: asap_types::workload::QueryTimeScope::RealTime, - lookback: Some(DurationMs(60_000)), - as_of: None, - }; - let plan = plan_summary_lifecycles( - summary(), - WorkloadDemand::new( - &workload(vec![], vec![entry], continuous(1_000, 60_000)), - &[0], - ), - 1_000, - Some(Horizon(10.0)), - LifecycleCapabilities::ALL, - &NoDelete, - ) - .unwrap(); - assert_eq!( - plan.deployments[0].alternatives[3].rejection, - Some(LifecycleRejection::SummaryDoesNotSupportDeletion) - ); - } - - #[test] - fn normalized_workload_drives_plan_space_recurrence_profiles() { - let root = query_root(); - let space = crate::replacement::search_workload(vec![("dashboard", Rc::clone(&root))]); - let workload = workload(vec![], vec![repeating()], continuous(1_000, 60_000)); - let profiles = space - .recurrence_profiles_from_workload(&workload, &[0], 1_000, Some(Horizon(10.0))) - .unwrap(); - // `search_workload` canonicalizes roots through CSE; recurrence - // profiles are keyed by that canonical post-CSE node. - let profile = profiles.for_target(&space.roots[0].1); - assert_eq!(profile.evaluation_rate, Some(EvaluationRate(1.0))); - assert_eq!(profile.update_rate, Some(UpdateRate(1.0))); - assert_eq!(profile.one_shot_consumers, 0); - } - - #[test] - fn recurrence_binding_is_explicit_when_root_order_differs_from_workload_order() { - let repeating_root = query_root_for("dashboard"); - let batch_root = query_root_for("batch"); - let space = crate::replacement::search_workload(vec![ - ("dashboard", repeating_root), - ("batch", batch_root), - ]); - let workload = workload( - vec![batch(Predictability::AdHoc)], - vec![repeating()], - at_rest(), - ); - let profiles = space - .recurrence_profiles_from_workload(&workload, &[1, 0], 1_000, Some(Horizon(10.0))) - .unwrap(); - let dashboard = profiles.for_target(&space.roots[0].1); - let batch = profiles.for_target(&space.roots[1].1); - assert_eq!(dashboard.evaluation_rate, Some(EvaluationRate(1.0))); - assert_eq!(dashboard.one_shot_consumers, 0); - assert_eq!(batch.evaluation_rate, None); - assert_eq!(batch.one_shot_consumers, 1); - } -} diff --git a/crates/asap-aware-mapping/src/replacement.rs b/crates/asap-aware-mapping/src/replacement.rs index fdbfdf7a..2c19330c 100644 --- a/crates/asap-aware-mapping/src/replacement.rs +++ b/crates/asap-aware-mapping/src/replacement.rs @@ -372,7 +372,7 @@ use crate::accuracy::{ KLL_RANK_ERROR_EXPONENT_99, }; use crate::accuracy_reconciliation::AccuracyReconciliationStrategy; -use crate::cost_model::{CostModel, CseCandidate, DefaultCostModel, ShareDecision}; +use crate::cost_model::{Cost, CostModel, CseCandidate, DefaultCostModel, ShareDecision}; use crate::grouping::HydraGroupingStrategy; use crate::recurrence::{ evaluation_rate_of, Horizon, RecurrenceError, RecurrenceProfile, RootRecurrence, UpdateRate, @@ -2181,6 +2181,30 @@ pub struct PlanSpace { order: Vec<*const QueryExpr>, } +/// Lifecycle-aware whole-subplan costs keyed by target and candidate identity. +#[derive(Default)] +pub(crate) struct CandidateCostOverrides { + costs: HashMap<(*const QueryExpr, *const ReplacementSubDAG), Cost>, +} + +impl CandidateCostOverrides { + pub(crate) fn insert( + &mut self, + target: &Rc, + candidate: &ReplacementSubDAG, + cost: Cost, + ) { + self.costs + .insert((Rc::as_ptr(target), candidate as *const _), cost); + } + + fn get(&self, target: &Rc, candidate: &ReplacementSubDAG) -> Option { + self.costs + .get(&(Rc::as_ptr(target), candidate as *const _)) + .copied() + } +} + impl PlanSpace { /// Every discovered group, in discovery order. pub fn groups(&self) -> impl Iterator { @@ -2599,6 +2623,54 @@ impl PlanSpace { .map(|rate| UpdateRate(rate.0)); self.recurrence_profiles(&recurrences, update_rate) } + + /// Associate every discovered target with the normalized workload entries + /// whose roots can reach it. + pub(crate) fn workload_entries_by_target( + &self, + workload: &QueryWorkload, + root_workload_entries: &[usize], + ) -> Result>, RecurrenceError> { + let entry_count = workload.entries().count(); + if root_workload_entries.len() != self.roots.len() { + return Err(RecurrenceError::RootCountMismatch { + expected: self.roots.len(), + got: root_workload_entries.len(), + }); + } + let mut bindings: HashMap<*const QueryExpr, HashSet> = HashMap::new(); + for ((_, root), &entry_index) in self.roots.iter().zip(root_workload_entries) { + if entry_index >= entry_count { + return Err(RecurrenceError::InvalidWorkloadEntry { + index: entry_index, + entry_count, + }); + } + let mut seen = HashSet::new(); + let mut queue = VecDeque::from([Rc::as_ptr(root)]); + while let Some(ptr) = queue.pop_front() { + if !seen.insert(ptr) { + continue; + } + bindings.entry(ptr).or_default().insert(entry_index); + if let Some(group) = self.groups.get(&ptr) { + queue.extend( + direct_child_counts(&group.target) + .into_iter() + .map(|(child, _)| child), + ); + } + } + } + Ok(bindings + .into_iter() + .map(|(ptr, entries)| { + let mut entries: Vec<_> = entries.into_iter().collect(); + entries.sort_unstable(); + (ptr, entries) + }) + .collect()) + } } /// Record `times` occurrences of `recurrence` against `ptr` — `times > 1` @@ -2889,6 +2961,23 @@ impl<'a> GlobalSelection<'a> { pub fn for_target(&self, target: &Rc) -> Option<&SelectedGroup<'a>> { self.groups.get(&Rc::as_ptr(target)) } + + /// Materialize the selected replacement at `target`. Exact operators + /// that remain in pre-ASAP IR are preserved by `KeepPreAsap`; logical + /// summary candidates are already fully bound post-ASAP nodes. + pub fn materialize( + &self, + target: &Rc, + ) -> Result>, ImplementError> { + let Some(selected) = self.for_target(target) else { + return Ok(None); + }; + match selected.chosen.map(|candidate| &candidate.replacement) { + Some(Replacement::Summary(node)) => Ok(Some(Rc::clone(node))), + Some(Replacement::Rewrite(rewritten)) => keep_pre_asap(rewritten).map(Some), + None => keep_pre_asap(target).map(Some), + } + } } impl PlanSpace { @@ -2900,7 +2989,7 @@ impl PlanSpace { /// [`Self::cost_sorted`], whose per-group ranking only ever sees a /// group's own raw [`MemoGroup::consumer_count`]. pub fn global_selection(&self, cost_model: &dyn CostModel) -> GlobalSelection<'_> { - self.global_selection_impl(cost_model, None, None) + self.global_selection_impl(cost_model, None, None, None) .expect("structural global selection cannot produce a recurrence error") } @@ -2914,7 +3003,17 @@ impl PlanSpace { profiles: &RecurrenceProfileMap, horizon: Option, ) -> Result, RecurrenceError> { - self.global_selection_impl(cost_model, Some(profiles), horizon) + self.global_selection_impl(cost_model, Some(profiles), horizon, None) + } + + pub(crate) fn global_selection_with_candidate_costs( + &self, + cost_model: &dyn CostModel, + profiles: &RecurrenceProfileMap, + horizon: Option, + costs: &CandidateCostOverrides, + ) -> Result, RecurrenceError> { + self.global_selection_impl(cost_model, Some(profiles), horizon, Some(costs)) } fn global_selection_impl( @@ -2922,6 +3021,7 @@ impl PlanSpace { cost_model: &dyn CostModel, profiles: Option<&RecurrenceProfileMap>, horizon: Option, + candidate_costs: Option<&CandidateCostOverrides>, ) -> Result, RecurrenceError> { let graph = reference_graph(self); let topo = topological_order(&self.order, &graph); @@ -2936,7 +3036,21 @@ impl PlanSpace { let effective = effective_uses.get(ptr).copied().unwrap_or(0); effective_uses.insert(*ptr, effective); - let chosen = if effective >= 2 && cse_candidate_pair(group).is_some() { + let lifecycle_choice = candidate_costs.and_then(|costs| { + group + .candidates + .iter() + .filter_map(|candidate| { + costs + .get(&group.target, candidate) + .map(|cost| (candidate, cost)) + }) + .min_by(|(_, a), (_, b)| a.0.total_cmp(&b.0)) + .map(|(candidate, _)| candidate) + }); + let chosen = if lifecycle_choice.is_some() { + lifecycle_choice + } else if effective >= 2 && cse_candidate_pair(group).is_some() { let decision = if let Some(profiles) = profiles { decide_group_with_recurrence( group, diff --git a/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs b/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs new file mode 100644 index 00000000..d69fd27e --- /dev/null +++ b/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs @@ -0,0 +1,1771 @@ +//! Workload-aware physical summary-maintenance lifecycle planning. +//! +//! This module answers how each unique `SummaryAgg` state is deployed for the +//! supplied query and data workloads. Unknown evidence stays unknown and +//! therefore cannot make a long-lived summary maintenance lifecycle win. + +use std::collections::{HashMap, HashSet}; +use std::rc::Rc; + +use asap_types::post_asap::{ + EvaluationSchedule, OutputRepresentation, SummaryExpr, SummaryMaintenanceLifecycle, + SummaryMaintenanceLifecycleGuarantee, SummaryMaintenanceMode, SummaryNode, +}; +use asap_types::pre_asap::QueryExpr; +use asap_types::workload::{ + DataArrival, Predictability, QueryRecurrence, QueryWorkload, RepeatedDemand, TimestampMs, + WorkloadError, +}; + +use crate::cost_model::{Cost, CostModel}; +use crate::recurrence::{ + CostRate, EvaluationRate, Horizon, RecurrenceError, RecurrenceProfile, UpdateRate, +}; +use crate::replacement::{ + CandidateCostOverrides, GlobalSelection, ImplementError, PlanSpace, Replacement, +}; + +/// Summary maintenance lifecycle shapes available to the runtime planner. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct SummaryMaintenanceLifecycleCapabilities { + pub ephemeral: bool, + pub prepared: bool, + pub shared: bool, + pub continuously_maintained: bool, +} + +/// Capabilities of one concrete summary family/state representation. +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +pub struct SummaryMaintenanceCapabilities { + pub incremental_update: bool, + pub merge: bool, + pub delete: bool, +} + +impl SummaryMaintenanceLifecycleCapabilities { + pub const ALL: Self = Self { + ephemeral: true, + prepared: true, + shared: true, + continuously_maintained: true, + }; +} + +impl Default for SummaryMaintenanceLifecycleCapabilities { + fn default() -> Self { + Self::ALL + } +} + +/// Primitive costs for one concrete summary state. Every field is optional: +/// missing statistics produce an uncosted alternative, never a zero. +#[derive(Debug, Clone, Default, PartialEq)] +pub struct SummaryMaintenanceLifecycleCostInputs { + pub build_cost: Option, + pub maintenance_cost_per_update: Option, + pub summary_read_cost: Option, + pub retention_cost_rate: Option, + pub retirement_cost: Option, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum SummaryMaintenanceLifecycleRejection { + UnsupportedByRuntime, + RequiresPredictableOneTimeQuery, + RequiresMultipleReads, + RequiresHorizon, + RequiresContinuousData, + MissingOrStaleIngestionRate, + SummaryDoesNotSupportIncrementalUpdates, + SummaryDoesNotSupportDeletion, + MissingCostEvidence, +} + +#[derive(Debug, Clone, PartialEq)] +pub struct SummaryMaintenanceLifecycleAlternative { + pub summary_maintenance_lifecycle: SummaryMaintenanceLifecycle, + pub total_cost: Option, + pub rejection: Option, + pub assumptions: Vec, +} + +impl SummaryMaintenanceLifecycleAlternative { + fn selectable(&self) -> bool { + self.rejection.is_none() && self.total_cost.is_some() + } +} + +/// One unique summary-state deployment. Shared `Rc` nodes are emitted once. +#[derive(Debug, Clone)] +pub struct SummaryMaintenanceDeployment { + pub summary_index: usize, + pub summary: Rc, + pub summary_maintenance_lifecycle_guarantee: Option, + pub alternatives: Vec, +} + +#[derive(Debug, Clone)] +pub struct SummaryMaintenanceLifecyclePlan { + pub root: Rc, + pub deployments: Vec, + pub horizon: Option, + pub evaluation_rate: Option, + pub update_rate: Option, + pub expected_reads: Option, + pub selected_raw_recompute: bool, + pub summary_total_cost: Option, + pub raw_recompute_total_cost: Option, +} + +/// Explicit association between a materialized target and the normalized +/// workload entries whose demand consumes it. +#[derive(Debug, Clone, Copy)] +pub struct WorkloadDemand<'a> { + pub workload: &'a QueryWorkload, + pub entry_indices: &'a [usize], +} + +impl<'a> WorkloadDemand<'a> { + pub const fn new(workload: &'a QueryWorkload, entry_indices: &'a [usize]) -> Self { + Self { + workload, + entry_indices, + } + } +} + +#[derive(Debug, thiserror::Error)] +pub enum SummaryMaintenanceLifecyclePlanError { + #[error(transparent)] + InvalidWorkload(#[from] WorkloadError), + #[error("optimization horizon must be finite and strictly positive")] + InvalidHorizon, + #[error("workload entry index {index} is out of bounds for {entry_count} entries")] + InvalidWorkloadEntry { index: usize, entry_count: usize }, + #[error("a workload-demand binding must contain at least one entry")] + EmptyWorkloadDemand, + #[error("workload entry index {index} appears more than once in one demand binding")] + DuplicateWorkloadEntry { index: usize }, +} + +#[derive(Debug, thiserror::Error)] +pub enum MaterializeSummaryMaintenanceLifecycleError { + #[error(transparent)] + Materialize(#[from] ImplementError), + #[error(transparent)] + SummaryMaintenance(#[from] SummaryMaintenanceLifecyclePlanError), +} + +/// Failure while deriving workload-aware candidate costs before global +/// selection. +#[derive(Debug, thiserror::Error)] +pub enum SummaryMaintenanceLifecycleSelectionError { + #[error(transparent)] + Recurrence(#[from] RecurrenceError), + #[error(transparent)] + SummaryMaintenance(#[from] SummaryMaintenanceLifecyclePlanError), +} + +#[derive(Debug)] +struct WorkloadFacts { + reads: Option, + one_time_invocations: u64, + evaluation_rate: Option, + update_rate: Option, + arrival: DataArrival, + prepared_window: Option<(TimestampMs, TimestampMs)>, + prepared_eligible: bool, + requires_deletion: bool, +} + +/// Validate a materialized plan, enumerate lifecycle alternatives for each +/// unique summary state, and select the cheapest legal alternative whose cost +/// is fully known. +pub fn plan_summary_maintenance_lifecycles( + root: Rc, + demand: WorkloadDemand<'_>, + now_ms: u64, + horizon: Option, + capabilities: SummaryMaintenanceLifecycleCapabilities, + cost_model: &dyn CostModel, +) -> Result { + plan_summary_maintenance_lifecycles_with_profile( + root, + demand, + now_ms, + horizon, + capabilities, + cost_model, + None, + ) +} + +/// Internal candidate-costing form. The workload binding supplies temporal +/// eligibility and data-arrival facts; `profile` supplies effective uses after +/// DAG path multiplicity has been propagated by `PlanSpace`. +fn plan_summary_maintenance_lifecycles_with_profile( + root: Rc, + demand: WorkloadDemand<'_>, + now_ms: u64, + horizon: Option, + capabilities: SummaryMaintenanceLifecycleCapabilities, + cost_model: &dyn CostModel, + profile: Option, +) -> Result { + demand.workload.validate()?; + if horizon.is_some_and(|h| !h.0.is_finite() || h.0 <= 0.0) { + return Err(SummaryMaintenanceLifecyclePlanError::InvalidHorizon); + } + let mut facts = workload_facts(demand.workload, demand.entry_indices, now_ms, horizon)?; + if let Some(profile) = profile { + facts.one_time_invocations = u64::try_from(profile.one_shot_consumers).unwrap_or(u64::MAX); + facts.evaluation_rate = profile.evaluation_rate; + facts.update_rate = profile.update_rate; + facts.reads = match (profile.evaluation_rate, horizon) { + (Some(rate), Some(horizon)) => { + Some(profile.one_shot_consumers as f64 + rate.0 * horizon.0) + } + (Some(_), None) => None, + (None, _) if profile.one_shot_consumers > 0 => Some(profile.one_shot_consumers as f64), + // Preserve unknown recurrence from the normalized workload. An + // empty profile does not prove that the target is never read. + (None, _) => facts.reads, + }; + } + let mut summaries = Vec::new(); + collect_summary_aggs(&root, &mut HashSet::new(), &mut summaries); + let components = summary_state_components(&summaries); + let mut deployments: Vec = summaries + .into_iter() + .enumerate() + .map(|(summary_index, summary)| { + let alternatives = alternatives_for( + &facts, + horizon, + capabilities, + cost_model.summary_maintenance_capabilities(&summary), + cost_model.summary_maintenance_lifecycle_cost_inputs(&summary), + ); + SummaryMaintenanceDeployment { + summary_index, + summary, + summary_maintenance_lifecycle_guarantee: None, + alternatives, + } + }) + .collect(); + select_compatible_lifecycles(&mut deployments, &components, facts.arrival); + let summary_total_cost = deployments.iter().try_fold(Cost::ZERO, |sum, deployment| { + let selected = &deployment + .summary_maintenance_lifecycle_guarantee + .as_ref()? + .summary_maintenance_lifecycle; + let cost = deployment + .alternatives + .iter() + .find(|alternative| &alternative.summary_maintenance_lifecycle == selected)? + .total_cost?; + Some(Cost(sum.0 + cost.0)) + }); + Ok(SummaryMaintenanceLifecyclePlan { + root, + deployments, + horizon, + evaluation_rate: facts.evaluation_rate, + update_rate: facts.update_rate, + expected_reads: facts.reads, + selected_raw_recompute: false, + summary_total_cost, + raw_recompute_total_cost: None, + }) +} + +/// Rank semantic summary siblings using the cheapest legal +/// summary-maintenance lifecycle for each candidate before final global +/// selection. The candidate space stays compact; only cost overrides are +/// attached, so shared `Rc` identity and exact-composition commitments remain +/// the responsibility of `GlobalSelection`. +pub fn global_selection_with_summary_maintenance_lifecycles<'a, Id>( + space: &'a PlanSpace, + workload: &QueryWorkload, + root_workload_entries: &[usize], + now_ms: u64, + horizon: Option, + capabilities: SummaryMaintenanceLifecycleCapabilities, + cost_model: &dyn CostModel, +) -> Result, SummaryMaintenanceLifecycleSelectionError> { + let profiles = space.recurrence_profiles_from_workload( + workload, + root_workload_entries, + now_ms, + horizon, + )?; + let bindings = space.workload_entries_by_target(workload, root_workload_entries)?; + let mut costs = CandidateCostOverrides::default(); + for group in space.groups() { + let Some(entry_indices) = bindings.get(&Rc::as_ptr(&group.target)) else { + continue; + }; + for candidate in &group.candidates { + let Replacement::Summary(summary) = &candidate.replacement else { + continue; + }; + let plan = plan_summary_maintenance_lifecycles_with_profile( + Rc::clone(summary), + WorkloadDemand::new(workload, entry_indices), + now_ms, + horizon, + capabilities, + cost_model, + Some(profiles.for_target(&group.target)), + )?; + if !plan.deployments.is_empty() { + if let Some(total) = plan.summary_total_cost { + costs.insert(&group.target, candidate, total); + } + } + } + } + Ok(space.global_selection_with_candidate_costs(cost_model, &profiles, horizon, &costs)?) +} + +/// Materialize a globally selected phase-valid DAG and immediately attach +/// workload-aware summary maintenance deployments. +pub fn materialize_with_summary_maintenance_lifecycles( + selection: &GlobalSelection<'_>, + target: &Rc, + demand: WorkloadDemand<'_>, + now_ms: u64, + horizon: Option, + capabilities: SummaryMaintenanceLifecycleCapabilities, + cost_model: &dyn CostModel, +) -> Result, MaterializeSummaryMaintenanceLifecycleError> { + selection + .materialize(target)? + .map(|root| { + let mut plan = plan_summary_maintenance_lifecycles( + root, + demand, + now_ms, + horizon, + capabilities, + cost_model, + )?; + plan.raw_recompute_total_cost = cost_model + .raw_query_recompute_cost(target) + .zip(plan.expected_reads) + .map(|(per_read, reads)| Cost(per_read.0 * reads)); + if plan.raw_recompute_total_cost.is_some_and(|raw| { + plan.summary_total_cost + .is_none_or(|summary| raw.0 <= summary.0) + }) { + plan.root = crate::replacement::keep_pre_asap(target)?; + plan.deployments.clear(); + plan.selected_raw_recompute = true; + } + Ok(plan) + }) + .transpose() +} + +fn workload_facts( + workload: &QueryWorkload, + workload_entry_indices: &[usize], + now_ms: u64, + horizon: Option, +) -> Result { + let mut one_time_invocations = 0u64; + let mut recurring_reads = 0.0; + let mut recurring_known = true; + let mut evaluation_rate = 0.0; + let mut has_evaluation_rate = false; + let mut prepared_start: Option = None; + let mut prepared_end: Option = None; + let mut prepared_eligible = true; + let mut requires_deletion = false; + + let entries: Vec<_> = workload.entries().collect(); + if workload_entry_indices.is_empty() { + return Err(SummaryMaintenanceLifecyclePlanError::EmptyWorkloadDemand); + } + let mut seen_indices = HashSet::new(); + for &index in workload_entry_indices { + if !seen_indices.insert(index) { + return Err(SummaryMaintenanceLifecyclePlanError::DuplicateWorkloadEntry { index }); + } + let entry = entries.get(index).ok_or( + SummaryMaintenanceLifecyclePlanError::InvalidWorkloadEntry { + index, + entry_count: entries.len(), + }, + )?; + requires_deletion |= entry.time_selection.lookback.is_some() + && entry.time_selection.as_of.is_none() + && matches!( + entry.time_selection.scope, + asap_types::workload::QueryTimeScope::RealTime + | asap_types::workload::QueryTimeScope::Mixed + ); + match &entry.recurrence { + QueryRecurrence::OneTime { + invocations, + execute_at, + } => { + one_time_invocations = one_time_invocations.saturating_add(*invocations); + let covered = if let ( + Predictability::Predictable { + known_at: Some(known), + }, + Some(execute), + ) = (&entry.predictability, execute_at) + { + if known < execute { + prepared_start = Some(prepared_start.map_or(*known, |old| old.min(*known))); + prepared_end = Some(prepared_end.map_or(*execute, |old| old.max(*execute))); + true + } else { + false + } + } else { + false + }; + prepared_eligible &= covered; + } + QueryRecurrence::Repeated(RepeatedDemand::FixedInterval(interval)) => { + prepared_eligible = false; + let rate = 1000.0 / f64::from(interval.0); + evaluation_rate += rate; + has_evaluation_rate = true; + if let Some(h) = horizon { + recurring_reads += h.0 * rate; + } else { + recurring_known = false; + } + } + QueryRecurrence::Repeated(RepeatedDemand::Scheduled(schedule)) => { + prepared_eligible = false; + if let Some(h) = horizon { + let end_ms = now_ms.saturating_add((h.0 * 1000.0) as u64); + let reads_in_horizon = schedule + .iter() + .filter(|at| at.0 >= now_ms && at.0 <= end_ms) + .count() as f64; + recurring_reads += reads_in_horizon; + evaluation_rate += reads_in_horizon / h.0; + has_evaluation_rate = true; + } else { + recurring_known = false; + } + } + QueryRecurrence::Repeated(RepeatedDemand::EstimatedRate(estimate)) => { + prepared_eligible = false; + if !estimate.is_fresh_at(now_ms) { + recurring_known = false; + continue; + } + let rate = match estimate.expected { + asap_types::workload::ExpectedDemand::AverageRate(rate) => Some(rate.0), + asap_types::workload::ExpectedDemand::InvocationCount(count) => { + let millis = estimate + .observation_window + .end + .0 + .saturating_sub(estimate.observation_window.start.0); + (millis > 0).then_some(count as f64 / (millis as f64 / 1000.0)) + } + }; + if let Some(rate) = rate { + evaluation_rate += rate; + has_evaluation_rate = true; + if let Some(h) = horizon { + recurring_reads += h.0 * rate; + } else { + recurring_known = false; + } + } else { + recurring_known = false; + } + } + QueryRecurrence::Unknown => { + prepared_eligible = false; + recurring_known = false; + } + } + } + + let data = workload.data_workload.as_ref(); + let arrival = data.map_or(DataArrival::Unknown, |data| data.arrival); + let update_rate = data + .and_then(|data| data.ingestion_rate.value_at(now_ms)) + .map(|rate| UpdateRate(rate.0)); + let reads = recurring_known.then_some(one_time_invocations as f64 + recurring_reads); + Ok(WorkloadFacts { + reads, + one_time_invocations, + evaluation_rate: has_evaluation_rate.then_some(EvaluationRate(evaluation_rate)), + update_rate, + arrival, + prepared_window: prepared_start.zip(prepared_end), + prepared_eligible, + requires_deletion, + }) +} + +fn alternatives_for( + facts: &WorkloadFacts, + horizon: Option, + capabilities: SummaryMaintenanceLifecycleCapabilities, + summary_capabilities: SummaryMaintenanceCapabilities, + costs: SummaryMaintenanceLifecycleCostInputs, +) -> Vec { + let alternatives = vec![ + ephemeral(facts, capabilities, &costs), + prepared(facts, capabilities, summary_capabilities, &costs), + shared(facts, horizon, capabilities, summary_capabilities, &costs), + continuous(facts, horizon, capabilities, summary_capabilities, &costs), + ]; + alternatives +} + +fn ephemeral( + facts: &WorkloadFacts, + capabilities: SummaryMaintenanceLifecycleCapabilities, + costs: &SummaryMaintenanceLifecycleCostInputs, +) -> SummaryMaintenanceLifecycleAlternative { + let lifecycle = SummaryMaintenanceLifecycle::Ephemeral; + if !capabilities.ephemeral { + return rejected( + lifecycle, + SummaryMaintenanceLifecycleRejection::UnsupportedByRuntime, + ); + } + let total_cost = zip_costs(&[ + costs.build_cost, + costs.summary_read_cost, + costs.retirement_cost, + ]) + .zip(facts.reads) + .map(|(per_read, reads)| Cost(per_read * reads)); + costed_or_unknown( + lifecycle, + total_cost, + vec!["state is rebuilt per invocation".into()], + ) +} + +fn prepared( + facts: &WorkloadFacts, + capabilities: SummaryMaintenanceLifecycleCapabilities, + summary_capabilities: SummaryMaintenanceCapabilities, + costs: &SummaryMaintenanceLifecycleCostInputs, +) -> SummaryMaintenanceLifecycleAlternative { + if !facts.prepared_eligible { + return rejected( + SummaryMaintenanceLifecycle::Prepared { + activate_at: TimestampMs(0), + retire_at: TimestampMs(0), + }, + SummaryMaintenanceLifecycleRejection::RequiresPredictableOneTimeQuery, + ); + } + let Some((activate_at, retire_at)) = facts.prepared_window else { + return rejected( + SummaryMaintenanceLifecycle::Prepared { + activate_at: TimestampMs(0), + retire_at: TimestampMs(0), + }, + SummaryMaintenanceLifecycleRejection::RequiresPredictableOneTimeQuery, + ); + }; + let lifecycle = SummaryMaintenanceLifecycle::Prepared { + activate_at, + retire_at, + }; + if !capabilities.prepared { + return rejected( + lifecycle, + SummaryMaintenanceLifecycleRejection::UnsupportedByRuntime, + ); + } + if let Some(rejection) = maintenance_capability_rejection(facts, summary_capabilities) { + return rejected(lifecycle, rejection); + } + let seconds = retire_at.0.saturating_sub(activate_at.0) as f64 / 1000.0; + let maintenance = maintenance_cost(facts, costs, seconds); + let total_cost = match ( + costs.build_cost, + costs.summary_read_cost, + costs.retention_cost_rate, + costs.retirement_cost, + maintenance, + ) { + (Some(build), Some(read), Some(retention), Some(retire), Some(maintenance)) => Some(Cost( + build.0 + + read.0 * facts.one_time_invocations as f64 + + retention.0 * seconds + + retire.0 + + maintenance, + )), + _ => None, + }; + costed_or_unknown( + lifecycle, + total_cost, + vec!["activation and retirement come from the declared schedule".into()], + ) +} + +fn shared( + facts: &WorkloadFacts, + horizon: Option, + capabilities: SummaryMaintenanceLifecycleCapabilities, + summary_capabilities: SummaryMaintenanceCapabilities, + costs: &SummaryMaintenanceLifecycleCostInputs, +) -> SummaryMaintenanceLifecycleAlternative { + let lifecycle = SummaryMaintenanceLifecycle::Shared { + retention: asap_types::workload::DurationMs(horizon.map_or(0, |h| (h.0 * 1000.0) as u64)), + }; + if !capabilities.shared { + return rejected( + lifecycle, + SummaryMaintenanceLifecycleRejection::UnsupportedByRuntime, + ); + } + if let Some(rejection) = maintenance_capability_rejection(facts, summary_capabilities) { + return rejected(lifecycle, rejection); + } + if facts.reads.is_none_or(|reads| reads <= 1.0) { + return rejected( + lifecycle, + SummaryMaintenanceLifecycleRejection::RequiresMultipleReads, + ); + } + let Some(horizon) = horizon else { + return rejected( + lifecycle, + SummaryMaintenanceLifecycleRejection::RequiresHorizon, + ); + }; + let total_cost = retained_cost(facts, costs, horizon.0); + costed_or_unknown( + lifecycle, + total_cost, + vec!["one state is shared across reads".into()], + ) +} + +fn continuous( + facts: &WorkloadFacts, + horizon: Option, + capabilities: SummaryMaintenanceLifecycleCapabilities, + summary_capabilities: SummaryMaintenanceCapabilities, + costs: &SummaryMaintenanceLifecycleCostInputs, +) -> SummaryMaintenanceLifecycleAlternative { + let lifecycle = SummaryMaintenanceLifecycle::ContinuouslyMaintained; + if !capabilities.continuously_maintained { + return rejected( + lifecycle, + SummaryMaintenanceLifecycleRejection::UnsupportedByRuntime, + ); + } + if !matches!( + facts.arrival, + DataArrival::ContinuouslyIngesting | DataArrival::Mixed + ) { + return rejected( + lifecycle, + SummaryMaintenanceLifecycleRejection::RequiresContinuousData, + ); + } + if facts.update_rate.is_none() { + return rejected( + lifecycle, + SummaryMaintenanceLifecycleRejection::MissingOrStaleIngestionRate, + ); + } + if let Some(rejection) = maintenance_capability_rejection(facts, summary_capabilities) { + return rejected(lifecycle, rejection); + } + let Some(horizon) = horizon else { + return rejected( + lifecycle, + SummaryMaintenanceLifecycleRejection::RequiresHorizon, + ); + }; + let total_cost = retained_cost(facts, costs, horizon.0); + costed_or_unknown( + lifecycle, + total_cost, + vec!["updates are applied for the optimization horizon".into()], + ) +} + +fn maintenance_capability_rejection( + facts: &WorkloadFacts, + capabilities: SummaryMaintenanceCapabilities, +) -> Option { + if matches!( + facts.arrival, + DataArrival::ContinuouslyIngesting | DataArrival::Mixed + ) && !capabilities.incremental_update + { + Some(SummaryMaintenanceLifecycleRejection::SummaryDoesNotSupportIncrementalUpdates) + } else if matches!( + facts.arrival, + DataArrival::ContinuouslyIngesting | DataArrival::Mixed + ) && facts.requires_deletion + && !capabilities.delete + { + Some(SummaryMaintenanceLifecycleRejection::SummaryDoesNotSupportDeletion) + } else { + None + } +} + +fn retained_cost( + facts: &WorkloadFacts, + costs: &SummaryMaintenanceLifecycleCostInputs, + seconds: f64, +) -> Option { + let reads = facts.reads?; + let maintenance = maintenance_cost(facts, costs, seconds)?; + Some(Cost( + costs.build_cost?.0 + + maintenance + + reads * costs.summary_read_cost?.0 + + seconds * costs.retention_cost_rate?.0 + + costs.retirement_cost?.0, + )) +} + +fn maintenance_cost( + facts: &WorkloadFacts, + costs: &SummaryMaintenanceLifecycleCostInputs, + seconds: f64, +) -> Option { + match facts.arrival { + DataArrival::AtRest => Some(0.0), + DataArrival::ContinuouslyIngesting | DataArrival::Mixed => { + Some(seconds * facts.update_rate?.0 * costs.maintenance_cost_per_update?.0) + } + DataArrival::Unknown => None, + } +} + +fn zip_costs(costs: &[Option]) -> Option { + costs + .iter() + .try_fold(0.0, |sum, cost| Some(sum + cost.as_ref()?.0)) +} + +fn costed_or_unknown( + summary_maintenance_lifecycle: SummaryMaintenanceLifecycle, + total_cost: Option, + assumptions: Vec, +) -> SummaryMaintenanceLifecycleAlternative { + SummaryMaintenanceLifecycleAlternative { + summary_maintenance_lifecycle, + total_cost, + rejection: total_cost + .is_none() + .then_some(SummaryMaintenanceLifecycleRejection::MissingCostEvidence), + assumptions, + } +} + +fn rejected( + summary_maintenance_lifecycle: SummaryMaintenanceLifecycle, + rejection: SummaryMaintenanceLifecycleRejection, +) -> SummaryMaintenanceLifecycleAlternative { + SummaryMaintenanceLifecycleAlternative { + summary_maintenance_lifecycle, + total_cost: None, + rejection: Some(rejection), + assumptions: Vec::new(), + } +} + +fn collect_summary_aggs( + node: &Rc, + seen: &mut HashSet<*const SummaryNode>, + output: &mut Vec>, +) { + if !seen.insert(Rc::as_ptr(node)) { + return; + } + match &node.expr { + SummaryExpr::SummaryAgg { child, .. } => { + output.push(Rc::clone(node)); + collect_summary_aggs(child, seen, output); + } + SummaryExpr::SummaryJoin { outer, inner, .. } + | SummaryExpr::SummarySubtract { + left: outer, + right: inner, + } => { + collect_summary_aggs(outer, seen, output); + collect_summary_aggs(inner, seen, output); + } + SummaryExpr::SummaryDelete { summary_input, .. } + | SummaryExpr::SummaryEstimate { summary_input, .. } => { + collect_summary_aggs(summary_input, seen, output) + } + SummaryExpr::SummaryMerge { children } => { + for child in children { + collect_summary_aggs(child, seen, output); + } + } + SummaryExpr::KeepPreAsap(_) => {} + } +} + +fn evaluation_schedule( + lifecycle: &SummaryMaintenanceLifecycle, + arrival: DataArrival, +) -> EvaluationSchedule { + match lifecycle { + SummaryMaintenanceLifecycle::Ephemeral => EvaluationSchedule::OneShot, + SummaryMaintenanceLifecycle::Prepared { .. } + | SummaryMaintenanceLifecycle::Shared { .. } + if matches!( + arrival, + DataArrival::ContinuouslyIngesting | DataArrival::Mixed + ) => + { + EvaluationSchedule::PerUpdate + } + SummaryMaintenanceLifecycle::Prepared { .. } => EvaluationSchedule::OneShot, + SummaryMaintenanceLifecycle::Shared { .. } => EvaluationSchedule::OnRead, + SummaryMaintenanceLifecycle::ContinuouslyMaintained => EvaluationSchedule::PerUpdate, + } +} + +/// Summary states composed on one maintenance path must be produced on the +/// same schedule. Return a component id for each collected `SummaryAgg`. +fn summary_state_components(summaries: &[Rc]) -> Vec { + let indices: HashMap<_, _> = summaries + .iter() + .enumerate() + .map(|(index, summary)| (Rc::as_ptr(summary), index)) + .collect(); + let mut parents: Vec<_> = (0..summaries.len()).collect(); + + fn find(parents: &mut [usize], index: usize) -> usize { + if parents[index] != index { + parents[index] = find(parents, parents[index]); + } + parents[index] + } + + for (parent_index, summary) in summaries.iter().enumerate() { + let SummaryExpr::SummaryAgg { child, .. } = &summary.expr else { + continue; + }; + if !matches!( + child.expr, + SummaryExpr::SummaryAgg { .. } + | SummaryExpr::SummaryJoin { .. } + | SummaryExpr::SummarySubtract { .. } + | SummaryExpr::SummaryDelete { .. } + | SummaryExpr::SummaryMerge { .. } + ) { + continue; + } + let mut descendants = Vec::new(); + collect_summary_aggs(child, &mut HashSet::new(), &mut descendants); + for descendant in descendants { + let child_index = indices[&Rc::as_ptr(&descendant)]; + let parent_root = find(&mut parents, parent_index); + let child_root = find(&mut parents, child_index); + parents[child_root] = parent_root; + } + } + (0..parents.len()) + .map(|index| find(&mut parents, index)) + .collect() +} + +fn select_compatible_lifecycles( + deployments: &mut [SummaryMaintenanceDeployment], + components: &[usize], + arrival: DataArrival, +) { + let component_ids: HashSet<_> = components.iter().copied().collect(); + for component in component_ids { + let members: Vec<_> = components + .iter() + .enumerate() + .filter_map(|(index, &id)| (id == component).then_some(index)) + .collect(); + let selected_schedule = [ + EvaluationSchedule::OneShot, + EvaluationSchedule::PerUpdate, + EvaluationSchedule::OnRead, + ] + .into_iter() + .filter_map(|schedule| { + members + .iter() + .try_fold(0.0, |sum, &index| { + deployments[index] + .alternatives + .iter() + .filter(|candidate| { + candidate.selectable() + && evaluation_schedule( + &candidate.summary_maintenance_lifecycle, + arrival, + ) == schedule + }) + .map(|candidate| candidate.total_cost.unwrap().0) + .min_by(f64::total_cmp) + .map(|cost| sum + cost) + }) + .map(|cost| (schedule, cost)) + }) + .min_by(|(_, a), (_, b)| a.total_cmp(b)) + .map(|(schedule, _)| schedule); + + let Some(schedule) = selected_schedule else { + continue; + }; + for index in members { + let selected = deployments[index] + .alternatives + .iter() + .filter(|candidate| { + candidate.selectable() + && evaluation_schedule(&candidate.summary_maintenance_lifecycle, arrival) + == schedule + }) + .min_by(|a, b| a.total_cost.unwrap().0.total_cmp(&b.total_cost.unwrap().0)); + deployments[index].summary_maintenance_lifecycle_guarantee = + selected.map(|candidate| SummaryMaintenanceLifecycleGuarantee { + summary_maintenance_lifecycle: candidate.summary_maintenance_lifecycle.clone(), + summary_maintenance_mode: maintenance_mode( + &candidate.summary_maintenance_lifecycle, + arrival, + ), + evaluation_schedule: schedule, + output_representation: OutputRepresentation::SummaryState, + }); + } + } +} + +fn maintenance_mode( + lifecycle: &SummaryMaintenanceLifecycle, + arrival: DataArrival, +) -> SummaryMaintenanceMode { + match lifecycle { + SummaryMaintenanceLifecycle::Ephemeral => SummaryMaintenanceMode::DirectBuild, + SummaryMaintenanceLifecycle::ContinuouslyMaintained => SummaryMaintenanceMode::Incremental, + SummaryMaintenanceLifecycle::Prepared { .. } + | SummaryMaintenanceLifecycle::Shared { .. } => match arrival { + DataArrival::ContinuouslyIngesting | DataArrival::Mixed => { + SummaryMaintenanceMode::Incremental + } + DataArrival::AtRest | DataArrival::Unknown => SummaryMaintenanceMode::DirectBuild, + }, + } +} + +#[cfg(test)] +mod tests { + use super::*; + use asap_types::post_asap::{ + ExactKind, ExactParams, GroupingStrategy, ResultGuarantee, SketchAlgorithm, + SummaryFamilyType, SummaryField, SummarySchema, + }; + use asap_types::pre_asap::AggIntent; + use asap_types::pre_asap::{Column, ColumnRef, DataType, QueryExpr, Reduction, Schema, Source}; + use asap_types::types::AccuracyTarget; + use asap_types::workload::{ + BatchEntry, DataWorkload, DurationMs, Evidence, EvidenceSource, Predictability, Query, + QueryLanguage, QueryRequirements, Rate, RepeatingEntry, RepetitionInterval, TimeSelection, + }; + + struct UnitCosts; + + impl CostModel for UnitCosts { + fn rank_candidates( + &self, + _intent: &asap_types::pre_asap::AggIntent, + candidates: &[asap_types::post_asap::SketchAlgorithm], + ) -> Vec { + candidates.to_vec() + } + + fn summary_maintenance_lifecycle_cost_inputs( + &self, + _summary: &SummaryNode, + ) -> SummaryMaintenanceLifecycleCostInputs { + SummaryMaintenanceLifecycleCostInputs { + build_cost: Some(Cost(10.0)), + maintenance_cost_per_update: Some(Cost(1.0)), + summary_read_cost: Some(Cost(1.0)), + retention_cost_rate: Some(CostRate(0.1)), + retirement_cost: Some(Cost(1.0)), + } + } + + fn summary_maintenance_capabilities( + &self, + _summary: &SummaryNode, + ) -> SummaryMaintenanceCapabilities { + SummaryMaintenanceCapabilities { + incremental_update: true, + merge: true, + delete: true, + } + } + } + + struct RawCheaper; + + impl CostModel for RawCheaper { + fn rank_candidates( + &self, + _intent: &asap_types::pre_asap::AggIntent, + candidates: &[asap_types::post_asap::SketchAlgorithm], + ) -> Vec { + candidates.to_vec() + } + + fn summary_maintenance_lifecycle_cost_inputs( + &self, + summary: &SummaryNode, + ) -> SummaryMaintenanceLifecycleCostInputs { + UnitCosts.summary_maintenance_lifecycle_cost_inputs(summary) + } + + fn summary_maintenance_capabilities( + &self, + summary: &SummaryNode, + ) -> SummaryMaintenanceCapabilities { + UnitCosts.summary_maintenance_capabilities(summary) + } + + fn raw_query_recompute_cost(&self, _target: &QueryExpr) -> Option { + Some(Cost(1.0)) + } + } + + struct NoDelete; + + impl CostModel for NoDelete { + fn rank_candidates( + &self, + _intent: &asap_types::pre_asap::AggIntent, + candidates: &[asap_types::post_asap::SketchAlgorithm], + ) -> Vec { + candidates.to_vec() + } + + fn summary_maintenance_lifecycle_cost_inputs( + &self, + summary: &SummaryNode, + ) -> SummaryMaintenanceLifecycleCostInputs { + UnitCosts.summary_maintenance_lifecycle_cost_inputs(summary) + } + + fn summary_maintenance_capabilities( + &self, + _summary: &SummaryNode, + ) -> SummaryMaintenanceCapabilities { + SummaryMaintenanceCapabilities { + incremental_update: true, + merge: true, + delete: false, + } + } + } + + struct SummaryMaintenancePrefersDdSketch; + + impl CostModel for SummaryMaintenancePrefersDdSketch { + fn rank_candidates( + &self, + _intent: &AggIntent, + candidates: &[SketchAlgorithm], + ) -> Vec { + // Preserve semantic mapping's KLL-first order. The lifecycle + // total below must be what changes the final choice. + candidates.to_vec() + } + + fn summary_maintenance_lifecycle_cost_inputs( + &self, + summary: &SummaryNode, + ) -> SummaryMaintenanceLifecycleCostInputs { + let build = match sketch_algorithm(summary) { + Some(SketchAlgorithm::Kll) => 100.0, + Some(SketchAlgorithm::DDSketch) => 1.0, + _ => 10.0, + }; + SummaryMaintenanceLifecycleCostInputs { + build_cost: Some(Cost(build)), + maintenance_cost_per_update: Some(Cost(1.0)), + summary_read_cost: Some(Cost(1.0)), + retention_cost_rate: Some(CostRate(0.1)), + retirement_cost: Some(Cost(1.0)), + } + } + } + + struct IncompatibleNestedCosts; + + impl CostModel for IncompatibleNestedCosts { + fn rank_candidates( + &self, + _intent: &AggIntent, + candidates: &[SketchAlgorithm], + ) -> Vec { + candidates.to_vec() + } + + fn summary_maintenance_lifecycle_cost_inputs( + &self, + summary: &SummaryNode, + ) -> SummaryMaintenanceLifecycleCostInputs { + let is_leaf = matches!( + summary.expr, + SummaryExpr::SummaryAgg { ref child, .. } + if matches!(child.expr, SummaryExpr::KeepPreAsap(_)) + ); + SummaryMaintenanceLifecycleCostInputs { + build_cost: Some(Cost(if is_leaf { 1.0 } else { 100.0 })), + maintenance_cost_per_update: Some(Cost(if is_leaf { 100.0 } else { 0.0 })), + summary_read_cost: Some(Cost::ZERO), + retention_cost_rate: Some(CostRate(0.0)), + retirement_cost: Some(Cost::ZERO), + } + } + + fn summary_maintenance_capabilities( + &self, + _summary: &SummaryNode, + ) -> SummaryMaintenanceCapabilities { + SummaryMaintenanceCapabilities { + incremental_update: true, + merge: true, + delete: true, + } + } + } + + fn sketch_algorithm(node: &SummaryNode) -> Option { + match &node.expr { + SummaryExpr::SummaryEstimate { summary_input, .. } => sketch_algorithm(summary_input), + SummaryExpr::SummaryAgg { + family: SummaryFamilyType::Sketch(kind, _), + .. + } => Some(kind.algorithm().clone()), + _ => None, + } + } + + fn query_root() -> Rc { + query_root_for("m") + } + + fn query_root_for(metric: &str) -> Rc { + Rc::new(QueryExpr::Scan { + source: Source::TimeSeries { + metric: metric.into(), + }, + predicates: vec![], + schema: Schema::with_time_index( + vec![ + Column::new("ts", DataType::Timestamp, false), + Column::new("value", DataType::Float64, false), + ], + 0, + vec![], + ), + }) + } + + fn sum_query() -> Rc { + Rc::new(QueryExpr::Aggregate { + reduction: Reduction::by(vec![]), + measures: vec![AggIntent::Sum { col: None }], + output_names: vec![], + having: None, + child: query_root(), + }) + } + + fn quantile_query() -> Rc { + Rc::new(QueryExpr::Aggregate { + reduction: Reduction::by(vec![]), + measures: vec![AggIntent::Quantile { + col: None, + q: 0.99, + accuracy: AccuracyTarget::Epsilon(0.1), + }], + output_names: vec![], + having: None, + child: query_root(), + }) + } + + fn summary() -> Rc { + let child = Rc::new(SummaryNode { + expr: SummaryExpr::KeepPreAsap(query_root()), + schema: SummarySchema { + fields: vec![], + time_index: None, + }, + guarantee: Some(ResultGuarantee::exact("raw")), + }); + let family = SummaryFamilyType::ExactAggregate(ExactKind::Sum, ExactParams::Sum); + Rc::new(SummaryNode { + expr: SummaryExpr::SummaryAgg { + child, + family: family.clone(), + col: ColumnRef::Named("value".into()), + reduction: Reduction::by(vec![]), + grouping: GroupingStrategy::default(), + }, + schema: SummarySchema { + fields: vec![SummaryField { + name: "state".into(), + dtype: family, + nullable: false, + }], + time_index: None, + }, + guarantee: Some(ResultGuarantee::exact("sum")), + }) + } + + fn nested_summary() -> Rc { + let child = summary(); + let family = SummaryFamilyType::ExactAggregate(ExactKind::Sum, ExactParams::Sum); + Rc::new(SummaryNode { + expr: SummaryExpr::SummaryAgg { + child, + family: family.clone(), + col: ColumnRef::Named("state".into()), + reduction: Reduction::by(vec![]), + grouping: GroupingStrategy::default(), + }, + schema: SummarySchema { + fields: vec![SummaryField { + name: "state".into(), + dtype: family, + nullable: false, + }], + time_index: None, + }, + guarantee: Some(ResultGuarantee::exact("nested sum")), + }) + } + + fn batch(predictability: Predictability) -> BatchEntry { + BatchEntry { + query: Query("sum(m)".into()), + requirements: QueryRequirements::default(), + predictability, + invocations: 1, + execute_at: None, + time_selection: TimeSelection::default(), + } + } + + fn workload( + batches: Vec, + repeating: Vec, + data: DataWorkload, + ) -> QueryWorkload { + QueryWorkload { + language: QueryLanguage::PromQL, + query_batch: (!batches.is_empty()).then_some(batches), + repeating_queries: (!repeating.is_empty()).then_some(repeating), + data_workload: Some(data), + } + } + + fn at_rest() -> DataWorkload { + DataWorkload { + arrival: DataArrival::AtRest, + ..Default::default() + } + } + + fn continuous(observed_at_ms: u64, valid_for_ms: u64) -> DataWorkload { + DataWorkload { + arrival: DataArrival::ContinuouslyIngesting, + ingestion_rate: Evidence { + value: Some(Rate(1.0)), + source: EvidenceSource::Observed, + observed_at_ms: Some(observed_at_ms), + valid_for_ms: Some(valid_for_ms), + }, + ..Default::default() + } + } + + fn repeating() -> RepeatingEntry { + RepeatingEntry { + query: Query("sum(m)".into()), + demand: RepeatedDemand::FixedInterval(RepetitionInterval(1_000)), + requirements: QueryRequirements::default(), + predictability: Predictability::Predictable { known_at: None }, + time_selection: TimeSelection::default(), + } + } + + fn selected_summary_maintenance_lifecycle( + deployment: &SummaryMaintenanceDeployment, + ) -> Option<&SummaryMaintenanceLifecycle> { + deployment + .summary_maintenance_lifecycle_guarantee + .as_ref() + .map(|guarantee| &guarantee.summary_maintenance_lifecycle) + } + + #[test] + fn unpredictable_one_time_at_rest_selects_ephemeral() { + let plan = plan_summary_maintenance_lifecycles( + summary(), + WorkloadDemand::new( + &workload(vec![batch(Predictability::AdHoc)], vec![], at_rest()), + &[0], + ), + 1_000, + None, + SummaryMaintenanceLifecycleCapabilities::ALL, + &UnitCosts, + ) + .unwrap(); + assert_eq!(plan.deployments.len(), 1); + assert_eq!( + selected_summary_maintenance_lifecycle(&plan.deployments[0]), + Some(&SummaryMaintenanceLifecycle::Ephemeral) + ); + let guarantee = plan.deployments[0] + .summary_maintenance_lifecycle_guarantee + .as_ref() + .unwrap(); + assert_eq!(guarantee.evaluation_schedule, EvaluationSchedule::OneShot); + assert_eq!( + guarantee.summary_maintenance_mode, + SummaryMaintenanceMode::DirectBuild + ); + assert_eq!( + guarantee.output_representation, + OutputRepresentation::SummaryState + ); + assert_eq!( + plan.deployments[0].alternatives[0].total_cost, + Some(Cost(12.0)) + ); + } + + #[test] + fn predictable_scheduled_one_time_offers_prepared_state() { + let mut entry = batch(Predictability::Predictable { + known_at: Some(TimestampMs(1_000)), + }); + entry.execute_at = Some(TimestampMs(11_000)); + let plan = plan_summary_maintenance_lifecycles( + summary(), + WorkloadDemand::new(&workload(vec![entry], vec![], at_rest()), &[0]), + 1_000, + None, + SummaryMaintenanceLifecycleCapabilities::ALL, + &UnitCosts, + ) + .unwrap(); + let prepared = &plan.deployments[0].alternatives[1]; + assert!(prepared.rejection.is_none()); + assert_eq!(prepared.total_cost, Some(Cost(13.0))); + } + + #[test] + fn nested_summary_lifecycles_have_compatible_evaluation_schedules() { + let workload = workload(vec![], vec![repeating()], continuous(1_000, 20_000)); + let plan = plan_summary_maintenance_lifecycles( + nested_summary(), + WorkloadDemand::new(&workload, &[0]), + 1_000, + Some(Horizon(10.0)), + SummaryMaintenanceLifecycleCapabilities::ALL, + &IncompatibleNestedCosts, + ) + .unwrap(); + + assert_eq!(plan.deployments.len(), 2); + let schedules: HashSet<_> = plan + .deployments + .iter() + .map(|deployment| { + deployment + .summary_maintenance_lifecycle_guarantee + .as_ref() + .unwrap() + .evaluation_schedule + }) + .collect(); + assert_eq!(schedules.len(), 1); + } + + #[test] + fn repeated_at_rest_selects_shared_without_inventing_updates() { + let plan = plan_summary_maintenance_lifecycles( + summary(), + WorkloadDemand::new(&workload(vec![], vec![repeating()], at_rest()), &[0]), + 1_000, + Some(Horizon(10.0)), + SummaryMaintenanceLifecycleCapabilities::ALL, + &UnitCosts, + ) + .unwrap(); + assert_eq!( + selected_summary_maintenance_lifecycle(&plan.deployments[0]), + Some(&SummaryMaintenanceLifecycle::Shared { + retention: DurationMs(10_000) + }) + ); + assert_eq!( + plan.deployments[0] + .summary_maintenance_lifecycle_guarantee + .as_ref() + .unwrap() + .summary_maintenance_mode, + SummaryMaintenanceMode::DirectBuild + ); + assert_eq!( + plan.deployments[0].alternatives[3].rejection, + Some(SummaryMaintenanceLifecycleRejection::RequiresContinuousData) + ); + assert_eq!(plan.update_rate, None); + } + + #[test] + fn repeated_continuous_workload_can_select_continuous_maintenance() { + let capabilities = SummaryMaintenanceLifecycleCapabilities { + shared: false, + ..SummaryMaintenanceLifecycleCapabilities::ALL + }; + let plan = plan_summary_maintenance_lifecycles( + summary(), + WorkloadDemand::new( + &workload(vec![], vec![repeating()], continuous(1_000, 60_000)), + &[0], + ), + 1_000, + Some(Horizon(10.0)), + capabilities, + &UnitCosts, + ) + .unwrap(); + assert_eq!( + selected_summary_maintenance_lifecycle(&plan.deployments[0]), + Some(&SummaryMaintenanceLifecycle::ContinuouslyMaintained) + ); + assert_eq!( + plan.deployments[0] + .summary_maintenance_lifecycle_guarantee + .as_ref() + .unwrap() + .summary_maintenance_mode, + SummaryMaintenanceMode::Incremental + ); + assert_eq!(plan.evaluation_rate, Some(EvaluationRate(1.0))); + assert_eq!(plan.update_rate, Some(UpdateRate(1.0))); + } + + #[test] + fn stale_ingestion_evidence_cannot_enable_continuous_maintenance() { + let plan = plan_summary_maintenance_lifecycles( + summary(), + WorkloadDemand::new( + &workload(vec![], vec![repeating()], continuous(1_000, 1_000)), + &[0], + ), + 3_000, + Some(Horizon(10.0)), + SummaryMaintenanceLifecycleCapabilities::ALL, + &UnitCosts, + ) + .unwrap(); + assert_eq!( + plan.deployments[0].alternatives[3].rejection, + Some(SummaryMaintenanceLifecycleRejection::MissingOrStaleIngestionRate) + ); + assert_eq!(plan.update_rate, None); + } + + #[test] + fn unknown_costs_do_not_make_a_long_lived_lifecycle_win() { + let plan = plan_summary_maintenance_lifecycles( + summary(), + WorkloadDemand::new( + &workload(vec![], vec![repeating()], continuous(1_000, 60_000)), + &[0], + ), + 1_000, + Some(Horizon(10.0)), + SummaryMaintenanceLifecycleCapabilities::ALL, + &crate::cost_model::DefaultCostModel, + ) + .unwrap(); + assert_eq!( + selected_summary_maintenance_lifecycle(&plan.deployments[0]), + None + ); + assert!(plan.deployments[0] + .alternatives + .iter() + .all(|alternative| alternative.rejection.is_some())); + } + + #[test] + fn unrelated_workload_entries_do_not_create_reuse_for_a_target() { + let plan = plan_summary_maintenance_lifecycles( + summary(), + WorkloadDemand::new( + &workload( + vec![batch(Predictability::AdHoc), batch(Predictability::AdHoc)], + vec![], + at_rest(), + ), + &[0], + ), + 1_000, + Some(Horizon(10.0)), + SummaryMaintenanceLifecycleCapabilities::ALL, + &UnitCosts, + ) + .unwrap(); + assert_eq!( + selected_summary_maintenance_lifecycle(&plan.deployments[0]), + Some(&SummaryMaintenanceLifecycle::Ephemeral) + ); + assert_eq!( + plan.deployments[0].alternatives[2].rejection, + Some(SummaryMaintenanceLifecycleRejection::RequiresMultipleReads) + ); + } + + #[test] + fn scheduled_rate_counts_only_executions_inside_the_horizon() { + let mut entry = repeating(); + entry.demand = RepeatedDemand::Scheduled(vec![ + TimestampMs(999), + TimestampMs(5_000), + TimestampMs(20_000), + ]); + let plan = plan_summary_maintenance_lifecycles( + summary(), + WorkloadDemand::new(&workload(vec![], vec![entry], at_rest()), &[0]), + 1_000, + Some(Horizon(10.0)), + SummaryMaintenanceLifecycleCapabilities::ALL, + &UnitCosts, + ) + .unwrap(); + assert_eq!(plan.evaluation_rate, Some(EvaluationRate(0.1))); + } + + #[test] + fn demand_binding_rejects_empty_and_duplicate_entries() { + let workload = workload(vec![batch(Predictability::AdHoc)], vec![], at_rest()); + assert!(matches!( + plan_summary_maintenance_lifecycles( + summary(), + WorkloadDemand::new(&workload, &[]), + 1_000, + None, + SummaryMaintenanceLifecycleCapabilities::ALL, + &UnitCosts, + ), + Err(SummaryMaintenanceLifecyclePlanError::EmptyWorkloadDemand) + )); + assert!(matches!( + plan_summary_maintenance_lifecycles( + summary(), + WorkloadDemand::new(&workload, &[0, 0]), + 1_000, + None, + SummaryMaintenanceLifecycleCapabilities::ALL, + &UnitCosts, + ), + Err(SummaryMaintenanceLifecyclePlanError::DuplicateWorkloadEntry { index: 0 }) + )); + } + + #[test] + fn prepared_requires_every_bound_consumer_to_be_scheduled_and_predictable() { + let mut predictable = batch(Predictability::Predictable { + known_at: Some(TimestampMs(1_000)), + }); + predictable.execute_at = Some(TimestampMs(2_000)); + let workload = workload( + vec![predictable, batch(Predictability::AdHoc)], + vec![], + at_rest(), + ); + let plan = plan_summary_maintenance_lifecycles( + summary(), + WorkloadDemand::new(&workload, &[0, 1]), + 1_000, + Some(Horizon(10.0)), + SummaryMaintenanceLifecycleCapabilities::ALL, + &UnitCosts, + ) + .unwrap(); + assert_eq!( + plan.deployments[0].alternatives[1].rejection, + Some(SummaryMaintenanceLifecycleRejection::RequiresPredictableOneTimeQuery) + ); + } + + #[test] + fn moving_realtime_maintenance_requires_summary_deletion_support() { + let mut entry = repeating(); + entry.time_selection = TimeSelection { + scope: asap_types::workload::QueryTimeScope::RealTime, + lookback: Some(DurationMs(60_000)), + as_of: None, + }; + let plan = plan_summary_maintenance_lifecycles( + summary(), + WorkloadDemand::new( + &workload(vec![], vec![entry], continuous(1_000, 60_000)), + &[0], + ), + 1_000, + Some(Horizon(10.0)), + SummaryMaintenanceLifecycleCapabilities::ALL, + &NoDelete, + ) + .unwrap(); + assert_eq!( + plan.deployments[0].alternatives[3].rejection, + Some(SummaryMaintenanceLifecycleRejection::SummaryDoesNotSupportDeletion) + ); + } + + #[test] + fn lifecycle_cost_can_fall_back_to_raw_recomputation() { + let target = sum_query(); + let space = crate::replacement::search_workload(vec![("q", Rc::clone(&target))]); + let selection = space.global_selection(&RawCheaper); + let workload = workload(vec![batch(Predictability::AdHoc)], vec![], at_rest()); + let plan = materialize_with_summary_maintenance_lifecycles( + &selection, + &space.roots[0].1, + WorkloadDemand::new(&workload, &[0]), + 1_000, + None, + SummaryMaintenanceLifecycleCapabilities::ALL, + &RawCheaper, + ) + .unwrap() + .unwrap(); + assert!(plan.selected_raw_recompute); + assert_eq!(plan.raw_recompute_total_cost, Some(Cost(1.0))); + assert!(plan.deployments.is_empty()); + assert!(matches!(plan.root.expr, SummaryExpr::KeepPreAsap(_))); + } + + #[test] + fn lifecycle_cost_reorders_semantic_summary_candidates_before_materialization() { + let target = quantile_query(); + let space = crate::replacement::search_workload(vec![("q", target)]); + let workload = workload(vec![batch(Predictability::AdHoc)], vec![], at_rest()); + + let selection = global_selection_with_summary_maintenance_lifecycles( + &space, + &workload, + &[0], + 1_000, + None, + SummaryMaintenanceLifecycleCapabilities::ALL, + &SummaryMaintenancePrefersDdSketch, + ) + .unwrap(); + let materialized = selection.materialize(&space.roots[0].1).unwrap().unwrap(); + + assert_eq!( + sketch_algorithm(&materialized), + Some(SketchAlgorithm::DDSketch) + ); + } + + #[test] + fn lifecycle_cost_counts_one_shared_summary_node_once() { + let shared = summary(); + let root = Rc::new(SummaryNode { + expr: SummaryExpr::SummaryMerge { + children: vec![Rc::clone(&shared), Rc::clone(&shared)], + }, + schema: shared.schema.clone(), + guarantee: None, + }); + let workload = workload( + vec![batch(Predictability::AdHoc), batch(Predictability::AdHoc)], + vec![], + at_rest(), + ); + let horizon = Some(Horizon(10.0)); + let plan = plan_summary_maintenance_lifecycles( + root, + WorkloadDemand::new(&workload, &[0, 1]), + 1_000, + horizon, + SummaryMaintenanceLifecycleCapabilities::ALL, + &UnitCosts, + ) + .unwrap(); + assert_eq!(plan.deployments.len(), 1); + assert!(matches!( + selected_summary_maintenance_lifecycle(&plan.deployments[0]), + Some(SummaryMaintenanceLifecycle::Shared { .. }) + )); + } + + #[test] + fn normalized_workload_drives_plan_space_recurrence_profiles() { + let root = query_root(); + let space = crate::replacement::search_workload(vec![("dashboard", Rc::clone(&root))]); + let workload = workload(vec![], vec![repeating()], continuous(1_000, 60_000)); + let profiles = space + .recurrence_profiles_from_workload(&workload, &[0], 1_000, Some(Horizon(10.0))) + .unwrap(); + // `search_workload` canonicalizes roots through CSE; recurrence + // profiles are keyed by that canonical post-CSE node. + let profile = profiles.for_target(&space.roots[0].1); + assert_eq!(profile.evaluation_rate, Some(EvaluationRate(1.0))); + assert_eq!(profile.update_rate, Some(UpdateRate(1.0))); + assert_eq!(profile.one_shot_consumers, 0); + } + + #[test] + fn recurrence_binding_is_explicit_when_root_order_differs_from_workload_order() { + let repeating_root = query_root_for("dashboard"); + let batch_root = query_root_for("batch"); + let space = crate::replacement::search_workload(vec![ + ("dashboard", repeating_root), + ("batch", batch_root), + ]); + let workload = workload( + vec![batch(Predictability::AdHoc)], + vec![repeating()], + at_rest(), + ); + let profiles = space + .recurrence_profiles_from_workload(&workload, &[1, 0], 1_000, Some(Horizon(10.0))) + .unwrap(); + let dashboard = profiles.for_target(&space.roots[0].1); + let batch = profiles.for_target(&space.roots[1].1); + assert_eq!(dashboard.evaluation_rate, Some(EvaluationRate(1.0))); + assert_eq!(dashboard.one_shot_consumers, 0); + assert_eq!(batch.evaluation_rate, None); + assert_eq!(batch.one_shot_consumers, 1); + } +} diff --git a/crates/types/src/post_asap/lifecycle.rs b/crates/types/src/post_asap/lifecycle.rs deleted file mode 100644 index 995c2f79..00000000 --- a/crates/types/src/post_asap/lifecycle.rs +++ /dev/null @@ -1,37 +0,0 @@ -//! Physical lifecycle vocabulary for summary state. -//! -//! These choices are attached by physical planning; a `SummaryAgg` does not -//! imply continuous maintenance by itself. - -use crate::workload::{DurationMs, TimestampMs}; - -/// When an operator is evaluated. This is independent of whether it owns -/// state and how long that state is retained. -#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] -pub enum EvaluationSchedule { - OneShot, - PerUpdate, - OnRead, -} - -/// The physical value crossing an execution boundary. -#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] -pub enum OutputRepresentation { - PlainRows, - SummaryState, - FinalizedValue, -} - -/// How long one planned summary state deployment exists. -#[derive(Debug, Clone, PartialEq, Eq, Hash)] -pub enum StateLifecycle { - Ephemeral, - Prepared { - activate_at: TimestampMs, - retire_at: TimestampMs, - }, - Shared { - retention: DurationMs, - }, - ContinuouslyMaintained, -} diff --git a/crates/types/src/post_asap/mod.rs b/crates/types/src/post_asap/mod.rs index 5c572b05..3b424bba 100644 --- a/crates/types/src/post_asap/mod.rs +++ b/crates/types/src/post_asap/mod.rs @@ -29,18 +29,17 @@ pub mod expr; pub mod guarantee; -pub mod lifecycle; pub mod query_time; pub mod schema; pub mod sketch; pub mod summary_maintenance; +pub mod summary_maintenance_lifecycle; pub use expr::{SummaryExpr, SummaryNode}; pub use guarantee::{ AccuracyError, BoundExpr, CompositionOperator, ErrorMetric, GuaranteeSource, ProbabilityExpr, ResultGuarantee, }; -pub use lifecycle::{EvaluationSchedule, OutputRepresentation, StateLifecycle}; pub use query_time::{ classic_cms_sizing, cms_posterior_error_bound, count_sketch_posterior_error_bound, cu_sketch_posterior_error_bound, traditional_a_priori_bound, @@ -52,3 +51,7 @@ pub use sketch::{ SketchParams, SketchQuery, StatModelKind, StatModelParams, WaveletKind, WaveletParams, }; pub use summary_maintenance::SummaryMaintenanceMode; +pub use summary_maintenance_lifecycle::{ + EvaluationSchedule, OutputRepresentation, SummaryMaintenanceLifecycle, + SummaryMaintenanceLifecycleGuarantee, +}; diff --git a/crates/types/src/post_asap/summary_maintenance_lifecycle.rs b/crates/types/src/post_asap/summary_maintenance_lifecycle.rs new file mode 100644 index 00000000..966b2cd5 --- /dev/null +++ b/crates/types/src/post_asap/summary_maintenance_lifecycle.rs @@ -0,0 +1,53 @@ +//! Physical summary-maintenance lifecycle vocabulary. +//! +//! These choices are attached by physical planning; a `SummaryAgg` does not +//! imply continuous maintenance by itself. "Summary maintenance lifecycle" +//! is deliberately narrower than the end-to-end data lifecycle (collection, +//! transmission, storage, and analytics). + +use super::SummaryMaintenanceMode; +use crate::workload::{DurationMs, TimestampMs}; + +/// When an operator is evaluated. This is independent of whether it owns +/// state and how long that state is retained. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] +pub enum EvaluationSchedule { + OneShot, + PerUpdate, + OnRead, +} + +/// The physical value crossing an execution boundary. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] +pub enum OutputRepresentation { + PlainRows, + SummaryState, + FinalizedValue, +} + +/// How long one planned summary state deployment exists. +#[derive(Debug, Clone, PartialEq, Eq, Hash)] +pub enum SummaryMaintenanceLifecycle { + Ephemeral, + Prepared { + activate_at: TimestampMs, + retire_at: TimestampMs, + }, + Shared { + retention: DurationMs, + }, + ContinuouslyMaintained, +} + +/// The lifecycle commitment emitted for one materialized summary deployment. +/// +/// This names the summary-maintenance promise explicitly so consumers do not +/// confuse it with guarantees about the broader data lifecycle. Accuracy is a +/// separate [`super::ResultGuarantee`]. +#[derive(Debug, Clone, PartialEq, Eq, Hash)] +pub struct SummaryMaintenanceLifecycleGuarantee { + pub summary_maintenance_lifecycle: SummaryMaintenanceLifecycle, + pub summary_maintenance_mode: SummaryMaintenanceMode, + pub evaluation_schedule: EvaluationSchedule, + pub output_representation: OutputRepresentation, +}