Skip to main content

subms_segment_reader/
lib.rs

1//! Read length-prefix framed records from a segment file.
2//!
3//! Frame format (big-endian on disk; matches the LSM SSTable and Java side):
4//!
5//! ```text
6//! u32 length
7//! u8  payload[length]
8//! ```
9//!
10//! The reader streams record-by-record. Truncated tails surface as a typed
11//! `Error::TruncatedFrame` instead of crashing - the recipe sees crash-recovery
12//! workloads as common and the API treats them as expected.
13//!
14//! ```
15//! use subms_segment_reader::{SegmentWriter, SegmentReader};
16//! let mut buf = Vec::new();
17//! { let mut w = SegmentWriter::new(&mut buf); w.write(b"alice").unwrap(); w.write(b"bob").unwrap(); }
18//! let mut r = SegmentReader::new(buf.as_slice());
19//! assert_eq!(r.next_record().unwrap().unwrap(), b"alice");
20//! assert_eq!(r.next_record().unwrap().unwrap(), b"bob");
21//! assert!(r.next_record().unwrap().is_none());
22//! ```
23
24use std::io::{self, Read, Write};
25
26#[derive(Debug)]
27pub enum Error {
28    /// Underlying IO error.
29    Io(io::Error),
30    /// Header or payload truncated at the tail of the segment.
31    TruncatedFrame,
32    /// Block trailer's checksum did not match the payload's checksum.
33    /// Raised by the `crc32` and `xxh3` opt-in readers.
34    ChecksumMismatch,
35    /// Block's algorithm tag is unknown or a decompress call failed.
36    /// Raised by the `lz4` opt-in reader.
37    DecompressionFailed,
38}
39
40impl From<io::Error> for Error {
41    fn from(e: io::Error) -> Self {
42        Error::Io(e)
43    }
44}
45
46impl std::fmt::Display for Error {
47    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
48        match self {
49            Error::Io(e) => write!(f, "io error: {e}"),
50            Error::TruncatedFrame => write!(f, "truncated frame at segment tail"),
51            Error::ChecksumMismatch => write!(f, "block checksum mismatch"),
52            Error::DecompressionFailed => write!(f, "block decompression failed"),
53        }
54    }
55}
56
57impl std::error::Error for Error {}
58
59pub struct SegmentReader<R: Read> {
60    reader: R,
61    buffer: Vec<u8>,
62}
63
64impl<R: Read> SegmentReader<R> {
65    pub fn new(reader: R) -> Self {
66        Self {
67            reader,
68            buffer: Vec::new(),
69        }
70    }
71
72    /// Read the next record. Returns `Ok(None)` at clean EOF;
73    /// `Err(TruncatedFrame)` if the segment is cut in the middle of a frame.
74    pub fn next_record(&mut self) -> Result<Option<&[u8]>, Error> {
75        let mut len_buf = [0u8; 4];
76        match self.reader.read(&mut len_buf)? {
77            0 => return Ok(None),
78            n if n < 4 => return Err(Error::TruncatedFrame),
79            _ => {}
80        }
81        let len = u32::from_be_bytes(len_buf) as usize;
82        self.buffer.resize(len, 0);
83        self.reader.read_exact(&mut self.buffer).map_err(|e| {
84            if e.kind() == io::ErrorKind::UnexpectedEof {
85                Error::TruncatedFrame
86            } else {
87                Error::Io(e)
88            }
89        })?;
90        Ok(Some(&self.buffer))
91    }
92}
93
94pub struct SegmentWriter<W: Write> {
95    writer: W,
96}
97
98impl<W: Write> SegmentWriter<W> {
99    pub fn new(writer: W) -> Self {
100        Self { writer }
101    }
102
103    pub fn write(&mut self, record: &[u8]) -> io::Result<()> {
104        let len = record.len() as u32;
105        self.writer.write_all(&len.to_be_bytes())?;
106        self.writer.write_all(record)?;
107        Ok(())
108    }
109
110    pub fn flush(&mut self) -> io::Result<()> {
111        self.writer.flush()
112    }
113}
114
115#[cfg(feature = "harness")]
116pub mod recipe;
117
118// Opt-in feature catalog. Each submodule is gated by its own Cargo
119// feature flag and adds a capability without bloating the core build.
120// See `Cargo.toml` `[features]` for the catalog.
121#[cfg(any(
122    feature = "mmap",
123    feature = "crc32",
124    feature = "xxh3",
125    feature = "lz4",
126    feature = "seek-index",
127    feature = "wal-cursor",
128))]
129pub mod features;
130
131#[cfg(feature = "crc32")]
132pub use features::crc32::Crc32SegmentReader;
133#[cfg(feature = "lz4")]
134pub use features::lz4::{Lz4BlockWriter, Lz4SegmentReader};
135#[cfg(feature = "mmap")]
136pub use features::mmap::MmapSegmentReader;
137#[cfg(feature = "seek-index")]
138pub use features::seek_index::IndexedSegmentReader;
139#[cfg(feature = "wal-cursor")]
140pub use features::wal_cursor::WalCursorReader;
141#[cfg(feature = "xxh3")]
142pub use features::xxh3::Xxh3SegmentReader;