From 4f182bff096c26691c6a62e947de7d0d854b8a0d Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 3 Sep 2026 18:11:20 -0600 Subject: [PATCH] feat(asapquery): gate warm routes on materialization readiness --- control_plane/src/query_plan.rs | 13 ++ data_plane/src/drivers/query/servers/http.rs | 9 +- .../query_engines/asap_query_engine/engine.rs | 181 ++++++++++++++- .../asap_query_engine/post_asap_planner.rs | 3 + .../types/hot_reload_config.rs | 212 +++++++++++++++++- docs/user_guide/asapquery-profile.md | 8 + 6 files changed, 415 insertions(+), 11 deletions(-) diff --git a/control_plane/src/query_plan.rs b/control_plane/src/query_plan.rs index 6f627a0a0..a0c4c003d 100644 --- a/control_plane/src/query_plan.rs +++ b/control_plane/src/query_plan.rs @@ -84,6 +84,19 @@ pub struct InstantExecution { } impl QueryPlanEntry { + /// Materializations this executable DAG reads, in stable node order. + /// Serving uses this set for readiness accounting; it never performs a + /// catalog candidate search to reconstruct dependencies. + pub fn materialization_bindings(&self) -> Vec<&MaterializationBinding> { + self.nodes + .values() + .filter_map(|node| match node { + QueryPlanNode::ReadMaterialization { binding } => Some(binding), + _ => None, + }) + .collect() + } + pub fn compile_bound( query_id: String, canonical_promql: String, diff --git a/data_plane/src/drivers/query/servers/http.rs b/data_plane/src/drivers/query/servers/http.rs index 7b0ad9a82..8eb17ea0f 100644 --- a/data_plane/src/drivers/query/servers/http.rs +++ b/data_plane/src/drivers/query/servers/http.rs @@ -5869,10 +5869,17 @@ async fn handle_physical_plan_status(State(state): State) -> axum::res ) .into_response(); }; + let materializations = state + .active_physical_plan + .as_ref() + .map(|active| active.materialization_statuses()) + .unwrap_or_default(); ( StatusCode::OK, axum::Json(serde_json::json!({ - "status": "success", "plans": lifecycle.statuses() + "status": "success", + "plans": lifecycle.statuses(), + "materializations": materializations })), ) .into_response() diff --git a/data_plane/src/query_engines/asap_query_engine/engine.rs b/data_plane/src/query_engines/asap_query_engine/engine.rs index 9e9592922..a26878b05 100644 --- a/data_plane/src/query_engines/asap_query_engine/engine.rs +++ b/data_plane/src/query_engines/asap_query_engine/engine.rs @@ -4,6 +4,66 @@ use std::sync::Arc; use asap_types::query_requirements::QueryRequirements; use asap_types::KeyByLabelNames; +#[derive(Clone)] +struct QueryReadinessRequirement { + materializations: Vec, + max_window_ms: u64, +} + +fn readiness_requirement( + entry: &control_plane::query_plan::QueryPlanEntry, +) -> QueryReadinessRequirement { + let bindings = entry.materialization_bindings(); + let mut materializations = bindings + .iter() + .map(|binding| binding.materialization) + .collect::>(); + materializations.sort_unstable(); + materializations.dedup(); + QueryReadinessRequirement { + materializations, + max_window_ms: bindings + .iter() + .map(|binding| binding.window_ms) + .max() + .unwrap_or(0), + } +} + +fn complete_window_coverage( + coverage: Option<(u64, u64)>, + t0_ms: u64, + t1_ms: u64, + window_ms: u64, +) -> bool { + let Some((coverage_start, coverage_end)) = coverage else { + return false; + }; + if window_ms == 0 || coverage_start > coverage_end || t0_ms > t1_ms { + return false; + } + coverage_start.saturating_sub(window_ms) <= t0_ms + && coverage_end >= t1_ms + && coverage_end + .saturating_sub(coverage_start) + .saturating_add(window_ms) + >= t1_ms.saturating_sub(t0_ms) +} + +#[cfg(test)] +mod readiness_coverage_tests { + use super::complete_window_coverage; + + #[test] + fn readiness_requires_complete_and_fresh_window_span() { + assert!(!complete_window_coverage(None, 100, 400, 100)); + assert!(!complete_window_coverage(Some((200, 300)), 100, 500, 100)); + assert!(!complete_window_coverage(Some((300, 500)), 100, 500, 100)); + assert!(complete_window_coverage(Some((200, 500)), 100, 500, 100)); + assert!(complete_window_coverage(Some((500, 500)), 400, 500, 100)); + } +} + #[cfg(test)] use crate::storage_engines::types::KeyByLabelValues; #[cfg(test)] @@ -309,9 +369,16 @@ impl ASAPQueryEngine { )); }; - let planned = match self.physical_plan_snapshot() { + let physical_plan = self.physical_plan_snapshot(); + let mut readiness = None; + let planned = match physical_plan.as_ref() { Some(physical_plan) => match physical_plan.query_plan.lookup(query) { Ok(query_entry) => { + readiness = Some(( + physical_plan.backend_plan.plan_id, + physical_plan.backend_plan.plan_version, + readiness_requirement(query_entry), + )); crate::query_engines::asap_query_engine::live_serve::serve_from_query_plan( idx, query_entry, @@ -337,7 +404,50 @@ impl ASAPQueryEngine { "no active physical QueryPlan".into(), )), }; - let result = planned.map_err(|reason| { + let result = planned.and_then(|result| { + if let Some((plan_id, plan_version, requirement)) = readiness.as_ref() { + let complete = !requirement.materializations.is_empty() + && complete_window_coverage( + result.coverage, + start_ms, + end_ms, + requirement.max_window_ms, + ); + let Some(active) = self.active_physical_plan.as_ref() else { + return Err(crate::query_engines::asap_query_engine::post_asap_planner::LoweringSkip::MaterializationNotReady( + "physical readiness registry is unavailable".into(), + )); + }; + if !complete { + active.mark_materializing( + *plan_id, + *plan_version, + &requirement.materializations, + result.coverage, + ); + return Err(crate::query_engines::asap_query_engine::post_asap_planner::LoweringSkip::MaterializationNotReady( + format!("coverage {:?} does not completely and freshly cover [{start_ms}, {end_ms}]", result.coverage), + )); + } + let coverage = result.coverage.expect("complete coverage checked above"); + if !active.mark_ready( + *plan_id, + *plan_version, + &requirement.materializations, + coverage, + ) || !active.mark_serving( + *plan_id, + *plan_version, + &requirement.materializations, + coverage, + ) { + return Err(crate::query_engines::asap_query_engine::post_asap_planner::LoweringSkip::MaterializationNotReady( + "physical generation changed while checking readiness".into(), + )); + } + } + Ok(result) + }).map_err(|reason| { if let Some(req) = Self::requirements_from_query_str(query) { crate::drivers::control_plane_client::spawn_capability_miss_notify( &self.control_plane_client, @@ -356,8 +466,9 @@ impl ASAPQueryEngine { // Matrix shape — the range_query wire format requires it. let warm_qr = asap_tier_result_to_query_result(result.clone(), end_ms, true); - // FIX 3 — coverage-aware warm+archive HYBRID STITCH for RANGE - // queries. The instant path (`execute`) already stitches when warm + // Complete coverage is a prerequisite above. Hybrid stitching is + // retained for legacy/test callers without an active QueryPlan only. + // The instant path (`execute`) already stitches when warm // coverage is narrower than the request; the range path historically // returned warm-only, so a request `[start_ms, end_ms]` whose warm // sketches only cover a suffix `[cov_lo, cov_hi]` lost the @@ -556,11 +667,20 @@ impl crate::query_engines::routing::query_engine_routing::QueryEngine for ASAPQu .duration_since(std::time::SystemTime::UNIX_EPOCH) .map(|d| d.as_millis() as u64) .unwrap_or(0); - let planned = match self.physical_plan_snapshot() { + let physical_plan = self.physical_plan_snapshot(); + let mut readiness = None; + let planned = match physical_plan.as_ref() { Some(physical_plan) => match physical_plan.query_plan.lookup(query) { - Ok(query_entry) => crate::query_engines::asap_query_engine::live_serve::serve_instant_from_query_plan( - idx, query_entry, now_ms, - ), + Ok(query_entry) => { + readiness = Some(( + physical_plan.backend_plan.plan_id, + physical_plan.backend_plan.plan_version, + readiness_requirement(query_entry), + )); + crate::query_engines::asap_query_engine::live_serve::serve_instant_from_query_plan( + idx, query_entry, now_ms, + ) + }, Err(reason) => Err(crate::query_engines::asap_query_engine::post_asap_planner::LoweringSkip::QueryNotPlanned(reason.to_string())), }, #[cfg(test)] @@ -572,7 +692,50 @@ impl crate::query_engines::routing::query_engine_routing::QueryEngine for ASAPQu "no active physical QueryPlan".into(), )), }; - let (result, t0_ms) = planned.map_err(|reason| { + let (result, t0_ms) = planned.and_then(|(result, t0_ms)| { + if let Some((plan_id, plan_version, requirement)) = readiness.as_ref() { + let complete = !requirement.materializations.is_empty() + && complete_window_coverage( + result.coverage, + t0_ms, + now_ms, + requirement.max_window_ms, + ); + let Some(active) = self.active_physical_plan.as_ref() else { + return Err(crate::query_engines::asap_query_engine::post_asap_planner::LoweringSkip::MaterializationNotReady( + "physical readiness registry is unavailable".into(), + )); + }; + if !complete { + active.mark_materializing( + *plan_id, + *plan_version, + &requirement.materializations, + result.coverage, + ); + return Err(crate::query_engines::asap_query_engine::post_asap_planner::LoweringSkip::MaterializationNotReady( + format!("coverage {:?} does not completely and freshly cover [{t0_ms}, {now_ms}]", result.coverage), + )); + } + let coverage = result.coverage.expect("complete coverage checked above"); + if !active.mark_ready( + *plan_id, + *plan_version, + &requirement.materializations, + coverage, + ) || !active.mark_serving( + *plan_id, + *plan_version, + &requirement.materializations, + coverage, + ) { + return Err(crate::query_engines::asap_query_engine::post_asap_planner::LoweringSkip::MaterializationNotReady( + "physical generation changed while checking readiness".into(), + )); + } + } + Ok((result, t0_ms)) + }).map_err(|reason| { if let Some(req) = Self::requirements_from_query_str(query) { crate::drivers::control_plane_client::spawn_capability_miss_notify( &self.control_plane_client, 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 68a4877b0..e2b351b87 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 @@ -68,6 +68,9 @@ pub enum LoweringSkip { /// The installed QueryPlan entry could not be reconstructed or failed /// its internal contract. The request must fail closed. InvalidQueryPlan(String), + /// The QueryPlan binding exists, but its warm materializations do not yet + /// cover the requested evaluation interval with complete, fresh windows. + MaterializationNotReady(String), /// `parse_query_expr_canonical` failed — same failure mode the legacy /// parsing path tolerates and reports as an archive fallback reason. ParseFailed(String), diff --git a/data_plane/src/storage_engines/types/hot_reload_config.rs b/data_plane/src/storage_engines/types/hot_reload_config.rs index 9ceab4638..959a123bf 100644 --- a/data_plane/src/storage_engines/types/hot_reload_config.rs +++ b/data_plane/src/storage_engines/types/hot_reload_config.rs @@ -97,6 +97,105 @@ pub struct ActivePhysicalPlan { #[derive(Clone)] pub struct HotReloadActivePhysicalPlan { inner: Arc>, + readiness: Arc>, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)] +#[serde(rename_all = "snake_case")] +pub enum MaterializationPhase { + Materializing, + Ready, + Serving, +} + +#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)] +pub struct MaterializationStatus { + pub plan_id: u64, + pub plan_version: u64, + pub materialization: u64, + pub phase: MaterializationPhase, + pub coverage_start_unix_ms: Option, + pub coverage_end_unix_ms: Option, +} + +struct MaterializationReadinessState { + plan_id: u64, + plan_version: u64, + statuses: BTreeMap, +} + +impl MaterializationReadinessState { + fn for_plan(plan: &ActivePhysicalPlan) -> Self { + let plan_id = plan.backend_plan.plan_id; + let plan_version = plan.backend_plan.plan_version; + let statuses = plan + .backend_plan + .materializations + .keys() + .copied() + .map(|materialization| { + ( + materialization, + MaterializationStatus { + plan_id, + plan_version, + materialization: materialization.0, + phase: MaterializationPhase::Materializing, + coverage_start_unix_ms: None, + coverage_end_unix_ms: None, + }, + ) + }) + .collect(); + Self { + plan_id, + plan_version, + statuses, + } + } + + fn update( + &mut self, + plan_id: u64, + plan_version: u64, + materializations: &[asap_types::PolicyFingerprint], + phase: MaterializationPhase, + coverage: Option<(u64, u64)>, + ) -> bool { + if self.plan_id != plan_id || self.plan_version != plan_version { + return false; + } + if materializations + .iter() + .any(|materialization| !self.statuses.contains_key(materialization)) + { + return false; + } + for materialization in materializations { + let status = self + .statuses + .get_mut(materialization) + .expect("materializations were validated above"); + if phase != MaterializationPhase::Materializing + || status.phase == MaterializationPhase::Materializing + { + status.phase = phase; + } + if let Some((start, end)) = coverage { + status.coverage_start_unix_ms = Some( + status + .coverage_start_unix_ms + .map_or(start, |current| current.min(start)), + ); + status.coverage_end_unix_ms = Some( + status + .coverage_end_unix_ms + .map_or(end, |current| current.max(end)), + ); + } + } + true + } } #[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)] @@ -304,8 +403,10 @@ fn status_for(plan: &ActivePhysicalPlan, phase: PhysicalPlanPhase) -> PhysicalPl impl HotReloadActivePhysicalPlan { pub fn new(initial: ActivePhysicalPlan) -> Self { + let readiness = MaterializationReadinessState::for_plan(&initial); Self { inner: Arc::new(ArcSwap::new(Arc::new(initial))), + readiness: Arc::new(std::sync::Mutex::new(readiness)), } } @@ -314,7 +415,85 @@ impl HotReloadActivePhysicalPlan { } pub fn swap(&self, next: ActivePhysicalPlan) -> Arc { - self.inner.swap(Arc::new(next)) + let next_readiness = MaterializationReadinessState::for_plan(&next); + let old = self.inner.swap(Arc::new(next)); + *self + .readiness + .lock() + .expect("materialization readiness lock poisoned") = next_readiness; + old + } + + pub fn materialization_statuses(&self) -> Vec { + self.readiness + .lock() + .expect("materialization readiness lock poisoned") + .statuses + .values() + .cloned() + .collect() + } + + pub fn mark_materializing( + &self, + plan_id: u64, + plan_version: u64, + materializations: &[asap_types::PolicyFingerprint], + coverage: Option<(u64, u64)>, + ) -> bool { + self.update_materializations( + plan_id, + plan_version, + materializations, + MaterializationPhase::Materializing, + coverage, + ) + } + + pub fn mark_ready( + &self, + plan_id: u64, + plan_version: u64, + materializations: &[asap_types::PolicyFingerprint], + coverage: (u64, u64), + ) -> bool { + self.update_materializations( + plan_id, + plan_version, + materializations, + MaterializationPhase::Ready, + Some(coverage), + ) + } + + pub fn mark_serving( + &self, + plan_id: u64, + plan_version: u64, + materializations: &[asap_types::PolicyFingerprint], + coverage: (u64, u64), + ) -> bool { + self.update_materializations( + plan_id, + plan_version, + materializations, + MaterializationPhase::Serving, + Some(coverage), + ) + } + + fn update_materializations( + &self, + plan_id: u64, + plan_version: u64, + materializations: &[asap_types::PolicyFingerprint], + phase: MaterializationPhase, + coverage: Option<(u64, u64)>, + ) -> bool { + self.readiness + .lock() + .expect("materialization readiness lock poisoned") + .update(plan_id, plan_version, materializations, phase, coverage) } } @@ -802,6 +981,37 @@ mod tests { assert_eq!(lifecycle.statuses()[0].phase, PhysicalPlanPhase::Retired); } + #[test] + fn materialization_readiness_is_generation_scoped_and_monotonic() { + let active = HotReloadActivePhysicalPlan::new(physical_plan(7, 1, 100, None)); + let fingerprint = asap_types::PolicyFingerprint(41); + active.readiness.lock().unwrap().statuses.insert( + fingerprint, + MaterializationStatus { + plan_id: 7, + plan_version: 1, + materialization: fingerprint.0, + phase: MaterializationPhase::Materializing, + coverage_start_unix_ms: None, + coverage_end_unix_ms: None, + }, + ); + + assert!(active.mark_ready(7, 1, &[fingerprint], (100, 200))); + assert!(active.mark_serving(7, 1, &[fingerprint], (100, 300))); + assert!(active.mark_materializing(7, 1, &[fingerprint], Some((200, 250)))); + let status = active.materialization_statuses().pop().unwrap(); + assert_eq!(status.phase, MaterializationPhase::Serving); + assert_eq!(status.coverage_start_unix_ms, Some(100)); + assert_eq!(status.coverage_end_unix_ms, Some(300)); + + assert!(!active.mark_ready(7, 2, &[fingerprint], (100, 400))); + assert_eq!( + active.materialization_statuses()[0].coverage_end_unix_ms, + Some(300) + ); + } + #[test] fn physical_plan_rejects_stale_and_expired_generations() { let active = HotReloadActivePhysicalPlan::new(physical_plan(7, 2, 100, None)); diff --git a/docs/user_guide/asapquery-profile.md b/docs/user_guide/asapquery-profile.md index 7e637e4e0..6f0b80646 100644 --- a/docs/user_guide/asapquery-profile.md +++ b/docs/user_guide/asapquery-profile.md @@ -84,6 +84,14 @@ Query serving uses only the installed `QueryPlan` DAG. A query absent from that DAG is a capability miss and goes to the exact Prometheus fallback; the backend does not search materialization candidates while serving. +Activation and warm readiness are deliberately separate. A newly activated +generation exposes each materialization as `materializing`; the QueryPlan path +promotes it through `ready` to `serving` only after its closed-window coverage +fully spans the requested interval. Missing, partial, stale, or generation-raced +coverage is a capability miss, never a partial warm success. Inspect the +per-materialization state and observed coverage through +`GET /api/v1/physical-plan/status`. + Snapshot schema version `1` currently accepts fixed-interval repeating PromQL queries with explicit whole-second lookbacks and fresh ingestion-rate evidence. Unsupported snapshot semantics fail startup rather than silently