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