From 28d3a1bb5a3cd44ff9809e05f992ab0067b157a3 Mon Sep 17 00:00:00 2001 From: zz_y Date: Wed, 2 Sep 2026 08:52:25 -0600 Subject: [PATCH 1/4] feat(cost): compare streaming lifecycle plans --- .../asap-aware-mapping/src/analytical_cost.rs | 2 + .../src/summary_maintenance_cost.rs | 660 +++++++++++++++++- 2 files changed, 649 insertions(+), 13 deletions(-) diff --git a/crates/asap-aware-mapping/src/analytical_cost.rs b/crates/asap-aware-mapping/src/analytical_cost.rs index d10fdacf..b4f3d44b 100644 --- a/crates/asap-aware-mapping/src/analytical_cost.rs +++ b/crates/asap-aware-mapping/src/analytical_cost.rs @@ -1946,6 +1946,8 @@ pub enum AnalyticalCostError { InvalidIngestionRate(f64), #[error("summary lifecycle, maintenance mode, and evaluation schedule are inconsistent")] IncompatibleLifecycleGuarantee, + #[error("bootstrap row and byte evidence must either both be zero or both be non-zero")] + InconsistentBootstrapEvidence, #[error("summary operation cost {0} must be finite and non-negative, got {1}")] InvalidOperationCost(&'static str, f64), #[error("required analytical input {0} is missing or zero")] diff --git a/crates/asap-aware-mapping/src/summary_maintenance_cost.rs b/crates/asap-aware-mapping/src/summary_maintenance_cost.rs index ced47e2e..55cba088 100644 --- a/crates/asap-aware-mapping/src/summary_maintenance_cost.rs +++ b/crates/asap-aware-mapping/src/summary_maintenance_cost.rs @@ -7,15 +7,24 @@ use std::collections::HashSet; use asap_types::post_asap::{ - SummaryExpr, SummaryMaintenanceLifecycle, SummaryMaintenanceLifecycleGuarantee, - SummaryMaintenanceMode, SummaryNode, + SketchAlgorithm, SummaryExpr, SummaryMaintenanceLifecycle, + SummaryMaintenanceLifecycleGuarantee, SummaryMaintenanceMode, SummaryNode, }; +use asap_types::pre_asap::{agg_intent::AggIntent, QueryExpr}; use asap_types::workload::{DataArrival, DataWorkload, QueryWorkloadEntry}; use serde::{Deserialize, Serialize}; -use crate::analytical_cost::{AnalyticalCostError, ResourceEstimate}; +use crate::analytical_cost::{ + AnalyticalCostError, ResourceCalibration, ResourceEstimate, +}; +use crate::cost_model::{Cost, CostModel, DefaultCostModel}; use crate::physical_operator_statistics::evaluations_in_horizon; -use crate::summary_maintenance_lifecycle::{evaluation_schedule, maintenance_mode}; +use crate::recurrence::CostRate; +use crate::replacement::{ReplacementSubDAG, TargetSubDAG}; +use crate::summary_maintenance_lifecycle::{ + evaluation_schedule, maintenance_mode, SummaryMaintenanceCapabilities, + SummaryMaintenanceLifecycleCostInputs, +}; pub const SUMMARY_MAINTENANCE_COST_MODEL_VERSION: &str = "summary-maintenance-resource-v1"; @@ -135,9 +144,7 @@ impl StreamingSummaryInputs { self.initial_input_bytes > 0 && self.initial_source_scan_bytes > 0 }; if !bootstrap_is_consistent { - return Err(AnalyticalCostError::InconsistentOperatorStatistics( - "streaming bootstrap rows, logical bytes, and source bytes disagree", - )); + return Err(AnalyticalCostError::InconsistentBootstrapEvidence); } if !self.ingestion_rate_per_second.is_finite() || self.ingestion_rate_per_second < 0.0 { return Err(AnalyticalCostError::InvalidIngestionRate( @@ -160,6 +167,119 @@ pub struct SummaryOperationCpuEvidence { pub readout_cpu_ops: Option, } +/// Physical evidence for one `SummaryJoin` implementation. Cardinality and +/// working memory cannot be inferred from the logical join key alone. +#[derive(Debug, Clone, Copy, PartialEq, Serialize, Deserialize)] +pub struct SummaryJoinEvidence { + pub matched_state_pairs_per_evaluation: u64, + pub cpu_ops_per_matched_pair: f64, + pub working_memory_bytes: u64, +} + +/// Complete physical work for one raw evaluation. It is deliberately +/// per-evaluation so the same normalized query recurrence/horizon can multiply +/// both raw and summary alternatives. +#[derive(Debug, Clone, Copy, PartialEq, Serialize, Deserialize)] +pub struct StreamingRawInputEvidence { + pub input_rows_per_evaluation: u64, + pub input_bytes_per_evaluation: u64, + pub source_scan_bytes_per_evaluation: u64, + pub cpu_ops_per_row: f64, + pub peak_memory_bytes: u64, +} + +/// Adapter that supplies the existing lifecycle planner with analytical +/// streaming costs. It does not define lifecycle policy: the planner's +/// existing enums and legality checks remain authoritative. +#[derive(Debug, Clone)] +pub struct StreamingAnalyticalCostModel { + pub summary_inputs: StreamingSummaryInputs, + pub raw: StreamingRawInputEvidence, + pub cpu: SummaryOperationCpuEvidence, + pub calibration: ResourceCalibration, + pub capabilities: SummaryMaintenanceCapabilities, +} + +impl StreamingAnalyticalCostModel { + fn calibrated(&self, estimate: ResourceEstimate) -> Option { + estimate.calibrated_cost(&self.calibration).ok().map(Cost) + } + + fn lifecycle_inputs(&self) -> Option { + let inputs = self.summary_inputs.validate().ok()?; + let insert = required_cpu("insert_cpu_ops", self.cpu.insert_cpu_ops).ok()?; + let readout = required_cpu("readout_cpu_ops", self.cpu.readout_cpu_ops).ok()?; + let build = self.calibrated(ResourceEstimate { + cpu_ops: inputs.initial_input_rows as f64 * insert, + peak_memory_bytes: 0, + scan_bytes: inputs.initial_source_scan_bytes, + })?; + let maintenance = self.calibrated(ResourceEstimate { + cpu_ops: inputs.active_window_count as f64 * insert, + peak_memory_bytes: 0, + scan_bytes: 0, + })?; + let read = self.calibrated(ResourceEstimate { + cpu_ops: inputs.physical_sketch_count as f64 * readout, + peak_memory_bytes: 0, + scan_bytes: 0, + })?; + let retained = inputs + .active_window_count + .checked_add(inputs.retained_window_count)? + .checked_mul(inputs.physical_sketch_count)? + .checked_mul(inputs.state_bytes_per_sketch)?; + let horizon_seconds = inputs.horizon_ms as f64 / 1_000.0; + let retention_total = self.calibrated(ResourceEstimate { + cpu_ops: 0.0, + peak_memory_bytes: retained, + scan_bytes: 0, + })?; + Some(SummaryMaintenanceLifecycleCostInputs { + build_cost: Some(build), + maintenance_cost_per_update: Some(maintenance), + summary_read_cost: Some(read), + retention_cost_rate: Some(CostRate(retention_total.0 / horizon_seconds)), + // Releasing memory has no modeled CPU or I/O. This is not an + // implicit expiration/rebuild policy; those require an explicit + // SummaryDelete or future authoritative lifecycle evidence. + retirement_cost: Some(Cost::ZERO), + }) + } +} + +impl CostModel for StreamingAnalyticalCostModel { + fn rank_candidates( + &self, + intent: &AggIntent, + candidates: &[SketchAlgorithm], + ) -> Vec { + DefaultCostModel.rank_candidates(intent, candidates) + } + + fn estimate_cost(&self, candidate: &ReplacementSubDAG, target: &TargetSubDAG<'_>) -> f64 { + DefaultCostModel.estimate_cost(candidate, target) + } + + fn summary_maintenance_lifecycle_cost_inputs( + &self, + _summary: &SummaryNode, + ) -> SummaryMaintenanceLifecycleCostInputs { + self.lifecycle_inputs().unwrap_or_default() + } + + fn summary_maintenance_capabilities( + &self, + _summary: &SummaryNode, + ) -> SummaryMaintenanceCapabilities { + self.capabilities + } + + fn raw_query_recompute_cost(&self, _target: &QueryExpr) -> Option { + self.calibrated(estimate_streaming_raw_recompute(self.raw, 1).ok()?) + } +} + #[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] struct SummaryOperationCounts { state_builds: u64, @@ -167,6 +287,49 @@ struct SummaryOperationCounts { subtracts_per_read: u64, deletes_per_update: u64, readouts_per_read: u64, + joins_per_read: u64, +} + +pub fn estimate_streaming_raw_recompute( + evidence: StreamingRawInputEvidence, + evaluation_count: u64, +) -> Result { + for (name, value) in [ + ( + "raw_input_rows_per_evaluation", + evidence.input_rows_per_evaluation, + ), + ( + "raw_input_bytes_per_evaluation", + evidence.input_bytes_per_evaluation, + ), + ("raw_peak_memory_bytes", evidence.peak_memory_bytes), + ("evaluation_count", evaluation_count), + ] { + if value == 0 { + return Err(AnalyticalCostError::MissingOrZero(name)); + } + } + if !evidence.cpu_ops_per_row.is_finite() || evidence.cpu_ops_per_row < 0.0 { + return Err(AnalyticalCostError::InvalidOperationCost( + "raw_cpu_ops_per_row", + evidence.cpu_ops_per_row, + )); + } + let cpu_ops = evidence.input_rows_per_evaluation as f64 + * evidence.cpu_ops_per_row + * evaluation_count as f64; + if !cpu_ops.is_finite() { + return Err(AnalyticalCostError::Overflow); + } + Ok(ResourceEstimate { + cpu_ops, + peak_memory_bytes: evidence.peak_memory_bytes, + scan_bytes: evidence + .source_scan_bytes_per_evaluation + .checked_mul(evaluation_count) + .ok_or(AnalyticalCostError::Overflow)?, + }) } /// Cost a selected incremental deployment without changing its lifecycle @@ -178,6 +341,16 @@ pub fn estimate_incremental_summary_maintenance( guarantee: &SummaryMaintenanceLifecycleGuarantee, inputs: StreamingSummaryInputs, cpu: SummaryOperationCpuEvidence, +) -> Result { + estimate_incremental_summary_maintenance_with_join(root, guarantee, inputs, cpu, None) +} + +pub fn estimate_incremental_summary_maintenance_with_join( + root: &SummaryNode, + guarantee: &SummaryMaintenanceLifecycleGuarantee, + inputs: StreamingSummaryInputs, + cpu: SummaryOperationCpuEvidence, + join: Option, ) -> Result { let inputs = inputs.validate()?; validate_guarantee(guarantee, inputs.data_arrival)?; @@ -207,6 +380,27 @@ pub fn estimate_incremental_summary_maintenance( "readout_cpu_ops", cpu.readout_cpu_ops, )?; + let join_cpu = match (counts.joins_per_read, join) { + (0, _) => 0.0, + (_, Some(evidence)) + if evidence.matched_state_pairs_per_evaluation > 0 + && evidence.cpu_ops_per_matched_pair.is_finite() + && evidence.cpu_ops_per_matched_pair >= 0.0 + && evidence.working_memory_bytes > 0 => + { + evidence.matched_state_pairs_per_evaluation as f64 * evidence.cpu_ops_per_matched_pair + } + (_, Some(evidence)) + if !evidence.cpu_ops_per_matched_pair.is_finite() + || evidence.cpu_ops_per_matched_pair < 0.0 => + { + return Err(AnalyticalCostError::InvalidOperationCost( + "summary_join_cpu_ops_per_matched_pair", + evidence.cpu_ops_per_matched_pair, + )); + } + _ => return Err(AnalyticalCostError::MissingOrStale("summary_join")), + }; let build_inserts = inputs .initial_input_rows @@ -223,7 +417,8 @@ pub fn estimate_incremental_summary_maintenance( + evaluations * counts.merges_per_read as f64 * instances * merge + evaluations * counts.subtracts_per_read as f64 * instances * subtract + update_inserts as f64 * counts.deletes_per_update as f64 * delete - + evaluations * counts.readouts_per_read as f64 * instances * readout; + + evaluations * counts.readouts_per_read as f64 * instances * readout + + evaluations * counts.joins_per_read as f64 * join_cpu; if !cpu_ops.is_finite() { return Err(AnalyticalCostError::Overflow); } @@ -248,6 +443,11 @@ pub fn estimate_incremental_summary_maintenance( } else { 0 }; + let join_bytes = match (counts.joins_per_read, join) { + (0, _) => 0, + (_, Some(evidence)) => evidence.working_memory_bytes, + _ => return Err(AnalyticalCostError::MissingOrStale("summary_join")), + }; let bootstrap_row_buffer = if inputs.initial_input_rows == 0 { 0 } else { @@ -259,6 +459,7 @@ pub fn estimate_incremental_summary_maintenance( cpu_ops, peak_memory_bytes: retained_bytes .checked_add(transient_bytes) + .and_then(|bytes| bytes.checked_add(join_bytes)) .ok_or(AnalyticalCostError::Overflow)? .max(bootstrap_row_buffer), scan_bytes: inputs.initial_source_scan_bytes, @@ -390,8 +591,13 @@ fn count_operations(root: &SummaryNode) -> Result { - return Err(AnalyticalCostError::UnsupportedSummaryOperation("join")); + SummaryExpr::SummaryJoin { outer, inner, .. } => { + counts.joins_per_read = counts + .joins_per_read + .checked_add(1) + .ok_or(AnalyticalCostError::Overflow)?; + visit(outer, seen, counts)?; + visit(inner, seen, counts)?; } } Ok(()) @@ -411,13 +617,22 @@ mod tests { SummaryExpr, SummaryFamilyType, SummaryField, SummaryMaintenanceLifecycle, SummaryMaintenanceLifecycleGuarantee, SummaryMaintenanceMode, SummarySchema, }; - use asap_types::pre_asap::{ColumnRef, QueryExpr, Reduction}; + use asap_types::pre_asap::{ + agg_intent::AggIntent, Column, ColumnRef, DataType, QueryExpr, Reduction, Schema, Source, + }; use asap_types::workload::{ - DataWorkload, Evidence, EvidenceSource, Predictability, Query, QueryRecurrence, - QueryRequirements, QueryTimeScope, Rate, RepeatedDemand, RepetitionInterval, TimeSelection, + DataWorkload, Evidence, EvidenceSource, Predictability, Query, QueryLanguage, + QueryRecurrence, QueryRequirements, QueryTimeScope, QueryWorkload, Rate, RepeatedDemand, + RepeatingEntry, RepetitionInterval, TimeSelection, }; use super::*; + use crate::recurrence::Horizon; + use crate::summary_maintenance_lifecycle::{ + global_selection_with_summary_maintenance_lifecycles, + materialize_with_summary_maintenance_lifecycles, plan_summary_maintenance_lifecycles, + SummaryMaintenanceLifecycleCapabilities, WorkloadDemand, + }; fn physical() -> StreamingPhysicalInputEvidence { StreamingPhysicalInputEvidence { @@ -484,6 +699,228 @@ mod tests { assert_eq!(inputs.evaluation_count, 5); } + #[test] + fn pure_streaming_can_bootstrap_from_an_empty_state() { + let data = DataWorkload { + arrival: DataArrival::ContinuouslyIngesting, + ingestion_rate: Evidence { + value: Some(Rate(2.0)), + source: EvidenceSource::Declared, + observed_at_ms: None, + valid_for_ms: None, + }, + input_cardinality: Evidence { + value: Some(0), + source: EvidenceSource::Declared, + observed_at_ms: None, + valid_for_ms: None, + }, + ..DataWorkload::default() + }; + let mut empty = physical(); + empty.initial_input_bytes = 0; + empty.initial_source_scan_bytes = 0; + let inputs = + StreamingSummaryInputs::from_workload(empty, &data, &query(), 0, 5_000).unwrap(); + let estimate = estimate_incremental_summary_maintenance( + &summary_with_operations(false, false, false), + &continuous_guarantee(), + inputs, + SummaryOperationCpuEvidence { + insert_cpu_ops: Some(2.0), + readout_cpu_ops: Some(0.0), + ..SummaryOperationCpuEvidence::default() + }, + ) + .unwrap(); + assert_eq!(estimate.cpu_ops, 40.0); // 10 arrivals * 2 active windows * 2 ops. + assert_eq!(estimate.scan_bytes, 0); + } + + #[test] + fn bootstrap_rows_and_bytes_must_be_present_together() { + let mut inputs = StreamingSummaryInputs { + data_arrival: DataArrival::ContinuouslyIngesting, + initial_input_rows: 0, + initial_input_bytes: 8, + initial_source_scan_bytes: 0, + ingestion_rate_per_second: 1.0, + planning_time_ms: 0, + horizon_ms: 1_000, + active_window_count: 1, + retained_window_count: 1, + physical_sketch_count: 1, + state_bytes_per_sketch: 8, + evaluation_count: 1, + }; + assert_eq!( + inputs.validate(), + Err(AnalyticalCostError::InconsistentBootstrapEvidence) + ); + inputs.initial_input_rows = 1; + inputs.initial_input_bytes = 0; + assert_eq!( + inputs.validate(), + Err(AnalyticalCostError::InconsistentBootstrapEvidence) + ); + } + + #[test] + fn existing_lifecycle_planner_selects_a_fully_costed_streaming_alternative() { + let inputs = StreamingSummaryInputs { + data_arrival: DataArrival::ContinuouslyIngesting, + initial_input_rows: 10, + initial_input_bytes: 640, + initial_source_scan_bytes: 640, + ingestion_rate_per_second: 2.0, + planning_time_ms: 0, + horizon_ms: 5_000, + active_window_count: 2, + retained_window_count: 3, + physical_sketch_count: 2, + state_bytes_per_sketch: 100, + evaluation_count: 5, + }; + let model = StreamingAnalyticalCostModel { + summary_inputs: inputs, + raw: StreamingRawInputEvidence { + input_rows_per_evaluation: 20, + input_bytes_per_evaluation: 1_280, + source_scan_bytes_per_evaluation: 1_280, + cpu_ops_per_row: 2.0, + peak_memory_bytes: 320, + }, + cpu: SummaryOperationCpuEvidence { + insert_cpu_ops: Some(2.0), + readout_cpu_ops: Some(3.0), + ..SummaryOperationCpuEvidence::default() + }, + calibration: ResourceCalibration { + cost_per_cpu_op: 1.0, + cost_per_scan_byte: 1.0, + cost_per_retained_byte: 1.0, + version: "test".into(), + }, + capabilities: SummaryMaintenanceCapabilities { + incremental_update: true, + merge: false, + delete: false, + }, + }; + let workload = QueryWorkload { + language: QueryLanguage::PromQL, + query_batch: None, + repeating_queries: Some(vec![RepeatingEntry { + query: Query("streaming count".into()), + demand: RepeatedDemand::FixedInterval(RepetitionInterval(1_000)), + requirements: QueryRequirements::default(), + predictability: Predictability::Predictable { known_at: None }, + time_selection: TimeSelection::default(), + }]), + data_workload: Some(DataWorkload { + arrival: DataArrival::ContinuouslyIngesting, + ingestion_rate: Evidence { + value: Some(Rate(2.0)), + source: EvidenceSource::Declared, + observed_at_ms: None, + valid_for_ms: None, + }, + ..DataWorkload::default() + }), + }; + let root = summary_with_operations(false, false, false); + let plan = plan_summary_maintenance_lifecycles( + root, + WorkloadDemand::new(&workload, &[0]), + 0, + Some(Horizon(5.0)), + SummaryMaintenanceLifecycleCapabilities::ALL, + &model, + ) + .unwrap(); + let selected = plan.deployments[0] + .summary_maintenance_lifecycle_guarantee + .as_ref() + .unwrap(); + assert_eq!( + selected.summary_maintenance_mode, + SummaryMaintenanceMode::Incremental + ); + assert!(matches!( + selected.summary_maintenance_lifecycle, + SummaryMaintenanceLifecycle::Shared { .. } + | SummaryMaintenanceLifecycle::ContinuouslyMaintained + )); + assert!(plan.summary_total_cost.is_some()); + assert_eq!( + model.raw_query_recompute_cost(&QueryExpr::promql_scalar(1.0)), + Some(Cost(1_640.0)) + ); + } + + #[test] + fn global_selection_compares_streaming_summary_and_raw_over_one_horizon() { + let target = streaming_sum_query(); + let space = crate::replacement::search_workload(vec![("q", Rc::clone(&target))]); + let workload = streaming_workload(); + let model = streaming_model(); + let selection = global_selection_with_summary_maintenance_lifecycles( + &space, + &workload, + &[0], + 0, + Some(Horizon(5.0)), + SummaryMaintenanceLifecycleCapabilities::ALL, + &model, + ) + .unwrap(); + let plan = materialize_with_summary_maintenance_lifecycles( + &selection, + &space.roots[0].1, + WorkloadDemand::new(&workload, &[0]), + 0, + Some(Horizon(5.0)), + SummaryMaintenanceLifecycleCapabilities::ALL, + &model, + ) + .unwrap() + .unwrap(); + assert!(!plan.selected_raw_recompute); + assert_eq!(plan.raw_recompute_total_cost, Some(Cost(8_200.0))); + + let mut raw_cheaper = model; + raw_cheaper.raw = StreamingRawInputEvidence { + input_rows_per_evaluation: 1, + input_bytes_per_evaluation: 1, + source_scan_bytes_per_evaluation: 0, + cpu_ops_per_row: 0.0, + peak_memory_bytes: 1, + }; + let cheap_selection = global_selection_with_summary_maintenance_lifecycles( + &space, + &workload, + &[0], + 0, + Some(Horizon(5.0)), + SummaryMaintenanceLifecycleCapabilities::ALL, + &raw_cheaper, + ) + .unwrap(); + let cheap_plan = materialize_with_summary_maintenance_lifecycles( + &cheap_selection, + &space.roots[0].1, + WorkloadDemand::new(&workload, &[0]), + 0, + Some(Horizon(5.0)), + SummaryMaintenanceLifecycleCapabilities::ALL, + &raw_cheaper, + ) + .unwrap() + .unwrap(); + assert!(cheap_plan.selected_raw_recompute); + assert_eq!(cheap_plan.raw_recompute_total_cost, Some(Cost(5.0))); + } + #[test] fn mixed_arrival_fails_closed_until_backlog_and_stream_are_separate() { let data = DataWorkload { @@ -766,6 +1203,72 @@ mod tests { assert_eq!(estimate.scan_bytes, 0); } + #[test] + fn raw_recompute_uses_the_same_evaluation_horizon() { + let estimate = estimate_streaming_raw_recompute( + StreamingRawInputEvidence { + input_rows_per_evaluation: 100, + input_bytes_per_evaluation: 6_400, + source_scan_bytes_per_evaluation: 6_400, + cpu_ops_per_row: 2.0, + peak_memory_bytes: 800, + }, + 5, + ) + .unwrap(); + assert_eq!(estimate.cpu_ops, 1_000.0); + assert_eq!(estimate.scan_bytes, 32_000); + assert_eq!(estimate.peak_memory_bytes, 800); + } + + #[test] + fn summary_join_requires_cardinality_and_working_memory_evidence() { + let joined = summary_join(); + let inputs = StreamingSummaryInputs { + data_arrival: DataArrival::ContinuouslyIngesting, + initial_input_rows: 1, + initial_input_bytes: 8, + initial_source_scan_bytes: 8, + ingestion_rate_per_second: 1.0, + planning_time_ms: 0, + horizon_ms: 1_000, + active_window_count: 1, + retained_window_count: 1, + physical_sketch_count: 1, + state_bytes_per_sketch: 8, + evaluation_count: 2, + }; + let cpu = SummaryOperationCpuEvidence { + insert_cpu_ops: Some(1.0), + readout_cpu_ops: Some(1.0), + ..SummaryOperationCpuEvidence::default() + }; + assert_eq!( + estimate_incremental_summary_maintenance_with_join( + &joined, + &continuous_guarantee(), + inputs, + cpu, + None, + ), + Err(AnalyticalCostError::MissingOrStale("summary_join")) + ); + let estimate = estimate_incremental_summary_maintenance_with_join( + &joined, + &continuous_guarantee(), + inputs, + cpu, + Some(SummaryJoinEvidence { + matched_state_pairs_per_evaluation: 3, + cpu_ops_per_matched_pair: 4.0, + working_memory_bytes: 32, + }), + ) + .unwrap(); + assert_eq!(estimate.cpu_ops, 30.0); // 2 inserts + 4 readouts + 24 join ops. + assert_eq!(estimate.peak_memory_bytes, 64); // 4 persistent states + join memory. + } + fn summary_with_operations(merge: bool, subtract: bool, delete: bool) -> Rc { let state_type = SummaryFamilyType::ExactAggregate(ExactKind::Count, ExactParams::Count); let schema = SummarySchema { @@ -834,4 +1337,135 @@ mod tests { guarantee: None, }) } + + fn summary_join() -> Rc { + let left = summary_with_operations(false, false, false); + let right = summary_with_operations(false, false, false); + let SummaryExpr::SummaryEstimate { + summary_input: left, + .. + } = &left.expr + else { + unreachable!() + }; + let SummaryExpr::SummaryEstimate { + summary_input: right, + .. + } = &right.expr + else { + unreachable!() + }; + let schema = left.schema.clone(); + let join = Rc::new(SummaryNode { + expr: SummaryExpr::SummaryJoin { + outer: Rc::clone(left), + inner: Rc::clone(right), + key: ColumnRef::Wildcard, + family: SummaryFamilyType::ExactAggregate(ExactKind::Count, ExactParams::Count), + }, + schema: schema.clone(), + guarantee: None, + }); + Rc::new(SummaryNode { + expr: SummaryExpr::SummaryEstimate { + summary_input: join, + query: asap_types::post_asap::SketchQuery::PointCount { + key: ColumnRef::Wildcard, + value: None, + }, + }, + schema, + guarantee: None, + }) + } + + fn streaming_sum_query() -> Rc { + let scan = Rc::new(QueryExpr::Scan { + source: Source::TimeSeries { + metric: "metrics".into(), + }, + predicates: vec![], + schema: Schema::with_time_index( + vec![ + Column::new("ts", DataType::Timestamp, false), + Column::new("value", DataType::Float64, false), + ], + 0, + vec![], + ), + }); + Rc::new(QueryExpr::Aggregate { + reduction: Reduction::by(vec![]), + measures: vec![AggIntent::Sum { col: None }], + output_names: vec![], + having: None, + child: scan, + }) + } + + fn streaming_workload() -> QueryWorkload { + QueryWorkload { + language: QueryLanguage::PromQL, + query_batch: None, + repeating_queries: Some(vec![RepeatingEntry { + query: Query("sum(metrics)".into()), + demand: RepeatedDemand::FixedInterval(RepetitionInterval(1_000)), + requirements: QueryRequirements::default(), + predictability: Predictability::Predictable { known_at: None }, + time_selection: TimeSelection::default(), + }]), + data_workload: Some(DataWorkload { + arrival: DataArrival::ContinuouslyIngesting, + ingestion_rate: Evidence { + value: Some(Rate(2.0)), + source: EvidenceSource::Declared, + observed_at_ms: None, + valid_for_ms: None, + }, + ..DataWorkload::default() + }), + } + } + + fn streaming_model() -> StreamingAnalyticalCostModel { + StreamingAnalyticalCostModel { + summary_inputs: StreamingSummaryInputs { + data_arrival: DataArrival::ContinuouslyIngesting, + initial_input_rows: 10, + initial_input_bytes: 640, + initial_source_scan_bytes: 640, + ingestion_rate_per_second: 2.0, + planning_time_ms: 0, + horizon_ms: 5_000, + active_window_count: 2, + retained_window_count: 3, + physical_sketch_count: 2, + state_bytes_per_sketch: 100, + evaluation_count: 5, + }, + raw: StreamingRawInputEvidence { + input_rows_per_evaluation: 20, + input_bytes_per_evaluation: 1_280, + source_scan_bytes_per_evaluation: 1_280, + cpu_ops_per_row: 2.0, + peak_memory_bytes: 320, + }, + cpu: SummaryOperationCpuEvidence { + insert_cpu_ops: Some(2.0), + readout_cpu_ops: Some(3.0), + ..SummaryOperationCpuEvidence::default() + }, + calibration: ResourceCalibration { + cost_per_cpu_op: 1.0, + cost_per_scan_byte: 1.0, + cost_per_retained_byte: 1.0, + version: "test".into(), + }, + capabilities: SummaryMaintenanceCapabilities { + incremental_update: true, + merge: false, + delete: false, + }, + } + } } From d4e9007ef03afd78806dcd8fae0ac6b46511f144 Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 3 Sep 2026 18:07:37 -0600 Subject: [PATCH 2/4] fix(cost): compare complete streaming alternatives --- .../src/summary_maintenance_cost.rs | 197 +++++------------- 1 file changed, 54 insertions(+), 143 deletions(-) diff --git a/crates/asap-aware-mapping/src/summary_maintenance_cost.rs b/crates/asap-aware-mapping/src/summary_maintenance_cost.rs index 55cba088..89d21c42 100644 --- a/crates/asap-aware-mapping/src/summary_maintenance_cost.rs +++ b/crates/asap-aware-mapping/src/summary_maintenance_cost.rs @@ -167,25 +167,14 @@ pub struct SummaryOperationCpuEvidence { pub readout_cpu_ops: Option, } -/// Physical evidence for one `SummaryJoin` implementation. Cardinality and -/// working memory cannot be inferred from the logical join key alone. -#[derive(Debug, Clone, Copy, PartialEq, Serialize, Deserialize)] -pub struct SummaryJoinEvidence { - pub matched_state_pairs_per_evaluation: u64, - pub cpu_ops_per_matched_pair: f64, - pub working_memory_bytes: u64, -} - /// Complete physical work for one raw evaluation. It is deliberately /// per-evaluation so the same normalized query recurrence/horizon can multiply -/// both raw and summary alternatives. +/// both raw and summary alternatives. The estimate comes from the complete raw +/// physical DAG; it is not reconstructed here from a query-shape-specific +/// rows-times-CPU shortcut. #[derive(Debug, Clone, Copy, PartialEq, Serialize, Deserialize)] pub struct StreamingRawInputEvidence { - pub input_rows_per_evaluation: u64, - pub input_bytes_per_evaluation: u64, - pub source_scan_bytes_per_evaluation: u64, - pub cpu_ops_per_row: f64, - pub peak_memory_bytes: u64, + pub per_evaluation: ResourceEstimate, } /// Adapter that supplies the existing lifecycle planner with analytical @@ -210,7 +199,7 @@ impl StreamingAnalyticalCostModel { let insert = required_cpu("insert_cpu_ops", self.cpu.insert_cpu_ops).ok()?; let readout = required_cpu("readout_cpu_ops", self.cpu.readout_cpu_ops).ok()?; let build = self.calibrated(ResourceEstimate { - cpu_ops: inputs.initial_input_rows as f64 * insert, + cpu_ops: inputs.initial_input_rows as f64 * inputs.active_window_count as f64 * insert, peak_memory_bytes: 0, scan_bytes: inputs.initial_source_scan_bytes, })?; @@ -220,15 +209,15 @@ impl StreamingAnalyticalCostModel { scan_bytes: 0, })?; let read = self.calibrated(ResourceEstimate { - cpu_ops: inputs.physical_sketch_count as f64 * readout, + cpu_ops: inputs.physical_summary_count as f64 * readout, peak_memory_bytes: 0, scan_bytes: 0, })?; let retained = inputs .active_window_count .checked_add(inputs.retained_window_count)? - .checked_mul(inputs.physical_sketch_count)? - .checked_mul(inputs.state_bytes_per_sketch)?; + .checked_mul(inputs.physical_summary_count)? + .checked_mul(inputs.state_bytes_per_summary)?; let horizon_seconds = inputs.horizon_ms as f64 / 1_000.0; let retention_total = self.calibrated(ResourceEstimate { cpu_ops: 0.0, @@ -287,46 +276,31 @@ struct SummaryOperationCounts { subtracts_per_read: u64, deletes_per_update: u64, readouts_per_read: u64, - joins_per_read: u64, } pub fn estimate_streaming_raw_recompute( evidence: StreamingRawInputEvidence, evaluation_count: u64, ) -> Result { - for (name, value) in [ - ( - "raw_input_rows_per_evaluation", - evidence.input_rows_per_evaluation, - ), - ( - "raw_input_bytes_per_evaluation", - evidence.input_bytes_per_evaluation, - ), - ("raw_peak_memory_bytes", evidence.peak_memory_bytes), - ("evaluation_count", evaluation_count), - ] { - if value == 0 { - return Err(AnalyticalCostError::MissingOrZero(name)); - } + if evaluation_count == 0 { + return Err(AnalyticalCostError::MissingOrZero("evaluation_count")); } - if !evidence.cpu_ops_per_row.is_finite() || evidence.cpu_ops_per_row < 0.0 { + if !evidence.per_evaluation.cpu_ops.is_finite() || evidence.per_evaluation.cpu_ops < 0.0 { return Err(AnalyticalCostError::InvalidOperationCost( - "raw_cpu_ops_per_row", - evidence.cpu_ops_per_row, + "raw_cpu_ops_per_evaluation", + evidence.per_evaluation.cpu_ops, )); } - let cpu_ops = evidence.input_rows_per_evaluation as f64 - * evidence.cpu_ops_per_row - * evaluation_count as f64; + let cpu_ops = evidence.per_evaluation.cpu_ops * evaluation_count as f64; if !cpu_ops.is_finite() { return Err(AnalyticalCostError::Overflow); } Ok(ResourceEstimate { cpu_ops, - peak_memory_bytes: evidence.peak_memory_bytes, + peak_memory_bytes: evidence.per_evaluation.peak_memory_bytes, scan_bytes: evidence - .source_scan_bytes_per_evaluation + .per_evaluation + .scan_bytes .checked_mul(evaluation_count) .ok_or(AnalyticalCostError::Overflow)?, }) @@ -341,16 +315,6 @@ pub fn estimate_incremental_summary_maintenance( guarantee: &SummaryMaintenanceLifecycleGuarantee, inputs: StreamingSummaryInputs, cpu: SummaryOperationCpuEvidence, -) -> Result { - estimate_incremental_summary_maintenance_with_join(root, guarantee, inputs, cpu, None) -} - -pub fn estimate_incremental_summary_maintenance_with_join( - root: &SummaryNode, - guarantee: &SummaryMaintenanceLifecycleGuarantee, - inputs: StreamingSummaryInputs, - cpu: SummaryOperationCpuEvidence, - join: Option, ) -> Result { let inputs = inputs.validate()?; validate_guarantee(guarantee, inputs.data_arrival)?; @@ -380,28 +344,6 @@ pub fn estimate_incremental_summary_maintenance_with_join( "readout_cpu_ops", cpu.readout_cpu_ops, )?; - let join_cpu = match (counts.joins_per_read, join) { - (0, _) => 0.0, - (_, Some(evidence)) - if evidence.matched_state_pairs_per_evaluation > 0 - && evidence.cpu_ops_per_matched_pair.is_finite() - && evidence.cpu_ops_per_matched_pair >= 0.0 - && evidence.working_memory_bytes > 0 => - { - evidence.matched_state_pairs_per_evaluation as f64 * evidence.cpu_ops_per_matched_pair - } - (_, Some(evidence)) - if !evidence.cpu_ops_per_matched_pair.is_finite() - || evidence.cpu_ops_per_matched_pair < 0.0 => - { - return Err(AnalyticalCostError::InvalidOperationCost( - "summary_join_cpu_ops_per_matched_pair", - evidence.cpu_ops_per_matched_pair, - )); - } - _ => return Err(AnalyticalCostError::MissingOrStale("summary_join")), - }; - let build_inserts = inputs .initial_input_rows .checked_mul(inputs.active_window_count) @@ -417,8 +359,7 @@ pub fn estimate_incremental_summary_maintenance_with_join( + evaluations * counts.merges_per_read as f64 * instances * merge + evaluations * counts.subtracts_per_read as f64 * instances * subtract + update_inserts as f64 * counts.deletes_per_update as f64 * delete - + evaluations * counts.readouts_per_read as f64 * instances * readout - + evaluations * counts.joins_per_read as f64 * join_cpu; + + evaluations * counts.readouts_per_read as f64 * instances * readout; if !cpu_ops.is_finite() { return Err(AnalyticalCostError::Overflow); } @@ -443,11 +384,6 @@ pub fn estimate_incremental_summary_maintenance_with_join( } else { 0 }; - let join_bytes = match (counts.joins_per_read, join) { - (0, _) => 0, - (_, Some(evidence)) => evidence.working_memory_bytes, - _ => return Err(AnalyticalCostError::MissingOrStale("summary_join")), - }; let bootstrap_row_buffer = if inputs.initial_input_rows == 0 { 0 } else { @@ -459,7 +395,6 @@ pub fn estimate_incremental_summary_maintenance_with_join( cpu_ops, peak_memory_bytes: retained_bytes .checked_add(transient_bytes) - .and_then(|bytes| bytes.checked_add(join_bytes)) .ok_or(AnalyticalCostError::Overflow)? .max(bootstrap_row_buffer), scan_bytes: inputs.initial_source_scan_bytes, @@ -591,13 +526,8 @@ fn count_operations(root: &SummaryNode) -> Result { - counts.joins_per_read = counts - .joins_per_read - .checked_add(1) - .ok_or(AnalyticalCostError::Overflow)?; - visit(outer, seen, counts)?; - visit(inner, seen, counts)?; + SummaryExpr::SummaryJoin { .. } => { + return Err(AnalyticalCostError::UnsupportedSummaryOperation("join")); } } Ok(()) @@ -728,12 +658,13 @@ mod tests { inputs, SummaryOperationCpuEvidence { insert_cpu_ops: Some(2.0), - readout_cpu_ops: Some(0.0), + readout_cpu_ops: Some(1.0), ..SummaryOperationCpuEvidence::default() }, ) .unwrap(); - assert_eq!(estimate.cpu_ops, 40.0); // 10 arrivals * 2 active windows * 2 ops. + // 10 arrivals * 2 active windows * 2 insert ops + 5 evaluations * 2 summaries. + assert_eq!(estimate.cpu_ops, 50.0); assert_eq!(estimate.scan_bytes, 0); } @@ -749,8 +680,8 @@ mod tests { horizon_ms: 1_000, active_window_count: 1, retained_window_count: 1, - physical_sketch_count: 1, - state_bytes_per_sketch: 8, + physical_summary_count: 1, + state_bytes_per_summary: 8, evaluation_count: 1, }; assert_eq!( @@ -777,18 +708,18 @@ mod tests { horizon_ms: 5_000, active_window_count: 2, retained_window_count: 3, - physical_sketch_count: 2, - state_bytes_per_sketch: 100, + physical_summary_count: 2, + state_bytes_per_summary: 100, evaluation_count: 5, }; let model = StreamingAnalyticalCostModel { summary_inputs: inputs, raw: StreamingRawInputEvidence { - input_rows_per_evaluation: 20, - input_bytes_per_evaluation: 1_280, - source_scan_bytes_per_evaluation: 1_280, - cpu_ops_per_row: 2.0, - peak_memory_bytes: 320, + per_evaluation: ResourceEstimate { + cpu_ops: 40.0, + peak_memory_bytes: 320, + scan_bytes: 1_280, + }, }, cpu: SummaryOperationCpuEvidence { insert_cpu_ops: Some(2.0), @@ -890,11 +821,11 @@ mod tests { let mut raw_cheaper = model; raw_cheaper.raw = StreamingRawInputEvidence { - input_rows_per_evaluation: 1, - input_bytes_per_evaluation: 1, - source_scan_bytes_per_evaluation: 0, - cpu_ops_per_row: 0.0, - peak_memory_bytes: 1, + per_evaluation: ResourceEstimate { + cpu_ops: 0.0, + peak_memory_bytes: 1, + scan_bytes: 0, + }, }; let cheap_selection = global_selection_with_summary_maintenance_lifecycles( &space, @@ -1207,11 +1138,11 @@ mod tests { fn raw_recompute_uses_the_same_evaluation_horizon() { let estimate = estimate_streaming_raw_recompute( StreamingRawInputEvidence { - input_rows_per_evaluation: 100, - input_bytes_per_evaluation: 6_400, - source_scan_bytes_per_evaluation: 6_400, - cpu_ops_per_row: 2.0, - peak_memory_bytes: 800, + per_evaluation: ResourceEstimate { + cpu_ops: 200.0, + peak_memory_bytes: 800, + scan_bytes: 6_400, + }, }, 5, ) @@ -1222,7 +1153,7 @@ mod tests { } #[test] - fn summary_join_requires_cardinality_and_working_memory_evidence() { + fn flat_summary_evidence_rejects_summary_join() { let joined = summary_join(); let inputs = StreamingSummaryInputs { data_arrival: DataArrival::ContinuouslyIngesting, @@ -1234,8 +1165,8 @@ mod tests { horizon_ms: 1_000, active_window_count: 1, retained_window_count: 1, - physical_sketch_count: 1, - state_bytes_per_sketch: 8, + physical_summary_count: 1, + state_bytes_per_summary: 8, evaluation_count: 2, }; let cpu = SummaryOperationCpuEvidence { @@ -1244,29 +1175,9 @@ mod tests { ..SummaryOperationCpuEvidence::default() }; assert_eq!( - estimate_incremental_summary_maintenance_with_join( - &joined, - &continuous_guarantee(), - inputs, - cpu, - None, - ), - Err(AnalyticalCostError::MissingOrStale("summary_join")) + estimate_incremental_summary_maintenance(&joined, &continuous_guarantee(), inputs, cpu,), + Err(AnalyticalCostError::UnsupportedSummaryOperation("join")) ); - let estimate = estimate_incremental_summary_maintenance_with_join( - &joined, - &continuous_guarantee(), - inputs, - cpu, - Some(SummaryJoinEvidence { - matched_state_pairs_per_evaluation: 3, - cpu_ops_per_matched_pair: 4.0, - working_memory_bytes: 32, - }), - ) - .unwrap(); - assert_eq!(estimate.cpu_ops, 30.0); // 2 inserts + 4 readouts + 24 join ops. - assert_eq!(estimate.peak_memory_bytes, 64); // 4 persistent states + join memory. } fn summary_with_operations(merge: bool, subtract: bool, delete: bool) -> Rc { @@ -1439,16 +1350,16 @@ mod tests { horizon_ms: 5_000, active_window_count: 2, retained_window_count: 3, - physical_sketch_count: 2, - state_bytes_per_sketch: 100, + physical_summary_count: 2, + state_bytes_per_summary: 100, evaluation_count: 5, }, raw: StreamingRawInputEvidence { - input_rows_per_evaluation: 20, - input_bytes_per_evaluation: 1_280, - source_scan_bytes_per_evaluation: 1_280, - cpu_ops_per_row: 2.0, - peak_memory_bytes: 320, + per_evaluation: ResourceEstimate { + cpu_ops: 40.0, + peak_memory_bytes: 320, + scan_bytes: 1_280, + }, }, cpu: SummaryOperationCpuEvidence { insert_cpu_ops: Some(2.0), From 103dc876434ef391f3b086a08b3b8b7144c15535 Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 3 Sep 2026 18:11:57 -0600 Subject: [PATCH 3/4] fix(cost): bind lifecycle evidence to physical identities --- .../src/summary_maintenance_cost.rs | 67 ++++++++++++++++--- .../analytical-resource-cost.md | 30 +++++++++ 2 files changed, 87 insertions(+), 10 deletions(-) diff --git a/crates/asap-aware-mapping/src/summary_maintenance_cost.rs b/crates/asap-aware-mapping/src/summary_maintenance_cost.rs index 89d21c42..fbbf1776 100644 --- a/crates/asap-aware-mapping/src/summary_maintenance_cost.rs +++ b/crates/asap-aware-mapping/src/summary_maintenance_cost.rs @@ -5,6 +5,7 @@ //! window counts, and per-operation CPU measurements or complexity estimates. use std::collections::HashSet; +use std::rc::Rc; use asap_types::post_asap::{ SketchAlgorithm, SummaryExpr, SummaryMaintenanceLifecycle, @@ -182,6 +183,10 @@ pub struct StreamingRawInputEvidence { /// existing enums and legality checks remain authoritative. #[derive(Debug, Clone)] pub struct StreamingAnalyticalCostModel { + /// Exact logical summary identity to which the flat evidence applies. + pub summary: Rc, + /// Exact raw target identity to which `raw` applies. + pub raw_target: Rc, pub summary_inputs: StreamingSummaryInputs, pub raw: StreamingRawInputEvidence, pub cpu: SummaryOperationCpuEvidence, @@ -252,19 +257,29 @@ impl CostModel for StreamingAnalyticalCostModel { fn summary_maintenance_lifecycle_cost_inputs( &self, - _summary: &SummaryNode, + summary: &SummaryNode, ) -> SummaryMaintenanceLifecycleCostInputs { - self.lifecycle_inputs().unwrap_or_default() + std::ptr::eq(self.summary.as_ref(), summary) + .then(|| self.lifecycle_inputs()) + .flatten() + .unwrap_or_default() } fn summary_maintenance_capabilities( &self, - _summary: &SummaryNode, + summary: &SummaryNode, ) -> SummaryMaintenanceCapabilities { - self.capabilities + if std::ptr::eq(self.summary.as_ref(), summary) { + self.capabilities + } else { + SummaryMaintenanceCapabilities::default() + } } - fn raw_query_recompute_cost(&self, _target: &QueryExpr) -> Option { + fn raw_query_recompute_cost(&self, target: &QueryExpr) -> Option { + if !std::ptr::eq(self.raw_target.as_ref(), target) { + return None; + } self.calibrated(estimate_streaming_raw_recompute(self.raw, 1).ok()?) } } @@ -698,6 +713,8 @@ mod tests { #[test] fn existing_lifecycle_planner_selects_a_fully_costed_streaming_alternative() { + let root = summary_with_operations(false, false, false); + let raw_target = Rc::new(QueryExpr::promql_scalar(1.0)); let inputs = StreamingSummaryInputs { data_arrival: DataArrival::ContinuouslyIngesting, initial_input_rows: 10, @@ -713,6 +730,8 @@ mod tests { evaluation_count: 5, }; let model = StreamingAnalyticalCostModel { + summary: single_summary_agg(&root), + raw_target: Rc::clone(&raw_target), summary_inputs: inputs, raw: StreamingRawInputEvidence { per_evaluation: ResourceEstimate { @@ -759,9 +778,8 @@ mod tests { ..DataWorkload::default() }), }; - let root = summary_with_operations(false, false, false); let plan = plan_summary_maintenance_lifecycles( - root, + Rc::clone(&root), WorkloadDemand::new(&workload, &[0]), 0, Some(Horizon(5.0)), @@ -784,9 +802,14 @@ mod tests { )); assert!(plan.summary_total_cost.is_some()); assert_eq!( - model.raw_query_recompute_cost(&QueryExpr::promql_scalar(1.0)), + model.raw_query_recompute_cost(&raw_target), Some(Cost(1_640.0)) ); + assert_eq!( + model.raw_query_recompute_cost(&QueryExpr::promql_scalar(1.0)), + None, + "structurally equal but independently allocated targets must not reuse evidence" + ); } #[test] @@ -794,7 +817,18 @@ mod tests { let target = streaming_sum_query(); let space = crate::replacement::search_workload(vec![("q", Rc::clone(&target))]); let workload = streaming_workload(); - let model = streaming_model(); + let raw_target = Rc::clone(&space.roots[0].1); + let summary = space + .group_for(&raw_target) + .unwrap() + .candidates + .iter() + .find_map(|candidate| match &candidate.replacement { + crate::replacement::Replacement::Summary(summary) => Some(Rc::clone(summary)), + crate::replacement::Replacement::Rewrite(_) => None, + }) + .unwrap(); + let model = streaming_model(single_summary_agg(&summary), raw_target); let selection = global_selection_with_summary_maintenance_lifecycles( &space, &workload, @@ -1338,8 +1372,13 @@ mod tests { } } - fn streaming_model() -> StreamingAnalyticalCostModel { + fn streaming_model( + summary: Rc, + raw_target: Rc, + ) -> StreamingAnalyticalCostModel { StreamingAnalyticalCostModel { + summary, + raw_target, summary_inputs: StreamingSummaryInputs { data_arrival: DataArrival::ContinuouslyIngesting, initial_input_rows: 10, @@ -1379,4 +1418,12 @@ mod tests { }, } } + + fn single_summary_agg(root: &Rc) -> Rc { + match &root.expr { + SummaryExpr::SummaryAgg { .. } => Rc::clone(root), + SummaryExpr::SummaryEstimate { summary_input, .. } => single_summary_agg(summary_input), + other => panic!("expected one SummaryAgg, got {other:?}"), + } + } } diff --git a/docs/design_docs/asap-aware-mapping/analytical-resource-cost.md b/docs/design_docs/asap-aware-mapping/analytical-resource-cost.md index dceeb4cd..f6379950 100644 --- a/docs/design_docs/asap-aware-mapping/analytical-resource-cost.md +++ b/docs/design_docs/asap-aware-mapping/analytical-resource-cost.md @@ -608,6 +608,36 @@ horizon, so the estimator does not require `retention >= horizon`. Lifecycle legality and query time-coverage checks establish whether the retained window can answer the query. +### Comparing single-summary lifecycle alternatives + +For one logical `SummaryAgg`, the analytical lifecycle adapter converts the +same physical evidence into the existing lifecycle planner's five cost terms: + +| Lifecycle term | Resource basis | +|---|---| +| Initial build | Bootstrap rows routed to every bootstrap-active window, plus the bootstrap source read. | +| Maintenance per update | One arriving row routed to every currently active window. | +| Summary read | Readout of every physical summary instance needed by one query evaluation. | +| Retention rate | All active and retained state bytes calibrated over the finite comparison horizon. | +| Retirement | Zero only for releasing modeled memory; an actual delete, expiration, or rebuild requires explicit operation evidence. | + +The existing lifecycle model—not this adapter—enumerates `Ephemeral`, +`Prepared`, `Shared`, and `ContinuouslyMaintained`, checks workload and runtime +legality, and multiplies per-update and per-read terms by the normalized +workload rates. Missing any required term leaves that alternative unavailable. + +The raw side is supplied as a complete `ResourceEstimate` for one execution of +the raw physical DAG. The lifecycle planner applies the same recurrence and +horizon. This deliberately avoids reconstructing raw work with a special-case +`input_rows × cpu_per_row` formula that would omit joins, windows, sorts, or +other operators. + +Flat single-summary evidence is bound to the exact `SummaryNode` and raw +`QueryExpr` identities for which it was produced. It cannot be reused for a +structurally similar node or for multiple summary states. A complete +multi-summary `SummaryExpr` DAG requires per-node physical evidence and +physical-identity deduplication. + A retained summary bootstraps every active window, consumes arriving rows, and serves later reads from state. For the simple build/update/readout shape: From b4c5480f9ae9b746c668533a9bca5fd3d899481d Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 3 Sep 2026 20:57:09 -0600 Subject: [PATCH 4/4] refactor(cost): name summary maintenance model --- .../src/summary_maintenance_cost.rs | 16 +++++++--------- 1 file changed, 7 insertions(+), 9 deletions(-) diff --git a/crates/asap-aware-mapping/src/summary_maintenance_cost.rs b/crates/asap-aware-mapping/src/summary_maintenance_cost.rs index fbbf1776..5dbee68a 100644 --- a/crates/asap-aware-mapping/src/summary_maintenance_cost.rs +++ b/crates/asap-aware-mapping/src/summary_maintenance_cost.rs @@ -15,9 +15,7 @@ use asap_types::pre_asap::{agg_intent::AggIntent, QueryExpr}; use asap_types::workload::{DataArrival, DataWorkload, QueryWorkloadEntry}; use serde::{Deserialize, Serialize}; -use crate::analytical_cost::{ - AnalyticalCostError, ResourceCalibration, ResourceEstimate, -}; +use crate::analytical_cost::{AnalyticalCostError, ResourceCalibration, ResourceEstimate}; use crate::cost_model::{Cost, CostModel, DefaultCostModel}; use crate::physical_operator_statistics::evaluations_in_horizon; use crate::recurrence::CostRate; @@ -182,7 +180,7 @@ pub struct StreamingRawInputEvidence { /// streaming costs. It does not define lifecycle policy: the planner's /// existing enums and legality checks remain authoritative. #[derive(Debug, Clone)] -pub struct StreamingAnalyticalCostModel { +pub struct SummaryMaintenanceCostModel { /// Exact logical summary identity to which the flat evidence applies. pub summary: Rc, /// Exact raw target identity to which `raw` applies. @@ -194,7 +192,7 @@ pub struct StreamingAnalyticalCostModel { pub capabilities: SummaryMaintenanceCapabilities, } -impl StreamingAnalyticalCostModel { +impl SummaryMaintenanceCostModel { fn calibrated(&self, estimate: ResourceEstimate) -> Option { estimate.calibrated_cost(&self.calibration).ok().map(Cost) } @@ -242,7 +240,7 @@ impl StreamingAnalyticalCostModel { } } -impl CostModel for StreamingAnalyticalCostModel { +impl CostModel for SummaryMaintenanceCostModel { fn rank_candidates( &self, intent: &AggIntent, @@ -729,7 +727,7 @@ mod tests { state_bytes_per_summary: 100, evaluation_count: 5, }; - let model = StreamingAnalyticalCostModel { + let model = SummaryMaintenanceCostModel { summary: single_summary_agg(&root), raw_target: Rc::clone(&raw_target), summary_inputs: inputs, @@ -1375,8 +1373,8 @@ mod tests { fn streaming_model( summary: Rc, raw_target: Rc, - ) -> StreamingAnalyticalCostModel { - StreamingAnalyticalCostModel { + ) -> SummaryMaintenanceCostModel { + SummaryMaintenanceCostModel { summary, raw_target, summary_inputs: StreamingSummaryInputs {