From 03a6ff680f87dd17d2c11bd227dba67d8ec5f294 Mon Sep 17 00:00:00 2001 From: zz_y Date: Wed, 2 Sep 2026 08:30:21 -0600 Subject: [PATCH 1/6] feat(cost): model incremental summary maintenance --- .../asap-aware-mapping/src/analytical_cost.rs | 8 +- .../src/analytical_streaming_cost.rs | 792 ++++++++++++++++++ crates/asap-aware-mapping/src/lib.rs | 1 + .../src/summary_maintenance_lifecycle.rs | 4 +- 4 files changed, 802 insertions(+), 3 deletions(-) create mode 100644 crates/asap-aware-mapping/src/analytical_streaming_cost.rs diff --git a/crates/asap-aware-mapping/src/analytical_cost.rs b/crates/asap-aware-mapping/src/analytical_cost.rs index 4b766082..d10fdacf 100644 --- a/crates/asap-aware-mapping/src/analytical_cost.rs +++ b/crates/asap-aware-mapping/src/analytical_cost.rs @@ -1940,8 +1940,14 @@ impl ResourceEstimate { #[derive(Debug, Clone, PartialEq, thiserror::Error)] pub enum AnalyticalCostError { - #[error("analytical resource model v1 supports only DataArrival::AtRest, got {0:?}")] + #[error("data arrival {0:?} is unsupported by the selected analytical adapter")] UnsupportedDataArrival(DataArrival), + #[error("ingestion rate must be finite and non-negative, got {0}")] + InvalidIngestionRate(f64), + #[error("summary lifecycle, maintenance mode, and evaluation schedule are inconsistent")] + IncompatibleLifecycleGuarantee, + #[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")] MissingOrZero(&'static str), #[error("required analytical evidence {0} is missing or stale")] diff --git a/crates/asap-aware-mapping/src/analytical_streaming_cost.rs b/crates/asap-aware-mapping/src/analytical_streaming_cost.rs new file mode 100644 index 00000000..973e67a2 --- /dev/null +++ b/crates/asap-aware-mapping/src/analytical_streaming_cost.rs @@ -0,0 +1,792 @@ +//! Analytical resource cost for incrementally maintained summary deployments. +//! +//! The canonical workload and lifecycle types own deployment semantics. This +//! module only adds physical evidence absent from those schemas: state size, +//! window counts, and per-operation CPU measurements or complexity estimates. + +use std::collections::HashSet; + +use asap_types::post_asap::{ + SummaryExpr, SummaryMaintenanceLifecycle, SummaryMaintenanceLifecycleGuarantee, + SummaryMaintenanceMode, SummaryNode, +}; +use asap_types::workload::{DataArrival, DataWorkload, QueryWorkloadEntry}; +use serde::{Deserialize, Serialize}; + +use crate::analytical_cost::{evaluations_in_horizon, AnalyticalCostError, ResourceEstimate}; +use crate::summary_maintenance_lifecycle::{evaluation_schedule, maintenance_mode}; + +pub const ANALYTICAL_STREAMING_MODEL_VERSION: &str = "analytical-summary-incremental-v1"; + +/// Physical evidence that is not represented by [`DataWorkload`] for one +/// incrementally maintained summary deployment. Window counts describe the +/// already-selected physical deployment; this layer does not define another +/// tumbling/sliding policy enum. +#[derive(Debug, Clone, Copy, PartialEq, Serialize, Deserialize)] +pub struct StreamingPhysicalInputEvidence { + /// Logical bytes in the snapshot used to bootstrap the state. + pub initial_input_bytes: u64, + /// Source bytes read while bootstrapping. Arriving stream bytes are not a + /// disk scan and are therefore excluded. + pub initial_source_scan_bytes: u64, + /// Simultaneously open windows receiving each arriving item. + pub active_window_count: u64, + /// Completed windows retained for query coverage. + pub retained_window_count: u64, + /// Independent state instances per window: one for shared + /// multi-subpopulation state, otherwise the resolved group count. + pub physical_sketch_count: u64, + /// Resident bytes of one concrete state instance. + pub state_bytes_per_sketch: u64, +} + +/// Workload-normalized inputs for incremental maintenance over one finite +/// comparison horizon. +#[derive(Debug, Clone, Copy, PartialEq, Serialize, Deserialize)] +pub struct StreamingSummaryInputs { + pub data_arrival: DataArrival, + pub initial_input_rows: u64, + pub initial_input_bytes: u64, + pub initial_source_scan_bytes: u64, + pub ingestion_rate_per_second: f64, + pub planning_time_ms: u64, + pub horizon_ms: u64, + pub active_window_count: u64, + pub retained_window_count: u64, + pub physical_sketch_count: u64, + pub state_bytes_per_sketch: u64, + pub evaluation_count: u64, +} + +impl StreamingSummaryInputs { + /// Resolve snapshot size, arriving rows, and reads from the canonical + /// workload. Positive fractional expected work rounds up conservatively. + /// + /// `Mixed` fails closed because today's workload schema cannot distinguish + /// its at-rest backlog from its continuing-arrival cardinality. + pub fn from_workload( + physical: StreamingPhysicalInputEvidence, + data: &DataWorkload, + query: &QueryWorkloadEntry, + planning_time_ms: u64, + horizon_ms: u64, + ) -> Result { + if data.arrival != DataArrival::ContinuouslyIngesting { + return Err(AnalyticalCostError::UnsupportedDataArrival(data.arrival)); + } + let initial_input_rows = data + .input_cardinality + .value_at(planning_time_ms) + .copied() + .ok_or(AnalyticalCostError::MissingOrStale("input_cardinality"))?; + let ingestion_rate = data + .ingestion_rate + .value_at(planning_time_ms) + .copied() + .ok_or(AnalyticalCostError::MissingOrStale("ingestion_rate"))?; + if !ingestion_rate.0.is_finite() || ingestion_rate.0 < 0.0 { + return Err(AnalyticalCostError::InvalidIngestionRate(ingestion_rate.0)); + } + if horizon_ms == 0 { + return Err(AnalyticalCostError::MissingOrZero("horizon_ms")); + } + Self { + data_arrival: data.arrival, + initial_input_rows, + initial_input_bytes: physical.initial_input_bytes, + initial_source_scan_bytes: physical.initial_source_scan_bytes, + ingestion_rate_per_second: ingestion_rate.0, + planning_time_ms, + horizon_ms, + active_window_count: physical.active_window_count, + retained_window_count: physical.retained_window_count, + physical_sketch_count: physical.physical_sketch_count, + state_bytes_per_sketch: physical.state_bytes_per_sketch, + evaluation_count: evaluations_in_horizon( + &query.recurrence, + planning_time_ms, + horizon_ms, + )?, + } + .validate() + } + + pub fn validate(self) -> Result { + if self.data_arrival != DataArrival::ContinuouslyIngesting { + return Err(AnalyticalCostError::UnsupportedDataArrival( + self.data_arrival, + )); + } + for (name, value) in [ + ("initial_input_rows", self.initial_input_rows), + ("initial_input_bytes", self.initial_input_bytes), + ("horizon_ms", self.horizon_ms), + ("active_window_count", self.active_window_count), + ("retained_window_count", self.retained_window_count), + ("physical_sketch_count", self.physical_sketch_count), + ("state_bytes_per_sketch", self.state_bytes_per_sketch), + ("evaluation_count", self.evaluation_count), + ] { + if value == 0 { + return Err(AnalyticalCostError::MissingOrZero(name)); + } + } + if !self.ingestion_rate_per_second.is_finite() || self.ingestion_rate_per_second < 0.0 { + return Err(AnalyticalCostError::InvalidIngestionRate( + self.ingestion_rate_per_second, + )); + } + Ok(self) + } +} + +/// CPU operations for one concrete state operation on one state instance. +/// Missing evidence is legal only when the selected summary DAG does not use +/// that operation. +#[derive(Debug, Clone, Copy, Default, PartialEq, Serialize, Deserialize)] +pub struct SummaryOperationCpuEvidence { + pub insert_cpu_ops: Option, + pub merge_cpu_ops: Option, + pub subtract_cpu_ops: Option, + pub delete_cpu_ops: Option, + pub readout_cpu_ops: Option, +} + +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +struct SummaryOperationCounts { + state_builds: u64, + merges_per_read: u64, + subtracts_per_read: u64, + deletes_per_update: u64, + readouts_per_read: u64, +} + +/// Cost a selected incremental deployment without changing its lifecycle +/// decision. Shared `Rc` nodes are visited once, so shared state is built and +/// retained once. Summary merge/subtract/readout run per query evaluation; +/// summary delete runs per arriving update. +pub fn estimate_incremental_summary_maintenance( + root: &SummaryNode, + guarantee: &SummaryMaintenanceLifecycleGuarantee, + inputs: StreamingSummaryInputs, + cpu: SummaryOperationCpuEvidence, +) -> Result { + let inputs = inputs.validate()?; + validate_guarantee(guarantee, inputs.data_arrival)?; + let arriving_input_rows = arriving_rows_for_lifecycle(inputs, guarantee)?; + let counts = count_operations(root)?; + if counts.state_builds == 0 { + return Err(AnalyticalCostError::UnsupportedCandidate); + } + + let insert = required_cpu("insert_cpu_ops", cpu.insert_cpu_ops)?; + let merge = required_cpu_when(counts.merges_per_read, "merge_cpu_ops", cpu.merge_cpu_ops)?; + let subtract = required_cpu_when( + counts.subtracts_per_read, + "subtract_cpu_ops", + cpu.subtract_cpu_ops, + )?; + let delete = required_cpu_when( + counts.deletes_per_update, + "delete_cpu_ops", + cpu.delete_cpu_ops, + )?; + let readout = required_cpu_when( + counts.readouts_per_read, + "readout_cpu_ops", + cpu.readout_cpu_ops, + )?; + + let build_inserts = inputs + .initial_input_rows + .checked_mul(counts.state_builds) + .ok_or(AnalyticalCostError::Overflow)?; + let update_inserts = arriving_input_rows + .checked_mul(inputs.active_window_count) + .and_then(|n| n.checked_mul(counts.state_builds)) + .ok_or(AnalyticalCostError::Overflow)?; + let instances = inputs.physical_sketch_count as f64; + let evaluations = inputs.evaluation_count as f64; + let updates = arriving_input_rows as f64; + let cpu_ops = (build_inserts as f64 + update_inserts as f64) * insert + + evaluations * counts.merges_per_read as f64 * instances * merge + + evaluations * counts.subtracts_per_read as f64 * instances * subtract + + updates * counts.deletes_per_update as f64 * instances * delete + + evaluations * counts.readouts_per_read as f64 * instances * readout; + if !cpu_ops.is_finite() { + return Err(AnalyticalCostError::Overflow); + } + + let state_instances = inputs + .active_window_count + .checked_add(inputs.retained_window_count) + .and_then(|n| n.checked_mul(inputs.physical_sketch_count)) + .and_then(|n| n.checked_mul(counts.state_builds)) + .ok_or(AnalyticalCostError::Overflow)?; + let retained_bytes = state_instances + .checked_mul(inputs.state_bytes_per_sketch) + .ok_or(AnalyticalCostError::Overflow)?; + // Merge/subtract may stream over persistent inputs but still needs one + // result state per physical instance. Persistent retained windows are + // already included above and are not loaded a second time. + let transient_bytes = if counts.merges_per_read > 0 || counts.subtracts_per_read > 0 { + inputs + .physical_sketch_count + .checked_mul(inputs.state_bytes_per_sketch) + .ok_or(AnalyticalCostError::Overflow)? + } else { + 0 + }; + let bootstrap_row_buffer = inputs + .initial_input_bytes + .div_ceil(inputs.initial_input_rows); + Ok(ResourceEstimate { + cpu_ops, + peak_memory_bytes: retained_bytes + .checked_add(transient_bytes) + .ok_or(AnalyticalCostError::Overflow)? + .max(bootstrap_row_buffer), + scan_bytes: inputs.initial_source_scan_bytes, + }) +} + +fn arriving_rows_for_lifecycle( + inputs: StreamingSummaryInputs, + guarantee: &SummaryMaintenanceLifecycleGuarantee, +) -> Result { + let horizon_end = inputs + .planning_time_ms + .checked_add(inputs.horizon_ms) + .ok_or(AnalyticalCostError::Overflow)?; + let active_ms = match guarantee.summary_maintenance_lifecycle { + SummaryMaintenanceLifecycle::Prepared { + activate_at, + retire_at, + } => { + if activate_at.0 >= retire_at.0 { + return Err(AnalyticalCostError::IncompatibleLifecycleGuarantee); + } + let start = activate_at.0.max(inputs.planning_time_ms); + let end = retire_at.0.min(horizon_end); + end.saturating_sub(start) + } + SummaryMaintenanceLifecycle::Shared { retention } => { + if retention.0 < inputs.horizon_ms { + return Err(AnalyticalCostError::IncompatibleLifecycleGuarantee); + } + inputs.horizon_ms + } + SummaryMaintenanceLifecycle::ContinuouslyMaintained => inputs.horizon_ms, + SummaryMaintenanceLifecycle::Ephemeral => { + return Err(AnalyticalCostError::IncompatibleLifecycleGuarantee) + } + }; + let rows = inputs.ingestion_rate_per_second * active_ms as f64 / 1000.0; + if !rows.is_finite() || rows > u64::MAX as f64 { + return Err(AnalyticalCostError::Overflow); + } + Ok(rows.ceil() as u64) +} + +fn validate_guarantee( + guarantee: &SummaryMaintenanceLifecycleGuarantee, + arrival: DataArrival, +) -> Result<(), AnalyticalCostError> { + if guarantee.summary_maintenance_mode != SummaryMaintenanceMode::Incremental + || guarantee.summary_maintenance_mode + != maintenance_mode(&guarantee.summary_maintenance_lifecycle, arrival) + || guarantee.evaluation_schedule + != evaluation_schedule(&guarantee.summary_maintenance_lifecycle, arrival) + { + return Err(AnalyticalCostError::IncompatibleLifecycleGuarantee); + } + Ok(()) +} + +fn required_cpu(name: &'static str, value: Option) -> Result { + let value = value.ok_or(AnalyticalCostError::MissingOrStale(name))?; + if !value.is_finite() || value < 0.0 { + return Err(AnalyticalCostError::InvalidOperationCost(name, value)); + } + Ok(value) +} + +fn required_cpu_when( + count: u64, + name: &'static str, + value: Option, +) -> Result { + if count == 0 { + return Ok(0.0); + } + required_cpu(name, value) +} + +fn count_operations(root: &SummaryNode) -> Result { + fn visit( + node: &SummaryNode, + seen: &mut HashSet<*const SummaryNode>, + counts: &mut SummaryOperationCounts, + ) -> Result<(), AnalyticalCostError> { + if !seen.insert(node as *const SummaryNode) { + return Ok(()); + } + match &node.expr { + SummaryExpr::KeepPreAsap(_) => {} + SummaryExpr::SummaryAgg { child, .. } => { + counts.state_builds = counts + .state_builds + .checked_add(1) + .ok_or(AnalyticalCostError::Overflow)?; + visit(child, seen, counts)?; + } + SummaryExpr::SummaryMerge { children } => { + if children.is_empty() { + return Err(AnalyticalCostError::InvalidPhysicalDag( + "summary merge has no children", + )); + } + counts.merges_per_read = counts + .merges_per_read + .checked_add(children.len().saturating_sub(1) as u64) + .ok_or(AnalyticalCostError::Overflow)?; + for child in children { + visit(child, seen, counts)?; + } + } + SummaryExpr::SummarySubtract { left, right } => { + counts.subtracts_per_read = counts + .subtracts_per_read + .checked_add(1) + .ok_or(AnalyticalCostError::Overflow)?; + visit(left, seen, counts)?; + visit(right, seen, counts)?; + } + SummaryExpr::SummaryDelete { summary_input, .. } => { + counts.deletes_per_update = counts + .deletes_per_update + .checked_add(1) + .ok_or(AnalyticalCostError::Overflow)?; + visit(summary_input, seen, counts)?; + } + SummaryExpr::SummaryEstimate { summary_input, .. } => { + counts.readouts_per_read = counts + .readouts_per_read + .checked_add(1) + .ok_or(AnalyticalCostError::Overflow)?; + visit(summary_input, seen, counts)?; + } + SummaryExpr::SummaryJoin { .. } => { + return Err(AnalyticalCostError::UnsupportedSummaryOperation("join")); + } + } + Ok(()) + } + + let mut counts = SummaryOperationCounts::default(); + visit(root, &mut HashSet::new(), &mut counts)?; + Ok(counts) +} + +#[cfg(test)] +mod tests { + use std::rc::Rc; + + use asap_types::post_asap::{ + EvaluationSchedule, ExactKind, ExactParams, GroupingStrategy, OutputRepresentation, + SummaryExpr, SummaryFamilyType, SummaryField, SummaryMaintenanceLifecycle, + SummaryMaintenanceLifecycleGuarantee, SummaryMaintenanceMode, SummarySchema, + }; + use asap_types::pre_asap::{ColumnRef, QueryExpr, Reduction}; + use asap_types::workload::{ + DataWorkload, Evidence, EvidenceSource, Predictability, Query, QueryRecurrence, + QueryRequirements, QueryTimeScope, Rate, RepeatedDemand, RepetitionInterval, TimeSelection, + }; + + use super::*; + + fn physical() -> StreamingPhysicalInputEvidence { + StreamingPhysicalInputEvidence { + initial_input_bytes: 640, + initial_source_scan_bytes: 640, + active_window_count: 2, + retained_window_count: 3, + physical_sketch_count: 2, + state_bytes_per_sketch: 100, + } + } + + fn query() -> QueryWorkloadEntry { + QueryWorkloadEntry { + query: Query("streaming count".into()), + requirements: QueryRequirements::default(), + predictability: Predictability::Unknown, + recurrence: QueryRecurrence::Repeated(RepeatedDemand::FixedInterval( + RepetitionInterval(1_000), + )), + time_selection: TimeSelection { + scope: QueryTimeScope::Unknown, + lookback: None, + as_of: None, + }, + } + } + + fn continuous_guarantee() -> SummaryMaintenanceLifecycleGuarantee { + SummaryMaintenanceLifecycleGuarantee { + summary_maintenance_lifecycle: SummaryMaintenanceLifecycle::ContinuouslyMaintained, + summary_maintenance_mode: SummaryMaintenanceMode::Incremental, + evaluation_schedule: EvaluationSchedule::PerUpdate, + output_representation: OutputRepresentation::SummaryState, + } + } + + #[test] + fn workload_adapter_derives_updates_and_reads_over_one_horizon() { + let data = DataWorkload { + arrival: DataArrival::ContinuouslyIngesting, + ingestion_rate: Evidence { + value: Some(Rate(2.0)), + source: EvidenceSource::Observed, + observed_at_ms: Some(100), + valid_for_ms: Some(10_000), + }, + input_cardinality: Evidence { + value: Some(10), + source: EvidenceSource::Observed, + observed_at_ms: Some(100), + valid_for_ms: Some(10_000), + }, + ..DataWorkload::default() + }; + + let inputs = + StreamingSummaryInputs::from_workload(physical(), &data, &query(), 100, 5_000).unwrap(); + assert_eq!(inputs.initial_input_rows, 10); + assert_eq!( + arriving_rows_for_lifecycle(inputs, &continuous_guarantee()).unwrap(), + 10 + ); + assert_eq!(inputs.evaluation_count, 5); + } + + #[test] + fn mixed_arrival_fails_closed_until_backlog_and_stream_are_separate() { + let data = DataWorkload { + arrival: DataArrival::Mixed, + ..DataWorkload::default() + }; + assert_eq!( + StreamingSummaryInputs::from_workload(physical(), &data, &query(), 0, 1_000), + Err(AnalyticalCostError::UnsupportedDataArrival( + DataArrival::Mixed + )) + ); + } + + #[test] + fn direct_read_costs_build_updates_windows_and_recurrence() { + let estimate = estimate_incremental_summary_maintenance( + &summary_with_operations(false, false, false), + &continuous_guarantee(), + 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, + }, + SummaryOperationCpuEvidence { + insert_cpu_ops: Some(2.0), + readout_cpu_ops: Some(3.0), + ..SummaryOperationCpuEvidence::default() + }, + ) + .unwrap(); + // 10 bootstrap + 10 arrivals into two active windows; two states read 5 times. + assert_eq!(estimate.cpu_ops, 90.0); + assert_eq!(estimate.peak_memory_bytes, 1_000); + assert_eq!(estimate.scan_bytes, 640); + } + + #[test] + fn operations_use_update_or_read_multiplicity_and_shared_state_once() { + let estimate = estimate_incremental_summary_maintenance( + &summary_with_operations(true, true, true), + &continuous_guarantee(), + StreamingSummaryInputs { + data_arrival: DataArrival::ContinuouslyIngesting, + initial_input_rows: 1, + initial_input_bytes: 8, + initial_source_scan_bytes: 8, + ingestion_rate_per_second: 4.0, + planning_time_ms: 0, + horizon_ms: 1_000, + active_window_count: 1, + retained_window_count: 2, + physical_sketch_count: 2, + state_bytes_per_sketch: 10, + evaluation_count: 3, + }, + SummaryOperationCpuEvidence { + insert_cpu_ops: Some(1.0), + merge_cpu_ops: Some(2.0), + subtract_cpu_ops: Some(3.0), + delete_cpu_ops: Some(5.0), + readout_cpu_ops: Some(7.0), + }, + ) + .unwrap(); + assert_eq!(estimate.cpu_ops, 5.0 + 40.0 + 12.0 + 18.0 + 42.0); + // Three persistent windows plus one transient result, for two instances. + assert_eq!(estimate.peak_memory_bytes, 80); + } + + #[test] + fn lifecycle_mode_and_schedule_must_match_existing_planner_semantics() { + let mut guarantee = continuous_guarantee(); + guarantee.evaluation_schedule = EvaluationSchedule::OnRead; + assert_eq!( + estimate_incremental_summary_maintenance( + &summary_with_operations(false, false, false), + &guarantee, + 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: 1, + }, + SummaryOperationCpuEvidence { + insert_cpu_ops: Some(1.0), + readout_cpu_ops: Some(1.0), + ..SummaryOperationCpuEvidence::default() + }, + ), + Err(AnalyticalCostError::IncompatibleLifecycleGuarantee) + ); + } + + #[test] + fn missing_cost_for_an_operation_in_the_dag_fails_closed() { + assert_eq!( + estimate_incremental_summary_maintenance( + &summary_with_operations(true, false, false), + &continuous_guarantee(), + 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: 1, + }, + SummaryOperationCpuEvidence { + insert_cpu_ops: Some(1.0), + readout_cpu_ops: Some(1.0), + ..SummaryOperationCpuEvidence::default() + }, + ), + Err(AnalyticalCostError::MissingOrStale("merge_cpu_ops")) + ); + } + + #[test] + fn direct_build_mode_is_not_mispriced_as_incremental_maintenance() { + let mut guarantee = continuous_guarantee(); + guarantee.summary_maintenance_lifecycle = SummaryMaintenanceLifecycle::Ephemeral; + guarantee.summary_maintenance_mode = SummaryMaintenanceMode::DirectBuild; + guarantee.evaluation_schedule = EvaluationSchedule::OneShot; + assert_eq!( + estimate_incremental_summary_maintenance( + &summary_with_operations(false, false, false), + &guarantee, + 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: 1, + }, + SummaryOperationCpuEvidence { + insert_cpu_ops: Some(1.0), + readout_cpu_ops: Some(1.0), + ..SummaryOperationCpuEvidence::default() + }, + ), + Err(AnalyticalCostError::IncompatibleLifecycleGuarantee) + ); + } + + #[test] + fn prepared_maintenance_charges_only_its_active_interval() { + let guarantee = SummaryMaintenanceLifecycleGuarantee { + summary_maintenance_lifecycle: SummaryMaintenanceLifecycle::Prepared { + activate_at: asap_types::workload::TimestampMs(1_000), + retire_at: asap_types::workload::TimestampMs(2_000), + }, + summary_maintenance_mode: SummaryMaintenanceMode::Incremental, + evaluation_schedule: EvaluationSchedule::PerUpdate, + output_representation: OutputRepresentation::SummaryState, + }; + let estimate = estimate_incremental_summary_maintenance( + &summary_with_operations(false, false, false), + &guarantee, + StreamingSummaryInputs { + data_arrival: DataArrival::ContinuouslyIngesting, + initial_input_rows: 10, + initial_input_bytes: 80, + initial_source_scan_bytes: 80, + ingestion_rate_per_second: 2.0, + planning_time_ms: 0, + horizon_ms: 5_000, + active_window_count: 1, + retained_window_count: 1, + physical_sketch_count: 1, + state_bytes_per_sketch: 8, + evaluation_count: 1, + }, + SummaryOperationCpuEvidence { + insert_cpu_ops: Some(1.0), + readout_cpu_ops: Some(0.0), + ..SummaryOperationCpuEvidence::default() + }, + ) + .unwrap(); + assert_eq!(estimate.cpu_ops, 12.0); // 10 bootstrap + 2 updates in [1s, 2s]. + } + + #[test] + fn shared_retention_must_cover_the_comparison_horizon() { + let guarantee = SummaryMaintenanceLifecycleGuarantee { + summary_maintenance_lifecycle: SummaryMaintenanceLifecycle::Shared { + retention: asap_types::workload::DurationMs(999), + }, + summary_maintenance_mode: SummaryMaintenanceMode::Incremental, + evaluation_schedule: EvaluationSchedule::PerUpdate, + output_representation: OutputRepresentation::SummaryState, + }; + assert_eq!( + estimate_incremental_summary_maintenance( + &summary_with_operations(false, false, false), + &guarantee, + 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: 1, + }, + SummaryOperationCpuEvidence { + insert_cpu_ops: Some(1.0), + readout_cpu_ops: Some(1.0), + ..SummaryOperationCpuEvidence::default() + }, + ), + Err(AnalyticalCostError::IncompatibleLifecycleGuarantee) + ); + } + + fn summary_with_operations(merge: bool, subtract: bool, delete: bool) -> Rc { + let state_type = SummaryFamilyType::ExactAggregate(ExactKind::Count, ExactParams::Count); + let schema = SummarySchema { + fields: vec![SummaryField { + name: "count".into(), + dtype: state_type.clone(), + nullable: false, + }], + time_index: None, + }; + let leaf = Rc::new(SummaryNode { + expr: SummaryExpr::KeepPreAsap(Rc::new(QueryExpr::promql_scalar(1.0))), + schema: schema.clone(), + guarantee: None, + }); + let agg = Rc::new(SummaryNode { + expr: SummaryExpr::SummaryAgg { + child: leaf, + family: state_type, + col: ColumnRef::Wildcard, + reduction: Reduction::by(vec![]), + grouping: GroupingStrategy::PerSubpopulationInstance, + }, + schema: schema.clone(), + guarantee: None, + }); + let mut root = Rc::clone(&agg); + if merge { + root = Rc::new(SummaryNode { + expr: SummaryExpr::SummaryMerge { + children: vec![Rc::clone(&agg), Rc::clone(&agg)], + }, + schema: schema.clone(), + guarantee: None, + }); + } + if subtract { + root = Rc::new(SummaryNode { + expr: SummaryExpr::SummarySubtract { + left: Rc::clone(&root), + right: Rc::clone(&agg), + }, + schema: schema.clone(), + guarantee: None, + }); + } + if delete { + root = Rc::new(SummaryNode { + expr: SummaryExpr::SummaryDelete { + summary_input: root, + key: ColumnRef::Wildcard, + }, + schema: schema.clone(), + guarantee: None, + }); + } + Rc::new(SummaryNode { + expr: SummaryExpr::SummaryEstimate { + summary_input: root, + query: asap_types::post_asap::SketchQuery::PointCount { + key: ColumnRef::Wildcard, + value: None, + }, + }, + schema, + guarantee: None, + }) + } +} diff --git a/crates/asap-aware-mapping/src/lib.rs b/crates/asap-aware-mapping/src/lib.rs index 3e4fa92b..a3071c28 100644 --- a/crates/asap-aware-mapping/src/lib.rs +++ b/crates/asap-aware-mapping/src/lib.rs @@ -184,6 +184,7 @@ pub mod accuracy; pub mod accuracy_reconciliation; pub mod analytical_cost; +pub mod analytical_streaming_cost; pub mod cost_model; pub mod explanation; pub mod grouping; diff --git a/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs b/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs index 476d9c25..c8ce8853 100644 --- a/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs +++ b/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs @@ -919,7 +919,7 @@ fn collect_summary_aggs( } } -fn evaluation_schedule( +pub(crate) fn evaluation_schedule( lifecycle: &SummaryMaintenanceLifecycle, arrival: DataArrival, ) -> EvaluationSchedule { @@ -1053,7 +1053,7 @@ fn select_compatible_lifecycles( } } -fn maintenance_mode( +pub(crate) fn maintenance_mode( lifecycle: &SummaryMaintenanceLifecycle, arrival: DataArrival, ) -> SummaryMaintenanceMode { From 853d7bad05e8f6dc3536ef6091eb5fd1d81b660d Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 3 Sep 2026 10:40:54 -0600 Subject: [PATCH 2/6] fix(cost): use shared recurrence evidence in streaming base --- crates/asap-aware-mapping/src/analytical_streaming_cost.rs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/crates/asap-aware-mapping/src/analytical_streaming_cost.rs b/crates/asap-aware-mapping/src/analytical_streaming_cost.rs index 973e67a2..cd6a00eb 100644 --- a/crates/asap-aware-mapping/src/analytical_streaming_cost.rs +++ b/crates/asap-aware-mapping/src/analytical_streaming_cost.rs @@ -13,7 +13,8 @@ use asap_types::post_asap::{ use asap_types::workload::{DataArrival, DataWorkload, QueryWorkloadEntry}; use serde::{Deserialize, Serialize}; -use crate::analytical_cost::{evaluations_in_horizon, AnalyticalCostError, ResourceEstimate}; +use crate::analytical_cost::{AnalyticalCostError, ResourceEstimate}; +use crate::analytical_statistics::evaluations_in_horizon; use crate::summary_maintenance_lifecycle::{evaluation_schedule, maintenance_mode}; pub const ANALYTICAL_STREAMING_MODEL_VERSION: &str = "analytical-summary-incremental-v1"; From 10d3922cc9d760cbebd38d515faee858fcd2f129 Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 3 Sep 2026 18:01:01 -0600 Subject: [PATCH 3/6] fix(cost): make streaming summary evidence conservative --- .../src/analytical_streaming_cost.rs | 196 +++++++++++------- .../analytical-resource-cost.md | 41 ++++ 2 files changed, 161 insertions(+), 76 deletions(-) diff --git a/crates/asap-aware-mapping/src/analytical_streaming_cost.rs b/crates/asap-aware-mapping/src/analytical_streaming_cost.rs index cd6a00eb..44e11fc2 100644 --- a/crates/asap-aware-mapping/src/analytical_streaming_cost.rs +++ b/crates/asap-aware-mapping/src/analytical_streaming_cost.rs @@ -34,11 +34,11 @@ pub struct StreamingPhysicalInputEvidence { pub active_window_count: u64, /// Completed windows retained for query coverage. pub retained_window_count: u64, - /// Independent state instances per window: one for shared + /// Independent summary-state instances per window: one for shared /// multi-subpopulation state, otherwise the resolved group count. - pub physical_sketch_count: u64, + pub physical_summary_count: u64, /// Resident bytes of one concrete state instance. - pub state_bytes_per_sketch: u64, + pub state_bytes_per_summary: u64, } /// Workload-normalized inputs for incremental maintenance over one finite @@ -54,8 +54,8 @@ pub struct StreamingSummaryInputs { pub horizon_ms: u64, pub active_window_count: u64, pub retained_window_count: u64, - pub physical_sketch_count: u64, - pub state_bytes_per_sketch: u64, + pub physical_summary_count: u64, + pub state_bytes_per_summary: u64, pub evaluation_count: u64, } @@ -101,8 +101,8 @@ impl StreamingSummaryInputs { horizon_ms, active_window_count: physical.active_window_count, retained_window_count: physical.retained_window_count, - physical_sketch_count: physical.physical_sketch_count, - state_bytes_per_sketch: physical.state_bytes_per_sketch, + physical_summary_count: physical.physical_summary_count, + state_bytes_per_summary: physical.state_bytes_per_summary, evaluation_count: evaluations_in_horizon( &query.recurrence, planning_time_ms, @@ -119,19 +119,26 @@ impl StreamingSummaryInputs { )); } for (name, value) in [ - ("initial_input_rows", self.initial_input_rows), - ("initial_input_bytes", self.initial_input_bytes), ("horizon_ms", self.horizon_ms), ("active_window_count", self.active_window_count), - ("retained_window_count", self.retained_window_count), - ("physical_sketch_count", self.physical_sketch_count), - ("state_bytes_per_sketch", self.state_bytes_per_sketch), + ("physical_summary_count", self.physical_summary_count), + ("state_bytes_per_summary", self.state_bytes_per_summary), ("evaluation_count", self.evaluation_count), ] { if value == 0 { return Err(AnalyticalCostError::MissingOrZero(name)); } } + let bootstrap_is_consistent = if self.initial_input_rows == 0 { + self.initial_input_bytes == 0 && self.initial_source_scan_bytes == 0 + } else { + 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", + )); + } if !self.ingestion_rate_per_second.is_finite() || self.ingestion_rate_per_second < 0.0 { return Err(AnalyticalCostError::InvalidIngestionRate( self.ingestion_rate_per_second, @@ -176,7 +183,10 @@ pub fn estimate_incremental_summary_maintenance( validate_guarantee(guarantee, inputs.data_arrival)?; let arriving_input_rows = arriving_rows_for_lifecycle(inputs, guarantee)?; let counts = count_operations(root)?; - if counts.state_builds == 0 { + if counts.state_builds != 1 { + // One flat evidence record cannot safely describe several summary + // nodes with different inputs, state sizes, or algorithms. Complete + // multi-node streaming DAGs use per-node evidence instead. return Err(AnalyticalCostError::UnsupportedCandidate); } @@ -200,19 +210,19 @@ pub fn estimate_incremental_summary_maintenance( let build_inserts = inputs .initial_input_rows - .checked_mul(counts.state_builds) + .checked_mul(inputs.active_window_count) + .and_then(|n| n.checked_mul(counts.state_builds)) .ok_or(AnalyticalCostError::Overflow)?; let update_inserts = arriving_input_rows .checked_mul(inputs.active_window_count) .and_then(|n| n.checked_mul(counts.state_builds)) .ok_or(AnalyticalCostError::Overflow)?; - let instances = inputs.physical_sketch_count as f64; + let instances = inputs.physical_summary_count as f64; let evaluations = inputs.evaluation_count as f64; - let updates = arriving_input_rows as f64; let cpu_ops = (build_inserts as f64 + update_inserts as f64) * insert + evaluations * counts.merges_per_read as f64 * instances * merge + evaluations * counts.subtracts_per_read as f64 * instances * subtract - + updates * counts.deletes_per_update as f64 * instances * delete + + update_inserts as f64 * counts.deletes_per_update as f64 * delete + evaluations * counts.readouts_per_read as f64 * instances * readout; if !cpu_ops.is_finite() { return Err(AnalyticalCostError::Overflow); @@ -221,26 +231,30 @@ pub fn estimate_incremental_summary_maintenance( let state_instances = inputs .active_window_count .checked_add(inputs.retained_window_count) - .and_then(|n| n.checked_mul(inputs.physical_sketch_count)) + .and_then(|n| n.checked_mul(inputs.physical_summary_count)) .and_then(|n| n.checked_mul(counts.state_builds)) .ok_or(AnalyticalCostError::Overflow)?; let retained_bytes = state_instances - .checked_mul(inputs.state_bytes_per_sketch) + .checked_mul(inputs.state_bytes_per_summary) .ok_or(AnalyticalCostError::Overflow)?; // Merge/subtract may stream over persistent inputs but still needs one // result state per physical instance. Persistent retained windows are // already included above and are not loaded a second time. let transient_bytes = if counts.merges_per_read > 0 || counts.subtracts_per_read > 0 { inputs - .physical_sketch_count - .checked_mul(inputs.state_bytes_per_sketch) + .physical_summary_count + .checked_mul(inputs.state_bytes_per_summary) .ok_or(AnalyticalCostError::Overflow)? } else { 0 }; - let bootstrap_row_buffer = inputs - .initial_input_bytes - .div_ceil(inputs.initial_input_rows); + let bootstrap_row_buffer = if inputs.initial_input_rows == 0 { + 0 + } else { + inputs + .initial_input_bytes + .div_ceil(inputs.initial_input_rows) + }; Ok(ResourceEstimate { cpu_ops, peak_memory_bytes: retained_bytes @@ -269,14 +283,13 @@ fn arriving_rows_for_lifecycle( } let start = activate_at.0.max(inputs.planning_time_ms); let end = retire_at.0.min(horizon_end); - end.saturating_sub(start) - } - SummaryMaintenanceLifecycle::Shared { retention } => { - if retention.0 < inputs.horizon_ms { + let active_ms = end.saturating_sub(start); + if active_ms == 0 { return Err(AnalyticalCostError::IncompatibleLifecycleGuarantee); } - inputs.horizon_ms + active_ms } + SummaryMaintenanceLifecycle::Shared { .. } => inputs.horizon_ms, SummaryMaintenanceLifecycle::ContinuouslyMaintained => inputs.horizon_ms, SummaryMaintenanceLifecycle::Ephemeral => { return Err(AnalyticalCostError::IncompatibleLifecycleGuarantee) @@ -306,7 +319,7 @@ fn validate_guarantee( fn required_cpu(name: &'static str, value: Option) -> Result { let value = value.ok_or(AnalyticalCostError::MissingOrStale(name))?; - if !value.is_finite() || value < 0.0 { + if !value.is_finite() || value <= 0.0 { return Err(AnalyticalCostError::InvalidOperationCost(name, value)); } Ok(value) @@ -412,8 +425,8 @@ mod tests { initial_source_scan_bytes: 640, 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, } } @@ -500,8 +513,8 @@ 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, }, SummaryOperationCpuEvidence { @@ -511,8 +524,9 @@ mod tests { }, ) .unwrap(); - // 10 bootstrap + 10 arrivals into two active windows; two states read 5 times. - assert_eq!(estimate.cpu_ops, 90.0); + // 10 bootstrap + 10 arrivals, each routed to two active windows; + // two physical summary instances are read 5 times. + assert_eq!(estimate.cpu_ops, 110.0); assert_eq!(estimate.peak_memory_bytes, 1_000); assert_eq!(estimate.scan_bytes, 640); } @@ -532,8 +546,8 @@ mod tests { horizon_ms: 1_000, active_window_count: 1, retained_window_count: 2, - physical_sketch_count: 2, - state_bytes_per_sketch: 10, + physical_summary_count: 2, + state_bytes_per_summary: 10, evaluation_count: 3, }, SummaryOperationCpuEvidence { @@ -545,7 +559,7 @@ mod tests { }, ) .unwrap(); - assert_eq!(estimate.cpu_ops, 5.0 + 40.0 + 12.0 + 18.0 + 42.0); + assert_eq!(estimate.cpu_ops, 5.0 + 12.0 + 18.0 + 20.0 + 42.0); // Three persistent windows plus one transient result, for two instances. assert_eq!(estimate.peak_memory_bytes, 80); } @@ -568,8 +582,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, }, SummaryOperationCpuEvidence { @@ -598,8 +612,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, }, SummaryOperationCpuEvidence { @@ -632,8 +646,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, }, SummaryOperationCpuEvidence { @@ -670,22 +684,22 @@ mod tests { horizon_ms: 5_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, }, SummaryOperationCpuEvidence { insert_cpu_ops: Some(1.0), - readout_cpu_ops: Some(0.0), + readout_cpu_ops: Some(1.0), ..SummaryOperationCpuEvidence::default() }, ) .unwrap(); - assert_eq!(estimate.cpu_ops, 12.0); // 10 bootstrap + 2 updates in [1s, 2s]. + assert_eq!(estimate.cpu_ops, 13.0); // 10 bootstrap + 2 updates + one read. } #[test] - fn shared_retention_must_cover_the_comparison_horizon() { + fn shared_retention_is_not_confused_with_the_planning_horizon() { let guarantee = SummaryMaintenanceLifecycleGuarantee { summary_maintenance_lifecycle: SummaryMaintenanceLifecycle::Shared { retention: asap_types::workload::DurationMs(999), @@ -694,32 +708,62 @@ mod tests { evaluation_schedule: EvaluationSchedule::PerUpdate, output_representation: OutputRepresentation::SummaryState, }; - assert_eq!( - estimate_incremental_summary_maintenance( - &summary_with_operations(false, false, false), - &guarantee, - 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: 1, - }, - SummaryOperationCpuEvidence { - insert_cpu_ops: Some(1.0), - readout_cpu_ops: Some(1.0), - ..SummaryOperationCpuEvidence::default() - }, - ), - Err(AnalyticalCostError::IncompatibleLifecycleGuarantee) - ); + assert!(estimate_incremental_summary_maintenance( + &summary_with_operations(false, false, false), + &guarantee, + 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_summary_count: 1, + state_bytes_per_summary: 8, + evaluation_count: 1, + }, + SummaryOperationCpuEvidence { + insert_cpu_ops: Some(1.0), + readout_cpu_ops: Some(1.0), + ..SummaryOperationCpuEvidence::default() + }, + ) + .is_ok()); + } + + #[test] + fn an_empty_bootstrap_is_valid_for_a_new_stream() { + let estimate = estimate_incremental_summary_maintenance( + &summary_with_operations(false, false, false), + &continuous_guarantee(), + StreamingSummaryInputs { + data_arrival: DataArrival::ContinuouslyIngesting, + initial_input_rows: 0, + initial_input_bytes: 0, + 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: 0, + physical_summary_count: 1, + state_bytes_per_summary: 8, + evaluation_count: 1, + }, + SummaryOperationCpuEvidence { + insert_cpu_ops: Some(1.0), + readout_cpu_ops: Some(1.0), + ..SummaryOperationCpuEvidence::default() + }, + ) + .unwrap(); + + assert_eq!(estimate.cpu_ops, 2.0); + assert_eq!(estimate.peak_memory_bytes, 8); + assert_eq!(estimate.scan_bytes, 0); } fn summary_with_operations(merge: bool, subtract: bool, delete: bool) -> Rc { 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 ec6756a3..c1ddc546 100644 --- a/docs/design_docs/asap-aware-mapping/analytical-resource-cost.md +++ b/docs/design_docs/asap-aware-mapping/analytical-resource-cost.md @@ -562,6 +562,47 @@ input. ## Summary operator formulas +### Incremental single-summary foundation + +For `DataArrival::ContinuouslyIngesting`, the incremental estimator accepts +one selected lifecycle and one unique logical `SummaryAgg`. This deliberately +narrow contract prevents one flat evidence record from being reused across +several summary nodes with different input cardinalities, algorithms, or state +sizes. Complete multi-node streaming alternatives require per-node physical +evidence. + +The canonical workload supplies fresh bootstrap cardinality, ingestion rate, +query recurrence, planning time, and a finite horizon. Physical evidence adds +logical/bootstrap bytes, physical bootstrap scan bytes, active and retained +window counts, the number of concrete summary-state instances per window, and +bytes per state instance. Names use `summary`, not `sketch`, because an exact +aggregate or another non-sketch state is equally valid. + +For bootstrap rows `B`, arrivals `U`, simultaneously updated windows `A`, +query evaluations `Q`, physical summary instances `P`, and state bytes `S`: + +```text +insert invocations = (B + U) × A +retained memory = (A + retained_windows) × P × S +``` + +Each input row is routed to its matching summary instance; it is not inserted +into every group. Merge, subtract, and readout work may operate over all `P` +instances. Delete work follows the same routed window updates rather than +multiplying every update by every possible group. + +An empty bootstrap is valid and has zero logical bytes and zero source reads. +A non-empty bootstrap requires positive logical and physical source bytes. +Active window count, summary-instance count, state width, horizon, and query +evaluation count must be positive; retained-window count may be zero for a new +stream. Required per-operation CPU evidence must be finite and positive. + +Lifecycle retention and the planning horizon are different quantities. A +short retained window may be maintained throughout a much longer planning +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. + A retained sketch performs one build and serves later reads from state: ```text From 59f17345e5c2fb504f94d0a283b096c0fe5c846e Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 3 Sep 2026 18:07:30 -0600 Subject: [PATCH 4/6] docs(cost): distinguish batch and streaming entry points --- .../analytical-resource-cost.md | 29 +++++++++++-------- 1 file changed, 17 insertions(+), 12 deletions(-) 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 c1ddc546..d95aaf57 100644 --- a/docs/design_docs/asap-aware-mapping/analytical-resource-cost.md +++ b/docs/design_docs/asap-aware-mapping/analytical-resource-cost.md @@ -10,11 +10,17 @@ plans. These are separate concerns: - the **analytical estimation layer** applies algorithmic formulas to that evidence to estimate CPU work, peak memory, and source/disk I/O. -The model version implemented here is explicitly for `DataArrival::AtRest`. -It replaces dimensionless plan-node counts with estimates derived from -operator complexity, cardinality, row width, and concrete summary parameters. -The estimates are predictions; they are not measurements reported by a -physical executor. +There are two arrival-specific entry points. The complete physical-DAG +comparison currently models `DataArrival::AtRest`. The streaming summary +extension models `DataArrival::ContinuouslyIngesting` over a finite horizon, +including bootstrap, arriving updates, retained state, and query readout. +Evidence from one arrival mode must not be reused for the other. `Mixed` and +`Unknown` remain unavailable until their distinct data regions are modeled. + +Both entry points replace dimensionless plan-node counts with estimates +derived from operator complexity, cardinality, row width, and concrete summary +parameters. The estimates are predictions; they are not measurements reported +by a physical executor. The model does not decide semantic or accuracy legality. Candidate generation and guarantee composition run first; costing ranks only the candidates that @@ -137,7 +143,7 @@ costed. ## Workload horizon and lifecycle Every alternative must cover the same source data and query horizon. The -implemented `DataArrival::AtRest` comparison is build-once, read-many: +`DataArrival::AtRest` physical-DAG comparison is build-once, read-many: ```text retained-summary builds = 1 @@ -153,12 +159,11 @@ An at-rest estimate must not be reused for `unknown`, `mixed`, or `continuously_ingesting` data. Callers fail closed instead of pretending that incremental updates are a one-time snapshot build. -The sketch alternative scans the selected source snapshot once and retains -state. The raw alternative recomputes from that snapshot for every query -read. Continuously ingesting data, rebuilds, deletions, expiration, and -retention duration belong to the summary-maintenance lifecycle model. They -must contribute update/build/delete work before a continuously maintained -plan is compared with raw execution. +The at-rest summary alternative scans the selected source snapshot once and +retains state. Its raw alternative recomputes from that snapshot for every +query read. The continuously-ingesting entry point separately charges +bootstrap, updates, summary operations, retained state, and raw evaluations; +its lifecycle rules are defined below. ### Comparable source and workload scope From 92c206edaf075ed08b8605024e49295bb90ef8dd Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 3 Sep 2026 18:08:08 -0600 Subject: [PATCH 5/6] docs(cost): correct incremental work formula --- .../asap-aware-mapping/analytical-resource-cost.md | 10 +++++++--- 1 file changed, 7 insertions(+), 3 deletions(-) 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 d95aaf57..dceeb4cd 100644 --- a/docs/design_docs/asap-aware-mapping/analytical-resource-cost.md +++ b/docs/design_docs/asap-aware-mapping/analytical-resource-cost.md @@ -608,16 +608,20 @@ 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. -A retained sketch performs one build and serves later reads from state: +A retained summary bootstraps every active window, consumes arriving rows, and +serves later reads from state. For the simple build/update/readout shape: ```text -cpu_ops = input_rows // build scan - + input_rows × update_ops(params) +cpu_ops = (bootstrap_rows + arriving_rows) + × active_window_count × insert_ops(params) + evaluation_count × physical_summary_count × read_ops(params) scan_bytes = source_read_bytes for the build ``` +Merge, subtract, and delete add their own invocation counts described above; +they are never folded into the simple formula implicitly. + Concrete accuracy-sized parameters determine state and work: | Summary | Update operations per row | Read operations | State bytes per physical instance | From e9cb271d1a5c20f698b69911cda7d907e7eeaa62 Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 3 Sep 2026 20:49:30 -0600 Subject: [PATCH 6/6] refactor(cost): name summary maintenance cost layer --- crates/asap-aware-mapping/src/lib.rs | 2 +- ...alytical_streaming_cost.rs => summary_maintenance_cost.rs} | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) rename crates/asap-aware-mapping/src/{analytical_streaming_cost.rs => summary_maintenance_cost.rs} (99%) diff --git a/crates/asap-aware-mapping/src/lib.rs b/crates/asap-aware-mapping/src/lib.rs index a3071c28..39a43014 100644 --- a/crates/asap-aware-mapping/src/lib.rs +++ b/crates/asap-aware-mapping/src/lib.rs @@ -184,7 +184,6 @@ pub mod accuracy; pub mod accuracy_reconciliation; pub mod analytical_cost; -pub mod analytical_streaming_cost; pub mod cost_model; pub mod explanation; pub mod grouping; @@ -195,6 +194,7 @@ pub mod recurrence; pub mod replacement; pub mod rewrite; pub mod rollup; +pub mod summary_maintenance_cost; pub mod summary_maintenance_dag_export; pub mod summary_maintenance_lifecycle; pub mod topk_reuse; diff --git a/crates/asap-aware-mapping/src/analytical_streaming_cost.rs b/crates/asap-aware-mapping/src/summary_maintenance_cost.rs similarity index 99% rename from crates/asap-aware-mapping/src/analytical_streaming_cost.rs rename to crates/asap-aware-mapping/src/summary_maintenance_cost.rs index 44e11fc2..ced47e2e 100644 --- a/crates/asap-aware-mapping/src/analytical_streaming_cost.rs +++ b/crates/asap-aware-mapping/src/summary_maintenance_cost.rs @@ -14,10 +14,10 @@ use asap_types::workload::{DataArrival, DataWorkload, QueryWorkloadEntry}; use serde::{Deserialize, Serialize}; use crate::analytical_cost::{AnalyticalCostError, ResourceEstimate}; -use crate::analytical_statistics::evaluations_in_horizon; +use crate::physical_operator_statistics::evaluations_in_horizon; use crate::summary_maintenance_lifecycle::{evaluation_schedule, maintenance_mode}; -pub const ANALYTICAL_STREAMING_MODEL_VERSION: &str = "analytical-summary-incremental-v1"; +pub const SUMMARY_MAINTENANCE_COST_MODEL_VERSION: &str = "summary-maintenance-resource-v1"; /// Physical evidence that is not represented by [`DataWorkload`] for one /// incrementally maintained summary deployment. Window counts describe the