use fs4::FileExt;
use std::fs::File;
use std::io::{BufReader, BufWriter, Read, Seek, SeekFrom, Write};
use std::ops::Bound;
use std::path::{Path, PathBuf};
const KEY_VAL_HEADER_LEN: u32 = 4;
const MERGE_FILE_EXT: &str = "merge";
const WRITE_BUF_TARGET: usize = 4 * 1024 * 1024;
const HINT_MAGIC_V1: &[u8; 8] = b"HINT0001";
const HINT_MAGIC_V2: &[u8; 8] = b"HINT0002";
const HINT_HEADER_SIZE_V1: usize = 32;
const HINT_HEADER_SIZE_V2: usize = 40;
const HINT_FLAG_LIVE: i32 = 0;
pub const DEFAULT_MAX_FILE_BYTES: u64 = 64 * 1024 * 1024;
const FLAG_HARD_DELETE: i32 = -1;
const FLAG_SOFT_TOMBSTONE: i32 = -2;
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum KeyLoc {
Live {
file_id: u32,
pos: u64,
len: u32,
},
SoftTombstone { file_id: u32 },
}
pub type KeyDir = std::collections::BTreeMap<Vec<u8>, KeyLoc>;
pub type Result<T> = std::result::Result<T, std::io::Error>;
fn apply_log_record(
key_dir: &mut KeyDir,
file_id: u32,
key: Vec<u8>,
value_pos: u64,
flag: i32,
value_len: u32,
) {
if flag >= 0 {
key_dir.insert(
key,
KeyLoc::Live {
file_id,
pos: value_pos,
len: value_len,
},
);
} else if flag == FLAG_SOFT_TOMBSTONE {
key_dir.insert(key, KeyLoc::SoftTombstone { file_id });
} else {
key_dir.remove(&key);
}
}
pub fn hint_path_for(data_path: &Path) -> PathBuf {
data_path.with_extension("hint")
}
pub fn segment_path(active_path: &Path, file_id: u32) -> PathBuf {
let stem = active_path
.file_stem()
.and_then(|s| s.to_str())
.unwrap_or("data");
let parent = active_path.parent().unwrap_or_else(|| Path::new("."));
parent.join(format!("{stem}.{file_id}.log"))
}
fn discover_sealed_ids(active_path: &Path) -> Vec<u32> {
let stem = active_path
.file_stem()
.and_then(|s| s.to_str())
.unwrap_or("data");
let parent = match active_path.parent() {
Some(p) => p,
None => return Vec::new(),
};
let mut ids = Vec::new();
let Ok(rd) = std::fs::read_dir(parent) else {
return ids;
};
let prefix = format!("{stem}.");
for ent in rd.flatten() {
let name = ent.file_name();
let name = name.to_string_lossy();
if let Some(rest) = name.strip_prefix(&prefix) {
if let Some(num) = rest.strip_suffix(".log") {
if let Ok(id) = num.parse::<u32>() {
ids.push(id);
}
}
}
}
ids.sort_unstable();
ids
}
pub fn write_hint_file(
path: &Path,
key_dir: &KeyDir,
active_id: u32,
active_log_size: u64,
) -> Result<()> {
let tmp = {
let mut p = path.as_os_str().to_os_string();
p.push(".tmp");
PathBuf::from(p)
};
if let Some(parent) = path.parent() {
if !parent.as_os_str().is_empty() {
let _ = std::fs::create_dir_all(parent);
}
}
{
let f = std::fs::OpenOptions::new()
.write(true)
.create(true)
.truncate(true)
.open(&tmp)?;
let mut w = BufWriter::new(f);
let count = key_dir.len() as u64;
let mut hdr = [0u8; HINT_HEADER_SIZE_V2];
hdr[0..8].copy_from_slice(HINT_MAGIC_V2);
hdr[8..12].copy_from_slice(&active_id.to_le_bytes());
hdr[16..24].copy_from_slice(&active_log_size.to_le_bytes());
hdr[24..32].copy_from_slice(&count.to_le_bytes());
w.write_all(&hdr)?;
for (key, loc) in key_dir.iter() {
let (flag, file_id, pos, len) = match loc {
KeyLoc::Live {
file_id,
pos,
len,
} => (HINT_FLAG_LIVE, *file_id, *pos, *len),
KeyLoc::SoftTombstone { file_id } => {
(FLAG_SOFT_TOMBSTONE, *file_id, 0u64, 0u32)
}
};
w.write_all(&(key.len() as u32).to_le_bytes())?;
w.write_all(&flag.to_le_bytes())?;
w.write_all(&file_id.to_le_bytes())?;
w.write_all(&pos.to_le_bytes())?;
w.write_all(&len.to_le_bytes())?;
w.write_all(key)?;
}
w.flush()?;
w.into_inner()
.map_err(|e| e.into_error())?
.sync_all()?;
}
if path.exists() {
let _ = std::fs::remove_file(path);
}
std::fs::rename(&tmp, path)?;
Ok(())
}
pub fn read_hint_file(path: &Path) -> Result<(KeyDir, u32, u64)> {
let mut f = std::fs::File::open(path)?;
let mut magic = [0u8; 8];
f.read_exact(&mut magic)?;
if &magic == HINT_MAGIC_V2 {
let mut rest = [0u8; HINT_HEADER_SIZE_V2 - 8];
f.read_exact(&mut rest)?;
let active_id = u32::from_le_bytes(rest[0..4].try_into().unwrap());
let active_log_size = u64::from_le_bytes(rest[8..16].try_into().unwrap());
let entry_count = u64::from_le_bytes(rest[16..24].try_into().unwrap());
let mut key_dir = KeyDir::new();
let mut r = BufReader::new(f);
for _ in 0..entry_count {
let mut nbuf = [0u8; 4];
r.read_exact(&mut nbuf)?;
let key_len = u32::from_le_bytes(nbuf) as usize;
r.read_exact(&mut nbuf)?;
let flag = i32::from_le_bytes(nbuf);
r.read_exact(&mut nbuf)?;
let file_id = u32::from_le_bytes(nbuf);
let mut pbuf = [0u8; 8];
r.read_exact(&mut pbuf)?;
let pos = u64::from_le_bytes(pbuf);
r.read_exact(&mut nbuf)?;
let len = u32::from_le_bytes(nbuf);
if key_len > 16 * 1024 * 1024 {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidData,
"hint key_len 异常",
));
}
let mut key = vec![0u8; key_len];
r.read_exact(&mut key)?;
match flag {
HINT_FLAG_LIVE => {
key_dir.insert(
key,
KeyLoc::Live {
file_id,
pos,
len,
},
);
}
FLAG_SOFT_TOMBSTONE => {
key_dir.insert(key, KeyLoc::SoftTombstone { file_id });
}
_ => {}
}
}
Ok((key_dir, active_id, active_log_size))
} else if &magic == HINT_MAGIC_V1 {
let mut rest = [0u8; HINT_HEADER_SIZE_V1 - 8];
f.read_exact(&mut rest)?;
let log_size = u64::from_le_bytes(rest[0..8].try_into().unwrap());
let entry_count = u64::from_le_bytes(rest[8..16].try_into().unwrap());
let mut key_dir = KeyDir::new();
let mut r = BufReader::new(f);
for _ in 0..entry_count {
let mut nbuf = [0u8; 4];
r.read_exact(&mut nbuf)?;
let key_len = u32::from_le_bytes(nbuf) as usize;
r.read_exact(&mut nbuf)?;
let flag = i32::from_le_bytes(nbuf);
let mut pbuf = [0u8; 8];
r.read_exact(&mut pbuf)?;
let pos = u64::from_le_bytes(pbuf);
r.read_exact(&mut nbuf)?;
let len = u32::from_le_bytes(nbuf);
let mut key = vec![0u8; key_len];
r.read_exact(&mut key)?;
match flag {
HINT_FLAG_LIVE => {
key_dir.insert(
key,
KeyLoc::Live {
file_id: 0,
pos,
len,
},
);
}
FLAG_SOFT_TOMBSTONE => {
key_dir.insert(key, KeyLoc::SoftTombstone { file_id: 0 });
}
_ => {}
}
}
Ok((key_dir, 0, log_size))
} else {
Err(std::io::Error::new(
std::io::ErrorKind::InvalidData,
"非法 hint 魔数",
))
}
}
pub struct Log {
path: PathBuf,
file: File,
end_pos: u64,
buf: Vec<u8>,
}
impl Drop for Log {
fn drop(&mut self) {
let _ = self.flush_buf();
}
}
impl Log {
pub 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)
.open(&path)?;
file.try_lock_exclusive()?;
let end_pos = file.metadata()?.len();
Ok(Self {
path,
file,
end_pos,
buf: Vec::with_capacity(WRITE_BUF_TARGET),
})
}
pub fn flush_buf(&mut self) -> Result<()> {
if self.buf.is_empty() {
return Ok(());
}
self.file.seek(SeekFrom::End(0))?;
self.file.write_all(&self.buf)?;
self.buf.clear();
Ok(())
}
fn scan_log_into(
&mut self,
key_dir: &mut KeyDir,
file_id: u32,
from: u64,
to: u64,
) -> Result<()> {
if from >= to {
return Ok(());
}
let mut len_buf = [0u8; KEY_VAL_HEADER_LEN as usize];
let mut r = BufReader::new(&mut self.file);
let mut pos: u64 = r.seek(SeekFrom::Start(from))?;
while pos < to {
let read_one = || -> Result<(Vec<u8>, u64, i32, u32)> {
r.read_exact(&mut len_buf)?;
let key_len = u32::from_be_bytes(len_buf);
r.read_exact(&mut len_buf)?;
let flag = i32::from_be_bytes(len_buf);
let value_pos = pos + KEY_VAL_HEADER_LEN as u64 * 2 + key_len as u64;
let mut key = vec![0; key_len as usize];
r.read_exact(&mut key)?;
let value_len = if flag >= 0 { flag as u32 } else { 0 };
if flag >= 0 {
r.seek_relative(value_len as i64)?;
}
Ok((key, value_pos, flag, value_len))
}();
match read_one {
Ok((key, value_pos, flag, value_len)) => {
apply_log_record(key_dir, file_id, key, value_pos, flag, value_len);
if flag >= 0 {
pos = value_pos + value_len as u64;
} else {
pos = value_pos;
}
}
Err(err) => return Err(err),
}
}
Ok(())
}
pub fn load_index(&mut self, file_id: u32) -> Result<KeyDir> {
self.flush_buf()?;
let file_len = self.file.metadata()?.len();
self.end_pos = file_len;
let mut key_dir = KeyDir::new();
self.scan_log_into(&mut key_dir, file_id, 0, file_len)?;
Ok(key_dir)
}
pub fn load_index_from(
&mut self,
key_dir: &mut KeyDir,
file_id: u32,
from_offset: u64,
) -> Result<()> {
self.flush_buf()?;
let file_len = self.file.metadata()?.len();
self.end_pos = file_len;
if from_offset > file_len {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidData,
format!("hint log_size={from_offset} > data.log len={file_len}"),
));
}
self.scan_log_into(key_dir, file_id, from_offset, file_len)
}
fn read_value(&mut self, value_pos: u64, value_len: u32) -> Result<Vec<u8>> {
self.flush_buf()?;
let mut value = vec![0; value_len as usize];
self.file.seek(SeekFrom::Start(value_pos))?;
self.file.read_exact(&mut value)?;
Ok(value)
}
fn encode_entry_to(
out: &mut Vec<u8>,
base_offset: u64,
key: &[u8],
value: Option<&[u8]>,
hard_delete: bool,
) -> (u64, u32, i32) {
let key_len = key.len() as u32;
let (flag, value_len, payload): (i32, u32, Option<&[u8]>) = if hard_delete {
(FLAG_HARD_DELETE, 0, None)
} else {
match value {
Some(v) => (v.len() as i32, v.len() as u32, Some(v)),
None => (FLAG_SOFT_TOMBSTONE, 0, None),
}
};
let value_pos = base_offset + KEY_VAL_HEADER_LEN as u64 * 2 + key_len as u64;
out.extend_from_slice(&key_len.to_be_bytes());
out.extend_from_slice(&flag.to_be_bytes());
out.extend_from_slice(key);
if let Some(v) = payload {
out.extend_from_slice(v);
}
(value_pos, value_len, flag)
}
fn append_entry(
&mut self,
key: &[u8],
value: Option<&[u8]>,
hard_delete: bool,
) -> Result<(u64, u32, i32)> {
let base = self.end_pos;
let before = self.buf.len();
let (value_pos, value_len, flag) =
Self::encode_entry_to(&mut self.buf, base, key, value, hard_delete);
let written = (self.buf.len() - before) as u64;
self.end_pos += written;
if self.buf.len() >= WRITE_BUF_TARGET {
self.flush_buf()?;
}
Ok((value_pos, value_len, flag))
}
pub fn write_entry(&mut self, key: &[u8], value: Option<&[u8]>) -> Result<(u64, u32)> {
match value {
Some(v) => {
let (pos, len, _) = self.append_entry(key, Some(v), false)?;
let entry_len = KEY_VAL_HEADER_LEN * 2 + key.len() as u32 + len;
let offset = pos - KEY_VAL_HEADER_LEN as u64 * 2 - key.len() as u64;
Ok((offset, entry_len))
}
None => self.write_hard_delete(key),
}
}
pub fn write_hard_delete(&mut self, key: &[u8]) -> Result<(u64, u32)> {
let (pos, _, _) = self.append_entry(key, None, true)?;
let entry_len = KEY_VAL_HEADER_LEN * 2 + key.len() as u32;
let offset = pos - KEY_VAL_HEADER_LEN as u64 * 2 - key.len() as u64;
Ok((offset, entry_len))
}
pub fn write_soft_tombstone(&mut self, key: &[u8]) -> Result<(u64, u32)> {
let (pos, _, _) = self.append_entry(key, None, false)?;
let entry_len = KEY_VAL_HEADER_LEN * 2 + key.len() as u32;
let offset = pos - KEY_VAL_HEADER_LEN as u64 * 2 - key.len() as u64;
Ok((offset, entry_len))
}
}
pub struct ScanIterator<'a> {
items: Vec<(Vec<u8>, KeyLoc)>,
idx: usize,
rev: bool,
eng: &'a mut Bitcask,
}
impl<'a> ScanIterator<'a> {
fn read_item(&mut self, key: Vec<u8>, loc: KeyLoc) -> Option<<Self as Iterator>::Item> {
match loc {
KeyLoc::SoftTombstone { .. } => None,
KeyLoc::Live {
file_id,
pos,
len,
} => {
let value = self.eng.read_live(file_id, pos, len);
Some(value.map(|v| (key, v)))
}
}
}
}
impl<'a> Iterator for ScanIterator<'a> {
type Item = Result<(Vec<u8>, Vec<u8>)>;
fn next(&mut self) -> Option<Self::Item> {
if self.rev {
while self.idx < self.items.len() {
let (key, loc) = self.items[self.idx].clone();
self.idx += 1;
if let Some(mapped) = self.read_item(key, loc) {
return Some(mapped);
}
}
return None;
}
while self.idx < self.items.len() {
let (key, loc) = self.items[self.idx].clone();
self.idx += 1;
if let Some(mapped) = self.read_item(key, loc) {
return Some(mapped);
}
}
None
}
}
impl<'a> DoubleEndedIterator for ScanIterator<'a> {
fn next_back(&mut self) -> Option<Self::Item> {
self.rev = true;
while !self.items.is_empty() && self.idx < self.items.len() {
let last = self.items.len() - 1;
if last < self.idx {
break;
}
let (key, loc) = self.items.remove(last);
if let Some(mapped) = self.read_item(key, loc) {
return Some(mapped);
}
}
None
}
}
pub struct Bitcask {
base_path: PathBuf,
active: Log,
active_id: u32,
key_dir: KeyDir,
max_file_bytes: u64,
sealed: std::collections::HashMap<u32, Log>,
}
impl Drop for Bitcask {
fn drop(&mut self) {
if let Err(e) = self.flush() {
log::error!("failed to flush file: {e:?}");
}
if let Err(e) = self.dump_hint() {
log::error!("dump_hint on drop failed: {e:?}");
}
}
}
impl Bitcask {
pub fn new(path: PathBuf) -> Result<Self> {
Self::with_max_file_bytes(path, DEFAULT_MAX_FILE_BYTES)
}
pub fn with_max_file_bytes(path: PathBuf, max_file_bytes: u64) -> Result<Self> {
let base_path = path;
let sealed_ids = discover_sealed_ids(&base_path);
let hint_path = hint_path_for(&base_path);
let mut from_hint: Option<(KeyDir, u32, u64)> = None;
if let Ok((kd, active_id, active_size)) = read_hint_file(&hint_path) {
from_hint = Some((kd, active_id, active_size));
}
let (key_dir, active_id, active) = if let Some((mut kd, active_id, active_size)) = from_hint
{
let mut active = Log::new(base_path.clone())?;
let file_len = active.file.metadata()?.len();
if active_size > file_len {
log::warn!(
"hint active_size={} > data.log len={},回退全量扫",
active_size,
file_len
);
drop(active);
Self::full_rebuild(&base_path, &sealed_ids)?
} else if let Err(e) = active.load_index_from(&mut kd, active_id, active_size) {
log::warn!("hint 增量扫失败: {e},回退全量扫");
drop(active);
Self::full_rebuild(&base_path, &sealed_ids)?
} else {
(kd, active_id, active)
}
} else {
Self::full_rebuild(&base_path, &sealed_ids)?
};
Ok(Self {
base_path,
active,
active_id,
key_dir,
max_file_bytes,
sealed: std::collections::HashMap::new(),
})
}
fn full_rebuild(base_path: &Path, sealed_ids: &[u32]) -> Result<(KeyDir, u32, Log)> {
let mut key_dir = KeyDir::new();
for &id in sealed_ids {
let p = segment_path(base_path, id);
if !p.exists() {
continue;
}
let mut log = Log::new(p)?;
let partial = log.load_index(id)?;
for (k, v) in partial {
key_dir.insert(k, v);
}
drop(log);
}
let active = Log::new(base_path.to_path_buf())?;
let active_id = if sealed_ids.is_empty() {
0
} else {
sealed_ids.iter().copied().max().unwrap_or(0).saturating_add(1)
};
let partial = {
let mut active_mut = active;
let kd = active_mut.load_index(active_id)?;
for (k, v) in kd {
key_dir.insert(k, v);
}
(key_dir, active_id, active_mut)
};
Ok(partial)
}
pub fn path(&self) -> &std::path::Path {
&self.base_path
}
pub fn active_id(&self) -> u32 {
self.active_id
}
pub fn hint_path(&self) -> PathBuf {
hint_path_for(&self.base_path)
}
pub fn dump_hint(&mut self) -> Result<()> {
self.active.flush_buf()?;
let log_size = self.active.end_pos;
let path = self.hint_path();
write_hint_file(&path, &self.key_dir, self.active_id, log_size)
}
fn read_live(&mut self, file_id: u32, pos: u64, len: u32) -> Result<Vec<u8>> {
if file_id == self.active_id {
return self.active.read_value(pos, len);
}
if !self.sealed.contains_key(&file_id) {
let p = segment_path(&self.base_path, file_id);
let log = Log::new(p)?;
self.sealed.insert(file_id, log);
}
let log = self
.sealed
.get_mut(&file_id)
.expect("sealed just inserted");
log.read_value(pos, len)
}
fn maybe_rotate(&mut self) -> Result<()> {
if self.active.end_pos < self.max_file_bytes {
return Ok(());
}
self.rotate()
}
fn rotate(&mut self) -> Result<()> {
self.active.flush_buf()?;
let _ = self.active.file.sync_all();
let sealed_path = segment_path(&self.base_path, self.active_id);
let old = std::mem::replace(
&mut self.active,
Log {
path: PathBuf::new(),
file: {
let hold = self.base_path.with_extension("rotate_hold");
std::fs::OpenOptions::new()
.read(true)
.write(true)
.create(true)
.open(&hold)?
},
end_pos: 0,
buf: Vec::new(),
},
);
let old_path = old.path.clone();
drop(old);
if sealed_path.exists() {
let _ = std::fs::remove_file(&sealed_path);
}
std::fs::rename(&old_path, &sealed_path)?;
self.active = Log::new(self.base_path.clone())?;
self.active_id = self.active_id.saturating_add(1);
let hold = self.base_path.with_extension("rotate_hold");
let _ = std::fs::remove_file(&hold);
Ok(())
}
pub fn merge(&mut self) -> Result<()> {
self.active.flush_buf()?;
let target = self.base_path.clone();
let mut merge_path = target.clone();
merge_path.set_extension(MERGE_FILE_EXT);
let snapshot: Vec<(Vec<u8>, KeyLoc)> = self
.key_dir
.iter()
.map(|(k, loc)| (k.clone(), loc.clone()))
.collect();
let mut new_log = Log::new(merge_path)?;
let mut new_key_dir = KeyDir::new();
let new_file_id = 0u32;
for (key, loc) in snapshot {
match loc {
KeyLoc::Live {
file_id,
pos,
len,
} => {
let value = self.read_live(file_id, pos, len)?;
let (value_pos, value_len, _) =
new_log.append_entry(&key, Some(&value), false)?;
new_key_dir.insert(
key,
KeyLoc::Live {
file_id: new_file_id,
pos: value_pos,
len: value_len,
},
);
}
KeyLoc::SoftTombstone { .. } => {
new_log.write_soft_tombstone(&key)?;
new_key_dir.insert(key, KeyLoc::SoftTombstone { file_id: new_file_id });
}
}
}
new_log.flush_buf()?;
new_log.file.sync_all()?;
let new_path = new_log.path.clone();
let new_end = new_log.end_pos;
drop(new_log);
self.sealed.clear();
let hold_path = target.with_extension("merge_hold");
let old_active = std::mem::replace(
&mut self.active,
Log {
path: PathBuf::new(),
file: {
let f = std::fs::OpenOptions::new()
.read(true)
.write(true)
.create(true)
.open(&hold_path)?;
f
},
end_pos: 0,
buf: Vec::new(),
},
);
drop(old_active);
for id in discover_sealed_ids(&target) {
let p = segment_path(&target, id);
let _ = std::fs::remove_file(&p);
}
if target.exists() {
let _ = std::fs::remove_file(&target);
}
std::fs::rename(&new_path, &target)?;
let file = std::fs::OpenOptions::new()
.read(true)
.write(true)
.open(&target)?;
file.try_lock_exclusive()?;
self.active = Log {
path: target,
file,
end_pos: new_end,
buf: Vec::with_capacity(WRITE_BUF_TARGET),
};
self.active_id = 0;
self.key_dir = new_key_dir;
let _ = std::fs::remove_file(&hold_path);
let _ = self.dump_hint();
Ok(())
}
pub fn set(&mut self, key: &[u8], value: Vec<u8>) -> Result<()> {
let fid = self.active_id;
let (value_pos, value_len, _) = self.active.append_entry(key, Some(&value), false)?;
self.key_dir.insert(
key.to_vec(),
KeyLoc::Live {
file_id: fid,
pos: value_pos,
len: value_len,
},
);
self.maybe_rotate()?;
Ok(())
}
pub fn get(&mut self, key: &[u8]) -> Result<Option<Vec<u8>>> {
match self.key_dir.get(key).cloned() {
Some(KeyLoc::Live {
file_id,
pos,
len,
}) => {
let val = self.read_live(file_id, pos, len)?;
Ok(Some(val))
}
Some(KeyLoc::SoftTombstone { .. }) | None => Ok(None),
}
}
pub fn delete(&mut self, key: &[u8]) -> Result<()> {
self.active.write_hard_delete(key)?;
self.key_dir.remove(key);
self.maybe_rotate()?;
Ok(())
}
pub fn flush(&mut self) -> Result<()> {
self.active.flush_buf()?;
self.active.file.sync_all()?;
Ok(())
}
pub fn flush_buf(&mut self) -> Result<()> {
self.active.flush_buf()
}
pub fn scan(&mut self, range: impl std::ops::RangeBounds<Vec<u8>>) -> ScanIterator<'_> {
let _ = self.active.flush_buf();
let items: Vec<(Vec<u8>, KeyLoc)> = self
.key_dir
.range(range)
.map(|(k, loc)| (k.clone(), loc.clone()))
.collect();
ScanIterator {
items,
idx: 0,
rev: false,
eng: self,
}
}
pub fn scan_prefix(&mut self, prefix: &[u8]) -> ScanIterator<'_> {
let start = Bound::Included(prefix.to_vec());
let mut bound_prefix = prefix.to_vec();
if let Some(last) = bound_prefix.iter_mut().last() {
*last = last.saturating_add(1);
}
let end = Bound::Excluded(bound_prefix);
self.scan((start, end))
}
pub fn insert(&mut self, key: Vec<u8>, value: Option<Vec<u8>>) -> Result<()> {
let fid = self.active_id;
match value {
Some(v) => {
let (value_pos, value_len, _) =
self.active.append_entry(&key, Some(&v), false)?;
self.key_dir.insert(
key,
KeyLoc::Live {
file_id: fid,
pos: value_pos,
len: value_len,
},
);
}
None => {
self.active.write_soft_tombstone(&key)?;
self.key_dir
.insert(key, KeyLoc::SoftTombstone { file_id: fid });
}
}
self.maybe_rotate()?;
Ok(())
}
pub fn insert_batch(
&mut self,
items: impl IntoIterator<Item = (Vec<u8>, Option<Vec<u8>>)>,
) -> Result<()> {
for (key, value) in items {
let fid = self.active_id;
match value {
Some(v) => {
let (value_pos, value_len, _) =
self.active.append_entry(&key, Some(&v), false)?;
self.key_dir.insert(
key,
KeyLoc::Live {
file_id: fid,
pos: value_pos,
len: value_len,
},
);
}
None => {
self.active.write_soft_tombstone(&key)?;
self.key_dir
.insert(key, KeyLoc::SoftTombstone { file_id: fid });
}
}
self.maybe_rotate()?;
}
Ok(())
}
pub fn remove(&mut self, key: &[u8]) -> Result<()> {
self.active.write_hard_delete(key)?;
self.key_dir.remove(key);
self.maybe_rotate()?;
Ok(())
}
pub fn get_entry(&mut self, key: &[u8]) -> Result<Option<Option<Vec<u8>>>> {
match self.key_dir.get(key).cloned() {
None => Ok(None),
Some(KeyLoc::SoftTombstone { .. }) => Ok(Some(None)),
Some(KeyLoc::Live {
file_id,
pos,
len,
}) => {
let v = self.read_live(file_id, pos, len)?;
Ok(Some(Some(v)))
}
}
}
pub fn range_scan(
&mut self,
low: Vec<u8>,
high: Vec<u8>,
) -> Result<Vec<(Vec<u8>, Option<Vec<u8>>)>> {
self.active.flush_buf()?;
let keys: Vec<(Vec<u8>, KeyLoc)> = self
.key_dir
.range(low..=high)
.map(|(k, loc)| (k.clone(), loc.clone()))
.collect();
let mut out = Vec::with_capacity(keys.len());
for (k, loc) in keys {
match loc {
KeyLoc::SoftTombstone { .. } => out.push((k, None)),
KeyLoc::Live {
file_id,
pos,
len,
} => {
let v = self.read_live(file_id, pos, len)?;
out.push((k, Some(v)));
}
}
}
Ok(out)
}
pub fn iter(&mut self) -> Result<Vec<(Vec<u8>, Option<Vec<u8>>)>> {
self.active.flush_buf()?;
let keys: Vec<(Vec<u8>, KeyLoc)> = self
.key_dir
.iter()
.map(|(k, loc)| (k.clone(), loc.clone()))
.collect();
let mut out = Vec::with_capacity(keys.len());
for (k, loc) in keys {
match loc {
KeyLoc::SoftTombstone { .. } => out.push((k, None)),
KeyLoc::Live {
file_id,
pos,
len,
} => {
let v = self.read_live(file_id, pos, len)?;
out.push((k, Some(v)));
}
}
}
Ok(out)
}
pub fn iter_meta(&self) -> impl Iterator<Item = (&Vec<u8>, &KeyLoc)> {
self.key_dir.iter()
}
pub fn len(&self) -> usize {
self.key_dir.len()
}
pub fn is_empty(&self) -> bool {
self.key_dir.is_empty()
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::time::{SystemTime, UNIX_EPOCH};
fn tmp_log(tag: &str) -> PathBuf {
let nanos = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap()
.as_nanos();
std::env::temp_dir().join(format!(
"bitcask_hint_{tag}_{}_{}.log",
std::process::id(),
nanos
))
}
fn cleanup(path: &Path) {
let _ = std::fs::remove_file(path);
let _ = std::fs::remove_file(hint_path_for(path));
let _ = std::fs::remove_file(path.with_extension("merge"));
let _ = std::fs::remove_file(path.with_extension("merge_hold"));
let _ = std::fs::remove_file(path.with_extension("rotate_hold"));
for id in discover_sealed_ids(path) {
let _ = std::fs::remove_file(segment_path(path, id));
}
}
#[test]
fn test_hint_roundtrip_and_fast_open() {
let path = tmp_log("rt");
{
let mut eng = Bitcask::new(path.clone()).unwrap();
eng.set(b"a", b"1".to_vec()).unwrap();
eng.set(b"b", b"2".to_vec()).unwrap();
eng.set(b"c", b"3".to_vec()).unwrap();
eng.delete(b"b").unwrap();
eng.flush().unwrap();
eng.dump_hint().unwrap();
assert!(eng.hint_path().exists());
assert_eq!(eng.len(), 2);
}
{
let mut eng = Bitcask::new(path.clone()).unwrap();
assert_eq!(eng.get(b"a").unwrap(), Some(b"1".to_vec()));
assert_eq!(eng.get(b"b").unwrap(), None);
assert_eq!(eng.get(b"c").unwrap(), Some(b"3".to_vec()));
assert_eq!(eng.len(), 2);
}
cleanup(&path);
}
#[test]
fn test_hint_incremental_tail_scan() {
let path = tmp_log("incr");
{
let mut eng = Bitcask::new(path.clone()).unwrap();
eng.set(b"k1", b"v1".to_vec()).unwrap();
eng.set(b"k2", b"v2".to_vec()).unwrap();
eng.flush().unwrap();
eng.dump_hint().unwrap();
}
{
let mut eng = Bitcask::new(path.clone()).unwrap();
eng.set(b"k3", b"v3".to_vec()).unwrap();
eng.set(b"k1", b"v1b".to_vec()).unwrap(); eng.flush().unwrap();
drop(eng);
}
{
let mut eng = Bitcask::new(path.clone()).unwrap();
assert_eq!(eng.get(b"k1").unwrap(), Some(b"v1b".to_vec()));
assert_eq!(eng.get(b"k2").unwrap(), Some(b"v2".to_vec()));
assert_eq!(eng.get(b"k3").unwrap(), Some(b"v3".to_vec()));
assert_eq!(eng.len(), 3);
}
cleanup(&path);
}
#[test]
fn test_merge_writes_hint() {
let path = tmp_log("merge");
{
let mut eng = Bitcask::new(path.clone()).unwrap();
eng.set(b"a", b"1".to_vec()).unwrap();
eng.set(b"a", b"2".to_vec()).unwrap();
eng.set(b"b", b"3".to_vec()).unwrap();
eng.delete(b"b").unwrap();
eng.merge().unwrap();
assert!(eng.hint_path().exists());
assert_eq!(eng.get(b"a").unwrap(), Some(b"2".to_vec()));
assert_eq!(eng.get(b"b").unwrap(), None);
}
{
let mut eng = Bitcask::new(path.clone()).unwrap();
assert_eq!(eng.get(b"a").unwrap(), Some(b"2".to_vec()));
assert_eq!(eng.len(), 1);
}
cleanup(&path);
}
#[test]
fn test_corrupt_hint_falls_back() {
let path = tmp_log("bad");
{
let mut eng = Bitcask::new(path.clone()).unwrap();
eng.set(b"x", b"y".to_vec()).unwrap();
eng.flush().unwrap();
eng.dump_hint().unwrap();
}
std::fs::write(hint_path_for(&path), b"GARBAGE!!").unwrap();
{
let mut eng = Bitcask::new(path.clone()).unwrap();
assert_eq!(eng.get(b"x").unwrap(), Some(b"y".to_vec()));
}
cleanup(&path);
}
#[test]
fn test_soft_tombstone_in_hint() {
let path = tmp_log("soft");
{
let mut eng = Bitcask::new(path.clone()).unwrap();
eng.insert(b"vkey".to_vec(), Some(b"val".to_vec()))
.unwrap();
eng.insert(b"tomb".to_vec(), None).unwrap(); eng.flush().unwrap();
eng.dump_hint().unwrap();
assert_eq!(eng.len(), 2);
assert_eq!(eng.get_entry(b"tomb").unwrap(), Some(None));
}
{
let mut eng = Bitcask::new(path.clone()).unwrap();
assert_eq!(eng.get(b"vkey").unwrap(), Some(b"val".to_vec()));
assert_eq!(eng.get_entry(b"tomb").unwrap(), Some(None));
assert_eq!(eng.len(), 2);
}
cleanup(&path);
}
#[test]
fn test_rotate_and_reopen() {
let path = tmp_log("rotate");
{
let mut eng = Bitcask::with_max_file_bytes(path.clone(), 256).unwrap();
for i in 0..40u32 {
let k = format!("k{i:03}");
let v = format!("value-padding-{:03}-xxxxxxxx", i);
eng.set(k.as_bytes(), v.into_bytes()).unwrap();
}
eng.flush().unwrap();
eng.dump_hint().unwrap();
assert!(
eng.active_id() > 0,
"expected rotate, active_id={}",
eng.active_id()
);
let sealed = discover_sealed_ids(&path);
assert!(!sealed.is_empty(), "expected sealed segments");
assert_eq!(eng.get(b"k000").unwrap(), Some(b"value-padding-000-xxxxxxxx".to_vec()));
assert_eq!(eng.get(b"k039").unwrap(), Some(b"value-padding-039-xxxxxxxx".to_vec()));
}
{
let mut eng = Bitcask::new(path.clone()).unwrap();
assert_eq!(eng.get(b"k000").unwrap(), Some(b"value-padding-000-xxxxxxxx".to_vec()));
assert_eq!(eng.get(b"k020").unwrap(), Some(b"value-padding-020-xxxxxxxx".to_vec()));
assert_eq!(eng.get(b"k039").unwrap(), Some(b"value-padding-039-xxxxxxxx".to_vec()));
assert_eq!(eng.len(), 40);
}
{
let mut eng = Bitcask::new(path.clone()).unwrap();
eng.merge().unwrap();
assert_eq!(eng.active_id(), 0);
assert!(discover_sealed_ids(&path).is_empty());
assert_eq!(eng.get(b"k010").unwrap(), Some(b"value-padding-010-xxxxxxxx".to_vec()));
}
cleanup(&path);
}
}