subms_segment_reader/
lib.rs1use std::io::{self, Read, Write};
25
26#[derive(Debug)]
27pub enum Error {
28 Io(io::Error),
30 TruncatedFrame,
32 ChecksumMismatch,
35 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 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#[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;