use std::{
fmt,
fs::create_dir_all,
ops::Bound,
path::Path,
sync::{
Arc,
atomic::{AtomicBool, Ordering},
},
thread::{Builder, JoinHandle, park_timeout},
time::Duration,
};
use fjall::{
CompressionType, Database, Keyspace, KeyspaceCreateOptions, PersistMode as FjallPersistMode,
};
use crate::{
api::key::cleanup_composite_data_raw,
conf::{Conf, DbConfig, PersistMode},
error::{Error, Result},
key_composer::{KeyComposer, KeyTag, SystemDomainTag},
keyspace::{self, DATA, META},
meta::{KeyMeta, current_now_ms, init_version_counter},
string::{decode_string_value, is_string_expired},
};
#[derive(Default, Debug)]
pub(crate) struct ExpireCursors {
pub meta: Option<Vec<u8>>,
pub meta_ns: Option<Vec<u8>>,
pub data: Option<Vec<u8>>,
pub data_ns: Option<Vec<u8>>,
}
pub(crate) struct FlusherGuard {
stop: Arc<AtomicBool>,
handle: parking_lot::Mutex<Option<JoinHandle<()>>>,
}
impl Drop for FlusherGuard {
fn drop(&mut self) {
self.stop.store(true, Ordering::Release);
let mut guard = self.handle.lock();
if let Some(h) = guard.take() {
h.thread().unpark();
let _ = h.join();
}
}
}
#[derive(Clone)]
pub struct WeDb {
pub db: Arc<Database>,
pub data: Keyspace,
pub meta: Keyspace,
pub data_ns: Keyspace,
pub meta_ns: Keyspace,
pub(crate) ns_cache: Arc<papaya::HashMap<hipstr::HipStr<'static>, u64>>,
pub(crate) ns_lock: Arc<parking_lot::Mutex<()>>,
pub(crate) expire_cursor: Arc<parking_lot::Mutex<ExpireCursors>>,
pub(crate) _flusher: Option<Arc<FlusherGuard>>,
}
#[inline]
fn sweep_ks<F>(
ks: &Keyspace,
cursor: &mut Option<Vec<u8>>,
sample_limit: usize,
mut on_entry: F,
) -> Result<()>
where
F: FnMut(&[u8], &[u8]) -> Result<()>,
{
let mut count = 0;
let mut last_key = None;
let iter = match cursor.as_ref() {
Some(cur) => ks.range::<&[u8], _>((Bound::Excluded(cur.as_slice()), Bound::Unbounded)),
None => ks.iter(),
};
for guard in iter.take(sample_limit) {
let (k, v) = guard.into_inner()?;
on_entry(&k, &v)?;
count += 1;
if count == sample_limit {
last_key = Some(k.to_vec());
}
}
if count < sample_limit {
*cursor = None;
} else {
*cursor = last_key;
}
Ok(())
}
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<I>(path: impl AsRef<Path>, conf_li: I) -> Result<Self>
where
I: IntoIterator<Item = Conf>,
{
let cfg = DbConfig::from_conf_li(conf_li);
let path = path.as_ref();
if let Some(parent) = path.parent() {
create_dir_all(parent)?;
}
init_version_counter();
let mut builder = Database::builder(path);
let journal_comp: CompressionType = cfg.journal_compression.into();
builder = builder.cache_size(cfg.cache_size as u64);
builder = builder.journal_compression(journal_comp);
builder = builder.manual_journal_persist(cfg.manual_journal_persist);
if let Some(worker_threads) = cfg.worker_threads
&& worker_threads > 0
{
builder = builder.worker_threads(worker_threads);
}
if let Some(max_journaling_size) = cfg.max_journaling_size {
builder = builder.max_journaling_size(max_journaling_size as u64);
}
if let Some(max_cached_files) = cfg.max_cached_files {
builder = builder.max_cached_files(Some(max_cached_files));
}
if cfg.temporary {
builder = builder.temporary(true);
}
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_with_cfg(&db, &cfg)?;
let persist_interval = cfg.persist_interval_ms;
let persist_mode: FjallPersistMode = cfg.persist_mode.into();
let _flusher = if cfg.manual_journal_persist && persist_interval > 0 {
let stop = Arc::new(AtomicBool::new(false));
let stop_clone = stop.clone();
let db_weak = Arc::downgrade(&db);
let handle = Builder::new()
.name("wedb-flusher".to_string())
.spawn(move || {
let interval = Duration::from_millis(persist_interval);
while !stop_clone.load(Ordering::Acquire) {
park_timeout(interval);
if stop_clone.load(Ordering::Acquire) {
break;
}
if let Some(d) = db_weak.upgrade() {
let _ = d.persist(persist_mode);
} else {
break;
}
}
if let Some(d) = db_weak.upgrade() {
let _ = d.persist(FjallPersistMode::SyncAll);
}
})
.ok();
Some(Arc::new(FlusherGuard {
stop,
handle: parking_lot::Mutex::new(handle),
}))
} else {
None
};
Ok(Self {
db,
data: keyspace.data,
meta: keyspace.meta,
data_ns: keyspace.data_ns,
meta_ns: keyspace.meta_ns,
ns_cache: Arc::new(papaya::HashMap::default()),
ns_lock: Arc::new(parking_lot::Mutex::new(())),
expire_cursor: Arc::new(parking_lot::Mutex::new(ExpireCursors::default())),
_flusher,
})
}
#[inline]
pub fn database(&self) -> &Arc<Database> {
&self.db
}
#[inline]
pub fn raw(&self) -> &Arc<Database> {
&self.db
}
#[inline]
pub fn data(&self) -> &Keyspace {
&self.data
}
#[inline]
pub fn meta(&self) -> &Keyspace {
&self.meta
}
#[inline]
pub fn data_ns(&self) -> &Keyspace {
&self.data_ns
}
#[inline]
pub fn meta_ns(&self) -> &Keyspace {
&self.meta_ns
}
#[inline(always)]
pub(crate) fn data_ks<'a>(&'a self, kc: &KeyComposer<'_>) -> &'a Keyspace {
if kc.is_default() {
&self.data
} else {
&self.data_ns
}
}
#[inline(always)]
pub(crate) fn meta_ks<'a>(&'a self, kc: &KeyComposer<'_>) -> &'a Keyspace {
if kc.is_default() {
&self.meta
} else {
&self.meta_ns
}
}
#[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.into())
.map_err(|e| Error::internal_with_source("Persist error", e))
}
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);
}
for item in self.data_ns.iter() {
let k = item.key()?;
batch.remove(&self.data_ns, k);
}
for item in self.meta_ns.iter() {
let k = item.key()?;
batch.remove(&self.meta_ns, k);
}
batch.commit()?;
let pin = self.ns_cache.pin();
pin.clear();
Ok(())
}
pub fn active_expire_cycle(&self, sample_limit: usize) -> Result<usize> {
if sample_limit == 0 {
return Ok(0);
}
let now_ms = current_now_ms();
let mut cleaned = 0;
let mut batch = self.db.batch();
let mut buf = Vec::new();
let mut cursors = self.expire_cursor.lock();
let kc_def = KeyComposer::new("default");
sweep_ks(&self.meta, &mut cursors.meta, sample_limit, |k, v| {
if let Some(base_meta) = KeyMeta::decode(v)
&& base_meta.is_expired(now_ms)
{
batch.remove(&self.meta, k);
if !k.is_empty() {
cleanup_composite_data_raw(
&self.data,
&self.meta,
&kc_def,
k[0],
&k[1..],
&mut batch,
&mut buf,
)?;
}
cleaned += 1;
}
Ok(())
})?;
sweep_ks(&self.meta_ns, &mut cursors.meta_ns, sample_limit, |k, v| {
if k.len() >= 2 && k[0] == 0 && SystemDomainTag::is_system_domain_byte(k[1]) {
return Ok(());
}
if let Some(base_meta) = KeyMeta::decode(v)
&& base_meta.is_expired(now_ms)
{
batch.remove(&self.meta_ns, k);
if let Some((kc, _, remain)) = KeyComposer::parse_scoped_prefix(k)
&& !remain.is_empty()
{
cleanup_composite_data_raw(
&self.data_ns,
&self.meta_ns,
&kc,
remain[0],
&remain[1..],
&mut batch,
&mut buf,
)?;
}
cleaned += 1;
}
Ok(())
})?;
sweep_ks(&self.data, &mut cursors.data, sample_limit, |k, v| {
if !k.is_empty() && k[0] == KeyTag::RawString as u8 {
let (expire_at, _) = decode_string_value(v);
if is_string_expired(expire_at, now_ms) {
batch.remove(&self.data, k);
cleaned += 1;
}
}
Ok(())
})?;
sweep_ks(&self.data_ns, &mut cursors.data_ns, sample_limit, |k, v| {
if let Some((_, _, remain)) = KeyComposer::parse_scoped_prefix(k)
&& !remain.is_empty()
&& remain[0] == KeyTag::RawString as u8
{
let (expire_at, _) = decode_string_value(v);
if is_string_expired(expire_at, now_ms) {
batch.remove(&self.data_ns, k);
cleaned += 1;
}
}
Ok(())
})?;
if cleaned > 0 {
batch.commit()?;
}
Ok(cleaned)
}
}