From a3311463c7b0605925dd3d7a643e837081e98008 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Sat, 3 Oct 2026 03:55:00 +0000 Subject: [PATCH] fix(planner): enforce response bounds on raw fallback --- crates/asap-aware-mapping/src/cost_model.rs | 7 ++ .../src/summary_maintenance_lifecycle.rs | 63 +++++++++++++++-- crates/planner/tests/summary_sharing.rs | 70 ++++++++++++++++++- .../operator-design-acceptance.md | 11 +++ 4 files changed, 143 insertions(+), 8 deletions(-) diff --git a/crates/asap-aware-mapping/src/cost_model.rs b/crates/asap-aware-mapping/src/cost_model.rs index 35685ede4..5b0bcbb76 100644 --- a/crates/asap-aware-mapping/src/cost_model.rs +++ b/crates/asap-aware-mapping/src/cost_model.rs @@ -888,6 +888,13 @@ pub trait CostModel { self.raw_query_recompute_cost(target) .map(|per_read| Cost(per_read.0 * expected_reads)) } + /// Response latency for executing the original ordinary query once. + /// Separate from amortized workload cost: cheap recomputation can still + /// miss a response deadline. `None` means the bound cannot be checked. + fn raw_query_response_latency_ms(&self, _target: &OperatorNode) -> Option { + None + } + /// Physical feasibility evidence for a complete summary candidate. /// `None` defers admission to physical/deployment compilation; `Some(false)` /// excludes the candidate without changing its computation or parameters. diff --git a/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs b/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs index 782f221af..de56c100d 100644 --- a/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs +++ b/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs @@ -33,8 +33,8 @@ use asap_types::post_asap::{ }; use asap_types::types::AccuracyTarget; use asap_types::workload::{ - DataArrival, DataWorkload, Predictability, QueryRecurrence, QueryWorkload, RepeatedDemand, - TimestampMs, WorkloadError, + DataArrival, DataWorkload, LatencyRequirement, Predictability, QueryRecurrence, QueryWorkload, + RepeatedDemand, TimestampMs, WorkloadError, }; use crate::analytical_cost::AnalyticalCostError; @@ -351,6 +351,8 @@ impl<'a> WorkloadDemand<'a> { #[derive(Debug, thiserror::Error)] pub enum SummaryMaintenanceLifecyclePlanError { + #[error("no latency-feasible plan: raw response estimate {estimate_ms} ms cannot meet {bound_ms} ms, and no costed summary alternative is available")] + NoLatencyFeasiblePlan { bound_ms: f64, estimate_ms: f64 }, #[error(transparent)] InvalidWorkload(#[from] WorkloadError), #[error("optimization horizon must be finite and strictly positive")] @@ -841,10 +843,22 @@ pub fn global_selection_with_summary_maintenance_lifecycles<'a, Id>( let raw = plan .expected_reads .and_then(|reads| cost_model.raw_query_recompute_total_cost(&group.target, reads)); - // Final comparison is atomic: without the raw side, no summary - // override is published even when that summary alone is costed. - if let Some(raw) = raw { - costs.insert_raw(&group.target, raw); + let raw_admitted = raw_response_latency_violation( + &group.target, + WorkloadDemand { + workload, + data_workload, + entry_indices, + }, + cost_model, + ) + .is_none(); + // Missing cost evidence keeps the atomic comparison, unless the + // response deadline has already ruled out raw execution. + if raw.is_some() || !raw_admitted { + if raw_admitted { + costs.insert_raw(&group.target, raw.expect("checked above")); + } if !plan.deployments.is_empty() { if let Some(total) = plan.summary_total_cost { costs.insert(&group.target, candidate, total); @@ -1046,7 +1060,20 @@ pub(crate) fn plan_assembled_dag( plan.raw_recompute_total_cost = plan .expected_reads .and_then(|reads| cost_model.raw_query_recompute_total_cost(target, reads)); - if !plan.selected_raw_recompute + let raw_violation = raw_response_latency_violation(target, demand, cost_model); + if let Some((bound_ms, estimate_ms)) = raw_violation { + if plan.selected_raw_recompute || plan.summary_total_cost.is_none() { + return Err( + SummaryMaintenanceLifecyclePlanError::NoLatencyFeasiblePlan { + bound_ms, + estimate_ms, + } + .into(), + ); + } + } + if raw_violation.is_none() + && !plan.selected_raw_recompute && plan.raw_recompute_total_cost.is_none_or(|raw| { plan.summary_total_cost .is_none_or(|summary| raw.0 <= summary.0) @@ -1062,6 +1089,28 @@ pub(crate) fn plan_assembled_dag( Ok(plan) } +/// A quote is per execution, so every consumer's bound applies even when +/// repeated reads make this alternative cheap over the planning horizon. +fn raw_response_latency_violation( + target: &OperatorNode, + demand: WorkloadDemand<'_>, + cost_model: &dyn CostModel, +) -> Option<(f64, f64)> { + let bound_ms = demand + .workload + .entries() + .enumerate() + .filter(|(index, _)| demand.entry_indices.contains(index)) + .filter_map(|(_, entry)| match entry.requirements.response_latency { + LatencyRequirement::ExplicitMaxMs(bound) => Some(bound), + LatencyRequirement::Unspecified => None, + }) + .min_by(f64::total_cmp)?; + let estimate_ms = cost_model.raw_query_response_latency_ms(target)?; + (!estimate_ms.is_finite() || estimate_ms < 0.0 || estimate_ms > bound_ms) + .then_some((bound_ms, estimate_ms)) +} + fn workload_facts( workload: &QueryWorkload, data_workload: Option<&DataWorkload>, diff --git a/crates/planner/tests/summary_sharing.rs b/crates/planner/tests/summary_sharing.rs index ce7505b35..58e497291 100644 --- a/crates/planner/tests/summary_sharing.rs +++ b/crates/planner/tests/summary_sharing.rs @@ -10,7 +10,8 @@ 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, SummaryMaintenanceLifecycleRejection, + SummaryMaintenanceLifecycleCapabilities, SummaryMaintenanceLifecycleCostInputs, + SummaryMaintenanceLifecycleRejection, }; use asap_frontend_sql::SqlCatalog; use asap_planner::{e2e_plan, FrontendInput, UserInput}; @@ -84,6 +85,10 @@ impl CostModel for FixedCosts { }) } + fn raw_query_response_latency_ms(&self, _target: &OperatorNode) -> Option { + self.latency_estimates.then_some(250.0) + } + fn raw_query_recompute_cost(&self, _target: &OperatorNode) -> Option { Some(Cost(self.raw_per_read)) } @@ -603,3 +608,66 @@ async fn e2e_certified_frequency_evaluations_share_one_univmon_state() { assert!(DefaultAccuracyModel.satisfies(guarantee, &AccuracyTarget::Epsilon(epsilon))); } } + +/// A cheap raw scan cannot beat a maintained state when it misses the response deadline. +#[tokio::test] +async fn slow_cheap_raw_recompute_cannot_bypass_the_response_bound() { + 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: 100.0, + raw_per_read: 0.01, + 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("a fast maintained alternative exists"); + assert!(!output.plans[0].plan.selected_raw_recompute); + assert!(output.plans[0].plan.summary_total_cost.is_some()); +} + +/// With only raw execution available, a known missed deadline fails planning. +#[tokio::test] +async fn slow_raw_only_query_reports_no_latency_feasible_plan() { + 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: 100.0, + raw_per_read: 0.01, + latency_estimates: true, + }; + let result = e2e_plan(UserInput::new( + &workload, + FrontendInput::Promql { + now_ms: NOW_MS, + histograms: None, + }, + PlanningModels::builtin() + .with_cost(&costs) + .with_capabilities(SummaryMaintenanceLifecycleCapabilities { + supports_ephemeral: true, + supports_prepared: false, + supports_shared: false, + supports_continuously_maintained: false, + }), + lifecycle(), + )) + .await; + let error = result.expect_err("a known slow raw fallback must not escape summary rejection"); + assert!( + error.to_string().contains("no latency-feasible plan"), + "{error}" + ); +} diff --git a/docs/develop_docs/operator-design-acceptance.md b/docs/develop_docs/operator-design-acceptance.md index daee31c74..c71c80051 100644 --- a/docs/develop_docs/operator-design-acceptance.md +++ b/docs/develop_docs/operator-design-acceptance.md @@ -131,3 +131,14 @@ states, merges them through unified export/native compilation, and checks p99 for both rebuilding raw inputs and reading materialized panes (Examples 3B/4B). This proves the merge building block; automatic window candidate generation and Exponential Histogram construction remain separate work. + +## Planner-layering follow-up: raw response latency + +`CostModel::raw_query_response_latency_ms` quotes one execution separately from +amortized workload cost. A known quote that exceeds any bound of its consumers +cannot win against a feasible summary, and cannot reappear during final raw +fallback. If no costed summary survives, planning returns `NoLatencyFeasiblePlan`. +The public planner's `slow_cheap_raw_recompute_cannot_bypass_the_response_bound` +and `slow_raw_only_query_reports_no_latency_feasible_plan` reproduce both paths. +As with summary latency quotes, missing raw latency evidence remains unchecked; +this does not establish a latency guarantee for an unmeasured deployment.