From 5baecb1057d5714f4956c9f0190fa12e666d2080 Mon Sep 17 00:00:00 2001 From: zz_y Date: Sat, 29 Aug 2026 07:22:17 -0600 Subject: [PATCH 01/11] feat(workload): plan summary state lifecycles --- crates/asap-aware-mapping/src/cost_model.rs | 9 + crates/asap-aware-mapping/src/lib.rs | 8 +- crates/asap-aware-mapping/src/lifecycle.rs | 835 ++++++++++++++++++ crates/types/src/post_asap/lifecycle.rs | 37 + crates/types/src/post_asap/mod.rs | 2 + .../workload-demand-and-summary-lifecycle.md | 32 +- 6 files changed, 913 insertions(+), 10 deletions(-) create mode 100644 crates/asap-aware-mapping/src/lifecycle.rs create mode 100644 crates/types/src/post_asap/lifecycle.rs diff --git a/crates/asap-aware-mapping/src/cost_model.rs b/crates/asap-aware-mapping/src/cost_model.rs index b1780c8d..4f9225f3 100644 --- a/crates/asap-aware-mapping/src/cost_model.rs +++ b/crates/asap-aware-mapping/src/cost_model.rs @@ -56,6 +56,7 @@ 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; use crate::recurrence::{ self, Horizon, RecurrenceCostExplanation, RecurrenceError, RecurrenceProfile, }; @@ -504,6 +505,14 @@ pub trait CostModel { let _ = (candidate, target); f64::NAN } + + /// Primitive build, update, read, retention, and retirement costs used to + /// 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 sketch_state( diff --git a/crates/asap-aware-mapping/src/lib.rs b/crates/asap-aware-mapping/src/lib.rs index 613d5061..b2afeb29 100644 --- a/crates/asap-aware-mapping/src/lib.rs +++ b/crates/asap-aware-mapping/src/lib.rs @@ -186,6 +186,7 @@ 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; @@ -195,7 +196,7 @@ pub mod topk_reuse; pub use accuracy::{ AccuracyAllocation, AccuracyBudgetAllocator, AccuracyEvidenceProvider, AccuracyModel, CompositionShape, DefaultAccuracyModel, EqualSplitAllocator, NoAccuracyEvidence, - PropagationStats, + PropagationStats, WorkloadAccuracyEvidence, }; pub use accuracy_reconciliation::AccuracyReconciliationStrategy; pub use cost_model::{CostModel, DefaultCostModel}; @@ -203,6 +204,11 @@ pub use explanation::{ explain_replacements, explain_replacements_with, ExplanationKind, ReplacementExplanation, }; pub use grouping::{has_subpopulations, HydraGroupingStrategy}; +pub use lifecycle::{ + materialize_with_lifecycles, plan_summary_lifecycles, LifecycleAlternative, + LifecycleCapabilities, LifecycleCostInputs, LifecyclePlan, LifecyclePlanError, + LifecycleRejection, MaterializeLifecycleError, StateDeployment, +}; pub use recurrence::{ evaluation_rate_of, total_cost, update_rate_from_data_workload, CostRate, EvaluationRate, Horizon, RecurrenceCostExplanation, RecurrenceError, RecurrenceProfile, RootRecurrence, diff --git a/crates/asap-aware-mapping/src/lifecycle.rs b/crates/asap-aware-mapping/src/lifecycle.rs new file mode 100644 index 00000000..12bed5c7 --- /dev/null +++ b/crates/asap-aware-mapping/src/lifecycle.rs @@ -0,0 +1,835 @@ +//! Workload-aware physical lifecycle planning for summary state. +//! +//! Phase validation from PR #300 answers whether a post-ASAP DAG can execute. +//! 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::{validate_execution_phases, StateLifecycle, SummaryExpr, SummaryNode}; +use asap_types::post_asap::{EvaluationSchedule, OutputRepresentation}; +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, UpdateRate}; +use crate::replacement::{GlobalSelection, ImplementError}; + +/// 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, +} + +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, + MissingCostEvidence, +} + +#[derive(Debug, Clone, PartialEq)] +pub struct LifecycleAlternative { + pub lifecycle: StateLifecycle, + 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 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, +} + +#[derive(Debug, thiserror::Error)] +pub enum LifecyclePlanError { + #[error(transparent)] + InvalidWorkload(#[from] WorkloadError), + #[error(transparent)] + InvalidExecutionPhases(#[from] asap_types::post_asap::PhaseError), + #[error("optimization horizon must be finite and strictly positive")] + InvalidHorizon, +} + +#[derive(Debug, thiserror::Error)] +pub enum MaterializeLifecycleError { + #[error(transparent)] + Materialize(#[from] ImplementError), + #[error(transparent)] + Lifecycle(#[from] LifecyclePlanError), +} + +#[derive(Debug)] +struct WorkloadFacts { + reads: Option, + one_time_invocations: u64, + evaluation_rate: Option, + update_rate: Option, + arrival: DataArrival, + prepared_window: Option<(TimestampMs, TimestampMs)>, +} + +/// 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, + workload: &QueryWorkload, + now_ms: u64, + horizon: Option, + capabilities: LifecycleCapabilities, + cost_model: &dyn CostModel, +) -> Result { + workload.validate()?; + validate_execution_phases(&root)?; + if horizon.is_some_and(|h| !h.0.is_finite() || h.0 <= 0.0) { + return Err(LifecyclePlanError::InvalidHorizon); + } + let facts = workload_facts(workload, now_ms, horizon); + let mut summaries = Vec::new(); + collect_summary_aggs(&root, &mut HashSet::new(), &mut summaries); + let deployments = summaries + .into_iter() + .enumerate() + .map(|(summary_index, summary)| { + let alternatives = alternatives_for( + &facts, + horizon, + capabilities, + 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()); + 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, + evaluation_schedule, + output_representation: OutputRepresentation::SummaryState, + alternatives, + } + }) + .collect(); + Ok(LifecyclePlan { + root, + deployments, + horizon, + evaluation_rate: facts.evaluation_rate, + update_rate: facts.update_rate, + }) +} + +/// Materialize PR #300's globally selected phase-valid DAG and immediately +/// attach workload-aware lifecycle deployments. +pub fn materialize_with_lifecycles( + selection: &GlobalSelection<'_>, + target: &Rc, + workload: &QueryWorkload, + now_ms: u64, + horizon: Option, + capabilities: LifecycleCapabilities, + cost_model: &dyn CostModel, +) -> Result, MaterializeLifecycleError> { + selection + .materialize(target)? + .map(|root| { + plan_summary_lifecycles(root, workload, now_ms, horizon, capabilities, cost_model) + }) + .transpose() + .map_err(Into::into) +} + +fn workload_facts( + workload: &QueryWorkload, + now_ms: u64, + horizon: Option, +) -> WorkloadFacts { + 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; + + for entry in workload.entries() { + match &entry.recurrence { + QueryRecurrence::OneTime { + invocations, + execute_at, + } => { + one_time_invocations = one_time_invocations.saturating_add(*invocations); + 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))); + } + } + } + QueryRecurrence::Repeated(RepeatedDemand::FixedInterval(interval)) => { + 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)) => { + if let Some(h) = horizon { + let end_ms = now_ms.saturating_add((h.0 * 1000.0) as u64); + recurring_reads += schedule + .iter() + .filter(|at| at.0 >= now_ms && at.0 <= end_ms) + .count() as f64; + evaluation_rate += schedule.len() as f64 / h.0; + has_evaluation_rate = true; + } else { + recurring_known = false; + } + } + QueryRecurrence::Repeated(RepeatedDemand::EstimatedRate(estimate)) => { + let fresh = match (estimate.observed_at, estimate.valid_for) { + (Some(observed), Some(valid_for)) => { + now_ms <= observed.0.saturating_add(valid_for.0) + } + (None, Some(_)) => false, + _ => true, + }; + if !fresh { + 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 => 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); + 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), + } +} + +fn alternatives_for( + facts: &WorkloadFacts, + horizon: Option, + capabilities: LifecycleCapabilities, + costs: LifecycleCostInputs, +) -> Vec { + let mut alternatives = Vec::with_capacity(4); + alternatives.push(ephemeral(facts, capabilities, &costs)); + alternatives.push(prepared(facts, capabilities, &costs)); + alternatives.push(shared(facts, horizon, capabilities, &costs)); + alternatives.push(continuous(facts, horizon, capabilities, &costs)); + alternatives +} + +fn ephemeral( + facts: &WorkloadFacts, + capabilities: LifecycleCapabilities, + costs: &LifecycleCostInputs, +) -> LifecycleAlternative { + let lifecycle = StateLifecycle::Ephemeral; + if !capabilities.ephemeral { + return rejected(lifecycle, 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, + total_cost, + vec!["state is rebuilt per invocation".into()], + ) +} + +fn prepared( + facts: &WorkloadFacts, + capabilities: LifecycleCapabilities, + costs: &LifecycleCostInputs, +) -> LifecycleAlternative { + let Some((activate_at, retire_at)) = facts.prepared_window else { + return rejected( + StateLifecycle::Prepared { + activate_at: TimestampMs(0), + retire_at: TimestampMs(0), + }, + LifecycleRejection::RequiresPredictableOneTimeQuery, + ); + }; + let lifecycle = StateLifecycle::Prepared { + activate_at, + retire_at, + }; + if !capabilities.prepared { + return rejected(lifecycle, LifecycleRejection::UnsupportedByRuntime); + } + 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: LifecycleCapabilities, + 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, LifecycleRejection::UnsupportedByRuntime); + } + if facts.reads.is_none_or(|reads| reads <= 1.0) { + return rejected(lifecycle, LifecycleRejection::RequiresMultipleReads); + } + let Some(horizon) = horizon else { + return rejected(lifecycle, LifecycleRejection::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: LifecycleCapabilities, + costs: &LifecycleCostInputs, +) -> LifecycleAlternative { + let lifecycle = StateLifecycle::ContinuouslyMaintained; + if !capabilities.continuously_maintained { + return rejected(lifecycle, LifecycleRejection::UnsupportedByRuntime); + } + if !matches!( + facts.arrival, + DataArrival::ContinuouslyIngesting | DataArrival::Mixed + ) { + return rejected(lifecycle, LifecycleRejection::RequiresContinuousData); + } + if facts.update_rate.is_none() { + return rejected(lifecycle, LifecycleRejection::MissingOrStaleIngestionRate); + } + let Some(horizon) = horizon else { + return rejected(lifecycle, LifecycleRejection::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 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, + total_cost: Option, + assumptions: Vec, +) -> LifecycleAlternative { + LifecycleAlternative { + lifecycle, + total_cost, + rejection: total_cost + .is_none() + .then_some(LifecycleRejection::MissingCostEvidence), + assumptions, + } +} + +fn rejected(lifecycle: StateLifecycle, rejection: LifecycleRejection) -> LifecycleAlternative { + LifecycleAlternative { + 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::ExactTransform { child, .. } | SummaryExpr::ExactPostProcess { child, .. } => { + 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 query_root() -> Rc { + Rc::new(QueryExpr::Scan { + source: Source::TimeSeries { metric: "m".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(), + &workload(vec![batch(Predictability::AdHoc)], vec![], at_rest()), + 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].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(), + &workload(vec![entry], vec![], at_rest()), + 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(), + &workload(vec![], vec![repeating()], at_rest()), + 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].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(), + &workload(vec![], vec![repeating()], continuous(1_000, 60_000)), + 1_000, + Some(Horizon(10.0)), + capabilities, + &UnitCosts, + ) + .unwrap(); + assert_eq!( + plan.deployments[0].selected, + Some(StateLifecycle::ContinuouslyMaintained) + ); + 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(), + &workload(vec![], vec![repeating()], continuous(1_000, 1_000)), + 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(), + &workload(vec![], vec![repeating()], continuous(1_000, 60_000)), + 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 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, 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); + } +} diff --git a/crates/types/src/post_asap/lifecycle.rs b/crates/types/src/post_asap/lifecycle.rs new file mode 100644 index 00000000..995c2f79 --- /dev/null +++ b/crates/types/src/post_asap/lifecycle.rs @@ -0,0 +1,37 @@ +//! 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 40c4a7a4..5c572b05 100644 --- a/crates/types/src/post_asap/mod.rs +++ b/crates/types/src/post_asap/mod.rs @@ -29,6 +29,7 @@ pub mod expr; pub mod guarantee; +pub mod lifecycle; pub mod query_time; pub mod schema; pub mod sketch; @@ -39,6 +40,7 @@ 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, diff --git a/docs/design_docs/asap-aware-mapping/workload-demand-and-summary-lifecycle.md b/docs/design_docs/asap-aware-mapping/workload-demand-and-summary-lifecycle.md index fef2aff1..5d98890d 100644 --- a/docs/design_docs/asap-aware-mapping/workload-demand-and-summary-lifecycle.md +++ b/docs/design_docs/asap-aware-mapping/workload-demand-and-summary-lifecycle.md @@ -5,8 +5,9 @@ This document is for ASAPPlanner designers, architects, researchers, and developers working on workload-aware plan selection. It defines how the planner should describe query workload, data workload, and the lifecycle of -summary state. It is a design contract, not a description of the current -public Rust API. +summary state. It is the design contract for the public Rust model and the +workload-to-lifecycle planning API; deployments still supply their own cost +statistics and runtime capabilities. The terminology follows the ProjectASAP [glossary](https://github.com/ProjectASAP/internal-docs/blob/03e1c70f5af3ae9221471898541067eee7f86338/glossary.md). @@ -41,13 +42,26 @@ may arrive unexpectedly during exploration, run once at a scheduled time, or repeat every ten seconds on a dashboard. Planning summary state from syntax alone either misses reuse or invents reuse that the workload does not justify. -The current normalized workload distinguishes a one-shot `query_batch` from -fixed-interval `repeating_queries`, and the recurrence cost model distinguishes -one-shot consumers from evaluation and update rates. This is a useful base, but -it does not represent predictability, uncertain demand, real-time versus -longitudinal scope, at-rest versus continuously ingesting data, or summary-state -lifecycle. It also risks treating "repeating query" and "streaming data" as the -same fact even though the glossary defines them on different axes. +The normalized workload preserves `query_batch` and `repeating_queries` as +compatibility-shaped inputs, then exposes both through `QueryWorkload::entries` +as recurrence, predictability, requirements, and time-selection axes. Data +arrival and fresh ingestion evidence remain a separate `DataWorkload`; a +repeating query therefore never implies streaming data. + +### Implementation map + +- `asap_types::workload` defines the normalized query/data workload and + evidence freshness contract. +- `PlanSpace::recurrence_profiles_from_workload` derives per-target read and + update recurrence without treating missing evidence as zero. +- `WorkloadAccuracyEvidence` supplies fresh cardinality and distribution to + accuracy models. +- `plan_summary_lifecycles` enumerates legal ephemeral, prepared, shared, and + continuously maintained alternatives and compares their costs over the + caller's explicit horizon. +- `materialize_with_lifecycles` attaches those state deployments to PR #300's + phase-validated global selection. Each deployment retains assumptions and + rejected alternatives for explanation. ## Inputs, outputs, and end-to-end behavior From 842ebb11bd0ea77411cc34da793f643481a04157 Mon Sep 17 00:00:00 2001 From: zz_y Date: Sat, 29 Aug 2026 16:43:19 -0600 Subject: [PATCH 02/11] refactor(lifecycle): use generic phase operations --- crates/asap-aware-mapping/src/lifecycle.rs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/crates/asap-aware-mapping/src/lifecycle.rs b/crates/asap-aware-mapping/src/lifecycle.rs index 12bed5c7..1acfc931 100644 --- a/crates/asap-aware-mapping/src/lifecycle.rs +++ b/crates/asap-aware-mapping/src/lifecycle.rs @@ -551,7 +551,8 @@ fn collect_summary_aggs( collect_summary_aggs(child, seen, output); } } - SummaryExpr::ExactTransform { child, .. } | SummaryExpr::ExactPostProcess { child, .. } => { + SummaryExpr::UpdateTransform { child, .. } + | SummaryExpr::ReadoutPostProcess { child, .. } => { collect_summary_aggs(child, seen, output) } SummaryExpr::KeepPreAsap(_) => {} From 7b8900abddaac7483cdfd87cca877f01a9635b31 Mon Sep 17 00:00:00 2001 From: zz_y Date: Sat, 29 Aug 2026 07:25:20 -0600 Subject: [PATCH 03/11] fix(workload): satisfy lifecycle lint --- crates/asap-aware-mapping/src/lifecycle.rs | 11 ++++++----- 1 file changed, 6 insertions(+), 5 deletions(-) diff --git a/crates/asap-aware-mapping/src/lifecycle.rs b/crates/asap-aware-mapping/src/lifecycle.rs index 1acfc931..1508d179 100644 --- a/crates/asap-aware-mapping/src/lifecycle.rs +++ b/crates/asap-aware-mapping/src/lifecycle.rs @@ -332,11 +332,12 @@ fn alternatives_for( capabilities: LifecycleCapabilities, costs: LifecycleCostInputs, ) -> Vec { - let mut alternatives = Vec::with_capacity(4); - alternatives.push(ephemeral(facts, capabilities, &costs)); - alternatives.push(prepared(facts, capabilities, &costs)); - alternatives.push(shared(facts, horizon, capabilities, &costs)); - alternatives.push(continuous(facts, horizon, capabilities, &costs)); + let alternatives = vec![ + ephemeral(facts, capabilities, &costs), + prepared(facts, capabilities, &costs), + shared(facts, horizon, capabilities, &costs), + continuous(facts, horizon, capabilities, &costs), + ]; alternatives } From 63c013721d42e70b5155efcd5e08fa830b5b7ae1 Mon Sep 17 00:00:00 2001 From: zz_y Date: Sat, 29 Aug 2026 16:45:49 -0600 Subject: [PATCH 04/11] chore(lifecycle): defer design documentation to docs PR --- .../workload-demand-and-summary-lifecycle.md | 32 ++++++------------- 1 file changed, 9 insertions(+), 23 deletions(-) diff --git a/docs/design_docs/asap-aware-mapping/workload-demand-and-summary-lifecycle.md b/docs/design_docs/asap-aware-mapping/workload-demand-and-summary-lifecycle.md index 5d98890d..fef2aff1 100644 --- a/docs/design_docs/asap-aware-mapping/workload-demand-and-summary-lifecycle.md +++ b/docs/design_docs/asap-aware-mapping/workload-demand-and-summary-lifecycle.md @@ -5,9 +5,8 @@ This document is for ASAPPlanner designers, architects, researchers, and developers working on workload-aware plan selection. It defines how the planner should describe query workload, data workload, and the lifecycle of -summary state. It is the design contract for the public Rust model and the -workload-to-lifecycle planning API; deployments still supply their own cost -statistics and runtime capabilities. +summary state. It is a design contract, not a description of the current +public Rust API. The terminology follows the ProjectASAP [glossary](https://github.com/ProjectASAP/internal-docs/blob/03e1c70f5af3ae9221471898541067eee7f86338/glossary.md). @@ -42,26 +41,13 @@ may arrive unexpectedly during exploration, run once at a scheduled time, or repeat every ten seconds on a dashboard. Planning summary state from syntax alone either misses reuse or invents reuse that the workload does not justify. -The normalized workload preserves `query_batch` and `repeating_queries` as -compatibility-shaped inputs, then exposes both through `QueryWorkload::entries` -as recurrence, predictability, requirements, and time-selection axes. Data -arrival and fresh ingestion evidence remain a separate `DataWorkload`; a -repeating query therefore never implies streaming data. - -### Implementation map - -- `asap_types::workload` defines the normalized query/data workload and - evidence freshness contract. -- `PlanSpace::recurrence_profiles_from_workload` derives per-target read and - update recurrence without treating missing evidence as zero. -- `WorkloadAccuracyEvidence` supplies fresh cardinality and distribution to - accuracy models. -- `plan_summary_lifecycles` enumerates legal ephemeral, prepared, shared, and - continuously maintained alternatives and compares their costs over the - caller's explicit horizon. -- `materialize_with_lifecycles` attaches those state deployments to PR #300's - phase-validated global selection. Each deployment retains assumptions and - rejected alternatives for explanation. +The current normalized workload distinguishes a one-shot `query_batch` from +fixed-interval `repeating_queries`, and the recurrence cost model distinguishes +one-shot consumers from evaluation and update rates. This is a useful base, but +it does not represent predictability, uncertain demand, real-time versus +longitudinal scope, at-rest versus continuously ingesting data, or summary-state +lifecycle. It also risks treating "repeating query" and "streaming data" as the +same fact even though the glossary defines them on different axes. ## Inputs, outputs, and end-to-end behavior From 7a2f2323d54e6a670dbdc70b7743fd837dce45e3 Mon Sep 17 00:00:00 2001 From: zz_y Date: Sun, 30 Aug 2026 15:13:15 -0600 Subject: [PATCH 05/11] refactor(lifecycle): attach summary maintenance modes --- crates/asap-aware-mapping/src/lib.rs | 5 +- crates/asap-aware-mapping/src/lifecycle.rs | 143 ++++++++++++------- crates/asap-aware-mapping/src/replacement.rs | 3 +- 3 files changed, 92 insertions(+), 59 deletions(-) diff --git a/crates/asap-aware-mapping/src/lib.rs b/crates/asap-aware-mapping/src/lib.rs index b2afeb29..773464b9 100644 --- a/crates/asap-aware-mapping/src/lib.rs +++ b/crates/asap-aware-mapping/src/lib.rs @@ -205,9 +205,8 @@ pub use explanation::{ }; pub use grouping::{has_subpopulations, HydraGroupingStrategy}; pub use lifecycle::{ - materialize_with_lifecycles, plan_summary_lifecycles, LifecycleAlternative, - LifecycleCapabilities, LifecycleCostInputs, LifecyclePlan, LifecyclePlanError, - LifecycleRejection, MaterializeLifecycleError, StateDeployment, + plan_summary_lifecycles, LifecycleAlternative, LifecycleCapabilities, LifecycleCostInputs, + LifecyclePlan, LifecyclePlanError, LifecycleRejection, StateDeployment, }; pub use recurrence::{ evaluation_rate_of, total_cost, update_rate_from_data_workload, CostRate, EvaluationRate, diff --git a/crates/asap-aware-mapping/src/lifecycle.rs b/crates/asap-aware-mapping/src/lifecycle.rs index 1508d179..08b1ab11 100644 --- a/crates/asap-aware-mapping/src/lifecycle.rs +++ b/crates/asap-aware-mapping/src/lifecycle.rs @@ -1,6 +1,5 @@ //! Workload-aware physical lifecycle planning for summary state. //! -//! Phase validation from PR #300 answers whether a post-ASAP DAG can execute. //! 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. @@ -8,9 +7,8 @@ use std::collections::HashSet; use std::rc::Rc; -use asap_types::post_asap::{validate_execution_phases, StateLifecycle, SummaryExpr, SummaryNode}; -use asap_types::post_asap::{EvaluationSchedule, OutputRepresentation}; -use asap_types::pre_asap::QueryExpr; +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, @@ -18,7 +16,6 @@ use asap_types::workload::{ use crate::cost_model::{Cost, CostModel}; use crate::recurrence::{CostRate, EvaluationRate, Horizon, UpdateRate}; -use crate::replacement::{GlobalSelection, ImplementError}; /// Runtime lifecycle shapes available to the planner. #[derive(Debug, Clone, Copy, PartialEq, Eq)] @@ -69,6 +66,8 @@ pub enum LifecycleRejection { #[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, @@ -86,6 +85,7 @@ 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, @@ -104,20 +104,10 @@ pub struct LifecyclePlan { pub enum LifecyclePlanError { #[error(transparent)] InvalidWorkload(#[from] WorkloadError), - #[error(transparent)] - InvalidExecutionPhases(#[from] asap_types::post_asap::PhaseError), #[error("optimization horizon must be finite and strictly positive")] InvalidHorizon, } -#[derive(Debug, thiserror::Error)] -pub enum MaterializeLifecycleError { - #[error(transparent)] - Materialize(#[from] ImplementError), - #[error(transparent)] - Lifecycle(#[from] LifecyclePlanError), -} - #[derive(Debug)] struct WorkloadFacts { reads: Option, @@ -140,7 +130,6 @@ pub fn plan_summary_lifecycles( cost_model: &dyn CostModel, ) -> Result { workload.validate()?; - validate_execution_phases(&root)?; if horizon.is_some_and(|h| !h.0.is_finite() || h.0 <= 0.0) { return Err(LifecyclePlanError::InvalidHorizon); } @@ -161,8 +150,8 @@ pub fn plan_summary_lifecycles( .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()); - let evaluation_schedule = selected.as_ref().map(|lifecycle| match lifecycle { + .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!( @@ -179,7 +168,8 @@ pub fn plan_summary_lifecycles( StateDeployment { summary_index, summary, - selected, + selected: selected.as_ref().map(|(lifecycle, _)| lifecycle.clone()), + selected_maintenance_mode: selected.map(|(_, mode)| mode), evaluation_schedule, output_representation: OutputRepresentation::SummaryState, alternatives, @@ -195,26 +185,6 @@ pub fn plan_summary_lifecycles( }) } -/// Materialize PR #300's globally selected phase-valid DAG and immediately -/// attach workload-aware lifecycle deployments. -pub fn materialize_with_lifecycles( - selection: &GlobalSelection<'_>, - target: &Rc, - workload: &QueryWorkload, - now_ms: u64, - horizon: Option, - capabilities: LifecycleCapabilities, - cost_model: &dyn CostModel, -) -> Result, MaterializeLifecycleError> { - selection - .materialize(target)? - .map(|root| { - plan_summary_lifecycles(root, workload, now_ms, horizon, capabilities, cost_model) - }) - .transpose() - .map_err(Into::into) -} - fn workload_facts( workload: &QueryWorkload, now_ms: u64, @@ -348,7 +318,11 @@ fn ephemeral( ) -> LifecycleAlternative { let lifecycle = StateLifecycle::Ephemeral; if !capabilities.ephemeral { - return rejected(lifecycle, LifecycleRejection::UnsupportedByRuntime); + return rejected( + lifecycle, + SummaryMaintenanceMode::DirectBuild, + LifecycleRejection::UnsupportedByRuntime, + ); } let total_cost = zip_costs(&[ costs.build_cost, @@ -359,6 +333,7 @@ fn ephemeral( .map(|(per_read, reads)| Cost(per_read * reads)); costed_or_unknown( lifecycle, + SummaryMaintenanceMode::DirectBuild, total_cost, vec!["state is rebuilt per invocation".into()], ) @@ -375,6 +350,7 @@ fn prepared( activate_at: TimestampMs(0), retire_at: TimestampMs(0), }, + retained_mode(facts), LifecycleRejection::RequiresPredictableOneTimeQuery, ); }; @@ -383,7 +359,11 @@ fn prepared( retire_at, }; if !capabilities.prepared { - return rejected(lifecycle, LifecycleRejection::UnsupportedByRuntime); + return rejected( + lifecycle, + retained_mode(facts), + LifecycleRejection::UnsupportedByRuntime, + ); } let seconds = retire_at.0.saturating_sub(activate_at.0) as f64 / 1000.0; let maintenance = maintenance_cost(facts, costs, seconds); @@ -405,6 +385,7 @@ fn prepared( }; costed_or_unknown( lifecycle, + retained_mode(facts), total_cost, vec!["activation and retirement come from the declared schedule".into()], ) @@ -420,17 +401,30 @@ fn shared( retention: asap_types::workload::DurationMs(horizon.map_or(0, |h| (h.0 * 1000.0) as u64)), }; if !capabilities.shared { - return rejected(lifecycle, LifecycleRejection::UnsupportedByRuntime); + return rejected( + lifecycle, + retained_mode(facts), + LifecycleRejection::UnsupportedByRuntime, + ); } if facts.reads.is_none_or(|reads| reads <= 1.0) { - return rejected(lifecycle, LifecycleRejection::RequiresMultipleReads); + return rejected( + lifecycle, + retained_mode(facts), + LifecycleRejection::RequiresMultipleReads, + ); } let Some(horizon) = horizon else { - return rejected(lifecycle, LifecycleRejection::RequiresHorizon); + 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()], ) @@ -444,23 +438,40 @@ fn continuous( ) -> LifecycleAlternative { let lifecycle = StateLifecycle::ContinuouslyMaintained; if !capabilities.continuously_maintained { - return rejected(lifecycle, LifecycleRejection::UnsupportedByRuntime); + return rejected( + lifecycle, + SummaryMaintenanceMode::Incremental, + LifecycleRejection::UnsupportedByRuntime, + ); } if !matches!( facts.arrival, DataArrival::ContinuouslyIngesting | DataArrival::Mixed ) { - return rejected(lifecycle, LifecycleRejection::RequiresContinuousData); + return rejected( + lifecycle, + SummaryMaintenanceMode::Incremental, + LifecycleRejection::RequiresContinuousData, + ); } if facts.update_rate.is_none() { - return rejected(lifecycle, LifecycleRejection::MissingOrStaleIngestionRate); + return rejected( + lifecycle, + SummaryMaintenanceMode::Incremental, + LifecycleRejection::MissingOrStaleIngestionRate, + ); } let Some(horizon) = horizon else { - return rejected(lifecycle, LifecycleRejection::RequiresHorizon); + 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()], ) @@ -500,11 +511,13 @@ fn zip_costs(costs: &[Option]) -> Option { 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() @@ -513,15 +526,29 @@ fn costed_or_unknown( } } -fn rejected(lifecycle: StateLifecycle, rejection: LifecycleRejection) -> LifecycleAlternative { +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>, @@ -552,10 +579,6 @@ fn collect_summary_aggs( collect_summary_aggs(child, seen, output); } } - SummaryExpr::UpdateTransform { child, .. } - | SummaryExpr::ReadoutPostProcess { child, .. } => { - collect_summary_aggs(child, seen, output) - } SummaryExpr::KeepPreAsap(_) => {} } } @@ -710,6 +733,10 @@ mod tests { 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)) @@ -753,6 +780,10 @@ mod tests { 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) @@ -779,6 +810,10 @@ mod tests { 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))); } @@ -825,7 +860,7 @@ mod tests { 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, 1_000, Some(Horizon(10.0))) + .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. diff --git a/crates/asap-aware-mapping/src/replacement.rs b/crates/asap-aware-mapping/src/replacement.rs index 1824a5ac..fdbfdf7a 100644 --- a/crates/asap-aware-mapping/src/replacement.rs +++ b/crates/asap-aware-mapping/src/replacement.rs @@ -375,8 +375,7 @@ use crate::accuracy_reconciliation::AccuracyReconciliationStrategy; use crate::cost_model::{CostModel, CseCandidate, DefaultCostModel, ShareDecision}; use crate::grouping::HydraGroupingStrategy; use crate::recurrence::{ - evaluation_rate_of, CostRate, Horizon, RecurrenceError, RecurrenceProfile, RootRecurrence, - UpdateRate, + evaluation_rate_of, Horizon, RecurrenceError, RecurrenceProfile, RootRecurrence, UpdateRate, }; use crate::rollup::RollupStrategy; use crate::topk_reuse::TopKLimitReuseStrategy; From 1c7a034740aae4bd78726a9ff7bf2723352b0d4f Mon Sep 17 00:00:00 2001 From: zz_y Date: Mon, 31 Aug 2026 06:08:29 -0600 Subject: [PATCH 06/11] refactor(lifecycle): name summary maintenance lifecycle explicitly --- crates/asap-aware-mapping/src/cost_model.rs | 9 +- crates/asap-aware-mapping/src/lib.rs | 12 +- ...le.rs => summary_maintenance_lifecycle.rs} | 231 ++++++++++-------- crates/types/src/post_asap/mod.rs | 6 +- ...le.rs => summary_maintenance_lifecycle.rs} | 13 +- 5 files changed, 155 insertions(+), 116 deletions(-) rename crates/asap-aware-mapping/src/{lifecycle.rs => summary_maintenance_lifecycle.rs} (76%) rename crates/types/src/post_asap/{lifecycle.rs => summary_maintenance_lifecycle.rs} (57%) diff --git a/crates/asap-aware-mapping/src/cost_model.rs b/crates/asap-aware-mapping/src/cost_model.rs index 4f9225f3..254d6964 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; use crate::recurrence::{ self, Horizon, RecurrenceCostExplanation, RecurrenceError, RecurrenceProfile, }; @@ -64,6 +63,7 @@ use crate::replacement::{ realize_child, Implementation, Replacement, ReplacementProvenance, ReplacementSubDAG, TargetSubDAG, }; +use crate::summary_maintenance_lifecycle::SummaryMaintenanceLifecycleCostInputs; /// A CSE-detected, legality-gated shared subtree with two or more consumers /// — the unit [`CostModel::cse_share_decision`] decides over. Built by @@ -510,8 +510,11 @@ 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() } } diff --git a/crates/asap-aware-mapping/src/lib.rs b/crates/asap-aware-mapping/src/lib.rs index 773464b9..5fd54655 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,10 +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, -}; pub use recurrence::{ evaluation_rate_of, total_cost, update_rate_from_data_workload, CostRate, EvaluationRate, Horizon, RecurrenceCostExplanation, RecurrenceError, RecurrenceProfile, RootRecurrence, @@ -222,4 +218,10 @@ pub use replacement::{ MAX_SEARCH_ITERATIONS, }; pub use rewrite::AvgToSumOverCountStrategy; +pub use summary_maintenance_lifecycle::{ + plan_summary_maintenance_lifecycles, SummaryMaintenanceDeployment, + SummaryMaintenanceLifecycleAlternative, SummaryMaintenanceLifecycleCapabilities, + SummaryMaintenanceLifecycleCostInputs, SummaryMaintenanceLifecyclePlan, + SummaryMaintenanceLifecyclePlanError, SummaryMaintenanceLifecycleRejection, +}; pub use topk_reuse::TopKLimitReuseStrategy; diff --git a/crates/asap-aware-mapping/src/lifecycle.rs b/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs similarity index 76% rename from crates/asap-aware-mapping/src/lifecycle.rs rename to crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs index 08b1ab11..503d6957 100644 --- a/crates/asap-aware-mapping/src/lifecycle.rs +++ b/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs @@ -1,14 +1,23 @@ -//! Workload-aware physical lifecycle planning for summary state. +//! 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 lifecycle win. +//! A **summary-maintenance lifecycle** is the physical policy for when one +//! materialized summary state is created, retained or shared, updated as data +//! arrives, and retired. It is deliberately narrower than the end-to-end data +//! lifecycle and independent of query recurrence: recurrence is workload +//! evidence used to choose a lifecycle, not a lifecycle itself. +//! +//! This module enumerates and costs `Ephemeral`, `Prepared`, `Shared`, and +//! `ContinuouslyMaintained` alternatives for every unique `SummaryAgg` in a +//! materialized plan. [`SummaryMaintenanceMode`] is an orthogonal detail of +//! the selected deployment: state is either built directly or updated +//! incrementally. Unknown evidence stays unknown and therefore cannot make a +//! long-lived alternative 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::post_asap::{SummaryExpr, SummaryMaintenanceLifecycle, SummaryNode}; use asap_types::workload::{ DataArrival, Predictability, QueryRecurrence, QueryWorkload, RepeatedDemand, TimestampMs, WorkloadError, @@ -17,16 +26,16 @@ use asap_types::workload::{ use crate::cost_model::{Cost, CostModel}; use crate::recurrence::{CostRate, EvaluationRate, Horizon, UpdateRate}; -/// Runtime lifecycle shapes available to the planner. +/// Summary-maintenance lifecycle shapes supported by the target runtime. #[derive(Debug, Clone, Copy, PartialEq, Eq)] -pub struct LifecycleCapabilities { +pub struct SummaryMaintenanceLifecycleCapabilities { pub ephemeral: bool, pub prepared: bool, pub shared: bool, pub continuously_maintained: bool, } -impl LifecycleCapabilities { +impl SummaryMaintenanceLifecycleCapabilities { pub const ALL: Self = Self { ephemeral: true, prepared: true, @@ -35,7 +44,7 @@ impl LifecycleCapabilities { }; } -impl Default for LifecycleCapabilities { +impl Default for SummaryMaintenanceLifecycleCapabilities { fn default() -> Self { Self::ALL } @@ -44,7 +53,7 @@ impl Default for LifecycleCapabilities { /// 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 struct SummaryMaintenanceLifecycleCostInputs { pub build_cost: Option, pub maintenance_cost_per_update: Option, pub summary_read_cost: Option, @@ -53,7 +62,7 @@ pub struct LifecycleCostInputs { } #[derive(Debug, Clone, PartialEq, Eq)] -pub enum LifecycleRejection { +pub enum SummaryMaintenanceLifecycleRejection { UnsupportedByRuntime, RequiresPredictableOneTimeQuery, RequiresMultipleReads, @@ -64,16 +73,16 @@ pub enum LifecycleRejection { } #[derive(Debug, Clone, PartialEq)] -pub struct LifecycleAlternative { - pub lifecycle: StateLifecycle, +pub struct SummaryMaintenanceLifecycleAlternative { + pub summary_maintenance_lifecycle: SummaryMaintenanceLifecycle, /// How this lifecycle obtains and refreshes its summary state. pub maintenance_mode: SummaryMaintenanceMode, pub total_cost: Option, - pub rejection: Option, + pub rejection: Option, pub assumptions: Vec, } -impl LifecycleAlternative { +impl SummaryMaintenanceLifecycleAlternative { fn selectable(&self) -> bool { self.rejection.is_none() && self.total_cost.is_some() } @@ -81,27 +90,27 @@ impl LifecycleAlternative { /// One unique summary-state deployment. Shared `Rc` nodes are emitted once. #[derive(Debug, Clone)] -pub struct StateDeployment { +pub struct SummaryMaintenanceDeployment { pub summary_index: usize, pub summary: Rc, - pub selected: Option, + pub selected_summary_maintenance_lifecycle: Option, pub selected_maintenance_mode: Option, pub evaluation_schedule: Option, pub output_representation: OutputRepresentation, - pub alternatives: Vec, + pub alternatives: Vec, } #[derive(Debug, Clone)] -pub struct LifecyclePlan { +pub struct SummaryMaintenanceLifecyclePlan { pub root: Rc, - pub deployments: Vec, + pub deployments: Vec, pub horizon: Option, pub evaluation_rate: Option, pub update_rate: Option, } #[derive(Debug, thiserror::Error)] -pub enum LifecyclePlanError { +pub enum SummaryMaintenanceLifecyclePlanError { #[error(transparent)] InvalidWorkload(#[from] WorkloadError), #[error("optimization horizon must be finite and strictly positive")] @@ -121,17 +130,17 @@ struct WorkloadFacts { /// 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( +pub fn plan_summary_maintenance_lifecycles( root: Rc, workload: &QueryWorkload, now_ms: u64, horizon: Option, - capabilities: LifecycleCapabilities, + capabilities: SummaryMaintenanceLifecycleCapabilities, cost_model: &dyn CostModel, -) -> Result { +) -> Result { workload.validate()?; if horizon.is_some_and(|h| !h.0.is_finite() || h.0 <= 0.0) { - return Err(LifecyclePlanError::InvalidHorizon); + return Err(SummaryMaintenanceLifecyclePlanError::InvalidHorizon); } let facts = workload_facts(workload, now_ms, horizon); let mut summaries = Vec::new(); @@ -144,16 +153,22 @@ pub fn plan_summary_lifecycles( &facts, horizon, capabilities, - cost_model.summary_lifecycle_cost_inputs(&summary), + cost_model.summary_maintenance_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)); + .map(|candidate| { + ( + candidate.summary_maintenance_lifecycle.clone(), + candidate.maintenance_mode, + ) + }); let evaluation_schedule = selected.as_ref().map(|(lifecycle, _)| match lifecycle { - StateLifecycle::Ephemeral => EvaluationSchedule::OneShot, - StateLifecycle::Prepared { .. } | StateLifecycle::Shared { .. } + SummaryMaintenanceLifecycle::Ephemeral => EvaluationSchedule::OneShot, + SummaryMaintenanceLifecycle::Prepared { .. } + | SummaryMaintenanceLifecycle::Shared { .. } if matches!( facts.arrival, DataArrival::ContinuouslyIngesting | DataArrival::Mixed @@ -161,14 +176,18 @@ pub fn plan_summary_lifecycles( { EvaluationSchedule::PerUpdate } - StateLifecycle::Prepared { .. } => EvaluationSchedule::OneShot, - StateLifecycle::Shared { .. } => EvaluationSchedule::OnRead, - StateLifecycle::ContinuouslyMaintained => EvaluationSchedule::PerUpdate, + SummaryMaintenanceLifecycle::Prepared { .. } => EvaluationSchedule::OneShot, + SummaryMaintenanceLifecycle::Shared { .. } => EvaluationSchedule::OnRead, + SummaryMaintenanceLifecycle::ContinuouslyMaintained => { + EvaluationSchedule::PerUpdate + } }); - StateDeployment { + SummaryMaintenanceDeployment { summary_index, summary, - selected: selected.as_ref().map(|(lifecycle, _)| lifecycle.clone()), + selected_summary_maintenance_lifecycle: selected + .as_ref() + .map(|(lifecycle, _)| lifecycle.clone()), selected_maintenance_mode: selected.map(|(_, mode)| mode), evaluation_schedule, output_representation: OutputRepresentation::SummaryState, @@ -176,7 +195,7 @@ pub fn plan_summary_lifecycles( } }) .collect(); - Ok(LifecyclePlan { + Ok(SummaryMaintenanceLifecyclePlan { root, deployments, horizon, @@ -299,9 +318,9 @@ fn workload_facts( fn alternatives_for( facts: &WorkloadFacts, horizon: Option, - capabilities: LifecycleCapabilities, - costs: LifecycleCostInputs, -) -> Vec { + capabilities: SummaryMaintenanceLifecycleCapabilities, + costs: SummaryMaintenanceLifecycleCostInputs, +) -> Vec { let alternatives = vec![ ephemeral(facts, capabilities, &costs), prepared(facts, capabilities, &costs), @@ -313,15 +332,15 @@ fn alternatives_for( fn ephemeral( facts: &WorkloadFacts, - capabilities: LifecycleCapabilities, - costs: &LifecycleCostInputs, -) -> LifecycleAlternative { - let lifecycle = StateLifecycle::Ephemeral; + capabilities: SummaryMaintenanceLifecycleCapabilities, + costs: &SummaryMaintenanceLifecycleCostInputs, +) -> SummaryMaintenanceLifecycleAlternative { + let lifecycle = SummaryMaintenanceLifecycle::Ephemeral; if !capabilities.ephemeral { return rejected( lifecycle, SummaryMaintenanceMode::DirectBuild, - LifecycleRejection::UnsupportedByRuntime, + SummaryMaintenanceLifecycleRejection::UnsupportedByRuntime, ); } let total_cost = zip_costs(&[ @@ -341,20 +360,20 @@ fn ephemeral( fn prepared( facts: &WorkloadFacts, - capabilities: LifecycleCapabilities, - costs: &LifecycleCostInputs, -) -> LifecycleAlternative { + capabilities: SummaryMaintenanceLifecycleCapabilities, + costs: &SummaryMaintenanceLifecycleCostInputs, +) -> SummaryMaintenanceLifecycleAlternative { let Some((activate_at, retire_at)) = facts.prepared_window else { return rejected( - StateLifecycle::Prepared { + SummaryMaintenanceLifecycle::Prepared { activate_at: TimestampMs(0), retire_at: TimestampMs(0), }, retained_mode(facts), - LifecycleRejection::RequiresPredictableOneTimeQuery, + SummaryMaintenanceLifecycleRejection::RequiresPredictableOneTimeQuery, ); }; - let lifecycle = StateLifecycle::Prepared { + let lifecycle = SummaryMaintenanceLifecycle::Prepared { activate_at, retire_at, }; @@ -362,7 +381,7 @@ fn prepared( return rejected( lifecycle, retained_mode(facts), - LifecycleRejection::UnsupportedByRuntime, + SummaryMaintenanceLifecycleRejection::UnsupportedByRuntime, ); } let seconds = retire_at.0.saturating_sub(activate_at.0) as f64 / 1000.0; @@ -394,31 +413,31 @@ fn prepared( fn shared( facts: &WorkloadFacts, horizon: Option, - capabilities: LifecycleCapabilities, - costs: &LifecycleCostInputs, -) -> LifecycleAlternative { - let lifecycle = StateLifecycle::Shared { + capabilities: SummaryMaintenanceLifecycleCapabilities, + 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, retained_mode(facts), - LifecycleRejection::UnsupportedByRuntime, + SummaryMaintenanceLifecycleRejection::UnsupportedByRuntime, ); } if facts.reads.is_none_or(|reads| reads <= 1.0) { return rejected( lifecycle, retained_mode(facts), - LifecycleRejection::RequiresMultipleReads, + SummaryMaintenanceLifecycleRejection::RequiresMultipleReads, ); } let Some(horizon) = horizon else { return rejected( lifecycle, retained_mode(facts), - LifecycleRejection::RequiresHorizon, + SummaryMaintenanceLifecycleRejection::RequiresHorizon, ); }; let total_cost = retained_cost(facts, costs, horizon.0); @@ -433,15 +452,15 @@ fn shared( fn continuous( facts: &WorkloadFacts, horizon: Option, - capabilities: LifecycleCapabilities, - costs: &LifecycleCostInputs, -) -> LifecycleAlternative { - let lifecycle = StateLifecycle::ContinuouslyMaintained; + capabilities: SummaryMaintenanceLifecycleCapabilities, + costs: &SummaryMaintenanceLifecycleCostInputs, +) -> SummaryMaintenanceLifecycleAlternative { + let lifecycle = SummaryMaintenanceLifecycle::ContinuouslyMaintained; if !capabilities.continuously_maintained { return rejected( lifecycle, SummaryMaintenanceMode::Incremental, - LifecycleRejection::UnsupportedByRuntime, + SummaryMaintenanceLifecycleRejection::UnsupportedByRuntime, ); } if !matches!( @@ -451,21 +470,21 @@ fn continuous( return rejected( lifecycle, SummaryMaintenanceMode::Incremental, - LifecycleRejection::RequiresContinuousData, + SummaryMaintenanceLifecycleRejection::RequiresContinuousData, ); } if facts.update_rate.is_none() { return rejected( lifecycle, SummaryMaintenanceMode::Incremental, - LifecycleRejection::MissingOrStaleIngestionRate, + SummaryMaintenanceLifecycleRejection::MissingOrStaleIngestionRate, ); } let Some(horizon) = horizon else { return rejected( lifecycle, SummaryMaintenanceMode::Incremental, - LifecycleRejection::RequiresHorizon, + SummaryMaintenanceLifecycleRejection::RequiresHorizon, ); }; let total_cost = retained_cost(facts, costs, horizon.0); @@ -477,7 +496,11 @@ fn continuous( ) } -fn retained_cost(facts: &WorkloadFacts, costs: &LifecycleCostInputs, seconds: f64) -> Option { +fn retained_cost( + facts: &WorkloadFacts, + costs: &SummaryMaintenanceLifecycleCostInputs, + seconds: f64, +) -> Option { let reads = facts.reads?; let maintenance = maintenance_cost(facts, costs, seconds)?; Some(Cost( @@ -491,7 +514,7 @@ fn retained_cost(facts: &WorkloadFacts, costs: &LifecycleCostInputs, seconds: f6 fn maintenance_cost( facts: &WorkloadFacts, - costs: &LifecycleCostInputs, + costs: &SummaryMaintenanceLifecycleCostInputs, seconds: f64, ) -> Option { match facts.arrival { @@ -510,29 +533,29 @@ fn zip_costs(costs: &[Option]) -> Option { } fn costed_or_unknown( - lifecycle: StateLifecycle, + lifecycle: SummaryMaintenanceLifecycle, maintenance_mode: SummaryMaintenanceMode, total_cost: Option, assumptions: Vec, -) -> LifecycleAlternative { - LifecycleAlternative { - lifecycle, +) -> SummaryMaintenanceLifecycleAlternative { + SummaryMaintenanceLifecycleAlternative { + summary_maintenance_lifecycle: lifecycle, maintenance_mode, total_cost, rejection: total_cost .is_none() - .then_some(LifecycleRejection::MissingCostEvidence), + .then_some(SummaryMaintenanceLifecycleRejection::MissingCostEvidence), assumptions, } } fn rejected( - lifecycle: StateLifecycle, + lifecycle: SummaryMaintenanceLifecycle, maintenance_mode: SummaryMaintenanceMode, - rejection: LifecycleRejection, -) -> LifecycleAlternative { - LifecycleAlternative { - lifecycle, + rejection: SummaryMaintenanceLifecycleRejection, +) -> SummaryMaintenanceLifecycleAlternative { + SummaryMaintenanceLifecycleAlternative { + summary_maintenance_lifecycle: lifecycle, maintenance_mode, total_cost: None, rejection: Some(rejection), @@ -607,8 +630,11 @@ mod tests { candidates.to_vec() } - fn summary_lifecycle_cost_inputs(&self, _summary: &SummaryNode) -> LifecycleCostInputs { - LifecycleCostInputs { + 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)), @@ -719,19 +745,19 @@ mod tests { #[test] fn unpredictable_one_time_at_rest_selects_ephemeral() { - let plan = plan_summary_lifecycles( + let plan = plan_summary_maintenance_lifecycles( summary(), &workload(vec![batch(Predictability::AdHoc)], vec![], at_rest()), 1_000, None, - LifecycleCapabilities::ALL, + SummaryMaintenanceLifecycleCapabilities::ALL, &UnitCosts, ) .unwrap(); assert_eq!(plan.deployments.len(), 1); assert_eq!( - plan.deployments[0].selected, - Some(StateLifecycle::Ephemeral) + plan.deployments[0].selected_summary_maintenance_lifecycle, + Some(SummaryMaintenanceLifecycle::Ephemeral) ); assert_eq!( plan.deployments[0].selected_maintenance_mode, @@ -749,12 +775,12 @@ mod tests { known_at: Some(TimestampMs(1_000)), }); entry.execute_at = Some(TimestampMs(11_000)); - let plan = plan_summary_lifecycles( + let plan = plan_summary_maintenance_lifecycles( summary(), &workload(vec![entry], vec![], at_rest()), 1_000, None, - LifecycleCapabilities::ALL, + SummaryMaintenanceLifecycleCapabilities::ALL, &UnitCosts, ) .unwrap(); @@ -765,18 +791,18 @@ mod tests { #[test] fn repeated_at_rest_selects_shared_without_inventing_updates() { - let plan = plan_summary_lifecycles( + let plan = plan_summary_maintenance_lifecycles( summary(), &workload(vec![], vec![repeating()], at_rest()), 1_000, Some(Horizon(10.0)), - LifecycleCapabilities::ALL, + SummaryMaintenanceLifecycleCapabilities::ALL, &UnitCosts, ) .unwrap(); assert_eq!( - plan.deployments[0].selected, - Some(StateLifecycle::Shared { + plan.deployments[0].selected_summary_maintenance_lifecycle, + Some(SummaryMaintenanceLifecycle::Shared { retention: DurationMs(10_000) }) ); @@ -786,18 +812,18 @@ mod tests { ); assert_eq!( plan.deployments[0].alternatives[3].rejection, - Some(LifecycleRejection::RequiresContinuousData) + Some(SummaryMaintenanceLifecycleRejection::RequiresContinuousData) ); assert_eq!(plan.update_rate, None); } #[test] fn repeated_continuous_workload_can_select_continuous_maintenance() { - let capabilities = LifecycleCapabilities { + let capabilities = SummaryMaintenanceLifecycleCapabilities { shared: false, - ..LifecycleCapabilities::ALL + ..SummaryMaintenanceLifecycleCapabilities::ALL }; - let plan = plan_summary_lifecycles( + let plan = plan_summary_maintenance_lifecycles( summary(), &workload(vec![], vec![repeating()], continuous(1_000, 60_000)), 1_000, @@ -807,8 +833,8 @@ mod tests { ) .unwrap(); assert_eq!( - plan.deployments[0].selected, - Some(StateLifecycle::ContinuouslyMaintained) + plan.deployments[0].selected_summary_maintenance_lifecycle, + Some(SummaryMaintenanceLifecycle::ContinuouslyMaintained) ); assert_eq!( plan.deployments[0].selected_maintenance_mode, @@ -820,34 +846,37 @@ mod tests { #[test] fn stale_ingestion_evidence_cannot_enable_continuous_maintenance() { - let plan = plan_summary_lifecycles( + let plan = plan_summary_maintenance_lifecycles( summary(), &workload(vec![], vec![repeating()], continuous(1_000, 1_000)), 3_000, Some(Horizon(10.0)), - LifecycleCapabilities::ALL, + SummaryMaintenanceLifecycleCapabilities::ALL, &UnitCosts, ) .unwrap(); assert_eq!( plan.deployments[0].alternatives[3].rejection, - Some(LifecycleRejection::MissingOrStaleIngestionRate) + 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_lifecycles( + let plan = plan_summary_maintenance_lifecycles( summary(), &workload(vec![], vec![repeating()], continuous(1_000, 60_000)), 1_000, Some(Horizon(10.0)), - LifecycleCapabilities::ALL, + SummaryMaintenanceLifecycleCapabilities::ALL, &crate::cost_model::DefaultCostModel, ) .unwrap(); - assert_eq!(plan.deployments[0].selected, None); + assert_eq!( + plan.deployments[0].selected_summary_maintenance_lifecycle, + None + ); assert!(plan.deployments[0] .alternatives .iter() diff --git a/crates/types/src/post_asap/mod.rs b/crates/types/src/post_asap/mod.rs index 5c572b05..4c5de298 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,6 @@ pub use sketch::{ SketchParams, SketchQuery, StatModelKind, StatModelParams, WaveletKind, WaveletParams, }; pub use summary_maintenance::SummaryMaintenanceMode; +pub use summary_maintenance_lifecycle::{ + EvaluationSchedule, OutputRepresentation, SummaryMaintenanceLifecycle, +}; diff --git a/crates/types/src/post_asap/lifecycle.rs b/crates/types/src/post_asap/summary_maintenance_lifecycle.rs similarity index 57% rename from crates/types/src/post_asap/lifecycle.rs rename to crates/types/src/post_asap/summary_maintenance_lifecycle.rs index 995c2f79..614c18f1 100644 --- a/crates/types/src/post_asap/lifecycle.rs +++ b/crates/types/src/post_asap/summary_maintenance_lifecycle.rs @@ -1,7 +1,10 @@ -//! Physical lifecycle vocabulary for summary state. +//! Physical summary-maintenance lifecycle vocabulary. //! -//! These choices are attached by physical planning; a `SummaryAgg` does not -//! imply continuous maintenance by itself. +//! A **summary-maintenance lifecycle** describes when one materialized summary +//! state is created, retained or shared, updated, and retired. It does not +//! describe the broader data lifecycle (collection, transport, and storage), +//! and it is not implied by a logical `SummaryAgg`. Physical planning chooses +//! it from workload evidence and runtime capabilities. use crate::workload::{DurationMs, TimestampMs}; @@ -22,9 +25,9 @@ pub enum OutputRepresentation { FinalizedValue, } -/// How long one planned summary state deployment exists. +/// Lifetime and reuse policy for one materialized summary-state deployment. #[derive(Debug, Clone, PartialEq, Eq, Hash)] -pub enum StateLifecycle { +pub enum SummaryMaintenanceLifecycle { Ephemeral, Prepared { activate_at: TimestampMs, From 81a9d928b73e6152cd3d08d92ed32ab6a0fb2343 Mon Sep 17 00:00:00 2001 From: zz_y Date: Mon, 31 Aug 2026 06:30:13 -0600 Subject: [PATCH 07/11] docs(lifecycle): define planning inputs and outputs --- .../src/summary_maintenance_lifecycle.rs | 98 +++++++++++++++---- 1 file changed, 79 insertions(+), 19 deletions(-) diff --git a/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs b/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs index 503d6957..dedda127 100644 --- a/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs +++ b/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs @@ -27,11 +27,21 @@ use crate::cost_model::{Cost, CostModel}; use crate::recurrence::{CostRate, EvaluationRate, Horizon, UpdateRate}; /// Summary-maintenance lifecycle shapes supported by the target runtime. +/// +/// These flags describe runtime implementation capabilities, not which +/// alternative the planner prefers. A supported alternative may still be +/// rejected because workload evidence is missing or its cost is unknown. #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub struct SummaryMaintenanceLifecycleCapabilities { + /// The runtime can build a fresh state for each invocation and retire it + /// after that invocation finishes. pub ephemeral: bool, + /// The runtime can build state before a predictable execution and retain + /// it until that scheduled execution window ends. pub prepared: bool, + /// The runtime can retain one state and reuse it across multiple reads. pub shared: bool, + /// The runtime can keep state current by applying arriving data updates. pub continuously_maintained: bool, } @@ -52,12 +62,21 @@ impl Default for SummaryMaintenanceLifecycleCapabilities { /// Primitive costs for one concrete summary state. Every field is optional: /// missing statistics produce an uncosted alternative, never a zero. +/// +/// The lifecycle planner combines these state-specific inputs with workload +/// rates, invocation counts, and the optimization horizon. All `Cost` fields +/// are one-time costs unless their name explicitly says otherwise. #[derive(Debug, Clone, Default, PartialEq)] pub struct SummaryMaintenanceLifecycleCostInputs { + /// One-time cost to construct the state from its input. pub build_cost: Option, + /// Cost to incorporate one arriving input update into existing state. pub maintenance_cost_per_update: Option, + /// Cost of one read or finalization from already-built summary state. pub summary_read_cost: Option, + /// Cost per second for retaining the state over a lifecycle window. pub retention_cost_rate: Option, + /// One-time cost to release or retire the state. pub retirement_cost: Option, } @@ -72,13 +91,23 @@ pub enum SummaryMaintenanceLifecycleRejection { MissingCostEvidence, } +/// One candidate physical policy for a particular summary deployment. +/// +/// `total_cost: None` never means zero: it means the planner lacks enough +/// evidence to cost the candidate. Such a candidate is not selectable and its +/// `rejection` explains why. #[derive(Debug, Clone, PartialEq)] pub struct SummaryMaintenanceLifecycleAlternative { + /// State creation, retention, sharing, update, and retirement policy. pub summary_maintenance_lifecycle: SummaryMaintenanceLifecycle, /// How this lifecycle obtains and refreshes its summary state. pub maintenance_mode: SummaryMaintenanceMode, + /// Complete cost over the requested horizon, when every input is known. pub total_cost: Option, + /// Why this alternative cannot be selected; `None` means it is legal and + /// fully costed. pub rejection: Option, + /// Human-readable premises used when deriving and costing the alternative. pub assumptions: Vec, } @@ -91,21 +120,41 @@ impl SummaryMaintenanceLifecycleAlternative { /// One unique summary-state deployment. Shared `Rc` nodes are emitted once. #[derive(Debug, Clone)] pub struct SummaryMaintenanceDeployment { + /// Traversal-local ordinal used to associate this deployment with exports; + /// it is not a persistent identity across independently planned DAGs. pub summary_index: usize, + /// The unique materialized `SummaryAgg` represented by this deployment. pub summary: Rc, + /// Cheapest legal, fully costed lifecycle, or `None` when none is + /// selectable. pub selected_summary_maintenance_lifecycle: Option, + /// Maintenance mechanism paired with the selected lifecycle. It is `None` + /// exactly when no lifecycle was selected. pub selected_maintenance_mode: Option, + /// When the selected deployment is evaluated. It is `None` when no + /// lifecycle was selected. pub evaluation_schedule: Option, + /// Physical value exposed by this deployment to its downstream consumer. pub output_representation: OutputRepresentation, + /// Every lifecycle shape considered, including rejected and uncosted ones. pub alternatives: Vec, } +/// Workload-aware physical deployment decisions for every unique summary +/// state reachable from one materialized post-ASAP root. #[derive(Debug, Clone)] pub struct SummaryMaintenanceLifecyclePlan { + /// Root of the materialized post-ASAP DAG being deployed. pub root: Rc, + /// One entry per unique reachable `SummaryAgg`; shared `Rc` nodes appear + /// only once. pub deployments: Vec, + /// Caller-supplied optimization horizon used to turn rates into total + /// costs. `None` keeps horizon-dependent alternatives unselectable. pub horizon: Option, + /// Aggregate recurring query-evaluation rate derived from the workload. pub evaluation_rate: Option, + /// Fresh source-data ingestion rate, when supplied by the workload. pub update_rate: Option, } @@ -117,13 +166,31 @@ pub enum SummaryMaintenanceLifecyclePlanError { InvalidHorizon, } +/// Workload-wide evidence derived specifically for summary-maintenance +/// lifecycle enumeration and costing. +/// +/// This is not another workload input model. [`QueryWorkload`] and its +/// normalized entries remain the source of truth. Unlike one +/// [`asap_types::workload::QueryWorkloadEntry`], these values aggregate all +/// entries at a particular planning time and optional horizon. It also cannot +/// reuse [`crate::recurrence::RecurrenceProfile`], which describes recurrence +/// for one candidate target and counts consumers rather than invocations. #[derive(Debug)] -struct WorkloadFacts { +struct SummaryMaintenanceWorkloadFacts { + /// Total one-time and recurring reads inside the horizon. `None` means a + /// recurrence or horizon was unknown, not zero reads. reads: Option, + /// Sum of declared invocations across all one-time workload entries. one_time_invocations: u64, + /// Sum of usable recurring query rates in evaluations per second. evaluation_rate: Option, + /// Fresh workload-level ingestion rate in updates per second. update_rate: Option, + /// Whether the workload's source data is static, arriving, mixed, or + /// unknown. arrival: DataArrival, + /// Earliest known activation and latest scheduled execution across + /// predictable one-time entries. `None` means no valid preparation window. prepared_window: Option<(TimestampMs, TimestampMs)>, } @@ -208,7 +275,7 @@ fn workload_facts( workload: &QueryWorkload, now_ms: u64, horizon: Option, -) -> WorkloadFacts { +) -> SummaryMaintenanceWorkloadFacts { let mut one_time_invocations = 0u64; let mut recurring_reads = 0.0; let mut recurring_known = true; @@ -261,14 +328,7 @@ fn workload_facts( } } QueryRecurrence::Repeated(RepeatedDemand::EstimatedRate(estimate)) => { - let fresh = match (estimate.observed_at, estimate.valid_for) { - (Some(observed), Some(valid_for)) => { - now_ms <= observed.0.saturating_add(valid_for.0) - } - (None, Some(_)) => false, - _ => true, - }; - if !fresh { + if !estimate.is_fresh_at(now_ms) { recurring_known = false; continue; } @@ -305,7 +365,7 @@ fn workload_facts( .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); - WorkloadFacts { + SummaryMaintenanceWorkloadFacts { reads, one_time_invocations, evaluation_rate: has_evaluation_rate.then_some(EvaluationRate(evaluation_rate)), @@ -316,7 +376,7 @@ fn workload_facts( } fn alternatives_for( - facts: &WorkloadFacts, + facts: &SummaryMaintenanceWorkloadFacts, horizon: Option, capabilities: SummaryMaintenanceLifecycleCapabilities, costs: SummaryMaintenanceLifecycleCostInputs, @@ -331,7 +391,7 @@ fn alternatives_for( } fn ephemeral( - facts: &WorkloadFacts, + facts: &SummaryMaintenanceWorkloadFacts, capabilities: SummaryMaintenanceLifecycleCapabilities, costs: &SummaryMaintenanceLifecycleCostInputs, ) -> SummaryMaintenanceLifecycleAlternative { @@ -359,7 +419,7 @@ fn ephemeral( } fn prepared( - facts: &WorkloadFacts, + facts: &SummaryMaintenanceWorkloadFacts, capabilities: SummaryMaintenanceLifecycleCapabilities, costs: &SummaryMaintenanceLifecycleCostInputs, ) -> SummaryMaintenanceLifecycleAlternative { @@ -411,7 +471,7 @@ fn prepared( } fn shared( - facts: &WorkloadFacts, + facts: &SummaryMaintenanceWorkloadFacts, horizon: Option, capabilities: SummaryMaintenanceLifecycleCapabilities, costs: &SummaryMaintenanceLifecycleCostInputs, @@ -450,7 +510,7 @@ fn shared( } fn continuous( - facts: &WorkloadFacts, + facts: &SummaryMaintenanceWorkloadFacts, horizon: Option, capabilities: SummaryMaintenanceLifecycleCapabilities, costs: &SummaryMaintenanceLifecycleCostInputs, @@ -497,7 +557,7 @@ fn continuous( } fn retained_cost( - facts: &WorkloadFacts, + facts: &SummaryMaintenanceWorkloadFacts, costs: &SummaryMaintenanceLifecycleCostInputs, seconds: f64, ) -> Option { @@ -513,7 +573,7 @@ fn retained_cost( } fn maintenance_cost( - facts: &WorkloadFacts, + facts: &SummaryMaintenanceWorkloadFacts, costs: &SummaryMaintenanceLifecycleCostInputs, seconds: f64, ) -> Option { @@ -563,7 +623,7 @@ fn rejected( } } -fn retained_mode(facts: &WorkloadFacts) -> SummaryMaintenanceMode { +fn retained_mode(facts: &SummaryMaintenanceWorkloadFacts) -> SummaryMaintenanceMode { match facts.arrival { DataArrival::ContinuouslyIngesting | DataArrival::Mixed => { SummaryMaintenanceMode::Incremental From 483e4bd5bc79fb396863e12c02f1f0d585ce4e5b Mon Sep 17 00:00:00 2001 From: zz_y Date: Mon, 31 Aug 2026 06:41:24 -0600 Subject: [PATCH 08/11] docs(lifecycle): explain how recurrence informs planning --- .../src/summary_maintenance_lifecycle.rs | 7 +++++-- .../types/src/post_asap/summary_maintenance_lifecycle.rs | 5 +++-- 2 files changed, 8 insertions(+), 4 deletions(-) diff --git a/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs b/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs index dedda127..20ddf92f 100644 --- a/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs +++ b/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs @@ -3,8 +3,11 @@ //! A **summary-maintenance lifecycle** is the physical policy for when one //! materialized summary state is created, retained or shared, updated as data //! arrives, and retired. It is deliberately narrower than the end-to-end data -//! lifecycle and independent of query recurrence: recurrence is workload -//! evidence used to choose a lifecycle, not a lifecycle itself. +//! lifecycle and independent of query recurrence. Recurrence says when and how +//! often queries will read the result. The planner converts that demand into +//! expected reads and an evaluation rate, then uses those quantities to compare +//! rebuilding per query with retaining or continuously maintaining state. +//! Recurrence does not itself prescribe a state-maintenance policy. //! //! This module enumerates and costs `Ephemeral`, `Prepared`, `Shared`, and //! `ContinuouslyMaintained` alternatives for every unique `SummaryAgg` in a diff --git a/crates/types/src/post_asap/summary_maintenance_lifecycle.rs b/crates/types/src/post_asap/summary_maintenance_lifecycle.rs index 614c18f1..89ce7e8c 100644 --- a/crates/types/src/post_asap/summary_maintenance_lifecycle.rs +++ b/crates/types/src/post_asap/summary_maintenance_lifecycle.rs @@ -3,8 +3,9 @@ //! A **summary-maintenance lifecycle** describes when one materialized summary //! state is created, retained or shared, updated, and retired. It does not //! describe the broader data lifecycle (collection, transport, and storage), -//! and it is not implied by a logical `SummaryAgg`. Physical planning chooses -//! it from workload evidence and runtime capabilities. +//! and it is not implied by a logical `SummaryAgg`. Physical planning compares +//! alternatives using the expected number and timing of reads, the source-data +//! arrival/update rate, state-operation costs, and runtime capabilities. use crate::workload::{DurationMs, TimestampMs}; From af0342638dc14ae8533acce12a4b00f3fefa3567 Mon Sep 17 00:00:00 2001 From: zz_y Date: Mon, 31 Aug 2026 12:47:54 -0600 Subject: [PATCH 09/11] refactor(lifecycle): clarify runtime capability flags --- .../src/summary_maintenance_lifecycle.rs | 36 ++++++++++--------- 1 file changed, 20 insertions(+), 16 deletions(-) diff --git a/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs b/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs index 20ddf92f..837a9754 100644 --- a/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs +++ b/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs @@ -31,29 +31,33 @@ use crate::recurrence::{CostRate, EvaluationRate, Horizon, UpdateRate}; /// Summary-maintenance lifecycle shapes supported by the target runtime. /// -/// These flags describe runtime implementation capabilities, not which -/// alternative the planner prefers. A supported alternative may still be -/// rejected because workload evidence is missing or its cost is unknown. +/// These independent flags describe the set of lifecycle alternatives the +/// runtime implements, not simultaneous states of one deployment. Multiple +/// flags may be `true` (a runtime can support both ephemeral and prepared +/// state, for example); the planner still selects exactly one mutually +/// exclusive [`SummaryMaintenanceLifecycle`] for each deployment. A supported +/// alternative may still be rejected because workload evidence is missing or +/// its cost is unknown. #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub struct SummaryMaintenanceLifecycleCapabilities { /// The runtime can build a fresh state for each invocation and retire it /// after that invocation finishes. - pub ephemeral: bool, + pub supports_ephemeral: bool, /// The runtime can build state before a predictable execution and retain /// it until that scheduled execution window ends. - pub prepared: bool, + pub supports_prepared: bool, /// The runtime can retain one state and reuse it across multiple reads. - pub shared: bool, + pub supports_shared: bool, /// The runtime can keep state current by applying arriving data updates. - pub continuously_maintained: bool, + pub supports_continuously_maintained: bool, } impl SummaryMaintenanceLifecycleCapabilities { pub const ALL: Self = Self { - ephemeral: true, - prepared: true, - shared: true, - continuously_maintained: true, + supports_ephemeral: true, + supports_prepared: true, + supports_shared: true, + supports_continuously_maintained: true, }; } @@ -399,7 +403,7 @@ fn ephemeral( costs: &SummaryMaintenanceLifecycleCostInputs, ) -> SummaryMaintenanceLifecycleAlternative { let lifecycle = SummaryMaintenanceLifecycle::Ephemeral; - if !capabilities.ephemeral { + if !capabilities.supports_ephemeral { return rejected( lifecycle, SummaryMaintenanceMode::DirectBuild, @@ -440,7 +444,7 @@ fn prepared( activate_at, retire_at, }; - if !capabilities.prepared { + if !capabilities.supports_prepared { return rejected( lifecycle, retained_mode(facts), @@ -482,7 +486,7 @@ fn shared( let lifecycle = SummaryMaintenanceLifecycle::Shared { retention: asap_types::workload::DurationMs(horizon.map_or(0, |h| (h.0 * 1000.0) as u64)), }; - if !capabilities.shared { + if !capabilities.supports_shared { return rejected( lifecycle, retained_mode(facts), @@ -519,7 +523,7 @@ fn continuous( costs: &SummaryMaintenanceLifecycleCostInputs, ) -> SummaryMaintenanceLifecycleAlternative { let lifecycle = SummaryMaintenanceLifecycle::ContinuouslyMaintained; - if !capabilities.continuously_maintained { + if !capabilities.supports_continuously_maintained { return rejected( lifecycle, SummaryMaintenanceMode::Incremental, @@ -883,7 +887,7 @@ mod tests { #[test] fn repeated_continuous_workload_can_select_continuous_maintenance() { let capabilities = SummaryMaintenanceLifecycleCapabilities { - shared: false, + supports_shared: false, ..SummaryMaintenanceLifecycleCapabilities::ALL }; let plan = plan_summary_maintenance_lifecycles( From c865af4c27dfe832a215e12100710b21c7cbc1e3 Mon Sep 17 00:00:00 2001 From: zz_y Date: Mon, 31 Aug 2026 12:50:51 -0600 Subject: [PATCH 10/11] docs(lifecycle): define deployment output boundary --- crates/types/src/post_asap/summary_maintenance_lifecycle.rs | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/crates/types/src/post_asap/summary_maintenance_lifecycle.rs b/crates/types/src/post_asap/summary_maintenance_lifecycle.rs index 89ce7e8c..1cf493d4 100644 --- a/crates/types/src/post_asap/summary_maintenance_lifecycle.rs +++ b/crates/types/src/post_asap/summary_maintenance_lifecycle.rs @@ -18,7 +18,9 @@ pub enum EvaluationSchedule { OnRead, } -/// The physical value crossing an execution boundary. +/// The value a summary-maintenance deployment provides to its downstream +/// query operator: ordinary rows, reusable summary state, or a finalized +/// scalar/result value. #[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] pub enum OutputRepresentation { PlainRows, From 58b1b63c0bb5008948a88b9c52a3a8deb3ba43d8 Mon Sep 17 00:00:00 2001 From: zz_y Date: Mon, 31 Aug 2026 12:54:11 -0600 Subject: [PATCH 11/11] docs(lifecycle): define summary output consumer --- .../types/src/post_asap/summary_maintenance_lifecycle.rs | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/crates/types/src/post_asap/summary_maintenance_lifecycle.rs b/crates/types/src/post_asap/summary_maintenance_lifecycle.rs index 1cf493d4..5cebdc23 100644 --- a/crates/types/src/post_asap/summary_maintenance_lifecycle.rs +++ b/crates/types/src/post_asap/summary_maintenance_lifecycle.rs @@ -18,9 +18,11 @@ pub enum EvaluationSchedule { OnRead, } -/// The value a summary-maintenance deployment provides to its downstream -/// query operator: ordinary rows, reusable summary state, or a finalized -/// scalar/result value. +/// The form in which this deployment exposes its result to its consumer. The +/// consumer is the next operator in the execution plan that reads the +/// summary's output; for example, `Estimate` is the consumer in +/// `SummaryAgg -> Estimate`. The exposed result is ordinary rows, reusable +/// summary state, or a finalized value. #[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] pub enum OutputRepresentation { PlainRows,