use core::str::from_utf8;
use crate::{
db::WeDb,
error::{Error, Result},
key::Key,
key_composer::{
KeyComposer,
ns::{
DEFAULT_NAMESPACE, NS_NAME_PREFIX, NS_NEXT_ID_KEY, catalog_db_key, catalog_ns_prefix,
is_default_namespace, ns_id_key, ns_name_key,
},
oppv::{decode_oppv_u64, encode_oppv_u64, encode_oppv_u64_fixed},
tag::ScopeModeTag,
},
};
#[derive(Clone, Copy)]
pub struct Namespace<'a> {
pub db: &'a WeDb,
pub kc: KeyComposer<'a>,
}
impl WeDb {
#[inline]
pub fn kc_for_ns<'a>(&self, name: &'a str) -> KeyComposer<'a> {
if is_default_namespace(name) {
KeyComposer::new("default")
} else {
let ns_id = self.ns_id(name).unwrap_or(0);
KeyComposer::new_named(name, ns_id, 0)
}
}
#[inline]
pub fn kc_for_ns_db<'a>(&self, name: &'a str, db: u64) -> KeyComposer<'a> {
if is_default_namespace(name) {
KeyComposer::new_db(db)
} else {
let ns_id = self.ns_id(name).unwrap_or(0);
KeyComposer::new_named(name, ns_id, db)
}
}
#[inline]
pub fn namespace<'a>(&'a self, name: &'a str) -> Namespace<'a> {
let kc = self.kc_for_ns(name);
Namespace { db: self, kc }
}
#[inline]
pub fn select_db(&self, db_index: u64) -> Result<Namespace<'_>> {
if db_index > 0 {
self.activate_db("default", db_index)?;
}
Ok(Namespace {
db: self,
kc: KeyComposer::new_db(db_index),
})
}
#[inline]
pub fn default_ns(&self) -> Namespace<'_> {
self.namespace("default")
}
#[inline]
pub fn iter(&self) -> Namespaces {
Namespaces {
iter: self.meta_ns.prefix(NS_NAME_PREFIX),
emitted_default: false,
}
}
pub fn activate_db(&self, ns: &str, db_idx: u64) -> Result<()> {
let key = catalog_db_key(ns, db_idx);
self.meta_ns.insert(key, b"")?;
Ok(())
}
pub fn ns_id(&self, name: &str) -> Result<u64> {
if is_default_namespace(name) {
return Ok(0);
}
let pin = self.ns_cache.pin();
if let Some(&id) = pin.get(name) {
return Ok(id);
}
let _guard = self.ns_lock.lock();
let pin = self.ns_cache.pin();
if let Some(&id) = pin.get(name) {
return Ok(id);
}
let key = ns_name_key(name);
if let Some(val) = self.meta_ns.get(&key)?
&& let Some((id, _)) = decode_oppv_u64(&val)
{
pin.insert(hipstr::HipStr::from(name), id);
return Ok(id);
}
let current_id = match self.meta_ns.get(NS_NEXT_ID_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 = [0u8; 9];
let val_len = encode_oppv_u64_fixed(new_id, &mut val_buf);
batch.insert(&self.meta_ns, &key, &val_buf[..val_len]);
let id_key = ns_id_key(new_id);
batch.insert(&self.meta_ns, id_key, name.as_bytes());
let mut next_val_buf = [0u8; 9];
let next_len = encode_oppv_u64_fixed(next_val, &mut next_val_buf);
batch.insert(&self.meta_ns, NS_NEXT_ID_KEY, &next_val_buf[..next_len]);
batch.commit()?;
pin.insert(hipstr::HipStr::from(name), new_id);
Ok(new_id)
}
pub fn query_ns_id(&self, name: &str) -> Result<Option<u64>> {
if is_default_namespace(name) {
return Ok(Some(0));
}
let pin = self.ns_cache.pin();
if let Some(&id) = pin.get(name) {
return Ok(Some(id));
}
let key = ns_name_key(name);
if let Some(val) = self.meta_ns.get(&key)?
&& let Some((id, _)) = decode_oppv_u64(&val)
{
pin.insert(hipstr::HipStr::from(name), id);
return Ok(Some(id));
}
Ok(None)
}
pub fn ns_name(&self, ns_id: u64) -> Result<Option<String>> {
if ns_id == 0 {
return Ok(Some("default".to_string()));
}
let id_key = ns_id_key(ns_id);
if let Some(val) = self.meta_ns.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 _guard = self.ns_lock.lock();
let ns_id = match self.query_ns_id(old_name)? {
Some(id) => id,
None => return Err(Error::invalid_data("namespace not found")),
};
if self.query_ns_id(new_name)?.is_some() {
return Err(Error::invalid_data("target namespace already exists"));
}
let mut batch = self.db.batch();
let old_key = ns_name_key(old_name);
batch.remove(&self.meta_ns, old_key);
let new_key = ns_name_key(new_name);
let mut val_buf = [0u8; 9];
let val_len = encode_oppv_u64_fixed(ns_id, &mut val_buf);
batch.insert(&self.meta_ns, &new_key, &val_buf[..val_len]);
let id_key = ns_id_key(ns_id);
batch.insert(&self.meta_ns, id_key, new_name.as_bytes());
let cat_old_prefix = catalog_ns_prefix(old_name);
for item in self.meta_ns.prefix(&cat_old_prefix) {
let k = item.key()?;
if !k.starts_with(&cat_old_prefix) {
break;
}
batch.remove(&self.meta_ns, &*k);
let remain = &k[cat_old_prefix.len()..];
let mut new_cat_k = catalog_ns_prefix(new_name);
new_cat_k.extend_from_slice(remain);
batch.insert(&self.meta_ns, new_cat_k, b"");
}
batch.commit()?;
let pin = self.ns_cache.pin();
pin.remove(old_name);
pin.insert(hipstr::HipStr::from(new_name), ns_id);
Ok(())
}
#[inline]
pub fn keys<P: AsRef<[u8]>>(&self, pattern: P) -> Result<Vec<Vec<u8>>> {
Key::keys(self, pattern)
}
#[inline]
pub fn key_count(&self) -> Result<usize> {
Key::key_count(self)
}
}
#[inline]
fn clear_ks_prefix(
ks: &fjall::Keyspace,
prefix: &[u8],
batch: &mut fjall::OwnedWriteBatch,
count: &mut u64,
) -> Result<()> {
for item in ks.prefix(prefix) {
let k = item.key()?;
if !k.starts_with(prefix) {
break;
}
batch.remove(ks, k);
*count += 1;
}
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 id(&self) -> Result<u64> {
self.db.ns_id(self.name())
}
#[inline]
pub fn rename(&self, new_name: &str) -> Result<()> {
self.db.rename_namespace(self.name(), new_name)
}
#[inline]
pub fn select_db(&self, db_index: u64) -> Result<Namespace<'a>> {
if db_index > 0 {
self.db.activate_db(self.kc.ns(), db_index)?;
}
Ok(Namespace {
db: self.db,
kc: KeyComposer::new_named(self.kc.ns(), self.kc.ns_id(), db_index),
})
}
#[inline]
pub fn db_index(&self) -> u64 {
self.kc.db()
}
#[inline]
pub fn keys<P: AsRef<[u8]>>(&self, pattern: P) -> Result<Vec<Vec<u8>>> {
Key::keys(self, pattern)
}
#[inline]
pub fn key_count(&self) -> Result<usize> {
Key::key_count(self)
}
#[inline]
pub fn del<K: AsRef<[u8]>>(&self, keys: &[K]) -> Result<usize> {
Key::del(self, keys)
}
#[inline]
pub fn exists<K: AsRef<[u8]>>(&self, keys: &[K]) -> Result<usize> {
Key::exists(self, keys)
}
pub fn clear(&self) -> Result<u64> {
let mut count = 0u64;
let mut batch = self.db.db.batch();
if self.is_default() {
for item in self.db.data.iter() {
batch.remove(&self.db.data, item.key()?);
count += 1;
}
for item in self.db.meta.iter() {
batch.remove(&self.db.meta, item.key()?);
count += 1;
}
let prefix_default_dbs = ScopeModeTag::DefaultDb.as_prefix_slice();
clear_ks_prefix(&self.db.data_ns, prefix_default_dbs, &mut batch, &mut count)?;
clear_ks_prefix(&self.db.meta_ns, prefix_default_dbs, &mut batch, &mut count)?;
let cat_prefix = catalog_ns_prefix(DEFAULT_NAMESPACE);
let mut _dummy = 0;
clear_ks_prefix(&self.db.meta_ns, &cat_prefix, &mut batch, &mut _dummy)?;
} else if self.kc.db() > 0 {
let prefix = self.kc.namespace_prefix();
clear_ks_prefix(&self.db.data_ns, &prefix, &mut batch, &mut count)?;
clear_ks_prefix(&self.db.meta_ns, &prefix, &mut batch, &mut count)?;
let cat_key = catalog_db_key(self.name(), self.kc.db());
batch.remove(&self.db.meta_ns, cat_key);
} else {
let prefix_single = self.kc.namespace_prefix();
clear_ks_prefix(&self.db.data_ns, &prefix_single, &mut batch, &mut count)?;
clear_ks_prefix(&self.db.meta_ns, &prefix_single, &mut batch, &mut count)?;
let mut prefix_multi = Vec::with_capacity(2 + 9);
prefix_multi.extend_from_slice(ScopeModeTag::TenantDb.as_prefix_slice());
encode_oppv_u64(self.kc.ns_id(), &mut prefix_multi);
clear_ks_prefix(&self.db.data_ns, &prefix_multi, &mut batch, &mut count)?;
clear_ks_prefix(&self.db.meta_ns, &prefix_multi, &mut batch, &mut count)?;
let cat_prefix = catalog_ns_prefix(self.name());
let mut _dummy = 0;
clear_ks_prefix(&self.db.meta_ns, &cat_prefix, &mut batch, &mut _dummy)?;
}
batch.commit()?;
Ok(count)
}
pub fn flushdb(&self) -> Result<u64> {
let mut count = 0u64;
let mut batch = self.db.db.batch();
if self.is_default() {
for item in self.db.data.iter() {
batch.remove(&self.db.data, item.key()?);
count += 1;
}
for item in self.db.meta.iter() {
batch.remove(&self.db.meta, item.key()?);
count += 1;
}
} else {
let prefix = self.kc.namespace_prefix();
clear_ks_prefix(&self.db.data_ns, &prefix, &mut batch, &mut count)?;
clear_ks_prefix(&self.db.meta_ns, &prefix, &mut batch, &mut count)?;
if self.kc.db() > 0 {
let cat_key = catalog_db_key(self.name(), self.kc.db());
batch.remove(&self.db.meta_ns, cat_key);
}
}
batch.commit()?;
Ok(count)
}
#[inline]
pub fn flushall(&self) -> Result<u64> {
if self.is_default() || self.kc.db() == 0 {
self.clear()
} else {
Namespace {
db: self.db,
kc: KeyComposer::new_named(self.kc.ns(), self.kc.ns_id(), 0),
}
.clear()
}
}
#[inline]
pub fn iter(&self) -> Dbs {
let cat_prefix = catalog_ns_prefix(self.name());
let iter = self.db.meta_ns.prefix(&cat_prefix);
Dbs {
prefix: cat_prefix,
iter,
emitted_zero: false,
}
}
}
pub struct Namespaces {
iter: fjall::Iter,
emitted_default: bool,
}
impl Iterator for Namespaces {
type Item = String;
fn next(&mut self) -> Option<Self::Item> {
if !self.emitted_default {
self.emitted_default = true;
return Some("default".to_string());
}
let prefix = NS_NAME_PREFIX;
for item in self.iter.by_ref() {
let k = item.key().ok()?;
if !k.starts_with(prefix) {
return None;
}
if let Ok(ns_name) = from_utf8(&k[prefix.len()..]) {
return Some(ns_name.to_string());
}
}
None
}
}
impl IntoIterator for &WeDb {
type Item = String;
type IntoIter = Namespaces;
#[inline]
fn into_iter(self) -> Self::IntoIter {
self.iter()
}
}
impl IntoIterator for WeDb {
type Item = String;
type IntoIter = Namespaces;
#[inline]
fn into_iter(self) -> Self::IntoIter {
self.iter()
}
}
pub struct Dbs {
prefix: Vec<u8>,
iter: fjall::Iter,
emitted_zero: bool,
}
impl Iterator for Dbs {
type Item = u64;
fn next(&mut self) -> Option<Self::Item> {
if !self.emitted_zero {
self.emitted_zero = true;
return Some(0);
}
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 > 0
{
return Some(db_idx);
}
}
None
}
}
impl IntoIterator for &Namespace<'_> {
type Item = u64;
type IntoIter = Dbs;
#[inline]
fn into_iter(self) -> Self::IntoIter {
self.iter()
}
}
impl IntoIterator for Namespace<'_> {
type Item = u64;
type IntoIter = Dbs;
#[inline]
fn into_iter(self) -> Self::IntoIter {
self.iter()
}
}