Skip to main content

vole_document/container/
record.rs

1//! Length-delimited record framing.
2//!
3//! ```text
4//! record := tag:u8 flags:u8 reserved:u16=0 length:u32 payload:[u8;length] crc32c:u32
5//! ```
6//!
7//! The CRC covers the 8-byte record header and the payload. `reserved` must be
8//! zero. The `flags` bit [`FLAG_OPTIONAL`] marks a record whose *tag* a decoder
9//! may skip if unknown; unknown non-optional tags fail closed.
10
11use std::io::{Read, Seek, SeekFrom};
12
13use crate::error::{Error, Result};
14use crate::integrity::crc32c;
15use crate::limits::Limits;
16
17/// Bytes of per-record framing before the payload.
18pub const RECORD_HEADER_LEN: usize = 8;
19/// Bytes of per-record framing after the payload (the CRC).
20pub const RECORD_TRAILER_LEN: usize = 4;
21/// Total framing overhead per record.
22pub const RECORD_OVERHEAD: usize = RECORD_HEADER_LEN + RECORD_TRAILER_LEN;
23
24/// Flag: this record's tag may be skipped by a decoder that does not know it.
25pub const FLAG_OPTIONAL: u8 = 0x01;
26
27/// Known record classes.
28#[derive(Debug, Clone, Copy, PartialEq, Eq)]
29#[repr(u8)]
30pub enum RecordTag {
31    /// Universe declaration (opcode/coder/limit semantics version).
32    Universe = 0x01,
33    /// Source-format descriptor.
34    Format = 0x02,
35    /// A raw byte object.
36    Object = 0x10,
37    /// The reconstruction graph (DRA program).
38    Graph = 0x20,
39    /// An entropy model (Phase 2+).
40    Model = 0x30,
41    /// An entropy channel capsule (Phase 2+).
42    EntropyChannel = 0x40,
43    /// A typed residual stream (later phases).
44    Residual = 0x50,
45    /// A random-access checkpoint (later phases).
46    Checkpoint = 0x60,
47    /// An optional observation index (Phase 7): op/selector/digest map.
48    ///
49    /// Reuses the reserved `0x70` slot. It is written with [`FLAG_OPTIONAL`]; a
50    /// decoder that does not implement partial decode skips it and still fully
51    /// materializes.
52    ObservationIndex = 0x70,
53    /// An optional seek directory (Phase 8): a bounded map from record class to
54    /// on-disk locator, written as the first record at fixed offset 64.
55    ///
56    /// It is written with [`FLAG_OPTIONAL`]; a decoder that does not implement
57    /// seek-based partial I/O skips it (the `None`/unknown arm in
58    /// `Descriptor::parse`) and still fully materializes the source, because the
59    /// reconstruction program alone is complete.
60    Directory = 0x71,
61    /// A reference to an external content-addressed object (Phase 9+).
62    ExternalRef = 0x80,
63    /// Whole-source integrity manifest.
64    Integrity = 0xF0,
65    /// Terminal record.
66    Trailer = 0xFF,
67}
68
69impl RecordTag {
70    /// Map a raw tag byte to a known tag, if any.
71    pub const fn from_u8(b: u8) -> Option<RecordTag> {
72        match b {
73            0x01 => Some(RecordTag::Universe),
74            0x02 => Some(RecordTag::Format),
75            0x10 => Some(RecordTag::Object),
76            0x20 => Some(RecordTag::Graph),
77            0x30 => Some(RecordTag::Model),
78            0x40 => Some(RecordTag::EntropyChannel),
79            0x50 => Some(RecordTag::Residual),
80            0x60 => Some(RecordTag::Checkpoint),
81            0x70 => Some(RecordTag::ObservationIndex),
82            0x71 => Some(RecordTag::Directory),
83            0x80 => Some(RecordTag::ExternalRef),
84            0xF0 => Some(RecordTag::Integrity),
85            0xFF => Some(RecordTag::Trailer),
86            _ => None,
87        }
88    }
89
90    /// Stable short name for diagnostics and receipts.
91    pub const fn name(self) -> &'static str {
92        match self {
93            RecordTag::Universe => "UNIVERSE",
94            RecordTag::Format => "FORMAT",
95            RecordTag::Object => "OBJECT",
96            RecordTag::Graph => "GRAPH",
97            RecordTag::Model => "MODEL",
98            RecordTag::EntropyChannel => "ENTROPY_CHANNEL",
99            RecordTag::Residual => "RESIDUAL",
100            RecordTag::Checkpoint => "CHECKPOINT",
101            RecordTag::ObservationIndex => "OBSERVATION_INDEX",
102            RecordTag::Directory => "DIRECTORY",
103            RecordTag::ExternalRef => "EXTERNAL_REF",
104            RecordTag::Integrity => "INTEGRITY",
105            RecordTag::Trailer => "TRAILER",
106        }
107    }
108}
109
110/// A single decoded record.
111#[derive(Debug, Clone, PartialEq, Eq)]
112pub struct Record {
113    /// Raw tag byte (kept raw so unknown optional tags survive a reparse).
114    pub tag: u8,
115    /// Raw flags byte.
116    pub flags: u8,
117    /// Payload bytes.
118    pub payload: Vec<u8>,
119}
120
121impl Record {
122    /// Construct a record with the given tag byte and payload.
123    pub fn new(tag: u8, payload: Vec<u8>) -> Self {
124        Record {
125            tag,
126            flags: 0,
127            payload,
128        }
129    }
130
131    /// Construct a record for a known tag.
132    pub fn known(tag: RecordTag, payload: Vec<u8>) -> Self {
133        Record {
134            tag: tag as u8,
135            flags: 0,
136            payload,
137        }
138    }
139
140    /// Whether the optional-skip flag is set.
141    pub fn is_optional(&self) -> bool {
142        self.flags & FLAG_OPTIONAL != 0
143    }
144}
145
146/// Append an encoded record to `out`. `payload` must fit in `u32`.
147pub fn write_record(out: &mut Vec<u8>, tag: u8, flags: u8, payload: &[u8]) -> Result<()> {
148    let len = u32::try_from(payload.len())
149        .map_err(|_| Error::resource_limit("record payload exceeds 4 GiB"))?;
150    let start = out.len();
151    out.push(tag);
152    out.push(flags);
153    out.extend_from_slice(&0u16.to_le_bytes());
154    out.extend_from_slice(&len.to_le_bytes());
155    out.extend_from_slice(payload);
156    let crc = crc32c(&out[start..]);
157    out.extend_from_slice(&crc.to_le_bytes());
158    Ok(())
159}
160
161/// Read and CRC-check exactly one record located at `offset` in a seekable
162/// source.
163///
164/// This is the seek-reader counterpart of
165/// [`RecordReader::next_record`]: it applies the *same* framing checks (reserved
166/// must be zero, `len <= max_record_len`, CRC32C over the 8-byte header plus the
167/// payload) but seeks to an absolute offset instead of walking sequentially, so a
168/// caller reads only the records a query needs. It returns the fully materialized
169/// [`Record`]; the tag byte, flags byte, and payload length are the framing's own
170/// claim and must still be cross-checked by the caller against any directory
171/// locator.
172pub fn read_record_at<R: Read + Seek>(
173    reader: &mut R,
174    offset: u64,
175    limits: Limits,
176) -> Result<Record> {
177    reader
178        .seek(SeekFrom::Start(offset))
179        .map_err(|e| Error::invalid_container(format!("seek to record at {offset} failed: {e}")))?;
180    let mut hdr = [0u8; RECORD_HEADER_LEN];
181    reader.read_exact(&mut hdr).map_err(|_| {
182        Error::invalid_container(format!("truncated record header at offset {offset}"))
183    })?;
184    let tag = hdr[0];
185    let flags = hdr[1];
186    let reserved = u16::from_le_bytes([hdr[2], hdr[3]]);
187    if reserved != 0 {
188        return Err(Error::invalid_container(
189            "record reserved field must be zero",
190        ));
191    }
192    let len = u32::from_le_bytes([hdr[4], hdr[5], hdr[6], hdr[7]]);
193    if len > limits.max_record_len {
194        return Err(Error::resource_limit(format!(
195            "record payload length {len} exceeds limit {}",
196            limits.max_record_len
197        )));
198    }
199    let body_len = RECORD_HEADER_LEN
200        .checked_add(len as usize)
201        .ok_or_else(|| Error::invalid_container("record length overflow"))?;
202    let mut body = vec![0u8; body_len];
203    body[..RECORD_HEADER_LEN].copy_from_slice(&hdr);
204    reader
205        .read_exact(&mut body[RECORD_HEADER_LEN..])
206        .map_err(|_| {
207            Error::invalid_container(format!("truncated record payload at offset {offset}"))
208        })?;
209    let mut crc = [0u8; RECORD_TRAILER_LEN];
210    reader.read_exact(&mut crc).map_err(|_| {
211        Error::invalid_container(format!("truncated record CRC at offset {offset}"))
212    })?;
213    let want_crc = u32::from_le_bytes(crc);
214    let got_crc = crc32c(&body);
215    if want_crc != got_crc {
216        return Err(Error::invalid_container(format!(
217            "record CRC32C mismatch at offset {offset}: declared {want_crc:#010x}, computed {got_crc:#010x}"
218        )));
219    }
220    let payload = body[RECORD_HEADER_LEN..].to_vec();
221    Ok(Record {
222        tag,
223        flags,
224        payload,
225    })
226}
227
228/// An iterator over the records of a byte slice.
229pub struct RecordReader<'a> {
230    data: &'a [u8],
231    pos: usize,
232    limits: Limits,
233    count: u32,
234}
235
236impl<'a> RecordReader<'a> {
237    /// Create a reader starting at `offset` within `data`.
238    pub fn new(data: &'a [u8], offset: usize, limits: Limits) -> Self {
239        RecordReader {
240            data,
241            pos: offset,
242            limits,
243            count: 0,
244        }
245    }
246
247    /// Current byte offset of the next record header.
248    pub fn position(&self) -> usize {
249        self.pos
250    }
251
252    /// Records successfully read so far.
253    pub fn records_read(&self) -> u32 {
254        self.count
255    }
256
257    /// Read the next record, or `None` at end of input.
258    pub fn next_record(&mut self) -> Result<Option<Record>> {
259        if self.pos == self.data.len() {
260            return Ok(None);
261        }
262        if self.count >= self.limits.max_record_count {
263            return Err(Error::resource_limit("record count limit exceeded"));
264        }
265        let remaining = self.data.len() - self.pos;
266        if remaining < RECORD_OVERHEAD {
267            return Err(Error::invalid_container(format!(
268                "truncated record header: {remaining} bytes remain"
269            )));
270        }
271        let hdr = &self.data[self.pos..self.pos + RECORD_HEADER_LEN];
272        let tag = hdr[0];
273        let flags = hdr[1];
274        let reserved = u16::from_le_bytes([hdr[2], hdr[3]]);
275        if reserved != 0 {
276            return Err(Error::invalid_container(
277                "record reserved field must be zero",
278            ));
279        }
280        let len = u32::from_le_bytes([hdr[4], hdr[5], hdr[6], hdr[7]]);
281        if len > self.limits.max_record_len {
282            return Err(Error::resource_limit(format!(
283                "record payload length {len} exceeds limit {}",
284                self.limits.max_record_len
285            )));
286        }
287        let total = RECORD_HEADER_LEN
288            .checked_add(len as usize)
289            .and_then(|n| n.checked_add(RECORD_TRAILER_LEN))
290            .ok_or_else(|| Error::invalid_container("record length overflow"))?;
291        if total > remaining {
292            return Err(Error::invalid_container(format!(
293                "truncated record payload: need {total}, have {remaining}"
294            )));
295        }
296        let body = &self.data[self.pos..self.pos + RECORD_HEADER_LEN + len as usize];
297        let want_crc = u32::from_le_bytes([
298            self.data[self.pos + RECORD_HEADER_LEN + len as usize],
299            self.data[self.pos + RECORD_HEADER_LEN + len as usize + 1],
300            self.data[self.pos + RECORD_HEADER_LEN + len as usize + 2],
301            self.data[self.pos + RECORD_HEADER_LEN + len as usize + 3],
302        ]);
303        let got_crc = crc32c(body);
304        if want_crc != got_crc {
305            return Err(Error::invalid_container(format!(
306                "record CRC32C mismatch at offset {}: declared {want_crc:#010x}, computed {got_crc:#010x}",
307                self.pos
308            )));
309        }
310        let payload = self.data
311            [self.pos + RECORD_HEADER_LEN..self.pos + RECORD_HEADER_LEN + len as usize]
312            .to_vec();
313        self.pos += total;
314        self.count += 1;
315        Ok(Some(Record {
316            tag,
317            flags,
318            payload,
319        }))
320    }
321}
322
323#[cfg(test)]
324mod tests {
325    use super::*;
326
327    #[test]
328    fn write_read_roundtrip() {
329        let mut buf = Vec::new();
330        write_record(&mut buf, RecordTag::Universe as u8, 0, b"hello").unwrap();
331        write_record(&mut buf, RecordTag::Object as u8, 0, &[1, 2, 3]).unwrap();
332        let mut r = RecordReader::new(&buf, 0, Limits::DEFAULT);
333        let a = r.next_record().unwrap().unwrap();
334        assert_eq!(a.tag, RecordTag::Universe as u8);
335        assert_eq!(a.payload, b"hello");
336        let b = r.next_record().unwrap().unwrap();
337        assert_eq!(b.tag, RecordTag::Object as u8);
338        assert_eq!(b.payload, vec![1, 2, 3]);
339        assert!(r.next_record().unwrap().is_none());
340    }
341
342    #[test]
343    fn detects_payload_corruption() {
344        let mut buf = Vec::new();
345        write_record(&mut buf, 0x10, 0, b"abcdef").unwrap();
346        buf[9] ^= 0xFF; // corrupt a payload byte
347        let mut r = RecordReader::new(&buf, 0, Limits::DEFAULT);
348        let e = r.next_record().unwrap_err();
349        assert_eq!(e.class(), crate::ErrorClass::InvalidContainer);
350    }
351
352    #[test]
353    fn detects_truncation() {
354        let mut buf = Vec::new();
355        write_record(&mut buf, 0x10, 0, b"abcdef").unwrap();
356        buf.truncate(buf.len() - 2);
357        let mut r = RecordReader::new(&buf, 0, Limits::DEFAULT);
358        let e = r.next_record().unwrap_err();
359        assert_eq!(e.class(), crate::ErrorClass::InvalidContainer);
360    }
361
362    #[test]
363    fn rejects_nonzero_reserved() {
364        let mut buf = Vec::new();
365        write_record(&mut buf, 0x10, 0, b"x").unwrap();
366        buf[2] = 1; // reserved
367        let mut r = RecordReader::new(&buf, 0, Limits::DEFAULT);
368        let e = r.next_record().unwrap_err();
369        assert_eq!(e.class(), crate::ErrorClass::InvalidContainer);
370    }
371
372    #[test]
373    fn enforces_record_length_limit() {
374        let mut buf = Vec::new();
375        write_record(&mut buf, 0x10, 0, &[0u8; 100]).unwrap();
376        let limits = Limits {
377            max_record_len: 10,
378            ..Limits::DEFAULT
379        };
380        let mut r = RecordReader::new(&buf, 0, limits);
381        let e = r.next_record().unwrap_err();
382        assert_eq!(e.class(), crate::ErrorClass::ResourceLimit);
383    }
384}