use rapidhash::{HashSetExt, RapidHashSet as HashSet};
use crate::{
api::sortedint::{
compose_si_item_key, compose_si_meta_key, compose_si_prefix, compose_si_prefix_stack,
conf::{BE_LEN, SortedintRangeSpec, encode_be_u64},
extract_id,
meta::SortedintMeta,
traits::SortedInt,
},
error::Result,
key::{clear_prefix_in_batch, get_meta_checked},
meta::{current_now_ms, normalize_range},
traits::DbLike,
};
fn prepare_si_meta_for_write(
db: &impl DbLike,
k_bytes: &[u8],
prefix: &[u8],
meta_k: &[u8],
now_ms: u64,
batch: &mut fjall::OwnedWriteBatch,
) -> Result<(SortedintMeta, bool)> {
match get_meta_checked::<SortedintMeta>(db, k_bytes, meta_k, now_ms)? {
Some(meta) => Ok((meta, false)),
None => {
clear_prefix_in_batch(db.data(), prefix, batch)?;
Ok((SortedintMeta::default(), true))
}
}
}
impl<T: DbLike> SortedInt for T {
fn si_iter<K: AsRef<[u8]>, F: FnMut(u64) -> bool>(&self, key: K, mut f: F) -> Result<()> {
let k_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_si_meta_key(&kc, k_bytes);
let now_ms = current_now_ms();
if get_meta_checked::<SortedintMeta>(self, k_bytes, &meta_k, now_ms)?.is_none() {
return Ok(());
}
let prefix = compose_si_prefix(&kc, k_bytes);
let prefix_len = prefix.len();
for g in self.data().prefix(&prefix) {
let (k, _) = g.into_inner()?;
if !k.starts_with(&prefix) {
break;
}
if let Some(id) = extract_id(&k, prefix_len)
&& !f(id)
{
break;
}
}
Ok(())
}
fn si_add<K: AsRef<[u8]>>(&self, key: K, ids: &[u64]) -> Result<usize> {
if ids.is_empty() {
return Ok(0);
}
let k_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_si_meta_key(&kc, k_bytes);
let prefix = compose_si_prefix(&kc, k_bytes);
let now_ms = current_now_ms();
let data_ks = self.data();
let meta_ks = self.meta();
let mut batch = self.batch();
let (mut meta, is_fresh) =
prepare_si_meta_for_write(self, k_bytes, &prefix, &meta_k, now_ms, &mut batch)?;
let mut added = 0usize;
let mut item_buf = compose_si_prefix_stack(&kc, k_bytes);
let prefix_len = item_buf.len();
item_buf.extend_from_slice(&[0u8; BE_LEN]);
if ids.len() == 1 {
let id = ids[0];
let be_bytes = encode_be_u64(id);
item_buf[prefix_len..].copy_from_slice(&be_bytes);
if is_fresh || !data_ks.contains_key(&item_buf)? {
added = 1;
meta.base.size = meta.base.size.saturating_add(1);
batch.insert(data_ks, &item_buf, b"");
}
} else {
let mut seen = HashSet::with_capacity(ids.len());
for &id in ids {
if !seen.insert(id) {
continue;
}
let be_bytes = encode_be_u64(id);
item_buf[prefix_len..].copy_from_slice(&be_bytes);
if is_fresh || !data_ks.contains_key(&item_buf)? {
added += 1;
meta.base.size = meta.base.size.saturating_add(1);
batch.insert(data_ks, &item_buf, b"");
}
}
}
if added > 0 || is_fresh {
batch.insert(meta_ks, &meta_k, meta.encode());
batch.commit()?;
}
Ok(added)
}
fn si_rem<K: AsRef<[u8]>>(&self, key: K, ids: &[u64]) -> Result<usize> {
if ids.is_empty() {
return Ok(0);
}
let k_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_si_meta_key(&kc, k_bytes);
let now_ms = current_now_ms();
let data_ks = self.data();
let meta_ks = self.meta();
let mut meta = match get_meta_checked::<SortedintMeta>(self, k_bytes, &meta_k, now_ms)? {
Some(m) => m,
None => return Ok(0),
};
let mut item_buf = compose_si_prefix_stack(&kc, k_bytes);
let prefix_len = item_buf.len();
item_buf.extend_from_slice(&[0u8; BE_LEN]);
let mut deleted = 0usize;
let mut batch = self.batch();
if ids.len() == 1 {
let id = ids[0];
let be_bytes = encode_be_u64(id);
item_buf[prefix_len..].copy_from_slice(&be_bytes);
if data_ks.contains_key(&item_buf)? {
deleted = 1;
meta.base.size = meta.base.size.saturating_sub(1);
batch.remove(data_ks, &item_buf);
}
} else {
let mut seen = HashSet::with_capacity(ids.len());
for &id in ids {
if !seen.insert(id) {
continue;
}
let be_bytes = encode_be_u64(id);
item_buf[prefix_len..].copy_from_slice(&be_bytes);
if data_ks.contains_key(&item_buf)? {
deleted += 1;
meta.base.size = meta.base.size.saturating_sub(1);
batch.remove(data_ks, &item_buf);
}
}
}
if deleted > 0 {
if meta.base.size == 0 {
batch.remove(meta_ks, &meta_k);
} else {
batch.insert(meta_ks, &meta_k, meta.encode());
}
batch.commit()?;
}
Ok(deleted)
}
fn si_card<K: AsRef<[u8]>>(&self, key: K) -> Result<u64> {
let k_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_si_meta_key(&kc, k_bytes);
let now_ms = current_now_ms();
Ok(
get_meta_checked::<SortedintMeta>(self, k_bytes, &meta_k, now_ms)?.map_or(0, |m| m.base.size),
)
}
fn si_exists<K: AsRef<[u8]>>(&self, key: K, id: u64) -> Result<bool> {
let k_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_si_meta_key(&kc, k_bytes);
let now_ms = current_now_ms();
if get_meta_checked::<SortedintMeta>(self, k_bytes, &meta_k, now_ms)?.is_none() {
return Ok(false);
}
let prefix = compose_si_prefix(&kc, k_bytes);
let item_k = compose_si_item_key(&prefix, id);
Ok(self.data().contains_key(&item_k)?)
}
fn si_mexist<K: AsRef<[u8]>>(&self, key: K, ids: &[u64]) -> Result<Vec<bool>> {
let mut results = Vec::with_capacity(ids.len());
if ids.is_empty() {
return Ok(results);
}
let k_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_si_meta_key(&kc, k_bytes);
let now_ms = current_now_ms();
if get_meta_checked::<SortedintMeta>(self, k_bytes, &meta_k, now_ms)?.is_none() {
results.resize(ids.len(), false);
return Ok(results);
}
let mut item_buf = compose_si_prefix_stack(&kc, k_bytes);
let prefix_len = item_buf.len();
item_buf.extend_from_slice(&[0u8; BE_LEN]);
let data_ks = self.data();
for &id in ids {
let be_bytes = encode_be_u64(id);
item_buf[prefix_len..].copy_from_slice(&be_bytes);
results.push(data_ks.contains_key(&item_buf)?);
}
Ok(results)
}
fn si_members<K: AsRef<[u8]>>(&self, key: K) -> Result<Vec<u64>> {
let k_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_si_meta_key(&kc, k_bytes);
let now_ms = current_now_ms();
let meta = match get_meta_checked::<SortedintMeta>(self, k_bytes, &meta_k, now_ms)? {
Some(m) => m,
None => return Ok(Vec::new()),
};
let mut results = Vec::with_capacity((meta.base.size as usize).min(4096));
self.si_iter(key, |id| {
results.push(id);
true
})?;
Ok(results)
}
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_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_si_meta_key(&kc, k_bytes);
let now_ms = current_now_ms();
if get_meta_checked::<SortedintMeta>(self, k_bytes, &meta_k, now_ms)?.is_none() {
return Ok(Vec::new());
}
let prefix = compose_si_prefix(&kc, k_bytes);
let prefix_len = prefix.len();
if !reversed {
let start_k = compose_si_item_key(&prefix, cursor);
let end_k = compose_si_item_key(&prefix, u64::MAX);
let mut results = Vec::with_capacity(limit.min(1024));
let mut pos = 0usize;
for g in self.data().range(start_k..=end_k) {
let (k, _) = g.into_inner()?;
if let Some(id) = extract_id(&k, prefix_len) {
if cursor > 0 && id == cursor {
continue;
}
if pos < offset {
pos += 1;
continue;
}
results.push(id);
if results.len() >= limit {
break;
}
}
}
Ok(results)
} else {
let start_k = compose_si_item_key(&prefix, 0);
let end_k = compose_si_item_key(&prefix, if cursor == 0 { u64::MAX } else { cursor });
let mut results = Vec::with_capacity(limit.min(1024));
let mut pos = 0usize;
for g in self.data().range(start_k..=end_k).rev() {
let (k, _) = g.into_inner()?;
if let Some(id) = extract_id(&k, prefix_len) {
if cursor > 0 && id == cursor {
continue;
}
if pos < offset {
pos += 1;
continue;
}
results.push(id);
if results.len() >= limit {
break;
}
}
}
Ok(results)
}
}
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(0) = spec.count {
return Ok(Vec::new());
}
let k_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_si_meta_key(&kc, k_bytes);
let now_ms = current_now_ms();
if get_meta_checked::<SortedintMeta>(self, k_bytes, &meta_k, now_ms)?.is_none() {
return Ok(Vec::new());
}
let prefix = compose_si_prefix(&kc, k_bytes);
let prefix_len = prefix.len();
let start_k = compose_si_item_key(&prefix, spec.min);
let end_k = compose_si_item_key(&prefix, spec.max);
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().range(start_k..=end_k) {
let (k, _) = g.into_inner()?;
if let Some(id) = extract_id(&k, prefix_len) {
if spec.minex && id == spec.min {
continue;
}
if spec.maxex && id == spec.max {
break;
}
if pos < spec.offset {
pos += 1;
continue;
}
results.push(id);
if let Some(cnt) = spec.count
&& results.len() >= cnt
{
break;
}
}
}
Ok(results)
} else {
let mut results = Vec::with_capacity(spec.count.unwrap_or(16).min(1024));
let mut pos = 0usize;
for g in self.data().range(start_k..=end_k).rev() {
let (k, _) = g.into_inner()?;
if let Some(id) = extract_id(&k, prefix_len) {
if spec.maxex && id == spec.max {
continue;
}
if spec.minex && id == spec.min {
break;
}
if pos < spec.offset {
pos += 1;
continue;
}
results.push(id);
if let Some(cnt) = spec.count
&& results.len() >= cnt
{
break;
}
}
}
Ok(results)
}
}
fn si_rem_range_by_value<K: AsRef<[u8]>>(
&self,
key: K,
spec: &SortedintRangeSpec,
) -> Result<usize> {
if spec.is_empty_range() {
return Ok(0);
}
let k_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_si_meta_key(&kc, k_bytes);
let now_ms = current_now_ms();
let mut meta = match get_meta_checked::<SortedintMeta>(self, k_bytes, &meta_k, now_ms)? {
Some(m) => m,
None => return Ok(0),
};
let prefix = compose_si_prefix(&kc, k_bytes);
let prefix_len = prefix.len();
let start_k = compose_si_item_key(&prefix, spec.min);
let end_k = compose_si_item_key(&prefix, spec.max);
let mut deleted = 0usize;
let mut batch = self.batch();
for g in self.data().range(start_k..=end_k) {
let (k, _) = g.into_inner()?;
if let Some(id) = extract_id(&k, prefix_len) {
if spec.minex && id == spec.min {
continue;
}
if spec.maxex && id == spec.max {
break;
}
deleted += 1;
batch.remove(self.data(), &*k);
}
}
if deleted > 0 {
meta.base.size = meta.base.size.saturating_sub(deleted as u64);
if meta.base.size == 0 {
batch.remove(self.meta(), &meta_k);
} else {
batch.insert(self.meta(), &meta_k, meta.encode());
}
batch.commit()?;
}
Ok(deleted)
}
fn si_rem_range_by_rank<K: AsRef<[u8]>>(&self, key: K, start: i64, stop: i64) -> Result<usize> {
let k_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_si_meta_key(&kc, k_bytes);
let now_ms = current_now_ms();
let data_ks = self.data();
let meta_ks = self.meta();
let mut meta = match get_meta_checked::<SortedintMeta>(self, k_bytes, &meta_k, now_ms)? {
Some(m) => m,
None => return Ok(0),
};
let card = meta.base.size 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 prefix = compose_si_prefix(&kc, k_bytes);
let prefix_len = prefix.len();
let mut skipped = 0usize;
let mut deleted = 0usize;
let mut batch = self.batch();
for g in data_ks.prefix(&prefix) {
let (k, _) = g.into_inner()?;
if !k.starts_with(&prefix) {
break;
}
if extract_id(&k, prefix_len).is_some() {
if skipped < offset {
skipped += 1;
continue;
}
deleted += 1;
batch.remove(data_ks, &*k);
if deleted >= limit {
break;
}
}
}
if deleted > 0 {
meta.base.size = meta.base.size.saturating_sub(deleted as u64);
if meta.base.size == 0 {
batch.remove(meta_ks, &meta_k);
} else {
batch.insert(meta_ks, &meta_k, meta.encode());
}
batch.commit()?;
}
Ok(deleted)
}
fn si_rank<K: AsRef<[u8]>>(&self, key: K, id: u64) -> Result<Option<usize>> {
let mut rank = 0usize;
let mut found = false;
self.si_iter(key, |cur_id| {
if cur_id == id {
found = true;
return false;
}
if cur_id > id {
return false;
}
rank += 1;
true
})?;
if found { Ok(Some(rank)) } else { Ok(None) }
}
fn si_revrank<K: AsRef<[u8]>>(&self, key: K, id: u64) -> Result<Option<usize>> {
let k_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_si_meta_key(&kc, k_bytes);
let now_ms = current_now_ms();
let meta = match get_meta_checked::<SortedintMeta>(self, k_bytes, &meta_k, now_ms)? {
Some(m) => m,
None => return Ok(None),
};
if let Some(rank) = self.si_rank(key, id)? {
Ok(Some(
(meta.base.size as usize)
.saturating_sub(1)
.saturating_sub(rank),
))
} else {
Ok(None)
}
}
fn si_count<K: AsRef<[u8]>>(&self, key: K, spec: &SortedintRangeSpec) -> Result<usize> {
if spec.is_empty_range() {
return Ok(0);
}
let k_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_si_meta_key(&kc, k_bytes);
let now_ms = current_now_ms();
if get_meta_checked::<SortedintMeta>(self, k_bytes, &meta_k, now_ms)?.is_none() {
return Ok(0);
}
let prefix = compose_si_prefix(&kc, k_bytes);
let prefix_len = prefix.len();
let start_k = compose_si_item_key(&prefix, spec.min);
let end_k = compose_si_item_key(&prefix, spec.max);
let mut count = 0usize;
for g in self.data().range(start_k..=end_k) {
let (k, _) = g.into_inner()?;
if let Some(id) = extract_id(&k, prefix_len) {
if spec.minex && id == spec.min {
continue;
}
if spec.maxex && id == spec.max {
break;
}
count += 1;
}
}
Ok(count)
}
fn si_clear<K: AsRef<[u8]>>(&self, key: K) -> Result<usize> {
let k_bytes = key.as_ref();
let kc = self.kc();
let meta_k = compose_si_meta_key(&kc, k_bytes);
let prefix = compose_si_prefix(&kc, k_bytes);
let data_ks = self.data();
let meta_ks = self.meta();
let mut deleted = 0usize;
let mut batch = self.batch();
for g in data_ks.prefix(&prefix) {
let (k, _) = g.into_inner()?;
if !k.starts_with(&prefix) {
break;
}
deleted += 1;
batch.remove(data_ks, &*k);
}
let has_meta = meta_ks.contains_key(&meta_k)?;
if has_meta {
batch.remove(meta_ks, &meta_k);
}
if deleted > 0 || has_meta {
batch.commit()?;
}
Ok(deleted)
}
}