Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
47 changes: 15 additions & 32 deletions controller/src/emit/asapquery_backend.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,35 +10,25 @@
//! `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;

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
Expand Down Expand Up @@ -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),
Expand Down Expand Up @@ -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]
Expand Down
17 changes: 10 additions & 7 deletions data_plane/src/query_engines/asap_query_engine/engine.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};
Expand Down Expand Up @@ -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<u64> = snap_after.aggregation_configs.keys().copied().collect();
assert_eq!(new_ids.len(), 1);
let id = new_ids[0];
Expand All @@ -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"
);
}

Expand Down