pub mod traits;
pub use traits::*;
use rapidhash::RapidHashSet as HashSet;
use std::str::from_utf8;
use crate::db::WeDb;
use crate::error::{Error, Result};
use crate::key_composer::oppv::{decode_oppv_u64, encode_oppv_u64};
use crate::key_composer::{KeyComposer, is_default_namespace, matches_glob_bytes, namespace_to_db};
use crate::string::conf::Set;
use crate::zset::RangeScoreSpec;
#[derive(Clone, Copy)]
pub struct Namespace<'a> {
pub db: &'a WeDb,
pub kc: KeyComposer<'a>,
}
impl WeDb {
#[inline]
pub fn namespace<'a>(&'a self, name: &'a str) -> Namespace<'a> {
Namespace {
db: self,
kc: KeyComposer::new(name),
}
}
#[inline]
pub fn select_db(&self, db_index: u64) -> Namespace<'_> {
if db_index > 0 {
let _ = self.activate_db("default", db_index);
}
self.namespace(match db_index {
0 => "default",
1 => "db1",
2 => "db2",
3 => "db3",
4 => "db4",
5 => "db5",
6 => "db6",
7 => "db7",
8 => "db8",
9 => "db9",
10 => "db10",
11 => "db11",
12 => "db12",
13 => "db13",
14 => "db14",
15 => "db15",
_ => "default",
})
}
#[inline]
pub fn default_ns(&self) -> Namespace<'_> {
self.namespace("default")
}
#[inline]
pub fn iter(&self) -> Namespaces<'_> {
Namespaces {
db: self,
cursor: None,
emitted_default: false,
}
}
pub fn activate_db(&self, ns: &str, db_idx: u64) -> Result<()> {
let mut key = Vec::with_capacity(5 + ns.len() + 4 + 9);
key.extend_from_slice(b"\x00cat:");
key.extend_from_slice(ns.as_bytes());
key.extend_from_slice(b":db:");
encode_oppv_u64(db_idx, &mut key); self.meta.insert(key, b"")?;
Ok(())
}
pub fn get_or_create_ns_id(&self, name: &str) -> Result<u64> {
if is_default_namespace(name) {
return Ok(0);
}
let mut key = Vec::with_capacity(9 + name.len());
key.extend_from_slice(b"\x00ns:name:");
key.extend_from_slice(name.as_bytes());
if let Some(val) = self.meta.get(&key)?
&& let Some((id, _)) = decode_oppv_u64(&val)
{
return Ok(id);
}
let next_key = b"\x00ns:next_id";
let current_id = match self.meta.get(next_key)? {
Some(val) => decode_oppv_u64(&val).map(|(v, _)| v).unwrap_or(1),
None => 1,
};
let new_id = current_id;
let next_val = current_id + 1;
let mut batch = self.db.batch();
let mut val_buf = Vec::new();
encode_oppv_u64(new_id, &mut val_buf);
batch.insert(&self.meta, &key, &val_buf);
let mut id_key = Vec::with_capacity(7 + 9);
id_key.extend_from_slice(b"\x00ns:id:");
encode_oppv_u64(new_id, &mut id_key);
batch.insert(&self.meta, id_key, name.as_bytes());
let mut next_val_buf = Vec::new();
encode_oppv_u64(next_val, &mut next_val_buf);
batch.insert(&self.meta, next_key, next_val_buf);
batch.commit()?;
Ok(new_id)
}
pub fn get_ns_id(&self, name: &str) -> Result<Option<u64>> {
if is_default_namespace(name) {
return Ok(Some(0));
}
let mut key = Vec::with_capacity(9 + name.len());
key.extend_from_slice(b"\x00ns:name:");
key.extend_from_slice(name.as_bytes());
if let Some(val) = self.meta.get(&key)? {
return Ok(decode_oppv_u64(&val).map(|(v, _)| v));
}
Ok(None)
}
pub fn get_ns_name(&self, ns_id: u64) -> Result<Option<String>> {
if ns_id == 0 {
return Ok(Some("default".to_string()));
}
let mut id_key = Vec::with_capacity(7 + 9);
id_key.extend_from_slice(b"\x00ns:id:");
encode_oppv_u64(ns_id, &mut id_key);
if let Some(val) = self.meta.get(&id_key)? {
return Ok(String::from_utf8(val.to_vec()).ok());
}
Ok(None)
}
pub fn rename_namespace(&self, old_name: &str, new_name: &str) -> Result<()> {
if is_default_namespace(old_name) || is_default_namespace(new_name) {
return Err(Error::invalid_data("cannot rename default namespace"));
}
let ns_id = match self.get_ns_id(old_name)? {
Some(id) => id,
None => return Err(Error::invalid_data("namespace not found")),
};
let mut batch = self.db.batch();
let mut old_key = Vec::with_capacity(9 + old_name.len());
old_key.extend_from_slice(b"\x00ns:name:");
old_key.extend_from_slice(old_name.as_bytes());
batch.remove(&self.meta, old_key);
let mut new_key = Vec::with_capacity(9 + new_name.len());
new_key.extend_from_slice(b"\x00ns:name:");
new_key.extend_from_slice(new_name.as_bytes());
let mut val_buf = Vec::new();
encode_oppv_u64(ns_id, &mut val_buf);
batch.insert(&self.meta, new_key, val_buf);
let mut id_key = Vec::with_capacity(7 + 9);
id_key.extend_from_slice(b"\x00ns:id:");
encode_oppv_u64(ns_id, &mut id_key);
batch.insert(&self.meta, id_key, new_name.as_bytes());
batch.commit()?;
Ok(())
}
}
impl<'a> Namespace<'a> {
#[inline]
pub fn name(&self) -> &str {
self.kc.ns()
}
#[inline]
pub fn is_default(&self) -> bool {
self.kc.is_default()
}
#[inline]
pub fn ns_id(&self) -> Result<u64> {
self.db.get_or_create_ns_id(self.name())
}
#[inline]
pub fn db_index(&self) -> u64 {
namespace_to_db(self.name())
}
pub fn clear(&self) -> Result<u64> {
let prefix = self.kc.namespace_prefix();
if prefix.is_empty() {
return Ok(0);
}
let mut count = 0u64;
let mut batch = self.db.db.batch();
for item in self.db.data.prefix(&prefix) {
batch.remove(&self.db.data, item.key()?);
count += 1;
}
for item in self.db.meta.prefix(&prefix) {
batch.remove(&self.db.meta, item.key()?);
count += 1;
}
let mut cat_prefix = Vec::with_capacity(5 + self.name().len() + 1);
cat_prefix.extend_from_slice(b"\x00cat:");
cat_prefix.extend_from_slice(self.name().as_bytes());
cat_prefix.push(b':');
for item in self.db.meta.prefix(&cat_prefix) {
batch.remove(&self.db.meta, item.key()?);
}
batch.commit()?;
Ok(count)
}
pub fn key_count(&self) -> Result<usize> {
let keys = self.keys("*")?;
Ok(keys.len())
}
pub fn keys(&self, pattern: &str) -> Result<Vec<Vec<u8>>> {
let pat_bytes = pattern.as_bytes();
let mut seen = HashSet::default();
let prefix = self.kc.namespace_prefix();
let mut scan_ks = |ks: &fjall::Keyspace| -> Result<()> {
if prefix.is_empty() {
for item in ks.iter() {
let k = item.key()?;
if let Some(user_k) = self.kc.extract_user_key(&k)
&& matches_glob_bytes(pat_bytes, user_k)
&& !seen.contains(user_k)
{
seen.insert(user_k.to_vec());
}
}
} else {
for item in ks.prefix(&prefix) {
let k = item.key()?;
if !k.starts_with(&prefix) {
break;
}
if let Some(user_k) = self.kc.extract_user_key(&k)
&& matches_glob_bytes(pat_bytes, user_k)
&& !seen.contains(user_k)
{
seen.insert(user_k.to_vec());
}
}
}
Ok(())
};
scan_ks(&self.db.data)?;
scan_ks(&self.db.meta)?;
let mut result: Vec<Vec<u8>> = seen.into_iter().collect();
result.sort();
Ok(result)
}
#[inline]
pub fn del<K: AsRef<[u8]>>(&self, keys: &[K]) -> Result<usize> {
self.db.del_with_kc(&self.kc, keys)
}
#[inline]
pub fn exists<K: AsRef<[u8]>>(&self, keys: &[K]) -> Result<usize> {
self.db.exists_with_kc(&self.kc, keys)
}
#[inline]
pub fn iter(&self) -> Dbs {
let mut cat_prefix = Vec::with_capacity(5 + self.name().len() + 4);
cat_prefix.extend_from_slice(b"\x00cat:");
cat_prefix.extend_from_slice(self.name().as_bytes());
cat_prefix.extend_from_slice(b":db:");
let iter = self.db.meta.prefix(&cat_prefix);
Dbs {
prefix: cat_prefix,
iter,
emitted_self: false,
self_db: self.db_index(),
}
}
}
impl<'a> KvOps for Namespace<'a> {
#[inline]
fn get(&self, key: impl AsRef<[u8]>) -> Result<Option<Vec<u8>>> {
let raw_key = self.kc.raw_key_bytes(key.as_ref());
self.db.get(&raw_key)
}
#[inline]
fn set<'b>(
&self,
key: impl AsRef<[u8]>,
val: impl AsRef<[u8]>,
opts: impl AsRef<[Set<'b>]>,
) -> Result<()> {
let raw_key = self.kc.raw_key_bytes(key.as_ref());
self.db.set(&raw_key, val, opts)?;
Ok(())
}
#[inline]
fn set_ex(&self, key: impl AsRef<[u8]>, val: impl AsRef<[u8]>, expire_ms: u64) -> Result<()> {
let raw_key = self.kc.raw_key_bytes(key.as_ref());
self.db.setex(&raw_key, val, expire_ms)?;
Ok(())
}
#[inline]
fn incr(&self, key: impl AsRef<[u8]>) -> Result<i64> {
let raw_key = self.kc.raw_key_bytes(key.as_ref());
self.db.incr(&raw_key)
}
#[inline]
fn incrby(&self, key: impl AsRef<[u8]>, step: i64) -> Result<i64> {
let raw_key = self.kc.raw_key_bytes(key.as_ref());
self.db.incrby(&raw_key, step)
}
#[inline]
fn decr(&self, key: impl AsRef<[u8]>) -> Result<i64> {
let raw_key = self.kc.raw_key_bytes(key.as_ref());
self.db.decr(&raw_key)
}
#[inline]
fn decrby(&self, key: impl AsRef<[u8]>, step: i64) -> Result<i64> {
let raw_key = self.kc.raw_key_bytes(key.as_ref());
self.db.decrby(&raw_key, step)
}
#[inline]
fn strlen(&self, key: impl AsRef<[u8]>) -> Result<usize> {
let raw_key = self.kc.raw_key_bytes(key.as_ref());
self.db.strlen(&raw_key)
}
#[inline]
fn getdel(&self, key: impl AsRef<[u8]>) -> Result<Option<Vec<u8>>> {
let raw_key = self.kc.raw_key_bytes(key.as_ref());
self.db.getdel(&raw_key)
}
#[inline]
fn getset(&self, key: impl AsRef<[u8]>, val: impl AsRef<[u8]>) -> Result<Option<Vec<u8>>> {
let raw_key = self.kc.raw_key_bytes(key.as_ref());
self.db.getset(&raw_key, val)
}
#[inline]
fn mget<K: AsRef<[u8]>>(&self, keys: &[K]) -> Result<Vec<Option<Vec<u8>>>> {
let raw_keys: Vec<_> = keys
.iter()
.map(|k| self.kc.raw_key_bytes(k.as_ref()))
.collect();
self.db.mget(&raw_keys)
}
#[inline]
fn mset<K: AsRef<[u8]>, V: AsRef<[u8]>>(&self, kvs: &[(K, V)]) -> Result<()> {
let raw_kvs: Vec<_> = kvs
.iter()
.map(|(k, v)| (self.kc.raw_key_bytes(k.as_ref()), v.as_ref()))
.collect();
self.db.mset(&raw_kvs)
}
}
impl<'a> HashOps for Namespace<'a> {
#[inline]
fn hget(&self, key: impl AsRef<[u8]>, field: impl AsRef<[u8]>) -> Result<Option<Vec<u8>>> {
self.db.hget_with_kc(&self.kc, key, field)
}
#[inline]
fn hset<K: AsRef<[u8]>, F: AsRef<[u8]>, V: AsRef<[u8]>>(
&self,
key: K,
fields: &[(F, V)],
) -> Result<usize> {
self.db.hset_with_kc(&self.kc, key, fields)
}
#[inline]
fn hdel<K: AsRef<[u8]>, F: AsRef<[u8]>>(&self, key: K, fields: &[F]) -> Result<usize> {
self.db.hdel_with_kc(&self.kc, key, fields)
}
#[inline]
fn hexists(&self, key: impl AsRef<[u8]>, field: impl AsRef<[u8]>) -> Result<bool> {
self.db.hexists_with_kc(&self.kc, key, field)
}
#[inline]
fn hlen(&self, key: impl AsRef<[u8]>) -> Result<usize> {
self.db.hlen_with_kc(&self.kc, key)
}
#[inline]
fn hgetall(&self, key: impl AsRef<[u8]>) -> Result<Vec<(Vec<u8>, Vec<u8>)>> {
self.db.hgetall_with_kc(&self.kc, key)
}
#[inline]
fn hkeys(&self, key: impl AsRef<[u8]>) -> Result<Vec<Vec<u8>>> {
self.db.hkeys_with_kc(&self.kc, key)
}
#[inline]
fn hvals(&self, key: impl AsRef<[u8]>) -> Result<Vec<Vec<u8>>> {
self.db.hvals_with_kc(&self.kc, key)
}
#[inline]
fn hincrby(&self, key: impl AsRef<[u8]>, field: impl AsRef<[u8]>, step: i64) -> Result<i64> {
self.db.hincrby_with_kc(&self.kc, key, field, step)
}
#[inline]
fn hincrbyfloat(
&self,
key: impl AsRef<[u8]>,
field: impl AsRef<[u8]>,
step: f64,
) -> Result<f64> {
self.db.hincrbyfloat_with_kc(&self.kc, key, field, step)
}
#[inline]
fn hmget<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
key: K,
fields: &[F],
) -> Result<Vec<Option<Vec<u8>>>> {
self.db.hmget_with_kc(&self.kc, key, fields)
}
#[inline]
fn hmset<K: AsRef<[u8]>, F: AsRef<[u8]>, V: AsRef<[u8]>>(
&self,
key: K,
fields: &[(F, V)],
) -> Result<()> {
self.db.hmset_with_kc(&self.kc, key, fields)
}
#[inline]
fn hstrlen(&self, key: impl AsRef<[u8]>, field: impl AsRef<[u8]>) -> Result<usize> {
self.db.hstrlen_with_kc(&self.kc, key, field)
}
}
impl<'a> ListOps for Namespace<'a> {
#[inline]
fn lpush<K: AsRef<[u8]>, V: AsRef<[u8]>>(&self, key: K, values: &[V]) -> Result<u64> {
self.db.lpush_with_kc(&self.kc, key, values)
}
#[inline]
fn rpush<K: AsRef<[u8]>, V: AsRef<[u8]>>(&self, key: K, values: &[V]) -> Result<u64> {
self.db.rpush_with_kc(&self.kc, key, values)
}
#[inline]
fn lpop(&self, key: impl AsRef<[u8]>, count: usize) -> Result<Vec<Vec<u8>>> {
self.db.lpop_with_kc(&self.kc, key, count)
}
#[inline]
fn rpop(&self, key: impl AsRef<[u8]>, count: usize) -> Result<Vec<Vec<u8>>> {
self.db.rpop_with_kc(&self.kc, key, count)
}
#[inline]
fn llen(&self, key: impl AsRef<[u8]>) -> Result<u64> {
self.db.llen_with_kc(&self.kc, key)
}
#[inline]
fn lindex(&self, key: impl AsRef<[u8]>, index: i64) -> Result<Option<Vec<u8>>> {
self.db.lindex_with_kc(&self.kc, key, index)
}
#[inline]
fn lrange(&self, key: impl AsRef<[u8]>, start: i64, stop: i64) -> Result<Vec<Vec<u8>>> {
self.db.lrange_with_kc(&self.kc, key, start, stop)
}
#[inline]
fn ltrim(&self, key: impl AsRef<[u8]>, start: i64, stop: i64) -> Result<()> {
self.db.ltrim_with_kc(&self.kc, key, start, stop)
}
#[inline]
fn lset(&self, key: impl AsRef<[u8]>, index: i64, value: impl AsRef<[u8]>) -> Result<()> {
self.db.lset_with_kc(&self.kc, key, index, value)
}
}
impl<'a> SetOps for Namespace<'a> {
#[inline]
fn sadd<K: AsRef<[u8]>, M: AsRef<[u8]>>(&self, key: K, members: &[M]) -> Result<usize> {
self.db.sadd_with_kc(&self.kc, key, members)
}
#[inline]
fn srem<K: AsRef<[u8]>, M: AsRef<[u8]>>(&self, key: K, members: &[M]) -> Result<usize> {
self.db.srem_with_kc(&self.kc, key, members)
}
#[inline]
fn smembers(&self, key: impl AsRef<[u8]>) -> Result<Vec<Vec<u8>>> {
self.db.smembers_with_kc(&self.kc, key)
}
#[inline]
fn sismember(&self, key: impl AsRef<[u8]>, member: impl AsRef<[u8]>) -> Result<bool> {
self.db.sismember_with_kc(&self.kc, key, member)
}
#[inline]
fn scard(&self, key: impl AsRef<[u8]>) -> Result<u64> {
self.db.scard_with_kc(&self.kc, key)
}
#[inline]
fn spop(&self, key: impl AsRef<[u8]>, count: usize) -> Result<Vec<Vec<u8>>> {
self.db.spop_with_kc(&self.kc, key, count)
}
}
impl<'a> ZSetOps for Namespace<'a> {
#[inline]
fn zadd<K: AsRef<[u8]>, M: AsRef<[u8]>>(&self, key: K, members: &[(f64, M)]) -> Result<usize> {
self.db.zadd_with_kc(&self.kc, key, members, [])
}
#[inline]
fn zrem<K: AsRef<[u8]>, M: AsRef<[u8]>>(&self, key: K, members: &[M]) -> Result<usize> {
self.db.zrem_with_kc(&self.kc, key, members)
}
#[inline]
fn zscore(&self, key: impl AsRef<[u8]>, member: impl AsRef<[u8]>) -> Result<Option<f64>> {
self.db.zscore_with_kc(&self.kc, key, member)
}
#[inline]
fn zcard(&self, key: impl AsRef<[u8]>) -> Result<u64> {
self.db.zcard_with_kc(&self.kc, key)
}
#[inline]
fn zrank(&self, key: impl AsRef<[u8]>, member: impl AsRef<[u8]>) -> Result<Option<u64>> {
self.db.zrank_with_kc(&self.kc, key, member)
}
#[inline]
fn zrevrank(&self, key: impl AsRef<[u8]>, member: impl AsRef<[u8]>) -> Result<Option<u64>> {
self.db.zrevrank_with_kc(&self.kc, key, member)
}
#[inline]
fn zrange(&self, key: impl AsRef<[u8]>, start: i64, stop: i64) -> Result<Vec<(Vec<u8>, f64)>> {
self.db.zrange_with_kc(&self.kc, key, start, stop)
}
#[inline]
fn zcount(&self, key: impl AsRef<[u8]>, spec: &RangeScoreSpec) -> Result<u64> {
self.db.zcount_with_kc(&self.kc, key, spec)
}
#[inline]
fn zincrby(&self, key: impl AsRef<[u8]>, step: f64, member: impl AsRef<[u8]>) -> Result<f64> {
self.db.zincrby_with_kc(&self.kc, key, step, member)
}
}
#[inline]
fn find_next_ns_in_ks(ks: &fjall::Keyspace, cursor: &[u8]) -> Option<String> {
for item in ks.range(cursor..) {
let k = item.key().ok()?;
if !k.starts_with(b"\x00ns:") {
return None;
}
let remain = &k[4..];
if let Some(colon_pos) = memchr::memchr(b':', remain)
&& let Ok(ns_str) = from_utf8(&remain[..colon_pos])
&& !ns_str.is_empty()
{
return Some(ns_str.to_string());
}
}
None
}
pub struct Namespaces<'a> {
db: &'a WeDb,
cursor: Option<Vec<u8>>,
emitted_default: bool,
}
impl<'a> Iterator for Namespaces<'a> {
type Item = String;
fn next(&mut self) -> Option<Self::Item> {
if !self.emitted_default {
self.emitted_default = true;
return Some("default".to_string());
}
let cur = self.cursor.as_deref().unwrap_or(b"\x00ns:");
let ns_data = find_next_ns_in_ks(&self.db.data, cur);
let ns_meta = find_next_ns_in_ks(&self.db.meta, cur);
let next_ns = match (ns_data, ns_meta) {
(Some(a), Some(b)) => Some(a.min(b)),
(Some(a), None) => Some(a),
(None, Some(b)) => Some(b),
(None, None) => None,
}?;
let mut next_cursor = Vec::with_capacity(4 + next_ns.len() + 2);
next_cursor.extend_from_slice(b"\x00ns:");
next_cursor.extend_from_slice(next_ns.as_bytes());
next_cursor.push(b':');
next_cursor.push(0xff);
self.cursor = Some(next_cursor);
Some(next_ns)
}
}
impl<'a> IntoIterator for &'a WeDb {
type Item = String;
type IntoIter = Namespaces<'a>;
#[inline]
fn into_iter(self) -> Self::IntoIter {
self.iter()
}
}
pub struct Dbs {
prefix: Vec<u8>,
iter: fjall::Iter,
emitted_self: bool,
self_db: u64,
}
impl Iterator for Dbs {
type Item = u64;
fn next(&mut self) -> Option<Self::Item> {
if !self.emitted_self {
self.emitted_self = true;
return Some(self.self_db);
}
for item in self.iter.by_ref() {
let k = item.key().ok()?;
if !k.starts_with(&self.prefix) {
return None;
}
let remain = &k[self.prefix.len()..];
if let Some((db_idx, _consumed)) = decode_oppv_u64(remain)
&& db_idx != self.self_db
{
return Some(db_idx);
}
}
None
}
}
impl<'a> IntoIterator for &'a Namespace<'a> {
type Item = u64;
type IntoIter = Dbs;
#[inline]
fn into_iter(self) -> Self::IntoIter {
self.iter()
}
}
impl<'a> IntoIterator for Namespace<'a> {
type Item = u64;
type IntoIter = Dbs;
#[inline]
fn into_iter(self) -> Self::IntoIter {
self.iter()
}
}