diff --git a/control_plane/src/backend_plan/mod.rs b/control_plane/src/backend_plan/mod.rs index 507dbaf4..6ab9b05b 100644 --- a/control_plane/src/backend_plan/mod.rs +++ b/control_plane/src/backend_plan/mod.rs @@ -60,6 +60,18 @@ pub enum DecodeError { UnsupportedLifecycle(String), } +#[derive(Debug, Error, PartialEq, Eq)] +pub enum ValidationError { + #[error("materialization map key {key} does not match embedded fingerprint {embedded}")] + FingerprintMismatch { key: u64, embedded: u64 }, + #[error("materialization {fingerprint} has a zero-sized window")] + ZeroWindow { fingerprint: u64 }, + #[error("materialization {fingerprint} has a zero slide")] + ZeroSlide { fingerprint: u64 }, + #[error("route references unknown materialization {fingerprint}")] + UnknownMaterialization { fingerprint: u64 }, +} + // ── WindowSpec ─────────────────────────────────────────────────────────────── /// `kind`/`size`/`slide` triple — mirrors `QueryExpr::Window`'s fields @@ -782,6 +794,33 @@ impl BackendPlan { pub fn decode(bytes: &[u8]) -> Result { BackendPlan::try_from(proto::BackendPlan::decode(bytes)?) } + + /// Validate cross-references and invariants required before a decoded + /// plan may become visible to ingest or query readers. + pub fn validate(&self) -> Result<(), ValidationError> { + for (key, materialization) in &self.materializations { + if *key != materialization.fingerprint { + return Err(ValidationError::FingerprintMismatch { + key: key.0, + embedded: materialization.fingerprint.0, + }); + } + if materialization.window.size_ms == 0 { + return Err(ValidationError::ZeroWindow { fingerprint: key.0 }); + } + if materialization.window.slide_ms == Some(0) { + return Err(ValidationError::ZeroSlide { fingerprint: key.0 }); + } + } + for route in &self.routing { + if !self.materializations.contains_key(&route.materialization) { + return Err(ValidationError::UnknownMaterialization { + fingerprint: route.materialization.0, + }); + } + } + Ok(()) + } } #[cfg(test)] diff --git a/data_plane/src/drivers/query/servers/http.rs b/data_plane/src/drivers/query/servers/http.rs index 08ac102d..f947f943 100644 --- a/data_plane/src/drivers/query/servers/http.rs +++ b/data_plane/src/drivers/query/servers/http.rs @@ -5615,7 +5615,12 @@ async fn handle_post_backend_plan( let materialization_count = new_plan.materializations.len(); let routing_count = new_plan.routing.len(); let plan_id = new_plan.plan_id; - handle.swap(new_plan); + if let Err(error) = handle.install(new_plan) { + let body = serde_json::json!({ + "status": "error", + "error": format!("BackendPlan validation error: {error}")}); + return (StatusCode::UNPROCESSABLE_ENTITY, axum::Json(body)).into_response(); + } let body = serde_json::json!({ "status": "success", 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 66b56372..9f184bc4 100644 --- a/data_plane/src/storage_engines/types/hot_reload_config.rs +++ b/data_plane/src/storage_engines/types/hot_reload_config.rs @@ -118,6 +118,19 @@ impl HotReloadBackendPlan { ) -> Arc { self.inner.swap(Arc::new(new)) } + + /// Validate then atomically install a plan. A rejected plan never becomes + /// observable and the previous snapshot remains active. + pub fn install( + &self, + new: control_plane::backend_plan::BackendPlan, + ) -> Result< + Arc, + control_plane::backend_plan::ValidationError, + > { + new.validate()?; + Ok(self.swap(new)) + } } impl Default for HotReloadBackendPlan { @@ -170,6 +183,23 @@ mod hot_reload_backend_plan_tests { hr.swap(plan(2)); assert_eq!(hr_clone.snapshot().plan_id, 2); } + + #[test] + fn invalid_plan_is_rejected_without_replacing_snapshot() { + let hr = HotReloadBackendPlan::new(plan(1)); + let mut invalid = plan(2); + invalid + .routing + .push(control_plane::backend_plan::RoutingEntry { + satisfies: control_plane::physical::runtime_capability::Capability::ExactAgg( + asap_types::AggregationType::Sum, + ), + materialization: asap_types::PolicyFingerprint(99), + storage_backend: control_plane::backend_plan::StorageBackend::SketchStore, + }); + assert!(hr.install(invalid).is_err()); + assert_eq!(hr.snapshot().plan_id, 1); + } } /// Thin wrapper around `ArcSwap` with ergonomic