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