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