From cb02df420414f8cb2f423ffa3b3204762e738d55 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Sun, 4 Oct 2026 22:48:08 +0000 Subject: [PATCH] chore(plan-selection): delete the legacy physical-plan cost model PhysicalPlanCostModel, PlannerPhysicalPlanProvider and PhysicalEvidenceSnapshot had no caller once dag_export was gone, and Stage 3 never used them. asap_types::cost keeps only CostUnit; the annotation and workload-summary types were dag_export's. cache_hit_ratios moves into analytical_cost's tests, its only caller. Co-Authored-By: Claude Opus 5.5 --- .../src/cost/analytical_cost.rs | 34 +- crates/plan-selection/src/cost/mod.rs | 3 +- .../src/cost/physical_plan_cost_model.rs | 1112 ----------------- crates/types/src/cost.rs | 679 +--------- crates/types/src/lib.rs | 2 +- 5 files changed, 10 insertions(+), 1820 deletions(-) delete mode 100644 crates/plan-selection/src/cost/physical_plan_cost_model.rs diff --git a/crates/plan-selection/src/cost/analytical_cost.rs b/crates/plan-selection/src/cost/analytical_cost.rs index 46c89f4eb..11138082d 100644 --- a/crates/plan-selection/src/cost/analytical_cost.rs +++ b/crates/plan-selection/src/cost/analytical_cost.rs @@ -30,8 +30,6 @@ struct ResolvedCacheProfile { integer_executions: Option, cpu_execution_factor: f64, scan_execution_factor: f64, - result_hit_ratio: f64, - buffer_hit_ratio: f64, buffer_miss_bytes: u64, buffer_working_set_bytes: u64, } @@ -51,8 +49,6 @@ fn resolve_cache_profile( integer_executions: Some(evaluation_count), cpu_execution_factor: evaluation_count as f64, scan_execution_factor: evaluation_count as f64, - result_hit_ratio: 0.0, - buffer_hit_ratio: 0.0, buffer_miss_bytes: 1, buffer_working_set_bytes: 1, }); @@ -114,8 +110,6 @@ fn resolve_cache_profile( let result_miss_ratio = invalidation + miss_fraction(evidence.result_cache)? * (1.0 - invalidation); let buffer_miss_ratio = miss_fraction(evidence.buffer_cache)?; - let result_hit_ratio = 1.0 - result_miss_ratio; - let buffer_hit_ratio = 1.0 - buffer_miss_ratio; let cpu_execution_factor = evidence.distinct_evaluations as f64 + evidence.repeated_identical_evaluations as f64 * result_miss_ratio; Ok(ResolvedCacheProfile { @@ -130,8 +124,6 @@ fn resolve_cache_profile( }, cpu_execution_factor, scan_execution_factor: cpu_execution_factor * buffer_miss_ratio, - result_hit_ratio, - buffer_hit_ratio, buffer_miss_bytes: evidence .buffer_cache .working_set_bytes @@ -140,16 +132,6 @@ fn resolve_cache_profile( }) } -/// Derive cache hit ratios from the shared assumptions for this workload. -pub fn cache_hit_ratios( - profile: &CacheProfile, - evaluation_count: u64, - data_arrival: DataArrival, -) -> Result<(f64, f64), AnalyticalCostError> { - let resolved = resolve_cache_profile(profile, evaluation_count, data_arrival)?; - Ok((resolved.result_hit_ratio, resolved.buffer_hit_ratio)) -} - /// Conversion from physical dimensions to one deployment-specific objective. /// Memory's coefficient means cost units per retained byte over this model's /// explicit comparison scope; it is not mixed with a rate implicitly. @@ -3686,10 +3668,6 @@ mod tests { fn accepts_shared(_: &asap_types::workload::resources::CacheProfile) {} accepts_shared(&legacy); assert_eq!(central, legacy); - assert_eq!( - cache_hit_ratios(¢ral, 6, DataArrival::AtRest).unwrap(), - (1.0, 0.5) - ); } #[test] @@ -3778,13 +3756,13 @@ mod tests { }; evidence.distinct_evaluations = 1; assert!(matches!( - cache_hit_ratios(&profile, 6, DataArrival::AtRest), + resolve_cache_profile(&profile, 6, DataArrival::AtRest), Err(AnalyticalCostError::InvalidCacheEvidence(_)) )); let profile = cache_profile(100, 100); assert!(matches!( - cache_hit_ratios(&profile, 6, DataArrival::ContinuouslyIngesting), + resolve_cache_profile(&profile, 6, DataArrival::ContinuouslyIngesting), Err(AnalyticalCostError::InvalidCacheEvidence(_)) )); } @@ -3837,14 +3815,14 @@ mod tests { }; inputs.distinct_evaluations = 0; inputs.repeated_identical_evaluations = 6; - assert!(cache_hit_ratios(&cache, 6, DataArrival::AtRest).is_err()); + assert!(resolve_cache_profile(&cache, 6, DataArrival::AtRest).is_err()); for ratio in [f64::NAN, f64::INFINITY, -1.0, 0.5] { let mut cache = cache_profile(100, 0); let CacheProfile::Evidence(inputs) = &mut cache else { unreachable!() }; inputs.result_invalidation_ratio = Some(ratio); - assert!(cache_hit_ratios(&cache, 6, DataArrival::AtRest).is_err()); + assert!(resolve_cache_profile(&cache, 6, DataArrival::AtRest).is_err()); } } @@ -3990,10 +3968,6 @@ mod tests { inputs.result_invalidation_ratio = Some(0.5); let estimate = estimate_physical_dag_with_cache(&nodes, "scan", &scope, &evidence, &profile).unwrap(); - assert_eq!( - cache_hit_ratios(&profile, 6, scope.data_arrival).unwrap(), - (0.5, 0.5) - ); assert_eq!(estimate.cpu_ops(), 400.0); assert_eq!(estimate.scan_bytes(), 2_000); } diff --git a/crates/plan-selection/src/cost/mod.rs b/crates/plan-selection/src/cost/mod.rs index f057ad296..6c1f427b2 100644 --- a/crates/plan-selection/src/cost/mod.rs +++ b/crates/plan-selection/src/cost/mod.rs @@ -2,7 +2,7 @@ //! deployment's cost-based selection plugs into, with the built-in //! [`DefaultCostModel`](cost_model::DefaultCostModel); recurring and one-shot //! cost rates ([`recurrence`]); analytical and evidence-based pricing -//! ([`analytical_cost`], [`physical_plan_cost_model`], [`empirical_cost`]); and +//! ([`analytical_cost`], [`empirical_cost`]); and //! the physical lowering and storage I/O profiles they price //! ([`query_physical_lowering`], [`storage_io`]). @@ -12,7 +12,6 @@ pub mod empirical_cost; pub mod empirical_resources; pub mod physical_handoff_cost; pub mod physical_operator_statistics; -pub mod physical_plan_cost_model; pub mod query_physical_lowering; pub mod recurrence; pub mod storage_io; diff --git a/crates/plan-selection/src/cost/physical_plan_cost_model.rs b/crates/plan-selection/src/cost/physical_plan_cost_model.rs deleted file mode 100644 index 5e1835f33..000000000 --- a/crates/plan-selection/src/cost/physical_plan_cost_model.rs +++ /dev/null @@ -1,1112 +0,0 @@ -//! Evidence-backed physical-plan costing at the planner selection handoff. - -use std::{cell::RefCell, rc::Rc}; - -use asap_types::ir::operator::AggIntent; -use asap_types::ir::schema::SketchAlgorithm; -use asap_types::ir::OperatorNode; -use asap_types::workload::resources::CacheProfile; - -use crate::cost::analytical_cost::{ - estimate_physical_dag_comparison, AnalyticalCostError, - EvidenceBackedPhysicalDAG as PhysicalDAG, PhysicalDAGComparisonEstimate, - PhysicalDAGEstimateRequest, PhysicalNodeEvidence, ResourceCalibration, -}; -use crate::cost::cost_model::{Cost, CostModel, DefaultCostModel}; -use crate::cost::physical_operator_statistics::ComparisonScope; -use crate::cost::query_physical_lowering::{ - lower_query_physical_dag, PhysicalNodeEvidenceProvider, PhysicalNodeRequest, -}; -use asap_logical_optimizer::pass1::replacement::{Replacement, ReplacementSubDAG, TargetSubDAG}; - -/// One immutable generation of deployment evidence for a planner target. -/// -/// The opaque version lets a provider bind every subsequent query and summary -/// lookup to the same catalog/runtime snapshot. A cost-model instance retains -/// this value for the target, so a caller that needs fresher evidence creates a -/// new model instead of mixing generations in one ranking decision. -#[derive(Debug, Clone, PartialEq)] -pub struct PhysicalEvidenceSnapshot { - pub version: String, - pub scope: ComparisonScope, - pub cache_profile: CacheProfile, - pub storage_io: Option, - pub handoffs: Option, -} - -/// Deployment evidence needed to price one planner alternative. -/// -/// The planner lowers raw queries and logical rewrites itself. A deployment -/// supplies the comparison scope and atomic evidence for each selected query -/// operator. Post-ASAP summary operators need a physical plan provider because their -/// implementation, placement, and retained-state layout are deployment -/// choices; that provider must return the complete summary DAG, including any -/// non-ASAP work kept inside it. -pub trait PlannerPhysicalPlanProvider { - /// Atomically captures the comparison scope and evidence generation. - fn capture_evidence_snapshot( - &self, - target: &TargetSubDAG<'_>, - ) -> Result; - - fn query_node_evidence( - &self, - snapshot: &PhysicalEvidenceSnapshot, - request: PhysicalNodeRequest<'_>, - ) -> Result; - - fn summary_physical_dag( - &self, - snapshot: &PhysicalEvidenceSnapshot, - summary: &Rc, - target: &TargetSubDAG<'_>, - ) -> Result; -} - -/// Dimensional comparison retained for explanations and verification. -#[derive(Debug, Clone, PartialEq)] -pub struct PhysicalPlanComparison { - pub resources: PhysicalDAGComparisonEstimate, - pub raw_cost: Cost, - pub candidate_cost: Cost, - pub storage_io: Option<( - crate::cost::storage_io::StorageEstimate, - crate::cost::storage_io::StorageEstimate, - )>, - pub handoffs: Option<( - crate::cost::physical_handoff_cost::PhysicalHandoffEstimate, - crate::cost::physical_handoff_cost::PhysicalHandoffEstimate, - )>, -} - -/// Planner cost model that admits only complete, cheaper physical plans. -/// -/// There is deliberately no structural or compact-formula fallback. Failure -/// to lower either alternative, missing statistics, an unknown physical -/// algorithm, or invalid source/horizon evidence makes the candidate -/// unavailable and leaves the target on its raw path. -pub struct PhysicalPlanCostModel<'a> { - provider: &'a dyn PlannerPhysicalPlanProvider, - calibration: ResourceCalibration, - target_evidence: RefCell>, -} - -struct CachedTargetEvidence { - root: Rc, - consumer_count: usize, - snapshot: PhysicalEvidenceSnapshot, - raw: PhysicalDAG, -} - -impl<'a> PhysicalPlanCostModel<'a> { - pub fn new( - provider: &'a dyn PlannerPhysicalPlanProvider, - calibration: ResourceCalibration, - ) -> Result { - // Exported scores retain base calibration provenance even when the - // objective is priced entirely by supplementary storage operations. - if calibration.version.trim().is_empty() { - return Err(AnalyticalCostError::MissingOrStale( - "resource_calibration.version", - )); - } - // Supplemental storage or handoff pricing may supply a zero-base objective. - match calibration.validate() { - Ok(()) | Err(AnalyticalCostError::ZeroCalibration) => {} - Err(error) => return Err(error), - } - Ok(Self { - provider, - calibration, - target_evidence: RefCell::new(Vec::new()), - }) - } - - fn target_evidence( - &self, - target: &TargetSubDAG<'_>, - ) -> Result<(PhysicalEvidenceSnapshot, PhysicalDAG), AnalyticalCostError> { - if let Some(cached) = self.target_evidence.borrow().iter().find(|cached| { - Rc::ptr_eq(&cached.root, target.root) && cached.consumer_count == target.consumer_count - }) { - return Ok((cached.snapshot.clone(), cached.raw.clone())); - } - - let snapshot = self.provider.capture_evidence_snapshot(target)?; - if snapshot.version.trim().is_empty() { - return Err(AnalyticalCostError::MissingOrStale( - "planner_evidence_snapshot.version", - )); - } - snapshot.scope.validate()?; - let evidence = QueryEvidence { - provider: self.provider, - snapshot: &snapshot, - }; - let raw = lower_query_physical_dag(target.root, &snapshot.scope, &evidence)?; - self.target_evidence - .borrow_mut() - .push(CachedTargetEvidence { - root: Rc::clone(target.root), - consumer_count: target.consumer_count, - snapshot: snapshot.clone(), - raw: raw.clone(), - }); - Ok((snapshot, raw)) - } - - pub fn estimate_candidate( - &self, - candidate: &ReplacementSubDAG, - target: &TargetSubDAG<'_>, - ) -> Result { - let (snapshot, raw) = self.target_evidence(target)?; - // Aggregate cache hit ratios do not identify which independently rounded - // extents issue requests. Do not mix cache-adjusted bytes/CPU with - // uncached request counts until cache-aware extent evidence is available. - if snapshot.storage_io.is_some() - && matches!(&snapshot.cache_profile, CacheProfile::Evidence(_)) - { - return Err(AnalyticalCostError::InvalidCacheEvidence( - "storage operation costs require an explicit no-cache profile; cache-aware extent evidence is unavailable", - )); - } - let scope = &snapshot.scope; - if snapshot.handoffs.is_some() - && matches!(snapshot.cache_profile, CacheProfile::Evidence(_)) - { - // Scope-based handoff multiplicity does not describe which actions - // cache hits skip; do not mix pre-cache byte work with discounted CPU. - return Err(AnalyticalCostError::MissingOrStale( - "cache-aware handoff execution evidence", - )); - } - let evidence = QueryEvidence { - provider: self.provider, - snapshot: &snapshot, - }; - let replacement = match &candidate.replacement { - Replacement::ExactComposition(_) => { - return Err(AnalyticalCostError::UnsupportedCandidate) - } - // A sub-DAG without summary state is the planner's own query - // lowering; anything with summary state is deployment-provided. - Replacement::SubDAG(sub_dag) if !sub_dag.contains_asap() => { - lower_query_physical_dag(sub_dag, scope, &evidence)? - } - Replacement::SubDAG(summary) => self - .provider - .summary_physical_dag(&snapshot, summary, target)?, - }; - let resources = estimate_physical_dag_comparison( - PhysicalDAGEstimateRequest { - nodes: &raw.nodes, - root: &raw.root, - scope, - statistics: &raw, - cache_profile: &snapshot.cache_profile, - }, - PhysicalDAGEstimateRequest { - nodes: &replacement.nodes, - root: &replacement.root, - scope, - statistics: &replacement, - cache_profile: &snapshot.cache_profile, - }, - )?; - let storage_io = snapshot - .storage_io - .as_ref() - .map(|profile| { - Ok(( - crate::cost::storage_io::estimate_storage_io( - &raw, - scope, - profile, - &snapshot.version, - )?, - crate::cost::storage_io::estimate_storage_io( - &replacement, - scope, - profile, - &snapshot.version, - )?, - )) - }) - .transpose()?; - let handoffs = snapshot - .handoffs - .as_ref() - .map(|profile| { - Ok(( - crate::cost::physical_handoff_cost::estimate_physical_handoffs( - &raw, - scope, - profile, - &snapshot.version, - )?, - crate::cost::physical_handoff_cost::estimate_physical_handoffs( - &replacement, - scope, - profile, - &snapshot.version, - )?, - )) - }) - .transpose()?; - let (mut raw_cost, mut candidate_cost) = match self.calibration.validate() { - Ok(()) => ( - Cost(resources.raw.calibrated_cost(&self.calibration)?), - Cost(resources.candidate.calibrated_cost(&self.calibration)?), - ), - Err(AnalyticalCostError::ZeroCalibration) => { - // Storage estimation above validates the coefficients and - // evidence. At least one priced dimension must remain. - let has_storage_objective = snapshot.storage_io.as_ref().is_some_and(|profile| { - let calibration = &profile.calibration; - [ - calibration.cost_per_disk_read, - calibration.cost_per_disk_write, - calibration.cost_per_object_get, - calibration.cost_per_object_put, - ] - .into_iter() - .any(|coefficient| coefficient > 0.0) - }); - // Boundary estimation above validates evidence and pricing. - // At least one dimension must have a positive coefficient. - let has_handoff_objective = snapshot.handoffs.as_ref().is_some_and(|profile| { - profile.calibration.cost_per_network_byte > 0.0 - || profile.calibration.cost_per_materialization_byte > 0.0 - }); - if !has_storage_objective && !has_handoff_objective { - return Err(AnalyticalCostError::ZeroCalibration); - } - (Cost(0.0), Cost(0.0)) - } - Err(error) => return Err(error), - }; - if let Some((raw, candidate)) = &storage_io { - raw_cost.0 += raw.cost; - candidate_cost.0 += candidate.cost; - } - if let Some((raw, candidate)) = &handoffs { - raw_cost.0 += raw.cost; - candidate_cost.0 += candidate.cost; - } - if !raw_cost.0.is_finite() || !candidate_cost.0.is_finite() { - return Err(AnalyticalCostError::Overflow); - } - Ok(PhysicalPlanComparison { - resources, - raw_cost, - candidate_cost, - storage_io, - handoffs, - }) - } -} - -struct QueryEvidence<'a> { - provider: &'a dyn PlannerPhysicalPlanProvider, - snapshot: &'a PhysicalEvidenceSnapshot, -} - -impl PhysicalNodeEvidenceProvider for QueryEvidence<'_> { - fn evidence( - &self, - request: PhysicalNodeRequest<'_>, - ) -> Result { - self.provider.query_node_evidence(self.snapshot, request) - } -} - -impl CostModel for PhysicalPlanCostModel<'_> { - fn candidate_cost_covers_complete_plan(&self) -> bool { - true - } - - fn candidate_cost( - &self, - candidate: &ReplacementSubDAG, - target: &TargetSubDAG<'_>, - ) -> Option { - self.estimate_candidate(candidate, target) - .ok() - .filter(|estimate| estimate.candidate_cost < estimate.raw_cost) - .map(|estimate| estimate.candidate_cost) - } - - fn rank_candidates( - &self, - intent: &AggIntent, - candidates: &[SketchAlgorithm], - ) -> Vec { - // Physical plans do not exist at algorithm enumeration time. Final - // ranking happens after binding, through `candidate_cost` above. - DefaultCostModel.rank_candidates(intent, candidates) - } - - fn estimate_cost(&self, candidate: &ReplacementSubDAG, target: &TargetSubDAG<'_>) -> f64 { - self.candidate_cost(candidate, target) - .map_or(f64::NAN, |cost| cost.0) - } -} - -#[cfg(test)] -mod tests { - use super::*; - use crate::candidate_selection::global_selection; - use std::cell::Cell; - use std::collections::HashMap; - - use asap_types::ir::operator::{Reduction, Source}; - use asap_types::ir::schema::{DataType, Field, Schema}; - use asap_types::ir::{NonASAPOp, OperatorNode}; - use asap_types::types::AccuracyTarget; - use asap_types::workload::{ - DataArrival, DurationMs, QueryRecurrence, QueryTimeScope, TimeSelection, TimestampMs, - }; - - use crate::cost::analytical_cost::{ExecutionMultiplicity, PhysicalDAGNode, PhysicalOperator}; - use crate::cost::physical_operator_statistics::{ - EdgeStatistics, OperatorStatistics, ScanSelection, UnaryEdgeStatistics, - }; - use asap_logical_optimizer::pass1::replacement::ReplacementStrategy; - - fn edge(rows: u64, bytes: u64) -> EdgeStatistics { - EdgeStatistics { rows, bytes } - } - - fn unary_edges(input: EdgeStatistics, output: EdgeStatistics) -> UnaryEdgeStatistics { - UnaryEdgeStatistics { - input, - output, - promql: None, - } - } - - fn scan_statistics(source_read_bytes: u64, edge: EdgeStatistics) -> OperatorStatistics { - OperatorStatistics::Scan { - edges: unary_edges(edge, edge), - source_read_bytes, - } - } - - fn aggregate_statistics(input: EdgeStatistics, output: EdgeStatistics) -> OperatorStatistics { - OperatorStatistics::HashAggregate { - edges: unary_edges(input, output), - group_count: 1, - key_bytes: 0, - accumulator_bytes_per_group: 8, - } - } - - fn pass_through_statistics(edge: EdgeStatistics) -> OperatorStatistics { - OperatorStatistics::PassThrough { - edges: unary_edges(edge, edge), - } - } - - fn query() -> Rc { - let scan = OperatorNode::new_shared(asap_types::ir::Operator::NonASAP(NonASAPOp::Scan { - source: Source::Table { - table_ref: "events".into(), - }, - predicates: vec![], - schema: Schema::new(vec![Field::plain("value", DataType::Float64, false)]), - })) - .unwrap(); - OperatorNode::new_shared(asap_types::ir::Operator::NonASAP(NonASAPOp::Aggregate { - reduction: Reduction::by(vec![]), - measures: vec![AggIntent::Count { - accuracy: AccuracyTarget::Epsilon(0.01), - }], - output_names: vec![], - filters: vec![], - having: None, - child: scan, - })) - .unwrap() - } - - fn scope() -> ComparisonScope { - ComparisonScope { - data_arrival: DataArrival::AtRest, - planning_time: TimestampMs(1_000), - horizon: DurationMs(10_000), - recurrence: QueryRecurrence::OneTime { - invocations: 10, - execute_at: None, - }, - time_selection: TimeSelection { - scope: QueryTimeScope::Longitudinal, - lookback: Some(DurationMs(10_000)), - as_of: Some(TimestampMs(1_000)), - }, - sources: vec![ScanSelection { - source: Source::Table { - table_ref: "events".into(), - }, - source_snapshot_id: "snapshot-1".into(), - predicates: vec![], - info_matchers: vec![], - }], - } - } - - struct TestProvider { - storage_io: Option, - summary_available: bool, - candidate_scan_bytes: u64, - handoffs: Option, - snapshot_calls: Cell, - raw_evidence_calls: Cell, - } - - impl TestProvider { - fn new(summary_available: bool, candidate_scan_bytes: u64) -> Self { - Self { - storage_io: None, - summary_available, - candidate_scan_bytes, - handoffs: None, - snapshot_calls: Cell::new(0), - raw_evidence_calls: Cell::new(0), - } - } - - fn summary_dag(&self, scope: &ComparisonScope) -> PhysicalDAG { - let scan_statistics = scan_statistics(self.candidate_scan_bytes, edge(100, 800)); - let aggregate_statistics = aggregate_statistics(edge(100, 800), edge(1, 8)); - let read_statistics = pass_through_statistics(edge(1, 8)); - let evidence = HashMap::from([ - ( - "candidate-scan".into(), - PhysicalNodeEvidence { - physical_id: "candidate-scan".into(), - statistics: scan_statistics, - output_buffer_bytes: 8, - }, - ), - ( - "candidate-state".into(), - PhysicalNodeEvidence { - physical_id: "candidate-state".into(), - statistics: aggregate_statistics, - output_buffer_bytes: 8, - }, - ), - ( - "candidate-read".into(), - PhysicalNodeEvidence { - physical_id: "candidate-read".into(), - statistics: read_statistics, - output_buffer_bytes: 8, - }, - ), - ]); - PhysicalDAG { - nodes: vec![ - PhysicalDAGNode { - id: "candidate-scan".into(), - operator: PhysicalOperator::Scan, - children: vec![], - scan_selection: Some(scope.sources[0].clone()), - output_buffer_bytes: 8, - retained_bytes: 0, - execution: ExecutionMultiplicity::Once, - }, - PhysicalDAGNode { - id: "candidate-state".into(), - operator: PhysicalOperator::HashAggregate { - grouping_key_count: 0, - accumulator_count: 1, - }, - children: vec!["candidate-scan".into()], - scan_selection: None, - output_buffer_bytes: 8, - retained_bytes: 8, - execution: ExecutionMultiplicity::Once, - }, - PhysicalDAGNode { - id: "candidate-read".into(), - operator: PhysicalOperator::PassThrough, - children: vec!["candidate-state".into()], - scan_selection: None, - output_buffer_bytes: 8, - retained_bytes: 0, - execution: ExecutionMultiplicity::PerEvaluation, - }, - ], - root: "candidate-read".into(), - evidence, - } - } - } - - impl PlannerPhysicalPlanProvider for TestProvider { - fn capture_evidence_snapshot( - &self, - _target: &TargetSubDAG<'_>, - ) -> Result { - self.snapshot_calls.set(self.snapshot_calls.get() + 1); - Ok(PhysicalEvidenceSnapshot { - version: "test-snapshot-1".into(), - scope: scope(), - cache_profile: CacheProfile::no_cache(), - storage_io: self.storage_io.clone(), - handoffs: self.handoffs.clone(), - }) - } - - fn query_node_evidence( - &self, - snapshot: &PhysicalEvidenceSnapshot, - request: PhysicalNodeRequest<'_>, - ) -> Result { - assert_eq!(snapshot.version, "test-snapshot-1"); - self.raw_evidence_calls - .set(self.raw_evidence_calls.get() + 1); - let (physical_id, statistics) = match request.operator { - PhysicalOperator::Scan => ("raw-scan", scan_statistics(800, edge(100, 800))), - PhysicalOperator::HashAggregate { - grouping_key_count: 0, - accumulator_count: 1, - } => ( - "raw-aggregate", - aggregate_statistics(edge(100, 800), edge(1, 8)), - ), - _ => return Err(AnalyticalCostError::UnsupportedQueryOperator), - }; - Ok(PhysicalNodeEvidence { - physical_id: physical_id.into(), - output_buffer_bytes: 8, - statistics, - }) - } - - fn summary_physical_dag( - &self, - snapshot: &PhysicalEvidenceSnapshot, - _summary: &Rc, - _target: &TargetSubDAG<'_>, - ) -> Result { - assert_eq!(snapshot.version, "test-snapshot-1"); - self.summary_available - .then(|| self.summary_dag(&snapshot.scope)) - .ok_or(AnalyticalCostError::MissingOrStale("summary_physical_plan")) - } - } - - fn calibration() -> ResourceCalibration { - ResourceCalibration { - cost_per_cpu_op: 1.0, - cost_per_scan_byte: 1.0, - cost_per_retained_byte: 1.0, - version: "test-v1".into(), - } - } - - /// Both base-priced and storage-only objectives need identifiable calibration. - #[test] - fn blank_base_calibration_version_is_rejected_before_snapshot_lookup() { - let provider = TestProvider::new(true, 800); - for version in ["", " \t\n"] { - for zero_base in [false, true] { - let mut base = calibration(); - base.version = version.into(); - if zero_base { - base.cost_per_cpu_op = 0.0; - base.cost_per_scan_byte = 0.0; - base.cost_per_retained_byte = 0.0; - } - assert!(matches!( - PhysicalPlanCostModel::new(&provider, base), - Err(AnalyticalCostError::MissingOrStale( - "resource_calibration.version" - )) - )); - } - } - assert_eq!(provider.snapshot_calls.get(), 0); - } - - // A request-only objective ranks complete evidence and rejects absent, - // zero, or invalid supplemental calibration. - #[test] - fn storage_only_objective_requires_positive_valid_storage_calibration() { - use crate::cost::storage_io::*; - let root = query(); - let target = TargetSubDAG::new(&root); - let mut provider = TestProvider::new(true, 800); - let snapshot = provider.capture_evidence_snapshot(&target).unwrap(); - let raw = lower_query_physical_dag( - &root, - &snapshot.scope, - &QueryEvidence { - provider: &provider, - snapshot: &snapshot, - }, - ) - .unwrap(); - let mut profile = StorageIoProfile { - evidence_version: snapshot.version, - observed_at_ms: 900, - valid_until_ms: 2000, - calibration: StorageCalibration { - version: "requests-v1".into(), - cost_per_disk_read: 0.0, - cost_per_disk_write: 0.0, - cost_per_object_get: 1.0, - cost_per_object_put: 0.0, - }, - nodes: HashMap::new(), - }; - for dag in [raw, provider.summary_dag(&snapshot.scope)] { - for node in dag.nodes { - let statistics = dag.evidence[&node.id].statistics.clone(); - let accesses = match statistics { - OperatorStatistics::Scan { - source_read_bytes, .. - } => vec![StorageAccess { - operation: StorageOperation::ObjectGet, - extent_bytes: vec![source_read_bytes], - bytes_per_request: 400, - }], - _ => vec![], - }; - profile.nodes.insert( - node.id.clone(), - StorageNodeEvidence { - node, - statistics, - accesses, - }, - ); - } - } - let base = ResourceCalibration { - cost_per_cpu_op: 0.0, - cost_per_scan_byte: 0.0, - cost_per_retained_byte: 0.0, - version: "unused-base-v1".into(), - }; - let candidates = asap_logical_optimizer::pass1::replacement::ASAPStrategies::default() - .replacements(&target); - provider.storage_io = Some(profile.clone()); - let model = PhysicalPlanCostModel::new(&provider, base.clone()).unwrap(); - let estimate = model.estimate_candidate(&candidates[0], &target).unwrap(); - assert_eq!(estimate.raw_cost.0, 20.0); - assert_eq!(estimate.candidate_cost.0, 2.0); - assert_eq!( - model.candidate_cost(&candidates[0], &target), - Some(Cost(2.0)) - ); - drop(model); - for coefficient in [0.0, -1.0, f64::NAN] { - profile.calibration.cost_per_object_get = coefficient; - provider.storage_io = Some(profile.clone()); - let model = PhysicalPlanCostModel::new(&provider, base.clone()).unwrap(); - assert!(model.estimate_candidate(&candidates[0], &target).is_err()); - } - provider.storage_io = None; - let model = PhysicalPlanCostModel::new(&provider, base).unwrap(); - assert_eq!( - model.estimate_candidate(&candidates[0], &target), - Err(AnalyticalCostError::ZeroCalibration) - ); - } - - /// Shared count defaults must not turn an omitted profile into measured zero I/O. - #[test] - fn missing_storage_profile_remains_unestimated() { - let root = query(); - let target = TargetSubDAG::new(&root); - let candidates = asap_logical_optimizer::pass1::replacement::ASAPStrategies::default() - .replacements(&target); - let provider = TestProvider::new(true, 800); - let model = PhysicalPlanCostModel::new(&provider, calibration()).unwrap(); - let estimate = model.estimate_candidate(&candidates[0], &target).unwrap(); - assert!(estimate.storage_io.is_none()); - } - - // Combined objectives must identify the base calibration even when its - // coefficients are zero and handoffs supply the entire objective. - #[test] - fn blank_base_calibration_version_is_rejected() { - let provider = TestProvider::new(true, 800); - for version in ["", " \t\n"] { - for handoff_only in [false, true] { - let mut calibration = calibration(); - calibration.version = version.into(); - if handoff_only { - calibration.cost_per_cpu_op = 0.0; - calibration.cost_per_scan_byte = 0.0; - calibration.cost_per_retained_byte = 0.0; - } - assert!( - PhysicalPlanCostModel::new(&provider, calibration).is_err(), - "blank base provenance accepted (handoff_only={handoff_only})" - ); - } - } - } - - // Explicit byte pricing can rank complete plans without pricing CPU or scans. - #[test] - fn handoff_only_objective_ranks_complete_plans() { - use crate::cost::physical_handoff_cost::*; - let root = query(); - let target = TargetSubDAG::new(&root); - let mut provider = TestProvider::new(true, 800); - let raw = PhysicalPlanCostModel::new(&provider, calibration()) - .unwrap() - .target_evidence(&target) - .unwrap() - .1; - let candidate = provider.summary_dag(&scope()); - provider.handoffs = Some(PhysicalHandoffProfile { - evidence_version: "test-snapshot-1".into(), - observed_at_ms: 900, - valid_until_ms: 2000, - calibration: PhysicalHandoffCalibration { - version: "network-only-v1".into(), - cost_per_network_byte: 1.0, - cost_per_materialization_byte: 0.0, - }, - plans: [&raw, &candidate] - .into_iter() - .map(|dag| PhysicalHandoffPlanEvidence { - root: dag.root.clone(), - nodes: dag - .nodes - .iter() - .map(|node| { - let statistics = dag.evidence[&node.id].statistics.clone(); - let bytes = statistics.output().bytes; - let handoffs = if matches!(node.operator, PhysicalOperator::Scan) { - vec![PhysicalHandoff { - id: "scan-transfer".into(), - consumer: None, - kind: PhysicalHandoffKind::Network { - source_location: "edge".into(), - destination_location: "backend".into(), - }, - logical_bytes: bytes, - encoded_bytes: bytes, - copies: 1, - }] - } else { - vec![] - }; - ( - node.id.clone(), - PhysicalHandoffNodeEvidence { - node: node.clone(), - statistics, - handoffs, - }, - ) - }) - .collect(), - }) - .collect(), - }); - let zero_base = ResourceCalibration { - cost_per_cpu_op: 0.0, - cost_per_scan_byte: 0.0, - cost_per_retained_byte: 0.0, - version: "handoff-only-v1".into(), - }; - let space = asap_logical_optimizer::pass1::replacement::search_workload_with( - vec![("q", Rc::clone(&root))], - &asap_logical_optimizer::pass1::replacement::default_strategies(), - ); - let model = PhysicalPlanCostModel::new(&provider, zero_base.clone()).unwrap(); - let selected = global_selection(&space, &model); - assert!(selected - .for_target(&space.roots[0].1) - .unwrap() - .chosen - .is_some()); - drop(model); - for coefficient in [0.0, -1.0, f64::NAN] { - provider - .handoffs - .as_mut() - .unwrap() - .calibration - .cost_per_network_byte = coefficient; - let model = PhysicalPlanCostModel::new(&provider, zero_base.clone()).unwrap(); - assert!(global_selection(&space, &model) - .for_target(&space.roots[0].1) - .unwrap() - .chosen - .is_none()); - } - provider.handoffs = None; - let model = PhysicalPlanCostModel::new(&provider, zero_base).unwrap(); - assert!(global_selection(&space, &model) - .for_target(&space.roots[0].1) - .unwrap() - .chosen - .is_none()); - } - - #[test] - fn global_selection_uses_complete_physical_comparison() { - let root = query(); - let space = asap_logical_optimizer::pass1::replacement::search_workload_with( - vec![("q", Rc::clone(&root))], - &asap_logical_optimizer::pass1::replacement::default_strategies(), - ); - let planned_root = Rc::clone(&space.roots[0].1); - let provider = TestProvider::new(true, 800); - let model = PhysicalPlanCostModel::new(&provider, calibration()).unwrap(); - - let selected = global_selection(&space, &model); - assert!( - selected.for_target(&planned_root).unwrap().chosen.is_some(), - "a fully bound build-once summary cheaper than ten raw scans must be selected" - ); - } - - // Ranking must retain the uncovered byte of a nearly resident buffer cache. - #[test] - fn tiny_buffer_misses_still_affect_global_selection() { - use crate::cost::analytical_cost::{CacheCapacityEvidence, CacheEvidence}; - const WORKING_SET: u64 = 1_u64 << 63; - struct AlmostResident(TestProvider); - impl PlannerPhysicalPlanProvider for AlmostResident { - fn capture_evidence_snapshot( - &self, - target: &TargetSubDAG<'_>, - ) -> Result { - let mut snapshot = self.0.capture_evidence_snapshot(target)?; - snapshot.cache_profile = CacheProfile::Evidence(CacheEvidence { - version: "almost-resident-v1".into(), - distinct_evaluations: 10, - repeated_identical_evaluations: 0, - result_invalidation_ratio: None, - result_cache: CacheCapacityEvidence { - working_set_bytes: 1, - capacity_bytes: 0, - }, - buffer_cache: CacheCapacityEvidence { - working_set_bytes: WORKING_SET, - capacity_bytes: WORKING_SET - 1, - }, - }); - Ok(snapshot) - } - fn query_node_evidence( - &self, - snapshot: &PhysicalEvidenceSnapshot, - request: PhysicalNodeRequest<'_>, - ) -> Result { - let mut evidence = self.0.query_node_evidence(snapshot, request)?; - if let OperatorStatistics::Scan { - source_read_bytes, .. - } = &mut evidence.statistics - { - *source_read_bytes = WORKING_SET; - } - Ok(evidence) - } - fn summary_physical_dag( - &self, - snapshot: &PhysicalEvidenceSnapshot, - summary: &Rc, - target: &TargetSubDAG<'_>, - ) -> Result { - self.0.summary_physical_dag(snapshot, summary, target) - } - } - let space = asap_logical_optimizer::pass1::replacement::search_workload_with( - vec![("q", query())], - &asap_logical_optimizer::pass1::replacement::default_strategies(), - ); - let provider = AlmostResident(TestProvider::new(true, WORKING_SET)); - let model = PhysicalPlanCostModel::new( - &provider, - ResourceCalibration { - cost_per_cpu_op: 0.0, - cost_per_scan_byte: 1.0, - cost_per_retained_byte: 0.0, - version: "disk-only-v1".into(), - }, - ) - .unwrap(); - let selected = global_selection(&space, &model); - assert!( - selected - .for_target(&space.roots[0].1) - .unwrap() - .chosen - .is_some(), - "one remaining build read must rank ahead of ten remaining raw reads" - ); - } - - // Versioned physical ranking cannot accept an anonymous calibration generation. - #[test] - fn blank_calibration_version_is_rejected_before_ranking() { - let provider = TestProvider::new(true, 800); - for version in ["", " \t\n"] { - let mut coefficients = calibration(); - coefficients.version = version.into(); - assert!(matches!( - PhysicalPlanCostModel::new(&provider, coefficients), - Err(AnalyticalCostError::MissingOrStale( - "resource_calibration.version" - )) - )); - } - } - - #[test] - fn missing_summary_evidence_keeps_the_raw_target() { - let root = query(); - let space = asap_logical_optimizer::pass1::replacement::search_workload_with( - vec![("q", Rc::clone(&root))], - &asap_logical_optimizer::pass1::replacement::default_strategies(), - ); - let planned_root = Rc::clone(&space.roots[0].1); - let provider = TestProvider::new(false, 800); - let model = PhysicalPlanCostModel::new(&provider, calibration()).unwrap(); - - let selected = global_selection(&space, &model); - assert!( - selected.for_target(&planned_root).unwrap().chosen.is_none(), - "missing physical summary evidence must not fall back to a structural estimate" - ); - } - - #[test] - fn candidate_with_incomplete_scope_is_unavailable() { - struct WrongScope(TestProvider); - impl PlannerPhysicalPlanProvider for WrongScope { - fn capture_evidence_snapshot( - &self, - target: &TargetSubDAG<'_>, - ) -> Result { - self.0.capture_evidence_snapshot(target) - } - - fn query_node_evidence( - &self, - snapshot: &PhysicalEvidenceSnapshot, - request: PhysicalNodeRequest<'_>, - ) -> Result { - self.0.query_node_evidence(snapshot, request) - } - - fn summary_physical_dag( - &self, - snapshot: &PhysicalEvidenceSnapshot, - summary: &Rc, - target: &TargetSubDAG<'_>, - ) -> Result { - let mut dag = self.0.summary_physical_dag(snapshot, summary, target)?; - dag.nodes[0] - .scan_selection - .as_mut() - .unwrap() - .source_snapshot_id = "other".into(); - Ok(dag) - } - } - - let root = query(); - let candidates = asap_logical_optimizer::pass1::replacement::ASAPStrategies::default() - .replacements(&TargetSubDAG::new(&root)); - let provider = WrongScope(TestProvider::new(true, 800)); - let model = PhysicalPlanCostModel::new(&provider, calibration()).unwrap(); - assert_eq!( - model.candidate_cost(&candidates[0], &TargetSubDAG::new(&root)), - None - ); - } - - #[test] - fn blank_snapshot_version_is_unavailable_before_evidence_lookup() { - struct BlankVersionProvider; - impl PlannerPhysicalPlanProvider for BlankVersionProvider { - fn capture_evidence_snapshot( - &self, - _target: &TargetSubDAG<'_>, - ) -> Result { - Ok(PhysicalEvidenceSnapshot { - version: " \t".into(), - scope: scope(), - cache_profile: CacheProfile::no_cache(), - storage_io: None, - handoffs: None, - }) - } - - fn query_node_evidence( - &self, - _snapshot: &PhysicalEvidenceSnapshot, - _request: PhysicalNodeRequest<'_>, - ) -> Result { - panic!("blank snapshot versions must fail before evidence lookup") - } - - fn summary_physical_dag( - &self, - _snapshot: &PhysicalEvidenceSnapshot, - _summary: &Rc, - _target: &TargetSubDAG<'_>, - ) -> Result { - panic!("blank snapshot versions must fail before summary binding") - } - } - - let root = query(); - let candidates = asap_logical_optimizer::pass1::replacement::ASAPStrategies::default() - .replacements(&TargetSubDAG::new(&root)); - let model = PhysicalPlanCostModel::new(&BlankVersionProvider, calibration()).unwrap(); - assert_eq!( - model.candidate_cost(&candidates[0], &TargetSubDAG::new(&root)), - None - ); - } - - #[test] - fn complete_candidate_that_costs_more_than_raw_is_not_selected() { - let root = query(); - let space = asap_logical_optimizer::pass1::replacement::search_workload_with( - vec![("q", Rc::clone(&root))], - &asap_logical_optimizer::pass1::replacement::default_strategies(), - ); - let planned_root = Rc::clone(&space.roots[0].1); - let provider = TestProvider::new(true, 100_000); - let model = PhysicalPlanCostModel::new(&provider, calibration()).unwrap(); - - let selected = global_selection(&space, &model); - assert!(selected.for_target(&planned_root).unwrap().chosen.is_none()); - } - - #[test] - fn sibling_candidates_share_one_scope_and_raw_baseline() { - let root = query(); - let candidates = asap_logical_optimizer::pass1::replacement::ASAPStrategies::default() - .replacements(&TargetSubDAG::new(&root)); - assert!(candidates.len() >= 2); - let provider = TestProvider::new(true, 800); - let model = PhysicalPlanCostModel::new(&provider, calibration()).unwrap(); - let target = TargetSubDAG::new(&root); - - model.estimate_candidate(&candidates[0], &target).unwrap(); - model.estimate_candidate(&candidates[1], &target).unwrap(); - - assert_eq!(provider.snapshot_calls.get(), 1); - assert_eq!( - provider.raw_evidence_calls.get(), - 2, - "the scan and aggregate evidence for the raw baseline are captured once" - ); - } -} diff --git a/crates/types/src/cost.rs b/crates/types/src/cost.rs index 725807b01..c7158329f 100644 --- a/crates/types/src/cost.rs +++ b/crates/types/src/cost.rs @@ -1,19 +1,9 @@ -//! Unit- and provenance-aware cost annotations for exported DAGs. -//! -//! This crate defines the interchange schema; it does not estimate costs. -//! Producers may report a horizon total, a recurring rate, or no value. -//! Values with different [`CostUnit`]s are -//! never combined. A missing estimate is represented by -//! [`CostAnnotation::unavailable`], never by zero. -//! -//! A baseline comparison uses `delta = baseline_value - selected_value` and, -//! when `baseline_value > 0`, `benefit_ratio = delta / baseline_value`. -//! Model inputs and a model version (or a benchmark identifier for measured -//! values) keep every emitted number auditable. +//! The unit a cost is expressed in. Values with different [`CostUnit`]s are +//! never combined. use serde::{Deserialize, Serialize}; -/// The unit one [`CostAnnotation::value`] (and its `delta`) is expressed in. +/// The unit a cost value is expressed in. /// Producers and consumers must not compare or aggregate different units. #[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)] pub enum CostUnit { @@ -23,7 +13,7 @@ pub enum CostUnit { /// A finite-run or one-shot total at some horizon `H` — `total_cost(H)` /// in the issue's formulas, or a standalone one-shot addend. Never /// aggregated with [`CostUnitsPerSecond`](CostUnit::CostUnitsPerSecond) - /// except through [`total_cost`], which keeps the two terms explicit + /// except through `total_cost`, which keeps the two terms explicit /// rather than silently adding a rate to a total. CostUnits, } @@ -46,664 +36,3 @@ impl std::fmt::Display for CostUnit { } } } - -/// Provenance of one [`CostAnnotation`]'s value — the issue's three states. -#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] -pub enum CostSource { - /// Generated by a named/versioned cost model - /// ([`CostAnnotation::model_version`]). - Modeled, - /// Loaded from a reproducible benchmark artifact - /// ([`CostAnnotation::benchmark_id`]) — issue #288. - Measured, - /// No value. `value` is always `None` for this source — never encode an - /// unknown value as `0` or another synthetic number. - Unavailable, -} - -/// What a [`CostAnnotation`]'s `baseline`/`delta`/`benefit_ratio` are -/// measured against. The issue names two examples explicitly; both are -/// first-class variants here rather than opaque strings so a renderer can -/// display them without guessing. [`Named`](BaselineRef::Named) covers any -/// other explicitly-chosen baseline a future caller introduces. -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -#[serde(tag = "kind", content = "detail")] -pub enum BaselineRef { - /// Recomputing the pre-ASAP target independently at every consumer site - /// — "do nothing" (never apply ASAP-aware replacement at all). - PreAsapRecomputation, - /// The best-ranked *non-selected* legal candidate for the same target - /// (`rank` into that target's own `candidate_selection::cost_sorted` ordering, - /// `0` = best; a baseline referencing this variant is always `rank >= - /// 1`, since `rank 0` is what got selected). - HighestRankedNonSelectedCandidate { rank: usize }, - /// Any other explicitly-named baseline. - Named(String), -} - -/// One raw input that fed a [`CostAnnotation`]'s `value` — surfaced so the -/// viewer's sidebar can show *why* a modeled number is what it is, not just -/// the number itself. -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -pub struct CostInput { - pub name: String, - pub value: f64, - #[serde(skip_serializing_if = "Option::is_none")] - pub unit: Option, -} - -impl CostInput { - pub fn new(name: impl Into, value: f64) -> Self { - CostInput { - name: name.into(), - value, - unit: None, - } - } -} - -/// A structured, optional cost/benefit annotation — issue #286's schema. -/// Every numeric value carries units and provenance; a missing value is -/// [`CostSource::Unavailable`] with `value: None`, never a fabricated -/// number. See the module doc for the exact formulas this backs. -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -pub struct CostAnnotation { - pub value: Option, - pub unit: CostUnit, - pub source: CostSource, - #[serde(skip_serializing_if = "Option::is_none")] - pub baseline: Option, - /// `baseline_value - value` under `baseline`, when both are known — the - /// corresponding benefit in this annotation's declared unit. - #[serde(skip_serializing_if = "Option::is_none")] - pub delta: Option, - /// `delta / baseline_value`, when `baseline_value > 0` — the issue's - /// `estimated_benefit_ratio`. Not part of the issue's own sketch schema - /// verbatim (the issue only asks that ratio be *derivable*), but kept - /// as its own explicit field rather than pushed onto the caller to - /// recompute from `delta` plus a `baseline` it would otherwise have no - /// value for. - #[serde(skip_serializing_if = "Option::is_none")] - pub benefit_ratio: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub model_version: Option, - /// Immutable catalog/runtime evidence generation used by the model. - /// Kept separate from `model_version`: changing evidence must remain - /// visible even when the analytical formulas are unchanged. - #[serde(skip_serializing_if = "Option::is_none")] - pub evidence_version: Option, - /// Named, versioned cache assumption used by the physical cost model. - #[serde(skip_serializing_if = "Option::is_none")] - pub cache_profile: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub benchmark_id: Option, - #[serde(default, skip_serializing_if = "Vec::is_empty")] - pub inputs: Vec, -} - -impl CostAnnotation { - /// No value at all — `value: None`, `source: Unavailable`. `unit` is - /// still required: it states what unit a value *would* have been in, - /// had one been available, so a renderer can group "not estimated" - /// alongside the kind of number it's missing. - pub fn unavailable(unit: CostUnit) -> Self { - CostAnnotation { - value: None, - unit, - source: CostSource::Unavailable, - baseline: None, - delta: None, - benefit_ratio: None, - model_version: None, - evidence_version: None, - cache_profile: None, - benchmark_id: None, - inputs: Vec::new(), - } - } - - /// A modeled value with no baseline comparison attached yet — chain - /// [`with_baseline`](Self::with_baseline) to add one. Non-finite values, - /// blank model versions, or non-finite inputs produce an unavailable - /// annotation rather than poisoned JSON. - pub fn modeled( - value: f64, - unit: CostUnit, - model_version: impl Into, - inputs: Vec, - ) -> Self { - let model_version = model_version.into(); - if !value.is_finite() - || model_version.trim().is_empty() - || inputs.iter().any(|input| !input.value.is_finite()) - { - return Self::unavailable(unit); - } - CostAnnotation { - value: Some(value), - unit, - source: CostSource::Modeled, - baseline: None, - delta: None, - benefit_ratio: None, - model_version: Some(model_version), - evidence_version: None, - cache_profile: None, - benchmark_id: None, - inputs, - } - } - - /// Attach `baseline` plus its own value, computing `delta` and - /// [`benefit_ratio`](Self::benefit_ratio) per the issue's formulas. - /// `self.value` must already be `Some` (an - /// [`unavailable`](Self::unavailable) annotation has nothing to - /// subtract a baseline from and is returned unchanged). - pub fn with_baseline(mut self, baseline: BaselineRef, baseline_value: f64) -> Self { - let Some(value) = self.value else { - return self; - }; - if !value.is_finite() || !baseline_value.is_finite() || baseline_value < 0.0 { - return Self::unavailable(self.unit); - } - let delta = baseline_value - value; - self.delta = Some(delta); - self.benefit_ratio = benefit_ratio(baseline_value, delta); - self.baseline = Some(baseline); - self - } - - /// Bind a modeled annotation to the immutable evidence generation used - /// to produce it. Blank identities make the annotation unavailable. - pub fn with_evidence_version(mut self, evidence_version: impl Into) -> Self { - if self.source != CostSource::Modeled || self.value.is_none() { - return self; - } - let evidence_version = evidence_version.into(); - if evidence_version.trim().is_empty() { - return Self::unavailable(self.unit); - } - self.evidence_version = Some(evidence_version); - self - } - - pub fn with_cache_profile(mut self, cache_profile: impl Into) -> Self { - if self.source != CostSource::Modeled || self.value.is_none() { - return self; - } - let cache_profile = cache_profile.into(); - if cache_profile.trim().is_empty() { - return Self::unavailable(self.unit); - } - self.cache_profile = Some(cache_profile); - self - } -} - -/// `delta / baseline_value`, or `None` when `baseline_value <= 0` — the -/// issue's explicit "ratio unavailable" guard (a zero or negative baseline -/// makes a ratio meaningless, not just numerically awkward). -pub fn benefit_ratio(baseline_value: f64, delta: f64) -> Option { - if baseline_value.is_finite() && baseline_value > 0.0 && delta.is_finite() { - Some(delta / baseline_value) - } else { - None - } -} - -/// `total_cost(H) = recurring_cost_rate * H + one_shot_cost` — finite-run or -/// one-shot totals require an explicit horizon. Returns `None` (never a -/// fabricated or poisoned total) when both terms are `None`, when `horizon` -/// isn't a finite, non-negative number for a `Some(rate)`, or when either -/// `recurring_cost_rate` or `one_shot_cost` is itself non-finite (`NaN` or -/// infinite) — a non-finite input must never silently produce `Some(NaN)`. -/// Rate and one-shot costs are passed as separate arguments specifically so -/// they can never be silently added by a caller before this function ever -/// sees them. -pub fn total_cost( - recurring_cost_rate: Option, - horizon: f64, - one_shot_cost: Option, -) -> Option { - // A non-finite input anywhere here (NaN/±inf — e.g. from a future #287 - // caller upstream) must never quietly poison the total into `Some(NaN)`: - // that would violate this module's own "never fabricate, never a - // poisoned total" rule as much as inventing a number from nothing would. - if let Some(rate) = recurring_cost_rate { - if !rate.is_finite() { - return None; - } - } - if let Some(one_shot) = one_shot_cost { - if !one_shot.is_finite() { - return None; - } - } - let recurring = match recurring_cost_rate { - Some(rate) => { - if !horizon.is_finite() || horizon < 0.0 { - return None; - } - rate * horizon - } - None => 0.0, - }; - match (recurring_cost_rate, one_shot_cost) { - (None, None) => None, - _ => Some(recurring + one_shot_cost.unwrap_or(0.0)), - } -} - -/// Two [`CostAnnotation`]s were summed by [`sum_workload_costs`] despite -/// disagreeing on [`CostUnit`] — unit-incompatible aggregation is rejected -/// rather than silently mixed. -#[derive(Debug, Clone, Copy, PartialEq, Eq)] -pub struct UnitMismatch { - pub first: CostUnit, - pub second: CostUnit, -} - -impl std::fmt::Display for UnitMismatch { - fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - write!( - f, - "cannot aggregate cost annotations with mismatched units: {} vs {}", - self.first, self.second - ) - } -} -impl std::error::Error for UnitMismatch {} - -/// Sum a workload's per-node cost annotations into one workload-wide total, -/// counting each distinct `workload_node_id` exactly once — the "shared -/// nodes must be counted once in workload totals" requirement. `entries` is -/// `(workload_node_id, annotation)` pairs; an entry with `workload_node_id: -/// None` is never deduplicated against anything else (each is its own, -/// always-unique contribution). If any distinct entry has no value, the -/// aggregate is unavailable as well: returning a numeric subtotal would -/// misrepresent it as the whole workload's cost. -/// -/// Rejects the sum with [`UnitMismatch`] the moment two distinct entries -/// disagree on [`CostUnit`] — "unit-incompatible aggregation is rejected". -pub fn sum_workload_costs<'a, I>(entries: I) -> Result -where - I: IntoIterator, &'a CostAnnotation)>, -{ - let mut seen_ids = std::collections::HashSet::new(); - let mut unit: Option = None; - let mut total = 0.0_f64; - let mut counted_any = false; - let mut missing_any = false; - let mut model_versions: Vec = Vec::new(); - let mut evidence_versions: Vec = Vec::new(); - let mut cache_profiles: Vec> = Vec::new(); - - for (workload_node_id, annotation) in entries { - if let Some(id) = workload_node_id { - if !seen_ids.insert(id) { - continue; // already counted this shared node once - } - } - match unit { - None => unit = Some(annotation.unit), - Some(existing) if existing != annotation.unit => { - return Err(UnitMismatch { - first: existing, - second: annotation.unit, - }); - } - _ => {} - } - let Some(value) = annotation.value.filter(|value| value.is_finite()) else { - missing_any = true; - continue; - }; - total += value; - if !total.is_finite() { - missing_any = true; - continue; - } - counted_any = true; - if let Some(version) = &annotation.model_version { - if !model_versions.contains(version) { - model_versions.push(version.clone()); - } - } - if let Some(version) = &annotation.evidence_version { - if !evidence_versions.contains(version) { - evidence_versions.push(version.clone()); - } - } - if !cache_profiles.contains(&annotation.cache_profile) { - cache_profiles.push(annotation.cache_profile.clone()); - } - } - - // One workload total must not silently combine different immutable - // catalog/runtime generations. Such a subtotal is not a comparable - // snapshot even though each component is individually numeric. - if evidence_versions.len() > 1 || cache_profiles.len() > 1 { - missing_any = true; - } - - Ok(if counted_any && !missing_any { - CostAnnotation { - value: Some(total), - unit: unit.expect("counted_any implies unit was set"), - source: CostSource::Modeled, - baseline: None, - delta: None, - benefit_ratio: None, - model_version: if model_versions.len() == 1 { - Some(model_versions.remove(0)) - } else { - None - }, - evidence_version: if evidence_versions.len() == 1 { - Some(evidence_versions.remove(0)) - } else { - None - }, - cache_profile: if cache_profiles.len() == 1 { - cache_profiles.remove(0) - } else { - None - }, - benchmark_id: None, - inputs: Vec::new(), - } - } else { - CostAnnotation::unavailable(unit.unwrap_or(CostUnit::CostUnits)) - }) -} - -/// Whole-selected-workload cost/benefit — one query's (or one workload -/// batch's) aggregate baseline, selected, and benefit, built from -/// [`sum_workload_costs`] over that scope's own per-decision node -/// annotations. -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -pub struct WorkloadCostSummary { - pub baseline_cost: CostAnnotation, - pub selected_cost: CostAnnotation, - pub benefit: CostAnnotation, -} - -/// Build a [`WorkloadCostSummary`] from one `(workload_node_id, baseline, -/// selected)` triple per decision node in scope — see -/// [`WorkloadCostSummary`]'s own doc. `model_version` labels the resulting -/// `benefit` annotation. -pub fn workload_cost_summary<'a, I>( - entries: I, - model_version: impl Into, -) -> Result -where - I: IntoIterator, &'a CostAnnotation, &'a CostAnnotation)> + Clone, -{ - let model_version = model_version.into(); - let baseline_cost = sum_workload_costs( - entries - .clone() - .into_iter() - .map(|(id, baseline, _)| (id, baseline)), - )?; - let selected_cost = - sum_workload_costs(entries.into_iter().map(|(id, _, selected)| (id, selected)))?; - - let benefit = match (baseline_cost.value, selected_cost.value) { - (Some(baseline_value), Some(selected_value)) - if baseline_cost.unit == selected_cost.unit - && baseline_cost.evidence_version == selected_cost.evidence_version - && baseline_cost.cache_profile == selected_cost.cache_profile - && !model_version.trim().is_empty() => - { - let delta = baseline_value - selected_value; - CostAnnotation { - value: Some(delta), - unit: selected_cost.unit, - source: CostSource::Modeled, - baseline: Some(BaselineRef::PreAsapRecomputation), - delta: None, - benefit_ratio: benefit_ratio(baseline_value, delta), - model_version: Some(model_version), - evidence_version: selected_cost.evidence_version.clone(), - cache_profile: selected_cost.cache_profile.clone(), - benchmark_id: None, - inputs: Vec::new(), - } - } - (Some(_), Some(_)) if baseline_cost.unit != selected_cost.unit => { - return Err(UnitMismatch { - first: baseline_cost.unit, - second: selected_cost.unit, - }) - } - (Some(_), Some(_)) => CostAnnotation::unavailable(selected_cost.unit), - _ => CostAnnotation::unavailable(selected_cost.unit), - }; - - Ok(WorkloadCostSummary { - baseline_cost, - selected_cost, - benefit, - }) -} - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn unavailable_has_no_value_and_no_synthetic_number() { - let a = CostAnnotation::unavailable(CostUnit::CostUnits); - assert_eq!(a.value, None); - assert_eq!(a.source, CostSource::Unavailable); - } - - #[test] - fn with_baseline_computes_delta_and_ratio() { - let a = CostAnnotation::modeled(3.0, CostUnit::CostUnits, "v1", vec![]) - .with_baseline(BaselineRef::PreAsapRecomputation, 10.0); - assert_eq!(a.delta, Some(7.0)); - assert_eq!(a.benefit_ratio, Some(0.7)); - assert_eq!(a.baseline, Some(BaselineRef::PreAsapRecomputation)); - } - - #[test] - fn modeled_annotations_reject_non_finite_values_inputs_and_blank_versions() { - for annotation in [ - CostAnnotation::modeled(f64::NAN, CostUnit::CostUnits, "v1", vec![]), - CostAnnotation::modeled(f64::INFINITY, CostUnit::CostUnits, "v1", vec![]), - CostAnnotation::modeled(1.0, CostUnit::CostUnits, " ", vec![]), - CostAnnotation::modeled( - 1.0, - CostUnit::CostUnits, - "v1", - vec![CostInput::new("rows", f64::NAN)], - ), - ] { - assert_eq!(annotation.source, CostSource::Unavailable); - assert_eq!(annotation.value, None); - } - } - - #[test] - fn modeled_annotations_retain_evidence_generation_separately() { - let annotation = ann(3.0, CostUnit::CostUnits).with_evidence_version("catalog-42"); - assert_eq!(annotation.model_version.as_deref(), Some("v1")); - assert_eq!(annotation.evidence_version.as_deref(), Some("catalog-42")); - - let unavailable = ann(3.0, CostUnit::CostUnits).with_evidence_version(" \t"); - assert_eq!(unavailable.source, CostSource::Unavailable); - assert_eq!(unavailable.value, None); - } - - #[test] - fn modeled_annotations_export_cache_profile_provenance() { - let annotation = ann(3.0, CostUnit::CostUnits).with_cache_profile("no-cache-v1"); - assert_eq!(annotation.cache_profile.as_deref(), Some("no-cache-v1")); - assert_eq!( - serde_json::to_value(annotation).unwrap()["cache_profile"], - "no-cache-v1" - ); - - let unavailable = ann(3.0, CostUnit::CostUnits).with_cache_profile(" \t"); - assert_eq!(unavailable.source, CostSource::Unavailable); - } - - // Totals and benefits must not combine known and incompatible/unknown cache assumptions. - #[test] - fn workload_cache_provenance_must_be_complete_and_equal() { - let raw = ann(10.0, CostUnit::CostUnits).with_cache_profile("no-cache-v1"); - let warm = ann(1.0, CostUnit::CostUnits).with_cache_profile("warm-v1"); - let unknown = ann(2.0, CostUnit::CostUnits); - assert_eq!( - sum_workload_costs([(None, &raw), (None, &unknown)]) - .unwrap() - .value, - None - ); - assert_eq!( - workload_cost_summary([(None, &raw, &warm)], "v1") - .unwrap() - .benefit - .value, - None - ); - assert_eq!( - sum_workload_costs([(None, &raw), (None, &raw)]) - .unwrap() - .cache_profile, - raw.cache_profile - ); - } - - #[test] - fn benefit_ratio_unavailable_when_baseline_not_positive() { - assert_eq!(benefit_ratio(0.0, 5.0), None); - assert_eq!(benefit_ratio(-1.0, 5.0), None); - assert_eq!(benefit_ratio(2.0, 1.0), Some(0.5)); - assert_eq!(benefit_ratio(f64::INFINITY, 1.0), None); - assert_eq!(benefit_ratio(2.0, f64::NAN), None); - } - - #[test] - fn with_baseline_on_unavailable_annotation_is_a_no_op() { - let a = CostAnnotation::unavailable(CostUnit::CostUnits) - .with_baseline(BaselineRef::PreAsapRecomputation, 10.0); - assert_eq!(a.value, None); - assert_eq!(a.delta, None); - assert_eq!(a.baseline, None); - } - - #[test] - fn total_cost_requires_a_horizon_for_a_recurring_rate() { - assert_eq!(total_cost(Some(2.0), 5.0, None), Some(10.0)); - assert_eq!(total_cost(Some(2.0), 5.0, Some(1.0)), Some(11.0)); - assert_eq!(total_cost(None, 5.0, Some(4.0)), Some(4.0)); - assert_eq!(total_cost(None, 5.0, None), None); - } - - #[test] - fn total_cost_rejects_a_non_finite_or_negative_horizon() { - assert_eq!(total_cost(Some(2.0), f64::NAN, None), None); - assert_eq!(total_cost(Some(2.0), f64::INFINITY, None), None); - assert_eq!(total_cost(Some(2.0), -1.0, None), None); - } - - /// A non-finite `recurring_cost_rate` (e.g. a stray `NaN` from a future - /// #287 caller) must never silently produce `Some(NaN)` — that's a - /// poisoned total, exactly the kind of fabricated-looking value this - /// module's "never fabricate" rule exists to prevent. - #[test] - fn total_cost_rejects_a_non_finite_recurring_rate() { - assert_eq!(total_cost(Some(f64::NAN), 5.0, None), None); - assert_eq!(total_cost(Some(f64::INFINITY), 5.0, None), None); - assert_eq!(total_cost(Some(f64::NEG_INFINITY), 5.0, Some(1.0)), None); - } - - /// Same guard on the one-shot addend — a `NaN`/infinite one-shot cost - /// must not poison the total either, even when the recurring side is - /// perfectly well-formed. - #[test] - fn total_cost_rejects_a_non_finite_one_shot_cost() { - assert_eq!(total_cost(Some(2.0), 5.0, Some(f64::NAN)), None); - assert_eq!(total_cost(None, 5.0, Some(f64::INFINITY)), None); - } - - fn ann(value: f64, unit: CostUnit) -> CostAnnotation { - CostAnnotation::modeled(value, unit, "v1", vec![]) - } - - #[test] - fn sum_workload_costs_counts_a_shared_node_once() { - let a = ann(5.0, CostUnit::CostUnits); - let b = ann(5.0, CostUnit::CostUnits); - let c = ann(2.0, CostUnit::CostUnits); - // Node id 1 shared by two queries (same decision, same cost) must - // only be counted once; node id 2 is a distinct contribution. - let total = sum_workload_costs(vec![(Some(1), &a), (Some(1), &b), (Some(2), &c)]).unwrap(); - assert_eq!(total.value, Some(7.0), "5.0 (once) + 2.0, not 5+5+2"); - } - - #[test] - fn sum_workload_costs_never_deduplicates_entries_with_no_workload_id() { - let a = ann(3.0, CostUnit::CostUnits); - let b = ann(3.0, CostUnit::CostUnits); - let total = sum_workload_costs(vec![(None, &a), (None, &b)]).unwrap(); - assert_eq!(total.value, Some(6.0)); - } - - #[test] - fn sum_workload_costs_is_unavailable_when_any_distinct_entry_is_unavailable() { - let known = ann(4.0, CostUnit::CostUnits); - let unavailable = CostAnnotation::unavailable(CostUnit::CostUnits); - let total = sum_workload_costs(vec![(None, &known), (None, &unavailable)]).unwrap(); - assert_eq!(total.value, None); - assert_eq!(total.source, CostSource::Unavailable); - } - - #[test] - fn sum_workload_costs_rejects_mismatched_units() { - let rate = ann(1.0, CostUnit::CostUnitsPerSecond); - let total = ann(1.0, CostUnit::CostUnits); - let err = sum_workload_costs(vec![(None, &rate), (None, &total)]).unwrap_err(); - assert_eq!(err.first, CostUnit::CostUnitsPerSecond); - assert_eq!(err.second, CostUnit::CostUnits); - } - - #[test] - fn workload_totals_reject_mixed_evidence_generations() { - let first = ann(1.0, CostUnit::CostUnits).with_evidence_version("snapshot-a"); - let second = ann(2.0, CostUnit::CostUnits).with_evidence_version("snapshot-b"); - let total = sum_workload_costs(vec![(None, &first), (None, &second)]).unwrap(); - assert_eq!(total.source, CostSource::Unavailable); - assert_eq!(total.value, None); - - let baseline = ann(5.0, CostUnit::CostUnits).with_evidence_version("snapshot-a"); - let selected = ann(2.0, CostUnit::CostUnits).with_evidence_version("snapshot-b"); - let summary = workload_cost_summary(vec![(None, &baseline, &selected)], "v1").unwrap(); - assert_eq!(summary.benefit.source, CostSource::Unavailable); - assert_eq!(summary.benefit.value, None); - } - - #[test] - fn workload_cost_summary_computes_benefit_from_deduplicated_totals() { - let baseline1 = ann(10.0, CostUnit::CostUnits); - let selected1 = ann(3.0, CostUnit::CostUnits); - let baseline2 = ann(10.0, CostUnit::CostUnits); // same shared node - let selected2 = ann(3.0, CostUnit::CostUnits); - let baseline3 = ann(4.0, CostUnit::CostUnits); - let selected3 = ann(1.0, CostUnit::CostUnits); - - let entries = vec![ - (Some(1_u32), &baseline1, &selected1), - (Some(1_u32), &baseline2, &selected2), // duplicate of node 1 — must not double count - (Some(2_u32), &baseline3, &selected3), - ]; - let summary = workload_cost_summary(entries, "test-v1").unwrap(); - assert_eq!(summary.baseline_cost.value, Some(14.0), "10 (once) + 4"); - assert_eq!(summary.selected_cost.value, Some(4.0), "3 (once) + 1"); - assert_eq!(summary.benefit.value, Some(10.0)); - assert_eq!(summary.benefit.benefit_ratio, Some(10.0 / 14.0)); - } -} diff --git a/crates/types/src/lib.rs b/crates/types/src/lib.rs index 92c9f1891..ab72b083a 100644 --- a/crates/types/src/lib.rs +++ b/crates/types/src/lib.rs @@ -11,7 +11,7 @@ //! beside the cost and accuracy models. //! - [`physical`] — #509 Stage 2 decision data: exact-operator schema helpers //! and window-summary pane primitives. -//! - [`types`] / [`cost`] — accuracy targets and cost annotations. +//! - [`types`] / [`cost`] — accuracy targets and cost units. pub mod cost; pub mod deployment; pub mod ir;