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