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")
}
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,
})
}
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,
})
}
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(())
}