saluki_components/destinations/dsd_debug_log/
mod.rs1use 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
27pub struct DogStatsDDebugLogConfiguration {
29 pub metrics_stats_enabled: Live<bool>,
33
34 pub log_file: PathBuf,
36
37 pub log_file_max_size: u64,
39
40 pub log_file_max_rolls: usize,
42}
43
44struct 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 .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}