From 15f03e6f964a9247e2b93f13f740328014cb859c Mon Sep 17 00:00:00 2001 From: zz_y Date: Wed, 2 Sep 2026 16:12:11 -0600 Subject: [PATCH] feat(query): require BackendPlan warm routes --- .../asap_query_engine/post_asap_planner.rs | 106 ++++++++++++++++-- .../asap_query_engine/summary_executor.rs | 4 + 2 files changed, 98 insertions(+), 12 deletions(-) diff --git a/data_plane/src/query_engines/asap_query_engine/post_asap_planner.rs b/data_plane/src/query_engines/asap_query_engine/post_asap_planner.rs index ab07579b6..9430386ec 100644 --- a/data_plane/src/query_engines/asap_query_engine/post_asap_planner.rs +++ b/data_plane/src/query_engines/asap_query_engine/post_asap_planner.rs @@ -94,6 +94,10 @@ pub enum LoweringSkip { /// outcome, just detected one step earlier so the caller can skip /// without even constructing a `QueryExecutionContext`. NotRealized, + /// ASAPPlanner produced a maintained-summary plan, but the installed + /// BackendPlan has no warm route to a materialization with the exact + /// source, family and parameters required by that post-ASAP tree. + NoWarmRoute(String), /// The tree lowered successfully, but `crate::query_engines::asap_query_engine::summary_exec::execute()` /// itself returned `Err` (`NoCandidates`, `MergeKindParamsMismatch`, /// a decode/merge failure surfaced from `summary_executor.rs`, ...). @@ -215,7 +219,18 @@ fn observed_family_for_metric_from_plan( plan: &control_plane::backend_plan::BackendPlan, metric: &str, ) -> Option<(SketchAlgorithm, SketchParams)> { - plan.materializations.values().find_map(|m| { + let mut warm_fingerprints: Vec<_> = plan + .routing + .iter() + .filter(|route| { + route.storage_backend == control_plane::backend_plan::StorageBackend::SketchStore + }) + .map(|route| route.materialization) + .collect(); + warm_fingerprints.sort_by_key(|fp| fp.0); + warm_fingerprints.dedup(); + warm_fingerprints.into_iter().find_map(|fingerprint| { + let m = plan.materializations.get(&fingerprint)?; if !matches!(&m.source, planner_types::pre_asap::Source::TimeSeries { metric: mm } if mm == metric) { return None; @@ -229,6 +244,66 @@ fn observed_family_for_metric_from_plan( }) } +fn backend_plan_covers_post_asap( + plan: &control_plane::backend_plan::BackendPlan, + node: &SummaryNode, +) -> Result<(), LoweringSkip> { + fn visit( + plan: &control_plane::backend_plan::BackendPlan, + node: &SummaryNode, + ) -> Result<(), LoweringSkip> { + match &node.expr { + SummaryExpr::SummaryAgg { child, family, .. } => { + let metric = super::summary_executor::find_metric_in_query_expr_from_summary(child) + .ok_or_else(|| { + LoweringSkip::NoWarmRoute("summary has no time-series source".into()) + })?; + let (kind, params): (asap_types::SummaryKind, asap_types::SummaryParams) = + match family { + planner_types::post_asap::SummaryFamilyType::ExactAggregate( + kind, + params, + ) => (kind.clone().into(), params.clone().into()), + planner_types::post_asap::SummaryFamilyType::Sketch(kind, _) => { + (kind.clone().into(), kind.params().clone().into()) + } + _ => { + return Err(LoweringSkip::NoWarmRoute(format!( + "unsupported maintained family for metric `{metric}`" + ))) + } + }; + let covered = plan.routing.iter().any(|route| { + route.storage_backend == control_plane::backend_plan::StorageBackend::SketchStore + && plan.materializations.get(&route.materialization).is_some_and(|m| { + matches!(&m.source, planner_types::pre_asap::Source::TimeSeries { metric: mm } if mm == &metric) + && m.kind == kind + && m.params == params + }) + }); + if !covered { + return Err(LoweringSkip::NoWarmRoute(format!( + "no warm BackendPlan route for metric `{metric}` and family `{kind:?}`" + ))); + } + visit(plan, child) + } + SummaryExpr::SummaryEstimate { summary_input, .. } => visit(plan, summary_input), + SummaryExpr::SummaryMerge { children } => { + for child in children { + visit(plan, child)?; + } + Ok(()) + } + SummaryExpr::KeepPreAsap(_) => Ok(()), + _ => Err(LoweringSkip::NoWarmRoute( + "post-ASAP operator is not executable by the warm runtime".into(), + )), + } + } + visit(plan, node) +} + /// Lower a raw PromQL query string to the `SummaryNode` tree /// `crate::query_engines::asap_query_engine::summary_exec::execute`/`SummaryExecutor` needs — the actual /// serving cutover (`live_serve.rs`). Returns `Err` for any shape serving @@ -288,6 +363,9 @@ pub fn plan_promql_to_post_asap( if matches!(node.expr, SummaryExpr::KeepPreAsap(_)) { Err(LoweringSkip::NotRealized) } else { + if let Some(plan) = backend_plan { + backend_plan_covers_post_asap(plan, &node)?; + } Ok(node) } } @@ -420,7 +498,9 @@ mod tests { AccuracyBound, Capability, SketchInstanceMetadata, SketchKindHandle, }; use asap_types::enums::WindowKind; - use control_plane::backend_plan::{BackendPlan, Materialization, WindowSpec}; + use control_plane::backend_plan::{ + BackendPlan, Materialization, RoutingEntry, StorageBackend, WindowSpec, + }; use planner_types::pre_asap::{ColumnRef, Source}; use std::collections::HashMap; @@ -478,7 +558,11 @@ mod tests { plan_id: 1, generated_at_unix_ms: 0, materializations, - routing: Vec::new(), + routing: vec![RoutingEntry { + satisfies: Capability::QuantileApprox(SketchKindHandle::DDSketch), + materialization: fingerprint, + storage_backend: StorageBackend::SketchStore, + }], monitors: Vec::new(), } } @@ -539,22 +623,20 @@ mod tests { } #[test] - fn plan_present_but_no_materialization_for_metric_falls_back_to_sketchstore() { - // The plan is installed but doesn't cover THIS metric -- - // `observed_family_for_metric_from_plan` returns `None` for - // it, so the lookup must fall through to SketchStore - // reconstruction, not silently fail to observe anything. + fn installed_plan_without_metric_route_fails_closed() { + // Once a BackendPlan is installed it is authoritative. Store + // contents that are absent from the plan must not make a query + // warm-eligible, even when a matching legacy SID still exists. let idx = SketchStore::new(); register_kll(&idx, "m"); let plan = plan_with_ddsketch_materialization("some_other_metric"); - let node = plan_promql_to_post_asap( + let result = plan_promql_to_post_asap( &idx, "quantile_over_time(0.99, m[1m])", accuracy(), Some(&plan), - ) - .expect("should lower"); - assert_eq!(bound_family(&node).0, SketchAlgorithm::Kll); + ); + assert!(matches!(result, Err(LoweringSkip::NoWarmRoute(_)))); } } } diff --git a/data_plane/src/query_engines/asap_query_engine/summary_executor.rs b/data_plane/src/query_engines/asap_query_engine/summary_executor.rs index b648fc6a8..4a900b473 100644 --- a/data_plane/src/query_engines/asap_query_engine/summary_executor.rs +++ b/data_plane/src/query_engines/asap_query_engine/summary_executor.rs @@ -951,6 +951,10 @@ fn find_metric(node: &SummaryNode) -> Option { } } +pub(crate) fn find_metric_in_query_expr_from_summary(node: &SummaryNode) -> Option { + find_metric(node) +} + /// Walk a canonical `QueryExpr` down to its first `Scan { /// source: Source::TimeSeries { metric }, .. }` to recover the target /// metric name. Shared with `post_asap_planner.rs`'s observed-family lookup