use std::collections::BTreeMap;
use std::collections::btree_map::Range;
use std::fs::File;
use std::io::{BufReader, BufWriter, Read as _, Seek as _, SeekFrom, Write as _};
use std::ops::{Bound, RangeBounds};
use std::path::PathBuf;
use std::result::Result as StdResult;
use fs4::fs_std::FileExt;
use log::{error, info};
use super::{Engine, Status};
use crate::error::{Error, Result};
pub struct BitCask {
log: Log,
keydir: KeyDir,
}
type KeyDir = BTreeMap<Vec<u8>, ValueLocation>;
#[derive(Clone, Copy)]
struct ValueLocation {
offset: u64,
length: usize,
}
impl ValueLocation {
fn end(&self) -> u64 {
self.offset + self.length as u64
}
}
impl BitCask {
pub fn new(path: PathBuf) -> Result<Self> {
let mut log = Log::new(path.clone())?;
let keydir = log.build_keydir()?;
info!("Opened {} with {} live keys", path.display(), keydir.len());
Ok(Self { log, keydir })
}
pub fn new_maybe_compact(
path: PathBuf,
garbage_min_fraction: f64,
garbage_min_bytes: u64,
) -> Result<Self> {
let mut engine = Self::new(path)?;
let status = engine.status()?;
let total_size = status.disk_size;
let garbage_size = status.garbage_disk_size();
let garbage_fraction = garbage_size as f64 / total_size as f64;
if garbage_size > 0
&& garbage_size >= garbage_min_bytes
&& garbage_fraction >= garbage_min_fraction
{
info!(
"Compacting {} to remove {:.0}% garbage ({:.1} MB out of {:.1} MB)",
engine.log.path.display(),
garbage_fraction * 100.0,
garbage_size as f64 / 1024.0 / 1024.0,
total_size as f64 / 1024.0 / 1024.0
);
engine.compact()?;
info!(
"Compacted {} to size {:.1} MB",
engine.log.path.display(),
(total_size - garbage_size) as f64 / 1024.0 / 1024.0
);
}
Ok(engine)
}
}
impl Engine for BitCask {
type ScanIterator<'a> = ScanIterator<'a>;
fn delete(&mut self, key: &[u8]) -> Result<()> {
self.log.write_entry(key, None)?;
self.keydir.remove(key);
Ok(())
}
fn flush(&mut self) -> Result<()> {
#[cfg(not(test))]
self.log.file.sync_all()?;
Ok(())
}
fn get(&mut self, key: &[u8]) -> Result<Option<Vec<u8>>> {
let Some(location) = self.keydir.get(key) else {
return Ok(None);
};
self.log.read_value(*location).map(Some)
}
fn scan(&mut self, range: impl RangeBounds<Vec<u8>>) -> Self::ScanIterator<'_> {
ScanIterator { inner: self.keydir.range(range), log: &mut self.log }
}
fn scan_dyn(
&mut self,
range: (Bound<Vec<u8>>, Bound<Vec<u8>>),
) -> Box<dyn super::ScanIterator + '_> {
Box::new(self.scan(range))
}
fn set(&mut self, key: &[u8], value: Vec<u8>) -> Result<()> {
let value_location = self.log.write_entry(key, Some(&*value))?;
self.keydir.insert(key.to_vec(), value_location);
Ok(())
}
fn status(&mut self) -> Result<Status> {
let keys = self.keydir.len() as u64;
let size =
self.keydir.iter().map(|(key, value_loc)| (key.len() + value_loc.length) as u64).sum();
let disk_size = self.log.file.metadata()?.len();
let live_disk_size = size + 8 * keys; Ok(Status { name: "bitcask".to_string(), keys, size, disk_size, live_disk_size })
}
}
impl BitCask {
pub fn compact(&mut self) -> Result<()> {
let new_path = self.log.path.with_extension("new");
let mut new_log = Log::new(new_path)?;
new_log.file.set_len(0)?;
let mut new_keydir = KeyDir::new();
for (key, value_loc) in &self.keydir {
let value = self.log.read_value(*value_loc)?;
let value_loc = new_log.write_entry(key, Some(&value))?;
new_keydir.insert(key.clone(), value_loc);
}
std::fs::rename(&new_log.path, &self.log.path)?;
new_log.path = self.log.path.clone();
self.log = new_log;
self.keydir = new_keydir;
Ok(())
}
}
impl Drop for BitCask {
fn drop(&mut self) {
if let Err(error) = self.flush() {
error!("failed to flush file: {}", error)
}
}
}
pub struct ScanIterator<'a> {
inner: Range<'a, Vec<u8>, ValueLocation>,
log: &'a mut Log,
}
impl ScanIterator<'_> {
fn map(&mut self, item: (&Vec<u8>, &ValueLocation)) -> <Self as Iterator>::Item {
let (key, value_loc) = item;
Ok((key.clone(), self.log.read_value(*value_loc)?))
}
}
impl Iterator for ScanIterator<'_> {
type Item = Result<(Vec<u8>, Vec<u8>)>;
fn next(&mut self) -> Option<Self::Item> {
self.inner.next().map(|item| self.map(item))
}
}
impl DoubleEndedIterator for ScanIterator<'_> {
fn next_back(&mut self) -> Option<Self::Item> {
self.inner.next_back().map(|item| self.map(item))
}
}
struct Log {
file: File,
path: PathBuf,
}
impl Log {
fn new(path: PathBuf) -> Result<Self> {
if let Some(dir) = path.parent() {
std::fs::create_dir_all(dir)?
}
let file = std::fs::OpenOptions::new()
.read(true)
.write(true)
.create(true)
.truncate(false)
.open(&path)?;
if !file.try_lock_exclusive()? {
return Err(Error::IO(format!("file {path:?} is already is use")));
}
Ok(Self { file, path })
}
fn build_keydir(&mut self) -> Result<KeyDir> {
let mut len_buf = [0u8; 4];
let mut keydir = KeyDir::new();
let file_len = self.file.metadata()?.len();
let mut r = BufReader::new(&mut self.file);
let mut offset = r.seek(SeekFrom::Start(0))?;
while offset < file_len {
let result = || -> StdResult<(Vec<u8>, Option<ValueLocation>), std::io::Error> {
r.read_exact(&mut len_buf)?;
let key_len = u32::from_be_bytes(len_buf);
r.read_exact(&mut len_buf)?;
let value_loc = match i32::from_be_bytes(len_buf) {
..0 => None, len => Some(ValueLocation {
offset: offset + 8 + key_len as u64,
length: len as usize,
}),
};
let mut key: Vec<u8> = vec![0; key_len as usize];
r.read_exact(&mut key)?;
if let Some(value_loc) = value_loc {
if value_loc.end() > file_len {
return Err(std::io::Error::new(
std::io::ErrorKind::UnexpectedEof,
"value extends beyond end of file",
));
}
r.seek_relative(value_loc.length as i64)?;
}
offset += 8 + key_len as u64 + value_loc.map_or(0, |v| v.length) as u64;
Ok((key, value_loc))
}();
match result {
Ok((key, Some(value_loc))) => keydir.insert(key, value_loc),
Ok((key, None)) => keydir.remove(&key),
Err(err) if err.kind() == std::io::ErrorKind::UnexpectedEof => {
error!("Found incomplete entry at offset {offset}, truncating file");
self.file.set_len(offset)?; break;
}
Err(err) => return Err(err.into()),
};
}
Ok(keydir)
}
fn read_value(&mut self, location: ValueLocation) -> Result<Vec<u8>> {
let mut value: Vec<u8> = vec![0; location.length];
self.file.seek(SeekFrom::Start(location.offset))?;
self.file.read_exact(&mut value)?;
Ok(value)
}
fn write_entry(&mut self, key: &[u8], value: Option<&[u8]>) -> Result<ValueLocation> {
let length = 8 + key.len() + value.map_or(0, |v| v.len());
let offset = self.file.seek(SeekFrom::End(0))?;
let mut w = BufWriter::with_capacity(length, &mut self.file);
w.write_all(&(key.len() as u32).to_be_bytes())?;
w.write_all(&value.map_or(-1, |v| v.len() as i32).to_be_bytes())?;
w.write_all(key)?;
w.write_all(value.unwrap_or_default())?;
w.flush()?;
Ok(ValueLocation {
offset: offset + 8 + key.len() as u64,
length: value.map_or(0, |v| v.len()),
})
}
}