journal.rs

← Back to explorer
src/ storage/ journal.rs
Raw
Rust 137 lines · UTF-8
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)
    }
}