diff --git a/crates/asap-aware-mapping/src/cost_model.rs b/crates/asap-aware-mapping/src/cost_model.rs index 9fc17066..35685ede 100644 --- a/crates/asap-aware-mapping/src/cost_model.rs +++ b/crates/asap-aware-mapping/src/cost_model.rs @@ -52,7 +52,8 @@ use crate::exact_composition::ExactOperation; use asap_types::ir::{ASAPOp, Operator, OperatorNode}; use asap_types::post_asap::{ FieldDataType, GroupingStrategy, HydraParams, ResultGuarantee, SketchAlgorithm, SketchParams, - SketchStatistic, SummaryMaintenanceLifecycleGuarantee, SummaryWindowFramework, + SketchStatistic, SummaryMaintenanceLifecycle, SummaryMaintenanceLifecycleGuarantee, + SummaryWindowFramework, }; use asap_types::pre_asap::agg_intent::AggIntent; use asap_types::pre_asap::expr_ir::ColumnRef; @@ -801,6 +802,17 @@ pub trait CostModel { SummaryMaintenanceCapabilities::default() } + /// Estimated response latency for reading `summary` under one selected + /// physical lifecycle. `None` means the deployment has no estimate; it + /// does not make the alternative invalid. + fn summary_read_latency_ms( + &self, + _summary: &OperatorNode, + _lifecycle: &SummaryMaintenanceLifecycle, + ) -> Option { + None + } + /// Replace the sum of selected per-state lifecycle costs with a complete /// root-DAG cost. The default preserves legacy models. Evidence-strict /// models return `None` when any root operation is unavailable; callers diff --git a/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs b/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs index 28d75863..782f221a 100644 --- a/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs +++ b/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs @@ -138,6 +138,7 @@ pub enum SummaryMaintenanceLifecycleRejection { SummaryDoesNotSupportIncrementalUpdates, SummaryDoesNotSupportDeletion, MissingCostEvidence, + ExceedsLatencyBound, } /// One candidate lifecycle policy for a particular summary deployment. @@ -559,6 +560,7 @@ impl SummaryMaintenanceLifecycleCandidates<'_> { #[derive(Debug)] struct SummaryMaintenanceWorkloadFacts { required_accuracy: Vec, + latency_bound_ms: Option, /// Total one-time and recurring reads inside the horizon. `None` means a /// recurrence or horizon was unknown, not zero reads. reads: Option, @@ -706,13 +708,35 @@ fn enumerate_with_profile<'a>( } else { capabilities }; - let alternatives = alternatives_for( + let mut alternatives = alternatives_for( &facts, horizon, capabilities, cost_model.summary_maintenance_capabilities(&summary), cost_model.summary_maintenance_lifecycle_cost_inputs_for_horizon(&summary, horizon), ); + if let Some(bound_ms) = facts.latency_bound_ms { + for alternative in &mut alternatives { + match cost_model.summary_read_latency_ms( + &summary, + &alternative.summary_maintenance_lifecycle, + ) { + Some(latency_ms) if latency_ms.is_finite() && latency_ms >= 0.0 => { + if latency_ms > bound_ms && alternative.rejection.is_none() { + alternative.rejection = Some( + SummaryMaintenanceLifecycleRejection::ExceedsLatencyBound, + ); + alternative.assumptions.push(format!( + "estimated response latency {latency_ms} ms exceeds {bound_ms} ms bound" + )); + } + } + _ => alternative + .assumptions + .push(format!("latency bound {bound_ms} ms unchecked: no estimate")), + } + } + } SummaryMaintenanceDeployment { post_asap_node_id: timing_memo .timed(&summary) @@ -1055,6 +1079,7 @@ fn workload_facts( let mut prepared_eligible = true; let mut requires_deletion = false; let mut required_accuracy = Vec::new(); + let mut latency_bound_ms: Option = None; let entries: Vec<_> = workload.entries().collect(); if workload_entry_indices.is_empty() { @@ -1072,6 +1097,11 @@ fn workload_facts( }, )?; required_accuracy.push(entry.requirements.accuracy.target()); + if let asap_types::workload::LatencyRequirement::ExplicitMaxMs(bound) = + entry.requirements.response_latency + { + latency_bound_ms = Some(latency_bound_ms.map_or(bound, |current| current.min(bound))); + } requires_deletion |= entry.time_selection.lookback.is_some() && entry.time_selection.as_of.is_none() && matches!( @@ -1186,6 +1216,7 @@ fn workload_facts( }; Ok(SummaryMaintenanceWorkloadFacts { required_accuracy, + latency_bound_ms, reads, one_time_invocations, evaluation_rate: has_evaluation_rate.then_some(EvaluationRate(evaluation_rate)), diff --git a/crates/planner/tests/summary_sharing.rs b/crates/planner/tests/summary_sharing.rs index 6cad6b8b..ce7505b3 100644 --- a/crates/planner/tests/summary_sharing.rs +++ b/crates/planner/tests/summary_sharing.rs @@ -10,13 +10,13 @@ use asap_aware_mapping::pass::{PlanOutput, PlanningModels}; use asap_aware_mapping::replacement::{default_size_params, DEFAULT_DELTA}; use asap_aware_mapping::{ CostModel, CostRate, DefaultCostModel, Horizon, LifecycleInput, SummaryMaintenanceCapabilities, - SummaryMaintenanceLifecycleCostInputs, + SummaryMaintenanceLifecycleCostInputs, SummaryMaintenanceLifecycleRejection, }; use asap_frontend_sql::SqlCatalog; use asap_planner::{e2e_plan, FrontendInput, UserInput}; use asap_types::post_asap::{ AccuracyError, BoundExpr, CompositionOperator, ErrorMetric, ProbabilityExpr, ResultGuarantee, - SketchStatistic, + SketchStatistic, SummaryMaintenanceLifecycle, }; use asap_types::post_asap::{FieldDataType, SketchAlgorithm, SketchParams}; use asap_types::pre_asap::agg_intent::default_quantile; @@ -37,6 +37,7 @@ const HORIZON_S: f64 = 3_600.0; struct FixedCosts { build: f64, raw_per_read: f64, + latency_estimates: bool, } impl CostModel for FixedCosts { @@ -72,6 +73,17 @@ impl CostModel for FixedCosts { } } + fn summary_read_latency_ms( + &self, + _summary: &OperatorNode, + lifecycle: &SummaryMaintenanceLifecycle, + ) -> Option { + self.latency_estimates.then_some(match lifecycle { + SummaryMaintenanceLifecycle::Ephemeral => 250.0, + _ => 50.0, + }) + } + fn raw_query_recompute_cost(&self, _target: &OperatorNode) -> Option { Some(Cost(self.raw_per_read)) } @@ -82,6 +94,7 @@ impl CostModel for FixedCosts { const CHEAP_SUMMARY: FixedCosts = FixedCosts { build: 1.0, raw_per_read: 1_000.0, + latency_estimates: false, }; fn requirements(epsilon: f64) -> QueryRequirements { @@ -147,6 +160,74 @@ async fn plan_promql(queries: &[(&str, f64)], costs: &FixedCosts) -> PlanOutput e2e_plan(input).await.expect("workload plans") } +#[tokio::test] +async fn latency_bound_rejects_slow_ephemeral_summary_and_keeps_fast_maintained_one() { + let mut workload = promql_workload(&[("quantile_over_time(0.99, lat[5m])", 0.01)]); + workload.query_workload.repeating_queries.as_mut().unwrap()[0] + .requirements + .response_latency = LatencyRequirement::ExplicitMaxMs(100.0); + let costs = FixedCosts { + build: 1.0, + raw_per_read: 1_000.0, + latency_estimates: true, + }; + let output = e2e_plan(UserInput::new( + &workload, + FrontendInput::Promql { + now_ms: NOW_MS, + histograms: None, + }, + PlanningModels::builtin().with_cost(&costs), + lifecycle(), + )) + .await + .expect("maintained summary satisfies the response bound"); + + let deployment = &output.plans[0].plan.deployments[0]; + assert!(deployment.alternatives.iter().any(|alternative| { + matches!( + alternative.summary_maintenance_lifecycle, + SummaryMaintenanceLifecycle::Ephemeral + ) && alternative.rejection + == Some(SummaryMaintenanceLifecycleRejection::ExceedsLatencyBound) + })); + assert!(deployment.alternatives.iter().any(|alternative| { + matches!( + alternative.summary_maintenance_lifecycle, + SummaryMaintenanceLifecycle::ContinuouslyMaintained + ) && alternative.rejection.is_none() + })); +} + +#[tokio::test] +async fn missing_latency_estimate_keeps_candidate_and_records_unchecked_reason() { + let mut workload = promql_workload(&[("quantile_over_time(0.99, lat[5m])", 0.01)]); + workload.query_workload.repeating_queries.as_mut().unwrap()[0] + .requirements + .response_latency = LatencyRequirement::ExplicitMaxMs(100.0); + let output = e2e_plan(UserInput::new( + &workload, + FrontendInput::Promql { + now_ms: NOW_MS, + histograms: None, + }, + PlanningModels::builtin().with_cost(&CHEAP_SUMMARY), + lifecycle(), + )) + .await + .expect("an unchecked bound keeps legal candidates"); + + assert!(output.plans[0].plan.deployments.iter().any(|deployment| { + deployment.alternatives.iter().any(|alternative| { + alternative.rejection.is_none() + && alternative + .assumptions + .iter() + .any(|assumption| assumption.contains("latency bound 100 ms unchecked")) + }) + })); +} + async fn plan_sql(queries: &[&str], costs: &FixedCosts) -> PlanOutput { let workload = PlanningWorkload { query_workload: QueryWorkload { @@ -430,6 +511,7 @@ async fn shared_amortization_alone_can_beat_raw_recompute() { let costs = FixedCosts { build: 100.0, raw_per_read: 10.0, + latency_estimates: false, }; let p50 = ("quantile_over_time(0.5, lat[5m])", 0.01); let p99 = ("quantile_over_time(0.99, lat[5m])", 0.01);