From 316660a4e58898fda2137cf0fedf18000372f63d Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 10 Sep 2026 08:14:50 -0600 Subject: [PATCH 1/2] Use summary definition IDs across physical plans --- control_plane/src/opamp/mod.rs | 4 +- control_plane/src/physical/compiler.rs | 44 +++++++++---------- .../src/physical/precompute_contract.rs | 2 +- control_plane/src/physical/workload_cost.rs | 2 +- crates/asap_types/src/sds.rs | 2 +- data_plane/src/drivers/ingest/otel.rs | 14 +++--- .../src/precompute_engine/frame_lineage.rs | 4 +- .../storage_engines/sketch_db/index/mod.rs | 2 +- 8 files changed, 37 insertions(+), 37 deletions(-) diff --git a/control_plane/src/opamp/mod.rs b/control_plane/src/opamp/mod.rs index b58d7f11..47817887 100644 --- a/control_plane/src/opamp/mod.rs +++ b/control_plane/src/opamp/mod.rs @@ -1018,7 +1018,7 @@ mod tests { }, materializations: vec![crate::physical::compiler::CollectorMaterialization { query_id: "q".into(), - materialization: asap_types::PolicyFingerprint(1), + materialization: asap_types::PolicyFingerprint(1).into(), metric: "requests".into(), algorithm: "hll".into(), parameters: serde_json::json!({"precision": 14}), @@ -1039,7 +1039,7 @@ mod tests { }, }], transmission_rules: vec![crate::physical::compiler::TransmissionRule { - materialization: asap_types::PolicyFingerprint(1), + materialization: asap_types::PolicyFingerprint(1).into(), producer_id: collector_id.into(), schema_id: "asap-query-backend.v1:summary-state:v1:1".into(), mode: crate::physical::compiler::TransmissionMode::Full, diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index 88bfbfc4..205ae19b 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -263,7 +263,7 @@ pub struct PlanEnvelope { #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] pub struct CollectorMaterialization { pub query_id: String, - pub materialization: asap_types::PolicyFingerprint, + pub materialization: asap_types::sds::SummaryDefinitionId, pub metric: String, pub algorithm: String, pub parameters: Value, @@ -412,7 +412,7 @@ impl TryFrom<&SummaryFamilyType> for StateFamilyContract { pub struct StateSchemaContract { pub schema_id: String, pub schema_version: u32, - pub materialization: asap_types::PolicyFingerprint, + pub materialization: asap_types::sds::SummaryDefinitionId, pub family: StateFamilyContract, pub source: Source, pub value_column: planner_types::pre_asap::ColumnRef, @@ -442,7 +442,7 @@ pub struct StateWindowContract { pub struct ProducerContract { pub producer_id: String, pub collector_id: String, - pub materialization: asap_types::PolicyFingerprint, + pub materialization: asap_types::sds::SummaryDefinitionId, pub schema_id: String, } @@ -501,7 +501,7 @@ impl PrecomputePlan { Ok(StateSchemaContract { schema_id: state_schema_id(fingerprint), schema_version: 1, - materialization: fingerprint, + materialization: fingerprint.into(), family, source, value_column, @@ -615,7 +615,7 @@ impl PrecomputePlan { materialization.policy_fp_u64(), )); } - if !materializations.insert(materialization.policy_fingerprint()) { + if !materializations.insert(materialization.policy_fingerprint().into()) { return Err(PrecomputePlanError::DuplicateMaterialization( materialization.policy_fp_u64(), )); @@ -650,13 +650,13 @@ impl PrecomputePlan { let materialization = self .materializations .iter() - .find(|candidate| candidate.policy_fingerprint() == schema.materialization) + .find(|candidate| candidate.policy_fingerprint() == schema.materialization.fingerprint()) .ok_or(PrecomputePlanError::SchemaSetMismatch)?; let accumulator = materialization .accumulator_spec() - .map_err(|_| PrecomputePlanError::UnsupportedFamily(schema.materialization.0))?; + .map_err(|_| PrecomputePlanError::UnsupportedFamily(schema.materialization.as_u64()))?; let family = StateFamilyContract::try_from(&accumulator.family) - .map_err(|_| PrecomputePlanError::UnsupportedFamily(schema.materialization.0))?; + .map_err(|_| PrecomputePlanError::UnsupportedFamily(schema.materialization.as_u64()))?; let source = materialization.table_name.as_ref().map_or_else( || Source::TimeSeries { metric: materialization.metric.clone(), @@ -670,7 +670,7 @@ impl PrecomputePlan { .clone() .map(planner_types::pre_asap::ColumnRef::Named) .unwrap_or(planner_types::pre_asap::ColumnRef::SampleValue); - if schema.schema_id != state_schema_id(schema.materialization) + if schema.schema_id != state_schema_id(schema.materialization.fingerprint()) || schema.family != family || schema.source != source || schema.value_column != value_column @@ -720,7 +720,7 @@ impl PrecomputePlan { } if self.ingest.require_registered_producer { if let Some(missing) = materializations.difference(&produced).next() { - return Err(PrecomputePlanError::MissingProducer(missing.0)); + return Err(PrecomputePlanError::MissingProducer(missing.as_u64())); } } Ok(()) @@ -744,7 +744,7 @@ pub struct PhysicalPlan { #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] pub struct MaterializationLifecycleEstimate { - pub materialization: asap_types::PolicyFingerprint, + pub materialization: asap_types::sds::SummaryDefinitionId, pub consumer_query_ids: Vec, pub window_implementation_id: String, pub horizon_seconds: f64, @@ -778,7 +778,7 @@ pub struct FrameIdentityContract { #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] #[serde(deny_unknown_fields)] pub struct TransmissionRule { - pub materialization: asap_types::PolicyFingerprint, + pub materialization: asap_types::sds::SummaryDefinitionId, pub producer_id: String, pub schema_id: String, pub mode: TransmissionMode, @@ -981,7 +981,7 @@ pub struct RuntimeRulePolicy { pub struct RuntimeAdaptationEvidence { pub plan_id: u64, pub plan_version: u64, - pub materialization: asap_types::PolicyFingerprint, + pub materialization: asap_types::sds::SummaryDefinitionId, pub producer_id: String, pub schema_id: String, pub producer_version: String, @@ -1017,7 +1017,7 @@ pub struct SummaryFrameIdentity { pub plan_id: u64, pub plan_version: u64, pub backend_compat: String, - pub materialization: asap_types::PolicyFingerprint, + pub materialization: asap_types::sds::SummaryDefinitionId, /// Canonical producer-side identity for one concrete retained-label group. pub series_identity: String, pub schema_id: String, @@ -1062,7 +1062,7 @@ pub enum TransmissionPlanError { fn validate_catalog_projection( reference: Option<&super::summary_catalog::SummaryCatalogReference>, envelope: &PlanEnvelope, - materializations: impl IntoIterator, + materializations: impl IntoIterator, catalog: &super::summary_catalog::SummaryCatalog, ) -> Result<(), TransmissionPlanError> { let expected = catalog @@ -1079,11 +1079,11 @@ fn validate_catalog_projection( for id in materializations { if !catalog .materializations - .contains_key(&asap_types::sds::SummaryDefinitionId::from(id)) + .contains_key(&id) { return Err(TransmissionPlanError::Catalog(format!( "unknown materialization {}", - id.0 + id.as_u64() ))); } } @@ -1143,10 +1143,10 @@ impl TransmissionPlan { let materialization = precompute .materializations .iter() - .find(|m| m.policy_fingerprint() == producer.materialization) + .find(|m| m.policy_fingerprint() == producer.materialization.fingerprint()) .expect("validated PrecomputePlan materialization binding"); let runtime_policy = runtime_policies - .get(&producer.materialization) + .get(&producer.materialization.fingerprint()) .cloned() .unwrap_or_default(); let mode = if runtime_policy.delta.is_some() { @@ -2371,7 +2371,7 @@ impl PhysicalCompiler { lifecycle_estimates .entry(materialization) .or_insert_with(|| MaterializationLifecycleEstimate { - materialization, + materialization: materialization.into(), consumer_query_ids, window_implementation_id: window_implementation.implementation_id.clone(), horizon_seconds: query.lifecycle.horizon_seconds, @@ -2414,7 +2414,7 @@ impl PhysicalCompiler { compiled_materializations.push(runtime_materialization.clone()); let collector_materialization = CollectorMaterialization { query_id: format!("state-{}", materialization.0), - materialization, + materialization: materialization.into(), metric: metric.clone(), algorithm: physical_algorithm, parameters: selected.parameters, @@ -5399,7 +5399,7 @@ mod tests { assert!(validate_catalog_projection( Some(&catalog.reference().unwrap()), &bundle.envelope, - [asap_types::PolicyFingerprint(u64::MAX)], + [asap_types::PolicyFingerprint(u64::MAX).into()], catalog ) .is_err()); diff --git a/control_plane/src/physical/precompute_contract.rs b/control_plane/src/physical/precompute_contract.rs index fa84eeb5..adcf6765 100644 --- a/control_plane/src/physical/precompute_contract.rs +++ b/control_plane/src/physical/precompute_contract.rs @@ -83,7 +83,7 @@ impl PrecomputePlan { let schema = self .schemas .iter() - .find(|s| s.materialization == id.fingerprint()) + .find(|s| s.materialization == id) .ok_or(PrecomputePlanError::SchemaSetMismatch)?; let family = config .accumulator_spec() diff --git a/control_plane/src/physical/workload_cost.rs b/control_plane/src/physical/workload_cost.rs index 4279485c..0d1e83e8 100644 --- a/control_plane/src/physical/workload_cost.rs +++ b/control_plane/src/physical/workload_cost.rs @@ -171,7 +171,7 @@ pub fn manifest( .precompute_plan .materializations .iter() - .find(|m| m.policy_fingerprint() == schema.materialization) + .find(|m| m.policy_fingerprint() == schema.materialization.fingerprint()) .ok_or_else(|| invalid("state has no physical implementation"))?; let identity = json!({"schema": schema, "location": location, "physical": physical, "window_implementation": plan.lifecycle_estimates.iter().find(|e| e.materialization == schema.materialization).map(|e| &e.window_implementation_id)}); diff --git a/crates/asap_types/src/sds.rs b/crates/asap_types/src/sds.rs index 551824e0..2e662778 100644 --- a/crates/asap_types/src/sds.rs +++ b/crates/asap_types/src/sds.rs @@ -28,7 +28,7 @@ macro_rules! descriptor_id { }; } /// Semantic materialization reference. Wire-compatible with PolicyFingerprint, -/// but distinct from descriptor IDs and runtime instance/SID identity. +/// but distinct from descriptor IDs and concrete [`SummaryInstanceId`] identity. #[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)] #[serde(transparent)] pub struct SummaryDefinitionId(pub crate::PolicyFingerprint); diff --git a/data_plane/src/drivers/ingest/otel.rs b/data_plane/src/drivers/ingest/otel.rs index 4a568a47..f8d7c8d8 100644 --- a/data_plane/src/drivers/ingest/otel.rs +++ b/data_plane/src/drivers/ingest/otel.rs @@ -1265,7 +1265,7 @@ async fn route_modified_otlp_sketches_to_precompute( &dp.attrs.keys().cloned().collect(), ) }); - if observed_policy != frame.materialization { + if observed_policy != frame.materialization.fingerprint() { unreachable!( "materialization changed after successful request preflight" ); @@ -1637,7 +1637,7 @@ async fn route_modified_otlp_sketches_to_precompute( debug!( plan_id = frame.plan_id, plan_version = frame.plan_version, - materialization = frame.materialization.0, + materialization = frame.materialization.as_u64(), producer = %frame.producer_id, producer_epoch = %frame.producer_epoch, sequence = frame.sequence, @@ -1652,7 +1652,7 @@ async fn route_modified_otlp_sketches_to_precompute( warn!( plan_id = frame.plan_id, plan_version = frame.plan_version, - materialization = frame.materialization.0, + materialization = frame.materialization.as_u64(), producer = %frame.producer_id, producer_epoch = %frame.producer_epoch, sequence = frame.sequence, @@ -2237,10 +2237,10 @@ fn preflight_summary_frames( &dp.attrs.keys().cloned().collect(), ) }; - if observed != frame.materialization { + if observed != frame.materialization.fingerprint() { return Err(format!( "summary frame for {metric_name} declares materialization {} but active schema resolves {}", - frame.materialization.0, observed.0 + frame.materialization.as_u64(), observed.0 )); } Ok(frame) @@ -2392,7 +2392,7 @@ fn take_summary_frame_identity( plan_id, plan_version, backend_compat, - materialization, + materialization: materialization.into(), series_identity, schema_id, producer_id, @@ -4695,7 +4695,7 @@ mod sid_bucketing_tests { let frame = take_summary_frame_identity(&mut attrs, 100, 200).expect("valid identity"); assert_eq!(frame.plan_id, 42); assert_eq!(frame.plan_version, 3); - assert_eq!(frame.materialization, asap_types::PolicyFingerprint(99)); + assert_eq!(frame.materialization, asap_types::PolicyFingerprint(99).into()); assert_eq!(attrs, HashMap::from([("service".into(), "api".into())])); } } diff --git a/data_plane/src/precompute_engine/frame_lineage.rs b/data_plane/src/precompute_engine/frame_lineage.rs index a400dda3..c283ef41 100644 --- a/data_plane/src/precompute_engine/frame_lineage.rs +++ b/data_plane/src/precompute_engine/frame_lineage.rs @@ -45,7 +45,7 @@ pub enum FrameLineageError { struct FrameLineageKey { plan_id: u64, plan_version: u64, - materialization: asap_types::PolicyFingerprint, + materialization: asap_types::sds::SummaryDefinitionId, series_identity: String, producer_id: String, producer_epoch: String, @@ -199,7 +199,7 @@ mod tests { plan_id: 7, plan_version: 3, backend_compat: "asap-query-backend.v1".into(), - materialization: asap_types::PolicyFingerprint(41), + materialization: asap_types::PolicyFingerprint(41).into(), series_identity: "service=checkout,zone=a".into(), schema_id: "schema-41".into(), producer_id: "edge-a".into(), diff --git a/data_plane/src/storage_engines/sketch_db/index/mod.rs b/data_plane/src/storage_engines/sketch_db/index/mod.rs index 43474bbe..fb0ee47f 100644 --- a/data_plane/src/storage_engines/sketch_db/index/mod.rs +++ b/data_plane/src/storage_engines/sketch_db/index/mod.rs @@ -3130,7 +3130,7 @@ mod tests { plan_id: 7, plan_version: 2, backend_compat: "asap-query-backend.v1".into(), - materialization: PolicyFingerprint(41), + materialization: PolicyFingerprint(41).into(), series_identity: "service=checkout,zone=a".into(), schema_id: "schema-41".into(), producer_id: "edge-a".into(), From 3b7b9967dce182d7a53de62891586372b8e7c2b4 Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 10 Sep 2026 09:59:50 -0600 Subject: [PATCH 2/2] Format typed summary definition references --- control_plane/src/physical/compiler.rs | 20 ++++++++++---------- data_plane/src/drivers/ingest/otel.rs | 5 ++++- 2 files changed, 14 insertions(+), 11 deletions(-) diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index 205ae19b..bb92d9ab 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -650,13 +650,16 @@ impl PrecomputePlan { let materialization = self .materializations .iter() - .find(|candidate| candidate.policy_fingerprint() == schema.materialization.fingerprint()) + .find(|candidate| { + candidate.policy_fingerprint() == schema.materialization.fingerprint() + }) .ok_or(PrecomputePlanError::SchemaSetMismatch)?; - let accumulator = materialization - .accumulator_spec() - .map_err(|_| PrecomputePlanError::UnsupportedFamily(schema.materialization.as_u64()))?; - let family = StateFamilyContract::try_from(&accumulator.family) - .map_err(|_| PrecomputePlanError::UnsupportedFamily(schema.materialization.as_u64()))?; + let accumulator = materialization.accumulator_spec().map_err(|_| { + PrecomputePlanError::UnsupportedFamily(schema.materialization.as_u64()) + })?; + let family = StateFamilyContract::try_from(&accumulator.family).map_err(|_| { + PrecomputePlanError::UnsupportedFamily(schema.materialization.as_u64()) + })?; let source = materialization.table_name.as_ref().map_or_else( || Source::TimeSeries { metric: materialization.metric.clone(), @@ -1077,10 +1080,7 @@ fn validate_catalog_projection( )); } for id in materializations { - if !catalog - .materializations - .contains_key(&id) - { + if !catalog.materializations.contains_key(&id) { return Err(TransmissionPlanError::Catalog(format!( "unknown materialization {}", id.as_u64() diff --git a/data_plane/src/drivers/ingest/otel.rs b/data_plane/src/drivers/ingest/otel.rs index f8d7c8d8..bcf3e308 100644 --- a/data_plane/src/drivers/ingest/otel.rs +++ b/data_plane/src/drivers/ingest/otel.rs @@ -4695,7 +4695,10 @@ mod sid_bucketing_tests { let frame = take_summary_frame_identity(&mut attrs, 100, 200).expect("valid identity"); assert_eq!(frame.plan_id, 42); assert_eq!(frame.plan_version, 3); - assert_eq!(frame.materialization, asap_types::PolicyFingerprint(99).into()); + assert_eq!( + frame.materialization, + asap_types::PolicyFingerprint(99).into() + ); assert_eq!(attrs, HashMap::from([("service".into(), "api".into())])); } }