saluki_components/sources/dogstatsd/replay/
reader.rs1use 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
19pub enum TimestampResolution {
20 Seconds,
22 Nanoseconds,
24}
25
26#[derive(Debug)]
28pub struct TrafficCaptureReader {
29 contents: Vec<u8>,
30 version: u8,
31 offset: usize,
32}
33
34impl TrafficCaptureReader {
35 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 pub fn version(&self) -> u8 {
68 self.version
69 }
70
71 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 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 if size == 0 || self.offset + size > self.contents.len() {
96 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 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 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 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 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 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}