#![warn(clippy::missing_docs_in_private_items)]
use std::future::Future;
use std::ops::Range;
use std::pin::Pin;
use anyhow::bail;
use super::cursor::{DefaultKeysCursor, DefaultValsCursor};
use super::direction::Direction;
use super::err::{Error, Result};
use super::util;
use crate::key::debug::Sprintable;
use crate::kvs::batch::Batch;
use crate::kvs::timestamp::IncTimeStamp;
use crate::kvs::{
BoxTimeStamp, BoxTimeStampImpl, COUNT_BATCH_SIZE, HlcTimeStamp, HlcTimeStampImpl,
IncTimeStampImpl, Key, NORMAL_BATCH_SIZE, Val,
};
#[cfg(target_family = "wasm")]
pub(crate) type BoxFut<'a, T> = Pin<Box<dyn Future<Output = T> + 'a>>;
#[cfg(not(target_family = "wasm"))]
pub(crate) type BoxFut<'a, T> = Pin<Box<dyn Future<Output = T> + Send + 'a>>;
#[derive(Debug, Default)]
pub struct KeysResult {
pub keys: Vec<Key>,
pub key_bytes: u64,
}
#[derive(Debug, Default)]
pub struct ScanResult {
pub values: Vec<(Key, Val)>,
pub key_bytes: u64,
pub value_bytes: u64,
}
#[derive(Debug, Default, Clone, Copy)]
pub struct ScanChunkStats {
pub rows: u64,
pub key_bytes: u64,
pub value_bytes: u64,
}
#[derive(Debug, Default)]
pub struct GetMultiResult {
pub values: Vec<Option<Val>>,
pub records: u64,
pub value_bytes: u64,
}
pub mod requirements {
#[cfg(target_family = "wasm")]
pub trait TransactionRequirements {}
#[cfg(target_family = "wasm")]
impl<T> TransactionRequirements for T {}
#[cfg(not(target_family = "wasm"))]
pub trait TransactionRequirements: Send + Sync {}
#[cfg(not(target_family = "wasm"))]
impl<T: Send + Sync> TransactionRequirements for T {}
}
#[derive(Clone, Copy, Debug)]
pub(crate) struct KeySpan {
pub(crate) offset: usize,
pub(crate) len: usize,
}
#[derive(Clone, Copy, Debug)]
pub(crate) struct KeyValSpan {
pub(crate) key_offset: usize,
pub(crate) key_len: usize,
pub(crate) val_offset: usize,
pub(crate) val_len: usize,
}
pub struct KeysBatch<'c> {
buf: &'c [u8],
spans: &'c [KeySpan],
pub key_bytes: u64,
}
impl<'c> KeysBatch<'c> {
#[inline]
pub(crate) fn from_parts(buf: &'c [u8], spans: &'c [KeySpan], key_bytes: u64) -> Self {
Self {
buf,
spans,
key_bytes,
}
}
#[inline]
pub fn len(&self) -> usize {
self.spans.len()
}
#[inline]
pub fn is_empty(&self) -> bool {
self.spans.is_empty()
}
#[inline]
pub fn get(&self, i: usize) -> Option<&[u8]> {
let span = self.spans.get(i)?;
Some(&self.buf[span.offset..span.offset + span.len])
}
#[inline]
pub fn iter(&self) -> KeysIter<'_> {
KeysIter {
buf: self.buf,
spans: self.spans.iter(),
}
}
}
impl<'a, 'c: 'a> IntoIterator for &'a KeysBatch<'c> {
type Item = &'a [u8];
type IntoIter = KeysIter<'a>;
#[inline]
fn into_iter(self) -> Self::IntoIter {
self.iter()
}
}
pub struct KeysIter<'a> {
buf: &'a [u8],
spans: std::slice::Iter<'a, KeySpan>,
}
impl<'a> Iterator for KeysIter<'a> {
type Item = &'a [u8];
#[inline]
fn next(&mut self) -> Option<&'a [u8]> {
let span = self.spans.next()?;
Some(&self.buf[span.offset..span.offset + span.len])
}
#[inline]
fn size_hint(&self) -> (usize, Option<usize>) {
self.spans.size_hint()
}
}
impl ExactSizeIterator for KeysIter<'_> {}
pub struct ValsBatch<'c> {
key_buf: &'c [u8],
val_buf: &'c [u8],
spans: &'c [KeyValSpan],
pub key_bytes: u64,
pub value_bytes: u64,
}
impl<'c> ValsBatch<'c> {
#[inline]
pub(crate) fn from_parts(
key_buf: &'c [u8],
val_buf: &'c [u8],
spans: &'c [KeyValSpan],
key_bytes: u64,
value_bytes: u64,
) -> Self {
Self {
key_buf,
val_buf,
spans,
key_bytes,
value_bytes,
}
}
#[inline]
pub fn len(&self) -> usize {
self.spans.len()
}
#[inline]
pub fn is_empty(&self) -> bool {
self.spans.is_empty()
}
#[inline]
pub fn get(&self, i: usize) -> Option<(&[u8], &[u8])> {
let span = self.spans.get(i)?;
let k = &self.key_buf[span.key_offset..span.key_offset + span.key_len];
let v = &self.val_buf[span.val_offset..span.val_offset + span.val_len];
Some((k, v))
}
#[inline]
pub fn iter(&self) -> ValsIter<'_> {
ValsIter {
key_buf: self.key_buf,
val_buf: self.val_buf,
spans: self.spans.iter(),
}
}
}
impl<'a, 'c: 'a> IntoIterator for &'a ValsBatch<'c> {
type Item = (&'a [u8], &'a [u8]);
type IntoIter = ValsIter<'a>;
#[inline]
fn into_iter(self) -> Self::IntoIter {
self.iter()
}
}
pub struct ValsIter<'a> {
key_buf: &'a [u8],
val_buf: &'a [u8],
spans: std::slice::Iter<'a, KeyValSpan>,
}
impl<'a> Iterator for ValsIter<'a> {
type Item = (&'a [u8], &'a [u8]);
#[inline]
fn next(&mut self) -> Option<Self::Item> {
let span = self.spans.next()?;
let k = &self.key_buf[span.key_offset..span.key_offset + span.key_len];
let v = &self.val_buf[span.val_offset..span.val_offset + span.val_len];
Some((k, v))
}
#[inline]
fn size_hint(&self) -> (usize, Option<usize>) {
self.spans.size_hint()
}
}
impl ExactSizeIterator for ValsIter<'_> {}
#[cfg(not(target_family = "wasm"))]
pub trait ValVisitor: FnMut(&[u8], &[u8]) -> Result<std::ops::ControlFlow<()>> + Send {}
#[cfg(not(target_family = "wasm"))]
impl<T: FnMut(&[u8], &[u8]) -> Result<std::ops::ControlFlow<()>> + Send> ValVisitor for T {}
#[cfg(target_family = "wasm")]
pub trait ValVisitor: FnMut(&[u8], &[u8]) -> Result<std::ops::ControlFlow<()>> {}
#[cfg(target_family = "wasm")]
impl<T: FnMut(&[u8], &[u8]) -> Result<std::ops::ControlFlow<()>>> ValVisitor for T {}
#[cfg(not(target_family = "wasm"))]
pub trait KeyVisitor: FnMut(&[u8]) -> Result<std::ops::ControlFlow<()>> + Send {}
#[cfg(not(target_family = "wasm"))]
impl<T: FnMut(&[u8]) -> Result<std::ops::ControlFlow<()>> + Send> KeyVisitor for T {}
#[cfg(target_family = "wasm")]
pub trait KeyVisitor: FnMut(&[u8]) -> Result<std::ops::ControlFlow<()>> {}
#[cfg(target_family = "wasm")]
impl<T: FnMut(&[u8]) -> Result<std::ops::ControlFlow<()>>> KeyVisitor for T {}
pub trait ScanCursorKeys: requirements::TransactionRequirements {
fn next_batch<'s>(&'s mut self, limit: u32) -> BoxFut<'s, Result<KeysBatch<'s>>>;
fn for_each<'s>(
&'s mut self,
limit: u32,
f: &'s mut dyn KeyVisitor,
) -> BoxFut<'s, Result<ScanChunkStats>>;
}
pub trait ScanCursorVals: requirements::TransactionRequirements {
fn next_batch<'s>(&'s mut self, limit: u32) -> BoxFut<'s, Result<ValsBatch<'s>>>;
fn for_each<'s>(
&'s mut self,
limit: u32,
f: &'s mut dyn ValVisitor,
) -> BoxFut<'s, Result<ScanChunkStats>>;
}
#[allow(dead_code, reason = "Not used when none of the storage backends are enabled.")]
pub trait Transactable: requirements::TransactionRequirements {
fn kind(&self) -> &'static str;
fn closed(&self) -> bool;
fn writeable(&self) -> bool;
fn cancel(&self) -> BoxFut<'_, Result<()>>;
fn commit(&self) -> BoxFut<'_, Result<()>>;
fn exists(&self, key: Key, version: Option<u64>) -> BoxFut<'_, Result<bool>>;
fn get(&self, key: Key, version: Option<u64>) -> BoxFut<'_, Result<Option<Val>>>;
fn set(&self, key: Key, val: Val) -> BoxFut<'_, Result<()>>;
fn put(&self, key: Key, val: Val) -> BoxFut<'_, Result<()>>;
fn putc(&self, key: Key, val: Val, chk: Option<Val>) -> BoxFut<'_, Result<()>>;
fn del(&self, key: Key) -> BoxFut<'_, Result<()>>;
fn delc(&self, key: Key, chk: Option<Val>) -> BoxFut<'_, Result<()>>;
fn keys(
&self,
rng: Range<Key>,
limit: u32,
skip: u32,
version: Option<u64>,
) -> BoxFut<'_, Result<KeysResult>>;
fn keysr(
&self,
rng: Range<Key>,
limit: u32,
skip: u32,
version: Option<u64>,
) -> BoxFut<'_, Result<KeysResult>>;
fn scan(
&self,
rng: Range<Key>,
limit: u32,
skip: u32,
version: Option<u64>,
) -> BoxFut<'_, Result<ScanResult>>;
fn scanr(
&self,
rng: Range<Key>,
limit: u32,
skip: u32,
version: Option<u64>,
) -> BoxFut<'_, Result<ScanResult>>;
fn open_keys_cursor<'a>(
&'a self,
rng: Range<Key>,
dir: Direction,
skip: u32,
version: Option<u64>,
) -> BoxFut<'a, Result<Box<dyn ScanCursorKeys + 'a>>> {
Box::pin(async move {
if self.closed() {
return Err(Error::TransactionFinished);
}
Ok(Box::new(DefaultKeysCursor::new(self, rng, dir, version, skip))
as Box<dyn ScanCursorKeys + 'a>)
})
}
fn open_vals_cursor<'a>(
&'a self,
rng: Range<Key>,
dir: Direction,
skip: u32,
version: Option<u64>,
) -> BoxFut<'a, Result<Box<dyn ScanCursorVals + 'a>>> {
Box::pin(async move {
if self.closed() {
return Err(Error::TransactionFinished);
}
Ok(Box::new(DefaultValsCursor::new(self, rng, dir, version, skip))
as Box<dyn ScanCursorVals + 'a>)
})
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::api", skip(self), fields(key = key.sprint()))]
fn replace(&self, key: Key, val: Val) -> BoxFut<'_, Result<()>> {
Box::pin(async move { self.set(key, val).await })
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::api", skip(self), fields(key = key.sprint()))]
fn clr(&self, key: Key) -> BoxFut<'_, Result<()>> {
Box::pin(async move { self.del(key).await })
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::api", skip(self), fields(key = key.sprint()))]
fn clrc(&self, key: Key, chk: Option<Val>) -> BoxFut<'_, Result<()>> {
Box::pin(async move { self.delc(key, chk).await })
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::api", skip(self), fields(keys = keys.sprint()))]
fn getm(&self, keys: Vec<Key>, version: Option<u64>) -> BoxFut<'_, Result<GetMultiResult>> {
Box::pin(async move {
if self.closed() {
return Err(Error::TransactionFinished);
}
let mut out = Vec::with_capacity(keys.len());
let mut records = 0u64;
let mut value_bytes = 0u64;
for key in keys {
if let Some(val) = self.get(key, version).await? {
records += 1;
value_bytes += val.len() as u64;
out.push(Some(val));
} else {
out.push(None);
}
}
Ok(GetMultiResult {
values: out,
records,
value_bytes,
})
})
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::api", skip(self), fields(key = key.sprint()))]
fn getp(&self, key: Key, version: Option<u64>) -> BoxFut<'_, Result<ScanResult>> {
Box::pin(async move {
if self.closed() {
return Err(Error::TransactionFinished);
}
let range = util::to_prefix_range(&key)?;
self.getr(range, version).await
})
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::api", skip(self), fields(rng = rng.sprint()))]
fn getr(&self, rng: Range<Key>, version: Option<u64>) -> BoxFut<'_, Result<ScanResult>> {
Box::pin(async move {
if self.closed() {
return Err(Error::TransactionFinished);
}
let mut out: Vec<(Key, Val)> = vec![];
let mut key_bytes = 0u64;
let mut value_bytes = 0u64;
let mut next = Some(rng);
while let Some(rng) = next {
let res = self.batch_keys_vals(rng, NORMAL_BATCH_SIZE, version).await?;
next = res.next;
for (k, v) in res.result {
key_bytes += k.len() as u64;
value_bytes += v.len() as u64;
out.push((k, v));
}
}
Ok(ScanResult {
values: out,
key_bytes,
value_bytes,
})
})
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::api", skip(self), fields(key = key.sprint()))]
fn delp(&self, key: Key) -> BoxFut<'_, Result<()>> {
Box::pin(async move {
if self.closed() {
return Err(Error::TransactionFinished);
}
if !self.writeable() {
return Err(Error::TransactionReadonly);
}
let range = util::to_prefix_range(&key)?;
self.delr(range).await
})
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::api", skip(self), fields(rng = rng.sprint()))]
fn delr(&self, rng: Range<Key>) -> BoxFut<'_, Result<()>> {
Box::pin(async move {
if self.closed() {
return Err(Error::TransactionFinished);
}
if !self.writeable() {
return Err(Error::TransactionReadonly);
}
let mut next = Some(rng);
while let Some(rng) = next {
let res = self.batch_keys(rng, NORMAL_BATCH_SIZE, None).await?;
next = res.next;
for k in res.result {
self.del(k).await?;
}
}
Ok(())
})
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::api", skip(self), fields(key = key.sprint()))]
fn clrp(&self, key: Key) -> BoxFut<'_, Result<()>> {
Box::pin(async move {
if self.closed() {
return Err(Error::TransactionFinished);
}
if !self.writeable() {
return Err(Error::TransactionReadonly);
}
let range = util::to_prefix_range(&key)?;
self.clrr(range).await
})
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::api", skip(self), fields(rng = rng.sprint()))]
fn clrr(&self, rng: Range<Key>) -> BoxFut<'_, Result<()>> {
Box::pin(async move {
if self.closed() {
return Err(Error::TransactionFinished);
}
if !self.writeable() {
return Err(Error::TransactionReadonly);
}
let mut next = Some(rng);
while let Some(rng) = next {
let res = self.batch_keys(rng, NORMAL_BATCH_SIZE, None).await?;
next = res.next;
for k in res.result {
self.clr(k).await?;
}
}
Ok(())
})
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::api", skip(self), fields(rng = rng.sprint()))]
fn count(&self, rng: Range<Key>, version: Option<u64>) -> BoxFut<'_, Result<usize>> {
Box::pin(async move {
if self.closed() {
return Err(Error::TransactionFinished);
}
let mut len = 0;
let mut next = Some(rng);
while let Some(rng) = next {
let res = self.batch_keys(rng, COUNT_BATCH_SIZE, version).await?;
next = res.next;
len += res.result.len();
}
Ok(len)
})
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::api", skip(self), fields(rng = rng.sprint()))]
fn batch_keys(
&self,
rng: Range<Key>,
batch: u32,
version: Option<u64>,
) -> BoxFut<'_, Result<Batch<Key>>> {
Box::pin(async move {
if self.closed() {
return Err(Error::TransactionFinished);
}
let end = rng.end.clone();
let res = self.keys(rng, batch, 0, version).await?.keys;
if res.len() < batch as usize && batch > 0 {
Ok(Batch::<Key>::new(None, res))
} else {
match res.last() {
Some(k) => {
let mut k = k.clone();
util::advance_key(&mut k);
Ok(Batch::<Key>::new(
Some(Range {
start: k,
end,
}),
res,
))
}
None => Ok(Batch::<Key>::new(None, res)),
}
}
})
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::api", skip(self), fields(rng = rng.sprint()))]
fn batch_keys_vals(
&self,
rng: Range<Key>,
batch: u32,
version: Option<u64>,
) -> BoxFut<'_, Result<Batch<(Key, Val)>>> {
Box::pin(async move {
if self.closed() {
return Err(Error::TransactionFinished);
}
let end = rng.end.clone();
let res = self.scan(rng, batch, 0, version).await?.values;
if res.len() < batch as usize && batch > 0 {
Ok(Batch::<(Key, Val)>::new(None, res))
} else {
match res.last() {
Some((k, _)) => {
let mut k = k.clone();
util::advance_key(&mut k);
Ok(Batch::<(Key, Val)>::new(
Some(Range {
start: k,
end,
}),
res,
))
}
None => Ok(Batch::<(Key, Val)>::new(None, res)),
}
}
})
}
fn new_save_point(&self) -> BoxFut<'_, Result<()>>;
fn release_last_save_point(&self) -> BoxFut<'_, Result<()>>;
fn rollback_to_save_point(&self) -> BoxFut<'_, Result<()>>;
fn timestamp(&self) -> BoxFut<'_, Result<BoxTimeStamp>> {
Box::pin(async move {
if cfg!(test) {
Ok(BoxTimeStamp::new(IncTimeStamp::next()))
} else {
Ok(BoxTimeStamp::new(HlcTimeStamp::next()))
}
})
}
fn safe_timestamp(&self) -> BoxFut<'_, Result<BoxTimeStamp>> {
self.timestamp()
}
fn timestamp_impl(&self) -> BoxTimeStampImpl {
if cfg!(test) {
Box::new(IncTimeStampImpl)
} else {
Box::new(HlcTimeStampImpl)
}
}
fn compact(&self, _range: Option<Range<Key>>) -> BoxFut<'_, anyhow::Result<()>> {
Box::pin(async move { bail!(Error::CompactionNotSupported) })
}
}