diff --git a/crates/asap-aware-mapping/src/cost_model.rs b/crates/asap-aware-mapping/src/cost_model.rs index 4f9225f..fae0243 100644 --- a/crates/asap-aware-mapping/src/cost_model.rs +++ b/crates/asap-aware-mapping/src/cost_model.rs @@ -56,7 +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::lifecycle::{LifecycleCostInputs, SummaryLifecycleCapabilities}; use crate::recurrence::{ self, Horizon, RecurrenceCostExplanation, RecurrenceError, RecurrenceProfile, }; @@ -513,6 +513,22 @@ pub trait CostModel { fn summary_lifecycle_cost_inputs(&self, _summary: &SummaryNode) -> LifecycleCostInputs { LifecycleCostInputs::default() } + + /// Physical update/merge/delete support for one concrete summary. The + /// conservative default advertises no long-lived maintenance capability. + fn summary_lifecycle_capabilities( + &self, + _summary: &SummaryNode, + ) -> SummaryLifecycleCapabilities { + SummaryLifecycleCapabilities::default() + } + + /// Cost of evaluating `target` directly from its logical/raw inputs once. + /// When known, lifecycle-aware materialization compares this fallback with + /// the aggregate cost of the selected summary deployments. + fn raw_query_recompute_cost(&self, _target: &QueryExpr) -> Option { + None + } } fn sketch_state( diff --git a/crates/asap-aware-mapping/src/lib.rs b/crates/asap-aware-mapping/src/lib.rs index 773464b..be0f012 100644 --- a/crates/asap-aware-mapping/src/lib.rs +++ b/crates/asap-aware-mapping/src/lib.rs @@ -207,6 +207,7 @@ pub use grouping::{has_subpopulations, HydraGroupingStrategy}; pub use lifecycle::{ plan_summary_lifecycles, LifecycleAlternative, LifecycleCapabilities, LifecycleCostInputs, LifecyclePlan, LifecyclePlanError, LifecycleRejection, StateDeployment, + SummaryLifecycleCapabilities, WorkloadDemand, }; pub use recurrence::{ evaluation_rate_of, total_cost, update_rate_from_data_workload, CostRate, EvaluationRate, diff --git a/crates/asap-aware-mapping/src/lifecycle.rs b/crates/asap-aware-mapping/src/lifecycle.rs index 08b1ab1..8a4a37f 100644 --- a/crates/asap-aware-mapping/src/lifecycle.rs +++ b/crates/asap-aware-mapping/src/lifecycle.rs @@ -26,6 +26,14 @@ pub struct LifecycleCapabilities { pub continuously_maintained: bool, } +/// Capabilities of one concrete summary family/state representation. +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +pub struct SummaryLifecycleCapabilities { + pub incremental_update: bool, + pub merge: bool, + pub delete: bool, +} + impl LifecycleCapabilities { pub const ALL: Self = Self { ephemeral: true, @@ -60,6 +68,8 @@ pub enum LifecycleRejection { RequiresHorizon, RequiresContinuousData, MissingOrStaleIngestionRate, + SummaryDoesNotSupportIncrementalUpdates, + SummaryDoesNotSupportDeletion, MissingCostEvidence, } @@ -98,6 +108,27 @@ pub struct LifecyclePlan { pub horizon: Option, pub evaluation_rate: Option, pub update_rate: Option, + pub expected_reads: Option, + pub selected_raw_recompute: bool, + pub summary_total_cost: Option, + pub raw_recompute_total_cost: Option, +} + +/// Explicit association between a materialized target and the normalized +/// workload entries whose demand consumes it. +#[derive(Debug, Clone, Copy)] +pub struct WorkloadDemand<'a> { + pub workload: &'a QueryWorkload, + pub entry_indices: &'a [usize], +} + +impl<'a> WorkloadDemand<'a> { + pub const fn new(workload: &'a QueryWorkload, entry_indices: &'a [usize]) -> Self { + Self { + workload, + entry_indices, + } + } } #[derive(Debug, thiserror::Error)] @@ -106,6 +137,12 @@ pub enum LifecyclePlanError { InvalidWorkload(#[from] WorkloadError), #[error("optimization horizon must be finite and strictly positive")] InvalidHorizon, + #[error("workload entry index {index} is out of bounds for {entry_count} entries")] + InvalidWorkloadEntry { index: usize, entry_count: usize }, + #[error("a workload-demand binding must contain at least one entry")] + EmptyWorkloadDemand, + #[error("workload entry index {index} appears more than once in one demand binding")] + DuplicateWorkloadEntry { index: usize }, } #[derive(Debug)] @@ -116,6 +153,8 @@ struct WorkloadFacts { update_rate: Option, arrival: DataArrival, prepared_window: Option<(TimestampMs, TimestampMs)>, + prepared_eligible: bool, + requires_deletion: bool, } /// Validate a materialized plan, enumerate lifecycle alternatives for each @@ -123,20 +162,20 @@ struct WorkloadFacts { /// is fully known. pub fn plan_summary_lifecycles( root: Rc, - workload: &QueryWorkload, + demand: WorkloadDemand<'_>, now_ms: u64, horizon: Option, capabilities: LifecycleCapabilities, cost_model: &dyn CostModel, ) -> Result { - workload.validate()?; + demand.workload.validate()?; if horizon.is_some_and(|h| !h.0.is_finite() || h.0 <= 0.0) { return Err(LifecyclePlanError::InvalidHorizon); } - let facts = workload_facts(workload, now_ms, horizon); + let facts = workload_facts(demand.workload, demand.entry_indices, now_ms, horizon)?; let mut summaries = Vec::new(); collect_summary_aggs(&root, &mut HashSet::new(), &mut summaries); - let deployments = summaries + let deployments: Vec = summaries .into_iter() .enumerate() .map(|(summary_index, summary)| { @@ -144,6 +183,7 @@ pub fn plan_summary_lifecycles( &facts, horizon, capabilities, + cost_model.summary_lifecycle_capabilities(&summary), cost_model.summary_lifecycle_cost_inputs(&summary), ); let selected = alternatives @@ -176,20 +216,34 @@ pub fn plan_summary_lifecycles( } }) .collect(); + let summary_total_cost = deployments.iter().try_fold(Cost::ZERO, |sum, deployment| { + let selected = deployment.selected.as_ref()?; + let cost = deployment + .alternatives + .iter() + .find(|alternative| &alternative.lifecycle == selected)? + .total_cost?; + Some(Cost(sum.0 + cost.0)) + }); Ok(LifecyclePlan { root, deployments, horizon, evaluation_rate: facts.evaluation_rate, update_rate: facts.update_rate, + expected_reads: facts.reads, + selected_raw_recompute: false, + summary_total_cost, + raw_recompute_total_cost: None, }) } fn workload_facts( workload: &QueryWorkload, + workload_entry_indices: &[usize], now_ms: u64, horizon: Option, -) -> WorkloadFacts { +) -> Result { let mut one_time_invocations = 0u64; let mut recurring_reads = 0.0; let mut recurring_known = true; @@ -197,15 +251,38 @@ fn workload_facts( let mut has_evaluation_rate = false; let mut prepared_start: Option = None; let mut prepared_end: Option = None; + let mut prepared_eligible = true; + let mut requires_deletion = false; - for entry in workload.entries() { + let entries: Vec<_> = workload.entries().collect(); + if workload_entry_indices.is_empty() { + return Err(LifecyclePlanError::EmptyWorkloadDemand); + } + let mut seen_indices = HashSet::new(); + for &index in workload_entry_indices { + if !seen_indices.insert(index) { + return Err(LifecyclePlanError::DuplicateWorkloadEntry { index }); + } + let entry = entries + .get(index) + .ok_or(LifecyclePlanError::InvalidWorkloadEntry { + index, + entry_count: entries.len(), + })?; + requires_deletion |= entry.time_selection.lookback.is_some() + && entry.time_selection.as_of.is_none() + && matches!( + entry.time_selection.scope, + asap_types::workload::QueryTimeScope::RealTime + | asap_types::workload::QueryTimeScope::Mixed + ); match &entry.recurrence { QueryRecurrence::OneTime { invocations, execute_at, } => { one_time_invocations = one_time_invocations.saturating_add(*invocations); - if let ( + let covered = if let ( Predictability::Predictable { known_at: Some(known), }, @@ -215,10 +292,17 @@ fn workload_facts( if known < execute { prepared_start = Some(prepared_start.map_or(*known, |old| old.min(*known))); prepared_end = Some(prepared_end.map_or(*execute, |old| old.max(*execute))); + true + } else { + false } - } + } else { + false + }; + prepared_eligible &= covered; } QueryRecurrence::Repeated(RepeatedDemand::FixedInterval(interval)) => { + prepared_eligible = false; let rate = 1000.0 / f64::from(interval.0); evaluation_rate += rate; has_evaluation_rate = true; @@ -229,27 +313,23 @@ fn workload_facts( } } QueryRecurrence::Repeated(RepeatedDemand::Scheduled(schedule)) => { + prepared_eligible = false; if let Some(h) = horizon { let end_ms = now_ms.saturating_add((h.0 * 1000.0) as u64); - recurring_reads += schedule + let reads_in_horizon = schedule .iter() .filter(|at| at.0 >= now_ms && at.0 <= end_ms) .count() as f64; - evaluation_rate += schedule.len() as f64 / h.0; + recurring_reads += reads_in_horizon; + evaluation_rate += reads_in_horizon / h.0; has_evaluation_rate = true; } else { recurring_known = false; } } QueryRecurrence::Repeated(RepeatedDemand::EstimatedRate(estimate)) => { - 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 { + prepared_eligible = false; + if !estimate.is_fresh_at(now_ms) { recurring_known = false; continue; } @@ -276,7 +356,10 @@ fn workload_facts( recurring_known = false; } } - QueryRecurrence::Unknown => recurring_known = false, + QueryRecurrence::Unknown => { + prepared_eligible = false; + recurring_known = false; + } } } @@ -286,27 +369,30 @@ 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 { + Ok(WorkloadFacts { reads, one_time_invocations, evaluation_rate: has_evaluation_rate.then_some(EvaluationRate(evaluation_rate)), update_rate, arrival, prepared_window: prepared_start.zip(prepared_end), - } + prepared_eligible, + requires_deletion, + }) } fn alternatives_for( facts: &WorkloadFacts, horizon: Option, capabilities: LifecycleCapabilities, + summary_capabilities: SummaryLifecycleCapabilities, costs: LifecycleCostInputs, ) -> Vec { let alternatives = vec![ ephemeral(facts, capabilities, &costs), - prepared(facts, capabilities, &costs), - shared(facts, horizon, capabilities, &costs), - continuous(facts, horizon, capabilities, &costs), + prepared(facts, capabilities, summary_capabilities, &costs), + shared(facts, horizon, capabilities, summary_capabilities, &costs), + continuous(facts, horizon, capabilities, summary_capabilities, &costs), ]; alternatives } @@ -342,8 +428,19 @@ fn ephemeral( fn prepared( facts: &WorkloadFacts, capabilities: LifecycleCapabilities, + summary_capabilities: SummaryLifecycleCapabilities, costs: &LifecycleCostInputs, ) -> LifecycleAlternative { + if !facts.prepared_eligible { + return rejected( + StateLifecycle::Prepared { + activate_at: TimestampMs(0), + retire_at: TimestampMs(0), + }, + retained_mode(facts), + LifecycleRejection::RequiresPredictableOneTimeQuery, + ); + } let Some((activate_at, retire_at)) = facts.prepared_window else { return rejected( StateLifecycle::Prepared { @@ -365,6 +462,9 @@ fn prepared( LifecycleRejection::UnsupportedByRuntime, ); } + if let Some(rejection) = maintenance_capability_rejection(facts, summary_capabilities) { + return rejected(lifecycle, retained_mode(facts), rejection); + } let seconds = retire_at.0.saturating_sub(activate_at.0) as f64 / 1000.0; let maintenance = maintenance_cost(facts, costs, seconds); let total_cost = match ( @@ -395,6 +495,7 @@ fn shared( facts: &WorkloadFacts, horizon: Option, capabilities: LifecycleCapabilities, + summary_capabilities: SummaryLifecycleCapabilities, costs: &LifecycleCostInputs, ) -> LifecycleAlternative { let lifecycle = StateLifecycle::Shared { @@ -407,6 +508,9 @@ fn shared( LifecycleRejection::UnsupportedByRuntime, ); } + if let Some(rejection) = maintenance_capability_rejection(facts, summary_capabilities) { + return rejected(lifecycle, retained_mode(facts), rejection); + } if facts.reads.is_none_or(|reads| reads <= 1.0) { return rejected( lifecycle, @@ -434,6 +538,7 @@ fn continuous( facts: &WorkloadFacts, horizon: Option, capabilities: LifecycleCapabilities, + summary_capabilities: SummaryLifecycleCapabilities, costs: &LifecycleCostInputs, ) -> LifecycleAlternative { let lifecycle = StateLifecycle::ContinuouslyMaintained; @@ -461,6 +566,9 @@ fn continuous( LifecycleRejection::MissingOrStaleIngestionRate, ); } + if let Some(rejection) = maintenance_capability_rejection(facts, summary_capabilities) { + return rejected(lifecycle, SummaryMaintenanceMode::Incremental, rejection); + } let Some(horizon) = horizon else { return rejected( lifecycle, @@ -477,6 +585,28 @@ fn continuous( ) } +fn maintenance_capability_rejection( + facts: &WorkloadFacts, + capabilities: SummaryLifecycleCapabilities, +) -> Option { + if matches!( + facts.arrival, + DataArrival::ContinuouslyIngesting | DataArrival::Mixed + ) && !capabilities.incremental_update + { + Some(LifecycleRejection::SummaryDoesNotSupportIncrementalUpdates) + } else if matches!( + facts.arrival, + DataArrival::ContinuouslyIngesting | DataArrival::Mixed + ) && facts.requires_deletion + && !capabilities.delete + { + Some(LifecycleRejection::SummaryDoesNotSupportDeletion) + } else { + None + } +} + fn retained_cost(facts: &WorkloadFacts, costs: &LifecycleCostInputs, seconds: f64) -> Option { let reads = facts.reads?; let maintenance = maintenance_cost(facts, costs, seconds)?; @@ -616,11 +746,55 @@ mod tests { retirement_cost: Some(Cost(1.0)), } } + + fn summary_lifecycle_capabilities( + &self, + _summary: &SummaryNode, + ) -> SummaryLifecycleCapabilities { + SummaryLifecycleCapabilities { + incremental_update: true, + merge: true, + delete: true, + } + } + } + + struct NoDelete; + + impl CostModel for NoDelete { + fn rank_candidates( + &self, + _intent: &asap_types::pre_asap::AggIntent, + candidates: &[asap_types::post_asap::SketchAlgorithm], + ) -> Vec { + candidates.to_vec() + } + + fn summary_lifecycle_cost_inputs(&self, summary: &SummaryNode) -> LifecycleCostInputs { + UnitCosts.summary_lifecycle_cost_inputs(summary) + } + + fn summary_lifecycle_capabilities( + &self, + _summary: &SummaryNode, + ) -> SummaryLifecycleCapabilities { + SummaryLifecycleCapabilities { + incremental_update: true, + merge: true, + delete: false, + } + } } fn query_root() -> Rc { + query_root_for("m") + } + + fn query_root_for(metric: &str) -> Rc { Rc::new(QueryExpr::Scan { - source: Source::TimeSeries { metric: "m".into() }, + source: Source::TimeSeries { + metric: metric.into(), + }, predicates: vec![], schema: Schema::with_time_index( vec![ @@ -721,7 +895,10 @@ mod tests { fn unpredictable_one_time_at_rest_selects_ephemeral() { let plan = plan_summary_lifecycles( summary(), - &workload(vec![batch(Predictability::AdHoc)], vec![], at_rest()), + WorkloadDemand::new( + &workload(vec![batch(Predictability::AdHoc)], vec![], at_rest()), + &[0], + ), 1_000, None, LifecycleCapabilities::ALL, @@ -751,7 +928,7 @@ mod tests { entry.execute_at = Some(TimestampMs(11_000)); let plan = plan_summary_lifecycles( summary(), - &workload(vec![entry], vec![], at_rest()), + WorkloadDemand::new(&workload(vec![entry], vec![], at_rest()), &[0]), 1_000, None, LifecycleCapabilities::ALL, @@ -767,7 +944,7 @@ mod tests { fn repeated_at_rest_selects_shared_without_inventing_updates() { let plan = plan_summary_lifecycles( summary(), - &workload(vec![], vec![repeating()], at_rest()), + WorkloadDemand::new(&workload(vec![], vec![repeating()], at_rest()), &[0]), 1_000, Some(Horizon(10.0)), LifecycleCapabilities::ALL, @@ -799,7 +976,10 @@ mod tests { }; let plan = plan_summary_lifecycles( summary(), - &workload(vec![], vec![repeating()], continuous(1_000, 60_000)), + WorkloadDemand::new( + &workload(vec![], vec![repeating()], continuous(1_000, 60_000)), + &[0], + ), 1_000, Some(Horizon(10.0)), capabilities, @@ -822,7 +1002,10 @@ mod tests { fn stale_ingestion_evidence_cannot_enable_continuous_maintenance() { let plan = plan_summary_lifecycles( summary(), - &workload(vec![], vec![repeating()], continuous(1_000, 1_000)), + WorkloadDemand::new( + &workload(vec![], vec![repeating()], continuous(1_000, 1_000)), + &[0], + ), 3_000, Some(Horizon(10.0)), LifecycleCapabilities::ALL, @@ -840,7 +1023,10 @@ mod tests { 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)), + WorkloadDemand::new( + &workload(vec![], vec![repeating()], continuous(1_000, 60_000)), + &[0], + ), 1_000, Some(Horizon(10.0)), LifecycleCapabilities::ALL, @@ -854,6 +1040,133 @@ mod tests { .all(|alternative| alternative.rejection.is_some())); } + #[test] + fn unrelated_workload_entries_do_not_create_reuse_for_a_target() { + let plan = plan_summary_lifecycles( + summary(), + WorkloadDemand::new( + &workload( + vec![batch(Predictability::AdHoc), batch(Predictability::AdHoc)], + vec![], + at_rest(), + ), + &[0], + ), + 1_000, + Some(Horizon(10.0)), + LifecycleCapabilities::ALL, + &UnitCosts, + ) + .unwrap(); + assert_eq!( + plan.deployments[0].selected, + Some(StateLifecycle::Ephemeral) + ); + assert_eq!( + plan.deployments[0].alternatives[2].rejection, + Some(LifecycleRejection::RequiresMultipleReads) + ); + } + + #[test] + fn scheduled_rate_counts_only_executions_inside_the_horizon() { + let mut entry = repeating(); + entry.demand = RepeatedDemand::Scheduled(vec![ + TimestampMs(999), + TimestampMs(5_000), + TimestampMs(20_000), + ]); + let plan = plan_summary_lifecycles( + summary(), + WorkloadDemand::new(&workload(vec![], vec![entry], at_rest()), &[0]), + 1_000, + Some(Horizon(10.0)), + LifecycleCapabilities::ALL, + &UnitCosts, + ) + .unwrap(); + assert_eq!(plan.evaluation_rate, Some(EvaluationRate(0.1))); + } + + #[test] + fn demand_binding_rejects_empty_and_duplicate_entries() { + let workload = workload(vec![batch(Predictability::AdHoc)], vec![], at_rest()); + assert!(matches!( + plan_summary_lifecycles( + summary(), + WorkloadDemand::new(&workload, &[]), + 1_000, + None, + LifecycleCapabilities::ALL, + &UnitCosts, + ), + Err(LifecyclePlanError::EmptyWorkloadDemand) + )); + assert!(matches!( + plan_summary_lifecycles( + summary(), + WorkloadDemand::new(&workload, &[0, 0]), + 1_000, + None, + LifecycleCapabilities::ALL, + &UnitCosts, + ), + Err(LifecyclePlanError::DuplicateWorkloadEntry { index: 0 }) + )); + } + + #[test] + fn prepared_requires_every_bound_consumer_to_be_scheduled_and_predictable() { + let mut predictable = batch(Predictability::Predictable { + known_at: Some(TimestampMs(1_000)), + }); + predictable.execute_at = Some(TimestampMs(2_000)); + let workload = workload( + vec![predictable, batch(Predictability::AdHoc)], + vec![], + at_rest(), + ); + let plan = plan_summary_lifecycles( + summary(), + WorkloadDemand::new(&workload, &[0, 1]), + 1_000, + Some(Horizon(10.0)), + LifecycleCapabilities::ALL, + &UnitCosts, + ) + .unwrap(); + assert_eq!( + plan.deployments[0].alternatives[1].rejection, + Some(LifecycleRejection::RequiresPredictableOneTimeQuery) + ); + } + + #[test] + fn moving_realtime_maintenance_requires_summary_deletion_support() { + let mut entry = repeating(); + entry.time_selection = TimeSelection { + scope: asap_types::workload::QueryTimeScope::RealTime, + lookback: Some(DurationMs(60_000)), + as_of: None, + }; + let plan = plan_summary_lifecycles( + summary(), + WorkloadDemand::new( + &workload(vec![], vec![entry], continuous(1_000, 60_000)), + &[0], + ), + 1_000, + Some(Horizon(10.0)), + LifecycleCapabilities::ALL, + &NoDelete, + ) + .unwrap(); + assert_eq!( + plan.deployments[0].alternatives[3].rejection, + Some(LifecycleRejection::SummaryDoesNotSupportDeletion) + ); + } + #[test] fn normalized_workload_drives_plan_space_recurrence_profiles() { let root = query_root(); @@ -869,4 +1182,28 @@ mod tests { assert_eq!(profile.update_rate, Some(UpdateRate(1.0))); assert_eq!(profile.one_shot_consumers, 0); } + + #[test] + fn recurrence_binding_is_explicit_when_root_order_differs_from_workload_order() { + let repeating_root = query_root_for("dashboard"); + let batch_root = query_root_for("batch"); + let space = crate::replacement::search_workload(vec![ + ("dashboard", repeating_root), + ("batch", batch_root), + ]); + let workload = workload( + vec![batch(Predictability::AdHoc)], + vec![repeating()], + at_rest(), + ); + let profiles = space + .recurrence_profiles_from_workload(&workload, &[1, 0], 1_000, Some(Horizon(10.0))) + .unwrap(); + let dashboard = profiles.for_target(&space.roots[0].1); + let batch = profiles.for_target(&space.roots[1].1); + assert_eq!(dashboard.evaluation_rate, Some(EvaluationRate(1.0))); + assert_eq!(dashboard.one_shot_consumers, 0); + assert_eq!(batch.evaluation_rate, None); + assert_eq!(batch.one_shot_consumers, 1); + } } diff --git a/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 fef2aff..0a65805 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 @@ -60,10 +60,12 @@ The planner receives four logically distinct inputs: distribution; 4. existing summaries and the lifecycle actions available to the deployment. -The output is a legal physical-plan choice plus explicit state deployments. A +The implemented output is a phase-valid selected summary plan (or a +cost-preferred raw-recomputation fallback) plus explicit state deployments. A state deployment states whether a summary is ephemeral, prepared, shared for a -bounded period, or continuously maintained. Its cost explanation identifies -the demand and data evidence used in the decision. +bounded period, or continuously maintained. It retains costs, assumptions, and +structured rejection reasons. Exporting full input provenance remains a later +integration. ```text logical queries ---+ @@ -92,7 +94,8 @@ normalize query and data workloads -> validate summary capabilities and phase constraints -> derive and check accuracy guarantees -> normalize one-time and rate costs over an explicit horizon - -> rank legal alternatives + -> rank legal alternatives and compare the selected summary deployment + with raw recomputation -> emit plan, deployments, assumptions, and rejected alternatives ``` @@ -191,8 +194,8 @@ arbitrary approximation: the current normalization policy makes it whether the caller chose exactness or inherited the default. An unspecified response-latency requirement imposes no response-time constraint; it is not a zero-duration bound or evidence that every latency is acceptable. Accuracy is -checked as a legality constraint, while response latency is used to reject -plans that cannot meet the bound. +checked as a legality constraint. The normalized model preserves response +latency, but the current planner does not yet reject plans against that bound. #### Classification axes @@ -393,7 +396,8 @@ struct Evidence { ``` This reuses the provenance and freshness principles from empirical summary -parameter configuration. A missing or stale value remains unknown. +parameter configuration. Missing, stale, or future-dated evidence remains +unknown. ### Output cardinality is a derived or evidenced cost input @@ -461,7 +465,9 @@ enum StateLifecycle { The summary family and its properties constrain which lifecycles are legal. For example, an append-only sketch may support continuous inserts but not a sliding-window lifecycle requiring deletion. Lifecycle legality is checked -before cost ranking, like accuracy legality. +before cost ranking, like accuracy legality. Deployments provide these +per-summary properties through `summary_lifecycle_capabilities`; moving +real-time windows require deletion support as well as incremental updates. ### Existing summaries are planning input @@ -509,6 +515,12 @@ For repeated raw recomputation: total(H) = reads(H) * raw_recompute_cost ``` +The current lifecycle-aware materialization sums the selected summary +deployments and can replace that plan with raw recomputation when the raw cost +is lower or the summary lifecycle is uncostable. Jointly reconsidering every +sibling semantic candidate under lifecycle costs remains a later optimizer +integration; this document does not claim that broader search is implemented. + For an ephemeral summary: ```text @@ -552,15 +564,15 @@ The glossary review found the following required coverage and current gaps. | Glossary concept | Current ASAPPlanner representation | Missing design support | | --- | --- | --- | -| Data at rest vs continuously ingesting | Continuous ingest characteristics are available; no explicit arrival mode | Add `DataArrival`; support at-rest statistics without inventing update rate | -| Ingestion volume | Not a first-class workload input | Add evidenced volume with a time basis | -| Ingestion rate | Derived from series count and sample rate | Preserve as evidenced rate; do not conflate with query evaluation rate | -| Input cardinality | Partial `series_count` and distinct-key inputs | Associate each estimate with its dataset, metric, columns, and observation window | -| Data distribution | Small built-in enum | Preserve source/freshness; permit deployment-specific distributions later | -| Ad-hoc vs predictable | Not represented | Add predictability independently from recurrence | -| One-time vs repeated | Batch entries and fixed-interval repeating entries | Add scheduled one-time, unknown recurrence, and estimated/scheduled repetition | -| Query volume and characteristics | Fixed interval or structural consumer count | Add observation window, peak/burst and concurrency evidence where latency or capacity models require it | -| Real-time vs longitudinal | Temporal IR can carry ranges; no workload classification | Add time scope plus concrete selection; avoid inferring scope from lookback alone | +| Data at rest vs continuously ingesting | `DataArrival` is explicit | Runtime/catalog-specific arrival discovery remains external | +| Ingestion volume | `DataWorkload::ingestion_volume` carries evidence | A concrete time basis for volume remains deployment-specific | +| Ingestion rate | Evidenced independently from query evaluation rate | Preserve richer unit/provenance metadata when integrations require it | +| Input cardinality | Evidenced workload-level cardinality feeds accuracy | Per-dataset/metric/column scoping remains future work | +| Data distribution | Evidenced built-in enum | Permit deployment-specific distributions later | +| Ad-hoc vs predictable | `Predictability` is independent from recurrence | Parameterized-template equivalence remains open | +| One-time vs repeated | One-time, fixed, scheduled, estimated, and unknown recurrence | Forecast-policy integration remains future work | +| Query volume and characteristics | Estimates preserve average/count, peak, concurrency, confidence, and freshness | Peak and concurrency are not yet consumed by cost or latency models | +| Real-time vs longitudinal | `TimeSelection` carries scope, lookback, and `as_of` | Conflict policy with temporal IR remains open | | Output cardinality | May be inferred locally; no common evidenced input | Add derived/evidenced value and provenance for costing | | Lookback window | Represented in temporal query shapes/frontends | Establish query IR as authority and expose it to workload costing | | CTSA pipeline | Not explicitly modeled | Keep as architectural context; planner consumes collect/store/analyze facts but does not model transmission topology in the MVP | @@ -655,9 +667,9 @@ and after aggregation. - **Understandability:** explanations use glossary terms and show each axis separately. Proxy: reviewers can distinguish repeated queries from continuous ingestion in exported plan evidence. -- **Debuggability:** selected and rejected lifecycle alternatives record demand, - horizon, data statistics, and provenance. Proxy: no lifecycle decision is - explained only as a scalar cost. +- **Debuggability:** selected and rejected lifecycle alternatives record costs, + horizon-derived decisions, assumptions, and typed rejection reasons. Full + demand/data provenance in exported explanations remains future work. - **Maintainability:** current recurrence types remain the cost authority; normalized workload types remain the source authority. No duplicate formula system is introduced.