use std::collections::{HashMap, HashSet};
use std::fs::{File, OpenOptions};
use std::io::{Read, Seek, SeekFrom, Write};
use std::path::{Path, PathBuf};
use crate::encode::{encode, SchemaMode};
use crate::error::{Error, Result};
use crate::hash::crc32;
use crate::message::Message;
use crate::registry::SchemaRegistry;
use crate::resolve::Resolver;
use crate::schema::Schema;
use crate::value::Value;
pub const FILE_MAGIC: &[u8; 4] = b"VRTF";
pub const FILE_VERSION: u8 = 1;
pub const FILE_HEADER_LEN: usize = 32;
pub const FOOTER_LEN: usize = 64;
pub const INDEX_ENTRY_LEN: usize = 40;
pub const FIRST_RECORD_ID: u64 = 1;
pub const OPT_RECORD_CRC: u32 = 1;
const ALIGN: u64 = 8;
#[inline]
fn align_up(x: u64) -> Result<u64> {
x.checked_add(ALIGN - 1)
.map(|v| v & !(ALIGN - 1))
.ok_or(Error::BadFile("offset overflow"))
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct Record {
pub id: u64,
pub offset: u64,
pub length: u64,
pub schema_id: u128,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct Footer {
pub generation: u64,
pub index_offset: u64,
pub schema_offset: u64,
pub schema_len: u32,
pub record_count: u32,
pub file_len: u64,
pub next_record_id: u64,
}
impl Footer {
fn encode(&self) -> [u8; FOOTER_LEN] {
let mut out = [0u8; FOOTER_LEN];
out[0..8].copy_from_slice(&self.generation.to_le_bytes());
out[8..16].copy_from_slice(&self.index_offset.to_le_bytes());
out[16..24].copy_from_slice(&self.schema_offset.to_le_bytes());
out[24..28].copy_from_slice(&self.schema_len.to_le_bytes());
out[28..32].copy_from_slice(&self.record_count.to_le_bytes());
out[32..40].copy_from_slice(&self.file_len.to_le_bytes());
out[40..48].copy_from_slice(&self.next_record_id.to_le_bytes());
let crc = crc32(&out[0..56]);
out[56..60].copy_from_slice(&crc.to_le_bytes());
out[60..64].copy_from_slice(FILE_MAGIC);
out
}
fn decode_at(bytes: &[u8; FOOTER_LEN], at: u64) -> Option<Footer> {
if &bytes[60..64] != FILE_MAGIC {
return None;
}
if u32::from_le_bytes(bytes[56..60].try_into().ok()?) != crc32(&bytes[0..56]) {
return None;
}
if u64::from_le_bytes(bytes[48..56].try_into().ok()?) != 0 {
return None;
}
let file_len = u64::from_le_bytes(bytes[32..40].try_into().ok()?);
if file_len != at.checked_add(FOOTER_LEN as u64)? {
return None;
}
Some(Footer {
generation: u64::from_le_bytes(bytes[0..8].try_into().ok()?),
index_offset: u64::from_le_bytes(bytes[8..16].try_into().ok()?),
schema_offset: u64::from_le_bytes(bytes[16..24].try_into().ok()?),
schema_len: u32::from_le_bytes(bytes[24..28].try_into().ok()?),
record_count: u32::from_le_bytes(bytes[28..32].try_into().ok()?),
file_len,
next_record_id: u64::from_le_bytes(bytes[40..48].try_into().ok()?),
})
}
}
fn header_bytes(optional_features: u32) -> [u8; FILE_HEADER_LEN] {
let mut out = [0u8; FILE_HEADER_LEN];
out[0..4].copy_from_slice(FILE_MAGIC);
out[4] = FILE_VERSION;
out[12..16].copy_from_slice(&optional_features.to_le_bytes());
out
}
fn check_header(buf: &[u8]) -> Result<()> {
if buf.len() < FILE_HEADER_LEN {
return Err(Error::Truncated);
}
if &buf[0..4] != FILE_MAGIC {
return Err(Error::BadFile("bad magic (not a .verit file)"));
}
if buf[4] != FILE_VERSION {
return Err(Error::BadFile("unsupported .verit file version"));
}
if buf[5] != 0 || u16::from_le_bytes(buf[6..8].try_into().unwrap()) != 0 {
return Err(Error::BadFile("nonzero reserved header field"));
}
let required = u32::from_le_bytes(buf[8..12].try_into().unwrap());
if required != 0 {
return Err(Error::UnsupportedFileFeature(required));
}
if u64::from_le_bytes(buf[16..24].try_into().unwrap()) != 0
|| u64::from_le_bytes(buf[24..32].try_into().unwrap()) != 0
{
return Err(Error::BadFile("nonzero reserved header field"));
}
Ok(())
}
fn commit_tail(
index: &[Record],
registry: &SchemaRegistry,
tail_start: u64,
generation: u64,
next_record_id: u64,
checksums: Option<&[u32]>,
) -> Result<(Vec<u8>, Footer)> {
let count = u32::try_from(index.len()).map_err(|_| Error::BadFile("too many records"))?;
let live: HashSet<u128> = index.iter().map(|r| r.schema_id).collect();
let mut pruned = SchemaRegistry::new();
for id in live {
let schema = registry.get(id).ok_or(Error::MissingSchema(id))?;
pruned.register(schema.clone());
}
let section = pruned.to_bundle();
let schema_len =
u32::try_from(section.len()).map_err(|_| Error::BadFile("schema section too large"))?;
let schema_offset = tail_start;
let crc_bytes = match checksums {
Some(c) => {
debug_assert_eq!(c.len(), index.len());
(c.len() as u64) * 4
}
None => 0,
};
let index_offset = align_up(schema_offset + section.len() as u64 + crc_bytes)?;
let index_bytes = (count as u64)
.checked_mul(INDEX_ENTRY_LEN as u64)
.ok_or(Error::BadFile("index size overflow"))?;
let file_len = index_offset
.checked_add(index_bytes)
.and_then(|v| v.checked_add(FOOTER_LEN as u64))
.ok_or(Error::BadFile("file length overflow"))?;
let mut tail = Vec::with_capacity((file_len - tail_start) as usize);
tail.extend_from_slice(§ion);
if let Some(c) = checksums {
tail.resize((index_offset - schema_offset) as usize - c.len() * 4, 0);
for crc in c {
tail.extend_from_slice(&crc.to_le_bytes());
}
}
tail.resize((index_offset - schema_offset) as usize, 0); for r in index {
tail.extend_from_slice(&r.id.to_le_bytes());
tail.extend_from_slice(&r.offset.to_le_bytes());
tail.extend_from_slice(&r.length.to_le_bytes());
tail.extend_from_slice(&r.schema_id.to_le_bytes());
}
Ok((
tail,
Footer {
generation,
index_offset,
schema_offset,
schema_len,
record_count: count,
file_len,
next_record_id,
},
))
}
fn schema_of(bytes: &[u8]) -> Result<Schema> {
Message::parse(bytes)?
.writer_schema()?
.ok_or(Error::NoInlineSchema)
}
#[derive(Clone, Debug)]
pub struct FileView<'a> {
buf: &'a [u8],
footer: Footer,
registry: SchemaRegistry,
}
impl<'a> FileView<'a> {
pub fn open(buf: &'a [u8]) -> Result<FileView<'a>> {
check_header(buf)?;
if buf.len() < FILE_HEADER_LEN + FOOTER_LEN {
return Err(Error::Truncated);
}
let footer = Self::find_footer(buf)?;
if footer.file_len > buf.len() as u64 {
return Err(Error::BadFile(
"footer claims more bytes than the image has",
));
}
if footer.schema_offset % ALIGN != 0 || footer.index_offset % ALIGN != 0 {
return Err(Error::BadFile("misaligned schema or index offset"));
}
if footer.schema_offset < FILE_HEADER_LEN as u64 {
return Err(Error::BadFile("schema section overlaps the header"));
}
let schema_end = footer
.schema_offset
.checked_add(footer.schema_len as u64)
.ok_or(Error::BadFile("schema section extent overflow"))?;
if schema_end > footer.index_offset {
return Err(Error::BadFile("schema section overlaps the index"));
}
if u32::from_le_bytes(buf[12..16].try_into().unwrap()) & OPT_RECORD_CRC != 0 {
let crc_bytes = (footer.record_count as u64)
.checked_mul(4)
.ok_or(Error::BadFile("checksum array size overflow"))?;
if footer.index_offset < crc_bytes || footer.index_offset - crc_bytes < schema_end {
return Err(Error::BadFile("checksum array overlaps the schema section"));
}
}
let index_bytes = (footer.record_count as u64)
.checked_mul(INDEX_ENTRY_LEN as u64)
.ok_or(Error::BadFile("index size overflow"))?;
let index_end = footer
.index_offset
.checked_add(index_bytes)
.ok_or(Error::BadFile("index extent overflow"))?;
if index_end != footer.file_len - FOOTER_LEN as u64 {
return Err(Error::BadFile("index does not end at the footer"));
}
let registry =
SchemaRegistry::from_bundle(&buf[footer.schema_offset as usize..schema_end as usize])?;
let view = FileView {
buf,
footer,
registry,
};
let mut prev_id = 0u64;
for i in 0..view.len() {
let r = view.record(i)?;
if r.id == 0 {
return Err(Error::BadFile("record id 0 is reserved"));
}
if r.id <= prev_id {
return Err(Error::BadFile("record ids are not strictly ascending"));
}
prev_id = r.id;
if r.offset % ALIGN != 0 {
return Err(Error::BadFile("misaligned record"));
}
if r.offset < FILE_HEADER_LEN as u64 {
return Err(Error::BadFile("record overlaps the header"));
}
let end = r
.offset
.checked_add(r.length)
.ok_or(Error::BadFile("record extent overflow"))?;
if end > footer.schema_offset {
return Err(Error::BadFile("record outside the record region"));
}
if !view.registry.contains(r.schema_id) {
return Err(Error::MissingSchema(r.schema_id));
}
}
if footer.next_record_id <= prev_id {
return Err(Error::BadFile(
"next_record_id does not exceed every record id",
));
}
Ok(view)
}
fn find_footer(buf: &[u8]) -> Result<Footer> {
let mut pos = ((buf.len() - FOOTER_LEN) as u64) & !(ALIGN - 1);
loop {
let at = pos as usize;
let bytes: [u8; FOOTER_LEN] = buf[at..at + FOOTER_LEN]
.try_into()
.map_err(|_| Error::Internal("footer slice width"))?;
if let Some(f) = Footer::decode_at(&bytes, pos) {
return Ok(f);
}
if pos < FILE_HEADER_LEN as u64 + ALIGN {
return Err(Error::NoValidFooter);
}
pos -= ALIGN;
}
}
pub fn len(&self) -> usize {
self.footer.record_count as usize
}
pub fn is_empty(&self) -> bool {
self.footer.record_count == 0
}
pub fn generation(&self) -> u64 {
self.footer.generation
}
pub fn footer(&self) -> Footer {
self.footer
}
pub fn file_len(&self) -> u64 {
self.footer.file_len
}
pub fn next_record_id(&self) -> u64 {
self.footer.next_record_id
}
pub fn has_record_checksums(&self) -> bool {
u32::from_le_bytes(self.buf[12..16].try_into().unwrap()) & OPT_RECORD_CRC != 0
}
pub fn record_checksum(&self, i: usize) -> Option<u32> {
if !self.has_record_checksums() || i >= self.len() {
return None;
}
let base = self.footer.index_offset as usize - self.len() * 4 + i * 4;
Some(u32::from_le_bytes(
self.buf[base..base + 4].try_into().ok()?,
))
}
pub fn verify_checksums(&self) -> Result<usize> {
if !self.has_record_checksums() {
return Ok(0);
}
for i in 0..self.len() {
let want = self.record_checksum(i).ok_or(Error::BadFile(
"checksum array is shorter than the record count",
))?;
let got = crc32(self.get(i)?);
if got != want {
return Err(Error::ChecksumMismatch {
id: self.record(i)?.id,
expected: want,
found: got,
});
}
}
Ok(self.len())
}
pub fn schemas(&self) -> &SchemaRegistry {
&self.registry
}
pub fn record(&self, i: usize) -> Result<Record> {
if i >= self.len() {
return Err(Error::IndexOutOfBounds);
}
let base = self.footer.index_offset as usize + i * INDEX_ENTRY_LEN;
Ok(Record {
id: u64::from_le_bytes(self.buf[base..base + 8].try_into().unwrap()),
offset: u64::from_le_bytes(self.buf[base + 8..base + 16].try_into().unwrap()),
length: u64::from_le_bytes(self.buf[base + 16..base + 24].try_into().unwrap()),
schema_id: u128::from_le_bytes(self.buf[base + 24..base + 40].try_into().unwrap()),
})
}
pub fn get(&self, i: usize) -> Result<&'a [u8]> {
let r = self.record(i)?;
Ok(&self.buf[r.offset as usize..(r.offset + r.length) as usize])
}
pub fn find_by_id(&self, id: u64) -> Option<usize> {
let (mut lo, mut hi) = (0usize, self.len());
while lo < hi {
let mid = lo + (hi - lo) / 2;
let mid_id = self.record(mid).ok()?.id;
match mid_id.cmp(&id) {
std::cmp::Ordering::Equal => return Some(mid),
std::cmp::Ordering::Less => lo = mid + 1,
std::cmp::Ordering::Greater => hi = mid,
}
}
None
}
pub fn get_by_id(&self, id: u64) -> Result<&'a [u8]> {
self.get(self.find_by_id(id).ok_or(Error::IndexOutOfBounds)?)
}
pub fn records_after(&self, id: u64) -> impl Iterator<Item = Record> + '_ {
let (mut lo, mut hi) = (0usize, self.len());
while lo < hi {
let mid = lo + (hi - lo) / 2;
match self.record(mid) {
Ok(r) if r.id <= id => lo = mid + 1,
_ => hi = mid,
}
}
(lo..self.len()).map(move |i| self.record(i).expect("index validated in open"))
}
pub fn schema_id(&self, i: usize) -> Result<u128> {
Ok(self.record(i)?.schema_id)
}
pub fn schema(&self, i: usize) -> Result<&Schema> {
let id = self.schema_id(i)?;
self.registry.get(id).ok_or(Error::MissingSchema(id))
}
pub fn message(&self, i: usize) -> Result<Message<'a>> {
Message::parse(self.get(i)?)
}
pub fn resolver_for(&self, i: usize, reader: &Schema) -> Result<Resolver> {
self.registry.resolver_for(self.schema_id(i)?, reader)
}
pub fn resolvers(&self, reader: &Schema) -> Resolvers {
let mut by_schema: HashMap<u128, Resolver> = HashMap::new();
for r in self.records() {
if by_schema.contains_key(&r.schema_id) {
continue;
}
if let Ok(resolver) = self.registry.resolver_for(r.schema_id, reader) {
by_schema.insert(r.schema_id, resolver);
}
}
Resolvers {
by_schema,
reader_id: reader.id(),
}
}
pub fn dump_json(&self, i: usize) -> Result<String> {
crate::dump::dump_json_with(self.schema(i)?, self.get(i)?)
}
pub fn iter(&self) -> impl Iterator<Item = &'a [u8]> + '_ {
(0..self.len()).map(move |i| self.get(i).expect("index validated in open"))
}
pub fn records(&self) -> impl Iterator<Item = Record> + '_ {
(0..self.len()).map(move |i| self.record(i).expect("index validated in open"))
}
}
#[derive(Clone, Debug)]
pub struct Resolvers {
by_schema: HashMap<u128, Resolver>,
reader_id: u128,
}
impl Resolvers {
pub fn get(&self, schema_id: u128) -> Option<&Resolver> {
self.by_schema.get(&schema_id)
}
pub fn for_record<'a>(&self, view: &FileView<'a>, i: usize) -> Result<&Resolver> {
let schema_id = view.schema_id(i)?;
self.by_schema.get(&schema_id).ok_or_else(|| {
Error::Incompatible(format!(
"record {i} was written under schema {schema_id:#034x}, \
which does not resolve into reader schema {:#034x}",
self.reader_id
))
})
}
pub fn len(&self) -> usize {
self.by_schema.len()
}
pub fn is_empty(&self) -> bool {
self.by_schema.is_empty()
}
}
#[derive(Clone, Debug)]
pub struct FileReader {
bytes: Vec<u8>,
}
impl FileReader {
pub fn open<P: AsRef<Path>>(path: P) -> Result<FileReader> {
FileReader::from_bytes(std::fs::read(path)?)
}
pub fn from_bytes(bytes: Vec<u8>) -> Result<FileReader> {
FileView::open(&bytes)?;
Ok(FileReader { bytes })
}
pub fn view(&self) -> Result<FileView<'_>> {
FileView::open(&self.bytes)
}
pub fn bytes(&self) -> &[u8] {
&self.bytes
}
}
pub struct FileBuilder {
buf: Vec<u8>,
index: Vec<Record>,
registry: SchemaRegistry,
next_id: u64,
checksums: Option<Vec<u32>>,
}
impl Default for FileBuilder {
fn default() -> FileBuilder {
FileBuilder::new()
}
}
impl FileBuilder {
pub fn new() -> FileBuilder {
FileBuilder {
buf: header_bytes(0).to_vec(),
index: Vec::new(),
registry: SchemaRegistry::new(),
next_id: FIRST_RECORD_ID,
checksums: None,
}
}
pub fn with_record_checksums(mut self) -> FileBuilder {
debug_assert!(self.index.is_empty(), "call before appending");
self.buf = header_bytes(OPT_RECORD_CRC).to_vec();
self.checksums = Some(Vec::new());
self
}
pub fn append(&mut self, schema: &Schema, value: &Value) -> Result<u64> {
let bytes = encode(schema, value, SchemaMode::HashOnly)?;
self.append_message(schema, &bytes)
}
pub fn append_message(&mut self, schema: &Schema, bytes: &[u8]) -> Result<u64> {
let id = self.next_id;
self.append_message_with_id(schema, bytes, id)
}
pub fn append_message_with_id(
&mut self,
schema: &Schema,
bytes: &[u8],
id: u64,
) -> Result<u64> {
if id < self.next_id {
return Err(Error::BadFile("record ids must be strictly ascending"));
}
let msg = Message::parse(bytes)?;
if msg.schema_id() != schema.id() {
return Err(Error::SchemaIdMismatch {
message: msg.schema_id(),
expected: schema.id(),
});
}
let schema_id = self.registry.register(schema.clone());
self.push_bytes(bytes, schema_id, id)
}
pub fn append_self_describing(&mut self, bytes: &[u8]) -> Result<u64> {
let schema = schema_of(bytes)?;
let schema_id = self.registry.register(schema);
let id = self.next_id;
self.push_bytes(bytes, schema_id, id)
}
fn push_bytes(&mut self, bytes: &[u8], schema_id: u128, id: u64) -> Result<u64> {
let offset = align_up(self.buf.len() as u64)?;
self.buf.resize(offset as usize, 0); self.buf.extend_from_slice(bytes);
if let Some(c) = &mut self.checksums {
c.push(crc32(bytes));
}
self.index.push(Record {
id,
offset,
length: bytes.len() as u64,
schema_id,
});
self.next_id = id
.checked_add(1)
.ok_or(Error::BadFile("record id space exhausted"))?;
Ok(id)
}
pub fn reserve_next_record_id(&mut self, n: u64) -> &mut Self {
self.next_id = self.next_id.max(n);
self
}
pub fn len(&self) -> usize {
self.index.len()
}
pub fn is_empty(&self) -> bool {
self.index.is_empty()
}
pub fn finish(mut self) -> Result<Vec<u8>> {
let tail_start = align_up(self.buf.len() as u64)?;
self.buf.resize(tail_start as usize, 0);
let (tail, footer) = commit_tail(
&self.index,
&self.registry,
tail_start,
1,
self.next_id,
self.checksums.as_deref(),
)?;
self.buf.extend_from_slice(&tail);
self.buf.extend_from_slice(&footer.encode());
debug_assert_eq!(self.buf.len() as u64, footer.file_len);
Ok(self.buf)
}
}
fn sync_parent_dir(path: &Path) -> Result<()> {
#[cfg(unix)]
{
let dir = path.parent().unwrap_or_else(|| Path::new("."));
File::open(dir)?.sync_all()?;
}
#[cfg(not(unix))]
{
let _ = path;
}
Ok(())
}
#[derive(Debug)]
struct WriterLock {
path: PathBuf,
}
impl WriterLock {
fn acquire(target: &Path) -> Result<WriterLock> {
let mut lock = target.as_os_str().to_os_string();
lock.push(".lock");
let path = PathBuf::from(lock);
match OpenOptions::new().write(true).create_new(true).open(&path) {
Ok(_) => Ok(WriterLock { path }),
Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
Err(Error::AlreadyLocked(path.display().to_string()))
}
Err(e) => Err(Error::Io(e.to_string())),
}
}
}
impl Drop for WriterLock {
fn drop(&mut self) {
let _ = std::fs::remove_file(&self.path);
}
}
enum Slot {
Committed {
id: u64,
offset: u64,
length: u64,
schema_id: u128,
crc: u32,
},
Staged {
id: u64,
staged: usize,
schema_id: u128,
crc: u32,
},
}
impl Slot {
fn id(&self) -> u64 {
match self {
Slot::Committed { id, .. } | Slot::Staged { id, .. } => *id,
}
}
fn crc(&self) -> u32 {
match self {
Slot::Committed { crc, .. } | Slot::Staged { crc, .. } => *crc,
}
}
fn schema_id(&self) -> u128 {
match self {
Slot::Committed { schema_id, .. } | Slot::Staged { schema_id, .. } => *schema_id,
}
}
}
pub struct FileWriter {
file: File,
path: PathBuf,
lock: Option<WriterLock>,
generation: u64,
committed_len: u64,
next_record_id: u64,
slots: Vec<Slot>,
staged: Vec<Vec<u8>>,
registry: SchemaRegistry,
checksums: bool,
}
impl FileWriter {
pub fn create<P: AsRef<Path>>(path: P) -> Result<FileWriter> {
FileWriter::create_with(path, 0)
}
pub fn create_checksummed<P: AsRef<Path>>(path: P) -> Result<FileWriter> {
FileWriter::create_with(path, OPT_RECORD_CRC)
}
fn create_with<P: AsRef<Path>>(path: P, optional_features: u32) -> Result<FileWriter> {
let path = path.as_ref().to_path_buf();
let mut file = OpenOptions::new()
.read(true)
.write(true)
.create(true)
.truncate(true)
.open(&path)?;
file.write_all(&header_bytes(optional_features))?;
let mut w = FileWriter {
file,
path,
lock: None,
generation: 0,
committed_len: FILE_HEADER_LEN as u64,
next_record_id: FIRST_RECORD_ID,
slots: Vec::new(),
staged: Vec::new(),
registry: SchemaRegistry::new(),
checksums: optional_features & OPT_RECORD_CRC != 0,
};
w.commit()?;
Ok(w)
}
pub fn open<P: AsRef<Path>>(path: P) -> Result<FileWriter> {
let path = path.as_ref().to_path_buf();
let mut file = OpenOptions::new().read(true).write(true).open(&path)?;
let mut image = Vec::new();
file.read_to_end(&mut image)?;
let view = FileView::open(&image)?;
let checksums = view.has_record_checksums();
let slots = (0..view.len())
.map(|i| {
let r = view.record(i).expect("index validated in open");
Slot::Committed {
id: r.id,
offset: r.offset,
length: r.length,
schema_id: r.schema_id,
crc: view.record_checksum(i).unwrap_or(0),
}
})
.collect();
Ok(FileWriter {
lock: None,
generation: view.generation(),
committed_len: view.file_len(),
next_record_id: view.next_record_id(),
registry: view.schemas().clone(),
slots,
staged: Vec::new(),
checksums,
file,
path,
})
}
pub fn open_locked<P: AsRef<Path>>(path: P) -> Result<FileWriter> {
let path = path.as_ref();
let lock = WriterLock::acquire(path)?;
let mut w = FileWriter::open(path)?;
w.lock = Some(lock);
Ok(w)
}
pub fn open_or_create_locked<P: AsRef<Path>>(path: P) -> Result<FileWriter> {
let path = path.as_ref();
let lock = WriterLock::acquire(path)?;
let mut w = FileWriter::open_or_create(path)?;
w.lock = Some(lock);
Ok(w)
}
pub fn is_locked(&self) -> bool {
self.lock.is_some()
}
pub fn open_or_create<P: AsRef<Path>>(path: P) -> Result<FileWriter> {
let path = path.as_ref();
if path.exists() {
FileWriter::open(path)
} else {
FileWriter::create(path)
}
}
pub fn append(&mut self, schema: &Schema, value: &Value) -> Result<u64> {
let bytes = encode(schema, value, SchemaMode::HashOnly)?;
self.append_message(schema, &bytes)
}
pub fn append_message(&mut self, schema: &Schema, bytes: &[u8]) -> Result<u64> {
let msg = Message::parse(bytes)?;
if msg.schema_id() != schema.id() {
return Err(Error::SchemaIdMismatch {
message: msg.schema_id(),
expected: schema.id(),
});
}
let schema_id = self.registry.register(schema.clone());
self.stage(bytes.to_vec(), schema_id)
}
pub fn append_self_describing(&mut self, bytes: &[u8]) -> Result<u64> {
let schema = schema_of(bytes)?;
let schema_id = self.registry.register(schema);
self.stage(bytes.to_vec(), schema_id)
}
fn stage(&mut self, bytes: Vec<u8>, schema_id: u128) -> Result<u64> {
let id = self.next_record_id;
self.next_record_id = id
.checked_add(1)
.ok_or(Error::BadFile("record id space exhausted"))?;
let crc = if self.checksums { crc32(&bytes) } else { 0 };
self.staged.push(bytes);
self.slots.push(Slot::Staged {
id,
staged: self.staged.len() - 1,
schema_id,
crc,
});
Ok(id)
}
pub fn remove_id(&mut self, id: u64) -> Result<&mut Self> {
let at = self
.slots
.iter()
.position(|s| s.id() == id)
.ok_or(Error::IndexOutOfBounds)?;
self.slots.remove(at);
Ok(self)
}
pub fn remove_ids(&mut self, ids: &[u64]) -> usize {
let doomed: HashSet<u64> = ids.iter().copied().collect();
let before = self.slots.len();
self.slots.retain(|s| !doomed.contains(&s.id()));
before - self.slots.len()
}
pub fn remove(&mut self, i: usize) -> Result<&mut Self> {
if i >= self.slots.len() {
return Err(Error::IndexOutOfBounds);
}
self.slots.remove(i);
Ok(self)
}
pub fn retain<F: FnMut(u64, u128) -> bool>(&mut self, mut keep: F) -> usize {
let before = self.slots.len();
self.slots.retain(|s| keep(s.id(), s.schema_id()));
before - self.slots.len()
}
pub fn purge_ids(&mut self, ids: &[u64]) -> Result<usize> {
let removed = self.remove_ids(ids);
self.commit()?;
self.compact()?;
Ok(removed)
}
pub fn len(&self) -> usize {
self.slots.len()
}
pub fn is_empty(&self) -> bool {
self.slots.is_empty()
}
pub fn pending(&self) -> usize {
self.staged.len()
}
pub fn generation(&self) -> u64 {
self.generation
}
pub fn next_record_id(&self) -> u64 {
self.next_record_id
}
pub fn ids(&self) -> impl Iterator<Item = u64> + '_ {
self.slots.iter().map(|s| s.id())
}
pub fn path(&self) -> &Path {
&self.path
}
pub fn commit(&mut self) -> Result<u64> {
self.file.set_len(self.committed_len)?;
self.file.seek(SeekFrom::Start(self.committed_len))?;
let mut cursor = self.committed_len;
let mut body = Vec::new();
let mut index = Vec::with_capacity(self.slots.len());
for slot in &self.slots {
match slot {
Slot::Committed {
id,
offset,
length,
schema_id,
..
} => index.push(Record {
id: *id,
offset: *offset,
length: *length,
schema_id: *schema_id,
}),
Slot::Staged {
id,
staged,
schema_id,
..
} => {
let bytes = &self.staged[*staged];
let offset = align_up(cursor)?;
body.resize((offset - self.committed_len) as usize, 0);
body.extend_from_slice(bytes);
cursor = offset + bytes.len() as u64;
index.push(Record {
id: *id,
offset,
length: bytes.len() as u64,
schema_id: *schema_id,
});
}
}
}
let tail_start = align_up(cursor)?;
body.resize((tail_start - self.committed_len) as usize, 0);
let generation = self.generation + 1;
let crcs: Option<Vec<u32>> = if self.checksums {
Some(self.slots.iter().map(|s| s.crc()).collect())
} else {
None
};
let (tail, footer) = commit_tail(
&index,
&self.registry,
tail_start,
generation,
self.next_record_id,
crcs.as_deref(),
)?;
body.extend_from_slice(&tail);
self.file.write_all(&body)?;
self.file.sync_all()?;
self.file.write_all(&footer.encode())?;
self.file.sync_all()?;
let kept: Vec<u32> = self.slots.iter().map(|s| s.crc()).collect();
self.slots = index
.iter()
.zip(kept)
.map(|(r, crc)| Slot::Committed {
id: r.id,
offset: r.offset,
length: r.length,
schema_id: r.schema_id,
crc,
})
.collect();
self.staged.clear();
self.generation = generation;
self.committed_len = footer.file_len;
Ok(generation)
}
pub fn compact(&mut self) -> Result<()> {
if !self.staged.is_empty() {
self.commit()?;
}
let mut builder = if self.checksums {
FileBuilder::new().with_record_checksums()
} else {
FileBuilder::new()
};
for i in 0..self.slots.len() {
let bytes = self.read_record(i)?;
let id = self.slots[i].id();
let schema_id = self.slots[i].schema_id();
let schema = self
.registry
.get(schema_id)
.ok_or(Error::MissingSchema(schema_id))?
.clone();
builder.append_message_with_id(&schema, &bytes, id)?;
}
builder.reserve_next_record_id(self.next_record_id);
let image = builder.finish()?;
let mut tmp = self.path.clone().into_os_string();
tmp.push(".compact");
let tmp = PathBuf::from(tmp);
{
let mut f = OpenOptions::new()
.read(true)
.write(true)
.create(true)
.truncate(true)
.open(&tmp)?;
f.write_all(&image)?;
f.sync_all()?;
}
std::fs::rename(&tmp, &self.path)?;
sync_parent_dir(&self.path)?;
let held = self.lock.take();
*self = FileWriter::open(&self.path)?;
self.lock = held;
Ok(())
}
pub fn read_record(&mut self, i: usize) -> Result<Vec<u8>> {
match self.slots.get(i).ok_or(Error::IndexOutOfBounds)? {
Slot::Staged { staged, .. } => Ok(self.staged[*staged].clone()),
Slot::Committed { offset, length, .. } => {
let (offset, length) = (*offset, *length as usize);
let mut buf = vec![0u8; length];
self.file.seek(SeekFrom::Start(offset))?;
self.file.read_exact(&mut buf)?;
Ok(buf)
}
}
}
pub fn read_record_by_id(&mut self, id: u64) -> Result<Vec<u8>> {
let at = self
.slots
.iter()
.position(|s| s.id() == id)
.ok_or(Error::IndexOutOfBounds)?;
self.read_record(at)
}
pub fn image(&mut self) -> Result<Vec<u8>> {
let mut buf = Vec::new();
self.file.seek(SeekFrom::Start(0))?;
self.file.read_to_end(&mut buf)?;
Ok(buf)
}
}