mod barrier;
mod ops;
mod snapshot;
use std::{
fs,
io::Read as _,
path::Path,
sync::{
Arc,
atomic::{AtomicBool, AtomicPtr, AtomicU8, AtomicUsize, Ordering},
},
};
use bf_tree::{BfTree, ConfigError};
pub use ops::SCAN_ALL_START_KEY;
use parking_lot::RwLock;
use wbase::backoff::backoff;
use crate::{
error::{Error, Result},
manager::CPR_MAGIC,
types::{BfTreeConfig, StorageBackendType},
};
const STACK_READ_BUF_SIZE: usize = 4096;
pub(crate) const STACK_SCAN_BUF_SIZE: usize = 8192;
const PRESET_LEAF_PAGE_SIZE: usize = 16384;
const PRESET_MAX_RECORD_SIZE: usize = 4096;
const PRESET_MAX_KEY_LEN: usize = 512;
const PRESET_MIN_RECORD_SIZE: usize = 4;
pub(crate) const MIN_MAX_RECORD_SIZE: usize = STACK_READ_BUF_SIZE;
#[inline]
pub fn file_has_cpr_magic(path: &Path) -> bool {
let Ok(mut file) = fs::File::open(path) else {
return false;
};
let mut magic = [0u8; CPR_MAGIC.len()];
file.read_exact(&mut magic).is_ok() && magic == *CPR_MAGIC
}
#[inline]
fn config_error_to_string(e: ConfigError) -> String {
match e {
ConfigError::MinimumRecordSize(s) => format!("MinimumRecordSize: {s}"),
ConfigError::MaximumRecordSize(s) => format!("MaximumRecordSize: {s}"),
ConfigError::LeafPageSize(s) => format!("LeafPageSize: {s}"),
ConfigError::MaxKeyLen(s) => format!("MaxKeyLen: {s}"),
ConfigError::CircularBufferSize(s) => format!("CircularBufferSize: {s}"),
ConfigError::SnapshotFileInvalid(s) => format!("SnapshotFileInvalid: {s}"),
ConfigError::SnapshotDisabled => "SnapshotDisabled".to_string(),
}
}
pub struct BfTreeService {
pub(crate) raw_tree: AtomicPtr<BfTree>,
pub(crate) arc_tree: RwLock<Option<Arc<BfTree>>>,
pub(crate) retired_trees: RwLock<Vec<Arc<BfTree>>>,
pub(crate) storage_backend: AtomicU8,
pub(crate) file_path: RwLock<Option<String>>,
pub(crate) max_record_size: AtomicUsize,
pub(crate) disposed: AtomicBool,
pub(crate) barriers: AtomicUsize,
}
unsafe impl Send for BfTreeService {}
unsafe impl Sync for BfTreeService {}
pub struct WriteBarrierGuard<'a> {
pub(crate) service: &'a BfTreeService,
}
impl Drop for WriteBarrierGuard<'_> {
fn drop(&mut self) {
self.service.barriers.fetch_sub(1, Ordering::Release);
}
}
impl BfTreeService {
#[cold]
#[inline(never)]
pub(crate) fn wait_for_barrier(&self) {
let mut spins = 0u32;
while self.barriers.load(Ordering::Acquire) != 0 {
backoff(spins);
spins = spins.wrapping_add(1);
}
}
pub(crate) fn new_with_backend(
config: impl Into<bf_tree::Config>,
storage_backend: StorageBackendType,
file_path: Option<String>,
) -> Result<Self> {
if storage_backend == StorageBackendType::Disk && file_path.is_none() {
return Err(Error::InvalidArgument(
"磁盘后端必须指定数据文件路径 (file_path)".into(),
));
}
let inner_cfg: bf_tree::Config = config.into();
let max_record_size = inner_cfg.get_cb_max_record_size().max(MIN_MAX_RECORD_SIZE);
let tree = match BfTree::with_config(inner_cfg, None) {
Ok(t) => Arc::new(t),
Err(e) => return Err(Error::InvalidConfig(config_error_to_string(e))),
};
let raw_ptr = Arc::as_ptr(&tree) as *mut BfTree;
Ok(Self {
raw_tree: AtomicPtr::new(raw_ptr),
arc_tree: RwLock::new(Some(tree)),
retired_trees: RwLock::new(Vec::new()),
storage_backend: AtomicU8::new(storage_backend as u8),
file_path: RwLock::new(file_path),
max_record_size: AtomicUsize::new(max_record_size),
disposed: AtomicBool::new(false),
barriers: AtomicUsize::new(0),
})
}
pub fn new(config: BfTreeConfig) -> Result<Self> {
let storage_backend = config.storage_backend;
let file_path = config.file_path;
Self::new_with_backend(config.inner, storage_backend, file_path)
}
pub fn preset_config(cb_min_record_size: usize) -> BfTreeConfig {
let mut config = BfTreeConfig::default();
config
.use_snapshot(true)
.leaf_page_size(PRESET_LEAF_PAGE_SIZE)
.cb_max_record_size(PRESET_MAX_RECORD_SIZE)
.cb_max_key_len(PRESET_MAX_KEY_LEN)
.cb_min_record_size(if cb_min_record_size > 0 {
cb_min_record_size
} else {
PRESET_MIN_RECORD_SIZE
})
.scan_promotion_rate(0);
config
}
pub fn open_disk(path: impl AsRef<Path>, cb_min_record_size: usize) -> Result<Self> {
let p = path.as_ref();
if p.as_os_str().is_empty() {
return Err(Error::InvalidArgument(
"磁盘后端必须指定有效的数据文件路径".into(),
));
}
if let Some(parent) = p.parent()
&& !parent.as_os_str().is_empty()
{
fs::create_dir_all(parent)?;
}
let mut config = Self::preset_config(cb_min_record_size);
config.file_path(p);
Self::new_with_backend(
config,
StorageBackendType::Disk,
Some(p.to_string_lossy().into_owned()),
)
}
pub fn open_memory(cb_min_record_size: usize) -> Result<Self> {
let mut config = Self::preset_config(cb_min_record_size);
config.cache_only(true);
Self::new_with_backend(config, StorageBackendType::Memory, None)
}
#[inline]
pub fn has_tree(&self) -> bool {
!self.raw_tree.load(Ordering::Acquire).is_null()
}
#[inline]
pub(crate) fn tree_arc(&self) -> Result<Arc<BfTree>> {
self.check_disposed()?;
self
.arc_tree
.read()
.as_ref()
.cloned()
.ok_or(Error::Disposed)
}
#[inline(always)]
pub(crate) fn tree_ref(&self) -> Result<&BfTree> {
let ptr = self.raw_tree.load(Ordering::Acquire);
if ptr.is_null() {
return Err(Error::Disposed);
}
Ok(unsafe { &*ptr })
}
#[inline]
pub fn native_ptr(&self) -> u64 {
self.raw_tree.load(Ordering::Relaxed) as usize as u64
}
#[inline]
pub fn file_path(&self) -> Option<String> {
self.file_path.read().clone()
}
#[inline]
pub fn storage_backend(&self) -> StorageBackendType {
StorageBackendType::from_u8(self.storage_backend.load(Ordering::Acquire))
}
#[inline]
pub(crate) fn max_record_size(&self) -> usize {
self.max_record_size.load(Ordering::Relaxed)
}
#[inline]
pub fn is_disposed(&self) -> bool {
self.disposed.load(Ordering::Acquire)
}
#[inline]
pub(crate) fn check_disposed(&self) -> Result<()> {
if self.is_disposed() {
Err(Error::Disposed)
} else {
Ok(())
}
}
}
impl Drop for BfTreeService {
fn drop(&mut self) {
self.dispose();
}
}
#[cfg(test)]
mod tests {
use std::{env, fs, process, sync::Arc, thread::spawn};
use super::*;
use crate::types::{
BfTreeConfig, BfTreeDeleteResult, BfTreeInsertResult, BfTreeReadResult, StorageBackendType,
};
fn mem_service() -> BfTreeService {
BfTreeService::open_memory(0).unwrap()
}
#[test]
fn test_insert_empty_value_rejected_without_engine_panic() {
let service = mem_service();
assert_eq!(
service.insert(b"long_enough_key", b""),
BfTreeInsertResult::InvalidKV
);
assert_eq!(service.insert(b"k", b""), BfTreeInsertResult::InvalidKV);
let (res, v) = service.read(b"long_enough_key");
assert_eq!(res, BfTreeReadResult::NotFound);
assert_eq!(v, None);
}
#[test]
fn test_dispose_lifecycle() {
let service = Arc::new(mem_service());
assert_eq!(
service.insert(b"key1", b"val1"),
BfTreeInsertResult::Success
);
assert!(!service.is_disposed());
assert!(service.has_tree());
service.dispose();
assert!(service.is_disposed());
assert!(!service.has_tree());
assert_eq!(
service.insert(b"key2", b"val2"),
BfTreeInsertResult::InvalidArguments
);
assert_eq!(
service.delete(b"key1"),
BfTreeDeleteResult::InvalidArguments
);
service.dispose();
assert!(service.is_disposed());
}
#[test]
fn test_cpr_snapshot_concurrent_with_writers() {
let service = Arc::new(mem_service());
for i in 0..100u32 {
let k = format!("base{i:04}");
assert_eq!(
service.insert(k.as_bytes(), b"base_value"),
BfTreeInsertResult::Success
);
}
let dir = env::temp_dir().join(format!(
"wbftree_cpr_concurrent_{}_{}",
process::id(),
fastrand::u64(..)
));
fs::create_dir_all(&dir).unwrap();
let snap = dir.join("snap.bftree");
let wsvc = Arc::clone(&service);
let writer = spawn(move || {
for i in 0..5000u32 {
let k = format!("live{i:05}");
assert_eq!(
wsvc.insert(k.as_bytes(), b"payload"),
BfTreeInsertResult::Success
);
}
});
service.cpr_snapshot(&snap).unwrap();
writer.join().unwrap();
let recovered =
BfTreeService::recover_from_cpr_snapshot(&snap, true, StorageBackendType::Disk).unwrap();
for i in 0..100u32 {
let k = format!("base{i:04}");
let (res, v) = recovered.read(k.as_bytes());
assert_eq!(res, BfTreeReadResult::Found, "基线键 {k} 必须在快照中");
assert_eq!(v.as_deref(), Some(&b"base_value"[..]));
}
fs::remove_dir_all(&dir).unwrap();
}
#[test]
fn test_recover_from_corrupt_snapshot_returns_err() {
let dir = env::temp_dir().join(format!(
"wbftree_corrupt_{}_{}",
process::id(),
fastrand::u64(..)
));
fs::create_dir_all(&dir).unwrap();
let bad_magic = dir.join("bad_magic.bftree");
fs::write(&bad_magic, b"garbage payload").unwrap();
let err = BfTreeService::recover_from_cpr_snapshot(&bad_magic, true, StorageBackendType::Disk);
assert!(matches!(err, Err(Error::Recovery(_))));
let truncated = dir.join("truncated.bftree");
fs::write(&truncated, b"BF-TREE").unwrap();
let err = BfTreeService::recover_from_cpr_snapshot(&truncated, true, StorageBackendType::Disk);
assert!(matches!(err, Err(Error::Recovery(_))));
let wrong_magic = dir.join("wrong_magic.bftree");
fs::write(&wrong_magic, b"XX-TREE-V0-BEGIN_PAYLOAD").unwrap();
let err =
BfTreeService::recover_from_cpr_snapshot(&wrong_magic, true, StorageBackendType::Disk);
assert!(matches!(err, Err(Error::Recovery(_))));
fs::remove_dir_all(&dir).unwrap();
}
#[test]
fn test_recover_in_place_republishes_max_record_size() {
let dir = env::temp_dir().join(format!(
"wbftree_swap_max_{}_{}",
process::id(),
fastrand::u64(..)
));
fs::create_dir_all(&dir).unwrap();
let src_work = dir.join("src.data.bftree");
let snap = dir.join("snap.bftree");
{
let mut config = BfTreeConfig::default();
config
.use_snapshot(true)
.leaf_page_size(32768)
.cb_max_record_size(8192)
.cb_max_key_len(PRESET_MAX_KEY_LEN)
.cb_min_record_size(8);
config.file_path(&src_work);
let src = BfTreeService::new(config).unwrap();
let big = [b'x'; 6000];
assert_eq!(src.insert(b"big_key", &big), BfTreeInsertResult::Success);
src.cpr_snapshot(&snap).unwrap();
}
let work = dir.join("work.data.bftree");
let target = BfTreeService::open_disk(&work, 0).unwrap();
assert_eq!(target.max_record_size(), PRESET_MAX_RECORD_SIZE);
let big = [b'x'; 6000];
assert_eq!(
target.insert(b"big_key", &big),
BfTreeInsertResult::InvalidKV
);
target.recover_in_place(&snap, &work).unwrap();
assert_eq!(target.max_record_size(), 8192);
let (res, v) = target.read(b"big_key");
assert_eq!(res, BfTreeReadResult::Found);
assert_eq!(v.as_deref(), Some(&big[..]));
assert_eq!(target.insert(b"fresh", &big), BfTreeInsertResult::Success);
let (res, v) = target.read(b"fresh");
assert_eq!(res, BfTreeReadResult::Found);
assert_eq!(v.as_deref(), Some(&big[..]));
let mut big_out = [0u8; 8192];
let (res, len) = target.read_into(b"big_key", &mut big_out);
assert_eq!(res, BfTreeReadResult::Found);
assert_eq!(&big_out[..len], &big[..]);
let mut small_out = [0u8; 64];
let (res, _) = target.read_into(b"big_key", &mut small_out);
assert_eq!(res, BfTreeReadResult::InvalidArguments);
fs::remove_dir_all(&dir).unwrap();
}
}