use std::sync::Arc;
use wbase::time::now_ms;
use wdev::Device;
use wval::{
CollectionType, CompactHashCodec, KeyTag, META_VALUE_SIZE, MetaValue, StorageEncoding,
ZSetSubKeyCodec,
};
use crate::{error::Result, range_index::range_index_blocking, session::StoreSession, ttl::TtlOpt};
pub const META_HAS_EXPIRE_MASK: u8 = 0x80;
const META_RESERVED1_OFFSET: usize = 10;
impl<D: Device> StoreSession<D> {
pub const CHUNK_CAPACITY: u16 = 128;
pub fn clear_bftree_zset(&self, key_id: u64, version: u64) -> Result<()> {
const BATCH_SIZE: usize = 256;
let prefixes = [
ZSetSubKeyCodec::encode_score_prefix(key_id, version),
ZSetSubKeyCodec::encode_member_header(key_id, version),
];
let mut batch = Vec::with_capacity(BATCH_SIZE);
for prefix in prefixes {
let mut end_key = prefix;
let mut has_next = false;
for i in (0..end_key.len()).rev() {
if end_key[i] < 0xff {
end_key[i] += 1;
has_next = true;
break;
}
end_key[i] = 0;
}
let end_slice: &[u8] = if has_next { &end_key } else { &[0xff; 18] };
let mut cursor = prefix.to_vec();
loop {
batch.clear();
let _ = self.store.bftree.scan_with_end_key_callback(
&cursor,
end_slice,
wbftree::ScanReturnField::Key,
|k, _| {
if k.starts_with(&prefix) {
batch.push(k.to_vec());
batch.len() < BATCH_SIZE
} else {
false
}
},
);
if batch.is_empty() {
break;
}
for k in &batch {
let _ = self.store.bftree.delete(k);
}
let Some(last) = batch.last() else {
break;
};
if last.as_slice() <= cursor.as_slice() {
break;
}
cursor.clear();
cursor.extend_from_slice(last);
let mut carried = false;
for b in cursor.iter_mut().rev() {
if *b < 0xff {
*b += 1;
carried = true;
break;
}
*b = 0;
}
if !carried || cursor.as_slice() > end_slice {
break;
}
}
}
Ok(())
}
pub async fn delete(&self, key: &[u8]) -> Result<bool> {
self.del_ttl(key).await?;
if let Some(false) = self.check_object_meta_fast(key)? {
let str_k = self.session_string_key(key);
return self.delete_raw(&str_k).await;
}
let mut meta_del = false;
if let Some(mut meta) = self.load_meta(key).await? {
let meta_k = self.session_meta_key(key);
if meta.size > 0 {
meta_del = true;
}
if matches!(
meta.collection_type,
CollectionType::ZSet | CollectionType::Hash | CollectionType::Set
) {
if meta.encoding() == StorageEncoding::Flattened {
if meta.collection_type == CollectionType::ZSet {
self.clear_bftree_zset(meta.key_id, meta.version)?;
}
meta.bump_version();
meta.size = 0;
Self::set_meta_chunk_info(&mut meta.reserved, 0, 0);
let bytes = meta.to_bytes();
self.upsert_raw(&meta_k, &bytes).await?;
self
.store
.update_key_id_meta(meta.key_id, meta.version, false);
} else {
self
.store
.update_key_id_meta(meta.key_id, meta.version, false);
self.delete_raw(&meta_k).await?;
}
} else if meta.collection_type == CollectionType::RangeIndex {
self.delete_raw(&meta_k).await?;
self
.store
.update_key_id_meta(meta.key_id, meta.version, false);
let mgr = Arc::clone(&self.store.range_index);
let del_key = key.to_vec();
let _deleted = range_index_blocking(move || mgr.delete_index(&del_key)).await?;
meta_del = true;
} else {
self
.store
.update_key_id_meta(meta.key_id, meta.version, false);
self.delete_raw(&meta_k).await?;
meta_del = true;
}
}
let str_k = self.session_string_key(key);
let key_del = self.delete_raw(&str_k).await?;
Ok(meta_del || key_del)
}
pub async fn contains_key(&self, key: &[u8]) -> Result<bool> {
if let Some(meta) = self.load_meta(key).await?
&& meta.size > 0
{
return Ok(true);
}
let str_k = self.session_string_key(key);
if self.contains_key_raw(&str_k).await? {
if self.has_ttl_tag(key)? && self.check_expired(key).await? {
return Ok(false);
}
return Ok(true);
}
Ok(false)
}
pub(crate) async fn contains_key_ignore_ttl(&self, key: &[u8]) -> Result<bool> {
let meta_k = self.session_meta_key(key);
if let Some(meta) = self.read_raw_with(&meta_k, MetaValue::from_slice).await? {
let meta = meta?;
if meta.size > 0 {
self
.store
.update_key_id_meta(meta.key_id, meta.version, true);
return Ok(true);
}
self
.store
.update_key_id_meta(meta.key_id, meta.version, false);
}
let str_k = self.session_string_key(key);
self.contains_key_raw(&str_k).await
}
#[inline(always)]
pub const fn get_meta_chunk_info(reserved: &[u8; 7]) -> (u32, u16) {
let chunk_id = u32::from_be_bytes([
reserved[1] & !META_HAS_EXPIRE_MASK,
reserved[2],
reserved[3],
reserved[4],
]);
let chunk_len = u16::from_be_bytes([reserved[5], reserved[6]]);
(chunk_id, chunk_len)
}
#[inline(always)]
pub const fn set_meta_chunk_info(reserved: &mut [u8; 7], chunk_id: u32, chunk_len: u16) {
let c = chunk_id.to_be_bytes();
reserved[1] = c[0] | (reserved[1] & META_HAS_EXPIRE_MASK);
reserved[2] = c[1];
reserved[3] = c[2];
reserved[4] = c[3];
let l = chunk_len.to_be_bytes();
reserved[5] = l[0];
reserved[6] = l[1];
}
#[inline(always)]
pub const fn get_meta_has_expire(reserved: &[u8; 7]) -> bool {
reserved[1] & META_HAS_EXPIRE_MASK != 0
}
#[inline(always)]
pub const fn set_meta_has_expire(reserved: &mut [u8; 7]) {
reserved[1] |= META_HAS_EXPIRE_MASK;
}
#[inline(always)]
pub(crate) const fn meta_value_has_expire(meta_bytes: &[u8]) -> bool {
meta_bytes.len() >= META_VALUE_SIZE
&& meta_bytes[META_RESERVED1_OFFSET] & META_HAS_EXPIRE_MASK != 0
}
#[inline]
pub(crate) fn probe_meta_has_expire(&self, user_key: &[u8]) -> Result<bool> {
let meta_k = self.session_meta_key(user_key);
match self.try_read_raw_in_memory(&meta_k, Self::meta_value_has_expire)? {
Some(Some(has)) => Ok(has),
Some(None) => Ok(false),
None => Ok(true),
}
}
#[inline]
pub fn check_object_meta_fast(&self, user_key: &[u8]) -> Result<Option<bool>> {
let meta_k = self.session_meta_key(user_key);
let _guard = self.participant.enter();
let Some(first_addr) = self.store.index.find_tag(&meta_k) else {
return Ok(Some(false));
};
match self.try_read_raw_in_memory_with_addr(&meta_k, Some(first_addr), |bytes| {
bytes.len() >= META_VALUE_SIZE && matches!(MetaValue::read_size(bytes), Ok(size) if size > 0)
})? {
Some(Some(is_active)) => Ok(Some(is_active)),
Some(None) => Ok(Some(false)),
None => Ok(None),
}
}
pub async fn load_meta(&self, user_key: &[u8]) -> Result<Option<MetaValue>> {
let meta_k = self.session_meta_key(user_key);
match self.read_raw_with(&meta_k, MetaValue::from_slice).await? {
Some(meta) => {
let meta = meta?;
if meta.size > 0 && self.has_ttl_tag(user_key)? && self.check_expired(user_key).await? {
return Ok(None);
}
self
.store
.update_key_id_meta(meta.key_id, meta.version, meta.size > 0);
Ok(Some(meta))
}
None => Ok(None),
}
}
pub async fn get_collection_chunk_info(&self, user_key: &[u8]) -> Result<Option<(u32, u16)>> {
if let Some(meta) = self.load_meta(user_key).await? {
Ok(Some(Self::get_meta_chunk_info(&meta.reserved)))
} else {
Ok(None)
}
}
pub async fn load_collection_raw_read(
&self,
key: &[u8],
expected: CollectionType,
) -> Result<Option<RawCollectionRead>> {
let Some((meta, bytes)) = self.load_live_meta_bytes(key).await? else {
return Ok(None);
};
if meta.collection_type != expected {
return Ok(None);
}
if meta.encoding() == StorageEncoding::Compact {
if meta.collection_type == CollectionType::Hash && Self::get_meta_has_expire(&meta.reserved) {
let now = now_ms();
let has_expired = bytes.len() > META_VALUE_SIZE
&& CompactHashCodec::iter_fields(&bytes[META_VALUE_SIZE..])
.any(|e| e.expire_at_ms.is_some_and(|exp| exp <= now));
if has_expired {
return self.purge_expired_hash_read(key).await;
}
}
return Ok(Some(RawCollectionRead::new(meta, Some(bytes))));
}
Ok(Some(RawCollectionRead::new(meta, None)))
}
async fn load_live_meta_bytes(&self, key: &[u8]) -> Result<Option<(MetaValue, Vec<u8>)>> {
let meta_k = self.session_meta_key(key);
let bytes = match self.read_raw(&meta_k).await? {
Some(b) => b,
None => {
if self.read(key).await?.is_some() {
return Err(wval::Error::InvalidCollectionType(0xFF).into());
}
return Ok(None);
}
};
let meta = MetaValue::from_slice(&bytes)?;
if meta.size == 0 {
if self.read(key).await?.is_some() {
return Err(wval::Error::InvalidCollectionType(0xFF).into());
}
return Ok(None);
}
if self.has_ttl_tag(key)? && self.check_expired(key).await? {
return Ok(None);
}
Ok(Some((meta, bytes)))
}
async fn purge_expired_hash_read(&self, key: &[u8]) -> Result<Option<RawCollectionRead>> {
let _key_lock = self.store.index.acquire_keys_lock_exclusive(&[key])?;
let Some((mut meta, Some(mut payload))) = self
.load_collection_raw_write(key, CollectionType::Hash)
.await?
else {
return Ok(None);
};
let purged = CompactHashCodec::purge_expired(&mut payload, now_ms())?;
if purged > 0 {
meta.dec_size(purged as u64);
self.save_compact_meta(key, &meta, &payload).await?;
if meta.size == 0 {
return Ok(None);
}
}
let mut fresh = Vec::with_capacity(META_VALUE_SIZE + payload.len());
fresh.extend_from_slice(&meta.to_bytes());
fresh.extend_from_slice(&payload);
Ok(Some(RawCollectionRead::new(meta, Some(fresh))))
}
pub async fn load_collection_raw_write(
&self,
key: &[u8],
expected: CollectionType,
) -> Result<Option<(MetaValue, Option<Vec<u8>>)>> {
let Some((meta, mut bytes)) = self.load_live_meta_bytes(key).await? else {
return Ok(None);
};
if meta.collection_type != expected {
return Err(wval::Error::InvalidCollectionType(meta.collection_type.as_u8()).into());
}
let payload = if meta.encoding() == StorageEncoding::Compact {
if bytes.len() > META_VALUE_SIZE {
bytes.drain(..META_VALUE_SIZE);
Some(bytes)
} else {
Some(Vec::new())
}
} else {
None
};
Ok(Some((meta, payload)))
}
pub async fn save_compact_meta(
&self,
user_key: &[u8],
meta: &MetaValue,
payload: &[u8],
) -> Result<()> {
self
.store
.update_key_id_meta(meta.key_id, meta.version, meta.size > 0);
let meta_k = self.session_meta_key(user_key);
if meta.size == 0 {
self.del_ttl(user_key).await?;
self.delete_raw(&meta_k).await?;
} else {
let total_len = META_VALUE_SIZE + payload.len();
const STACK_LIMIT: usize = 512;
if total_len <= STACK_LIMIT {
let mut buf = [0u8; STACK_LIMIT];
buf[..META_VALUE_SIZE].copy_from_slice(&meta.to_bytes());
buf[META_VALUE_SIZE..total_len].copy_from_slice(payload);
self.upsert_raw(&meta_k, &buf[..total_len]).await?;
} else {
let mut buf = Vec::with_capacity(total_len);
buf.extend_from_slice(&meta.to_bytes());
buf.extend_from_slice(payload);
self.upsert_raw(&meta_k, &buf).await?;
}
}
Ok(())
}
pub async fn save_meta(&self, user_key: &[u8], meta: &MetaValue) -> Result<()> {
self
.store
.update_key_id_meta(meta.key_id, meta.version, meta.size > 0);
let meta_k = self.session_meta_key(user_key);
if meta.size == 0 {
self.del_ttl(user_key).await?;
self.delete_raw(&meta_k).await?;
} else {
let bytes = meta.to_bytes();
self.upsert_raw(&meta_k, &bytes).await?;
}
Ok(())
}
pub async fn hexpire_at(
&self,
key: &[u8],
field: &[u8],
expire_at_ms: u64,
opt: TtlOpt,
) -> Result<i32> {
self
.hash_field_ttl(key, field, FieldTtlCmd::Expire { expire_at_ms, opt })
.await
}
pub async fn hpersist(&self, key: &[u8], field: &[u8]) -> Result<i32> {
self.hash_field_ttl(key, field, FieldTtlCmd::Persist).await
}
pub async fn collect_expired_hash_fields(&self, user_key: &[u8], now: u64) -> Result<u64> {
let _key_lock = self.store.index.acquire_keys_lock_exclusive(&[user_key])?;
let Some((mut meta, Some(mut payload))) = self
.load_collection_raw_write(user_key, CollectionType::Hash)
.await?
else {
return Ok(0);
};
if !Self::get_meta_has_expire(&meta.reserved) {
return Ok(0);
}
let purged = CompactHashCodec::purge_expired(&mut payload, now)?;
if purged > 0 {
meta.dec_size(purged as u64);
self.save_compact_meta(user_key, &meta, &payload).await?;
}
Ok(purged as u64)
}
async fn hash_field_ttl(&self, key: &[u8], field: &[u8], cmd: FieldTtlCmd) -> Result<i32> {
let _key_lock = self.store.index.acquire_keys_lock_exclusive(&[key])?;
let Some((mut meta, Some(mut payload))) = self
.load_collection_raw_write(key, CollectionType::Hash)
.await?
else {
return Ok(-2);
};
let purged = if Self::get_meta_has_expire(&meta.reserved) {
CompactHashCodec::purge_expired(&mut payload, now_ms())?
} else {
0
};
if purged > 0 {
meta.dec_size(purged as u64);
self.save_compact_meta(key, &meta, &payload).await?;
if meta.size == 0 {
return Ok(-1);
}
}
let Some(fv) = CompactHashCodec::find(&payload, field) else {
return Ok(-1);
};
match cmd {
FieldTtlCmd::Expire { expire_at_ms, opt } => {
let cond_fail = match fv.expire_at_ms {
Some(c) => opt.nx || (opt.gt && expire_at_ms <= c) || (opt.lt && expire_at_ms >= c),
None => opt.xx || opt.gt,
};
if cond_fail {
return Ok(0);
}
let now = now_ms();
if expire_at_ms <= now {
if CompactHashCodec::delete_field(&mut payload, field)? {
meta.dec_size(1);
}
self.save_compact_meta(key, &meta, &payload).await?;
return Ok(1);
}
let value = fv.value.to_vec();
CompactHashCodec::set_field(&mut payload, field, &value, Some(expire_at_ms))?;
Self::set_meta_has_expire(&mut meta.reserved);
self.save_compact_meta(key, &meta, &payload).await?;
Ok(1)
}
FieldTtlCmd::Persist => {
if fv.expire_at_ms.is_none() {
return Ok(0);
}
let value = fv.value.to_vec();
CompactHashCodec::set_field(&mut payload, field, &value, None)?;
self.save_compact_meta(key, &meta, &payload).await?;
Ok(1)
}
}
}
pub async fn append_hash_field(&self, meta: &mut MetaValue, field: &[u8]) -> Result<()> {
self.append_hash_fields_batch(meta, &[field]).await
}
async fn append_collection_chunks_batch(
&self,
tag: KeyTag,
meta: &mut MetaValue,
items: &[impl AsRef<[u8]>],
) -> Result<()> {
if items.is_empty() {
return Ok(());
}
let (mut max_chunk_id, mut curr_chunk_len) = Self::get_meta_chunk_info(&meta.reserved);
let mut chunk_buf: Option<Vec<u8>> = None;
for (idx, item) in items.iter().enumerate() {
let it = item.as_ref();
if curr_chunk_len == 0 || curr_chunk_len >= Self::CHUNK_CAPACITY {
if let Some(buf) = chunk_buf.take() {
let chunk_k = self.chunk_key(tag, meta.key_id, meta.version, max_chunk_id);
self.upsert_raw(&chunk_k, &buf).await?;
}
if curr_chunk_len >= Self::CHUNK_CAPACITY {
max_chunk_id = max_chunk_id.saturating_add(1);
}
let items_left = items.len().saturating_sub(idx);
let chunk_slots = (Self::CHUNK_CAPACITY as usize).min(items_left);
let est_cap = (4 + it.len()).saturating_mul(chunk_slots);
let mut buf = Vec::with_capacity(est_cap);
wrecord::ChunkCodec::append(it, &mut buf)?;
chunk_buf = Some(buf);
curr_chunk_len = 1;
} else {
let buf = match chunk_buf {
Some(ref mut b) => b,
None => {
let chunk_k = self.chunk_key(tag, meta.key_id, meta.version, max_chunk_id);
let mut b = self.read_raw(&chunk_k).await?.unwrap_or_default();
let items_left = items.len().saturating_sub(idx);
let chunk_slots =
(Self::CHUNK_CAPACITY.saturating_sub(curr_chunk_len) as usize).min(items_left);
let est_add_cap = (4 + it.len()).saturating_mul(chunk_slots);
b.reserve(est_add_cap);
chunk_buf.insert(b)
}
};
wrecord::ChunkCodec::append(it, buf)?;
curr_chunk_len = curr_chunk_len.saturating_add(1);
}
}
if let Some(buf) = chunk_buf {
let chunk_k = self.chunk_key(tag, meta.key_id, meta.version, max_chunk_id);
self.upsert_raw(&chunk_k, &buf).await?;
}
Self::set_meta_chunk_info(&mut meta.reserved, max_chunk_id, curr_chunk_len);
Ok(())
}
pub async fn append_hash_fields_batch(
&self,
meta: &mut MetaValue,
fields: &[impl AsRef<[u8]>],
) -> Result<()> {
self
.append_collection_chunks_batch(KeyTag::HashChunk, meta, fields)
.await
}
pub async fn append_set_member(&self, meta: &mut MetaValue, member: &[u8]) -> Result<()> {
self.append_set_members_batch(meta, &[member]).await
}
pub async fn append_set_members_batch(
&self,
meta: &mut MetaValue,
members: &[impl AsRef<[u8]>],
) -> Result<()> {
self
.append_collection_chunks_batch(KeyTag::SetChunk, meta, members)
.await
}
}
#[derive(Debug, Clone)]
pub struct RawCollectionRead {
pub meta: MetaValue,
raw: Option<Vec<u8>>,
}
impl RawCollectionRead {
#[inline(always)]
pub const fn new(meta: MetaValue, raw: Option<Vec<u8>>) -> Self {
Self { meta, raw }
}
#[inline(always)]
pub fn compact_payload(&self) -> Option<&[u8]> {
self.raw.as_deref().map(|b| {
if b.len() > META_VALUE_SIZE {
&b[META_VALUE_SIZE..]
} else {
&[]
}
})
}
}
#[derive(Debug, Clone, Copy)]
enum FieldTtlCmd {
Expire {
expire_at_ms: u64,
opt: TtlOpt,
},
Persist,
}