use crate::storage::file_manager::FileHandle;
use crate::{Result, StorageError};
use lru::LruCache;
use parking_lot::{Mutex, RwLock};
use serde::{Deserialize, Serialize};
use std::fs::{File, OpenOptions};
use std::io::{Read, Seek, SeekFrom, Write};
use std::num::NonZeroUsize;
use std::path::PathBuf;
use std::sync::Arc;
pub const BTREE_ORDER: usize = 255;
pub const MAX_PAGE_SIZE: usize = 15 + BTREE_ORDER * 8 * 2;
const PAGE_HEADER_SIZE: usize = 15;
pub const DEFAULT_PAGE_CACHE: usize = 1024;
const INVALID_PAGE_ID: u64 = u64::MAX;
const BTREE_MAGIC: u32 = 0x42545245;
const BTREE_VERSION: u32 = 2;
type InsertResult = Result<(Option<u64>, Option<(u64, u64)>)>;
#[derive(Serialize, Deserialize, Debug, Clone)]
struct SuperBlock {
magic: u32,
version: u32,
root_page_id: u64,
next_page_id: u64,
total_keys: usize,
total_pages: usize,
leaf_pages: usize,
internal_pages: usize,
tree_height: usize,
page_offsets: Vec<u64>,
}
pub struct BTree {
root_page_id: Arc<RwLock<u64>>,
page_cache: Arc<RwLock<LruCache<u64, Arc<RwLock<Page>>>>>,
next_page_id: Arc<RwLock<u64>>,
storage_file: Arc<RwLock<File>>,
flush_lock: Arc<Mutex<()>>,
_storage_path: PathBuf,
config: BTreeConfig,
stats: Arc<RwLock<BTreeStats>>,
page_offsets: Arc<RwLock<Vec<u64>>>,
_file_handle: Option<FileHandle>,
}
#[derive(Clone)]
pub struct BTreeConfig {
pub order: usize,
pub cache_size: usize,
pub unique_keys: bool,
pub allow_updates: bool,
pub immediate_sync: bool,
}
impl Default for BTreeConfig {
fn default() -> Self {
Self {
order: BTREE_ORDER,
cache_size: DEFAULT_PAGE_CACHE,
unique_keys: false,
allow_updates: true,
immediate_sync: false,
}
}
}
#[derive(Clone)]
struct Page {
page_id: u64,
is_leaf: bool,
num_keys: usize,
keys: Vec<u64>,
values: Vec<u64>,
children: Vec<u64>,
next_leaf: u64,
dirty: bool,
}
impl Page {
fn new_leaf(page_id: u64) -> Self {
Self {
page_id,
is_leaf: true,
num_keys: 0,
keys: Vec::with_capacity(BTREE_ORDER),
values: Vec::with_capacity(BTREE_ORDER),
children: Vec::new(),
next_leaf: INVALID_PAGE_ID,
dirty: true,
}
}
fn new_internal(page_id: u64) -> Self {
Self {
page_id,
is_leaf: false,
num_keys: 0,
keys: Vec::with_capacity(BTREE_ORDER),
values: Vec::new(),
children: Vec::with_capacity(BTREE_ORDER + 1),
next_leaf: INVALID_PAGE_ID,
dirty: true,
}
}
fn serialize_compact(&self) -> Result<Vec<u8>> {
let data_len = PAGE_HEADER_SIZE
+ self.num_keys * 8
+ if self.is_leaf {
self.num_keys * 8
} else {
(self.num_keys + 1) * 8
};
let mut buf = vec![0u8; data_len];
let mut offset = 0;
buf[offset] = if self.is_leaf { 1 } else { 0 };
offset += 1;
buf[offset..offset + 4].copy_from_slice(&(self.num_keys as u32).to_le_bytes());
offset += 4;
buf[offset..offset + 8].copy_from_slice(&self.next_leaf.to_le_bytes());
offset += 8;
buf[offset..offset + 2].copy_from_slice(&(data_len as u16).to_le_bytes());
offset += 2;
for &key in &self.keys {
buf[offset..offset + 8].copy_from_slice(&key.to_le_bytes());
offset += 8;
}
if self.is_leaf {
for &value in &self.values {
buf[offset..offset + 8].copy_from_slice(&value.to_le_bytes());
offset += 8;
}
} else {
for &child in &self.children {
buf[offset..offset + 8].copy_from_slice(&child.to_le_bytes());
offset += 8;
}
}
Ok(buf)
}
fn deserialize_compact(page_id: u64, buf: &[u8]) -> Result<Self> {
if buf.len() < PAGE_HEADER_SIZE {
return Err(StorageError::InvalidData(format!(
"Page buffer too small: {} < header size {}",
buf.len(),
PAGE_HEADER_SIZE
)));
}
let mut offset = 0;
let is_leaf = buf[offset] == 1;
offset += 1;
let num_keys = u32::from_le_bytes([
buf[offset],
buf[offset + 1],
buf[offset + 2],
buf[offset + 3],
]) as usize;
offset += 4;
if num_keys > BTREE_ORDER {
return Err(StorageError::Corruption(format!(
"Invalid num_keys in page {}: {} exceeds max {}",
page_id, num_keys, BTREE_ORDER
)));
}
let next_leaf = u64::from_le_bytes([
buf[offset],
buf[offset + 1],
buf[offset + 2],
buf[offset + 3],
buf[offset + 4],
buf[offset + 5],
buf[offset + 6],
buf[offset + 7],
]);
offset += 8;
let _content_len = u16::from_le_bytes([buf[offset], buf[offset + 1]]) as usize;
offset += 2;
let mut keys = Vec::with_capacity(num_keys);
for _ in 0..num_keys {
let key = u64::from_le_bytes([
buf[offset],
buf[offset + 1],
buf[offset + 2],
buf[offset + 3],
buf[offset + 4],
buf[offset + 5],
buf[offset + 6],
buf[offset + 7],
]);
keys.push(key);
offset += 8;
}
let mut values = Vec::new();
let mut children = Vec::new();
if is_leaf {
for _ in 0..num_keys {
let value = u64::from_le_bytes([
buf[offset],
buf[offset + 1],
buf[offset + 2],
buf[offset + 3],
buf[offset + 4],
buf[offset + 5],
buf[offset + 6],
buf[offset + 7],
]);
values.push(value);
offset += 8;
}
} else if num_keys > 0 {
for _ in 0..=num_keys {
let child = u64::from_le_bytes([
buf[offset],
buf[offset + 1],
buf[offset + 2],
buf[offset + 3],
buf[offset + 4],
buf[offset + 5],
buf[offset + 6],
buf[offset + 7],
]);
children.push(child);
offset += 8;
}
}
Ok(Self {
page_id,
is_leaf,
num_keys,
keys,
values,
children,
next_leaf,
dirty: false,
})
}
fn validate(&self) -> Result<()> {
if self.is_leaf {
if self.keys.len() != self.values.len() {
return Err(StorageError::Corruption(format!(
"Leaf page {} has mismatched keys ({}) and values ({})",
self.page_id,
self.keys.len(),
self.values.len()
)));
}
} else {
if self.num_keys > 0 && self.children.len() != self.num_keys + 1 {
return Err(StorageError::Corruption(format!(
"Internal page {} has mismatched keys ({}) and children ({})",
self.page_id,
self.num_keys,
self.children.len()
)));
}
if self.num_keys == 0 {
return Err(StorageError::Corruption(format!(
"Internal page {} has num_keys=0 (invalid state)",
self.page_id
)));
}
}
Ok(())
}
}
#[derive(Default, Debug, Clone)]
pub struct BTreeStats {
pub total_keys: usize,
pub total_pages: usize,
pub leaf_pages: usize,
pub internal_pages: usize,
pub tree_height: usize,
pub page_cache_hits: u64,
pub page_cache_misses: u64,
}
#[derive(Default, Debug, Clone)]
pub struct RangeQueryProfile {
pub find_leaf_us: u64, pub scan_us: u64, pub total_us: u64, pub pages_scanned: usize, pub keys_examined: usize, pub results_found: usize, }
impl BTree {
pub fn new(storage_path: PathBuf) -> Result<Self> {
Self::with_config(storage_path, BTreeConfig::default())
}
pub fn with_config(storage_path: PathBuf, config: BTreeConfig) -> Result<Self> {
if let Some(parent) = storage_path.parent() {
std::fs::create_dir_all(parent)?;
}
let mut file = OpenOptions::new()
.read(true)
.write(true)
.create(true)
.truncate(false)
.open(&storage_path)?;
let metadata = file.metadata()?;
let file_size = metadata.len();
let is_new_file = file_size == 0;
let (_superblock, root_page_id, next_page_id, stats, page_offsets) = if is_new_file {
let superblock = SuperBlock {
magic: BTREE_MAGIC,
version: BTREE_VERSION,
root_page_id: 0,
next_page_id: 1,
total_keys: 0,
total_pages: 0,
leaf_pages: 0,
internal_pages: 0,
tree_height: 0,
page_offsets: vec![0], };
Self::write_superblock(&mut file, &superblock)?;
let page_offsets = vec![0u64]; (superblock, 0, 1, BTreeStats::default(), page_offsets)
} else {
let superblock = Self::read_superblock(&mut file)?;
if superblock.magic != BTREE_MAGIC {
return Err(StorageError::Corruption(format!(
"Invalid B+Tree magic number: expected 0x{:08X}, got 0x{:08X}",
BTREE_MAGIC, superblock.magic
)));
}
if superblock.version != BTREE_VERSION {
drop(file);
let _ = std::fs::remove_file(&storage_path);
return Self::with_config(storage_path, config);
}
let root_id = superblock.root_page_id;
let next_id = superblock.next_page_id;
let stats = BTreeStats {
total_keys: superblock.total_keys,
total_pages: superblock.total_pages,
leaf_pages: superblock.leaf_pages,
internal_pages: superblock.internal_pages,
tree_height: superblock.tree_height,
page_cache_hits: 0,
page_cache_misses: 0,
};
let page_offsets = superblock.page_offsets.clone();
(superblock, root_id, next_id, stats, page_offsets)
};
Ok(Self {
root_page_id: Arc::new(RwLock::new(root_page_id)),
page_cache: Arc::new(RwLock::new(LruCache::new(
NonZeroUsize::new(config.cache_size.max(1)).unwrap(),
))),
next_page_id: Arc::new(RwLock::new(next_page_id)),
storage_file: Arc::new(RwLock::new(file)),
flush_lock: Arc::new(Mutex::new(())),
_storage_path: storage_path,
config,
stats: Arc::new(RwLock::new(stats)),
page_offsets: Arc::new(RwLock::new(page_offsets)),
_file_handle: None,
})
}
const SUPERBLOCK_SIZE: usize = 4096;
fn read_superblock(file: &mut File) -> Result<SuperBlock> {
file.seek(SeekFrom::Start(0))?;
let mut buf = vec![0u8; Self::SUPERBLOCK_SIZE];
file.read_exact(&mut buf)?;
bincode::deserialize(&buf).map_err(|e| {
StorageError::Corruption(format!("Failed to deserialize SuperBlock: {}", e))
})
}
fn write_superblock(file: &mut File, superblock: &SuperBlock) -> Result<()> {
file.seek(SeekFrom::Start(0))?;
let data = bincode::serialize(superblock)
.map_err(|e| StorageError::Index(format!("Failed to serialize SuperBlock: {}", e)))?;
if data.len() > Self::SUPERBLOCK_SIZE {
return Err(StorageError::Index(format!(
"SuperBlock too large: {} bytes exceeds {} reservation",
data.len(),
Self::SUPERBLOCK_SIZE
)));
}
let mut buf = vec![0u8; Self::SUPERBLOCK_SIZE];
buf[..data.len()].copy_from_slice(&data);
file.write_all(&buf)?;
Ok(())
}
fn sync_superblock(&self) -> Result<()> {
let root_page_id = *self.root_page_id.read();
let next_page_id = *self.next_page_id.read();
let stats = self.stats.read();
let page_offsets = self.page_offsets.read();
let superblock = SuperBlock {
magic: BTREE_MAGIC,
version: BTREE_VERSION,
root_page_id,
next_page_id,
total_keys: stats.total_keys,
total_pages: stats.total_pages,
leaf_pages: stats.leaf_pages,
internal_pages: stats.internal_pages,
tree_height: stats.tree_height,
page_offsets: page_offsets.clone(),
};
let mut file = self.storage_file.write();
Self::write_superblock(&mut file, &superblock)
}
fn load_page(&self, page_id: u64) -> Result<Arc<RwLock<Page>>> {
if page_id == 0 {
let root_id = *self.root_page_id.read();
return Err(StorageError::Corruption(format!(
"Cannot load Page 0: reserved for SuperBlock (root_id={})",
root_id
)));
}
{
let mut cache = self.page_cache.write();
if let Some(page) = cache.get(&page_id) {
let mut stats = self.stats.write();
stats.page_cache_hits += 1;
return Ok(Arc::clone(page));
}
}
let mut stats = self.stats.write();
stats.page_cache_misses += 1;
drop(stats);
let file_offset = {
let offsets = self.page_offsets.read();
let idx = page_id as usize;
if idx >= offsets.len() || offsets[idx] == 0 {
return Err(StorageError::Corruption(format!(
"Page {} not found in page table",
page_id
)));
}
offsets[idx]
};
let file = self.storage_file.read();
use std::os::unix::fs::FileExt;
let mut header_buf = [0u8; 15];
file.read_exact_at(&mut header_buf, file_offset)?;
let content_len = u16::from_le_bytes([header_buf[13], header_buf[14]]) as usize;
if !(15..=MAX_PAGE_SIZE).contains(&content_len) {
return Err(StorageError::Corruption(format!(
"Invalid content_len {} for page {} at offset {}",
content_len, page_id, file_offset
)));
}
let mut buf = vec![0u8; content_len];
file.read_exact_at(&mut buf, file_offset)?;
let page = Page::deserialize_compact(page_id, &buf)?;
page.validate()?;
let page_arc = Arc::new(RwLock::new(page));
let mut cache = self.page_cache.write();
cache.put(page_id, Arc::clone(&page_arc));
Ok(page_arc)
}
fn flush_page(&self, page: &Page) -> Result<()> {
if !page.dirty {
return Ok(());
}
if page.page_id == 0 {
return Err(StorageError::Corruption(
"Cannot flush Page 0: reserved for SuperBlock".into(),
));
}
let buf = page.serialize_compact()?;
let _flush_guard = self.flush_lock.lock();
let mut file = self.storage_file.write();
let file_end = file.metadata()?.len().max(Self::SUPERBLOCK_SIZE as u64);
file.seek(SeekFrom::Start(file_end))?;
file.write_all(&buf)?;
{
let mut offsets = self.page_offsets.write();
let idx = page.page_id as usize;
if idx >= offsets.len() {
offsets.resize(idx + 1, 0);
}
offsets[idx] = file_end;
}
if self.config.immediate_sync {
file.sync_all()?;
}
Ok(())
}
fn alloc_page(&self, is_leaf: bool) -> Result<Arc<RwLock<Page>>> {
let page_id = {
let mut next_id = self.next_page_id.write();
let id = *next_id;
*next_id += 1;
id
};
let page = if is_leaf {
Page::new_leaf(page_id)
} else {
Page::new_internal(page_id)
};
let page_arc = Arc::new(RwLock::new(page));
let mut cache = self.page_cache.write();
cache.put(page_id, Arc::clone(&page_arc));
Ok(page_arc)
}
fn search_internal(&self, page_id: u64, key: u64) -> Result<Option<u64>> {
let page_arc = self.load_page(page_id)?;
let page = page_arc.read();
if page.is_leaf {
match page.keys.binary_search(&key) {
Ok(idx) => Ok(Some(page.values[idx])),
Err(_) => Ok(None),
}
} else {
let child_idx = match page.keys.binary_search(&key) {
Ok(idx) => idx + 1, Err(idx) => idx, };
let child_page_id = page.children[child_idx];
drop(page);
self.search_internal(child_page_id, key)
}
}
pub fn insert(&mut self, key: u64, value: u64) -> Result<Option<u64>> {
let root_id = *self.root_page_id.read();
if root_id == 0 {
let root_page = self.alloc_page(true)?;
let new_root_id = {
let mut page = root_page.write();
page.keys.push(key);
page.values.push(value);
page.num_keys = 1;
page.dirty = true;
page.page_id
};
{
let page_ref = root_page.read();
self.flush_page(&page_ref)?;
}
{
let mut root = self.root_page_id.write();
*root = new_root_id;
}
self.sync_superblock()?;
let mut stats = self.stats.write();
stats.total_keys = 1;
stats.total_pages = 1;
stats.leaf_pages = 1;
return Ok(None);
}
let (old_value, split_info) = self.insert_internal(root_id, key, value)?;
if let Some((split_key, new_page_id)) = split_info {
let new_root = self.alloc_page(false)?;
{
let mut root = new_root.write();
root.keys.push(split_key);
root.children.push(root_id);
root.children.push(new_page_id);
root.num_keys = 1;
root.dirty = true;
}
{
let page_ref = new_root.read();
let new_root_id = page_ref.page_id;
self.flush_page(&page_ref)?;
drop(page_ref);
let mut root_page_id = self.root_page_id.write();
*root_page_id = new_root_id;
}
self.sync_superblock()?;
let mut stats = self.stats.write();
stats.total_pages += 1;
stats.internal_pages += 1;
stats.tree_height += 1;
}
if old_value.is_none() {
let mut stats = self.stats.write();
stats.total_keys += 1;
}
Ok(old_value)
}
fn insert_internal(&mut self, page_id: u64, key: u64, value: u64) -> InsertResult {
let page_arc = self.load_page(page_id)?;
let is_leaf = {
let page = page_arc.read();
page.is_leaf
};
if is_leaf {
let mut page = page_arc.write();
let search_result = page.keys.binary_search(&key);
let old_value = match search_result {
Ok(idx) => {
if !self.config.allow_updates {
return Err(StorageError::InvalidData(
"Key already exists and updates are disabled".into(),
));
}
let old = Some(page.values[idx]);
page.values[idx] = value;
page.dirty = true;
return Ok((old, None));
}
Err(idx) => {
page.keys.insert(idx, key);
page.values.insert(idx, value);
page.num_keys += 1;
page.dirty = true;
None
}
};
if page.num_keys >= self.config.order {
let split_info = self.split_leaf(&mut page)?;
drop(page);
Ok((old_value, Some(split_info)))
} else {
drop(page);
Ok((old_value, None))
}
} else {
let child_idx = {
let page = page_arc.read();
page.keys.binary_search(&key).unwrap_or_else(|idx| idx)
};
let child_page_id = {
let page = page_arc.read();
page.children[child_idx]
};
let (old_value, child_split) = self.insert_internal(child_page_id, key, value)?;
if let Some((split_key, new_child_id)) = child_split {
let mut page = page_arc.write();
let insert_idx = page
.keys
.binary_search(&split_key)
.unwrap_or_else(|idx| idx);
page.keys.insert(insert_idx, split_key);
page.children.insert(insert_idx + 1, new_child_id);
page.num_keys += 1;
page.dirty = true;
if page.num_keys >= self.config.order {
let split_info = self.split_internal(&mut page)?;
drop(page);
Ok((old_value, Some(split_info)))
} else {
drop(page);
Ok((old_value, None))
}
} else {
Ok((old_value, None))
}
}
}
fn split_leaf(&mut self, page: &mut Page) -> Result<(u64, u64)> {
let mid = page.num_keys * 7 / 10;
let new_page_arc = self.alloc_page(true)?;
let new_page_id = {
let mut new_page = new_page_arc.write();
new_page.keys = page.keys.split_off(mid);
new_page.values = page.values.split_off(mid);
new_page.num_keys = new_page.keys.len();
new_page.dirty = true;
new_page.next_leaf = page.next_leaf;
page.next_leaf = new_page.page_id;
let split_key = new_page.keys[0];
let new_id = new_page.page_id;
drop(new_page);
(split_key, new_id)
};
page.num_keys = page.keys.len();
page.dirty = true;
let mut stats = self.stats.write();
stats.total_pages += 1;
stats.leaf_pages += 1;
Ok(new_page_id)
}
fn split_internal(&mut self, page: &mut Page) -> Result<(u64, u64)> {
let original_num_keys = page.num_keys;
let original_num_children = page.children.len();
if page.num_keys < 1 {
return Err(StorageError::Index(format!(
"Cannot split internal node with {} keys",
page.num_keys
)));
}
let mid = page.num_keys * 7 / 10;
if mid >= page.num_keys {
return Err(StorageError::Index(format!(
"Invalid split mid={} for num_keys={}",
mid, page.num_keys
)));
}
let split_key = page.keys[mid];
let new_page_arc = self.alloc_page(false)?;
let new_page_id = {
let mut new_page = new_page_arc.write();
new_page.keys = page.keys.split_off(mid + 1);
new_page.children = page.children.split_off(mid + 1);
new_page.num_keys = new_page.keys.len();
new_page.dirty = true;
if new_page.num_keys == 0 {
return Err(StorageError::Corruption(
format!("Split internal node: right child has 0 keys (original_keys={}, mid={}, keys_after_splitoff={})",
original_num_keys, mid, page.keys.len())
));
}
if new_page.children.len() != new_page.num_keys + 1 {
return Err(StorageError::Corruption(format!(
"Split internal node: right child has {} keys but {} children",
new_page.num_keys,
new_page.children.len()
)));
}
let new_id = new_page.page_id;
drop(new_page);
(split_key, new_id)
};
if !page.keys.is_empty() {
page.keys.pop(); }
page.num_keys = page.keys.len();
page.dirty = true;
if page.num_keys == 0 {
return Err(StorageError::Corruption(
format!("Split internal node: left child has 0 keys after removing mid (original_keys={}, mid={})",
original_num_keys, mid)
));
}
if page.children.len() != page.num_keys + 1 {
return Err(StorageError::Corruption(
format!("Split internal node: left child has {} keys but {} children (original had {} keys, {} children)",
page.num_keys, page.children.len(), original_num_keys, original_num_children)
));
}
let mut stats = self.stats.write();
stats.total_pages += 1;
stats.internal_pages += 1;
Ok(new_page_id)
}
pub fn get(&self, key: &u64) -> Result<Option<u64>> {
let root_id = *self.root_page_id.read();
if root_id == 0 {
return Ok(None);
}
self.search_internal(root_id, *key)
}
pub fn remove(&mut self, key: &u64) -> Result<Option<u64>> {
let root_id = *self.root_page_id.read();
if root_id == 0 {
return Ok(None);
}
let leaf_id = self.find_leaf_for_key(root_id, *key)?;
let leaf_arc = self.load_page(leaf_id)?;
let mut leaf = leaf_arc.write();
if !leaf.is_leaf {
return Err(StorageError::Index(
"find_leaf_for_key returned non-leaf page".into(),
));
}
match leaf.keys.binary_search(key) {
Ok(idx) => {
let old_value = leaf.values[idx];
leaf.keys.remove(idx);
leaf.values.remove(idx);
leaf.num_keys -= 1;
leaf.dirty = true;
drop(leaf);
let mut stats = self.stats.write();
stats.total_keys = stats.total_keys.saturating_sub(1);
Ok(Some(old_value))
}
Err(_) => Ok(None),
}
}
pub fn contains_key(&self, key: &u64) -> Result<bool> {
Ok(self.get(key)?.is_some())
}
pub fn len(&self) -> usize {
self.stats.read().total_keys
}
pub fn is_empty(&self) -> bool {
self.len() == 0
}
pub fn stats(&self) -> BTreeStats {
self.stats.read().clone()
}
pub fn range_with_limit(
&self,
start: &u64,
end: &u64,
limit: usize,
) -> Result<Vec<(u64, u64)>> {
let root_id = *self.root_page_id.read();
if root_id == 0 || limit == 0 {
return Ok(Vec::new());
}
let first_leaf_id = self.find_leaf_for_key(root_id, *start)?;
let mut results = Vec::with_capacity(limit.min(64));
self.scan_leaf_chain_limited(first_leaf_id, *start, *end, &mut results, limit)?;
Ok(results)
}
pub fn range(&self, start: &u64, end: &u64) -> Result<Vec<(u64, u64)>> {
let root_id = *self.root_page_id.read();
if root_id == 0 {
return Ok(Vec::new());
}
let first_leaf_id = self.find_leaf_for_key(root_id, *start)?;
let mut results = Vec::new();
self.scan_leaf_chain(first_leaf_id, *start, *end, &mut results)?;
Ok(results)
}
fn find_leaf_for_key(&self, page_id: u64, key: u64) -> Result<u64> {
let page_arc = self.load_page(page_id)?;
let page = page_arc.read();
if page.is_leaf {
return Ok(page_id);
}
let mut child_idx = 0;
for i in 0..page.num_keys {
if key < page.keys[i] {
break;
}
child_idx = i + 1;
}
if child_idx >= page.children.len() {
return Err(StorageError::Index(format!(
"Child index {} out of bounds (num_children={}, num_keys={}, page_id={})",
child_idx,
page.children.len(),
page.num_keys,
page_id
)));
}
let child_id = page.children[child_idx];
if child_id == 0 {
return Err(StorageError::Corruption(format!(
"Invalid child_id=0 at page_id={}, child_idx={}, num_keys={}",
page_id, child_idx, page.num_keys
)));
}
drop(page);
self.find_leaf_for_key(child_id, key)
}
fn scan_leaf_chain(
&self,
start_leaf_id: u64,
start: u64,
end: u64,
results: &mut Vec<(u64, u64)>,
) -> Result<()> {
let mut current_leaf_id = start_leaf_id;
while current_leaf_id != INVALID_PAGE_ID {
let page_arc = self.load_page(current_leaf_id)?;
let page = page_arc.read();
if !page.is_leaf {
return Err(StorageError::Index("Expected leaf node".into()));
}
let mut found_end = false;
for i in 0..page.num_keys {
let key = page.keys[i];
if key > end {
found_end = true;
break;
}
if key >= start {
results.push((key, page.values[i]));
}
}
if found_end {
break;
}
current_leaf_id = page.next_leaf;
}
Ok(())
}
fn scan_leaf_chain_limited(
&self,
start_leaf_id: u64,
start: u64,
end: u64,
results: &mut Vec<(u64, u64)>,
limit: usize,
) -> Result<()> {
let mut current_leaf_id = start_leaf_id;
while current_leaf_id != INVALID_PAGE_ID && results.len() < limit {
let page_arc = self.load_page(current_leaf_id)?;
let page = page_arc.read();
if !page.is_leaf {
return Err(StorageError::Index("Expected leaf node".into()));
}
for i in 0..page.num_keys {
if results.len() >= limit {
return Ok(());
}
let key = page.keys[i];
if key > end {
return Ok(());
}
if key >= start {
results.push((key, page.values[i]));
}
}
current_leaf_id = page.next_leaf;
}
Ok(())
}
pub fn flush(&self) -> Result<()> {
let cache = self.page_cache.read();
let mut pages: Vec<(u64, Arc<RwLock<Page>>)> = cache
.iter()
.map(|(id, arc)| (*id, Arc::clone(arc)))
.collect();
pages.sort_by_key(|(id, _)| *id);
drop(cache);
let _flush_guard = self.flush_lock.lock();
let mut file = self.storage_file.write();
let mut offset = Self::SUPERBLOCK_SIZE as u64;
let mut new_offsets = vec![0u64];
for (page_id, page_arc) in &pages {
let page = page_arc.read();
let buf = page.serialize_compact()?;
file.seek(SeekFrom::Start(offset))?;
file.write_all(&buf)?;
let idx = *page_id as usize;
if idx >= new_offsets.len() {
new_offsets.resize(idx + 1, 0);
}
new_offsets[idx] = offset;
offset += buf.len() as u64;
}
let mut offsets = self.page_offsets.write();
*offsets = new_offsets;
drop(offsets);
for (_, page_arc) in &pages {
let mut page = page_arc.write();
page.dirty = false;
}
drop(file);
self.sync_superblock()?;
let file = self.storage_file.write();
file.set_len(offset)?;
file.sync_all()?;
Ok(())
}
pub fn scan(&self) -> Result<Vec<(u64, u64)>> {
let root_id = *self.root_page_id.read();
if root_id == 0 {
return Ok(Vec::new());
}
let mut results = Vec::new();
self.scan_internal(root_id, &mut results)?;
Ok(results)
}
fn scan_internal(&self, page_id: u64, results: &mut Vec<(u64, u64)>) -> Result<()> {
let page_arc = self.load_page(page_id)?;
let page = page_arc.read();
if page.is_leaf {
for i in 0..page.num_keys {
results.push((page.keys[i], page.values[i]));
}
if page.next_leaf != INVALID_PAGE_ID {
let next_id = page.next_leaf;
drop(page);
self.scan_internal(next_id, results)?;
}
} else {
if !page.children.is_empty() {
let first_child = page.children[0];
drop(page);
self.scan_internal(first_child, results)?;
}
}
Ok(())
}
pub fn min_key(&self) -> Result<Option<u64>> {
let root_id = *self.root_page_id.read();
if root_id == 0 {
return Ok(None);
}
let mut page_id = root_id;
loop {
let page_arc = self.load_page(page_id)?;
let page = page_arc.read();
if page.is_leaf {
return Ok(page.keys.first().copied());
}
page_id = match page.children.first() {
Some(&id) => id,
None => return Ok(page.keys.first().copied()),
};
}
}
pub fn max_key(&self) -> Result<Option<u64>> {
let root_id = *self.root_page_id.read();
if root_id == 0 {
return Ok(None);
}
let mut page_id = root_id;
loop {
let page_arc = self.load_page(page_id)?;
let page = page_arc.read();
if page.is_leaf {
return Ok(page.keys.last().copied());
}
page_id = match page.children.last() {
Some(&id) => id,
None => return Ok(page.keys.last().copied()),
};
}
}
}
impl Drop for BTree {
fn drop(&mut self) {
let _ = self.flush();
}
}
#[cfg(test)]
mod tests {
use super::*;
use tempfile::TempDir;
fn create_test_btree() -> (BTree, TempDir) {
let temp_dir = TempDir::new().unwrap();
let path = temp_dir.path().join("test.btree");
let btree = BTree::new(path).unwrap();
(btree, temp_dir)
}
#[test]
fn test_basic_operations() {
let (mut btree, _temp) = create_test_btree();
assert!(btree.insert(1, 100).unwrap().is_none());
assert!(btree.insert(2, 200).unwrap().is_none());
assert!(btree.insert(3, 300).unwrap().is_none());
assert_eq!(btree.get(&1).unwrap(), Some(100));
assert_eq!(btree.get(&2).unwrap(), Some(200));
assert_eq!(btree.get(&999).unwrap(), None);
assert_eq!(btree.len(), 3);
assert!(btree.contains_key(&1).unwrap());
assert!(!btree.contains_key(&999).unwrap());
}
#[test]
fn test_persistence() {
let temp_dir = TempDir::new().unwrap();
let path = temp_dir.path().join("persist.btree");
{
let mut btree = BTree::new(path.clone()).unwrap();
btree.insert(1, 100).unwrap();
btree.insert(2, 200).unwrap();
btree.flush().unwrap();
}
{
let btree = BTree::new(path).unwrap();
assert_eq!(btree.get(&1).unwrap(), Some(100));
assert_eq!(btree.get(&2).unwrap(), Some(200));
}
}
#[test]
fn test_superblock_persistence() {
let temp_dir = TempDir::new().unwrap();
let path = temp_dir.path().join("superblock.btree");
{
let mut btree = BTree::new(path.clone()).unwrap();
for i in 1..=100 {
btree.insert(i, i * 10).unwrap();
}
let stats_before = btree.stats();
assert_eq!(stats_before.total_keys, 100);
btree.flush().unwrap();
}
{
let btree = BTree::new(path).unwrap();
let stats_after = btree.stats();
assert_eq!(stats_after.total_keys, 100);
assert!(stats_after.total_pages > 0);
let root_id = *btree.root_page_id.read();
assert!(
root_id > 0,
"Root should be at Page 1 or higher (Page 0 is SuperBlock)"
);
assert_eq!(btree.get(&1).unwrap(), Some(10));
assert_eq!(btree.get(&50).unwrap(), Some(500));
assert_eq!(btree.get(&100).unwrap(), Some(1000));
assert_eq!(btree.len(), 100);
}
}
#[test]
fn test_range_query() {
let (mut btree, _temp) = create_test_btree();
for i in 1..=10 {
btree.insert(i, i * 100).unwrap();
}
let results = btree.range(&3, &7).unwrap();
assert_eq!(results.len(), 5);
assert_eq!(results[0], (3, 300));
assert_eq!(results[4], (7, 700));
}
#[test]
fn test_unique_constraint() {
let temp_dir = TempDir::new().unwrap();
let path = temp_dir.path().join("unique.btree");
let config = BTreeConfig {
unique_keys: true,
allow_updates: false, ..Default::default()
};
let mut btree = BTree::with_config(path, config).unwrap();
btree.insert(1, 100).unwrap();
let result = btree.insert(1, 200);
assert!(result.is_err());
}
#[test]
fn test_remove() {
let (mut btree, _temp) = create_test_btree();
btree.insert(1, 100).unwrap();
btree.insert(2, 200).unwrap();
assert_eq!(btree.remove(&1).unwrap(), Some(100));
assert_eq!(btree.len(), 1);
assert_eq!(btree.remove(&1).unwrap(), None);
}
#[test]
fn test_min_max_key() {
let (mut btree, _temp) = create_test_btree();
btree.insert(5, 50).unwrap();
btree.insert(1, 10).unwrap();
btree.insert(10, 100).unwrap();
assert_eq!(btree.min_key().unwrap(), Some(1));
assert_eq!(btree.max_key().unwrap(), Some(10));
}
#[test]
fn test_scan() {
let (mut btree, _temp) = create_test_btree();
for i in 1..=5 {
btree.insert(i, i * 10).unwrap();
}
let all = btree.scan().unwrap();
assert_eq!(all.len(), 5);
assert_eq!(all[0], (1, 10));
assert_eq!(all[4], (5, 50));
}
#[test]
fn test_update() {
let (mut btree, _temp) = create_test_btree();
btree.insert(1, 100).unwrap();
assert_eq!(btree.get(&1).unwrap(), Some(100));
btree.insert(1, 200).unwrap();
assert_eq!(btree.get(&1).unwrap(), Some(200));
assert_eq!(btree.len(), 1);
}
#[test]
fn test_simple_split() {
let (mut btree, _temp) = create_test_btree();
for i in 0..256 {
btree.insert(i, i * 10).unwrap();
}
for i in 0..256 {
let result = btree.get(&i).unwrap();
assert_eq!(result, Some(i * 10), "Key {} missing or wrong", i);
}
debug_log!("Stats: {:?}", btree.stats());
}
#[test]
fn test_node_split() {
let (mut btree, _temp) = create_test_btree();
for i in 0..1000 {
btree.insert(i, i * 10).unwrap();
}
for i in 0..1000 {
let result = btree.get(&i).unwrap();
if result != Some(i * 10) {
panic!(
"Key {} not found or has wrong value. Expected: {}, Got: {:?}",
i,
i * 10,
result
);
}
}
let stats = btree.stats();
assert_eq!(stats.total_keys, 1000);
assert!(stats.total_pages > 1); assert!(stats.tree_height > 0);
let all = btree.scan().unwrap();
assert_eq!(all.len(), 1000);
for (i, &(k, v)) in all.iter().enumerate() {
assert_eq!((k, v), (i as u64, (i * 10) as u64));
}
}
#[test]
fn test_large_dataset() {
let (mut btree, _temp) = create_test_btree();
let count = 5000;
for i in 0..count {
btree.insert(i, i).unwrap();
}
assert_eq!(btree.len(), count as usize);
assert_eq!(btree.get(&2500).unwrap(), Some(2500));
assert_eq!(btree.get(&4999).unwrap(), Some(4999));
assert_eq!(btree.get(&0).unwrap(), Some(0));
let results = btree.range(&1000, &1010).unwrap();
assert_eq!(results.len(), 11);
}
#[test]
fn test_range_with_limit() {
let (mut btree, _temp) = create_test_btree();
for i in 0..500 {
btree.insert(i, i * 10).unwrap();
}
let results = btree.range_with_limit(&0, &499, 5).unwrap();
assert_eq!(results.len(), 5);
assert_eq!(results[0], (0, 0));
assert_eq!(results[4], (4, 40));
let results = btree.range_with_limit(&100, &110, 50).unwrap();
assert_eq!(results.len(), 11);
let results = btree.range_with_limit(&0, &499, 0).unwrap();
assert!(results.is_empty());
let results = btree.range_with_limit(&600, &700, 10).unwrap();
assert!(results.is_empty());
}
}