From a996bc72b26e64a716074707bb3c8d4ca9b0a32d Mon Sep 17 00:00:00 2001 From: zz_y Date: Mon, 25 May 2026 17:14:08 -0600 Subject: [PATCH] feat(control_plane): plumb per-metric sample_p from workload to agent emit Add an optional per-metric sketch sampling probability so an operator can turn on the warm-sketch sampling per metric. WorkloadEntry gains a sample_p field (default 1.0, sampling disabled). The planner threads it into EdgeStageConfig.metric_to_sample_p via collect_metric_to_sample_p (mirroring the metric_to_grouping_labels / cumulative_counter_metrics stitch in both the bootstrap and replan paths), and the L5 edge emitter writes a sample_p knob onto the CMS / HLL sketch-processor block in build_edge_processor_block (and the fused asap_edge per-metric entries), guarded so it is only emitted when p < 1. Default 1.0 everywhere keeps the map empty and emits no sample_p key, so the agent config is byte-identical to today when unset; out-of-range values are validated (logged and skipped) at collection and re-guarded at emit. Tests cover the workload field default/parse and that a configured p<1 reaches the emitted agent YAML for CMS + HLL while 1.0/unset emits nothing. This is a static operator-set knob; an optimizer-driven dynamic p is a follow-up. Requires the agent-side sample_p support (CMS/HLL processor + precompute WithSampleP) to be merged first. Co-Authored-By: Claude Opus 4.7 (1M context) --- control_plane/src/emit/mod.rs | 60 +++++++++ control_plane/src/emit/otap.rs | 3 + control_plane/src/emit/stage_config.rs | 114 +++++++++++++++++- control_plane/src/emit/telegraf.rs | 4 + control_plane/src/emit/trait_def.rs | 1 + control_plane/src/main.rs | 11 ++ .../src/physical/colored_dag/emitter.rs | 23 ++++ control_plane/src/replan.rs | 5 + control_plane/src/store/workload.rs | 3 + control_plane/src/workload.rs | 56 +++++++++ 10 files changed, 277 insertions(+), 3 deletions(-) diff --git a/control_plane/src/emit/mod.rs b/control_plane/src/emit/mod.rs index acb9fd0bd..de9512a54 100644 --- a/control_plane/src/emit/mod.rs +++ b/control_plane/src/emit/mod.rs @@ -367,6 +367,62 @@ pub fn collect_metric_to_grouping_labels( out } +/// Sibling of [`collect_metric_to_grouping_labels`]: walk every registry +/// entry and return the per-metric **sampling probability** map the L5 +/// edge emitter drops into [`crate::physical::colored_dag::emitter::EdgeStageConfig::metric_to_sample_p`]. +/// +/// Only metrics whose workload sets a `sample_p` in `(0, 1)` are +/// included — `1.0` (the default / sampling-disabled) and out-of-range +/// values are skipped, so the map stays empty when no metric requests +/// sampling and the emitted agent config (hence the on-wire sketch bytes) +/// is byte-identical to the pre-sampling format. The edge emitter's +/// `insert_sample_p` re-guards the range defensively. +/// +/// As with the sibling collectors, an entry is only honoured when its +/// metric was successfully pre-populated into the workload store, keeping +/// the emit aligned with what the backend knows about. When a metric +/// carries multiple roles the FIRST registered entry's `sample_p` wins +/// (in practice all share it, since the field lives on the WorkloadEntry). +/// +/// A `p <= 0` or `p > 1` value is logged and skipped rather than emitted, +/// so a typo degrades to "no sampling" instead of a mis-scaled sketch. +/// +/// NOTE: this is a static operator-set knob. A dynamic, optimizer-driven +/// `p` (tuned online against an accuracy/bandwidth budget from runtime +/// samples) is a deliberate follow-up and is out of scope here. +pub fn collect_metric_to_sample_p( + registry: &WorkloadRegistry, + workload_store: &WorkloadStore, +) -> std::collections::HashMap { + let mut out = std::collections::HashMap::new(); + for entry in registry.entries() { + if workload_store + .get_all_for_metric(&entry.metric_name) + .into_iter() + .next() + .is_none() + { + continue; + } + let p = entry.sample_p; + if p >= 1.0 { + // Sampling disabled (the default) — emit nothing so the wire + // bytes stay byte-identical. + continue; + } + if p <= 0.0 || !p.is_finite() { + tracing::warn!( + metric = %entry.metric_name, + sample_p = p, + "ignoring out-of-range sample_p (must be in (0, 1]); treating metric as unsampled" + ); + continue; + } + out.entry(entry.metric_name.clone()).or_insert(p); + } + out +} + /// Issue #298 — sibling of [`collect_metric_to_family`] / /// [`collect_metric_to_grouping_labels`]: walk every registry entry and /// return the deduped list of metrics whose workload(s) classify as @@ -488,6 +544,7 @@ mod runtime_tests { cumulative_counter_metrics: Vec::new(), cold_ship_endpoint: None, cold_external_labels: Vec::new(), + metric_to_sample_p: std::collections::HashMap::new(), }; let collector = emit_for_runtime( @@ -525,6 +582,7 @@ mod runtime_tests { cumulative_counter_metrics: Vec::new(), cold_ship_endpoint: None, cold_external_labels: Vec::new(), + metric_to_sample_p: std::collections::HashMap::new(), }; let yaml = emit_for_runtime( AgentRuntime::AsapOtap, @@ -560,6 +618,7 @@ mod runtime_tests { cumulative_counter_metrics: Vec::new(), cold_ship_endpoint: None, cold_external_labels: Vec::new(), + metric_to_sample_p: std::collections::HashMap::new(), }; let toml = emit_for_runtime( AgentRuntime::AsapTelegraf, @@ -898,6 +957,7 @@ mod runtime_tests { cumulative_counter_metrics: Vec::new(), cold_ship_endpoint: None, cold_external_labels: Vec::new(), + metric_to_sample_p: std::collections::HashMap::new(), }; edge_cfg.metric_to_grouping_labels = collect_metric_to_grouping_labels(®istry, &store); diff --git a/control_plane/src/emit/otap.rs b/control_plane/src/emit/otap.rs index 58ac03b25..e5a3213c7 100644 --- a/control_plane/src/emit/otap.rs +++ b/control_plane/src/emit/otap.rs @@ -398,6 +398,7 @@ mod tests { cumulative_counter_metrics: Vec::new(), cold_ship_endpoint: None, cold_external_labels: Vec::new(), + metric_to_sample_p: std::collections::HashMap::new(), } } @@ -416,6 +417,7 @@ mod tests { cumulative_counter_metrics: Vec::new(), cold_ship_endpoint: None, cold_external_labels: Vec::new(), + metric_to_sample_p: std::collections::HashMap::new(), } } @@ -438,6 +440,7 @@ mod tests { cumulative_counter_metrics: Vec::new(), cold_ship_endpoint: None, cold_external_labels: Vec::new(), + metric_to_sample_p: std::collections::HashMap::new(), } } diff --git a/control_plane/src/emit/stage_config.rs b/control_plane/src/emit/stage_config.rs index 4d7aef8c2..109d2f371 100644 --- a/control_plane/src/emit/stage_config.rs +++ b/control_plane/src/emit/stage_config.rs @@ -250,6 +250,9 @@ pub fn emit_edge_yaml( clamp_window_secs(cfg.window_secs), &cfg.label_filters, cfg.source_metric.as_deref(), + cfg.source_metric + .as_deref() + .and_then(|m| cfg.metric_to_sample_p.get(m).copied()), ); // Use the processor_name verbatim as the YAML key — matches the // factory `Type` strings the patched OTel-contrib build registers @@ -1064,10 +1067,19 @@ fn emit_edge_yaml_5sketch_routing( // processor in the routed YAML inherits the unclamped 300s // window from `[5m]` queries. let clamped_window = clamp_window_secs(cfg.window_secs); + // Per-metric sampling probability for the metric routed to this + // family (keyed by the same `metric_name_hint` used above). + let sample_p = metric_name_hint.and_then(|m| cfg.metric_to_sample_p.get(m).copied()); let block = if let Some(sp) = family_to_proc.get(&kind) { - build_edge_processor_block(sp, clamped_window, &cfg.label_filters, metric_name_hint) + build_edge_processor_block( + sp, + clamped_window, + &cfg.label_filters, + metric_name_hint, + sample_p, + ) } else { - build_default_edge_processor_block(&kind, clamped_window, metric_name_hint) + build_default_edge_processor_block(&kind, clamped_window, metric_name_hint, sample_p) }; processors.insert(processor_name.to_string(), block); } @@ -1830,6 +1842,13 @@ fn emit_edge_yaml_asap_edge( } } } + // Per-metric sampling: emit `sample_p` for the sampling-aware + // families (CMS / HLL) only when `p < 1.0`. Mirrors + // `build_edge_processor_block`'s guarded emit so an unset / + // 1.0 probability keeps the fused entry byte-identical. + if matches!(kind, SketchKind::Cms | SketchKind::Hll) { + insert_sample_p(&mut e, cfg.metric_to_sample_p.get(*metric).copied()); + } // Sketch family IS a warm entry → warm signal = true. e.insert( "tier".into(), @@ -2161,6 +2180,7 @@ fn build_default_edge_processor_block( kind: &SketchKind, window_secs: Option, metric_name_hint: Option<&str>, + sample_p: Option, ) -> Value { use crate::sketch_algebra::params::{ CmsParams, CountSketchParams, DDSketchParams, HllParams, KllParams, @@ -2182,7 +2202,7 @@ fn build_default_edge_processor_block( sketch_params: params, aggregation_id: format!("agg_default_{}", sketch_kind_tag(kind)), }; - build_edge_processor_block(&synthetic, window_secs, &[], metric_name_hint) + build_edge_processor_block(&synthetic, window_secs, &[], metric_name_hint, sample_p) } /// Resolve an `ExportTarget` to a concrete `endpoint:port` string. Phase @@ -2220,6 +2240,7 @@ fn build_edge_processor_block( window_secs: Option, label_filters: &[(String, String)], metric_name_hint: Option<&str>, + sample_p: Option, ) -> Value { let mut m = Mapping::new(); @@ -2279,6 +2300,9 @@ fn build_edge_processor_block( // patched build hard-codes p=14); nothing further to set. m.insert("encoding".into(), Value::String("msgpack".into())); m.insert("delta_transmission".into(), Value::Bool(true)); + // Per-metric sampling: HLL's processor honours `sample_p` + // (hash-threshold element sampling in sketchlib-go). + insert_sample_p(&mut m, sample_p); } SketchParams::Cms(p) => { m.insert( @@ -2293,6 +2317,9 @@ fn build_edge_processor_block( m.insert("columns".into(), Value::Number((p.w as u64).into())); m.insert("encoding".into(), Value::String("msgpack".into())); m.insert("delta_transmission".into(), Value::Bool(true)); + // Per-metric sampling: the CMS processor honours `sample_p` + // (geometric admission sampling in sketchlib-go). + insert_sample_p(&mut m, sample_p); } SketchParams::CountSketch(p) => { // Translate (w, d) to the legacy (epsilon, delta) surface @@ -2310,6 +2337,23 @@ fn build_edge_processor_block( Value::Mapping(m) } +/// Write the per-metric `sample_p` knob onto a sketch-processor block, +/// but ONLY when sampling is actually requested (`p < 1.0`). +/// +/// `None` or `p >= 1.0` (the default / disabled state) emits no key, so +/// the agent processor's `Config.Validate` normalises the unset field to +/// `1.0` (sampling disabled) and the emitted YAML — hence the on-wire +/// sketch bytes — stays byte-identical to the pre-sampling format. Values +/// outside `(0, 1]` are dropped here too (the planner validates the range +/// before populating `metric_to_sample_p`, so this is a defensive guard). +fn insert_sample_p(m: &mut Mapping, sample_p: Option) { + if let Some(p) = sample_p { + if p > 0.0 && p < 1.0 { + m.insert("sample_p".into(), Value::Number(p.into())); + } + } +} + /// Compute the gateway-side merge processor name for a `GatewayMergeProcessor`. /// /// Today the typed emitter populates every entry's `processor_name` @@ -2528,6 +2572,7 @@ mod tests { cumulative_counter_metrics: Vec::new(), cold_ship_endpoint: None, cold_external_labels: Vec::new(), + metric_to_sample_p: HashMap::new(), } } @@ -3509,6 +3554,7 @@ mod tests { cumulative_counter_metrics: Vec::new(), cold_ship_endpoint: None, cold_external_labels: Vec::new(), + metric_to_sample_p: HashMap::new(), }; let yaml = emit_edge_yaml(&cfg, "ws://c/", "test-agent").expect("emit ok"); @@ -3856,6 +3902,7 @@ mod tests { cumulative_counter_metrics: Vec::new(), cold_ship_endpoint: None, cold_external_labels: Vec::new(), + metric_to_sample_p: HashMap::new(), } } @@ -3872,6 +3919,64 @@ mod tests { } } + #[test] + fn mvp46_default_emits_no_sample_p() { + // Default fixture (metric_to_sample_p empty) must NOT emit any + // `sample_p` knob — keeps the agent config byte-identical to the + // pre-sampling format when no metric requests sampling. + let _env = crate::test_support::env_lock(); + let cfg = five_sketch_edge_cfg(); + let yaml = emit_edge_yaml(&cfg, "ws://c/", "test-agent").expect("emit ok"); + assert!( + !yaml.contains("sample_p"), + "default (unset) sample_p must not appear in the emitted YAML\n{yaml}" + ); + } + + #[test] + fn mvp46_configured_sample_p_reaches_cms_and_hll_blocks() { + // A configured per-metric `sample_p < 1` must be threaded into the + // emitted agent sketch-processor config for the sampling-aware + // families (CMS / HLL). This is the control-plane half of the + // end-to-end path: workload `sample_p` → EdgeStageConfig. + // metric_to_sample_p → build_edge_processor_block → agent YAML → + // processor Config.SampleP → sketchlib-go WithSampleP. + let _env = crate::test_support::env_lock(); + let mut cfg = five_sketch_edge_cfg(); + // endpoint_request_freq → CMS, unique_users_per_min → HLL. + cfg.metric_to_sample_p + .insert("endpoint_request_freq".into(), 0.1); + cfg.metric_to_sample_p + .insert("unique_users_per_min".into(), 0.25); + let yaml = emit_edge_yaml(&cfg, "ws://c/", "test-agent").expect("emit ok"); + + // The CMS block carries sample_p: 0.1. + assert!( + yaml.contains("sample_p: 0.1"), + "CMS sample_p 0.1 did not reach the emitted YAML\n{yaml}" + ); + // The HLL block carries sample_p: 0.25. + assert!( + yaml.contains("sample_p: 0.25"), + "HLL sample_p 0.25 did not reach the emitted YAML\n{yaml}" + ); + } + + #[test] + fn mvp46_sample_p_of_one_emits_nothing() { + // sample_p == 1.0 is the disabled state — even when present in the + // map it must emit no knob (insert_sample_p guards on `< 1.0`). + let _env = crate::test_support::env_lock(); + let mut cfg = five_sketch_edge_cfg(); + cfg.metric_to_sample_p + .insert("endpoint_request_freq".into(), 1.0); + let yaml = emit_edge_yaml(&cfg, "ws://c/", "test-agent").expect("emit ok"); + assert!( + !yaml.contains("sample_p"), + "sample_p == 1.0 must not be emitted\n{yaml}" + ); + } + #[test] fn mvp46_routing_lives_in_connectors_not_processors() { let _env = crate::test_support::env_lock(); @@ -4147,6 +4252,7 @@ mod tests { cumulative_counter_metrics: Vec::new(), cold_ship_endpoint: None, cold_external_labels: Vec::new(), + metric_to_sample_p: HashMap::new(), } } @@ -5307,6 +5413,7 @@ mod tests { "http://gorilla-merger:10908/ingest/gorilla".into(), ), cold_external_labels: vec![("cluster".into(), "asap-mvp".into())], + metric_to_sample_p: HashMap::new(), } } @@ -5573,6 +5680,7 @@ mod tests { cumulative_counter_metrics: Vec::new(), cold_ship_endpoint: None, cold_external_labels: Vec::new(), + metric_to_sample_p: HashMap::new(), }; let yaml = emit_edge_yaml(&cfg, "ws://c/", "agent-1").expect("emit ok"); diff --git a/control_plane/src/emit/telegraf.rs b/control_plane/src/emit/telegraf.rs index 071b485dd..2b583af24 100644 --- a/control_plane/src/emit/telegraf.rs +++ b/control_plane/src/emit/telegraf.rs @@ -319,6 +319,7 @@ mod tests { cumulative_counter_metrics: Vec::new(), cold_ship_endpoint: None, cold_external_labels: Vec::new(), + metric_to_sample_p: std::collections::HashMap::new(), } } @@ -337,6 +338,7 @@ mod tests { cumulative_counter_metrics: Vec::new(), cold_ship_endpoint: None, cold_external_labels: Vec::new(), + metric_to_sample_p: std::collections::HashMap::new(), } } @@ -359,6 +361,7 @@ mod tests { cumulative_counter_metrics: Vec::new(), cold_ship_endpoint: None, cold_external_labels: Vec::new(), + metric_to_sample_p: std::collections::HashMap::new(), } } @@ -513,6 +516,7 @@ mod tests { cumulative_counter_metrics: Vec::new(), cold_ship_endpoint: None, cold_external_labels: Vec::new(), + metric_to_sample_p: std::collections::HashMap::new(), }; let toml = emit_telegraf_toml(&cfg, None).expect("emit ok"); assert!(toml.contains("k = 200"), "k not propagated\n{toml}"); diff --git a/control_plane/src/emit/trait_def.rs b/control_plane/src/emit/trait_def.rs index 48586fdfa..cdca2fb75 100644 --- a/control_plane/src/emit/trait_def.rs +++ b/control_plane/src/emit/trait_def.rs @@ -212,6 +212,7 @@ mod tests { cumulative_counter_metrics: Vec::new(), cold_ship_endpoint: None, cold_external_labels: Vec::new(), + metric_to_sample_p: HashMap::new(), } } diff --git a/control_plane/src/main.rs b/control_plane/src/main.rs index 51e44d2b5..8b6b538e3 100644 --- a/control_plane/src/main.rs +++ b/control_plane/src/main.rs @@ -615,6 +615,8 @@ async fn handle_plan( sketch_family_override: workload.sketch_type_override.clone(), target_path: None, grouping_labels: workload.group_by_labels.clone(), + // Role derivation does not depend on sampling; default 1.0. + sample_p: 1.0, }; control_plane::workload::derive_agg_role(&entry) }; @@ -1280,6 +1282,15 @@ async fn emit_bootstrap_typed( &st.workload_registry, &st.workload_store, ); + // Per-metric sketch sampling probability — companion stitch: maps + // each metric whose workload set `sample_p < 1` to its probability so + // the L5 edge emitter writes a `sample_p` knob onto the metric's + // CMS / HLL sketch-processor block. Empty when nothing is sampled + // (the default) ⇒ byte-identical agent config. + edge_cfg.metric_to_sample_p = emit::collect_metric_to_sample_p( + &st.workload_registry, + &st.workload_store, + ); // Issue #2: thread X-Agent-ID into the opamp block. Bootstrap GET // is per-agent when `pinned_agent_id` is set (the agent's own diff --git a/control_plane/src/physical/colored_dag/emitter.rs b/control_plane/src/physical/colored_dag/emitter.rs index 92af97f99..75df3a1c2 100644 --- a/control_plane/src/physical/colored_dag/emitter.rs +++ b/control_plane/src/physical/colored_dag/emitter.rs @@ -268,6 +268,28 @@ pub struct EdgeStageConfig { /// `ASAP_CLUSTER`, defaulting to `asap-mvp`). #[serde(default, skip_serializing_if = "Vec::is_empty")] pub cold_external_labels: Vec<(String, String)>, + /// Per-metric sketch **sampling probability** `p` in `(0, 1]`, + /// populated by the planner from each workload entry's + /// [`crate::workload::WorkloadEntry::sample_p`]. + /// + /// The L5 edge emitter reads this in `build_edge_processor_block` and + /// writes a `sample_p:

