Skip to main content

dial9_core/
sealed.rs

1//! Sealed-file detection for the worker pipeline.
2//!
3//! Finds `.bin` files produced by `DiskBuffer` rename-on-seal,
4//! ignoring `.active` files that are still being written.
5
6// The segment types here are the public pipeline API, exposed via the
7// `pipeline` module's re-export. Without that feature they stay crate-internal.
8#![cfg_attr(not(feature = "pipeline"), allow(unreachable_pub))]
9
10use std::path::{Path, PathBuf};
11
12use crate::primitives::fs;
13
14/// A sealed trace segment ready for processing (disk-backed).
15#[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    /// Build a disk segment directly.
24    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    /// Path to the sealed segment file on disk.
32    pub fn path(&self) -> &Path {
33        &self.path
34    }
35
36    /// Segment index (e.g. `3` for `trace.3.bin`).
37    pub fn index(&self) -> u32 {
38        self.index
39    }
40}
41
42/// A sealed trace segment backed by in-process memory.
43#[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    /// Segment index.
52    pub fn index(&self) -> u32 {
53        self.index
54    }
55
56    /// Encoded segment size in bytes.
57    pub fn size(&self) -> u64 {
58        self.size
59    }
60}
61
62/// A sealed trace segment, either disk-backed or memory-backed.
63#[derive(Debug, Clone, PartialEq, Eq)]
64#[non_exhaustive]
65pub enum SegmentRef {
66    /// On-disk segment file
67    Disk(SealedSegment),
68    /// In-process segment
69    Memory(MemorySegment),
70}
71
72impl SegmentRef {
73    /// Segment index.
74    pub fn index(&self) -> u32 {
75        match self {
76            SegmentRef::Disk(s) => s.index,
77            SegmentRef::Memory(m) => m.index,
78        }
79    }
80
81    /// Returns the on-disk path for disk-backed segments.
82    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
99/// Segment creation time as epoch seconds, parsed from the first clock
100/// anchor in the trace. Returns `(secs, true)` on a successful parse, or
101/// falls back to file mtime / current time with `(secs, false)`.
102pub(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/// Best-effort seal epoch for a disk segment: the file's mtime (the last
119/// write before seal), falling back to now. Together with
120/// [`creation_epoch_secs`] it gives the span the triggered worker matches
121/// against dump windows.
122#[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
141/// Legacy epoch-vs-monotonic discriminator for ambiguous
142/// `SegmentMetadataEvent.timestamp_ns` values.
143///
144/// Pre-clock-sync traces may store epoch nanoseconds there, while new traces
145/// use monotonic nanoseconds. Epoch timestamps are ~1e18, while monotonic uptime is much smaller.
146/// So we use `2020-01-01` ns as a floor, monotonic would need decades
147/// of continuous runtime to reach that range.
148///
149/// Keep in sync with `LEGACY_EPOCH_NS_FLOOR` in `dial9-viewer/ui/trace_parser.js`.
150pub(crate) const LEGACY_EPOCH_NS_FLOOR: u64 = 1_577_836_800_000_000_000;
151
152/// Parse wall-clock creation time from the first `ClockSyncEvent`,
153/// or from a legacy `SegmentMetadataEvent.timestamp_ns` that predates clock-sync support.
154fn 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(_)) => {} // schema/pool frame, keep going
189            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/// Find sealed `.bin` segments in `dir`, sorted oldest-first by index.
229///
230/// Matches files named `{stem}.{index}.bin` where `stem` matches the
231/// given base path's file stem. Ignores `.active` files and any files
232/// that don't match the expected naming pattern.
233#[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        // Guard against .bin.active being misread — those have extension "active"
243        // so the check above already excludes them, but be explicit.
244        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/// Trace-segment artifact found in the trace directory: a retained segment
257/// (`.bin` or write-back sibling like `.bin.gz`) or a stale `.bin.active`.
258#[derive(Debug, Clone, Copy, PartialEq, Eq)]
259pub(crate) enum SegmentArtifact {
260    Retained { index: u32 },
261    Active,
262}
263
264/// Classify a filename against the `{stem}.{index}.bin*` family.
265pub(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/// Parse segment index from a filename like `trace.3.bin`.
282/// Returns `None` if the filename doesn't match `{stem}.{index}.bin`.
283#[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// These tests cover the worker-facing segment-discovery path.
291#[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    /// Build a single-`ClockSyncEvent` trace carrying `realtime_ns` and parse it
350    /// back. Mirrors the anchor the writer emits at the start of every segment.
351    #[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        // Header parse succeeds, so the path is never stat'd.
384        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        // Should fall back to mtime, which should be recent
400        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    /// A legacy pre-clock-sync trace: wall clock lives in
408    /// `SegmentMetadataEvent.timestamp_ns`, no `ClockSyncEvent`.
409    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    /// Legacy pre-clock-sync files stored wall clock in `SegmentMetadataEvent.timestamp_ns`,
420    /// parser fallback should recover it.
421    #[test]
422    fn test_parse_segment_timestamp_legacy_segment_metadata_fallback() {
423        // 2024-06-01T00:00:00Z: comfortably epoch-scale.
424        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    /// A small (monotonic-looking) SegmentMeta.timestamp_ns must NOT be
444    /// mistaken for a legacy wall-clock value.
445    #[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}