#![cfg(feature = "kv-fdb")]
use crate::err::Error;
use crate::kvs::Check;
use crate::kvs::Key;
use crate::kvs::Val;
use crate::vs::{u64_to_versionstamp, Versionstamp};
use foundationdb::options;
use futures::TryStreamExt;
use std::ops::Range;
use std::sync::Arc;
use crate::key::error::KeyCategory;
use foundationdb::options::MutationType;
use futures::lock::Mutex;
use once_cell::sync::Lazy;
pub struct Datastore {
db: foundationdb::Database,
_fdbnet: Arc<foundationdb::api::NetworkAutoStop>,
}
pub struct Transaction {
done: bool,
lock: bool,
write: bool,
check: Check,
inner: Arc<Mutex<Option<foundationdb::Transaction>>>,
}
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> {
static FDBNET: Lazy<Arc<foundationdb::api::NetworkAutoStop>> =
Lazy::new(|| Arc::new(unsafe { foundationdb::boot() }));
let _fdbnet = (*FDBNET).clone();
match foundationdb::Database::from_path(path) {
Ok(db) => {
db.set_option(options::DatabaseOption::TransactionRetryLimit(5)).map_err(|e| {
Error::Ds(format!("Unable to set transaction retry limit: {}", e))
})?;
db.set_option(options::DatabaseOption::TransactionTimeout(5000))
.map_err(|e| Error::Ds(format!("Unable to set transaction timeout: {}", e)))?;
db.set_option(options::DatabaseOption::TransactionMaxRetryDelay(500)).map_err(
|e| Error::Ds(format!("Unable to set transaction max retry delay: {}", e)),
)?;
Ok(Datastore {
db,
_fdbnet,
})
}
Err(e) => Err(Error::Ds(e.to_string())),
}
}
pub(crate) async fn transaction(&self, write: bool, lock: bool) -> Result<Transaction, Error> {
#[cfg(not(debug_assertions))]
let check = Check::Warn;
#[cfg(debug_assertions)]
let check = Check::Panic;
match self.db.create_trx() {
Ok(inner) => Ok(Transaction {
done: false,
check,
write,
lock,
inner: Arc::new(Mutex::new(Some(inner))),
}),
Err(e) => Err(Error::Tx(e.to_string())),
}
}
}
impl Transaction {
pub(crate) fn check_level(&mut self, check: Check) {
self.check = check;
}
pub(crate) fn closed(&self) -> bool {
self.done
}
fn snapshot(&self) -> bool {
!self.write && !self.lock
}
pub(crate) async fn cancel(&mut self) -> Result<(), Error> {
if self.done {
return Err(Error::TxFinished);
}
self.done = true;
let inner = match self.inner.lock().await.take() {
Some(inner) => {
let tc = inner.cancel();
tc.reset()
}
_ => return Err(Error::Ds("Unexpected error".to_string())),
};
self.inner = Arc::new(Mutex::new(Some(inner)));
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;
let r = match self.inner.lock().await.take() {
Some(inner) => inner.commit().await,
_ => return Err(Error::Ds("Unexpected error".to_string())),
};
match r {
Ok(_r) => {}
Err(e) => {
return Err(Error::Tx(format!("Transaction commit error: {}", e)));
}
}
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 key: Vec<u8> = key.into();
let key: &[u8] = &key[..];
let inner = self.inner.lock().await;
let inner = inner.as_ref().unwrap();
inner
.get(key, self.snapshot())
.await
.map(|v| v.is_some())
.map_err(|e| Error::Tx(format!("Unable to get kv from FoundationDB: {}", e)))
}
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 key: Vec<u8> = key.into();
let key = &key[..];
let inner = self.inner.lock().await;
let inner = inner.as_ref().unwrap();
inner
.get(key, self.snapshot())
.await
.map(|v| v.as_ref().map(|v| v.to_vec()))
.map_err(|e| Error::Tx(format!("Unable to get kv from FoundationDB: {}", e)))
}
#[allow(unused)]
pub(crate) async fn get_timestamp(&mut self) -> Result<Versionstamp, Error> {
if self.done {
return Err(Error::TxFinished);
}
let inner = self.inner.lock().await;
let inner = inner.as_ref().unwrap();
let res = inner
.get_read_version()
.await
.map_err(|e| Error::Tx(format!("Unable to get read version from FDB: {}", e)))?;
let res: u64 = res.try_into().unwrap();
let res = u64_to_versionstamp(res);
Ok(res)
}
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);
}
let key: Vec<u8> = key.into();
let key = &key[..];
let val: Vec<u8> = val.into();
let val = &val[..];
let inner = self.inner.lock().await;
let inner = inner.as_ref().unwrap();
inner.set(key, val);
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 key: Vec<u8> = key.into();
if self.exi(key.clone().as_slice()).await? {
return Err(Error::TxKeyAlreadyExistsCategory(category));
}
let key: &[u8] = &key[..];
let val: Vec<u8> = val.into();
let val: &[u8] = &val[..];
let inner = self.inner.lock().await;
let inner = inner.as_ref().unwrap();
inner.set(key, val);
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 key: Vec<u8> = key.into();
let key: &[u8] = key.as_slice();
let val: Vec<u8> = val.into();
let val: &[u8] = val.as_slice();
let chk = chk.map(Into::into);
let inner = self.inner.lock().await;
let inner = inner.as_ref().unwrap();
let res = inner.get(key, false).await;
let res = res.map_err(|e| Error::Tx(format!("Unable to get kv from FoundationDB: {}", e)));
match (res, chk) {
(Ok(Some(v)), Some(w)) if *v.as_ref() == w => inner.set(key, val),
(Ok(None), None) => inner.set(key, val),
(Err(e), _) => return Err(e),
_ => return Err(Error::TxConditionNotMet),
};
Ok(())
}
pub(crate) async fn set_versionstamped_key<K, V>(
&mut self,
prefix: K,
suffix: 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 mut k: Vec<u8> = prefix.into();
let pos = k.len();
let pos: u32 = pos.try_into().unwrap();
let mut ts_placeholder: Vec<u8> =
vec![0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00];
k.append(&mut ts_placeholder);
k.append(&mut suffix.into());
let mut posbs: Vec<u8> = pos.to_le_bytes().to_vec();
k.append(&mut posbs);
let key: &[u8] = &k[..];
let val: Vec<u8> = val.into();
let val: &[u8] = &val[..];
let inner = self.inner.lock().await;
let inner = inner.as_ref().unwrap();
inner.atomic_op(key, val, MutationType::SetVersionstampedKey);
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);
}
let key: Vec<u8> = key.into();
let key: &[u8] = key.as_slice();
let inner = self.inner.lock().await;
let inner = inner.as_ref().unwrap();
inner.clear(key);
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 key: Vec<u8> = key.into();
let key: &[u8] = key.as_slice();
let chk: Option<Val> = chk.map(Into::into);
let inner = self.inner.lock().await;
let inner = inner.as_ref().unwrap();
let res = inner
.get(key, false)
.await
.map_err(|e| Error::Tx(format!("FoundationDB inner failure: {}", e)));
match (res, chk) {
(Ok(Some(v)), Some(w)) if *v.as_ref() == w => inner.clear(key),
(Ok(None), None) => inner.clear(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 rng: Range<Key> = Range {
start: rng.start.into(),
end: rng.end.into(),
};
let begin: Vec<u8> = rng.start;
let end: Vec<u8> = rng.end;
let opt = foundationdb::RangeOption {
limit: Some(limit.try_into().unwrap()),
..foundationdb::RangeOption::from((begin.as_slice(), end.as_slice()))
};
let inner = self.inner.lock().await;
let inner = inner.as_ref().unwrap();
let mut stream = inner.get_ranges_keyvalues(opt, self.snapshot());
let mut res: Vec<(Key, Val)> = vec![];
loop {
let x = stream.try_next().await;
match x {
Ok(Some(v)) => {
let x = (Key::from(v.key()), Val::from(v.value()));
res.push(x)
}
Ok(None) => break,
Err(e) => return Err(Error::Tx(format!("GetRanges failed: {}", e))),
}
}
Ok(res)
}
pub(crate) async fn delr<K>(&mut self, rng: Range<K>) -> Result<(), Error>
where
K: Into<Key>,
{
if self.done {
return Err(Error::TxFinished);
}
if !self.write {
return Err(Error::TxReadonly);
}
let begin: &[u8] = &rng.start.into();
let end: &[u8] = &rng.end.into();
let inner = self.inner.lock().await;
let inner = inner.as_ref().unwrap();
inner.clear_range(begin, end);
Ok(())
}
}