1use std::{
10 collections::VecDeque,
11 future::Future,
12 num::NonZeroUsize,
13 path::PathBuf,
14 pin::Pin,
15 sync::{Arc, LazyLock},
16 time::{Duration, SystemTime, UNIX_EPOCH},
17};
18
19use async_trait::async_trait;
20use bytes::{Buf, BufMut};
21use bytesize::ByteSize;
22use saluki_common::{
23 sync::shutdown::{ShutdownCoordinator, ShutdownHandle},
24 task::spawn_traced_named,
25};
26use saluki_config::{deserialize_space_separated_or_seq, GenericConfiguration};
27use saluki_context::{
28 origin::RawOrigin,
29 tags::{RawTags, RawTagsFilter},
30 TagsResolver,
31};
32use saluki_core::accounting::{MemoryBounds, MemoryBoundsBuilder, UsageExpr};
33use saluki_core::data_model::event::{
34 eventd::EventD,
35 metric::{Metric, MetricMetadata, MetricOrigin},
36 service_check::ServiceCheck,
37 Event, EventType,
38};
39use saluki_core::{
40 components::{sources::*, ComponentContext},
41 pooling::ElasticObjectPool,
42 topology::{interconnect::EventBufferManager, EventsBuffer, OutputDefinition},
43};
44use saluki_env::{workload::CaptureEntityResolver, WorkloadProvider};
45use saluki_error::{generic_error, ErrorContext as _, GenericError};
46use saluki_io::{
47 buf::{BytesBuffer, ClearableIoBuffer as _, FixedSizeVec},
48 deser::{
49 codec::dogstatsd::*,
50 framing::{Framer as _, FramingError, LengthDelimitedFramer},
51 },
52 net::{
53 listener::{Listener, ListenerError},
54 ConnectionAddress, ListenAddress, ProcessCredentials, ProcessIdentity, Stream,
55 },
56};
57use serde::{Deserialize, Deserializer};
58use serde_with::{serde_as, NoneAsEmptyString};
59use snafu::{ResultExt as _, Snafu};
60use stringtheory::MetaString;
61use tokio::{
62 pin, select,
63 time::{interval, MissedTickBehavior},
64};
65use tracing::{debug, error, info, trace, warn};
66
67mod forwarder;
68use self::forwarder::{PacketForwarder, PacketForwarderTarget};
69
70mod framer;
71use self::framer::{get_framer, DsdFramer};
72use crate::sources::dogstatsd::tags::{WellKnownTags, WellKnownTagsFilterPredicate};
73
74mod filters;
75use self::filters::EnablePayloadsFilter;
76
77mod io_buffer;
78use self::io_buffer::IoBufferManager;
79
80mod metrics;
81use self::metrics::{build_metrics, Metrics};
82
83mod replay;
84use self::replay::{CaptureRecord, CapturedTaggerHandle, TrafficCapture};
85pub use self::replay::{
86 DogStatsDCaptureAPIHandler, DogStatsDCaptureControl, DogStatsDReplayAPIHandler, DogStatsDReplayControl,
87 ReplaySession, TimestampResolution, TrafficCaptureReader, DEFAULT_REPLAY_LOOPS, REPLAY_CREDENTIALS_GID,
88};
89
90mod origin;
91use self::origin::{
92 mark_replay_process_id, origin_from_event_packet, origin_from_metric_packet, origin_from_service_check_packet,
93 DogStatsDOriginTagResolver, OriginEnrichmentConfiguration,
94};
95
96mod resolver;
97use self::resolver::ContextResolvers;
98
99mod tags;
100
101#[derive(Debug, Snafu)]
102#[snafu(context(suffix(false)))]
103enum Error {
104 #[snafu(display("Failed to create {} listener: {}", listener_type, source))]
105 FailedToCreateListener {
106 listener_type: &'static str,
107 source: ListenerError,
108 },
109
110 #[snafu(display("No listeners configured. Please specify a port (`dogstatsd_port`) or a socket path (`dogstatsd_socket` or `dogstatsd_stream_socket`) to enable a listener."))]
111 NoListenersConfigured,
112
113 #[snafu(display("Could not resolve bind_host '{}': {}", host, source))]
114 UnresolvableBindHost { host: String, source: std::io::Error },
115
116 #[snafu(display("bind_host '{}' resolved to zero IP addresses.", host))]
117 BindHostHasNoAddresses { host: String },
118}
119
120const INTERNER_BASELINE_BYTES_PER_ENTRY: u64 = 512;
125
126const fn default_buffer_size() -> usize {
127 8192
128}
129
130const fn default_buffer_count() -> usize {
131 128
132}
133
134const fn default_buffer_count_max() -> usize {
135 256
136}
137
138const fn default_port() -> u16 {
139 8125
140}
141
142const fn default_tcp_port() -> u16 {
143 0
144}
145
146const fn default_statsd_forward_port() -> u16 {
147 0
148}
149
150const fn default_socket_receive_buffer_size() -> usize {
151 0
152}
153
154const fn default_allow_context_heap_allocations() -> bool {
155 true
156}
157
158const fn default_no_aggregation_pipeline_support() -> bool {
159 true
160}
161
162const fn default_context_string_interner_entry_count() -> u64 {
163 4096
164}
165
166const fn default_cached_contexts_limit() -> usize {
167 500_000
168}
169
170const fn default_cached_tagsets_limit() -> usize {
171 500_000
172}
173
174const fn default_context_expiry_seconds() -> u64 {
175 20
176}
177
178const fn default_dogstatsd_permissive_decoding() -> bool {
179 true
180}
181
182const fn default_dogstatsd_minimum_sample_rate() -> f64 {
183 0.000000003845
184}
185
186const fn default_true() -> bool {
187 true
188}
189
190const fn default_windows_pipe_security_descriptor() -> &'static str {
192 "D:AI(A;;GA;;;WD)"
193}
194
195fn default_windows_pipe_security_descriptor_string() -> String {
196 default_windows_pipe_security_descriptor().to_string()
197}
198
199#[derive(Deserialize)]
201#[cfg_attr(test, derive(PartialEq, serde::Serialize))]
202pub struct EnablePayloadsConfiguration {
203 #[serde(default = "default_true")]
207 pub series: bool,
208
209 #[serde(default = "default_true")]
213 pub sketches: bool,
214
215 #[serde(default = "default_true")]
219 pub events: bool,
220
221 #[serde(default = "default_true")]
225 pub service_checks: bool,
226}
227
228impl Default for EnablePayloadsConfiguration {
229 fn default() -> Self {
230 Self {
231 series: true,
232 sketches: true,
233 events: true,
234 service_checks: true,
235 }
236 }
237}
238
239const MIN_CAPTURE_DEPTH: usize = 1024;
240
241const fn default_capture_depth() -> usize {
242 MIN_CAPTURE_DEPTH
243}
244
245const DOGSTATSD_CAPTURE_DIR: &str = "dsd_capture";
246
247fn deserialize_empty_metastring_as_none<'de, D>(deserializer: D) -> Result<Option<MetaString>, D::Error>
248where
249 D: Deserializer<'de>,
250{
251 let value = Option::<MetaString>::deserialize(deserializer)?;
252 Ok(value.filter(|host| !host.is_empty()))
253}
254
255#[derive(Deserialize, Default)]
256#[cfg_attr(test, derive(PartialEq, serde::Serialize))]
257struct DogStatsDTelemetryConfiguration {
258 #[serde(default)]
265 dogstatsd_origin: bool,
266}
267
268#[serde_as]
272#[derive(Deserialize, Default)]
273#[cfg_attr(test, derive(derive_where::DeriveWhere, serde::Serialize))]
274#[cfg_attr(test, derive_where(PartialEq))]
275pub struct DogStatsDConfiguration {
276 #[serde(skip)]
278 default_hostname: MetaString,
279
280 #[serde(rename = "dogstatsd_buffer_size", default = "default_buffer_size")]
286 buffer_size: usize,
287
288 #[serde(rename = "dogstatsd_buffer_count", default = "default_buffer_count")]
296 buffer_count: usize,
297
298 #[serde(rename = "dogstatsd_buffer_count_max", default = "default_buffer_count_max")]
310 buffer_count_max: usize,
311
312 #[serde(rename = "dogstatsd_port", default = "default_port")]
318 port: u16,
319
320 #[serde(rename = "dogstatsd_so_rcvbuf", default = "default_socket_receive_buffer_size")]
326 socket_receive_buffer_size: usize,
327
328 #[serde(rename = "dogstatsd_tcp_port", default = "default_tcp_port")]
334 tcp_port: u16,
335
336 #[serde(
343 rename = "statsd_forward_host",
344 default,
345 deserialize_with = "deserialize_empty_metastring_as_none"
346 )]
347 statsd_forward_host: Option<MetaString>,
348
349 #[serde(rename = "statsd_forward_port", default = "default_statsd_forward_port")]
355 statsd_forward_port: u16,
356
357 #[serde(rename = "dogstatsd_socket", default)]
363 #[serde_as(as = "NoneAsEmptyString")]
364 socket_path: Option<String>,
365
366 #[serde(rename = "dogstatsd_stream_socket", default)]
372 #[serde_as(as = "NoneAsEmptyString")]
373 socket_stream_path: Option<String>,
374
375 #[serde(rename = "dogstatsd_stream_log_too_big", default)]
384 stream_log_too_big: bool,
385
386 #[serde(rename = "dogstatsd_pipe_name", default)]
393 #[serde_as(as = "NoneAsEmptyString")]
394 pipe_name: Option<String>,
395
396 #[serde(
402 rename = "dogstatsd_windows_pipe_security_descriptor",
403 default = "default_windows_pipe_security_descriptor_string"
404 )]
405 windows_pipe_security_descriptor: String,
406
407 #[serde(rename = "dogstatsd_disable_verbose_logs", default)]
415 disable_verbose_logs: bool,
416
417 #[serde(
425 rename = "dogstatsd_eol_required",
426 default,
427 deserialize_with = "deserialize_space_separated_or_seq"
428 )]
429 eol_required: Vec<String>,
430
431 #[serde(rename = "bind_host", default)]
439 #[serde_as(as = "NoneAsEmptyString")]
440 bind_host: Option<String>,
441
442 #[serde(rename = "dogstatsd_non_local_traffic", default)]
449 non_local_traffic: bool,
450
451 #[serde(rename = "dogstatsd_autoscale_udp_listeners", default)]
466 autoscale_udp_listeners: bool,
467
468 #[serde(
478 rename = "dogstatsd_allow_context_heap_allocs",
479 default = "default_allow_context_heap_allocations"
480 )]
481 allow_context_heap_allocations: bool,
482
483 #[serde(
491 rename = "dogstatsd_no_aggregation_pipeline",
492 default = "default_no_aggregation_pipeline_support"
493 )]
494 no_aggregation_pipeline_support: bool,
495
496 #[serde(
504 rename = "dogstatsd_string_interner_size",
505 default = "default_context_string_interner_entry_count"
506 )]
507 context_string_interner_entry_count: u64,
508
509 #[serde(rename = "dogstatsd_string_interner_size_bytes", default)]
515 context_string_interner_size_bytes: Option<ByteSize>,
516
517 #[serde(
525 rename = "dogstatsd_cached_contexts_limit",
526 default = "default_cached_contexts_limit"
527 )]
528 cached_contexts_limit: usize,
529
530 #[serde(rename = "dogstatsd_cached_tagsets_limit", default = "default_cached_tagsets_limit")]
538 cached_tagsets_limit: usize,
539
540 #[serde(
546 rename = "dogstatsd_context_expiry_seconds",
547 default = "default_context_expiry_seconds"
548 )]
549 context_expiry_seconds: u64,
550
551 #[serde(
558 rename = "dogstatsd_permissive_decoding",
559 default = "default_dogstatsd_permissive_decoding"
560 )]
561 permissive_decoding: bool,
562
563 #[serde(
573 rename = "dogstatsd_minimum_sample_rate",
574 default = "default_dogstatsd_minimum_sample_rate"
575 )]
576 minimum_sample_rate: f64,
577
578 #[serde(rename = "enable_payloads", default)]
580 enable_payloads: EnablePayloadsConfiguration,
581
582 #[serde(flatten, default)]
584 origin_enrichment: OriginEnrichmentConfiguration,
585
586 #[serde(default)]
588 telemetry: DogStatsDTelemetryConfiguration,
589
590 #[serde(skip)]
592 #[cfg_attr(test, derive_where(skip))]
593 workload_provider: Option<Arc<dyn WorkloadProvider + Send + Sync>>,
594
595 #[serde(skip, default)]
597 #[cfg_attr(test, derive_where(skip))]
598 capture_entity_resolver: Option<Arc<dyn CaptureEntityResolver + Send + Sync>>,
599
600 #[serde(rename = "dogstatsd_tags", default)]
602 additional_tags: Vec<String>,
603
604 #[serde(rename = "dogstatsd_capture_path", default)]
612 capture_path: PathBuf,
613
614 #[serde(rename = "dogstatsd_capture_depth", default = "default_capture_depth")]
622 capture_depth: usize,
623
624 #[serde(skip, default)]
625 #[cfg_attr(test, derive_where(skip))]
626 capture_control: DogStatsDCaptureControl,
627
628 #[serde(skip, default)]
629 #[cfg_attr(test, derive_where(skip))]
630 replay_control: DogStatsDReplayControl,
631
632 #[serde(default)]
639 provider_kind: String,
640}
641
642#[derive(Clone, Copy, Default)]
643struct EolRequired {
644 udp: bool,
645 uds: bool,
646 named_pipe: bool,
647}
648
649impl EolRequired {
650 fn from_config_values(values: &[String]) -> Self {
651 let mut eol_required = Self::default();
652
653 for value in values {
654 match value.as_str() {
655 "udp" => eol_required.udp = true,
656 "uds" => eol_required.uds = true,
657 "named_pipe" => eol_required.named_pipe = true,
658 _ => warn!(
659 value,
660 "Invalid dogstatsd_eol_required value. Expected 'udp', 'uds', or 'named_pipe'."
661 ),
662 }
663 }
664
665 eol_required
666 }
667
668 fn for_listener(&self, listen_addr: &ListenAddress) -> bool {
669 match listen_addr {
670 ListenAddress::Udp(_) => self.udp,
671 ListenAddress::Tcp(_) => false,
672 ListenAddress::Unixgram(_) | ListenAddress::Unix(_) => self.uds,
673 ListenAddress::NamedPipe { .. } => self.named_pipe,
674 }
675 }
676}
677
678async fn resolve_bind_host(host: &str) -> Result<std::net::IpAddr, Error> {
684 let mut addrs = tokio::net::lookup_host((host, 0u16))
685 .await
686 .context(UnresolvableBindHost { host: host.to_string() })?;
687 addrs
688 .next()
689 .map(|sa| sa.ip())
690 .ok_or_else(|| Error::BindHostHasNoAddresses { host: host.to_string() })
691}
692
693impl DogStatsDConfiguration {
694 pub fn from_configuration(config: &GenericConfiguration) -> Result<Self, GenericError> {
696 let mut dogstatsd_config: Self = config.as_typed()?;
697 dogstatsd_config.fix_empty_capture_path(config);
698 dogstatsd_config.fix_capture_depth();
699 Ok(dogstatsd_config)
700 }
701
702 fn additional_tags(&self) -> Vec<String> {
704 if self.provider_kind.is_empty() {
705 return self.additional_tags.clone();
706 }
707
708 let mut tags = self.additional_tags.clone();
709 tags.push(format!("provider_kind:{}", self.provider_kind.clone()));
710 tags
711 }
712
713 fn fix_capture_depth(&mut self) {
714 self.capture_depth = self.capture_depth.max(MIN_CAPTURE_DEPTH);
715 }
716
717 fn effective_context_string_interner_bytes(&self) -> ByteSize {
723 match self.context_string_interner_size_bytes {
724 Some(explicit_bytes) => explicit_bytes,
725 None => {
726 saluki_antithesis::always_le!(
727 self.context_string_interner_entry_count,
728 u64::MAX / INTERNER_BASELINE_BYTES_PER_ENTRY,
729 "dogstatsd interner byte-size multiply does not overflow",
730 { "entry_count": self.context_string_interner_entry_count }
731 );
732 ByteSize::b(
733 self.context_string_interner_entry_count
734 .saturating_mul(INTERNER_BASELINE_BYTES_PER_ENTRY),
735 )
736 }
737 }
738 }
739
740 fn eol_required(&self) -> EolRequired {
741 EolRequired::from_config_values(&self.eol_required)
742 }
743
744 fn statsd_forward_target(&self) -> Option<(&MetaString, u16)> {
745 let host = self.statsd_forward_host.as_ref()?;
746 if self.statsd_forward_port == 0 {
747 return None;
748 }
749
750 Some((host, self.statsd_forward_port))
751 }
752
753 fn packet_forwarder_target(&self) -> Option<PacketForwarderTarget> {
754 let (host, port) = self.statsd_forward_target()?;
755 Some(PacketForwarderTarget::new(host.clone(), port))
756 }
757
758 fn udp_streams_to_yield(&self) -> Option<NonZeroUsize> {
764 if !self.autoscale_udp_listeners {
765 return None;
766 }
767
768 #[cfg(not(target_os = "linux"))]
769 if self.autoscale_udp_listeners {
770 warn!("UDP stream handler autoscaling not supported on non-Linux platforms. Default to single stream handler.");
771 return None;
772 }
773
774 let vcpus = std::thread::available_parallelism().map(NonZeroUsize::get).unwrap_or(1);
775 let streams = (1 + vcpus / 8).min(4);
776 NonZeroUsize::new(streams)
777 }
778
779 fn effective_max_buffer_count(&self) -> usize {
785 self.buffer_count_max.max(self.buffer_count)
786 }
787
788 pub fn with_default_hostname(mut self, hostname: impl Into<MetaString>) -> Self {
790 self.default_hostname = hostname.into();
791 self
792 }
793
794 pub fn with_workload_provider<W>(mut self, workload_provider: W) -> Self
800 where
801 W: WorkloadProvider + Send + Sync + 'static,
802 {
803 self.workload_provider = Some(Arc::new(workload_provider));
804 self
805 }
806
807 pub fn with_capture_entity_resolver<R>(mut self, capture_entity_resolver: R) -> Self
814 where
815 R: CaptureEntityResolver + Send + Sync + 'static,
816 {
817 self.capture_entity_resolver = Some(Arc::new(capture_entity_resolver));
818 self
819 }
820
821 pub fn capture_control(&self) -> DogStatsDCaptureControl {
823 self.capture_control.clone()
824 }
825
826 pub fn capture_api_handler(&self) -> DogStatsDCaptureAPIHandler {
828 DogStatsDCaptureAPIHandler::new(self.capture_control.clone())
829 }
830
831 pub fn replay_control(&self) -> DogStatsDReplayControl {
833 self.replay_control.clone()
834 }
835
836 pub fn replay_api_handler(&self) -> DogStatsDReplayAPIHandler {
838 DogStatsDReplayAPIHandler::new(self.replay_control.clone())
839 }
840
841 fn fix_empty_capture_path(&mut self, config: &GenericConfiguration) {
842 if self.capture_path.parent().is_some() {
843 return;
844 }
845
846 let capture_path = match config.try_get_typed::<PathBuf>("run_path") {
847 Ok(Some(mut run_path)) => {
848 run_path.push(DOGSTATSD_CAPTURE_DIR);
849 run_path
850 }
851 Ok(None) => {
852 debug!(
853 "`dogstatsd_capture_path` and `run_path` were empty. Default DogStatsD capture path is unavailable."
854 );
855 return;
856 }
857 Err(e) => {
858 debug!(
859 error = %e,
860 "Failed to read `run_path` from configuration. Default DogStatsD capture path is unavailable."
861 );
862 return;
863 }
864 };
865
866 self.capture_path = capture_path;
867 }
868
869 fn build_addresses(&self, bind_host: Option<std::net::IpAddr>) -> Vec<ListenAddress> {
879 let bind_ip: std::net::IpAddr = if self.non_local_traffic {
880 [0, 0, 0, 0].into()
881 } else {
882 bind_host.unwrap_or_else(|| [127, 0, 0, 1].into())
883 };
884
885 let mut addresses: Vec<ListenAddress> = Vec::new();
886
887 if self.port != 0 {
888 addresses.push(ListenAddress::Udp(std::net::SocketAddr::new(bind_ip, self.port)));
889 }
890
891 if self.tcp_port != 0 {
892 addresses.push(ListenAddress::Tcp(std::net::SocketAddr::new(bind_ip, self.tcp_port)));
893 }
894
895 if let Some(socket_path) = &self.socket_path {
896 addresses.push(ListenAddress::Unixgram(socket_path.into()));
897 }
898
899 if let Some(socket_stream_path) = &self.socket_stream_path {
900 addresses.push(ListenAddress::Unix(socket_stream_path.into()));
901 }
902
903 if let Some(pipe_name) = &self.pipe_name {
904 addresses.push(ListenAddress::named_pipe_with_input_buffer_size(
905 pipe_name,
906 &self.windows_pipe_security_descriptor,
907 self.buffer_size as u32,
908 ));
909 }
910
911 addresses
912 }
913
914 fn uds_origin_detection_unsupported_on_platform(&self, addresses: &[ListenAddress]) -> bool {
915 self.origin_enrichment.enabled()
916 && cfg!(not(target_os = "linux"))
917 && addresses
918 .iter()
919 .any(|address| matches!(address, ListenAddress::Unixgram(_) | ListenAddress::Unix(_)))
920 }
921
922 fn warn_if_uds_origin_detection_unsupported(&self, addresses: &[ListenAddress]) {
923 if self.uds_origin_detection_unsupported_on_platform(addresses) {
924 warn!(
925 "DogStatsD UDS origin detection is enabled, but PID-based Unix socket credentials are unsupported on \
926 this platform. Metrics are accepted without PID-based origin enrichment."
927 );
928 }
929 }
930
931 async fn build_listeners(&self) -> Result<Vec<Listener>, Error> {
933 let bind_host: Option<std::net::IpAddr> = if self.non_local_traffic {
937 None
938 } else {
939 match &self.bind_host {
940 Some(host) => Some(resolve_bind_host(host).await?),
941 None => None,
942 }
943 };
944
945 let addresses = self.build_addresses(bind_host);
946 self.warn_if_uds_origin_detection_unsupported(&addresses);
947 let mut listeners = Vec::new();
948 let socket_receive_buffer_size =
949 (self.socket_receive_buffer_size != 0).then_some(self.socket_receive_buffer_size);
950 let udp_streams_to_yield = self.udp_streams_to_yield();
951 for address in addresses {
952 let listener_type = address.listener_type();
953 let listener_streams = matches!(address, ListenAddress::Udp(_))
954 .then_some(udp_streams_to_yield)
955 .flatten();
956 let listener = Listener::from_listen_address(address, listener_streams)
957 .await
958 .context(FailedToCreateListener { listener_type })?
959 .with_receive_buffer_size(socket_receive_buffer_size);
960
961 listeners.push(listener);
962 }
963 Ok(listeners)
964 }
965}
966
967#[async_trait]
968impl SourceBuilder for DogStatsDConfiguration {
969 async fn build(&self, context: ComponentContext) -> Result<Box<dyn Source + Send>, GenericError> {
970 let listeners = self.build_listeners().await?;
971 if listeners.is_empty() {
972 return Err(Error::NoListenersConfigured.into());
973 }
974
975 let min_buffers: usize = listeners.iter().map(Listener::min_buffer_reservation).sum();
979 let max_buffers = self.effective_max_buffer_count();
980 if max_buffers < min_buffers {
981 return Err(generic_error!(
982 "The maximum I/O buffer count ({}) must be at least {} to service all configured listeners.",
983 max_buffers,
984 min_buffers,
985 ));
986 }
987
988 let origin_detection_enabled = self.origin_enrichment.enabled();
989 let captured_tagger = CapturedTaggerHandle::new();
992
993 let maybe_origin_tags_resolver = self.workload_provider.clone().map(|provider| {
994 DogStatsDOriginTagResolver::new(self.origin_enrichment.clone(), provider, captured_tagger.clone())
995 });
996 let context_resolvers = ContextResolvers::new(self, &context, maybe_origin_tags_resolver)
997 .error_context("Failed to create context resolvers.")?;
998
999 let codec_config = DogStatsDCodecConfiguration::default()
1000 .with_timestamps(self.no_aggregation_pipeline_support)
1001 .with_permissive_mode(self.permissive_decoding)
1002 .with_minimum_sample_rate(self.minimum_sample_rate)
1003 .with_client_origin_detection(self.origin_enrichment.origin_detection_client);
1004
1005 let codec = DogStatsDCodec::from_configuration(codec_config);
1006 let eol_required = self.eol_required();
1007
1008 let enable_payloads_filter = EnablePayloadsFilter::default()
1009 .with_allow_series(self.enable_payloads.series)
1010 .with_allow_sketches(self.enable_payloads.sketches)
1011 .with_allow_events(self.enable_payloads.events)
1012 .with_allow_service_checks(self.enable_payloads.service_checks);
1013 let traffic_capture = TrafficCapture::with_workload_provider(
1014 self.capture_path.clone(),
1015 self.capture_depth.max(MIN_CAPTURE_DEPTH),
1016 self.workload_provider.clone(),
1017 );
1018 self.capture_control.bind(traffic_capture.clone());
1019 let packet_forwarder_target = self.packet_forwarder_target();
1020
1021 self.replay_control.bind(captured_tagger);
1022
1023 let (io_buffer_pool, io_buffer_pool_shrinker) =
1027 build_io_buffer_pool(self.buffer_count, max_buffers, self.buffer_size);
1028
1029 Ok(Box::new(DogStatsD {
1030 listeners,
1031 io_buffer_pool,
1032 io_buffer_pool_shrinker: Box::pin(io_buffer_pool_shrinker),
1033 codec,
1034 context_resolvers,
1035 default_hostname: self.default_hostname.clone(),
1036 enabled_filter: enable_payloads_filter,
1037 origin_detection_enabled,
1038 origin_telemetry_enabled: self.telemetry.dogstatsd_origin,
1039 stream_log_too_big: self.stream_log_too_big,
1040 disable_verbose_logs: self.disable_verbose_logs,
1041 eol_required,
1042 additional_tags: self.additional_tags().into(),
1043 capture_entity_resolver: self.capture_entity_resolver.clone(),
1044 traffic_capture,
1045 packet_forwarder_target,
1046 }))
1047 }
1048
1049 fn outputs(&self) -> &[OutputDefinition<EventType>] {
1050 static OUTPUTS: LazyLock<Vec<OutputDefinition<EventType>>> = LazyLock::new(|| {
1051 vec![
1052 OutputDefinition::named_output("metrics", EventType::Metric),
1053 OutputDefinition::named_output("events", EventType::EventD),
1054 OutputDefinition::named_output("service_checks", EventType::ServiceCheck),
1055 ]
1056 });
1057 &OUTPUTS
1058 }
1059}
1060
1061impl MemoryBounds for DogStatsDConfiguration {
1062 fn specify_bounds(&self, builder: &mut MemoryBoundsBuilder) {
1063 let additional_buffers = self.effective_max_buffer_count().saturating_sub(self.buffer_count);
1064 let adjusted_buffer_size = get_adjusted_buffer_size(self.buffer_size);
1065
1066 builder
1067 .minimum()
1068 .with_single_value::<DogStatsD>("source struct")
1070 .with_expr(UsageExpr::product(
1072 "buffers",
1073 UsageExpr::config("dogstatsd_buffer_count", self.buffer_count),
1074 UsageExpr::config("dogstatsd_buffer_size", adjusted_buffer_size),
1075 ))
1076 .with_expr(UsageExpr::config(
1079 "dogstatsd_string_interner_size_bytes",
1080 self.effective_context_string_interner_bytes().as_u64() as usize,
1081 ));
1082
1083 builder.firm().with_expr(UsageExpr::product(
1085 "elastic buffers",
1086 UsageExpr::constant("dogstatsd_buffer_count_max_extra", additional_buffers),
1087 UsageExpr::config("dogstatsd_buffer_size", adjusted_buffer_size),
1088 ));
1089 }
1090}
1091
1092pub struct DogStatsD {
1094 listeners: Vec<Listener>,
1095 io_buffer_pool: ElasticObjectPool<BytesBuffer>,
1096 io_buffer_pool_shrinker: Pin<Box<dyn Future<Output = ()> + Send>>,
1097 codec: DogStatsDCodec,
1098 context_resolvers: ContextResolvers,
1099 default_hostname: MetaString,
1100 enabled_filter: EnablePayloadsFilter,
1101 origin_detection_enabled: bool,
1102 origin_telemetry_enabled: bool,
1103 stream_log_too_big: bool,
1104 disable_verbose_logs: bool,
1105 eol_required: EolRequired,
1106 additional_tags: Arc<[String]>,
1107 capture_entity_resolver: Option<Arc<dyn CaptureEntityResolver + Send + Sync>>,
1108 traffic_capture: TrafficCapture,
1109 packet_forwarder_target: Option<PacketForwarderTarget>,
1110}
1111
1112struct ListenerContext {
1113 shutdown_handle: ShutdownHandle,
1114 listener: Listener,
1115 io_buffer_pool: ElasticObjectPool<BytesBuffer>,
1116 codec: DogStatsDCodec,
1117 context_resolvers: ContextResolvers,
1118 default_hostname: MetaString,
1119 origin_detection_enabled: bool,
1120 origin_telemetry_enabled: bool,
1121 stream_log_too_big: bool,
1122 disable_verbose_logs: bool,
1123 eol_required: EolRequired,
1124 additional_tags: Arc<[String]>,
1125 capture_entity_resolver: Option<Arc<dyn CaptureEntityResolver + Send + Sync>>,
1126 traffic_capture: TrafficCapture,
1127 packet_forwarder_target: Option<PacketForwarderTarget>,
1128}
1129
1130struct HandlerContext {
1131 listen_addr: ListenAddress,
1132 framer: DsdFramer,
1133 codec: DogStatsDCodec,
1134 io_buffer_pool: ElasticObjectPool<BytesBuffer>,
1135 metrics: Metrics,
1136 context_resolvers: ContextResolvers,
1137 default_hostname: MetaString,
1138 origin_detection_enabled: bool,
1139 stream_log_too_big: bool,
1140 disable_verbose_logs: bool,
1141 additional_tags: Arc<[String]>,
1142 capture_entity_resolver: Option<Arc<dyn CaptureEntityResolver + Send + Sync>>,
1143 traffic_capture: TrafficCapture,
1144 packet_forwarder: Option<PacketForwarder>,
1145}
1146
1147#[async_trait]
1148impl Source for DogStatsD {
1149 async fn run(mut self: Box<Self>, mut context: SourceContext) -> Result<(), GenericError> {
1150 let global_shutdown = context.take_shutdown_handle();
1151 pin!(global_shutdown);
1152
1153 let mut health = context.take_health_handle();
1154
1155 let mut listener_shutdown_coordinator = ShutdownCoordinator::default();
1156 spawn_traced_named(
1157 "dogstatsd-io-buffer-pool-shrinker",
1158 process_io_buffer_pool_shrinker(self.io_buffer_pool_shrinker, listener_shutdown_coordinator.register()),
1159 );
1160
1161 for listener in self.listeners {
1163 let task_name = format!("dogstatsd-listener-{}", listener.listen_address().listener_type());
1164
1165 let listener_context = ListenerContext {
1172 shutdown_handle: listener_shutdown_coordinator.register(),
1173 listener,
1174 io_buffer_pool: self.io_buffer_pool.clone(),
1175 codec: self.codec.clone(),
1176 context_resolvers: self.context_resolvers.clone(),
1177 default_hostname: self.default_hostname.clone(),
1178 origin_detection_enabled: self.origin_detection_enabled,
1179 origin_telemetry_enabled: self.origin_telemetry_enabled,
1180 stream_log_too_big: self.stream_log_too_big,
1181 disable_verbose_logs: self.disable_verbose_logs,
1182 eol_required: self.eol_required,
1183 additional_tags: self.additional_tags.clone(),
1184 capture_entity_resolver: self.capture_entity_resolver.clone(),
1185 traffic_capture: self.traffic_capture.clone(),
1186 packet_forwarder_target: self.packet_forwarder_target.clone(),
1187 };
1188
1189 spawn_traced_named(
1190 task_name,
1191 process_listener(context.clone(), listener_context, self.enabled_filter),
1192 );
1193 }
1194
1195 health.mark_ready();
1196 debug!("DogStatsD source started.");
1197
1198 loop {
1203 select! {
1204 _ = &mut global_shutdown => {
1205 debug!("Received shutdown signal.");
1206 break
1207 },
1208 _ = health.live() => continue,
1209 }
1210 }
1211
1212 debug!("Stopping DogStatsD source...");
1213
1214 listener_shutdown_coordinator.shutdown_and_wait().await;
1215
1216 debug!("DogStatsD source stopped.");
1217
1218 Ok(())
1219 }
1220}
1221
1222async fn process_io_buffer_pool_shrinker(
1223 io_buffer_pool_shrinker: Pin<Box<dyn Future<Output = ()> + Send>>, shutdown_handle: ShutdownHandle,
1224) {
1225 pin!(shutdown_handle);
1226
1227 select! {
1228 _ = &mut shutdown_handle => {
1229 debug!("I/O buffer pool shrinker received shutdown signal.");
1230 },
1231 _ = io_buffer_pool_shrinker => {
1232 debug!("I/O buffer pool shrinker stopped.");
1233 },
1234 }
1235}
1236
1237fn build_io_buffer_pool(
1238 min_buffers: usize, max_buffers: usize, buffer_size: usize,
1239) -> (ElasticObjectPool<BytesBuffer>, impl Future<Output = ()> + Send) {
1240 saluki_antithesis::always_le!(
1241 buffer_size,
1242 usize::MAX - 4,
1243 "dogstatsd buffer size add does not overflow",
1244 { "buffer_size": buffer_size }
1245 );
1246 let adjusted_buffer_size = get_adjusted_buffer_size(buffer_size);
1247 ElasticObjectPool::with_builder("dsd_packet_bufs", min_buffers, max_buffers, move || {
1248 FixedSizeVec::with_capacity(adjusted_buffer_size)
1249 })
1250}
1251
1252async fn process_listener(
1253 source_context: SourceContext, listener_context: ListenerContext, enabled_filter: EnablePayloadsFilter,
1254) {
1255 let ListenerContext {
1256 shutdown_handle,
1257 mut listener,
1258 io_buffer_pool,
1259 codec,
1260 context_resolvers,
1261 default_hostname,
1262 origin_detection_enabled,
1263 origin_telemetry_enabled,
1264 stream_log_too_big,
1265 disable_verbose_logs,
1266 eol_required,
1267 additional_tags,
1268 capture_entity_resolver,
1269 traffic_capture,
1270 packet_forwarder_target,
1271 } = listener_context;
1272
1273 pin!(shutdown_handle);
1274
1275 let listen_addr = listener.listen_address().clone();
1276 let metrics = build_metrics(
1277 &listen_addr,
1278 source_context.component_context(),
1279 origin_telemetry_enabled,
1280 );
1281 let packet_forwarder = packet_forwarder_target
1282 .as_ref()
1283 .map(|target| target.to_forwarder(metrics.clone()));
1284 if let Some(packet_forwarder) = &packet_forwarder {
1285 packet_forwarder.spawn_connect();
1286 }
1287
1288 let mut stream_shutdown_coordinator = ShutdownCoordinator::default();
1289
1290 info!(%listen_addr, "DogStatsD listener started.");
1291
1292 loop {
1293 select! {
1294 _ = &mut shutdown_handle => {
1295 debug!(%listen_addr, "Received shutdown signal. Waiting for existing stream handlers to finish...");
1296 break;
1297 }
1298 result = listener.accept() => match result {
1299 Ok(stream) => {
1300 debug!(%listen_addr, "Spawning new stream handler.");
1301
1302 let handler_context = HandlerContext {
1303 listen_addr: listen_addr.clone(),
1304 framer: get_framer(&listen_addr, eol_required.for_listener(&listen_addr)),
1305 codec: codec.clone(),
1306 io_buffer_pool: io_buffer_pool.clone(),
1307 metrics: metrics.clone(),
1308 context_resolvers: context_resolvers.clone(),
1309 default_hostname: default_hostname.clone(),
1310 origin_detection_enabled,
1311 stream_log_too_big,
1312 disable_verbose_logs,
1313 additional_tags: additional_tags.clone(),
1314 capture_entity_resolver: capture_entity_resolver.clone(),
1315 traffic_capture: traffic_capture.clone(),
1316 packet_forwarder: packet_forwarder.clone(),
1317 };
1318
1319 let task_name = format!(
1320 "dogstatsd-stream-handler-{}",
1321 listen_addr.listener_type(),
1322 );
1323 spawn_traced_named(task_name, process_stream(stream, source_context.clone(), handler_context, stream_shutdown_coordinator.register(), enabled_filter));
1324 }
1325 Err(e) => {
1326 error!(%listen_addr, error = %e, "Failed to accept connection. Stopping listener.");
1327 break
1328 }
1329 }
1330 }
1331 }
1332
1333 stream_shutdown_coordinator.shutdown_and_wait().await;
1334
1335 info!(%listen_addr, "DogStatsD listener stopped.");
1336}
1337
1338async fn process_stream(
1339 stream: Stream, source_context: SourceContext, handler_context: HandlerContext, shutdown_handle: ShutdownHandle,
1340 enabled_filter: EnablePayloadsFilter,
1341) {
1342 select! {
1343 _ = shutdown_handle => {
1344 debug!("Stream handler received shutdown signal.");
1345 },
1346 _ = drive_stream(stream, source_context, handler_context, enabled_filter) => {},
1347 }
1348}
1349
1350fn origin_detection_failed_for_telemetry(
1351 origin_detection_enabled: bool, bytes_read: usize, peer_addr: &ConnectionAddress,
1352) -> bool {
1353 origin_detection_enabled && bytes_read > 0 && peer_addr.has_process_credential_telemetry_error()
1354}
1355
1356async fn drive_stream(
1357 mut stream: Stream, source_context: SourceContext, handler_context: HandlerContext,
1358 enabled_filter: EnablePayloadsFilter,
1359) {
1360 let HandlerContext {
1361 listen_addr,
1362 mut framer,
1363 codec,
1364 io_buffer_pool,
1365 metrics,
1366 mut context_resolvers,
1367 default_hostname,
1368 origin_detection_enabled,
1369 stream_log_too_big,
1370 disable_verbose_logs,
1371 additional_tags,
1372 capture_entity_resolver,
1373 traffic_capture,
1374 packet_forwarder,
1375 } = handler_context;
1376
1377 debug!(%listen_addr, "Stream handler started.");
1378
1379 if !stream.is_connectionless() {
1380 metrics.connections_active().increment(1);
1381 }
1382
1383 let mut stream_capture = StreamCaptureState::new();
1384 let mut buffer_flush = interval(Duration::from_millis(100));
1387 buffer_flush.set_missed_tick_behavior(MissedTickBehavior::Delay);
1388
1389 let mut event_buffer_manager = EventBufferManager::default();
1390 let mut io_buffer_manager = IoBufferManager::new(&io_buffer_pool, &stream);
1391 let memory_limiter = source_context.topology_context().memory_limiter();
1392
1393 'read: loop {
1394 let mut eof = false;
1395
1396 let mut io_buffer = io_buffer_manager.get_buffer_mut().await;
1397
1398 memory_limiter.wait_for_capacity().await;
1399
1400 select! {
1401 read_result = stream.receive(&mut io_buffer) => match read_result {
1403 Ok((bytes_read, peer_addr)) => {
1404 if bytes_read == 0 {
1405 eof = true;
1406 }
1407
1408 let is_connectionless = stream.is_connectionless();
1409 let payload = received_payload(io_buffer, bytes_read);
1410
1411 capture_uds_traffic(
1412 &listen_addr,
1413 &traffic_capture,
1414 capture_entity_resolver.as_deref(),
1415 &peer_addr,
1416 payload,
1417 &mut stream_capture,
1418 );
1419
1420 if is_connectionless {
1421 metrics.packet_receive_success().increment(1);
1422 }
1423 metrics.bytes_received().increment(bytes_read as u64);
1424 metrics.bytes_received_size().record(bytes_read as f64);
1425 let origin_detection_failed =
1426 origin_detection_failed_for_telemetry(origin_detection_enabled, bytes_read, &peer_addr);
1427 if origin_detection_failed && is_connectionless {
1428 metrics.origin_detection_errors().increment(1);
1429 }
1430
1431 let reached_eof = eof || is_connectionless;
1437
1438 trace!(
1439 buffer_len = io_buffer.remaining(),
1440 buffer_cap = io_buffer.remaining_mut(),
1441 eof = reached_eof,
1442 %listen_addr,
1443 %peer_addr,
1444 "Received {} bytes from stream.",
1445 bytes_read
1446 );
1447
1448 if should_drop_oversized_named_pipe_frame(&listen_addr, io_buffer) {
1449 metrics.framing_errors().increment(1);
1450 debug!(%listen_addr, %peer_addr, "DogStatsD named pipe frame exceeded the configured buffer size. Dropping frame.");
1451 io_buffer.clear();
1452 continue 'read;
1453 }
1454
1455 'frame: loop {
1456 let frame_result = framer.next_frame(io_buffer, reached_eof);
1457 let completed_outer_frames = framer.take_completed_outer_frames();
1458 if !is_connectionless && completed_outer_frames > 0 {
1459 metrics.packet_receive_success().increment(completed_outer_frames as u64);
1460 }
1461 if origin_detection_failed && completed_outer_frames > 0 {
1462 metrics.origin_detection_errors().increment(completed_outer_frames as u64);
1463 }
1464
1465 match frame_result {
1466 Ok(Some(frame)) => {
1467 if matches!(listen_addr, ListenAddress::NamedPipe { .. }) {
1468 metrics.packet_receive_success().increment(1);
1469 }
1470 trace!(%listen_addr, %peer_addr, ?frame, "Decoded frame.");
1471 if let Some(forwarder) = &packet_forwarder {
1472 forwarder.forward(frame.clone()).await;
1473 }
1474 match handle_frame(
1475 &frame[..],
1476 &codec,
1477 &mut context_resolvers,
1478 &metrics,
1479 capture_entity_resolver.as_deref(),
1480 origin_detection_enabled,
1481 &peer_addr,
1482 enabled_filter,
1483 &additional_tags,
1484 &default_hostname,
1485 ) {
1486 Ok(Some(event)) => {
1487 if let Some(event_buffer) = event_buffer_manager.try_push(event) {
1488 debug!(%listen_addr, %peer_addr, "Event buffer is full. Forwarding events.");
1489 dispatch_events(event_buffer, &source_context, &listen_addr).await;
1490 }
1491 },
1492 Ok(None) => {
1493 continue
1498 },
1499 Err(e) => {
1500 log_parse_failure(disable_verbose_logs, &listen_addr, &peer_addr, &frame, &e);
1501 },
1502 }
1503 }
1504 Err(e) => {
1505 metrics.framing_errors().increment(1);
1506 if should_warn_stream_log_too_big(&listen_addr, &e, stream_log_too_big) {
1507 warn!(
1508 %listen_addr,
1509 %peer_addr,
1510 error = %e,
1511 "DogStatsD stream frame exceeded the configured buffer size."
1512 );
1513 }
1514
1515 if stream.is_connectionless() {
1516 io_buffer.clear();
1517 debug!(%listen_addr, %peer_addr, error = %e, "Error decoding frame. Continuing stream.");
1520 continue 'read;
1521 } else {
1522 debug!(%listen_addr, %peer_addr, error = %e, "Error decoding frame. Stopping stream.");
1523 break 'read;
1524 }
1525 }
1526 Ok(None) => {
1527 trace!(%listen_addr, %peer_addr, "Not enough data to decode another frame.");
1528 if eof && !stream.is_connectionless() {
1529 debug!(%listen_addr, %peer_addr, "Stream received EOF. Shutting down handler.");
1530 break 'read;
1531 } else {
1532 break 'frame;
1533 }
1534 }
1535 }
1536 }
1537 },
1538 Err(e) => {
1539 metrics.packet_receive_failure().increment(1);
1540
1541 if stream.is_connectionless() {
1542 warn!(%listen_addr, error = %e, "I/O error while decoding. Continuing stream.");
1545 continue 'read;
1546 } else {
1547 warn!(%listen_addr, error = %e, "I/O error while decoding. Stopping stream.");
1548 break 'read;
1549 }
1550 }
1551 },
1552
1553 _ = buffer_flush.tick() => {
1554 if let Some(event_buffer) = event_buffer_manager.consume() {
1555 dispatch_events(event_buffer, &source_context, &listen_addr).await;
1556 }
1557 },
1558
1559 }
1560 }
1561
1562 if let Some(event_buffer) = event_buffer_manager.consume() {
1563 dispatch_events(event_buffer, &source_context, &listen_addr).await;
1564 }
1565
1566 metrics.connections_active().decrement(1);
1567
1568 debug!(%listen_addr, "Stream handler stopped.");
1569}
1570
1571fn should_drop_oversized_named_pipe_frame(listen_addr: &ListenAddress, buffer: &BytesBuffer) -> bool {
1572 matches!(listen_addr, ListenAddress::NamedPipe { .. })
1573 && buffer.remaining_mut() == 0
1574 && memchr::memchr(b'\n', buffer.chunk()).is_none()
1575}
1576
1577fn should_warn_stream_log_too_big(listen_addr: &ListenAddress, error: &FramingError, stream_log_too_big: bool) -> bool {
1578 stream_log_too_big
1579 && matches!(listen_addr, ListenAddress::Unix(_))
1580 && matches!(error, FramingError::InvalidFrame { .. })
1581}
1582
1583fn log_parse_failure(
1584 disable_verbose_logs: bool, listen_addr: &ListenAddress, peer_addr: &ConnectionAddress, frame: &[u8],
1585 error: &ParseError,
1586) {
1587 let frame = String::from_utf8_lossy(frame);
1588 if disable_verbose_logs {
1589 debug!(%listen_addr, %peer_addr, %frame, %error, "Failed to parse frame.");
1590 } else {
1591 warn!(%listen_addr, %peer_addr, %frame, %error, "Failed to parse frame.");
1592 }
1593}
1594
1595fn capture_uds_traffic(
1596 listen_addr: &ListenAddress, traffic_capture: &TrafficCapture,
1597 capture_entity_resolver: Option<&(dyn CaptureEntityResolver + Send + Sync)>, peer_addr: &ConnectionAddress,
1598 payload: &[u8], stream_capture: &mut StreamCaptureState,
1599) {
1600 if payload.is_empty() || !traffic_capture.is_ongoing() {
1601 return;
1602 }
1603
1604 match listen_addr {
1605 ListenAddress::Unixgram(_) => {
1606 let _ = traffic_capture.enqueue(build_capture_record(
1607 capture_entity_resolver,
1608 process_id_from_peer_addr(peer_addr),
1609 payload,
1610 ));
1611 }
1612 ListenAddress::Unix(_) => {
1613 stream_capture.update_peer_metadata(peer_addr);
1614 stream_capture.pending.extend(payload);
1615
1616 while let Ok(Some(outer_payload)) = stream_capture
1617 .outer_framer
1618 .next_frame(&mut stream_capture.pending, false)
1619 {
1620 let _ = traffic_capture.enqueue(build_capture_record(
1621 capture_entity_resolver,
1622 stream_capture.last_pid,
1623 &outer_payload,
1624 ));
1625 }
1626 }
1627 _ => {}
1628 }
1629}
1630
1631struct StreamCaptureState {
1632 outer_framer: LengthDelimitedFramer,
1633 pending: VecDeque<u8>,
1634 last_pid: Option<i32>,
1635}
1636
1637impl StreamCaptureState {
1638 fn new() -> Self {
1639 Self {
1640 outer_framer: LengthDelimitedFramer,
1641 pending: VecDeque::new(),
1642 last_pid: None,
1643 }
1644 }
1645
1646 fn update_peer_metadata(&mut self, peer_addr: &ConnectionAddress) {
1647 if let Some(process_id) = process_id_from_peer_addr(peer_addr) {
1648 self.last_pid = Some(process_id);
1649 }
1650 }
1651}
1652
1653fn build_capture_record(
1654 capture_entity_resolver: Option<&(dyn CaptureEntityResolver + Send + Sync)>, process_id: Option<i32>,
1655 payload: &[u8],
1656) -> CaptureRecord {
1657 CaptureRecord {
1658 timestamp_ns: capture_timestamp_ns(),
1659 payload: payload.to_vec(),
1660 pid: process_id,
1661 ancillary: Vec::new(),
1662 container_id: resolve_capture_container_id(capture_entity_resolver, process_id),
1663 }
1664}
1665
1666fn resolve_capture_container_id(
1667 capture_entity_resolver: Option<&(dyn CaptureEntityResolver + Send + Sync)>, process_id: Option<i32>,
1668) -> Option<String> {
1669 let process_id = u32::try_from(process_id?).ok()?;
1670 capture_entity_resolver
1671 .and_then(|resolver| resolver.resolve_container_entity_for_live_pid(process_id))
1672 .map(|entity_id| entity_id.to_string())
1673}
1674
1675fn process_id_from_peer_addr(peer_addr: &ConnectionAddress) -> Option<i32> {
1676 match peer_addr {
1677 ConnectionAddress::ProcessLike(ProcessIdentity::Credentials(creds)) => Some(creds.pid),
1678 _ => None,
1679 }
1680}
1681
1682fn apply_credentials_to_origin(origin: &mut RawOrigin<'_>, creds: &ProcessCredentials) {
1688 if creds.gid == REPLAY_CREDENTIALS_GID {
1689 origin.set_process_id(mark_replay_process_id(creds.uid));
1690 } else {
1691 origin.set_process_id(creds.pid as u32);
1692 }
1693}
1694
1695fn received_payload(buffer: &BytesBuffer, bytes_read: usize) -> &[u8] {
1696 let chunk = buffer.chunk();
1697 let start = chunk.len().saturating_sub(bytes_read);
1698 &chunk[start..]
1699}
1700
1701fn capture_timestamp_ns() -> i64 {
1702 SystemTime::now()
1703 .duration_since(UNIX_EPOCH)
1704 .map(|duration| duration.as_nanos().min(i64::MAX as u128) as i64)
1705 .unwrap_or_default()
1706}
1707
1708#[allow(clippy::too_many_arguments)]
1709fn handle_frame(
1710 frame: &[u8], codec: &DogStatsDCodec, context_resolvers: &mut ContextResolvers, source_metrics: &Metrics,
1711 capture_entity_resolver: Option<&(dyn CaptureEntityResolver + Send + Sync)>, origin_detection_enabled: bool,
1712 peer_addr: &ConnectionAddress, enabled_filter: EnablePayloadsFilter, additional_tags: &[String],
1713 default_hostname: &MetaString,
1714) -> Result<Option<Event>, ParseError> {
1715 let resolve_telemetry_origin = || {
1718 (source_metrics.origin_telemetry_enabled() && origin_detection_enabled)
1719 .then(|| resolve_capture_container_id(capture_entity_resolver, process_id_from_peer_addr(peer_addr)))
1720 .flatten()
1721 };
1722
1723 let parsed = match codec.decode_packet(frame) {
1724 Ok(parsed) => parsed,
1725 Err(e) => {
1726 match parse_message_type(frame) {
1728 MessageType::MetricSample => {
1729 source_metrics.record_metric_parse_failed(resolve_telemetry_origin().as_deref())
1730 }
1731 MessageType::Event => source_metrics.event_decode_failed().increment(1),
1732 MessageType::ServiceCheck => source_metrics.service_check_decode_failed().increment(1),
1733 }
1734
1735 return Err(e);
1736 }
1737 };
1738
1739 let event = match parsed {
1740 ParsedPacket::Metric(metric_packet) => {
1741 if metric_packet.num_points == 0 {
1742 return Ok(None);
1743 }
1744 let events_len = metric_packet.num_points;
1745 if !enabled_filter.allow_metric(&metric_packet) {
1746 trace!(
1747 metric.name = metric_packet.metric_name,
1748 "Skipping metric due to filter configuration."
1749 );
1750 return Ok(None);
1751 }
1752
1753 match handle_metric_packet(
1754 metric_packet,
1755 context_resolvers,
1756 peer_addr,
1757 additional_tags,
1758 default_hostname,
1759 ) {
1760 Some(metric) => {
1761 source_metrics.record_metrics_received(events_len, resolve_telemetry_origin().as_deref());
1762 Event::Metric(metric)
1763 }
1764 None => {
1765 source_metrics.failed_context_resolve_total().increment(1);
1767 return Ok(None);
1768 }
1769 }
1770 }
1771 ParsedPacket::Event(event) => {
1772 if !enabled_filter.allow_event(&event) {
1773 trace!("Skipping event {} due to filter configuration.", event.title);
1774 return Ok(None);
1775 }
1776 let tags_resolver = context_resolvers.tags();
1777 match handle_event_packet(event, tags_resolver, peer_addr, additional_tags) {
1778 Some(event) => {
1779 source_metrics.events_received().increment(1);
1780 Event::EventD(event)
1781 }
1782 None => {
1783 source_metrics.failed_context_resolve_total().increment(1);
1784 return Ok(None);
1785 }
1786 }
1787 }
1788 ParsedPacket::ServiceCheck(service_check) => {
1789 if !enabled_filter.allow_service_check(&service_check) {
1790 trace!(
1791 "Skipping service check {} due to filter configuration.",
1792 service_check.name
1793 );
1794 return Ok(None);
1795 }
1796 let tags_resolver = context_resolvers.tags();
1797 match handle_service_check_packet(service_check, tags_resolver, peer_addr, additional_tags) {
1798 Some(service_check) => {
1799 source_metrics.service_checks_received().increment(1);
1800 Event::ServiceCheck(service_check)
1801 }
1802 None => {
1803 source_metrics.failed_context_resolve_total().increment(1);
1804 return Ok(None);
1805 }
1806 }
1807 }
1808 };
1809
1810 Ok(Some(event))
1811}
1812
1813fn handle_metric_packet(
1814 packet: MetricPacket, context_resolvers: &mut ContextResolvers, peer_addr: &ConnectionAddress,
1815 additional_tags: &[String], default_hostname: &MetaString,
1816) -> Option<Metric> {
1817 let well_known_tags = WellKnownTags::from_raw_tags(packet.tags.clone());
1818
1819 let mut origin = origin_from_metric_packet(&packet, &well_known_tags);
1820 if let Some(creds) = peer_addr.process_credentials() {
1821 apply_credentials_to_origin(&mut origin, creds);
1822 }
1823
1824 let context_resolver = if packet.timestamp.is_some() {
1826 context_resolvers.no_agg()
1827 } else {
1828 context_resolvers.primary()
1829 };
1830
1831 let tags = get_filtered_tags_iterator(packet.tags, additional_tags);
1832
1833 let hostname = well_known_tags.hostname.unwrap_or(default_hostname);
1834
1835 let maybe_context = context_resolver.resolve_with_host(packet.metric_name, hostname, tags, Some(origin));
1837
1838 match maybe_context {
1839 Some(context) => {
1840 let metric_origin = well_known_tags
1841 .jmx_check_name
1842 .map(MetricOrigin::jmx_check)
1843 .unwrap_or_else(MetricOrigin::dogstatsd);
1844 let metadata = MetricMetadata::default()
1845 .with_origin(metric_origin)
1846 .with_unit(packet.unit.map_or_else(MetaString::empty, MetaString::from_static));
1847
1848 Some(Metric::from_parts(context, packet.values, metadata))
1849 }
1850 None => None,
1852 }
1853}
1854
1855fn handle_event_packet(
1856 packet: EventPacket, tags_resolver: &mut TagsResolver, peer_addr: &ConnectionAddress, additional_tags: &[String],
1857) -> Option<EventD> {
1858 let well_known_tags = WellKnownTags::from_raw_tags(packet.tags.clone());
1859
1860 let mut origin = origin_from_event_packet(&packet, &well_known_tags);
1861 if let Some(creds) = peer_addr.process_credentials() {
1862 apply_credentials_to_origin(&mut origin, creds);
1863 }
1864 let origin_tags = tags_resolver.resolve_origin_tags(Some(origin));
1865
1866 let tags = get_filtered_tags_iterator(packet.tags, additional_tags);
1867 let tags = tags_resolver.create_tag_set(tags)?;
1868
1869 let timestamp = packet
1873 .timestamp
1874 .or_else(|| SystemTime::now().duration_since(UNIX_EPOCH).ok().map(|d| d.as_secs()));
1875
1876 let eventd = EventD::new(packet.title, packet.text)
1877 .with_timestamp(timestamp)
1878 .with_hostname(packet.hostname.map(|s| s.into()))
1879 .with_aggregation_key(packet.aggregation_key.map(|s| s.into()))
1880 .with_alert_type(packet.alert_type)
1881 .with_priority(packet.priority)
1882 .with_source_type_name(Some(
1887 packet
1888 .source_type_name
1889 .map(|s| s.into())
1890 .unwrap_or_else(|| "api".into()),
1891 ))
1892 .with_alert_type(packet.alert_type)
1893 .with_tags(tags)
1894 .with_origin_tags(origin_tags);
1895
1896 Some(eventd)
1897}
1898
1899fn handle_service_check_packet(
1900 packet: ServiceCheckPacket, tags_resolver: &mut TagsResolver, peer_addr: &ConnectionAddress,
1901 additional_tags: &[String],
1902) -> Option<ServiceCheck> {
1903 let well_known_tags = WellKnownTags::from_raw_tags(packet.tags.clone());
1904
1905 let mut origin = origin_from_service_check_packet(&packet, &well_known_tags);
1906 if let Some(creds) = peer_addr.process_credentials() {
1907 apply_credentials_to_origin(&mut origin, creds);
1908 }
1909 let origin_tags = tags_resolver.resolve_origin_tags(Some(origin));
1910
1911 let tags = get_filtered_tags_iterator(packet.tags, additional_tags);
1912 let tags = tags_resolver.create_tag_set(tags)?;
1913
1914 let timestamp = packet
1918 .timestamp
1919 .or_else(|| SystemTime::now().duration_since(UNIX_EPOCH).ok().map(|d| d.as_secs()));
1920
1921 let service_check = ServiceCheck::new(packet.name, packet.status)
1922 .with_timestamp(timestamp)
1923 .with_hostname(packet.hostname.map(|s| s.into()))
1924 .with_tags(tags)
1925 .with_origin_tags(origin_tags)
1926 .with_message(packet.message.map(|s| s.into()));
1927
1928 Some(service_check)
1929}
1930
1931fn get_filtered_tags_iterator<'a>(
1932 raw_tags: RawTags<'a>, additional_tags: &'a [String],
1933) -> impl Iterator<Item = &'a str> + Clone {
1934 RawTagsFilter::exclude(raw_tags, WellKnownTagsFilterPredicate).chain(additional_tags.iter().map(|s| s.as_str()))
1937}
1938
1939async fn dispatch_events(mut event_buffer: EventsBuffer, source_context: &SourceContext, listen_addr: &ListenAddress) {
1940 debug!(%listen_addr, events_len = event_buffer.len(), "Forwarding events.");
1941
1942 if event_buffer.has_event_type(EventType::EventD) {
1952 let eventd_events = event_buffer.extract(Event::is_eventd);
1953 let events_output = source_context.dispatcher().buffered_named("events");
1954
1955 if events_output.is_err() {
1958 saluki_antithesis::unreachable!("dsd 'events' output missing at dispatch");
1959 }
1960
1961 if let Err(e) = events_output
1962 .expect("events output should always exist")
1963 .send_all(eventd_events)
1964 .await
1965 {
1966 error!(%listen_addr, error = %e, "Failed to dispatch eventd events.");
1967
1968 saluki_antithesis::unreachable!("dsd dispatch failed mid-buffer", { "stream": "events" });
1969 }
1970 }
1971
1972 if event_buffer.has_event_type(EventType::ServiceCheck) {
1974 let service_check_events = event_buffer.extract(Event::is_service_check);
1975 let service_checks_output = source_context.dispatcher().buffered_named("service_checks");
1976
1977 if service_checks_output.is_err() {
1978 saluki_antithesis::unreachable!("dsd 'service_checks' output missing at dispatch");
1979 }
1980
1981 if let Err(e) = service_checks_output
1982 .expect("service checks output should always exist")
1983 .send_all(service_check_events)
1984 .await
1985 {
1986 error!(%listen_addr, error = %e, "Failed to dispatch service check events.");
1987
1988 saluki_antithesis::unreachable!("dsd dispatch failed mid-buffer", { "stream": "service_checks" });
1989 }
1990 }
1991
1992 if !event_buffer.is_empty() {
1994 if let Err(e) = source_context
1995 .dispatcher()
1996 .dispatch_named("metrics", event_buffer)
1997 .await
1998 {
1999 error!(%listen_addr, error = %e, "Failed to dispatch metric events.");
2000
2001 saluki_antithesis::unreachable!("dsd dispatch failed mid-buffer", { "stream": "metrics" });
2002 }
2003 }
2004}
2005
2006const fn get_adjusted_buffer_size(buffer_size: usize) -> usize {
2007 buffer_size + 4
2022}
2023
2024#[cfg(test)]
2025mod tests {
2026 use std::{
2027 collections::HashMap,
2028 io::ErrorKind,
2029 net::{IpAddr, Ipv4Addr, SocketAddr, SocketAddrV4},
2030 path::PathBuf,
2031 sync::{Arc, OnceLock},
2032 time::Duration,
2033 };
2034
2035 use bytes::{BufMut as _, Bytes};
2036 use bytesize::ByteSize;
2037 use metrics::{Key, Label};
2038 use saluki_config::ConfigurationLoader;
2039 use saluki_context::{origin::RawOrigin, ContextResolverBuilder, TagsResolverBuilder};
2040 use saluki_core::{
2041 components::ComponentContext,
2042 pooling::{helpers::get_pooled_object_via_builder, ObjectPool as _},
2043 };
2044 use saluki_env::workload::{CaptureEntityResolver, EntityId};
2045 use saluki_io::{
2046 buf::{BytesBuffer, FixedSizeVec},
2047 deser::codec::dogstatsd::{DogStatsDCodec, DogStatsDCodecConfiguration, ParsedPacket},
2048 net::{ConnectionAddress, ListenAddress, ProcessCredentials, ProcessIdentity},
2049 };
2050 use saluki_metrics::test::TestRecorder;
2051 use serde_json::json;
2052 use stringtheory::MetaString;
2053 use tokio::{net::UdpSocket, sync::mpsc, time::timeout};
2054
2055 use super::{
2056 build_io_buffer_pool, default_buffer_size, default_windows_pipe_security_descriptor,
2057 filters::EnablePayloadsFilter,
2058 forwarder::{
2059 ConnectedPacketForwarder, ForwardPacket, PacketForwarder, PacketForwarderTarget, FORWARDER_QUEUE_CAPACITY,
2060 },
2061 handle_frame, handle_metric_packet,
2062 metrics::build_metrics,
2063 origin_detection_failed_for_telemetry, resolve_capture_container_id, ContextResolvers, DogStatsDConfiguration,
2064 DOGSTATSD_CAPTURE_DIR, MIN_CAPTURE_DEPTH,
2065 };
2066
2067 const LINUX_EAFNOSUPPORT: i32 = 97;
2068 const MACOS_EAFNOSUPPORT: i32 = 47;
2069
2070 fn is_ipv6_unavailable_error(error: &std::io::Error) -> bool {
2071 matches!(error.kind(), ErrorKind::AddrNotAvailable | ErrorKind::Unsupported)
2072 || matches!(error.raw_os_error(), Some(LINUX_EAFNOSUPPORT | MACOS_EAFNOSUPPORT))
2073 }
2074
2075 fn test_component_context() -> ComponentContext {
2076 ComponentContext::test_source("dogstatsd_test")
2077 }
2078
2079 #[derive(Default)]
2080 struct CaptureTestEntityResolver {
2081 pid_map: HashMap<u32, EntityId>,
2082 }
2083
2084 impl CaptureTestEntityResolver {
2085 fn with_pid_mapping(process_id: u32, entity_id: EntityId) -> Self {
2086 let mut pid_map = HashMap::new();
2087 pid_map.insert(process_id, entity_id);
2088 Self { pid_map }
2089 }
2090 }
2091
2092 impl CaptureEntityResolver for CaptureTestEntityResolver {
2093 fn resolve_container_entity_for_live_pid(&self, process_id: u32) -> Option<EntityId> {
2094 self.pid_map.get(&process_id).cloned()
2095 }
2096 }
2097
2098 fn packet_forwarder_from_sender(
2099 target_port: u16, packets_tx: mpsc::Sender<ForwardPacket>, metrics: super::metrics::Metrics,
2100 ) -> PacketForwarder {
2101 let mut forwarder =
2102 PacketForwarderTarget::new(MetaString::from_static("127.0.0.1"), target_port).to_forwarder(metrics);
2103 forwarder.connected = Arc::new(OnceLock::from(packets_tx));
2104 forwarder
2105 }
2106
2107 fn processed_metric_key(listener_type: &'static str, origin: Option<&str>) -> Key {
2108 let mut labels = vec![
2109 Label::from_static_parts("component_id", "dogstatsd_test"),
2110 Label::from_static_parts("component_type", "source"),
2111 Label::from_static_parts("listener_type", listener_type),
2112 Label::from_static_parts("message_type", "metrics"),
2113 ];
2114 if let Some(origin) = origin {
2115 labels.push(Label::new("origin", origin.to_string()));
2116 }
2117
2118 Key::from_parts("component_events_received_total", labels)
2119 }
2120
2121 fn test_context_resolvers() -> ContextResolvers {
2122 let tags_resolver = TagsResolverBuilder::for_tests().build();
2123 let context_resolver = ContextResolverBuilder::for_tests()
2124 .with_tags_resolver(Some(tags_resolver.clone()))
2125 .build();
2126 ContextResolvers::manual(context_resolver.clone(), context_resolver, tags_resolver)
2127 }
2128
2129 #[test]
2130 fn origin_telemetry_does_not_resolve_origin_when_origin_detection_is_disabled() {
2131 let recorder = TestRecorder::default();
2132 let _recorder_guard = metrics::set_default_local_recorder(&recorder);
2133 let listen_addr = ListenAddress::Unixgram("/tmp/dsd.sock".into());
2134 let context = test_component_context();
2135 let metrics = build_metrics(&listen_addr, &context, true);
2136 let codec = DogStatsDCodec::from_configuration(DogStatsDCodecConfiguration::default());
2137 let mut context_resolvers = test_context_resolvers();
2138 let capture_entity_resolver = CaptureTestEntityResolver::with_pid_mapping(
2139 42,
2140 EntityId::from_local_data("ci-pid-container").expect("container entity"),
2141 );
2142 let peer_addr = ConnectionAddress::ProcessLike(ProcessIdentity::Credentials(ProcessCredentials {
2143 pid: 42,
2144 uid: 0,
2145 gid: 0,
2146 }));
2147
2148 let event = handle_frame(
2149 b"test_metric:1|c",
2150 &codec,
2151 &mut context_resolvers,
2152 &metrics,
2153 Some(&capture_entity_resolver),
2154 false,
2155 &peer_addr,
2156 EnablePayloadsFilter::default(),
2157 &[],
2158 &MetaString::from_static("default-host"),
2159 )
2160 .expect("frame should parse");
2161
2162 assert!(event.is_some());
2163 assert_eq!(
2164 recorder.counter(processed_metric_key("unixgram", Some("container_id://pid-container"))),
2165 None
2166 );
2167 assert_eq!(recorder.counter(processed_metric_key("unixgram", Some(""))), Some(1));
2168 }
2169
2170 #[test]
2171 fn origin_telemetry_records_resolved_origin_when_origin_detection_is_enabled() {
2172 let recorder = TestRecorder::default();
2173 let _recorder_guard = metrics::set_default_local_recorder(&recorder);
2174 let listen_addr = ListenAddress::Unixgram("/tmp/dsd.sock".into());
2175 let context = test_component_context();
2176 let metrics = build_metrics(&listen_addr, &context, true);
2177 let codec = DogStatsDCodec::from_configuration(DogStatsDCodecConfiguration::default());
2178 let mut context_resolvers = test_context_resolvers();
2179 let capture_entity_resolver = CaptureTestEntityResolver::with_pid_mapping(
2180 42,
2181 EntityId::from_local_data("ci-pid-container").expect("container entity"),
2182 );
2183 let peer_addr = ConnectionAddress::ProcessLike(ProcessIdentity::Credentials(ProcessCredentials {
2184 pid: 42,
2185 uid: 0,
2186 gid: 0,
2187 }));
2188
2189 let event = handle_frame(
2190 b"test_metric:1|c",
2191 &codec,
2192 &mut context_resolvers,
2193 &metrics,
2194 Some(&capture_entity_resolver),
2195 true,
2196 &peer_addr,
2197 EnablePayloadsFilter::default(),
2198 &[],
2199 &MetaString::from_static("default-host"),
2200 )
2201 .expect("frame should parse");
2202
2203 assert!(event.is_some());
2204 assert_eq!(
2205 recorder.counter(processed_metric_key("unixgram", Some("container_id://pid-container"))),
2206 Some(1)
2207 );
2208 assert_eq!(recorder.counter(processed_metric_key("unixgram", Some(""))), Some(0));
2209 }
2210
2211 #[test]
2212 fn no_metrics_when_interner_full_allocations_disallowed() {
2213 let codec = DogStatsDCodec::from_configuration(DogStatsDCodecConfiguration::default());
2221 let tags_resolver = TagsResolverBuilder::for_tests().build();
2222 let context_resolver = ContextResolverBuilder::for_tests()
2223 .with_heap_allocations(false)
2224 .with_tags_resolver(Some(tags_resolver.clone()))
2225 .build();
2226 let mut context_resolvers = ContextResolvers::manual(context_resolver.clone(), context_resolver, tags_resolver);
2227 let peer_addr = ConnectionAddress::from("1.1.1.1:1234".parse::<SocketAddr>().unwrap());
2228
2229 let input = "big_metric_name_that_cant_possibly_be_inlined:1|c|#tag1:value1,tag2:value2,tag3:value3";
2230
2231 let Ok(ParsedPacket::Metric(packet)) = codec.decode_packet(input.as_bytes()) else {
2232 panic!("Failed to parse packet.");
2233 };
2234
2235 let maybe_metric = handle_metric_packet(
2236 packet,
2237 &mut context_resolvers,
2238 &peer_addr,
2239 &[],
2240 &MetaString::from_static("default-host"),
2241 );
2242 assert!(maybe_metric.is_none());
2243 }
2244
2245 #[test]
2246 fn metric_host_tag_disambiguates_contexts_without_remaining_tag() {
2247 let codec = DogStatsDCodec::from_configuration(DogStatsDCodecConfiguration::default());
2248 let mut context_resolvers = test_context_resolvers();
2249 let peer_addr = ConnectionAddress::from("1.1.1.1:1234".parse::<SocketAddr>().unwrap());
2250 let default_hostname = MetaString::from_static("default-host");
2251
2252 let packets = [
2253 ("unset", b"test_metric_name:1|g".as_slice(), "default-host"),
2254 ("empty", b"test_metric_name:2|g|#host:".as_slice(), ""),
2255 (
2256 "explicit_default",
2257 b"test_metric_name:3|g|#host:default-host".as_slice(),
2258 "default-host",
2259 ),
2260 (
2261 "custom",
2262 b"test_metric_name:4|g|#host:custom-host".as_slice(),
2263 "custom-host",
2264 ),
2265 ];
2266
2267 let mut metrics = Vec::new();
2268 for (case, raw, expected_host) in packets {
2269 let Ok(ParsedPacket::Metric(packet)) = codec.decode_packet(raw) else {
2270 panic!("Failed to parse {case} packet.");
2271 };
2272 let metric = handle_metric_packet(packet, &mut context_resolvers, &peer_addr, &[], &default_hostname)
2273 .unwrap_or_else(|| panic!("{case} metric should resolve"));
2274
2275 assert_eq!(metric.context().host(), Some(expected_host), "{case} context host");
2276 assert!(metric.context().tags().into_iter().all(|tag| tag.name() != "host"));
2277 metrics.push(metric);
2278 }
2279
2280 assert_eq!(metrics[0].context(), metrics[2].context());
2281 assert_ne!(metrics[0].context(), metrics[1].context());
2282 assert_ne!(metrics[0].context(), metrics[3].context());
2283 assert_ne!(metrics[1].context(), metrics[3].context());
2284 }
2285
2286 #[test]
2287 fn metric_with_additional_tags() {
2288 let codec = DogStatsDCodec::from_configuration(DogStatsDCodecConfiguration::default());
2289 let tags_resolver = TagsResolverBuilder::for_tests().build();
2290 let context_resolver = ContextResolverBuilder::for_tests()
2291 .with_heap_allocations(false)
2292 .with_tags_resolver(Some(tags_resolver.clone()))
2293 .build();
2294 let mut context_resolvers = ContextResolvers::manual(context_resolver.clone(), context_resolver, tags_resolver);
2295 let peer_addr = ConnectionAddress::from("1.1.1.1:1234".parse::<SocketAddr>().unwrap());
2296
2297 let existing_tags = ["tag1:value1", "tag2:value2", "tag3:value3"];
2298 let existing_tags_str = existing_tags.join(",");
2299
2300 let input = format!("test_metric_name:1|c|#{}", existing_tags_str);
2301 let additional_tags = [
2302 "tag4:value4".to_string(),
2303 "tag5:value5".to_string(),
2304 "tag6:value6".to_string(),
2305 ];
2306
2307 let Ok(ParsedPacket::Metric(packet)) = codec.decode_packet(input.as_bytes()) else {
2308 panic!("Failed to parse packet.");
2309 };
2310 let maybe_metric = handle_metric_packet(
2311 packet,
2312 &mut context_resolvers,
2313 &peer_addr,
2314 &additional_tags,
2315 &MetaString::from_static("default-host"),
2316 );
2317 assert!(maybe_metric.is_some());
2318
2319 let metric = maybe_metric.unwrap();
2320 let context = metric.context();
2321
2322 for tag in existing_tags {
2323 assert!(context.tags().has_tag(tag));
2324 }
2325
2326 for tag in additional_tags {
2327 assert!(context.tags().has_tag(tag));
2328 }
2329 }
2330
2331 fn deser_config(json: &str) -> DogStatsDConfiguration {
2332 serde_json::from_str(json).expect("failed to deserialize config")
2333 }
2334
2335 fn udp_listen_address() -> ListenAddress {
2336 ListenAddress::Udp(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 8125)))
2337 }
2338
2339 fn tcp_listen_address() -> ListenAddress {
2340 ListenAddress::Tcp(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 8125)))
2341 }
2342
2343 fn named_pipe_listen_address() -> ListenAddress {
2344 ListenAddress::named_pipe_with_input_buffer_size(
2345 "datadog-dogstatsd",
2346 default_windows_pipe_security_descriptor(),
2347 default_buffer_size() as u32,
2348 )
2349 }
2350
2351 #[test]
2352 fn build_addresses_includes_named_pipe_when_configured() {
2353 let config = deser_config(
2354 r#"{
2355 "dogstatsd_port": 0,
2356 "dogstatsd_pipe_name": "datadog-dogstatsd"
2357 }"#,
2358 );
2359
2360 let addresses = config.build_addresses(None);
2361
2362 assert_eq!(addresses, vec![named_pipe_listen_address()]);
2363 }
2364
2365 #[test]
2366 fn build_addresses_uses_dogstatsd_buffer_size_for_named_pipe_input_buffer() {
2367 let config = deser_config(
2368 r#"{
2369 "dogstatsd_port": 0,
2370 "dogstatsd_pipe_name": "datadog-dogstatsd",
2371 "dogstatsd_buffer_size": 16384
2372 }"#,
2373 );
2374
2375 let addresses = config.build_addresses(None);
2376
2377 let [ListenAddress::NamedPipe { input_buffer_size, .. }] = addresses.as_slice() else {
2378 panic!("expected only a named pipe listen address, got {addresses:?}");
2379 };
2380 assert_eq!(*input_buffer_size, Some(16_384));
2381 }
2382
2383 #[test]
2384 fn eol_required_matches_named_pipe_listener_type() {
2385 let config = deser_config(r#"{"dogstatsd_eol_required": ["named_pipe"]}"#);
2386 let eol_required = config.eol_required();
2387
2388 assert!(eol_required.for_listener(&named_pipe_listen_address()));
2389 assert!(!eol_required.for_listener(&udp_listen_address()));
2390 assert!(!eol_required.for_listener(&tcp_listen_address()));
2391 }
2392
2393 #[test]
2394 fn interner_size_defaults_to_2mib() {
2395 let config = deser_config("{}");
2396 assert_eq!(config.effective_context_string_interner_bytes(), ByteSize::mib(2));
2397 }
2398
2399 #[test]
2400 fn socket_receive_buffer_size_defaults_to_zero() {
2401 let config = deser_config("{}");
2402 assert_eq!(config.socket_receive_buffer_size, 0);
2403 }
2404
2405 #[test]
2406 fn socket_receive_buffer_size_from_config() {
2407 let config = deser_config(r#"{"dogstatsd_so_rcvbuf": 131072}"#);
2408 assert_eq!(config.socket_receive_buffer_size, 131_072);
2409 }
2410
2411 #[test]
2412 fn stream_log_too_big_defaults_to_false() {
2413 let config = deser_config("{}");
2414 assert!(!config.stream_log_too_big);
2415 }
2416
2417 #[test]
2418 fn stream_log_too_big_from_config() {
2419 let config = deser_config(r#"{"dogstatsd_stream_log_too_big": true}"#);
2420 assert!(config.stream_log_too_big);
2421 }
2422
2423 #[test]
2424 fn disable_verbose_logs_defaults_to_false() {
2425 let config = deser_config("{}");
2426 assert!(!config.disable_verbose_logs);
2427 }
2428
2429 #[test]
2430 fn disable_verbose_logs_from_config() {
2431 let config = deser_config(r#"{"dogstatsd_disable_verbose_logs": true}"#);
2432 assert!(config.disable_verbose_logs);
2433 }
2434
2435 #[test]
2436 fn statsd_forward_defaults_disabled() {
2437 let config = deser_config("{}");
2438 assert!(config.statsd_forward_host.is_none());
2439 assert_eq!(config.statsd_forward_port, 0);
2440 assert!(config.statsd_forward_target().is_none());
2441 }
2442
2443 #[test]
2444 fn statsd_forward_empty_host_disabled() {
2445 let config = deser_config(r#"{"statsd_forward_host": "", "statsd_forward_port": 9125}"#);
2446 assert!(config.statsd_forward_host.is_none());
2447 assert!(config.statsd_forward_target().is_none());
2448 }
2449
2450 #[test]
2451 fn statsd_forward_zero_port_disabled() {
2452 let config = deser_config(r#"{"statsd_forward_host": "127.0.0.1", "statsd_forward_port": 0}"#);
2453 assert_eq!(config.statsd_forward_host.as_deref(), Some("127.0.0.1"));
2454 assert!(config.statsd_forward_target().is_none());
2455 }
2456
2457 #[test]
2458 fn statsd_forward_host_and_port_enabled() {
2459 let config = deser_config(r#"{"statsd_forward_host": "127.0.0.1", "statsd_forward_port": 9125}"#);
2460 let (host, port) = config.statsd_forward_target().expect("forwarding should be enabled");
2461 assert_eq!(host.as_ref(), "127.0.0.1");
2462 assert_eq!(port, 9125);
2463 }
2464
2465 #[test]
2466 fn statsd_forward_invalid_target_still_builds_forwarder_handle() {
2467 let config = deser_config(r#"{"statsd_forward_host": "not a valid host", "statsd_forward_port": 9125}"#);
2468 assert!(config.packet_forwarder_target().is_some());
2469 }
2470
2471 #[tokio::test]
2472 async fn packet_forwarder_sends_payload_bytes() {
2473 let receiver = UdpSocket::bind("127.0.0.1:0").await.expect("receiver should bind");
2474 let receiver_addr = receiver.local_addr().expect("receiver should have an address");
2475 let forwarder = ConnectedPacketForwarder::connect("127.0.0.1", receiver_addr.port())
2476 .await
2477 .expect("forwarder should connect");
2478 let payload = b"daemon:666|g|#sometag1:somevalue1,sometag2:somevalue2";
2479
2480 let recorder = TestRecorder::default();
2481 let _recorder_guard = metrics::set_default_local_recorder(&recorder);
2482 let listen_addr = ListenAddress::Udp(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 8125)));
2483 let context = test_component_context();
2484 let metrics = build_metrics(&listen_addr, &context, false);
2485 let (packets_tx, packets_rx) = mpsc::channel(1);
2486 let worker = tokio::spawn(forwarder.run(packets_rx, metrics.clone()));
2487 let packet_forwarder = packet_forwarder_from_sender(receiver_addr.port(), packets_tx, metrics);
2488
2489 packet_forwarder.forward(Bytes::copy_from_slice(payload)).await;
2490
2491 let mut actual = [0u8; 128];
2492 let (received_len, _) = timeout(Duration::from_secs(1), receiver.recv_from(&mut actual))
2493 .await
2494 .expect("receive should not time out")
2495 .expect("receiver should receive payload");
2496
2497 assert_eq!(&actual[..received_len], payload);
2498 assert_eq!(
2499 recorder.counter((
2500 "component_packets_forwarded_total",
2501 &[
2502 ("component_id", "dogstatsd_test"),
2503 ("component_type", "source"),
2504 ("listener_type", "udp"),
2505 ("state", "ok"),
2506 ]
2507 )),
2508 Some(1)
2509 );
2510 assert_eq!(
2511 recorder.counter((
2512 "component_bytes_forwarded_total",
2513 &[
2514 ("component_id", "dogstatsd_test"),
2515 ("component_type", "source"),
2516 ("listener_type", "udp"),
2517 ]
2518 )),
2519 Some(payload.len() as u64)
2520 );
2521 worker.abort();
2522 }
2523
2524 #[tokio::test]
2525 async fn packet_forwarder_sends_payload_bytes_to_ipv6_target() {
2526 let receiver = match UdpSocket::bind("[::1]:0").await {
2527 Ok(receiver) => receiver,
2528 Err(e) if is_ipv6_unavailable_error(&e) => return,
2529 Err(e) => panic!("receiver should bind: {e}"),
2530 };
2531 let receiver_addr = receiver.local_addr().expect("receiver should have an address");
2532 let forwarder = ConnectedPacketForwarder::connect("::1", receiver_addr.port())
2533 .await
2534 .expect("forwarder should connect");
2535 let payload = b"daemon:666|g|#ip:6";
2536
2537 let recorder = TestRecorder::default();
2538 let _recorder_guard = metrics::set_default_local_recorder(&recorder);
2539 let listen_addr = ListenAddress::Udp(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 8125)));
2540 let context = test_component_context();
2541 let metrics = build_metrics(&listen_addr, &context, false);
2542 let (packets_tx, packets_rx) = mpsc::channel(1);
2543 let worker = tokio::spawn(forwarder.run(packets_rx, metrics.clone()));
2544 let packet_forwarder = packet_forwarder_from_sender(receiver_addr.port(), packets_tx, metrics);
2545
2546 packet_forwarder.forward(Bytes::copy_from_slice(payload)).await;
2547
2548 let mut actual = [0u8; 128];
2549 let (received_len, _) = timeout(Duration::from_secs(1), receiver.recv_from(&mut actual))
2550 .await
2551 .expect("receive should not time out")
2552 .expect("receiver should receive payload");
2553
2554 assert_eq!(&actual[..received_len], payload);
2555 worker.abort();
2556 }
2557
2558 #[tokio::test]
2559 async fn packet_forwarder_waits_when_queue_is_full() {
2560 let recorder = TestRecorder::default();
2561 let _recorder_guard = metrics::set_default_local_recorder(&recorder);
2562 let listen_addr = ListenAddress::Udp(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 8125)));
2563 let context = test_component_context();
2564 let metrics = build_metrics(&listen_addr, &context, false);
2565 let (packets_tx, _packets_rx) = mpsc::channel(FORWARDER_QUEUE_CAPACITY);
2566 let packet_forwarder = packet_forwarder_from_sender(9125, packets_tx, metrics);
2567
2568 for _ in 0..FORWARDER_QUEUE_CAPACITY {
2569 packet_forwarder.forward(Bytes::from_static(b"queued:1|c")).await;
2570 }
2571
2572 assert!(
2573 timeout(
2574 Duration::from_millis(100),
2575 packet_forwarder.forward(Bytes::from_static(b"blocked:1|c")),
2576 )
2577 .await
2578 .is_err(),
2579 "forwarding should wait for queue capacity instead of dropping"
2580 );
2581 }
2582
2583 #[tokio::test]
2584 async fn packet_forwarder_send_error_increments_error_telemetry() {
2585 let recorder = TestRecorder::default();
2586 let _recorder_guard = metrics::set_default_local_recorder(&recorder);
2587 let listen_addr = ListenAddress::Udp(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 8125)));
2588 let context = test_component_context();
2589 let metrics = build_metrics(&listen_addr, &context, false);
2590 let socket = UdpSocket::bind("127.0.0.1:0").await.expect("socket should bind");
2591 let forwarder = ConnectedPacketForwarder {
2592 socket,
2593 target: SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 9125)),
2594 };
2595 let (packets_tx, packets_rx) = mpsc::channel(1);
2596 let worker = tokio::spawn(forwarder.run(packets_rx, metrics.clone()));
2597 let packet_forwarder = packet_forwarder_from_sender(9125, packets_tx, metrics);
2598
2599 packet_forwarder.forward(Bytes::from_static(b"daemon:666|g")).await;
2600
2601 let deadline = tokio::time::Instant::now() + Duration::from_secs(1);
2602 loop {
2603 if recorder.counter((
2604 "component_packets_forwarded_total",
2605 &[
2606 ("component_id", "dogstatsd_test"),
2607 ("component_type", "source"),
2608 ("listener_type", "udp"),
2609 ("state", "error"),
2610 ],
2611 )) == Some(1)
2612 {
2613 break;
2614 }
2615
2616 assert!(
2617 tokio::time::Instant::now() < deadline,
2618 "forwarding error telemetry should be recorded"
2619 );
2620 tokio::time::sleep(Duration::from_millis(10)).await;
2621 }
2622 worker.abort();
2623 }
2624
2625 #[test]
2626 fn unsupported_platform_process_credentials_do_not_count_as_origin_detection_telemetry_errors() {
2627 let peer_addr = ConnectionAddress::ProcessLike(ProcessIdentity::Error(
2628 saluki_io::net::ProcessCredentialsError::UnsupportedPlatform,
2629 ));
2630
2631 assert!(!origin_detection_failed_for_telemetry(true, 1, &peer_addr));
2632 }
2633
2634 #[test]
2635 fn invalid_process_credentials_count_as_origin_detection_telemetry_errors() {
2636 let peer_addr = ConnectionAddress::ProcessLike(ProcessIdentity::Error(
2637 saluki_io::net::ProcessCredentialsError::InvalidCredentials,
2638 ));
2639
2640 assert!(origin_detection_failed_for_telemetry(true, 1, &peer_addr));
2641 }
2642
2643 #[test]
2644 fn autoscale_udp_listeners_defaults_to_false() {
2645 let config = deser_config("{}");
2646 assert!(!config.autoscale_udp_listeners);
2647 assert!(config.udp_streams_to_yield().is_none());
2648 }
2649
2650 #[test]
2651 fn effective_max_buffer_count_never_below_baseline() {
2652 let legacy = deser_config(r#"{"dogstatsd_buffer_count": 1024}"#);
2655 assert_eq!(legacy.effective_max_buffer_count(), 1024);
2656
2657 let explicit = deser_config(r#"{"dogstatsd_buffer_count": 128, "dogstatsd_buffer_count_max": 512}"#);
2659 assert_eq!(explicit.effective_max_buffer_count(), 512);
2660
2661 let below = deser_config(r#"{"dogstatsd_buffer_count": 200, "dogstatsd_buffer_count_max": 64}"#);
2663 assert_eq!(below.effective_max_buffer_count(), 200);
2664 }
2665
2666 #[tokio::test]
2667 async fn dogstatsd_io_buffer_pool_grows_on_demand_until_limit() {
2668 let min_buffers = 2;
2669 let max_buffers = 3;
2670 let (pool, shrinker) = build_io_buffer_pool(min_buffers, max_buffers, default_buffer_size());
2671
2672 let mut initial_buffers = Vec::with_capacity(min_buffers);
2673 for _ in 0..min_buffers {
2674 initial_buffers.push(
2675 timeout(Duration::from_secs(1), pool.acquire())
2676 .await
2677 .expect("initial buffer should be available"),
2678 );
2679 }
2680 let on_demand_buffer = timeout(Duration::from_secs(1), pool.acquire())
2681 .await
2682 .expect("pool should grow on demand before hitting the limit");
2683
2684 let capped_acquire = timeout(Duration::from_millis(25), pool.acquire()).await;
2685 assert!(capped_acquire.is_err(), "pool should wait once it reaches the limit");
2686
2687 drop(initial_buffers.pop().expect("initial buffer should still be held"));
2688 timeout(Duration::from_secs(1), pool.acquire())
2689 .await
2690 .expect("returned buffer should unblock acquisition");
2691
2692 drop(on_demand_buffer);
2693 drop(shrinker);
2694 }
2695
2696 #[test]
2697 #[cfg(target_os = "linux")]
2698 fn autoscale_udp_listeners_from_config_linux() {
2699 let config = deser_config(r#"{"dogstatsd_autoscale_udp_listeners": true}"#);
2700 assert!(config.autoscale_udp_listeners);
2701
2702 let streams = config
2703 .udp_streams_to_yield()
2704 .expect("autoscale yields at least 1 stream");
2705 let n = streams.get();
2706 assert!(
2707 (1..=4).contains(&n),
2708 "expected 1..=4 streams from vCPU formula, got {n}"
2709 );
2710 }
2711
2712 #[test]
2713 #[cfg(not(target_os = "linux"))]
2714 fn warns_for_uds_origin_detection_on_non_linux() {
2715 let config = deser_config(
2716 r#"{
2717 "dogstatsd_origin_detection": true,
2718 "dogstatsd_port": 0,
2719 "dogstatsd_socket": "/tmp/dsd.sock"
2720 }"#,
2721 );
2722 let addresses = config.build_addresses(None);
2723
2724 assert!(config.uds_origin_detection_unsupported_on_platform(&addresses));
2725 }
2726
2727 #[test]
2728 #[cfg(not(target_os = "linux"))]
2729 fn does_not_warn_for_udp_origin_detection_on_non_linux() {
2730 let config = deser_config(r#"{"dogstatsd_origin_detection": true}"#);
2731 let addresses = config.build_addresses(None);
2732
2733 assert!(!config.uds_origin_detection_unsupported_on_platform(&addresses));
2734 }
2735
2736 #[test]
2737 #[cfg(not(target_os = "linux"))]
2738 fn autoscale_udp_listeners_from_config_non_linux() {
2739 let config = deser_config(r#"{"dogstatsd_autoscale_udp_listeners": true}"#);
2740 assert!(config.autoscale_udp_listeners);
2741
2742 assert_eq!(None, config.udp_streams_to_yield());
2743 }
2744
2745 #[test]
2746 fn eol_required_defaults_to_no_listeners() {
2747 let config = deser_config("{}");
2748 let eol_required = config.eol_required();
2749
2750 assert!(!eol_required.for_listener(&udp_listen_address()));
2751 assert!(!eol_required.for_listener(&tcp_listen_address()));
2752 }
2753
2754 #[test]
2755 fn eol_required_matches_configured_listener_types() {
2756 let config = deser_config(r#"{"dogstatsd_eol_required": ["udp", "uds"]}"#);
2757 let eol_required = config.eol_required();
2758
2759 assert!(eol_required.for_listener(&udp_listen_address()));
2760 assert!(!eol_required.for_listener(&tcp_listen_address()));
2761
2762 #[cfg(unix)]
2763 {
2764 assert!(eol_required.for_listener(&ListenAddress::Unixgram("/tmp/dsd.sock".into())));
2765 assert!(eol_required.for_listener(&ListenAddress::Unix("/tmp/dsd-stream.sock".into())));
2766 }
2767 }
2768
2769 #[test]
2770 fn eol_required_accepts_space_separated_string() {
2771 let config = deser_config(r#"{"dogstatsd_eol_required": "udp uds"}"#);
2772 let eol_required = config.eol_required();
2773
2774 assert!(eol_required.for_listener(&udp_listen_address()));
2775 }
2776
2777 #[test]
2778 fn drops_full_named_pipe_buffer_without_newline() {
2779 let named_pipe_stream = named_pipe_listen_address();
2780 let mut buffer = get_pooled_object_via_builder::<_, BytesBuffer>(|| FixedSizeVec::with_capacity(8));
2781 buffer.put_slice(b"12345678");
2782
2783 assert!(super::should_drop_oversized_named_pipe_frame(
2784 &named_pipe_stream,
2785 &buffer
2786 ));
2787 }
2788
2789 #[test]
2790 fn keeps_named_pipe_partial_frame_when_buffer_has_capacity() {
2791 let named_pipe_stream = named_pipe_listen_address();
2792 let mut buffer = get_pooled_object_via_builder::<_, BytesBuffer>(|| FixedSizeVec::with_capacity(9));
2793 buffer.put_slice(b"12345678");
2794
2795 assert!(!super::should_drop_oversized_named_pipe_frame(
2796 &named_pipe_stream,
2797 &buffer
2798 ));
2799 }
2800
2801 #[test]
2802 fn keeps_full_named_pipe_buffer_with_newline() {
2803 let named_pipe_stream = named_pipe_listen_address();
2804 let mut buffer = get_pooled_object_via_builder::<_, BytesBuffer>(|| FixedSizeVec::with_capacity(8));
2805 buffer.put_slice(b"1234567\n");
2806
2807 assert!(!super::should_drop_oversized_named_pipe_frame(
2808 &named_pipe_stream,
2809 &buffer
2810 ));
2811 }
2812
2813 #[test]
2814 fn stream_log_too_big_warns_for_enabled_length_delimited_stream_invalid_frames() {
2815 let uds_stream = ListenAddress::Unix("/tmp/dsd-stream.sock".into());
2816 let named_pipe_stream = named_pipe_listen_address();
2817 let tcp_stream = ListenAddress::Tcp(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::LOCALHOST, 8125)));
2818 let error = saluki_io::deser::framing::FramingError::InvalidFrame {
2819 frame_len: 8193,
2820 reason: "frame length exceeds buffer capacity",
2821 };
2822
2823 assert!(super::should_warn_stream_log_too_big(&uds_stream, &error, true));
2824 assert!(!super::should_warn_stream_log_too_big(&uds_stream, &error, false));
2825 assert!(!super::should_warn_stream_log_too_big(&named_pipe_stream, &error, true));
2826 assert!(!super::should_warn_stream_log_too_big(&tcp_stream, &error, true));
2827 }
2828
2829 #[test]
2830 fn interner_size_from_entry_count() {
2831 let config = deser_config(r#"{"dogstatsd_string_interner_size": 4096}"#);
2833 assert_eq!(config.effective_context_string_interner_bytes(), ByteSize::mib(2));
2834 }
2835
2836 #[test]
2837 fn interner_size_from_explicit_bytes() {
2838 let config = deser_config(r#"{"dogstatsd_string_interner_size_bytes": 4194304}"#);
2839 assert_eq!(config.effective_context_string_interner_bytes(), ByteSize::b(4194304));
2840 }
2841
2842 #[test]
2843 fn interner_size_explicit_bytes_takes_priority() {
2844 let config = deser_config(
2845 r#"{"dogstatsd_string_interner_size": 4096, "dogstatsd_string_interner_size_bytes": 8388608}"#,
2846 );
2847 assert_eq!(config.effective_context_string_interner_bytes(), ByteSize::b(8388608));
2849 }
2850
2851 #[test]
2852 fn interner_size_custom_entry_count() {
2853 let config = deser_config(r#"{"dogstatsd_string_interner_size": 8192}"#);
2854 assert_eq!(config.effective_context_string_interner_bytes(), ByteSize::mib(4));
2856 }
2857
2858 fn address_list_eq(expected: &mut [ListenAddress], actual: &mut [ListenAddress]) -> Result<(), String> {
2860 if expected.len() != actual.len() {
2861 return Err(format!(
2862 "length mismatch: expected {} addresses, got {}",
2863 expected.len(),
2864 actual.len()
2865 ));
2866 }
2867
2868 expected.sort_by_key(|a| a.to_string());
2869 actual.sort_by_key(|a| a.to_string());
2870
2871 for (e, a) in expected.iter().zip(actual.iter()) {
2872 let (es, as_) = (e.to_string(), a.to_string());
2873 if es != as_ {
2874 return Err(format!("address mismatch: expected {}, got {}", es, as_));
2875 }
2876 }
2877
2878 Ok(())
2879 }
2880
2881 #[test]
2884 fn build_addresses_assertion_function_works() {
2885 let config = DogStatsDConfiguration {
2886 port: 0,
2887 tcp_port: 123,
2888 socket_path: None,
2889 socket_stream_path: None,
2890 non_local_traffic: false,
2891 ..Default::default()
2892 };
2893 let mut expected = vec![ListenAddress::Tcp(SocketAddr::V4(SocketAddrV4::new(
2894 Ipv4Addr::new(127, 0, 0, 2),
2896 123,
2897 )))];
2898 let mut actual = config.build_addresses(None);
2899 assert!(address_list_eq(&mut expected, &mut actual).is_err())
2900 }
2901
2902 #[test]
2904 fn build_addresses_no_listeners() {
2905 let config = DogStatsDConfiguration {
2906 port: 0,
2907 tcp_port: 0,
2908 socket_path: None,
2909 socket_stream_path: None,
2910 non_local_traffic: false,
2911 ..Default::default()
2912 };
2913 let mut expected = vec![];
2914 let mut actual = config.build_addresses(None);
2915 address_list_eq(&mut expected, &mut actual).unwrap();
2916 }
2917
2918 #[test]
2920 fn build_addresses_udp_local_only() {
2921 let config = DogStatsDConfiguration {
2922 port: 8125,
2923 tcp_port: 0,
2924 socket_path: None,
2925 socket_stream_path: None,
2926 non_local_traffic: false,
2927 ..Default::default()
2928 };
2929 let mut expected = vec![ListenAddress::Udp(SocketAddr::V4(SocketAddrV4::new(
2930 Ipv4Addr::new(127, 0, 0, 1),
2931 8125,
2932 )))];
2933 let mut actual = config.build_addresses(None);
2934 address_list_eq(&mut expected, &mut actual).unwrap();
2935 }
2936
2937 #[test]
2939 fn build_addresses_udp_non_local_only() {
2940 let config = DogStatsDConfiguration {
2941 port: 8125,
2942 tcp_port: 0,
2943 socket_path: None,
2944 socket_stream_path: None,
2945 non_local_traffic: true,
2946 ..Default::default()
2947 };
2948 let mut expected = vec![ListenAddress::Udp(SocketAddr::V4(SocketAddrV4::new(
2949 Ipv4Addr::new(0, 0, 0, 0),
2950 8125,
2951 )))];
2952 let mut actual = config.build_addresses(None);
2953 address_list_eq(&mut expected, &mut actual).unwrap();
2954 }
2955
2956 #[test]
2958 fn build_addresses_tcp_local_only() {
2959 let config = DogStatsDConfiguration {
2960 port: 0,
2961 tcp_port: 9000,
2962 socket_path: None,
2963 socket_stream_path: None,
2964 non_local_traffic: false,
2965 ..Default::default()
2966 };
2967 let mut expected = vec![ListenAddress::Tcp(SocketAddr::V4(SocketAddrV4::new(
2968 Ipv4Addr::new(127, 0, 0, 1),
2969 9000,
2970 )))];
2971 let mut actual = config.build_addresses(None);
2972 address_list_eq(&mut expected, &mut actual).unwrap();
2973 }
2974
2975 #[test]
2977 fn build_addresses_tcp_non_local_only() {
2978 let config = DogStatsDConfiguration {
2979 port: 0,
2980 tcp_port: 9000,
2981 socket_path: None,
2982 socket_stream_path: None,
2983 non_local_traffic: true,
2984 ..Default::default()
2985 };
2986 let mut expected = vec![ListenAddress::Tcp(SocketAddr::V4(SocketAddrV4::new(
2987 Ipv4Addr::new(0, 0, 0, 0),
2988 9000,
2989 )))];
2990 let mut actual = config.build_addresses(None);
2991 address_list_eq(&mut expected, &mut actual).unwrap();
2992 }
2993
2994 #[test]
2996 fn build_addresses_unixgram_only() {
2997 let config = DogStatsDConfiguration {
2998 port: 0,
2999 tcp_port: 0,
3000 socket_path: Some("/tmp/dsd.sock".to_string()),
3001 socket_stream_path: None,
3002 non_local_traffic: false,
3003 ..Default::default()
3004 };
3005 let mut expected = vec![ListenAddress::Unixgram("/tmp/dsd.sock".into())];
3006 let mut actual = config.build_addresses(None);
3007 address_list_eq(&mut expected, &mut actual).unwrap();
3008 }
3009
3010 #[test]
3012 fn build_addresses_unix_stream_only() {
3013 let config = DogStatsDConfiguration {
3014 port: 0,
3015 tcp_port: 0,
3016 socket_path: None,
3017 socket_stream_path: Some("/tmp/dsd-stream.sock".to_string()),
3018 non_local_traffic: false,
3019 ..Default::default()
3020 };
3021 let mut expected = vec![ListenAddress::Unix("/tmp/dsd-stream.sock".into())];
3022 let mut actual = config.build_addresses(None);
3023 address_list_eq(&mut expected, &mut actual).unwrap();
3024 }
3025
3026 #[test]
3028 fn build_addresses_all_four_non_local() {
3029 let config = DogStatsDConfiguration {
3030 port: 8125,
3031 tcp_port: 9000,
3032 socket_path: Some("/tmp/dsd.sock".to_string()),
3033 socket_stream_path: Some("/tmp/dsd-stream.sock".to_string()),
3034 non_local_traffic: true,
3035 ..Default::default()
3036 };
3037 let mut expected = vec![
3038 ListenAddress::Udp(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::new(0, 0, 0, 0), 8125))),
3039 ListenAddress::Tcp(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::new(0, 0, 0, 0), 9000))),
3040 ListenAddress::Unixgram("/tmp/dsd.sock".into()),
3041 ListenAddress::Unix("/tmp/dsd-stream.sock".into()),
3042 ];
3043 let mut actual = config.build_addresses(None);
3044 address_list_eq(&mut expected, &mut actual).unwrap();
3045 }
3046
3047 #[test]
3049 fn build_addresses_all_four_local() {
3050 let config = DogStatsDConfiguration {
3051 port: 8125,
3052 tcp_port: 9000,
3053 socket_path: Some("/tmp/dsd.sock".to_string()),
3054 socket_stream_path: Some("/tmp/dsd-stream.sock".to_string()),
3055 non_local_traffic: false,
3056 ..Default::default()
3057 };
3058 let mut expected = vec![
3059 ListenAddress::Udp(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::new(127, 0, 0, 1), 8125))),
3060 ListenAddress::Tcp(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::new(127, 0, 0, 1), 9000))),
3061 ListenAddress::Unixgram("/tmp/dsd.sock".into()),
3062 ListenAddress::Unix("/tmp/dsd-stream.sock".into()),
3063 ];
3064 let mut actual = config.build_addresses(None);
3065 address_list_eq(&mut expected, &mut actual).unwrap();
3066 }
3067
3068 #[test]
3071 fn build_addresses_bind_host_applies_to_udp_and_tcp() {
3072 let config = DogStatsDConfiguration {
3073 port: 8125,
3074 tcp_port: 9000,
3075 socket_path: Some("/tmp/dsd.sock".to_string()),
3076 socket_stream_path: None,
3077 non_local_traffic: false,
3078 ..Default::default()
3079 };
3080 let bind_host = Some(IpAddr::V4(Ipv4Addr::new(192, 168, 1, 50)));
3081 let mut expected = vec![
3082 ListenAddress::Udp(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::new(192, 168, 1, 50), 8125))),
3083 ListenAddress::Tcp(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::new(192, 168, 1, 50), 9000))),
3084 ListenAddress::Unixgram("/tmp/dsd.sock".into()),
3085 ];
3086 let mut actual = config.build_addresses(bind_host);
3087 address_list_eq(&mut expected, &mut actual).unwrap();
3088 }
3089
3090 #[test]
3094 fn build_addresses_non_local_clobbers_bind_host() {
3095 let config = DogStatsDConfiguration {
3096 port: 8125,
3097 tcp_port: 9000,
3098 socket_path: None,
3099 socket_stream_path: Some("/tmp/dsd-stream.sock".to_string()),
3100 non_local_traffic: true,
3101 ..Default::default()
3102 };
3103 let bind_host = Some(IpAddr::V4(Ipv4Addr::new(192, 168, 1, 50)));
3104 let mut expected = vec![
3105 ListenAddress::Udp(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::new(0, 0, 0, 0), 8125))),
3106 ListenAddress::Tcp(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::new(0, 0, 0, 0), 9000))),
3107 ListenAddress::Unix("/tmp/dsd-stream.sock".into()),
3108 ];
3109 let mut actual = config.build_addresses(bind_host);
3110 address_list_eq(&mut expected, &mut actual).unwrap();
3111 }
3112
3113 #[test]
3114 fn non_finite_metric_values_are_silently_dropped() {
3115 let codec = DogStatsDCodec::from_configuration(DogStatsDCodecConfiguration::default());
3120 for input in &[b"my.gauge:NaN|g" as &[u8], b"my.gauge:inf|g", b"my.gauge:-inf|g"] {
3121 match codec.decode_packet(input).expect("should decode without error") {
3122 ParsedPacket::Metric(packet) => assert_eq!(
3123 packet.num_points, 0,
3124 "non-finite value should be dropped, leaving 0 valid points"
3125 ),
3126 _ => panic!("expected Metric packet"),
3127 }
3128 }
3129 }
3130
3131 #[tokio::test]
3132 async fn fix_empty_capture_path_sets_path_from_run_path() {
3133 const RUN_PATH: &str = "/my/little/run_path";
3134
3135 let base_config_values = json!({ "run_path": RUN_PATH });
3136 let (config, _) = ConfigurationLoader::for_tests(Some(base_config_values), None, false).await;
3137
3138 let dogstatsd_config = DogStatsDConfiguration::from_configuration(&config).expect("should deserialize");
3139
3140 let expected = PathBuf::from(RUN_PATH).join(DOGSTATSD_CAPTURE_DIR);
3141 assert_eq!(expected, dogstatsd_config.capture_path);
3142 }
3143
3144 #[tokio::test]
3145 async fn fix_empty_capture_path_keeps_explicit_path() {
3146 const RUN_PATH: &str = "/my/little/run_path";
3147 const CAPTURE_PATH: &str = "/custom/path/to/capture";
3148
3149 let base_config_values = json!({ "run_path": RUN_PATH, "dogstatsd_capture_path": CAPTURE_PATH });
3150 let (config, _) = ConfigurationLoader::for_tests(Some(base_config_values), None, false).await;
3151
3152 let dogstatsd_config = DogStatsDConfiguration::from_configuration(&config).expect("should deserialize");
3153
3154 assert_eq!(PathBuf::from(CAPTURE_PATH), dogstatsd_config.capture_path);
3155 }
3156
3157 #[tokio::test]
3158 async fn from_configuration_normalizes_capture_depth() {
3159 let cases = [
3160 (json!({}), MIN_CAPTURE_DEPTH),
3161 (json!({ "dogstatsd_capture_depth": 0 }), MIN_CAPTURE_DEPTH),
3162 (json!({ "dogstatsd_capture_depth": 2048 }), 2048),
3163 ];
3164
3165 for (base_config_values, expected_depth) in cases {
3166 let (config, _) = ConfigurationLoader::for_tests(Some(base_config_values), None, false).await;
3167 let dogstatsd_config = DogStatsDConfiguration::from_configuration(&config).expect("should deserialize");
3168
3169 assert_eq!(expected_depth, dogstatsd_config.capture_depth);
3170 }
3171 }
3172
3173 #[test]
3174 fn capture_entity_resolver_is_configured_separately_from_workload_provider() {
3175 let config =
3176 DogStatsDConfiguration::default().with_capture_entity_resolver(CaptureTestEntityResolver::default());
3177
3178 assert!(config.capture_entity_resolver.is_some());
3179 assert!(config.workload_provider.is_none());
3180 }
3181
3182 #[test]
3183 fn resolve_capture_container_id_uses_live_pid_mapping() {
3184 let capture_entity_resolver = CaptureTestEntityResolver::with_pid_mapping(
3185 42,
3186 EntityId::from_local_data("ci-pid-container").expect("container entity"),
3187 );
3188
3189 assert_eq!(
3190 resolve_capture_container_id(Some(&capture_entity_resolver), Some(42)),
3191 Some("container_id://pid-container".to_string())
3192 );
3193 }
3194
3195 #[test]
3196 fn build_capture_record_ignores_payload_local_data() {
3197 let record = super::build_capture_record(None, None, b"test.metric:1|c|c:ci-local-container\n");
3198
3199 assert_eq!(record.container_id, None);
3200 assert!(record.ancillary.is_empty());
3201 }
3202
3203 #[test]
3204 fn stream_capture_state_preserves_last_pid_without_new_creds() {
3205 let mut stream_capture = super::StreamCaptureState::new();
3206
3207 stream_capture.update_peer_metadata(&ConnectionAddress::ProcessLike(ProcessIdentity::Credentials(
3208 ProcessCredentials {
3209 pid: 42,
3210 uid: 0,
3211 gid: 0,
3212 },
3213 )));
3214 stream_capture.update_peer_metadata(&ConnectionAddress::ProcessLike(ProcessIdentity::Unavailable));
3215
3216 assert_eq!(stream_capture.last_pid, Some(42));
3217 }
3218
3219 #[test]
3220 fn apply_credentials_uses_live_pid_for_normal_packet() {
3221 let mut origin = RawOrigin::default();
3222 let creds = ProcessCredentials {
3223 pid: 12345,
3224 uid: 1000,
3225 gid: 1000,
3226 };
3227 super::apply_credentials_to_origin(&mut origin, &creds);
3228
3229 assert_eq!(origin.process_id(), Some(12345));
3230 }
3231
3232 #[test]
3233 fn apply_credentials_unpacks_captured_pid_when_replay_gid_present() {
3234 let mut origin = RawOrigin::default();
3235 let captured_pid: u32 = 99887766;
3236 let creds = ProcessCredentials {
3237 pid: 12345, uid: captured_pid, gid: super::REPLAY_CREDENTIALS_GID,
3240 };
3241 super::apply_credentials_to_origin(&mut origin, &creds);
3242
3243 assert_eq!(
3244 origin.process_id(),
3245 Some(super::origin::mark_replay_process_id(captured_pid))
3246 );
3247 }
3248}
3249
3250#[cfg(test)]
3251mod config_smoke {
3252 use datadog_agent_config_testing::config_registry::structs;
3253 use datadog_agent_config_testing::run_config_smoke_tests;
3254 use serde_json::json;
3255
3256 use super::DogStatsDConfiguration;
3257 use crate::config::{DatadogRemapper, KEY_ALIASES};
3258
3259 #[tokio::test]
3260 async fn smoke_test() {
3261 run_config_smoke_tests(
3262 structs::DOGSTATSD_CONFIGURATION,
3263 &[],
3264 json!({}),
3265 |cfg| {
3266 cfg.as_typed::<DogStatsDConfiguration>()
3267 .expect("DogStatsDConfiguration should deserialize")
3268 },
3269 KEY_ALIASES,
3270 DatadogRemapper::from_env_vars,
3271 )
3272 .await
3273 }
3274}