#![cfg(feature = "kv-speedb")]
mod cnf;
use crate::err::Error;
use crate::key::error::KeyCategory;
use crate::kvs::Check;
use crate::kvs::Key;
use crate::kvs::Val;
use crate::vs::{try_to_u64_be, u64_to_versionstamp, Versionstamp};
use futures::lock::Mutex;
use speedb::{
DBCompactionStyle, DBCompressionType, LogLevel, OptimisticTransactionDB,
OptimisticTransactionOptions, Options, ReadOptions, WriteOptions,
};
use std::ops::Range;
use std::pin::Pin;
use std::sync::Arc;
#[derive(Clone)]
pub struct Datastore {
db: Pin<Arc<OptimisticTransactionDB>>,
}
pub struct Transaction {
done: bool,
write: bool,
check: Check,
inner: Arc<Mutex<Option<speedb::Transaction<'static, OptimisticTransactionDB>>>>,
ro: ReadOptions,
_db: Pin<Arc<OptimisticTransactionDB>>,
}
impl Drop for Transaction {
fn drop(&mut self) {
if !self.done && self.write {
if std::thread::panicking() {
return;
}
match self.check {
Check::None => {
trace!("A transaction was dropped without being committed or cancelled");
}
Check::Warn => {
warn!("A transaction was dropped without being committed or cancelled");
}
Check::Panic => {
#[cfg(debug_assertions)]
{
let backtrace = std::backtrace::Backtrace::force_capture();
if let std::backtrace::BacktraceStatus::Captured = backtrace.status() {
println!("{}", backtrace);
}
}
panic!("A transaction was dropped without being committed or cancelled");
}
}
}
}
}
impl Datastore {
pub(crate) async fn new(path: &str) -> Result<Datastore, Error> {
let mut opts = Options::default();
opts.set_use_fsync(false);
opts.set_log_level(LogLevel::Warn);
opts.set_keep_log_file_num(*cnf::SPEEDB_KEEP_LOG_FILE_NUM);
opts.create_if_missing(true);
opts.create_missing_column_families(true);
opts.set_compaction_style(DBCompactionStyle::Level);
opts.increase_parallelism(*cnf::SPEEDB_THREAD_COUNT);
opts.set_max_write_buffer_number(*cnf::SPEEDB_MAX_WRITE_BUFFER_NUMBER);
opts.set_write_buffer_size(*cnf::SPEEDB_WRITE_BUFFER_SIZE);
opts.set_target_file_size_base(*cnf::SPEEDB_TARGET_FILE_SIZE_BASE);
opts.set_min_write_buffer_number_to_merge(*cnf::SPEEDB_MIN_WRITE_BUFFER_NUMBER_TO_MERGE);
opts.set_enable_pipelined_write(*cnf::SPEEDB_ENABLE_PIPELINED_WRITES);
opts.set_enable_blob_files(*cnf::SPEEDB_ENABLE_BLOB_FILES);
opts.set_min_blob_size(*cnf::SPEEDB_MIN_BLOB_SIZE);
opts.set_compression_per_level(&[
DBCompressionType::None,
DBCompressionType::None,
DBCompressionType::Lz4hc,
DBCompressionType::Lz4hc,
DBCompressionType::Lz4hc,
]);
Ok(Datastore {
db: Arc::pin(OptimisticTransactionDB::open(&opts, path)?),
})
}
pub(crate) async fn transaction(&self, write: bool, _: bool) -> Result<Transaction, Error> {
let mut to = OptimisticTransactionOptions::default();
to.set_snapshot(true);
let mut wo = WriteOptions::default();
wo.set_sync(false);
let inner = self.db.transaction_opt(&wo, &to);
let inner = unsafe {
std::mem::transmute::<
speedb::Transaction<'_, OptimisticTransactionDB>,
speedb::Transaction<'static, OptimisticTransactionDB>,
>(inner)
};
let mut ro = ReadOptions::default();
ro.set_snapshot(&inner.snapshot());
ro.fill_cache(true);
#[cfg(not(debug_assertions))]
let check = Check::Warn;
#[cfg(debug_assertions)]
let check = Check::Panic;
Ok(Transaction {
done: false,
check,
write,
inner: Arc::new(Mutex::new(Some(inner))),
ro,
_db: self.db.clone(),
})
}
}
impl Transaction {
pub(crate) fn check_level(&mut self, check: Check) {
self.check = check;
}
pub(crate) fn closed(&self) -> bool {
self.done
}
pub(crate) async fn cancel(&mut self) -> Result<(), Error> {
if self.done {
return Err(Error::TxFinished);
}
self.done = true;
match self.inner.lock().await.take() {
Some(inner) => inner.rollback()?,
None => unreachable!(),
};
Ok(())
}
pub(crate) async fn commit(&mut self) -> Result<(), Error> {
if self.done {
return Err(Error::TxFinished);
}
if !self.write {
return Err(Error::TxReadonly);
}
self.done = true;
match self.inner.lock().await.take() {
Some(inner) => inner.commit()?,
None => unreachable!(),
};
Ok(())
}
pub(crate) async fn exi<K>(&mut self, key: K) -> Result<bool, Error>
where
K: Into<Key>,
{
if self.done {
return Err(Error::TxFinished);
}
let res =
self.inner.lock().await.as_ref().unwrap().get_opt(key.into(), &self.ro)?.is_some();
Ok(res)
}
pub(crate) async fn get<K>(&mut self, key: K) -> Result<Option<Val>, Error>
where
K: Into<Key>,
{
if self.done {
return Err(Error::TxFinished);
}
let res = self.inner.lock().await.as_ref().unwrap().get_opt(key.into(), &self.ro)?;
Ok(res)
}
#[allow(unused)]
pub(crate) async fn get_timestamp<K>(&mut self, key: K) -> Result<Versionstamp, Error>
where
K: Into<Key>,
{
if self.done {
return Err(Error::TxFinished);
}
let k: Key = key.into();
let prev = self.inner.lock().await.as_ref().unwrap().get_opt(k.clone(), &self.ro)?;
let ver = match prev {
Some(prev) => {
let slice = prev.as_slice();
let res: Result<[u8; 10], Error> = match slice.try_into() {
Ok(ba) => Ok(ba),
Err(e) => Err(Error::Ds(e.to_string())),
};
let array = res?;
let prev = try_to_u64_be(array)?;
prev + 1
}
None => 1,
};
let verbytes = u64_to_versionstamp(ver);
self.inner.lock().await.as_ref().unwrap().put(k, verbytes)?;
Ok(verbytes)
}
pub(crate) async fn get_versionstamped_key<K>(
&mut self,
ts_key: K,
prefix: K,
suffix: K,
) -> Result<Vec<u8>, Error>
where
K: Into<Key>,
{
if self.done {
return Err(Error::TxFinished);
}
if !self.write {
return Err(Error::TxReadonly);
}
let ts = self.get_timestamp(ts_key).await?;
let mut k: Vec<u8> = prefix.into();
k.append(&mut ts.to_vec());
k.append(&mut suffix.into());
Ok(k)
}
pub(crate) async fn set<K, V>(&mut self, key: K, val: V) -> Result<(), Error>
where
K: Into<Key>,
V: Into<Val>,
{
if self.done {
return Err(Error::TxFinished);
}
if !self.write {
return Err(Error::TxReadonly);
}
self.inner.lock().await.as_ref().unwrap().put(key.into(), val.into())?;
Ok(())
}
pub(crate) async fn put<K, V>(
&mut self,
category: KeyCategory,
key: K,
val: V,
) -> Result<(), Error>
where
K: Into<Key>,
V: Into<Val>,
{
if self.done {
return Err(Error::TxFinished);
}
if !self.write {
return Err(Error::TxReadonly);
}
let inner = self.inner.lock().await;
let inner = inner.as_ref().unwrap();
let key = key.into();
let val = val.into();
match inner.get_opt(&key, &self.ro)? {
None => inner.put(key, val)?,
_ => return Err(Error::TxKeyAlreadyExistsCategory(category)),
};
Ok(())
}
pub(crate) async fn putc<K, V>(&mut self, key: K, val: V, chk: Option<V>) -> Result<(), Error>
where
K: Into<Key>,
V: Into<Val>,
{
if self.done {
return Err(Error::TxFinished);
}
if !self.write {
return Err(Error::TxReadonly);
}
let inner = self.inner.lock().await;
let inner = inner.as_ref().unwrap();
let key = key.into();
let val = val.into();
let chk = chk.map(Into::into);
match (inner.get_opt(&key, &self.ro)?, chk) {
(Some(v), Some(w)) if v == w => inner.put(key, val)?,
(None, None) => inner.put(key, val)?,
_ => return Err(Error::TxConditionNotMet),
};
Ok(())
}
pub(crate) async fn del<K>(&mut self, key: K) -> Result<(), Error>
where
K: Into<Key>,
{
if self.done {
return Err(Error::TxFinished);
}
if !self.write {
return Err(Error::TxReadonly);
}
self.inner.lock().await.as_ref().unwrap().delete(key.into())?;
Ok(())
}
pub(crate) async fn delc<K, V>(&mut self, key: K, chk: Option<V>) -> Result<(), Error>
where
K: Into<Key>,
V: Into<Val>,
{
if self.done {
return Err(Error::TxFinished);
}
if !self.write {
return Err(Error::TxReadonly);
}
let inner = self.inner.lock().await;
let inner = inner.as_ref().unwrap();
let key = key.into();
let chk = chk.map(Into::into);
match (inner.get_opt(&key, &self.ro)?, chk) {
(Some(v), Some(w)) if v == w => inner.delete(key)?,
(None, None) => inner.delete(key)?,
_ => return Err(Error::TxConditionNotMet),
};
Ok(())
}
pub(crate) async fn scan<K>(
&mut self,
rng: Range<K>,
limit: u32,
) -> Result<Vec<(Key, Val)>, Error>
where
K: Into<Key>,
{
if self.done {
return Err(Error::TxFinished);
}
let inner = self.inner.lock().await;
let inner = inner.as_ref().unwrap();
let rng: Range<Key> = Range {
start: rng.start.into(),
end: rng.end.into(),
};
let mut res = vec![];
let beg = rng.start.as_slice();
let end = rng.end.as_slice();
let mut ro = ReadOptions::default();
ro.set_snapshot(&inner.snapshot());
let mut iter = inner.raw_iterator_opt(ro);
iter.seek(&rng.start);
while iter.valid() {
if res.len() < limit as usize {
let (k, v) = (iter.key(), iter.value());
if let (Some(k), Some(v)) = (k, v) {
if k >= beg && k < end {
res.push((k.to_vec(), v.to_vec()));
iter.next();
continue;
}
}
}
break;
}
Ok(res)
}
}