use std::io::{self, Read, Write};
#[derive(Debug)]
pub enum Error {
Io(io::Error),
TruncatedFrame,
ChecksumMismatch,
DecompressionFailed,
}
impl From<io::Error> for Error {
fn from(e: io::Error) -> Self {
Error::Io(e)
}
}
impl std::fmt::Display for Error {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Error::Io(e) => write!(f, "io error: {e}"),
Error::TruncatedFrame => write!(f, "truncated frame at segment tail"),
Error::ChecksumMismatch => write!(f, "block checksum mismatch"),
Error::DecompressionFailed => write!(f, "block decompression failed"),
}
}
}
impl std::error::Error for Error {}
pub struct SegmentReader<R: Read> {
reader: R,
buffer: Vec<u8>,
}
impl<R: Read> SegmentReader<R> {
pub fn new(reader: R) -> Self {
Self {
reader,
buffer: Vec::new(),
}
}
pub fn next_record(&mut self) -> Result<Option<&[u8]>, Error> {
let mut len_buf = [0u8; 4];
match self.reader.read(&mut len_buf)? {
0 => return Ok(None),
n if n < 4 => return Err(Error::TruncatedFrame),
_ => {}
}
let len = u32::from_be_bytes(len_buf) as usize;
self.buffer.resize(len, 0);
self.reader.read_exact(&mut self.buffer).map_err(|e| {
if e.kind() == io::ErrorKind::UnexpectedEof {
Error::TruncatedFrame
} else {
Error::Io(e)
}
})?;
Ok(Some(&self.buffer))
}
}
pub struct SegmentWriter<W: Write> {
writer: W,
}
impl<W: Write> SegmentWriter<W> {
pub fn new(writer: W) -> Self {
Self { writer }
}
pub fn write(&mut self, record: &[u8]) -> io::Result<()> {
let len = record.len() as u32;
self.writer.write_all(&len.to_be_bytes())?;
self.writer.write_all(record)?;
Ok(())
}
pub fn flush(&mut self) -> io::Result<()> {
self.writer.flush()
}
}
#[cfg(feature = "harness")]
pub mod recipe;
#[cfg(any(
feature = "mmap",
feature = "crc32",
feature = "xxh3",
feature = "lz4",
feature = "seek-index",
feature = "wal-cursor",
))]
pub mod features;
#[cfg(feature = "crc32")]
pub use features::crc32::Crc32SegmentReader;
#[cfg(feature = "lz4")]
pub use features::lz4::{Lz4BlockWriter, Lz4SegmentReader};
#[cfg(feature = "mmap")]
pub use features::mmap::MmapSegmentReader;
#[cfg(feature = "seek-index")]
pub use features::seek_index::IndexedSegmentReader;
#[cfg(feature = "wal-cursor")]
pub use features::wal_cursor::WalCursorReader;
#[cfg(feature = "xxh3")]
pub use features::xxh3::Xxh3SegmentReader;