use core::result;
use std::sync::Arc;
use thiserror::Error as ThisError;
use waof::WalLog;
use wdev::Device;
use wkv::{RangeIndexError, StorageBackend, StoreSession, TreeTuning, WedbStore};
use crate::{
aof::{self, AofEntryRef, AofOp, Replay, TtlPurgePayload, encode_entry},
resp,
};
#[derive(Debug, ThisError)]
pub enum Error {
#[error(transparent)]
Store(#[from] wkv::Error),
#[error(transparent)]
RangeIndex(#[from] RangeIndexError),
#[error(transparent)]
Wal(#[from] waof::Error),
#[error(transparent)]
Aof(#[from] aof::Error),
#[error("wal overwritten during replay: {skipped} record(s) lost")]
WalOverwritten { skipped: u64 },
}
pub type Result<T> = result::Result<T, Error>;
pub type SharedStore<D> = Arc<WedbStore<D>>;
pub struct NodeService<D: Device> {
session: StoreSession<D>,
wal: Arc<WalLog<D>>,
}
impl<D: Device> NodeService<D> {
pub fn new(store: SharedStore<D>, wal: Arc<WalLog<D>>) -> Result<Self>
where
D: 'static,
{
let sink = Arc::clone(&wal);
let _ = store.set_ttl_purge_listener(Arc::new(move |ns, db, key, expire_at_ms| {
let payload = TtlPurgePayload {
ns,
db,
expire_at_ms,
};
let _ = sink.enqueue(&encode_entry(AofOp::TtlPurge, 0, key, &payload.encode()));
}));
store.start_gc();
let session = store.new_session()?;
Ok(Self { session, wal })
}
#[inline]
pub fn session(&self) -> &StoreSession<D> {
&self.session
}
#[inline]
pub fn store(&self) -> &SharedStore<D> {
&self.session.store
}
#[inline]
pub fn wal(&self) -> &Arc<WalLog<D>> {
&self.wal
}
pub async fn ri_create(
&self,
key: &[u8],
storage_backend: StorageBackend,
tuning: TreeTuning,
) -> Result<()> {
let frame = resp::encode_ri_create(key, &storage_backend, tuning);
self
.session
.range_index_create(key, storage_backend, tuning)
.await?;
self.append(AofOp::RiCreate, key, &frame)?;
Ok(())
}
pub async fn ri_set(&self, key: &[u8], field: &[u8], value: &[u8]) -> Result<()> {
let frame = resp::encode_ri_set(key, field, value);
self.session.range_index_set(key, field, value).await?;
self.append(AofOp::RiSet, key, &frame)?;
Ok(())
}
pub async fn ri_del(&self, key: &[u8], field: &[u8]) -> Result<bool> {
let deleted = self.session.range_index_del(key, field).await?;
let frame = resp::encode_ri_del(key, field);
self.append(AofOp::RiDel, key, &frame)?;
Ok(deleted)
}
pub async fn replay(&self, replay: &mut impl Replay) -> Result<u64> {
let mut iter = self.wal.scan_committed();
let mut count = 0u64;
while let Some(record) = iter.next().await? {
let entry = AofEntryRef::decode(&record.payload)?;
if entry.op == AofOp::TtlPurge {
let payload = TtlPurgePayload::decode(entry.blob)?;
replay.on_ttl_purge(entry.key, payload)?;
} else {
replay.on_entry(entry)?;
}
count += 1;
}
let overwritten = iter.overwritten_skips();
if overwritten != 0 {
return Err(Error::WalOverwritten {
skipped: overwritten,
});
}
Ok(count)
}
pub async fn apply_ttl_purge(&self, ns: u64, db: u64, user_key: &[u8]) -> Result<()> {
let (prev_ns, prev_db) = (self.session.namespace(), self.session.active_db());
self.session.set_context(ns, db);
let deleted = self.session.delete(user_key).await;
self.session.set_context(prev_ns, prev_db);
deleted?;
Ok(())
}
#[inline]
fn append(&self, op: AofOp, key: &[u8], blob: &[u8]) -> Result<()> {
self.wal.enqueue(&encode_entry(op, 0, key, blob))?;
Ok(())
}
}