From 1f15b8158189994f96d1f38cf3b5439bb5e54c20 Mon Sep 17 00:00:00 2001 From: zz_y Date: Wed, 2 Sep 2026 16:16:18 -0600 Subject: [PATCH 1/2] feat(query): bind SID reads to active materializations --- .../asap_query_engine/post_asap_readout.rs | 10 +++ .../asap_query_engine/summary_executor.rs | 63 ++++++++++++++++--- 2 files changed, 64 insertions(+), 9 deletions(-) diff --git a/data_plane/src/query_engines/asap_query_engine/post_asap_readout.rs b/data_plane/src/query_engines/asap_query_engine/post_asap_readout.rs index 9fc8b6c7..86dade24 100644 --- a/data_plane/src/query_engines/asap_query_engine/post_asap_readout.rs +++ b/data_plane/src/query_engines/asap_query_engine/post_asap_readout.rs @@ -79,6 +79,16 @@ pub fn execute_post_asap_readout( t0_ms, t1_ms, is_cumulative, + allowed_materializations: backend_plan.map(|plan| { + plan.routing + .iter() + .filter(|route| { + route.storage_backend + == control_plane::backend_plan::StorageBackend::SketchStore + }) + .map(|route| route.materialization) + .collect() + }), }; match execute(&node, &ctx) { 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 4a900b47..0be2b334 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 @@ -101,6 +101,10 @@ pub struct QueryExecutionContext<'a> { /// `readout_cumulative`); `false` for a per-window matrix (one merged /// answer per window, via `readout_per_window`). pub is_cumulative: bool, + /// Materializations authorized by the active BackendPlan's warm routes. + /// `None` is the explicit legacy/no-plan mode; `Some` fails closed and + /// excludes stale or unrelated SIDs even when their metric/family match. + pub allowed_materializations: Option>, } /// One candidate sid, already carrying its `[t0, t1]` data and decode @@ -388,16 +392,25 @@ impl<'a> SummaryExecutor for QueryExecutionContext<'a> { for sid in candidate_sids { let candidate = self .index - .with_instance(sid, |m| match &m.agg_kind { - AggKind::Sketch { kind, config, .. } => { - summary_params_match(sketch, params, *kind, config) - .then(|| to_delta_kind(*kind, config)) - .flatten() - .map(Candidate::Sketch) + .with_instance(sid, |m| { + if self + .allowed_materializations + .as_ref() + .is_some_and(|allowed| !allowed.contains(&m.policy_fp)) + { + return None; } - AggKind::ExactAgg { agg_type, .. } => { - exact_agg_kind_match(sketch, params, *agg_type) - .then_some(Candidate::ExactAgg(*agg_type)) + match &m.agg_kind { + AggKind::Sketch { kind, config, .. } => { + summary_params_match(sketch, params, *kind, config) + .then(|| to_delta_kind(*kind, config)) + .flatten() + .map(Candidate::Sketch) + } + AggKind::ExactAgg { agg_type, .. } => { + exact_agg_kind_match(sketch, params, *agg_type) + .then_some(Candidate::ExactAgg(*agg_type)) + } } }) .flatten(); @@ -1358,6 +1371,7 @@ mod tests { t0_ms: T0, t1_ms: T1, is_cumulative: true, + allowed_materializations: None, } } @@ -1367,9 +1381,40 @@ mod tests { t0_ms: T0, t1_ms: T1, is_cumulative: false, + allowed_materializations: None, } } + #[test] + fn backend_plan_materialization_filter_excludes_stale_sid() { + let idx = SketchStore::new(); + let mut stale = kll_meta(1, "latency_ms", &[]); + stale.policy_fp = asap_types::PolicyFingerprint(10); + idx.register(stale); + let exec = QueryExecutionContext { + index: &idx, + t0_ms: T0, + t1_ms: T1, + is_cumulative: true, + allowed_materializations: Some( + [asap_types::PolicyFingerprint(20)].into_iter().collect(), + ), + }; + let family = sketch_family( + planner_types::post_asap::SketchAlgorithm::Kll, + planner_types::post_asap::SketchParams::Kll { k: 200 }, + ); + let result = exec + .find_candidates( + &family, + &ColumnRef::SampleValue, + &Reduction::PerEntity, + &scan_node("latency_ms", None), + ) + .expect("candidate lookup"); + assert!(result.is_empty()); + } + #[test] fn single_kll_sid_quantile_readout() { let idx = SketchStore::new(); From c663c9ead376a6148953ba655428383cffd2c874 Mon Sep 17 00:00:00 2001 From: zz_y Date: Wed, 2 Sep 2026 16:20:36 -0600 Subject: [PATCH 2/2] fix(query): authorize uniquely routed legacy SID --- .../src/query_engines/asap_query_engine/summary_executor.rs | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) 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 0be2b334..4846ae73 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 @@ -396,7 +396,11 @@ impl<'a> SummaryExecutor for QueryExecutionContext<'a> { if self .allowed_materializations .as_ref() - .is_some_and(|allowed| !allowed.contains(&m.policy_fp)) + .is_some_and(|allowed| { + !allowed.contains(&m.policy_fp) + && !(m.policy_fp == asap_types::PolicyFingerprint::UNSET + && allowed.len() == 1) + }) { return None; }