saluki_components/sources/dogstatsd/replay/
reader.rs

1//! Reader for Datadog DogStatsD capture files.
2//!
3//! Decodes a `.dog` or `.dog.zstd` capture file into the sequence of `UnixDogstatsdMsg` records it
4//! contains, plus the optional `TaggerState` trailer.
5
6use std::{fs, path::Path};
7
8use datadog_protos::agent::{TaggerState, UnixDogstatsdMsg};
9use prost::Message;
10use saluki_error::{generic_error, GenericError};
11
12use super::file_header::{file_version, valid_header, DATADOG_HEADER, MIN_NANO_VERSION, MIN_STATE_VERSION};
13
14const ZSTD_MAGIC: [u8; 4] = [0x28, 0xB5, 0x2F, 0xFD];
15const LENGTH_PREFIX_SIZE: usize = 4;
16
17/// Timestamp resolution recorded in a capture file.
18#[derive(Debug, Clone, Copy, PartialEq, Eq)]
19pub enum TimestampResolution {
20    /// Timestamps are recorded in whole seconds (file version < 3).
21    Seconds,
22    /// Timestamps are recorded in nanoseconds (file version >= 3).
23    Nanoseconds,
24}
25
26/// Reads back a DogStatsD traffic capture file.
27#[derive(Debug)]
28pub struct TrafficCaptureReader {
29    contents: Vec<u8>,
30    version: u8,
31    offset: usize,
32}
33
34impl TrafficCaptureReader {
35    /// Opens a capture file at the given path.
36    ///
37    /// Detects zstd-compressed inputs by magic bytes and decompresses transparently. Validates the
38    /// Datadog capture header and parses the file version.
39    pub fn from_path(path: &Path) -> Result<Self, GenericError> {
40        let raw =
41            fs::read(path).map_err(|e| generic_error!("Failed to read capture file '{}': {}", path.display(), e))?;
42
43        let contents = if has_zstd_magic(&raw) {
44            zstd::stream::decode_all(raw.as_slice())
45                .map_err(|e| generic_error!("Failed to decompress capture file '{}': {}", path.display(), e))?
46        } else {
47            raw
48        };
49
50        if !valid_header(&contents) {
51            return Err(generic_error!(
52                "Capture file '{}' does not begin with a valid Datadog capture header.",
53                path.display()
54            ));
55        }
56
57        let version = file_version(&contents)?;
58
59        Ok(Self {
60            contents,
61            version,
62            offset: DATADOG_HEADER.len(),
63        })
64    }
65
66    /// Returns the capture file version.
67    pub fn version(&self) -> u8 {
68        self.version
69    }
70
71    /// Returns the timestamp resolution implied by the file version.
72    pub fn timestamp_resolution(&self) -> TimestampResolution {
73        if self.version < MIN_NANO_VERSION {
74            TimestampResolution::Seconds
75        } else {
76            TimestampResolution::Nanoseconds
77        }
78    }
79
80    /// Reads the next captured DogStatsD record from the file.
81    ///
82    /// Returns `Ok(None)` when the stream of records is exhausted: at EOF, at the four-byte state
83    /// separator that introduces the tagger trailer, or at a truncated record boundary.
84    pub fn read_next(&mut self) -> Result<Option<UnixDogstatsdMsg>, GenericError> {
85        if self.offset + LENGTH_PREFIX_SIZE > self.contents.len() {
86            return Ok(None);
87        }
88
89        let size_bytes = &self.contents[self.offset..self.offset + LENGTH_PREFIX_SIZE];
90        let size = u32::from_le_bytes(size_bytes.try_into().expect("length prefix is 4 bytes")) as usize;
91        self.offset += LENGTH_PREFIX_SIZE;
92
93        // The writer emits a zero-length prefix to mark the start of the tagger state trailer; treat
94        // that (and any size that would overrun the buffer) as the end of the record stream.
95        if size == 0 || self.offset + size > self.contents.len() {
96            // A zero-length prefix is the legitimate trailer marker. A non-zero `size` that overruns the buffer is a
97            // corrupt/oversized length prefix being silently read as clean EOF, which drops every following
98            // well-formed record. Surface the corrupt case as distinct from a real trailer.
99            saluki_antithesis::always_or_unreachable!(
100                size == 0,
101                "replay read_next stopped at the real trailer, not on a corrupt length prefix",
102                { "size": size, "offset": self.offset, "len": self.contents.len() }
103            );
104
105            return Ok(None);
106        }
107
108        let msg = UnixDogstatsdMsg::decode(&self.contents[self.offset..self.offset + size])
109            .map_err(|e| generic_error!("Failed to decode captured DogStatsD record: {}", e))?;
110        self.offset += size;
111
112        Ok(Some(msg))
113    }
114
115    /// Reads the tagger state trailer from the end of the capture file.
116    ///
117    /// Returns `Ok(None)` if the file version predates state support, or if the trailer is empty.
118    /// Does not modify the read offset, so it can be called independently of `read_next`.
119    pub fn read_state(&self) -> Result<Option<TaggerState>, GenericError> {
120        if self.version < MIN_STATE_VERSION {
121            return Ok(None);
122        }
123
124        let len = self.contents.len();
125        if len < LENGTH_PREFIX_SIZE {
126            return Ok(None);
127        }
128
129        let size_bytes = &self.contents[len - LENGTH_PREFIX_SIZE..len];
130        let size = u32::from_le_bytes(size_bytes.try_into().expect("length suffix is 4 bytes")) as usize;
131        if size == 0 {
132            return Ok(None);
133        }
134
135        if size + LENGTH_PREFIX_SIZE > len {
136            return Err(generic_error!(
137                "Tagger state trailer size ({}) exceeds capture file length ({}).",
138                size,
139                len
140            ));
141        }
142
143        let state_start = len - LENGTH_PREFIX_SIZE - size;
144        let state = TaggerState::decode(&self.contents[state_start..len - LENGTH_PREFIX_SIZE])
145            .map_err(|e| generic_error!("Failed to decode tagger state trailer: {}", e))?;
146        Ok(Some(state))
147    }
148}
149
150fn has_zstd_magic(buf: &[u8]) -> bool {
151    buf.len() >= ZSTD_MAGIC.len() && buf[..ZSTD_MAGIC.len()] == ZSTD_MAGIC
152}
153
154#[cfg(test)]
155mod tests {
156    use std::{fs, path::PathBuf, sync::Arc, time::Duration};
157
158    use saluki_env::workload::{providers::TestWorkloadProvider, EntityId};
159
160    use super::super::test_support::{unique_dir, unique_path, wait_until_inactive};
161    use super::*;
162    use crate::sources::dogstatsd::replay::writer::{CaptureRecord, CaptureTargetDir, TrafficCaptureWriter};
163
164    #[test]
165    fn plain_capture_round_trip() {
166        let (path, _dir_guard) = run_capture(1, false, &[sample_record(100, b"metric.a:1|c", 11)]);
167
168        let mut reader = TrafficCaptureReader::from_path(&path).expect("reader should open");
169        assert_eq!(reader.version(), 3);
170        assert_eq!(reader.timestamp_resolution(), TimestampResolution::Nanoseconds);
171
172        let msg = reader.read_next().expect("read should succeed").expect("first record");
173        assert_eq!(msg.timestamp, 100);
174        assert_eq!(msg.payload, b"metric.a:1|c");
175        assert_eq!(msg.payload_size, msg.payload.len() as i32);
176        assert_eq!(msg.pid, 11);
177
178        assert!(reader.read_next().expect("read should succeed").is_none());
179    }
180
181    #[test]
182    fn compressed_capture_round_trip() {
183        let (path, _dir_guard) = run_capture(
184            1,
185            true,
186            &[
187                sample_record(1, b"metric.a:1|c", 1),
188                sample_record(2, b"metric.b:2|c", 1),
189                sample_record(3, b"metric.c:3|c", 1),
190            ],
191        );
192
193        let mut reader = TrafficCaptureReader::from_path(&path).expect("reader should open");
194
195        for expected_ts in [1, 2, 3] {
196            let msg = reader.read_next().expect("read should succeed").expect("record");
197            assert_eq!(msg.timestamp, expected_ts);
198        }
199
200        assert!(reader.read_next().expect("read should succeed").is_none());
201    }
202
203    #[test]
204    fn read_next_stops_at_state_separator() {
205        let (path, _dir_guard) = run_capture(1, false, &[sample_record(7, b"x:1|c", 5)]);
206
207        let mut reader = TrafficCaptureReader::from_path(&path).expect("reader should open");
208        assert!(reader.read_next().expect("first record").is_some());
209
210        // Subsequent calls must yield None rather than attempting to decode the trailer as a record.
211        assert!(reader.read_next().expect("trailer boundary").is_none());
212        assert!(reader.read_next().expect("idempotent EOF").is_none());
213    }
214
215    #[test]
216    fn read_next_reads_records_then_stops_cleanly_when_trailer_is_truncated() {
217        // The previous name ("truncated_record_returns_none") mischaracterized this test: dropping the final bytes
218        // truncates the tagger-state *trailer*, not a record. The single record stays intact, so `read_next` must
219        // still return it and then terminate cleanly at the zero-length separator, rather than erroring on the
220        // damaged trailer.
221        let (path, _dir_guard) = run_capture(1, false, &[sample_record(1, b"metric.a:1|c", 7)]);
222
223        let bytes = fs::read(&path).expect("capture readable");
224        let truncated_path = path.with_extension("truncated");
225        fs::write(&truncated_path, &bytes[..bytes.len().saturating_sub(8)]).expect("write truncated");
226
227        let mut reader = TrafficCaptureReader::from_path(&truncated_path).expect("reader should open");
228        let record = reader
229            .read_next()
230            .expect("record read should succeed")
231            .expect("intact record should still be recovered");
232        assert_eq!(record.payload, b"metric.a:1|c");
233        assert_eq!(record.pid, 7);
234        assert!(reader.read_next().expect("clean EOF after truncated trailer").is_none());
235
236        let _ = fs::remove_file(&truncated_path);
237    }
238
239    #[test]
240    fn read_next_treats_corrupt_length_prefix_as_end_of_stream() {
241        // Regression coverage for read_next's documented corrupt/oversized-length branch: a *non-zero* length prefix
242        // that overruns the buffer must be treated as a clean end-of-stream (`Ok(None)`) rather than decoding past the
243        // end of the buffer. This is distinct from the legitimate zero-length trailer marker (covered by
244        // `read_next_stops_at_state_separator`).
245        let (path, _dir_guard) = run_capture(1, false, &[sample_record(1, b"metric.a:1|c", 1)]);
246
247        let mut bytes = fs::read(&path).expect("capture readable");
248        // Overwrite the first record's 4-byte little-endian length prefix (immediately after the 8-byte header) with a
249        // value far larger than the remaining file.
250        let prefix_start = DATADOG_HEADER.len();
251        bytes[prefix_start..prefix_start + LENGTH_PREFIX_SIZE].copy_from_slice(&u32::MAX.to_le_bytes());
252        let corrupt_path = path.with_extension("corrupt");
253        fs::write(&corrupt_path, &bytes).expect("write corrupt capture");
254
255        let mut reader = TrafficCaptureReader::from_path(&corrupt_path).expect("reader should open");
256        assert!(
257            reader
258                .read_next()
259                .expect("corrupt length prefix should read as clean EOF")
260                .is_none(),
261            "an oversized length prefix must terminate the record stream, not decode past the buffer"
262        );
263
264        let _ = fs::remove_file(&corrupt_path);
265    }
266
267    #[test]
268    fn bad_header_is_rejected() {
269        let tmp = unique_path("reader-bad-header");
270        fs::write(&tmp, b"this is not a capture file").expect("write garbage");
271
272        let err = TrafficCaptureReader::from_path(&tmp).expect_err("bad header should fail");
273        assert!(err.to_string().contains("Datadog capture header"));
274
275        let _ = fs::remove_file(&tmp);
276    }
277
278    #[test]
279    fn read_state_recovers_entity_tags() {
280        let target_dir = unique_dir("reader-state");
281        let workload = Arc::new(TestWorkloadProvider::with_entity(
282            EntityId::Container("container-xyz".into()),
283            &["env:prod", "service:api"],
284        ));
285        let writer = TrafficCaptureWriter::with_workload_provider(1, Some(workload));
286
287        let path = writer
288            .start_capture(
289                CaptureTargetDir::Explicit(target_dir.clone()),
290                Duration::from_millis(250),
291                false,
292            )
293            .expect("capture should start");
294        assert!(writer.enqueue(CaptureRecord {
295            timestamp_ns: 1,
296            payload: b"metric:1|c".to_vec(),
297            pid: Some(99),
298            ancillary: Vec::new(),
299            container_id: Some("container_id://container-xyz".to_string()),
300        }));
301        writer.stop_capture();
302        wait_until_inactive(|| writer.is_ongoing());
303
304        let reader = TrafficCaptureReader::from_path(&path).expect("reader should open");
305        let state = reader
306            .read_state()
307            .expect("state should decode")
308            .expect("state present");
309        let entity = state
310            .state
311            .get("container_id://container-xyz")
312            .expect("captured entity present");
313        assert_eq!(
314            entity.low_cardinality_tags,
315            vec!["env:prod".to_string(), "service:api".to_string()]
316        );
317        assert_eq!(
318            state.pid_map.get(&99).map(String::as_str),
319            Some("container_id://container-xyz")
320        );
321
322        let _ = fs::remove_dir_all(target_dir);
323    }
324
325    fn run_capture(queue_depth: usize, compressed: bool, records: &[CaptureRecord]) -> (PathBuf, DirGuard) {
326        let target_dir = unique_dir("reader-capture");
327        let writer = TrafficCaptureWriter::new(queue_depth);
328
329        let path = writer
330            .start_capture(
331                CaptureTargetDir::Explicit(target_dir.clone()),
332                Duration::from_millis(500),
333                compressed,
334            )
335            .expect("capture should start");
336
337        for record in records {
338            assert!(writer.enqueue(clone_record(record)));
339        }
340        writer.stop_capture();
341        wait_until_inactive(|| writer.is_ongoing());
342
343        (path, DirGuard { path: target_dir })
344    }
345
346    fn clone_record(record: &CaptureRecord) -> CaptureRecord {
347        CaptureRecord {
348            timestamp_ns: record.timestamp_ns,
349            payload: record.payload.clone(),
350            pid: record.pid,
351            ancillary: record.ancillary.clone(),
352            container_id: record.container_id.clone(),
353        }
354    }
355
356    fn sample_record(timestamp_ns: i64, payload: &[u8], pid: i32) -> CaptureRecord {
357        CaptureRecord {
358            timestamp_ns,
359            payload: payload.to_vec(),
360            pid: Some(pid),
361            ancillary: Vec::new(),
362            container_id: Some(format!("container_id://container-{}", pid)),
363        }
364    }
365
366    struct DirGuard {
367        path: PathBuf,
368    }
369
370    impl Drop for DirGuard {
371        fn drop(&mut self) {
372            let _ = fs::remove_dir_all(&self.path);
373        }
374    }
375}