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    io, mem,
13    num::NonZeroUsize,
14    path::PathBuf,
15    pin::Pin,
16    sync::{Arc, LazyLock, Mutex as StdMutex},
17    time::{Duration, SystemTime, UNIX_EPOCH},
18};
19
20use async_trait::async_trait;
21use bytes::{Buf, BufMut, Bytes};
22use bytesize::ByteSize;
23use saluki_common::sync::shutdown::{ShutdownCoordinator, ShutdownHandle};
24use saluki_core::{
25    accounting::{MemoryBounds, MemoryBoundsBuilder, MemoryLimiter, UsageExpr},
26    components::{sources::*, BuildContext},
27    data_model::{
28        event::{
29            eventd::EventD,
30            metric::{Metric, MetricMetadata, MetricOrigin},
31            service_check::ServiceCheck,
32            Event, EventType,
33        },
34        tags::{RawTags, RawTagsFilter},
35    },
36    pooling::{ElasticObjectPool, ObjectPool as _},
37    runtime::{
38        self,
39        state::{AcquireError, ResourceLease},
40        InitializationError, ShutdownStrategy, Supervisable, SupervisorFuture,
41    },
42    topology::{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 _, CollapsibleReadWriteIoBuffer as _, FixedSizeVec, ReadIoBuffer as _},
48    deser::{
49        codec::dogstatsd::*,
50        framing::{Framer as _, FramingError, LengthDelimitedFramer},
51    },
52    net::{listener::Listener, ConnectionAddress, ListenAddress, ProcessIdentity, SocketSpecification, Stream},
53};
54use snafu::{ResultExt as _, Snafu};
55use stringtheory::MetaString;
56use tokio::{
57    pin, select,
58    sync::{mpsc, oneshot, Mutex},
59    time::{interval, MissedTickBehavior},
60};
61use tracing::{debug, error, info, trace, warn};
62
63mod forwarder;
64use self::forwarder::{PacketForwarder, PacketForwarderTarget};
65
66mod framer;
67use self::framer::{get_framer, DsdFramer};
68use crate::sources::dogstatsd::tags::{WellKnownTags, WellKnownTagsFilterPredicate};
69
70mod filters;
71use self::filters::EnablePayloadsFilter;
72
73mod metrics;
74use self::metrics::{build_metrics, Metrics};
75
76mod replay;
77use self::replay::{CaptureRecord, CapturedTaggerHandle, TrafficCapture};
78pub use self::replay::{
79    DogStatsDCaptureAPIHandler, DogStatsDCaptureControl, DogStatsDReplayAPIHandler, DogStatsDReplayControl,
80    ReplaySession, TimestampResolution, TrafficCaptureReader, DEFAULT_REPLAY_LOOPS, REPLAY_CREDENTIALS_GID,
81};
82
83mod origin;
84pub use self::origin::OriginEnrichmentConfiguration;
85use self::origin::{
86    origin_from_event_packet, origin_from_metric_packet, origin_from_service_check_packet, DogStatsDOriginTagResolver,
87    ProcessOrigin,
88};
89
90mod resolver;
91use self::resolver::ContextResolvers;
92
93mod tags;
94
95#[derive(Debug, Snafu)]
96#[snafu(context(suffix(false)))]
97enum Error {
98    #[snafu(display("Failed to acquire {} listener: {}", listener_type, source))]
99    FailedToAcquireListener {
100        listener_type: &'static str,
101        source: AcquireError,
102    },
103
104    #[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."))]
105    NoListenersConfigured,
106
107    #[snafu(display("Could not resolve bind_host '{}': {}", host, source))]
108    UnresolvableBindHost { host: String, source: std::io::Error },
109
110    #[snafu(display("bind_host '{}' resolved to zero IP addresses.", host))]
111    BindHostHasNoAddresses { host: String },
112}
113
114/// Baseline byte cost per interner entry, used to convert the Core Agent's entry-count-based
115/// `dogstatsd_string_interner_size` to a byte size.
116///
117/// 4096 entries × 512 bytes = 2 MiB, matching ADP's previous default.
118const INTERNER_BASELINE_BYTES_PER_ENTRY: u64 = 512;
119const DOGSTATSD_LISTENER_WORKER_COUNT: usize = 1;
120const DOGSTATSD_PIPELINE_COUNT: usize = 1;
121const MIN_DOGSTATSD_WORKER_COUNT: usize = 2;
122
123fn default_decoder_worker_count(vcpus: usize) -> usize {
124    vcpus
125        .saturating_sub(DOGSTATSD_LISTENER_WORKER_COUNT + DOGSTATSD_PIPELINE_COUNT)
126        .max(MIN_DOGSTATSD_WORKER_COUNT)
127}
128
129/// Controls which payload types are forwarded to the backend.
130pub struct EnablePayloadsConfiguration {
131    /// Whether or not to enable sending series (counter/gauge/rate) payloads.
132    pub series: bool,
133
134    /// Whether or not to enable sending sketch (distribution) payloads.
135    pub sketches: bool,
136
137    /// Whether or not to enable sending event payloads.
138    pub events: bool,
139
140    /// Whether or not to enable sending service check payloads.
141    pub service_checks: bool,
142}
143
144const MIN_CAPTURE_DEPTH: usize = 1024;
145
146/// DogStatsD source.
147///
148/// Accepts metrics over TCP, UDP, or Unix Domain Sockets in the StatsD/DogStatsD format.
149pub struct DogStatsDConfiguration {
150    /// Hostname used when DogStatsD metrics do not carry an explicit `host:` tag.
151    pub default_hostname: MetaString,
152
153    /// The size of the buffer used to receive messages into, in bytes.
154    ///
155    /// Payloads can't exceed this size, or they will be truncated, leading to discarded messages.
156    pub buffer_size: usize,
157
158    /// The number of message buffers to allocate up front.
159    ///
160    /// This is the baseline pool size allocated at startup. The pool then grows on demand up to
161    /// `buffer_count_max` as active stream connections and datagram queues need additional buffers.
162    /// Higher values allocate more memory at startup but reduce on-demand allocations during bursts.
163    pub buffer_count: usize,
164
165    /// The maximum number of message buffers to allocate overall.
166    ///
167    /// The global pool starts at `buffer_count` buffers and grows on demand up to this limit. Active stream
168    /// connections use these buffers for reads, while connectionless listeners use them to queue received packets for
169    /// decoding. Increasing this value lets datagram listeners absorb larger bursts at the cost of up to one additional
170    /// `buffer_size` allocation per buffer. This limit bounds payload buffers, but not per-connection task and channel
171    /// bookkeeping. After a short grace period without pool growth, ADP releases idle buffers until the pool returns to
172    /// `buffer_count`.
173    ///
174    /// The pool never holds fewer buffers than `buffer_count`, so a value below the baseline is treated as equal to it.
175    pub buffer_count_max: usize,
176
177    /// The number of workers in the global pool that decodes connectionless DogStatsD packets.
178    ///
179    /// If set to `0`, the worker count is derived from the available vCPUs using the Core Agent's default formula.
180    /// Positive values force an exact worker count. Higher values can improve throughput when decoding is CPU-bound,
181    /// but also increase task scheduling and per-worker event buffering overhead.
182    pub workers_count: usize,
183
184    /// The port to listen on in UDP mode.
185    ///
186    /// If set to `0`, UDP isn't used.
187    pub port: u16,
188
189    /// The size of the receive buffer requested for each DogStatsD socket, in bytes.
190    ///
191    /// Applies to the UDP, TCP, and UDS sockets alike. If set to `0`, the OS default is used.
192    pub socket_receive_buffer_size: usize,
193
194    /// The port to listen on in TCP mode.
195    ///
196    /// If set to `0`, TCP isn't used.
197    pub tcp_port: u16,
198
199    /// The host to forward framed DogStatsD messages to over UDP.
200    ///
201    /// Forwarding is enabled only when this value is set and `statsd_forward_port` is non-zero. Setup failures
202    /// are logged, and send failures are tracked through telemetry.
203    pub statsd_forward_host: Option<MetaString>,
204
205    /// The port to forward framed DogStatsD messages to over UDP.
206    ///
207    /// Forwarding is enabled only when this value is non-zero and `statsd_forward_host` is set.
208    pub statsd_forward_port: u16,
209
210    /// The Unix domain socket path to listen on, in datagram mode.
211    ///
212    /// If not set, UDS (in datagram mode) isn't used.
213    pub socket_path: Option<String>,
214
215    /// The Unix domain socket path to listen on, in stream mode.
216    ///
217    /// If not set, UDS (in stream mode) isn't used.
218    pub socket_stream_path: Option<String>,
219
220    /// Controls whether ADP logs oversized DogStatsD stream frames.
221    ///
222    /// When set to `true`, ADP emits a warning when a UDS stream frame exceeds the
223    /// configured DogStatsD buffer size. The frame is still rejected either way.
224    ///
225    /// Enable this when diagnosing clients that send oversized UDS stream frames.
226    pub stream_log_too_big: bool,
227
228    /// The Windows named pipe name to listen on.
229    ///
230    /// If set, ADP listens for DogStatsD stream traffic on `\\.\pipe\<name>` on Windows.
231    /// The listener is unsupported on non-Windows platforms.
232    pub pipe_name: Option<String>,
233
234    /// Windows named pipe security descriptor.
235    ///
236    /// This SDDL descriptor is applied when creating the named pipe listener.
237    pub windows_pipe_security_descriptor: String,
238
239    /// Whether ADP lowers DogStatsD parse-failure logs to debug level.
240    ///
241    /// When set to `true`, invalid metrics, events, and service checks still increment decode-failure telemetry, but
242    /// their parse-failure logs are emitted at debug level instead of warning level. Enable this to suppress noisy
243    /// parse-error logs from misbehaving clients.
244    pub disable_verbose_logs: bool,
245
246    /// Listener types that require DogStatsD messages to be newline-terminated.
247    ///
248    /// Valid values are `udp`, `uds`, and `named_pipe`. Invalid values are ignored.
249    ///
250    /// An empty list accepts the final message without a newline.
251    pub eol_required: Vec<String>,
252
253    /// The host address to bind DogStatsD UDP and TCP listeners to.
254    ///
255    /// When set, UDP and TCP listeners bind to this address. Accepts either an IP literal (for example,
256    /// `192.168.1.50`, `::1`) or a hostname that resolves via DNS (for example, `agent.internal`).
257    /// Ignored when `non_local_traffic` is `true`, and binds to `127.0.0.1` when unset.
258    pub bind_host: Option<String>,
259
260    /// Whether or not to listen for non-local traffic.
261    ///
262    /// If set to `true`, the UDP and TCP listeners bind to `0.0.0.0` and accept traffic from any interface. Otherwise,
263    /// they bind to `bind_host`, or `127.0.0.1` if `bind_host` isn't set.
264    pub non_local_traffic: bool,
265
266    /// Whether to autoscale UDP stream handlers using `SO_REUSEPORT`.
267    ///
268    /// When enabled on Linux, the DogStatsD source binds multiple UDP sockets to the configured port with
269    /// `SO_REUSEPORT`, allowing the kernel to load-balance incoming datagrams across independent stream handler
270    /// tasks. The number of sockets scales with available vCPUs: one stream handler base, plus one additional
271    /// per 8 vCPUs, capped at 4 total.
272    ///
273    /// Has no effect on non-Linux platforms because `SO_REUSEPORT` doesn't provide kernel-level load balancing
274    /// there; a warning is logged at startup if enabled outside of Linux.
275    ///
276    /// Enable this on multi-vCPU Linux deployments where UDP DogStatsD throughput is bottlenecked on a single
277    /// receive task.
278    pub autoscale_udp_listeners: bool,
279
280    /// Whether or not to allow heap allocations when resolving contexts.
281    ///
282    /// When resolving contexts during parsing, the metric name and tags are interned to reduce memory usage. The
283    /// interner has a fixed size, however, which means some strings can fail to be interned if the interner is full.
284    /// When set to `true`, we allow these strings to be allocated on the heap like normal, but this can lead to
285    /// increased (unbounded) memory usage. When set to `false`, if the metric name and all of its tags can't be
286    /// interned, the metric is skipped.
287    pub allow_context_heap_allocations: bool,
288
289    /// Whether or not to enable support for no-aggregation pipelines.
290    ///
291    /// When enabled, this influences how metrics are parsed, specifically around user-provided metric timestamps. When
292    /// metric timestamps are present, it's used as a signal to any aggregation transforms that the metric shouldn't
293    /// be aggregated.
294    pub no_aggregation_pipeline_support: bool,
295
296    /// Number of entries for the string interner, as interpreted by the Core Datadog Agent.
297    ///
298    /// When `context_string_interner_size_bytes` isn't set, this value is multiplied by 512 bytes per entry to
299    /// derive the interner byte size. This provides backwards compatibility for customers migrating configurations
300    /// from the Core Agent, where this setting represents an entry count rather than a byte size.
301    pub context_string_interner_entry_count: u64,
302
303    /// Total size of the string interner used for contexts, in bytes.
304    ///
305    /// When set, this takes priority over `context_string_interner_entry_count`. This controls the amount of memory
306    /// that can be used to intern metric names and tags. If the interner is full, metrics with contexts that haven't
307    /// already been resolved may or may not be dropped, depending on the value of `allow_context_heap_allocations`.
308    pub context_string_interner_size_bytes: Option<ByteSize>,
309
310    /// The maximum number of cached contexts to allow.
311    ///
312    /// This is the maximum number of resolved contexts that can be cached at any given time. This limit doesn't affect
313    /// the total number of contexts that can be _alive_ at any given time, which is dependent on the interner capacity
314    /// and whether or not heap allocations are allowed.
315    pub cached_contexts_limit: usize,
316
317    /// The maximum number of cached tagsets to allow.
318    ///
319    /// This is the maximum number of resolved tagsets that can be cached at any given time. This limit doesn't affect
320    /// the total number of tagsets that can be _alive_ at any given time, which is dependent on the interner capacity
321    /// and whether or not heap allocations are allowed.
322    pub cached_tagsets_limit: usize,
323
324    /// The number of seconds after which cached contexts will expire.
325    ///
326    /// Higher values allow for more effective caching for sparse metrics at the cost of increased memory usage.
327    pub context_expiry_seconds: u64,
328
329    /// Whether or not to enable permissive mode in the decoder.
330    ///
331    /// Permissive mode allows the decoder to relax its strictness around the allowed payloads, which lets it match the
332    /// decoding behavior of the Datadog Agent.
333    pub permissive_decoding: bool,
334
335    /// The minimum sample rate allowed for metrics.
336    ///
337    /// When metrics are sent with a sample rate _lower_ than this value then it will be clamped to this value. This is
338    /// done in order to ensure an upper bound on how many equivalent samples are tracked for the metric, as high sample
339    /// rates (very small numbers, such as `0.00000001`) can lead to large memory growth.
340    ///
341    /// A warning log will be emitted when clamping occurs, as this represents an effective loss of metric samples.
342    pub minimum_sample_rate: f64,
343
344    /// Which payload types to forward to the backend.
345    pub enable_payloads: EnablePayloadsConfiguration,
346
347    /// Configuration related to origin detection and enrichment.
348    pub origin_enrichment: OriginEnrichmentConfiguration,
349
350    /// Whether to break down DogStatsD processed-metric telemetry by UDS origin.
351    ///
352    /// When enabled, metric-message `dogstatsd.processed` telemetry includes an `origin` label derived from the
353    /// sender's UDS origin. This can add one telemetry series per origin and should primarily be used for diagnostics.
354    pub origin_telemetry_enabled: bool,
355
356    /// Workload provider to utilize for origin enrichment.
357    ///
358    /// Detected origins are only resolved to workload tags when a provider is set. Origin detection itself is
359    /// controlled by `origin_enrichment`.
360    pub workload_provider: Option<Arc<dyn WorkloadProvider + Send + Sync>>,
361
362    /// Resolver to use for mapping live sender PIDs to container entities before deferred processing.
363    ///
364    /// This resolver pins the sender entity while socket credentials are current so origin enrichment and traffic
365    /// capture do not resolve a stale or reused PID later. It is set separately from the workload provider because it
366    /// only needs a narrow live-PID lookup.
367    pub capture_entity_resolver: Option<Arc<dyn CaptureEntityResolver + Send + Sync>>,
368
369    /// Additional tags to add to all metrics.
370    ///
371    /// These are sorted and deduplicated before use, so callers may append tags required by the running environment.
372    pub additional_tags: Vec<String>,
373
374    /// The directory where DogStatsD capture files are written by default.
375    ///
376    /// When set to an empty path, callers must provide an explicit capture path when starting a capture session.
377    pub capture_path: PathBuf,
378
379    /// The maximum number of captured packets that can be queued for persistence.
380    ///
381    /// This controls the depth of the in-process capture queue. Values below `1024` are raised to `1024` before the
382    /// capture writer starts, preventing a zero-depth rendezvous channel from serializing DogStatsD stream handlers
383    /// behind capture persistence.
384    pub capture_depth: usize,
385
386    /// Control surface that starts and stops DogStatsD traffic capture.
387    ///
388    /// This is runtime wiring rather than configuration: the source binds the capture it creates to this handle, so
389    /// the caller creates the handle and keeps a clone to serve the capture API.
390    pub capture_control: DogStatsDCaptureControl,
391
392    /// Control surface that starts and stops DogStatsD traffic replay.
393    ///
394    /// This is runtime wiring rather than configuration: the source binds the captured tagger to this handle, so the
395    /// caller creates the handle and keeps a clone to serve the replay API.
396    pub replay_control: DogStatsDReplayControl,
397}
398
399#[derive(Clone, Copy, Default)]
400struct EolRequired {
401    udp: bool,
402    uds: bool,
403    named_pipe: bool,
404}
405
406impl EolRequired {
407    fn from_config_values(values: &[String]) -> Self {
408        let mut eol_required = Self::default();
409
410        for value in values {
411            match value.as_str() {
412                "udp" => eol_required.udp = true,
413                "uds" => eol_required.uds = true,
414                "named_pipe" => eol_required.named_pipe = true,
415                _ => warn!(
416                    value,
417                    "Invalid dogstatsd_eol_required value. Expected 'udp', 'uds', or 'named_pipe'."
418                ),
419            }
420        }
421
422        eol_required
423    }
424
425    fn for_listener(&self, listen_addr: &ListenAddress) -> bool {
426        match listen_addr {
427            ListenAddress::Udp(_) => self.udp,
428            ListenAddress::Tcp(_) => false,
429            ListenAddress::Unixgram(_) | ListenAddress::Unix(_) => self.uds,
430            ListenAddress::NamedPipe { .. } => self.named_pipe,
431        }
432    }
433}
434
435/// Resolves a `bind_host` string to an `IpAddr`.
436///
437/// Accepts either an IP literal (no DNS required) or a hostname (resolved via async DNS). Returns
438/// `UnresolvableBindHost` if the lookup fails, or `BindHostHasNoAddresses` if it succeeds but
439/// returns no addresses.
440async fn resolve_bind_host(host: &str) -> Result<std::net::IpAddr, Error> {
441    let mut addrs = tokio::net::lookup_host((host, 0u16))
442        .await
443        .context(UnresolvableBindHost { host: host.to_string() })?;
444    addrs
445        .next()
446        .map(|sa| sa.ip())
447        .ok_or_else(|| Error::BindHostHasNoAddresses { host: host.to_string() })
448}
449
450#[cfg(test)]
451impl DogStatsDConfiguration {
452    /// Creates a fixture configuration, for tests that exercise source behavior rather than configuration.
453    ///
454    /// Values match the effective defaults the configuration system resolves, so a test only states the settings it
455    /// cares about. `capture_depth` is the exception: a default configuration leaves it at `0`, which the capture
456    /// writer raises to its minimum, so the fixture states that minimum directly.
457    fn for_test() -> Self {
458        Self {
459            default_hostname: MetaString::default(),
460            buffer_size: 8192,
461            buffer_count: 128,
462            buffer_count_max: 32_768,
463            workers_count: 0,
464            port: 8125,
465            socket_receive_buffer_size: 0,
466            tcp_port: 0,
467            statsd_forward_host: None,
468            statsd_forward_port: 0,
469            socket_path: None,
470            socket_stream_path: None,
471            stream_log_too_big: false,
472            pipe_name: None,
473            windows_pipe_security_descriptor: "D:AI(A;;GA;;;WD)".to_string(),
474            disable_verbose_logs: false,
475            eol_required: Vec::new(),
476            bind_host: None,
477            non_local_traffic: false,
478            autoscale_udp_listeners: false,
479            allow_context_heap_allocations: true,
480            no_aggregation_pipeline_support: true,
481            context_string_interner_entry_count: 4096,
482            context_string_interner_size_bytes: None,
483            cached_contexts_limit: 500_000,
484            cached_tagsets_limit: 500_000,
485            context_expiry_seconds: 20,
486            permissive_decoding: true,
487            minimum_sample_rate: 0.000000003845,
488            enable_payloads: EnablePayloadsConfiguration {
489                series: true,
490                sketches: true,
491                events: true,
492                service_checks: true,
493            },
494            origin_enrichment: OriginEnrichmentConfiguration::for_test(),
495            origin_telemetry_enabled: false,
496            workload_provider: None,
497            capture_entity_resolver: None,
498            additional_tags: Vec::new(),
499            capture_path: PathBuf::new(),
500            capture_depth: MIN_CAPTURE_DEPTH,
501            capture_control: DogStatsDCaptureControl::default(),
502            replay_control: DogStatsDReplayControl::default(),
503        }
504    }
505}
506
507impl DogStatsDConfiguration {
508    /// Gets the effective source-wide DogStatsD tags.
509    fn additional_tags(&self) -> Vec<String> {
510        let mut tags = self.additional_tags.clone();
511        tags.sort_unstable();
512        tags.dedup();
513        tags
514    }
515
516    /// Returns the effective string interner size in bytes.
517    ///
518    /// If `context_string_interner_size_bytes` is set, it's used directly. Otherwise,
519    /// `context_string_interner_entry_count` is multiplied by 512 bytes per entry to derive the byte size.
520    fn effective_context_string_interner_bytes(&self) -> ByteSize {
521        match self.context_string_interner_size_bytes {
522            Some(explicit_bytes) => explicit_bytes,
523            None => {
524                saluki_antithesis::always_le!(
525                    self.context_string_interner_entry_count,
526                    u64::MAX / INTERNER_BASELINE_BYTES_PER_ENTRY,
527                    "dogstatsd interner byte-size multiply does not overflow",
528                    { "entry_count": self.context_string_interner_entry_count }
529                );
530                ByteSize::b(
531                    self.context_string_interner_entry_count
532                        .saturating_mul(INTERNER_BASELINE_BYTES_PER_ENTRY),
533                )
534            }
535        }
536    }
537
538    fn eol_required(&self) -> EolRequired {
539        EolRequired::from_config_values(&self.eol_required)
540    }
541
542    fn statsd_forward_target(&self) -> Option<(&MetaString, u16)> {
543        let host = self.statsd_forward_host.as_ref()?;
544        if self.statsd_forward_port == 0 {
545            return None;
546        }
547
548        Some((host, self.statsd_forward_port))
549    }
550
551    fn packet_forwarder_target(&self) -> Option<PacketForwarderTarget> {
552        let (host, port) = self.statsd_forward_target()?;
553        Some(PacketForwarderTarget::new(host.clone(), port))
554    }
555
556    /// Returns the number of UDP stream handlers to spawn, derived from `dogstatsd_autoscale_udp_listeners` and
557    /// the number of available vCPUs.
558    ///
559    /// Returns `None` when autoscaling is disabled, which keeps the legacy single-socket behavior. The platform
560    /// gate for `SO_REUSEPORT` lives inside the listener—this method intentionally stays platform-agnostic.
561    fn udp_streams_to_yield(&self) -> Option<NonZeroUsize> {
562        if !self.autoscale_udp_listeners {
563            return None;
564        }
565
566        #[cfg(not(target_os = "linux"))]
567        if self.autoscale_udp_listeners {
568            warn!("UDP stream handler autoscaling not supported on non-Linux platforms. Default to single stream handler.");
569            return None;
570        }
571
572        let vcpus = std::thread::available_parallelism().map(NonZeroUsize::get).unwrap_or(1);
573        let streams = (1 + vcpus / 8).min(4);
574        NonZeroUsize::new(streams)
575    }
576
577    /// Returns the effective maximum size of the I/O buffer pool.
578    ///
579    /// The pool can never hold fewer buffers than the configured baseline, so a `dogstatsd_buffer_count_max` below
580    /// `dogstatsd_buffer_count` (including a legacy config that only raised `dogstatsd_buffer_count`) is treated as
581    /// equal to the baseline rather than reducing capacity.
582    fn effective_max_buffer_count(&self) -> usize {
583        self.buffer_count_max.max(self.buffer_count)
584    }
585
586    fn decoder_worker_count(&self) -> NonZeroUsize {
587        let worker_count = if self.workers_count == 0 {
588            let vcpus = std::thread::available_parallelism().map(NonZeroUsize::get).unwrap_or(1);
589            default_decoder_worker_count(vcpus)
590        } else {
591            self.workers_count
592        };
593
594        NonZeroUsize::new(worker_count).expect("DogStatsD decoder worker count must be non-zero")
595    }
596
597    /// Using the current configuration, determines which listeners should be created and adds an address for each into
598    /// a `Vec<ListenAddress>`. This function has no side effects so that it can be unit tested whereas build_listeners`
599    /// actually binds the listeners on the system.
600    ///
601    /// `bind_host` is the pre-resolved IP that UDP and TCP listeners should bind to (provided by
602    /// `resolve_bind_host`). Precedence matches the Agent:
603    ///   - `non_local_traffic=true` → `0.0.0.0` (`bind_host` ignored)
604    ///   - `bind_host=Some(ip)`     → `ip`
605    ///   - `bind_host=None`         → `127.0.0.1`
606    fn build_addresses(&self, bind_host: Option<std::net::IpAddr>) -> Vec<ListenAddress> {
607        let bind_ip: std::net::IpAddr = if self.non_local_traffic {
608            [0, 0, 0, 0].into()
609        } else {
610            bind_host.unwrap_or_else(|| [127, 0, 0, 1].into())
611        };
612
613        let mut addresses: Vec<ListenAddress> = Vec::new();
614
615        if self.port != 0 {
616            addresses.push(ListenAddress::Udp(std::net::SocketAddr::new(bind_ip, self.port)));
617        }
618
619        if self.tcp_port != 0 {
620            addresses.push(ListenAddress::Tcp(std::net::SocketAddr::new(bind_ip, self.tcp_port)));
621        }
622
623        if let Some(socket_path) = &self.socket_path {
624            addresses.push(ListenAddress::Unixgram(socket_path.into()));
625        }
626
627        if let Some(socket_stream_path) = &self.socket_stream_path {
628            addresses.push(ListenAddress::Unix(socket_stream_path.into()));
629        }
630
631        if let Some(pipe_name) = &self.pipe_name {
632            addresses.push(ListenAddress::named_pipe_with_input_buffer_size(
633                pipe_name,
634                &self.windows_pipe_security_descriptor,
635                self.buffer_size as u32,
636            ));
637        }
638
639        addresses
640    }
641
642    fn uds_origin_detection_unsupported_on_platform(&self, addresses: &[ListenAddress]) -> bool {
643        self.origin_enrichment.enabled
644            && cfg!(not(target_os = "linux"))
645            && addresses
646                .iter()
647                .any(|address| matches!(address, ListenAddress::Unixgram(_) | ListenAddress::Unix(_)))
648    }
649
650    fn warn_if_uds_origin_detection_unsupported(&self, addresses: &[ListenAddress]) {
651        if self.uds_origin_detection_unsupported_on_platform(addresses) {
652            warn!(
653                "DogStatsD UDS origin detection is enabled, but PID-based Unix socket credentials are unsupported on \
654                 this platform. Metrics are accepted without PID-based origin enrichment."
655            );
656        }
657    }
658
659    /// Acquires the appropriate `Listener` objects from the resource registry.
660    ///
661    /// The listeners are leased, not owned: the registry keeps the bound sockets when this source is torn down, so a
662    /// rebuilt source gets the same sockets back rather than racing to re-bind addresses it just released.
663    async fn build_listeners(&self, context: &BuildContext) -> Result<Vec<ResourceLease<Listener>>, Error> {
664        // Resolve `bind_host` to an IP (via DNS if needed). Skip the lookup when
665        // `non_local_traffic=true` since `bind_host` is ignored in that branch—matches Go's
666        // laziness and avoids failing startup on an unresolvable hostname that wouldn't be used.
667        let bind_host: Option<std::net::IpAddr> = if self.non_local_traffic {
668            None
669        } else {
670            match &self.bind_host {
671                Some(host) => Some(resolve_bind_host(host).await?),
672                None => None,
673            }
674        };
675
676        let addresses = self.build_addresses(bind_host);
677        self.warn_if_uds_origin_detection_unsupported(&addresses);
678        let mut listeners = Vec::new();
679        let socket_receive_buffer_size =
680            (self.socket_receive_buffer_size != 0).then_some(self.socket_receive_buffer_size);
681        let udp_streams_to_yield = self.udp_streams_to_yield();
682        for address in addresses {
683            let listener_type = address.listener_type();
684            let listener_streams = matches!(address, ListenAddress::Udp(_))
685                .then_some(udp_streams_to_yield)
686                .flatten();
687            let spec = SocketSpecification::new(address)
688                .with_udp_streams(listener_streams)
689                .with_receive_buffer_size(socket_receive_buffer_size);
690            let listener = context
691                .acquire_resource(spec)
692                .await
693                .context(FailedToAcquireListener { listener_type })?;
694
695            listeners.push(listener);
696        }
697        Ok(listeners)
698    }
699}
700
701#[async_trait]
702impl SourceBuilder for DogStatsDConfiguration {
703    async fn build(&self, context: BuildContext) -> Result<Box<dyn Source + Send>, GenericError> {
704        let listeners = self.build_listeners(&context).await?;
705        if listeners.is_empty() {
706            return Err(Error::NoListenersConfigured.into());
707        }
708
709        // Every listener requires at least one I/O buffer to ensure that all listeners can be serviced without
710        // deadlocking any of the others. Multi-socket connectionless listeners require one buffer per yielded socket.
711        let min_buffers: usize = listeners.iter().map(|listener| listener.min_buffer_reservation()).sum();
712        let max_buffers = self.effective_max_buffer_count();
713        if max_buffers < min_buffers {
714            return Err(generic_error!(
715                "The maximum I/O buffer count ({}) must be at least {} to service all configured listeners.",
716                max_buffers,
717                min_buffers,
718            ));
719        }
720
721        let origin_detection_enabled = self.origin_enrichment.enabled;
722        // Single CapturedTaggerHandle is cloned to both the resolver (reader of the captured store) and the replay
723        // control surface (writer). Both sides reference the same atomic slot.
724        let captured_tagger = CapturedTaggerHandle::new();
725
726        let maybe_origin_tags_resolver = self.workload_provider.clone().map(|provider| {
727            DogStatsDOriginTagResolver::new(self.origin_enrichment.clone(), provider, captured_tagger.clone())
728        });
729        let context_resolvers = ContextResolvers::new(self, context.component_context(), maybe_origin_tags_resolver)
730            .error_context("Failed to create context resolvers.")?;
731
732        let codec_config = DogStatsDCodecConfiguration::default()
733            .with_timestamps(self.no_aggregation_pipeline_support)
734            .with_permissive_mode(self.permissive_decoding)
735            .with_minimum_sample_rate(self.minimum_sample_rate)
736            .with_client_origin_detection(self.origin_enrichment.origin_detection_client);
737
738        let codec = DogStatsDCodec::from_configuration(codec_config);
739        let eol_required = self.eol_required();
740
741        let enable_payloads_filter = EnablePayloadsFilter::default()
742            .with_allow_series(self.enable_payloads.series)
743            .with_allow_sketches(self.enable_payloads.sketches)
744            .with_allow_events(self.enable_payloads.events)
745            .with_allow_service_checks(self.enable_payloads.service_checks);
746        let traffic_capture = TrafficCapture::with_workload_provider(
747            self.capture_path.clone(),
748            self.capture_depth,
749            self.workload_provider.clone(),
750        );
751        self.capture_control.bind(traffic_capture.clone());
752        let packet_forwarder_target = self.packet_forwarder_target();
753
754        self.replay_control.bind(captured_tagger);
755
756        // The pool allocates `buffer_count` buffers up front and may grow on demand up to `max_buffers`. The effective
757        // maximum is never below the baseline, so configs that only raise `dogstatsd_buffer_count` keep their full
758        // capacity instead of being silently reduced to the `dogstatsd_buffer_count_max` default.
759        let (io_buffer_pool, io_buffer_pool_shrinker) =
760            build_io_buffer_pool(self.buffer_count, max_buffers, self.buffer_size);
761        Ok(Box::new(DogStatsD {
762            listeners,
763            decoder_worker_count: self.decoder_worker_count(),
764            io_buffer_pool,
765            io_buffer_queue_capacity: max_buffers,
766            io_buffer_pool_shrinker: Box::pin(io_buffer_pool_shrinker),
767            codec,
768            context_resolvers,
769            default_hostname: self.default_hostname.clone(),
770            enabled_filter: enable_payloads_filter,
771            origin_detection_enabled,
772            origin_telemetry_enabled: self.origin_telemetry_enabled,
773            stream_log_too_big: self.stream_log_too_big,
774            disable_verbose_logs: self.disable_verbose_logs,
775            eol_required,
776            additional_tags: self.additional_tags().into(),
777            capture_entity_resolver: self.capture_entity_resolver.clone(),
778            traffic_capture,
779            packet_forwarder_target,
780        }))
781    }
782
783    fn outputs(&self) -> &[OutputDefinition<EventType>] {
784        static OUTPUTS: LazyLock<Vec<OutputDefinition<EventType>>> = LazyLock::new(|| {
785            vec![
786                OutputDefinition::named_output("metrics", EventType::Metric),
787                OutputDefinition::named_output("events", EventType::EventD),
788                OutputDefinition::named_output("service_checks", EventType::ServiceCheck),
789            ]
790        });
791        &OUTPUTS
792    }
793}
794
795impl MemoryBounds for DogStatsDConfiguration {
796    fn specify_bounds(&self, builder: &mut MemoryBoundsBuilder) {
797        let additional_buffers = self.effective_max_buffer_count().saturating_sub(self.buffer_count);
798        let adjusted_buffer_size = get_adjusted_buffer_size(self.buffer_size);
799
800        builder
801            .minimum()
802            // Capture the size of the heap allocation when the component is built.
803            .with_single_value::<DogStatsD>("source struct")
804            // We allocate the baseline buffer pool up front.
805            .with_expr(UsageExpr::product(
806                "buffers",
807                UsageExpr::config("dogstatsd_buffer_count", self.buffer_count),
808                UsageExpr::config("dogstatsd_buffer_size", adjusted_buffer_size),
809            ))
810            // We also allocate the backing storage for the string interner up front, which is used by our context
811            // resolver.
812            .with_expr(UsageExpr::config(
813                "dogstatsd_string_interner_size_bytes",
814                self.effective_context_string_interner_bytes().as_u64() as usize,
815            ));
816
817        // The pool can grow on demand up to its maximum, so account for the additional headroom as firm usage.
818        builder.firm().with_expr(UsageExpr::product(
819            "elastic buffers",
820            UsageExpr::constant("dogstatsd_buffer_count_max_extra", additional_buffers),
821            UsageExpr::config("dogstatsd_buffer_size", adjusted_buffer_size),
822        ));
823    }
824}
825
826/// DogStatsD source.
827pub struct DogStatsD {
828    listeners: Vec<ResourceLease<Listener>>,
829    decoder_worker_count: NonZeroUsize,
830    io_buffer_pool: ElasticObjectPool<BytesBuffer>,
831    io_buffer_queue_capacity: usize,
832    io_buffer_pool_shrinker: Pin<Box<dyn Future<Output = ()> + Send>>,
833    codec: DogStatsDCodec,
834    context_resolvers: ContextResolvers,
835    default_hostname: MetaString,
836    enabled_filter: EnablePayloadsFilter,
837    origin_detection_enabled: bool,
838    origin_telemetry_enabled: bool,
839    stream_log_too_big: bool,
840    disable_verbose_logs: bool,
841    eol_required: EolRequired,
842    additional_tags: Arc<[String]>,
843    capture_entity_resolver: Option<Arc<dyn CaptureEntityResolver + Send + Sync>>,
844    traffic_capture: TrafficCapture,
845    packet_forwarder_target: Option<PacketForwarderTarget>,
846}
847
848struct ListenerContext {
849    shutdown_handle: ShutdownHandle,
850    listener: ResourceLease<Listener>,
851    datagram_sender: mpsc::Sender<QueuedDatagram>,
852    io_buffer_pool: ElasticObjectPool<BytesBuffer>,
853    origin_telemetry_enabled: bool,
854    eol_required: EolRequired,
855    decoder_context: DecoderContext,
856    capture_entity_resolver: Option<Arc<dyn CaptureEntityResolver + Send + Sync>>,
857    packet_forwarder_target: Option<PacketForwarderTarget>,
858}
859
860#[derive(Clone)]
861struct HandlerContext {
862    listen_addr: ListenAddress,
863    eol_required: bool,
864    io_buffer_pool: ElasticObjectPool<BytesBuffer>,
865    datagram_sender: mpsc::Sender<QueuedDatagram>,
866    datagram_context: Option<Arc<DatagramSocketContext>>,
867    metrics: Metrics,
868    decoder_context: DecoderContext,
869    capture_entity_resolver: Option<Arc<dyn CaptureEntityResolver + Send + Sync>>,
870    packet_forwarder: Option<PacketForwarder>,
871}
872
873#[derive(Clone)]
874struct DecoderContext {
875    codec: DogStatsDCodec,
876    context_resolvers: ContextResolvers,
877    default_hostname: MetaString,
878    enabled_filter: EnablePayloadsFilter,
879    origin_detection_enabled: bool,
880    stream_log_too_big: bool,
881    disable_verbose_logs: bool,
882    additional_tags: Arc<[String]>,
883    traffic_capture: TrafficCapture,
884}
885
886struct DatagramSocketContext {
887    listen_addr: ListenAddress,
888    eol_required: bool,
889    metrics: Metrics,
890    packet_forwarder: Option<PacketForwarder>,
891}
892
893struct QueuedDatagram {
894    result: io::Result<ReceivedBuffer>,
895    socket_context: Arc<DatagramSocketContext>,
896}
897
898struct DogStatsDDecoder {
899    source_context: SourceContext,
900    codec: DogStatsDCodec,
901    context_resolvers: ContextResolvers,
902    default_hostname: MetaString,
903    enabled_filter: EnablePayloadsFilter,
904    origin_detection_enabled: bool,
905    stream_log_too_big: bool,
906    disable_verbose_logs: bool,
907    additional_tags: Arc<[String]>,
908    traffic_capture: TrafficCapture,
909    event_buffer: Option<EventsBuffer>,
910}
911
912#[derive(Clone, Copy)]
913enum BufferDecodeMode {
914    Connectionless,
915    Connected,
916}
917
918impl BufferDecodeMode {
919    fn is_eof(self, bytes_read: usize) -> bool {
920        match self {
921            Self::Connectionless => true,
922            Self::Connected => bytes_read == 0,
923        }
924    }
925
926    fn should_stop_on_eof(self, eof: bool) -> bool {
927        matches!(self, Self::Connected) && eof
928    }
929
930    fn should_stop_on_framing_error(self) -> bool {
931        matches!(self, Self::Connected)
932    }
933}
934
935struct BufferDecodeContext<'a> {
936    listen_addr: &'a ListenAddress,
937    metrics: &'a Metrics,
938    packet_forwarder: Option<&'a PacketForwarder>,
939    mode: BufferDecodeMode,
940    framer: DsdFramer,
941    stream_capture: StreamCaptureState,
942}
943
944impl<'a> BufferDecodeContext<'a> {
945    fn new(
946        listen_addr: &'a ListenAddress, eol_required: bool, metrics: &'a Metrics,
947        packet_forwarder: Option<&'a PacketForwarder>, mode: BufferDecodeMode,
948    ) -> Self {
949        Self {
950            listen_addr,
951            metrics,
952            packet_forwarder,
953            mode,
954            framer: get_framer(listen_addr, eol_required),
955            stream_capture: StreamCaptureState::new(),
956        }
957    }
958}
959
960#[derive(Clone, Copy, Debug, Eq, PartialEq)]
961enum DecodeOutcome {
962    Continue,
963    Stop,
964}
965
966#[async_trait]
967impl Source for DogStatsD {
968    async fn run(mut self: Box<Self>, mut context: SourceContext) -> Result<(), GenericError> {
969        let global_shutdown = context.take_shutdown_handle();
970        pin!(global_shutdown);
971
972        let mut health = context.take_health_handle();
973
974        // Brutal: the shrinker loops forever with no terminal condition of its own, so waiting for it would hold
975        // every shutdown open until the component's budget elapsed.
976        runtime::worker("io_buffer_pool_shrinker", self.io_buffer_pool_shrinker)
977            .with_shutdown_strategy(ShutdownStrategy::Brutal)
978            .spawn();
979
980        let (datagram_sender, datagram_receiver) = mpsc::channel(self.io_buffer_queue_capacity);
981        let datagram_receiver = Arc::new(Mutex::new(datagram_receiver));
982        let decoder_context = DecoderContext {
983            codec: self.codec.clone(),
984            context_resolvers: self.context_resolvers.clone(),
985            default_hostname: self.default_hostname.clone(),
986            enabled_filter: self.enabled_filter,
987            origin_detection_enabled: self.origin_detection_enabled,
988            stream_log_too_big: self.stream_log_too_big,
989            disable_verbose_logs: self.disable_verbose_logs,
990            additional_tags: self.additional_tags.clone(),
991            traffic_capture: self.traffic_capture.clone(),
992        };
993
994        // Decoders must drain their queue to completion, so they deliberately ignore the shutdown signal and stop only
995        // once the datagram channel closes -- which happens when every listener and stream handler has dropped its
996        // sender. The coordinator here is only used for its handle-drop accounting, so we can wait for that draining
997        // to finish; see `shutdown_listeners_and_drain_datagram_decoders`.
998        let mut decoder_shutdown_coordinator = ShutdownCoordinator::default();
999        for worker_id in 0..self.decoder_worker_count.get() {
1000            let decoder_shutdown = decoder_shutdown_coordinator.register();
1001            let datagram_receiver = datagram_receiver.clone();
1002            let decoder_source_context = context.clone();
1003            let decoder_context = decoder_context.clone();
1004
1005            runtime::worker(format!("datagram_decoder_{worker_id}"), async move {
1006                let _decoder_shutdown = decoder_shutdown;
1007                process_datagram_decoder(datagram_receiver, decoder_source_context, decoder_context).await;
1008            })
1009            .spawn();
1010        }
1011        drop(datagram_receiver);
1012
1013        let mut listener_shutdown_coordinator = ShutdownCoordinator::default();
1014        // For each listener, spawn a dedicated task to run it.
1015        for listener in self.listeners {
1016            let task_name = format!("listener_{}", listener.listen_address().listener_type());
1017            let listener_source_context = context.clone();
1018
1019            // TODO: Create a health handle for each listener.
1020            //
1021            // We need to rework `HealthRegistry` to look a little more like `ComponentRegistry` so that we can have it
1022            // already be scoped properly, otherwise all we can do here at present is either have a relative name, like
1023            // `uds-stream`, or try and hardcode the full component name, which we will inevitably forget to update if
1024            // we tweak the topology configuration, etc.
1025            let listener_context = ListenerContext {
1026                shutdown_handle: listener_shutdown_coordinator.register(),
1027                listener,
1028                datagram_sender: datagram_sender.clone(),
1029                io_buffer_pool: self.io_buffer_pool.clone(),
1030                origin_telemetry_enabled: self.origin_telemetry_enabled,
1031                eol_required: self.eol_required,
1032                decoder_context: decoder_context.clone(),
1033                capture_entity_resolver: self.capture_entity_resolver.clone(),
1034                packet_forwarder_target: self.packet_forwarder_target.clone(),
1035            };
1036
1037            runtime::supervisable(ListenerWorker::new(
1038                task_name,
1039                listener_source_context,
1040                listener_context,
1041            ))
1042            .temporary()
1043            .with_significant(true)
1044            .spawn();
1045        }
1046        drop(datagram_sender);
1047
1048        health.mark_ready();
1049        debug!("DogStatsD source started.");
1050
1051        // Wait for the global shutdown signal, then notify listeners to shutdown.
1052        //
1053        // We also handle liveness here, which doesn't really matter for _this_ task, since the real work is happening
1054        // in the listeners, but we need to satisfy the health checker.
1055        loop {
1056            select! {
1057                _ = &mut global_shutdown => {
1058                    debug!("Received shutdown signal.");
1059                    break
1060                },
1061                _ = health.live() => continue,
1062            }
1063        }
1064
1065        debug!("Stopping DogStatsD source...");
1066
1067        shutdown_listeners_and_drain_datagram_decoders(listener_shutdown_coordinator, decoder_shutdown_coordinator)
1068            .await;
1069
1070        debug!("DogStatsD source stopped.");
1071
1072        Ok(())
1073    }
1074}
1075
1076fn build_io_buffer_pool(
1077    min_buffers: usize, max_buffers: usize, buffer_size: usize,
1078) -> (ElasticObjectPool<BytesBuffer>, impl Future<Output = ()> + Send) {
1079    saluki_antithesis::always_le!(
1080        buffer_size,
1081        usize::MAX - 4,
1082        "dogstatsd buffer size add does not overflow",
1083        { "buffer_size": buffer_size }
1084    );
1085    let adjusted_buffer_size = get_adjusted_buffer_size(buffer_size);
1086    ElasticObjectPool::with_builder("dsd_packet_bufs", min_buffers, max_buffers, move || {
1087        FixedSizeVec::with_capacity(adjusted_buffer_size)
1088    })
1089}
1090
1091fn is_connectionless_listen_address(listen_addr: &ListenAddress) -> bool {
1092    match listen_addr {
1093        ListenAddress::Udp(_) => true,
1094        #[cfg(unix)]
1095        ListenAddress::Unixgram(_) => true,
1096        _ => false,
1097    }
1098}
1099
1100/// A DogStatsD listener, as a [`Supervisable`].
1101///
1102/// An accept loop has no terminal condition of its own -- it runs until something tells it to stop -- so unlike the
1103/// other children this source spawns, it genuinely needs the supervisor's shutdown signal. That is what makes it a
1104/// `Supervisable` rather than a plain worker. See [`process_listener`] for the two paths that end the loop.
1105struct ListenerWorker {
1106    name: String,
1107    inner: StdMutex<Option<(SourceContext, ListenerContext)>>,
1108}
1109
1110impl ListenerWorker {
1111    fn new(name: String, source_context: SourceContext, listener_context: ListenerContext) -> Self {
1112        Self {
1113            name,
1114            inner: StdMutex::new(Some((source_context, listener_context))),
1115        }
1116    }
1117}
1118
1119#[async_trait]
1120impl Supervisable for ListenerWorker {
1121    fn name(&self) -> &str {
1122        &self.name
1123    }
1124
1125    fn shutdown_strategy(&self) -> ShutdownStrategy {
1126        // No deadline of its own: the component's supervisor holds a single budget covering the source and everything
1127        // it spawned.
1128        ShutdownStrategy::Graceful(Duration::MAX)
1129    }
1130
1131    async fn initialize(&self, process_shutdown: ShutdownHandle) -> Result<SupervisorFuture, InitializationError> {
1132        let (source_context, listener_context) = self
1133            .inner
1134            .lock()
1135            .expect("listener worker mutex poisoned")
1136            .take()
1137            .ok_or_else(|| InitializationError::from(generic_error!("listener worker already initialized")))?;
1138
1139        Ok(Box::pin(async move {
1140            process_listener(source_context, listener_context, process_shutdown).await;
1141            Ok(())
1142        }))
1143    }
1144}
1145
1146async fn process_listener(
1147    source_context: SourceContext, listener_context: ListenerContext, process_shutdown: ShutdownHandle,
1148) {
1149    let ListenerContext {
1150        shutdown_handle,
1151        mut listener,
1152        datagram_sender,
1153        io_buffer_pool,
1154        origin_telemetry_enabled,
1155        eol_required,
1156        decoder_context,
1157        capture_entity_resolver,
1158        packet_forwarder_target,
1159    } = listener_context;
1160
1161    pin!(shutdown_handle, process_shutdown);
1162
1163    let listen_addr = listener.listen_address().clone();
1164    let metrics = build_metrics(
1165        &listen_addr,
1166        source_context.component_context(),
1167        origin_telemetry_enabled,
1168    );
1169    let packet_forwarder = packet_forwarder_target
1170        .as_ref()
1171        .map(|target| target.to_forwarder(metrics.clone()));
1172    if let Some(packet_forwarder) = &packet_forwarder {
1173        packet_forwarder.spawn();
1174    }
1175    let datagram_context = is_connectionless_listen_address(&listen_addr).then(|| {
1176        Arc::new(DatagramSocketContext {
1177            listen_addr: listen_addr.clone(),
1178            eol_required: eol_required.for_listener(&listen_addr),
1179            metrics: metrics.clone(),
1180            packet_forwarder: packet_forwarder.clone(),
1181        })
1182    });
1183
1184    let mut stream_shutdown_coordinator = ShutdownCoordinator::default();
1185
1186    info!(%listen_addr, "DogStatsD listener started.");
1187
1188    loop {
1189        select! {
1190            _ = &mut shutdown_handle => {
1191                debug!(%listen_addr, "Received shutdown signal. Waiting for existing stream handlers to finish...");
1192                break;
1193            }
1194            // The other way in: the supervisor tearing the component down while the source's own `run` is still going
1195            // and so isn't driving the coordinator above. Without this the accept loop would keep going -- and keep
1196            // holding a datagram sender, keeping the decoders waiting on a queue that never closes -- until the
1197            // shutdown budget elapsed and aborted the lot.
1198            _ = &mut process_shutdown => {
1199                debug!(%listen_addr, "Supervisor signalled shutdown. Waiting for existing stream handlers to finish...");
1200                break;
1201            }
1202            result = listener.accept() => match result {
1203                Ok(stream) => {
1204                    debug!(%listen_addr, "Spawning new stream handler.");
1205
1206                    let handler_context = HandlerContext {
1207                        listen_addr: listen_addr.clone(),
1208                        eol_required: eol_required.for_listener(&listen_addr),
1209                        io_buffer_pool: io_buffer_pool.clone(),
1210                        datagram_sender: datagram_sender.clone(),
1211                        datagram_context: datagram_context.clone(),
1212                        metrics: metrics.clone(),
1213                        decoder_context: decoder_context.clone(),
1214                        capture_entity_resolver: capture_entity_resolver.clone(),
1215                        packet_forwarder: packet_forwarder.clone(),
1216                    };
1217
1218                    let task_name = format!("conn_{}", listen_addr.listener_type());
1219
1220                    // The coordinator handle stays even though this is now a supervised child. Supervision makes the
1221                    // handler a sibling of the listener, not a descendant, so it alone wouldn't keep "the listener
1222                    // waits for its own streams" true -- the handle-drop accounting below is what does.
1223                    let stream_shutdown = stream_shutdown_coordinator.register();
1224                    let handler_source_context = source_context.clone();
1225                    let handler = process_stream(stream, handler_source_context, handler_context, stream_shutdown);
1226
1227                    runtime::worker(task_name, handler).spawn();
1228                }
1229                Err(e) => {
1230                    error!(%listen_addr, error = %e, "Failed to accept connection. Stopping listener.");
1231                    break
1232                }
1233            }
1234        }
1235    }
1236
1237    stream_shutdown_coordinator.shutdown_and_wait().await;
1238
1239    info!(%listen_addr, "DogStatsD listener stopped.");
1240}
1241
1242async fn process_stream(
1243    stream: Stream, source_context: SourceContext, handler_context: HandlerContext, shutdown_handle: ShutdownHandle,
1244) {
1245    select! {
1246        _ = shutdown_handle => {
1247            debug!("Stream handler received shutdown signal.");
1248        },
1249        _ = drive_stream(stream, source_context, handler_context) => {},
1250    }
1251}
1252
1253fn origin_detection_failed_for_telemetry(
1254    origin_detection_enabled: bool, bytes_read: usize, peer_addr: &ConnectionAddress,
1255) -> bool {
1256    origin_detection_enabled && bytes_read > 0 && peer_addr.has_process_credential_telemetry_error()
1257}
1258
1259struct ReceivedBuffer {
1260    buffer: Option<BytesBuffer>,
1261    bytes_read: usize,
1262    peer_addr: ConnectionAddress,
1263    process_origin: Option<ProcessOrigin>,
1264    buffer_sender: Option<oneshot::Sender<BytesBuffer>>,
1265}
1266
1267impl ReceivedBuffer {
1268    fn with_return(
1269        buffer: BytesBuffer, bytes_read: usize, peer_addr: ConnectionAddress, process_origin: Option<ProcessOrigin>,
1270    ) -> (Self, oneshot::Receiver<BytesBuffer>) {
1271        let (buffer_sender, returned_buffer) = oneshot::channel();
1272        (
1273            Self {
1274                buffer: Some(buffer),
1275                bytes_read,
1276                peer_addr,
1277                process_origin,
1278                buffer_sender: Some(buffer_sender),
1279            },
1280            returned_buffer,
1281        )
1282    }
1283
1284    fn without_return(
1285        buffer: BytesBuffer, bytes_read: usize, peer_addr: ConnectionAddress, process_origin: Option<ProcessOrigin>,
1286    ) -> Self {
1287        Self {
1288            buffer: Some(buffer),
1289            bytes_read,
1290            peer_addr,
1291            process_origin,
1292            buffer_sender: None,
1293        }
1294    }
1295
1296    fn parts_mut(&mut self) -> (&mut BytesBuffer, usize, &ConnectionAddress, Option<&ProcessOrigin>) {
1297        let Self {
1298            buffer,
1299            bytes_read,
1300            peer_addr,
1301            process_origin,
1302            ..
1303        } = self;
1304        (
1305            buffer.as_mut().expect("Received buffer already taken."),
1306            *bytes_read,
1307            peer_addr,
1308            process_origin.as_ref(),
1309        )
1310    }
1311
1312    #[cfg(all(test, unix))]
1313    fn buffer(&self) -> &BytesBuffer {
1314        self.buffer.as_ref().expect("Received buffer already taken.")
1315    }
1316
1317    #[cfg(all(test, unix))]
1318    fn buffer_mut(&mut self) -> &mut BytesBuffer {
1319        self.buffer.as_mut().expect("Received buffer already taken.")
1320    }
1321}
1322
1323impl Drop for ReceivedBuffer {
1324    fn drop(&mut self) {
1325        if let Some(sender) = self.buffer_sender.take() {
1326            let buffer = self.buffer.take().expect("Received buffer already taken.");
1327            let _ = sender.send(buffer);
1328        }
1329    }
1330}
1331
1332struct BufferedStreamReader {
1333    receiver: Option<mpsc::Receiver<io::Result<ReceivedBuffer>>>,
1334
1335    // Held but never fired explicitly, so that dropping this reader signals the worker to stop.
1336    _shutdown_coordinator: ShutdownCoordinator,
1337}
1338
1339impl BufferedStreamReader {
1340    fn new(
1341        stream: Stream, io_buffer_pool: ElasticObjectPool<BytesBuffer>, memory_limiter: MemoryLimiter,
1342        origin_detection_enabled: bool, traffic_capture: TrafficCapture,
1343        capture_entity_resolver: Option<Arc<dyn CaptureEntityResolver + Send + Sync>>,
1344    ) -> Self {
1345        debug_assert!(!stream.is_connectionless());
1346        let (packets_tx, receiver) = mpsc::channel(1);
1347        let (shutdown_coordinator, shutdown_handle) = ShutdownHandle::paired();
1348
1349        let reader = receive_connected_stream(
1350            stream,
1351            io_buffer_pool,
1352            memory_limiter,
1353            origin_detection_enabled,
1354            traffic_capture,
1355            capture_entity_resolver,
1356            packets_tx,
1357        );
1358
1359        runtime::worker("stream_reader", async move {
1360            select! {
1361                _ = shutdown_handle => debug!("Stream reader cancelled."),
1362                _ = reader => {},
1363            }
1364        })
1365        .spawn();
1366
1367        Self {
1368            receiver: Some(receiver),
1369            _shutdown_coordinator: shutdown_coordinator,
1370        }
1371    }
1372
1373    fn take_receiver(&mut self) -> mpsc::Receiver<io::Result<ReceivedBuffer>> {
1374        self.receiver.take().expect("Buffered stream receiver already taken")
1375    }
1376}
1377
1378async fn receive_connected_stream(
1379    mut stream: Stream, io_buffer_pool: ElasticObjectPool<BytesBuffer>, memory_limiter: MemoryLimiter,
1380    origin_detection_enabled: bool, traffic_capture: TrafficCapture,
1381    capture_entity_resolver: Option<Arc<dyn CaptureEntityResolver + Send + Sync>>,
1382    packets_tx: mpsc::Sender<io::Result<ReceivedBuffer>>,
1383) {
1384    debug!("Stream reader started.");
1385
1386    let mut retained_buffer: Option<BytesBuffer> = None;
1387    loop {
1388        memory_limiter.wait_for_capacity().await;
1389
1390        let mut buffer = match retained_buffer.take() {
1391            Some(mut buffer) => {
1392                buffer.collapse();
1393                buffer
1394            }
1395            None => acquire_io_buffer(&io_buffer_pool).await,
1396        };
1397        let (bytes_read, peer_addr) = match stream.receive(&mut buffer).await {
1398            Ok(received) => received,
1399            Err(error) => {
1400                let _ = packets_tx.send(Err(error)).await;
1401                break;
1402            }
1403        };
1404        let process_origin = resolve_process_origin_if_needed(
1405            origin_detection_enabled,
1406            &traffic_capture,
1407            capture_entity_resolver.as_deref(),
1408            &peer_addr,
1409        );
1410
1411        let (received, returned_buffer) = ReceivedBuffer::with_return(buffer, bytes_read, peer_addr, process_origin);
1412
1413        if packets_tx.send(Ok(received)).await.is_err() {
1414            debug!("Failed to enqueue DogStatsD packet for decoding: receiver dropped.");
1415            break;
1416        }
1417
1418        match returned_buffer.await {
1419            Ok(buffer) if buffer.has_remaining() => retained_buffer = Some(buffer),
1420            Ok(buffer) => drop(buffer),
1421            Err(_) => break,
1422        }
1423    }
1424
1425    debug!("Stream reader stopped.");
1426}
1427
1428async fn receive_connectionless_stream(
1429    mut stream: Stream, io_buffer_pool: ElasticObjectPool<BytesBuffer>, memory_limiter: MemoryLimiter,
1430    origin_detection_enabled: bool, traffic_capture: TrafficCapture,
1431    capture_entity_resolver: Option<Arc<dyn CaptureEntityResolver + Send + Sync>>,
1432    datagram_sender: mpsc::Sender<QueuedDatagram>, socket_context: Arc<DatagramSocketContext>,
1433) {
1434    debug!(listen_addr = %socket_context.listen_addr, "Datagram reader started.");
1435
1436    loop {
1437        memory_limiter.wait_for_capacity().await;
1438
1439        let mut buffer = acquire_io_buffer(&io_buffer_pool).await;
1440        let result = match stream.receive(&mut buffer).await {
1441            Ok((bytes_read, peer_addr)) => {
1442                let process_origin = resolve_process_origin_if_needed(
1443                    origin_detection_enabled,
1444                    &traffic_capture,
1445                    capture_entity_resolver.as_deref(),
1446                    &peer_addr,
1447                );
1448                Ok(ReceivedBuffer::without_return(
1449                    buffer,
1450                    bytes_read,
1451                    peer_addr,
1452                    process_origin,
1453                ))
1454            }
1455            Err(error) => Err(error),
1456        };
1457
1458        let receive_failed = result.is_err();
1459        let queued = QueuedDatagram {
1460            result,
1461            socket_context: socket_context.clone(),
1462        };
1463        if datagram_sender.send(queued).await.is_err() {
1464            debug!(
1465                listen_addr = %socket_context.listen_addr,
1466                "Failed to enqueue DogStatsD packet for decoding: receiver dropped."
1467            );
1468            break;
1469        }
1470        if receive_failed {
1471            continue;
1472        }
1473    }
1474
1475    debug!(listen_addr = %socket_context.listen_addr, "Datagram reader stopped.");
1476}
1477
1478async fn acquire_io_buffer(io_buffer_pool: &ElasticObjectPool<BytesBuffer>) -> BytesBuffer {
1479    let buffer = io_buffer_pool.acquire().await;
1480    trace!(
1481        remaining = buffer.remaining(),
1482        capacity = buffer.capacity(),
1483        "Acquired new buffer from pool."
1484    );
1485    buffer
1486}
1487
1488async fn process_datagram_decoder(
1489    datagram_receiver: Arc<Mutex<mpsc::Receiver<QueuedDatagram>>>, source_context: SourceContext,
1490    decoder_context: DecoderContext,
1491) {
1492    drive_datagram_decoder(datagram_receiver, source_context, decoder_context).await;
1493    debug!("Datagram decoder drained its queue.");
1494}
1495
1496/// Stops the listeners and then waits for the datagram decoders to finish draining.
1497///
1498/// The order matters and is what makes shutdown lossless: the listeners (and, transitively, their stream handlers) are
1499/// the only holders of a datagram sender, so waiting for them first is what closes the decoders' queue. Only then can
1500/// the decoders observe the close, drain what's left, and drop their handles.
1501async fn shutdown_listeners_and_drain_datagram_decoders(
1502    listener_shutdown_coordinator: ShutdownCoordinator, decoder_shutdown_coordinator: ShutdownCoordinator,
1503) {
1504    listener_shutdown_coordinator.shutdown_and_wait().await;
1505    decoder_shutdown_coordinator.shutdown_and_wait().await;
1506}
1507
1508impl DogStatsDDecoder {
1509    fn new(source_context: SourceContext, decoder_context: DecoderContext) -> Self {
1510        let DecoderContext {
1511            codec,
1512            context_resolvers,
1513            default_hostname,
1514            enabled_filter,
1515            origin_detection_enabled,
1516            stream_log_too_big,
1517            disable_verbose_logs,
1518            additional_tags,
1519            traffic_capture,
1520        } = decoder_context;
1521
1522        Self {
1523            source_context,
1524            codec,
1525            context_resolvers,
1526            default_hostname,
1527            enabled_filter,
1528            origin_detection_enabled,
1529            stream_log_too_big,
1530            disable_verbose_logs,
1531            additional_tags,
1532            traffic_capture,
1533            event_buffer: None,
1534        }
1535    }
1536
1537    async fn decode_buffer(
1538        &mut self, context: &mut BufferDecodeContext<'_>, mut received: ReceivedBuffer,
1539    ) -> DecodeOutcome {
1540        let (buffer, bytes_read, peer_addr, process_origin) = received.parts_mut();
1541        self.decode_buffer_contents(context, buffer, bytes_read, peer_addr, process_origin)
1542            .await
1543    }
1544
1545    async fn decode_buffer_contents(
1546        &mut self, context: &mut BufferDecodeContext<'_>, io_buffer: &mut BytesBuffer, bytes_read: usize,
1547        peer_addr: &ConnectionAddress, process_origin: Option<&ProcessOrigin>,
1548    ) -> DecodeOutcome {
1549        let listen_addr = context.listen_addr;
1550        let metrics = context.metrics;
1551        let packet_forwarder = context.packet_forwarder;
1552        let mode = context.mode;
1553
1554        let payload = received_payload(io_buffer, bytes_read);
1555        capture_uds_traffic(
1556            listen_addr,
1557            &self.traffic_capture,
1558            peer_addr,
1559            process_origin,
1560            payload,
1561            &mut context.stream_capture,
1562        );
1563
1564        metrics.bytes_received().increment(bytes_read as u64);
1565        metrics.bytes_received_size().record(bytes_read as f64);
1566        let origin_detection_failed =
1567            origin_detection_failed_for_telemetry(self.origin_detection_enabled, bytes_read, peer_addr);
1568
1569        if matches!(mode, BufferDecodeMode::Connectionless) {
1570            metrics.packet_receive_success().increment(1);
1571            if origin_detection_failed {
1572                metrics.origin_detection_errors().increment(1);
1573            }
1574        }
1575
1576        let eof = mode.is_eof(bytes_read);
1577        trace!(
1578            buffer_len = io_buffer.remaining(),
1579            buffer_cap = io_buffer.remaining_mut(),
1580            %listen_addr,
1581            %peer_addr,
1582            eof,
1583            "Received {} bytes from socket.",
1584            bytes_read
1585        );
1586
1587        if should_drop_oversized_named_pipe_frame(listen_addr, io_buffer) {
1588            metrics.framing_errors().increment(1);
1589            debug!(%listen_addr, %peer_addr, "DogStatsD named pipe frame exceeded the configured buffer size. Dropping frame.");
1590            io_buffer.clear();
1591            return DecodeOutcome::Continue;
1592        }
1593
1594        loop {
1595            let frame_result = context.framer.next_frame(io_buffer, eof);
1596            let completed_outer_frames = context.framer.take_completed_outer_frames();
1597            if completed_outer_frames > 0 {
1598                metrics
1599                    .packet_receive_success()
1600                    .increment(completed_outer_frames as u64);
1601            }
1602            if origin_detection_failed && completed_outer_frames > 0 {
1603                metrics
1604                    .origin_detection_errors()
1605                    .increment(completed_outer_frames as u64);
1606            }
1607
1608            match frame_result {
1609                Ok(Some(frame)) => {
1610                    if matches!(listen_addr, ListenAddress::NamedPipe { .. }) {
1611                        capture_named_pipe_frame(&self.traffic_capture, &frame);
1612                        metrics.packet_receive_success().increment(1);
1613                    }
1614                    self.decode_frame(frame, listen_addr, peer_addr, process_origin, metrics, packet_forwarder)
1615                        .await;
1616                }
1617                Ok(None) => {
1618                    if mode.should_stop_on_eof(eof) {
1619                        debug!(%listen_addr, %peer_addr, "Stream received EOF. Shutting down handler.");
1620                        return DecodeOutcome::Stop;
1621                    }
1622                    return DecodeOutcome::Continue;
1623                }
1624                Err(error) => {
1625                    metrics.framing_errors().increment(1);
1626                    if should_warn_stream_log_too_big(listen_addr, &error, self.stream_log_too_big) {
1627                        warn!(
1628                            %listen_addr,
1629                            %peer_addr,
1630                            error = %error,
1631                            "DogStatsD stream frame exceeded the configured buffer size."
1632                        );
1633                    }
1634
1635                    if mode.should_stop_on_framing_error() {
1636                        debug!(%listen_addr, %peer_addr, %error, "Error decoding frame. Stopping stream.");
1637                        return DecodeOutcome::Stop;
1638                    }
1639
1640                    debug!(%listen_addr, %peer_addr, %error, "Error decoding datagram frame. Continuing listener.");
1641                    return DecodeOutcome::Continue;
1642                }
1643            }
1644        }
1645    }
1646
1647    async fn decode_frame(
1648        &mut self, frame: Bytes, listen_addr: &ListenAddress, peer_addr: &ConnectionAddress,
1649        process_origin: Option<&ProcessOrigin>, metrics: &Metrics, packet_forwarder: Option<&PacketForwarder>,
1650    ) {
1651        trace!(%listen_addr, %peer_addr, ?frame, "Decoded frame.");
1652        if let Some(forwarder) = packet_forwarder {
1653            forwarder.forward(frame.clone()).await;
1654        }
1655
1656        match handle_frame(
1657            &frame[..],
1658            &self.codec,
1659            &mut self.context_resolvers,
1660            metrics,
1661            self.origin_detection_enabled,
1662            process_origin,
1663            self.enabled_filter,
1664            &self.additional_tags,
1665            &self.default_hostname,
1666        ) {
1667            Ok(Some(event)) => {
1668                if let Some(event_buffer) = self.buffer_event(event) {
1669                    debug!(%listen_addr, %peer_addr, "Event buffer is full. Forwarding events.");
1670                    dispatch_events(event_buffer, &self.source_context).await;
1671                }
1672            }
1673            Ok(None) => {}
1674            Err(error) => {
1675                log_parse_failure(self.disable_verbose_logs, listen_addr, peer_addr, &frame, &error);
1676            }
1677        }
1678    }
1679
1680    fn buffer_event(&mut self, event: Event) -> Option<EventsBuffer> {
1681        let event_buffer = self.event_buffer.get_or_insert_default();
1682        match event_buffer.try_push(event) {
1683            Some(event) => {
1684                let full_event_buffer = mem::take(event_buffer);
1685                assert!(
1686                    event_buffer.try_push(event).is_none(),
1687                    "New event buffer is unexpectedly full."
1688                );
1689                Some(full_event_buffer)
1690            }
1691            None => None,
1692        }
1693    }
1694
1695    async fn flush_events(&mut self) {
1696        if let Some(event_buffer) = self.event_buffer.take() {
1697            dispatch_events(event_buffer, &self.source_context).await;
1698        }
1699    }
1700}
1701
1702async fn drive_datagram_decoder(
1703    datagram_receiver: Arc<Mutex<mpsc::Receiver<QueuedDatagram>>>, source_context: SourceContext,
1704    decoder_context: DecoderContext,
1705) {
1706    let mut decoder = DogStatsDDecoder::new(source_context, decoder_context);
1707    let mut buffer_flush = interval(Duration::from_millis(100));
1708    buffer_flush.set_missed_tick_behavior(MissedTickBehavior::Delay);
1709
1710    loop {
1711        select! {
1712            maybe_datagram = async {
1713                datagram_receiver.lock().await.recv().await
1714            } => {
1715                let Some(QueuedDatagram { result, socket_context }) = maybe_datagram else {
1716                    break;
1717                };
1718                {
1719                    let DatagramSocketContext {
1720                        listen_addr,
1721                        eol_required,
1722                        metrics,
1723                        packet_forwarder,
1724                    } = socket_context.as_ref();
1725
1726                    let received = match result {
1727                        Ok(received) => received,
1728                        Err(error) => {
1729                            metrics.packet_receive_failure().increment(1);
1730                            warn!(%listen_addr, %error, "I/O error while reading datagram. Continuing listener.");
1731                            continue;
1732                        }
1733                    };
1734
1735                    let mut buffer_decode_context = BufferDecodeContext::new(
1736                        listen_addr,
1737                        *eol_required,
1738                        metrics,
1739                        packet_forwarder.as_ref(),
1740                        BufferDecodeMode::Connectionless,
1741                    );
1742                    let outcome = decoder.decode_buffer(&mut buffer_decode_context, received).await;
1743                    debug_assert_eq!(outcome, DecodeOutcome::Continue);
1744                }
1745            }
1746            _ = buffer_flush.tick() => {
1747                decoder.flush_events().await;
1748            }
1749        }
1750    }
1751
1752    decoder.flush_events().await;
1753}
1754
1755async fn drive_stream(stream: Stream, source_context: SourceContext, handler_context: HandlerContext) {
1756    if stream.is_connectionless() {
1757        let memory_limiter = source_context.topology_context().memory_limiter().clone();
1758        let socket_context = handler_context
1759            .datagram_context
1760            .clone()
1761            .expect("connectionless stream must have a datagram context");
1762        receive_connectionless_stream(
1763            stream,
1764            handler_context.io_buffer_pool,
1765            memory_limiter,
1766            handler_context.decoder_context.origin_detection_enabled,
1767            handler_context.decoder_context.traffic_capture.clone(),
1768            handler_context.capture_entity_resolver,
1769            handler_context.datagram_sender,
1770            socket_context,
1771        )
1772        .await;
1773        return;
1774    }
1775
1776    drive_connected_stream(stream, source_context, handler_context).await;
1777}
1778
1779async fn drive_connected_stream(stream: Stream, source_context: SourceContext, handler_context: HandlerContext) {
1780    let listen_addr = handler_context.listen_addr.clone();
1781    let metrics = handler_context.metrics.clone();
1782    let memory_limiter = source_context.topology_context().memory_limiter().clone();
1783    let mut stream_reader = BufferedStreamReader::new(
1784        stream,
1785        handler_context.io_buffer_pool.clone(),
1786        memory_limiter,
1787        handler_context.decoder_context.origin_detection_enabled,
1788        handler_context.decoder_context.traffic_capture.clone(),
1789        handler_context.capture_entity_resolver.clone(),
1790    );
1791    let receiver = stream_reader.take_receiver();
1792
1793    debug!(%listen_addr, "Stream handler started.");
1794
1795    metrics.connections_active().increment(1);
1796    drive_decoder(receiver, source_context, handler_context).await;
1797    metrics.connections_active().decrement(1);
1798
1799    debug!(%listen_addr, "Stream handler stopped.");
1800}
1801
1802async fn drive_decoder(
1803    mut stream_receiver: mpsc::Receiver<io::Result<ReceivedBuffer>>, source_context: SourceContext,
1804    handler_context: HandlerContext,
1805) {
1806    let HandlerContext {
1807        listen_addr,
1808        eol_required,
1809        metrics,
1810        decoder_context,
1811        packet_forwarder,
1812        ..
1813    } = handler_context;
1814    let mut decoder = DogStatsDDecoder::new(source_context, decoder_context);
1815    // Set a buffer flush interval of 100ms, which will ensure we always flush buffered events at least every 100ms if
1816    // we're otherwise idle and not receiving packets from the client.
1817    let mut buffer_flush = interval(Duration::from_millis(100));
1818    buffer_flush.set_missed_tick_behavior(MissedTickBehavior::Delay);
1819    let mut buffer_decode_context = BufferDecodeContext::new(
1820        &listen_addr,
1821        eol_required,
1822        &metrics,
1823        packet_forwarder.as_ref(),
1824        BufferDecodeMode::Connected,
1825    );
1826
1827    'read: loop {
1828        select! {
1829            // We read from the stream.
1830            maybe_read_result = stream_receiver.recv() => match maybe_read_result {
1831                Some(Ok(received)) => {
1832                    let outcome = decoder.decode_buffer(&mut buffer_decode_context, received).await;
1833                    if outcome == DecodeOutcome::Stop {
1834                        break 'read;
1835                    }
1836                },
1837                Some(Err(e)) => {
1838                    metrics.packet_receive_failure().increment(1);
1839                    warn!(%listen_addr, error = %e, "I/O error while decoding. Stopping stream.");
1840                    break 'read;
1841                },
1842                None => {
1843                    warn!(%listen_addr, "Buffered stream reader stopped unexpectedly. Stopping stream.");
1844                    break 'read;
1845                }
1846            },
1847
1848            _ = buffer_flush.tick() => {
1849                decoder.flush_events().await;
1850            },
1851
1852        }
1853    }
1854
1855    decoder.flush_events().await;
1856}
1857
1858fn should_drop_oversized_named_pipe_frame(listen_addr: &ListenAddress, buffer: &BytesBuffer) -> bool {
1859    matches!(listen_addr, ListenAddress::NamedPipe { .. })
1860        && buffer.remaining_mut() == 0
1861        && memchr::memchr(b'\n', buffer.chunk()).is_none()
1862}
1863
1864fn should_warn_stream_log_too_big(listen_addr: &ListenAddress, error: &FramingError, stream_log_too_big: bool) -> bool {
1865    stream_log_too_big
1866        && matches!(listen_addr, ListenAddress::Unix(_))
1867        && matches!(error, FramingError::InvalidFrame { .. })
1868}
1869
1870fn log_parse_failure(
1871    disable_verbose_logs: bool, listen_addr: &ListenAddress, peer_addr: &ConnectionAddress, frame: &[u8],
1872    error: &ParseError,
1873) {
1874    let frame = String::from_utf8_lossy(frame);
1875    if disable_verbose_logs {
1876        debug!(%listen_addr, %peer_addr, %frame, %error, "Failed to parse frame.");
1877    } else {
1878        warn!(%listen_addr, %peer_addr, %frame, %error, "Failed to parse frame.");
1879    }
1880}
1881
1882fn capture_named_pipe_frame(traffic_capture: &TrafficCapture, frame: &[u8]) {
1883    if !frame.is_empty() && traffic_capture.is_ongoing() {
1884        let _ = traffic_capture.enqueue(build_capture_record(None, None, frame));
1885    }
1886}
1887
1888fn capture_uds_traffic(
1889    listen_addr: &ListenAddress, traffic_capture: &TrafficCapture, peer_addr: &ConnectionAddress,
1890    process_origin: Option<&ProcessOrigin>, payload: &[u8], stream_capture: &mut StreamCaptureState,
1891) {
1892    if payload.is_empty() || !traffic_capture.is_ongoing() {
1893        return;
1894    }
1895
1896    match listen_addr {
1897        ListenAddress::Unixgram(_) => {
1898            let _ = traffic_capture.enqueue(build_capture_record(
1899                process_id_from_peer_addr(peer_addr),
1900                process_origin,
1901                payload,
1902            ));
1903        }
1904        ListenAddress::Unix(_) => {
1905            stream_capture.update_peer_metadata(peer_addr);
1906            stream_capture.pending.extend(payload);
1907
1908            while let Ok(Some(outer_payload)) = stream_capture
1909                .outer_framer
1910                .next_frame(&mut stream_capture.pending, false)
1911            {
1912                let _ = traffic_capture.enqueue(build_capture_record(
1913                    stream_capture.last_pid,
1914                    process_origin,
1915                    &outer_payload,
1916                ));
1917            }
1918        }
1919        _ => {}
1920    }
1921}
1922
1923struct StreamCaptureState {
1924    outer_framer: LengthDelimitedFramer,
1925    pending: VecDeque<u8>,
1926    last_pid: Option<i32>,
1927}
1928
1929impl StreamCaptureState {
1930    fn new() -> Self {
1931        Self {
1932            outer_framer: LengthDelimitedFramer,
1933            pending: VecDeque::new(),
1934            last_pid: None,
1935        }
1936    }
1937
1938    fn update_peer_metadata(&mut self, peer_addr: &ConnectionAddress) {
1939        if let Some(process_id) = process_id_from_peer_addr(peer_addr) {
1940            self.last_pid = Some(process_id);
1941        }
1942    }
1943}
1944
1945fn build_capture_record(
1946    process_id: Option<i32>, process_origin: Option<&ProcessOrigin>, payload: &[u8],
1947) -> CaptureRecord {
1948    CaptureRecord {
1949        timestamp_ns: capture_timestamp_ns(),
1950        payload: payload.to_vec(),
1951        pid: process_id,
1952        ancillary: Vec::new(),
1953        container_id: process_origin
1954            .and_then(ProcessOrigin::container_entity_id)
1955            .map(ToString::to_string),
1956    }
1957}
1958
1959fn process_id_from_peer_addr(peer_addr: &ConnectionAddress) -> Option<i32> {
1960    match peer_addr {
1961        ConnectionAddress::ProcessLike(ProcessIdentity::Credentials(creds)) => Some(creds.pid),
1962        _ => None,
1963    }
1964}
1965
1966fn resolve_process_origin(
1967    capture_entity_resolver: Option<&(dyn CaptureEntityResolver + Send + Sync)>, peer_addr: &ConnectionAddress,
1968) -> Option<ProcessOrigin> {
1969    let creds = peer_addr.process_credentials()?;
1970    if creds.gid == REPLAY_CREDENTIALS_GID {
1971        return Some(ProcessOrigin::Replay(creds.uid));
1972    }
1973
1974    let process_id = u32::try_from(creds.pid).ok()?;
1975    Some(match capture_entity_resolver {
1976        Some(resolver) => ProcessOrigin::Pinned(resolver.resolve_container_entity_for_live_pid(process_id)),
1977        None => ProcessOrigin::Unpinned(process_id),
1978    })
1979}
1980
1981fn resolve_process_origin_if_needed(
1982    origin_detection_enabled: bool, traffic_capture: &TrafficCapture,
1983    capture_entity_resolver: Option<&(dyn CaptureEntityResolver + Send + Sync)>, peer_addr: &ConnectionAddress,
1984) -> Option<ProcessOrigin> {
1985    if !origin_detection_enabled && !traffic_capture.is_ongoing() {
1986        return None;
1987    }
1988
1989    resolve_process_origin(capture_entity_resolver, peer_addr)
1990}
1991
1992fn received_payload(buffer: &BytesBuffer, bytes_read: usize) -> &[u8] {
1993    let chunk = buffer.chunk();
1994    let start = chunk.len().saturating_sub(bytes_read);
1995    &chunk[start..]
1996}
1997
1998fn capture_timestamp_ns() -> i64 {
1999    SystemTime::now()
2000        .duration_since(UNIX_EPOCH)
2001        .map(|duration| duration.as_nanos().min(i64::MAX as u128) as i64)
2002        .unwrap_or_default()
2003}
2004
2005#[allow(clippy::too_many_arguments)]
2006fn handle_frame(
2007    frame: &[u8], codec: &DogStatsDCodec, context_resolvers: &mut ContextResolvers, source_metrics: &Metrics,
2008    origin_detection_enabled: bool, process_origin: Option<&ProcessOrigin>, enabled_filter: EnablePayloadsFilter,
2009    additional_tags: &[String], default_hostname: &MetaString,
2010) -> Result<Option<Event>, ParseError> {
2011    let resolve_telemetry_origin = || {
2012        (source_metrics.origin_telemetry_enabled() && origin_detection_enabled)
2013            .then(|| {
2014                process_origin
2015                    .and_then(ProcessOrigin::container_entity_id)
2016                    .map(ToString::to_string)
2017            })
2018            .flatten()
2019    };
2020
2021    let parsed = match codec.decode_packet(frame) {
2022        Ok(parsed) => parsed,
2023        Err(e) => {
2024            // Try and determine what the message type was, if possible, to increment the correct error counter.
2025            match parse_message_type(frame) {
2026                MessageType::MetricSample => {
2027                    source_metrics.record_metric_parse_failed(resolve_telemetry_origin().as_deref())
2028                }
2029                MessageType::Event => source_metrics.event_decode_failed().increment(1),
2030                MessageType::ServiceCheck => source_metrics.service_check_decode_failed().increment(1),
2031            }
2032
2033            return Err(e);
2034        }
2035    };
2036
2037    let event = match parsed {
2038        ParsedPacket::Metric(metric_packet) => {
2039            if metric_packet.num_points == 0 {
2040                return Ok(None);
2041            }
2042            let events_len = metric_packet.num_points;
2043            if !enabled_filter.allow_metric(&metric_packet) {
2044                trace!(
2045                    metric.name = metric_packet.metric_name,
2046                    "Skipping metric due to filter configuration."
2047                );
2048                return Ok(None);
2049            }
2050
2051            match handle_metric_packet(
2052                metric_packet,
2053                context_resolvers,
2054                process_origin,
2055                additional_tags,
2056                default_hostname,
2057            ) {
2058                Some(metric) => {
2059                    source_metrics.record_metrics_received(events_len, resolve_telemetry_origin().as_deref());
2060                    Event::Metric(metric)
2061                }
2062                None => {
2063                    // We can only fail to get a metric back if we failed to resolve the context.
2064                    source_metrics.failed_context_resolve_total().increment(1);
2065                    return Ok(None);
2066                }
2067            }
2068        }
2069        ParsedPacket::Event(event) => {
2070            if !enabled_filter.allow_event(&event) {
2071                trace!("Skipping event {} due to filter configuration.", event.title);
2072                return Ok(None);
2073            }
2074            match handle_event_packet(event, context_resolvers, process_origin, additional_tags) {
2075                Some(event) => {
2076                    source_metrics.events_received().increment(1);
2077                    Event::EventD(event)
2078                }
2079                None => {
2080                    source_metrics.failed_context_resolve_total().increment(1);
2081                    return Ok(None);
2082                }
2083            }
2084        }
2085        ParsedPacket::ServiceCheck(service_check) => {
2086            if !enabled_filter.allow_service_check(&service_check) {
2087                trace!(
2088                    "Skipping service check {} due to filter configuration.",
2089                    service_check.name
2090                );
2091                return Ok(None);
2092            }
2093            match handle_service_check_packet(service_check, context_resolvers, process_origin, additional_tags) {
2094                Some(service_check) => {
2095                    source_metrics.service_checks_received().increment(1);
2096                    Event::ServiceCheck(service_check)
2097                }
2098                None => {
2099                    source_metrics.failed_context_resolve_total().increment(1);
2100                    return Ok(None);
2101                }
2102            }
2103        }
2104    };
2105
2106    Ok(Some(event))
2107}
2108
2109fn handle_metric_packet(
2110    packet: MetricPacket, context_resolvers: &mut ContextResolvers, process_origin: Option<&ProcessOrigin>,
2111    additional_tags: &[String], default_hostname: &MetaString,
2112) -> Option<Metric> {
2113    let well_known_tags = WellKnownTags::from_raw_tags(&packet.tags);
2114
2115    let origin = origin_from_metric_packet(&packet, &well_known_tags);
2116    let origin_tags = context_resolvers.resolve_origin_tags(origin, process_origin);
2117
2118    // Choose the right context resolver based on whether or not this metric is pre-aggregated.
2119    let context_resolver = if packet.timestamp.is_some() {
2120        context_resolvers.no_agg()
2121    } else {
2122        context_resolvers.primary()
2123    };
2124
2125    let tags = get_filtered_tags_iterator(&packet.tags, additional_tags);
2126
2127    let hostname = well_known_tags.hostname.unwrap_or(default_hostname);
2128
2129    // Try to resolve the context for this metric.
2130    let maybe_context =
2131        context_resolver.resolve_with_host_and_origin_tags(packet.metric_name, hostname, tags, origin_tags);
2132
2133    match maybe_context {
2134        Some(context) => {
2135            let metric_origin = well_known_tags
2136                .jmx_check_name
2137                .map(MetricOrigin::jmx_check)
2138                .unwrap_or_else(MetricOrigin::dogstatsd);
2139            let metadata = MetricMetadata::default()
2140                .with_origin(metric_origin)
2141                .with_unit(packet.unit.map_or_else(MetaString::empty, MetaString::from_static));
2142
2143            Some(Metric::from_parts(context, packet.values, metadata))
2144        }
2145        // We failed to resolve the context, likely due to not having enough interner capacity.
2146        None => None,
2147    }
2148}
2149
2150fn handle_event_packet(
2151    packet: EventPacket, context_resolvers: &mut ContextResolvers, process_origin: Option<&ProcessOrigin>,
2152    additional_tags: &[String],
2153) -> Option<EventD> {
2154    let well_known_tags = WellKnownTags::from_raw_tags(&packet.tags);
2155
2156    let origin = origin_from_event_packet(&packet, &well_known_tags);
2157    let origin_tags = context_resolvers.resolve_origin_tags(origin, process_origin);
2158
2159    let tags = get_filtered_tags_iterator(&packet.tags, additional_tags);
2160    let tags_resolver = context_resolvers.tags();
2161    let tags = tags_resolver.create_tag_set(tags)?;
2162
2163    // When no d: field is present, backfill the current time—matching the stock Datadog Agent's
2164    // behavior in pkg/aggregator/aggregator.go (addEvent), which sets e.Ts = time.Now().Unix()
2165    // for any event with Ts == 0.
2166    let timestamp = packet
2167        .timestamp
2168        .or_else(|| SystemTime::now().duration_since(UNIX_EPOCH).ok().map(|d| d.as_secs()));
2169
2170    let eventd = EventD::new(packet.title, packet.text)
2171        .with_timestamp(timestamp)
2172        .with_hostname(packet.hostname.map(|s| s.into()))
2173        .with_aggregation_key(packet.aggregation_key.map(|s| s.into()))
2174        .with_alert_type(packet.alert_type)
2175        .with_priority(packet.priority)
2176        // When no source type is provided, default to "api"—the same default the stock Datadog
2177        // Agent applies when serializing DogStatsD events to the intake JSON format. The agent
2178        // groups events by source type name and uses "api" as the key for events without an
2179        // explicit `s:` field. See: pkg/serializer/internal/metrics/events.go (writeItem).
2180        .with_source_type_name(Some(
2181            packet
2182                .source_type_name
2183                .map(|s| s.into())
2184                .unwrap_or_else(|| "api".into()),
2185        ))
2186        .with_alert_type(packet.alert_type)
2187        .with_tags(tags)
2188        .with_origin_tags(origin_tags);
2189
2190    Some(eventd)
2191}
2192
2193fn handle_service_check_packet(
2194    packet: ServiceCheckPacket, context_resolvers: &mut ContextResolvers, process_origin: Option<&ProcessOrigin>,
2195    additional_tags: &[String],
2196) -> Option<ServiceCheck> {
2197    let well_known_tags = WellKnownTags::from_raw_tags(&packet.tags);
2198
2199    let origin = origin_from_service_check_packet(&packet, &well_known_tags);
2200    let origin_tags = context_resolvers.resolve_origin_tags(origin, process_origin);
2201
2202    let tags = get_filtered_tags_iterator(&packet.tags, additional_tags);
2203    let tags_resolver = context_resolvers.tags();
2204    let tags = tags_resolver.create_tag_set(tags)?;
2205
2206    // When no d: field is present, backfill the current time—matching the stock Datadog Agent's
2207    // behavior, which sets the timestamp to time.Now().Unix() for any service check with a zero
2208    // timestamp.
2209    let timestamp = packet
2210        .timestamp
2211        .or_else(|| SystemTime::now().duration_since(UNIX_EPOCH).ok().map(|d| d.as_secs()));
2212
2213    let service_check = ServiceCheck::new(packet.name, packet.status)
2214        .with_timestamp(timestamp)
2215        .with_hostname(packet.hostname.map(|s| s.into()))
2216        .with_tags(tags)
2217        .with_origin_tags(origin_tags)
2218        .with_message(packet.message.map(|s| s.into()));
2219
2220    Some(service_check)
2221}
2222
2223fn get_filtered_tags_iterator<'a>(
2224    raw_tags: &'a RawTags<'_>, additional_tags: &'a [String],
2225) -> impl Iterator<Item = &'a str> + Clone {
2226    // This filters out "well-known" tags from the raw tags in the DogStatsD packet, and then chains on any additional tags
2227    // that were configured on the source.
2228    RawTagsFilter::exclude(raw_tags, WellKnownTagsFilterPredicate).chain(additional_tags.iter().map(|s| s.as_str()))
2229}
2230
2231async fn dispatch_events(mut event_buffer: EventsBuffer, source_context: &SourceContext) {
2232    debug!(events_len = event_buffer.len(), "Forwarding events.");
2233
2234    // TODO: This is maybe a little dicey because if we fail to dispatch the events, we may not have iterated over all of
2235    // them, so there might still be eventd events when get to the service checks point, and eventd events and/or service
2236    // check events when we get to the metrics point, and so on.
2237    //
2238    // There's probably something to be said for erroring out fully if this happens, since we should only fail to
2239    // dispatch if the downstream component fails entirely... and unless we have a way to restart the component, then
2240    // we're going to continue to fail to dispatch any more events until the process is restarted anyways.
2241
2242    // Dispatch any eventd events, if present.
2243    if event_buffer.has_event_type(EventType::EventD) {
2244        let eventd_events = event_buffer.extract(Event::is_eventd);
2245        let events_output = source_context.dispatcher().buffered_named("events");
2246
2247        // The `events` output is always wired in the DSD topology, so a missing output is an invariant violation that
2248        // crashes this component.
2249        if events_output.is_err() {
2250            saluki_antithesis::unreachable!("dsd 'events' output missing at dispatch");
2251        }
2252
2253        if let Err(e) = events_output
2254            .expect("events output should always exist")
2255            .send_all(eventd_events)
2256            .await
2257        {
2258            error!(error = %e, "Failed to dispatch eventd events.");
2259
2260            saluki_antithesis::unreachable!("dsd dispatch failed mid-buffer", { "stream": "events" });
2261        }
2262    }
2263
2264    // Dispatch any service check events, if present.
2265    if event_buffer.has_event_type(EventType::ServiceCheck) {
2266        let service_check_events = event_buffer.extract(Event::is_service_check);
2267        let service_checks_output = source_context.dispatcher().buffered_named("service_checks");
2268
2269        if service_checks_output.is_err() {
2270            saluki_antithesis::unreachable!("dsd 'service_checks' output missing at dispatch");
2271        }
2272
2273        if let Err(e) = service_checks_output
2274            .expect("service checks output should always exist")
2275            .send_all(service_check_events)
2276            .await
2277        {
2278            error!(error = %e, "Failed to dispatch service check events.");
2279
2280            saluki_antithesis::unreachable!("dsd dispatch failed mid-buffer", { "stream": "service_checks" });
2281        }
2282    }
2283
2284    // Finally, if there are events left, they'll be metrics, so dispatch them.
2285    if !event_buffer.is_empty() {
2286        if let Err(e) = source_context
2287            .dispatcher()
2288            .dispatch_named("metrics", event_buffer)
2289            .await
2290        {
2291            error!(error = %e, "Failed to dispatch metric events.");
2292
2293            saluki_antithesis::unreachable!("dsd dispatch failed mid-buffer", { "stream": "metrics" });
2294        }
2295    }
2296}
2297
2298const fn get_adjusted_buffer_size(buffer_size: usize) -> usize {
2299    // This is a little goofy, but hear me out:
2300    //
2301    // In the Datadog Agent, the way the UDS listener works is that if it's in stream mode, it will do a standalone
2302    // socket read to get _just_ the length delimiter, which is 4 bytes. After that, it will do a read to get the packet
2303    // data itself, up to the limit of `dogstatsd_buffer_size`. This means that a _full_ UDS stream packet can be up to
2304    // `dogstatsd_buffer_size + 4` bytes.
2305    //
2306    // 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
2307    // able to get an entire frame in a single buffer for the purpose of decoding the frame. Rather than rewriting our
2308    // read loop such that we have to change the logic depending on UDP/UDS datagram vs UDS stream, we simply increase
2309    // the buffer size by 4 bytes to account for the length delimiter.
2310    //
2311    // We do it this way so that we don't have to change the buffer size in the configuration, since if you just ported
2312    // over a Datadog Agent configuration, the value would be too small, and vise versa.
2313    buffer_size + 4
2314}
2315
2316#[cfg(test)]
2317mod tests {
2318    use std::{
2319        collections::HashMap,
2320        net::{IpAddr, Ipv4Addr, SocketAddr, SocketAddrV4},
2321        path::PathBuf,
2322        sync::{
2323            atomic::{AtomicUsize, Ordering},
2324            Arc, Mutex as StdMutex,
2325        },
2326        time::{Duration, Instant},
2327    };
2328
2329    use bytes::{Buf as _, BufMut as _};
2330    use bytesize::ByteSize;
2331    use metrics::{Key, Label};
2332    use saluki_common::sync::shutdown::ShutdownCoordinator;
2333    use saluki_core::{
2334        accounting::{ComponentRegistry, MemoryLimiter},
2335        components::{sources::SourceContext, ComponentContext},
2336        data_model::event::metric::context::{ContextResolverBuilder, TagsResolverBuilder},
2337        health::HealthRegistry,
2338        pooling::{helpers::get_pooled_object_via_builder, ObjectPool as _},
2339        runtime::state::DataspaceRegistry,
2340        support::SubsystemIdentifier,
2341        topology::{EventsBuffer, EventsDispatcher, OutputName, TopologyContext},
2342    };
2343    #[cfg(target_os = "linux")]
2344    use saluki_env::workload::providers::TestWorkloadProvider;
2345    use saluki_env::workload::{CaptureEntityResolver, EntityId};
2346    #[cfg(unix)]
2347    use saluki_io::net::Stream;
2348    use saluki_io::{
2349        buf::{BytesBuffer, FixedSizeVec},
2350        deser::codec::dogstatsd::{DogStatsDCodec, DogStatsDCodecConfiguration, ParsedPacket},
2351        net::{ConnectionAddress, ListenAddress, ProcessCredentials, ProcessIdentity},
2352    };
2353    use saluki_metrics::test::TestRecorder;
2354    use stringtheory::MetaString;
2355    #[cfg(unix)]
2356    use tokio::{
2357        io::AsyncWriteExt as _,
2358        net::{UnixDatagram, UnixStream},
2359    };
2360    use tokio::{
2361        runtime::Handle,
2362        sync::{mpsc, Mutex},
2363        task::yield_now,
2364        time::timeout,
2365    };
2366
2367    use super::{
2368        build_io_buffer_pool, capture_named_pipe_frame, default_decoder_worker_count, filters::EnablePayloadsFilter,
2369        handle_frame, handle_metric_packet, metrics::build_metrics, origin_detection_failed_for_telemetry,
2370        resolve_process_origin, resolve_process_origin_if_needed, shutdown_listeners_and_drain_datagram_decoders,
2371        BufferDecodeContext, BufferDecodeMode, ContextResolvers, DatagramSocketContext, DecodeOutcome, DecoderContext,
2372        DogStatsDConfiguration, DogStatsDDecoder, OriginEnrichmentConfiguration, ProcessOrigin, QueuedDatagram,
2373        ReceivedBuffer, TrafficCapture, TrafficCaptureReader,
2374    };
2375    #[cfg(unix)]
2376    use super::{receive_connected_stream, receive_connectionless_stream, received_payload};
2377    #[cfg(target_os = "linux")]
2378    use super::{DogStatsDOriginTagResolver, Listener};
2379    use crate::sources::dogstatsd::PacketForwarder;
2380
2381    const TEST_WINDOWS_PIPE_SECURITY_DESCRIPTOR: &str = "D:AI(A;;GA;;;WD)";
2382    const TEST_BUFFER_SIZE: usize = 8192;
2383
2384    fn test_component_context() -> ComponentContext {
2385        ComponentContext::test_source("dogstatsd_test")
2386    }
2387
2388    fn test_datagram_socket_context(listen_addr: ListenAddress) -> Arc<DatagramSocketContext> {
2389        Arc::new(DatagramSocketContext {
2390            metrics: build_metrics(&listen_addr, &test_component_context(), false),
2391            listen_addr,
2392            eol_required: false,
2393            packet_forwarder: None,
2394        })
2395    }
2396
2397    #[derive(Default)]
2398    struct CaptureTestEntityResolver {
2399        pid_map: StdMutex<HashMap<u32, EntityId>>,
2400        resolution_count: AtomicUsize,
2401    }
2402
2403    impl CaptureTestEntityResolver {
2404        fn with_pid_mapping(process_id: u32, entity_id: EntityId) -> Self {
2405            let mut pid_map = HashMap::new();
2406            pid_map.insert(process_id, entity_id);
2407            Self {
2408                pid_map: StdMutex::new(pid_map),
2409                resolution_count: AtomicUsize::new(0),
2410            }
2411        }
2412
2413        fn resolution_count(&self) -> usize {
2414            self.resolution_count.load(Ordering::Relaxed)
2415        }
2416
2417        #[cfg(target_os = "linux")]
2418        fn set_pid_mapping(&self, process_id: u32, entity_id: EntityId) {
2419            self.pid_map
2420                .lock()
2421                .expect("PID map lock should not be poisoned")
2422                .insert(process_id, entity_id);
2423        }
2424    }
2425
2426    impl CaptureEntityResolver for CaptureTestEntityResolver {
2427        fn resolve_container_entity_for_live_pid(&self, process_id: u32) -> Option<EntityId> {
2428            self.resolution_count.fetch_add(1, Ordering::Relaxed);
2429            self.pid_map
2430                .lock()
2431                .expect("PID map lock should not be poisoned")
2432                .get(&process_id)
2433                .cloned()
2434        }
2435    }
2436
2437    fn processed_metric_key(listener_type: &'static str, origin: Option<&str>) -> Key {
2438        let mut labels = vec![
2439            Label::from_static_parts("component_id", "dogstatsd_test"),
2440            Label::from_static_parts("component_type", "source"),
2441            Label::from_static_parts("listener_type", listener_type),
2442            Label::from_static_parts("message_type", "metrics"),
2443        ];
2444        if let Some(origin) = origin {
2445            labels.push(Label::new("origin", origin.to_string()));
2446        }
2447
2448        Key::from_parts("component_events_received_total", labels)
2449    }
2450
2451    fn test_context_resolvers() -> ContextResolvers {
2452        let tags_resolver = TagsResolverBuilder::for_tests().build();
2453        let context_resolver = ContextResolverBuilder::for_tests()
2454            .with_tags_resolver(Some(tags_resolver.clone()))
2455            .build();
2456        ContextResolvers::manual(context_resolver.clone(), context_resolver, tags_resolver)
2457    }
2458
2459    /// Builds a source context, plus the receiver for its `metrics` output.
2460    ///
2461    /// Only useful for tests that drive pieces of the source directly. A test that runs the source itself needs a
2462    /// running supervisor for it to spawn its children onto -- see `build_source`.
2463    fn test_source_context() -> (SourceContext, mpsc::Receiver<EventsBuffer>) {
2464        let component_context = test_component_context();
2465        let mut dispatcher = EventsDispatcher::new(component_context.clone());
2466        let metrics_output = OutputName::Given("metrics".into());
2467        dispatcher
2468            .add_output(metrics_output.clone())
2469            .expect("metrics output should be added");
2470        let (metrics_tx, metrics_rx) = mpsc::channel(4);
2471        dispatcher
2472            .attach_sender_to_output(&metrics_output, metrics_tx)
2473            .expect("metrics output should accept a sender");
2474
2475        let health_registry = HealthRegistry::new();
2476        let topology_context = TopologyContext::new(
2477            Arc::from("test"),
2478            MemoryLimiter::noop(),
2479            health_registry.clone(),
2480            Handle::current(),
2481            DataspaceRegistry::new(),
2482        );
2483        let health = health_registry
2484            .register_component(&SubsystemIdentifier::from_dotted("test.decoder"))
2485            .expect("test decoder should have a health handle");
2486        let source_context = SourceContext::new(
2487            &topology_context,
2488            &component_context,
2489            ComponentRegistry::default(),
2490            health,
2491            dispatcher,
2492        );
2493
2494        (source_context, metrics_rx)
2495    }
2496
2497    fn test_decoder_context(origin_detection_enabled: bool) -> DecoderContext {
2498        DecoderContext {
2499            codec: DogStatsDCodec::from_configuration(DogStatsDCodecConfiguration::default()),
2500            context_resolvers: test_context_resolvers(),
2501            default_hostname: MetaString::from_static("default-hostname"),
2502            enabled_filter: EnablePayloadsFilter::default(),
2503            origin_detection_enabled,
2504            stream_log_too_big: false,
2505            disable_verbose_logs: false,
2506            additional_tags: Vec::<String>::new().into(),
2507            traffic_capture: TrafficCapture::new(PathBuf::new(), 1),
2508        }
2509    }
2510
2511    fn test_io_buffer(payload: &[u8], capacity: usize) -> BytesBuffer {
2512        let mut buffer: BytesBuffer = get_pooled_object_via_builder(|| FixedSizeVec::with_capacity(capacity));
2513        buffer.put_slice(payload);
2514        buffer
2515    }
2516
2517    #[test]
2518    fn named_pipe_frames_are_written_to_an_active_capture() {
2519        let capture_directory = tempfile::tempdir().expect("temporary capture directory should be created");
2520        let capture = TrafficCapture::new(capture_directory.path().to_path_buf(), 1);
2521        let capture_path = capture
2522            .start_capture(None, Duration::from_secs(30), false)
2523            .expect("capture should start");
2524
2525        capture_named_pipe_frame(&capture, b"captured.named_pipe.one:1|c");
2526        capture_named_pipe_frame(&capture, b"captured.named_pipe.two:1|c");
2527        capture.stop_capture();
2528
2529        let deadline = Instant::now() + Duration::from_secs(2);
2530        while capture.is_ongoing() && Instant::now() < deadline {
2531            std::thread::sleep(Duration::from_millis(10));
2532        }
2533        assert!(!capture.is_ongoing(), "capture should stop after its sender is dropped");
2534
2535        let mut reader = TrafficCaptureReader::from_path(&capture_path).expect("capture should be readable");
2536        let first_record = reader
2537            .read_next()
2538            .expect("first capture record should decode")
2539            .expect("first named-pipe frame should be captured");
2540        let second_record = reader
2541            .read_next()
2542            .expect("second capture record should decode")
2543            .expect("second named-pipe frame should be captured");
2544        assert_eq!(first_record.payload, b"captured.named_pipe.one:1|c");
2545        assert_eq!(first_record.pid, 0);
2546        assert_eq!(second_record.payload, b"captured.named_pipe.two:1|c");
2547        assert_eq!(second_record.pid, 0);
2548        assert!(reader.read_next().expect("capture should terminate cleanly").is_none());
2549    }
2550
2551    #[tokio::test]
2552    async fn connectionless_decoder_dispatches_full_and_flushed_buffers_and_forwards_frames() {
2553        let recorder = TestRecorder::default();
2554        let _recorder_guard = metrics::set_default_local_recorder(&recorder);
2555        let listen_addr = ListenAddress::Unixgram("/tmp/dsd.sock".into());
2556        let peer_addr = ConnectionAddress::ProcessLike(ProcessIdentity::Unavailable);
2557        let metrics = build_metrics(&listen_addr, &test_component_context(), false);
2558        let (source_context, mut metrics_rx) = test_source_context();
2559        let mut decoder = DogStatsDDecoder::new(source_context, test_decoder_context(false));
2560        let event_buffer_capacity = EventsBuffer::default().capacity();
2561        let (packet_forwarder, mut packets_rx) = PacketForwarder::for_tests(event_buffer_capacity + 1);
2562        let mut buffer_decode_context = BufferDecodeContext::new(
2563            &listen_addr,
2564            false,
2565            &metrics,
2566            Some(&packet_forwarder),
2567            BufferDecodeMode::Connectionless,
2568        );
2569
2570        let mut payload = b"decoder.metric:1|c\n".repeat(event_buffer_capacity);
2571        payload.extend_from_slice(b"decoder.metric:1|c");
2572        let io_buffer = test_io_buffer(&payload, payload.len());
2573        let (received, returned_buffer) = ReceivedBuffer::with_return(io_buffer, payload.len(), peer_addr, None);
2574        let outcome = decoder.decode_buffer(&mut buffer_decode_context, received).await;
2575
2576        assert_eq!(outcome, DecodeOutcome::Continue);
2577        let io_buffer = returned_buffer
2578            .await
2579            .expect("connectionless decoder should return the I/O buffer");
2580        assert_eq!(io_buffer.remaining(), 0);
2581        let full_buffer = timeout(Duration::from_secs(1), metrics_rx.recv())
2582            .await
2583            .expect("full event buffer dispatch should not time out")
2584            .expect("metrics output should remain connected");
2585        assert_eq!(full_buffer.len(), event_buffer_capacity);
2586        assert!(
2587            timeout(Duration::from_secs(1), packets_rx.recv())
2588                .await
2589                .expect("forwarded packet should not time out")
2590                .is_some(),
2591            "forwarded packet should be queued"
2592        );
2593        assert_eq!(packets_rx.len(), event_buffer_capacity);
2594
2595        decoder.flush_events().await;
2596        let flushed_buffer = timeout(Duration::from_secs(1), metrics_rx.recv())
2597            .await
2598            .expect("partial event buffer flush should not time out")
2599            .expect("metrics output should remain connected");
2600        assert_eq!(flushed_buffer.len(), 1);
2601        decoder.flush_events().await;
2602        assert!(
2603            metrics_rx.try_recv().is_err(),
2604            "empty flush should not dispatch another buffer"
2605        );
2606
2607        assert_eq!(
2608            recorder.counter((
2609                "component_packets_received_total",
2610                &[
2611                    ("component_id", "dogstatsd_test"),
2612                    ("component_type", "source"),
2613                    ("listener_type", "unixgram"),
2614                    ("state", "ok"),
2615                ],
2616            )),
2617            Some(1)
2618        );
2619        assert_eq!(
2620            recorder.counter((
2621                "component_bytes_received_total",
2622                &[
2623                    ("component_id", "dogstatsd_test"),
2624                    ("component_type", "source"),
2625                    ("listener_type", "unixgram"),
2626                ],
2627            )),
2628            Some(payload.len() as u64)
2629        );
2630        assert_eq!(
2631            recorder.counter(processed_metric_key("unixgram", None)),
2632            Some((event_buffer_capacity + 1) as u64)
2633        );
2634    }
2635
2636    #[tokio::test]
2637    async fn connected_decoder_waits_for_complete_outer_frame_and_stops_on_eof() {
2638        let recorder = TestRecorder::default();
2639        let _recorder_guard = metrics::set_default_local_recorder(&recorder);
2640        let listen_addr = ListenAddress::Unix("/tmp/dsd.socket".into());
2641        let peer_addr = ConnectionAddress::ProcessLike(ProcessIdentity::Error(
2642            saluki_io::net::ProcessCredentialsError::InvalidCredentials,
2643        ));
2644        let metrics = build_metrics(&listen_addr, &test_component_context(), false);
2645        let (source_context, mut metrics_rx) = test_source_context();
2646        let mut decoder = DogStatsDDecoder::new(source_context, test_decoder_context(true));
2647        let (packet_forwarder, mut packets_rx) = PacketForwarder::for_tests(1);
2648        let mut buffer_decode_context = BufferDecodeContext::new(
2649            &listen_addr,
2650            false,
2651            &metrics,
2652            Some(&packet_forwarder),
2653            BufferDecodeMode::Connected,
2654        );
2655
2656        let frame = b"stream.metric:1|c\n";
2657        let mut payload = Vec::with_capacity(frame.len() + 4);
2658        payload.extend_from_slice(&(frame.len() as u32).to_le_bytes());
2659        payload.extend_from_slice(frame);
2660        let io_buffer = test_io_buffer(&payload[..2], payload.len() + 16);
2661        let (received, returned_buffer) = ReceivedBuffer::with_return(io_buffer, 2, peer_addr.clone(), None);
2662
2663        let partial_outcome = decoder.decode_buffer(&mut buffer_decode_context, received).await;
2664        assert_eq!(partial_outcome, DecodeOutcome::Continue);
2665        let mut io_buffer = returned_buffer
2666            .await
2667            .expect("connected decoder should return the partial I/O buffer");
2668        assert_eq!(io_buffer.remaining(), 2);
2669        assert_eq!(
2670            recorder.counter((
2671                "component_packets_received_total",
2672                &[
2673                    ("component_id", "dogstatsd_test"),
2674                    ("component_type", "source"),
2675                    ("listener_type", "unix"),
2676                    ("state", "ok"),
2677                ],
2678            )),
2679            Some(0)
2680        );
2681        assert_eq!(
2682            recorder.counter((
2683                "component_errors_total",
2684                &[
2685                    ("component_id", "dogstatsd_test"),
2686                    ("component_type", "source"),
2687                    ("error_type", "origin_detection"),
2688                ],
2689            )),
2690            Some(0)
2691        );
2692
2693        io_buffer.put_slice(&payload[2..]);
2694        let (received, returned_buffer) =
2695            ReceivedBuffer::with_return(io_buffer, payload.len() - 2, peer_addr.clone(), None);
2696        let complete_outcome = decoder.decode_buffer(&mut buffer_decode_context, received).await;
2697        assert_eq!(complete_outcome, DecodeOutcome::Continue);
2698        let io_buffer = returned_buffer
2699            .await
2700            .expect("connected decoder should return the consumed I/O buffer");
2701        assert_eq!(io_buffer.remaining(), 0);
2702        assert!(
2703            timeout(Duration::from_secs(1), packets_rx.recv())
2704                .await
2705                .expect("forwarded stream packet should not time out")
2706                .is_some(),
2707            "complete outer frame should be forwarded"
2708        );
2709
2710        let (received, returned_buffer) = ReceivedBuffer::with_return(io_buffer, 0, peer_addr, None);
2711        let eof_outcome = decoder.decode_buffer(&mut buffer_decode_context, received).await;
2712        assert_eq!(eof_outcome, DecodeOutcome::Stop);
2713        let returned_buffer = returned_buffer
2714            .await
2715            .expect("stopped decoder should return the I/O buffer");
2716        assert!(!returned_buffer.has_remaining());
2717
2718        decoder.flush_events().await;
2719        let flushed_buffer = timeout(Duration::from_secs(1), metrics_rx.recv())
2720            .await
2721            .expect("connected event buffer flush should not time out")
2722            .expect("metrics output should remain connected");
2723        assert_eq!(flushed_buffer.len(), 1);
2724        assert_eq!(
2725            recorder.counter((
2726                "component_packets_received_total",
2727                &[
2728                    ("component_id", "dogstatsd_test"),
2729                    ("component_type", "source"),
2730                    ("listener_type", "unix"),
2731                    ("state", "ok"),
2732                ],
2733            )),
2734            Some(1)
2735        );
2736        assert_eq!(
2737            recorder.counter((
2738                "component_bytes_received_total",
2739                &[
2740                    ("component_id", "dogstatsd_test"),
2741                    ("component_type", "source"),
2742                    ("listener_type", "unix"),
2743                ],
2744            )),
2745            Some(payload.len() as u64)
2746        );
2747        assert_eq!(
2748            recorder.counter((
2749                "component_errors_total",
2750                &[
2751                    ("component_id", "dogstatsd_test"),
2752                    ("component_type", "source"),
2753                    ("error_type", "origin_detection"),
2754                ],
2755            )),
2756            Some(1)
2757        );
2758        assert_eq!(recorder.counter(processed_metric_key("unix", None)), Some(1));
2759    }
2760
2761    #[tokio::test]
2762    async fn decoder_continues_connectionless_but_stops_connected_on_framing_error() {
2763        let recorder = TestRecorder::default();
2764        let _recorder_guard = metrics::set_default_local_recorder(&recorder);
2765        let peer_addr = ConnectionAddress::ProcessLike(ProcessIdentity::Unavailable);
2766        let (source_context, _metrics_rx) = test_source_context();
2767        let mut decoder = DogStatsDDecoder::new(source_context, test_decoder_context(false));
2768
2769        let datagram_addr = ListenAddress::Unixgram("/tmp/dsd.sock".into());
2770        let datagram_metrics = build_metrics(&datagram_addr, &test_component_context(), false);
2771        let payload = b"missing.newline:1|c";
2772        let datagram_buffer = test_io_buffer(payload, payload.len());
2773        let mut datagram_context = BufferDecodeContext::new(
2774            &datagram_addr,
2775            true,
2776            &datagram_metrics,
2777            None,
2778            BufferDecodeMode::Connectionless,
2779        );
2780        let (received, _returned_datagram_buffer) =
2781            ReceivedBuffer::with_return(datagram_buffer, payload.len(), peer_addr.clone(), None);
2782        let datagram_outcome = decoder.decode_buffer(&mut datagram_context, received).await;
2783        assert_eq!(datagram_outcome, DecodeOutcome::Continue);
2784
2785        let stream_addr = ListenAddress::Unix("/tmp/dsd.socket".into());
2786        let stream_metrics = build_metrics(&stream_addr, &test_component_context(), false);
2787        let stream_buffer_capacity = 64;
2788        let oversized_frame = (stream_buffer_capacity as u32).to_le_bytes();
2789        let stream_buffer = test_io_buffer(&oversized_frame, stream_buffer_capacity);
2790        let mut stream_context =
2791            BufferDecodeContext::new(&stream_addr, false, &stream_metrics, None, BufferDecodeMode::Connected);
2792        let (received, _returned_stream_buffer) =
2793            ReceivedBuffer::with_return(stream_buffer, oversized_frame.len(), peer_addr, None);
2794        let stream_outcome = decoder.decode_buffer(&mut stream_context, received).await;
2795        assert_eq!(stream_outcome, DecodeOutcome::Stop);
2796
2797        for listener_type in ["unixgram", "unix"] {
2798            assert_eq!(
2799                recorder.counter((
2800                    "component_errors_total",
2801                    &[
2802                        ("component_id", "dogstatsd_test"),
2803                        ("component_type", "source"),
2804                        ("listener_type", listener_type),
2805                        ("error_type", "framing"),
2806                    ],
2807                )),
2808                Some(1)
2809            );
2810        }
2811        assert_eq!(
2812            recorder.counter((
2813                "component_packets_received_total",
2814                &[
2815                    ("component_id", "dogstatsd_test"),
2816                    ("component_type", "source"),
2817                    ("listener_type", "unixgram"),
2818                    ("state", "ok"),
2819                ],
2820            )),
2821            Some(1)
2822        );
2823        assert_eq!(
2824            recorder.counter((
2825                "component_packets_received_total",
2826                &[
2827                    ("component_id", "dogstatsd_test"),
2828                    ("component_type", "source"),
2829                    ("listener_type", "unix"),
2830                    ("state", "ok"),
2831                ],
2832            )),
2833            Some(0)
2834        );
2835    }
2836
2837    #[test]
2838    fn origin_telemetry_does_not_resolve_origin_when_origin_detection_is_disabled() {
2839        let recorder = TestRecorder::default();
2840        let _recorder_guard = metrics::set_default_local_recorder(&recorder);
2841        let listen_addr = ListenAddress::Unixgram("/tmp/dsd.sock".into());
2842        let context = test_component_context();
2843        let metrics = build_metrics(&listen_addr, &context, true);
2844        let codec = DogStatsDCodec::from_configuration(DogStatsDCodecConfiguration::default());
2845        let mut context_resolvers = test_context_resolvers();
2846        let capture_entity_resolver = CaptureTestEntityResolver::with_pid_mapping(
2847            42,
2848            EntityId::from_local_data("ci-pid-container").expect("container entity"),
2849        );
2850        let peer_addr = ConnectionAddress::ProcessLike(ProcessIdentity::Credentials(ProcessCredentials {
2851            pid: 42,
2852            uid: 0,
2853            gid: 0,
2854        }));
2855        let process_origin = resolve_process_origin(Some(&capture_entity_resolver), &peer_addr);
2856
2857        let event = handle_frame(
2858            b"test_metric:1|c",
2859            &codec,
2860            &mut context_resolvers,
2861            &metrics,
2862            false,
2863            process_origin.as_ref(),
2864            EnablePayloadsFilter::default(),
2865            &[],
2866            &MetaString::from_static("default-host"),
2867        )
2868        .expect("frame should parse");
2869
2870        assert!(event.is_some());
2871        assert_eq!(
2872            recorder.counter(processed_metric_key("unixgram", Some("container_id://pid-container"))),
2873            None
2874        );
2875        assert_eq!(recorder.counter(processed_metric_key("unixgram", Some(""))), Some(1));
2876    }
2877
2878    #[test]
2879    fn origin_telemetry_records_resolved_origin_when_origin_detection_is_enabled() {
2880        let recorder = TestRecorder::default();
2881        let _recorder_guard = metrics::set_default_local_recorder(&recorder);
2882        let listen_addr = ListenAddress::Unixgram("/tmp/dsd.sock".into());
2883        let context = test_component_context();
2884        let metrics = build_metrics(&listen_addr, &context, true);
2885        let codec = DogStatsDCodec::from_configuration(DogStatsDCodecConfiguration::default());
2886        let mut context_resolvers = test_context_resolvers();
2887        let capture_entity_resolver = CaptureTestEntityResolver::with_pid_mapping(
2888            42,
2889            EntityId::from_local_data("ci-pid-container").expect("container entity"),
2890        );
2891        let peer_addr = ConnectionAddress::ProcessLike(ProcessIdentity::Credentials(ProcessCredentials {
2892            pid: 42,
2893            uid: 0,
2894            gid: 0,
2895        }));
2896        let process_origin = resolve_process_origin(Some(&capture_entity_resolver), &peer_addr);
2897
2898        let event = handle_frame(
2899            b"test_metric:1|c",
2900            &codec,
2901            &mut context_resolvers,
2902            &metrics,
2903            true,
2904            process_origin.as_ref(),
2905            EnablePayloadsFilter::default(),
2906            &[],
2907            &MetaString::from_static("default-host"),
2908        )
2909        .expect("frame should parse");
2910
2911        assert!(event.is_some());
2912        assert_eq!(
2913            recorder.counter(processed_metric_key("unixgram", Some("container_id://pid-container"))),
2914            Some(1)
2915        );
2916        assert_eq!(recorder.counter(processed_metric_key("unixgram", Some(""))), Some(0));
2917    }
2918
2919    #[test]
2920    fn no_metrics_when_interner_full_allocations_disallowed() {
2921        // We're specifically testing here that when we don't allow outside allocations, we should not be able to
2922        // resolve a context if the interner is full. A no-op interner has the smallest possible size, so that's going
2923        // to assure we can't intern anything... but we also need a string (name or one of the tags) that can't be
2924        // _inlined_ either, since that will get around the interner being full.
2925        //
2926        // We set our metric name to be longer than 31 bytes (the inlining limit) to ensure this.
2927
2928        let codec = DogStatsDCodec::from_configuration(DogStatsDCodecConfiguration::default());
2929        let tags_resolver = TagsResolverBuilder::for_tests().build();
2930        let context_resolver = ContextResolverBuilder::for_tests()
2931            .with_heap_allocations(false)
2932            .with_tags_resolver(Some(tags_resolver.clone()))
2933            .build();
2934        let mut context_resolvers = ContextResolvers::manual(context_resolver.clone(), context_resolver, tags_resolver);
2935        let input = "big_metric_name_that_cant_possibly_be_inlined:1|c|#tag1:value1,tag2:value2,tag3:value3";
2936
2937        let Ok(ParsedPacket::Metric(packet)) = codec.decode_packet(input.as_bytes()) else {
2938            panic!("Failed to parse packet.");
2939        };
2940
2941        let maybe_metric = handle_metric_packet(
2942            packet,
2943            &mut context_resolvers,
2944            None,
2945            &[],
2946            &MetaString::from_static("default-host"),
2947        );
2948        assert!(maybe_metric.is_none());
2949    }
2950
2951    #[test]
2952    fn metric_host_tag_disambiguates_contexts_without_remaining_tag() {
2953        let codec = DogStatsDCodec::from_configuration(DogStatsDCodecConfiguration::default());
2954        let mut context_resolvers = test_context_resolvers();
2955        let default_hostname = MetaString::from_static("default-host");
2956
2957        let packets = [
2958            ("unset", b"test_metric_name:1|g".as_slice(), "default-host"),
2959            ("empty", b"test_metric_name:2|g|#host:".as_slice(), ""),
2960            (
2961                "explicit_default",
2962                b"test_metric_name:3|g|#host:default-host".as_slice(),
2963                "default-host",
2964            ),
2965            (
2966                "custom",
2967                b"test_metric_name:4|g|#host:custom-host".as_slice(),
2968                "custom-host",
2969            ),
2970        ];
2971
2972        let mut metrics = Vec::new();
2973        for (case, raw, expected_host) in packets {
2974            let Ok(ParsedPacket::Metric(packet)) = codec.decode_packet(raw) else {
2975                panic!("Failed to parse {case} packet.");
2976            };
2977            let metric = handle_metric_packet(packet, &mut context_resolvers, None, &[], &default_hostname)
2978                .unwrap_or_else(|| panic!("{case} metric should resolve"));
2979
2980            assert_eq!(metric.context().host(), Some(expected_host), "{case} context host");
2981            assert!(metric.context().tags().into_iter().all(|tag| tag.name() != "host"));
2982            metrics.push(metric);
2983        }
2984
2985        assert_eq!(metrics[0].context(), metrics[2].context());
2986        assert_ne!(metrics[0].context(), metrics[1].context());
2987        assert_ne!(metrics[0].context(), metrics[3].context());
2988        assert_ne!(metrics[1].context(), metrics[3].context());
2989    }
2990
2991    #[test]
2992    fn additional_tags_are_sorted_and_deduplicated() {
2993        let config = DogStatsDConfiguration {
2994            additional_tags: vec![
2995                "dogstatsd:configured".to_string(),
2996                "env:prod".to_string(),
2997                "provider_kind:autopilot".to_string(),
2998                "kube_distribution:eks".to_string(),
2999                "env:prod".to_string(),
3000            ],
3001            ..DogStatsDConfiguration::for_test()
3002        };
3003
3004        assert_eq!(
3005            config.additional_tags(),
3006            [
3007                "dogstatsd:configured",
3008                "env:prod",
3009                "kube_distribution:eks",
3010                "provider_kind:autopilot",
3011            ]
3012        );
3013    }
3014
3015    #[test]
3016    fn metric_with_additional_tags() {
3017        let codec = DogStatsDCodec::from_configuration(DogStatsDCodecConfiguration::default());
3018        let tags_resolver = TagsResolverBuilder::for_tests().build();
3019        let context_resolver = ContextResolverBuilder::for_tests()
3020            .with_heap_allocations(false)
3021            .with_tags_resolver(Some(tags_resolver.clone()))
3022            .build();
3023        let mut context_resolvers = ContextResolvers::manual(context_resolver.clone(), context_resolver, tags_resolver);
3024        let existing_tags = ["tag1:value1", "tag2:value2", "tag3:value3"];
3025        let existing_tags_str = existing_tags.join(",");
3026
3027        let input = format!("test_metric_name:1|c|#{}", existing_tags_str);
3028        let additional_tags = [
3029            "tag4:value4".to_string(),
3030            "tag5:value5".to_string(),
3031            "tag6:value6".to_string(),
3032        ];
3033
3034        let Ok(ParsedPacket::Metric(packet)) = codec.decode_packet(input.as_bytes()) else {
3035            panic!("Failed to parse packet.");
3036        };
3037        let maybe_metric = handle_metric_packet(
3038            packet,
3039            &mut context_resolvers,
3040            None,
3041            &additional_tags,
3042            &MetaString::from_static("default-host"),
3043        );
3044        assert!(maybe_metric.is_some());
3045
3046        let metric = maybe_metric.unwrap();
3047        let context = metric.context();
3048
3049        for tag in existing_tags {
3050            assert!(context.tags().has_tag(tag));
3051        }
3052
3053        for tag in additional_tags {
3054            assert!(context.tags().has_tag(tag));
3055        }
3056    }
3057
3058    fn udp_listen_address() -> ListenAddress {
3059        ListenAddress::udp_loopback(8125)
3060    }
3061
3062    fn tcp_listen_address() -> ListenAddress {
3063        ListenAddress::tcp_loopback(8125)
3064    }
3065
3066    fn named_pipe_listen_address() -> ListenAddress {
3067        ListenAddress::named_pipe_with_input_buffer_size(
3068            "datadog-dogstatsd",
3069            TEST_WINDOWS_PIPE_SECURITY_DESCRIPTOR,
3070            TEST_BUFFER_SIZE as u32,
3071        )
3072    }
3073
3074    #[test]
3075    fn build_addresses_includes_named_pipe_when_configured() {
3076        let config = DogStatsDConfiguration {
3077            port: 0,
3078            pipe_name: Some("datadog-dogstatsd".to_string()),
3079            ..DogStatsDConfiguration::for_test()
3080        };
3081
3082        let addresses = config.build_addresses(None);
3083
3084        assert_eq!(addresses, vec![named_pipe_listen_address()]);
3085    }
3086
3087    #[test]
3088    fn build_addresses_uses_buffer_size_for_named_pipe_input_buffer() {
3089        let config = DogStatsDConfiguration {
3090            port: 0,
3091            pipe_name: Some("datadog-dogstatsd".to_string()),
3092            buffer_size: 16384,
3093            ..DogStatsDConfiguration::for_test()
3094        };
3095
3096        let addresses = config.build_addresses(None);
3097
3098        let [ListenAddress::NamedPipe { input_buffer_size, .. }] = addresses.as_slice() else {
3099            panic!("expected only a named pipe listen address, got {addresses:?}");
3100        };
3101        assert_eq!(*input_buffer_size, Some(16_384));
3102    }
3103
3104    #[test]
3105    fn eol_required_matches_named_pipe_listener_type() {
3106        let config = DogStatsDConfiguration {
3107            eol_required: vec!["named_pipe".to_string()],
3108            ..DogStatsDConfiguration::for_test()
3109        };
3110        let eol_required = config.eol_required();
3111
3112        assert!(eol_required.for_listener(&named_pipe_listen_address()));
3113        assert!(!eol_required.for_listener(&udp_listen_address()));
3114        assert!(!eol_required.for_listener(&tcp_listen_address()));
3115    }
3116
3117    #[test]
3118    fn statsd_forward_no_host_disabled() {
3119        let config = DogStatsDConfiguration {
3120            statsd_forward_port: 9125,
3121            ..DogStatsDConfiguration::for_test()
3122        };
3123        assert!(config.statsd_forward_target().is_none());
3124    }
3125
3126    #[test]
3127    fn statsd_forward_zero_port_disabled() {
3128        let config = DogStatsDConfiguration {
3129            statsd_forward_host: Some(MetaString::from("127.0.0.1")),
3130            statsd_forward_port: 0,
3131            ..DogStatsDConfiguration::for_test()
3132        };
3133        assert!(config.statsd_forward_target().is_none());
3134    }
3135
3136    #[test]
3137    fn statsd_forward_host_and_port_enabled() {
3138        let config = DogStatsDConfiguration {
3139            statsd_forward_host: Some(MetaString::from("127.0.0.1")),
3140            statsd_forward_port: 9125,
3141            ..DogStatsDConfiguration::for_test()
3142        };
3143        let (host, port) = config.statsd_forward_target().expect("forwarding should be enabled");
3144        assert_eq!(host.as_ref(), "127.0.0.1");
3145        assert_eq!(port, 9125);
3146    }
3147
3148    #[test]
3149    fn statsd_forward_invalid_target_still_builds_forwarder_handle() {
3150        let config = DogStatsDConfiguration {
3151            statsd_forward_host: Some(MetaString::from("not a valid host")),
3152            statsd_forward_port: 9125,
3153            ..DogStatsDConfiguration::for_test()
3154        };
3155        assert!(config.packet_forwarder_target().is_some());
3156    }
3157
3158    #[test]
3159    fn unsupported_platform_process_credentials_do_not_count_as_origin_detection_telemetry_errors() {
3160        let peer_addr = ConnectionAddress::ProcessLike(ProcessIdentity::Error(
3161            saluki_io::net::ProcessCredentialsError::UnsupportedPlatform,
3162        ));
3163
3164        assert!(!origin_detection_failed_for_telemetry(true, 1, &peer_addr));
3165    }
3166
3167    #[test]
3168    fn invalid_process_credentials_count_as_origin_detection_telemetry_errors() {
3169        let peer_addr = ConnectionAddress::ProcessLike(ProcessIdentity::Error(
3170            saluki_io::net::ProcessCredentialsError::InvalidCredentials,
3171        ));
3172
3173        assert!(origin_detection_failed_for_telemetry(true, 1, &peer_addr));
3174    }
3175
3176    #[test]
3177    fn autoscaling_disabled_yields_a_single_udp_stream() {
3178        let config = DogStatsDConfiguration {
3179            autoscale_udp_listeners: false,
3180            ..DogStatsDConfiguration::for_test()
3181        };
3182        assert!(config.udp_streams_to_yield().is_none());
3183    }
3184
3185    #[test]
3186    fn effective_max_buffer_count_never_below_baseline() {
3187        fn buffer_counts(buffer_count: usize, buffer_count_max: usize) -> DogStatsDConfiguration {
3188            DogStatsDConfiguration {
3189                buffer_count,
3190                buffer_count_max,
3191                ..DogStatsDConfiguration::for_test()
3192            }
3193        }
3194
3195        // A config that only raised the baseline keeps its full capacity rather than being capped to the maximum.
3196        assert_eq!(buffer_counts(65536, 32_768).effective_max_buffer_count(), 65536);
3197
3198        // An explicit maximum above the baseline is honored as-is.
3199        assert_eq!(buffer_counts(128, 512).effective_max_buffer_count(), 512);
3200
3201        // A maximum below the baseline is treated as equal to the baseline.
3202        assert_eq!(buffer_counts(200, 64).effective_max_buffer_count(), 200);
3203    }
3204
3205    #[test]
3206    fn decoder_worker_count_matches_core_agent_defaults() {
3207        assert_eq!(default_decoder_worker_count(1), 2);
3208        assert_eq!(default_decoder_worker_count(4), 2);
3209        assert_eq!(default_decoder_worker_count(8), 6);
3210    }
3211
3212    #[test]
3213    fn decoder_worker_count_honors_explicit_override() {
3214        let config = DogStatsDConfiguration {
3215            workers_count: 1,
3216            ..DogStatsDConfiguration::for_test()
3217        };
3218
3219        assert_eq!(config.decoder_worker_count().get(), 1);
3220    }
3221
3222    #[tokio::test]
3223    async fn global_datagram_receiver_distributes_packets_to_workers() {
3224        let (sender, receiver) = mpsc::channel(2);
3225        let receiver = Arc::new(Mutex::new(receiver));
3226        let first_worker = receiver.clone();
3227        let second_worker = receiver;
3228        let socket_context = test_datagram_socket_context(udp_listen_address());
3229
3230        sender
3231            .send(QueuedDatagram {
3232                result: Err(std::io::Error::other("first")),
3233                socket_context: socket_context.clone(),
3234            })
3235            .await
3236            .expect("first packet should be queued");
3237        sender
3238            .send(QueuedDatagram {
3239                result: Err(std::io::Error::other("second")),
3240                socket_context,
3241            })
3242            .await
3243            .expect("second packet should be queued");
3244
3245        let (first, second) = tokio::join!(async { first_worker.lock().await.recv().await }, async {
3246            second_worker.lock().await.recv().await
3247        },);
3248
3249        assert!(first.expect("first worker should receive a packet").result.is_err());
3250        assert!(second.expect("second worker should receive a packet").result.is_err());
3251    }
3252
3253    #[tokio::test]
3254    async fn shutdown_drains_queued_datagrams_after_listeners_stop() {
3255        let mut listener_shutdown_coordinator = ShutdownCoordinator::default();
3256        let listener_shutdown = listener_shutdown_coordinator.register();
3257        let (sender, mut receiver) = mpsc::channel(2);
3258        sender.send(()).await.expect("first datagram should be queued");
3259        sender.send(()).await.expect("second datagram should be queued");
3260
3261        let listener_task = tokio::spawn(async move {
3262            listener_shutdown.await;
3263            drop(sender);
3264        });
3265        let decoded = Arc::new(AtomicUsize::new(0));
3266        let decoder_count = decoded.clone();
3267
3268        // The decoder reports completion by dropping its handle, exactly as the supervised children do.
3269        let mut decoder_shutdown_coordinator = ShutdownCoordinator::default();
3270        let decoder_shutdown = decoder_shutdown_coordinator.register();
3271        let decoder_task = tokio::spawn(async move {
3272            let _decoder_shutdown = decoder_shutdown;
3273            while receiver.recv().await.is_some() {
3274                decoder_count.fetch_add(1, Ordering::Relaxed);
3275            }
3276        });
3277
3278        shutdown_listeners_and_drain_datagram_decoders(listener_shutdown_coordinator, decoder_shutdown_coordinator)
3279            .await;
3280        listener_task.await.expect("listener task should stop cleanly");
3281        decoder_task.await.expect("decoder task should stop cleanly");
3282
3283        assert_eq!(decoded.load(Ordering::Relaxed), 2);
3284    }
3285
3286    #[tokio::test]
3287    async fn dogstatsd_io_buffer_pool_grows_on_demand_until_limit() {
3288        let min_buffers = 2;
3289        let max_buffers = 3;
3290        let (pool, shrinker) = build_io_buffer_pool(min_buffers, max_buffers, TEST_BUFFER_SIZE);
3291
3292        let mut initial_buffers = Vec::with_capacity(min_buffers);
3293        for _ in 0..min_buffers {
3294            initial_buffers.push(
3295                timeout(Duration::from_secs(1), pool.acquire())
3296                    .await
3297                    .expect("initial buffer should be available"),
3298            );
3299        }
3300        let on_demand_buffer = timeout(Duration::from_secs(1), pool.acquire())
3301            .await
3302            .expect("pool should grow on demand before hitting the limit");
3303
3304        let capped_acquire = timeout(Duration::from_millis(25), pool.acquire()).await;
3305        assert!(capped_acquire.is_err(), "pool should wait once it reaches the limit");
3306
3307        drop(initial_buffers.pop().expect("initial buffer should still be held"));
3308        timeout(Duration::from_secs(1), pool.acquire())
3309            .await
3310            .expect("returned buffer should unblock acquisition");
3311
3312        drop(on_demand_buffer);
3313        drop(shrinker);
3314    }
3315
3316    #[cfg(unix)]
3317    #[tokio::test]
3318    async fn uds_datagram_reader_is_bounded_by_io_buffer_pool() {
3319        let temp_dir = tempfile::tempdir().expect("temp directory should be created");
3320        let socket_path = temp_dir.path().join("dogstatsd.socket");
3321        let receiver = UnixDatagram::bind(&socket_path).expect("receiver should bind");
3322        let sender = UnixDatagram::unbound().expect("sender should be created");
3323        let (pool, shrinker) = build_io_buffer_pool(2, 2, TEST_BUFFER_SIZE);
3324        let (packets_tx, mut packets_rx) = mpsc::channel(3);
3325        let listen_addr = ListenAddress::Unixgram(socket_path.clone());
3326        let socket_context = Arc::new(DatagramSocketContext {
3327            metrics: build_metrics(&listen_addr, &test_component_context(), false),
3328            listen_addr,
3329            eol_required: false,
3330            packet_forwarder: None,
3331        });
3332        let reader = tokio::spawn(receive_connectionless_stream(
3333            Stream::from(receiver),
3334            pool,
3335            MemoryLimiter::noop(),
3336            false,
3337            TrafficCapture::new(PathBuf::new(), 1),
3338            None,
3339            packets_tx,
3340            socket_context,
3341        ));
3342        let payloads: [&[u8]; 3] = [b"first", b"second", b"third"];
3343
3344        for payload in payloads {
3345            sender
3346                .send_to(payload, &socket_path)
3347                .await
3348                .expect("payload should send");
3349        }
3350
3351        timeout(Duration::from_secs(1), async {
3352            while packets_rx.len() < 2 {
3353                yield_now().await;
3354            }
3355        })
3356        .await
3357        .expect("reader should fill the two-buffer pool");
3358        assert_eq!(packets_rx.len(), 2);
3359
3360        let first = packets_rx
3361            .recv()
3362            .await
3363            .expect("first packet should be queued")
3364            .result
3365            .expect("first receive should succeed");
3366        assert_eq!(received_payload(first.buffer(), first.bytes_read), payloads[0]);
3367        drop(first);
3368
3369        timeout(Duration::from_secs(1), async {
3370            while packets_rx.len() < 2 {
3371                yield_now().await;
3372            }
3373        })
3374        .await
3375        .expect("returning a buffer should allow the third packet to be read");
3376
3377        for expected in &payloads[1..] {
3378            let received = packets_rx
3379                .recv()
3380                .await
3381                .expect("packet should be queued")
3382                .result
3383                .expect("receive should succeed");
3384            assert_eq!(received_payload(received.buffer(), received.bytes_read), *expected);
3385        }
3386
3387        drop(packets_rx);
3388        sender
3389            .send_to(b"shutdown", &socket_path)
3390            .await
3391            .expect("shutdown payload should send");
3392        timeout(Duration::from_secs(1), reader)
3393            .await
3394            .expect("reader should stop after observing the closed queue")
3395            .expect("reader task should not panic");
3396        drop(shrinker);
3397    }
3398
3399    #[cfg(target_os = "linux")]
3400    #[tokio::test]
3401    async fn uds_datagram_reader_pins_origin_before_decode() {
3402        let temp_dir = tempfile::tempdir().expect("temp directory should be created");
3403        let socket_path = temp_dir.path().join("dogstatsd.socket");
3404        let listen_addr = ListenAddress::Unixgram(socket_path.clone());
3405        let mut listener = Listener::from_listen_address(listen_addr.clone(), None)
3406            .await
3407            .expect("listener should bind");
3408        let stream = listener.accept().await.expect("listener should yield its socket");
3409        let sender = UnixDatagram::unbound().expect("sender should be created");
3410        let process_id = std::process::id();
3411        let original_entity = EntityId::from_local_data("ci-original-container").expect("container entity");
3412        let reused_entity = EntityId::from_local_data("ci-reused-container").expect("container entity");
3413        let capture_entity_resolver = Arc::new(CaptureTestEntityResolver::with_pid_mapping(
3414            process_id,
3415            original_entity.clone(),
3416        ));
3417        let (pool, shrinker) = build_io_buffer_pool(1, 1, TEST_BUFFER_SIZE);
3418        let (packets_tx, mut packets_rx) = mpsc::channel(1);
3419        let reader = tokio::spawn(receive_connectionless_stream(
3420            stream,
3421            pool,
3422            MemoryLimiter::noop(),
3423            true,
3424            TrafficCapture::new(PathBuf::new(), 1),
3425            Some(capture_entity_resolver.clone()),
3426            packets_tx,
3427            test_datagram_socket_context(listen_addr),
3428        ));
3429
3430        sender
3431            .send_to(b"test.metric:1|c", &socket_path)
3432            .await
3433            .expect("payload should send");
3434        let received = timeout(Duration::from_secs(1), packets_rx.recv())
3435            .await
3436            .expect("packet should be received")
3437            .expect("reader should remain active")
3438            .result
3439            .expect("receive should succeed");
3440        assert_eq!(
3441            received.process_origin,
3442            Some(ProcessOrigin::Pinned(Some(original_entity.clone())))
3443        );
3444        assert_eq!(capture_entity_resolver.resolution_count(), 1);
3445
3446        // Simulate the sender exiting and its PID being reused before a decoder worker reaches the queued packet.
3447        capture_entity_resolver.set_pid_mapping(process_id, reused_entity.clone());
3448
3449        let mut workload_provider = TestWorkloadProvider::new();
3450        workload_provider.add_entity(original_entity, &["container:original"]);
3451        workload_provider.add_entity(reused_entity, &["container:reused"]);
3452        let origin_config = OriginEnrichmentConfiguration {
3453            enabled: true,
3454            ..OriginEnrichmentConfiguration::for_test()
3455        };
3456        let origin_resolver = DogStatsDOriginTagResolver::new(
3457            origin_config,
3458            Arc::new(workload_provider),
3459            super::CapturedTaggerHandle::new(),
3460        );
3461        let tags_resolver = TagsResolverBuilder::for_tests().build();
3462        let context_resolver = ContextResolverBuilder::for_tests()
3463            .with_tags_resolver(Some(tags_resolver.clone()))
3464            .build();
3465        let mut context_resolvers = ContextResolvers::manual_with_origin(
3466            context_resolver.clone(),
3467            context_resolver,
3468            tags_resolver,
3469            origin_resolver,
3470        );
3471        let codec = DogStatsDCodec::from_configuration(DogStatsDCodecConfiguration::default());
3472        let Ok(ParsedPacket::Metric(packet)) = codec.decode_packet(b"test.metric:1|c") else {
3473            panic!("metric should parse");
3474        };
3475        let metric = handle_metric_packet(
3476            packet,
3477            &mut context_resolvers,
3478            received.process_origin.as_ref(),
3479            &[],
3480            &MetaString::from_static("default-host"),
3481        )
3482        .expect("metric context should resolve");
3483
3484        assert!(metric.context().origin_tags().has_tag("container:original"));
3485        assert!(!metric.context().origin_tags().has_tag("container:reused"));
3486
3487        drop(received);
3488        drop(packets_rx);
3489        sender
3490            .send_to(b"shutdown", &socket_path)
3491            .await
3492            .expect("shutdown payload should send");
3493        timeout(Duration::from_secs(1), reader)
3494            .await
3495            .expect("reader should stop after observing the closed queue")
3496            .expect("reader task should not panic");
3497        drop(shrinker);
3498    }
3499
3500    #[cfg(unix)]
3501    #[tokio::test]
3502    async fn connection_oriented_reader_preserves_partial_frames() {
3503        let (mut sender, receiver) = UnixStream::pair().expect("stream pair should be created");
3504        let (pool, shrinker) = build_io_buffer_pool(1, 1, TEST_BUFFER_SIZE);
3505        let (packets_tx, mut packets_rx) = mpsc::channel(1);
3506        let reader = tokio::spawn(receive_connected_stream(
3507            Stream::from(receiver),
3508            pool,
3509            MemoryLimiter::noop(),
3510            false,
3511            TrafficCapture::new(PathBuf::new(), 1),
3512            None,
3513            packets_tx,
3514        ));
3515
3516        sender.write_all(b"partial").await.expect("first payload should send");
3517        let first = timeout(Duration::from_secs(1), packets_rx.recv())
3518            .await
3519            .expect("first read should finish")
3520            .expect("reader should remain active")
3521            .expect("first read should succeed");
3522        assert_eq!(first.buffer().chunk(), b"partial");
3523        drop(first);
3524
3525        sender.write_all(b"-frame").await.expect("second payload should send");
3526        let second = timeout(Duration::from_secs(1), packets_rx.recv())
3527            .await
3528            .expect("second read should finish")
3529            .expect("reader should remain active")
3530            .expect("second read should succeed");
3531        assert_eq!(second.bytes_read, b"-frame".len());
3532        assert_eq!(second.buffer().chunk(), b"partial-frame");
3533        drop(second);
3534
3535        drop(packets_rx);
3536        sender
3537            .write_all(b"shutdown")
3538            .await
3539            .expect("shutdown payload should send");
3540        timeout(Duration::from_secs(1), reader)
3541            .await
3542            .expect("reader should stop after observing the closed queue")
3543            .expect("reader task should not panic");
3544        drop(shrinker);
3545    }
3546
3547    #[cfg(unix)]
3548    #[tokio::test]
3549    async fn connection_oriented_reader_releases_drained_buffer_before_reacquiring() {
3550        let (mut sender, receiver) = UnixStream::pair().expect("stream pair should be created");
3551        let (pool, shrinker) = build_io_buffer_pool(1, 1, TEST_BUFFER_SIZE);
3552        let (packets_tx, mut packets_rx) = mpsc::channel(1);
3553        let reader = tokio::spawn(receive_connected_stream(
3554            Stream::from(receiver),
3555            pool,
3556            MemoryLimiter::noop(),
3557            false,
3558            TrafficCapture::new(PathBuf::new(), 1),
3559            None,
3560            packets_tx,
3561        ));
3562
3563        sender.write_all(b"first").await.expect("first payload should send");
3564        let mut first = timeout(Duration::from_secs(1), packets_rx.recv())
3565            .await
3566            .expect("first read should finish")
3567            .expect("reader should remain active")
3568            .expect("first read should succeed");
3569        let bytes_read = first.bytes_read;
3570        first.buffer_mut().advance(bytes_read);
3571        drop(first);
3572
3573        sender.write_all(b"second").await.expect("second payload should send");
3574        let mut second = timeout(Duration::from_secs(1), packets_rx.recv())
3575            .await
3576            .expect("reader should reacquire the released buffer")
3577            .expect("reader should remain active")
3578            .expect("second read should succeed");
3579        let bytes_read = second.bytes_read;
3580        assert_eq!(received_payload(second.buffer(), bytes_read), b"second");
3581        second.buffer_mut().advance(bytes_read);
3582        drop(second);
3583
3584        drop(packets_rx);
3585        sender
3586            .write_all(b"shutdown")
3587            .await
3588            .expect("shutdown payload should send");
3589        timeout(Duration::from_secs(1), reader)
3590            .await
3591            .expect("reader should stop after observing the closed queue")
3592            .expect("reader task should not panic");
3593        drop(shrinker);
3594    }
3595
3596    #[test]
3597    #[cfg(target_os = "linux")]
3598    fn autoscale_udp_listeners_yields_multiple_streams_on_linux() {
3599        let config = DogStatsDConfiguration {
3600            autoscale_udp_listeners: true,
3601            ..DogStatsDConfiguration::for_test()
3602        };
3603
3604        let streams = config
3605            .udp_streams_to_yield()
3606            .expect("autoscale yields at least 1 stream");
3607        let n = streams.get();
3608        assert!(
3609            (1..=4).contains(&n),
3610            "expected 1..=4 streams from vCPU formula, got {n}"
3611        );
3612    }
3613
3614    #[test]
3615    #[cfg(not(target_os = "linux"))]
3616    fn warns_for_uds_origin_detection_on_non_linux() {
3617        let config = DogStatsDConfiguration {
3618            port: 0,
3619            socket_path: Some("/tmp/dsd.sock".to_string()),
3620            origin_enrichment: OriginEnrichmentConfiguration {
3621                enabled: true,
3622                ..OriginEnrichmentConfiguration::for_test()
3623            },
3624            ..DogStatsDConfiguration::for_test()
3625        };
3626        let addresses = config.build_addresses(None);
3627
3628        assert!(config.uds_origin_detection_unsupported_on_platform(&addresses));
3629    }
3630
3631    #[test]
3632    #[cfg(not(target_os = "linux"))]
3633    fn does_not_warn_for_udp_origin_detection_on_non_linux() {
3634        let config = DogStatsDConfiguration {
3635            origin_enrichment: OriginEnrichmentConfiguration {
3636                enabled: true,
3637                ..OriginEnrichmentConfiguration::for_test()
3638            },
3639            ..DogStatsDConfiguration::for_test()
3640        };
3641        let addresses = config.build_addresses(None);
3642
3643        assert!(!config.uds_origin_detection_unsupported_on_platform(&addresses));
3644    }
3645
3646    #[test]
3647    #[cfg(not(target_os = "linux"))]
3648    fn autoscale_udp_listeners_yields_a_single_stream_on_non_linux() {
3649        let config = DogStatsDConfiguration {
3650            autoscale_udp_listeners: true,
3651            ..DogStatsDConfiguration::for_test()
3652        };
3653
3654        assert_eq!(None, config.udp_streams_to_yield());
3655    }
3656
3657    #[test]
3658    fn no_eol_required_listener_types_requires_no_newline() {
3659        let config = DogStatsDConfiguration {
3660            eol_required: Vec::new(),
3661            ..DogStatsDConfiguration::for_test()
3662        };
3663        let eol_required = config.eol_required();
3664
3665        assert!(!eol_required.for_listener(&udp_listen_address()));
3666        assert!(!eol_required.for_listener(&tcp_listen_address()));
3667    }
3668
3669    #[test]
3670    fn eol_required_matches_configured_listener_types() {
3671        let config = DogStatsDConfiguration {
3672            eol_required: vec!["udp".to_string(), "uds".to_string()],
3673            ..DogStatsDConfiguration::for_test()
3674        };
3675        let eol_required = config.eol_required();
3676
3677        assert!(eol_required.for_listener(&udp_listen_address()));
3678        assert!(!eol_required.for_listener(&tcp_listen_address()));
3679
3680        #[cfg(unix)]
3681        {
3682            assert!(eol_required.for_listener(&ListenAddress::Unixgram("/tmp/dsd.sock".into())));
3683            assert!(eol_required.for_listener(&ListenAddress::Unix("/tmp/dsd-stream.sock".into())));
3684        }
3685    }
3686
3687    #[test]
3688    fn eol_required_ignores_unrecognized_listener_types() {
3689        let config = DogStatsDConfiguration {
3690            eol_required: vec!["udp".to_string(), "carrier_pigeon".to_string()],
3691            ..DogStatsDConfiguration::for_test()
3692        };
3693        let eol_required = config.eol_required();
3694
3695        assert!(eol_required.for_listener(&udp_listen_address()));
3696        assert!(!eol_required.for_listener(&named_pipe_listen_address()));
3697    }
3698
3699    #[test]
3700    fn drops_full_named_pipe_buffer_without_newline() {
3701        let named_pipe_stream = named_pipe_listen_address();
3702        let mut buffer = get_pooled_object_via_builder::<_, BytesBuffer>(|| FixedSizeVec::with_capacity(8));
3703        buffer.put_slice(b"12345678");
3704
3705        assert!(super::should_drop_oversized_named_pipe_frame(
3706            &named_pipe_stream,
3707            &buffer
3708        ));
3709    }
3710
3711    #[test]
3712    fn keeps_named_pipe_partial_frame_when_buffer_has_capacity() {
3713        let named_pipe_stream = named_pipe_listen_address();
3714        let mut buffer = get_pooled_object_via_builder::<_, BytesBuffer>(|| FixedSizeVec::with_capacity(9));
3715        buffer.put_slice(b"12345678");
3716
3717        assert!(!super::should_drop_oversized_named_pipe_frame(
3718            &named_pipe_stream,
3719            &buffer
3720        ));
3721    }
3722
3723    #[test]
3724    fn keeps_full_named_pipe_buffer_with_newline() {
3725        let named_pipe_stream = named_pipe_listen_address();
3726        let mut buffer = get_pooled_object_via_builder::<_, BytesBuffer>(|| FixedSizeVec::with_capacity(8));
3727        buffer.put_slice(b"1234567\n");
3728
3729        assert!(!super::should_drop_oversized_named_pipe_frame(
3730            &named_pipe_stream,
3731            &buffer
3732        ));
3733    }
3734
3735    #[test]
3736    fn stream_log_too_big_warns_for_enabled_length_delimited_stream_invalid_frames() {
3737        let uds_stream = ListenAddress::Unix("/tmp/dsd-stream.sock".into());
3738        let named_pipe_stream = named_pipe_listen_address();
3739        let tcp_stream = ListenAddress::Tcp(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 8125)));
3740        let error = saluki_io::deser::framing::FramingError::InvalidFrame {
3741            frame_len: 8193,
3742            reason: "frame length exceeds buffer capacity",
3743        };
3744
3745        assert!(super::should_warn_stream_log_too_big(&uds_stream, &error, true));
3746        assert!(!super::should_warn_stream_log_too_big(&uds_stream, &error, false));
3747        assert!(!super::should_warn_stream_log_too_big(&named_pipe_stream, &error, true));
3748        assert!(!super::should_warn_stream_log_too_big(&tcp_stream, &error, true));
3749    }
3750
3751    #[test]
3752    fn interner_size_from_entry_count() {
3753        // A Core Agent migration config with entry count 4096 should yield 2 MiB, not 4096 bytes.
3754        let config = DogStatsDConfiguration {
3755            context_string_interner_entry_count: 4096,
3756            ..DogStatsDConfiguration::for_test()
3757        };
3758        assert_eq!(config.effective_context_string_interner_bytes(), ByteSize::mib(2));
3759    }
3760
3761    #[test]
3762    fn interner_size_from_explicit_bytes() {
3763        let config = DogStatsDConfiguration {
3764            context_string_interner_size_bytes: Some(ByteSize::b(4194304)),
3765            ..DogStatsDConfiguration::for_test()
3766        };
3767        assert_eq!(config.effective_context_string_interner_bytes(), ByteSize::b(4194304));
3768    }
3769
3770    #[test]
3771    fn interner_size_explicit_bytes_takes_priority() {
3772        let config = DogStatsDConfiguration {
3773            context_string_interner_entry_count: 4096,
3774            context_string_interner_size_bytes: Some(ByteSize::b(8388608)),
3775            ..DogStatsDConfiguration::for_test()
3776        };
3777        // The explicit byte size (8 MiB) takes priority over the entry count.
3778        assert_eq!(config.effective_context_string_interner_bytes(), ByteSize::b(8388608));
3779    }
3780
3781    #[test]
3782    fn interner_size_custom_entry_count() {
3783        let config = DogStatsDConfiguration {
3784            context_string_interner_entry_count: 8192,
3785            ..DogStatsDConfiguration::for_test()
3786        };
3787        // 8192 entries * 512 bytes = 4 MiB
3788        assert_eq!(config.effective_context_string_interner_bytes(), ByteSize::mib(4));
3789    }
3790
3791    /// Asserts that two lists of ListenAddress are equivalent.
3792    fn address_list_eq(expected: &mut [ListenAddress], actual: &mut [ListenAddress]) -> Result<(), String> {
3793        if expected.len() != actual.len() {
3794            return Err(format!(
3795                "length mismatch: expected {} addresses, got {}",
3796                expected.len(),
3797                actual.len()
3798            ));
3799        }
3800
3801        expected.sort_by_key(|a| a.to_string());
3802        actual.sort_by_key(|a| a.to_string());
3803
3804        for (e, a) in expected.iter().zip(actual.iter()) {
3805            let (es, as_) = (e.to_string(), a.to_string());
3806            if es != as_ {
3807                return Err(format!("address mismatch: expected {}, got {}", es, as_));
3808            }
3809        }
3810
3811        Ok(())
3812    }
3813
3814    /// This test verifies that we didn't accidentally break the `build_addresses_no_listeners` helper function which
3815    /// would render all further tests useless.
3816    #[test]
3817    fn build_addresses_assertion_function_works() {
3818        let config = DogStatsDConfiguration {
3819            port: 0,
3820            tcp_port: 123,
3821            socket_path: None,
3822            socket_stream_path: None,
3823            non_local_traffic: false,
3824            ..DogStatsDConfiguration::for_test()
3825        };
3826        let mut expected = vec![ListenAddress::Tcp(SocketAddr::V4(SocketAddrV4::new(
3827            // Close, but not quite! This is intentionally *not* 127.0.0.1 to test that the assertion will fail
3828            Ipv4Addr::new(127, 0, 0, 2),
3829            123,
3830        )))];
3831        let mut actual = config.build_addresses(None);
3832        assert!(address_list_eq(&mut expected, &mut actual).is_err())
3833    }
3834
3835    /// With all four listener gates off, `build_addresses` returns an empty Vec.
3836    #[test]
3837    fn build_addresses_no_listeners() {
3838        let config = DogStatsDConfiguration {
3839            port: 0,
3840            tcp_port: 0,
3841            socket_path: None,
3842            socket_stream_path: None,
3843            non_local_traffic: false,
3844            ..DogStatsDConfiguration::for_test()
3845        };
3846        let mut expected = vec![];
3847        let mut actual = config.build_addresses(None);
3848        address_list_eq(&mut expected, &mut actual).unwrap();
3849    }
3850
3851    /// UDP port set, `non_local_traffic=false` -> UDP listener bound to `127.0.0.1`.
3852    #[test]
3853    fn build_addresses_udp_local_only() {
3854        let config = DogStatsDConfiguration {
3855            port: 8125,
3856            tcp_port: 0,
3857            socket_path: None,
3858            socket_stream_path: None,
3859            non_local_traffic: false,
3860            ..DogStatsDConfiguration::for_test()
3861        };
3862        let mut expected = vec![ListenAddress::udp_loopback(8125)];
3863        let mut actual = config.build_addresses(None);
3864        address_list_eq(&mut expected, &mut actual).unwrap();
3865    }
3866
3867    /// UDP port set, `non_local_traffic=true` -> UDP listener bound to `0.0.0.0`.
3868    #[test]
3869    fn build_addresses_udp_non_local_only() {
3870        let config = DogStatsDConfiguration {
3871            port: 8125,
3872            tcp_port: 0,
3873            socket_path: None,
3874            socket_stream_path: None,
3875            non_local_traffic: true,
3876            ..DogStatsDConfiguration::for_test()
3877        };
3878        let mut expected = vec![ListenAddress::udp_any(8125)];
3879        let mut actual = config.build_addresses(None);
3880        address_list_eq(&mut expected, &mut actual).unwrap();
3881    }
3882
3883    /// TCP port set, `non_local_traffic=false` -> TCP listener bound to `127.0.0.1`.
3884    #[test]
3885    fn build_addresses_tcp_local_only() {
3886        let config = DogStatsDConfiguration {
3887            port: 0,
3888            tcp_port: 9000,
3889            socket_path: None,
3890            socket_stream_path: None,
3891            non_local_traffic: false,
3892            ..DogStatsDConfiguration::for_test()
3893        };
3894        let mut expected = vec![ListenAddress::Tcp(SocketAddr::V4(SocketAddrV4::new(
3895            Ipv4Addr::new(127, 0, 0, 1),
3896            9000,
3897        )))];
3898        let mut actual = config.build_addresses(None);
3899        address_list_eq(&mut expected, &mut actual).unwrap();
3900    }
3901
3902    /// TCP port set, `non_local_traffic=true` -> TCP listener bound to `0.0.0.0`.
3903    #[test]
3904    fn build_addresses_tcp_non_local_only() {
3905        let config = DogStatsDConfiguration {
3906            port: 0,
3907            tcp_port: 9000,
3908            socket_path: None,
3909            socket_stream_path: None,
3910            non_local_traffic: true,
3911            ..DogStatsDConfiguration::for_test()
3912        };
3913        let mut expected = vec![ListenAddress::Tcp(SocketAddr::V4(SocketAddrV4::new(
3914            Ipv4Addr::new(0, 0, 0, 0),
3915            9000,
3916        )))];
3917        let mut actual = config.build_addresses(None);
3918        address_list_eq(&mut expected, &mut actual).unwrap();
3919    }
3920
3921    /// `socket_path` set -> a `Unixgram` address is produced with that path.
3922    #[test]
3923    fn build_addresses_unixgram_only() {
3924        let config = DogStatsDConfiguration {
3925            port: 0,
3926            tcp_port: 0,
3927            socket_path: Some("/tmp/dsd.sock".to_string()),
3928            socket_stream_path: None,
3929            non_local_traffic: false,
3930            ..DogStatsDConfiguration::for_test()
3931        };
3932        let mut expected = vec![ListenAddress::Unixgram("/tmp/dsd.sock".into())];
3933        let mut actual = config.build_addresses(None);
3934        address_list_eq(&mut expected, &mut actual).unwrap();
3935    }
3936
3937    /// `socket_stream_path` set -> a `Unix` (stream) address is produced with that path.
3938    #[test]
3939    fn build_addresses_unix_stream_only() {
3940        let config = DogStatsDConfiguration {
3941            port: 0,
3942            tcp_port: 0,
3943            socket_path: None,
3944            socket_stream_path: Some("/tmp/dsd-stream.sock".to_string()),
3945            non_local_traffic: false,
3946            ..DogStatsDConfiguration::for_test()
3947        };
3948        let mut expected = vec![ListenAddress::Unix("/tmp/dsd-stream.sock".into())];
3949        let mut actual = config.build_addresses(None);
3950        address_list_eq(&mut expected, &mut actual).unwrap();
3951    }
3952
3953    /// All four listener types enabled at once, with `non_local_traffic=true`.
3954    #[test]
3955    fn build_addresses_all_four_non_local() {
3956        let config = DogStatsDConfiguration {
3957            port: 8125,
3958            tcp_port: 9000,
3959            socket_path: Some("/tmp/dsd.sock".to_string()),
3960            socket_stream_path: Some("/tmp/dsd-stream.sock".to_string()),
3961            non_local_traffic: true,
3962            ..DogStatsDConfiguration::for_test()
3963        };
3964        let mut expected = vec![
3965            ListenAddress::udp_any(8125),
3966            ListenAddress::tcp_any(9000),
3967            ListenAddress::Unixgram("/tmp/dsd.sock".into()),
3968            ListenAddress::Unix("/tmp/dsd-stream.sock".into()),
3969        ];
3970        let mut actual = config.build_addresses(None);
3971        address_list_eq(&mut expected, &mut actual).unwrap();
3972    }
3973
3974    /// All four listener types enabled at once, with `non_local_traffic=false`.
3975    #[test]
3976    fn build_addresses_all_four_local() {
3977        let config = DogStatsDConfiguration {
3978            port: 8125,
3979            tcp_port: 9000,
3980            socket_path: Some("/tmp/dsd.sock".to_string()),
3981            socket_stream_path: Some("/tmp/dsd-stream.sock".to_string()),
3982            non_local_traffic: false,
3983            ..DogStatsDConfiguration::for_test()
3984        };
3985        let mut expected = vec![
3986            ListenAddress::udp_loopback(8125),
3987            ListenAddress::tcp_loopback(9000),
3988            ListenAddress::Unixgram("/tmp/dsd.sock".into()),
3989            ListenAddress::Unix("/tmp/dsd-stream.sock".into()),
3990        ];
3991        let mut actual = config.build_addresses(None);
3992        address_list_eq(&mut expected, &mut actual).unwrap();
3993    }
3994
3995    /// Passing `Some(ip)` to `build_addresses` with `non_local_traffic=false` -> both UDP and TCP
3996    /// bind to that IP. Includes a UDS datagram socket to confirm `bind_host` doesn't affect it.
3997    #[test]
3998    fn build_addresses_bind_host_applies_to_udp_and_tcp() {
3999        let config = DogStatsDConfiguration {
4000            port: 8125,
4001            tcp_port: 9000,
4002            socket_path: Some("/tmp/dsd.sock".to_string()),
4003            socket_stream_path: None,
4004            non_local_traffic: false,
4005            ..DogStatsDConfiguration::for_test()
4006        };
4007        let bind_host = Some(IpAddr::V4(Ipv4Addr::new(192, 168, 1, 50)));
4008        let mut expected = vec![
4009            ListenAddress::Udp(([192, 168, 1, 50], 8125).into()),
4010            ListenAddress::Tcp(([192, 168, 1, 50], 9000).into()),
4011            ListenAddress::Unixgram("/tmp/dsd.sock".into()),
4012        ];
4013        let mut actual = config.build_addresses(bind_host);
4014        address_list_eq(&mut expected, &mut actual).unwrap();
4015    }
4016
4017    /// Passing `Some(ip)` to `build_addresses` with `non_local_traffic=true` -> both UDP and TCP
4018    /// bind to `0.0.0.0`; the `bind_host` parameter is ignored (precedence matches the Agent).
4019    /// Includes a UDS stream socket to confirm `bind_host` doesn't affect it.
4020    #[test]
4021    fn build_addresses_non_local_clobbers_bind_host() {
4022        let config = DogStatsDConfiguration {
4023            port: 8125,
4024            tcp_port: 9000,
4025            socket_path: None,
4026            socket_stream_path: Some("/tmp/dsd-stream.sock".to_string()),
4027            non_local_traffic: true,
4028            ..DogStatsDConfiguration::for_test()
4029        };
4030        let bind_host = Some(IpAddr::V4(Ipv4Addr::new(192, 168, 1, 50)));
4031        let mut expected = vec![
4032            ListenAddress::udp_any(8125),
4033            ListenAddress::tcp_any(9000),
4034            ListenAddress::Unix("/tmp/dsd-stream.sock".into()),
4035        ];
4036        let mut actual = config.build_addresses(bind_host);
4037        address_list_eq(&mut expected, &mut actual).unwrap();
4038    }
4039
4040    #[test]
4041    fn non_finite_metric_values_are_silently_dropped() {
4042        // The Datadog Agent sends NaN gauges (for example, encode_ms.avg computed as 0.0/0.0 in Go).
4043        // FloatIter skips non-finite values with a debug log, so decode_packet returns Ok with
4044        // num_points == 0. handle_frame then returns Ok(None) for zero-point packets, which is
4045        // the existing silent-drop path (no warning emitted).
4046        let codec = DogStatsDCodec::from_configuration(DogStatsDCodecConfiguration::default());
4047        for input in &[b"my.gauge:NaN|g" as &[u8], b"my.gauge:inf|g", b"my.gauge:-inf|g"] {
4048            match codec.decode_packet(input).expect("should decode without error") {
4049                ParsedPacket::Metric(packet) => assert_eq!(
4050                    packet.num_points, 0,
4051                    "non-finite value should be dropped, leaving 0 valid points"
4052                ),
4053                _ => panic!("expected Metric packet"),
4054            }
4055        }
4056    }
4057
4058    #[test]
4059    fn resolve_process_origin_pins_live_entity() {
4060        let capture_entity_resolver = CaptureTestEntityResolver::with_pid_mapping(
4061            42,
4062            EntityId::from_local_data("ci-pid-container").expect("container entity"),
4063        );
4064        let peer_addr = ConnectionAddress::ProcessLike(ProcessIdentity::Credentials(ProcessCredentials {
4065            pid: 42,
4066            uid: 0,
4067            gid: 0,
4068        }));
4069
4070        assert_eq!(
4071            resolve_process_origin(Some(&capture_entity_resolver), &peer_addr),
4072            Some(ProcessOrigin::Pinned(Some(
4073                EntityId::from_local_data("ci-pid-container").expect("container entity")
4074            )))
4075        );
4076    }
4077
4078    #[test]
4079    fn resolve_process_origin_skips_live_lookup_when_unused() {
4080        let capture_entity_resolver = CaptureTestEntityResolver::with_pid_mapping(
4081            42,
4082            EntityId::from_local_data("ci-pid-container").expect("container entity"),
4083        );
4084        let peer_addr = ConnectionAddress::ProcessLike(ProcessIdentity::Credentials(ProcessCredentials {
4085            pid: 42,
4086            uid: 0,
4087            gid: 0,
4088        }));
4089        let traffic_capture = TrafficCapture::new(PathBuf::new(), 1);
4090
4091        assert_eq!(
4092            resolve_process_origin_if_needed(false, &traffic_capture, Some(&capture_entity_resolver), &peer_addr),
4093            None
4094        );
4095        assert_eq!(capture_entity_resolver.resolution_count(), 0);
4096    }
4097
4098    #[tokio::test]
4099    async fn resolve_process_origin_pins_live_entity_during_capture() {
4100        let capture_entity_resolver = CaptureTestEntityResolver::with_pid_mapping(
4101            42,
4102            EntityId::from_local_data("ci-pid-container").expect("container entity"),
4103        );
4104        let peer_addr = ConnectionAddress::ProcessLike(ProcessIdentity::Credentials(ProcessCredentials {
4105            pid: 42,
4106            uid: 0,
4107            gid: 0,
4108        }));
4109        let capture_dir = tempfile::tempdir().expect("capture directory should be created");
4110        let traffic_capture = TrafficCapture::new(capture_dir.path().to_path_buf(), 1);
4111        traffic_capture
4112            .start_capture(None, Duration::from_secs(30), false)
4113            .expect("capture should start");
4114
4115        assert_eq!(
4116            resolve_process_origin_if_needed(false, &traffic_capture, Some(&capture_entity_resolver), &peer_addr),
4117            Some(ProcessOrigin::Pinned(Some(
4118                EntityId::from_local_data("ci-pid-container").expect("container entity")
4119            )))
4120        );
4121        assert_eq!(capture_entity_resolver.resolution_count(), 1);
4122
4123        traffic_capture.stop_capture();
4124        timeout(Duration::from_secs(1), async {
4125            while traffic_capture.is_ongoing() {
4126                yield_now().await;
4127            }
4128        })
4129        .await
4130        .expect("capture should stop");
4131    }
4132
4133    #[test]
4134    fn build_capture_record_ignores_payload_local_data() {
4135        let record = super::build_capture_record(None, None, b"test.metric:1|c|c:ci-local-container\n");
4136
4137        assert_eq!(record.container_id, None);
4138        assert!(record.ancillary.is_empty());
4139    }
4140
4141    #[test]
4142    fn stream_capture_state_preserves_last_pid_without_new_creds() {
4143        let mut stream_capture = super::StreamCaptureState::new();
4144
4145        stream_capture.update_peer_metadata(&ConnectionAddress::ProcessLike(ProcessIdentity::Credentials(
4146            ProcessCredentials {
4147                pid: 42,
4148                uid: 0,
4149                gid: 0,
4150            },
4151        )));
4152        stream_capture.update_peer_metadata(&ConnectionAddress::ProcessLike(ProcessIdentity::Unavailable));
4153
4154        assert_eq!(stream_capture.last_pid, Some(42));
4155    }
4156
4157    #[test]
4158    fn resolve_process_origin_preserves_live_pid_without_entity_resolver() {
4159        let peer_addr = ConnectionAddress::ProcessLike(ProcessIdentity::Credentials(ProcessCredentials {
4160            pid: 12345,
4161            uid: 1000,
4162            gid: 1000,
4163        }));
4164
4165        assert_eq!(
4166            resolve_process_origin(None, &peer_addr),
4167            Some(ProcessOrigin::Unpinned(12345))
4168        );
4169    }
4170
4171    #[test]
4172    fn resolve_process_origin_unpacks_captured_pid_when_replay_gid_present() {
4173        let captured_pid: u32 = 99887766;
4174        let peer_addr = ConnectionAddress::ProcessLike(ProcessIdentity::Credentials(ProcessCredentials {
4175            pid: 12345,        // our PID (irrelevant for replay)
4176            uid: captured_pid, // captured PID packed by the sender
4177            gid: super::REPLAY_CREDENTIALS_GID,
4178        }));
4179
4180        assert_eq!(
4181            resolve_process_origin(None, &peer_addr),
4182            Some(ProcessOrigin::Replay(captured_pid))
4183        );
4184    }
4185}
4186
4187/// Tests covering the source's background work running as supervised children.
4188///
4189/// The `tests` module above exercises the decode path in isolation. What matters here is the wiring: that `run` puts
4190/// the pool shrinker, datagram decoders, listeners, and per-connection handlers under the component's supervisor, and
4191/// that the whole subtree comes down cleanly and without dropping queued work.
4192#[cfg(test)]
4193mod supervision {
4194    use std::net::SocketAddr;
4195    use std::num::NonZeroUsize;
4196    use std::path::PathBuf;
4197    use std::sync::Arc;
4198    use std::time::Duration;
4199
4200    use bytes::BytesMut;
4201    use saluki_common::sync::shutdown::{ShutdownCoordinator, ShutdownHandle};
4202    use saluki_core::accounting::{ComponentRegistry, MemoryLimiter};
4203    use saluki_core::components::test_util::TestComponentSupervisor;
4204    #[cfg(unix)]
4205    use saluki_core::components::{sources::SourceBuilder as _, BuildContext};
4206    use saluki_core::{
4207        components::{
4208            sources::{Source as _, SourceContext},
4209            ComponentContext,
4210        },
4211        data_model::event::metric::context::{ContextResolverBuilder, TagsResolverBuilder},
4212        health::HealthRegistry,
4213        runtime::{
4214            state::{DataspaceRegistry, ResourceRegistry},
4215            SupervisorError,
4216        },
4217        support::SubsystemIdentifier,
4218        topology::{EventsBuffer, EventsDispatcher, OutputName, TopologyContext},
4219    };
4220    use saluki_io::net::listener::Listener;
4221    use saluki_io::net::{BoundListenAddress, ListenAddress};
4222    use stringtheory::MetaString;
4223    use tokio::net::UdpSocket;
4224    use tokio::runtime::Handle;
4225    use tokio::sync::mpsc;
4226    use tokio::time::timeout;
4227
4228    use super::*;
4229
4230    /// Bound on driving the source to completion, so a hang fails rather than stalling the suite.
4231    const RUN_TIMEOUT: Duration = Duration::from_secs(10);
4232
4233    /// Decoder workers the test source runs with.
4234    ///
4235    /// More than one, so the drain covers the real case of several decoders sharing the queue.
4236    const DECODER_WORKERS: usize = 2;
4237
4238    /// How long to let datagrams make their way through to the (single-slot) output before asserting on backpressure.
4239    const BACKPRESSURE_SETTLE: Duration = Duration::from_millis(250);
4240
4241    /// Children a running source with one UDP listener should have.
4242    ///
4243    /// The pool shrinker, one child per decoder worker, the listener itself, and one stream handler: a connectionless
4244    /// listener's `accept` yields its socket straight away, so its handler exists from the start rather than only once
4245    /// a peer shows up.
4246    const fn udp_child_count() -> usize {
4247        1 + DECODER_WORKERS + 1 + 1
4248    }
4249
4250    /// Everything a test needs to drive `DogStatsD::run` and observe the result.
4251    struct Harness {
4252        source: Box<DogStatsD>,
4253        context: SourceContext,
4254        /// Events dispatched to the `metrics` output.
4255        metrics_rx: mpsc::Receiver<EventsBuffer>,
4256        /// The address the source is listening on.
4257        listen_addr: SocketAddr,
4258    }
4259
4260    /// Builds a source with a single UDP listener bound to an ephemeral port.
4261    ///
4262    /// `health_registry` **MUST** outlive the run: `Health::live` resolves immediately once its registry is dropped,
4263    /// and the source's run loop polls liveness in a `select!`, so a dropped registry spins it instead of waiting for
4264    /// shutdown.
4265    /// Builds a source whose `metrics` output holds a single buffer.
4266    ///
4267    /// One slot makes backpressure trivial to arrange: one dispatched buffer fills it and the next blocks a decoder.
4268    async fn build_source(health_registry: &HealthRegistry, shutdown: ShutdownHandle) -> Harness {
4269        build_source_with_output_capacity(health_registry, shutdown, 1).await
4270    }
4271
4272    async fn build_source_with_output_capacity(
4273        health_registry: &HealthRegistry, shutdown: ShutdownHandle, output_capacity: usize,
4274    ) -> Harness {
4275        build_source_with(health_registry, shutdown, output_capacity, 16).await
4276    }
4277
4278    async fn build_source_with(
4279        health_registry: &HealthRegistry, shutdown: ShutdownHandle, output_capacity: usize, io_buffer_count: usize,
4280    ) -> Harness {
4281        let component_context = ComponentContext::test_source("dogstatsd");
4282
4283        let listener = leased_listener(ListenAddress::udp_loopback(0)).await;
4284        let listen_addr = match listener.bound_listen_address() {
4285            BoundListenAddress::Udp(addr) => addr,
4286            other => panic!("expected a bound UDP address, got {other:?}"),
4287        };
4288
4289        let (io_buffer_pool, io_buffer_pool_shrinker) = build_io_buffer_pool(
4290            io_buffer_count,
4291            io_buffer_count,
4292            DogStatsDConfiguration::for_test().buffer_size,
4293        );
4294
4295        let tags_resolver = TagsResolverBuilder::for_tests().build();
4296        let context_resolver = ContextResolverBuilder::for_tests()
4297            .with_tags_resolver(Some(tags_resolver.clone()))
4298            .build();
4299
4300        let source = Box::new(DogStatsD {
4301            listeners: vec![listener],
4302            decoder_worker_count: NonZeroUsize::new(DECODER_WORKERS).expect("decoder count should be non-zero"),
4303            io_buffer_pool,
4304            io_buffer_queue_capacity: 16,
4305            io_buffer_pool_shrinker: Box::pin(io_buffer_pool_shrinker),
4306            codec: DogStatsDCodec::from_configuration(DogStatsDCodecConfiguration::default()),
4307            context_resolvers: ContextResolvers::manual(context_resolver.clone(), context_resolver, tags_resolver),
4308            default_hostname: MetaString::from_static("test-host"),
4309            enabled_filter: EnablePayloadsFilter::default(),
4310            origin_detection_enabled: false,
4311            origin_telemetry_enabled: false,
4312            stream_log_too_big: false,
4313            disable_verbose_logs: false,
4314            eol_required: EolRequired::default(),
4315            additional_tags: Vec::<String>::new().into(),
4316            capture_entity_resolver: None,
4317            traffic_capture: TrafficCapture::new(PathBuf::new(), 1),
4318            packet_forwarder_target: None,
4319        });
4320
4321        // Only the `metrics` output is wired up: these tests feed counters, so nothing reaches the other two.
4322        let mut dispatcher = EventsDispatcher::new(component_context.clone());
4323        let metrics_output = OutputName::Given("metrics".into());
4324        dispatcher
4325            .add_output(metrics_output.clone())
4326            .expect("metrics output should be added");
4327        let (metrics_tx, metrics_rx) = mpsc::channel(output_capacity);
4328        dispatcher
4329            .attach_sender_to_output(&metrics_output, metrics_tx)
4330            .expect("metrics output should accept a sender");
4331
4332        let topology_context = TopologyContext::new(
4333            Arc::from("test"),
4334            MemoryLimiter::noop(),
4335            health_registry.clone(),
4336            Handle::current(),
4337            DataspaceRegistry::new(),
4338        );
4339        let health = health_registry
4340            .register_component(&SubsystemIdentifier::from_dotted("test.dogstatsd"))
4341            .expect("component was not previously registered");
4342
4343        let mut context = SourceContext::new(
4344            &topology_context,
4345            &component_context,
4346            ComponentRegistry::default(),
4347            health,
4348            dispatcher,
4349        );
4350
4351        context.set_shutdown_handle_for_test(shutdown);
4352
4353        Harness {
4354            source,
4355            context,
4356            metrics_rx,
4357            listen_addr,
4358        }
4359    }
4360
4361    /// Acquires a listener the way a real build does, so the source holds a lease rather than a bare listener.
4362    ///
4363    /// A registry per call is enough for a test that only needs one holder: a `ResourceLease` keeps its own clone of
4364    /// the registry, so the registry outlives the lease without the caller holding it.
4365    async fn leased_listener(address: ListenAddress) -> ResourceLease<Listener> {
4366        leased_listener_from(&ResourceRegistry::new(), address).await
4367    }
4368
4369    /// Acquires a listener from a specific registry, for a test that cares about what the registry does afterwards.
4370    async fn leased_listener_from(registry: &ResourceRegistry, address: ListenAddress) -> ResourceLease<Listener> {
4371        registry
4372            .acquire(
4373                &SubsystemIdentifier::from_segments(["test", "dogstatsd"]),
4374                SocketSpecification::new(address),
4375            )
4376            .await
4377            .expect("listener should bind")
4378    }
4379
4380    #[tokio::test]
4381    async fn background_work_runs_as_supervised_children() {
4382        // The pool shrinker, each datagram decoder, and each listener are supervised children rather than detached
4383        // tasks, so they are all accounted for while the source runs and all gone once it stops.
4384        let mut supervisor = TestComponentSupervisor::start("dogstatsd").await;
4385        let health_registry = HealthRegistry::new();
4386        let harness = build_source(&health_registry, supervisor.component_shutdown_handle()).await;
4387        let Harness { source, context, .. } = harness;
4388
4389        let run = tokio::spawn(supervisor.handle().scope(async move { source.run(context).await }));
4390
4391        supervisor.wait_for_children(udp_child_count()).await;
4392
4393        supervisor.signal_shutdown();
4394        timeout(RUN_TIMEOUT, run)
4395            .await
4396            .expect("source should stop on shutdown")
4397            .expect("source task should not panic")
4398            .expect("source should stop cleanly");
4399
4400        // The supervisor is draining by now -- its shutdown and the source's are one signal, as in production -- and a
4401        // drain freezes the child roster rather than emptying it, so there is no count left to sample here. What the
4402        // drain does still show is that nothing had to be forced: `ShutdownTimedOut` would mean a child ignored
4403        // shutdown and was aborted, which is the same "everything stopped on its own" property the count stood in for.
4404        // The brutally-stopped pool shrinker is the deliberate exception, and doesn't count as forced.
4405        let result = supervisor.wait().await;
4406        assert!(result.is_ok(), "every child should have stopped on its own: {result:?}");
4407    }
4408
4409    #[tokio::test]
4410    async fn run_does_not_return_until_the_decoders_have_drained() {
4411        // The drain guarantee that matters: work already queued when shutdown is signalled must still be decoded and
4412        // dispatched. The decoders deliberately ignore shutdown for exactly this reason, exiting only once the
4413        // listeners have dropped their senders and the queue is empty -- and `run` waits for that.
4414        //
4415        // Rather than race a datagram against shutdown, this pins a decoder open: the metrics output holds a single
4416        // buffer, so one dispatched buffer fills it and the next blocks a decoder mid-dispatch. `run` must not return
4417        // while that is true, no matter that shutdown has been signalled.
4418        let mut supervisor = TestComponentSupervisor::start("dogstatsd").await;
4419        let health_registry = HealthRegistry::new();
4420        let harness = build_source(&health_registry, supervisor.component_shutdown_handle()).await;
4421        let Harness {
4422            source,
4423            context,
4424            mut metrics_rx,
4425            listen_addr,
4426        } = harness;
4427
4428        let mut run = tokio::spawn(supervisor.handle().scope(async move { source.run(context).await }));
4429        supervisor.wait_for_children(udp_child_count()).await;
4430
4431        let client = UdpSocket::bind("127.0.0.1:0").await.expect("client should bind");
4432
4433        // Establish that the pipeline works, and leave the output empty again.
4434        client
4435            .send_to(b"first.metric:1|c", listen_addr)
4436            .await
4437            .expect("client should send");
4438        let buffer = timeout(RUN_TIMEOUT, metrics_rx.recv())
4439            .await
4440            .expect("the first datagram should be decoded and dispatched")
4441            .expect("the metrics output should stay open");
4442        assert_eq!(buffer.len(), 1);
4443
4444        // Now back the output up: the first of these fills the single slot, the second blocks a decoder.
4445        for payload in [b"second.metric:1|c".as_slice(), b"third.metric:1|c".as_slice()] {
4446            client.send_to(payload, listen_addr).await.expect("client should send");
4447        }
4448        tokio::time::sleep(BACKPRESSURE_SETTLE).await;
4449
4450        supervisor.signal_shutdown();
4451
4452        // `run` must still be waiting on the blocked decoder. Without the drain it would return here.
4453        tokio::time::sleep(BACKPRESSURE_SETTLE).await;
4454        assert!(
4455            !run.is_finished(),
4456            "`run` returned while a decoder was still mid-dispatch"
4457        );
4458
4459        // Draining the output unblocks the decoder, which lets the drain -- and so `run` -- complete.
4460        let mut drained = 1;
4461        let result = loop {
4462            select! {
4463                run_result = &mut run => break run_result,
4464                maybe_buffer = metrics_rx.recv() => match maybe_buffer {
4465                    Some(buffer) => drained += buffer.len(),
4466                    None => break (&mut run).await,
4467                },
4468            }
4469        };
4470        result
4471            .expect("source task should not panic")
4472            .expect("source should stop cleanly");
4473
4474        // `run` returning and the last buffer arriving are concurrent, so collect whatever is still queued. The
4475        // source's dispatcher is gone by now, so this terminates as soon as the channel is empty.
4476        while let Some(buffer) = metrics_rx.recv().await {
4477            drained += buffer.len();
4478        }
4479
4480        // Nothing accepted before shutdown should have been dropped on the way out.
4481        assert_eq!(drained, 3, "every queued datagram should have been dispatched");
4482
4483        let supervisor_result = supervisor.wait().await;
4484        assert!(
4485            supervisor_result.is_ok(),
4486            "every child should have stopped on its own: {supervisor_result:?}"
4487        );
4488    }
4489
4490    #[tokio::test]
4491    async fn supervisor_shutdown_stops_the_subtree_without_the_source_driving_it() {
4492        // The orderly path is the source's own `run` signalling its coordinators. This is the other one: the
4493        // supervisor tears the component down while `run` is still going, so nothing is driving those coordinators.
4494        //
4495        // The listener has to observe the supervisor's own signal for this to terminate. If it only watched the
4496        // source's coordinator it would keep accepting, keep holding a datagram sender, and so keep the decoders
4497        // waiting on a queue that never closes -- until the shutdown budget elapsed and aborted the lot.
4498        let supervisor = TestComponentSupervisor::start("dogstatsd").await;
4499        let health_registry = HealthRegistry::new();
4500
4501        // A signal of the test's own, deliberately never fired, so `run` stays in its loop for the whole test. This is
4502        // the one case production can't produce -- there a source's shutdown signal *is* its supervisor's -- and
4503        // separating them is what isolates the listener's own `process_shutdown` arm.
4504        let mut source_shutdown = ShutdownCoordinator::default();
4505        let harness = build_source(&health_registry, source_shutdown.register()).await;
4506        let Harness { source, context, .. } = harness;
4507
4508        let _run = tokio::spawn(supervisor.handle().scope(async move { source.run(context).await }));
4509        supervisor.wait_for_children(udp_child_count()).await;
4510
4511        let result = supervisor.shutdown().await;
4512        assert!(
4513            result.is_ok(),
4514            "the subtree should have stopped on the supervisor's signal alone: {result:?}"
4515        );
4516    }
4517
4518    #[tokio::test]
4519    async fn stream_handlers_are_supervised_children() {
4520        // Per-connection handlers are supervised too, under one fixed name per listener type. Accepting a connection
4521        // adds children, and tearing the source down takes them with everything else.
4522        let mut supervisor = TestComponentSupervisor::start("dogstatsd").await;
4523        let health_registry = HealthRegistry::new();
4524
4525        // A TCP listener, so there are real connections to accept.
4526        let tcp_listener = leased_listener(ListenAddress::Tcp("127.0.0.1:0".parse().unwrap())).await;
4527        let listen_addr = match tcp_listener.bound_listen_address() {
4528            BoundListenAddress::Tcp(addr) => addr,
4529            other => panic!("expected a bound TCP address, got {other:?}"),
4530        };
4531
4532        let mut harness = build_source(&health_registry, supervisor.component_shutdown_handle()).await;
4533        harness.source.listeners = vec![tcp_listener];
4534
4535        let Harness { source, context, .. } = harness;
4536
4537        let run = tokio::spawn(supervisor.handle().scope(async move { source.run(context).await }));
4538
4539        let baseline = 1 + DECODER_WORKERS + 1;
4540        supervisor.wait_for_children(baseline).await;
4541
4542        let _client = tokio::net::TcpStream::connect(listen_addr)
4543            .await
4544            .expect("client should connect");
4545
4546        // Two children per connection, not one: the handler, and the stream reader it starts. Spawning on the
4547        // ambient supervisor makes the reader the handler's *sibling* rather than its descendant, which is why it
4548        // shows up in the supervisor's own child count.
4549        supervisor.wait_for_children(baseline + 2).await;
4550
4551        supervisor.signal_shutdown();
4552        timeout(RUN_TIMEOUT, run)
4553            .await
4554            .expect("source should stop on shutdown")
4555            .expect("source task should not panic")
4556            .expect("source should stop cleanly");
4557
4558        // As above, the drain is what attests that the handler stopped on its own rather than being aborted.
4559        let result = supervisor.wait().await;
4560        assert!(result.is_ok(), "every child should have stopped on its own: {result:?}");
4561    }
4562
4563    #[tokio::test]
4564    async fn the_registry_keeps_the_listener_bound_after_the_source_stops() {
4565        // The property the whole migration is for. The source only ever borrows its listener, so stopping it hands the
4566        // socket back to the registry still bound, rather than releasing the port to the operating system and leaving
4567        // a rebuilt source to race whatever else wants it.
4568        let mut supervisor = TestComponentSupervisor::start("dogstatsd").await;
4569        let health_registry = HealthRegistry::new();
4570
4571        // A registry the test keeps, so it can ask for the same address once the source is gone.
4572        //
4573        // Port 0 is what makes the assertion below mean anything: the specification is keyed on the address as asked
4574        // for, so re-acquiring it can only come back on the same port if the registry really held the socket. A fixed
4575        // port could be re-bound from scratch and look identical.
4576        let registry = ResourceRegistry::new();
4577        let address = ListenAddress::udp_loopback(0);
4578
4579        let mut harness = build_source(&health_registry, supervisor.component_shutdown_handle()).await;
4580        harness.source.listeners = vec![leased_listener_from(&registry, address.clone()).await];
4581        let bound = harness.source.listeners[0].bound_listen_address();
4582        let listen_addr = match bound {
4583            BoundListenAddress::Udp(addr) => addr,
4584            other => panic!("expected a bound UDP address, got {other:?}"),
4585        };
4586
4587        let Harness { source, context, .. } = harness;
4588        let run = tokio::spawn(supervisor.handle().scope(async move { source.run(context).await }));
4589        supervisor.wait_for_children(udp_child_count()).await;
4590
4591        supervisor.signal_shutdown();
4592        timeout(RUN_TIMEOUT, run)
4593            .await
4594            .expect("source should stop on shutdown")
4595            .expect("source task should not panic")
4596            .expect("source should stop cleanly");
4597        supervisor.wait().await.expect("supervisor should drain cleanly");
4598
4599        // The same socket, not a lucky re-bind: the registry never let go of it.
4600        let mut listener = leased_listener_from(&registry, address).await;
4601        assert_eq!(listener.bound_listen_address(), bound);
4602
4603        // And it is still usable. A listener whose socket had been moved out by the previous holder rather than
4604        // borrowed would come back exhausted and hang here instead of yielding anything.
4605        let mut stream = timeout(RUN_TIMEOUT, listener.accept())
4606            .await
4607            .expect("accept should not hang once the source has released its lease")
4608            .expect("listener should yield its socket");
4609
4610        let client = UdpSocket::bind("127.0.0.1:0").await.expect("client should bind");
4611        client
4612            .send_to(b"after.teardown:1|c", listen_addr)
4613            .await
4614            .expect("client should send");
4615
4616        let mut buf = BytesMut::with_capacity(64);
4617        let (read, _) = timeout(RUN_TIMEOUT, stream.receive(&mut buf))
4618            .await
4619            .expect("receive should not time out")
4620            .expect("receive should succeed");
4621        assert_eq!(&buf[..read], b"after.teardown:1|c");
4622    }
4623
4624    #[cfg(unix)]
4625    #[tokio::test]
4626    async fn building_against_an_already_leased_address_says_so() {
4627        // Two holders of one address is a conflict the registry can name, rather than the bind failure an operator
4628        // would otherwise have to work backwards from.
4629        //
4630        // A socket path rather than a port: the configuration has to name the address concretely for the build to ask
4631        // for the same one, and a path in a temporary directory is unique to this test by construction, so there is no
4632        // ephemeral port to discover and race over.
4633        let temp_dir = tempfile::tempdir().expect("temp directory should be created");
4634        let socket_path = temp_dir.path().join("dogstatsd.socket");
4635
4636        let registry = ResourceRegistry::new();
4637        let _held = leased_listener_from(&registry, ListenAddress::Unixgram(socket_path.clone())).await;
4638
4639        let mut config = DogStatsDConfiguration::for_test();
4640        config.socket_path = Some(socket_path.to_string_lossy().into_owned());
4641
4642        // The fixture defaults to the real DogStatsD port, which this test has no business binding.
4643        config.port = 0;
4644
4645        // The build and the existing holder have to share a registry for the conflict to exist at all;
4646        // `BuildContext::test_source` would give the build one of its own.
4647        let context = BuildContext::new(ComponentContext::test_source("dogstatsd"), registry);
4648        let Err(error) = config.build(context).await else {
4649            panic!("building against an address that is already leased should be refused");
4650        };
4651        let error = error.to_string();
4652        assert!(error.contains("already leased"), "unexpected error: {error}");
4653    }
4654
4655    #[tokio::test]
4656    async fn a_listener_stopping_takes_the_component_down() {
4657        // The listener is a significant child, so it stopping is a component-level event rather than a source that
4658        // carries on with nothing left to listen on.
4659        //
4660        // A fatal accept error is what causes this in production, which a test can't readily provoke. Signalling only
4661        // the source's own shutdown produces the same shape: the listener exits while the supervisor around it is
4662        // still running normally, which is exactly the case significance exists to catch.
4663        let supervisor = TestComponentSupervisor::start("dogstatsd").await;
4664        let health_registry = HealthRegistry::new();
4665
4666        let mut source_shutdown = ShutdownCoordinator::default();
4667        let harness = build_source(&health_registry, source_shutdown.register()).await;
4668        let Harness { source, context, .. } = harness;
4669
4670        let run = tokio::spawn(supervisor.handle().scope(async move { source.run(context).await }));
4671        supervisor.wait_for_children(udp_child_count()).await;
4672
4673        source_shutdown.shutdown();
4674        timeout(RUN_TIMEOUT, run)
4675            .await
4676            .expect("source should stop on shutdown")
4677            .expect("source task should not panic")
4678            .expect("source should stop cleanly");
4679
4680        // The supervisor was never told to stop, so it stopping at all is the assertion. The timeout is what turns a
4681        // listener that was left non-significant into a failure rather than a hang.
4682        let result = timeout(RUN_TIMEOUT, supervisor.wait())
4683            .await
4684            .expect("the supervisor should have stopped on its own");
4685        assert!(
4686            matches!(result, Err(SupervisorError::SignificantChildExited)),
4687            "a stopped listener should have brought the component down, got {result:?}"
4688        );
4689    }
4690}