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}