1#![cfg_attr(not(feature = "pipeline"), allow(unreachable_pub))]
9
10use std::path::{Path, PathBuf};
11
12use crate::primitives::fs;
13
14#[derive(Debug, Clone, PartialEq, Eq)]
16pub struct SealedSegment {
17 pub(crate) path: PathBuf,
18 pub(crate) index: u32,
19}
20
21#[cfg(all(feature = "test-util", feature = "pipeline"))]
22impl SealedSegment {
23 pub(crate) fn new_for_test(path: PathBuf, index: u32) -> Self {
25 Self { path, index }
26 }
27}
28
29#[cfg(feature = "pipeline")]
30impl SealedSegment {
31 pub fn path(&self) -> &Path {
33 &self.path
34 }
35
36 pub fn index(&self) -> u32 {
38 self.index
39 }
40}
41
42#[derive(Debug, Clone, PartialEq, Eq)]
44pub struct MemorySegment {
45 pub(crate) index: u32,
46 pub(crate) size: u64,
47}
48
49#[cfg(feature = "pipeline")]
50impl MemorySegment {
51 pub fn index(&self) -> u32 {
53 self.index
54 }
55
56 pub fn size(&self) -> u64 {
58 self.size
59 }
60}
61
62#[derive(Debug, Clone, PartialEq, Eq)]
64#[non_exhaustive]
65pub enum SegmentRef {
66 Disk(SealedSegment),
68 Memory(MemorySegment),
70}
71
72impl SegmentRef {
73 pub fn index(&self) -> u32 {
75 match self {
76 SegmentRef::Disk(s) => s.index,
77 SegmentRef::Memory(m) => m.index,
78 }
79 }
80
81 pub fn disk_path(&self) -> Option<&Path> {
83 match self {
84 SegmentRef::Disk(s) => Some(&s.path),
85 SegmentRef::Memory(_) => None,
86 }
87 }
88}
89
90impl std::fmt::Display for SegmentRef {
91 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
92 match self {
93 SegmentRef::Disk(s) => write!(f, "{}", s.path.display()),
94 SegmentRef::Memory(m) => write!(f, "mem://{}", m.index),
95 }
96 }
97}
98
99pub(crate) fn creation_epoch_secs(data: &[u8], path: &Path) -> (u64, bool) {
103 match parse_segment_timestamp(data) {
104 Ok(ts) => return (ts / 1_000_000_000, true),
105 Err(e) => {
106 crate::rate_limit::rate_limited!(std::time::Duration::from_secs(60), {
107 tracing::warn!(
108 path = %path.display(),
109 error = %e,
110 "failed to parse segment timestamp, falling back to mtime"
111 );
112 });
113 }
114 }
115 (mtime_or_now_secs(path), false)
116}
117
118#[cfg(feature = "pipeline")]
123pub(crate) fn seal_epoch_secs(path: &Path) -> u64 {
124 mtime_or_now_secs(path)
125}
126
127fn mtime_or_now_secs(path: &Path) -> u64 {
128 fs::metadata(path)
129 .and_then(|m| m.modified())
130 .ok()
131 .and_then(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
132 .map(|d| d.as_secs())
133 .unwrap_or_else(|| {
134 std::time::SystemTime::now()
135 .duration_since(std::time::UNIX_EPOCH)
136 .unwrap_or_default()
137 .as_secs()
138 })
139}
140
141pub(crate) const LEGACY_EPOCH_NS_FLOOR: u64 = 1_577_836_800_000_000_000;
151
152fn parse_segment_timestamp(data: &[u8]) -> Result<u64, ParseTimestampError> {
155 use dial9_trace_format::decoder::{DecodedFrameRef, Decoder};
156 use dial9_trace_format::types::FieldValueRef;
157
158 let mut dec = Decoder::new(data).ok_or(ParseTimestampError::InvalidHeader)?;
159 let mut events_seen = 0;
160 let mut legacy_fallback: Option<u64> = None;
161 loop {
162 match dec.next_frame_ref() {
163 Ok(Some(DecodedFrameRef::Event {
164 type_id,
165 timestamp_ns,
166 values,
167 ..
168 })) => {
169 events_seen += 1;
170 let name = dec
171 .registry()
172 .get(type_id)
173 .map(|s| s.name())
174 .ok_or(ParseTimestampError::UnknownTypeId(type_id.0))?;
175 if name == "ClockSyncEvent" {
176 return match values.first() {
177 Some(FieldValueRef::Varint(v)) => Ok(*v),
178 _ => Err(ParseTimestampError::MissingRealtimeField),
179 };
180 }
181 if name == "SegmentMetadataEvent" && timestamp_ns >= LEGACY_EPOCH_NS_FLOOR {
182 legacy_fallback = Some(timestamp_ns);
183 }
184 if events_seen >= 10 {
185 return legacy_fallback.ok_or(ParseTimestampError::NoAnchorInFirst10Events);
186 }
187 }
188 Ok(Some(_)) => {} Ok(None) => {
190 return legacy_fallback.ok_or(ParseTimestampError::EndOfStream { events_seen });
191 }
192 Err(e) => {
193 return Err(ParseTimestampError::DecodeError(e.to_string()));
194 }
195 }
196 }
197}
198
199#[derive(Debug)]
200enum ParseTimestampError {
201 InvalidHeader,
202 UnknownTypeId(u16),
203 MissingRealtimeField,
204 NoAnchorInFirst10Events,
205 EndOfStream { events_seen: u32 },
206 DecodeError(String),
207}
208
209impl std::fmt::Display for ParseTimestampError {
210 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
211 match self {
212 Self::InvalidHeader => write!(f, "invalid trace header"),
213 Self::UnknownTypeId(id) => write!(f, "unknown type_id {id} not in registry"),
214 Self::MissingRealtimeField => write!(f, "ClockSyncEvent had no realtime_ns field"),
215 Self::NoAnchorInFirst10Events => write!(
216 f,
217 "no ClockSyncEvent or legacy wall-clock SegmentMetadataEvent in first 10 events"
218 ),
219 Self::EndOfStream { events_seen } => write!(
220 f,
221 "end of stream after {events_seen} events without a clock anchor"
222 ),
223 Self::DecodeError(e) => write!(f, "decode error: {e}"),
224 }
225 }
226}
227
228#[cfg(feature = "pipeline")]
234pub(crate) fn find_sealed_segments(dir: &Path, stem: &str) -> std::io::Result<Vec<SealedSegment>> {
235 let mut segments = Vec::new();
236 for entry in fs::read_dir(dir)? {
237 let entry = entry?;
238 let path = entry.path();
239 if path.extension().is_none_or(|ext| ext != "bin") {
240 continue;
241 }
242 let file_name = match path.file_name().and_then(|n| n.to_str()) {
245 Some(n) => n,
246 None => continue,
247 };
248 if let Some(index) = parse_segment_index(file_name, stem) {
249 segments.push(SealedSegment { path, index });
250 }
251 }
252 segments.sort_by_key(|s| s.index);
253 Ok(segments)
254}
255
256#[derive(Debug, Clone, Copy, PartialEq, Eq)]
259pub(crate) enum SegmentArtifact {
260 Retained { index: u32 },
261 Active,
262}
263
264pub(crate) fn parse_segment_artifact(file_name: &str, stem: &str) -> Option<SegmentArtifact> {
266 let rest = file_name.strip_prefix(stem)?.strip_prefix('.')?;
267 if let Some(idx) = rest.strip_suffix(".bin.active")
268 && idx.parse::<u32>().is_ok()
269 {
270 return Some(SegmentArtifact::Active);
271 }
272 let (index, suffix) = rest.split_once(".bin")?;
273 if !suffix.is_empty() && !suffix.starts_with('.') {
274 return None;
275 }
276 Some(SegmentArtifact::Retained {
277 index: index.parse().ok()?,
278 })
279}
280
281#[cfg(feature = "pipeline")]
284fn parse_segment_index(file_name: &str, stem: &str) -> Option<u32> {
285 let rest = file_name.strip_prefix(stem)?.strip_prefix('.')?;
286 let index_str = rest.strip_suffix(".bin")?;
287 index_str.parse().ok()
288}
289
290#[cfg(all(test, feature = "pipeline"))]
292mod tests {
293 use super::*;
294 use crate::format::{ClockSyncEvent, SegmentMetadataEvent};
295 use assert2::check;
296 use dial9_trace_format::encoder::Encoder;
297 use std::fs::File;
298 use tempfile::TempDir;
299
300 fn touch(dir: &Path, name: &str) {
301 File::create(dir.join(name)).unwrap();
302 }
303
304 #[test]
305 fn finds_sealed_ignores_active() {
306 let dir = TempDir::new().unwrap();
307 touch(dir.path(), "trace.0.bin");
308 touch(dir.path(), "trace.1.bin");
309 touch(dir.path(), "trace.2.bin.active");
310
311 let segments = find_sealed_segments(dir.path(), "trace").unwrap();
312 check!(segments.len() == 2);
313 check!(segments[0].index == 0);
314 check!(segments[1].index == 1);
315 }
316
317 #[test]
318 fn sorted_oldest_first() {
319 let dir = TempDir::new().unwrap();
320 touch(dir.path(), "trace.5.bin");
321 touch(dir.path(), "trace.2.bin");
322 touch(dir.path(), "trace.9.bin");
323
324 let segments = find_sealed_segments(dir.path(), "trace").unwrap();
325 let indices: Vec<u32> = segments.iter().map(|s| s.index).collect();
326 check!(indices == [2, 5, 9]);
327 }
328
329 #[test]
330 fn ignores_unrelated_files() {
331 let dir = TempDir::new().unwrap();
332 touch(dir.path(), "trace.0.bin");
333 touch(dir.path(), "other.0.bin");
334 touch(dir.path(), "trace.txt");
335 touch(dir.path(), "readme.md");
336
337 let segments = find_sealed_segments(dir.path(), "trace").unwrap();
338 check!(segments.len() == 1);
339 check!(segments[0].index == 0);
340 }
341
342 #[test]
343 fn empty_directory() {
344 let dir = TempDir::new().unwrap();
345 let segments = find_sealed_segments(dir.path(), "trace").unwrap();
346 check!(segments.is_empty());
347 }
348
349 #[test]
352 fn test_parse_segment_timestamp() {
353 let now_nanos = std::time::SystemTime::now()
354 .duration_since(std::time::UNIX_EPOCH)
355 .unwrap()
356 .as_nanos() as u64;
357 let mut enc = Encoder::new_to(Vec::new()).unwrap();
358 enc.write(&ClockSyncEvent {
359 timestamp_ns: 1_000_000_000,
360 realtime_ns: now_nanos,
361 })
362 .unwrap();
363 let data = enc.into_inner();
364
365 let parsed = parse_segment_timestamp(&data).unwrap();
366 check!(parsed == now_nanos);
367 }
368
369 #[test]
370 fn test_creation_epoch_secs_uses_parsed_timestamp() {
371 let now = std::time::SystemTime::now()
372 .duration_since(std::time::UNIX_EPOCH)
373 .unwrap();
374 let mut enc = Encoder::new_to(Vec::new()).unwrap();
375 enc.write(&ClockSyncEvent {
376 timestamp_ns: 1_000_000_000,
377 realtime_ns: now.as_nanos() as u64,
378 })
379 .unwrap();
380 let data = enc.into_inner();
381
382 let dir = TempDir::new().unwrap();
383 let path = dir.path().join("trace.0.bin");
385 let (epoch_secs, header_valid) = creation_epoch_secs(&data, &path);
386 check!(header_valid);
387 check!(epoch_secs == now.as_secs());
388 }
389
390 #[test]
391 fn test_creation_epoch_secs_invalid_data_falls_back() {
392 let dir = tempfile::TempDir::new().unwrap();
393 let path = dir.path().join("trace.0.bin");
394 std::fs::write(&path, b"not a valid trace").unwrap();
395
396 let data = std::fs::read(&path).unwrap();
397 let (epoch_secs, header_valid) = creation_epoch_secs(&data, &path);
398 check!(!header_valid);
399 let now = std::time::SystemTime::now()
401 .duration_since(std::time::UNIX_EPOCH)
402 .unwrap()
403 .as_secs();
404 check!(now.abs_diff(epoch_secs) < 60);
405 }
406
407 fn legacy_segment_metadata_trace(timestamp_ns: u64) -> Vec<u8> {
410 let mut enc = Encoder::new_to(Vec::new()).unwrap();
411 enc.write(&SegmentMetadataEvent {
412 timestamp_ns,
413 entries: vec![("k".into(), "v".into())],
414 })
415 .unwrap();
416 enc.into_inner()
417 }
418
419 #[test]
422 fn test_parse_segment_timestamp_legacy_segment_metadata_fallback() {
423 const LEGACY_WALL_NS: u64 = 1_717_200_000_000_000_000;
425 let data = legacy_segment_metadata_trace(LEGACY_WALL_NS);
426
427 let parsed = parse_segment_timestamp(&data).unwrap();
428 check!(parsed == LEGACY_WALL_NS);
429 }
430
431 #[test]
432 fn test_creation_epoch_secs_uses_legacy_segment_metadata_fallback() {
433 const LEGACY_WALL_NS: u64 = 1_717_200_000_000_000_000;
434 let data = legacy_segment_metadata_trace(LEGACY_WALL_NS);
435 let dir = TempDir::new().unwrap();
436 let path = dir.path().join("legacy.0.bin");
437 std::fs::write(&path, &data).unwrap();
438 let (epoch_secs, header_valid) = creation_epoch_secs(&data, &path);
439 check!(header_valid);
440 check!(epoch_secs == LEGACY_WALL_NS / 1_000_000_000);
441 }
442
443 #[test]
446 fn test_parse_segment_timestamp_monotonic_segment_metadata_is_not_legacy() {
447 let data = legacy_segment_metadata_trace(12_345);
448
449 let result = parse_segment_timestamp(&data);
450 check!(matches!(
451 result,
452 Err(ParseTimestampError::EndOfStream { .. })
453 | Err(ParseTimestampError::NoAnchorInFirst10Events)
454 ));
455 }
456
457 #[test]
458 fn parse_segment_artifact_classifies_family() {
459 check!(
460 parse_segment_artifact("trace.0.bin", "trace")
461 == Some(SegmentArtifact::Retained { index: 0 })
462 );
463 check!(
464 parse_segment_artifact("trace.42.bin.gz", "trace")
465 == Some(SegmentArtifact::Retained { index: 42 })
466 );
467 check!(
468 parse_segment_artifact("trace.5.bin.active", "trace") == Some(SegmentArtifact::Active)
469 );
470 check!(parse_segment_artifact("trace.bin", "trace").is_none());
471 check!(parse_segment_artifact("other.0.bin", "trace").is_none());
472 check!(parse_segment_artifact("trace.0.binx", "trace").is_none());
473 }
474
475 #[test]
476 #[cfg(feature = "pipeline")]
477 fn parse_segment_index_valid() {
478 check!(parse_segment_index("trace.0.bin", "trace") == Some(0));
479 check!(parse_segment_index("trace.42.bin", "trace") == Some(42));
480 check!(parse_segment_index("my-app.100.bin", "my-app") == Some(100));
481 }
482
483 #[test]
484 #[cfg(feature = "pipeline")]
485 fn parse_segment_index_invalid() {
486 check!(parse_segment_index("trace.0.bin.active", "trace") == None);
487 check!(parse_segment_index("trace.bin", "trace") == None);
488 check!(parse_segment_index("other.0.bin", "trace") == None);
489 check!(parse_segment_index("trace.abc.bin", "trace") == None);
490 }
491}