use std::{ffi::OsString, fs, path::Path, sync::Arc};
use wbase::base32::Base32Buf128;
use super::{RangeIndexManager, TreeEntry};
use crate::{
error::{Error, Result},
service::{BfTreeService, file_has_cpr_magic},
stub::RangeIndexStub,
types::{BfTreeConfig, StorageBackend, StorageBackendType, TreeTuning},
};
impl RangeIndexManager {
pub fn create_bftree(
&self,
key: &[u8],
storage_backend: StorageBackend,
tuning: TreeTuning,
) -> Result<Arc<BfTreeService>> {
let key_hash = Self::key_hash_of(key);
let _stripe_lock = self.locks.write(key_hash);
let key_id = Self::key_id_of(key);
let hash_prefix = Self::base32_prefix_of(key);
self.create_bftree_internal(key_id, key_hash, hash_prefix, storage_backend, tuning)
}
fn instantiate_tree(
&self,
hash_prefix: &str,
storage_backend: StorageBackend,
tuning: TreeTuning,
) -> Result<Arc<BfTreeService>> {
let mut config = BfTreeConfig::default();
let (file_path_str, backend_type) = if storage_backend == StorageBackend::Memory {
config.cache_only(true);
(None, StorageBackendType::Memory)
} else {
let data_path = self.data_file_path(hash_prefix);
config.file_path(&data_path);
(
Some(data_path.to_string_lossy().into_owned()),
StorageBackendType::Disk,
)
};
if tuning.cache_size > 0 {
config.cb_size_byte(tuning.cache_size);
}
if tuning.min_record_size > 0 {
config.cb_min_record_size(tuning.min_record_size);
}
if tuning.max_record_size > 0 {
config.cb_max_record_size(tuning.max_record_size);
}
if tuning.max_key_len > 0 {
config.cb_max_key_len(tuning.max_key_len);
}
let actual_leaf_page_size = if tuning.leaf_page_size > 0 {
tuning.leaf_page_size
} else if tuning.max_record_size > 0 {
Self::compute_leaf_page_size(tuning.max_record_size)
} else {
0
};
if actual_leaf_page_size > 0 {
config.leaf_page_size(actual_leaf_page_size);
}
config.use_snapshot(true);
Ok(Arc::new(BfTreeService::new_with_backend(
config,
backend_type,
file_path_str,
)?))
}
fn create_bftree_internal(
&self,
key_id: u128,
key_hash: u64,
hash_prefix: Base32Buf128,
storage_backend: StorageBackend,
tuning: TreeTuning,
) -> Result<Arc<BfTreeService>> {
let pin = self.live_indexes.pin();
if pin.contains_key(&key_id) {
return Err(Error::IndexExists);
}
let _ = fs::remove_file(self.data_file_path(&hash_prefix));
let tree = self.instantiate_tree(&hash_prefix, storage_backend, tuning)?;
let entry = Arc::new(TreeEntry::new(
Some(Arc::clone(&tree)),
key_hash,
key_id,
hash_prefix,
));
pin.insert(key_id, entry);
Ok(tree)
}
pub fn get_or_open_tree(&self, key: &[u8], stub: &RangeIndexStub) -> Result<Arc<BfTreeService>> {
let key_hash = Self::key_hash_of(key);
let _stripe_lock = self.locks.write(key_hash);
let key_id = Self::key_id_of(key);
{
let pin = self.live_indexes.pin();
if let Some(entry) = pin.get(&key_id) {
let tree_guard = entry.tree.read();
if let Some(t) = tree_guard.as_ref() {
return Ok(Arc::clone(t));
}
}
}
let hash_prefix = Self::base32_prefix_of(key);
let backend = StorageBackendType::from_u8(stub.storage_backend);
let data_path = self.data_file_path(&hash_prefix);
let flush_path = self.bare_flush_path(&hash_prefix);
if backend == StorageBackendType::Disk {
if flush_path.exists() {
fs::copy(&flush_path, &data_path)?;
} else if self.addr_flush_scan_pending() {
let scan_token = self.addr_flush_scan_token();
let mut found_flush = false;
if let Ok(entries) = fs::read_dir(&self.ri_log_root) {
let mut latest: Option<(i64, OsString)> = None;
for entry in entries.flatten() {
let name = entry.file_name();
if let Some(name_str) = name.to_str()
&& let Some((prefix, addr)) = Self::parse_flush_file_name(name_str)
&& prefix == hash_prefix
&& latest.as_ref().is_none_or(|(max_addr, _)| addr > *max_addr)
{
latest = Some((addr, name));
}
}
if let Some((_, name)) = latest {
fs::copy(self.ri_log_root.join(name), &data_path)?;
found_flush = true;
}
}
if !found_flush {
self.settle_addr_flush_scan(scan_token);
}
}
if !data_path.exists() {
let mut msg = String::from("数据文件缺失且无可用刷盘快照: ");
msg.push_str(&data_path.display().to_string());
return Err(Error::Recovery(msg));
}
}
let is_cpr =
(backend == StorageBackendType::Disk || data_path.exists()) && file_has_cpr_magic(&data_path);
let tree = if is_cpr {
Arc::new(BfTreeService::recover_from_cpr_snapshot(
&data_path, true, backend,
)?)
} else {
self.instantiate_tree(&hash_prefix, backend.into(), TreeTuning::from(stub))?
};
let pin = self.live_indexes.pin();
if let Some(entry) = pin.get(&key_id) {
*entry.tree.write() = Some(Arc::clone(&tree));
} else {
let entry = Arc::new(TreeEntry::new(
Some(Arc::clone(&tree)),
key_hash,
key_id,
hash_prefix,
));
pin.insert(key_id, entry);
}
Ok(tree)
}
pub fn pre_stage_and_register_pending(&self, key: &[u8], src_flush_address: i64) -> Result<()> {
let key_hash = Self::key_hash_of(key);
let _stripe_lock = self.locks.write(key_hash);
let key_id = Self::key_id_of(key);
let hash_prefix = Self::base32_prefix_of(key);
let snapshot_path = self.log_flush_path(&hash_prefix, src_flush_address);
self.notice_addr_flush_files();
if !snapshot_path.exists() {
return Ok(());
}
let data_path = self.data_file_path(&hash_prefix);
fs::copy(&snapshot_path, &data_path)?;
let entry = Arc::new(TreeEntry::new(None, key_hash, key_id, hash_prefix));
self.live_indexes.pin().insert(key_id, entry);
Ok(())
}
#[inline]
fn remove_and_take_tree(&self, key_id: u128) -> Option<Arc<BfTreeService>> {
let entry = self.live_indexes.pin().remove(&key_id)?.clone();
entry.tree.write().take()
}
pub fn unregister_index(&self, key: &[u8]) -> Result<bool> {
let key_hash = Self::key_hash_of(key);
let _stripe_lock = self.locks.write(key_hash);
let Some(tree) = self.remove_and_take_tree(Self::key_id_of(key)) else {
return Ok(false);
};
drop(_stripe_lock);
tree.dispose_quiesced()?;
Ok(true)
}
pub fn delete_index(&self, key: &[u8]) -> Result<bool> {
let key_hash = Self::key_hash_of(key);
let _stripe_lock = self.locks.write(key_hash);
let Some(tree) = self.remove_and_take_tree(Self::key_id_of(key)) else {
let data_path = self.data_file_path_for_key(key);
if data_path.exists() {
let _ = fs::remove_file(data_path);
}
return Ok(false);
};
drop(_stripe_lock);
tree.dispose_quiesced()?;
let data_path = self.data_file_path_for_key(key);
if data_path.exists() {
let _ = fs::remove_file(data_path);
}
Ok(true)
}
pub fn dispose_tree(&self, key: &[u8], delete_file: bool) -> Result<bool> {
if delete_file {
self.delete_index(key)
} else {
self.unregister_index(key)
}
}
pub fn register_tree(&self, key: &[u8], tree: Arc<BfTreeService>) {
let key_hash = Self::key_hash_of(key);
let _stripe_lock = self.locks.write(key_hash);
let key_id = Self::key_id_of(key);
let pin = self.live_indexes.pin();
if let Some(existing) = pin.get(&key_id) {
*existing.tree.write() = Some(tree);
} else {
let hash_prefix = Self::base32_prefix_of(key);
let entry = Arc::new(TreeEntry::new(Some(tree), key_hash, key_id, hash_prefix));
pin.insert(key_id, entry);
}
}
pub fn publish_tree_from_snapshot_locked(
&self,
key: &[u8],
snapshot_path: &Path,
replace: bool,
) -> Result<Arc<BfTreeService>> {
let key_hash = Self::key_hash_of(key);
let key_id = Self::key_id_of(key);
let hash_prefix = Self::base32_prefix_of(key);
let data_path = self.data_file_path(&hash_prefix);
let exists = self.live_indexes.pin().contains_key(&key_id);
if exists && !replace {
return Err(Error::IndexExists);
}
if exists {
if let Some(old) = self.remove_and_take_tree(key_id) {
old.dispose_quiesced()?;
}
}
if !snapshot_path.exists() {
let mut msg = String::from("迁移快照文件不存在: ");
msg.push_str(&snapshot_path.display().to_string());
return Err(Error::Recovery(msg));
}
if let Some(parent) = data_path.parent() {
fs::create_dir_all(parent)?;
}
if fs::rename(snapshot_path, &data_path).is_err() {
if data_path.exists() {
fs::remove_file(&data_path)?;
}
fs::copy(snapshot_path, &data_path)?;
let _ = fs::remove_file(snapshot_path);
}
let tree = Arc::new(BfTreeService::recover_from_cpr_snapshot(
&data_path,
true,
StorageBackendType::Disk,
)?);
let entry = Arc::new(TreeEntry::new(
Some(Arc::clone(&tree)),
key_hash,
key_id,
hash_prefix,
));
self.live_indexes.pin().insert(key_id, entry);
Ok(tree)
}
}