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}