Skip to main content

moqtap_trace/
writer.rs

1use std::io::{BufWriter, Write};
2
3use ciborium::Value;
4
5use crate::error::MoqTraceError;
6use crate::event::TraceEvent;
7use crate::header::TraceHeader;
8
9/// Magic bytes identifying a `.moqtrace` file, and each segment within one.
10pub const MOQTRACE_MAGIC: &[u8; 8] = b"MOQTRACE";
11
12/// Format version this crate writes.
13pub const MOQTRACE_VERSION: u32 = 2;
14
15/// Format versions this crate can read.
16///
17/// Version 1 is a version 2 file that happens to carry none of the keys this
18/// revision adds and is never segmented, so reading it costs nothing and
19/// keeps every capture taken before the bump openable.
20pub const MOQTRACE_VERSIONS_SUPPORTED: &[u32] = &[1, MOQTRACE_VERSION];
21
22/// Streaming writer for `.moqtrace` files.
23///
24/// Writes a preamble (magic, version, header) on construction, then accepts
25/// events one at a time via [`write_event`](Self::write_event).
26///
27/// For a segmented trace — on-disk capture rotation, or a live stream carried
28/// over MoQT — call [`start_segment`](Self::start_segment) to begin a new
29/// segment with a fresh preamble of its own. A plain file is one segment.
30///
31/// The inner writer is wrapped in a [`BufWriter`] so events can be appended
32/// at line rate without incurring one syscall per event. Call
33/// [`flush`](Self::flush) or [`into_inner`](Self::into_inner) to drain the
34/// buffer.
35#[derive(Debug)]
36pub struct MoqTraceWriter<W: Write> {
37    inner: BufWriter<W>,
38}
39
40impl<W: Write> MoqTraceWriter<W> {
41    /// Create a new writer, writing the preamble and header of the file (or
42    /// of its first segment).
43    pub fn new(writer: W, header: &TraceHeader) -> Result<Self, MoqTraceError> {
44        let mut writer = BufWriter::new(writer);
45        write_preamble(&mut writer, header)?;
46        Ok(Self { inner: writer })
47    }
48
49    /// Append a single event to the current segment.
50    pub fn write_event(&mut self, event: &TraceEvent) -> Result<(), MoqTraceError> {
51        ciborium::into_writer(event, &mut self.inner)?;
52        Ok(())
53    }
54
55    /// Begin a new segment, writing a fresh preamble immediately after the
56    /// previous segment's last event.
57    ///
58    /// The header must carry [`segment`](TraceHeader::segment) metadata, and
59    /// this returns an error if it does not. That field is what tells a
60    /// reader the sequence numbers and timestamps it is about to see restart
61    /// at zero; without it a reader takes each segment for a complete file
62    /// and reconstructs a timeline that repeatedly jumps backwards.
63    ///
64    /// The new header should keep the same `protocol`, `perspective`,
65    /// `detail` and `session_id` as the segments before it, and increment
66    /// `segment.sequence`.
67    pub fn start_segment(&mut self, header: &TraceHeader) -> Result<(), MoqTraceError> {
68        if header.segment.is_none() {
69            return Err(MoqTraceError::InvalidHeader(
70                "a segment header must carry 'segment' metadata".into(),
71            ));
72        }
73        write_preamble(&mut self.inner, header)
74    }
75
76    /// Flush the underlying writer.
77    pub fn flush(&mut self) -> Result<(), MoqTraceError> {
78        self.inner.flush()?;
79        Ok(())
80    }
81
82    /// Consume the writer and return the inner writer.
83    ///
84    /// Flushes any buffered bytes. Returns an error if the flush fails.
85    pub fn into_inner(self) -> Result<W, MoqTraceError> {
86        self.inner.into_inner().map_err(|e| MoqTraceError::Io(e.into_error()))
87    }
88}
89
90fn write_preamble<W: Write>(writer: &mut W, header: &TraceHeader) -> Result<(), MoqTraceError> {
91    // 1. Magic bytes
92    writer.write_all(MOQTRACE_MAGIC)?;
93
94    // 2. Format version (u32 LE)
95    writer.write_all(&MOQTRACE_VERSION.to_le_bytes())?;
96
97    // 3. CBOR-encode header
98    let header_value: Value = header.into();
99    let mut header_bytes = Vec::with_capacity(128);
100    ciborium::into_writer(&header_value, &mut header_bytes)?;
101
102    // 4. Header length (u32 LE)
103    let header_len = header_bytes.len() as u32;
104    writer.write_all(&header_len.to_le_bytes())?;
105
106    // 5. Header CBOR bytes
107    writer.write_all(&header_bytes)?;
108
109    Ok(())
110}