pub mod conf;
pub mod meta;
pub use conf::{
HEX_CHARS, HEX_LEN, SortedintRangeSpec, decode_be_u64, decode_hex_u64, encode_be_u64,
encode_hex_u64, parse_range_spec,
};
pub use meta::SortedintMeta;
use rapidhash::RapidHashSet;
use std::str;
use crate::db::WeDb;
use crate::error::Result;
use crate::key_composer::KeyComposer;
use crate::meta::normalize_range;
const DEFAULT_NS: &str = "default";
#[inline(always)]
fn extract_id(key_bytes: &[u8], prefix_len: usize) -> Option<u64> {
let sub = key_bytes.get(prefix_len..)?;
if sub.len() >= HEX_LEN
&& let Some(id) = decode_hex_u64(&sub[..HEX_LEN])
{
return Some(id);
}
if sub.len() >= 8 {
decode_be_u64(&sub[..8])
} else {
None
}
}
#[derive(Debug, Clone)]
struct SizedRing<T> {
buf: Vec<T>,
cap: usize,
pos: usize,
len: usize,
}
impl<T: Copy> SizedRing<T> {
#[inline]
fn new(cap: usize) -> Self {
Self {
buf: Vec::with_capacity(cap),
cap,
pos: 0,
len: 0,
}
}
#[inline]
fn push(&mut self, item: T) {
if self.cap == 0 {
return;
}
if self.len < self.cap {
self.buf.push(item);
self.len += 1;
} else {
self.buf[self.pos] = item;
self.pos = (self.pos + 1) % self.cap;
}
}
#[inline]
fn into_rev_vec_with_offset_limit(self, offset: usize, limit: usize) -> Vec<T> {
if self.len <= offset || limit == 0 {
return Vec::new();
}
let available = self.len - offset;
let take_count = available.min(limit);
let mut out = Vec::with_capacity(take_count);
if self.len < self.cap {
for i in 0..take_count {
let idx = self.len - 1 - offset - i;
out.push(self.buf[idx]);
}
} else {
for i in 0..take_count {
let logical_rev_idx = offset + i;
let idx = (self.pos + self.cap * 2 - 1 - (logical_rev_idx % self.cap)) % self.cap;
out.push(self.buf[idx]);
}
}
out
}
}
impl WeDb {
#[inline]
fn get_active_sortedint_meta(
&self,
key_str: &str,
now_ms: u64,
) -> Result<Option<SortedintMeta>> {
let kc = KeyComposer::new(DEFAULT_NS);
let meta_k = kc.si_meta(key_str);
match self.meta_ks.get(meta_k.as_bytes())? {
Some(m_bytes) => {
if let Some(meta) = SortedintMeta::decode(&m_bytes)
&& !meta.is_expired(now_ms)
&& !meta.is_empty()
{
return Ok(Some(meta));
}
Ok(None)
}
None => Ok(None),
}
}
pub fn si_add<K: AsRef<[u8]>>(&self, key: K, ids: &[u64]) -> Result<usize> {
if ids.is_empty() {
return Ok(0);
}
let kc = KeyComposer::new(DEFAULT_NS);
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.si_meta(k_str);
let now_ms = ts_::sec() * 1000;
let mut meta = match self.meta_ks.get(meta_k.as_bytes())? {
Some(m_bytes) => {
if let Some(m) = SortedintMeta::decode(&m_bytes) {
if m.is_expired(now_ms) {
let prefix = kc.si_prefix(k_str);
let mut clear_batch = self.db.batch();
for g in self.data_ks.prefix(&prefix) {
let (k, _) = g.into_inner()?;
if !k.starts_with(&prefix) {
break;
}
clear_batch.remove(&self.data_ks, &*k);
}
clear_batch.remove(&self.meta_ks, meta_k.as_bytes());
clear_batch.commit()?;
SortedintMeta::default()
} else {
m
}
} else {
SortedintMeta::default()
}
}
None => SortedintMeta::default(),
};
let mut added = 0;
let mut batch = self.db.batch();
let mut seen = RapidHashSet::with_capacity_and_hasher(ids.len(), Default::default());
let mut item_buf = kc.si_prefix(k_str);
let prefix_len = item_buf.len();
item_buf.resize(prefix_len + HEX_LEN, 0);
for &id in ids {
if !seen.insert(id) {
continue;
}
let hex_bytes = encode_hex_u64(id);
item_buf[prefix_len..].copy_from_slice(&hex_bytes);
if !self.data_ks.contains_key(&item_buf)? {
added += 1;
meta.base.size = meta.base.size.saturating_add(1);
batch.insert(&self.data_ks, &item_buf, b"");
}
}
if added > 0 {
batch.insert(&self.meta_ks, meta_k.as_bytes(), meta.encode());
batch.commit()?;
}
Ok(added)
}
pub fn si_rem<K: AsRef<[u8]>>(&self, key: K, ids: &[u64]) -> Result<usize> {
if ids.is_empty() {
return Ok(0);
}
let kc = KeyComposer::new(DEFAULT_NS);
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let now_ms = ts_::sec() * 1000;
let mut meta = match self.get_active_sortedint_meta(k_str, now_ms)? {
Some(m) => m,
None => return Ok(0),
};
let meta_k = kc.si_meta(k_str);
let mut deleted = 0;
let mut batch = self.db.batch();
let mut seen = RapidHashSet::with_capacity_and_hasher(ids.len(), Default::default());
let mut item_buf = kc.si_prefix(k_str);
let prefix_len = item_buf.len();
item_buf.resize(prefix_len + HEX_LEN, 0);
for &id in ids {
if !seen.insert(id) {
continue;
}
let hex_bytes = encode_hex_u64(id);
item_buf[prefix_len..].copy_from_slice(&hex_bytes);
if self.data_ks.contains_key(&item_buf)? {
deleted += 1;
meta.base.size = meta.base.size.saturating_sub(1);
batch.remove(&self.data_ks, &item_buf);
}
}
if deleted > 0 {
if meta.base.size == 0 {
batch.remove(&self.meta_ks, meta_k.as_bytes());
} else {
batch.insert(&self.meta_ks, meta_k.as_bytes(), meta.encode());
}
batch.commit()?;
}
Ok(deleted)
}
pub fn si_card<K: AsRef<[u8]>>(&self, key: K) -> Result<u64> {
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let now_ms = ts_::sec() * 1000;
Ok(self
.get_active_sortedint_meta(k_str, now_ms)?
.map_or(0, |m| m.base.size))
}
pub fn si_exists<K: AsRef<[u8]>>(&self, key: K, id: u64) -> Result<bool> {
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let now_ms = ts_::sec() * 1000;
if self.get_active_sortedint_meta(k_str, now_ms)?.is_none() {
return Ok(false);
}
let kc = KeyComposer::new(DEFAULT_NS);
let mut item_buf = kc.si_prefix(k_str);
let hex_bytes = encode_hex_u64(id);
item_buf.extend_from_slice(&hex_bytes);
Ok(self.data_ks.contains_key(&item_buf)?)
}
pub fn si_mexist<K: AsRef<[u8]>>(&self, key: K, ids: &[u64]) -> Result<Vec<bool>> {
let mut results = Vec::with_capacity(ids.len());
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let now_ms = ts_::sec() * 1000;
if self.get_active_sortedint_meta(k_str, now_ms)?.is_none() {
results.resize(ids.len(), false);
return Ok(results);
}
let kc = KeyComposer::new(DEFAULT_NS);
let mut item_buf = kc.si_prefix(k_str);
let prefix_len = item_buf.len();
item_buf.resize(prefix_len + HEX_LEN, 0);
for &id in ids {
let hex_bytes = encode_hex_u64(id);
item_buf[prefix_len..].copy_from_slice(&hex_bytes);
results.push(self.data_ks.contains_key(&item_buf)?);
}
Ok(results)
}
pub fn si_range<K: AsRef<[u8]>>(
&self,
key: K,
cursor: u64,
offset: usize,
limit: usize,
reversed: bool,
) -> Result<Vec<u64>> {
if limit == 0 {
return Ok(Vec::new());
}
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let now_ms = ts_::sec() * 1000;
if self.get_active_sortedint_meta(k_str, now_ms)?.is_none() {
return Ok(Vec::new());
}
let kc = KeyComposer::new(DEFAULT_NS);
let prefix = kc.si_prefix(k_str);
let prefix_len = prefix.len();
if !reversed {
let mut results = Vec::with_capacity(limit.min(1024));
let mut skipped = 0usize;
for g in self.data_ks.prefix(&prefix) {
let (k, _) = g.into_inner()?;
if !k.starts_with(&prefix) {
break;
}
if let Some(id) = extract_id(&k, prefix_len) {
if cursor > 0 && id <= cursor {
continue;
}
if skipped < offset {
skipped += 1;
continue;
}
results.push(id);
if results.len() >= limit {
break;
}
}
}
Ok(results)
} else {
let needed = offset.saturating_add(limit);
let mut ring = SizedRing::new(needed);
for g in self.data_ks.prefix(&prefix) {
let (k, _) = g.into_inner()?;
if !k.starts_with(&prefix) {
break;
}
if let Some(id) = extract_id(&k, prefix_len) {
if cursor > 0 && id >= cursor {
break;
}
ring.push(id);
}
}
Ok(ring.into_rev_vec_with_offset_limit(offset, limit))
}
}
#[inline]
pub fn si_rev_range<K: AsRef<[u8]>>(
&self,
key: K,
cursor: u64,
offset: usize,
limit: usize,
) -> Result<Vec<u64>> {
self.si_range(key, cursor, offset, limit, true)
}
pub fn si_range_by_value<K: AsRef<[u8]>>(
&self,
key: K,
spec: &SortedintRangeSpec,
) -> Result<Vec<u64>> {
if spec.is_empty_range() {
return Ok(Vec::new());
}
if let Some(cnt) = spec.count
&& cnt == 0
{
return Ok(Vec::new());
}
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let now_ms = ts_::sec() * 1000;
if self.get_active_sortedint_meta(k_str, now_ms)?.is_none() {
return Ok(Vec::new());
}
let kc = KeyComposer::new(DEFAULT_NS);
let prefix = kc.si_prefix(k_str);
let prefix_len = prefix.len();
if !spec.reversed {
let mut results = Vec::with_capacity(spec.count.unwrap_or(16).min(1024));
let mut pos = 0usize;
for g in self.data_ks.prefix(&prefix) {
let (k, _) = g.into_inner()?;
if !k.starts_with(&prefix) {
break;
}
if let Some(id) = extract_id(&k, prefix_len) {
if (spec.minex && id == spec.min) || id < spec.min {
continue;
}
if (spec.maxex && id == spec.max) || id > spec.max {
break;
}
if pos < spec.offset {
pos += 1;
continue;
}
results.push(id);
if let Some(cnt) = spec.count
&& cnt > 0
&& results.len() >= cnt
{
break;
}
}
}
Ok(results)
} else {
if let Some(cnt) = spec.count {
let needed = spec.offset.saturating_add(cnt);
let mut ring = SizedRing::new(needed);
for g in self.data_ks.prefix(&prefix) {
let (k, _) = g.into_inner()?;
if !k.starts_with(&prefix) {
break;
}
if let Some(id) = extract_id(&k, prefix_len) {
if (spec.minex && id == spec.min) || id < spec.min {
continue;
}
if (spec.maxex && id == spec.max) || id > spec.max {
break;
}
ring.push(id);
}
}
Ok(ring.into_rev_vec_with_offset_limit(spec.offset, cnt))
} else {
let mut matched_ids = Vec::new();
for g in self.data_ks.prefix(&prefix) {
let (k, _) = g.into_inner()?;
if !k.starts_with(&prefix) {
break;
}
if let Some(id) = extract_id(&k, prefix_len) {
if (spec.minex && id == spec.min) || id < spec.min {
continue;
}
if (spec.maxex && id == spec.max) || id > spec.max {
break;
}
matched_ids.push(id);
}
}
let total = matched_ids.len();
if total <= spec.offset {
return Ok(Vec::new());
}
matched_ids.truncate(total - spec.offset);
matched_ids.reverse();
Ok(matched_ids)
}
}
}
#[inline]
pub fn si_rev_range_by_value<K: AsRef<[u8]>>(
&self,
key: K,
spec: &SortedintRangeSpec,
) -> Result<Vec<u64>> {
let mut rev_spec = spec.clone();
rev_spec.reversed = true;
self.si_range_by_value(key, &rev_spec)
}
#[inline]
pub fn si_range_by_score<K: AsRef<[u8]>>(
&self,
key: K,
spec: &SortedintRangeSpec,
) -> Result<Vec<u64>> {
self.si_range_by_value(key, spec)
}
#[inline]
pub fn si_rev_range_by_score<K: AsRef<[u8]>>(
&self,
key: K,
spec: &SortedintRangeSpec,
) -> Result<Vec<u64>> {
self.si_rev_range_by_value(key, spec)
}
pub fn si_rem_range_by_value<K: AsRef<[u8]>>(
&self,
key: K,
spec: &SortedintRangeSpec,
) -> Result<usize> {
let ids = self.si_range_by_value(&key, spec)?;
if ids.is_empty() {
return Ok(0);
}
self.si_rem(key, &ids)
}
#[inline]
pub fn si_rem_range_by_score<K: AsRef<[u8]>>(
&self,
key: K,
spec: &SortedintRangeSpec,
) -> Result<usize> {
self.si_rem_range_by_value(key, spec)
}
pub fn si_rem_range_by_rank<K: AsRef<[u8]>>(
&self,
key: K,
start: i64,
stop: i64,
) -> Result<usize> {
let card = self.si_card(&key)? as i64;
if card == 0 {
return Ok(0);
}
let (s, e) = normalize_range(start, stop, card);
if s > e {
return Ok(0);
}
let offset = s as usize;
let limit = (e - s + 1) as usize;
let ids = self.si_range(&key, 0, offset, limit, false)?;
if ids.is_empty() {
return Ok(0);
}
self.si_rem(key, &ids)
}
pub fn si_rank<K: AsRef<[u8]>>(&self, key: K, id: u64) -> Result<Option<usize>> {
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let now_ms = ts_::sec() * 1000;
if self.get_active_sortedint_meta(k_str, now_ms)?.is_none() {
return Ok(None);
}
let kc = KeyComposer::new(DEFAULT_NS);
let prefix = kc.si_prefix(k_str);
let prefix_len = prefix.len();
let mut rank = 0usize;
for g in self.data_ks.prefix(&prefix) {
let (k, _) = g.into_inner()?;
if !k.starts_with(&prefix) {
break;
}
if let Some(cur_id) = extract_id(&k, prefix_len) {
if cur_id == id {
return Ok(Some(rank));
}
if cur_id > id {
return Ok(None);
}
rank += 1;
}
}
Ok(None)
}
pub fn si_revrank<K: AsRef<[u8]>>(&self, key: K, id: u64) -> Result<Option<usize>> {
let card = self.si_card(&key)?;
if card == 0 {
return Ok(None);
}
if let Some(rank) = self.si_rank(key, id)? {
Ok(Some((card as usize) - 1 - rank))
} else {
Ok(None)
}
}
#[inline]
pub fn si_count<K: AsRef<[u8]>>(&self, key: K, spec: &SortedintRangeSpec) -> Result<usize> {
let results = self.si_range_by_value(key, spec)?;
Ok(results.len())
}
pub fn si_clear<K: AsRef<[u8]>>(&self, key: K) -> Result<usize> {
let kc = KeyComposer::new(DEFAULT_NS);
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.si_meta(k_str);
let prefix = kc.si_prefix(k_str);
let mut deleted = 0;
let mut batch = self.db.batch();
for g in self.data_ks.prefix(&prefix) {
let (k, _) = g.into_inner()?;
if !k.starts_with(&prefix) {
break;
}
deleted += 1;
batch.remove(&self.data_ks, &*k);
}
let has_meta = self.meta_ks.contains_key(meta_k.as_bytes())?;
if has_meta {
batch.remove(&self.meta_ks, meta_k.as_bytes());
}
if deleted > 0 || has_meta {
batch.commit()?;
}
Ok(deleted)
}
}