use std::{
fs,
path::Path,
result::Result as StdResult,
sync::{Arc, atomic::Ordering},
};
use compio::runtime::spawn_blocking;
use thiserror::Error as ThisError;
use wbftree::{
BfTreeDeleteResult, BfTreeInsertResult, BfTreeReadResult, BfTreeService, RANGE_INDEX_STUB_SIZE,
RangeIndexManager, RangeIndexStub, ScanRecord, ScanReturnField, StorageBackend,
StorageBackendType, TreeTuning,
};
use wdev::Device;
use wval::{CollectionType, META_VALUE_SIZE, MetaValue};
use crate::{
error::{Error, Result},
session::StoreSession,
};
#[derive(ThisError, Debug, Clone, PartialEq, Eq)]
pub enum RangeIndexError {
#[error("ERR index already exists")]
AlreadyExists,
#[error("ERR range index not found")]
NotFound,
#[error("WRONGTYPE Operation against a key holding the wrong kind of value")]
WrongType,
#[error(
"ERR key+value size must be between {min_record_size} and {max_record_size} bytes (got {total_len}), max key length {max_key_len} (got {key_len})"
)]
InvalidKV {
min_record_size: u32,
max_record_size: u32,
max_key_len: u32,
total_len: usize,
key_len: usize,
},
#[error("ERR RI.SCAN is not supported for MEMORY-mode indexes")]
MemoryModeNotSupported,
#[error("ERR {0}")]
Internal(String),
}
impl From<Error> for RangeIndexError {
fn from(err: Error) -> Self {
Self::Internal(err.to_string())
}
}
impl From<wbftree::Error> for RangeIndexError {
fn from(err: wbftree::Error) -> Self {
match err {
wbftree::Error::IndexExists => Self::AlreadyExists,
other => Self::Internal(other.to_string()),
}
}
}
const DEFAULT_CACHE_SIZE: usize = 16 * 1024 * 1024;
const DEFAULT_MIN_RECORD_SIZE: usize = 64;
const DEFAULT_MAX_RECORD_SIZE: usize = 1024;
const DEFAULT_MAX_KEY_LEN: usize = 128;
#[inline]
const fn nz_or(v: usize, d: usize) -> usize {
if v > 0 { v } else { d }
}
pub(crate) async fn range_index_blocking<T: Send + 'static>(
op: impl FnOnce() -> T + Send + 'static,
) -> StdResult<T, RangeIndexError> {
spawn_blocking(op)
.await
.map_err(|e| RangeIndexError::Internal(format!("RangeIndex 阻塞任务异常退出: {e}")))
}
impl<D: Device> StoreSession<D> {
pub async fn range_index_create(
&self,
key: &[u8],
storage_backend: StorageBackend,
tuning: TreeTuning,
) -> StdResult<(), RangeIndexError> {
if self.read(key).await?.is_some() {
return Err(RangeIndexError::AlreadyExists);
}
if let Some(meta) = self.load_meta(key).await?
&& meta.size > 0
{
return Err(RangeIndexError::AlreadyExists);
}
let cache_size = nz_or(tuning.cache_size, DEFAULT_CACHE_SIZE);
let min_record_size = nz_or(tuning.min_record_size, DEFAULT_MIN_RECORD_SIZE);
let max_record_size = nz_or(tuning.max_record_size, DEFAULT_MAX_RECORD_SIZE);
let max_key_len = nz_or(tuning.max_key_len, DEFAULT_MAX_KEY_LEN);
let actual_leaf_page_size = if tuning.leaf_page_size > 0 {
tuning.leaf_page_size
} else {
RangeIndexManager::compute_leaf_page_size(max_record_size)
};
let mgr = Arc::clone(&self.store.range_index);
let create_key = key.to_vec();
let create_tuning = TreeTuning {
cache_size,
min_record_size,
max_record_size,
max_key_len,
leaf_page_size: actual_leaf_page_size,
};
let create_backend = storage_backend.clone();
let tree =
range_index_blocking(move || mgr.create_bftree(&create_key, create_backend, create_tuning))
.await?
.map_err(RangeIndexError::from)?;
let stub = RangeIndexStub::new(
tree.native_ptr(),
cache_size as u64,
min_record_size as u32,
max_record_size as u32,
max_key_len as u32,
actual_leaf_page_size as u32,
storage_backend,
);
let meta_k = self.session_meta_key(key);
let key_id = self.store.next_key_id.fetch_add(1, Ordering::Relaxed);
let meta = MetaValue::new(key_id, CollectionType::RangeIndex, 1, 1);
let val = encode_meta_stub_record(&meta, &stub);
if let Err(e) = self.upsert_raw(&meta_k, &val).await {
let _ = self.store.range_index.delete_index(key);
return Err(RangeIndexError::Internal(e.to_string()));
}
Ok(())
}
pub async fn load_range_index_stub(
&self,
key: &[u8],
) -> StdResult<Option<(MetaValue, RangeIndexStub)>, RangeIndexError> {
let meta_k = self.session_meta_key(key);
let Some(bytes) = self.read_raw(&meta_k).await? else {
if self.read(key).await?.is_some() {
return Err(RangeIndexError::WrongType);
}
return Ok(None);
};
if bytes.len() < META_VALUE_SIZE {
return Ok(None);
}
let meta = MetaValue::from_slice(&bytes[..META_VALUE_SIZE])
.map_err(|e| RangeIndexError::Internal(e.to_string()))?;
if meta.size == 0 {
self
.store
.update_key_id_meta(meta.key_id, meta.version, false);
return Ok(None);
}
if self.has_ttl_tag(key)? && self.check_expired(key).await? {
return Ok(None);
}
self
.store
.update_key_id_meta(meta.key_id, meta.version, true);
if meta.collection_type != CollectionType::RangeIndex {
return Err(RangeIndexError::WrongType);
}
if bytes.len() < META_VALUE_SIZE + RANGE_INDEX_STUB_SIZE {
return Ok(None);
}
let stub =
RangeIndexStub::decode(&bytes[META_VALUE_SIZE..META_VALUE_SIZE + RANGE_INDEX_STUB_SIZE])
.map_err(|e| RangeIndexError::Internal(e.to_string()))?;
Ok(Some((meta, stub)))
}
pub async fn acquire_tree_read(
&self,
key: &[u8],
stub: &RangeIndexStub,
) -> StdResult<(Arc<BfTreeService>, parking_lot::RwLockReadGuard<'_, ()>), RangeIndexError> {
let key_hash = RangeIndexManager::key_hash_of(key);
loop {
match self.store.range_index.wait_for_tree_checkpoint(key) {
Ok(true) => continue,
Ok(false) => {}
Err(e) => return Err(RangeIndexError::Internal(e.to_string())),
}
let read_lock = self.store.range_index.locks().read(key_hash);
if let Some(tree) = self.store.range_index.get_tree(key) {
return Ok((tree, read_lock));
}
drop(read_lock);
let mgr = Arc::clone(&self.store.range_index);
let restore_key = key.to_vec();
let restore_stub = *stub;
range_index_blocking(move || mgr.get_or_open_tree(&restore_key, &restore_stub))
.await?
.map_err(|e| RangeIndexError::Internal(e.to_string()))?;
}
}
pub async fn range_index_set(
&self,
key: &[u8],
field: &[u8],
value: &[u8],
) -> StdResult<(), RangeIndexError> {
let (_, stub) = self
.load_range_index_stub(key)
.await?
.ok_or(RangeIndexError::NotFound)?;
let total_len = field.len() + value.len();
if field.len() > stub.max_key_len as usize
|| total_len < stub.min_record_size as usize
|| total_len > stub.max_record_size as usize
{
return Err(RangeIndexError::InvalidKV {
min_record_size: stub.min_record_size,
max_record_size: stub.max_record_size,
max_key_len: stub.max_key_len,
total_len,
key_len: field.len(),
});
}
let (tree, _read_lock) = self.acquire_tree_read(key, &stub).await?;
match tree.insert(field, value) {
BfTreeInsertResult::Success => {
if let Some(listener) = self.store.range_listener() {
listener(key, field, value, false);
}
Ok(())
}
BfTreeInsertResult::InvalidArguments => {
Err(RangeIndexError::Internal("invalid arguments".to_string()))
}
BfTreeInsertResult::InvalidKV => Err(RangeIndexError::InvalidKV {
min_record_size: stub.min_record_size,
max_record_size: stub.max_record_size,
max_key_len: stub.max_key_len,
total_len,
key_len: field.len(),
}),
}
}
pub async fn range_index_get(
&self,
key: &[u8],
field: &[u8],
) -> StdResult<Option<Vec<u8>>, RangeIndexError> {
let (_, stub) = self
.load_range_index_stub(key)
.await?
.ok_or(RangeIndexError::NotFound)?;
let (tree, _read_lock) = self.acquire_tree_read(key, &stub).await?;
let (res, val) = tree.read(field);
match res {
BfTreeReadResult::Found => Ok(val),
BfTreeReadResult::NotFound | BfTreeReadResult::Deleted => Ok(None),
BfTreeReadResult::InvalidArguments | BfTreeReadResult::InvalidKey => {
Err(RangeIndexError::Internal("invalid arguments".to_string()))
}
}
}
pub async fn range_index_del(
&self,
key: &[u8],
field: &[u8],
) -> StdResult<bool, RangeIndexError> {
let (_, stub) = self
.load_range_index_stub(key)
.await?
.ok_or(RangeIndexError::NotFound)?;
let (tree, _read_lock) = self.acquire_tree_read(key, &stub).await?;
match tree.delete(field) {
BfTreeDeleteResult::Success => {
if let Some(listener) = self.store.range_listener() {
listener(key, field, &[], true);
}
Ok(true)
}
BfTreeDeleteResult::InvalidArguments => {
Err(RangeIndexError::Internal("invalid arguments".to_string()))
}
}
}
pub async fn range_index_scan_stream<F>(
&self,
key: &[u8],
start: &[u8],
count: usize,
return_field: ScanReturnField,
on_record: F,
) -> StdResult<usize, RangeIndexError>
where
F: FnMut(&[u8], &[u8]) -> bool,
{
let (_, stub) = self
.load_range_index_stub(key)
.await?
.ok_or(RangeIndexError::NotFound)?;
if stub.storage_backend == StorageBackendType::Memory.to_u8() {
return Err(RangeIndexError::MemoryModeNotSupported);
}
let (tree, _read_lock) = self.acquire_tree_read(key, &stub).await?;
Ok(tree.scan_with_count_callback(start, count, return_field, on_record)?)
}
pub async fn range_index_scan(
&self,
key: &[u8],
start: &[u8],
count: usize,
return_field: ScanReturnField,
) -> StdResult<Vec<ScanRecord>, RangeIndexError> {
let mut records = Vec::with_capacity(count.min(1024));
self
.range_index_scan_stream(
key,
start,
count,
return_field,
ScanRecord::sink(&mut records),
)
.await?;
Ok(records)
}
pub async fn range_index_range_stream<F>(
&self,
key: &[u8],
start: &[u8],
end: &[u8],
return_field: ScanReturnField,
on_record: F,
) -> StdResult<usize, RangeIndexError>
where
F: FnMut(&[u8], &[u8]) -> bool,
{
let (_, stub) = self
.load_range_index_stub(key)
.await?
.ok_or(RangeIndexError::NotFound)?;
if stub.storage_backend == StorageBackendType::Memory.to_u8() {
return Err(RangeIndexError::MemoryModeNotSupported);
}
let (tree, _read_lock) = self.acquire_tree_read(key, &stub).await?;
Ok(tree.scan_with_end_key_callback(start, end, return_field, on_record)?)
}
pub async fn range_index_range(
&self,
key: &[u8],
start: &[u8],
end: &[u8],
return_field: ScanReturnField,
) -> StdResult<Vec<ScanRecord>, RangeIndexError> {
let mut records = Vec::with_capacity(32);
self
.range_index_range_stream(
key,
start,
end,
return_field,
ScanRecord::sink(&mut records),
)
.await?;
Ok(records)
}
pub async fn range_index_exists(&self, key: &[u8]) -> Result<bool> {
if let Some(meta) = self.load_meta(key).await?
&& meta.size > 0
&& meta.collection_type == CollectionType::RangeIndex
{
return Ok(true);
}
Ok(false)
}
pub async fn range_index_config(&self, key: &[u8]) -> StdResult<RangeIndexStub, RangeIndexError> {
let (_, stub) = self
.load_range_index_stub(key)
.await?
.ok_or(RangeIndexError::NotFound)?;
Ok(stub)
}
pub async fn range_index_metrics(
&self,
key: &[u8],
) -> StdResult<(u64, bool, bool, bool), RangeIndexError> {
let (_, stub) = self
.load_range_index_stub(key)
.await?
.ok_or(RangeIndexError::NotFound)?;
let (tree_handle, is_live) = match self.store.range_index.get_tree(key) {
Some(tree) => (tree.native_ptr(), true),
None => (0, false),
};
Ok((tree_handle, is_live, stub.is_flushed(), stub.is_recovered()))
}
#[allow(clippy::await_holding_lock)]
pub async fn publish_migrated_range_index(
&self,
key: &[u8],
stub_bytes: &[u8],
temp_path: &Path,
replace: bool,
) -> StdResult<(), RangeIndexError> {
let key_hash = RangeIndexManager::key_hash_of(key);
let _xlock = self.store.range_index.locks().write(key_hash);
if self.range_index_exists(key).await? && !replace {
return Err(RangeIndexError::AlreadyExists);
}
let mgr = Arc::clone(&self.store.range_index);
let pub_key = key.to_vec();
let pub_src = temp_path.to_path_buf();
let tree = range_index_blocking(move || {
mgr.publish_tree_from_snapshot_locked(&pub_key, &pub_src, replace)
})
.await?
.map_err(|e| RangeIndexError::Internal(e.to_string()))?;
let mut stub =
RangeIndexStub::decode(stub_bytes).map_err(|e| RangeIndexError::Internal(e.to_string()))?;
stub.tree_handle = tree.native_ptr();
stub.reset_flags();
let meta_k = self.session_meta_key(key);
let key_id = self.store.next_key_id.fetch_add(1, Ordering::Relaxed);
let meta = MetaValue::new(key_id, CollectionType::RangeIndex, 1, 1);
let val = encode_meta_stub_record(&meta, &stub);
self
.upsert_raw(&meta_k, &val)
.await
.map_err(|e| RangeIndexError::Internal(e.to_string()))?;
Ok(())
}
pub async fn rename_range_index(&self, old_key: &[u8], new_key: &[u8]) -> Result<()> {
let old_meta_k = self.session_meta_key(old_key);
let Some(bytes) = self.read_raw(&old_meta_k).await? else {
return Ok(());
};
if bytes.len() < META_VALUE_SIZE + RANGE_INDEX_STUB_SIZE {
return Ok(());
}
let mut stub =
RangeIndexStub::decode(&bytes[META_VALUE_SIZE..META_VALUE_SIZE + RANGE_INDEX_STUB_SIZE])?;
while self.store.range_index.wait_for_tree_checkpoint(old_key)? {}
let mgr = Arc::clone(&self.store.range_index);
let restore_key = old_key.to_vec();
let snap_key = old_key.to_vec();
let new_key_owned = new_key.to_vec();
let restore_stub = stub;
let new_tree =
range_index_blocking(move || -> StdResult<Arc<BfTreeService>, RangeIndexError> {
let old_tree = match mgr.get_tree(&restore_key) {
Some(t) => t,
None => mgr.get_or_open_tree(&restore_key, &restore_stub)?,
};
let new_path = mgr.data_file_path_for_key(&new_key_owned);
if let Some(parent) = new_path.parent() {
let _ = fs::create_dir_all(parent);
}
let _ = fs::remove_file(&new_path);
let old_hash = RangeIndexManager::key_hash_of(&snap_key);
let _xlock = mgr.locks().write(old_hash);
mgr.snapshot_tree_to_path_locked(&snap_key, &old_tree, &new_path)?;
let backend = StorageBackendType::from_u8(restore_stub.storage_backend);
BfTreeService::recover_from_cpr_snapshot(&new_path, true, backend)
.map(Arc::new)
.map_err(RangeIndexError::from)
})
.await??;
stub.tree_handle = new_tree.native_ptr();
stub.reset_flags();
self.store.range_index.register_tree(new_key, new_tree);
let new_meta_k = self.session_meta_key(new_key);
let key_id = self.store.next_key_id.fetch_add(1, Ordering::Relaxed);
let meta = MetaValue::new(key_id, CollectionType::RangeIndex, 1, 1);
let val = encode_meta_stub_record(&meta, &stub);
self.upsert_raw(&new_meta_k, &val).await?;
Ok(())
}
}
#[inline]
pub(crate) fn encode_meta_stub_record(
meta: &MetaValue,
stub: &RangeIndexStub,
) -> [u8; META_VALUE_SIZE + RANGE_INDEX_STUB_SIZE] {
let mut val = [0u8; META_VALUE_SIZE + RANGE_INDEX_STUB_SIZE];
val[..META_VALUE_SIZE].copy_from_slice(&meta.to_bytes());
val[META_VALUE_SIZE..].copy_from_slice(&stub.encode());
val
}