quatzal-storage 0.1.0

Sharded LSM row-storage engine for Quatzal: WAL, snapshots, and crash recovery on io_uring (Linux only).
// SPDX-License-Identifier: Apache-2.0
//! Immutable, sorted, on-disk table. Phase 1 keeps a full in-memory key index built at
//! open/flush time (documented simplification in `CLAUDE.md` — no sparse index / mmap yet;
//! fine at Phase 1 data sizes, revisit before Phase 2 scale targets).

use std::collections::BTreeMap;
use std::path::{Path, PathBuf};

use glommio::io::BufferedFile;
use quatzal_schema::{UaceError, UaceResult};

use crate::framing::{for_each_framed, frame};
use crate::memtable::MemTable;
use crate::pax::PaxRecord;

pub struct SsTable {
    pub id: u64,
    pub level: u32,
    pub path: PathBuf,
    index: BTreeMap<Vec<u8>, PaxRecord>,
}

impl SsTable {
    pub fn file_name(level: u32, id: u64) -> String {
        format!("l{level}-{id:016x}.sst")
    }

    /// Write every record in `mem` (already key-sorted) to a new SSTable file.
    pub async fn flush_memtable(
        dir: &Path,
        id: u64,
        level: u32,
        mem: &MemTable,
    ) -> UaceResult<Self> {
        let path = dir.join(Self::file_name(level, id));
        let mut index = BTreeMap::new();
        let mut buf = Vec::new();
        for record in mem.iter_sorted() {
            buf.extend_from_slice(&frame(&record.to_bytes()));
            index.insert(record.key.clone(), record.clone());
        }
        write_whole_file(&path, buf).await?;
        Ok(SsTable {
            id,
            level,
            path,
            index,
        })
    }

    /// Write a pre-merged, key-sorted record list (used by compaction).
    pub async fn write_records(
        dir: &Path,
        id: u64,
        level: u32,
        records: impl Iterator<Item = PaxRecord>,
    ) -> UaceResult<Self> {
        let path = dir.join(Self::file_name(level, id));
        let mut index = BTreeMap::new();
        let mut buf = Vec::new();
        for record in records {
            buf.extend_from_slice(&frame(&record.to_bytes()));
            index.insert(record.key.clone(), record);
        }
        write_whole_file(&path, buf).await?;
        Ok(SsTable {
            id,
            level,
            path,
            index,
        })
    }

    /// Reopen an SSTable file that already exists on disk (used during `Shard::open`).
    pub async fn open(path: PathBuf, id: u64, level: u32) -> UaceResult<Self> {
        let file = BufferedFile::open(&path)
            .await
            .map_err(|e| UaceError::Io(format!("sstable open {path:?}: {e}")))?;
        let size = file
            .file_size()
            .await
            .map_err(|e| UaceError::Io(format!("sstable size {path:?}: {e}")))?;
        let result = file
            .read_at(0, size as usize)
            .await
            .map_err(|e| UaceError::Io(format!("sstable read {path:?}: {e}")))?;
        let bytes: &[u8] = &result;

        let mut index = BTreeMap::new();
        for_each_framed(bytes, |payload| {
            let (record, used) = PaxRecord::from_bytes(payload)?;
            if used != payload.len() {
                return Err(UaceError::Codec("sstable record framing mismatch".into()));
            }
            index.insert(record.key.clone(), record);
            Ok(())
        })?;

        Ok(SsTable {
            id,
            level,
            path,
            index,
        })
    }

    pub fn get(&self, key: &[u8]) -> Option<&PaxRecord> {
        self.index.get(key)
    }

    pub fn iter(&self) -> impl Iterator<Item = &PaxRecord> {
        self.index.values()
    }

    pub fn len(&self) -> usize {
        self.index.len()
    }

    pub fn is_empty(&self) -> bool {
        self.index.is_empty()
    }
}

async fn write_whole_file(path: &Path, bytes: Vec<u8>) -> UaceResult<()> {
    let file = BufferedFile::create(path)
        .await
        .map_err(|e| UaceError::Io(format!("sstable create {path:?}: {e}")))?;
    file.write_at(bytes, 0)
        .await
        .map_err(|e| UaceError::Io(format!("sstable write {path:?}: {e}")))?;
    file.fdatasync()
        .await
        .map_err(|e| UaceError::Io(format!("sstable fsync {path:?}: {e}")))?;
    file.close()
        .await
        .map_err(|e| UaceError::Io(format!("sstable close {path:?}: {e}")))?;
    Ok(())
}