pub mod conf;
pub mod meta;
pub use conf::{ERR_WRONG_TYPE, SetScanResult};
pub use meta::SetMeta;
use rapidhash::RapidHashSet as HashSet;
use crate::db::WeDb;
use crate::error::Result;
use crate::key_composer::{KeyComposer, KeyTag, SmallKey, SubkeyComposer, matches_glob_bytes};
use crate::meta::current_now_ms;
#[inline]
fn compose_set_meta_key(kc: &KeyComposer<'_>, key: &[u8]) -> Vec<u8> {
kc.compose_meta_key(KeyTag::SetMeta.as_slice(), key)
}
#[inline]
fn compose_set_prefix(kc: &KeyComposer<'_>, key: &[u8]) -> Vec<u8> {
kc.compose_prefix(KeyTag::SetData.as_slice(), key)
}
#[inline]
fn compose_set_prefix_stack(kc: &KeyComposer<'_>, key: &[u8]) -> SmallKey {
kc.compose_prefix_stack(KeyTag::SetData.as_slice(), key)
}
#[derive(Debug, Clone)]
pub struct SetItemKeyComposer {
composer: SubkeyComposer,
}
impl SetItemKeyComposer {
#[inline]
pub fn new(kc: &KeyComposer<'_>, key: &[u8]) -> Self {
let prefix = compose_set_prefix_stack(kc, key);
Self {
composer: SubkeyComposer::from_slice(&prefix),
}
}
#[inline(always)]
pub fn key_for_member<'a>(&'a mut self, member: &[u8]) -> &'a [u8] {
self.composer.key_for(member)
}
#[inline(always)]
pub fn prefix(&self) -> &[u8] {
self.composer.prefix()
}
}
impl WeDb {
#[inline]
fn get_set_meta_checked(
&self,
kc: &KeyComposer<'_>,
k: &[u8],
mk: &[u8],
now: u64,
) -> Result<Option<SetMeta>> {
self.get_meta_checked(kc, k, mk, now)
}
#[inline]
fn prepare_set_meta_for_write(
&self,
kc: &KeyComposer<'_>,
k_bytes: &[u8],
prefix: &[u8],
meta_k: &[u8],
now_ms: u64,
batch: &mut fjall::OwnedWriteBatch,
) -> Result<(SetMeta, bool)> {
match self.meta.get(meta_k)? {
Some(m_bytes) => match SetMeta::decode(&m_bytes) {
Some(meta) if !meta.is_expired(now_ms) => Ok((meta, false)),
Some(_) => {
self.check_key_not_other_type(kc, k_bytes, KeyTag::SetMeta.as_slice(), now_ms)?;
self.clear_prefix_in_batch(prefix, batch)?;
Ok((SetMeta::new_with_version(0, 0), true))
}
None => {
self.check_key_not_other_type(kc, k_bytes, KeyTag::SetMeta.as_slice(), now_ms)?;
Ok((SetMeta::new_with_version(0, 0), true))
}
},
None => {
self.check_key_not_other_type(kc, k_bytes, KeyTag::SetMeta.as_slice(), now_ms)?;
Ok((SetMeta::new_with_version(0, 0), true))
}
}
}
#[inline]
pub fn siter<K: AsRef<[u8]>, F>(&self, key: K, f: F) -> Result<()>
where
F: FnMut(&[u8]) -> bool,
{
self.siter_with_kc(&KeyComposer::new("default"), key, f)
}
pub fn siter_with_kc<K: AsRef<[u8]>, F>(
&self,
kc: &KeyComposer<'_>,
key: K,
mut f: F,
) -> Result<()>
where
F: FnMut(&[u8]) -> bool,
{
let key_bytes = key.as_ref();
let meta_k = compose_set_meta_key(kc, key_bytes);
let now_ms = current_now_ms();
match self.get_set_meta_checked(kc, key_bytes, &meta_k, now_ms)? {
Some(m) if !m.is_empty() => m,
_ => return Ok(()),
};
let prefix = compose_set_prefix(kc, key_bytes);
for g in self.data.prefix(&prefix) {
let (k, _) = g.into_inner()?;
if !k.starts_with(&prefix) {
break;
}
let member = &k[prefix.len()..];
if !f(member) {
break;
}
}
Ok(())
}
pub fn sadd<K: AsRef<[u8]>, M: AsRef<[u8]>>(&self, key: K, members: &[M]) -> Result<usize> {
self.sadd_with_kc(&KeyComposer::new("default"), key, members)
}
pub fn sadd_with_kc<K: AsRef<[u8]>, M: AsRef<[u8]>>(
&self,
kc: &KeyComposer<'_>,
key: K,
members: &[M],
) -> Result<usize> {
if members.is_empty() {
return Ok(0);
}
let key_bytes = key.as_ref();
let meta_k = compose_set_meta_key(kc, key_bytes);
let mut composer = SetItemKeyComposer::new(kc, key_bytes);
let now_ms = current_now_ms();
let mut batch = self.db.batch();
let (mut meta, is_fresh) = self.prepare_set_meta_for_write(
kc,
key_bytes,
composer.prefix(),
&meta_k,
now_ms,
&mut batch,
)?;
let mut added = 0usize;
if members.len() == 1 {
let m_bytes = members[0].as_ref();
let item_k = composer.key_for_member(m_bytes);
if is_fresh || !self.data.contains_key(item_k)? {
added = 1;
meta.base.size = meta.base.size.saturating_add(1);
batch.insert(&self.data, item_k, b"");
}
} else {
let mut seen = HashSet::with_capacity_and_hasher(members.len(), Default::default());
for m in members {
let m_bytes = m.as_ref();
if !seen.insert(m_bytes) {
continue;
}
let item_k = composer.key_for_member(m_bytes);
if is_fresh || !self.data.contains_key(item_k)? {
added += 1;
meta.base.size = meta.base.size.saturating_add(1);
batch.insert(&self.data, item_k, b"");
}
}
}
if added > 0 {
batch.insert(&self.meta, &meta_k, meta.encode());
batch.commit()?;
}
Ok(added)
}
pub fn srem<K: AsRef<[u8]>, M: AsRef<[u8]>>(&self, key: K, members: &[M]) -> Result<usize> {
self.srem_with_kc(&KeyComposer::new("default"), key, members)
}
pub fn srem_with_kc<K: AsRef<[u8]>, M: AsRef<[u8]>>(
&self,
kc: &KeyComposer<'_>,
key: K,
members: &[M],
) -> Result<usize> {
if members.is_empty() {
return Ok(0);
}
let key_bytes = key.as_ref();
let meta_k = compose_set_meta_key(kc, key_bytes);
let now_ms = current_now_ms();
let mut meta = match self.get_set_meta_checked(kc, key_bytes, &meta_k, now_ms)? {
Some(m) if !m.is_empty() => m,
_ => return Ok(0),
};
let mut deleted = 0usize;
let mut batch = self.db.batch();
let mut seen = HashSet::with_capacity_and_hasher(members.len(), Default::default());
let mut composer = SetItemKeyComposer::new(kc, key_bytes);
for m in members {
let m_bytes = m.as_ref();
if !seen.insert(m_bytes) {
continue;
}
let item_k = composer.key_for_member(m_bytes);
if self.data.contains_key(item_k)? {
deleted += 1;
batch.remove(&self.data, item_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)
}
pub fn smembers<K: AsRef<[u8]>>(&self, key: K) -> Result<Vec<Vec<u8>>> {
self.smembers_with_kc(&KeyComposer::new("default"), key)
}
pub fn smembers_with_kc<K: AsRef<[u8]>>(
&self,
kc: &KeyComposer<'_>,
key: K,
) -> Result<Vec<Vec<u8>>> {
let key_bytes = key.as_ref();
let meta_k = compose_set_meta_key(kc, key_bytes);
let now_ms = current_now_ms();
let meta = match self.get_set_meta_checked(kc, key_bytes, &meta_k, now_ms)? {
Some(m) if !m.is_empty() => m,
_ => return Ok(Vec::new()),
};
let mut members = Vec::with_capacity(meta.size() as usize);
self.siter_with_kc(kc, key_bytes, |m| {
members.push(m.to_vec());
true
})?;
Ok(members)
}
pub fn sismember<K: AsRef<[u8]>, M: AsRef<[u8]>>(&self, key: K, member: M) -> Result<bool> {
self.sismember_with_kc(&KeyComposer::new("default"), key, member)
}
pub fn sismember_with_kc<K: AsRef<[u8]>, M: AsRef<[u8]>>(
&self,
kc: &KeyComposer<'_>,
key: K,
member: M,
) -> Result<bool> {
let key_bytes = key.as_ref();
let meta_k = compose_set_meta_key(kc, key_bytes);
let now_ms = current_now_ms();
if match self.get_set_meta_checked(kc, key_bytes, &meta_k, now_ms)? {
Some(m) => m.is_empty(),
None => true,
} {
return Ok(false);
}
let mut composer = SetItemKeyComposer::new(kc, key_bytes);
let item_k = composer.key_for_member(member.as_ref());
Ok(self.data.contains_key(item_k)?)
}
pub fn smismember<K: AsRef<[u8]>, M: AsRef<[u8]>>(
&self,
key: K,
members: &[M],
) -> Result<Vec<bool>> {
if members.is_empty() {
return Ok(Vec::new());
}
let key_bytes = key.as_ref();
let kc = KeyComposer::new("default");
let meta_k = compose_set_meta_key(&kc, key_bytes);
let now_ms = current_now_ms();
if match self.get_set_meta_checked(&kc, key_bytes, &meta_k, now_ms)? {
Some(m) => m.is_empty(),
None => true,
} {
return Ok(vec![false; members.len()]);
}
let mut composer = SetItemKeyComposer::new(&kc, key_bytes);
let mut results = Vec::with_capacity(members.len());
for m in members {
let item_k = composer.key_for_member(m.as_ref());
results.push(self.data.contains_key(item_k)?);
}
Ok(results)
}
pub fn scard<K: AsRef<[u8]>>(&self, key: K) -> Result<u64> {
self.scard_with_kc(&KeyComposer::new("default"), key)
}
pub fn scard_with_kc<K: AsRef<[u8]>>(&self, kc: &KeyComposer<'_>, key: K) -> Result<u64> {
let key_bytes = key.as_ref();
let meta_k = compose_set_meta_key(kc, key_bytes);
let now_ms = current_now_ms();
match self.get_set_meta_checked(kc, key_bytes, &meta_k, now_ms)? {
Some(meta) => Ok(meta.base.size),
None => Ok(0),
}
}
pub fn smove<S: AsRef<[u8]>, D: AsRef<[u8]>, M: AsRef<[u8]>>(
&self,
src: S,
dst: D,
member: M,
) -> Result<bool> {
let src_bytes = src.as_ref();
let dst_bytes = dst.as_ref();
let m_bytes = member.as_ref();
if src_bytes == dst_bytes {
return self.sismember(src_bytes, m_bytes);
}
let kc = KeyComposer::new("default");
let now_ms = current_now_ms();
self.check_key_not_other_type(&kc, dst_bytes, KeyTag::SetMeta.as_slice(), now_ms)?;
let src_meta_k = compose_set_meta_key(&kc, src_bytes);
let mut src_meta = match self.get_set_meta_checked(&kc, src_bytes, &src_meta_k, now_ms)? {
Some(m) if !m.is_empty() => m,
_ => return Ok(false),
};
let mut src_composer = SetItemKeyComposer::new(&kc, src_bytes);
let src_item_k = src_composer.key_for_member(m_bytes);
if !self.data.contains_key(src_item_k)? {
return Ok(false);
}
let mut batch = self.db.batch();
src_meta.base.size = src_meta.base.size.saturating_sub(1);
batch.remove(&self.data, src_item_k);
if src_meta.base.size == 0 {
batch.remove(&self.meta, &src_meta_k);
} else {
batch.insert(&self.meta, &src_meta_k, src_meta.encode());
}
let dst_meta_k = compose_set_meta_key(&kc, dst_bytes);
let mut dst_composer = SetItemKeyComposer::new(&kc, dst_bytes);
let (mut dst_meta, dst_is_fresh) = self.prepare_set_meta_for_write(
&kc,
dst_bytes,
dst_composer.prefix(),
&dst_meta_k,
now_ms,
&mut batch,
)?;
let dst_item_k = dst_composer.key_for_member(m_bytes);
let dst_needs_insert = dst_is_fresh || !self.data.contains_key(dst_item_k)?;
if dst_needs_insert {
dst_meta.base.size = dst_meta.base.size.saturating_add(1);
batch.insert(&self.data, dst_item_k, b"");
batch.insert(&self.meta, &dst_meta_k, dst_meta.encode());
}
batch.commit()?;
Ok(true)
}
pub fn spop<K: AsRef<[u8]>>(&self, key: K, count: usize) -> Result<Vec<Vec<u8>>> {
self.spop_with_kc(&KeyComposer::new("default"), key, count)
}
pub fn spop_with_kc<K: AsRef<[u8]>>(
&self,
kc: &KeyComposer<'_>,
key: K,
count: usize,
) -> Result<Vec<Vec<u8>>> {
if count == 0 {
return Ok(Vec::new());
}
let key_bytes = key.as_ref();
let meta_k = compose_set_meta_key(kc, key_bytes);
let now_ms = current_now_ms();
let mut meta = match self.get_set_meta_checked(kc, key_bytes, &meta_k, now_ms)? {
Some(m) if !m.is_empty() => m,
_ => return Ok(Vec::new()),
};
let mut all = self.smembers_with_kc(kc, key_bytes)?;
if all.is_empty() {
return Ok(Vec::new());
}
let n = all.len();
let num_pop = count.min(n);
let mut composer = SetItemKeyComposer::new(kc, key_bytes);
if num_pop == n {
let mut batch = self.db.batch();
for m in &all {
let item_k = composer.key_for_member(m);
batch.remove(&self.data, item_k);
}
batch.remove(&self.meta, &meta_k);
batch.commit()?;
return Ok(all);
}
for i in 0..num_pop {
let j = fastrand::usize(i..n);
all.swap(i, j);
}
let mut batch = self.db.batch();
for m in &all[..num_pop] {
let item_k = composer.key_for_member(m);
batch.remove(&self.data, item_k);
}
meta.base.size = meta.base.size.saturating_sub(num_pop as u64);
batch.insert(&self.meta, &meta_k, meta.encode());
batch.commit()?;
all.truncate(num_pop);
Ok(all)
}
pub fn srandmember<K: AsRef<[u8]>>(&self, key: K, count: i64) -> Result<Vec<Vec<u8>>> {
if count == 0 {
return Ok(Vec::new());
}
let mut all = self.smembers(key)?;
if all.is_empty() {
return Ok(Vec::new());
}
let n = all.len();
if count > 0 {
let k = (count as usize).min(n);
if k == n {
return Ok(all);
}
for i in 0..k {
let j = fastrand::usize(i..n);
all.swap(i, j);
}
all.truncate(k);
Ok(all)
} else {
let num = count.unsigned_abs() as usize;
let mut results = Vec::with_capacity(num);
for _ in 0..num {
let idx = fastrand::usize(..n);
results.push(all[idx].clone());
}
Ok(results)
}
}
pub fn overwrite_set<K: AsRef<[u8]>, M: AsRef<[u8]>>(
&self,
key: K,
members: &[M],
) -> Result<usize> {
let key_bytes = key.as_ref();
let kc = KeyComposer::new("default");
let meta_k = compose_set_meta_key(&kc, key_bytes);
let mut composer = SetItemKeyComposer::new(&kc, key_bytes);
let prefix = composer.prefix().to_vec();
let now_ms = current_now_ms();
self.check_key_not_other_type(&kc, key_bytes, KeyTag::SetMeta.as_slice(), now_ms)?;
let mut batch = self.db.batch();
for g in self.data.prefix(&prefix) {
let (k, _) = g.into_inner()?;
if k.starts_with(&prefix) {
batch.remove(&self.data, &*k);
} else {
break;
}
}
batch.remove(&self.meta, &meta_k);
let raw_k = kc.raw_key_bytes(key_bytes);
batch.remove(&self.data, &*raw_k);
let mut seen = HashSet::with_capacity_and_hasher(members.len(), Default::default());
let mut count = 0u64;
for m in members {
let m_bytes = m.as_ref();
if !seen.insert(m_bytes) {
continue;
}
let item_k = composer.key_for_member(m_bytes);
batch.insert(&self.data, item_k, b"");
count += 1;
}
if count > 0 {
let meta = SetMeta::new_with_version(0, count);
batch.insert(&self.meta, &meta_k, meta.encode());
}
batch.commit()?;
Ok(count as usize)
}
pub fn sdiff<K: AsRef<[u8]>>(&self, keys: &[K]) -> Result<Vec<Vec<u8>>> {
if keys.is_empty() {
return Ok(Vec::new());
}
let source = self.smembers(&keys[0])?;
if source.is_empty() || keys.len() == 1 {
return Ok(source);
}
let mut source_set: HashSet<Vec<u8>> =
HashSet::with_capacity_and_hasher(source.len(), Default::default());
source_set.extend(source);
for k in &keys[1..] {
if source_set.is_empty() {
break;
}
let card = self.scard(k)?;
if card == 0 {
continue;
}
if (source_set.len() as u64) * 4 < card {
let mut err = None;
source_set.retain(|m| match self.sismember(k, m) {
Ok(found) => !found,
Err(e) => {
err = Some(e);
false
}
});
if let Some(e) = err {
return Err(e);
}
} else {
self.siter(k, |m| {
source_set.remove(m);
!source_set.is_empty()
})?;
}
}
Ok(source_set.into_iter().collect())
}
pub fn sdiffstore<D: AsRef<[u8]>, K: AsRef<[u8]>>(&self, dst: D, keys: &[K]) -> Result<usize> {
let diff = self.sdiff(keys)?;
self.overwrite_set(dst, &diff)
}
pub fn sdiffcard<K: AsRef<[u8]>>(&self, keys: &[K], limit: usize) -> Result<usize> {
if keys.is_empty() {
return Ok(0);
}
let first_card = self.scard(&keys[0])? as usize;
if first_card == 0 {
return Ok(0);
}
if keys.len() == 1 {
return Ok(if limit == 0 {
first_card
} else {
first_card.min(limit)
});
}
let diff = self.sdiff(keys)?;
Ok(if limit == 0 {
diff.len()
} else {
diff.len().min(limit)
})
}
pub fn sunion<K: AsRef<[u8]>>(&self, keys: &[K]) -> Result<Vec<Vec<u8>>> {
if keys.is_empty() {
return Ok(Vec::new());
}
if keys.len() == 1 {
return self.smembers(&keys[0]);
}
let first_card = self.scard(&keys[0])? as usize;
let mut union_set = HashSet::with_capacity_and_hasher(first_card, Default::default());
for k in keys {
self.siter(k, |m| {
if !union_set.contains(m) {
union_set.insert(m.to_vec());
}
true
})?;
}
Ok(union_set.into_iter().collect())
}
pub fn sunionstore<D: AsRef<[u8]>, K: AsRef<[u8]>>(&self, dst: D, keys: &[K]) -> Result<usize> {
let union_res = self.sunion(keys)?;
self.overwrite_set(dst, &union_res)
}
pub fn sunioncard<K: AsRef<[u8]>>(&self, keys: &[K], limit: usize) -> Result<usize> {
if keys.is_empty() {
return Ok(0);
}
if keys.len() == 1 {
let card = self.scard(&keys[0])? as usize;
return Ok(if limit == 0 { card } else { card.min(limit) });
}
let mut union_set = HashSet::default();
if limit > 0 {
for k in keys {
self.siter(k, |m| {
if !union_set.contains(m) {
union_set.insert(m.to_vec());
}
union_set.len() < limit
})?;
if union_set.len() >= limit {
return Ok(limit);
}
}
Ok(union_set.len().min(limit))
} else {
for k in keys {
self.siter(k, |m| {
if !union_set.contains(m) {
union_set.insert(m.to_vec());
}
true
})?;
}
Ok(union_set.len())
}
}
pub fn sinter<K: AsRef<[u8]>>(&self, keys: &[K]) -> Result<Vec<Vec<u8>>> {
if keys.is_empty() {
return Ok(Vec::new());
}
if keys.len() == 1 {
return self.smembers(&keys[0]);
}
let mut key_cards: Vec<(&K, u64)> = Vec::with_capacity(keys.len());
for k in keys {
let card = self.scard(k)?;
if card == 0 {
return Ok(Vec::new());
}
key_cards.push((k, card));
}
key_cards.sort_unstable_by_key(|&(_, card)| card);
let smallest_key = key_cards[0].0;
let mut current: HashSet<Vec<u8>> =
HashSet::with_capacity_and_hasher(key_cards[0].1 as usize, Default::default());
self.siter(smallest_key, |m| {
current.insert(m.to_vec());
true
})?;
if current.is_empty() {
return Ok(Vec::new());
}
for &(k, next_card) in &key_cards[1..] {
if current.is_empty() {
return Ok(Vec::new());
}
if (current.len() as u64) * 4 < next_card {
let mut err = None;
current.retain(|m| match self.sismember(k, m) {
Ok(found) => found,
Err(e) => {
err = Some(e);
false
}
});
if let Some(e) = err {
return Err(e);
}
} else {
let mut next_set =
HashSet::with_capacity_and_hasher(current.len(), Default::default());
self.siter(k, |m| {
if current.contains(m) {
next_set.insert(m.to_vec());
}
next_set.len() < current.len()
})?;
current = next_set;
}
}
Ok(current.into_iter().collect())
}
pub fn sinterstore<D: AsRef<[u8]>, K: AsRef<[u8]>>(&self, dst: D, keys: &[K]) -> Result<usize> {
let inter_res = self.sinter(keys)?;
self.overwrite_set(dst, &inter_res)
}
pub fn sintercard<K: AsRef<[u8]>>(&self, keys: &[K], limit: usize) -> Result<usize> {
if keys.is_empty() {
return Ok(0);
}
if keys.len() == 1 {
let card = self.scard(&keys[0])? as usize;
return Ok(if limit == 0 { card } else { card.min(limit) });
}
let mut key_cards: Vec<(&K, u64)> = Vec::with_capacity(keys.len());
for k in keys {
let card = self.scard(k)?;
if card == 0 {
return Ok(0);
}
key_cards.push((k, card));
}
key_cards.sort_unstable_by_key(|&(_, card)| card);
let smallest_key = key_cards[0].0;
let mut current: HashSet<Vec<u8>> =
HashSet::with_capacity_and_hasher(key_cards[0].1 as usize, Default::default());
self.siter(smallest_key, |m| {
current.insert(m.to_vec());
true
})?;
if current.is_empty() {
return Ok(0);
}
for &(k, next_card) in &key_cards[1..] {
if current.is_empty() {
return Ok(0);
}
if (current.len() as u64) * 4 < next_card {
let mut err = None;
current.retain(|m| match self.sismember(k, m) {
Ok(found) => found,
Err(e) => {
err = Some(e);
false
}
});
if let Some(e) = err {
return Err(e);
}
} else {
let mut next_set =
HashSet::with_capacity_and_hasher(current.len(), Default::default());
self.siter(k, |m| {
if current.contains(m) {
next_set.insert(m.to_vec());
}
next_set.len() < current.len()
})?;
current = next_set;
}
}
let total = current.len();
Ok(if limit == 0 { total } else { total.min(limit) })
}
pub fn sscan<K: AsRef<[u8]>>(
&self,
key: K,
cursor: u64,
pattern: Option<&[u8]>,
count: Option<usize>,
) -> Result<SetScanResult> {
let total = self.scard(&key)? as usize;
let start = cursor as usize;
if start >= total || total == 0 {
return Ok((0, Vec::new()));
}
let step = count.unwrap_or(10).max(1);
let end = (start + step).min(total);
let next_cursor = if end >= total { 0 } else { end as u64 };
let is_match_all = match pattern {
Some(p) => p == b"*",
None => true,
};
let mut results = Vec::new();
let mut current_idx = 0usize;
self.siter(key, |member| {
if current_idx >= start && current_idx < end {
let matches = is_match_all
|| match pattern {
Some(pat) => matches_glob_bytes(pat, member),
None => true,
};
if matches {
results.push(member.to_vec());
}
}
current_idx += 1;
current_idx < end
})?;
Ok((next_cursor, results))
}
impl_expire_ops! {
prefix = s,
meta_type = SetMeta
}
}