saluki_components/destinations/dsd_debug_log/
mod.rs

1use std::{io::Write, path::PathBuf};
2
3use agent_data_plane_config::Live;
4use async_trait::async_trait;
5use chrono::{DateTime, Utc};
6use saluki_common::collections::FastHashMap;
7use saluki_core::{
8    accounting::{MemoryBounds, MemoryBoundsBuilder},
9    components::{
10        destinations::{Destination, DestinationBuilder, DestinationContext},
11        BuildContext,
12    },
13    data_model::{
14        event::{metric::Metric, Event, EventType},
15        tags::TagSet,
16    },
17};
18use saluki_error::{generic_error, GenericError};
19use stringtheory::MetaString;
20use tokio::select;
21use tracing::{debug, warn};
22use tracing_appender::non_blocking::{NonBlocking, NonBlockingBuilder, WorkerGuard};
23use tracing_rolling_file::{RollingConditionBase, RollingFileAppenderBase};
24
25const DEBUG_LOG_WRITER_BUFFER_LINES: usize = 4096;
26
27/// Configuration for the DogStatsD debug log destination.
28pub struct DogStatsDDebugLogConfiguration {
29    /// Whether DogStatsD metric-level statistics are enabled.
30    ///
31    /// The destination drops metrics while this runtime setting is `false`.
32    pub metrics_stats_enabled: Live<bool>,
33
34    /// Path to the DogStatsD debug log file.
35    pub log_file: PathBuf,
36
37    /// Maximum size of the active debug log file before rotation, in bytes.
38    pub log_file_max_size: u64,
39
40    /// Number of rotated debug log files to keep.
41    pub log_file_max_rolls: usize,
42}
43
44/// DogStatsD destination that writes metric debug lines to a rotating file.
45struct DogStatsDDebugLog {
46    log_file: PathBuf,
47    log_file_max_size: u64,
48    log_file_max_rolls: usize,
49    writer: Option<DebugLogWriter>,
50    metrics_stats_enabled: Live<bool>,
51    stats: FastHashMap<ContextNoOrigin, MetricSample>,
52}
53
54struct DebugLogWriter {
55    writer: NonBlocking,
56    _guard: WorkerGuard,
57}
58
59#[derive(Debug, Default)]
60struct MetricSample {
61    count: u64,
62    last_seen: u64,
63}
64
65#[derive(Eq, Hash, PartialEq)]
66struct ContextNoOrigin {
67    name: MetaString,
68    tags: TagSet,
69}
70
71impl DogStatsDDebugLog {
72    fn new(config: &DogStatsDDebugLogConfiguration) -> Result<Self, GenericError> {
73        let mut destination = Self {
74            log_file: config.log_file.clone(),
75            log_file_max_size: config.log_file_max_size,
76            log_file_max_rolls: config.log_file_max_rolls,
77            writer: None,
78            metrics_stats_enabled: config.metrics_stats_enabled.clone(),
79            stats: FastHashMap::default(),
80        };
81
82        if *destination.metrics_stats_enabled {
83            destination.ensure_writer()?;
84        }
85
86        Ok(destination)
87    }
88
89    fn process_metric(&mut self, metric: &Metric) -> Result<(), GenericError> {
90        if !*self.metrics_stats_enabled {
91            return Ok(());
92        }
93
94        self.write_metric(metric)
95    }
96
97    fn write_metric(&mut self, metric: &Metric) -> Result<(), GenericError> {
98        self.ensure_writer()?;
99
100        let context = metric.context();
101        let metric_context = ContextNoOrigin {
102            name: context.name().clone(),
103            tags: context.tags().clone(),
104        };
105
106        let timestamp = saluki_common::time::get_coarse_unix_timestamp();
107        let sample = self.stats.entry(metric_context).or_default();
108        sample.count += 1;
109        sample.last_seen = timestamp;
110
111        let writer = self.writer.as_mut().expect("writer should be initialized");
112        writeln!(
113            writer.writer,
114            "Metric Name: {} | Tags: {{{}}} | Count: {} | Last Seen: {}",
115            context.name(),
116            format_tags(context.tags()),
117            sample.count,
118            format_timestamp(sample.last_seen)
119        )
120        .map_err(|e| {
121            generic_error!(
122                "Failed to write to DogStatsD debug log file '{}': {}",
123                self.log_file.display(),
124                e
125            )
126        })
127    }
128
129    fn ensure_writer(&mut self) -> Result<(), GenericError> {
130        if self.writer.is_some() {
131            return Ok(());
132        }
133
134        let appender = RollingFileAppenderBase::new(
135            &self.log_file,
136            RollingConditionBase::new().max_size(self.log_file_max_size),
137            self.log_file_max_rolls,
138        )
139        .map_err(|e| generic_error!("Failed to open dogstatsd_log_file '{}': {}", self.log_file.display(), e))?;
140
141        let (writer, guard) = NonBlockingBuilder::default()
142            .thread_name("dsd-dbg-writer")
143            .buffered_lines_limit(DEBUG_LOG_WRITER_BUFFER_LINES)
144            // Drop debug log lines rather than slow DogStatsD metric ingestion.
145            .lossy(true)
146            .finish(appender);
147
148        self.writer = Some(DebugLogWriter { writer, _guard: guard });
149
150        Ok(())
151    }
152}
153
154#[async_trait]
155impl Destination for DogStatsDDebugLog {
156    async fn run(mut self: Box<Self>, mut context: DestinationContext) -> Result<(), GenericError> {
157        let mut health = context.take_health_handle();
158        health.mark_ready();
159
160        loop {
161            select! {
162                _ = health.live() => continue,
163                maybe_events = context.events().next() => match maybe_events {
164                    Some(events) => {
165                        for event in events {
166                            if let Event::Metric(metric) = event {
167                                if let Err(error) = self.process_metric(&metric) {
168                                    warn!(error = %error, "Failed to write DogStatsD debug log line; continuing.");
169                                }
170                            }
171                        }
172                    },
173                    None => break,
174                },
175                metrics_stats_enabled = self.metrics_stats_enabled.changed() => {
176                    debug!(metrics_stats_enabled, "Updated DogStatsD metrics stats debug logging gate.");
177                },
178            }
179        }
180
181        Ok(())
182    }
183}
184
185#[async_trait]
186impl DestinationBuilder for DogStatsDDebugLogConfiguration {
187    fn input_event_type(&self) -> EventType {
188        EventType::Metric
189    }
190
191    async fn build(&self, _context: BuildContext) -> Result<Box<dyn Destination + Send>, GenericError> {
192        DogStatsDDebugLog::new(self).map(|destination| Box::new(destination) as Box<dyn Destination + Send>)
193    }
194}
195
196impl MemoryBounds for DogStatsDDebugLogConfiguration {
197    fn specify_bounds(&self, builder: &mut MemoryBoundsBuilder) {
198        builder
199            .minimum()
200            .with_single_value::<DogStatsDDebugLog>("component struct");
201    }
202}
203
204fn format_tags(tags: &TagSet) -> String {
205    let mut formatted = String::new();
206
207    for tag in tags {
208        if !formatted.is_empty() {
209            formatted.push(' ');
210        }
211        formatted.push_str(tag.as_str());
212    }
213
214    formatted
215}
216
217fn format_timestamp(timestamp: u64) -> String {
218    i64::try_from(timestamp)
219        .ok()
220        .and_then(|ts| DateTime::<Utc>::from_timestamp(ts, 0))
221        .map(|dt| dt.format("%Y-%m-%d %H:%M:%S +0000 UTC").to_string())
222        .unwrap_or_else(|| timestamp.to_string())
223}
224
225#[cfg(test)]
226mod tests {
227    use std::{
228        fs,
229        path::{Path, PathBuf},
230    };
231    use std::{sync::Arc, time::Duration};
232
233    use agent_data_plane_config::{Live, SalukiConfiguration};
234    use saluki_core::{
235        accounting::{ComponentRegistry, MemoryLimiter},
236        components::{destinations::DestinationContext, ComponentContext},
237        data_model::event::{
238            metric::{context::Context, Metric},
239            Event,
240        },
241        health::HealthRegistry,
242        runtime::state::DataspaceRegistry,
243        topology::{interconnect::Consumer, EventsBuffer, TopologyContext},
244    };
245    use tempfile::tempdir;
246    use tokio::{runtime::Handle, sync::mpsc};
247
248    use super::{Destination, DogStatsDDebugLog, DogStatsDDebugLogConfiguration};
249
250    fn test_config(log_file: PathBuf, max_size: u64, max_rolls: usize) -> DogStatsDDebugLogConfiguration {
251        DogStatsDDebugLogConfiguration {
252            metrics_stats_enabled: Live::new_fixed(true),
253            log_file,
254            log_file_max_size: max_size,
255            log_file_max_rolls: max_rolls,
256        }
257    }
258
259    fn read_log_files(log_file: &Path, max_rolls: usize) -> String {
260        let mut output = String::new();
261
262        for roll in (0..=max_rolls).rev() {
263            let path = rolled_path(log_file, roll);
264            if path.exists() {
265                output.push_str(&fs::read_to_string(&path).expect("debug log file should be readable"));
266            }
267        }
268
269        output
270    }
271
272    fn rolled_path(log_file: &Path, roll: usize) -> PathBuf {
273        if roll == 0 {
274            log_file.to_path_buf()
275        } else {
276            PathBuf::from(format!("{}.{}", log_file.display(), roll))
277        }
278    }
279
280    fn tagged_metric() -> Metric {
281        let context = Context::from_static_parts("custom.metric", &["env:prod", "service:web"]);
282        Metric::counter(context, 1.0)
283    }
284
285    #[tokio::test]
286    async fn writes_metric_debug_lines_and_updates_count() {
287        let tempdir = tempdir().expect("temporary directory should be created");
288        let log_file = tempdir.path().join("dogstatsd-stats.log");
289        let config = test_config(log_file.clone(), 64_000, 3);
290        let metric = tagged_metric();
291
292        let mut destination = DogStatsDDebugLog::new(&config).expect("debug log destination should be built");
293        destination
294            .write_metric(&metric)
295            .expect("first metric should be written");
296        destination
297            .write_metric(&metric)
298            .expect("second metric should be written");
299        drop(destination);
300
301        let output = read_log_files(&log_file, config.log_file_max_rolls);
302        let lines = output.lines().collect::<Vec<_>>();
303
304        assert_eq!(lines.len(), 2);
305        assert!(lines[0].contains("Metric Name: custom.metric"));
306        assert!(lines[0].contains("Tags: {env:prod service:web}"));
307        assert!(lines[0].contains("Count: 1"));
308        assert!(lines[0].contains("Last Seen: "));
309        assert!(lines[1].contains("Count: 2"));
310    }
311
312    #[tokio::test]
313    async fn run_starts_and_stops_logging_with_metrics_stats_setting() {
314        let tempdir = tempdir().expect("temporary directory should be created");
315        let log_file = tempdir.path().join("dogstatsd-stats.log");
316        let cell = Arc::new(arc_swap::ArcSwap::from_pointee(SalukiConfiguration::default()));
317        let (tick_tx, tick_rx) = tokio::sync::watch::channel(());
318        let mut config = test_config(log_file.clone(), 64_000, 3);
319        config.metrics_stats_enabled = Live::new_dynamic(Arc::clone(&cell), tick_rx, |config| {
320            &config.domains.dogstatsd.debug_log.metrics_stats_enable
321        });
322        let destination = DogStatsDDebugLog::new(&config).expect("debug log destination should be built");
323
324        let component_context = ComponentContext::test_destination("test");
325        let (events_tx, events_rx) = mpsc::channel::<EventsBuffer>(4);
326        let consumer = Consumer::new(component_context.clone(), events_rx);
327        let topology_context = TopologyContext::new(
328            Arc::from("test"),
329            MemoryLimiter::noop(),
330            HealthRegistry::new(),
331            Handle::current(),
332            DataspaceRegistry::new(),
333        );
334        let health = HealthRegistry::new()
335            .register_component(&saluki_core::support::SubsystemIdentifier::from_dotted("test"))
336            .expect("component was not previously registered");
337        let context = DestinationContext::new(
338            &topology_context,
339            &component_context,
340            ComponentRegistry::default(),
341            health,
342            consumer,
343        );
344        let run_handle = tokio::spawn(async move { Box::new(destination).run(context).await });
345
346        let mut events = EventsBuffer::default();
347        assert!(events.try_push(Event::Metric(tagged_metric())).is_none());
348        events_tx
349            .send(events)
350            .await
351            .expect("disabled metric should be accepted");
352        tokio::time::timeout(Duration::from_secs(2), async {
353            while events_tx.capacity() != 4 {
354                tokio::task::yield_now().await;
355            }
356        })
357        .await
358        .expect("disabled metric should be consumed");
359        assert!(!log_file.exists());
360
361        let mut updated = (*cell.load_full()).clone();
362        updated.domains.dogstatsd.debug_log.metrics_stats_enable = true;
363        cell.store(Arc::new(updated));
364        tick_tx.send_replace(());
365
366        tokio::time::timeout(Duration::from_secs(2), async {
367            loop {
368                let mut events = EventsBuffer::default();
369                assert!(events.try_push(Event::Metric(tagged_metric())).is_none());
370                events_tx.send(events).await.expect("enabled metric should be accepted");
371                tokio::time::sleep(Duration::from_millis(10)).await;
372                if fs::read_to_string(&log_file).is_ok_and(|output| output.contains("Metric Name: custom.metric")) {
373                    break;
374                }
375            }
376        })
377        .await
378        .expect("metrics should be logged after the runtime setting is enabled");
379
380        let mut updated = (*cell.load_full()).clone();
381        updated.domains.dogstatsd.debug_log.metrics_stats_enable = false;
382        cell.store(Arc::new(updated));
383        tick_tx.send_replace(());
384
385        let line_count_after_disable = tokio::time::timeout(Duration::from_secs(2), async {
386            let mut previous_line_count = read_log_files(&log_file, config.log_file_max_rolls).lines().count();
387            let mut unchanged_samples = 0;
388
389            loop {
390                let mut events = EventsBuffer::default();
391                assert!(events.try_push(Event::Metric(tagged_metric())).is_none());
392                events_tx
393                    .send(events)
394                    .await
395                    .expect("metric should be accepted while disabling");
396                while events_tx.capacity() != 4 {
397                    tokio::task::yield_now().await;
398                }
399                tokio::time::sleep(Duration::from_millis(20)).await;
400
401                let current_line_count = read_log_files(&log_file, config.log_file_max_rolls).lines().count();
402                if current_line_count == previous_line_count {
403                    unchanged_samples += 1;
404                    if unchanged_samples == 5 {
405                        break current_line_count;
406                    }
407                } else {
408                    previous_line_count = current_line_count;
409                    unchanged_samples = 0;
410                }
411            }
412        })
413        .await
414        .expect("metrics should stop being logged after the runtime setting is disabled");
415
416        for _ in 0..3 {
417            let mut events = EventsBuffer::default();
418            assert!(events.try_push(Event::Metric(tagged_metric())).is_none());
419            events_tx
420                .send(events)
421                .await
422                .expect("disabled metric should be accepted");
423        }
424        while events_tx.capacity() != 4 {
425            tokio::task::yield_now().await;
426        }
427        tokio::time::sleep(Duration::from_millis(100)).await;
428
429        let output = read_log_files(&log_file, config.log_file_max_rolls);
430        assert_eq!(output.lines().count(), line_count_after_disable);
431
432        drop(events_tx);
433        run_handle
434            .await
435            .expect("destination task should not panic")
436            .expect("destination should stop cleanly");
437    }
438
439    #[tokio::test]
440    async fn rotates_log_file_at_configured_size() {
441        let tempdir = tempdir().expect("temporary directory should be created");
442        let log_file = tempdir.path().join("dogstatsd-stats.log");
443        let min_debug_line_len =
444            "Metric Name: custom.metric | Tags: {env:prod service:web} | Count: 1 | Last Seen: ".len();
445        let config = test_config(log_file.clone(), min_debug_line_len as u64, 2);
446        let metric = tagged_metric();
447
448        let mut destination = DogStatsDDebugLog::new(&config).expect("debug log destination should be built");
449        for _ in 0..12 {
450            destination.write_metric(&metric).expect("metric should be written");
451        }
452        drop(destination);
453
454        assert!(log_file.exists());
455        assert!(rolled_path(&log_file, 1).exists());
456        assert!(rolled_path(&log_file, 2).exists());
457        assert!(!rolled_path(&log_file, 3).exists());
458
459        let output = read_log_files(&log_file, config.log_file_max_rolls);
460        assert!(output.contains("Metric Name: custom.metric"));
461    }
462
463    #[tokio::test]
464    async fn build_error_mentions_log_file_config_key_and_path() {
465        let tempdir = tempdir().expect("temporary directory should be created");
466        let blocked_parent = tempdir.path().join("not-a-directory");
467        fs::write(&blocked_parent, "not a directory").expect("blocking file should be written");
468        let log_file = blocked_parent.join("dogstatsd-stats.log");
469        let config = test_config(log_file.clone(), 64_000, 3);
470
471        let err = match DogStatsDDebugLog::new(&config) {
472            Ok(_) => panic!("build should fail"),
473            Err(err) => err,
474        };
475        let err = err.to_string();
476
477        assert!(err.contains("dogstatsd_log_file"));
478        assert!(err.contains(&log_file.display().to_string()));
479    }
480}