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