diff --git a/crates/asap_types/src/aggregation_config.rs b/crates/asap_types/src/aggregation_config.rs index 8254b65b..15144fdb 100644 --- a/crates/asap_types/src/aggregation_config.rs +++ b/crates/asap_types/src/aggregation_config.rs @@ -101,6 +101,11 @@ pub struct PrecomputeMaterialization { pub parameters: HashMap, #[serde(serialize_with = "crate::grouping_projection::serialize_config_grouping")] pub grouping_labels: crate::GroupingProjection, + #[serde( + default, + skip_serializing_if = "crate::grouping_projection::PopulationKeyEncoding::is_legacy" + )] + pub population_key_encoding: crate::grouping_projection::PopulationKeyEncoding, #[serde(default, skip_serializing_if = "Option::is_none")] pub partitioning: Option, #[serde(default, skip_serializing_if = "Option::is_none")] @@ -286,6 +291,7 @@ impl PrecomputeMaterialization { aggregation_sub_type, parameters, grouping_labels: grouping_labels.into(), + population_key_encoding: Default::default(), partitioning: None, derived_input: None, aggregated_labels, @@ -433,6 +439,11 @@ impl PrecomputeMaterialization { table_name, value_column, ); + config.population_key_encoding = data + .get("population_key_encoding") + .map(|value| serde_json::from_value(value.clone())) + .transpose()? + .unwrap_or_default(); config.derived_input = data .get("derived_input") .filter(|v| !v.is_null()) @@ -633,6 +644,11 @@ impl PrecomputeMaterialization { table_name, value_column, ); + config.population_key_encoding = aggregation_data + .get("population_key_encoding") + .map(|value| serde_yaml::from_value(value.clone())) + .transpose()? + .unwrap_or_default(); config.derived_input = aggregation_data .get("derived_input") .filter(|v| !v.is_null()) @@ -683,6 +699,9 @@ impl SerializableToSink for PrecomputeMaterialization { "metric": self.metric, }); + if !self.population_key_encoding.is_legacy() { + json["population_key_encoding"] = serde_json::json!(self.population_key_encoding); + } if let Some(input) = &self.derived_input { json["derived_input"] = serde_json::json!(input); } diff --git a/crates/asap_types/src/grouping_projection.rs b/crates/asap_types/src/grouping_projection.rs index a64018e3..8bb955cb 100644 --- a/crates/asap_types/src/grouping_projection.rs +++ b/crates/asap_types/src/grouping_projection.rs @@ -3,6 +3,21 @@ use crate::KeyByLabelNames; use planner_types::pre_asap::{Column, DataType}; use serde::{Deserialize, Deserializer, Serialize, Serializer}; +/// Version of the population-key routing contract. Choosing a new version +/// changes materialization and data identity; it never migrates old SIDs. +#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum PopulationKeyEncoding { + #[default] + LegacyDelimited, + CanonicalLabelsV1, +} +impl PopulationKeyEncoding { + pub fn is_legacy(&self) -> bool { + *self == Self::LegacyDelimited + } +} + #[derive(Debug, Clone, PartialEq, Eq)] pub struct GroupingProjection(Vec); diff --git a/crates/asap_types/src/lib.rs b/crates/asap_types/src/lib.rs index e61e299d..0cccd182 100644 --- a/crates/asap_types/src/lib.rs +++ b/crates/asap_types/src/lib.rs @@ -27,7 +27,7 @@ pub use accuracy::{AccuracyKind, AccuracyProfile}; pub use aggregation_config::*; pub use aggregation_type::AggregationType; pub use enums::*; -pub use grouping_projection::GroupingProjection; +pub use grouping_projection::{GroupingProjection, PopulationKeyEncoding}; pub use key_by_label_names::KeyByLabelNames; pub use monitor_spec::{MonitorFunctional, MonitorSpec}; pub use policy_fingerprint::PolicyFingerprint; diff --git a/crates/asap_types/src/policy_fingerprint.rs b/crates/asap_types/src/policy_fingerprint.rs index 38afc45f..448c7742 100644 --- a/crates/asap_types/src/policy_fingerprint.rs +++ b/crates/asap_types/src/policy_fingerprint.rs @@ -89,6 +89,10 @@ impl PolicyFingerprint { pub fn from_config(cfg: &AggregationConfig) -> Self { let mut buf: Vec = Vec::with_capacity(512); + if !cfg.population_key_encoding.is_legacy() { + // UTF-8 raw metrics cannot alias this nonlegacy domain prefix. + buf.extend_from_slice(b"\xffpopulation-key-canonical-labels-v1\0"); + } // 1. metric name if cfg.derived_input.is_some() { // Raw policies start with UTF-8 metric bytes; 0xff is impossible @@ -285,6 +289,45 @@ mod tests { ) } + #[test] + fn population_encoding_preserves_legacy_wire_and_separates_identity() { + use crate::grouping_projection::PopulationKeyEncoding; + let legacy = cfg( + "m", + AggregationType::Sum, + HashMap::new(), + vec!["host"], + 60, + "", + ); + let wire = serde_json::to_value(&legacy).unwrap(); + assert!(wire.get("population_key_encoding").is_none()); + let decoded: AggregationConfig = serde_json::from_value(wire).unwrap(); + assert!(decoded.population_key_encoding.is_legacy()); + assert_eq!(legacy.policy_fingerprint(), decoded.policy_fingerprint()); + let mut canonical = legacy.clone(); + canonical.population_key_encoding = PopulationKeyEncoding::CanonicalLabelsV1; + assert_ne!(legacy.policy_fingerprint(), canonical.policy_fingerprint()); + let wire = serde_json::to_value(&canonical).unwrap(); + assert_eq!(wire["population_key_encoding"], "canonical_labels_v1"); + let decoded: AggregationConfig = serde_json::from_value(wire).unwrap(); + assert_eq!(decoded.policy_fingerprint(), canonical.policy_fingerprint()); + use crate::traits::SerializableToSink; + let mut sink = canonical.serialize_to_json(); + // Transport wrappers supply the three label projections separately. + sink["groupingLabels"] = serde_json::to_value(&canonical.grouping_labels).unwrap(); + sink["aggregatedLabels"] = + serde_json::to_value(&canonical.aggregated_labels.labels).unwrap(); + sink["rollupLabels"] = serde_json::to_value(&canonical.rollup_labels.labels).unwrap(); + let decoded = AggregationConfig::deserialize_from_json(&sink).unwrap(); + assert_eq!( + decoded.population_key_encoding, + canonical.population_key_encoding + ); + + assert!(serde_json::from_str::("\"canonical_labels_v2\"").is_err()); + } + #[test] fn same_config_yields_same_fingerprint() { let a = cfg( diff --git a/crates/asap_types/src/precompute_plan.rs b/crates/asap_types/src/precompute_plan.rs index ef4621ba..ebc9db31 100644 --- a/crates/asap_types/src/precompute_plan.rs +++ b/crates/asap_types/src/precompute_plan.rs @@ -410,6 +410,11 @@ impl PrecomputePlan { return Err(PrecomputePlanError::UnsupportedIngestEndpoint); } for config in &self.materializations { + if !config.population_key_encoding.is_legacy() { + return Err(PrecomputePlanError::CatalogContract( + "population key encoding is not supported by the installed runtime".into(), + )); + } let Some(derived) = &config.derived_input else { continue; }; @@ -824,6 +829,29 @@ mod source_window_cohort_tests { })) .unwrap() } + #[test] + fn canonical_population_encoding_is_not_yet_installable() { + let mut config = full_window(); + config.slide_interval = config.window_size; + config.population_key_encoding = + crate::grouping_projection::PopulationKeyEncoding::CanonicalLabelsV1; + let envelope = PlanEnvelope { + plan_id: 1, + plan_version: 1, + generated_at_unix_ms: 0, + activation_unix_ms: 0, + expiry_unix_ms: None, + backend_compat: "test".into(), + planner_revision: "test".into(), + capability_snapshot_id: "test".into(), + }; + let error = PrecomputePlan::build_backend_local(envelope, vec![config]).unwrap_err(); + assert!( + error.to_string().contains("population key encoding"), + "{error}" + ); + } + #[test] fn full_sliding_cohort_preserves_explicit_windows_and_identity() { let target = full_window(); diff --git a/crates/asap_types/src/precompute_plan/catalog.rs b/crates/asap_types/src/precompute_plan/catalog.rs index e9dd2a59..201b9a5f 100644 --- a/crates/asap_types/src/precompute_plan/catalog.rs +++ b/crates/asap_types/src/precompute_plan/catalog.rs @@ -87,7 +87,8 @@ impl PrecomputePlan { .map_err(invalid)?; } let expected_projection = config.effective_value_projection(); - if data.partitioning != config.partitioning + if data.population_key_encoding != config.population_key_encoding + || data.partitioning != config.partitioning || data.timestamp_column != config.table_timestamp_column || data.source != expected_source || &data.value_projection != expected_projection diff --git a/crates/asap_types/src/sds.rs b/crates/asap_types/src/sds.rs index 6b313369..0ef82f9d 100644 --- a/crates/asap_types/src/sds.rs +++ b/crates/asap_types/src/sds.rs @@ -842,6 +842,11 @@ pub struct DataDescriptor { pub timestamp_column: Option, pub population_filter_canonical: String, pub group_by_keys: crate::GroupingProjection, + #[serde( + default, + skip_serializing_if = "crate::grouping_projection::PopulationKeyEncoding::is_legacy" + )] + pub population_key_encoding: crate::grouping_projection::PopulationKeyEncoding, #[serde(default, skip_serializing_if = "Option::is_none")] pub partitioning: Option, /// Versioned contract for timestamp interpretation and @@ -905,6 +910,7 @@ impl DataDescriptor { &observation_semantics, None, None, + Default::default(), ); Self { id, @@ -914,6 +920,7 @@ impl DataDescriptor { population_filter_canonical, group_by_keys, partitioning: None, + population_key_encoding: Default::default(), observation_semantics, } } @@ -927,6 +934,7 @@ impl DataDescriptor { &self.observation_semantics, self.partitioning, self.timestamp_column.as_deref(), + self.population_key_encoding, ); self } @@ -940,9 +948,18 @@ impl DataDescriptor { &self.observation_semantics, partitioning, self.timestamp_column.as_deref(), + self.population_key_encoding, ); self } + pub fn with_population_key_encoding( + mut self, + encoding: crate::grouping_projection::PopulationKeyEncoding, + ) -> Self { + self.population_key_encoding = encoding; + let partitioning = self.partitioning; + self.with_partitioning(partitioning) + } pub fn with_timestamp_column(mut self, column: Option) -> Self { self.timestamp_column = column; let partitioning = self.partitioning; @@ -987,6 +1004,7 @@ impl DataDescriptor { &self.observation_semantics, self.partitioning, self.timestamp_column.as_deref(), + self.population_key_encoding, ) { return Err(SdsError("data descriptor ID/content mismatch".into())); @@ -994,6 +1012,7 @@ impl DataDescriptor { Ok(()) } } +#[allow(clippy::too_many_arguments)] fn data_descriptor_id( source: &DataSourceIdentity, value_projection: &ValueProjectionIdentity, @@ -1002,6 +1021,7 @@ fn data_descriptor_id( observation_semantics: &str, partitioning: Option, timestamp_column: Option<&str>, + population_key_encoding: crate::grouping_projection::PopulationKeyEncoding, ) -> DataDescriptorId { // Length framing keeps distinct typed sources, projections, predicates, // and grouping keys collision-free in the content identity. @@ -1032,6 +1052,9 @@ fn data_descriptor_id( "|{}:{observation_semantics}", observation_semantics.len() )); + if !population_key_encoding.is_legacy() { + key = format!("data:population-key:canonical-labels-v1|{key}"); + } DataDescriptorId(key) } @@ -1040,6 +1063,24 @@ mod tests { use super::*; use std::collections::BTreeSet; + #[test] + fn population_encoding_is_part_of_data_identity_with_legacy_default() { + use crate::grouping_projection::PopulationKeyEncoding; + let legacy = DataDescriptor::new("m", "", ["host".to_string()]); + let wire = serde_json::to_value(&legacy).unwrap(); + assert!(wire.get("population_key_encoding").is_none()); + let decoded: DataDescriptor = serde_json::from_value(wire).unwrap(); + assert_eq!(legacy.id(), decoded.id()); + let canonical = legacy + .clone() + .with_population_key_encoding(PopulationKeyEncoding::CanonicalLabelsV1); + assert_ne!(legacy.id(), canonical.id()); + canonical.validate().unwrap(); + let reverted = + canonical.with_population_key_encoding(PopulationKeyEncoding::LegacyDelimited); + assert_eq!(legacy.id(), reverted.id()); + } + fn completion(epoch: u64) -> SummaryWindowCompletion { SummaryWindowCompletion { catalog_generation: CatalogGeneration { diff --git a/crates/asap_types/src/summary_catalog.rs b/crates/asap_types/src/summary_catalog.rs index d6d404b0..ce2c7e48 100644 --- a/crates/asap_types/src/summary_catalog.rs +++ b/crates/asap_types/src/summary_catalog.rs @@ -120,6 +120,7 @@ impl SummaryCatalog { ) .with_grouping_projection(config.grouping_labels.clone()) .with_partitioning(config.partitioning) + .with_population_key_encoding(config.population_key_encoding) .with_timestamp_column(config.table_timestamp_column.clone()); Ok(( config.policy_fingerprint(), diff --git a/data_plane/src/drivers/ingest/prometheus_remote_write.rs b/data_plane/src/drivers/ingest/prometheus_remote_write.rs index 4fc18d2a..ea06ebb8 100644 --- a/data_plane/src/drivers/ingest/prometheus_remote_write.rs +++ b/data_plane/src/drivers/ingest/prometheus_remote_write.rs @@ -1013,6 +1013,7 @@ mod tests { use asap_types::enums::WindowKind; use asap_types::{AggregationConfig, AggregationType, KeyByLabelNames}; let aggregation = AggregationConfig { + population_key_encoding: Default::default(), aggregation_type: AggregationType::Sum, aggregation_sub_type: String::new(), parameters: HashMap::new(), @@ -1075,6 +1076,7 @@ mod tests { let config = |aggregation_type, grouping: Vec, aggregated: Vec| AggregationConfig { + population_key_encoding: Default::default(), aggregation_type, aggregation_sub_type: String::new(), parameters: match aggregation_type { diff --git a/data_plane/src/drivers/query/servers/http.rs b/data_plane/src/drivers/query/servers/http.rs index 914ec4c3..4071ffef 100644 --- a/data_plane/src/drivers/query/servers/http.rs +++ b/data_plane/src/drivers/query/servers/http.rs @@ -3704,6 +3704,7 @@ aggregations: for marker in active_agg_ids { let metric = format!("metric_{marker}"); let cfg = AggregationConfig { + population_key_encoding: Default::default(), aggregation_type: AggregationType::Sum, aggregation_sub_type: String::new(), parameters: HashMap::new(), diff --git a/data_plane/src/precompute_engine/output_sink.rs b/data_plane/src/precompute_engine/output_sink.rs index baf2684b..538940d8 100644 --- a/data_plane/src/precompute_engine/output_sink.rs +++ b/data_plane/src/precompute_engine/output_sink.rs @@ -404,6 +404,7 @@ mod tests { // via `PolicyFingerprint::from_config`. Callers obtain the id // via `config.policy_fp_u64()`. AggregationConfig { + population_key_encoding: Default::default(), aggregation_type: AggregationType::Sum, aggregation_sub_type: String::new(), parameters: HashMap::new(), diff --git a/data_plane/src/storage_engines/sketch_db/lifecycle/eviction.rs b/data_plane/src/storage_engines/sketch_db/lifecycle/eviction.rs index 30713e75..469db04a 100644 --- a/data_plane/src/storage_engines/sketch_db/lifecycle/eviction.rs +++ b/data_plane/src/storage_engines/sketch_db/lifecycle/eviction.rs @@ -257,6 +257,7 @@ mod tests { fn sum_agg_config(id: u64) -> AggregationConfig { AggregationConfig { + population_key_encoding: Default::default(), aggregation_type: AggregationType::Sum, aggregation_sub_type: String::new(), parameters: HashMap::new(), diff --git a/data_plane/src/storage_engines/types/streaming_config.rs b/data_plane/src/storage_engines/types/streaming_config.rs index 2298955d..5c8b9ff4 100644 --- a/data_plane/src/storage_engines/types/streaming_config.rs +++ b/data_plane/src/storage_engines/types/streaming_config.rs @@ -133,6 +133,11 @@ impl StreamingConfig { num_aggregates_to_retain, QueryLanguage::PromQl, )?; + if !config.population_key_encoding.is_legacy() { + anyhow::bail!( + "legacy streaming input does not support this population key encoding" + ); + } if config.derived_input.is_some() { anyhow::bail!( "legacy streaming input cannot execute a derived summary program" @@ -207,6 +212,33 @@ mod tests { .contains("legacy streaming input cannot execute")); } + #[test] + fn legacy_yaml_rejects_canonical_population_key_encoding() { + let data: Value = serde_yaml::from_str( + r#" +aggregations: +- aggregationType: Sum + aggregationSubType: '' + metric: m + population_key_encoding: canonical_labels_v1 + labels: + grouping: [host] + rollup: [] + aggregated: [] + parameters: {} + windowSize: 60 + windowType: tumbling + spatialFilter: '' +"#, + ) + .unwrap(); + let error = StreamingConfig::from_yaml_data(&data).unwrap_err(); + assert!( + error.to_string().contains("population key encoding"), + "{error}" + ); + } + /// PR 5: a streaming-config YAML that omits `aggregationId` /// parses correctly — the backend derives identity from content /// via `PolicyFingerprint::from_config`. The map key is the diff --git a/data_plane/src/tests/test_utilities/engine_factories.rs b/data_plane/src/tests/test_utilities/engine_factories.rs index 564d21d7..dd88cab6 100644 --- a/data_plane/src/tests/test_utilities/engine_factories.rs +++ b/data_plane/src/tests/test_utilities/engine_factories.rs @@ -90,6 +90,7 @@ pub fn create_engine_single_pop_with_aggregated( let mut aggregation_configs = HashMap::new(); let agg_config = AggregationConfig { + population_key_encoding: Default::default(), aggregation_type, aggregation_sub_type: String::new(), parameters: HashMap::new(), @@ -188,6 +189,7 @@ pub fn create_engine_dual_input( // Value aggregation let value_agg_config = AggregationConfig { + population_key_encoding: Default::default(), aggregation_type: value_agg_type, aggregation_sub_type: String::new(), parameters: HashMap::new(), @@ -216,6 +218,7 @@ pub fn create_engine_dual_input( // Keys aggregation let keys_agg_config = AggregationConfig { + population_key_encoding: Default::default(), aggregation_type: key_agg_type, aggregation_sub_type: String::new(), parameters: HashMap::new(), @@ -309,6 +312,7 @@ pub fn create_engine_two_metrics( let mut aggregation_configs = HashMap::new(); let agg_config_a = AggregationConfig { + population_key_encoding: Default::default(), aggregation_type: aggregation_type_a, aggregation_sub_type: String::new(), parameters: HashMap::new(), @@ -336,6 +340,7 @@ pub fn create_engine_two_metrics( aggregation_configs.insert(id_a, agg_config_a); let agg_config_b = AggregationConfig { + population_key_encoding: Default::default(), aggregation_type: aggregation_type_b, aggregation_sub_type: String::new(), parameters: HashMap::new(), @@ -439,6 +444,7 @@ pub fn create_engine_three_metrics( (aggregation_type_c, &labels_c, metric_c), ] { let cfg = AggregationConfig { + population_key_encoding: Default::default(), aggregation_type: agg_type, aggregation_sub_type: String::new(), parameters: HashMap::new(), @@ -519,6 +525,7 @@ pub fn create_engine_multi_timestamp( let mut aggregation_configs = HashMap::new(); let agg_config = AggregationConfig { + population_key_encoding: Default::default(), aggregation_type, aggregation_sub_type: String::new(), parameters: HashMap::new(), @@ -591,6 +598,7 @@ pub fn create_engine_multi_timestamp_with_window( let mut aggregation_configs = HashMap::new(); let agg_config = AggregationConfig { + population_key_encoding: Default::default(), aggregation_type, aggregation_sub_type: String::new(), parameters: HashMap::new(),