From 5c162535b4c758cd5458f38a893ac3ecd3be2792 Mon Sep 17 00:00:00 2001 From: zz_y Date: Wed, 2 Sep 2026 20:37:59 -0600 Subject: [PATCH] refactor(runtime): publish one active physical plan snapshot --- .../src/backend_plan/from_stage_config.rs | 10 +-- control_plane/src/physical/compiler.rs | 2 +- crates/asap_types/src/aggregation_config.rs | 10 ++- data_plane/src/drivers/query/servers/http.rs | 50 +++++++----- data_plane/src/main.rs | 52 +++++++++++-- .../routing/backend_storage_routing.rs | 19 +++++ .../types/hot_reload_config.rs | 77 ++++++++++++++++++- 7 files changed, 182 insertions(+), 38 deletions(-) diff --git a/control_plane/src/backend_plan/from_stage_config.rs b/control_plane/src/backend_plan/from_stage_config.rs index d54bd98f5..aa154a160 100644 --- a/control_plane/src/backend_plan/from_stage_config.rs +++ b/control_plane/src/backend_plan/from_stage_config.rs @@ -21,7 +21,7 @@ use std::collections::HashMap; use anyhow::{Context, Result}; use asap_types::SummaryKind; -use asap_types::{AggregationConfig, MonitorSpec, PolicyFingerprint, QueryLanguage}; +use asap_types::{MonitorSpec, PolicyFingerprint, PrecomputeMaterialization, QueryLanguage}; use crate::emit::monitor::{agg_id_for_metric, MonitorIntent}; use crate::emit::stage_config::build_backend_aggregation_json; @@ -138,12 +138,12 @@ pub fn from_stage_config( /// JSON is valid YAML) — rather than re-deriving the field mapping here. pub fn aggregation_config_for_materialization( agg: &BackendAggregation, -) -> Result { +) -> Result { let json = build_backend_aggregation_json(agg); let text = serde_json::to_string(&json).context("serialize synthesized aggregation JSON")?; let yaml_value: serde_yaml::Value = serde_yaml::from_str(&text).context("parse synthesized aggregation JSON as YAML")?; - AggregationConfig::from_yaml_data(&yaml_value, None, QueryLanguage::promql) + PrecomputeMaterialization::from_yaml_data(&yaml_value, None, QueryLanguage::promql) .context("build AggregationConfig from synthesized aggregation JSON") } @@ -247,7 +247,7 @@ mod tests { /// Independently construct the `AggregationConfig` a hand-written /// (non-JSON-round-trip) reader would build for this fixture, so the /// parity test doesn't just check the implementation against itself. - fn hand_built_config(agg: &BackendAggregation) -> AggregationConfig { + fn hand_built_config(agg: &BackendAggregation) -> PrecomputeMaterialization { let parameters: StdHashMap = match &agg.sketch_params { SummaryParams::DDSketch { alpha } => { StdHashMap::from([("alpha".to_string(), serde_json::json!(alpha))]) @@ -266,7 +266,7 @@ mod tests { ]), other => unreachable!("fixture doesn't exercise {other:?}"), }; - AggregationConfig::new( + PrecomputeMaterialization::new( match agg.sketch_kind { SummaryKind::DDSketch => AggregationType::DDSketch, SummaryKind::Kll => AggregationType::DatasketchesKLL, diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index 65540589f..ebbbe4752 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -142,7 +142,7 @@ pub struct CollectorPlan { #[derive(Debug, Clone, Serialize, Deserialize)] pub struct PrecomputePlan { pub envelope: PlanEnvelope, - pub materializations: Vec, + pub materializations: Vec, } /// Complete physical projection of one post-ASAP planning decision. diff --git a/crates/asap_types/src/aggregation_config.rs b/crates/asap_types/src/aggregation_config.rs index e8d39a551..95c7f5d07 100644 --- a/crates/asap_types/src/aggregation_config.rs +++ b/crates/asap_types/src/aggregation_config.rs @@ -22,7 +22,7 @@ use crate::KeyByLabelNames; /// stop reading it). Existing fixtures that still spell out /// `aggregationId: N` parse cleanly — the field is silently dropped. #[derive(Debug, Clone, Serialize, Deserialize)] -pub struct AggregationConfig { +pub struct PrecomputeMaterialization { pub aggregation_type: AggregationType, pub aggregation_sub_type: String, pub parameters: HashMap, @@ -74,7 +74,11 @@ impl AggregationIdInfo { } } -impl AggregationConfig { +/// Compatibility name for legacy streaming-config and precompute call sites. +/// New PhysicalPlan code should use [`PrecomputeMaterialization`]. +pub type AggregationConfig = PrecomputeMaterialization; + +impl PrecomputeMaterialization { #[allow(clippy::too_many_arguments)] pub fn new( aggregation_type: AggregationType, @@ -340,7 +344,7 @@ impl AggregationConfig { } } -impl SerializableToSink for AggregationConfig { +impl SerializableToSink for PrecomputeMaterialization { fn serialize_to_json(&self) -> Value { // PR 5: `aggregationId` is no longer emitted — readers derive it // from content via `PolicyFingerprint::from_config(...).as_u64()`. diff --git a/data_plane/src/drivers/query/servers/http.rs b/data_plane/src/drivers/query/servers/http.rs index be50c40e1..14e8a7412 100644 --- a/data_plane/src/drivers/query/servers/http.rs +++ b/data_plane/src/drivers/query/servers/http.rs @@ -149,6 +149,7 @@ pub struct HttpServer { /// Serializes multi-document physical-plan publication so two control /// plane generations cannot interleave their config and BackendPlan. physical_plan_lock: Arc>, + active_physical_plan: Option, } #[derive(Clone)] @@ -175,6 +176,7 @@ struct AppState { /// See [`HttpServer::probe_cache`]. probe_cache: Option>, physical_plan_lock: Arc>, + active_physical_plan: Option, } impl HttpServer { @@ -200,6 +202,7 @@ impl HttpServer { data_retention_ms: None, probe_cache: None, physical_plan_lock: Arc::new(tokio::sync::Mutex::new(())), + active_physical_plan: None, } } @@ -257,6 +260,14 @@ impl HttpServer { self } + pub fn with_active_physical_plan( + mut self, + handle: crate::storage_engines::types::HotReloadActivePhysicalPlan, + ) -> Self { + self.active_physical_plan = Some(handle); + self + } + /// Attach a per-metric storage-backend routing table. The table is /// wrapped in a hot-reload handle internally so the /// `POST /api/v1/storage_routing` endpoint (Phase α) can swap it @@ -364,6 +375,7 @@ impl HttpServer { data_retention_ms: self.data_retention_ms, probe_cache: self.probe_cache.clone(), physical_plan_lock: self.physical_plan_lock.clone(), + active_physical_plan: self.active_physical_plan.clone(), }; let range_query_endpoint = adapter.get_range_query_endpoint(); @@ -451,6 +463,7 @@ impl HttpServer { data_retention_ms: self.data_retention_ms, probe_cache: self.probe_cache.clone(), physical_plan_lock: self.physical_plan_lock.clone(), + active_physical_plan: self.active_physical_plan.clone(), }; let range_query_endpoint = adapter.get_range_query_endpoint(); @@ -5363,10 +5376,7 @@ async fn handle_post_physical_plan( use axum::response::IntoResponse; use std::collections::BTreeSet; - let (Some(config_handle), Some(plan_handle)) = ( - state.hot_reload_config.as_ref(), - state.hot_reload_backend_plan.as_ref(), - ) else { + let Some(active_handle) = state.active_physical_plan.as_ref() else { return ( StatusCode::SERVICE_UNAVAILABLE, axum::Json(serde_json::json!({ @@ -5426,33 +5436,31 @@ async fn handle_post_physical_plan( }, None => None, }; - if new_routing.is_some() && state.backend_storage_routing.is_none() { - return ( - StatusCode::SERVICE_UNAVAILABLE, - axum::Json(serde_json::json!({ - "status": "error", "error": "storage-routing hot-reload handle is not attached" - })), - ) - .into_response(); - } + let new_routing = new_routing + .map(Arc::new) + .unwrap_or_else(|| active_handle.snapshot().storage_routing.clone()); let _guard = state.physical_plan_lock.lock().await; let plan_id = new_plan.plan_id; - if let Err(error) = plan_handle.install(new_plan) { + let active = crate::storage_engines::types::ActivePhysicalPlan { + precompute_plan: request.precompute_plan, + runtime_config: Arc::new(new_config), + backend_plan: Arc::new(new_plan), + storage_routing: new_routing, + }; + let generated = active.backend_plan.generated_at_unix_ms; + let current = active_handle.snapshot(); + if generated < current.backend_plan.generated_at_unix_ms { return ( StatusCode::CONFLICT, axum::Json(serde_json::json!({ - "status": "error", "error": format!("BackendPlan install error: {error}") + "status": "error", "error": "stale physical-plan generation" })), ) .into_response(); } - config_handle.swap(new_config); - if let (Some(handle), Some(routing)) = (state.backend_storage_routing.as_ref(), new_routing) { - let tenant = routing.tenant().to_string(); - handle.swap_tenant(&tenant, routing); - } - let snap = config_handle.snapshot(); + active_handle.swap(active); + let snap = active_handle.snapshot().runtime_config.clone(); let sid_summary = crate::storage_engines::sketch_db::lifecycle::reconcile_from_streaming_config( state.sketch_index.as_ref(), snap.as_ref(), diff --git a/data_plane/src/main.rs b/data_plane/src/main.rs index 9788ce0bc..c69d15f4b 100644 --- a/data_plane/src/main.rs +++ b/data_plane/src/main.rs @@ -318,9 +318,6 @@ async fn main() -> Result<()> { // startup snapshot; hot-reload currently only affects the // control-plane GET/POST endpoint. Phase 2 will extend the swap // to query execution and ingest routing. - let hot_reload_config = data_plane::storage_engines::types::HotReloadStreamingConfig::from_arc( - streaming_config.clone(), - ); // M2.3.6g — the legacy `SketchStore` construction is gone. // Production data lives in `SketchStore` (allocated below); the @@ -426,9 +423,38 @@ async fn main() -> Result<()> { // engine (serving-time cutover, Phase 4) and the HTTP server (the // push target) so a POST is observable by the next query, same // sharing contract as `hot_reload_config`. - let hot_reload_backend_plan = data_plane::storage_engines::types::HotReloadBackendPlan::new( - control_plane::backend_plan::BackendPlan::default(), + let initial_backend_plan = control_plane::backend_plan::BackendPlan::default(); + let initial_precompute_plan = control_plane::physical::compiler::PrecomputePlan { + envelope: control_plane::physical::compiler::PlanEnvelope { + plan_id: 0, + generated_at_unix_ms: 0, + planner_revision: control_plane::physical::compiler::PLANNER_REVISION.into(), + capability_snapshot_id: "bootstrap".into(), + }, + materializations: streaming_config + .aggregation_configs + .values() + .cloned() + .collect(), + }; + let active_physical_plan = data_plane::storage_engines::types::HotReloadActivePhysicalPlan::new( + data_plane::storage_engines::types::ActivePhysicalPlan { + precompute_plan: initial_precompute_plan, + runtime_config: streaming_config.clone(), + backend_plan: Arc::new(initial_backend_plan), + storage_routing: Arc::new( + data_plane::storage_engines::types::BackendStorageRouting::empty(), + ), + }, ); + let hot_reload_config = + data_plane::storage_engines::types::HotReloadStreamingConfig::from_active( + active_physical_plan.clone(), + ); + let hot_reload_backend_plan = + data_plane::storage_engines::types::HotReloadBackendPlan::from_active( + active_physical_plan.clone(), + ); // Setup query engine. ASAPQueryEngine shares the same // HotReloadStreamingConfig handle as the HTTP server, so a POST @@ -730,6 +756,7 @@ async fn main() -> Result<()> { let mut server = HttpServer::new(http_config, engine, sketch_index.clone()) .with_hot_reload_config(hot_reload_config.clone()) .with_hot_reload_backend_plan(hot_reload_backend_plan.clone()) + .with_active_physical_plan(active_physical_plan.clone()) .with_probe_cache(probe_cache.clone()); // Per-metric storage-backend routing table (issue #46 @@ -772,7 +799,20 @@ async fn main() -> Result<()> { ); data_plane::storage_engines::types::BackendStorageRouting::empty() }; - server = server.with_backend_storage_routing(Arc::new(bootstrap_routing)); + if active_physical_plan.snapshot().backend_plan.plan_id == 0 { + let current = active_physical_plan.snapshot(); + active_physical_plan.swap(data_plane::storage_engines::types::ActivePhysicalPlan { + precompute_plan: current.precompute_plan.clone(), + runtime_config: current.runtime_config.clone(), + backend_plan: current.backend_plan.clone(), + storage_routing: Arc::new(bootstrap_routing), + }); + } + server = server.with_hot_reload_backend_storage_routing( + data_plane::query_engines::routing::HotReloadBackendStorageRouting::from_active( + active_physical_plan.clone(), + ), + ); // Phase-5/6 + Step-2.3: register the Thanos query engine on the // capability router. Path A2 is the only archive path now: when diff --git a/data_plane/src/query_engines/routing/backend_storage_routing.rs b/data_plane/src/query_engines/routing/backend_storage_routing.rs index 3e2674498..e72712ba1 100644 --- a/data_plane/src/query_engines/routing/backend_storage_routing.rs +++ b/data_plane/src/query_engines/routing/backend_storage_routing.rs @@ -944,6 +944,7 @@ pub struct HotReloadBackendStorageRouting { /// once and pick the tenant's `Arc`. inner: std::sync::Arc>>>, + active: Option, } impl HotReloadBackendStorageRouting { @@ -957,6 +958,7 @@ impl HotReloadBackendStorageRouting { map.insert(initial.tenant.clone(), std::sync::Arc::new(initial)); Self { inner: std::sync::Arc::new(arc_swap::ArcSwap::new(std::sync::Arc::new(map))), + active: None, } } @@ -979,6 +981,17 @@ impl HotReloadBackendStorageRouting { map.insert(initial.tenant.clone(), initial); Self { inner: std::sync::Arc::new(arc_swap::ArcSwap::new(std::sync::Arc::new(map))), + active: None, + } + } + + pub fn from_active(active: crate::storage_engines::types::HotReloadActivePhysicalPlan) -> Self { + let initial = active.snapshot().storage_routing.clone(); + let mut map = HashMap::new(); + map.insert(initial.tenant().to_string(), initial); + Self { + inner: std::sync::Arc::new(arc_swap::ArcSwap::new(std::sync::Arc::new(map))), + active: Some(active), } } @@ -998,6 +1011,12 @@ impl HotReloadBackendStorageRouting { /// table when even the default tenant is missing. Stable for the /// caller's lifetime; concurrent swaps don't invalidate it. pub fn snapshot_for_tenant(&self, tenant: &str) -> std::sync::Arc { + if let Some(active) = &self.active { + let routing = active.snapshot().storage_routing.clone(); + if routing.tenant() == tenant || tenant == DEFAULT_TENANT { + return routing; + } + } let map = self.inner.load_full(); if let Some(t) = map.get(tenant) { return t.clone(); 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 3a8c316d0..62322bb3d 100644 --- a/data_plane/src/storage_engines/types/hot_reload_config.rs +++ b/data_plane/src/storage_engines/types/hot_reload_config.rs @@ -81,6 +81,50 @@ use arc_swap::ArcSwap; use crate::storage_engines::types::StreamingConfig; +/// One immutable, generation-consistent runtime snapshot. Every execution +/// subsystem must project its view from the same `Arc`. +#[derive(Debug, Clone)] +pub struct ActivePhysicalPlan { + pub precompute_plan: control_plane::physical::compiler::PrecomputePlan, + pub runtime_config: Arc, + pub backend_plan: Arc, + pub storage_routing: Arc, +} + +#[derive(Clone)] +pub struct HotReloadActivePhysicalPlan { + inner: Arc>, +} + +impl HotReloadActivePhysicalPlan { + pub fn new(initial: ActivePhysicalPlan) -> Self { + Self { + inner: Arc::new(ArcSwap::new(Arc::new(initial))), + } + } + + pub fn snapshot(&self) -> Arc { + self.inner.load_full() + } + + pub fn swap(&self, next: ActivePhysicalPlan) -> Arc { + self.inner.swap(Arc::new(next)) + } +} + +impl std::fmt::Debug for HotReloadActivePhysicalPlan { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + let snapshot = self.snapshot(); + f.debug_struct("HotReloadActivePhysicalPlan") + .field("plan_id", &snapshot.backend_plan.plan_id) + .field( + "materializations", + &snapshot.precompute_plan.materializations.len(), + ) + .finish() + } +} + /// Hot-reloadable `BackendPlan` state — same `ArcSwap` shape as /// [`HotReloadStreamingConfig`], applied to /// `control_plane::backend_plan::BackendPlan` (see @@ -94,6 +138,7 @@ use crate::storage_engines::types::StreamingConfig; pub struct HotReloadBackendPlan { inner: Arc>, install_lock: Arc>, + active: Option, } impl HotReloadBackendPlan { @@ -101,6 +146,7 @@ impl HotReloadBackendPlan { Self { inner: Arc::new(ArcSwap::new(Arc::new(initial))), install_lock: Arc::new(std::sync::Mutex::new(())), + active: None, } } @@ -108,11 +154,24 @@ impl HotReloadBackendPlan { Self { inner: Arc::new(ArcSwap::new(initial)), install_lock: Arc::new(std::sync::Mutex::new(())), + active: None, + } + } + + pub fn from_active(active: HotReloadActivePhysicalPlan) -> Self { + let initial = active.snapshot().backend_plan.clone(); + Self { + inner: Arc::new(ArcSwap::new(initial)), + install_lock: Arc::new(std::sync::Mutex::new(())), + active: Some(active), } } pub fn snapshot(&self) -> Arc { - self.inner.load_full() + self.active + .as_ref() + .map(|a| a.snapshot().backend_plan.clone()) + .unwrap_or_else(|| self.inner.load_full()) } pub fn swap( @@ -238,6 +297,7 @@ mod hot_reload_backend_plan_tests { #[derive(Clone)] pub struct HotReloadStreamingConfig { inner: Arc>, + active: Option, } impl HotReloadStreamingConfig { @@ -247,6 +307,7 @@ impl HotReloadStreamingConfig { pub fn new(initial: StreamingConfig) -> Self { Self { inner: Arc::new(ArcSwap::new(Arc::new(initial))), + active: None, } } @@ -256,6 +317,15 @@ impl HotReloadStreamingConfig { pub fn from_arc(initial: Arc) -> Self { Self { inner: Arc::new(ArcSwap::new(initial)), + active: None, + } + } + + pub fn from_active(active: HotReloadActivePhysicalPlan) -> Self { + let initial = active.snapshot().runtime_config.clone(); + Self { + inner: Arc::new(ArcSwap::new(initial)), + active: Some(active), } } @@ -263,7 +333,10 @@ impl HotReloadStreamingConfig { /// returned `Arc` is stable for the caller's lifetime — a /// concurrent swap produces a new `Arc` and leaves this one alone. pub fn snapshot(&self) -> Arc { - self.inner.load_full() + self.active + .as_ref() + .map(|a| a.snapshot().runtime_config.clone()) + .unwrap_or_else(|| self.inner.load_full()) } /// Atomically replace the current config. The previous `Arc` is