#![cfg(feature = "kv-tikv")]
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 std::ops::Range;
use tikv::CheckLevel;
use tikv::TimestampExt;
use tikv::TransactionOptions;
pub struct Datastore {
db: tikv::TransactionClient,
}
pub struct Transaction {
done: bool,
write: bool,
check: Check,
inner: tikv::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> {
match tikv::TransactionClient::new(vec![path]).await {
Ok(db) => Ok(Datastore {
db,
}),
Err(e) => Err(Error::Ds(e.to_string())),
}
}
pub(crate) async fn transaction(&self, write: bool, lock: bool) -> Result<Transaction, Error> {
#[cfg(debug_assertions)]
if lock {
panic!("There are issues with pessimistic locking in TiKV");
}
let mut opt = if lock {
TransactionOptions::new_pessimistic()
} else {
TransactionOptions::new_optimistic()
};
opt = opt.drop_check(CheckLevel::Warn);
if !write {
opt = opt.read_only();
}
#[cfg(not(debug_assertions))]
let check = Check::Warn;
#[cfg(debug_assertions)]
let check = Check::Panic;
match self.db.begin_with_options(opt).await {
Ok(inner) => Ok(Transaction {
done: false,
check,
write,
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
}
pub(crate) async fn cancel(&mut self) -> Result<(), Error> {
if self.done {
return Err(Error::TxFinished);
}
self.done = true;
if self.write {
self.inner.rollback().await?;
}
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;
if let Err(err) = self.inner.commit().await {
if let Err(inner_err) = self.inner.rollback().await {
error!("Transaction commit failed {} and rollback failed: {}", err, inner_err);
}
return Err(err.into());
}
Ok(())
}
#[allow(unused)]
pub(crate) async fn get_timestamp<K>(
&mut self,
key: K,
lock: bool,
) -> Result<Versionstamp, Error>
where
K: Into<Key>,
{
if self.done {
return Err(Error::TxFinished);
}
let res = self.inner.get_current_timestamp().await?;
let ver = res.version();
let verbytes = u64_to_versionstamp(ver);
let k: Key = key.into();
if lock {
let prev = self.inner.get(k.clone()).await?;
if let Some(prev) = 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)?;
if prev >= ver {
return Err(Error::TxFailure);
}
}
self.inner.put(k, verbytes.to_vec()).await?;
}
Ok(u64_to_versionstamp(ver))
}
#[allow(unused)]
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, false).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 exi<K>(&mut self, key: K) -> Result<bool, Error>
where
K: Into<Key>,
{
if self.done {
return Err(Error::TxFinished);
}
let res = self.inner.key_exists(key.into()).await?;
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.get(key.into()).await?;
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);
}
self.inner.put(key.into(), val.into()).await?;
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 = key.into();
let val = val.into();
match self.inner.key_exists(key.clone()).await? {
false => self.inner.put(key, val).await?,
_ => 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 key = key.into();
let val = val.into();
let chk = chk.map(Into::into);
match (self.inner.get(key.clone()).await?, chk) {
(Some(v), Some(w)) if v == w => self.inner.put(key, val).await?,
(None, None) => self.inner.put(key, val).await?,
_ => 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.delete(key.into()).await?;
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 = key.into();
let chk = chk.map(Into::into);
match (self.inner.get(key.clone()).await?, chk) {
(Some(v), Some(w)) if v == w => self.inner.delete(key).await?,
(None, None) => self.inner.delete(key).await?,
_ => 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 res = self.inner.scan(rng, limit).await?;
let res = res.map(|kv| (Key::from(kv.0), kv.1)).collect();
Ok(res)
}
pub(crate) async fn delr<K>(&mut self, rng: Range<K>, limit: u32) -> Result<(), Error>
where
K: Into<Key>,
{
if self.done {
return Err(Error::TxFinished);
}
if !self.write {
return Err(Error::TxReadonly);
}
let rng: Range<Key> = Range {
start: rng.start.into(),
end: rng.end.into(),
};
let res = self.inner.scan_keys(rng, limit).await?;
for key in res {
self.inner.delete(key).await?;
}
Ok(())
}
}