` knob onto the matching metric's + /// sketch-processor block ONLY when `p < 1.0`. A metric absent from + /// this map (or mapped to `1.0`) emits no `sample_p` key, so the + /// agent's processor `Config.Validate` normalises the unset field to + /// `1.0` (sampling disabled) and the wire bytes stay byte-identical to + /// the pre-sampling format. + /// + /// Activates the warm-sketch sampling layer the agent's sketch + /// processors carry (sketchlib-go geometric / hash-threshold sampling): + /// the encoder admits a `p` fraction of updates, stores the RAW + /// sampled state + `p`, and the backend rescales count-like estimates + /// by `1/p` at query time. This is a static operator-set knob; an + /// optimizer-driven dynamic `p` is a follow-up (out of scope here). + /// + /// Empty map (default) ⇒ no metric carries sampling — backward-compat. + #[serde(default, skip_serializing_if = "HashMap::is_empty")] + pub metric_to_sample_p: HashMap, } /// Named default for [`EdgeStageConfig::cold_ship_endpoint`]. @@ -573,6 +595,7 @@ impl Emitter for ThreeStageEmitter { // post-emit (same pattern as `exporter_target`). cold_ship_endpoint: Some(default_cold_ship_endpoint()), cold_external_labels: default_cold_external_labels(), + metric_to_sample_p: HashMap::new(), }; let mut backend_aggregations: Vec = Vec::new(); let mut gateway_processors: Vec = Vec::new(); diff --git a/control_plane/src/replan.rs b/control_plane/src/replan.rs index b5c7a6674..249b7e0c1 100644 --- a/control_plane/src/replan.rs +++ b/control_plane/src/replan.rs @@ -289,6 +289,11 @@ impl Replanner { // bootstrap YAML did. edge_cfg.cumulative_counter_metrics = crate::emit::collect_cumulative_counter_metrics(registry, &self.workload_store); + // Per-metric sketch sampling probability — companion stitch, + // mirrors the bootstrap path so OpAMP-pushed re-plans carry the + // same `sample_p` knob the first-connect bootstrap YAML did. + edge_cfg.metric_to_sample_p = + crate::emit::collect_metric_to_sample_p(registry, &self.workload_store); } // OpAMP `on_connect` doesn't expose the agent's runtime diff --git a/control_plane/src/store/workload.rs b/control_plane/src/store/workload.rs index 36767d5f0..7c1a4e0fd 100644 --- a/control_plane/src/store/workload.rs +++ b/control_plane/src/store/workload.rs @@ -240,6 +240,7 @@ mod tests { sketch_family_override: None, target_path: None, grouping_labels: vec!["zone".into()], + sample_p: 1.0, }, WorkloadEntry { metric_name: "http_requests_total".into(), @@ -249,6 +250,7 @@ mod tests { sketch_family_override: None, target_path: None, grouping_labels: vec!["zone".into()], + sample_p: 1.0, }, WorkloadEntry { metric_name: "http_requests_total".into(), @@ -258,6 +260,7 @@ mod tests { sketch_family_override: None, target_path: None, grouping_labels: vec!["zone".into()], + sample_p: 1.0, }, ]; diff --git a/control_plane/src/workload.rs b/control_plane/src/workload.rs index 94d373bda..4179df300 100644 --- a/control_plane/src/workload.rs +++ b/control_plane/src/workload.rs @@ -214,11 +214,40 @@ pub struct WorkloadEntry { /// pre-B3 (no allowlist injected). #[serde(default)] pub grouping_labels: Vec, + /// Optional per-metric sketch sampling probability `p` in `(0, 1]`. + /// + /// Activates the warm-sketch sampling layer (geometric admission for + /// the frequency families / hash-threshold for HLL) the agent's + /// sketch processors carry: the encoder admits a `p` fraction of + /// updates, stores the RAW sampled state, stamps `p` on the + /// `SketchEnvelope`, and the backend rescales count-like estimates by + /// `1/p` at query time. `1.0` (the default — also the value `0`/unset + /// normalises to) disables sampling so the emitted agent config and + /// wire bytes are byte-identical to today. + /// + /// Today this is a static operator-set knob; a dynamic + /// optimizer-driven `p` (tuned from runtime samples against an + /// accuracy/bandwidth budget) is a follow-up and is intentionally out + /// of scope here. + /// + /// Threaded into `EdgeStageConfig::metric_to_sample_p` by the registry + /// pre-pop loop in `main`, which the L5 edge emitter reads in + /// `build_edge_processor_block` and writes onto the per-metric + /// sketch-processor block (`sample_p`) only when `< 1.0`. + #[serde(default = "default_sample_p")] + pub sample_p: f64, } fn default_accuracy_sla() -> f64 { 0.01 } + +/// Default per-metric sampling probability — `1.0` (sampling disabled, +/// exact). Keeps the emitted config byte-identical to pre-sampling when a +/// workload entry omits `sample_p`. +fn default_sample_p() -> f64 { + 1.0 +} fn default_role() -> String { "agent".into() } @@ -361,6 +390,7 @@ mod tests { sketch_family_override: None, target_path: None, grouping_labels: vec![], + sample_p: 1.0, }, WorkloadEntry { metric_name: "b".into(), @@ -370,6 +400,7 @@ mod tests { sketch_family_override: None, target_path: None, grouping_labels: vec![], + sample_p: 1.0, }, WorkloadEntry { metric_name: "c".into(), @@ -379,6 +410,7 @@ mod tests { sketch_family_override: None, target_path: None, grouping_labels: vec![], + sample_p: 1.0, }, ], }; @@ -387,6 +419,29 @@ mod tests { assert_eq!(reg.first_for_role("agent").unwrap().metric_name, "a"); } + #[test] + fn deserialize_sample_p_default_and_explicit() { + // sample_p is optional and defaults to 1.0 (sampling disabled); + // an explicit value round-trips. This is the operator-facing + // per-metric sampling knob. + let yaml = r#" +- metric_name: freq_metric + sketch_family_override: CountMinSketch + sample_p: 0.1 +- metric_name: card_metric + sketch_family_override: HLL + sample_p: 0.25 +- metric_name: unset_metric + sketch_family_override: HLL +"#; + let entries: Vec = serde_yaml::from_str(yaml).unwrap(); + assert_eq!(entries.len(), 3); + assert!((entries[0].sample_p - 0.1).abs() < 1e-12); + assert!((entries[1].sample_p - 0.25).abs() < 1e-12); + // Unset ⇒ default 1.0 (sampling disabled / byte-identical). + assert!((entries[2].sample_p - 1.0).abs() < 1e-12); + } + #[test] fn deserialize_sketch_family_override_mixed_case() { // The live wire YAML in `deploy/configs/mvp-workload.yaml` spells @@ -437,6 +492,7 @@ mod tests { sketch_family_override: override_, target_path: None, grouping_labels: vec![], + sample_p: 1.0, } }