diff --git a/controller/src/emit/asapquery_backend.rs b/controller/src/emit/asapquery_backend.rs index 61d1d16ef..16eed5549 100644 --- a/controller/src/emit/asapquery_backend.rs +++ b/controller/src/emit/asapquery_backend.rs @@ -10,13 +10,18 @@ //! `processors: { ddsketch_merge: {...} }` + `service.pipelines`. //! * `config::asapquery_backend` (this module) — ASAPQuery-backend //! query engine, expects -//! `aggregations: [{ aggregationId, aggregationType, metric, labels, -//! parameters, windowSize, windowType, spatialFilter }]`. +//! `aggregations: [{ aggregationType, metric, labels, parameters, +//! windowSize, windowType, spatialFilter }]`. //! //! Both are generated from the same `CollectionPlan` fields but target //! different services. The replanner pushes the OTel YAML via OpAMP to //! backend-role collectors and pushes this one via HTTP to the //! ASAPQuery-backend's `/api/v1/streaming-config` endpoint. +//! +//! Phase 5 M2.2: this emitter no longer writes `aggregationId`. The +//! backend's `AggregationConfig::from_yaml_data` derives one +//! deterministically via `compute_agg_config_id` from the same set of +//! fields we emit, so the explicit field is redundant. use std::time::Duration; @@ -24,21 +29,6 @@ use anyhow::{Context, Result}; use crate::types::{AgentCollectorConfig, CollectionPlan, SketchType}; -/// Stable aggregation ID used when the planner has no explicit id to -/// assign. The ASAPQuery-backend uses `u64` agg IDs; we derive one -/// deterministically from the metric name so the same metric always -/// maps to the same id across successive pushes (otherwise the backend -/// would grow unbounded as each replan introduces a new agg_id). -pub fn deterministic_agg_id(metric: &str) -> u64 { - use std::collections::hash_map::DefaultHasher; - use std::hash::{Hash, Hasher}; - let mut h = DefaultHasher::new(); - metric.hash(&mut h); - // Bias away from 0 so the id space is [1, u64::MAX]; 0 is reserved - // in some of the backend's existing test fixtures as a sentinel. - h.finish().saturating_add(1) -} - /// Generate the `StreamingConfig` YAML for the ASAPQuery-backend from a /// single-metric `CollectionPlan`. Produces a one-element `aggregations` /// list — the backend's endpoint will merge this into its active config @@ -73,10 +63,6 @@ pub fn generate_streaming_config_yaml(metric: &str, plan: &CollectionPlan) -> Re let agg_type_str = map_sketch_type_to_agg_type(&agg.sketch_type); let aggregation = serde_yaml::Mapping::from_iter([ - ( - serde_yaml::Value::from("aggregationId"), - serde_yaml::Value::from(deterministic_agg_id(metric)), - ), ( serde_yaml::Value::from("aggregationType"), serde_yaml::Value::from(agg_type_str), @@ -218,18 +204,15 @@ mod tests { } #[test] - fn deterministic_id_is_stable_across_calls() { - assert_eq!( - deterministic_agg_id("cpu_usage"), - deterministic_agg_id("cpu_usage") - ); - assert_ne!( - deterministic_agg_id("cpu_usage"), - deterministic_agg_id("mem_usage") + fn emitted_yaml_omits_aggregation_id() { + let plan = dummy_plan(SketchType::DDSketch); + let yaml = generate_streaming_config_yaml("cpu_usage", &plan).expect("yaml ok"); + let parsed: serde_yaml::Value = serde_yaml::from_str(&yaml).expect("re-parse ok"); + let a = &parsed["aggregations"][0]; + assert!( + a["aggregationId"].is_null(), + "controller must not emit aggregationId — backend derives it from content (M2.2)" ); - // Id is biased away from 0 so test fixtures that use 0 as a - // sentinel don't accidentally collide. - assert_ne!(deterministic_agg_id("any"), 0); } #[test] diff --git a/data_plane/src/query_engines/asap_query_engine/engine.rs b/data_plane/src/query_engines/asap_query_engine/engine.rs index d255b1591..0b29f2d54 100644 --- a/data_plane/src/query_engines/asap_query_engine/engine.rs +++ b/data_plane/src/query_engines/asap_query_engine/engine.rs @@ -5003,9 +5003,12 @@ mod e2e_feedback_loop_tests { // that covers the requested metric. This mirrors DC's // replanner running and POSTing via its BackendClient. let mock = Arc::new(InProcessMockController::new(hot_reload.clone(), |req| { - // Use a deterministic agg_id derived from the metric - // name (same strategy as DC's - // asapquery_backend::deterministic_agg_id from PR #156). + // Mock controller mints an explicit id here just to keep + // the test self-contained. In production, the + // `controller::emit::asapquery_backend` emitter no longer + // writes `aggregationId` (M2.2) and the backend derives + // one via `compute_agg_config_id`; explicit ids in the + // YAML are still honored for backwards compatibility. let id: u64 = { use std::collections::hash_map::DefaultHasher; use std::hash::{Hash, Hasher}; @@ -5092,9 +5095,9 @@ mod e2e_feedback_loop_tests { assert_eq!(recorded[0].statistics, vec![Statistic::Sum]); assert_eq!(recorded[0].data_range_ms, Some(60_000)); - // 9. And the aggregation_id is the deterministic hash the - // planner produced — not a random value. This pins the - // DC #156 deterministic_agg_id contract. + // 9. And the aggregation_id is the deterministic value the + // mock controller produced — not a random one. Explicit + // ids in YAML are still honored after M2.2. let new_ids: Vec = snap_after.aggregation_configs.keys().copied().collect(); assert_eq!(new_ids.len(), 1); let id = new_ids[0]; @@ -5109,7 +5112,7 @@ mod e2e_feedback_loop_tests { }; assert_eq!( id, expected_id, - "deterministic_agg_id contract: same metric → same id" + "explicit YAML aggregationId honored unchanged" ); }