use std::fmt;
use std::fs::create_dir_all;
use std::path::Path;
use std::str;
use std::sync::Arc;
use crate::conf::{Conf, parse_compression_type};
use crate::error::{ERR_WRONG_TYPE, Error, Result};
use crate::key_composer::{ALL_COMPOSITE_META_TAGS, KeyComposer, KeyTag};
use crate::keyspace::{self, DATA, META};
use crate::meta::{KeyMeta, MetaOps, current_now_ms, init_version_counter};
use crate::string::conf::Set;
use crate::string::{
ERR_STRING_EXCEEDS_MAX_SIZE, MAX_STRING_SIZE, decode_string_value, encode_string_value,
is_string_expired,
};
use fjall::{CompressionType, Database, Keyspace, KeyspaceCreateOptions, PersistMode};
#[derive(Clone)]
pub struct WeDb {
pub db: Arc<Database>,
pub data: Keyspace,
pub meta: Keyspace,
}
impl fmt::Debug for WeDb {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("WeDb").finish()
}
}
impl WeDb {
pub const DEFAULT_KEYSPACE: &'static str = DATA;
pub const META_KEYSPACE: &'static str = META;
pub fn open(path: impl AsRef<Path>) -> Result<Self> {
Self::open_with_conf(&Conf {
data_path: path.as_ref().to_string_lossy().to_string(),
..Default::default()
})
}
pub fn open_with_conf(conf: &Conf) -> Result<Self> {
init_version_counter();
let path = Path::new(&conf.data_path);
if let Some(parent) = path.parent() {
create_dir_all(parent)?;
}
let mut builder = Database::builder(path);
let journal_comp =
parse_compression_type(conf.journal_compression.as_deref(), CompressionType::None);
if let Some(cache_size) = conf.cache_size {
builder = builder.cache_size(cache_size as u64);
}
builder = builder.journal_compression(journal_comp);
if let Some(manual_persist) = conf.manual_journal_persist {
builder = builder.manual_journal_persist(manual_persist);
}
if let Some(worker_threads) = conf.worker_threads
&& worker_threads > 0
{
builder = builder.worker_threads(worker_threads);
}
if let Some(max_journaling_size) = conf.max_journaling_size {
builder = builder.max_journaling_size(max_journaling_size as u64);
}
if let Some(max_cached_files) = conf.max_cached_files {
builder = builder.max_cached_files(Some(max_cached_files));
}
let db = Arc::new(builder.open().map_err(|e| {
Error::internal_with_source(format!("Failed to open Fjall at {path:?}"), e)
})?);
let keyspace = keyspace::Keyspace::open(&db, conf)?;
Ok(Self {
db,
data: keyspace.data,
meta: keyspace.meta,
})
}
#[inline]
pub fn database(&self) -> &Arc<Database> {
&self.db
}
#[inline]
pub fn data(&self) -> &Keyspace {
&self.data
}
#[inline]
pub fn meta(&self) -> &Keyspace {
&self.meta
}
#[inline]
pub fn keyspace(&self, name: &str) -> Result<Keyspace> {
self.db
.keyspace(name, KeyspaceCreateOptions::default)
.map_err(|e| Error::internal_with_source(format!("Keyspace '{name}' error"), e))
}
#[inline]
pub fn persist(&self, mode: PersistMode) -> Result<()> {
self.db
.persist(mode)
.map_err(|e| Error::internal_with_source("Persist error", e))
}
#[inline]
pub fn get(&self, key: impl AsRef<[u8]>) -> Result<Option<Vec<u8>>> {
let key_bytes = key.as_ref();
let kc = KeyComposer::new("default");
let raw_k = kc.raw_key_bytes(key_bytes);
let now_ms = current_now_ms();
if let Some(raw) = self.data.get(&raw_k)? {
let (expire_at, payload) = decode_string_value(&raw);
if !is_string_expired(expire_at, now_ms) {
return Ok(Some(payload.to_vec()));
}
}
if self.meta.is_empty()? {
return Ok(None);
}
let mut buf = Vec::with_capacity(32 + key_bytes.len());
for &tag in ALL_COMPOSITE_META_TAGS {
kc.compose_meta_key_into(tag, key_bytes, &mut buf);
if let Some(m_bytes) = self.meta.get(&buf)?
&& let Some(base_meta) = KeyMeta::decode(&m_bytes)
&& !base_meta.is_expired(now_ms)
{
return Err(Error::wrong_type(ERR_WRONG_TYPE));
}
}
Ok(None)
}
#[inline]
pub fn get_with_expire(&self, key: impl AsRef<[u8]>) -> Result<(Option<Vec<u8>>, u64)> {
let key_bytes = key.as_ref();
let kc = KeyComposer::new("default");
let raw_k = kc.raw_key_bytes(key_bytes);
let now_ms = current_now_ms();
if let Some(raw) = self.data.get(&raw_k)? {
let (expire_at, payload) = decode_string_value(&raw);
if !is_string_expired(expire_at, now_ms) {
return Ok((Some(payload.to_vec()), expire_at));
}
}
if self.meta.is_empty()? {
return Ok((None, 0));
}
let mut buf = Vec::with_capacity(32 + key_bytes.len());
for &tag in ALL_COMPOSITE_META_TAGS {
kc.compose_meta_key_into(tag, key_bytes, &mut buf);
if let Some(m_bytes) = self.meta.get(&buf)?
&& let Some(base_meta) = KeyMeta::decode(&m_bytes)
&& !base_meta.is_expired(now_ms)
{
return Err(Error::wrong_type(ERR_WRONG_TYPE));
}
}
Ok((None, 0))
}
pub fn set<'a>(
&self,
key: impl AsRef<[u8]>,
val: impl AsRef<[u8]>,
conf_li: impl AsRef<[Set<'a>]>,
) -> Result<Option<Vec<u8>>> {
let confs = conf_li.as_ref();
if confs.is_empty() {
let key_bytes = key.as_ref();
let val_bytes = val.as_ref();
if val_bytes.len() > MAX_STRING_SIZE {
return Err(Error::invalid_data(ERR_STRING_EXCEEDS_MAX_SIZE));
}
let kc = KeyComposer::new("default");
let raw_k = kc.raw_key_bytes(key_bytes);
let enc_val = encode_string_value(val_bytes, 0);
self.data.insert(&*raw_k, enc_val)?;
return Ok(Some(Vec::new()));
}
let now_ms = current_now_ms();
let args = Set::parse_options(confs, now_ms);
self.set_args(key, val, &args)
}
#[inline]
pub fn del(&self, keys: &[impl AsRef<[u8]>]) -> Result<usize> {
self.del_with_kc(&KeyComposer::new("default"), keys)
}
pub fn del_with_kc(&self, kc: &KeyComposer<'_>, keys: &[impl AsRef<[u8]>]) -> Result<usize> {
let mut deleted = 0;
let now_ms = current_now_ms();
let mut batch = self.db.batch();
let mut buf = Vec::new();
let meta_empty = self.meta.is_empty()?;
for k in keys {
let k_bytes = k.as_ref();
let mut key_deleted = false;
let raw_k = kc.raw_key_bytes(k_bytes);
if let Some(raw) = self.data.get(&raw_k)? {
let (expire_at, _) = decode_string_value(&raw);
if !is_string_expired(expire_at, now_ms) {
key_deleted = true;
}
batch.remove(&self.data, raw_k.as_ref());
}
if !meta_empty {
for &meta_tag in ALL_COMPOSITE_META_TAGS {
kc.compose_meta_key_into(meta_tag, k_bytes, &mut buf);
if let Some(m_bytes) = self.meta.get(&buf)? {
if let Some(base_meta) = KeyMeta::decode(&m_bytes)
&& !base_meta.is_expired(now_ms)
{
key_deleted = true;
}
batch.remove(&self.meta, buf.as_slice());
self.cleanup_composite_data(kc, meta_tag, k_bytes, &mut batch, &mut buf)?;
}
}
}
if key_deleted {
deleted += 1;
}
}
batch.commit()?;
Ok(deleted)
}
#[inline]
pub fn clear_prefix_in_batch(
&self,
prefix: &[u8],
batch: &mut fjall::OwnedWriteBatch,
) -> Result<()> {
for item in self.data.prefix(prefix) {
batch.remove(&self.data, item.key()?);
}
Ok(())
}
#[inline]
pub fn cleanup_composite_data(
&self,
kc: &KeyComposer<'_>,
meta_tag: &[u8],
k_bytes: &[u8],
batch: &mut fjall::OwnedWriteBatch,
buf: &mut Vec<u8>,
) -> Result<()> {
if meta_tag.is_empty() {
return Ok(());
}
if let Some(tag) = KeyTag::from_u8(meta_tag[0]) {
match tag {
KeyTag::HashMeta => {
kc.compose_prefix_into(KeyTag::HashData.as_slice(), k_bytes, buf);
self.clear_prefix_in_batch(buf, batch)?;
}
KeyTag::ListMeta => {
kc.compose_prefix_into(KeyTag::ListData.as_slice(), k_bytes, buf);
self.clear_prefix_in_batch(buf, batch)?;
}
KeyTag::SetMeta => {
kc.compose_prefix_into(KeyTag::SetData.as_slice(), k_bytes, buf);
self.clear_prefix_in_batch(buf, batch)?;
}
KeyTag::ZSetMeta => {
kc.compose_prefix_into(KeyTag::ZSetData.as_slice(), k_bytes, buf);
self.clear_prefix_in_batch(buf, batch)?;
kc.compose_prefix_into(KeyTag::ZSetScore.as_slice(), k_bytes, buf);
self.clear_prefix_in_batch(buf, batch)?;
}
KeyTag::BitmapMeta => {
kc.compose_prefix_into(KeyTag::BitmapData.as_slice(), k_bytes, buf);
self.clear_prefix_in_batch(buf, batch)?;
}
KeyTag::BloomMeta => {
kc.compose_prefix_into(KeyTag::BloomData.as_slice(), k_bytes, buf);
self.clear_prefix_in_batch(buf, batch)?;
}
KeyTag::CuckooMeta => {
kc.compose_prefix_into(KeyTag::CuckooData.as_slice(), k_bytes, buf);
self.clear_prefix_in_batch(buf, batch)?;
}
KeyTag::SortedIntMeta => {
kc.compose_prefix_into(KeyTag::SortedIntData.as_slice(), k_bytes, buf);
self.clear_prefix_in_batch(buf, batch)?;
}
KeyTag::TimeSeriesMeta => {
kc.compose_prefix_into(KeyTag::TimeSeriesData.as_slice(), k_bytes, buf);
self.clear_prefix_in_batch(buf, batch)?;
}
KeyTag::StreamMeta => {
for prefix_tag in [
KeyTag::StreamData.as_slice(),
KeyTag::StreamGroup.as_slice(),
KeyTag::StreamConsumer.as_slice(),
KeyTag::StreamPel.as_slice(),
] {
kc.compose_prefix_into(prefix_tag, k_bytes, buf);
self.clear_prefix_in_batch(buf, batch)?;
}
}
KeyTag::HllMeta => {
kc.compose_meta_key_into(KeyTag::HllRaw.as_slice(), k_bytes, buf);
batch.remove(&self.data, buf.as_slice());
}
_ => {}
}
}
Ok(())
}
#[inline]
pub fn cleanup_all_composite_data(
&self,
kc: &KeyComposer<'_>,
k_bytes: &[u8],
batch: &mut fjall::OwnedWriteBatch,
) -> Result<()> {
if self.meta.is_empty()? {
return Ok(());
}
let mut buf = Vec::with_capacity(32 + k_bytes.len());
for &meta_tag in ALL_COMPOSITE_META_TAGS {
kc.compose_meta_key_into(meta_tag, k_bytes, &mut buf);
if self.meta.contains_key(&buf)? {
batch.remove(&self.meta, buf.as_slice());
self.cleanup_composite_data(kc, meta_tag, k_bytes, batch, &mut buf)?;
}
}
Ok(())
}
#[inline]
pub fn exists(&self, keys: &[impl AsRef<[u8]>]) -> Result<usize> {
self.exists_with_kc(&KeyComposer::new("default"), keys)
}
pub fn exists_with_kc(&self, kc: &KeyComposer<'_>, keys: &[impl AsRef<[u8]>]) -> Result<usize> {
let mut count = 0;
let now_ms = current_now_ms();
let mut buf = Vec::new();
let meta_empty = self.meta.is_empty()?;
for k in keys {
let k_bytes = k.as_ref();
let raw_k = kc.raw_key_bytes(k_bytes);
if let Some(raw) = self.data.get(&raw_k)? {
let (expire_at, _) = decode_string_value(&raw);
if !is_string_expired(expire_at, now_ms) {
count += 1;
continue;
}
}
if meta_empty {
continue;
}
let mut found = false;
for &meta_tag in ALL_COMPOSITE_META_TAGS {
kc.compose_meta_key_into(meta_tag, k_bytes, &mut buf);
if let Some(raw_meta) = self.meta.get(&buf)?
&& let Some(meta) = KeyMeta::decode(&raw_meta)
&& !meta.is_expired(now_ms)
{
found = true;
break;
}
}
if found {
count += 1;
}
}
Ok(count)
}
#[inline]
pub fn check_key_not_other_type(
&self,
kc: &KeyComposer<'_>,
k_bytes: &[u8],
current_meta_tag: &[u8],
now_ms: u64,
) -> Result<()> {
let raw_k = kc.raw_key_bytes(k_bytes);
if let Some(raw) = self.data.get(&raw_k)? {
let (expire_at, _) = decode_string_value(&raw);
if !is_string_expired(expire_at, now_ms) {
return Err(Error::wrong_type(ERR_WRONG_TYPE));
}
}
if self.meta.is_empty()? {
return Ok(());
}
let mut buf = Vec::with_capacity(32 + k_bytes.len());
for &tag in ALL_COMPOSITE_META_TAGS {
if tag == current_meta_tag {
continue;
}
kc.compose_meta_key_into(tag, k_bytes, &mut buf);
if let Some(m_bytes) = self.meta.get(&buf)?
&& let Some(base_meta) = KeyMeta::decode(&m_bytes)
&& !base_meta.is_expired(now_ms)
{
return Err(Error::wrong_type(ERR_WRONG_TYPE));
}
}
Ok(())
}
#[inline]
pub fn get_meta_checked<M: MetaOps>(
&self,
kc: &KeyComposer<'_>,
k_bytes: &[u8],
meta_k: &[u8],
now_ms: u64,
) -> Result<Option<M>> {
let raw_k = kc.raw_key_bytes(k_bytes);
if let Some(raw) = self.data.get(&raw_k)? {
let (expire_at, _) = decode_string_value(&raw);
if !is_string_expired(expire_at, now_ms) {
return Err(Error::wrong_type(ERR_WRONG_TYPE));
}
}
if let Some(m_bytes) = self.meta.get(meta_k)?
&& let Some(meta) = M::decode(&m_bytes)
&& !meta.is_expired(now_ms)
{
return Ok(Some(meta));
}
if !self.meta.is_empty()? {
let mut buf = Vec::with_capacity(32 + k_bytes.len());
for &tag in ALL_COMPOSITE_META_TAGS {
if tag == M::TAG {
continue;
}
kc.compose_meta_key_into(tag, k_bytes, &mut buf);
if let Some(m_bytes) = self.meta.get(&buf)?
&& let Some(base_meta) = KeyMeta::decode(&m_bytes)
&& !base_meta.is_expired(now_ms)
{
return Err(Error::wrong_type(ERR_WRONG_TYPE));
}
}
}
Ok(None)
}
pub fn expireat_generic<M: MetaOps>(&self, key: &[u8], expire_at_ms: u64) -> Result<bool> {
let kc = KeyComposer::new("default");
let mk = kc.compose_meta_key_stack(M::TAG, key);
let now = current_now_ms();
let mut meta: M = match self.get_meta_checked(&kc, key, &mk, now)? {
Some(m) => m,
None => return Ok(false),
};
meta.base_mut().expire_at = expire_at_ms;
let mut batch = self.db.batch();
batch.insert(&self.meta, &*mk, meta.encode_bytes().as_ref());
batch.commit()?;
Ok(true)
}
pub fn ttl_generic<M: MetaOps>(&self, key: &[u8]) -> Result<i64> {
let kc = KeyComposer::new("default");
let mk = kc.compose_meta_key_stack(M::TAG, key);
let now = current_now_ms();
match self.get_meta_checked::<M>(&kc, key, &mk, now)? {
Some(m) => Ok(m.base().ttl_sec(now)),
None => Ok(-2),
}
}
pub fn pttl_generic<M: MetaOps>(&self, key: &[u8]) -> Result<i64> {
let kc = KeyComposer::new("default");
let mk = kc.compose_meta_key_stack(M::TAG, key);
let now = current_now_ms();
match self.get_meta_checked::<M>(&kc, key, &mk, now)? {
Some(m) => Ok(m.base().ttl_ms(now)),
None => Ok(-2),
}
}
pub fn expiretime_generic<M: MetaOps>(&self, key: &[u8]) -> Result<i64> {
let kc = KeyComposer::new("default");
let mk = kc.compose_meta_key_stack(M::TAG, key);
let now = current_now_ms();
match self.get_meta_checked::<M>(&kc, key, &mk, now)? {
Some(m) if m.base().expire_at > 0 => {
Ok(KeyMeta::expire_at_ms_to_sec(m.base().expire_at) as i64)
}
Some(_) => Ok(-1),
None => Ok(-2),
}
}
pub fn pexpiretime_generic<M: MetaOps>(&self, key: &[u8]) -> Result<i64> {
let kc = KeyComposer::new("default");
let mk = kc.compose_meta_key_stack(M::TAG, key);
let now = current_now_ms();
match self.get_meta_checked::<M>(&kc, key, &mk, now)? {
Some(m) if m.base().expire_at > 0 => Ok(m.base().expire_at as i64),
Some(_) => Ok(-1),
None => Ok(-2),
}
}
pub fn persist_generic<M: MetaOps>(&self, key: &[u8]) -> Result<bool> {
let kc = KeyComposer::new("default");
let mk = kc.compose_meta_key_stack(M::TAG, key);
let now = current_now_ms();
let mut meta: M = match self.get_meta_checked::<M>(&kc, key, &mk, now)? {
Some(m) if m.base().expire_at > 0 => m,
_ => return Ok(false),
};
meta.base_mut().expire_at = 0;
let mut batch = self.db.batch();
batch.insert(&self.meta, &*mk, meta.encode_bytes().as_ref());
batch.commit()?;
Ok(true)
}
pub fn flushall(&self) -> Result<()> {
let mut batch = self.db.batch();
for item in self.data.iter() {
let k = item.key()?;
batch.remove(&self.data, k);
}
for item in self.meta.iter() {
let k = item.key()?;
batch.remove(&self.meta, k);
}
batch.commit()?;
Ok(())
}
pub fn active_expire_cycle(&self, sample_limit: usize) -> Result<usize> {
let now_ms = current_now_ms();
let mut cleaned = 0;
let mut batch = self.db.batch();
for guard in self.meta.iter().take(sample_limit) {
let (k, v) = guard.into_inner()?;
if let Some(base_meta) = KeyMeta::decode(&v)
&& base_meta.is_expired(now_ms)
{
batch.remove(&self.meta, &*k);
cleaned += 1;
}
}
for guard in self.data.iter().take(sample_limit) {
let (k, v) = guard.into_inner()?;
let (expire_at, _) = decode_string_value(&v);
if is_string_expired(expire_at, now_ms) {
batch.remove(&self.data, &*k);
cleaned += 1;
}
}
if cleaned > 0 {
batch.commit()?;
}
Ok(cleaned)
}
}