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}