agent_data_plane_config/domains/dogstatsd.rs
1//! DogStatsD domain: source listeners, parsing, origin detection, aggregation, mapping, filters
2//! (some dynamic-capable), and debug logging.
3
4use std::collections::HashMap;
5use std::fmt;
6use std::num::{NonZeroU64, NonZeroUsize};
7use std::path::PathBuf;
8use std::str::FromStr;
9use std::time::Duration;
10
11use serde::{Deserialize, Serialize};
12
13use crate::defaults::{
14 DEFAULT_AGGREGATE_CONTEXT_LIMIT, DEFAULT_AGGREGATE_FLUSH_INTERVAL,
15 DEFAULT_AGGREGATE_PASSTHROUGH_IDLE_FLUSH_TIMEOUT, DEFAULT_AGGREGATE_WINDOW_DURATION_SECONDS,
16 DEFAULT_DOGSTATSD_ALLOW_CONTEXT_HEAP_ALLOCS, DEFAULT_DOGSTATSD_AUTOSCALE_UDP_LISTENERS,
17 DEFAULT_DOGSTATSD_BUFFER_COUNT, DEFAULT_DOGSTATSD_BUFFER_COUNT_MAX, DEFAULT_DOGSTATSD_CACHED_CONTEXTS_LIMIT,
18 DEFAULT_DOGSTATSD_CACHED_TAGSETS_LIMIT, DEFAULT_DOGSTATSD_MAPPER_STRING_INTERNER_SIZE_BYTES,
19 DEFAULT_DOGSTATSD_MINIMUM_SAMPLE_RATE, DEFAULT_DOGSTATSD_PERMISSIVE_DECODING, DEFAULT_DOGSTATSD_TCP_PORT,
20};
21use crate::Error;
22
23// TODO: better name than Domain? Pipeline? Topology? BlueprintConfig?
24/// Resolved DogStatsD configuration.
25#[derive(Clone, Debug, Default, PartialEq, Serialize)]
26pub struct Domain {
27 /// Source listeners and packet-decoding options.
28 pub listeners: Listeners,
29
30 /// Origin detection and tag cardinality.
31 pub origin: OriginDetection,
32
33 /// Context cache sizing and the sample-rate floor.
34 pub contexts: Contexts,
35
36 /// Metric aggregation window and flush behavior.
37 pub aggregation: Aggregation,
38
39 /// Metric-name mapper.
40 pub mapper: Mapper,
41
42 /// Which payload types are emitted.
43 pub enable_payloads: EnablePayloads,
44
45 /// Metric-name filtering.
46 pub metric_filter: MetricFilter,
47
48 /// Metric namespace prefixing.
49 pub prefix_filter: PrefixFilter,
50
51 /// Per-metric tag include/exclude rules.
52 pub tag_filterlist: Vec<MetricTagFilterEntry>,
53
54 /// Per-metric tag value allow-list rules.
55 pub tag_value_allowlist: Vec<MetricTagValueAllowlistEntry>,
56
57 /// Extra tags added to every metric.
58 pub tags: Vec<String>,
59
60 /// Telemetry emitted by the DogStatsD source.
61 pub telemetry: Telemetry,
62
63 /// Debug-log and verbose-log settings for the DogStatsD source.
64 pub debug_log: DebugLog,
65}
66
67/// Source listeners and packet-decoding options.
68#[derive(Clone, Debug, PartialEq, Serialize)]
69pub struct Listeners {
70 /// UDP port DogStatsD listens on.
71 pub port: u16,
72
73 /// TCP port DogStatsD listens on. (not in Datadog Agent config schema)
74 ///
75 /// Defaults to `0`, which disables TCP.
76 pub tcp_port: u16,
77
78 /// Path of the Unix datagram socket DogStatsD listens on.
79 pub socket: Option<String>,
80
81 /// Path of the Unix stream socket DogStatsD listens on.
82 pub stream_socket: Option<String>,
83
84 /// Windows named pipe name DogStatsD listens on. Unset when no named pipe is configured.
85 pub pipe_name: Option<String>,
86
87 /// SDDL security descriptor applied to the Windows named pipe listener.
88 pub windows_pipe_security_descriptor: String,
89
90 /// Whether the UDP listener accepts traffic from non-local addresses.
91 pub non_local_traffic: bool,
92
93 /// Host the UDP listener binds to.
94 pub bind_host: Option<String>,
95
96 /// Size, in bytes, requested for the socket receive buffer.
97 pub so_rcvbuf: usize,
98
99 /// Size, in bytes, of each packet receive buffer.
100 pub buffer_size: usize,
101
102 /// Number of receive buffers allocated. (not in Datadog Agent config schema)
103 pub buffer_count: usize,
104
105 /// Maximum number of receive buffers. (not in Datadog Agent config schema)
106 pub buffer_count_max: usize,
107
108 /// Number of connectionless packet decoder workers.
109 pub workers_count: usize,
110
111 /// Whether to bind multiple UDP sockets via `SO_REUSEPORT`. (not in Datadog Agent config
112 /// schema)
113 pub autoscale_udp_listeners: bool,
114
115 /// Path a traffic capture is written to or replayed from.
116 pub capture_path: PathBuf,
117
118 /// Maximum number of captured packets queued for persistence.
119 ///
120 /// The capture writer raises a depth below its own minimum, so the default, 0, selects that minimum.
121 pub capture_depth: usize,
122
123 /// End-of-line markers required to terminate a stream-socket message.
124 pub eol_required: Vec<String>,
125
126 /// Whether to log stream messages that exceed the buffer size.
127 pub stream_log_too_big: bool,
128
129 /// Whether to relax decoder strictness on malformed packets. (not in Datadog Agent config
130 /// schema)
131 pub permissive_decoding: bool,
132
133 /// Host that received metrics are additionally forwarded to.
134 pub forward_host: Option<String>,
135
136 /// Port that received metrics are additionally forwarded to.
137 pub forward_port: u16,
138}
139
140impl Default for Listeners {
141 fn default() -> Self {
142 Self {
143 // Saluki-only settings retain these defaults when unset.
144 buffer_count: DEFAULT_DOGSTATSD_BUFFER_COUNT,
145 buffer_count_max: DEFAULT_DOGSTATSD_BUFFER_COUNT_MAX,
146 permissive_decoding: DEFAULT_DOGSTATSD_PERMISSIVE_DECODING,
147 tcp_port: DEFAULT_DOGSTATSD_TCP_PORT,
148 autoscale_udp_listeners: DEFAULT_DOGSTATSD_AUTOSCALE_UDP_LISTENERS,
149 // Witnessed settings are overwritten during translation.
150 port: 0,
151 socket: None,
152 stream_socket: None,
153 pipe_name: None,
154 windows_pipe_security_descriptor: String::new(),
155 non_local_traffic: false,
156 bind_host: None,
157 so_rcvbuf: 0,
158 buffer_size: 0,
159 workers_count: 0,
160 capture_path: PathBuf::new(),
161 capture_depth: 0,
162 eol_required: Vec::new(),
163 stream_log_too_big: false,
164 forward_host: None,
165 forward_port: 0,
166 }
167 }
168}
169
170/// Origin detection and tag cardinality.
171#[derive(Clone, Debug, Default, PartialEq, Serialize)]
172pub struct OriginDetection {
173 /// Whether origin detection tags metrics with their source workload.
174 pub detection: bool,
175
176 /// Whether client-supplied origin information is honored.
177 pub detection_client: bool,
178
179 /// Whether the unified origin-detection scheme is used.
180 pub unified: bool,
181
182 /// Whether a client may opt out of origin detection per metric.
183 pub optout_enabled: bool,
184
185 /// Whether a client-supplied entity ID takes precedence over the detected origin.
186 pub entity_id_precedence: bool,
187
188 /// Tag cardinality applied to origin-detected tags, when a payload doesn't specify one itself.
189 pub tag_cardinality: OriginTagCardinality,
190}
191
192/// Tag cardinality applied during origin detection.
193#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize)]
194pub enum OriginTagCardinality {
195 #[default]
196 Low,
197 Orchestrator,
198 High,
199 None,
200}
201
202impl FromStr for OriginTagCardinality {
203 type Err = Error;
204
205 fn from_str(value: &str) -> Result<Self, Self::Err> {
206 // The Datadog Agent lower-cases the value before matching, and accepts `orch` as a shorthand
207 // for `orchestrator`.
208 match value.to_ascii_lowercase().as_str() {
209 "low" => Ok(Self::Low),
210 "orch" | "orchestrator" => Ok(Self::Orchestrator),
211 "high" => Ok(Self::High),
212 "none" => Ok(Self::None),
213 other => Err(Error::new_without_source(format!(
214 "unknown tag cardinality `{other}`; expected low, orch, orchestrator, high, or none"
215 ))),
216 }
217 }
218}
219
220/// Telemetry emitted by the DogStatsD source.
221#[derive(Clone, Debug, Default, PartialEq, Serialize)]
222pub struct Telemetry {
223 /// Whether processed-metric telemetry is broken down by detected origin.
224 pub origin_breakdown: bool,
225}
226
227/// Context cache sizing and sample-rate floor.
228#[derive(Clone, Debug, PartialEq, Serialize)]
229pub struct Contexts {
230 /// Maximum number of metric contexts held in the cache. (not in Datadog Agent config schema)
231 pub cached_contexts_limit: usize,
232
233 /// Maximum number of tagsets held in the cache. (not in Datadog Agent config schema)
234 pub cached_tagsets_limit: usize,
235
236 /// Number of entries the context string interner holds.
237 pub string_interner_size: u64,
238
239 /// Byte budget for the context string interner, overriding the entry count when set. (not in
240 /// Datadog Agent config schema)
241 pub string_interner_size_bytes: Option<u64>,
242
243 /// Whether contexts may be heap-allocated when the interner is full. (not in Datadog Agent
244 /// config schema)
245 pub allow_context_heap_allocs: bool,
246
247 /// Lowest sample rate honored; the decoder warns and clamps anything below it. (not in Datadog
248 /// Agent config schema)
249 pub minimum_sample_rate: f64,
250}
251
252impl Default for Contexts {
253 fn default() -> Self {
254 Self {
255 // Saluki-only settings retain these defaults when unset.
256 cached_contexts_limit: DEFAULT_DOGSTATSD_CACHED_CONTEXTS_LIMIT,
257 cached_tagsets_limit: DEFAULT_DOGSTATSD_CACHED_TAGSETS_LIMIT,
258 string_interner_size_bytes: None,
259 allow_context_heap_allocs: DEFAULT_DOGSTATSD_ALLOW_CONTEXT_HEAP_ALLOCS,
260 minimum_sample_rate: DEFAULT_DOGSTATSD_MINIMUM_SAMPLE_RATE,
261 // Witnessed settings are overwritten during translation.
262 string_interner_size: 0,
263 }
264 }
265}
266
267/// Metric aggregation window and flush behavior.
268#[derive(Clone, Debug, PartialEq, Serialize)]
269pub struct Aggregation {
270 /// Length, in seconds, of each aggregation window. (not in Datadog Agent config schema)
271 pub window_duration_seconds: NonZeroU64,
272
273 /// Maximum number of contexts held per aggregation window. (not in Datadog Agent config schema)
274 pub context_limit: usize,
275
276 /// How often aggregated metrics are flushed. (not in Datadog Agent config schema)
277 pub flush_interval: Duration,
278
279 /// Whether windows that are still open are flushed on shutdown.
280 ///
281 /// Set by the Datadog `dogstatsd_flush_incomplete_buckets` key.
282 pub flush_open_windows: bool,
283
284 /// How long the no-aggregation passthrough waits before flushing while idle. (not in Datadog
285 /// Agent config schema)
286 pub passthrough_idle_flush_timeout: Duration,
287
288 /// How long, in seconds, a counter value is retained after its last update before expiring.
289 ///
290 /// Set by the Datadog `dogstatsd_expiry_seconds` key. A value of `0` disables zero-value counter
291 /// emission.
292 pub counter_expiry_seconds: Option<u64>,
293
294 /// How long, in seconds, a context is retained after its last update before expiring.
295 pub context_expiry_seconds: u64,
296
297 /// Whether metrics bypass aggregation and are forwarded directly.
298 pub no_aggregation_pipeline: bool,
299
300 /// Capacity of the aggregator's tag-filter result cache.
301 pub aggregator_tag_filter_cache_capacity: usize,
302}
303
304impl Default for Aggregation {
305 fn default() -> Self {
306 Self {
307 // Saluki-schema-only knobs: the Datadog Agent schema does not publish these, so they are
308 // seeded only when set; absent that, these defaults stand.
309 window_duration_seconds: DEFAULT_AGGREGATE_WINDOW_DURATION_SECONDS,
310 context_limit: DEFAULT_AGGREGATE_CONTEXT_LIMIT,
311 flush_interval: DEFAULT_AGGREGATE_FLUSH_INTERVAL,
312 passthrough_idle_flush_timeout: DEFAULT_AGGREGATE_PASSTHROUGH_IDLE_FLUSH_TIMEOUT,
313 // Datadog-schema knobs: always written by the witness driver, so these values are
314 // placeholders that never survive translation.
315 flush_open_windows: false,
316 counter_expiry_seconds: None,
317 context_expiry_seconds: 0,
318 no_aggregation_pipeline: false,
319 aggregator_tag_filter_cache_capacity: 0,
320 }
321 }
322}
323
324/// DogStatsD metric mapper.
325#[derive(Clone, Debug, PartialEq, Serialize)]
326pub struct Mapper {
327 /// Mapper profiles that rewrite matching metric names and tags.
328 pub profiles: Vec<MapperProfile>,
329
330 /// Number of mapper match results cached.
331 pub cache_size: usize,
332
333 /// Byte capacity of the mapper's string interner. (not in Datadog Agent config schema)
334 pub string_interner_size_bytes: NonZeroUsize,
335}
336
337impl Default for Mapper {
338 fn default() -> Self {
339 Self {
340 profiles: Vec::new(),
341 // Written by the Datadog witness driver.
342 cache_size: 0,
343 string_interner_size_bytes: DEFAULT_DOGSTATSD_MAPPER_STRING_INTERNER_SIZE_BYTES,
344 }
345 }
346}
347
348/// One mapper profile: a name, a metric prefix, and the mappings under it.
349#[derive(Clone, Debug, Default, PartialEq, Serialize)]
350pub struct MapperProfile {
351 /// Profile name, for diagnostics.
352 pub name: String,
353
354 /// Metric-name prefix the profile's mappings apply to.
355 pub prefix: String,
356
357 /// The name/tag mappings under this profile.
358 pub mappings: Vec<MetricMapping>,
359}
360
361/// A single metric-name mapping within a [`MapperProfile`].
362#[derive(Clone, Debug, Default, PartialEq, Serialize)]
363pub struct MetricMapping {
364 /// Pattern a metric name must match.
365 pub metric_match: String,
366
367 /// How `metric_match` is interpreted (for example, `wildcard` or `regex`).
368 pub match_type: String,
369
370 /// Replacement name emitted for a matching metric.
371 pub name: String,
372
373 /// Tags added to a matching metric, with values captured from the match.
374 pub tags: HashMap<String, String>,
375}
376
377/// Which payload types are emitted.
378#[derive(Clone, Debug, Default, PartialEq, Serialize)]
379pub struct EnablePayloads {
380 /// Whether event payloads are emitted.
381 pub events: bool,
382
383 /// Whether series (metric) payloads are emitted.
384 pub series: bool,
385
386 /// Whether service-check payloads are emitted.
387 pub service_checks: bool,
388
389 /// Whether sketch (distribution) payloads are emitted.
390 pub sketches: bool,
391}
392
393/// Metric-name filtering (dynamic-capable).
394#[derive(Clone, Debug, Default, PartialEq, Serialize)]
395pub struct MetricFilter {
396 /// Metric names or prefixes that are dropped.
397 ///
398 /// Defaults to an empty list, which disables metric-name filtering.
399 pub values: Vec<String>,
400
401 /// Whether entries match by prefix rather than exact name.
402 ///
403 /// Defaults to `false`, which requires exact matches.
404 pub match_prefix: bool,
405}
406
407/// Metric namespace prefixing.
408#[derive(Clone, Debug, Default, PartialEq, Serialize)]
409pub struct PrefixFilter {
410 /// Namespace prepended to every metric name.
411 pub metric_namespace: String,
412
413 /// Namespaces excluded from the metric-namespace prefixing.
414 pub metric_namespace_blocklist: Vec<String>,
415}
416
417/// One tag-filterlist entry (dynamic-capable).
418#[derive(Clone, Debug, Default, PartialEq, Serialize)]
419pub struct MetricTagFilterEntry {
420 /// Metric name the entry applies to.
421 pub metric_name: String,
422
423 /// Whether the listed tags are included or excluded.
424 pub action: FilterAction,
425
426 /// Tags the action applies to.
427 pub tags: Vec<String>,
428}
429
430/// One tag value allow-list entry.
431///
432/// Rules apply to counters and sketch-backed metrics after mapper rewrites and metric namespace prefixing, and only
433/// to metrics that are aggregated: metrics carrying an explicit client timestamp bypass tag filtering entirely.
434/// Distinct prefixes must not overlap. Multiple rules may use the same prefix when they target different tags.
435#[derive(Clone, Debug, PartialEq, Deserialize, Serialize)]
436pub struct MetricTagValueAllowlistEntry {
437 /// Non-empty metric-name prefix the entry applies to.
438 ///
439 /// Matching is exact and case-sensitive, including any whitespace. Empty prefixes and overlapping distinct
440 /// prefixes are invalid. Multiple rules may use the same prefix when they target different tags.
441 pub metric_prefix: String,
442
443 /// Non-empty tag key whose values are constrained.
444 ///
445 /// Bare tags have no value and are not changed. Key/value tags with an empty value are processed normally. Empty
446 /// names and names containing `:` are invalid. Matching is exact and preserves whitespace.
447 pub tag_name: String,
448
449 /// Tag values retained unchanged.
450 ///
451 /// The default is an empty list, which treats every key/value tag as a mismatch. The empty string is a valid list
452 /// member and retains tags with an empty value. Matching is exact and preserves whitespace.
453 #[serde(default)]
454 pub values: Vec<String>,
455
456 /// Action applied when a tag value is absent from [`values`][Self::values].
457 ///
458 /// The default is [`Remove`][TagValueMismatchAction::Remove].
459 #[serde(default)]
460 pub on_miss: TagValueMismatchAction,
461
462 /// Replacement value used when `on_miss` is [`Replace`][TagValueMismatchAction::Replace].
463 ///
464 /// The default is `other`. This field has no effect when `on_miss` is
465 /// [`Remove`][TagValueMismatchAction::Remove]. The replacement is emitted exactly as configured, including
466 /// whitespace.
467 #[serde(default = "default_tag_value_replacement")]
468 pub replacement: String,
469}
470
471fn default_tag_value_replacement() -> String {
472 "other".to_string()
473}
474
475impl Default for MetricTagValueAllowlistEntry {
476 fn default() -> Self {
477 Self {
478 metric_prefix: String::new(),
479 tag_name: String::new(),
480 values: Vec::new(),
481 on_miss: TagValueMismatchAction::Remove,
482 replacement: default_tag_value_replacement(),
483 }
484 }
485}
486
487/// Reports why a metric tag value allow-list cannot be represented.
488#[derive(Clone, Debug, PartialEq, Eq)]
489pub struct InvalidMetricTagValueAllowlist(String);
490
491impl fmt::Display for InvalidMetricTagValueAllowlist {
492 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
493 f.write_str(&self.0)
494 }
495}
496
497impl std::error::Error for InvalidMetricTagValueAllowlist {}
498
499/// Validates metric tag value allow-list entries.
500///
501/// # Errors
502///
503/// Returns an error for empty prefixes or tag names, tag names containing `:`, overlapping distinct prefixes, or a
504/// duplicate prefix and tag pair.
505pub fn validate_metric_tag_value_allowlists(
506 entries: &[MetricTagValueAllowlistEntry],
507) -> Result<(), InvalidMetricTagValueAllowlist> {
508 for (index, entry) in entries.iter().enumerate() {
509 let rule = index + 1;
510 if entry.metric_prefix.is_empty() {
511 return Err(InvalidMetricTagValueAllowlist(format!(
512 "metric tag value allow-list rule {rule} has an empty `metric_prefix`; configure a non-empty metric-name prefix"
513 )));
514 }
515 if entry.tag_name.is_empty() {
516 return Err(InvalidMetricTagValueAllowlist(format!(
517 "metric tag value allow-list rule {rule} for prefix '{}' has an empty `tag_name`; configure a non-empty tag name",
518 entry.metric_prefix
519 )));
520 }
521 if entry.tag_name.contains(':') {
522 return Err(InvalidMetricTagValueAllowlist(format!(
523 "metric tag value allow-list tag name '{}' contains ':'; configure only the tag key, without a colon or value",
524 entry.tag_name
525 )));
526 }
527 }
528
529 let mut sorted_entries = entries.iter().collect::<Vec<_>>();
530 sorted_entries.sort_unstable_by(|left, right| {
531 left.metric_prefix
532 .cmp(&right.metric_prefix)
533 .then_with(|| left.tag_name.cmp(&right.tag_name))
534 });
535
536 // After sorting by prefix and then tag, duplicate prefix/tag pairs are adjacent. Any distinct prefix that extends
537 // another prefix follows the complete group for the shorter prefix, so one adjacent pair also exposes that overlap.
538 for pair in sorted_entries.windows(2) {
539 let [left, right] = pair else {
540 unreachable!("a two-entry window must contain two entries");
541 };
542 if left.metric_prefix == right.metric_prefix {
543 if left.tag_name == right.tag_name {
544 return Err(InvalidMetricTagValueAllowlist(format!(
545 "metric prefix '{}' is configured more than once for tag '{}'; configure each prefix and tag pair only once",
546 left.metric_prefix, left.tag_name
547 )));
548 }
549 } else if right.metric_prefix.starts_with(&left.metric_prefix) {
550 return Err(InvalidMetricTagValueAllowlist(format!(
551 "overlapping metric prefixes '{}' and '{}' are configured; configure distinct prefixes that do not overlap",
552 left.metric_prefix, right.metric_prefix
553 )));
554 }
555 }
556
557 Ok(())
558}
559
560/// Action applied when a tag value is absent from its allow-list.
561#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Deserialize, Serialize)]
562#[serde(rename_all = "snake_case")]
563pub enum TagValueMismatchAction {
564 /// Removes the tag.
565 #[default]
566 Remove,
567 /// Replaces the tag value with the configured sentinel.
568 Replace,
569}
570
571/// Whether a tag-filterlist entry includes or excludes the listed tags.
572#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize)]
573pub enum FilterAction {
574 Include,
575 #[default]
576 Exclude,
577}
578
579/// DogStatsD debug-log and verbose-log settings.
580#[derive(Clone, Debug, Default, PartialEq, Serialize)]
581pub struct DebugLog {
582 /// Whether the metric debug-log destination is added at startup.
583 ///
584 /// Defaults to `true`. Set this to `false` to remove the destination from the topology.
585 pub logging_enabled: bool,
586
587 /// Path of the DogStatsD metric debug log.
588 ///
589 /// Defaults to `None`, which selects the platform-specific Agent log path at startup.
590 /// Set a path to write the diagnostic log elsewhere.
591 pub log_file: Option<PathBuf>,
592
593 /// Number of rotated debug log files retained.
594 ///
595 /// Defaults to `3`. The file writer retains one rotated file when this is `0`.
596 pub log_file_max_rolls: usize,
597
598 /// Maximum size, in bytes, of the active debug log file before rotation.
599 ///
600 /// Defaults to `10,000,000` bytes. A value of `0` is accepted as the rotation threshold.
601 /// Set this based on the diagnostic data and disk space to retain.
602 pub log_file_max_size: u64,
603
604 /// Whether per-metric processing statistics are written to the debug log.
605 ///
606 /// Defaults to `false` and can change at runtime. Enable this when collecting metric-level diagnostics.
607 pub metrics_stats_enable: bool,
608
609 /// Whether per-packet parse errors are logged at debug rather than error level.
610 ///
611 /// Defaults to `false`. Enable this to reduce error-level log volume from malformed packets.
612 pub disable_verbose_logs: bool,
613}