From b68b23108c515d1b0605fc0931c0fa22ee8b82e3 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Fri, 2 Oct 2026 17:01:39 +0000 Subject: [PATCH] fix: generalize summary cost evidence across arrival modes --- .../src/summary_maintenance_cost/estimator.rs | 32 +- .../src/summary_maintenance_cost/evidence.rs | 123 ++++--- .../src/summary_maintenance_cost/mod.rs | 4 +- .../src/summary_maintenance_cost/model.rs | 341 +++++++++++++----- .../src/summary_maintenance_cost/window.rs | 16 +- .../architecture/physical-plan-integration.md | 34 +- .../architecture/planner-runtime-contract.md | 2 +- 7 files changed, 381 insertions(+), 171 deletions(-) diff --git a/crates/asap-aware-mapping/src/summary_maintenance_cost/estimator.rs b/crates/asap-aware-mapping/src/summary_maintenance_cost/estimator.rs index 4af4dabd..e3160640 100644 --- a/crates/asap-aware-mapping/src/summary_maintenance_cost/estimator.rs +++ b/crates/asap-aware-mapping/src/summary_maintenance_cost/estimator.rs @@ -3,9 +3,9 @@ use super::*; pub(super) fn estimate_heterogeneous_summary( root: &SummaryNode, deployments: &[CostedSummaryDeployment<'_>], - evidence: &StreamingNodeEvidence, + evidence: &SummaryNodeEvidence, scope: &ComparisonScope, - raw: &StreamingRawInputEvidence, + raw: &RawInputEvidence, window_frameworks: &[Option], ) -> Result { if window_frameworks.len() != deployments.len() { @@ -71,7 +71,7 @@ pub(super) fn estimate_heterogeneous_summary( let mut physical_states = HashMap::< String, ( - StreamingAggregateEvidence, + SummaryAggregateEvidence, SummaryMaintenanceLifecycleGuarantee, String, Option, @@ -85,6 +85,7 @@ pub(super) fn estimate_heterogeneous_summary( return Err(AnalyticalCostError::UnsupportedCandidate); }; let inputs = node_evidence.inputs.validate()?; + validate_arrival_rate(scope.data_arrival, inputs.ingestion_rate_per_second)?; match node_evidence.source_coverage_index { Some(index) => { let declared = @@ -255,7 +256,7 @@ pub(super) fn estimate_heterogeneous_summary( node: &SummaryNode, seen: &mut HashSet, by_node: &HashMap<*const SummaryNode, &CostedSummaryDeployment<'_>>, - evidence: &StreamingNodeEvidence, + evidence: &SummaryNodeEvidence, scope: &ComparisonScope, evaluation_count: u64, cpu_ops: &mut f64, @@ -392,7 +393,7 @@ pub(super) fn estimate_heterogeneous_summary( } SummaryExpr::SummaryDelete { summary_input, .. } => { let delete = summary_operation_evidence(node, evidence)?; - let StreamingSummaryOperatorEvidence::Delete { + let SummaryOperatorEvidence::Delete { resource: operation, events_per_second, routing_fanout, @@ -613,7 +614,7 @@ fn add_operator_io( fn validate_summary_edges_and_physical_ids( root: &SummaryNode, - evidence: &StreamingNodeEvidence, + evidence: &SummaryNodeEvidence, frameworks_by_node: &HashMap<*const SummaryNode, &Option>, ) -> Result<(), AnalyticalCostError> { fn children(node: &SummaryNode) -> Vec<&SummaryNode> { @@ -643,7 +644,7 @@ fn validate_summary_edges_and_physical_ids( } fn metadata( node: &SummaryNode, - evidence: &StreamingNodeEvidence, + evidence: &SummaryNodeEvidence, ) -> Result<(String, Vec, EdgeStatistics), AnalyticalCostError> { match &node.expr { SummaryExpr::KeepPreAsap(_) => { @@ -683,7 +684,7 @@ fn validate_summary_edges_and_physical_ids( } fn visit( node: &SummaryNode, - evidence: &StreamingNodeEvidence, + evidence: &SummaryNodeEvidence, frameworks_by_node: &HashMap<*const SummaryNode, &Option>, seen: &mut HashSet<*const SummaryNode>, physical: &mut HashMap, EdgeStatistics, String)>, @@ -757,7 +758,7 @@ fn validate_summary_edges_and_physical_ids( fn summary_physical_id( node: &SummaryNode, - evidence: &StreamingNodeEvidence, + evidence: &SummaryNodeEvidence, ) -> Result { match &node.expr { SummaryExpr::KeepPreAsap(_) => evidence @@ -786,7 +787,7 @@ fn summary_physical_id( /// operator workspace and its output buffer coexist during that execution. pub(super) fn estimate_transient_liveness( root: &SummaryNode, - evidence: &StreamingNodeEvidence, + evidence: &SummaryNodeEvidence, ) -> Result { fn children(node: &SummaryNode) -> Vec<&SummaryNode> { match &node.expr { @@ -815,7 +816,7 @@ pub(super) fn estimate_transient_liveness( } fn visit<'a>( node: &'a SummaryNode, - evidence: &StreamingNodeEvidence, + evidence: &SummaryNodeEvidence, seen: &mut HashSet, uses: &mut HashMap, order: &mut Vec<&'a SummaryNode>, @@ -834,7 +835,7 @@ pub(super) fn estimate_transient_liveness( } fn memory( node: &SummaryNode, - evidence: &StreamingNodeEvidence, + evidence: &SummaryNodeEvidence, ) -> Result<(u64, u64), AnalyticalCostError> { match &node.expr { SummaryExpr::KeepPreAsap(_) => evidence @@ -971,7 +972,7 @@ struct SummaryOperationCounts { pub(super) fn estimate_incremental_summary_maintenance( root: &SummaryNode, guarantee: &SummaryMaintenanceLifecycleGuarantee, - inputs: StreamingSummaryInputs, + inputs: SummaryMaintenanceInputs, cpu: SummaryOperationCpuEvidence, scope: &ComparisonScope, ) -> Result { @@ -981,7 +982,7 @@ pub(super) fn estimate_incremental_summary_maintenance( pub(super) fn estimate_incremental_summary_maintenance_with_join( root: &SummaryNode, guarantee: &SummaryMaintenanceLifecycleGuarantee, - inputs: StreamingSummaryInputs, + inputs: SummaryMaintenanceInputs, cpu: SummaryOperationCpuEvidence, join: Option, scope: &ComparisonScope, @@ -1120,10 +1121,11 @@ pub(super) fn estimate_incremental_summary_maintenance_with_join( } pub(super) fn lifecycle_row_counts( - inputs: StreamingSummaryInputs, + inputs: SummaryMaintenanceInputs, guarantee: &SummaryMaintenanceLifecycleGuarantee, scope: &ComparisonScope, ) -> Result<(u64, u64, u64), AnalyticalCostError> { + validate_arrival_rate(scope.data_arrival, inputs.ingestion_rate_per_second)?; let horizon_end = scope .planning_time .0 diff --git a/crates/asap-aware-mapping/src/summary_maintenance_cost/evidence.rs b/crates/asap-aware-mapping/src/summary_maintenance_cost/evidence.rs index ce47f7ae..de20a94a 100644 --- a/crates/asap-aware-mapping/src/summary_maintenance_cost/evidence.rs +++ b/crates/asap-aware-mapping/src/summary_maintenance_cost/evidence.rs @@ -1,11 +1,11 @@ use super::*; /// Physical evidence that is not represented by [`DataWorkload`] for one -/// incrementally maintained summary deployment. Window counts describe the +/// 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 { +pub struct SummaryPhysicalInputEvidence { /// 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 @@ -24,10 +24,10 @@ pub struct StreamingPhysicalInputEvidence { pub state_bytes_per_summary: u64, } -/// Workload-normalized inputs for incremental maintenance over one finite +/// Workload-normalized inputs for summary construction and maintenance over one finite /// comparison horizon. #[derive(Debug, Clone, Copy, PartialEq, Serialize, Deserialize)] -pub struct StreamingSummaryInputs { +pub struct SummaryMaintenanceInputs { pub initial_input_rows: u64, pub initial_input_bytes: u64, pub initial_source_scan_bytes: u64, @@ -39,40 +39,52 @@ pub struct StreamingSummaryInputs { pub state_bytes_per_summary: u64, } -impl StreamingSummaryInputs { +impl SummaryMaintenanceInputs { /// Resolve snapshot size, arriving rows, and reads from the canonical /// workload. Positive fractional expected work rounds up conservatively. /// + /// `AtRest` needs snapshot cardinality and implies zero arrivals; continuous + /// ingestion additionally requires fresh rate evidence. /// `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, + physical: SummaryPhysicalInputEvidence, data: &DataWorkload, scope: &ComparisonScope, ) -> Result { let _ = scope.validate()?; - if data.arrival != DataArrival::ContinuouslyIngesting || scope.data_arrival != data.arrival - { - return Err(AnalyticalCostError::UnsupportedDataArrival(data.arrival)); + if scope.data_arrival != data.arrival { + return Err(AnalyticalCostError::ComparisonScopeMismatch("data arrival")); } let initial_input_rows = data .input_cardinality .value_at(scope.planning_time.0) .copied() .ok_or(AnalyticalCostError::MissingOrStale("input_cardinality"))?; - let ingestion_rate = data - .ingestion_rate - .value_at(scope.planning_time.0) - .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)); - } + let ingestion_rate = match data.arrival { + DataArrival::AtRest => { + // A declared snapshot has no arrivals. Reject contradictory fresh + // evidence rather than silently pricing the wrong workload. + if let Some(rate) = data.ingestion_rate.value_at(scope.planning_time.0) { + validate_arrival_rate(data.arrival, rate.0)?; + } + 0.0 + } + DataArrival::ContinuouslyIngesting => { + data.ingestion_rate + .value_at(scope.planning_time.0) + .copied() + .ok_or(AnalyticalCostError::MissingOrStale("ingestion_rate"))? + .0 + } + arrival => return Err(AnalyticalCostError::UnsupportedDataArrival(arrival)), + }; + validate_arrival_rate(data.arrival, ingestion_rate)?; Self { 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, + ingestion_rate_per_second: ingestion_rate, active_window_count: physical.active_window_count, bootstrap_window_count: physical.bootstrap_window_count, retained_window_count: physical.retained_window_count, @@ -142,7 +154,7 @@ pub struct SummaryJoinEvidence { } #[derive(Debug, Clone, PartialEq)] -pub struct StreamingAggregateEvidence { +pub struct SummaryAggregateEvidence { pub physical_id: String, pub input: EdgeStatistics, pub output: EdgeStatistics, @@ -153,7 +165,7 @@ pub struct StreamingAggregateEvidence { /// Provider-owned identity of the physical bootstrap read. Equal source /// coverage alone does not prove two independent builds share I/O. pub bootstrap_read_identity: String, - pub inputs: StreamingSummaryInputs, + pub inputs: SummaryMaintenanceInputs, /// CPU operations to insert one routed row into one state instance. pub insert_cpu_ops: f64, } @@ -175,7 +187,7 @@ pub struct SummaryOperatorResourceEvidence { /// Evidence is structured by logical summary operation so delete-only facts /// cannot be attached to merge, subtract, or readout nodes. #[derive(Debug, Clone, PartialEq)] -pub enum StreamingSummaryOperatorEvidence { +pub enum SummaryOperatorEvidence { /// Exact query-time arithmetic over two independently realized operands. Binary(SummaryOperatorResourceEvidence), /// Query-time or maintenance-time plain-value work. For `Sort`/`Limit`, @@ -192,7 +204,7 @@ pub enum StreamingSummaryOperatorEvidence { Readout(SummaryOperatorResourceEvidence), } -impl StreamingSummaryOperatorEvidence { +impl SummaryOperatorEvidence { pub(super) fn resource(&self) -> &SummaryOperatorResourceEvidence { match self { Self::Binary(resource) @@ -221,7 +233,7 @@ impl StreamingSummaryOperatorEvidence { /// horizon. Bootstrap/source I/O belongs exclusively to the owning aggregate, /// and summary insertion belongs exclusively to its insert evidence. #[derive(Debug, Clone, PartialEq)] -pub struct StreamingRetainedQueryEvidence { +pub struct RetainedSubDagEvidence { pub physical_id: String, /// Logical output edge consumed by the parent summary operator. pub output: EdgeStatistics, @@ -234,19 +246,19 @@ pub struct StreamingRetainedQueryEvidence { /// Physical evidence bound to the selected DAG's `Rc` identity. A copied, /// structurally equal node is not silently treated as the same deployment. #[derive(Debug, Clone, Default)] -pub struct StreamingNodeEvidence { - pub(super) aggregations: HashMap<*const SummaryNode, StreamingAggregateEvidence>, +pub struct SummaryNodeEvidence { + pub(super) aggregations: HashMap<*const SummaryNode, SummaryAggregateEvidence>, pub(super) joins: HashMap<*const SummaryNode, SummaryJoinEvidence>, - pub(super) operations: HashMap<*const SummaryNode, StreamingSummaryOperatorEvidence>, + pub(super) operations: HashMap<*const SummaryNode, SummaryOperatorEvidence>, pub(super) operation_state_owners: HashMap<*const SummaryNode, *const SummaryNode>, - pub(super) retained_queries: HashMap<*const SummaryNode, StreamingRetainedQueryEvidence>, + pub(super) retained_queries: HashMap<*const SummaryNode, RetainedSubDagEvidence>, } -impl StreamingNodeEvidence { +impl SummaryNodeEvidence { pub fn insert_aggregation( &mut self, node: &Rc, - evidence: StreamingAggregateEvidence, + evidence: SummaryAggregateEvidence, ) { self.aggregations.insert(Rc::as_ptr(node), evidence); } @@ -255,11 +267,7 @@ impl StreamingNodeEvidence { self.joins.insert(Rc::as_ptr(node), evidence); } - pub fn insert_operation( - &mut self, - node: &Rc, - evidence: StreamingSummaryOperatorEvidence, - ) { + pub fn insert_operation(&mut self, node: &Rc, evidence: SummaryOperatorEvidence) { self.operations.insert(Rc::as_ptr(node), evidence); } @@ -269,7 +277,7 @@ impl StreamingNodeEvidence { &mut self, node: &Rc, state: &Rc, - evidence: StreamingSummaryOperatorEvidence, + evidence: SummaryOperatorEvidence, ) { self.operations.insert(Rc::as_ptr(node), evidence); self.operation_state_owners @@ -279,20 +287,20 @@ impl StreamingNodeEvidence { pub fn insert_retained_query( &mut self, node: &Rc, - evidence: StreamingRetainedQueryEvidence, + evidence: RetainedSubDagEvidence, ) { self.retained_queries.insert(Rc::as_ptr(node), evidence); } - pub(super) fn aggregation(&self, node: &SummaryNode) -> Option { + pub(super) fn aggregation(&self, node: &SummaryNode) -> Option { self.aggregations.get(&(node as *const _)).cloned() } } pub(super) fn summary_operation_evidence<'a>( node: &SummaryNode, - evidence: &'a StreamingNodeEvidence, -) -> Result<&'a StreamingSummaryOperatorEvidence, AnalyticalCostError> { + evidence: &'a SummaryNodeEvidence, +) -> Result<&'a SummaryOperatorEvidence, AnalyticalCostError> { let operation = evidence .operations .get(&(node as *const _)) @@ -301,22 +309,22 @@ pub(super) fn summary_operation_evidence<'a>( (&node.expr, operation), ( SummaryExpr::BinaryOp { .. }, - StreamingSummaryOperatorEvidence::Binary(_) + SummaryOperatorEvidence::Binary(_) ) | ( SummaryExpr::ValueOperation { .. }, - StreamingSummaryOperatorEvidence::ValueOperation(_) + SummaryOperatorEvidence::ValueOperation(_) ) | ( SummaryExpr::SummaryMerge { .. }, - StreamingSummaryOperatorEvidence::Merge(_) + SummaryOperatorEvidence::Merge(_) ) | ( SummaryExpr::SummarySubtract { .. }, - StreamingSummaryOperatorEvidence::Subtract(_) + SummaryOperatorEvidence::Subtract(_) ) | ( SummaryExpr::SummaryDelete { .. }, - StreamingSummaryOperatorEvidence::Delete { .. } + SummaryOperatorEvidence::Delete { .. } ) | ( SummaryExpr::SummaryEstimate { .. }, - StreamingSummaryOperatorEvidence::Readout(_) + SummaryOperatorEvidence::Readout(_) ) ); if matches { @@ -333,7 +341,7 @@ pub(super) fn summary_operation_evidence<'a>( /// evaluation adds arrivals since planning time; `physical_dag` is therefore /// a once-counted DAG whose edge statistics already aggregate all evaluations. #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -pub struct StreamingRawInputEvidence { +pub struct RawInputEvidence { pub planning_time_input_rows: u64, pub planning_time_input_bytes: u64, pub planning_time_source_scan_bytes: u64, @@ -347,10 +355,27 @@ pub struct StreamingRawInputEvidence { } /// One complete provider-enumerated physical implementation of the selected -/// streaming summary DAG. The identifier is stable provenance; concrete +/// summary DAG. The identifier is stable provenance; concrete /// framework selection is performed by ranking these complete alternatives. #[derive(Debug, Clone)] -pub struct StreamingPhysicalPlanAlternative { +pub struct SummaryPhysicalPlanAlternative { pub physical_plan_id: String, - pub node_evidence: StreamingNodeEvidence, + pub node_evidence: SummaryNodeEvidence, +} + +/// Apply arrival semantics to both workload-derived and directly bound evidence. +pub(super) fn validate_arrival_rate( + arrival: DataArrival, + rate: f64, +) -> Result<(), AnalyticalCostError> { + if !rate.is_finite() || rate < 0.0 { + return Err(AnalyticalCostError::InvalidIngestionRate(rate)); + } + match arrival { + DataArrival::AtRest if rate != 0.0 => Err(AnalyticalCostError::ComparisonScopeMismatch( + "at-rest ingestion rate", + )), + DataArrival::AtRest | DataArrival::ContinuouslyIngesting => Ok(()), + other => Err(AnalyticalCostError::UnsupportedDataArrival(other)), + } } diff --git a/crates/asap-aware-mapping/src/summary_maintenance_cost/mod.rs b/crates/asap-aware-mapping/src/summary_maintenance_cost/mod.rs index 23dd7e52..6e4901f4 100644 --- a/crates/asap-aware-mapping/src/summary_maintenance_cost/mod.rs +++ b/crates/asap-aware-mapping/src/summary_maintenance_cost/mod.rs @@ -1,4 +1,4 @@ -//! Analytical resource cost for incrementally maintained summary deployments. +//! Analytical resource cost for at-rest and 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, @@ -40,7 +40,7 @@ use crate::summary_maintenance_lifecycle::{ SummaryMaintenanceLifecycleCostInputs, }; -pub const SUMMARY_MAINTENANCE_COST_MODEL_VERSION: &str = "summary-maintenance-resource-v1"; +pub const SUMMARY_MAINTENANCE_COST_MODEL_VERSION: &str = "summary-maintenance-resource-v2"; mod estimator; mod evidence; diff --git a/crates/asap-aware-mapping/src/summary_maintenance_cost/model.rs b/crates/asap-aware-mapping/src/summary_maintenance_cost/model.rs index 46753516..8bbca6cc 100644 --- a/crates/asap-aware-mapping/src/summary_maintenance_cost/model.rs +++ b/crates/asap-aware-mapping/src/summary_maintenance_cost/model.rs @@ -1,18 +1,18 @@ use super::*; /// 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. +/// summary costs across at-rest and continuously ingesting workloads. The +/// planner's existing lifecycle enums and legality checks remain authoritative. #[derive(Debug, Clone)] pub struct SummaryMaintenanceCostModel { - pub node_evidence: StreamingNodeEvidence, + pub node_evidence: SummaryNodeEvidence, pub calibration: ResourceCalibration, pub capabilities: SummaryMaintenanceCapabilities, - target_comparisons: HashMap<*const QueryExpr, StreamingTargetComparison>, + target_comparisons: HashMap<*const QueryExpr, SummaryTargetComparison>, candidate_comparisons: HashMap, physical_plan_alternatives: - HashMap>, + HashMap>, window_framework_candidates: - HashMap>, + HashMap>, } type CandidateComparisonKey = (*const QueryExpr, *const SummaryNode); @@ -24,10 +24,10 @@ struct BoundCandidateIdentity { } #[derive(Debug, Clone)] -struct StreamingTargetComparison { +struct SummaryTargetComparison { _target: Rc, scope: ComparisonScope, - raw: StreamingRawInputEvidence, + raw: RawInputEvidence, } pub(super) type LogicalSourceSelection = (Source, Vec, Vec); @@ -231,9 +231,10 @@ fn reachable_physical_nodes( } fn validate_raw_snapshot_dimensions( - raw: &StreamingRawInputEvidence, + raw: &RawInputEvidence, scope: &ComparisonScope, ) -> Result<(), AnalyticalCostError> { + validate_arrival_rate(scope.data_arrival, raw.ingestion_rate_per_second)?; if scope.sources.len() != 1 { return Err(AnalyticalCostError::MissingComparisonScope( "single-source raw evolution", @@ -330,7 +331,7 @@ fn validate_raw_snapshot_dimensions( } pub(super) fn ephemeral_rows_over_horizon( - inputs: StreamingSummaryInputs, + inputs: SummaryMaintenanceInputs, scope: &ComparisonScope, ) -> Result { evaluation_offsets_ms(scope)? @@ -348,8 +349,8 @@ pub(super) fn ephemeral_rows_over_horizon( } pub(super) fn ephemeral_scan_bytes_over_horizon( - inputs: StreamingSummaryInputs, - raw: &StreamingRawInputEvidence, + inputs: SummaryMaintenanceInputs, + raw: &RawInputEvidence, scope: &ComparisonScope, ) -> Result { if inputs.initial_source_scan_bytes == 0 { @@ -412,7 +413,7 @@ impl SummaryMaintenanceCostModel { capabilities: SummaryMaintenanceCapabilities, ) -> Self { Self { - node_evidence: StreamingNodeEvidence::default(), + node_evidence: SummaryNodeEvidence::default(), calibration, capabilities, target_comparisons: HashMap::new(), @@ -430,7 +431,7 @@ impl SummaryMaintenanceCostModel { target: &Rc, root: &Rc, scope: ComparisonScope, - raw: StreamingRawInputEvidence, + raw: RawInputEvidence, ) -> Result<(), AnalyticalCostError> { scope.validate()?; validate_query_scope(target, &scope)?; @@ -454,7 +455,7 @@ impl SummaryMaintenanceCostModel { // not carry one owning target; context identity is `(target, root)`. self.target_comparisons .entry(target_ptr) - .or_insert(StreamingTargetComparison { + .or_insert(SummaryTargetComparison { _target: Rc::clone(target), scope, raw, @@ -475,7 +476,7 @@ impl SummaryMaintenanceCostModel { &mut self, target: &Rc, root: &Rc, - alternative: StreamingPhysicalPlanAlternative, + alternative: SummaryPhysicalPlanAlternative, ) -> Result<(), AnalyticalCostError> { let key = (Rc::as_ptr(target), Rc::as_ptr(root)); if !self.candidate_comparisons.contains_key(&key) { @@ -513,7 +514,7 @@ impl SummaryMaintenanceCostModel { &mut self, target: &Rc, root: &Rc, - candidate: StreamingWindowFrameworkCandidate, + candidate: SummaryWindowFrameworkCandidate, ) -> Result<(), AnalyticalCostError> { let key = (Rc::as_ptr(target), Rc::as_ptr(root)); if !self.candidate_comparisons.contains_key(&key) { @@ -568,7 +569,7 @@ impl SummaryMaintenanceCostModel { target: Option<&QueryExpr>, horizon: Option, expected_reads: Option, - ) -> Option<(CandidateComparisonKey, &StreamingTargetComparison)> { + ) -> Option<(CandidateComparisonKey, &SummaryTargetComparison)> { let root_ptr = root as *const _; let target_ptr = match target { Some(target) => target as *const _, @@ -601,8 +602,8 @@ impl SummaryMaintenanceCostModel { &self, root: &SummaryNode, deployments: &[CostedSummaryDeployment<'_>], - comparison: &StreamingTargetComparison, - evidence: &StreamingNodeEvidence, + comparison: &SummaryTargetComparison, + evidence: &SummaryNodeEvidence, window_frameworks: &[Option], ) -> Option { self.calibrated( @@ -618,7 +619,7 @@ impl SummaryMaintenanceCostModel { ) } - fn canonical_inputs(&self, summary: &SummaryNode) -> Option { + fn canonical_inputs(&self, summary: &SummaryNode) -> Option { let evidence = self.node_evidence.aggregation(summary)?; evidence.inputs.validate().ok()?; Some(evidence) @@ -904,7 +905,7 @@ mod tests { fn estimate_test( root: &SummaryNode, guarantee: &SummaryMaintenanceLifecycleGuarantee, - inputs: StreamingSummaryInputs, + inputs: SummaryMaintenanceInputs, cpu: SummaryOperationCpuEvidence, ) -> Result { estimate_incremental_summary_maintenance(root, guarantee, inputs, cpu, &streaming_scope()) @@ -913,7 +914,7 @@ mod tests { fn estimate_join_test( root: &SummaryNode, guarantee: &SummaryMaintenanceLifecycleGuarantee, - inputs: StreamingSummaryInputs, + inputs: SummaryMaintenanceInputs, cpu: SummaryOperationCpuEvidence, join: Option, ) -> Result { @@ -950,8 +951,8 @@ mod tests { .unwrap() } - fn physical() -> StreamingPhysicalInputEvidence { - StreamingPhysicalInputEvidence { + fn physical() -> SummaryPhysicalInputEvidence { + SummaryPhysicalInputEvidence { initial_input_bytes: 640, initial_source_scan_bytes: 640, active_window_count: 2, @@ -987,6 +988,158 @@ mod tests { } } + /// A fixed snapshot needs cardinality evidence, but no stream-rate evidence. + #[test] + fn at_rest_workload_adapter_builds_once_without_arrivals() { + let mut data = streaming_data_workload(); + data.arrival = DataArrival::AtRest; + data.ingestion_rate = Evidence::default(); + data.input_cardinality = Evidence { + value: Some(10), + source: EvidenceSource::Declared, + ..Default::default() + }; + let scope = scope_for(&data, &query(), 0, 5_000); + let inputs = SummaryMaintenanceInputs::from_workload(physical(), &data, &scope).unwrap(); + assert_eq!(inputs.initial_input_rows, 10); + assert_eq!(inputs.ingestion_rate_per_second, 0.0); + let guarantee = SummaryMaintenanceLifecycleGuarantee { + summary_maintenance_lifecycle: SummaryMaintenanceLifecycle::Shared { + retention: asap_types::workload::DurationMs(5_000), + }, + summary_maintenance_mode: SummaryMaintenanceMode::DirectBuild, + evaluation_schedule: EvaluationSchedule::OnRead, + output_representation: OutputRepresentation::SummaryState, + }; + assert_eq!( + lifecycle_row_counts(inputs, &guarantee, &scope).unwrap(), + (10, 0, 5_000) + ); + } + + /// Arrival semantics cannot be overridden by missing or contradictory rate evidence. + #[test] + fn workload_adapter_checks_arrival_scope_and_rate_evidence() { + let mut data = streaming_data_workload(); + data.input_cardinality = Evidence { + value: Some(10), + source: EvidenceSource::Declared, + ..Default::default() + }; + let mut scope = scope_for(&data, &query(), 0, 5_000); + data.ingestion_rate = Evidence::default(); + assert_eq!( + SummaryMaintenanceInputs::from_workload(physical(), &data, &scope), + Err(AnalyticalCostError::MissingOrStale("ingestion_rate")) + ); + data.arrival = DataArrival::AtRest; + assert_eq!( + SummaryMaintenanceInputs::from_workload(physical(), &data, &scope), + Err(AnalyticalCostError::ComparisonScopeMismatch("data arrival")) + ); + scope.data_arrival = DataArrival::AtRest; + for rate in [1.0, -1.0, f64::INFINITY, f64::NAN] { + data.ingestion_rate = Evidence { + value: Some(Rate(rate)), + source: EvidenceSource::Declared, + ..Default::default() + }; + assert!(SummaryMaintenanceInputs::from_workload(physical(), &data, &scope).is_err()); + } + data.ingestion_rate = Evidence::default(); + data.input_cardinality = Evidence::default(); + assert_eq!( + SummaryMaintenanceInputs::from_workload(physical(), &data, &scope), + Err(AnalyticalCostError::MissingOrStale("input_cardinality")) + ); + } + + /// The real lifecycle planner costs a fixed snapshot with the same node evidence API. + #[test] + fn lifecycle_planner_selects_fully_costed_at_rest_summary() { + let workload = streaming_workload(); + let mut data = streaming_data_workload(); + data.arrival = DataArrival::AtRest; + data.ingestion_rate = Evidence::default(); + data.input_cardinality = Evidence { + value: Some(10), + source: EvidenceSource::Declared, + ..Default::default() + }; + let mut scope = streaming_scope(); + scope.data_arrival = DataArrival::AtRest; + let inputs = SummaryMaintenanceInputs::from_workload(physical(), &data, &scope).unwrap(); + let target = streaming_sum_query(); + let root = summary_with_operations(false, false, false); + let mut provider = streaming_model(); + bind_aggregations(&mut provider, &target, &root, inputs, streaming_cpu()); + let mut model = streaming_model(); + model.node_evidence = provider.node_evidence; + let mut raw = streaming_raw(); + raw.ingestion_rate_per_second = 0.0; + let edge = EdgeStatistics { + rows: 50, + bytes: 3_200, + }; + raw.physical_dag + .evidence + .get_mut("raw-scan") + .unwrap() + .statistics = OperatorStatistics::Scan { + source_read_bytes: 3_200, + edges: UnaryEdgeStatistics { + input: edge, + output: edge, + promql: None, + }, + }; + model + .bind_candidate_comparison(&target, &root, scope.clone(), raw.clone()) + .unwrap(); + let plan = plan_summary_maintenance_lifecycles( + Rc::clone(&root), + WorkloadDemand::new_with_data(&workload, &data, &[0]), + 0, + Some(Horizon(5.0)), + SummaryMaintenanceLifecycleCapabilities::ALL, + &model, + ) + .unwrap(); + assert!(plan.summary_total_cost.is_some()); + assert!(!plan.selected_raw_recompute); + assert_eq!( + plan.deployments[0] + .summary_maintenance_lifecycle_guarantee + .as_ref() + .unwrap() + .summary_maintenance_mode, + SummaryMaintenanceMode::DirectBuild + ); + + // Directly supplied raw evidence must not bypass the workload invariant. + raw.ingestion_rate_per_second = 1.0; + assert_eq!( + streaming_model().bind_candidate_comparison(&target, &root, scope, raw), + Err(AnalyticalCostError::ComparisonScopeMismatch( + "at-rest ingestion rate" + )) + ); + // Nor may a provider hide arrivals on a summary edge. + for aggregation in model.node_evidence.aggregations.values_mut() { + aggregation.inputs.ingestion_rate_per_second = 1.0; + } + let invalid = plan_summary_maintenance_lifecycles( + root, + WorkloadDemand::new_with_data(&workload, &data, &[0]), + 0, + Some(Horizon(5.0)), + SummaryMaintenanceLifecycleCapabilities::ALL, + &model, + ) + .unwrap(); + assert_eq!(invalid.summary_total_cost, None); + } + #[test] fn workload_adapter_derives_updates_and_reads_over_one_horizon() { let data = DataWorkload { @@ -1007,7 +1160,7 @@ mod tests { }; let scope = scope_for(&data, &query(), 100, 5_000); - let inputs = StreamingSummaryInputs::from_workload(physical(), &data, &scope).unwrap(); + let inputs = SummaryMaintenanceInputs::from_workload(physical(), &data, &scope).unwrap(); assert_eq!(inputs.initial_input_rows, 10); assert_eq!( lifecycle_row_counts(inputs, &continuous_guarantee(), &scope) @@ -1040,7 +1193,7 @@ mod tests { empty.initial_input_bytes = 0; empty.initial_source_scan_bytes = 0; let scope = scope_for(&data, &query(), 0, 5_000); - let inputs = StreamingSummaryInputs::from_workload(empty, &data, &scope).unwrap(); + let inputs = SummaryMaintenanceInputs::from_workload(empty, &data, &scope).unwrap(); let estimate = estimate_test( &summary_with_operations(false, false, false), &continuous_guarantee(), @@ -1059,7 +1212,7 @@ mod tests { #[test] fn bootstrap_rows_and_bytes_must_be_present_together() { - let mut inputs = StreamingSummaryInputs { + let mut inputs = SummaryMaintenanceInputs { initial_input_rows: 0, initial_input_bytes: 8, initial_source_scan_bytes: 0, @@ -1126,7 +1279,7 @@ mod tests { #[test] fn existing_lifecycle_planner_selects_a_fully_costed_streaming_alternative() { - let inputs = StreamingSummaryInputs { + let inputs = SummaryMaintenanceInputs { initial_input_rows: 10, initial_input_bytes: 640, initial_source_scan_bytes: 640, @@ -1240,11 +1393,11 @@ mod tests { aggregate.inputs.retained_window_count = 2; } for alternative in [ - StreamingPhysicalPlanAlternative { + SummaryPhysicalPlanAlternative { physical_plan_id: "high-retention-layout".into(), node_evidence: high_retention, }, - StreamingPhysicalPlanAlternative { + SummaryPhysicalPlanAlternative { physical_plan_id: "low-retention-layout".into(), node_evidence: low_retention, }, @@ -1900,7 +2053,7 @@ mod tests { bind_comparison(&mut model, &target, &root); model.node_evidence.aggregations.insert( aggregations[0] as *const _, - StreamingAggregateEvidence { + SummaryAggregateEvidence { physical_id: "left-state".into(), input: test_edge(), output: test_edge(), @@ -1927,7 +2080,7 @@ mod tests { second_cpu.insert_cpu_ops = Some(5.0); model.node_evidence.aggregations.insert( aggregations[1] as *const _, - StreamingAggregateEvidence { + SummaryAggregateEvidence { physical_id: "right-state".into(), input: test_edge(), output: test_edge(), @@ -1952,7 +2105,7 @@ mod tests { ); model.node_evidence.operations.insert( Rc::as_ptr(&root), - StreamingSummaryOperatorEvidence::Readout(SummaryOperatorResourceEvidence { + SummaryOperatorEvidence::Readout(SummaryOperatorResourceEvidence { physical_id: "root-readout".into(), inputs: vec![test_edge()], output: test_edge(), @@ -2060,31 +2213,31 @@ mod tests { aggregate.inputs.retained_window_count = 2; } for candidate in [ - StreamingWindowFrameworkCandidate { + SummaryWindowFrameworkCandidate { physical_plan_id: "tumbling-v1".into(), - assignments: vec![StreamingWindowFrameworkAssignment { + assignments: vec![SummaryWindowFrameworkAssignment { summary: Rc::clone(&windowed_summary), framework: Some(SummaryWindowFramework::Tumbling), }], - accuracy: StreamingWindowAccuracyEvidence::Exact, + accuracy: SummaryWindowAccuracyEvidence::Exact, node_evidence: tumbling, }, - StreamingWindowFrameworkCandidate { + SummaryWindowFrameworkCandidate { physical_plan_id: "sliding-v1".into(), - assignments: vec![StreamingWindowFrameworkAssignment { + assignments: vec![SummaryWindowFrameworkAssignment { summary: Rc::clone(&windowed_summary), framework: Some(SummaryWindowFramework::Sliding), }], - accuracy: StreamingWindowAccuracyEvidence::Exact, + accuracy: SummaryWindowAccuracyEvidence::Exact, node_evidence: sliding, }, - StreamingWindowFrameworkCandidate { + SummaryWindowFrameworkCandidate { physical_plan_id: "eh-v1".into(), - assignments: vec![StreamingWindowFrameworkAssignment { + assignments: vec![SummaryWindowFrameworkAssignment { summary: Rc::clone(&windowed_summary), framework: Some(SummaryWindowFramework::ExponentialHistogram), }], - accuracy: StreamingWindowAccuracyEvidence::ExponentialHistogram( + accuracy: SummaryWindowAccuracyEvidence::ExponentialHistogram( ExponentialHistogramAccuracyEvidence::UniversalGsum { epsilon: 0.05, failure_probability: 0.01, @@ -2101,7 +2254,7 @@ mod tests { // Framework candidates are authoritative. Selection must not depend // on duplicating one arbitrary implementation into the legacy global // evidence map. - model.node_evidence = StreamingNodeEvidence::default(); + model.node_evidence = SummaryNodeEvidence::default(); let plan = plan_summary_maintenance_lifecycles( Rc::clone(&root), @@ -2183,22 +2336,22 @@ mod tests { let empty = model.bind_window_framework_candidate( &target, &root, - StreamingWindowFrameworkCandidate { + SummaryWindowFrameworkCandidate { physical_plan_id: "empty-assignments".into(), assignments: vec![], - accuracy: StreamingWindowAccuracyEvidence::Exact, + accuracy: SummaryWindowAccuracyEvidence::Exact, node_evidence: model.node_evidence.clone(), }, ); assert!(matches!(empty, Err(AnalyticalCostError::MissingOrZero(_)))); - let candidate = StreamingWindowFrameworkCandidate { + let candidate = SummaryWindowFrameworkCandidate { physical_plan_id: "tumbling-v1".into(), - assignments: vec![StreamingWindowFrameworkAssignment { + assignments: vec![SummaryWindowFrameworkAssignment { summary: windowed_summary, framework: Some(SummaryWindowFramework::Tumbling), }], - accuracy: StreamingWindowAccuracyEvidence::Exact, + accuracy: SummaryWindowAccuracyEvidence::Exact, node_evidence: model.node_evidence.clone(), }; model @@ -2274,19 +2427,19 @@ mod tests { .insert(Rc::as_ptr(child), shared_retained.clone()); } - let candidate = StreamingWindowFrameworkCandidate { + let candidate = SummaryWindowFrameworkCandidate { physical_plan_id: "mixed-framework-join".into(), assignments: vec![ - StreamingWindowFrameworkAssignment { + SummaryWindowFrameworkAssignment { summary: Rc::clone(&aggregation_nodes[0]), framework: Some(SummaryWindowFramework::Tumbling), }, - StreamingWindowFrameworkAssignment { + SummaryWindowFrameworkAssignment { summary: Rc::clone(&aggregation_nodes[1]), framework: Some(SummaryWindowFramework::Sliding), }, ], - accuracy: StreamingWindowAccuracyEvidence::Exact, + accuracy: SummaryWindowAccuracyEvidence::Exact, node_evidence: model.node_evidence.clone(), }; model @@ -2307,7 +2460,7 @@ mod tests { #[test] fn promsketch_eh_accuracy_composes_registered_full_and_subwindow_bounds() { - let full = StreamingWindowAccuracyEvidence::ExponentialHistogram( + let full = SummaryWindowAccuracyEvidence::ExponentialHistogram( ExponentialHistogramAccuracyEvidence::KllRank { eh_epsilon: 0.01, kll_epsilon: 0.02, @@ -2320,7 +2473,7 @@ mod tests { assert_eq!(full.metric, ErrorMetric::Rank); assert!((full.bound.evaluate().unwrap() - 0.04).abs() < f64::EPSILON); - let subwindow = StreamingWindowAccuracyEvidence::ExponentialHistogram( + let subwindow = SummaryWindowAccuracyEvidence::ExponentialHistogram( ExponentialHistogramAccuracyEvidence::KllRank { eh_epsilon: 0.01, kll_epsilon: 0.02, @@ -2335,7 +2488,7 @@ mod tests { .unwrap(); assert!((subwindow.bound.evaluate().unwrap() - 0.10).abs() < f64::EPSILON); - let gsum = StreamingWindowAccuracyEvidence::ExponentialHistogram( + let gsum = SummaryWindowAccuracyEvidence::ExponentialHistogram( ExponentialHistogramAccuracyEvidence::UniversalGsum { epsilon: 0.05, failure_probability: 0.30, @@ -2353,7 +2506,7 @@ mod tests { #[test] fn eh_accuracy_rejects_negative_components_and_mismatched_summary_guarantees() { - let evidence = StreamingWindowAccuracyEvidence::ExponentialHistogram( + let evidence = SummaryWindowAccuracyEvidence::ExponentialHistogram( ExponentialHistogramAccuracyEvidence::KllRank { eh_epsilon: -0.01, kll_epsilon: 0.03, @@ -2363,7 +2516,7 @@ mod tests { ); assert!(evidence.guarantee(true).is_none()); - let evidence = StreamingWindowAccuracyEvidence::ExponentialHistogram( + let evidence = SummaryWindowAccuracyEvidence::ExponentialHistogram( ExponentialHistogramAccuracyEvidence::KllRank { eh_epsilon: 0.01, kll_epsilon: 0.02, @@ -2406,13 +2559,13 @@ mod tests { .bind_window_framework_candidate( &target, &root, - StreamingWindowFrameworkCandidate { + SummaryWindowFrameworkCandidate { physical_plan_id: "invalid-eh".into(), - assignments: vec![StreamingWindowFrameworkAssignment { + assignments: vec![SummaryWindowFrameworkAssignment { summary: Rc::clone(summary_input), framework: Some(SummaryWindowFramework::ExponentialHistogram), }], - accuracy: StreamingWindowAccuracyEvidence::Exact, + accuracy: SummaryWindowAccuracyEvidence::Exact, node_evidence: model.node_evidence.clone(), }, ) @@ -2566,7 +2719,7 @@ mod tests { let mut scope = streaming_scope(); scope.data_arrival = DataArrival::Mixed; assert_eq!( - StreamingSummaryInputs::from_workload(physical(), &data, &scope), + SummaryMaintenanceInputs::from_workload(physical(), &data, &scope), Err(AnalyticalCostError::UnsupportedDataArrival( DataArrival::Mixed )) @@ -2578,7 +2731,7 @@ mod tests { let estimate = estimate_test( &summary_with_operations(false, false, false), &continuous_guarantee(), - StreamingSummaryInputs { + SummaryMaintenanceInputs { initial_input_rows: 10, initial_input_bytes: 640, initial_source_scan_bytes: 640, @@ -2607,7 +2760,7 @@ mod tests { let estimate = estimate_test( &summary_with_operations(true, true, true), &continuous_guarantee(), - StreamingSummaryInputs { + SummaryMaintenanceInputs { initial_input_rows: 1, initial_input_bytes: 8, initial_source_scan_bytes: 8, @@ -2642,7 +2795,7 @@ mod tests { estimate_test( &summary_with_operations(false, false, false), &guarantee, - StreamingSummaryInputs { + SummaryMaintenanceInputs { initial_input_rows: 1, initial_input_bytes: 8, initial_source_scan_bytes: 8, @@ -2669,7 +2822,7 @@ mod tests { estimate_test( &summary_with_operations(true, false, false), &continuous_guarantee(), - StreamingSummaryInputs { + SummaryMaintenanceInputs { initial_input_rows: 1, initial_input_bytes: 8, initial_source_scan_bytes: 8, @@ -2700,7 +2853,7 @@ mod tests { estimate_test( &summary_with_operations(false, false, false), &guarantee, - StreamingSummaryInputs { + SummaryMaintenanceInputs { initial_input_rows: 1, initial_input_bytes: 8, initial_source_scan_bytes: 8, @@ -2735,7 +2888,7 @@ mod tests { let estimate = estimate_test( &summary_with_operations(false, false, false), &guarantee, - StreamingSummaryInputs { + SummaryMaintenanceInputs { initial_input_rows: 10, initial_input_bytes: 80, initial_source_scan_bytes: 80, @@ -2771,7 +2924,7 @@ mod tests { assert!(estimate_test( &summary_with_operations(false, false, false), &guarantee, - StreamingSummaryInputs { + SummaryMaintenanceInputs { initial_input_rows: 1, initial_input_bytes: 8, initial_source_scan_bytes: 8, @@ -2815,7 +2968,7 @@ mod tests { #[test] fn summary_join_requires_cardinality_and_working_memory_evidence() { let joined = summary_join(); - let inputs = StreamingSummaryInputs { + let inputs = SummaryMaintenanceInputs { initial_input_rows: 1, initial_input_bytes: 8, initial_source_scan_bytes: 8, @@ -3112,7 +3265,7 @@ mod tests { ) } - fn streaming_raw() -> StreamingRawInputEvidence { + fn streaming_raw() -> RawInputEvidence { let scope = streaming_scope(); let node = PhysicalDagNode { id: "raw-scan".into(), @@ -3135,7 +3288,7 @@ mod tests { promql: None, }, }; - StreamingRawInputEvidence { + RawInputEvidence { planning_time_input_rows: 10, planning_time_input_bytes: 640, planning_time_source_scan_bytes: 640, @@ -3177,7 +3330,7 @@ mod tests { SummaryExpr::KeepPreAsap(_) => { model.node_evidence.insert_retained_query( node, - StreamingRetainedQueryEvidence { + RetainedSubDagEvidence { physical_id: format!("retained-{node:p}"), output: test_edge(), preprocessing_cpu_ops_over_horizon: 1.0, @@ -3217,8 +3370,8 @@ mod tests { retained(model, root, &mut HashSet::new()); } - fn streaming_inputs() -> StreamingSummaryInputs { - StreamingSummaryInputs { + fn streaming_inputs() -> SummaryMaintenanceInputs { + SummaryMaintenanceInputs { initial_input_rows: 10, initial_input_bytes: 640, initial_source_scan_bytes: 640, @@ -3247,7 +3400,7 @@ mod tests { model: &mut SummaryMaintenanceCostModel, target: &Rc, root: &Rc, - inputs: StreamingSummaryInputs, + inputs: SummaryMaintenanceInputs, cpu: SummaryOperationCpuEvidence, ) { bind_comparison(model, target, root); @@ -3265,7 +3418,7 @@ mod tests { } model.node_evidence.aggregations.insert( node as *const _, - StreamingAggregateEvidence { + SummaryAggregateEvidence { physical_id: format!("agg-{node:p}"), input: test_edge(), output: test_edge(), @@ -3284,7 +3437,7 @@ mod tests { model: &mut SummaryMaintenanceCostModel, node: &SummaryNode, seen: &mut HashSet<*const SummaryNode>, - inputs: StreamingSummaryInputs, + inputs: SummaryMaintenanceInputs, cpu: SummaryOperationCpuEvidence, ) { if !seen.insert(node as *const _) { @@ -3292,7 +3445,7 @@ mod tests { } let operation = match &node.expr { SummaryExpr::BinaryOp { .. } => cpu.readout_cpu_ops.map(|cpu_ops| { - StreamingSummaryOperatorEvidence::Binary(SummaryOperatorResourceEvidence { + SummaryOperatorEvidence::Binary(SummaryOperatorResourceEvidence { physical_id: format!("binary-{node:p}"), inputs: vec![test_edge(), test_edge()], output: test_edge(), @@ -3304,21 +3457,19 @@ mod tests { }) }), SummaryExpr::ValueOperation { .. } => cpu.readout_cpu_ops.map(|cpu_ops| { - StreamingSummaryOperatorEvidence::ValueOperation( - SummaryOperatorResourceEvidence { - physical_id: format!("value-operation-{node:p}"), - inputs: vec![test_edge()], - output: test_edge(), - cpu_ops, - working_memory_bytes: 0, - output_buffer_bytes: 0, - executions_per_evaluation: 1, - io_bytes_per_execution: Some(0), - }, - ) + SummaryOperatorEvidence::ValueOperation(SummaryOperatorResourceEvidence { + physical_id: format!("value-operation-{node:p}"), + inputs: vec![test_edge()], + output: test_edge(), + cpu_ops, + working_memory_bytes: 0, + output_buffer_bytes: 0, + executions_per_evaluation: 1, + io_bytes_per_execution: Some(0), + }) }), SummaryExpr::SummaryMerge { .. } => cpu.merge_cpu_ops.map(|cpu_ops| { - StreamingSummaryOperatorEvidence::Merge(SummaryOperatorResourceEvidence { + SummaryOperatorEvidence::Merge(SummaryOperatorResourceEvidence { physical_id: format!("merge-{node:p}"), inputs: match &node.expr { SummaryExpr::SummaryMerge { children, .. } => { @@ -3335,7 +3486,7 @@ mod tests { }) }), SummaryExpr::SummarySubtract { .. } => cpu.subtract_cpu_ops.map(|cpu_ops| { - StreamingSummaryOperatorEvidence::Subtract(SummaryOperatorResourceEvidence { + SummaryOperatorEvidence::Subtract(SummaryOperatorResourceEvidence { physical_id: format!("subtract-{node:p}"), inputs: vec![test_edge(), test_edge()], output: test_edge(), @@ -3347,7 +3498,7 @@ mod tests { }) }), SummaryExpr::SummaryDelete { .. } => cpu.delete_cpu_ops.and_then(|cpu_ops| { - Some(StreamingSummaryOperatorEvidence::Delete { + Some(SummaryOperatorEvidence::Delete { resource: SummaryOperatorResourceEvidence { physical_id: format!("delete-{node:p}"), inputs: vec![test_edge()], @@ -3363,7 +3514,7 @@ mod tests { }) }), SummaryExpr::SummaryEstimate { .. } => cpu.readout_cpu_ops.map(|cpu_ops| { - StreamingSummaryOperatorEvidence::Readout(SummaryOperatorResourceEvidence { + SummaryOperatorEvidence::Readout(SummaryOperatorResourceEvidence { physical_id: format!("readout-{node:p}"), inputs: vec![test_edge()], output: test_edge(), diff --git a/crates/asap-aware-mapping/src/summary_maintenance_cost/window.rs b/crates/asap-aware-mapping/src/summary_maintenance_cost/window.rs index eff12cf3..867a60e5 100644 --- a/crates/asap-aware-mapping/src/summary_maintenance_cost/window.rs +++ b/crates/asap-aware-mapping/src/summary_maintenance_cost/window.rs @@ -2,7 +2,7 @@ use super::*; /// One per-state window choice within a complete Planner candidate. #[derive(Debug, Clone)] -pub struct StreamingWindowFrameworkAssignment { +pub struct SummaryWindowFrameworkAssignment { pub summary: Rc, /// `None` explicitly means that this state is not window-organized. pub framework: Option, @@ -16,17 +16,17 @@ pub struct StreamingWindowFrameworkAssignment { /// provenance for the chosen implementation, while deployment placement and /// runtime configuration remain downstream concerns. #[derive(Debug, Clone)] -pub struct StreamingWindowFrameworkCandidate { +pub struct SummaryWindowFrameworkCandidate { /// Stable identity of the complete provider implementation whose evidence /// is bound to this planner-visible framework assignment. pub physical_plan_id: String, /// Exactly one assignment for every summary deployment in the DAG. - pub assignments: Vec, + pub assignments: Vec, /// Registered end-to-end accuracy composition for this complete window /// assignment. EH combinations must use one of the specialized proofs; /// unknown combinations fail closed. - pub accuracy: StreamingWindowAccuracyEvidence, - pub node_evidence: StreamingNodeEvidence, + pub accuracy: SummaryWindowAccuracyEvidence, + pub node_evidence: SummaryNodeEvidence, } pub(super) fn summary_aggregation_identities(root: &SummaryNode) -> HashSet<*const SummaryNode> { @@ -105,7 +105,7 @@ pub enum ExponentialHistogramAccuracyEvidence { } #[derive(Debug, Clone, PartialEq)] -pub enum StreamingWindowAccuracyEvidence { +pub enum SummaryWindowAccuracyEvidence { /// The window implementation preserves exact query-time coverage and adds /// no error. Used for exact tumbling/sliding realizations. Exact, @@ -127,10 +127,10 @@ impl ExponentialHistogramQueryRange { } } -impl StreamingWindowAccuracyEvidence { +impl SummaryWindowAccuracyEvidence { pub(super) fn matches_assignments( &self, - assignments: &[StreamingWindowFrameworkAssignment], + assignments: &[SummaryWindowFrameworkAssignment], ) -> bool { let eh_summaries: Vec<_> = assignments .iter() diff --git a/docs/design_docs/architecture/physical-plan-integration.md b/docs/design_docs/architecture/physical-plan-integration.md index fee2cfd2..47fc3a1b 100644 --- a/docs/design_docs/architecture/physical-plan-integration.md +++ b/docs/design_docs/architecture/physical-plan-integration.md @@ -120,7 +120,7 @@ statistics contract, validation rule, and resource formula for an operation, a candidate containing it is unavailable. The streaming integration can consume a complete binding through -`StreamingNodeEvidence`. That binding is keyed to exact `SummaryNode` +`SummaryNodeEvidence`. That binding is keyed to exact `SummaryNode` identities and uses structured evidence for aggregate state, join, merge, subtract, delete, readout, and retained pre-ASAP work. It is a physical evidence boundary, not automatic physical lowering: a deployment must still @@ -468,3 +468,35 @@ candidate completeness; that evidence belongs to the semi-join's pruning step. The semantic `SummaryExpr` constructors still propose an initial execution layout. Uniform phase assignment applies to the exported post-ASAP DAG; it is not a claim that every deployment has implemented every placement. + + +## Summary cost evidence across data-arrival modes + +`SummaryMaintenanceCostModel` binds `SummaryNodeEvidence` and +`SummaryOperatorEvidence` independently of data-arrival mode. `ComparisonScope` +and the canonical `DataWorkload` determine arrival semantics; individual operator +resource records do not define another workload model. + +`SummaryMaintenanceInputs::from_workload` requires fresh snapshot cardinality. +For `AtRest`, it derives zero arrivals without requiring ingestion-rate evidence; +a fresh nonzero or invalid rate contradicts that declaration and is rejected. +For `ContinuouslyIngesting`, fresh, finite, nonnegative rate evidence remains +mandatory. Missing continuous rate evidence is never treated as zero. +Raw and summary evidence supplied directly by a provider obey the same arrival +invariant. Their source lineage, horizon, evaluation count, and snapshot dimensions +must still match. The existing lifecycle planner selects direct builds for a fixed +snapshot and charges bootstrap work, result evaluation, and retention; it charges +no arrival updates. This does not add computation-placement policy. + +`Mixed` and `Unknown` remain unsupported for analytical comparisons: the current +workload schema cannot identify separate backlog and arrival populations. The +adapter fails explicitly rather than guessing a split. The estimator version is +`summary-maintenance-resource-v2`; evidence type names drop the `Streaming` prefix +(`SummaryMaintenanceInputs`, `SummaryPhysicalInputEvidence`, `SummaryAggregateEvidence`, +`RetainedSubDagEvidence`, `RawInputEvidence`, and the summary window/alternative +types). Update source imports; no legacy-name aliases are provided. + +Regressions cover a fixed snapshot with no rate evidence, contradictory arrival +rates, scope mismatches, missing continuous-rate/cardinality evidence, and actual +lifecycle selection of a completely costed at-rest summary against its raw scan. +The existing continuous-ingestion and mixed-arrival rejection tests remain. diff --git a/docs/design_docs/architecture/planner-runtime-contract.md b/docs/design_docs/architecture/planner-runtime-contract.md index e567e124..2d071e65 100644 --- a/docs/design_docs/architecture/planner-runtime-contract.md +++ b/docs/design_docs/architecture/planner-runtime-contract.md @@ -106,7 +106,7 @@ such as cache behavior, serialization overhead, compression, spill I/O, or data-distribution-dependent sketch error. Provenance and version information must accompany those facts so stale observations fail closed. -`StreamingPhysicalPlanAlternative` is the current integration point for a +`SummaryPhysicalPlanAlternative` is the current integration point for a complete provider-enumerated implementation. Its identity is returned with the winning lifecycle combination. More structured planner-owned realization contracts can refine the candidate space without moving executor