saluki_components/destinations/
dogstatsd_client_telemetry.rs

1use async_trait::async_trait;
2use metrics::Counter;
3use saluki_core::{
4    accounting::{MemoryBounds, MemoryBoundsBuilder},
5    components::{destinations::*, ComponentContext},
6    data_model::event::{
7        metric::{Metric, MetricValues},
8        Event, EventType,
9    },
10    observability::ComponentMetricsExt as _,
11};
12use saluki_error::GenericError;
13use saluki_metrics::MetricsBuilder;
14use tokio::select;
15use tracing::debug;
16
17/// Configuration for the DogStatsD client telemetry destination.
18#[derive(Default)]
19pub struct DogStatsDClientTelemetryConfiguration;
20
21#[async_trait]
22impl DestinationBuilder for DogStatsDClientTelemetryConfiguration {
23    fn input_event_type(&self) -> EventType {
24        EventType::Metric
25    }
26
27    async fn build(&self, context: ComponentContext) -> Result<Box<dyn Destination + Send>, GenericError> {
28        Ok(Box::new(DogStatsDClientTelemetry::new(
29            MetricsBuilder::from_component_context(&context),
30        )))
31    }
32}
33
34impl MemoryBounds for DogStatsDClientTelemetryConfiguration {
35    fn specify_bounds(&self, builder: &mut MemoryBoundsBuilder) {
36        builder
37            .minimum()
38            .with_single_value::<DogStatsDClientTelemetry>("destination struct");
39    }
40}
41
42/// Mirrors supported DogStatsD client telemetry metrics into ADP internal telemetry.
43pub struct DogStatsDClientTelemetry {
44    bytes_sent: Counter,
45    bytes_dropped: Counter,
46    bytes_dropped_queue: Counter,
47    bytes_dropped_writer: Counter,
48}
49
50impl DogStatsDClientTelemetry {
51    pub(super) fn new(metrics_builder: MetricsBuilder) -> Self {
52        Self {
53            bytes_sent: metrics_builder.register_counter("dogstatsd_client_telemetry_bytes_sent"),
54            bytes_dropped: metrics_builder.register_counter("dogstatsd_client_telemetry_bytes_dropped"),
55            bytes_dropped_queue: metrics_builder.register_counter("dogstatsd_client_telemetry_bytes_dropped_queue"),
56            bytes_dropped_writer: metrics_builder.register_counter("dogstatsd_client_telemetry_bytes_dropped_writer"),
57        }
58    }
59
60    pub(super) fn record_metric(&self, metric: &Metric) {
61        let counter = match metric.context().name().as_ref() {
62            "datadog.dogstatsd.client.bytes_sent" => &self.bytes_sent,
63            "datadog.dogstatsd.client.bytes_dropped" => &self.bytes_dropped,
64            "datadog.dogstatsd.client.bytes_dropped_queue" => &self.bytes_dropped_queue,
65            "datadog.dogstatsd.client.bytes_dropped_writer" => &self.bytes_dropped_writer,
66            _ => return,
67        };
68
69        if let MetricValues::Rate(values, _) = metric.values() {
70            // A delayed aggregate flush can contain several closed time buckets in one metric. Separate tag contexts
71            // arrive as separate metrics and intentionally accumulate into this dimensionless COAT counter.
72            for (_, value) in values {
73                if value.is_finite() && value >= 0.0 && value.fract() == 0.0 && value <= u64::MAX as f64 {
74                    counter.increment(value as u64);
75                }
76            }
77        }
78    }
79}
80
81#[async_trait]
82impl Destination for DogStatsDClientTelemetry {
83    async fn run(self: Box<Self>, mut context: DestinationContext) -> Result<(), GenericError> {
84        let mut health = context.take_health_handle();
85        health.mark_ready();
86        debug!("DogStatsD client telemetry destination started.");
87
88        loop {
89            select! {
90                _ = health.live() => continue,
91                maybe_events = context.events().next() => match maybe_events {
92                    Some(events) => {
93                        for event in events {
94                            if let Event::Metric(metric) = event {
95                                self.record_metric(&metric);
96                            }
97                        }
98                    },
99                    None => break,
100                },
101            }
102        }
103
104        debug!("DogStatsD client telemetry destination stopped.");
105        Ok(())
106    }
107}