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
60 changes: 60 additions & 0 deletions control_plane/src/emit/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<String, f64> {
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
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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(&registry, &store);

Expand Down
3 changes: 3 additions & 0 deletions control_plane/src/emit/otap.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(),
}
}

Expand All @@ -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(),
}
}

Expand All @@ -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(),
}
}

Expand Down
114 changes: 111 additions & 3 deletions control_plane/src/emit/stage_config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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);
}
Expand Down Expand Up @@ -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(),
Expand Down Expand Up @@ -2161,6 +2180,7 @@ fn build_default_edge_processor_block(
kind: &SketchKind,
window_secs: Option<u64>,
metric_name_hint: Option<&str>,
sample_p: Option<f64>,
) -> Value {
use crate::sketch_algebra::params::{
CmsParams, CountSketchParams, DDSketchParams, HllParams, KllParams,
Expand All @@ -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
Expand Down Expand Up @@ -2220,6 +2240,7 @@ fn build_edge_processor_block(
window_secs: Option<u64>,
label_filters: &[(String, String)],
metric_name_hint: Option<&str>,
sample_p: Option<f64>,
) -> Value {
let mut m = Mapping::new();

Expand Down Expand Up @@ -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(
Expand All @@ -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
Expand All @@ -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<f64>) {
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`
Expand Down Expand Up @@ -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(),
}
}

Expand Down Expand Up @@ -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");

Expand Down Expand Up @@ -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(),
}
}

Expand All @@ -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();
Expand Down Expand Up @@ -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(),
}
}

Expand Down Expand Up @@ -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(),
}
}

Expand Down Expand Up @@ -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");
Expand Down
4 changes: 4 additions & 0 deletions control_plane/src/emit/telegraf.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(),
}
}

Expand All @@ -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(),
}
}

Expand All @@ -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(),
}
}

Expand Down Expand Up @@ -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}");
Expand Down
1 change: 1 addition & 0 deletions control_plane/src/emit/trait_def.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(),
}
}

Expand Down
Loading