#![cfg(feature = "kv-indxdb")]
use crate::err::Error;
use crate::key::debug::Sprintable;
use crate::kvs::savepoint::{SavePointImpl, SavePoints};
use crate::kvs::{Check, Key, KeyEncode, Val};
use std::fmt::Debug;
use std::ops::Range;
pub struct Datastore {
db: indxdb::Db,
}
pub struct Transaction {
done: bool,
write: bool,
check: Check,
inner: indxdb::Tx,
save_points: SavePoints,
}
impl Drop for Transaction {
fn drop(&mut self) {
if !self.done && self.write {
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::Error => {
error!("A transaction was dropped without being committed or cancelled");
}
}
}
}
}
impl Datastore {
pub async fn new(path: &str) -> Result<Datastore, Error> {
match indxdb::db::new(path).await {
Ok(db) => Ok(Datastore {
db,
}),
Err(e) => Err(Error::Ds(e.to_string())),
}
}
pub(crate) async fn shutdown(&self) -> Result<(), Error> {
Ok(())
}
pub async fn transaction(&self, write: bool, _: bool) -> Result<Transaction, Error> {
#[cfg(not(debug_assertions))]
let check = Check::Warn;
#[cfg(debug_assertions)]
let check = Check::Error;
match self.db.begin(write).await {
Ok(inner) => Ok(Transaction {
done: false,
check,
write,
inner,
save_points: Default::default(),
}),
Err(e) => Err(Error::Tx(e.to_string())),
}
}
}
impl super::api::Transaction for Transaction {
fn supports_reverse_scan(&self) -> bool {
false
}
fn check_level(&mut self, check: Check) {
self.check = check;
}
fn closed(&self) -> bool {
self.done
}
fn writeable(&self) -> bool {
self.write
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::api", skip(self))]
async fn cancel(&mut self) -> Result<(), Error> {
if self.done {
return Err(Error::TxFinished);
}
self.done = true;
self.inner.cancel().await?;
Ok(())
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::api", skip(self))]
async fn commit(&mut self) -> Result<(), Error> {
if self.done {
return Err(Error::TxFinished);
}
if !self.write {
return Err(Error::TxReadonly);
}
self.done = true;
self.inner.commit().await?;
Ok(())
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::api", skip(self), fields(key = key.sprint()))]
async fn exists<K>(&mut self, key: K, version: Option<u64>) -> Result<bool, Error>
where
K: KeyEncode + Sprintable + Debug,
{
if version.is_some() {
return Err(Error::UnsupportedVersionedQueries);
}
if self.done {
return Err(Error::TxFinished);
}
let res = self.inner.exi(key.encode_owned()?).await?;
Ok(res)
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::api", skip(self), fields(key = key.sprint()))]
async fn get<K>(&mut self, key: K, version: Option<u64>) -> Result<Option<Val>, Error>
where
K: KeyEncode + Sprintable + Debug,
{
if version.is_some() {
return Err(Error::UnsupportedVersionedQueries);
}
if self.done {
return Err(Error::TxFinished);
}
let res = self.inner.get(key.encode_owned()?).await?;
Ok(res)
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::api", skip(self), fields(key = key.sprint()))]
async fn set<K, V>(&mut self, key: K, val: V, version: Option<u64>) -> Result<(), Error>
where
K: KeyEncode + Sprintable + Debug,
V: Into<Val> + Debug,
{
if version.is_some() {
return Err(Error::UnsupportedVersionedQueries);
}
if self.done {
return Err(Error::TxFinished);
}
if !self.write {
return Err(Error::TxReadonly);
}
self.inner.set(key.encode_owned()?, val.into()).await?;
Ok(())
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::api", skip(self), fields(key = key.sprint()))]
async fn put<K, V>(&mut self, key: K, val: V, version: Option<u64>) -> Result<(), Error>
where
K: KeyEncode + Sprintable + Debug,
V: Into<Val> + Debug,
{
if version.is_some() {
return Err(Error::UnsupportedVersionedQueries);
}
if self.done {
return Err(Error::TxFinished);
}
if !self.write {
return Err(Error::TxReadonly);
}
self.inner.put(key.encode_owned()?, val.into()).await?;
Ok(())
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::api", skip(self), fields(key = key.sprint()))]
async fn putc<K, V>(&mut self, key: K, val: V, chk: Option<V>) -> Result<(), Error>
where
K: KeyEncode + Sprintable + Debug,
V: Into<Val> + Debug,
{
if self.done {
return Err(Error::TxFinished);
}
if !self.write {
return Err(Error::TxReadonly);
}
self.inner.putc(key.encode_owned()?, val.into(), chk.map(Into::into)).await?;
Ok(())
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::api", skip(self), fields(key = key.sprint()))]
async fn del<K>(&mut self, key: K) -> Result<(), Error>
where
K: KeyEncode + Sprintable + Debug,
{
if self.done {
return Err(Error::TxFinished);
}
if !self.write {
return Err(Error::TxReadonly);
}
let res = self.inner.del(key.encode_owned()?).await?;
Ok(res)
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::api", skip(self), fields(key = key.sprint()))]
async fn delc<K, V>(&mut self, key: K, chk: Option<V>) -> Result<(), Error>
where
K: KeyEncode + Sprintable + Debug,
V: Into<Val> + Debug,
{
if self.done {
return Err(Error::TxFinished);
}
if !self.write {
return Err(Error::TxReadonly);
}
let res = self.inner.delc(key.encode_owned()?, chk.map(Into::into)).await?;
Ok(res)
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::api", skip(self), fields(rng = rng.sprint()))]
async fn keys<K>(
&mut self,
rng: Range<K>,
limit: u32,
version: Option<u64>,
) -> Result<Vec<Key>, Error>
where
K: KeyEncode + Sprintable + Debug,
{
if version.is_some() {
return Err(Error::UnsupportedVersionedQueries);
}
if self.done {
return Err(Error::TxFinished);
}
let rng: Range<Key> = Range {
start: rng.start.encode_owned()?,
end: rng.end.encode_owned()?,
};
let res = self.inner.keys(rng, limit).await?;
Ok(res)
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::api", skip(self), fields(rng = rng.sprint()))]
async fn scan<K>(
&mut self,
rng: Range<K>,
limit: u32,
version: Option<u64>,
) -> Result<Vec<(Key, Val)>, Error>
where
K: KeyEncode + Sprintable + Debug,
{
if version.is_some() {
return Err(Error::UnsupportedVersionedQueries);
}
if self.done {
return Err(Error::TxFinished);
}
let rng: Range<Key> = Range {
start: rng.start.encode_owned()?,
end: rng.end.encode_owned()?,
};
let res = self.inner.scan(rng, limit).await?;
Ok(res)
}
}
impl SavePointImpl for Transaction {
fn get_save_points(&mut self) -> &mut SavePoints {
&mut self.save_points
}
}