saluki_components/sources/dogstatsd/
mod.rs

1//! DogStatsD source.
2//!
3//! # Missing
4//!
5//! - Create a health handle for each listener.
6//! - Handle UDS stream framing without treating EOF the same way as UDP and UDS datagram framing.
7//! - Track dispatch failures without depending on whether all events were already iterated.
8
9use std::{
10    collections::VecDeque,
11    future::Future,
12    num::NonZeroUsize,
13    path::PathBuf,
14    pin::Pin,
15    sync::{Arc, LazyLock},
16    time::{Duration, SystemTime, UNIX_EPOCH},
17};
18
19use async_trait::async_trait;
20use bytes::{Buf, BufMut};
21use bytesize::ByteSize;
22use saluki_common::{
23    sync::shutdown::{ShutdownCoordinator, ShutdownHandle},
24    task::spawn_traced_named,
25};
26use saluki_config::{deserialize_space_separated_or_seq, GenericConfiguration};
27use saluki_context::{
28    origin::RawOrigin,
29    tags::{RawTags, RawTagsFilter},
30    TagsResolver,
31};
32use saluki_core::accounting::{MemoryBounds, MemoryBoundsBuilder, UsageExpr};
33use saluki_core::data_model::event::{
34    eventd::EventD,
35    metric::{Metric, MetricMetadata, MetricOrigin},
36    service_check::ServiceCheck,
37    Event, EventType,
38};
39use saluki_core::{
40    components::{sources::*, ComponentContext},
41    pooling::ElasticObjectPool,
42    topology::{interconnect::EventBufferManager, EventsBuffer, OutputDefinition},
43};
44use saluki_env::{workload::CaptureEntityResolver, WorkloadProvider};
45use saluki_error::{generic_error, ErrorContext as _, GenericError};
46use saluki_io::{
47    buf::{BytesBuffer, ClearableIoBuffer as _, FixedSizeVec},
48    deser::{
49        codec::dogstatsd::*,
50        framing::{Framer as _, FramingError, LengthDelimitedFramer},
51    },
52    net::{
53        listener::{Listener, ListenerError},
54        ConnectionAddress, ListenAddress, ProcessCredentials, ProcessIdentity, Stream,
55    },
56};
57use serde::{Deserialize, Deserializer};
58use serde_with::{serde_as, NoneAsEmptyString};
59use snafu::{ResultExt as _, Snafu};
60use stringtheory::MetaString;
61use tokio::{
62    pin, select,
63    time::{interval, MissedTickBehavior},
64};
65use tracing::{debug, error, info, trace, warn};
66
67mod forwarder;
68use self::forwarder::{PacketForwarder, PacketForwarderTarget};
69
70mod framer;
71use self::framer::{get_framer, DsdFramer};
72use crate::sources::dogstatsd::tags::{WellKnownTags, WellKnownTagsFilterPredicate};
73
74mod filters;
75use self::filters::EnablePayloadsFilter;
76
77mod io_buffer;
78use self::io_buffer::IoBufferManager;
79
80mod metrics;
81use self::metrics::{build_metrics, Metrics};
82
83mod replay;
84use self::replay::{CaptureRecord, CapturedTaggerHandle, TrafficCapture};
85pub use self::replay::{
86    DogStatsDCaptureAPIHandler, DogStatsDCaptureControl, DogStatsDReplayAPIHandler, DogStatsDReplayControl,
87    ReplaySession, TimestampResolution, TrafficCaptureReader, DEFAULT_REPLAY_LOOPS, REPLAY_CREDENTIALS_GID,
88};
89
90mod origin;
91use self::origin::{
92    mark_replay_process_id, origin_from_event_packet, origin_from_metric_packet, origin_from_service_check_packet,
93    DogStatsDOriginTagResolver, OriginEnrichmentConfiguration,
94};
95
96mod resolver;
97use self::resolver::ContextResolvers;
98
99mod tags;
100
101#[derive(Debug, Snafu)]
102#[snafu(context(suffix(false)))]
103enum Error {
104    #[snafu(display("Failed to create {} listener: {}", listener_type, source))]
105    FailedToCreateListener {
106        listener_type: &'static str,
107        source: ListenerError,
108    },
109
110    #[snafu(display("No listeners configured. Please specify a port (`dogstatsd_port`) or a socket path (`dogstatsd_socket` or `dogstatsd_stream_socket`) to enable a listener."))]
111    NoListenersConfigured,
112
113    #[snafu(display("Could not resolve bind_host '{}': {}", host, source))]
114    UnresolvableBindHost { host: String, source: std::io::Error },
115
116    #[snafu(display("bind_host '{}' resolved to zero IP addresses.", host))]
117    BindHostHasNoAddresses { host: String },
118}
119
120/// Baseline byte cost per interner entry, used to convert the Core Agent's entry-count-based
121/// `dogstatsd_string_interner_size` to a byte size.
122///
123/// 4096 entries × 512 bytes = 2 MiB, matching ADP's previous default.
124const INTERNER_BASELINE_BYTES_PER_ENTRY: u64 = 512;
125
126const fn default_buffer_size() -> usize {
127    8192
128}
129
130const fn default_buffer_count() -> usize {
131    128
132}
133
134const fn default_buffer_count_max() -> usize {
135    256
136}
137
138const fn default_port() -> u16 {
139    8125
140}
141
142const fn default_tcp_port() -> u16 {
143    0
144}
145
146const fn default_statsd_forward_port() -> u16 {
147    0
148}
149
150const fn default_socket_receive_buffer_size() -> usize {
151    0
152}
153
154const fn default_allow_context_heap_allocations() -> bool {
155    true
156}
157
158const fn default_no_aggregation_pipeline_support() -> bool {
159    true
160}
161
162const fn default_context_string_interner_entry_count() -> u64 {
163    4096
164}
165
166const fn default_cached_contexts_limit() -> usize {
167    500_000
168}
169
170const fn default_cached_tagsets_limit() -> usize {
171    500_000
172}
173
174const fn default_context_expiry_seconds() -> u64 {
175    20
176}
177
178const fn default_dogstatsd_permissive_decoding() -> bool {
179    true
180}
181
182const fn default_dogstatsd_minimum_sample_rate() -> f64 {
183    0.000000003845
184}
185
186const fn default_true() -> bool {
187    true
188}
189
190/// Returns the core Agent default SDDL applied to DogStatsD Windows named pipes.
191const fn default_windows_pipe_security_descriptor() -> &'static str {
192    "D:AI(A;;GA;;;WD)"
193}
194
195fn default_windows_pipe_security_descriptor_string() -> String {
196    default_windows_pipe_security_descriptor().to_string()
197}
198
199/// Controls which payload types are forwarded to the backend.
200#[derive(Deserialize)]
201#[cfg_attr(test, derive(PartialEq, serde::Serialize))]
202pub struct EnablePayloadsConfiguration {
203    /// Whether or not to enable sending series (counter/gauge/rate) payloads.
204    ///
205    /// Defaults to `true`.
206    #[serde(default = "default_true")]
207    pub series: bool,
208
209    /// Whether or not to enable sending sketch (distribution) payloads.
210    ///
211    /// Defaults to `true`.
212    #[serde(default = "default_true")]
213    pub sketches: bool,
214
215    /// Whether or not to enable sending event payloads.
216    ///
217    /// Defaults to `true`.
218    #[serde(default = "default_true")]
219    pub events: bool,
220
221    /// Whether or not to enable sending service check payloads.
222    ///
223    /// Defaults to `true`.
224    #[serde(default = "default_true")]
225    pub service_checks: bool,
226}
227
228impl Default for EnablePayloadsConfiguration {
229    fn default() -> Self {
230        Self {
231            series: true,
232            sketches: true,
233            events: true,
234            service_checks: true,
235        }
236    }
237}
238
239const MIN_CAPTURE_DEPTH: usize = 1024;
240
241const fn default_capture_depth() -> usize {
242    MIN_CAPTURE_DEPTH
243}
244
245const DOGSTATSD_CAPTURE_DIR: &str = "dsd_capture";
246
247fn deserialize_empty_metastring_as_none<'de, D>(deserializer: D) -> Result<Option<MetaString>, D::Error>
248where
249    D: Deserializer<'de>,
250{
251    let value = Option::<MetaString>::deserialize(deserializer)?;
252    Ok(value.filter(|host| !host.is_empty()))
253}
254
255#[derive(Deserialize, Default)]
256#[cfg_attr(test, derive(PartialEq, serde::Serialize))]
257struct DogStatsDTelemetryConfiguration {
258    /// Whether to break down DogStatsD processed-metric telemetry by UDS origin.
259    ///
260    /// When enabled, metric-message `dogstatsd.processed` telemetry includes an `origin` label derived from the
261    /// sender's UDS origin. This can add one telemetry series per origin and should primarily be used for diagnostics.
262    ///
263    /// Defaults to `false`.
264    #[serde(default)]
265    dogstatsd_origin: bool,
266}
267
268/// DogStatsD source.
269///
270/// Accepts metrics over TCP, UDP, or Unix Domain Sockets in the StatsD/DogStatsD format.
271#[serde_as]
272#[derive(Deserialize, Default)]
273#[cfg_attr(test, derive(derive_where::DeriveWhere, serde::Serialize))]
274#[cfg_attr(test, derive_where(PartialEq))]
275pub struct DogStatsDConfiguration {
276    /// Hostname used when DogStatsD metrics do not carry an explicit `host:` tag.
277    #[serde(skip)]
278    default_hostname: MetaString,
279
280    /// The size of the buffer used to receive messages into, in bytes.
281    ///
282    /// Payloads can't exceed this size, or they will be truncated, leading to discarded messages.
283    ///
284    /// Defaults to 8192 bytes.
285    #[serde(rename = "dogstatsd_buffer_size", default = "default_buffer_size")]
286    buffer_size: usize,
287
288    /// The number of message buffers to allocate up front.
289    ///
290    /// This is the baseline pool size allocated at startup. The pool then grows on demand up to
291    /// `dogstatsd_buffer_count_max` as active stream connections need additional buffers. The default value should be
292    /// suitable for the majority of workloads.
293    ///
294    /// Defaults to 128.
295    #[serde(rename = "dogstatsd_buffer_count", default = "default_buffer_count")]
296    buffer_count: usize,
297
298    /// The maximum number of message buffers to allocate overall.
299    ///
300    /// The pool starts at `dogstatsd_buffer_count` buffers and grows on demand up to this limit as active stream
301    /// connections need them, which loosely correlates with how many messages can be received per second. This caps
302    /// memory growth. The default value should be suitable for the majority of workloads, but high-throughput or
303    /// high-fan-in workloads may consider increasing this value.
304    ///
305    /// The pool never holds fewer buffers than `dogstatsd_buffer_count`, so a value below the baseline is treated as
306    /// equal to it.
307    ///
308    /// Defaults to 256, or `dogstatsd_buffer_count` if that is larger.
309    #[serde(rename = "dogstatsd_buffer_count_max", default = "default_buffer_count_max")]
310    buffer_count_max: usize,
311
312    /// The port to listen on in UDP mode.
313    ///
314    /// If set to `0`, UDP isn't used.
315    ///
316    /// Defaults to 8125.
317    #[serde(rename = "dogstatsd_port", default = "default_port")]
318    port: u16,
319
320    /// The size of the DogStatsD UDP/UDS socket receive buffer, in bytes.
321    ///
322    /// If set to `0`, the OS default is used.
323    ///
324    /// Defaults to 0.
325    #[serde(rename = "dogstatsd_so_rcvbuf", default = "default_socket_receive_buffer_size")]
326    socket_receive_buffer_size: usize,
327
328    /// The port to listen on in TCP mode.
329    ///
330    /// If set to `0`, TCP isn't used.
331    ///
332    /// Defaults to 0.
333    #[serde(rename = "dogstatsd_tcp_port", default = "default_tcp_port")]
334    tcp_port: u16,
335
336    /// The host to forward framed DogStatsD messages to over UDP.
337    ///
338    /// Forwarding is enabled only when this value is non-empty and `statsd_forward_port` is non-zero. Setup failures
339    /// are logged, and send failures are tracked through telemetry.
340    ///
341    /// Defaults to unset.
342    #[serde(
343        rename = "statsd_forward_host",
344        default,
345        deserialize_with = "deserialize_empty_metastring_as_none"
346    )]
347    statsd_forward_host: Option<MetaString>,
348
349    /// The port to forward framed DogStatsD messages to over UDP.
350    ///
351    /// Forwarding is enabled only when this value is non-zero and `statsd_forward_host` is non-empty.
352    ///
353    /// Defaults to 0.
354    #[serde(rename = "statsd_forward_port", default = "default_statsd_forward_port")]
355    statsd_forward_port: u16,
356
357    /// The Unix domain socket path to listen on, in datagram mode.
358    ///
359    /// If not set, UDS (in datagram mode) isn't used.
360    ///
361    /// Defaults to unset.
362    #[serde(rename = "dogstatsd_socket", default)]
363    #[serde_as(as = "NoneAsEmptyString")]
364    socket_path: Option<String>,
365
366    /// The Unix domain socket path to listen on, in stream mode.
367    ///
368    /// If not set, UDS (in stream mode) isn't used.
369    ///
370    /// Defaults to unset.
371    #[serde(rename = "dogstatsd_stream_socket", default)]
372    #[serde_as(as = "NoneAsEmptyString")]
373    socket_stream_path: Option<String>,
374
375    /// Controls whether ADP logs oversized DogStatsD stream frames.
376    ///
377    /// When set to `true`, ADP emits a warning when a UDS stream frame exceeds the
378    /// configured DogStatsD buffer size. The frame is still rejected either way.
379    ///
380    /// Enable this when diagnosing clients that send oversized UDS stream frames.
381    ///
382    /// Defaults to `false`.
383    #[serde(rename = "dogstatsd_stream_log_too_big", default)]
384    stream_log_too_big: bool,
385
386    /// The Windows named pipe name to listen on.
387    ///
388    /// If set, ADP listens for DogStatsD stream traffic on `\\.\pipe\<name>` on Windows.
389    /// The listener is unsupported on non-Windows platforms.
390    ///
391    /// Defaults to unset.
392    #[serde(rename = "dogstatsd_pipe_name", default)]
393    #[serde_as(as = "NoneAsEmptyString")]
394    pipe_name: Option<String>,
395
396    /// Windows named pipe security descriptor.
397    ///
398    /// This SDDL descriptor is applied when creating the named pipe listener.
399    ///
400    /// Defaults to `D:AI(A;;GA;;;WD)`.
401    #[serde(
402        rename = "dogstatsd_windows_pipe_security_descriptor",
403        default = "default_windows_pipe_security_descriptor_string"
404    )]
405    windows_pipe_security_descriptor: String,
406
407    /// Whether ADP lowers DogStatsD parse-failure logs to debug level.
408    ///
409    /// When set to `true`, invalid metrics, events, and service checks still increment decode-failure telemetry, but
410    /// their parse-failure logs are emitted at debug level instead of warning level. Enable this to suppress noisy
411    /// parse-error logs from misbehaving clients.
412    ///
413    /// Defaults to `false`.
414    #[serde(rename = "dogstatsd_disable_verbose_logs", default)]
415    disable_verbose_logs: bool,
416
417    /// Listener types that require DogStatsD messages to be newline-terminated.
418    ///
419    /// Valid values are `udp`, `uds`, and `named_pipe`. Invalid values are ignored.
420    ///
421    /// Enable this when DogStatsD clients must reject packets or stream frames that don't end with a newline.
422    ///
423    /// Defaults to unset, which accepts the final message without a newline.
424    #[serde(
425        rename = "dogstatsd_eol_required",
426        default,
427        deserialize_with = "deserialize_space_separated_or_seq"
428    )]
429    eol_required: Vec<String>,
430
431    /// The host address to bind DogStatsD UDP and TCP listeners to.
432    ///
433    /// When set, UDP and TCP listeners bind to this address. Accepts either an IP literal (for example,
434    /// `192.168.1.50`, `::1`) or a hostname that resolves via DNS (for example, `agent.internal`).
435    /// Ignored when `dogstatsd_non_local_traffic` is `true`.
436    ///
437    /// Defaults to unset, which binds to `127.0.0.1`.
438    #[serde(rename = "bind_host", default)]
439    #[serde_as(as = "NoneAsEmptyString")]
440    bind_host: Option<String>,
441
442    /// Whether or not to listen for non-local traffic in UDP mode.
443    ///
444    /// If set to `true`, the listener will accept packets from any interface/address. Otherwise, the source will only
445    /// listen on the address specified by `bind_host`, or `127.0.0.1` if `bind_host` isn't set.
446    ///
447    /// Defaults to `false`.
448    #[serde(rename = "dogstatsd_non_local_traffic", default)]
449    non_local_traffic: bool,
450
451    /// Whether to autoscale UDP stream handlers using `SO_REUSEPORT`.
452    ///
453    /// When enabled on Linux, the DogStatsD source binds multiple UDP sockets to the configured port with
454    /// `SO_REUSEPORT`, allowing the kernel to load-balance incoming datagrams across independent stream handler
455    /// tasks. The number of sockets scales with available vCPUs: one stream handler base, plus one additional
456    /// per 8 vCPUs, capped at 4 total.
457    ///
458    /// Has no effect on non-Linux platforms because `SO_REUSEPORT` doesn't provide kernel-level load balancing
459    /// there; a warning is logged at startup if enabled outside of Linux.
460    ///
461    /// Enable this on multi-vCPU Linux deployments where UDP DogStatsD throughput is bottlenecked on a single
462    /// receive task.
463    ///
464    /// Defaults to `false`.
465    #[serde(rename = "dogstatsd_autoscale_udp_listeners", default)]
466    autoscale_udp_listeners: bool,
467
468    /// Whether or not to allow heap allocations when resolving contexts.
469    ///
470    /// When resolving contexts during parsing, the metric name and tags are interned to reduce memory usage. The
471    /// interner has a fixed size, however, which means some strings can fail to be interned if the interner is full.
472    /// When set to `true`, we allow these strings to be allocated on the heap like normal, but this can lead to
473    /// increased (unbounded) memory usage. When set to `false`, if the metric name and all of its tags can't be
474    /// interned, the metric is skipped.
475    ///
476    /// Defaults to `true`.
477    #[serde(
478        rename = "dogstatsd_allow_context_heap_allocs",
479        default = "default_allow_context_heap_allocations"
480    )]
481    allow_context_heap_allocations: bool,
482
483    /// Whether or not to enable support for no-aggregation pipelines.
484    ///
485    /// When enabled, this influences how metrics are parsed, specifically around user-provided metric timestamps. When
486    /// metric timestamps are present, it's used as a signal to any aggregation transforms that the metric shouldn't
487    /// be aggregated.
488    ///
489    /// Defaults to `true`.
490    #[serde(
491        rename = "dogstatsd_no_aggregation_pipeline",
492        default = "default_no_aggregation_pipeline_support"
493    )]
494    no_aggregation_pipeline_support: bool,
495
496    /// Number of entries for the string interner, as interpreted by the Core Datadog Agent.
497    ///
498    /// When `dogstatsd_string_interner_size_bytes` isn't set, this value is multiplied by 512 bytes per entry to
499    /// derive the interner byte size. This provides backwards compatibility for customers migrating configurations
500    /// from the Core Agent, where this setting represents an entry count rather than a byte size.
501    ///
502    /// Defaults to 4096 entries, which yields 2 MiB when converted.
503    #[serde(
504        rename = "dogstatsd_string_interner_size",
505        default = "default_context_string_interner_entry_count"
506    )]
507    context_string_interner_entry_count: u64,
508
509    /// Total size of the string interner used for contexts, in bytes.
510    ///
511    /// When set, this takes priority over `dogstatsd_string_interner_size`. This controls the amount of memory that
512    /// can be used to intern metric names and tags. If the interner is full, metrics with contexts that haven't
513    /// already been resolved may or may not be dropped, depending on the value of `allow_context_heap_allocations`.
514    #[serde(rename = "dogstatsd_string_interner_size_bytes", default)]
515    context_string_interner_size_bytes: Option<ByteSize>,
516
517    /// The maximum number of cached contexts to allow.
518    ///
519    /// This is the maximum number of resolved contexts that can be cached at any given time. This limit doesn't affect
520    /// the total number of contexts that can be _alive_ at any given time, which is dependent on the interner capacity
521    /// and whether or not heap allocations are allowed.
522    ///
523    /// Defaults to 500,000.
524    #[serde(
525        rename = "dogstatsd_cached_contexts_limit",
526        default = "default_cached_contexts_limit"
527    )]
528    cached_contexts_limit: usize,
529
530    /// The maximum number of cached tagsets to allow.
531    ///
532    /// This is the maximum number of resolved tagsets that can be cached at any given time. This limit doesn't affect
533    /// the total number of tagsets that can be _alive_ at any given time, which is dependent on the interner capacity
534    /// and whether or not heap allocations are allowed.
535    ///
536    /// Defaults to 500,000.
537    #[serde(rename = "dogstatsd_cached_tagsets_limit", default = "default_cached_tagsets_limit")]
538    cached_tagsets_limit: usize,
539
540    /// The number of seconds after which cached contexts will expire.
541    ///
542    /// Higher values allow for more effective caching for sparse metrics at the cost of increased memory usage.
543    ///
544    /// Defaults to 20 seconds.
545    #[serde(
546        rename = "dogstatsd_context_expiry_seconds",
547        default = "default_context_expiry_seconds"
548    )]
549    context_expiry_seconds: u64,
550
551    /// Whether or not to enable permissive mode in the decoder.
552    ///
553    /// Permissive mode allows the decoder to relax its strictness around the allowed payloads, which lets it match the
554    /// decoding behavior of the Datadog Agent.
555    ///
556    /// Defaults to `true`.
557    #[serde(
558        rename = "dogstatsd_permissive_decoding",
559        default = "default_dogstatsd_permissive_decoding"
560    )]
561    permissive_decoding: bool,
562
563    /// The minimum sample rate allowed for metrics.
564    ///
565    /// When metrics are sent with a sample rate _lower_ than this value then it will be clamped to this value. This is
566    /// done in order to ensure an upper bound on how many equivalent samples are tracked for the metric, as high sample
567    /// rates (very small numbers, such as `0.00000001`) can lead to large memory growth.
568    ///
569    /// A warning log will be emitted when clamping occurs, as this represents an effective loss of metric samples.
570    ///
571    /// Defaults to `0.000000003845`. (~260M samples)
572    #[serde(
573        rename = "dogstatsd_minimum_sample_rate",
574        default = "default_dogstatsd_minimum_sample_rate"
575    )]
576    minimum_sample_rate: f64,
577
578    /// Which payload types to forward to the backend.
579    #[serde(rename = "enable_payloads", default)]
580    enable_payloads: EnablePayloadsConfiguration,
581
582    /// Configuration related to origin detection and enrichment.
583    #[serde(flatten, default)]
584    origin_enrichment: OriginEnrichmentConfiguration,
585
586    /// Configuration related to DogStatsD telemetry.
587    #[serde(default)]
588    telemetry: DogStatsDTelemetryConfiguration,
589
590    /// Workload provider to utilize for origin detection/enrichment.
591    #[serde(skip)]
592    #[cfg_attr(test, derive_where(skip))]
593    workload_provider: Option<Arc<dyn WorkloadProvider + Send + Sync>>,
594
595    /// Resolver to use for mapping live sender PIDs to container entities during traffic capture.
596    #[serde(skip, default)]
597    #[cfg_attr(test, derive_where(skip))]
598    capture_entity_resolver: Option<Arc<dyn CaptureEntityResolver + Send + Sync>>,
599
600    /// Additional tags to add to all metrics.
601    #[serde(rename = "dogstatsd_tags", default)]
602    additional_tags: Vec<String>,
603
604    /// The directory where DogStatsD capture files are written by default.
605    ///
606    /// When set to an empty path, the source attempts to derive the directory from `run_path` by appending
607    /// `dsd_capture`. If neither value is available, callers must provide an explicit capture path when starting a
608    /// capture session.
609    ///
610    /// Defaults to empty.
611    #[serde(rename = "dogstatsd_capture_path", default)]
612    capture_path: PathBuf,
613
614    /// The maximum number of captured packets that can be queued for persistence.
615    ///
616    /// This controls the depth of the in-process capture queue. Values below `1024` are raised to `1024` before the
617    /// capture writer starts, preventing a zero-depth rendezvous channel from serializing DogStatsD stream handlers
618    /// behind capture persistence.
619    ///
620    /// Defaults to `1024`.
621    #[serde(rename = "dogstatsd_capture_depth", default = "default_capture_depth")]
622    capture_depth: usize,
623
624    #[serde(skip, default)]
625    #[cfg_attr(test, derive_where(skip))]
626    capture_control: DogStatsDCaptureControl,
627
628    #[serde(skip, default)]
629    #[cfg_attr(test, derive_where(skip))]
630    replay_control: DogStatsDReplayControl,
631
632    /// Provider kind tag appended to all metrics as `provider_kind:<value>`.
633    ///
634    /// Set via `DD_PROVIDER_KIND` by the Helm chart on GKE Autopilot (`gke-autopilot`) and GKE on
635    /// Google Distributed Cloud (`gke-gdc`). When empty or absent, no tag is added.
636    ///
637    /// Defaults to `""` (disabled).
638    #[serde(default)]
639    provider_kind: String,
640}
641
642#[derive(Clone, Copy, Default)]
643struct EolRequired {
644    udp: bool,
645    uds: bool,
646    named_pipe: bool,
647}
648
649impl EolRequired {
650    fn from_config_values(values: &[String]) -> Self {
651        let mut eol_required = Self::default();
652
653        for value in values {
654            match value.as_str() {
655                "udp" => eol_required.udp = true,
656                "uds" => eol_required.uds = true,
657                "named_pipe" => eol_required.named_pipe = true,
658                _ => warn!(
659                    value,
660                    "Invalid dogstatsd_eol_required value. Expected 'udp', 'uds', or 'named_pipe'."
661                ),
662            }
663        }
664
665        eol_required
666    }
667
668    fn for_listener(&self, listen_addr: &ListenAddress) -> bool {
669        match listen_addr {
670            ListenAddress::Udp(_) => self.udp,
671            ListenAddress::Tcp(_) => false,
672            ListenAddress::Unixgram(_) | ListenAddress::Unix(_) => self.uds,
673            ListenAddress::NamedPipe { .. } => self.named_pipe,
674        }
675    }
676}
677
678/// Resolves a `bind_host` string to an `IpAddr`.
679///
680/// Accepts either an IP literal (no DNS required) or a hostname (resolved via async DNS). Returns
681/// `UnresolvableBindHost` if the lookup fails, or `BindHostHasNoAddresses` if it succeeds but
682/// returns no addresses.
683async fn resolve_bind_host(host: &str) -> Result<std::net::IpAddr, Error> {
684    let mut addrs = tokio::net::lookup_host((host, 0u16))
685        .await
686        .context(UnresolvableBindHost { host: host.to_string() })?;
687    addrs
688        .next()
689        .map(|sa| sa.ip())
690        .ok_or_else(|| Error::BindHostHasNoAddresses { host: host.to_string() })
691}
692
693impl DogStatsDConfiguration {
694    /// Creates a new `DogStatsDConfiguration` from the given configuration.
695    pub fn from_configuration(config: &GenericConfiguration) -> Result<Self, GenericError> {
696        let mut dogstatsd_config: Self = config.as_typed()?;
697        dogstatsd_config.fix_empty_capture_path(config);
698        dogstatsd_config.fix_capture_depth();
699        Ok(dogstatsd_config)
700    }
701
702    /// Gets both the `additional_tags` and any others specified by other configuration fields, such as `provider_kind`.
703    fn additional_tags(&self) -> Vec<String> {
704        if self.provider_kind.is_empty() {
705            return self.additional_tags.clone();
706        }
707
708        let mut tags = self.additional_tags.clone();
709        tags.push(format!("provider_kind:{}", self.provider_kind.clone()));
710        tags
711    }
712
713    fn fix_capture_depth(&mut self) {
714        self.capture_depth = self.capture_depth.max(MIN_CAPTURE_DEPTH);
715    }
716
717    /// Returns the effective string interner size in bytes.
718    ///
719    /// If `dogstatsd_string_interner_size_bytes` is set, it's used directly. Otherwise,
720    /// `dogstatsd_string_interner_size` (an entry count) is multiplied by 512 bytes per entry to derive the byte
721    /// size.
722    fn effective_context_string_interner_bytes(&self) -> ByteSize {
723        match self.context_string_interner_size_bytes {
724            Some(explicit_bytes) => explicit_bytes,
725            None => {
726                saluki_antithesis::always_le!(
727                    self.context_string_interner_entry_count,
728                    u64::MAX / INTERNER_BASELINE_BYTES_PER_ENTRY,
729                    "dogstatsd interner byte-size multiply does not overflow",
730                    { "entry_count": self.context_string_interner_entry_count }
731                );
732                ByteSize::b(
733                    self.context_string_interner_entry_count
734                        .saturating_mul(INTERNER_BASELINE_BYTES_PER_ENTRY),
735                )
736            }
737        }
738    }
739
740    fn eol_required(&self) -> EolRequired {
741        EolRequired::from_config_values(&self.eol_required)
742    }
743
744    fn statsd_forward_target(&self) -> Option<(&MetaString, u16)> {
745        let host = self.statsd_forward_host.as_ref()?;
746        if self.statsd_forward_port == 0 {
747            return None;
748        }
749
750        Some((host, self.statsd_forward_port))
751    }
752
753    fn packet_forwarder_target(&self) -> Option<PacketForwarderTarget> {
754        let (host, port) = self.statsd_forward_target()?;
755        Some(PacketForwarderTarget::new(host.clone(), port))
756    }
757
758    /// Returns the number of UDP stream handlers to spawn, derived from `dogstatsd_autoscale_udp_listeners` and
759    /// the number of available vCPUs.
760    ///
761    /// Returns `None` when autoscaling is disabled, which keeps the legacy single-socket behavior. The platform
762    /// gate for `SO_REUSEPORT` lives inside the listener—this method intentionally stays platform-agnostic.
763    fn udp_streams_to_yield(&self) -> Option<NonZeroUsize> {
764        if !self.autoscale_udp_listeners {
765            return None;
766        }
767
768        #[cfg(not(target_os = "linux"))]
769        if self.autoscale_udp_listeners {
770            warn!("UDP stream handler autoscaling not supported on non-Linux platforms. Default to single stream handler.");
771            return None;
772        }
773
774        let vcpus = std::thread::available_parallelism().map(NonZeroUsize::get).unwrap_or(1);
775        let streams = (1 + vcpus / 8).min(4);
776        NonZeroUsize::new(streams)
777    }
778
779    /// Returns the effective maximum size of the I/O buffer pool.
780    ///
781    /// The pool can never hold fewer buffers than the configured baseline, so a `dogstatsd_buffer_count_max` below
782    /// `dogstatsd_buffer_count` (including a legacy config that only raised `dogstatsd_buffer_count`) is treated as
783    /// equal to the baseline rather than reducing capacity.
784    fn effective_max_buffer_count(&self) -> usize {
785        self.buffer_count_max.max(self.buffer_count)
786    }
787
788    /// Sets the default hostname used when DogStatsD metrics do not carry an explicit `host:` tag.
789    pub fn with_default_hostname(mut self, hostname: impl Into<MetaString>) -> Self {
790        self.default_hostname = hostname.into();
791        self
792    }
793
794    /// Sets the workload provider to use for configuring origin detection/enrichment.
795    ///
796    /// A workload provider must be set otherwise origin detection/enrichment won't be enabled.
797    ///
798    /// Defaults to unset.
799    pub fn with_workload_provider<W>(mut self, workload_provider: W) -> Self
800    where
801        W: WorkloadProvider + Send + Sync + 'static,
802    {
803        self.workload_provider = Some(Arc::new(workload_provider));
804        self
805    }
806
807    /// Sets the resolver to use for mapping live sender PIDs while capturing DogStatsD traffic.
808    ///
809    /// This resolver is intentionally configured separately from the workload provider because capture only needs a
810    /// narrow live-PID lookup, while normal origin enrichment uses the broader workload provider contract.
811    ///
812    /// Defaults to unset.
813    pub fn with_capture_entity_resolver<R>(mut self, capture_entity_resolver: R) -> Self
814    where
815        R: CaptureEntityResolver + Send + Sync + 'static,
816    {
817        self.capture_entity_resolver = Some(Arc::new(capture_entity_resolver));
818        self
819    }
820
821    /// Returns the shared control handle for DogStatsD traffic capture.
822    pub fn capture_control(&self) -> DogStatsDCaptureControl {
823        self.capture_control.clone()
824    }
825
826    /// Returns an HTTP API handler exposing the DogStatsD capture control surface.
827    pub fn capture_api_handler(&self) -> DogStatsDCaptureAPIHandler {
828        DogStatsDCaptureAPIHandler::new(self.capture_control.clone())
829    }
830
831    /// Returns the shared control handle for DogStatsD traffic replay.
832    pub fn replay_control(&self) -> DogStatsDReplayControl {
833        self.replay_control.clone()
834    }
835
836    /// Returns an HTTP API handler exposing the DogStatsD replay control surface.
837    pub fn replay_api_handler(&self) -> DogStatsDReplayAPIHandler {
838        DogStatsDReplayAPIHandler::new(self.replay_control.clone())
839    }
840
841    fn fix_empty_capture_path(&mut self, config: &GenericConfiguration) {
842        if self.capture_path.parent().is_some() {
843            return;
844        }
845
846        let capture_path = match config.try_get_typed::<PathBuf>("run_path") {
847            Ok(Some(mut run_path)) => {
848                run_path.push(DOGSTATSD_CAPTURE_DIR);
849                run_path
850            }
851            Ok(None) => {
852                debug!(
853                    "`dogstatsd_capture_path` and `run_path` were empty. Default DogStatsD capture path is unavailable."
854                );
855                return;
856            }
857            Err(e) => {
858                debug!(
859                    error = %e,
860                    "Failed to read `run_path` from configuration. Default DogStatsD capture path is unavailable."
861                );
862                return;
863            }
864        };
865
866        self.capture_path = capture_path;
867    }
868
869    /// Using the current configuration, determines which listeners should be created and adds an address for each into
870    /// a `Vec<ListenAddress>`. This function has no side effects so that it can be unit tested whereas build_listeners`
871    /// actually binds the listeners on the system.
872    ///
873    /// `bind_host` is the pre-resolved IP that UDP and TCP listeners should bind to (provided by
874    /// `resolve_bind_host`). Precedence matches the Agent:
875    ///   - `non_local_traffic=true` → `0.0.0.0` (`bind_host` ignored)
876    ///   - `bind_host=Some(ip)`     → `ip`
877    ///   - `bind_host=None`         → `127.0.0.1`
878    fn build_addresses(&self, bind_host: Option<std::net::IpAddr>) -> Vec<ListenAddress> {
879        let bind_ip: std::net::IpAddr = if self.non_local_traffic {
880            [0, 0, 0, 0].into()
881        } else {
882            bind_host.unwrap_or_else(|| [127, 0, 0, 1].into())
883        };
884
885        let mut addresses: Vec<ListenAddress> = Vec::new();
886
887        if self.port != 0 {
888            addresses.push(ListenAddress::Udp(std::net::SocketAddr::new(bind_ip, self.port)));
889        }
890
891        if self.tcp_port != 0 {
892            addresses.push(ListenAddress::Tcp(std::net::SocketAddr::new(bind_ip, self.tcp_port)));
893        }
894
895        if let Some(socket_path) = &self.socket_path {
896            addresses.push(ListenAddress::Unixgram(socket_path.into()));
897        }
898
899        if let Some(socket_stream_path) = &self.socket_stream_path {
900            addresses.push(ListenAddress::Unix(socket_stream_path.into()));
901        }
902
903        if let Some(pipe_name) = &self.pipe_name {
904            addresses.push(ListenAddress::named_pipe_with_input_buffer_size(
905                pipe_name,
906                &self.windows_pipe_security_descriptor,
907                self.buffer_size as u32,
908            ));
909        }
910
911        addresses
912    }
913
914    fn uds_origin_detection_unsupported_on_platform(&self, addresses: &[ListenAddress]) -> bool {
915        self.origin_enrichment.enabled()
916            && cfg!(not(target_os = "linux"))
917            && addresses
918                .iter()
919                .any(|address| matches!(address, ListenAddress::Unixgram(_) | ListenAddress::Unix(_)))
920    }
921
922    fn warn_if_uds_origin_detection_unsupported(&self, addresses: &[ListenAddress]) {
923        if self.uds_origin_detection_unsupported_on_platform(addresses) {
924            warn!(
925                "DogStatsD UDS origin detection is enabled, but PID-based Unix socket credentials are unsupported on \
926                 this platform. Metrics are accepted without PID-based origin enrichment."
927            );
928        }
929    }
930
931    /// Builds the appropriate `Listener` objects.
932    async fn build_listeners(&self) -> Result<Vec<Listener>, Error> {
933        // Resolve `bind_host` to an IP (via DNS if needed). Skip the lookup when
934        // `non_local_traffic=true` since `bind_host` is ignored in that branch—matches Go's
935        // laziness and avoids failing startup on an unresolvable hostname that wouldn't be used.
936        let bind_host: Option<std::net::IpAddr> = if self.non_local_traffic {
937            None
938        } else {
939            match &self.bind_host {
940                Some(host) => Some(resolve_bind_host(host).await?),
941                None => None,
942            }
943        };
944
945        let addresses = self.build_addresses(bind_host);
946        self.warn_if_uds_origin_detection_unsupported(&addresses);
947        let mut listeners = Vec::new();
948        let socket_receive_buffer_size =
949            (self.socket_receive_buffer_size != 0).then_some(self.socket_receive_buffer_size);
950        let udp_streams_to_yield = self.udp_streams_to_yield();
951        for address in addresses {
952            let listener_type = address.listener_type();
953            let listener_streams = matches!(address, ListenAddress::Udp(_))
954                .then_some(udp_streams_to_yield)
955                .flatten();
956            let listener = Listener::from_listen_address(address, listener_streams)
957                .await
958                .context(FailedToCreateListener { listener_type })?
959                .with_receive_buffer_size(socket_receive_buffer_size);
960
961            listeners.push(listener);
962        }
963        Ok(listeners)
964    }
965}
966
967#[async_trait]
968impl SourceBuilder for DogStatsDConfiguration {
969    async fn build(&self, context: ComponentContext) -> Result<Box<dyn Source + Send>, GenericError> {
970        let listeners = self.build_listeners().await?;
971        if listeners.is_empty() {
972            return Err(Error::NoListenersConfigured.into());
973        }
974
975        // Every listener requires at least one I/O buffer to ensure that all listeners can be serviced without
976        // deadlocking any of the others. Connectionless listeners retain their buffer for the lifetime of the stream,
977        // so multi-socket UDP listeners require one buffer per yielded socket.
978        let min_buffers: usize = listeners.iter().map(Listener::min_buffer_reservation).sum();
979        let max_buffers = self.effective_max_buffer_count();
980        if max_buffers < min_buffers {
981            return Err(generic_error!(
982                "The maximum I/O buffer count ({}) must be at least {} to service all configured listeners.",
983                max_buffers,
984                min_buffers,
985            ));
986        }
987
988        let origin_detection_enabled = self.origin_enrichment.enabled();
989        // Single CapturedTaggerHandle is cloned to both the resolver (reader of the captured store) and the replay
990        // control surface (writer). Both sides reference the same atomic slot.
991        let captured_tagger = CapturedTaggerHandle::new();
992
993        let maybe_origin_tags_resolver = self.workload_provider.clone().map(|provider| {
994            DogStatsDOriginTagResolver::new(self.origin_enrichment.clone(), provider, captured_tagger.clone())
995        });
996        let context_resolvers = ContextResolvers::new(self, &context, maybe_origin_tags_resolver)
997            .error_context("Failed to create context resolvers.")?;
998
999        let codec_config = DogStatsDCodecConfiguration::default()
1000            .with_timestamps(self.no_aggregation_pipeline_support)
1001            .with_permissive_mode(self.permissive_decoding)
1002            .with_minimum_sample_rate(self.minimum_sample_rate)
1003            .with_client_origin_detection(self.origin_enrichment.origin_detection_client);
1004
1005        let codec = DogStatsDCodec::from_configuration(codec_config);
1006        let eol_required = self.eol_required();
1007
1008        let enable_payloads_filter = EnablePayloadsFilter::default()
1009            .with_allow_series(self.enable_payloads.series)
1010            .with_allow_sketches(self.enable_payloads.sketches)
1011            .with_allow_events(self.enable_payloads.events)
1012            .with_allow_service_checks(self.enable_payloads.service_checks);
1013        let traffic_capture = TrafficCapture::with_workload_provider(
1014            self.capture_path.clone(),
1015            self.capture_depth.max(MIN_CAPTURE_DEPTH),
1016            self.workload_provider.clone(),
1017        );
1018        self.capture_control.bind(traffic_capture.clone());
1019        let packet_forwarder_target = self.packet_forwarder_target();
1020
1021        self.replay_control.bind(captured_tagger);
1022
1023        // The pool allocates `buffer_count` buffers up front and may grow on demand up to `max_buffers`. The effective
1024        // maximum is never below the baseline, so configs that only raise `dogstatsd_buffer_count` keep their full
1025        // capacity instead of being silently reduced to the `dogstatsd_buffer_count_max` default.
1026        let (io_buffer_pool, io_buffer_pool_shrinker) =
1027            build_io_buffer_pool(self.buffer_count, max_buffers, self.buffer_size);
1028
1029        Ok(Box::new(DogStatsD {
1030            listeners,
1031            io_buffer_pool,
1032            io_buffer_pool_shrinker: Box::pin(io_buffer_pool_shrinker),
1033            codec,
1034            context_resolvers,
1035            default_hostname: self.default_hostname.clone(),
1036            enabled_filter: enable_payloads_filter,
1037            origin_detection_enabled,
1038            origin_telemetry_enabled: self.telemetry.dogstatsd_origin,
1039            stream_log_too_big: self.stream_log_too_big,
1040            disable_verbose_logs: self.disable_verbose_logs,
1041            eol_required,
1042            additional_tags: self.additional_tags().into(),
1043            capture_entity_resolver: self.capture_entity_resolver.clone(),
1044            traffic_capture,
1045            packet_forwarder_target,
1046        }))
1047    }
1048
1049    fn outputs(&self) -> &[OutputDefinition<EventType>] {
1050        static OUTPUTS: LazyLock<Vec<OutputDefinition<EventType>>> = LazyLock::new(|| {
1051            vec![
1052                OutputDefinition::named_output("metrics", EventType::Metric),
1053                OutputDefinition::named_output("events", EventType::EventD),
1054                OutputDefinition::named_output("service_checks", EventType::ServiceCheck),
1055            ]
1056        });
1057        &OUTPUTS
1058    }
1059}
1060
1061impl MemoryBounds for DogStatsDConfiguration {
1062    fn specify_bounds(&self, builder: &mut MemoryBoundsBuilder) {
1063        let additional_buffers = self.effective_max_buffer_count().saturating_sub(self.buffer_count);
1064        let adjusted_buffer_size = get_adjusted_buffer_size(self.buffer_size);
1065
1066        builder
1067            .minimum()
1068            // Capture the size of the heap allocation when the component is built.
1069            .with_single_value::<DogStatsD>("source struct")
1070            // We allocate the baseline buffer pool up front.
1071            .with_expr(UsageExpr::product(
1072                "buffers",
1073                UsageExpr::config("dogstatsd_buffer_count", self.buffer_count),
1074                UsageExpr::config("dogstatsd_buffer_size", adjusted_buffer_size),
1075            ))
1076            // We also allocate the backing storage for the string interner up front, which is used by our context
1077            // resolver.
1078            .with_expr(UsageExpr::config(
1079                "dogstatsd_string_interner_size_bytes",
1080                self.effective_context_string_interner_bytes().as_u64() as usize,
1081            ));
1082
1083        // The pool can grow on demand up to its maximum, so account for the additional headroom as firm usage.
1084        builder.firm().with_expr(UsageExpr::product(
1085            "elastic buffers",
1086            UsageExpr::constant("dogstatsd_buffer_count_max_extra", additional_buffers),
1087            UsageExpr::config("dogstatsd_buffer_size", adjusted_buffer_size),
1088        ));
1089    }
1090}
1091
1092/// DogStatsD source.
1093pub struct DogStatsD {
1094    listeners: Vec<Listener>,
1095    io_buffer_pool: ElasticObjectPool<BytesBuffer>,
1096    io_buffer_pool_shrinker: Pin<Box<dyn Future<Output = ()> + Send>>,
1097    codec: DogStatsDCodec,
1098    context_resolvers: ContextResolvers,
1099    default_hostname: MetaString,
1100    enabled_filter: EnablePayloadsFilter,
1101    origin_detection_enabled: bool,
1102    origin_telemetry_enabled: bool,
1103    stream_log_too_big: bool,
1104    disable_verbose_logs: bool,
1105    eol_required: EolRequired,
1106    additional_tags: Arc<[String]>,
1107    capture_entity_resolver: Option<Arc<dyn CaptureEntityResolver + Send + Sync>>,
1108    traffic_capture: TrafficCapture,
1109    packet_forwarder_target: Option<PacketForwarderTarget>,
1110}
1111
1112struct ListenerContext {
1113    shutdown_handle: ShutdownHandle,
1114    listener: Listener,
1115    io_buffer_pool: ElasticObjectPool<BytesBuffer>,
1116    codec: DogStatsDCodec,
1117    context_resolvers: ContextResolvers,
1118    default_hostname: MetaString,
1119    origin_detection_enabled: bool,
1120    origin_telemetry_enabled: bool,
1121    stream_log_too_big: bool,
1122    disable_verbose_logs: bool,
1123    eol_required: EolRequired,
1124    additional_tags: Arc<[String]>,
1125    capture_entity_resolver: Option<Arc<dyn CaptureEntityResolver + Send + Sync>>,
1126    traffic_capture: TrafficCapture,
1127    packet_forwarder_target: Option<PacketForwarderTarget>,
1128}
1129
1130struct HandlerContext {
1131    listen_addr: ListenAddress,
1132    framer: DsdFramer,
1133    codec: DogStatsDCodec,
1134    io_buffer_pool: ElasticObjectPool<BytesBuffer>,
1135    metrics: Metrics,
1136    context_resolvers: ContextResolvers,
1137    default_hostname: MetaString,
1138    origin_detection_enabled: bool,
1139    stream_log_too_big: bool,
1140    disable_verbose_logs: bool,
1141    additional_tags: Arc<[String]>,
1142    capture_entity_resolver: Option<Arc<dyn CaptureEntityResolver + Send + Sync>>,
1143    traffic_capture: TrafficCapture,
1144    packet_forwarder: Option<PacketForwarder>,
1145}
1146
1147#[async_trait]
1148impl Source for DogStatsD {
1149    async fn run(mut self: Box<Self>, mut context: SourceContext) -> Result<(), GenericError> {
1150        let global_shutdown = context.take_shutdown_handle();
1151        pin!(global_shutdown);
1152
1153        let mut health = context.take_health_handle();
1154
1155        let mut listener_shutdown_coordinator = ShutdownCoordinator::default();
1156        spawn_traced_named(
1157            "dogstatsd-io-buffer-pool-shrinker",
1158            process_io_buffer_pool_shrinker(self.io_buffer_pool_shrinker, listener_shutdown_coordinator.register()),
1159        );
1160
1161        // For each listener, spawn a dedicated task to run it.
1162        for listener in self.listeners {
1163            let task_name = format!("dogstatsd-listener-{}", listener.listen_address().listener_type());
1164
1165            // TODO: Create a health handle for each listener.
1166            //
1167            // We need to rework `HealthRegistry` to look a little more like `ComponentRegistry` so that we can have it
1168            // already be scoped properly, otherwise all we can do here at present is either have a relative name, like
1169            // `uds-stream`, or try and hardcode the full component name, which we will inevitably forget to update if
1170            // we tweak the topology configuration, etc.
1171            let listener_context = ListenerContext {
1172                shutdown_handle: listener_shutdown_coordinator.register(),
1173                listener,
1174                io_buffer_pool: self.io_buffer_pool.clone(),
1175                codec: self.codec.clone(),
1176                context_resolvers: self.context_resolvers.clone(),
1177                default_hostname: self.default_hostname.clone(),
1178                origin_detection_enabled: self.origin_detection_enabled,
1179                origin_telemetry_enabled: self.origin_telemetry_enabled,
1180                stream_log_too_big: self.stream_log_too_big,
1181                disable_verbose_logs: self.disable_verbose_logs,
1182                eol_required: self.eol_required,
1183                additional_tags: self.additional_tags.clone(),
1184                capture_entity_resolver: self.capture_entity_resolver.clone(),
1185                traffic_capture: self.traffic_capture.clone(),
1186                packet_forwarder_target: self.packet_forwarder_target.clone(),
1187            };
1188
1189            spawn_traced_named(
1190                task_name,
1191                process_listener(context.clone(), listener_context, self.enabled_filter),
1192            );
1193        }
1194
1195        health.mark_ready();
1196        debug!("DogStatsD source started.");
1197
1198        // Wait for the global shutdown signal, then notify listeners to shutdown.
1199        //
1200        // We also handle liveness here, which doesn't really matter for _this_ task, since the real work is happening
1201        // in the listeners, but we need to satisfy the health checker.
1202        loop {
1203            select! {
1204                _ = &mut global_shutdown => {
1205                    debug!("Received shutdown signal.");
1206                    break
1207                },
1208                _ = health.live() => continue,
1209            }
1210        }
1211
1212        debug!("Stopping DogStatsD source...");
1213
1214        listener_shutdown_coordinator.shutdown_and_wait().await;
1215
1216        debug!("DogStatsD source stopped.");
1217
1218        Ok(())
1219    }
1220}
1221
1222async fn process_io_buffer_pool_shrinker(
1223    io_buffer_pool_shrinker: Pin<Box<dyn Future<Output = ()> + Send>>, shutdown_handle: ShutdownHandle,
1224) {
1225    pin!(shutdown_handle);
1226
1227    select! {
1228        _ = &mut shutdown_handle => {
1229            debug!("I/O buffer pool shrinker received shutdown signal.");
1230        },
1231        _ = io_buffer_pool_shrinker => {
1232            debug!("I/O buffer pool shrinker stopped.");
1233        },
1234    }
1235}
1236
1237fn build_io_buffer_pool(
1238    min_buffers: usize, max_buffers: usize, buffer_size: usize,
1239) -> (ElasticObjectPool<BytesBuffer>, impl Future<Output = ()> + Send) {
1240    saluki_antithesis::always_le!(
1241        buffer_size,
1242        usize::MAX - 4,
1243        "dogstatsd buffer size add does not overflow",
1244        { "buffer_size": buffer_size }
1245    );
1246    let adjusted_buffer_size = get_adjusted_buffer_size(buffer_size);
1247    ElasticObjectPool::with_builder("dsd_packet_bufs", min_buffers, max_buffers, move || {
1248        FixedSizeVec::with_capacity(adjusted_buffer_size)
1249    })
1250}
1251
1252async fn process_listener(
1253    source_context: SourceContext, listener_context: ListenerContext, enabled_filter: EnablePayloadsFilter,
1254) {
1255    let ListenerContext {
1256        shutdown_handle,
1257        mut listener,
1258        io_buffer_pool,
1259        codec,
1260        context_resolvers,
1261        default_hostname,
1262        origin_detection_enabled,
1263        origin_telemetry_enabled,
1264        stream_log_too_big,
1265        disable_verbose_logs,
1266        eol_required,
1267        additional_tags,
1268        capture_entity_resolver,
1269        traffic_capture,
1270        packet_forwarder_target,
1271    } = listener_context;
1272
1273    pin!(shutdown_handle);
1274
1275    let listen_addr = listener.listen_address().clone();
1276    let metrics = build_metrics(
1277        &listen_addr,
1278        source_context.component_context(),
1279        origin_telemetry_enabled,
1280    );
1281    let packet_forwarder = packet_forwarder_target
1282        .as_ref()
1283        .map(|target| target.to_forwarder(metrics.clone()));
1284    if let Some(packet_forwarder) = &packet_forwarder {
1285        packet_forwarder.spawn_connect();
1286    }
1287
1288    let mut stream_shutdown_coordinator = ShutdownCoordinator::default();
1289
1290    info!(%listen_addr, "DogStatsD listener started.");
1291
1292    loop {
1293        select! {
1294            _ = &mut shutdown_handle => {
1295                debug!(%listen_addr, "Received shutdown signal. Waiting for existing stream handlers to finish...");
1296                break;
1297            }
1298            result = listener.accept() => match result {
1299                Ok(stream) => {
1300                    debug!(%listen_addr, "Spawning new stream handler.");
1301
1302                    let handler_context = HandlerContext {
1303                        listen_addr: listen_addr.clone(),
1304                        framer: get_framer(&listen_addr, eol_required.for_listener(&listen_addr)),
1305                        codec: codec.clone(),
1306                        io_buffer_pool: io_buffer_pool.clone(),
1307                        metrics: metrics.clone(),
1308                        context_resolvers: context_resolvers.clone(),
1309                        default_hostname: default_hostname.clone(),
1310                        origin_detection_enabled,
1311                        stream_log_too_big,
1312                        disable_verbose_logs,
1313                        additional_tags: additional_tags.clone(),
1314                        capture_entity_resolver: capture_entity_resolver.clone(),
1315                        traffic_capture: traffic_capture.clone(),
1316                        packet_forwarder: packet_forwarder.clone(),
1317                    };
1318
1319                    let task_name = format!(
1320                        "dogstatsd-stream-handler-{}",
1321                        listen_addr.listener_type(),
1322                    );
1323                    spawn_traced_named(task_name, process_stream(stream, source_context.clone(), handler_context, stream_shutdown_coordinator.register(), enabled_filter));
1324                }
1325                Err(e) => {
1326                    error!(%listen_addr, error = %e, "Failed to accept connection. Stopping listener.");
1327                    break
1328                }
1329            }
1330        }
1331    }
1332
1333    stream_shutdown_coordinator.shutdown_and_wait().await;
1334
1335    info!(%listen_addr, "DogStatsD listener stopped.");
1336}
1337
1338async fn process_stream(
1339    stream: Stream, source_context: SourceContext, handler_context: HandlerContext, shutdown_handle: ShutdownHandle,
1340    enabled_filter: EnablePayloadsFilter,
1341) {
1342    select! {
1343        _ = shutdown_handle => {
1344            debug!("Stream handler received shutdown signal.");
1345        },
1346        _ = drive_stream(stream, source_context, handler_context, enabled_filter) => {},
1347    }
1348}
1349
1350fn origin_detection_failed_for_telemetry(
1351    origin_detection_enabled: bool, bytes_read: usize, peer_addr: &ConnectionAddress,
1352) -> bool {
1353    origin_detection_enabled && bytes_read > 0 && peer_addr.has_process_credential_telemetry_error()
1354}
1355
1356async fn drive_stream(
1357    mut stream: Stream, source_context: SourceContext, handler_context: HandlerContext,
1358    enabled_filter: EnablePayloadsFilter,
1359) {
1360    let HandlerContext {
1361        listen_addr,
1362        mut framer,
1363        codec,
1364        io_buffer_pool,
1365        metrics,
1366        mut context_resolvers,
1367        default_hostname,
1368        origin_detection_enabled,
1369        stream_log_too_big,
1370        disable_verbose_logs,
1371        additional_tags,
1372        capture_entity_resolver,
1373        traffic_capture,
1374        packet_forwarder,
1375    } = handler_context;
1376
1377    debug!(%listen_addr, "Stream handler started.");
1378
1379    if !stream.is_connectionless() {
1380        metrics.connections_active().increment(1);
1381    }
1382
1383    let mut stream_capture = StreamCaptureState::new();
1384    // Set a buffer flush interval of 100ms, which will ensure we always flush buffered events at least every 100ms if
1385    // we're otherwise idle and not receiving packets from the client.
1386    let mut buffer_flush = interval(Duration::from_millis(100));
1387    buffer_flush.set_missed_tick_behavior(MissedTickBehavior::Delay);
1388
1389    let mut event_buffer_manager = EventBufferManager::default();
1390    let mut io_buffer_manager = IoBufferManager::new(&io_buffer_pool, &stream);
1391    let memory_limiter = source_context.topology_context().memory_limiter();
1392
1393    'read: loop {
1394        let mut eof = false;
1395
1396        let mut io_buffer = io_buffer_manager.get_buffer_mut().await;
1397
1398        memory_limiter.wait_for_capacity().await;
1399
1400        select! {
1401            // We read from the stream.
1402            read_result = stream.receive(&mut io_buffer) => match read_result {
1403                Ok((bytes_read, peer_addr)) => {
1404                    if bytes_read == 0 {
1405                        eof = true;
1406                    }
1407
1408                    let is_connectionless = stream.is_connectionless();
1409                    let payload = received_payload(io_buffer, bytes_read);
1410
1411                    capture_uds_traffic(
1412                        &listen_addr,
1413                        &traffic_capture,
1414                        capture_entity_resolver.as_deref(),
1415                        &peer_addr,
1416                        payload,
1417                        &mut stream_capture,
1418                    );
1419
1420                    if is_connectionless {
1421                        metrics.packet_receive_success().increment(1);
1422                    }
1423                    metrics.bytes_received().increment(bytes_read as u64);
1424                    metrics.bytes_received_size().record(bytes_read as f64);
1425                    let origin_detection_failed =
1426                        origin_detection_failed_for_telemetry(origin_detection_enabled, bytes_read, &peer_addr);
1427                    if origin_detection_failed && is_connectionless {
1428                        metrics.origin_detection_errors().increment(1);
1429                    }
1430
1431                    // When we're actually at EOF, or we're dealing with a connectionless stream, we try to decode in EOF mode.
1432                    //
1433                    // For connectionless streams, we always try to decode the buffer as if it's EOF, since it effectively _is_
1434                    // always the end of file after a receive. For connection-oriented streams, we only want to do this once we've
1435                    // actually hit true EOF.
1436                    let reached_eof = eof || is_connectionless;
1437
1438                    trace!(
1439                        buffer_len = io_buffer.remaining(),
1440                        buffer_cap = io_buffer.remaining_mut(),
1441                        eof = reached_eof,
1442                        %listen_addr,
1443                        %peer_addr,
1444                        "Received {} bytes from stream.",
1445                        bytes_read
1446                    );
1447
1448                    if should_drop_oversized_named_pipe_frame(&listen_addr, io_buffer) {
1449                        metrics.framing_errors().increment(1);
1450                        debug!(%listen_addr, %peer_addr, "DogStatsD named pipe frame exceeded the configured buffer size. Dropping frame.");
1451                        io_buffer.clear();
1452                        continue 'read;
1453                    }
1454
1455                    'frame: loop {
1456                        let frame_result = framer.next_frame(io_buffer, reached_eof);
1457                        let completed_outer_frames = framer.take_completed_outer_frames();
1458                        if !is_connectionless && completed_outer_frames > 0 {
1459                            metrics.packet_receive_success().increment(completed_outer_frames as u64);
1460                        }
1461                        if origin_detection_failed && completed_outer_frames > 0 {
1462                            metrics.origin_detection_errors().increment(completed_outer_frames as u64);
1463                        }
1464
1465                        match frame_result {
1466                            Ok(Some(frame)) => {
1467                                if matches!(listen_addr, ListenAddress::NamedPipe { .. }) {
1468                                    metrics.packet_receive_success().increment(1);
1469                                }
1470                                trace!(%listen_addr, %peer_addr, ?frame, "Decoded frame.");
1471                                if let Some(forwarder) = &packet_forwarder {
1472                                    forwarder.forward(frame.clone()).await;
1473                                }
1474                                match handle_frame(
1475                                    &frame[..],
1476                                    &codec,
1477                                    &mut context_resolvers,
1478                                    &metrics,
1479                                    capture_entity_resolver.as_deref(),
1480                                    origin_detection_enabled,
1481                                    &peer_addr,
1482                                    enabled_filter,
1483                                    &additional_tags,
1484                                    &default_hostname,
1485                                ) {
1486                                    Ok(Some(event)) => {
1487                                        if let Some(event_buffer) = event_buffer_manager.try_push(event) {
1488                                            debug!(%listen_addr, %peer_addr, "Event buffer is full. Forwarding events.");
1489                                            dispatch_events(event_buffer, &source_context, &listen_addr).await;
1490                                        }
1491                                    },
1492                                    Ok(None) => {
1493                                        // We didn't decode an event, but there was no inherent error. This is likely
1494                                        // due to hitting resource limits, etc.
1495                                        //
1496                                        // Simply continue on.
1497                                        continue
1498                                    },
1499                                    Err(e) => {
1500                                        log_parse_failure(disable_verbose_logs, &listen_addr, &peer_addr, &frame, &e);
1501                                    },
1502                                }
1503                            }
1504                            Err(e) => {
1505                                metrics.framing_errors().increment(1);
1506                                if should_warn_stream_log_too_big(&listen_addr, &e, stream_log_too_big) {
1507                                    warn!(
1508                                        %listen_addr,
1509                                        %peer_addr,
1510                                        error = %e,
1511                                        "DogStatsD stream frame exceeded the configured buffer size."
1512                                    );
1513                                }
1514
1515                                if stream.is_connectionless() {
1516                                    io_buffer.clear();
1517                                    // For connectionless streams, we don't want to shutdown the stream since we can just keep
1518                                    // reading more packets.
1519                                    debug!(%listen_addr, %peer_addr, error = %e, "Error decoding frame. Continuing stream.");
1520                                    continue 'read;
1521                                } else {
1522                                    debug!(%listen_addr, %peer_addr, error = %e, "Error decoding frame. Stopping stream.");
1523                                    break 'read;
1524                                }
1525                            }
1526                            Ok(None) => {
1527                                trace!(%listen_addr, %peer_addr, "Not enough data to decode another frame.");
1528                                if eof && !stream.is_connectionless() {
1529                                    debug!(%listen_addr, %peer_addr, "Stream received EOF. Shutting down handler.");
1530                                    break 'read;
1531                                } else {
1532                                    break 'frame;
1533                                }
1534                            }
1535                        }
1536                    }
1537                },
1538                Err(e) => {
1539                    metrics.packet_receive_failure().increment(1);
1540
1541                    if stream.is_connectionless() {
1542                        // For connectionless streams, we don't want to shutdown the stream since we can just keep
1543                        // reading more packets.
1544                        warn!(%listen_addr, error = %e, "I/O error while decoding. Continuing stream.");
1545                        continue 'read;
1546                    } else {
1547                        warn!(%listen_addr, error = %e, "I/O error while decoding. Stopping stream.");
1548                        break 'read;
1549                    }
1550                }
1551            },
1552
1553            _ = buffer_flush.tick() => {
1554                if let Some(event_buffer) = event_buffer_manager.consume() {
1555                    dispatch_events(event_buffer, &source_context, &listen_addr).await;
1556                }
1557            },
1558
1559        }
1560    }
1561
1562    if let Some(event_buffer) = event_buffer_manager.consume() {
1563        dispatch_events(event_buffer, &source_context, &listen_addr).await;
1564    }
1565
1566    metrics.connections_active().decrement(1);
1567
1568    debug!(%listen_addr, "Stream handler stopped.");
1569}
1570
1571fn should_drop_oversized_named_pipe_frame(listen_addr: &ListenAddress, buffer: &BytesBuffer) -> bool {
1572    matches!(listen_addr, ListenAddress::NamedPipe { .. })
1573        && buffer.remaining_mut() == 0
1574        && memchr::memchr(b'\n', buffer.chunk()).is_none()
1575}
1576
1577fn should_warn_stream_log_too_big(listen_addr: &ListenAddress, error: &FramingError, stream_log_too_big: bool) -> bool {
1578    stream_log_too_big
1579        && matches!(listen_addr, ListenAddress::Unix(_))
1580        && matches!(error, FramingError::InvalidFrame { .. })
1581}
1582
1583fn log_parse_failure(
1584    disable_verbose_logs: bool, listen_addr: &ListenAddress, peer_addr: &ConnectionAddress, frame: &[u8],
1585    error: &ParseError,
1586) {
1587    let frame = String::from_utf8_lossy(frame);
1588    if disable_verbose_logs {
1589        debug!(%listen_addr, %peer_addr, %frame, %error, "Failed to parse frame.");
1590    } else {
1591        warn!(%listen_addr, %peer_addr, %frame, %error, "Failed to parse frame.");
1592    }
1593}
1594
1595fn capture_uds_traffic(
1596    listen_addr: &ListenAddress, traffic_capture: &TrafficCapture,
1597    capture_entity_resolver: Option<&(dyn CaptureEntityResolver + Send + Sync)>, peer_addr: &ConnectionAddress,
1598    payload: &[u8], stream_capture: &mut StreamCaptureState,
1599) {
1600    if payload.is_empty() || !traffic_capture.is_ongoing() {
1601        return;
1602    }
1603
1604    match listen_addr {
1605        ListenAddress::Unixgram(_) => {
1606            let _ = traffic_capture.enqueue(build_capture_record(
1607                capture_entity_resolver,
1608                process_id_from_peer_addr(peer_addr),
1609                payload,
1610            ));
1611        }
1612        ListenAddress::Unix(_) => {
1613            stream_capture.update_peer_metadata(peer_addr);
1614            stream_capture.pending.extend(payload);
1615
1616            while let Ok(Some(outer_payload)) = stream_capture
1617                .outer_framer
1618                .next_frame(&mut stream_capture.pending, false)
1619            {
1620                let _ = traffic_capture.enqueue(build_capture_record(
1621                    capture_entity_resolver,
1622                    stream_capture.last_pid,
1623                    &outer_payload,
1624                ));
1625            }
1626        }
1627        _ => {}
1628    }
1629}
1630
1631struct StreamCaptureState {
1632    outer_framer: LengthDelimitedFramer,
1633    pending: VecDeque<u8>,
1634    last_pid: Option<i32>,
1635}
1636
1637impl StreamCaptureState {
1638    fn new() -> Self {
1639        Self {
1640            outer_framer: LengthDelimitedFramer,
1641            pending: VecDeque::new(),
1642            last_pid: None,
1643        }
1644    }
1645
1646    fn update_peer_metadata(&mut self, peer_addr: &ConnectionAddress) {
1647        if let Some(process_id) = process_id_from_peer_addr(peer_addr) {
1648            self.last_pid = Some(process_id);
1649        }
1650    }
1651}
1652
1653fn build_capture_record(
1654    capture_entity_resolver: Option<&(dyn CaptureEntityResolver + Send + Sync)>, process_id: Option<i32>,
1655    payload: &[u8],
1656) -> CaptureRecord {
1657    CaptureRecord {
1658        timestamp_ns: capture_timestamp_ns(),
1659        payload: payload.to_vec(),
1660        pid: process_id,
1661        ancillary: Vec::new(),
1662        container_id: resolve_capture_container_id(capture_entity_resolver, process_id),
1663    }
1664}
1665
1666fn resolve_capture_container_id(
1667    capture_entity_resolver: Option<&(dyn CaptureEntityResolver + Send + Sync)>, process_id: Option<i32>,
1668) -> Option<String> {
1669    let process_id = u32::try_from(process_id?).ok()?;
1670    capture_entity_resolver
1671        .and_then(|resolver| resolver.resolve_container_entity_for_live_pid(process_id))
1672        .map(|entity_id| entity_id.to_string())
1673}
1674
1675fn process_id_from_peer_addr(peer_addr: &ConnectionAddress) -> Option<i32> {
1676    match peer_addr {
1677        ConnectionAddress::ProcessLike(ProcessIdentity::Credentials(creds)) => Some(creds.pid),
1678        _ => None,
1679    }
1680}
1681
1682/// Applies SCM_CREDENTIALS to the origin, dispatching on the replay marker GID.
1683///
1684/// Live packets carry the sender process's real PID/UID/GID; we use `creds.pid` as the origin's process ID. Replay
1685/// packets carry `gid == REPLAY_CREDENTIALS_GID` and pack the captured (original) PID into `creds.uid`; we recover
1686/// that PID with an internal marker so downstream tag resolution consults the captured tagger store.
1687fn apply_credentials_to_origin(origin: &mut RawOrigin<'_>, creds: &ProcessCredentials) {
1688    if creds.gid == REPLAY_CREDENTIALS_GID {
1689        origin.set_process_id(mark_replay_process_id(creds.uid));
1690    } else {
1691        origin.set_process_id(creds.pid as u32);
1692    }
1693}
1694
1695fn received_payload(buffer: &BytesBuffer, bytes_read: usize) -> &[u8] {
1696    let chunk = buffer.chunk();
1697    let start = chunk.len().saturating_sub(bytes_read);
1698    &chunk[start..]
1699}
1700
1701fn capture_timestamp_ns() -> i64 {
1702    SystemTime::now()
1703        .duration_since(UNIX_EPOCH)
1704        .map(|duration| duration.as_nanos().min(i64::MAX as u128) as i64)
1705        .unwrap_or_default()
1706}
1707
1708#[allow(clippy::too_many_arguments)]
1709fn handle_frame(
1710    frame: &[u8], codec: &DogStatsDCodec, context_resolvers: &mut ContextResolvers, source_metrics: &Metrics,
1711    capture_entity_resolver: Option<&(dyn CaptureEntityResolver + Send + Sync)>, origin_detection_enabled: bool,
1712    peer_addr: &ConnectionAddress, enabled_filter: EnablePayloadsFilter, additional_tags: &[String],
1713    default_hostname: &MetaString,
1714) -> Result<Option<Event>, ParseError> {
1715    // Resolving the origin requires a (potentially uncached) PID-to-container lookup, so only do it for metric frames,
1716    // which are the only frames that record per-origin telemetry.
1717    let resolve_telemetry_origin = || {
1718        (source_metrics.origin_telemetry_enabled() && origin_detection_enabled)
1719            .then(|| resolve_capture_container_id(capture_entity_resolver, process_id_from_peer_addr(peer_addr)))
1720            .flatten()
1721    };
1722
1723    let parsed = match codec.decode_packet(frame) {
1724        Ok(parsed) => parsed,
1725        Err(e) => {
1726            // Try and determine what the message type was, if possible, to increment the correct error counter.
1727            match parse_message_type(frame) {
1728                MessageType::MetricSample => {
1729                    source_metrics.record_metric_parse_failed(resolve_telemetry_origin().as_deref())
1730                }
1731                MessageType::Event => source_metrics.event_decode_failed().increment(1),
1732                MessageType::ServiceCheck => source_metrics.service_check_decode_failed().increment(1),
1733            }
1734
1735            return Err(e);
1736        }
1737    };
1738
1739    let event = match parsed {
1740        ParsedPacket::Metric(metric_packet) => {
1741            if metric_packet.num_points == 0 {
1742                return Ok(None);
1743            }
1744            let events_len = metric_packet.num_points;
1745            if !enabled_filter.allow_metric(&metric_packet) {
1746                trace!(
1747                    metric.name = metric_packet.metric_name,
1748                    "Skipping metric due to filter configuration."
1749                );
1750                return Ok(None);
1751            }
1752
1753            match handle_metric_packet(
1754                metric_packet,
1755                context_resolvers,
1756                peer_addr,
1757                additional_tags,
1758                default_hostname,
1759            ) {
1760                Some(metric) => {
1761                    source_metrics.record_metrics_received(events_len, resolve_telemetry_origin().as_deref());
1762                    Event::Metric(metric)
1763                }
1764                None => {
1765                    // We can only fail to get a metric back if we failed to resolve the context.
1766                    source_metrics.failed_context_resolve_total().increment(1);
1767                    return Ok(None);
1768                }
1769            }
1770        }
1771        ParsedPacket::Event(event) => {
1772            if !enabled_filter.allow_event(&event) {
1773                trace!("Skipping event {} due to filter configuration.", event.title);
1774                return Ok(None);
1775            }
1776            let tags_resolver = context_resolvers.tags();
1777            match handle_event_packet(event, tags_resolver, peer_addr, additional_tags) {
1778                Some(event) => {
1779                    source_metrics.events_received().increment(1);
1780                    Event::EventD(event)
1781                }
1782                None => {
1783                    source_metrics.failed_context_resolve_total().increment(1);
1784                    return Ok(None);
1785                }
1786            }
1787        }
1788        ParsedPacket::ServiceCheck(service_check) => {
1789            if !enabled_filter.allow_service_check(&service_check) {
1790                trace!(
1791                    "Skipping service check {} due to filter configuration.",
1792                    service_check.name
1793                );
1794                return Ok(None);
1795            }
1796            let tags_resolver = context_resolvers.tags();
1797            match handle_service_check_packet(service_check, tags_resolver, peer_addr, additional_tags) {
1798                Some(service_check) => {
1799                    source_metrics.service_checks_received().increment(1);
1800                    Event::ServiceCheck(service_check)
1801                }
1802                None => {
1803                    source_metrics.failed_context_resolve_total().increment(1);
1804                    return Ok(None);
1805                }
1806            }
1807        }
1808    };
1809
1810    Ok(Some(event))
1811}
1812
1813fn handle_metric_packet(
1814    packet: MetricPacket, context_resolvers: &mut ContextResolvers, peer_addr: &ConnectionAddress,
1815    additional_tags: &[String], default_hostname: &MetaString,
1816) -> Option<Metric> {
1817    let well_known_tags = WellKnownTags::from_raw_tags(packet.tags.clone());
1818
1819    let mut origin = origin_from_metric_packet(&packet, &well_known_tags);
1820    if let Some(creds) = peer_addr.process_credentials() {
1821        apply_credentials_to_origin(&mut origin, creds);
1822    }
1823
1824    // Choose the right context resolver based on whether or not this metric is pre-aggregated.
1825    let context_resolver = if packet.timestamp.is_some() {
1826        context_resolvers.no_agg()
1827    } else {
1828        context_resolvers.primary()
1829    };
1830
1831    let tags = get_filtered_tags_iterator(packet.tags, additional_tags);
1832
1833    let hostname = well_known_tags.hostname.unwrap_or(default_hostname);
1834
1835    // Try to resolve the context for this metric.
1836    let maybe_context = context_resolver.resolve_with_host(packet.metric_name, hostname, tags, Some(origin));
1837
1838    match maybe_context {
1839        Some(context) => {
1840            let metric_origin = well_known_tags
1841                .jmx_check_name
1842                .map(MetricOrigin::jmx_check)
1843                .unwrap_or_else(MetricOrigin::dogstatsd);
1844            let metadata = MetricMetadata::default()
1845                .with_origin(metric_origin)
1846                .with_unit(packet.unit.map_or_else(MetaString::empty, MetaString::from_static));
1847
1848            Some(Metric::from_parts(context, packet.values, metadata))
1849        }
1850        // We failed to resolve the context, likely due to not having enough interner capacity.
1851        None => None,
1852    }
1853}
1854
1855fn handle_event_packet(
1856    packet: EventPacket, tags_resolver: &mut TagsResolver, peer_addr: &ConnectionAddress, additional_tags: &[String],
1857) -> Option<EventD> {
1858    let well_known_tags = WellKnownTags::from_raw_tags(packet.tags.clone());
1859
1860    let mut origin = origin_from_event_packet(&packet, &well_known_tags);
1861    if let Some(creds) = peer_addr.process_credentials() {
1862        apply_credentials_to_origin(&mut origin, creds);
1863    }
1864    let origin_tags = tags_resolver.resolve_origin_tags(Some(origin));
1865
1866    let tags = get_filtered_tags_iterator(packet.tags, additional_tags);
1867    let tags = tags_resolver.create_tag_set(tags)?;
1868
1869    // When no d: field is present, backfill the current time—matching the stock Datadog Agent's
1870    // behavior in pkg/aggregator/aggregator.go (addEvent), which sets e.Ts = time.Now().Unix()
1871    // for any event with Ts == 0.
1872    let timestamp = packet
1873        .timestamp
1874        .or_else(|| SystemTime::now().duration_since(UNIX_EPOCH).ok().map(|d| d.as_secs()));
1875
1876    let eventd = EventD::new(packet.title, packet.text)
1877        .with_timestamp(timestamp)
1878        .with_hostname(packet.hostname.map(|s| s.into()))
1879        .with_aggregation_key(packet.aggregation_key.map(|s| s.into()))
1880        .with_alert_type(packet.alert_type)
1881        .with_priority(packet.priority)
1882        // When no source type is provided, default to "api"—the same default the stock Datadog
1883        // Agent applies when serializing DogStatsD events to the intake JSON format. The agent
1884        // groups events by source type name and uses "api" as the key for events without an
1885        // explicit `s:` field. See: pkg/serializer/internal/metrics/events.go (writeItem).
1886        .with_source_type_name(Some(
1887            packet
1888                .source_type_name
1889                .map(|s| s.into())
1890                .unwrap_or_else(|| "api".into()),
1891        ))
1892        .with_alert_type(packet.alert_type)
1893        .with_tags(tags)
1894        .with_origin_tags(origin_tags);
1895
1896    Some(eventd)
1897}
1898
1899fn handle_service_check_packet(
1900    packet: ServiceCheckPacket, tags_resolver: &mut TagsResolver, peer_addr: &ConnectionAddress,
1901    additional_tags: &[String],
1902) -> Option<ServiceCheck> {
1903    let well_known_tags = WellKnownTags::from_raw_tags(packet.tags.clone());
1904
1905    let mut origin = origin_from_service_check_packet(&packet, &well_known_tags);
1906    if let Some(creds) = peer_addr.process_credentials() {
1907        apply_credentials_to_origin(&mut origin, creds);
1908    }
1909    let origin_tags = tags_resolver.resolve_origin_tags(Some(origin));
1910
1911    let tags = get_filtered_tags_iterator(packet.tags, additional_tags);
1912    let tags = tags_resolver.create_tag_set(tags)?;
1913
1914    // When no d: field is present, backfill the current time—matching the stock Datadog Agent's
1915    // behavior, which sets the timestamp to time.Now().Unix() for any service check with a zero
1916    // timestamp.
1917    let timestamp = packet
1918        .timestamp
1919        .or_else(|| SystemTime::now().duration_since(UNIX_EPOCH).ok().map(|d| d.as_secs()));
1920
1921    let service_check = ServiceCheck::new(packet.name, packet.status)
1922        .with_timestamp(timestamp)
1923        .with_hostname(packet.hostname.map(|s| s.into()))
1924        .with_tags(tags)
1925        .with_origin_tags(origin_tags)
1926        .with_message(packet.message.map(|s| s.into()));
1927
1928    Some(service_check)
1929}
1930
1931fn get_filtered_tags_iterator<'a>(
1932    raw_tags: RawTags<'a>, additional_tags: &'a [String],
1933) -> impl Iterator<Item = &'a str> + Clone {
1934    // This filters out "well-known" tags from the raw tags in the DogStatsD packet, and then chains on any additional tags
1935    // that were configured on the source.
1936    RawTagsFilter::exclude(raw_tags, WellKnownTagsFilterPredicate).chain(additional_tags.iter().map(|s| s.as_str()))
1937}
1938
1939async fn dispatch_events(mut event_buffer: EventsBuffer, source_context: &SourceContext, listen_addr: &ListenAddress) {
1940    debug!(%listen_addr, events_len = event_buffer.len(), "Forwarding events.");
1941
1942    // TODO: This is maybe a little dicey because if we fail to dispatch the events, we may not have iterated over all of
1943    // them, so there might still be eventd events when get to the service checks point, and eventd events and/or service
1944    // check events when we get to the metrics point, and so on.
1945    //
1946    // There's probably something to be said for erroring out fully if this happens, since we should only fail to
1947    // dispatch if the downstream component fails entirely... and unless we have a way to restart the component, then
1948    // we're going to continue to fail to dispatch any more events until the process is restarted anyways.
1949
1950    // Dispatch any eventd events, if present.
1951    if event_buffer.has_event_type(EventType::EventD) {
1952        let eventd_events = event_buffer.extract(Event::is_eventd);
1953        let events_output = source_context.dispatcher().buffered_named("events");
1954
1955        // The `events` output is always wired in the DSD topology, so a missing output is an invariant violation that
1956        // crashes this component.
1957        if events_output.is_err() {
1958            saluki_antithesis::unreachable!("dsd 'events' output missing at dispatch");
1959        }
1960
1961        if let Err(e) = events_output
1962            .expect("events output should always exist")
1963            .send_all(eventd_events)
1964            .await
1965        {
1966            error!(%listen_addr, error = %e, "Failed to dispatch eventd events.");
1967
1968            saluki_antithesis::unreachable!("dsd dispatch failed mid-buffer", { "stream": "events" });
1969        }
1970    }
1971
1972    // Dispatch any service check events, if present.
1973    if event_buffer.has_event_type(EventType::ServiceCheck) {
1974        let service_check_events = event_buffer.extract(Event::is_service_check);
1975        let service_checks_output = source_context.dispatcher().buffered_named("service_checks");
1976
1977        if service_checks_output.is_err() {
1978            saluki_antithesis::unreachable!("dsd 'service_checks' output missing at dispatch");
1979        }
1980
1981        if let Err(e) = service_checks_output
1982            .expect("service checks output should always exist")
1983            .send_all(service_check_events)
1984            .await
1985        {
1986            error!(%listen_addr, error = %e, "Failed to dispatch service check events.");
1987
1988            saluki_antithesis::unreachable!("dsd dispatch failed mid-buffer", { "stream": "service_checks" });
1989        }
1990    }
1991
1992    // Finally, if there are events left, they'll be metrics, so dispatch them.
1993    if !event_buffer.is_empty() {
1994        if let Err(e) = source_context
1995            .dispatcher()
1996            .dispatch_named("metrics", event_buffer)
1997            .await
1998        {
1999            error!(%listen_addr, error = %e, "Failed to dispatch metric events.");
2000
2001            saluki_antithesis::unreachable!("dsd dispatch failed mid-buffer", { "stream": "metrics" });
2002        }
2003    }
2004}
2005
2006const fn get_adjusted_buffer_size(buffer_size: usize) -> usize {
2007    // This is a little goofy, but hear me out:
2008    //
2009    // In the Datadog Agent, the way the UDS listener works is that if it's in stream mode, it will do a standalone
2010    // socket read to get _just_ the length delimiter, which is 4 bytes. After that, it will do a read to get the packet
2011    // data itself, up to the limit of `dogstatsd_buffer_size`. This means that a _full_ UDS stream packet can be up to
2012    // `dogstatsd_buffer_size + 4` bytes.
2013    //
2014    // This isn't a problem in the Agent due to how it does the reads, but it's a problem for us because we want to be
2015    // able to get an entire frame in a single buffer for the purpose of decoding the frame. Rather than rewriting our
2016    // read loop such that we have to change the logic depending on UDP/UDS datagram vs UDS stream, we simply increase
2017    // the buffer size by 4 bytes to account for the length delimiter.
2018    //
2019    // We do it this way so that we don't have to change the buffer size in the configuration, since if you just ported
2020    // over a Datadog Agent configuration, the value would be too small, and vise versa.
2021    buffer_size + 4
2022}
2023
2024#[cfg(test)]
2025mod tests {
2026    use std::{
2027        collections::HashMap,
2028        io::ErrorKind,
2029        net::{IpAddr, Ipv4Addr, SocketAddr, SocketAddrV4},
2030        path::PathBuf,
2031        sync::{Arc, OnceLock},
2032        time::Duration,
2033    };
2034
2035    use bytes::{BufMut as _, Bytes};
2036    use bytesize::ByteSize;
2037    use metrics::{Key, Label};
2038    use saluki_config::ConfigurationLoader;
2039    use saluki_context::{origin::RawOrigin, ContextResolverBuilder, TagsResolverBuilder};
2040    use saluki_core::{
2041        components::ComponentContext,
2042        pooling::{helpers::get_pooled_object_via_builder, ObjectPool as _},
2043    };
2044    use saluki_env::workload::{CaptureEntityResolver, EntityId};
2045    use saluki_io::{
2046        buf::{BytesBuffer, FixedSizeVec},
2047        deser::codec::dogstatsd::{DogStatsDCodec, DogStatsDCodecConfiguration, ParsedPacket},
2048        net::{ConnectionAddress, ListenAddress, ProcessCredentials, ProcessIdentity},
2049    };
2050    use saluki_metrics::test::TestRecorder;
2051    use serde_json::json;
2052    use stringtheory::MetaString;
2053    use tokio::{net::UdpSocket, sync::mpsc, time::timeout};
2054
2055    use super::{
2056        build_io_buffer_pool, default_buffer_size, default_windows_pipe_security_descriptor,
2057        filters::EnablePayloadsFilter,
2058        forwarder::{
2059            ConnectedPacketForwarder, ForwardPacket, PacketForwarder, PacketForwarderTarget, FORWARDER_QUEUE_CAPACITY,
2060        },
2061        handle_frame, handle_metric_packet,
2062        metrics::build_metrics,
2063        origin_detection_failed_for_telemetry, resolve_capture_container_id, ContextResolvers, DogStatsDConfiguration,
2064        DOGSTATSD_CAPTURE_DIR, MIN_CAPTURE_DEPTH,
2065    };
2066
2067    const LINUX_EAFNOSUPPORT: i32 = 97;
2068    const MACOS_EAFNOSUPPORT: i32 = 47;
2069
2070    fn is_ipv6_unavailable_error(error: &std::io::Error) -> bool {
2071        matches!(error.kind(), ErrorKind::AddrNotAvailable | ErrorKind::Unsupported)
2072            || matches!(error.raw_os_error(), Some(LINUX_EAFNOSUPPORT | MACOS_EAFNOSUPPORT))
2073    }
2074
2075    fn test_component_context() -> ComponentContext {
2076        ComponentContext::test_source("dogstatsd_test")
2077    }
2078
2079    #[derive(Default)]
2080    struct CaptureTestEntityResolver {
2081        pid_map: HashMap<u32, EntityId>,
2082    }
2083
2084    impl CaptureTestEntityResolver {
2085        fn with_pid_mapping(process_id: u32, entity_id: EntityId) -> Self {
2086            let mut pid_map = HashMap::new();
2087            pid_map.insert(process_id, entity_id);
2088            Self { pid_map }
2089        }
2090    }
2091
2092    impl CaptureEntityResolver for CaptureTestEntityResolver {
2093        fn resolve_container_entity_for_live_pid(&self, process_id: u32) -> Option<EntityId> {
2094            self.pid_map.get(&process_id).cloned()
2095        }
2096    }
2097
2098    fn packet_forwarder_from_sender(
2099        target_port: u16, packets_tx: mpsc::Sender<ForwardPacket>, metrics: super::metrics::Metrics,
2100    ) -> PacketForwarder {
2101        let mut forwarder =
2102            PacketForwarderTarget::new(MetaString::from_static("127.0.0.1"), target_port).to_forwarder(metrics);
2103        forwarder.connected = Arc::new(OnceLock::from(packets_tx));
2104        forwarder
2105    }
2106
2107    fn processed_metric_key(listener_type: &'static str, origin: Option<&str>) -> Key {
2108        let mut labels = vec![
2109            Label::from_static_parts("component_id", "dogstatsd_test"),
2110            Label::from_static_parts("component_type", "source"),
2111            Label::from_static_parts("listener_type", listener_type),
2112            Label::from_static_parts("message_type", "metrics"),
2113        ];
2114        if let Some(origin) = origin {
2115            labels.push(Label::new("origin", origin.to_string()));
2116        }
2117
2118        Key::from_parts("component_events_received_total", labels)
2119    }
2120
2121    fn test_context_resolvers() -> ContextResolvers {
2122        let tags_resolver = TagsResolverBuilder::for_tests().build();
2123        let context_resolver = ContextResolverBuilder::for_tests()
2124            .with_tags_resolver(Some(tags_resolver.clone()))
2125            .build();
2126        ContextResolvers::manual(context_resolver.clone(), context_resolver, tags_resolver)
2127    }
2128
2129    #[test]
2130    fn origin_telemetry_does_not_resolve_origin_when_origin_detection_is_disabled() {
2131        let recorder = TestRecorder::default();
2132        let _recorder_guard = metrics::set_default_local_recorder(&recorder);
2133        let listen_addr = ListenAddress::Unixgram("/tmp/dsd.sock".into());
2134        let context = test_component_context();
2135        let metrics = build_metrics(&listen_addr, &context, true);
2136        let codec = DogStatsDCodec::from_configuration(DogStatsDCodecConfiguration::default());
2137        let mut context_resolvers = test_context_resolvers();
2138        let capture_entity_resolver = CaptureTestEntityResolver::with_pid_mapping(
2139            42,
2140            EntityId::from_local_data("ci-pid-container").expect("container entity"),
2141        );
2142        let peer_addr = ConnectionAddress::ProcessLike(ProcessIdentity::Credentials(ProcessCredentials {
2143            pid: 42,
2144            uid: 0,
2145            gid: 0,
2146        }));
2147
2148        let event = handle_frame(
2149            b"test_metric:1|c",
2150            &codec,
2151            &mut context_resolvers,
2152            &metrics,
2153            Some(&capture_entity_resolver),
2154            false,
2155            &peer_addr,
2156            EnablePayloadsFilter::default(),
2157            &[],
2158            &MetaString::from_static("default-host"),
2159        )
2160        .expect("frame should parse");
2161
2162        assert!(event.is_some());
2163        assert_eq!(
2164            recorder.counter(processed_metric_key("unixgram", Some("container_id://pid-container"))),
2165            None
2166        );
2167        assert_eq!(recorder.counter(processed_metric_key("unixgram", Some(""))), Some(1));
2168    }
2169
2170    #[test]
2171    fn origin_telemetry_records_resolved_origin_when_origin_detection_is_enabled() {
2172        let recorder = TestRecorder::default();
2173        let _recorder_guard = metrics::set_default_local_recorder(&recorder);
2174        let listen_addr = ListenAddress::Unixgram("/tmp/dsd.sock".into());
2175        let context = test_component_context();
2176        let metrics = build_metrics(&listen_addr, &context, true);
2177        let codec = DogStatsDCodec::from_configuration(DogStatsDCodecConfiguration::default());
2178        let mut context_resolvers = test_context_resolvers();
2179        let capture_entity_resolver = CaptureTestEntityResolver::with_pid_mapping(
2180            42,
2181            EntityId::from_local_data("ci-pid-container").expect("container entity"),
2182        );
2183        let peer_addr = ConnectionAddress::ProcessLike(ProcessIdentity::Credentials(ProcessCredentials {
2184            pid: 42,
2185            uid: 0,
2186            gid: 0,
2187        }));
2188
2189        let event = handle_frame(
2190            b"test_metric:1|c",
2191            &codec,
2192            &mut context_resolvers,
2193            &metrics,
2194            Some(&capture_entity_resolver),
2195            true,
2196            &peer_addr,
2197            EnablePayloadsFilter::default(),
2198            &[],
2199            &MetaString::from_static("default-host"),
2200        )
2201        .expect("frame should parse");
2202
2203        assert!(event.is_some());
2204        assert_eq!(
2205            recorder.counter(processed_metric_key("unixgram", Some("container_id://pid-container"))),
2206            Some(1)
2207        );
2208        assert_eq!(recorder.counter(processed_metric_key("unixgram", Some(""))), Some(0));
2209    }
2210
2211    #[test]
2212    fn no_metrics_when_interner_full_allocations_disallowed() {
2213        // We're specifically testing here that when we don't allow outside allocations, we should not be able to
2214        // resolve a context if the interner is full. A no-op interner has the smallest possible size, so that's going
2215        // to assure we can't intern anything... but we also need a string (name or one of the tags) that can't be
2216        // _inlined_ either, since that will get around the interner being full.
2217        //
2218        // We set our metric name to be longer than 31 bytes (the inlining limit) to ensure this.
2219
2220        let codec = DogStatsDCodec::from_configuration(DogStatsDCodecConfiguration::default());
2221        let tags_resolver = TagsResolverBuilder::for_tests().build();
2222        let context_resolver = ContextResolverBuilder::for_tests()
2223            .with_heap_allocations(false)
2224            .with_tags_resolver(Some(tags_resolver.clone()))
2225            .build();
2226        let mut context_resolvers = ContextResolvers::manual(context_resolver.clone(), context_resolver, tags_resolver);
2227        let peer_addr = ConnectionAddress::from("1.1.1.1:1234".parse::<SocketAddr>().unwrap());
2228
2229        let input = "big_metric_name_that_cant_possibly_be_inlined:1|c|#tag1:value1,tag2:value2,tag3:value3";
2230
2231        let Ok(ParsedPacket::Metric(packet)) = codec.decode_packet(input.as_bytes()) else {
2232            panic!("Failed to parse packet.");
2233        };
2234
2235        let maybe_metric = handle_metric_packet(
2236            packet,
2237            &mut context_resolvers,
2238            &peer_addr,
2239            &[],
2240            &MetaString::from_static("default-host"),
2241        );
2242        assert!(maybe_metric.is_none());
2243    }
2244
2245    #[test]
2246    fn metric_host_tag_disambiguates_contexts_without_remaining_tag() {
2247        let codec = DogStatsDCodec::from_configuration(DogStatsDCodecConfiguration::default());
2248        let mut context_resolvers = test_context_resolvers();
2249        let peer_addr = ConnectionAddress::from("1.1.1.1:1234".parse::<SocketAddr>().unwrap());
2250        let default_hostname = MetaString::from_static("default-host");
2251
2252        let packets = [
2253            ("unset", b"test_metric_name:1|g".as_slice(), "default-host"),
2254            ("empty", b"test_metric_name:2|g|#host:".as_slice(), ""),
2255            (
2256                "explicit_default",
2257                b"test_metric_name:3|g|#host:default-host".as_slice(),
2258                "default-host",
2259            ),
2260            (
2261                "custom",
2262                b"test_metric_name:4|g|#host:custom-host".as_slice(),
2263                "custom-host",
2264            ),
2265        ];
2266
2267        let mut metrics = Vec::new();
2268        for (case, raw, expected_host) in packets {
2269            let Ok(ParsedPacket::Metric(packet)) = codec.decode_packet(raw) else {
2270                panic!("Failed to parse {case} packet.");
2271            };
2272            let metric = handle_metric_packet(packet, &mut context_resolvers, &peer_addr, &[], &default_hostname)
2273                .unwrap_or_else(|| panic!("{case} metric should resolve"));
2274
2275            assert_eq!(metric.context().host(), Some(expected_host), "{case} context host");
2276            assert!(metric.context().tags().into_iter().all(|tag| tag.name() != "host"));
2277            metrics.push(metric);
2278        }
2279
2280        assert_eq!(metrics[0].context(), metrics[2].context());
2281        assert_ne!(metrics[0].context(), metrics[1].context());
2282        assert_ne!(metrics[0].context(), metrics[3].context());
2283        assert_ne!(metrics[1].context(), metrics[3].context());
2284    }
2285
2286    #[test]
2287    fn metric_with_additional_tags() {
2288        let codec = DogStatsDCodec::from_configuration(DogStatsDCodecConfiguration::default());
2289        let tags_resolver = TagsResolverBuilder::for_tests().build();
2290        let context_resolver = ContextResolverBuilder::for_tests()
2291            .with_heap_allocations(false)
2292            .with_tags_resolver(Some(tags_resolver.clone()))
2293            .build();
2294        let mut context_resolvers = ContextResolvers::manual(context_resolver.clone(), context_resolver, tags_resolver);
2295        let peer_addr = ConnectionAddress::from("1.1.1.1:1234".parse::<SocketAddr>().unwrap());
2296
2297        let existing_tags = ["tag1:value1", "tag2:value2", "tag3:value3"];
2298        let existing_tags_str = existing_tags.join(",");
2299
2300        let input = format!("test_metric_name:1|c|#{}", existing_tags_str);
2301        let additional_tags = [
2302            "tag4:value4".to_string(),
2303            "tag5:value5".to_string(),
2304            "tag6:value6".to_string(),
2305        ];
2306
2307        let Ok(ParsedPacket::Metric(packet)) = codec.decode_packet(input.as_bytes()) else {
2308            panic!("Failed to parse packet.");
2309        };
2310        let maybe_metric = handle_metric_packet(
2311            packet,
2312            &mut context_resolvers,
2313            &peer_addr,
2314            &additional_tags,
2315            &MetaString::from_static("default-host"),
2316        );
2317        assert!(maybe_metric.is_some());
2318
2319        let metric = maybe_metric.unwrap();
2320        let context = metric.context();
2321
2322        for tag in existing_tags {
2323            assert!(context.tags().has_tag(tag));
2324        }
2325
2326        for tag in additional_tags {
2327            assert!(context.tags().has_tag(tag));
2328        }
2329    }
2330
2331    fn deser_config(json: &str) -> DogStatsDConfiguration {
2332        serde_json::from_str(json).expect("failed to deserialize config")
2333    }
2334
2335    fn udp_listen_address() -> ListenAddress {
2336        ListenAddress::Udp(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 8125)))
2337    }
2338
2339    fn tcp_listen_address() -> ListenAddress {
2340        ListenAddress::Tcp(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 8125)))
2341    }
2342
2343    fn named_pipe_listen_address() -> ListenAddress {
2344        ListenAddress::named_pipe_with_input_buffer_size(
2345            "datadog-dogstatsd",
2346            default_windows_pipe_security_descriptor(),
2347            default_buffer_size() as u32,
2348        )
2349    }
2350
2351    #[test]
2352    fn build_addresses_includes_named_pipe_when_configured() {
2353        let config = deser_config(
2354            r#"{
2355                "dogstatsd_port": 0,
2356                "dogstatsd_pipe_name": "datadog-dogstatsd"
2357            }"#,
2358        );
2359
2360        let addresses = config.build_addresses(None);
2361
2362        assert_eq!(addresses, vec![named_pipe_listen_address()]);
2363    }
2364
2365    #[test]
2366    fn build_addresses_uses_dogstatsd_buffer_size_for_named_pipe_input_buffer() {
2367        let config = deser_config(
2368            r#"{
2369                "dogstatsd_port": 0,
2370                "dogstatsd_pipe_name": "datadog-dogstatsd",
2371                "dogstatsd_buffer_size": 16384
2372            }"#,
2373        );
2374
2375        let addresses = config.build_addresses(None);
2376
2377        let [ListenAddress::NamedPipe { input_buffer_size, .. }] = addresses.as_slice() else {
2378            panic!("expected only a named pipe listen address, got {addresses:?}");
2379        };
2380        assert_eq!(*input_buffer_size, Some(16_384));
2381    }
2382
2383    #[test]
2384    fn eol_required_matches_named_pipe_listener_type() {
2385        let config = deser_config(r#"{"dogstatsd_eol_required": ["named_pipe"]}"#);
2386        let eol_required = config.eol_required();
2387
2388        assert!(eol_required.for_listener(&named_pipe_listen_address()));
2389        assert!(!eol_required.for_listener(&udp_listen_address()));
2390        assert!(!eol_required.for_listener(&tcp_listen_address()));
2391    }
2392
2393    #[test]
2394    fn interner_size_defaults_to_2mib() {
2395        let config = deser_config("{}");
2396        assert_eq!(config.effective_context_string_interner_bytes(), ByteSize::mib(2));
2397    }
2398
2399    #[test]
2400    fn socket_receive_buffer_size_defaults_to_zero() {
2401        let config = deser_config("{}");
2402        assert_eq!(config.socket_receive_buffer_size, 0);
2403    }
2404
2405    #[test]
2406    fn socket_receive_buffer_size_from_config() {
2407        let config = deser_config(r#"{"dogstatsd_so_rcvbuf": 131072}"#);
2408        assert_eq!(config.socket_receive_buffer_size, 131_072);
2409    }
2410
2411    #[test]
2412    fn stream_log_too_big_defaults_to_false() {
2413        let config = deser_config("{}");
2414        assert!(!config.stream_log_too_big);
2415    }
2416
2417    #[test]
2418    fn stream_log_too_big_from_config() {
2419        let config = deser_config(r#"{"dogstatsd_stream_log_too_big": true}"#);
2420        assert!(config.stream_log_too_big);
2421    }
2422
2423    #[test]
2424    fn disable_verbose_logs_defaults_to_false() {
2425        let config = deser_config("{}");
2426        assert!(!config.disable_verbose_logs);
2427    }
2428
2429    #[test]
2430    fn disable_verbose_logs_from_config() {
2431        let config = deser_config(r#"{"dogstatsd_disable_verbose_logs": true}"#);
2432        assert!(config.disable_verbose_logs);
2433    }
2434
2435    #[test]
2436    fn statsd_forward_defaults_disabled() {
2437        let config = deser_config("{}");
2438        assert!(config.statsd_forward_host.is_none());
2439        assert_eq!(config.statsd_forward_port, 0);
2440        assert!(config.statsd_forward_target().is_none());
2441    }
2442
2443    #[test]
2444    fn statsd_forward_empty_host_disabled() {
2445        let config = deser_config(r#"{"statsd_forward_host": "", "statsd_forward_port": 9125}"#);
2446        assert!(config.statsd_forward_host.is_none());
2447        assert!(config.statsd_forward_target().is_none());
2448    }
2449
2450    #[test]
2451    fn statsd_forward_zero_port_disabled() {
2452        let config = deser_config(r#"{"statsd_forward_host": "127.0.0.1", "statsd_forward_port": 0}"#);
2453        assert_eq!(config.statsd_forward_host.as_deref(), Some("127.0.0.1"));
2454        assert!(config.statsd_forward_target().is_none());
2455    }
2456
2457    #[test]
2458    fn statsd_forward_host_and_port_enabled() {
2459        let config = deser_config(r#"{"statsd_forward_host": "127.0.0.1", "statsd_forward_port": 9125}"#);
2460        let (host, port) = config.statsd_forward_target().expect("forwarding should be enabled");
2461        assert_eq!(host.as_ref(), "127.0.0.1");
2462        assert_eq!(port, 9125);
2463    }
2464
2465    #[test]
2466    fn statsd_forward_invalid_target_still_builds_forwarder_handle() {
2467        let config = deser_config(r#"{"statsd_forward_host": "not a valid host", "statsd_forward_port": 9125}"#);
2468        assert!(config.packet_forwarder_target().is_some());
2469    }
2470
2471    #[tokio::test]
2472    async fn packet_forwarder_sends_payload_bytes() {
2473        let receiver = UdpSocket::bind("127.0.0.1:0").await.expect("receiver should bind");
2474        let receiver_addr = receiver.local_addr().expect("receiver should have an address");
2475        let forwarder = ConnectedPacketForwarder::connect("127.0.0.1", receiver_addr.port())
2476            .await
2477            .expect("forwarder should connect");
2478        let payload = b"daemon:666|g|#sometag1:somevalue1,sometag2:somevalue2";
2479
2480        let recorder = TestRecorder::default();
2481        let _recorder_guard = metrics::set_default_local_recorder(&recorder);
2482        let listen_addr = ListenAddress::Udp(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 8125)));
2483        let context = test_component_context();
2484        let metrics = build_metrics(&listen_addr, &context, false);
2485        let (packets_tx, packets_rx) = mpsc::channel(1);
2486        let worker = tokio::spawn(forwarder.run(packets_rx, metrics.clone()));
2487        let packet_forwarder = packet_forwarder_from_sender(receiver_addr.port(), packets_tx, metrics);
2488
2489        packet_forwarder.forward(Bytes::copy_from_slice(payload)).await;
2490
2491        let mut actual = [0u8; 128];
2492        let (received_len, _) = timeout(Duration::from_secs(1), receiver.recv_from(&mut actual))
2493            .await
2494            .expect("receive should not time out")
2495            .expect("receiver should receive payload");
2496
2497        assert_eq!(&actual[..received_len], payload);
2498        assert_eq!(
2499            recorder.counter((
2500                "component_packets_forwarded_total",
2501                &[
2502                    ("component_id", "dogstatsd_test"),
2503                    ("component_type", "source"),
2504                    ("listener_type", "udp"),
2505                    ("state", "ok"),
2506                ]
2507            )),
2508            Some(1)
2509        );
2510        assert_eq!(
2511            recorder.counter((
2512                "component_bytes_forwarded_total",
2513                &[
2514                    ("component_id", "dogstatsd_test"),
2515                    ("component_type", "source"),
2516                    ("listener_type", "udp"),
2517                ]
2518            )),
2519            Some(payload.len() as u64)
2520        );
2521        worker.abort();
2522    }
2523
2524    #[tokio::test]
2525    async fn packet_forwarder_sends_payload_bytes_to_ipv6_target() {
2526        let receiver = match UdpSocket::bind("[::1]:0").await {
2527            Ok(receiver) => receiver,
2528            Err(e) if is_ipv6_unavailable_error(&e) => return,
2529            Err(e) => panic!("receiver should bind: {e}"),
2530        };
2531        let receiver_addr = receiver.local_addr().expect("receiver should have an address");
2532        let forwarder = ConnectedPacketForwarder::connect("::1", receiver_addr.port())
2533            .await
2534            .expect("forwarder should connect");
2535        let payload = b"daemon:666|g|#ip:6";
2536
2537        let recorder = TestRecorder::default();
2538        let _recorder_guard = metrics::set_default_local_recorder(&recorder);
2539        let listen_addr = ListenAddress::Udp(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 8125)));
2540        let context = test_component_context();
2541        let metrics = build_metrics(&listen_addr, &context, false);
2542        let (packets_tx, packets_rx) = mpsc::channel(1);
2543        let worker = tokio::spawn(forwarder.run(packets_rx, metrics.clone()));
2544        let packet_forwarder = packet_forwarder_from_sender(receiver_addr.port(), packets_tx, metrics);
2545
2546        packet_forwarder.forward(Bytes::copy_from_slice(payload)).await;
2547
2548        let mut actual = [0u8; 128];
2549        let (received_len, _) = timeout(Duration::from_secs(1), receiver.recv_from(&mut actual))
2550            .await
2551            .expect("receive should not time out")
2552            .expect("receiver should receive payload");
2553
2554        assert_eq!(&actual[..received_len], payload);
2555        worker.abort();
2556    }
2557
2558    #[tokio::test]
2559    async fn packet_forwarder_waits_when_queue_is_full() {
2560        let recorder = TestRecorder::default();
2561        let _recorder_guard = metrics::set_default_local_recorder(&recorder);
2562        let listen_addr = ListenAddress::Udp(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 8125)));
2563        let context = test_component_context();
2564        let metrics = build_metrics(&listen_addr, &context, false);
2565        let (packets_tx, _packets_rx) = mpsc::channel(FORWARDER_QUEUE_CAPACITY);
2566        let packet_forwarder = packet_forwarder_from_sender(9125, packets_tx, metrics);
2567
2568        for _ in 0..FORWARDER_QUEUE_CAPACITY {
2569            packet_forwarder.forward(Bytes::from_static(b"queued:1|c")).await;
2570        }
2571
2572        assert!(
2573            timeout(
2574                Duration::from_millis(100),
2575                packet_forwarder.forward(Bytes::from_static(b"blocked:1|c")),
2576            )
2577            .await
2578            .is_err(),
2579            "forwarding should wait for queue capacity instead of dropping"
2580        );
2581    }
2582
2583    #[tokio::test]
2584    async fn packet_forwarder_send_error_increments_error_telemetry() {
2585        let recorder = TestRecorder::default();
2586        let _recorder_guard = metrics::set_default_local_recorder(&recorder);
2587        let listen_addr = ListenAddress::Udp(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 8125)));
2588        let context = test_component_context();
2589        let metrics = build_metrics(&listen_addr, &context, false);
2590        let socket = UdpSocket::bind("127.0.0.1:0").await.expect("socket should bind");
2591        let forwarder = ConnectedPacketForwarder {
2592            socket,
2593            target: SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 9125)),
2594        };
2595        let (packets_tx, packets_rx) = mpsc::channel(1);
2596        let worker = tokio::spawn(forwarder.run(packets_rx, metrics.clone()));
2597        let packet_forwarder = packet_forwarder_from_sender(9125, packets_tx, metrics);
2598
2599        packet_forwarder.forward(Bytes::from_static(b"daemon:666|g")).await;
2600
2601        let deadline = tokio::time::Instant::now() + Duration::from_secs(1);
2602        loop {
2603            if recorder.counter((
2604                "component_packets_forwarded_total",
2605                &[
2606                    ("component_id", "dogstatsd_test"),
2607                    ("component_type", "source"),
2608                    ("listener_type", "udp"),
2609                    ("state", "error"),
2610                ],
2611            )) == Some(1)
2612            {
2613                break;
2614            }
2615
2616            assert!(
2617                tokio::time::Instant::now() < deadline,
2618                "forwarding error telemetry should be recorded"
2619            );
2620            tokio::time::sleep(Duration::from_millis(10)).await;
2621        }
2622        worker.abort();
2623    }
2624
2625    #[test]
2626    fn unsupported_platform_process_credentials_do_not_count_as_origin_detection_telemetry_errors() {
2627        let peer_addr = ConnectionAddress::ProcessLike(ProcessIdentity::Error(
2628            saluki_io::net::ProcessCredentialsError::UnsupportedPlatform,
2629        ));
2630
2631        assert!(!origin_detection_failed_for_telemetry(true, 1, &peer_addr));
2632    }
2633
2634    #[test]
2635    fn invalid_process_credentials_count_as_origin_detection_telemetry_errors() {
2636        let peer_addr = ConnectionAddress::ProcessLike(ProcessIdentity::Error(
2637            saluki_io::net::ProcessCredentialsError::InvalidCredentials,
2638        ));
2639
2640        assert!(origin_detection_failed_for_telemetry(true, 1, &peer_addr));
2641    }
2642
2643    #[test]
2644    fn autoscale_udp_listeners_defaults_to_false() {
2645        let config = deser_config("{}");
2646        assert!(!config.autoscale_udp_listeners);
2647        assert!(config.udp_streams_to_yield().is_none());
2648    }
2649
2650    #[test]
2651    fn effective_max_buffer_count_never_below_baseline() {
2652        // A legacy config that only raised `dogstatsd_buffer_count` keeps its full capacity rather than being capped
2653        // to the `dogstatsd_buffer_count_max` default.
2654        let legacy = deser_config(r#"{"dogstatsd_buffer_count": 1024}"#);
2655        assert_eq!(legacy.effective_max_buffer_count(), 1024);
2656
2657        // An explicit maximum above the baseline is honored as-is.
2658        let explicit = deser_config(r#"{"dogstatsd_buffer_count": 128, "dogstatsd_buffer_count_max": 512}"#);
2659        assert_eq!(explicit.effective_max_buffer_count(), 512);
2660
2661        // A maximum below the baseline is treated as equal to the baseline.
2662        let below = deser_config(r#"{"dogstatsd_buffer_count": 200, "dogstatsd_buffer_count_max": 64}"#);
2663        assert_eq!(below.effective_max_buffer_count(), 200);
2664    }
2665
2666    #[tokio::test]
2667    async fn dogstatsd_io_buffer_pool_grows_on_demand_until_limit() {
2668        let min_buffers = 2;
2669        let max_buffers = 3;
2670        let (pool, shrinker) = build_io_buffer_pool(min_buffers, max_buffers, default_buffer_size());
2671
2672        let mut initial_buffers = Vec::with_capacity(min_buffers);
2673        for _ in 0..min_buffers {
2674            initial_buffers.push(
2675                timeout(Duration::from_secs(1), pool.acquire())
2676                    .await
2677                    .expect("initial buffer should be available"),
2678            );
2679        }
2680        let on_demand_buffer = timeout(Duration::from_secs(1), pool.acquire())
2681            .await
2682            .expect("pool should grow on demand before hitting the limit");
2683
2684        let capped_acquire = timeout(Duration::from_millis(25), pool.acquire()).await;
2685        assert!(capped_acquire.is_err(), "pool should wait once it reaches the limit");
2686
2687        drop(initial_buffers.pop().expect("initial buffer should still be held"));
2688        timeout(Duration::from_secs(1), pool.acquire())
2689            .await
2690            .expect("returned buffer should unblock acquisition");
2691
2692        drop(on_demand_buffer);
2693        drop(shrinker);
2694    }
2695
2696    #[test]
2697    #[cfg(target_os = "linux")]
2698    fn autoscale_udp_listeners_from_config_linux() {
2699        let config = deser_config(r#"{"dogstatsd_autoscale_udp_listeners": true}"#);
2700        assert!(config.autoscale_udp_listeners);
2701
2702        let streams = config
2703            .udp_streams_to_yield()
2704            .expect("autoscale yields at least 1 stream");
2705        let n = streams.get();
2706        assert!(
2707            (1..=4).contains(&n),
2708            "expected 1..=4 streams from vCPU formula, got {n}"
2709        );
2710    }
2711
2712    #[test]
2713    #[cfg(not(target_os = "linux"))]
2714    fn warns_for_uds_origin_detection_on_non_linux() {
2715        let config = deser_config(
2716            r#"{
2717                "dogstatsd_origin_detection": true,
2718                "dogstatsd_port": 0,
2719                "dogstatsd_socket": "/tmp/dsd.sock"
2720            }"#,
2721        );
2722        let addresses = config.build_addresses(None);
2723
2724        assert!(config.uds_origin_detection_unsupported_on_platform(&addresses));
2725    }
2726
2727    #[test]
2728    #[cfg(not(target_os = "linux"))]
2729    fn does_not_warn_for_udp_origin_detection_on_non_linux() {
2730        let config = deser_config(r#"{"dogstatsd_origin_detection": true}"#);
2731        let addresses = config.build_addresses(None);
2732
2733        assert!(!config.uds_origin_detection_unsupported_on_platform(&addresses));
2734    }
2735
2736    #[test]
2737    #[cfg(not(target_os = "linux"))]
2738    fn autoscale_udp_listeners_from_config_non_linux() {
2739        let config = deser_config(r#"{"dogstatsd_autoscale_udp_listeners": true}"#);
2740        assert!(config.autoscale_udp_listeners);
2741
2742        assert_eq!(None, config.udp_streams_to_yield());
2743    }
2744
2745    #[test]
2746    fn eol_required_defaults_to_no_listeners() {
2747        let config = deser_config("{}");
2748        let eol_required = config.eol_required();
2749
2750        assert!(!eol_required.for_listener(&udp_listen_address()));
2751        assert!(!eol_required.for_listener(&tcp_listen_address()));
2752    }
2753
2754    #[test]
2755    fn eol_required_matches_configured_listener_types() {
2756        let config = deser_config(r#"{"dogstatsd_eol_required": ["udp", "uds"]}"#);
2757        let eol_required = config.eol_required();
2758
2759        assert!(eol_required.for_listener(&udp_listen_address()));
2760        assert!(!eol_required.for_listener(&tcp_listen_address()));
2761
2762        #[cfg(unix)]
2763        {
2764            assert!(eol_required.for_listener(&ListenAddress::Unixgram("/tmp/dsd.sock".into())));
2765            assert!(eol_required.for_listener(&ListenAddress::Unix("/tmp/dsd-stream.sock".into())));
2766        }
2767    }
2768
2769    #[test]
2770    fn eol_required_accepts_space_separated_string() {
2771        let config = deser_config(r#"{"dogstatsd_eol_required": "udp uds"}"#);
2772        let eol_required = config.eol_required();
2773
2774        assert!(eol_required.for_listener(&udp_listen_address()));
2775    }
2776
2777    #[test]
2778    fn drops_full_named_pipe_buffer_without_newline() {
2779        let named_pipe_stream = named_pipe_listen_address();
2780        let mut buffer = get_pooled_object_via_builder::<_, BytesBuffer>(|| FixedSizeVec::with_capacity(8));
2781        buffer.put_slice(b"12345678");
2782
2783        assert!(super::should_drop_oversized_named_pipe_frame(
2784            &named_pipe_stream,
2785            &buffer
2786        ));
2787    }
2788
2789    #[test]
2790    fn keeps_named_pipe_partial_frame_when_buffer_has_capacity() {
2791        let named_pipe_stream = named_pipe_listen_address();
2792        let mut buffer = get_pooled_object_via_builder::<_, BytesBuffer>(|| FixedSizeVec::with_capacity(9));
2793        buffer.put_slice(b"12345678");
2794
2795        assert!(!super::should_drop_oversized_named_pipe_frame(
2796            &named_pipe_stream,
2797            &buffer
2798        ));
2799    }
2800
2801    #[test]
2802    fn keeps_full_named_pipe_buffer_with_newline() {
2803        let named_pipe_stream = named_pipe_listen_address();
2804        let mut buffer = get_pooled_object_via_builder::<_, BytesBuffer>(|| FixedSizeVec::with_capacity(8));
2805        buffer.put_slice(b"1234567\n");
2806
2807        assert!(!super::should_drop_oversized_named_pipe_frame(
2808            &named_pipe_stream,
2809            &buffer
2810        ));
2811    }
2812
2813    #[test]
2814    fn stream_log_too_big_warns_for_enabled_length_delimited_stream_invalid_frames() {
2815        let uds_stream = ListenAddress::Unix("/tmp/dsd-stream.sock".into());
2816        let named_pipe_stream = named_pipe_listen_address();
2817        let tcp_stream = ListenAddress::Tcp(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 8125)));
2818        let error = saluki_io::deser::framing::FramingError::InvalidFrame {
2819            frame_len: 8193,
2820            reason: "frame length exceeds buffer capacity",
2821        };
2822
2823        assert!(super::should_warn_stream_log_too_big(&uds_stream, &error, true));
2824        assert!(!super::should_warn_stream_log_too_big(&uds_stream, &error, false));
2825        assert!(!super::should_warn_stream_log_too_big(&named_pipe_stream, &error, true));
2826        assert!(!super::should_warn_stream_log_too_big(&tcp_stream, &error, true));
2827    }
2828
2829    #[test]
2830    fn interner_size_from_entry_count() {
2831        // A Core Agent migration config with entry count 4096 should yield 2 MiB, not 4096 bytes.
2832        let config = deser_config(r#"{"dogstatsd_string_interner_size": 4096}"#);
2833        assert_eq!(config.effective_context_string_interner_bytes(), ByteSize::mib(2));
2834    }
2835
2836    #[test]
2837    fn interner_size_from_explicit_bytes() {
2838        let config = deser_config(r#"{"dogstatsd_string_interner_size_bytes": 4194304}"#);
2839        assert_eq!(config.effective_context_string_interner_bytes(), ByteSize::b(4194304));
2840    }
2841
2842    #[test]
2843    fn interner_size_explicit_bytes_takes_priority() {
2844        let config = deser_config(
2845            r#"{"dogstatsd_string_interner_size": 4096, "dogstatsd_string_interner_size_bytes": 8388608}"#,
2846        );
2847        // The _bytes key (8 MiB) takes priority over the entry count.
2848        assert_eq!(config.effective_context_string_interner_bytes(), ByteSize::b(8388608));
2849    }
2850
2851    #[test]
2852    fn interner_size_custom_entry_count() {
2853        let config = deser_config(r#"{"dogstatsd_string_interner_size": 8192}"#);
2854        // 8192 entries * 512 bytes = 4 MiB
2855        assert_eq!(config.effective_context_string_interner_bytes(), ByteSize::mib(4));
2856    }
2857
2858    /// Asserts that two lists of ListenAddress are equivalent.
2859    fn address_list_eq(expected: &mut [ListenAddress], actual: &mut [ListenAddress]) -> Result<(), String> {
2860        if expected.len() != actual.len() {
2861            return Err(format!(
2862                "length mismatch: expected {} addresses, got {}",
2863                expected.len(),
2864                actual.len()
2865            ));
2866        }
2867
2868        expected.sort_by_key(|a| a.to_string());
2869        actual.sort_by_key(|a| a.to_string());
2870
2871        for (e, a) in expected.iter().zip(actual.iter()) {
2872            let (es, as_) = (e.to_string(), a.to_string());
2873            if es != as_ {
2874                return Err(format!("address mismatch: expected {}, got {}", es, as_));
2875            }
2876        }
2877
2878        Ok(())
2879    }
2880
2881    /// This test verifies that we didn't accidentally break the `build_addresses_no_listeners` helper function which
2882    /// would render all further tests useless.
2883    #[test]
2884    fn build_addresses_assertion_function_works() {
2885        let config = DogStatsDConfiguration {
2886            port: 0,
2887            tcp_port: 123,
2888            socket_path: None,
2889            socket_stream_path: None,
2890            non_local_traffic: false,
2891            ..Default::default()
2892        };
2893        let mut expected = vec![ListenAddress::Tcp(SocketAddr::V4(SocketAddrV4::new(
2894            // Close, but not quite! This is intentionally *not* 127.0.0.1 to test that the assertion will fail
2895            Ipv4Addr::new(127, 0, 0, 2),
2896            123,
2897        )))];
2898        let mut actual = config.build_addresses(None);
2899        assert!(address_list_eq(&mut expected, &mut actual).is_err())
2900    }
2901
2902    /// With all four listener gates off, `build_addresses` returns an empty Vec.
2903    #[test]
2904    fn build_addresses_no_listeners() {
2905        let config = DogStatsDConfiguration {
2906            port: 0,
2907            tcp_port: 0,
2908            socket_path: None,
2909            socket_stream_path: None,
2910            non_local_traffic: false,
2911            ..Default::default()
2912        };
2913        let mut expected = vec![];
2914        let mut actual = config.build_addresses(None);
2915        address_list_eq(&mut expected, &mut actual).unwrap();
2916    }
2917
2918    /// UDP port set, `non_local_traffic=false` -> UDP listener bound to `127.0.0.1`.
2919    #[test]
2920    fn build_addresses_udp_local_only() {
2921        let config = DogStatsDConfiguration {
2922            port: 8125,
2923            tcp_port: 0,
2924            socket_path: None,
2925            socket_stream_path: None,
2926            non_local_traffic: false,
2927            ..Default::default()
2928        };
2929        let mut expected = vec![ListenAddress::Udp(SocketAddr::V4(SocketAddrV4::new(
2930            Ipv4Addr::new(127, 0, 0, 1),
2931            8125,
2932        )))];
2933        let mut actual = config.build_addresses(None);
2934        address_list_eq(&mut expected, &mut actual).unwrap();
2935    }
2936
2937    /// UDP port set, `non_local_traffic=true` -> UDP listener bound to `0.0.0.0`.
2938    #[test]
2939    fn build_addresses_udp_non_local_only() {
2940        let config = DogStatsDConfiguration {
2941            port: 8125,
2942            tcp_port: 0,
2943            socket_path: None,
2944            socket_stream_path: None,
2945            non_local_traffic: true,
2946            ..Default::default()
2947        };
2948        let mut expected = vec![ListenAddress::Udp(SocketAddr::V4(SocketAddrV4::new(
2949            Ipv4Addr::new(0, 0, 0, 0),
2950            8125,
2951        )))];
2952        let mut actual = config.build_addresses(None);
2953        address_list_eq(&mut expected, &mut actual).unwrap();
2954    }
2955
2956    /// TCP port set, `non_local_traffic=false` -> TCP listener bound to `127.0.0.1`.
2957    #[test]
2958    fn build_addresses_tcp_local_only() {
2959        let config = DogStatsDConfiguration {
2960            port: 0,
2961            tcp_port: 9000,
2962            socket_path: None,
2963            socket_stream_path: None,
2964            non_local_traffic: false,
2965            ..Default::default()
2966        };
2967        let mut expected = vec![ListenAddress::Tcp(SocketAddr::V4(SocketAddrV4::new(
2968            Ipv4Addr::new(127, 0, 0, 1),
2969            9000,
2970        )))];
2971        let mut actual = config.build_addresses(None);
2972        address_list_eq(&mut expected, &mut actual).unwrap();
2973    }
2974
2975    /// TCP port set, `non_local_traffic=true` -> TCP listener bound to `0.0.0.0`.
2976    #[test]
2977    fn build_addresses_tcp_non_local_only() {
2978        let config = DogStatsDConfiguration {
2979            port: 0,
2980            tcp_port: 9000,
2981            socket_path: None,
2982            socket_stream_path: None,
2983            non_local_traffic: true,
2984            ..Default::default()
2985        };
2986        let mut expected = vec![ListenAddress::Tcp(SocketAddr::V4(SocketAddrV4::new(
2987            Ipv4Addr::new(0, 0, 0, 0),
2988            9000,
2989        )))];
2990        let mut actual = config.build_addresses(None);
2991        address_list_eq(&mut expected, &mut actual).unwrap();
2992    }
2993
2994    /// `socket_path` set -> a `Unixgram` address is produced with that path.
2995    #[test]
2996    fn build_addresses_unixgram_only() {
2997        let config = DogStatsDConfiguration {
2998            port: 0,
2999            tcp_port: 0,
3000            socket_path: Some("/tmp/dsd.sock".to_string()),
3001            socket_stream_path: None,
3002            non_local_traffic: false,
3003            ..Default::default()
3004        };
3005        let mut expected = vec![ListenAddress::Unixgram("/tmp/dsd.sock".into())];
3006        let mut actual = config.build_addresses(None);
3007        address_list_eq(&mut expected, &mut actual).unwrap();
3008    }
3009
3010    /// `socket_stream_path` set -> a `Unix` (stream) address is produced with that path.
3011    #[test]
3012    fn build_addresses_unix_stream_only() {
3013        let config = DogStatsDConfiguration {
3014            port: 0,
3015            tcp_port: 0,
3016            socket_path: None,
3017            socket_stream_path: Some("/tmp/dsd-stream.sock".to_string()),
3018            non_local_traffic: false,
3019            ..Default::default()
3020        };
3021        let mut expected = vec![ListenAddress::Unix("/tmp/dsd-stream.sock".into())];
3022        let mut actual = config.build_addresses(None);
3023        address_list_eq(&mut expected, &mut actual).unwrap();
3024    }
3025
3026    /// All four listener types enabled at once, with `non_local_traffic=true`.
3027    #[test]
3028    fn build_addresses_all_four_non_local() {
3029        let config = DogStatsDConfiguration {
3030            port: 8125,
3031            tcp_port: 9000,
3032            socket_path: Some("/tmp/dsd.sock".to_string()),
3033            socket_stream_path: Some("/tmp/dsd-stream.sock".to_string()),
3034            non_local_traffic: true,
3035            ..Default::default()
3036        };
3037        let mut expected = vec![
3038            ListenAddress::Udp(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::new(0, 0, 0, 0), 8125))),
3039            ListenAddress::Tcp(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::new(0, 0, 0, 0), 9000))),
3040            ListenAddress::Unixgram("/tmp/dsd.sock".into()),
3041            ListenAddress::Unix("/tmp/dsd-stream.sock".into()),
3042        ];
3043        let mut actual = config.build_addresses(None);
3044        address_list_eq(&mut expected, &mut actual).unwrap();
3045    }
3046
3047    /// All four listener types enabled at once, with `non_local_traffic=false`.
3048    #[test]
3049    fn build_addresses_all_four_local() {
3050        let config = DogStatsDConfiguration {
3051            port: 8125,
3052            tcp_port: 9000,
3053            socket_path: Some("/tmp/dsd.sock".to_string()),
3054            socket_stream_path: Some("/tmp/dsd-stream.sock".to_string()),
3055            non_local_traffic: false,
3056            ..Default::default()
3057        };
3058        let mut expected = vec![
3059            ListenAddress::Udp(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::new(127, 0, 0, 1), 8125))),
3060            ListenAddress::Tcp(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::new(127, 0, 0, 1), 9000))),
3061            ListenAddress::Unixgram("/tmp/dsd.sock".into()),
3062            ListenAddress::Unix("/tmp/dsd-stream.sock".into()),
3063        ];
3064        let mut actual = config.build_addresses(None);
3065        address_list_eq(&mut expected, &mut actual).unwrap();
3066    }
3067
3068    /// Passing `Some(ip)` to `build_addresses` with `non_local_traffic=false` -> both UDP and TCP
3069    /// bind to that IP. Includes a UDS datagram socket to confirm `bind_host` doesn't affect it.
3070    #[test]
3071    fn build_addresses_bind_host_applies_to_udp_and_tcp() {
3072        let config = DogStatsDConfiguration {
3073            port: 8125,
3074            tcp_port: 9000,
3075            socket_path: Some("/tmp/dsd.sock".to_string()),
3076            socket_stream_path: None,
3077            non_local_traffic: false,
3078            ..Default::default()
3079        };
3080        let bind_host = Some(IpAddr::V4(Ipv4Addr::new(192, 168, 1, 50)));
3081        let mut expected = vec![
3082            ListenAddress::Udp(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::new(192, 168, 1, 50), 8125))),
3083            ListenAddress::Tcp(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::new(192, 168, 1, 50), 9000))),
3084            ListenAddress::Unixgram("/tmp/dsd.sock".into()),
3085        ];
3086        let mut actual = config.build_addresses(bind_host);
3087        address_list_eq(&mut expected, &mut actual).unwrap();
3088    }
3089
3090    /// Passing `Some(ip)` to `build_addresses` with `non_local_traffic=true` -> both UDP and TCP
3091    /// bind to `0.0.0.0`; the `bind_host` parameter is ignored (precedence matches the Agent).
3092    /// Includes a UDS stream socket to confirm `bind_host` doesn't affect it.
3093    #[test]
3094    fn build_addresses_non_local_clobbers_bind_host() {
3095        let config = DogStatsDConfiguration {
3096            port: 8125,
3097            tcp_port: 9000,
3098            socket_path: None,
3099            socket_stream_path: Some("/tmp/dsd-stream.sock".to_string()),
3100            non_local_traffic: true,
3101            ..Default::default()
3102        };
3103        let bind_host = Some(IpAddr::V4(Ipv4Addr::new(192, 168, 1, 50)));
3104        let mut expected = vec![
3105            ListenAddress::Udp(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::new(0, 0, 0, 0), 8125))),
3106            ListenAddress::Tcp(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::new(0, 0, 0, 0), 9000))),
3107            ListenAddress::Unix("/tmp/dsd-stream.sock".into()),
3108        ];
3109        let mut actual = config.build_addresses(bind_host);
3110        address_list_eq(&mut expected, &mut actual).unwrap();
3111    }
3112
3113    #[test]
3114    fn non_finite_metric_values_are_silently_dropped() {
3115        // The Datadog Agent sends NaN gauges (for example, encode_ms.avg computed as 0.0/0.0 in Go).
3116        // FloatIter skips non-finite values with a debug log, so decode_packet returns Ok with
3117        // num_points == 0. handle_frame then returns Ok(None) for zero-point packets, which is
3118        // the existing silent-drop path (no warning emitted).
3119        let codec = DogStatsDCodec::from_configuration(DogStatsDCodecConfiguration::default());
3120        for input in &[b"my.gauge:NaN|g" as &[u8], b"my.gauge:inf|g", b"my.gauge:-inf|g"] {
3121            match codec.decode_packet(input).expect("should decode without error") {
3122                ParsedPacket::Metric(packet) => assert_eq!(
3123                    packet.num_points, 0,
3124                    "non-finite value should be dropped, leaving 0 valid points"
3125                ),
3126                _ => panic!("expected Metric packet"),
3127            }
3128        }
3129    }
3130
3131    #[tokio::test]
3132    async fn fix_empty_capture_path_sets_path_from_run_path() {
3133        const RUN_PATH: &str = "/my/little/run_path";
3134
3135        let base_config_values = json!({ "run_path": RUN_PATH });
3136        let (config, _) = ConfigurationLoader::for_tests(Some(base_config_values), None, false).await;
3137
3138        let dogstatsd_config = DogStatsDConfiguration::from_configuration(&config).expect("should deserialize");
3139
3140        let expected = PathBuf::from(RUN_PATH).join(DOGSTATSD_CAPTURE_DIR);
3141        assert_eq!(expected, dogstatsd_config.capture_path);
3142    }
3143
3144    #[tokio::test]
3145    async fn fix_empty_capture_path_keeps_explicit_path() {
3146        const RUN_PATH: &str = "/my/little/run_path";
3147        const CAPTURE_PATH: &str = "/custom/path/to/capture";
3148
3149        let base_config_values = json!({ "run_path": RUN_PATH, "dogstatsd_capture_path": CAPTURE_PATH });
3150        let (config, _) = ConfigurationLoader::for_tests(Some(base_config_values), None, false).await;
3151
3152        let dogstatsd_config = DogStatsDConfiguration::from_configuration(&config).expect("should deserialize");
3153
3154        assert_eq!(PathBuf::from(CAPTURE_PATH), dogstatsd_config.capture_path);
3155    }
3156
3157    #[tokio::test]
3158    async fn from_configuration_normalizes_capture_depth() {
3159        let cases = [
3160            (json!({}), MIN_CAPTURE_DEPTH),
3161            (json!({ "dogstatsd_capture_depth": 0 }), MIN_CAPTURE_DEPTH),
3162            (json!({ "dogstatsd_capture_depth": 2048 }), 2048),
3163        ];
3164
3165        for (base_config_values, expected_depth) in cases {
3166            let (config, _) = ConfigurationLoader::for_tests(Some(base_config_values), None, false).await;
3167            let dogstatsd_config = DogStatsDConfiguration::from_configuration(&config).expect("should deserialize");
3168
3169            assert_eq!(expected_depth, dogstatsd_config.capture_depth);
3170        }
3171    }
3172
3173    #[test]
3174    fn capture_entity_resolver_is_configured_separately_from_workload_provider() {
3175        let config =
3176            DogStatsDConfiguration::default().with_capture_entity_resolver(CaptureTestEntityResolver::default());
3177
3178        assert!(config.capture_entity_resolver.is_some());
3179        assert!(config.workload_provider.is_none());
3180    }
3181
3182    #[test]
3183    fn resolve_capture_container_id_uses_live_pid_mapping() {
3184        let capture_entity_resolver = CaptureTestEntityResolver::with_pid_mapping(
3185            42,
3186            EntityId::from_local_data("ci-pid-container").expect("container entity"),
3187        );
3188
3189        assert_eq!(
3190            resolve_capture_container_id(Some(&capture_entity_resolver), Some(42)),
3191            Some("container_id://pid-container".to_string())
3192        );
3193    }
3194
3195    #[test]
3196    fn build_capture_record_ignores_payload_local_data() {
3197        let record = super::build_capture_record(None, None, b"test.metric:1|c|c:ci-local-container\n");
3198
3199        assert_eq!(record.container_id, None);
3200        assert!(record.ancillary.is_empty());
3201    }
3202
3203    #[test]
3204    fn stream_capture_state_preserves_last_pid_without_new_creds() {
3205        let mut stream_capture = super::StreamCaptureState::new();
3206
3207        stream_capture.update_peer_metadata(&ConnectionAddress::ProcessLike(ProcessIdentity::Credentials(
3208            ProcessCredentials {
3209                pid: 42,
3210                uid: 0,
3211                gid: 0,
3212            },
3213        )));
3214        stream_capture.update_peer_metadata(&ConnectionAddress::ProcessLike(ProcessIdentity::Unavailable));
3215
3216        assert_eq!(stream_capture.last_pid, Some(42));
3217    }
3218
3219    #[test]
3220    fn apply_credentials_uses_live_pid_for_normal_packet() {
3221        let mut origin = RawOrigin::default();
3222        let creds = ProcessCredentials {
3223            pid: 12345,
3224            uid: 1000,
3225            gid: 1000,
3226        };
3227        super::apply_credentials_to_origin(&mut origin, &creds);
3228
3229        assert_eq!(origin.process_id(), Some(12345));
3230    }
3231
3232    #[test]
3233    fn apply_credentials_unpacks_captured_pid_when_replay_gid_present() {
3234        let mut origin = RawOrigin::default();
3235        let captured_pid: u32 = 99887766;
3236        let creds = ProcessCredentials {
3237            pid: 12345,        // our PID (irrelevant for replay)
3238            uid: captured_pid, // captured PID packed by the sender
3239            gid: super::REPLAY_CREDENTIALS_GID,
3240        };
3241        super::apply_credentials_to_origin(&mut origin, &creds);
3242
3243        assert_eq!(
3244            origin.process_id(),
3245            Some(super::origin::mark_replay_process_id(captured_pid))
3246        );
3247    }
3248}
3249
3250#[cfg(test)]
3251mod config_smoke {
3252    use datadog_agent_config_testing::config_registry::structs;
3253    use datadog_agent_config_testing::run_config_smoke_tests;
3254    use serde_json::json;
3255
3256    use super::DogStatsDConfiguration;
3257    use crate::config::{DatadogRemapper, KEY_ALIASES};
3258
3259    #[tokio::test]
3260    async fn smoke_test() {
3261        run_config_smoke_tests(
3262            structs::DOGSTATSD_CONFIGURATION,
3263            &[],
3264            json!({}),
3265            |cfg| {
3266                cfg.as_typed::<DogStatsDConfiguration>()
3267                    .expect("DogStatsDConfiguration should deserialize")
3268            },
3269            KEY_ALIASES,
3270            DatadogRemapper::from_env_vars,
3271        )
3272        .await
3273    }
3274}