journal.rs
← Back to explorer
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137
// Created by AG on 15-08-2026
use crate::types::{SlimeDBError, Result, DBRecordVersioned};
use parking_lot::Mutex;
use std::fs::{File, OpenOptions};
use std::io::{BufReader, BufWriter, Read, Seek, SeekFrom, Write};
use std::path::PathBuf;
use std::sync::atomic::{AtomicU64, Ordering};
const JOURNAL_SIGNATURE: &[u8; 5] = b"SLIME";
const JOURNAL_RECORD_HEADER_SIZE: usize = 12;
pub struct Journal {
file_path: PathBuf,
writer: Mutex<BufWriter<File>>,
sequence: AtomicU64,
sync_on_write: bool,
}
impl Journal {
pub fn open(file_path: PathBuf, sync_on_write: bool) -> Result<Self> {
let file = OpenOptions::new()
.create(true)
.read(true)
.append(true)
.open(&file_path)?;
let mut file_reader = BufReader::new(file.try_clone()?);
let read_sequence = Self::recover_write_sequence(&mut file_reader)?;
let journal_writer = BufWriter::with_capacity(64 * 1024, file);
Ok(
Self {
file_path,
writer: Mutex::new(journal_writer),
sequence: AtomicU64::new(read_sequence),
sync_on_write,
}
)
}
fn recover_write_sequence(reader: &mut BufReader<File>) -> Result <u64> {
let file_length = reader.seek(SeekFrom::End(0))?;
if file_length == 0 {
return Ok(0);
}
reader.seek(SeekFrom::Start(0))?;
let mut journal_sig = [0u8; 5];
if reader.read_exact(&mut journal_sig).is_ok() && &journal_sig != JOURNAL_SIGNATURE {
return Err(SlimeDBError::Corruption("Invalid Journal Signature".to_string()));
}
let mut max_sequence = 0u64;
reader.seek(SeekFrom::Start(4))?;
loop {
let mut file_header = [0u8; JOURNAL_RECORD_HEADER_SIZE];
match reader.read_exact(&mut file_header) {
Ok(_) => {}
Err(err) if err.kind() == std::io::ErrorKind::UnexpectedEof => break,
Err(err) => return Err(err.into()),
}
let sequence = u64::from_le_bytes(file_header[0..8].try_into().unwrap());
let length = u32::from_le_bytes(file_header[0..12].try_into().unwrap()) as u64;
max_sequence = max_sequence.max(sequence);
let cursor_pos = reader.stream_position()?;
if cursor_pos + length + 4 > file_length {
break;
}
reader.seek(SeekFrom::Current(length as i64 + 4))?;
}
Ok(max_sequence + 1)
}
pub fn append(&self, record: &DBRecordVersioned) -> Result<u64> {
let sequence = self.sequence.fetch_add(1, Ordering::SeqCst);
let data = postcard::to_allocvec(record)
.map_err(|err| SlimeDBError::Serialization(err.to_string()))?;
let db_checksum = crc32fast::hash(&data);
let mut writer = self.writer.lock();
if sequence == 0 {
writer.write_all(JOURNAL_SIGNATURE)?;
}
writer.write_all(&sequence.to_le_bytes())?;
writer.write_all(&(data.len() as u32).to_le_bytes())?;
writer.write_all(&data)?;
writer.write_all(&db_checksum.to_le_bytes())?;
if self.sync_on_write {
writer.flush()?;
writer.get_ref().sync_data()?;
}
Ok(sequence)
}
pub fn batch_append(&self, records: &[DBRecordVersioned]) -> Result<u64> {
let mut writer = self.writer.lock();
let mut last_sequence = 0;
for record in records {
let sequence = self.sequence.fetch_add(1, Ordering::SeqCst);
let data = postcard::to_allocvec(record)
.map_err(|err| SlimeDBError::Serialization(err.to_string()))?;
let db_batch_checksum = crc32fast::hash(&data);
if sequence == 0 {
writer.write_all(JOURNAL_SIGNATURE)?;
}
writer.write_all(&sequence.to_le_bytes())?;
writer.write_all(&(data.len() as u32).to_le_bytes())?;
writer.write_all(&data)?;
writer.write_all(&db_batch_checksum.to_le_bytes())?;
last_sequence = sequence;
}
writer.flush()?;
if self.sync_on_write {
writer.get_ref().sync_data()?;
}
Ok(last_sequence)
}
}