Skip to main content

leviath_core/run_archive/
codec.rs

1//! The on-disk shape of `run.lvr`: magic, version, framing, and the reader
2//! that tolerates what it does not understand.
3//!
4//! Separate from the rest of the archive module because it answers a different
5//! question. Everything else decides what a record *means* - how a delta folds,
6//! what a window looks like after it. This decides only where one record ends
7//! and the next begins, which is the part that has to stay stable while the
8//! record set keeps growing.
9
10use std::io::{self, Read, Write};
11
12use super::RunRecord;
13
14/// File magic identifying a leviath run archive (`b"LVR1"`).
15pub const RUN_ARCHIVE_MAGIC: &[u8; 4] = b"LVR1";
16
17/// The archive format version this build writes.
18///
19/// Bump this only for a change to the *framing* - the preamble, the length
20/// prefix, or the payload encoding. Adding a record kind is not that: frames
21/// are length-prefixed and a reader skips a payload it cannot parse, so a new
22/// kind is readable by an older build (it just does not know what it says).
23/// Adding a field to an existing record is not that either, as long as the
24/// field is `#[serde(default)]`.
25///
26/// What a bump means is that an older build cannot find the record boundaries
27/// at all, which is why [`read_archive_start`] refuses a newer archive rather
28/// than trying.
29pub const RUN_ARCHIVE_VERSION: u16 = 1;
30
31// ─── codec ──────────────────────────────────────────────────────────────────
32
33/// Write the archive preamble (magic + version). Call once at file start.
34pub fn write_archive_start(w: &mut dyn Write, version: u16) -> io::Result<()> {
35    w.write_all(RUN_ARCHIVE_MAGIC)?;
36    w.write_all(&version.to_be_bytes())?;
37    Ok(())
38}
39
40/// Read + validate the archive preamble, returning the format version.
41///
42/// A version *newer* than this build understands is refused here rather than
43/// read. The version marks a framing change (see [`RUN_ARCHIVE_VERSION`]), so a
44/// newer archive is one whose record boundaries this build cannot find - and
45/// reading it anyway would not fail cleanly, it would produce nonsense from
46/// whatever the length prefixes happened to say. An older version is read
47/// normally: framing has not changed under it, and unknown record kinds are
48/// skipped rather than fatal.
49pub fn read_archive_start(r: &mut dyn Read) -> io::Result<u16> {
50    let mut magic = [0u8; 4];
51    r.read_exact(&mut magic)?;
52    if &magic != RUN_ARCHIVE_MAGIC {
53        return Err(io::Error::new(
54            io::ErrorKind::InvalidData,
55            "not a leviath run archive (bad magic)",
56        ));
57    }
58    let mut version = [0u8; 2];
59    r.read_exact(&mut version)?;
60    let version = u16::from_be_bytes(version);
61    if version > RUN_ARCHIVE_VERSION {
62        return Err(io::Error::new(
63            io::ErrorKind::InvalidData,
64            format!(
65                "run archive is format version {version}, but this build reads up to \
66                 {RUN_ARCHIVE_VERSION} - upgrade leviath to read it"
67            ),
68        ));
69    }
70    Ok(version)
71}
72
73/// Append one framed record. The frame length is a `u64` so it can never
74/// overflow the prefix (a `RunRecord` always serializes to JSON).
75pub fn write_record(w: &mut dyn Write, record: &RunRecord) -> io::Result<()> {
76    let payload = serde_json::to_vec(record).expect("a RunRecord always serializes to JSON");
77    let len = payload.len() as u64;
78    w.write_all(&len.to_be_bytes())?;
79    w.write_all(&payload)?;
80    Ok(())
81}
82
83/// Fill `buf` from `r`, returning `false` on a clean end-of-stream (zero bytes
84/// available at the call) and erroring only on a *partial* read (truncation).
85fn read_exact_or_eof(r: &mut dyn Read, buf: &mut [u8]) -> io::Result<bool> {
86    let mut filled = 0;
87    while filled < buf.len() {
88        match r.read(&mut buf[filled..])? {
89            0 => {
90                if filled == 0 {
91                    return Ok(false); // clean EOF at a record boundary
92                }
93                return Err(io::Error::new(
94                    io::ErrorKind::UnexpectedEof,
95                    "truncated run-archive frame",
96                ));
97            }
98            n => filled += n,
99        }
100    }
101    Ok(true)
102}
103
104/// The largest a single archive frame may claim to be.
105///
106/// Generous by design - a record holds one context snapshot, and 256 MiB is far
107/// past anything a real run writes - because this is a sanity bound on a length
108/// prefix, not a size policy. What it rules out is a torn or corrupt prefix
109/// being taken at its word and turned straight into an allocation.
110const MAX_RECORD_BYTES: u64 = 256 * 1024 * 1024;
111
112/// One frame off the wire: a record this build understands, or a complete
113/// payload it does not.
114///
115/// The distinction is the whole point of the length prefix. A frame whose
116/// *content* is unreadable - a record kind added by a later build - is a
117/// well-formed frame whose bytes can be stepped over, and everything after it
118/// is still readable. A frame that is *torn* cannot be stepped over, because
119/// its length is unknown or its payload ran out, and the file ends there.
120///
121/// Conflating the two is what made adding a record kind a breaking change:
122/// `read_archive_lenient` stopped at the first unknown record and returned the
123/// prefix, silently dropping every readable record after it.
124#[derive(Debug, Clone, PartialEq)]
125pub enum Frame {
126    /// A record this build knows how to read.
127    Record(Box<RunRecord>),
128    /// A complete frame whose payload this build cannot parse, stepped over.
129    /// Carries its size so a caller can report what it skipped.
130    Unreadable {
131        /// Payload length in bytes.
132        bytes: usize,
133    },
134}
135
136/// Read the next frame, or `None` at a clean end-of-stream.
137///
138/// Errors only on a *torn* frame. An unparseable-but-complete payload comes
139/// back as [`Frame::Unreadable`] with the stream positioned after it, so
140/// reading can continue.
141pub fn read_frame(r: &mut dyn Read) -> io::Result<Option<Frame>> {
142    let Some(payload) = read_framed_payload(r)? else {
143        return Ok(None);
144    };
145    match serde_json::from_slice(&payload) {
146        Ok(record) => Ok(Some(Frame::Record(Box::new(record)))),
147        // Deliberately not an error: the frame was intact and has been
148        // consumed, so the only question is whether the caller wants to know.
149        Err(_) => Ok(Some(Frame::Unreadable {
150            bytes: payload.len(),
151        })),
152    }
153}
154
155/// Read the next framed record, or `None` at a clean end-of-stream.
156///
157/// Strict about content: a payload this build cannot parse is an error. Prefer
158/// [`read_frame`] anywhere an archive written by a *different* build might be
159/// read, which is every path that loads a run from disk.
160pub fn read_record(r: &mut dyn Read) -> io::Result<Option<RunRecord>> {
161    let Some(payload) = read_framed_payload(r)? else {
162        return Ok(None);
163    };
164    let record = serde_json::from_slice(&payload)
165        .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
166    Ok(Some(record))
167}
168
169/// Read one length-prefixed payload, or `None` at a clean end-of-stream.
170///
171/// The framing half of a frame read, shared so the strict and skipping readers
172/// cannot disagree about where a record ends.
173fn read_framed_payload(r: &mut dyn Read) -> io::Result<Option<Vec<u8>>> {
174    let mut len_bytes = [0u8; 8];
175    if !read_exact_or_eof(r, &mut len_bytes)? {
176        return Ok(None);
177    }
178    let len = u64::from_be_bytes(len_bytes);
179    // A torn tail is the reason `read_archive_lenient` exists, and a torn
180    // *length prefix* is exactly where a nonsense `u64` comes from. Allocating
181    // it first would abort the process on a crash-truncated archive - during
182    // daemon recovery, which is the one moment the lenient reader is there to
183    // survive. Rejecting it makes the frame an ordinary error, so recovery
184    // folds back to the last intact record instead.
185    if len > MAX_RECORD_BYTES {
186        return Err(io::Error::new(
187            io::ErrorKind::InvalidData,
188            format!("run-archive frame claims {len} bytes, over the {MAX_RECORD_BYTES} cap"),
189        ));
190    }
191    let mut payload = vec![0u8; len as usize];
192    if !read_exact_or_eof(r, &mut payload)? {
193        return Err(io::Error::new(
194            io::ErrorKind::UnexpectedEof,
195            "truncated run-archive frame",
196        ));
197    }
198    Ok(Some(payload))
199}
200
201/// Read the whole archive: validate the preamble, then read every record.
202pub fn read_archive(r: &mut dyn Read) -> io::Result<(u16, Vec<RunRecord>)> {
203    let version = read_archive_start(r)?;
204    let mut records = Vec::new();
205    while let Some(record) = read_record(r)? {
206        records.push(record);
207    }
208    Ok((version, records))
209}
210
211/// Read the archive tolerantly: validate the preamble strictly, then read records
212/// until a clean end-of-stream **or the first unreadable frame**, returning the
213/// records collected so far.
214///
215/// A crash while the persistence lane is appending a record can leave a partial
216/// final frame (a truncated length prefix or payload). The strict [`read_archive`]
217/// would reject the whole file for that torn tail - and once a fallback-resume
218/// appends fresh records *past* the torn bytes, the archive would stay unreadable
219/// forever. This variant instead stops at the torn tail and keeps everything valid
220/// before it, so recovery can still fold the archive to its last intact point. The
221/// preamble is still validated strictly, so a file that isn't a run archive at all
222/// still errors rather than folding to nothing.
223pub fn read_archive_lenient(r: &mut dyn Read) -> io::Result<(u16, Vec<RunRecord>)> {
224    let version = read_archive_start(r)?;
225    let mut records = Vec::new();
226    let mut skipped = 0usize;
227    // A torn frame ends the read with whatever preceded it. A frame this build
228    // simply does not understand is stepped over instead: it was written by a
229    // later version, and everything after it is still ours to read.
230    while let Ok(Some(frame)) = read_frame(r) {
231        match frame {
232            Frame::Record(record) => records.push(*record),
233            Frame::Unreadable { .. } => skipped += 1,
234        }
235    }
236    if skipped > 0 {
237        tracing::debug!(
238            skipped,
239            "run archive holds record kinds this build does not know; skipped them"
240        );
241    }
242    Ok((version, records))
243}