From f3fde6122fa672a0a18d71ce18de6b225f17e896 Mon Sep 17 00:00:00 2001 From: zz_y Date: Wed, 13 May 2026 21:25:18 -0600 Subject: [PATCH] feat(sid): derive policy_fp at OTel sketch ingest via content lookup MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Wires up the OTel sketch-ingest path to populate `SketchInstanceMetadata.policy_fp` instead of leaving it `UNSET`. Sketch-backed sids now participate in the `policy_fp → {sids}` reverse index (PR #203), so the analyzer's `find_matching_policies` (PR #204) returns fingerprints whose sids are O(1)-reachable. ## What - New `control_plane::warm_tier_analysis::find_policy_by_content( registry, metric, group_by_keys, agg_type, expected_params)` — finds the policy whose contents match a freshly-ingested sketch's shape. Returns `Some(fp)` on a unique match, `None` on zero or multiple matches (ambiguous → stay UNSET). - Two data-plane helpers in `data_plane/src/drivers/ingest/otel.rs`: - `aggregation_type_for_sketch_handle(SketchKindHandle) → Option` - `sketch_config_to_params(&SketchConfig) → HashMap` Both lock in the data-plane → control-plane wire-shape mapping so drift surfaces as test failures, not silent lookup misses. - New `derive_sketch_policy_fp(ingest_state, metric, kind, cfg, group_by_keys)` threads the content match. Called from the OTel sketch ingest registration site (`route_modified_otlp_sketches_to_precompute`). Returns `PolicyFingerprint::UNSET` when no policy matches — sids stay reachable via the legacy `instances_matching(metric, gbk)` walk. ## Matching shape Match requires all of: 1. `policy.metric == metric` 2. `policy.aggregation_type == agg_type` (mapped from `SketchKindHandle`) 3. `policy.grouping_labels` (as a set) == `group_by_keys` 4. For every key in `expected_params`, `policy.parameters` has the same `serde_json::Value` (extra policy params tolerated) 5. `policy.spatial_filter_normalized.is_empty()` — OTLP sketches don't carry a filter context Ambiguous match (multiple policies → same shape) returns `None` intentionally. Such policies would have collided on sid identity anyway — surfacing as UNSET is the honest signal of a control-plane bug. ## Param-key vocabulary | `SketchConfig` | params key(s) | |---|---| | `DDSketch { relative_accuracy }` | `relative_accuracy` | | `Kll { k }` | `k` | | `Hll { precision }` | `precision` | | `CountSketch { rows, cols }` | `rows`, `cols` | | `CountMin { rows, cols }` | `rows`, `cols` | Names must stay in sync with `AggregationConfig::from_yaml_data` in `asap_types/src/aggregation_config.rs`. The new test `sketch_config_to_params_uses_canonical_keys` locks them. ## What's still standing - OTel sketch-ingest doesn't surface a spatial-filter context, so filtered policies remain unreachable from this path. When the agent's processor emits a filter shape in the OTLP DP, plumb it through to the lookup. - Keyed-group ExactAgg (`MultipleSum` / `MultipleIncrease`) policies don't appear in `policy_capability` either; if/when those land an L4 binder, both this lookup and `find_matching_policies` need matching arms. ## Test plan - [x] 2 new unit tests in `otel::policy_fp_lookup_tests` covering the `SketchKindHandle → AggregationType` mapping and the `SketchConfig → params` rendering - [x] `cargo check --workspace` clean - [x] `cargo test --workspace --lib --bins` green (pre-existing `avg_finds_sum_and_count` HashMap-iteration flake unchanged) 🤖 Generated with [Claude Code](https://claude.com/claude-code) Co-Authored-By: Claude Opus 4.7 (1M context) --- control_plane/src/warm_tier_analysis.rs | 72 ++++++++ data_plane/src/drivers/ingest/otel.rs | 211 +++++++++++++++++++++++- 2 files changed, 275 insertions(+), 8 deletions(-) diff --git a/control_plane/src/warm_tier_analysis.rs b/control_plane/src/warm_tier_analysis.rs index 1200e9b5..adfad022 100644 --- a/control_plane/src/warm_tier_analysis.rs +++ b/control_plane/src/warm_tier_analysis.rs @@ -435,6 +435,78 @@ pub fn policy_capability(cfg: &asap_types::AggregationConfig) -> Option, + agg_type: promql_utilities::query_logics::enums::AggregationType, + expected_params: &std::collections::HashMap, +) -> Option { + let mut hit: Option = None; + for (fp, cfg) in registry.iter() { + if cfg.metric != metric { + continue; + } + if cfg.aggregation_type != agg_type { + continue; + } + let policy_keys: BTreeSet = + cfg.grouping_labels.labels.iter().cloned().collect(); + if &policy_keys != group_by_keys { + continue; + } + if !cfg.spatial_filter_normalized.is_empty() { + continue; + } + // Param subset match — every key the caller named must appear + // in policy.parameters with an equal value. We don't require + // the reverse direction (policy may have extra params the DP + // didn't surface). + let params_ok = expected_params + .iter() + .all(|(k, v)| cfg.parameters.get(k).is_some_and(|pv| pv == v)); + if !params_ok { + continue; + } + // Track unique-match invariant. + if hit.is_some() { + // Ambiguous — multiple policies match the same shape. Skip. + return None; + } + hit = Some(*fp); + } + hit +} + /// Find every policy in `registry` whose contents satisfy `candidate`. /// The result is empty when no policy fits — caller routes the query /// to the archive engine (cold tier) in that case. Multiple matches diff --git a/data_plane/src/drivers/ingest/otel.rs b/data_plane/src/drivers/ingest/otel.rs index 080038d0..49cc613f 100644 --- a/data_plane/src/drivers/ingest/otel.rs +++ b/data_plane/src/drivers/ingest/otel.rs @@ -1052,6 +1052,27 @@ async fn route_modified_otlp_sketches_to_precompute( let group_by_keys: BTreeSet = dp.attrs.keys().cloned().collect(); let cfg = dp.container_config.clone(); + // Derive the policy fingerprint by content- + // matching the OTLP DP's shape against the + // streaming-config registry. Sketches arrive + // with `(kind, config)` embedded but no policy + // reference; we find the policy whose contents + // produce the same shape. Lookup returns + // `Some(fp)` on a unique match, `None` + // when zero policies match (sketch ingested + // before the streaming-config caught up) or + // when multiple policies match the same shape + // (would have been a sid-collision bug — + // surfaces as an UNSET registration so the + // legacy `instances_matching` walk still + // covers it). + let policy_fp = derive_sketch_policy_fp( + ingest_state, + &metric.name, + kind, + &cfg, + &group_by_keys, + ); ingest_state.sketch_index.register(SketchInstanceMetadata { sid, metric_name: metric.name.clone(), @@ -1066,14 +1087,7 @@ async fn route_modified_otlp_sketches_to_precompute( first_seen_unix_ms: ts_ms, retired_at_ms: None, expires_at_ms: None, - // OTel sketch ingest path doesn't have a - // source `AggregationConfig` here — sketches - // arrive with their shape (kind + config) - // embedded in the OTLP DP, not a policy - // reference. Leave UNSET; the reverse - // index skips these. Sketch sids stay - // reachable via `instances_matching`. - policy_fp: asap_types::PolicyFingerprint::UNSET, + policy_fp, }); } @@ -1237,6 +1251,111 @@ async fn route_modified_otlp_sketches_to_precompute( } } +/// Map `SketchKindHandle` to the corresponding wire-format +/// `AggregationType`. Inverse direction is in +/// `sketch_kind_handle_for` above. Used by +/// [`derive_sketch_policy_fp`] to find the policy whose +/// `AggregationConfig.aggregation_type` matches a freshly-ingested +/// sketch. +/// +/// `Any` is a control-plane analysis-time wildcard — it doesn't +/// appear on the ingest path. Returns `None` so the policy lookup +/// fails the (rare) defensive path explicitly. +fn aggregation_type_for_sketch_handle( + handle: crate::storage_engines::sketch_db::index::SketchKindHandle, +) -> Option { + use crate::storage_engines::sketch_db::index::SketchKindHandle; + use promql_utilities::query_logics::enums::AggregationType; + match handle { + SketchKindHandle::DDSketch => Some(AggregationType::DDSketch), + SketchKindHandle::Kll => Some(AggregationType::DatasketchesKLL), + SketchKindHandle::Hll => Some(AggregationType::HLL), + SketchKindHandle::CountSketch => Some(AggregationType::CountSketch), + SketchKindHandle::CountSketchWithHeap => Some(AggregationType::CountSketch), + SketchKindHandle::CountMin => Some(AggregationType::CountMinSketch), + SketchKindHandle::CmsWithHeap => Some(AggregationType::CountMinSketchWithHeap), + SketchKindHandle::Any => None, + } +} + +/// Render a `SketchConfig` into the param map the streaming-config +/// stores. The control plane authors these as +/// `parameters: {: }` JSON; the data plane has the +/// parameters typed in `SketchConfig`. This function converts. +/// +/// Keys MUST match what the control plane emits (see +/// `crates/asap_types/src/aggregation_config.rs::from_yaml_data` for +/// the canonical names). Drift here surfaces as policy lookups that +/// silently miss. +fn sketch_config_to_params( + cfg: &crate::storage_engines::sketch_db::data::SketchConfig, +) -> std::collections::HashMap { + use crate::storage_engines::sketch_db::data::SketchConfig; + let mut params = std::collections::HashMap::new(); + match cfg { + SketchConfig::DDSketch { relative_accuracy } => { + params.insert( + "relative_accuracy".to_string(), + serde_json::json!(*relative_accuracy), + ); + } + SketchConfig::Kll { k } => { + params.insert("k".to_string(), serde_json::json!(*k)); + } + SketchConfig::Hll { precision } => { + params.insert("precision".to_string(), serde_json::json!(*precision)); + } + SketchConfig::CountSketch { rows, cols } + | SketchConfig::CountMin { rows, cols } => { + params.insert("rows".to_string(), serde_json::json!(*rows)); + params.insert("cols".to_string(), serde_json::json!(*cols)); + } + } + params +} + +/// Look up the policy fingerprint for a freshly-ingested OTLP sketch +/// by content-matching against the streaming-config registry. +/// +/// Sketches arrive with `(metric, attrs, sketch_kind, sketch_config)` +/// embedded in the DP but no policy reference. The matching pass: +/// snapshots the current streaming config, derives a +/// `PolicyRegistry`, and asks `find_policy_by_content` for the +/// fingerprint of a policy whose contents match. Returns +/// `PolicyFingerprint::UNSET` when: +/// 1. The `SketchKindHandle::Any` wildcard reached this path +/// (defensive — shouldn't happen). +/// 2. No policy in the registry matches. +/// 3. Multiple policies match (would-have-been-a-bug case; +/// `find_policy_by_content` returns `None` on ambiguity). +/// +/// Callers register the sid with the returned fp regardless of +/// success — UNSET sids are simply absent from the policy_fp → +/// {sids} reverse index, and remain reachable via the legacy +/// `instances_matching(metric, gbk)` walk. +fn derive_sketch_policy_fp( + ingest_state: &IngestState, + metric: &str, + kind: crate::storage_engines::sketch_db::index::SketchKindHandle, + cfg: &crate::storage_engines::sketch_db::data::SketchConfig, + group_by_keys: &std::collections::BTreeSet, +) -> asap_types::PolicyFingerprint { + let Some(agg_type) = aggregation_type_for_sketch_handle(kind) else { + return asap_types::PolicyFingerprint::UNSET; + }; + let params = sketch_config_to_params(cfg); + let snap = ingest_state.config_snapshot(); + let registry = snap.policy_registry(); + control_plane::warm_tier_analysis::find_policy_by_content( + ®istry, + metric, + group_by_keys, + agg_type, + ¶ms, + ) + .unwrap_or(asap_types::PolicyFingerprint::UNSET) +} + /// Phase 5 helper — map a `ModifiedOtlpSketchDp` to the matching /// `SketchKindHandle` so registration and capability classification /// share one source of truth. @@ -1898,6 +2017,82 @@ fn attributes_to_map( m } +#[cfg(test)] +mod policy_fp_lookup_tests { + use super::*; + use crate::storage_engines::sketch_db::data::SketchConfig; + use crate::storage_engines::sketch_db::index::SketchKindHandle; + use promql_utilities::query_logics::enums::AggregationType; + + #[test] + fn handle_to_agg_type_round_trips_canonical_kinds() { + // Locks in the data-plane → control-plane name mapping. + // Drift surfaces as policy lookups that silently miss because + // the handle resolves to an `AggregationType` no policy uses. + assert_eq!( + aggregation_type_for_sketch_handle(SketchKindHandle::DDSketch), + Some(AggregationType::DDSketch) + ); + assert_eq!( + aggregation_type_for_sketch_handle(SketchKindHandle::Kll), + Some(AggregationType::DatasketchesKLL) + ); + assert_eq!( + aggregation_type_for_sketch_handle(SketchKindHandle::Hll), + Some(AggregationType::HLL) + ); + assert_eq!( + aggregation_type_for_sketch_handle(SketchKindHandle::CountMin), + Some(AggregationType::CountMinSketch) + ); + assert_eq!( + aggregation_type_for_sketch_handle(SketchKindHandle::CmsWithHeap), + Some(AggregationType::CountMinSketchWithHeap) + ); + assert_eq!( + aggregation_type_for_sketch_handle(SketchKindHandle::CountSketch), + Some(AggregationType::CountSketch) + ); + // `Any` is a control-plane wildcard, not a real DP shape. + assert_eq!( + aggregation_type_for_sketch_handle(SketchKindHandle::Any), + None + ); + } + + #[test] + fn sketch_config_to_params_uses_canonical_keys() { + // The param-name vocabulary must match what the control plane + // writes in streaming-config YAML (see + // `asap_types::aggregation_config::AggregationConfig::from_yaml_data`). + // Drift surfaces as `find_policy_by_content` missing matches. + let dd = sketch_config_to_params(&SketchConfig::DDSketch { + relative_accuracy: 0.01, + }); + assert_eq!(dd.get("relative_accuracy"), Some(&serde_json::json!(0.01))); + + let kll = sketch_config_to_params(&SketchConfig::Kll { k: 200 }); + assert_eq!(kll.get("k"), Some(&serde_json::json!(200))); + + let hll = sketch_config_to_params(&SketchConfig::Hll { precision: 14 }); + assert_eq!(hll.get("precision"), Some(&serde_json::json!(14))); + + let cs = sketch_config_to_params(&SketchConfig::CountSketch { + rows: 4, + cols: 256, + }); + assert_eq!(cs.get("rows"), Some(&serde_json::json!(4))); + assert_eq!(cs.get("cols"), Some(&serde_json::json!(256))); + + let cm = sketch_config_to_params(&SketchConfig::CountMin { + rows: 4, + cols: 256, + }); + assert_eq!(cm.get("rows"), Some(&serde_json::json!(4))); + assert_eq!(cm.get("cols"), Some(&serde_json::json!(256))); + } +} + #[cfg(test)] mod dispatcher_tests { use super::*;