diff --git a/control_plane/proto/backend_plan.proto b/control_plane/proto/backend_plan.proto index 91e2a109b..5c56e4476 100644 --- a/control_plane/proto/backend_plan.proto +++ b/control_plane/proto/backend_plan.proto @@ -191,4 +191,8 @@ message BackendPlan { map materializations = 3; repeated RoutingEntry routing = 4; repeated MonitorSpec monitors = 5; + uint64 plan_version = 6; + uint64 activation_unix_ms = 7; + optional uint64 expiry_unix_ms = 8; + string backend_compat = 9; } diff --git a/control_plane/src/backend_client.rs b/control_plane/src/backend_client.rs index 028fef543..c0d05909b 100644 --- a/control_plane/src/backend_client.rs +++ b/control_plane/src/backend_client.rs @@ -391,6 +391,35 @@ impl BackendClient { Err(classify_http_status(status, body, "PhysicalPlan POST")) } } + + pub async fn activate_physical_plan( + &self, + plan_id: u64, + plan_version: u64, + ) -> std::result::Result<(), BackendPostError> { + let url = format!("{}/activate", derive_physical_plan_url(&self.endpoint)); + let response = self + .http + .post(&url) + .json(&serde_json::json!({ + "plan_id": plan_id, + "plan_version": plan_version, + })) + .send() + .await + .map_err(classify_reqwest_error)?; + let status = response.status(); + if status.is_success() { + Ok(()) + } else { + let body = response.text().await.unwrap_or_default(); + Err(classify_http_status( + status, + body, + "PhysicalPlan activation POST", + )) + } + } } fn derive_physical_plan_url(endpoint: &str) -> String { diff --git a/control_plane/src/backend_plan/from_stage_config.rs b/control_plane/src/backend_plan/from_stage_config.rs index 73237768d..02f4f81d1 100644 --- a/control_plane/src/backend_plan/from_stage_config.rs +++ b/control_plane/src/backend_plan/from_stage_config.rs @@ -119,6 +119,10 @@ pub fn from_stage_config( Ok(BackendPlan { plan_id, generated_at_unix_ms, + plan_version: 1, + activation_unix_ms: generated_at_unix_ms, + expiry_unix_ms: None, + backend_compat: super::BACKEND_COMPAT.into(), materializations, routing, monitors: plan_monitors, diff --git a/control_plane/src/backend_plan/mod.rs b/control_plane/src/backend_plan/mod.rs index ca546ee1c..a834920b9 100644 --- a/control_plane/src/backend_plan/mod.rs +++ b/control_plane/src/backend_plan/mod.rs @@ -39,6 +39,8 @@ use asap_types::{AggregationType, MonitorSpec, PolicyFingerprint}; use prost::Message as _; use thiserror::Error; +pub const BACKEND_COMPAT: &str = "asap-query-backend.v1"; + use crate::physical::runtime_capability::Capability; use asap_types::enums::WindowKind; use planner_types::post_asap::{ @@ -82,6 +84,27 @@ pub enum ValidationError { IncompatibleRoute { fingerprint: u64 }, #[error("stale plan generation: incoming={incoming}, active={active}")] StaleGeneration { incoming: u64, active: u64 }, + #[error("stale plan version for plan {plan_id}: incoming={incoming}, active={active}")] + StalePlanVersion { + plan_id: u64, + incoming: u64, + active: u64, + }, + #[error("plan {plan_id} version {plan_version} was reused with different content")] + ReusedPlanVersion { plan_id: u64, plan_version: u64 }, + #[error("non-bootstrap plan must have a non-zero plan version")] + ZeroPlanVersion, + #[error("non-bootstrap plan must have an activation time")] + MissingActivation, + #[error("plan expiry {expiry} is not after activation {activation}")] + InvalidExpiry { activation: u64, expiry: u64 }, + #[error("non-bootstrap plan must declare backend compatibility")] + MissingBackendCompat, + #[error("unsupported backend compatibility `{actual}`; expected `{expected}`")] + UnsupportedBackendCompat { + actual: String, + expected: &'static str, + }, } // ── WindowSpec ─────────────────────────────────────────────────────────────── @@ -772,6 +795,10 @@ pub struct BackendPlan { /// existing convention for the pre-`BackendPlan` wire format). pub plan_id: u64, pub generated_at_unix_ms: u64, + pub plan_version: u64, + pub activation_unix_ms: u64, + pub expiry_unix_ms: Option, + pub backend_compat: String, pub materializations: HashMap, pub routing: Vec, pub monitors: Vec, @@ -782,6 +809,10 @@ impl From<&BackendPlan> for proto::BackendPlan { proto::BackendPlan { plan_id: p.plan_id, generated_at_unix_ms: p.generated_at_unix_ms, + plan_version: p.plan_version, + activation_unix_ms: p.activation_unix_ms, + expiry_unix_ms: p.expiry_unix_ms, + backend_compat: p.backend_compat.clone(), materializations: p .materializations .iter() @@ -809,6 +840,10 @@ impl TryFrom for BackendPlan { Ok(BackendPlan { plan_id: p.plan_id, generated_at_unix_ms: p.generated_at_unix_ms, + plan_version: p.plan_version, + activation_unix_ms: p.activation_unix_ms, + expiry_unix_ms: p.expiry_unix_ms, + backend_compat: p.backend_compat, materializations, routing, monitors: p.monitors.into_iter().map(Into::into).collect(), @@ -832,6 +867,31 @@ impl BackendPlan { /// Validate cross-references and invariants required before a decoded /// plan may become visible to ingest or query readers. pub fn validate(&self) -> Result<(), ValidationError> { + if self.plan_id != 0 { + if self.plan_version == 0 { + return Err(ValidationError::ZeroPlanVersion); + } + if self.activation_unix_ms == 0 { + return Err(ValidationError::MissingActivation); + } + if self.backend_compat.trim().is_empty() { + return Err(ValidationError::MissingBackendCompat); + } + if self.backend_compat != BACKEND_COMPAT { + return Err(ValidationError::UnsupportedBackendCompat { + actual: self.backend_compat.clone(), + expected: BACKEND_COMPAT, + }); + } + if let Some(expiry) = self.expiry_unix_ms { + if expiry <= self.activation_unix_ms { + return Err(ValidationError::InvalidExpiry { + activation: self.activation_unix_ms, + expiry, + }); + } + } + } for (key, materialization) in &self.materializations { if *key != materialization.fingerprint { return Err(ValidationError::FingerprintMismatch { @@ -976,6 +1036,10 @@ mod tests { BackendPlan { plan_id: 42, generated_at_unix_ms: 1_735_000_000_000, + plan_version: 1, + activation_unix_ms: 1_735_000_000_000, + expiry_unix_ms: None, + backend_compat: "asap-query-backend.v1".into(), materializations, routing: vec![ RoutingEntry { diff --git a/control_plane/src/emit/backend_push.rs b/control_plane/src/emit/backend_push.rs index a824d0641..7167cdf9b 100644 --- a/control_plane/src/emit/backend_push.rs +++ b/control_plane/src/emit/backend_push.rs @@ -247,6 +247,9 @@ async fn push_documents_coupled( plan_id: crate::backend_plan::BackendPlan::decode(&plan_bytes) .map(|plan| plan.plan_id) .unwrap_or_default(), + plan_version: crate::backend_plan::BackendPlan::decode(&plan_bytes) + .map(|plan| plan.plan_version) + .unwrap_or_default(), entries: Default::default(), }; @@ -454,7 +457,11 @@ async fn push_cumulative_entries( let precompute_plan = PrecomputePlan { envelope: PlanEnvelope { plan_id, + plan_version: 1, generated_at_unix_ms, + activation_unix_ms: generated_at_unix_ms, + expiry_unix_ms: None, + backend_compat: "asap-query-backend.v1".into(), planner_revision: crate::physical::compiler::PLANNER_REVISION.into(), capability_snapshot_id: "replanner".into(), }, diff --git a/control_plane/src/main.rs b/control_plane/src/main.rs index 0157126f0..602826745 100644 --- a/control_plane/src/main.rs +++ b/control_plane/src/main.rs @@ -595,6 +595,10 @@ struct CompileAndPublishPhysicalPlanRequest { evidence: HashMap, planner_revision: String, max_evidence_age_ms: u64, + plan_version: u64, + activation_unix_ms: u64, + expiry_unix_ms: Option, + backend_compat: String, #[serde(default = "default_physical_plan_timeout_ms")] apply_timeout_ms: u64, } @@ -606,6 +610,8 @@ fn default_physical_plan_timeout_ms() -> u64 { #[derive(Debug, Serialize)] struct CompileAndPublishPhysicalPlanResponse { plan_id: u64, + plan_version: u64, + status: &'static str, generated_at_unix_ms: u64, collector_ids: Vec, } @@ -667,9 +673,36 @@ async fn handle_compile_and_publish_physical_plan( ) .into_response(); } + let now = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap_or_default() + .as_millis() as u64; + let activation_wait = bundle.envelope.activation_unix_ms.saturating_sub(now); + if activation_wait > apply_timeout.as_millis() as u64 { + return ( + StatusCode::GATEWAY_TIMEOUT, + "activation time exceeds apply_timeout_ms; backend remains staged".to_string(), + ) + .into_response(); + } + if activation_wait > 0 { + tokio::time::sleep(Duration::from_millis(activation_wait)).await; + } + if let Err(error) = backend + .activate_physical_plan(bundle.envelope.plan_id, bundle.envelope.plan_version) + .await + { + return ( + StatusCode::BAD_GATEWAY, + format!("backend physical-plan activation failed: {error}"), + ) + .into_response(); + } Json(CompileAndPublishPhysicalPlanResponse { plan_id: bundle.envelope.plan_id, + plan_version: bundle.envelope.plan_version, + status: "active", generated_at_unix_ms: bundle.envelope.generated_at_unix_ms, collector_ids, }) @@ -693,6 +726,18 @@ fn compile_physical_plan_request( "max_evidence_age_ms and apply_timeout_ms must be non-zero".to_string(), )); } + if request.plan_version == 0 + || request.activation_unix_ms == 0 + || request.backend_compat != control_plane::backend_plan::BACKEND_COMPAT + || request + .expiry_unix_ms + .is_some_and(|expiry| expiry <= request.activation_unix_ms) + { + return Err(( + StatusCode::UNPROCESSABLE_ENTITY, + "plan_version, activation_unix_ms, backend_compat and expiry are invalid".to_string(), + )); + } let now = std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) @@ -748,6 +793,10 @@ fn compile_physical_plan_request( capability_snapshot_id: request.capability_snapshot_id, observed_at_unix_ms: now, max_evidence_age_ms: request.max_evidence_age_ms, + plan_version: request.plan_version, + activation_unix_ms: request.activation_unix_ms, + expiry_unix_ms: request.expiry_unix_ms, + backend_compat: request.backend_compat, }, ) { Ok(bundle) => bundle, diff --git a/control_plane/src/opamp/mod.rs b/control_plane/src/opamp/mod.rs index 954341ab0..4c08647f6 100644 --- a/control_plane/src/opamp/mod.rs +++ b/control_plane/src/opamp/mod.rs @@ -62,6 +62,7 @@ pub enum CollectorPlanStatusKind { #[serde(deny_unknown_fields)] pub struct CollectorPlanStatus { pub plan_id: u64, + pub plan_version: u64, pub status: CollectorPlanStatusKind, #[serde(default)] pub error: Option, @@ -81,12 +82,19 @@ pub enum CollectorPlanPublishError { #[source] source: serde_json::Error, }, - #[error("collector {collector_id} did not report plan {plan_id} before timeout")] - StatusTimeout { collector_id: String, plan_id: u64 }, - #[error("collector {collector_id} rejected plan {plan_id}: {error}")] + #[error( + "collector {collector_id} did not report plan {plan_id}/{plan_version} before timeout" + )] + StatusTimeout { + collector_id: String, + plan_id: u64, + plan_version: u64, + }, + #[error("collector {collector_id} rejected plan {plan_id}/{plan_version}: {error}")] Rejected { collector_id: String, plan_id: u64, + plan_version: u64, error: String, }, } @@ -137,7 +145,7 @@ struct AgentConnection { } type AgentMap = HashMap; -type PlanStatusMap = HashMap<(String, u64), CollectorPlanStatus>; +type PlanStatusMap = HashMap<(String, u64, u64), CollectorPlanStatus>; pub type OnConnectFn = Arc; pub type OnDisconnectFn = Arc; @@ -291,10 +299,11 @@ impl OpampServer { source, } })?; - self.plan_statuses - .write() - .await - .remove(&(plan.collector_id.clone(), plan.envelope.plan_id)); + self.plan_statuses.write().await.remove(&( + plan.collector_id.clone(), + plan.envelope.plan_id, + plan.envelope.plan_version, + )); let sent = self.send_collector_plan(&plan.collector_id, body).await; if !sent { return Err(CollectorPlanPublishError::Disconnected { @@ -306,16 +315,23 @@ impl OpampServer { let mut reports = Vec::with_capacity(plans.len()); for plan in plans { let report = self - .wait_for_plan_status(&plan.collector_id, plan.envelope.plan_id, timeout) + .wait_for_plan_status( + &plan.collector_id, + plan.envelope.plan_id, + plan.envelope.plan_version, + timeout, + ) .await .ok_or_else(|| CollectorPlanPublishError::StatusTimeout { collector_id: plan.collector_id.clone(), plan_id: plan.envelope.plan_id, + plan_version: plan.envelope.plan_version, })?; if report.status != CollectorPlanStatusKind::Applied { return Err(CollectorPlanPublishError::Rejected { collector_id: plan.collector_id.clone(), plan_id: report.plan_id, + plan_version: report.plan_version, error: report.error.clone().unwrap_or_else(|| "unspecified".into()), }); } @@ -399,9 +415,10 @@ impl OpampServer { &self, agent_id: &str, plan_id: u64, + plan_version: u64, timeout: Duration, ) -> Option { - let key = (agent_id.to_string(), plan_id); + let key = (agent_id.to_string(), plan_id, plan_version); let deadline = tokio::time::Instant::now() + timeout; loop { let changed = self.state_changed.notified(); @@ -526,10 +543,10 @@ async fn handle_socket( { match serde_json::from_slice::(&message.data) { Ok(status) => { - srv.plan_statuses - .write() - .await - .insert((agent_id.clone(), status.plan_id), status); + srv.plan_statuses.write().await.insert( + (agent_id.clone(), status.plan_id, status.plan_version), + status, + ); srv.state_changed.notify_waiters(); } Err(error) => warn!( @@ -989,7 +1006,11 @@ mod tests { collector_id: collector_id.into(), envelope: crate::physical::compiler::PlanEnvelope { plan_id, + plan_version: 1, generated_at_unix_ms: 1, + activation_unix_ms: 1, + expiry_unix_ms: None, + backend_compat: "asap-query-backend.v1".into(), planner_revision: crate::physical::compiler::PLANNER_REVISION.into(), capability_snapshot_id: "caps-1".into(), }, @@ -1072,6 +1093,7 @@ mod tests { let status = serde_json::to_vec(&CollectorPlanStatus { plan_id: 42, + plan_version: 1, status: CollectorPlanStatusKind::Applied, error: None, }) @@ -1092,6 +1114,7 @@ mod tests { let reports = publish.await.unwrap().unwrap(); assert_eq!(reports.len(), 1); assert_eq!(reports[0].plan_id, 42); + assert_eq!(reports[0].plan_version, 1); assert_eq!(reports[0].status, CollectorPlanStatusKind::Applied); } diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index 676b1bfba..ff5217bb6 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -138,12 +138,20 @@ pub struct DeploymentEnvironment { pub capability_snapshot_id: String, pub observed_at_unix_ms: u64, pub max_evidence_age_ms: u64, + pub plan_version: u64, + pub activation_unix_ms: u64, + pub expiry_unix_ms: Option, + pub backend_compat: String, } #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] pub struct PlanEnvelope { pub plan_id: u64, + pub plan_version: u64, pub generated_at_unix_ms: u64, + pub activation_unix_ms: u64, + pub expiry_unix_ms: Option, + pub backend_compat: String, pub planner_revision: String, pub capability_snapshot_id: String, } @@ -349,7 +357,11 @@ impl PhysicalCompiler { let plan_id = stable_plan_id(&collector_materializations); let envelope = PlanEnvelope { plan_id, + plan_version: environment.plan_version, generated_at_unix_ms: environment.observed_at_unix_ms, + activation_unix_ms: environment.activation_unix_ms, + expiry_unix_ms: environment.expiry_unix_ms, + backend_compat: environment.backend_compat.clone(), planner_revision: PLANNER_REVISION.into(), capability_snapshot_id: environment.capability_snapshot_id, }; @@ -363,6 +375,16 @@ impl PhysicalCompiler { plan_id, environment.observed_at_unix_ms, )?; + backend_plan.plan_version = envelope.plan_version; + backend_plan.activation_unix_ms = envelope.activation_unix_ms; + backend_plan.expiry_unix_ms = envelope.expiry_unix_ms; + backend_plan.backend_compat = envelope.backend_compat.clone(); + backend_plan + .validate() + .map_err(|error| CompileError::Query { + query_id: "physical-plan-envelope".into(), + reason: error.to_string(), + })?; for materialization in backend_plan.materializations.values_mut() { materialization.lifecycle = Some(SummaryMaintenanceLifecycleGuarantee { summary_maintenance_lifecycle: SummaryMaintenanceLifecycle::ContinuouslyMaintained, @@ -478,6 +500,7 @@ impl PhysicalCompiler { } let query_plan = QueryPlan { plan_id, + plan_version: envelope.plan_version, entries: query_entries, }; query_plan.validate(&materialization_fingerprints)?; @@ -843,6 +866,10 @@ mod tests { capability_snapshot_id: "caps-7".into(), observed_at_unix_ms: now, max_evidence_age_ms: 60_000, + plan_version: 1, + activation_unix_ms: now, + expiry_unix_ms: None, + backend_compat: "asap-query-backend.v1".into(), } } diff --git a/control_plane/src/query_plan.rs b/control_plane/src/query_plan.rs index c80d07fad..f50078dba 100644 --- a/control_plane/src/query_plan.rs +++ b/control_plane/src/query_plan.rs @@ -19,6 +19,7 @@ use asap_types::PolicyFingerprint; #[serde(deny_unknown_fields)] pub struct QueryPlan { pub plan_id: u64, + pub plan_version: u64, pub entries: BTreeMap, } @@ -26,6 +27,7 @@ impl QueryPlan { pub fn empty() -> Self { Self { plan_id: 0, + plan_version: 0, entries: BTreeMap::new(), } } @@ -38,6 +40,11 @@ impl QueryPlan { } pub fn validate(&self, available: &BTreeSet) -> Result<(), QueryPlanError> { + if self.plan_id != 0 && self.plan_version == 0 { + return Err(QueryPlanError::Invalid( + "non-bootstrap QueryPlan has zero plan_version".into(), + )); + } for (identity, entry) in &self.entries { if identity != &entry.canonical_promql { return Err(QueryPlanError::Invalid(format!( diff --git a/data_plane/src/drivers/query/servers/http.rs b/data_plane/src/drivers/query/servers/http.rs index 48a35403f..66c896e08 100644 --- a/data_plane/src/drivers/query/servers/http.rs +++ b/data_plane/src/drivers/query/servers/http.rs @@ -10,7 +10,7 @@ use axum::{ use serde_json::Value; use std::collections::HashMap; use std::sync::Arc; -use std::time::Instant; +use std::time::{Duration, Instant}; use tokio::net::TcpListener; use tracing::{debug, info, warn}; @@ -150,6 +150,7 @@ pub struct HttpServer { /// plane generations cannot interleave their config and BackendPlan. physical_plan_lock: Arc>, active_physical_plan: Option, + physical_plan_lifecycle: Option, } #[derive(Clone)] @@ -177,6 +178,7 @@ struct AppState { probe_cache: Option>, physical_plan_lock: Arc>, active_physical_plan: Option, + physical_plan_lifecycle: Option, } impl HttpServer { @@ -203,6 +205,7 @@ impl HttpServer { probe_cache: None, physical_plan_lock: Arc::new(tokio::sync::Mutex::new(())), active_physical_plan: None, + physical_plan_lifecycle: None, } } @@ -264,6 +267,9 @@ impl HttpServer { mut self, handle: crate::storage_engines::types::HotReloadActivePhysicalPlan, ) -> Self { + self.physical_plan_lifecycle = Some( + crate::storage_engines::types::PhysicalPlanLifecycle::new(handle.clone()), + ); self.active_physical_plan = Some(handle); self } @@ -376,6 +382,7 @@ impl HttpServer { probe_cache: self.probe_cache.clone(), physical_plan_lock: self.physical_plan_lock.clone(), active_physical_plan: self.active_physical_plan.clone(), + physical_plan_lifecycle: self.physical_plan_lifecycle.clone(), }; let range_query_endpoint = adapter.get_range_query_endpoint(); @@ -402,6 +409,14 @@ impl HttpServer { get(handle_get_backend_plan).post(handle_post_backend_plan), ) .route("/api/v1/physical-plan", post(handle_post_physical_plan)) + .route( + "/api/v1/physical-plan/activate", + post(handle_activate_physical_plan), + ) + .route( + "/api/v1/physical-plan/status", + get(handle_physical_plan_status), + ) // Phase α (MVP): control-plane-pushed `BackendStorageRouting` // table. POST replaces the current table atomically; GET // returns a JSON snapshot for operator diagnostics. @@ -464,6 +479,7 @@ impl HttpServer { probe_cache: self.probe_cache.clone(), physical_plan_lock: self.physical_plan_lock.clone(), active_physical_plan: self.active_physical_plan.clone(), + physical_plan_lifecycle: self.physical_plan_lifecycle.clone(), }; let range_query_endpoint = adapter.get_range_query_endpoint(); @@ -486,6 +502,14 @@ impl HttpServer { get(handle_get_backend_plan).post(handle_post_backend_plan), ) .route("/api/v1/physical-plan", post(handle_post_physical_plan)) + .route( + "/api/v1/physical-plan/activate", + post(handle_activate_physical_plan), + ) + .route( + "/api/v1/physical-plan/status", + get(handle_physical_plan_status), + ) // Phase α (MVP): control-plane-pushed `BackendStorageRouting` // table. POST replaces the current table atomically; GET // returns a JSON snapshot for operator diagnostics. @@ -2428,6 +2452,9 @@ aggregations: let new_plan = BackendPlan { plan_id: 7, generated_at_unix_ms: 123, + plan_version: 1, + activation_unix_ms: 123, + backend_compat: "asap-query-backend.v1".into(), ..Default::default() }; let bytes = new_plan.encode_to_vec(); @@ -5275,10 +5302,8 @@ struct PhysicalPlanInstallRequest { storage_routing: Option, } -/// Install the two backend views of one PhysicalPlan as one validated -/// publication. The BackendPlan is installed first, so readers racing the -/// short swap interval fail closed on missing fingerprints rather than using -/// a new streaming configuration with stale routing authority. +/// Validate and stage all backend views. Staging never changes query routing; +/// a separate activation request performs the single snapshot swap. async fn handle_post_physical_plan( State(state): State, axum::Json(request): axum::Json, @@ -5296,6 +5321,15 @@ async fn handle_post_physical_plan( ) .into_response(); }; + let Some(lifecycle) = state.physical_plan_lifecycle.as_ref() else { + return ( + StatusCode::SERVICE_UNAVAILABLE, + axum::Json(serde_json::json!({ + "status": "error", "error": "physical-plan lifecycle is not attached" + })), + ) + .into_response(); + }; let new_config = crate::storage_engines::types::StreamingConfig::new( request .precompute_plan @@ -5326,11 +5360,18 @@ async fn handle_post_physical_plan( ) .into_response(); } - if request.query_plan.plan_id != new_plan.plan_id { + if request.query_plan.plan_id != new_plan.plan_id + || request.query_plan.plan_version != new_plan.plan_version + || request.precompute_plan.envelope.plan_id != new_plan.plan_id + || request.precompute_plan.envelope.plan_version != new_plan.plan_version + || request.precompute_plan.envelope.activation_unix_ms != new_plan.activation_unix_ms + || request.precompute_plan.envelope.expiry_unix_ms != new_plan.expiry_unix_ms + || request.precompute_plan.envelope.backend_compat != new_plan.backend_compat + { return ( StatusCode::UNPROCESSABLE_ENTITY, axum::Json(serde_json::json!({ - "status": "error", "error": "QueryPlan and BackendPlan plan_id differ" + "status": "error", "error": "physical subplans have different plan identity/version" })), ) .into_response(); @@ -5379,20 +5420,78 @@ async fn handle_post_physical_plan( query_plan: Arc::new(request.query_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 { + let plan_version = active.backend_plan.plan_version; + let now = unix_time_ms(); + if let Err(error) = lifecycle.stage(active, now) { return ( StatusCode::CONFLICT, + axum::Json(serde_json::json!({"status": "error", "error": error.to_string()})), + ) + .into_response(); + } + ( + StatusCode::ACCEPTED, + axum::Json(serde_json::json!({ + "status": "staged", "plan_id": plan_id, "plan_version": plan_version, + "materialization_count": plan_fps.len() + })), + ) + .into_response() +} + +#[derive(serde::Deserialize)] +#[serde(deny_unknown_fields)] +struct ActivatePhysicalPlanRequest { + plan_id: u64, + plan_version: u64, +} + +async fn handle_activate_physical_plan( + State(state): State, + axum::Json(request): axum::Json, +) -> axum::response::Response { + use axum::response::IntoResponse; + let (Some(lifecycle), Some(active_handle)) = ( + state.physical_plan_lifecycle.as_ref(), + state.active_physical_plan.as_ref(), + ) else { + return ( + StatusCode::SERVICE_UNAVAILABLE, axum::Json(serde_json::json!({ - "status": "error", "error": "stale physical-plan generation" + "status": "error", "error": "physical-plan lifecycle is not attached" })), ) .into_response(); + }; + let _guard = state.physical_plan_lock.lock().await; + let old = match lifecycle.activate(request.plan_id, request.plan_version, unix_time_ms()) { + Ok(old) => old, + Err(error) => { + return ( + StatusCode::CONFLICT, + axum::Json(serde_json::json!({ + "status": "error", "error": error.to_string() + })), + ) + .into_response() + } + }; + if old.backend_plan.plan_id != 0 { + let draining_id = old.backend_plan.plan_id; + let draining_version = old.backend_plan.plan_version; + let lifecycle = lifecycle.clone(); + tokio::spawn(async move { + loop { + if Arc::strong_count(&old) == 1 { + lifecycle.retire_drained(draining_id, draining_version); + break; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }); } - 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( + let retired = crate::storage_engines::sketch_db::lifecycle::reconcile_from_streaming_config( state.sketch_index.as_ref(), snap.as_ref(), crate::storage_engines::sketch_db::DEFAULT_RETIREMENT_RETENTION, @@ -5400,13 +5499,40 @@ async fn handle_post_physical_plan( ( StatusCode::OK, axum::Json(serde_json::json!({ - "status": "success", "plan_id": plan_id, - "materialization_count": plan_fps.len(), "sids_retired": sid_summary.retired + "status": "active", "plan_id": request.plan_id, "plan_version": request.plan_version, + "sids_retired": retired.retired + })), + ) + .into_response() +} + +async fn handle_physical_plan_status(State(state): State) -> axum::response::Response { + use axum::response::IntoResponse; + let Some(lifecycle) = state.physical_plan_lifecycle.as_ref() else { + return ( + StatusCode::SERVICE_UNAVAILABLE, + axum::Json(serde_json::json!({ + "status": "error", "error": "physical-plan lifecycle is not attached" + })), + ) + .into_response(); + }; + ( + StatusCode::OK, + axum::Json(serde_json::json!({ + "status": "success", "plans": lifecycle.statuses() })), ) .into_response() } +fn unix_time_ms() -> u64 { + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap_or_default() + .as_millis() as u64 +} + // ── Phase α: BackendStorageRouting hot-reload endpoints ──────────── /// `GET /api/v1/storage_routing` — return a JSON snapshot of the diff --git a/data_plane/src/main.rs b/data_plane/src/main.rs index 5951742f8..554c2a071 100644 --- a/data_plane/src/main.rs +++ b/data_plane/src/main.rs @@ -427,7 +427,11 @@ async fn main() -> Result<()> { let initial_precompute_plan = control_plane::physical::compiler::PrecomputePlan { envelope: control_plane::physical::compiler::PlanEnvelope { plan_id: 0, + plan_version: 0, generated_at_unix_ms: 0, + activation_unix_ms: 0, + expiry_unix_ms: None, + backend_compat: "bootstrap".into(), planner_revision: control_plane::physical::compiler::PLANNER_REVISION.into(), capability_snapshot_id: "bootstrap".into(), }, 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 0fb403b32..68a4877b0 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 @@ -964,6 +964,10 @@ mod tests { BackendPlan { plan_id: 1, generated_at_unix_ms: 0, + plan_version: 1, + activation_unix_ms: 1, + expiry_unix_ms: None, + backend_compat: "asap-query-backend.v1".into(), materializations, routing: vec![RoutingEntry { satisfies: Capability::QuantileApprox(Some(SketchAlgorithm::DDSketch)), 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 8d2f60f43..abc4065bd 100644 --- a/data_plane/src/storage_engines/types/hot_reload_config.rs +++ b/data_plane/src/storage_engines/types/hot_reload_config.rs @@ -75,6 +75,7 @@ //! produce confusing store states where data under the same agg_id //! spans multiple parameter generations. +use std::collections::BTreeMap; use std::sync::Arc; use arc_swap::ArcSwap; @@ -97,6 +98,209 @@ pub struct HotReloadActivePhysicalPlan { inner: Arc>, } +#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)] +#[serde(rename_all = "snake_case")] +pub enum PhysicalPlanPhase { + Staged, + Active, + Draining, + Retired, +} + +#[derive(Debug, Clone, serde::Serialize)] +pub struct PhysicalPlanStatus { + pub plan_id: u64, + pub plan_version: u64, + pub phase: PhysicalPlanPhase, + pub activation_unix_ms: u64, + pub expiry_unix_ms: Option, +} + +#[derive(Debug, thiserror::Error, PartialEq, Eq)] +pub enum PhysicalPlanLifecycleError { + #[error("plan {plan_id}/{plan_version} is already staged or active")] + Duplicate { plan_id: u64, plan_version: u64 }, + #[error("plan {plan_id}/{plan_version} is not staged")] + NotStaged { plan_id: u64, plan_version: u64 }, + #[error("plan activation {activation} is later than now {now}")] + ActivationNotReached { activation: u64, now: u64 }, + #[error("plan expired at {expiry}; now is {now}")] + Expired { expiry: u64, now: u64 }, + #[error( + "plan version {incoming} is not newer than active version {active} for plan {plan_id}" + )] + StaleVersion { + plan_id: u64, + incoming: u64, + active: u64, + }, +} + +#[derive(Clone)] +pub struct PhysicalPlanLifecycle { + active: HotReloadActivePhysicalPlan, + state: Arc>, +} + +struct PhysicalPlanLifecycleState { + staged: BTreeMap<(u64, u64), ActivePhysicalPlan>, + statuses: BTreeMap<(u64, u64), PhysicalPlanStatus>, +} + +impl PhysicalPlanLifecycle { + pub fn new(active: HotReloadActivePhysicalPlan) -> Self { + let snapshot = active.snapshot(); + let mut statuses = BTreeMap::new(); + if snapshot.backend_plan.plan_id != 0 { + statuses.insert( + ( + snapshot.backend_plan.plan_id, + snapshot.backend_plan.plan_version, + ), + status_for(&snapshot, PhysicalPlanPhase::Active), + ); + } + Self { + active, + state: Arc::new(std::sync::Mutex::new(PhysicalPlanLifecycleState { + staged: BTreeMap::new(), + statuses, + })), + } + } + + pub fn stage( + &self, + plan: ActivePhysicalPlan, + now: u64, + ) -> Result<(), PhysicalPlanLifecycleError> { + let key = (plan.backend_plan.plan_id, plan.backend_plan.plan_version); + if let Some(expiry) = plan.backend_plan.expiry_unix_ms { + if expiry <= now { + return Err(PhysicalPlanLifecycleError::Expired { expiry, now }); + } + } + let active = self.active.snapshot(); + if active.backend_plan.plan_id != 0 && key.1 <= active.backend_plan.plan_version { + return Err(PhysicalPlanLifecycleError::StaleVersion { + plan_id: key.0, + incoming: key.1, + active: active.backend_plan.plan_version, + }); + } + let mut state = self + .state + .lock() + .expect("physical-plan lifecycle lock poisoned"); + if state.staged.contains_key(&key) + || state.statuses.values().any(|status| { + status.plan_version == key.1 && !matches!(status.phase, PhysicalPlanPhase::Retired) + }) + { + return Err(PhysicalPlanLifecycleError::Duplicate { + plan_id: key.0, + plan_version: key.1, + }); + } + state + .statuses + .insert(key, status_for(&plan, PhysicalPlanPhase::Staged)); + state.staged.insert(key, plan); + Ok(()) + } + + pub fn activate( + &self, + plan_id: u64, + plan_version: u64, + now: u64, + ) -> Result, PhysicalPlanLifecycleError> { + let key = (plan_id, plan_version); + let mut state = self + .state + .lock() + .expect("physical-plan lifecycle lock poisoned"); + let plan = state + .staged + .get(&key) + .ok_or(PhysicalPlanLifecycleError::NotStaged { + plan_id, + plan_version, + })?; + if now < plan.backend_plan.activation_unix_ms { + return Err(PhysicalPlanLifecycleError::ActivationNotReached { + activation: plan.backend_plan.activation_unix_ms, + now, + }); + } + if let Some(expiry) = plan.backend_plan.expiry_unix_ms { + if now >= expiry { + return Err(PhysicalPlanLifecycleError::Expired { expiry, now }); + } + } + let current = self.active.snapshot(); + if current.backend_plan.plan_id != 0 + && plan.backend_plan.plan_version <= current.backend_plan.plan_version + { + return Err(PhysicalPlanLifecycleError::StaleVersion { + plan_id, + incoming: plan_version, + active: current.backend_plan.plan_version, + }); + } + let plan = state.staged.remove(&key).expect("staged plan disappeared"); + let old = self.active.swap(plan.clone()); + if old.backend_plan.plan_id != 0 { + if let Some(status) = state + .statuses + .get_mut(&(old.backend_plan.plan_id, old.backend_plan.plan_version)) + { + status.phase = PhysicalPlanPhase::Draining; + } + } + state + .statuses + .insert(key, status_for(&plan, PhysicalPlanPhase::Active)); + Ok(old) + } + + pub fn statuses(&self) -> Vec { + self.state + .lock() + .expect("physical-plan lifecycle lock poisoned") + .statuses + .values() + .cloned() + .collect() + } + + /// Mark a superseded generation retired after all readers of the old + /// immutable snapshot have drained. + pub fn retire_drained(&self, plan_id: u64, plan_version: u64) { + if let Some(status) = self + .state + .lock() + .expect("physical-plan lifecycle lock poisoned") + .statuses + .get_mut(&(plan_id, plan_version)) + { + if status.phase == PhysicalPlanPhase::Draining { + status.phase = PhysicalPlanPhase::Retired; + } + } + } +} + +fn status_for(plan: &ActivePhysicalPlan, phase: PhysicalPlanPhase) -> PhysicalPlanStatus { + PhysicalPlanStatus { + plan_id: plan.backend_plan.plan_id, + plan_version: plan.backend_plan.plan_version, + phase, + activation_unix_ms: plan.backend_plan.activation_unix_ms, + expiry_unix_ms: plan.backend_plan.expiry_unix_ms, + } +} + impl HotReloadActivePhysicalPlan { pub fn new(initial: ActivePhysicalPlan) -> Self { Self { @@ -198,6 +402,28 @@ impl HotReloadBackendPlan { .expect("backend-plan install lock poisoned"); new.validate()?; let active = self.snapshot(); + if new.plan_id == active.plan_id && new.plan_id != 0 { + if new.plan_version < active.plan_version { + return Err( + control_plane::backend_plan::ValidationError::StalePlanVersion { + plan_id: new.plan_id, + incoming: new.plan_version, + active: active.plan_version, + }, + ); + } + if new.plan_version == active.plan_version { + if new == *active { + return Ok(active); + } + return Err( + control_plane::backend_plan::ValidationError::ReusedPlanVersion { + plan_id: new.plan_id, + plan_version: new.plan_version, + }, + ); + } + } if new.generated_at_unix_ms < active.generated_at_unix_ms { return Err( control_plane::backend_plan::ValidationError::StaleGeneration { @@ -235,6 +461,9 @@ mod hot_reload_backend_plan_tests { fn plan(plan_id: u64) -> BackendPlan { BackendPlan { plan_id, + plan_version: 1, + activation_unix_ms: 1, + backend_compat: "asap-query-backend.v1".into(), ..Default::default() } } @@ -370,6 +599,46 @@ mod tests { use std::collections::HashMap; use std::thread; + fn physical_plan( + plan_id: u64, + plan_version: u64, + activation_unix_ms: u64, + expiry_unix_ms: Option, + ) -> ActivePhysicalPlan { + let envelope = control_plane::physical::compiler::PlanEnvelope { + plan_id, + plan_version, + generated_at_unix_ms: activation_unix_ms, + activation_unix_ms, + expiry_unix_ms, + backend_compat: "asap-query-backend.v1".into(), + planner_revision: control_plane::physical::compiler::PLANNER_REVISION.into(), + capability_snapshot_id: "test".into(), + }; + ActivePhysicalPlan { + precompute_plan: control_plane::physical::compiler::PrecomputePlan { + envelope, + materializations: Vec::new(), + }, + runtime_config: Arc::new(StreamingConfig::new(HashMap::new())), + backend_plan: Arc::new(control_plane::backend_plan::BackendPlan { + plan_id, + plan_version, + generated_at_unix_ms: activation_unix_ms, + activation_unix_ms, + expiry_unix_ms, + backend_compat: "asap-query-backend.v1".into(), + ..Default::default() + }), + query_plan: Arc::new(control_plane::query_plan::QueryPlan { + plan_id, + plan_version, + entries: Default::default(), + }), + storage_routing: Arc::new(crate::storage_engines::types::BackendStorageRouting::empty()), + } + } + fn dummy_agg(id: u64) -> AggregationConfig { AggregationConfig::new( AggregationType::Sum, @@ -471,4 +740,66 @@ mod tests { writer.join().unwrap(); reader.join().unwrap(); } + + #[test] + fn physical_plan_stages_activates_drains_and_retires() { + let active = HotReloadActivePhysicalPlan::new(physical_plan(7, 1, 100, None)); + let lifecycle = PhysicalPlanLifecycle::new(active.clone()); + lifecycle + .stage(physical_plan(7, 2, 200, Some(500)), 150) + .unwrap(); + + assert_eq!(active.snapshot().backend_plan.plan_version, 1); + assert!(matches!( + lifecycle.activate(7, 2, 199), + Err(PhysicalPlanLifecycleError::ActivationNotReached { .. }) + )); + let old = lifecycle.activate(7, 2, 200).unwrap(); + assert_eq!(old.backend_plan.plan_version, 1); + assert_eq!(active.snapshot().backend_plan.plan_version, 2); + + let statuses = lifecycle.statuses(); + assert_eq!(statuses.len(), 2); + assert_eq!(statuses[0].phase, PhysicalPlanPhase::Draining); + assert_eq!(statuses[1].phase, PhysicalPlanPhase::Active); + lifecycle.retire_drained(7, 1); + assert_eq!(lifecycle.statuses()[0].phase, PhysicalPlanPhase::Retired); + } + + #[test] + fn physical_plan_rejects_stale_and_expired_generations() { + let active = HotReloadActivePhysicalPlan::new(physical_plan(7, 2, 100, None)); + let lifecycle = PhysicalPlanLifecycle::new(active); + assert!(matches!( + lifecycle.stage(physical_plan(7, 1, 100, None), 150), + Err(PhysicalPlanLifecycleError::StaleVersion { .. }) + )); + assert!(matches!( + lifecycle.stage(physical_plan(8, 1, 100, Some(150)), 150), + Err(PhysicalPlanLifecycleError::Expired { .. }) + )); + } + + #[test] + fn activation_rechecks_version_and_cannot_downgrade_across_plan_ids() { + let active = HotReloadActivePhysicalPlan::new(physical_plan(7, 1, 100, None)); + let lifecycle = PhysicalPlanLifecycle::new(active.clone()); + lifecycle + .stage(physical_plan(8, 3, 100, None), 100) + .unwrap(); + lifecycle + .stage(physical_plan(9, 2, 100, None), 100) + .unwrap(); + lifecycle.activate(8, 3, 100).unwrap(); + + assert!(matches!( + lifecycle.activate(9, 2, 100), + Err(PhysicalPlanLifecycleError::StaleVersion { + incoming: 2, + active: 3, + .. + }) + )); + assert_eq!(active.snapshot().backend_plan.plan_version, 3); + } } diff --git a/data_plane/tests/backend_process_e2e.rs b/data_plane/tests/backend_process_e2e.rs index 23c98eb3d..a0beba656 100644 --- a/data_plane/tests/backend_process_e2e.rs +++ b/data_plane/tests/backend_process_e2e.rs @@ -210,9 +210,13 @@ async fn apply_next_collector_plan(address: String) -> serde_json::Value { let plan_id = plan["envelope"]["plan_id"] .as_u64() .expect("collector plan ID"); + let plan_version = plan["envelope"]["plan_version"] + .as_u64() + .expect("collector plan version"); let status = serde_json::to_vec(&CollectorPlanStatus { plan_id, + plan_version, status: CollectorPlanStatusKind::Applied, error: None, })