wedb_embed 0.1.2

Embedded database engine providing Redis-like APIs, built on fjall / 嵌入式数据库引擎,提供类似 Redis 的接口,底层基于 fjall 开发
Documentation
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>,
  /// 默认命名空间业务数据表(0 字节纯裸键)
  pub data: Keyspace,
  /// 默认命名空间复合元数据表
  pub meta: Keyspace,
  /// 多租户命名空间业务数据表(带租户前缀)
  pub data_ns: Keyspace,
  /// 多租户命名空间复合元数据、Catalog 目录与 ID 映射表
  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;

  /// 统一打开数据库接口,支持传入配置项列表 conf_li
  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))
  }

  /// 清空数据库中的全部数据与元数据(FLUSHALL / FLUSHDB)
  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();

    // 1. 扫描 meta (默认租户 db 0 复合元数据)
    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(())
    })?;

    // 2. 扫描 meta_ns (多租户 / 多库 复合元数据)
    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(())
    })?;

    // 3. 针对原生 data 空间 (默认 db 0):仅扫描 String 键 (Tag == 0x00)
    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(())
    })?;

    // 4. 针对 data_ns 空间 (多租户 / 默认多库):仅扫描 String 键 (Tag == 0x00)
    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)
  }
}