From 4c5b1496210b8b9d4cf38b0090da4ff5b0ffb379 Mon Sep 17 00:00:00 2001 From: zz_y Date: Wed, 9 Sep 2026 13:25:08 -0600 Subject: [PATCH 1/3] Reference shared summary catalog from collector and transmission plans --- Cargo.lock | 12 ++ control_plane/Cargo.toml | 1 + control_plane/src/opamp/mod.rs | 1 + control_plane/src/physical/compiler.rs | 155 ++++++++++++++++++ control_plane/src/physical/summary_catalog.rs | 36 ++++ .../drivers/ingest/prometheus_remote_write.rs | 1 + data_plane/src/drivers/query/servers/http.rs | 1 + data_plane/src/main.rs | 1 + .../types/hot_reload_config.rs | 1 + 9 files changed, 209 insertions(+) diff --git a/Cargo.lock b/Cargo.lock index 855777f2..24c229e7 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -801,6 +801,7 @@ dependencies = [ "serde", "serde_json", "serde_yaml", + "sha2", "thiserror 1.0.69", "tokio", "tokio-stream", @@ -3209,6 +3210,17 @@ dependencies = [ "digest", ] +[[package]] +name = "sha2" +version = "0.10.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a7507d819769d01a365ab707794a4084392c824f54a7a6a7862f8c3d0892b283" +dependencies = [ + "cfg-if", + "cpufeatures", + "digest", +] + [[package]] name = "sharded-slab" version = "0.1.7" diff --git a/control_plane/Cargo.toml b/control_plane/Cargo.toml index b96327c7..00c8fe2c 100644 --- a/control_plane/Cargo.toml +++ b/control_plane/Cargo.toml @@ -17,6 +17,7 @@ axum = { version = "0.7", features = ["ws"] } futures-util = "0.3" serde = { version = "1", features = ["derive"] } serde_json = "1" +sha2 = "0.10" serde_yaml = "0.9" reqwest = { version = "0.12", default-features = false, features = ["json", "rustls-tls"] } anyhow = "1" diff --git a/control_plane/src/opamp/mod.rs b/control_plane/src/opamp/mod.rs index 74e73550..4cbd60f3 100644 --- a/control_plane/src/opamp/mod.rs +++ b/control_plane/src/opamp/mod.rs @@ -1004,6 +1004,7 @@ mod tests { plan_id: u64, ) -> crate::physical::compiler::CollectorPlan { crate::physical::compiler::CollectorPlan { + summary_catalog: None, collector_id: collector_id.into(), envelope: crate::physical::compiler::PlanEnvelope { plan_id, diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index 343454a3..c30df006 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -257,6 +257,9 @@ pub struct CollectorLifecycle { #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] pub struct CollectorPlan { + /// Absent only in legacy artifacts; catalog-aware validation requires it. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub summary_catalog: Option, pub collector_id: String, pub envelope: PlanEnvelope, pub materializations: Vec, @@ -904,6 +907,9 @@ pub struct RuntimeAdaptationEvidence { #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] #[serde(deny_unknown_fields)] pub struct TransmissionPlan { + /// Absent only in legacy artifacts; catalog-aware validation requires it. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub summary_catalog: Option, pub envelope: PlanEnvelope, pub frame_identity: FrameIdentityContract, pub rules: Vec, @@ -943,6 +949,8 @@ pub struct SummaryFrameIdentity { #[derive(Debug, Error, PartialEq, Eq)] pub enum TransmissionPlanError { + #[error("summary catalog mismatch: {0}")] + Catalog(String), #[error("TransmissionPlan envelope differs from PrecomputePlan")] EnvelopeMismatch, #[error("transmission rules do not exactly match precompute producer bindings")] @@ -966,7 +974,64 @@ pub enum TransmissionPlanError { }, } +fn validate_catalog_projection( + reference: Option<&super::summary_catalog::SummaryCatalogReference>, + envelope: &PlanEnvelope, + materializations: impl IntoIterator, + catalog: &super::summary_catalog::SummaryCatalog, +) -> Result<(), TransmissionPlanError> { + let expected = catalog + .reference() + .map_err(|error| TransmissionPlanError::Catalog(error.to_string()))?; + if reference != Some(&expected) + || envelope.plan_id != catalog.plan_id + || envelope.plan_version != catalog.plan_version + { + return Err(TransmissionPlanError::Catalog( + "missing or different snapshot reference".into(), + )); + } + for id in materializations { + if !catalog.materializations.contains_key(&id) { + return Err(TransmissionPlanError::Catalog(format!( + "unknown materialization {}", + id.0 + ))); + } + } + Ok(()) +} + +impl CollectorPlan { + pub fn validate_against_catalog( + &self, + catalog: &super::summary_catalog::SummaryCatalog, + ) -> Result<(), TransmissionPlanError> { + validate_catalog_projection( + self.summary_catalog.as_ref(), + &self.envelope, + self.materializations + .iter() + .map(|m| m.materialization) + .chain(self.transmission_rules.iter().map(|r| r.materialization)), + catalog, + ) + } +} + impl TransmissionPlan { + pub fn validate_against_catalog( + &self, + catalog: &super::summary_catalog::SummaryCatalog, + ) -> Result<(), TransmissionPlanError> { + validate_catalog_projection( + self.summary_catalog.as_ref(), + &self.envelope, + self.rules.iter().map(|r| r.materialization), + catalog, + ) + } + pub fn build( envelope: PlanEnvelope, precompute: &PrecomputePlan, @@ -1016,7 +1081,18 @@ impl TransmissionPlan { } }) .collect(); + let catalog = super::summary_catalog::SummaryCatalog::from_materializations( + envelope.plan_id, + envelope.plan_version, + &precompute.materializations, + ) + .map_err(|error| TransmissionPlanError::Catalog(error.to_string()))?; let plan = Self { + summary_catalog: Some( + catalog + .reference() + .map_err(|error| TransmissionPlanError::Catalog(error.to_string()))?, + ), envelope, frame_identity: FrameIdentityContract { identity_version: 1, @@ -1031,6 +1107,15 @@ impl TransmissionPlan { } pub fn validate(&self, precompute: &PrecomputePlan) -> Result<(), TransmissionPlanError> { + if self.summary_catalog.is_some() { + let catalog = super::summary_catalog::SummaryCatalog::from_materializations( + precompute.envelope.plan_id, + precompute.envelope.plan_version, + &precompute.materializations, + ) + .map_err(|error| TransmissionPlanError::Catalog(error.to_string()))?; + self.validate_against_catalog(&catalog)?; + } if self.envelope != precompute.envelope { return Err(TransmissionPlanError::EnvelopeMismatch); } @@ -2270,6 +2355,7 @@ impl PhysicalCompiler { let collector_plans = producer_ids .into_iter() .map(|collector_id| CollectorPlan { + summary_catalog: transmission_plan.summary_catalog.clone(), transmission_rules: transmission_plan .rules .iter() @@ -2411,6 +2497,20 @@ impl PhysicalCompiler { query_id: "summary-catalog".into(), reason: error.to_string(), })?; + transmission_plan + .validate_against_catalog(&summary_catalog) + .map_err(|error| CompileError::Query { + query_id: "summary-catalog".into(), + reason: error.to_string(), + })?; + for collector in &collector_plans { + collector + .validate_against_catalog(&summary_catalog) + .map_err(|error| CompileError::Query { + query_id: "summary-catalog".into(), + reason: error.to_string(), + })?; + } Ok(PhysicalPlan { envelope, summary_catalog, @@ -4018,6 +4118,61 @@ mod tests { } } + // Projections reference one immutable snapshot and reject drift or foreign state. + #[test] + fn catalog_projection_rejects_missing_stale_and_foreign_references() { + let snapshot: BackendLocalPlanningSnapshot = serde_json::from_str(include_str!( + "../../../docs/examples/asapquery-planning-snapshot.json" + )) + .unwrap(); + let bundle = snapshot.compile().unwrap(); + let catalog = &bundle.summary_catalog; + let mut transmission = bundle.transmission_plan.clone(); + transmission.validate_against_catalog(catalog).unwrap(); + assert_eq!( + transmission.summary_catalog.as_ref(), + Some(&catalog.reference().unwrap()) + ); + transmission + .summary_catalog + .as_mut() + .unwrap() + .snapshot_sha256 + .push('0'); + assert!(transmission.validate_against_catalog(catalog).is_err()); + assert!(transmission.validate(&bundle.precompute_plan).is_err()); + transmission.summary_catalog = None; + assert!(transmission.validate_against_catalog(catalog).is_err()); + let mut collector = CollectorPlan { + summary_catalog: Some(catalog.reference().unwrap()), + collector_id: "collector-test".into(), + envelope: bundle.envelope.clone(), + materializations: vec![], + transmission_rules: vec![], + }; + collector.validate_against_catalog(catalog).unwrap(); + collector.envelope.plan_version += 1; + assert!(collector.validate_against_catalog(catalog).is_err()); + assert!(validate_catalog_projection( + Some(&catalog.reference().unwrap()), + &bundle.envelope, + [asap_types::PolicyFingerprint(u64::MAX)], + catalog + ) + .is_err()); + let encoded = serde_json::to_value(&collector).unwrap(); + assert!(encoded["summary_catalog"] + .get("summary_descriptors") + .is_none()); + for actual in &bundle.collector_plans { + actual.validate_against_catalog(catalog).unwrap(); + assert_eq!( + actual.summary_catalog, + bundle.transmission_plan.summary_catalog + ); + } + } + #[test] fn canonical_snapshot_preserves_shared_bindings_after_serialization() { // Two different registered readouts survive publication with one state. diff --git a/control_plane/src/physical/summary_catalog.rs b/control_plane/src/physical/summary_catalog.rs index dfee576b..caea096b 100644 --- a/control_plane/src/physical/summary_catalog.rs +++ b/control_plane/src/physical/summary_catalog.rs @@ -13,6 +13,16 @@ use serde::{Deserialize, Serialize}; pub const SUMMARY_CATALOG_SCHEMA_VERSION: u32 = 1; +/// Identifies one immutable catalog snapshot without duplicating descriptors. +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[serde(deny_unknown_fields)] +pub struct SummaryCatalogReference { + pub schema_version: u32, + pub plan_id: u64, + pub plan_version: u64, + pub snapshot_sha256: String, +} + /// Stable materialization identity binds operator and population descriptors. /// Concrete intervals, groups and completeness belong to runtime instances. #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] @@ -46,6 +56,20 @@ pub enum SummaryCatalogError { } impl SummaryCatalog { + pub fn reference(&self) -> Result { + use sha2::{Digest, Sha256}; + self.validate()?; + // BTreeMap tables and typed descriptor fields serialize deterministically. + let bytes = serde_json::to_vec(self) + .map_err(|error| SummaryCatalogError::Descriptor(error.to_string()))?; + Ok(SummaryCatalogReference { + schema_version: self.schema_version, + plan_id: self.plan_id, + plan_version: self.plan_version, + snapshot_sha256: format!("{:x}", Sha256::digest(bytes)), + }) + } + pub fn from_materializations( plan_id: u64, plan_version: u64, @@ -176,6 +200,18 @@ mod tests { ) } + // Content changes invalidate references even when plan/version are reused. + #[test] + fn snapshot_reference_is_deterministic_and_content_sensitive() { + let a = config("requests", "", 60); + let b = config("errors", "", 60); + let left = SummaryCatalog::from_materializations(1, 2, &[a.clone(), b.clone()]).unwrap(); + let reordered = SummaryCatalog::from_materializations(1, 2, &[b, a.clone()]).unwrap(); + assert_eq!(left.reference().unwrap(), reordered.reference().unwrap()); + let changed = SummaryCatalog::from_materializations(1, 2, &[a]).unwrap(); + assert_ne!(left.reference().unwrap(), changed.reference().unwrap()); + } + // Panes share definitions but retain distinct physical materialization IDs. #[test] fn shares_descriptors_across_windows_and_deduplicates_materializations() { diff --git a/data_plane/src/drivers/ingest/prometheus_remote_write.rs b/data_plane/src/drivers/ingest/prometheus_remote_write.rs index f63156f6..58122622 100644 --- a/data_plane/src/drivers/ingest/prometheus_remote_write.rs +++ b/data_plane/src/drivers/ingest/prometheus_remote_write.rs @@ -788,6 +788,7 @@ mod tests { materializations: streaming.aggregation_configs.values().cloned().collect(), }, transmission_plan: TransmissionPlan { + summary_catalog: None, envelope, frame_identity: FrameIdentityContract { identity_version: 1, diff --git a/data_plane/src/drivers/query/servers/http.rs b/data_plane/src/drivers/query/servers/http.rs index 226ca617..fd4c59a6 100644 --- a/data_plane/src/drivers/query/servers/http.rs +++ b/data_plane/src/drivers/query/servers/http.rs @@ -2537,6 +2537,7 @@ mod tests { materializations: Vec::new(), }, transmission_plan: TransmissionPlan { + summary_catalog: None, envelope, frame_identity: FrameIdentityContract { identity_version: 1, diff --git a/data_plane/src/main.rs b/data_plane/src/main.rs index 7832356f..0235405e 100644 --- a/data_plane/src/main.rs +++ b/data_plane/src/main.rs @@ -701,6 +701,7 @@ async fn main() -> Result<()> { .collect(), }; let initial_transmission_plan = control_plane::physical::compiler::TransmissionPlan { + summary_catalog: None, envelope: initial_precompute_plan.envelope.clone(), frame_identity: control_plane::physical::compiler::FrameIdentityContract { identity_version: 1, 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 5d75f2b6..7c8e41ac 100644 --- a/data_plane/src/storage_engines/types/hot_reload_config.rs +++ b/data_plane/src/storage_engines/types/hot_reload_config.rs @@ -838,6 +838,7 @@ mod tests { materializations: Vec::new(), }, transmission_plan: control_plane::physical::compiler::TransmissionPlan { + summary_catalog: None, envelope: control_plane::physical::compiler::PlanEnvelope { plan_id, plan_version, From 69228887f93f84935cd70ae6f179ff437058dbad Mon Sep 17 00:00:00 2001 From: zz_y Date: Wed, 9 Sep 2026 13:30:44 -0600 Subject: [PATCH 2/3] Reference SummaryCatalog from execution plans --- control_plane/src/physical/compiler.rs | 111 +++++++- control_plane/src/physical/mod.rs | 1 + .../src/physical/precompute_contract.rs | 264 ++++++++++++++++++ .../drivers/ingest/prometheus_remote_write.rs | 2 + data_plane/src/drivers/query/servers/http.rs | 2 + data_plane/src/main.rs | 2 + .../types/hot_reload_config.rs | 2 + 7 files changed, 380 insertions(+), 4 deletions(-) create mode 100644 control_plane/src/physical/precompute_contract.rs diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index c30df006..d68e9092 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -273,6 +273,13 @@ pub struct CollectorPlan { /// series, maintains windows, and writes content-addressed materializations. #[derive(Debug, Clone, Serialize, Deserialize)] pub struct PrecomputePlan { + #[serde(default, skip_serializing_if = "Option::is_none")] + pub summary_catalog: Option, + #[serde(default, skip_serializing_if = "BTreeMap::is_empty")] + pub materialization_contracts: BTreeMap< + asap_types::sds::MaterializationId, + super::precompute_contract::PrecomputeMaterializationContract, + >, pub envelope: PlanEnvelope, pub ingest: IngestContract, pub schemas: Vec, @@ -399,6 +406,8 @@ pub struct ProducerContract { #[derive(Debug, Error, PartialEq, Eq)] pub enum PrecomputePlanError { + #[error("invalid precompute catalog contract: {0}")] + CatalogContract(String), #[error("PrecomputePlan envelope does not match BackendPlan identity/lifecycle")] PlanIdentityMismatch, #[error("unsupported precompute ingest protocol/endpoint/identity contract")] @@ -470,6 +479,8 @@ impl PrecomputePlan { }) .collect(); let plan = Self { + summary_catalog: None, + materialization_contracts: BTreeMap::new(), envelope, ingest: IngestContract { protocol: IngestProtocol::ModifiedOtlpMetricsV1, @@ -608,7 +619,7 @@ impl PrecomputePlan { return Err(PrecomputePlanError::MissingProducer(missing.0)); } } - Ok(()) + self.validate_catalog_contract() } pub fn validate_against_backend( @@ -992,7 +1003,10 @@ fn validate_catalog_projection( )); } for id in materializations { - if !catalog.materializations.contains_key(&id) { + if !catalog + .materializations + .contains_key(&asap_types::sds::MaterializationId::from(id)) + { return Err(TransmissionPlanError::Catalog(format!( "unknown materialization {}", id.0 @@ -2326,6 +2340,20 @@ impl PhysicalCompiler { .entry(materialization.policy_fingerprint()) .or_insert(materialization); } + // The compatibility emitter may clamp pane sizes. Publish the actual + // executable configuration in both projections, never the pre-clamp size. + for (id, config) in &materializations_by_fingerprint { + if let Some(state) = backend_plan.materializations.get_mut(id) { + state.window.size_ms = + config + .window_size + .checked_mul(1000) + .ok_or_else(|| CompileError::Query { + query_id: "precompute-window".into(), + reason: "window overflow".into(), + })?; + } + } let materializations = materializations_by_fingerprint.into_values().collect(); let mut precompute_plan = match environment.target { PhysicalDeploymentTarget::DistributedCollectors => PrecomputePlan::build( @@ -2497,6 +2525,12 @@ impl PhysicalCompiler { query_id: "summary-catalog".into(), reason: error.to_string(), })?; + precompute_plan + .bind_catalog(&summary_catalog) + .map_err(|error| CompileError::Query { + query_id: "precompute-catalog".into(), + reason: error.to_string(), + })?; transmission_plan .validate_against_catalog(&summary_catalog) .map_err(|error| CompileError::Query { @@ -2643,7 +2677,7 @@ pub fn select_post_asap( ) } -fn state_schema_id(fingerprint: asap_types::PolicyFingerprint) -> String { +pub(super) fn state_schema_id(fingerprint: asap_types::PolicyFingerprint) -> String { format!( "{}:summary-state:v1:{}", backend_plan::BACKEND_COMPAT, @@ -2651,7 +2685,7 @@ fn state_schema_id(fingerprint: asap_types::PolicyFingerprint) -> String { ) } -fn state_encodings(family: &SummaryFamilyType) -> Vec { +pub(super) fn state_encodings(family: &SummaryFamilyType) -> Vec { match family { SummaryFamilyType::ExactAggregate(..) => vec![StateEncoding::ExactAccumulatorV1], SummaryFamilyType::Sketch(kind, _) @@ -3460,6 +3494,7 @@ mod tests { .compile(request("counter", "rate(m[1m])"), environment(10_000)) .unwrap(); plan.precompute_plan.validate().unwrap(); + let catalog = plan.summary_catalog.clone(); let mut raw = plan.precompute_plan; raw.ingest.protocol = IngestProtocol::PrometheusRemoteWriteV1; raw.ingest.endpoint_path = "/api/v1/write".into(); @@ -3468,6 +3503,7 @@ mod tests { raw.ingest.require_materialization_identity = false; raw.ingest.require_registered_producer = false; raw.producers.clear(); + raw.bind_catalog(&catalog).unwrap(); raw.validate().unwrap(); } @@ -4200,6 +4236,73 @@ mod tests { query_plan.validate(&bindings).unwrap(); } + #[test] + fn precompute_catalog_validates_without_backend_projection() { + let bundle = PhysicalCompiler + .compile( + request("catalog", "quantile_over_time(0.99, m[1m])"), + environment(10_000), + ) + .unwrap(); + let original = &bundle.precompute_plan; + let catalog = &bundle.summary_catalog; + let roundtrip: PrecomputePlan = + serde_json::from_slice(&serde_json::to_vec(original).unwrap()).unwrap(); + roundtrip.validate_against_catalog(catalog).unwrap(); + let reject = + |mutated: PrecomputePlan| assert!(mutated.validate_against_catalog(catalog).is_err()); + let mut bad = original.clone(); + bad.schemas[0].source = planner_types::pre_asap::Source::TimeSeries { + metric: "other".into(), + }; + reject(bad); + let mut bad = original.clone(); + bad.schemas[0].value_column = planner_types::pre_asap::ColumnRef::Named("other".into()); + reject(bad); + let mut bad = original.clone(); + bad.schemas[0].group_by.push("other".into()); + reject(bad); + let mut bad = original.clone(); + bad.schemas[0].window.size_ms += 1; + reject(bad); + let mut bad = original.clone(); + bad.materialization_contracts + .values_mut() + .next() + .unwrap() + .activation_unix_ms += 1; + reject(bad); + let mut bad = original.clone(); + bad.materialization_contracts + .values_mut() + .next() + .unwrap() + .retained_windows = Some(0); + reject(bad); + let mut bad = original.clone(); + bad.materialization_contracts + .values_mut() + .next() + .unwrap() + .update = super::super::precompute_contract::PrecomputeUpdate::RawSamples; + reject(bad); + let mut bad = original.clone(); + bad.summary_catalog + .as_mut() + .unwrap() + .snapshot_sha256 + .push('0'); + reject(bad); + let mut corrupt_catalog = catalog.clone(); + corrupt_catalog.data_descriptors.clear(); + assert!(original.validate_against_catalog(&corrupt_catalog).is_err()); + let mut legacy = original.clone(); + legacy.summary_catalog = None; + legacy.materialization_contracts.clear(); + legacy.validate().unwrap(); + assert!(legacy.validate_against_catalog(catalog).is_err()); + } + #[test] fn compiles_one_decision_into_matching_collector_and_backend_views() { let bundle = PhysicalCompiler diff --git a/control_plane/src/physical/mod.rs b/control_plane/src/physical/mod.rs index b3c168da..b347524b 100644 --- a/control_plane/src/physical/mod.rs +++ b/control_plane/src/physical/mod.rs @@ -29,6 +29,7 @@ pub mod plan; pub mod plan_cache; pub mod planner; pub mod post_asap; +pub mod precompute_contract; pub mod runtime_capability; pub mod sketch_catalog; pub mod stage_split; diff --git a/control_plane/src/physical/precompute_contract.rs b/control_plane/src/physical/precompute_contract.rs new file mode 100644 index 00000000..302999cd --- /dev/null +++ b/control_plane/src/physical/precompute_contract.rs @@ -0,0 +1,264 @@ +//! Self-contained precompute execution contracts over a SummaryCatalog snapshot. +//! BackendPlan remains a compatibility projection, not a validation authority. +use super::compiler::*; +use super::summary_catalog::{MaterializationIdentity, SummaryCatalog}; +use asap_types::sds::{MaterializationId, SummaryDescriptor}; +use planner_types::pre_asap::{ColumnRef, Source}; +use serde::{Deserialize, Serialize}; +use std::collections::{BTreeMap, BTreeSet}; + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[serde(rename_all = "snake_case")] +pub enum PrecomputeUpdate { + RawSamples, + SummaryFrames, +} +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[serde(tag = "kind", rename_all = "snake_case", deny_unknown_fields)] +pub enum PrecomputePlacement { + BackendLocal, + Collectors { producer_ids: BTreeSet }, +} +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[serde(deny_unknown_fields)] +pub struct PrecomputeMaterializationContract { + pub descriptors: MaterializationIdentity, + pub schema_id: String, + pub update: PrecomputeUpdate, + pub placement: PrecomputePlacement, + pub activation_unix_ms: u64, + pub expiry_unix_ms: Option, + pub retained_windows: Option, +} +fn invalid(reason: impl Into) -> PrecomputePlanError { + PrecomputePlanError::CatalogContract(reason.into()) +} + +impl PrecomputePlan { + /// Bind references after retention decisions are finalized. The plan keeps + /// only the immutable snapshot reference; descriptors are installed once. + pub fn bind_catalog(&mut self, catalog: &SummaryCatalog) -> Result<(), PrecomputePlanError> { + let mut contracts = BTreeMap::new(); + for config in &self.materializations { + let id = MaterializationId::from(config.policy_fingerprint()); + let descriptors = catalog + .materializations + .get(&id) + .ok_or_else(|| invalid(format!("missing catalog materialization {}", id.as_u64())))? + .clone(); + let schema = self + .schemas + .iter() + .find(|s| s.materialization == id.fingerprint()) + .ok_or(PrecomputePlanError::SchemaSetMismatch)?; + let (update, placement) = self.expected_placement(id); + contracts.insert( + id, + PrecomputeMaterializationContract { + descriptors, + schema_id: schema.schema_id.clone(), + update, + placement, + activation_unix_ms: self.envelope.activation_unix_ms, + expiry_unix_ms: self.envelope.expiry_unix_ms, + retained_windows: config.num_aggregates_to_retain, + }, + ); + } + self.summary_catalog = Some( + catalog + .reference() + .map_err(|error| invalid(error.to_string()))?, + ); + self.materialization_contracts = contracts; + self.validate_against_catalog(catalog) + } + fn expected_placement(&self, id: MaterializationId) -> (PrecomputeUpdate, PrecomputePlacement) { + match self.ingest.protocol { + IngestProtocol::PrometheusRemoteWriteV1 => ( + PrecomputeUpdate::RawSamples, + PrecomputePlacement::BackendLocal, + ), + IngestProtocol::ModifiedOtlpMetricsV1 => ( + PrecomputeUpdate::SummaryFrames, + PrecomputePlacement::Collectors { + producer_ids: self + .producers + .iter() + .filter(|p| p.materialization == id.fingerprint()) + .map(|p| p.producer_id.clone()) + .collect(), + }, + ), + } + } + /// Strict validation resolves this plan's reference against the one catalog + /// snapshot carried by the enclosing physical-plan installation. + pub fn validate_against_catalog( + &self, + catalog: &SummaryCatalog, + ) -> Result<(), PrecomputePlanError> { + self.validate()?; + self.validate_catalog_contents(catalog) + } + pub(super) fn validate_catalog_contract(&self) -> Result<(), PrecomputePlanError> { + match ( + &self.summary_catalog, + self.materialization_contracts.is_empty(), + ) { + (None, true) | (Some(_), false) => Ok(()), + (None, false) => Err(invalid("descriptor references without catalog")), + (Some(_), true) if self.materializations.is_empty() => Ok(()), + (Some(_), true) => Err(invalid( + "catalog reference without materialization contracts", + )), + } + } + + fn validate_catalog_contents( + &self, + catalog: &SummaryCatalog, + ) -> Result<(), PrecomputePlanError> { + catalog.validate().map_err(|e| invalid(e.to_string()))?; + let expected_reference = catalog + .reference() + .map_err(|error| invalid(error.to_string()))?; + if self.summary_catalog.as_ref() != Some(&expected_reference) + || catalog.plan_id != self.envelope.plan_id + || catalog.plan_version != self.envelope.plan_version + { + return Err(PrecomputePlanError::PlanIdentityMismatch); + } + if self.envelope.activation_unix_ms < self.envelope.generated_at_unix_ms + || self + .envelope + .expiry_unix_ms + .is_some_and(|expiry| expiry <= self.envelope.activation_unix_ms) + { + return Err(invalid("invalid plan activation/expiry lifecycle")); + } + let ids: BTreeSet<_> = self + .materializations + .iter() + .map(|m| MaterializationId::from(m.policy_fingerprint())) + .collect(); + if ids != catalog.materializations.keys().copied().collect() + || ids != self.materialization_contracts.keys().copied().collect() + { + return Err(invalid("catalog/reference/materialization sets differ")); + } + for config in &self.materializations { + let id = MaterializationId::from(config.policy_fingerprint()); + let binding = &catalog.materializations[&id]; + let contract = &self.materialization_contracts[&id]; + if &contract.descriptors != binding { + return Err(invalid("materialization descriptor reference drift")); + } + let expected = + SummaryDescriptor::from_config(config).map_err(|e| invalid(e.to_string()))?; + if binding.summary_descriptor_id != expected.id { + return Err(invalid( + "summary operator/update contract differs from catalog", + )); + } + let data = &catalog.data_descriptors[&binding.data_descriptor_id]; + if data.metric_name != config.metric + || data.population_filter_canonical + != asap_types::utils::normalize_spatial_filter(&config.spatial_filter) + || data.group_by_keys != config.grouping_labels.labels.iter().cloned().collect() + { + return Err(invalid("source/population/grouping differs from catalog")); + } + if config.spatial_filter_normalized != data.population_filter_canonical { + return Err(invalid("normalized population predicate drift")); + } + let schema = self + .schemas + .iter() + .find(|s| s.materialization == id.fingerprint()) + .ok_or(PrecomputePlanError::SchemaSetMismatch)?; + let family = config + .accumulator_spec() + .map_err(|e| invalid(e.to_string()))? + .family; + let expected_family = StateFamilyContract::try_from(&family) + .map_err(|_| PrecomputePlanError::UnsupportedFamily(id.as_u64()))?; + let source = if let Some(table) = &config.table_name { + Source::Table { + table_ref: table.clone(), + } + } else { + Source::TimeSeries { + metric: config.metric.clone(), + } + }; + let col = if let Some(column) = &config.value_column { + ColumnRef::Named(column.clone()) + } else { + ColumnRef::SampleValue + }; + let size = config + .window_size + .checked_mul(1000) + .filter(|v| *v > 0) + .ok_or_else(|| invalid("invalid window size"))?; + let slide = config + .slide_interval + .checked_mul(1000) + .filter(|v| *v > 0) + .ok_or_else(|| invalid("invalid slide interval"))?; + let expected_slide = match config.window_type { + asap_types::WindowKind::Tumbling => None, + asap_types::WindowKind::Sliding => Some(slide), + asap_types::WindowKind::Session => { + return Err(invalid("session lifecycle is not supported")) + } + }; + if schema.schema_id != contract.schema_id + || schema.schema_id != state_schema_id(id.fingerprint()) + || schema.family != expected_family + || schema.source != source + || schema.value_column != col + || schema.group_by != config.grouping_labels.labels + || schema.window.kind != config.window_type + || schema.window.size_ms != size + || schema.window.slide_ms != expected_slide + { + return Err(PrecomputePlanError::InvalidSchema { + schema_id: schema.schema_id.clone(), + }); + } + if schema.encodings != state_encodings(&family) { + return Err(invalid("encoding does not match state family")); + } + let (update, placement) = self.expected_placement(id); + if contract.update != update || contract.placement != placement { + return Err(invalid( + "update/placement does not match ingress and producer bindings", + )); + } + if contract.activation_unix_ms != self.envelope.activation_unix_ms + || contract.expiry_unix_ms != self.envelope.expiry_unix_ms + || contract.retained_windows != config.num_aggregates_to_retain + || contract.retained_windows == Some(0) + { + return Err(invalid("materialization lifecycle/retention mismatch")); + } + } + if self.ingest.protocol == IngestProtocol::PrometheusRemoteWriteV1 + && !self.producers.is_empty() + { + return Err(invalid( + "backend-local plan cannot declare collector placement", + )); + } + if self.producers.iter().any(|p| { + p.producer_id.trim().is_empty() + || p.collector_id.trim().is_empty() + || p.producer_id != p.collector_id + }) { + return Err(invalid("invalid collector placement identity")); + } + Ok(()) + } +} diff --git a/data_plane/src/drivers/ingest/prometheus_remote_write.rs b/data_plane/src/drivers/ingest/prometheus_remote_write.rs index 58122622..961b9585 100644 --- a/data_plane/src/drivers/ingest/prometheus_remote_write.rs +++ b/data_plane/src/drivers/ingest/prometheus_remote_write.rs @@ -774,6 +774,8 @@ mod tests { }; let active = ActivePhysicalPlan { precompute_plan: PrecomputePlan { + summary_catalog: None, + materialization_contracts: Default::default(), envelope: envelope.clone(), ingest: IngestContract { protocol: IngestProtocol::PrometheusRemoteWriteV1, diff --git a/data_plane/src/drivers/query/servers/http.rs b/data_plane/src/drivers/query/servers/http.rs index fd4c59a6..91daf1f1 100644 --- a/data_plane/src/drivers/query/servers/http.rs +++ b/data_plane/src/drivers/query/servers/http.rs @@ -2523,6 +2523,8 @@ mod tests { let active = crate::storage_engines::types::HotReloadActivePhysicalPlan::new( crate::storage_engines::types::ActivePhysicalPlan { precompute_plan: PrecomputePlan { + summary_catalog: None, + materialization_contracts: Default::default(), envelope: envelope.clone(), ingest: IngestContract { protocol: IngestProtocol::PrometheusRemoteWriteV1, diff --git a/data_plane/src/main.rs b/data_plane/src/main.rs index 0235405e..703f9f19 100644 --- a/data_plane/src/main.rs +++ b/data_plane/src/main.rs @@ -674,6 +674,8 @@ async fn main() -> Result<()> { // sharing contract as `hot_reload_config`. let initial_backend_plan = control_plane::backend_plan::BackendPlan::default(); let initial_precompute_plan = control_plane::physical::compiler::PrecomputePlan { + summary_catalog: None, + materialization_contracts: Default::default(), envelope: control_plane::physical::compiler::PlanEnvelope { plan_id: 0, plan_version: 0, 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 7c8e41ac..0f53ff2a 100644 --- a/data_plane/src/storage_engines/types/hot_reload_config.rs +++ b/data_plane/src/storage_engines/types/hot_reload_config.rs @@ -822,6 +822,8 @@ mod tests { }; ActivePhysicalPlan { precompute_plan: control_plane::physical::compiler::PrecomputePlan { + summary_catalog: None, + materialization_contracts: Default::default(), envelope, ingest: control_plane::physical::compiler::IngestContract { protocol: From a374c0deab358e9e660624148290567d704504e6 Mon Sep 17 00:00:00 2001 From: zz_y Date: Wed, 9 Sep 2026 14:48:12 -0600 Subject: [PATCH 3/3] Format catalog plan projections --- data_plane/src/drivers/query/servers/http.rs | 4 ++-- data_plane/src/main.rs | 4 ++-- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/data_plane/src/drivers/query/servers/http.rs b/data_plane/src/drivers/query/servers/http.rs index 91daf1f1..acea438a 100644 --- a/data_plane/src/drivers/query/servers/http.rs +++ b/data_plane/src/drivers/query/servers/http.rs @@ -2523,8 +2523,8 @@ mod tests { let active = crate::storage_engines::types::HotReloadActivePhysicalPlan::new( crate::storage_engines::types::ActivePhysicalPlan { precompute_plan: PrecomputePlan { - summary_catalog: None, - materialization_contracts: Default::default(), + summary_catalog: None, + materialization_contracts: Default::default(), envelope: envelope.clone(), ingest: IngestContract { protocol: IngestProtocol::PrometheusRemoteWriteV1, diff --git a/data_plane/src/main.rs b/data_plane/src/main.rs index 703f9f19..f50c05a3 100644 --- a/data_plane/src/main.rs +++ b/data_plane/src/main.rs @@ -674,8 +674,8 @@ async fn main() -> Result<()> { // sharing contract as `hot_reload_config`. let initial_backend_plan = control_plane::backend_plan::BackendPlan::default(); let initial_precompute_plan = control_plane::physical::compiler::PrecomputePlan { - summary_catalog: None, - materialization_contracts: Default::default(), + summary_catalog: None, + materialization_contracts: Default::default(), envelope: control_plane::physical::compiler::PlanEnvelope { plan_id: 0, plan_version: 0,