#![cfg(feature = "kv-rocksdb")]
mod cnf;
use super::KeyEncode;
use crate::err::Error;
use crate::key::debug::Sprintable;
use crate::kvs::ds::{Metric, Metrics};
use crate::kvs::{Check, Key, Val};
use rocksdb::{
properties, BlockBasedOptions, Cache, DBCompactionStyle, DBCompressionType, Env, FlushOptions,
LogLevel, OptimisticTransactionDB, OptimisticTransactionOptions, Options, ReadOptions,
SstFileManager, WriteOptions,
};
use std::fmt::Debug;
use std::ops::Range;
use std::pin::Pin;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use std::thread;
use std::time::Duration;
const TARGET: &str = "surrealdb::core::kvs::rocksdb";
use crate::kvs::rocksdb::cnf::ROCKSDB_SST_MAX_ALLOWED_SPACE_USAGE;
pub struct Datastore {
db: Pin<Arc<OptimisticTransactionDB>>,
disk_space_manager: Option<DiskSpaceManager>,
}
#[derive(Clone)]
struct DiskSpaceManager {
sst_file_manager: Arc<SstFileManager>,
read_and_deletion_limit: u64,
limit_80: u64,
warn_80_percent_logged: Arc<AtomicBool>,
}
pub struct Transaction {
done: bool,
write: bool,
check: Check,
inner: Option<rocksdb::Transaction<'static, OptimisticTransactionDB>>,
ro: ReadOptions,
db: Pin<Arc<OptimisticTransactionDB>>,
deletion_only: bool,
contains_only_deletions: Option<bool>,
disk_space_manager: Option<DiskSpaceManager>,
}
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 DiskSpaceManager {
fn new(limit: u64, opts: &mut Options) -> Result<Self, Error> {
let env = Env::new()?;
let sst_file_manager = SstFileManager::new(&env)?;
sst_file_manager.set_max_allowed_space_usage(0);
opts.set_sst_file_manager(&sst_file_manager);
Ok(Self {
sst_file_manager: Arc::new(sst_file_manager),
read_and_deletion_limit: limit,
limit_80: (limit as f64 * 0.8) as u64,
warn_80_percent_logged: Arc::new(AtomicBool::new(false)),
})
}
fn is_deletion_only(&self) -> bool {
let current_size = self.sst_file_manager.get_total_size();
if current_size < self.limit_80 {
self.warn_80_percent_logged.store(false, Ordering::Relaxed);
return false;
}
if self
.warn_80_percent_logged
.compare_exchange(false, true, Ordering::Relaxed, Ordering::Relaxed)
.is_ok()
{
warn!(target: TARGET, "SST file space usage is at 80% of the limit ({})", current_size);
}
if current_size < self.read_and_deletion_limit {
return false;
}
warn!(
target: TARGET,
"Transitioning to read-and-deletion-only mode due to primary limit ({}) being reached",
current_size
);
true
}
}
impl Datastore {
pub(crate) async fn new(path: &str) -> Result<Datastore, Error> {
let mut opts = Options::default();
opts.set_use_fsync(false);
opts.create_if_missing(true);
opts.create_missing_column_families(true);
info!(target: TARGET, "Background thread count: {}", *cnf::ROCKSDB_THREAD_COUNT);
opts.increase_parallelism(*cnf::ROCKSDB_THREAD_COUNT);
info!(target: TARGET, "Maximum background jobs count: {}", *cnf::ROCKSDB_JOBS_COUNT);
opts.set_max_background_jobs(*cnf::ROCKSDB_JOBS_COUNT);
info!(target: TARGET, "Maximum number of open files: {}", *cnf::ROCKSDB_MAX_OPEN_FILES);
opts.set_max_open_files(*cnf::ROCKSDB_MAX_OPEN_FILES);
info!(target: TARGET, "Number of log files to keep: {}", *cnf::ROCKSDB_KEEP_LOG_FILE_NUM);
opts.set_keep_log_file_num(*cnf::ROCKSDB_KEEP_LOG_FILE_NUM);
info!(target: TARGET, "Maximum write buffers: {}", *cnf::ROCKSDB_MAX_WRITE_BUFFER_NUMBER);
opts.set_max_write_buffer_number(*cnf::ROCKSDB_MAX_WRITE_BUFFER_NUMBER);
info!(target: TARGET, "Write buffer size: {}", *cnf::ROCKSDB_WRITE_BUFFER_SIZE);
opts.set_write_buffer_size(*cnf::ROCKSDB_WRITE_BUFFER_SIZE);
info!(target: TARGET, "Target file size for compaction: {}", *cnf::ROCKSDB_TARGET_FILE_SIZE_BASE);
opts.set_target_file_size_base(*cnf::ROCKSDB_TARGET_FILE_SIZE_BASE);
info!(target: TARGET, "Target file size compaction multiplier: {}", *cnf::ROCKSDB_TARGET_FILE_SIZE_MULTIPLIER);
opts.set_target_file_size_multiplier(*cnf::ROCKSDB_TARGET_FILE_SIZE_MULTIPLIER);
info!(target: TARGET, "Minimum write buffers to merge: {}", *cnf::ROCKSDB_MIN_WRITE_BUFFER_NUMBER_TO_MERGE);
opts.set_min_write_buffer_number_to_merge(*cnf::ROCKSDB_MIN_WRITE_BUFFER_NUMBER_TO_MERGE);
info!(target: TARGET, "Number of files to trigger compaction: {}", *cnf::ROCKSDB_FILE_COMPACTION_TRIGGER);
opts.set_level_zero_file_num_compaction_trigger(*cnf::ROCKSDB_FILE_COMPACTION_TRIGGER);
info!(target: TARGET, "Compaction readahead size: {}", *cnf::ROCKSDB_COMPACTION_READAHEAD_SIZE);
opts.set_compaction_readahead_size(*cnf::ROCKSDB_COMPACTION_READAHEAD_SIZE);
info!(target: TARGET, "Maximum concurrent subcompactions: {}", *cnf::ROCKSDB_MAX_CONCURRENT_SUBCOMPACTIONS);
opts.set_max_subcompactions(*cnf::ROCKSDB_MAX_CONCURRENT_SUBCOMPACTIONS);
info!(target: TARGET, "Use separate thread queues: {}", *cnf::ROCKSDB_ENABLE_PIPELINED_WRITES);
opts.set_enable_pipelined_write(*cnf::ROCKSDB_ENABLE_PIPELINED_WRITES);
info!(target: TARGET, "Enable separation of keys and values: {}", *cnf::ROCKSDB_ENABLE_BLOB_FILES);
opts.set_enable_blob_files(*cnf::ROCKSDB_ENABLE_BLOB_FILES);
info!(target: TARGET, "Minimum blob value size: {}", *cnf::ROCKSDB_MIN_BLOB_SIZE);
opts.set_min_blob_size(*cnf::ROCKSDB_MIN_BLOB_SIZE);
info!(target: TARGET, "Target blob file size: {}", *cnf::ROCKSDB_BLOB_FILE_SIZE);
opts.set_blob_file_size(*cnf::ROCKSDB_BLOB_FILE_SIZE);
if let Some(c) = cnf::ROCKSDB_BLOB_COMPRESSION_TYPE.as_ref() {
info!(target: TARGET, "Blob compression type: {c}");
opts.set_blob_compression_type(match c.as_str() {
"none" => DBCompressionType::None,
"snappy" => DBCompressionType::Snappy,
"lz4" => DBCompressionType::Lz4,
"zstd" => DBCompressionType::Zstd,
l => {
return Err(Error::Ds(format!("Invalid compression type: {l}")));
}
});
}
info!(target: TARGET, "Enable blob garbage collection: {}", *cnf::ROCKSDB_ENABLE_BLOB_GC);
opts.set_enable_blob_gc(*cnf::ROCKSDB_ENABLE_BLOB_GC);
info!(target: TARGET, "Blob GC age cutoff: {}", *cnf::ROCKSDB_BLOB_GC_AGE_CUTOFF);
opts.set_blob_gc_age_cutoff(*cnf::ROCKSDB_BLOB_GC_AGE_CUTOFF);
info!(target: TARGET, "Blob GC force threshold: {}", *cnf::ROCKSDB_BLOB_GC_FORCE_THRESHOLD);
opts.set_blob_gc_force_threshold(*cnf::ROCKSDB_BLOB_GC_FORCE_THRESHOLD);
info!(target: TARGET, "Blob compaction readahead size: {}", *cnf::ROCKSDB_BLOB_COMPACTION_READAHEAD_SIZE);
opts.set_blob_compaction_readahead_size(
(*cnf::ROCKSDB_BLOB_COMPACTION_READAHEAD_SIZE) as u64,
);
info!(target: TARGET, "Write-ahead-log file size limit: {}MB", *cnf::ROCKSDB_WAL_SIZE_LIMIT);
opts.set_wal_size_limit_mb(*cnf::ROCKSDB_WAL_SIZE_LIMIT);
info!(target: TARGET, "Allow concurrent memtable writes: true");
opts.set_allow_concurrent_memtable_write(true);
info!(target: TARGET, "Avoid unnecessary blocking IO: true");
opts.set_avoid_unnecessary_blocking_io(true);
info!(target: TARGET, "Allow adaptive write thread yielding: true");
opts.set_enable_write_thread_adaptive_yield(true);
info!(target: TARGET, "Wait for disk sync acknowledgement: {}", *cnf::SYNC_DATA);
info!(target: TARGET, "Block cache size: {}", *cnf::ROCKSDB_BLOCK_CACHE_SIZE);
let cache = Cache::new_lru_cache(*cnf::ROCKSDB_BLOCK_CACHE_SIZE);
let mut block_opts = BlockBasedOptions::default();
block_opts.set_pin_l0_filter_and_index_blocks_in_cache(true);
block_opts.set_pin_top_level_index_and_filter(true);
block_opts.set_bloom_filter(10.0, false);
block_opts.set_block_size(*cnf::ROCKSDB_BLOCK_SIZE);
block_opts.set_block_cache(&cache);
opts.set_block_based_table_factory(&block_opts);
opts.set_blob_cache(&cache);
opts.set_row_cache(&cache);
info!(target: TARGET, "Enable memory-mapped reads: {}", *cnf::ROCKSDB_ENABLE_MEMORY_MAPPED_READS);
opts.set_allow_mmap_reads(*cnf::ROCKSDB_ENABLE_MEMORY_MAPPED_READS);
info!(target: TARGET, "Enable memory-mapped writes: {}", *cnf::ROCKSDB_ENABLE_MEMORY_MAPPED_WRITES);
opts.set_allow_mmap_writes(*cnf::ROCKSDB_ENABLE_MEMORY_MAPPED_WRITES);
info!(target: TARGET, "Setting delete compaction factory: {} / {} ({})",
*cnf::ROCKSDB_DELETION_FACTORY_WINDOW_SIZE,
*cnf::ROCKSDB_DELETION_FACTORY_DELETE_COUNT,
*cnf::ROCKSDB_DELETION_FACTORY_RATIO,
);
opts.add_compact_on_deletion_collector_factory(
*cnf::ROCKSDB_DELETION_FACTORY_WINDOW_SIZE,
*cnf::ROCKSDB_DELETION_FACTORY_DELETE_COUNT,
*cnf::ROCKSDB_DELETION_FACTORY_RATIO,
);
info!(target: TARGET, "Setting compaction style: {}", *cnf::ROCKSDB_COMPACTION_STYLE);
opts.set_compaction_style(
match cnf::ROCKSDB_COMPACTION_STYLE.to_ascii_lowercase().as_str() {
"universal" => DBCompactionStyle::Universal,
_ => DBCompactionStyle::Level,
},
);
info!(target: TARGET, "Setting compression level");
opts.set_compression_per_level(&[
DBCompressionType::None,
DBCompressionType::None,
DBCompressionType::Snappy,
DBCompressionType::Snappy,
DBCompressionType::Snappy,
]);
info!(target: TARGET, "Setting storage engine log level: {}", *cnf::ROCKSDB_STORAGE_LOG_LEVEL);
opts.set_log_level(match cnf::ROCKSDB_STORAGE_LOG_LEVEL.to_ascii_lowercase().as_str() {
"debug" => LogLevel::Debug,
"info" => LogLevel::Info,
"warn" => LogLevel::Warn,
"error" => LogLevel::Error,
"fatal" => LogLevel::Fatal,
l => {
return Err(Error::Ds(format!("Invalid storage engine log level specified: {l}")));
}
});
let disk_space_manager = if *ROCKSDB_SST_MAX_ALLOWED_SPACE_USAGE > 0 {
Some(DiskSpaceManager::new(*ROCKSDB_SST_MAX_ALLOWED_SPACE_USAGE, &mut opts)?)
} else {
None
};
let db = match *cnf::ROCKSDB_BACKGROUND_FLUSH {
false => {
info!(target: TARGET, "Background write-ahead-log flushing: disabled");
opts.set_manual_wal_flush(false);
Arc::pin(OptimisticTransactionDB::open(&opts, path)?)
}
true => {
info!(target: TARGET, "Background write-ahead-log flushing: enabled every {}ms", *cnf::ROCKSDB_BACKGROUND_FLUSH_INTERVAL);
opts.set_manual_wal_flush(true);
let db = Arc::pin(OptimisticTransactionDB::open(&opts, path)?);
let dbc = db.clone();
thread::spawn(move || loop {
let wait = *cnf::ROCKSDB_BACKGROUND_FLUSH_INTERVAL;
thread::sleep(Duration::from_millis(wait));
if let Err(err) = dbc.flush_wal(*cnf::SYNC_DATA) {
error!("Failed to flush WAL: {err}");
}
});
db
}
};
Ok(Datastore {
db,
disk_space_manager,
})
}
const BLOCK_CACHE_USAGE: &str = "rocksdb.block_cache_usage";
const BLOCK_CACHE_PINNED_USAGE: &str = "rocksdb.block_cache_pinned_usage";
const ESTIMATE_TABLE_READERS_MEM: &str = "rocksdb.estimate_table_readers_mem";
const CUR_SIZE_ALL_MEM_TABLES: &str = "rocksdb.cur_size_all_mem_tables";
pub(crate) fn register_metrics(&self) -> Metrics {
Metrics {
name: "surrealdb.rocksdb",
u64_metrics: vec![
Metric {
name: Self::BLOCK_CACHE_USAGE,
description: "Returns the memory size for the entries residing in block cache.",
},
Metric {
name: Self::BLOCK_CACHE_PINNED_USAGE,
description: "Returns the memory size for the entries being pinned.",
},
Metric {
name: Self::ESTIMATE_TABLE_READERS_MEM,
description: "Returns estimated memory used for reading SST tables, excluding memory used in block cache (e.g., filter and index blocks).",
},
Metric {
name: Self::CUR_SIZE_ALL_MEM_TABLES,
description: "Returns approximate size of active and unflushed immutable memtables",
},
],
}
}
pub(crate) fn collect_u64_metric(&self, metric: &str) -> Option<u64> {
let metric = match metric {
Self::BLOCK_CACHE_USAGE => Some(properties::BLOCK_CACHE_USAGE),
Self::BLOCK_CACHE_PINNED_USAGE => Some(properties::BLOCK_CACHE_PINNED_USAGE),
Self::ESTIMATE_TABLE_READERS_MEM => Some(properties::ESTIMATE_TABLE_READERS_MEM),
Self::CUR_SIZE_ALL_MEM_TABLES => Some(properties::CUR_SIZE_ALL_MEM_TABLES),
_ => None,
};
metric.map(|metric| {
self.db.property_int_value(metric).unwrap_or_default().unwrap_or_default()
})
}
pub(crate) async fn shutdown(&self) -> Result<(), Error> {
let mut opts = FlushOptions::default();
opts.set_wait(true);
if let Err(e) = self.db.flush_wal(true) {
error!("An error occured flushing the WAL buffer to disk: {e}");
}
if let Err(e) = self.db.flush_opt(&opts) {
error!("An error occured flushing memtables to SST files: {e}");
}
Ok(())
}
fn is_deletion_only(&self) -> bool {
self.disk_space_manager.as_ref().map(|dsm| dsm.is_deletion_only()).unwrap_or(false)
}
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(*cnf::SYNC_DATA);
let inner = self.db.transaction_opt(&wo, &to);
let inner = unsafe {
std::mem::transmute::<
rocksdb::Transaction<'_, OptimisticTransactionDB>,
rocksdb::Transaction<'static, OptimisticTransactionDB>,
>(inner)
};
let mut ro = ReadOptions::default();
ro.set_snapshot(&inner.snapshot());
ro.set_async_io(true);
ro.fill_cache(true);
#[cfg(not(debug_assertions))]
let check = Check::Warn;
#[cfg(debug_assertions)]
let check = Check::Error;
Ok(Transaction {
done: false,
write,
check,
inner: Some(inner),
ro,
db: self.db.clone(),
deletion_only: self.is_deletion_only(),
contains_only_deletions: None,
disk_space_manager: self.disk_space_manager.clone(),
})
}
}
impl Transaction {
fn ensure_write(&mut self, version: Option<u64>) -> Result<(), Error> {
if self.deletion_only {
return Err(Error::DbReadAndDeleteOnly);
}
if version.is_some() {
return Err(Error::UnsupportedVersionedQueries);
}
if self.done {
return Err(Error::TxFinished);
}
if !self.write {
return Err(Error::TxReadonly);
}
self.contains_only_deletions = Some(false);
Ok(())
}
fn ensure_deletion(&mut self, version: Option<u64>) -> Result<(), Error> {
if version.is_some() {
return Err(Error::UnsupportedVersionedQueries);
}
if self.done {
return Err(Error::TxFinished);
}
if !self.write {
return Err(Error::TxReadonly);
}
if self.contains_only_deletions.is_none() {
self.contains_only_deletions = Some(true);
}
Ok(())
}
fn ensure_read(&self, version: Option<u64>) -> Result<(), Error> {
if version.is_some() {
return Err(Error::UnsupportedVersionedQueries);
}
if self.done {
return Err(Error::TxFinished);
}
Ok(())
}
}
impl super::api::Transaction for Transaction {
fn supports_reverse_scan(&self) -> bool {
true
}
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.as_ref().unwrap().rollback()?;
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;
if let Some(disk_space_manager) = self.disk_space_manager.as_ref() {
if disk_space_manager.is_deletion_only() && self.contains_only_deletions == Some(false)
{
return Err(Error::DbReadAndDeleteOnly);
}
}
self.inner.take().unwrap().commit()?;
if self.deletion_only {
self.db.compact_range::<&[u8], &[u8]>(None, None);
}
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,
{
self.ensure_read(version)?;
let key = key.encode_owned()?;
let res = self.inner.as_ref().unwrap().get_pinned_opt(key, &self.ro)?.is_some();
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,
{
self.ensure_read(version)?;
let key = key.encode_owned()?;
let res = self.inner.as_ref().unwrap().get_opt(key, &self.ro)?;
Ok(res)
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::api", skip(self), fields(keys = keys.sprint()))]
async fn getm<K>(&mut self, keys: Vec<K>) -> Result<Vec<Option<Val>>, Error>
where
K: KeyEncode + Sprintable + Debug,
{
self.ensure_read(None)?;
let keys: Vec<Key> = keys.into_iter().map(K::encode_owned).collect::<Result<_, _>>()?;
let res = self.inner.as_ref().unwrap().multi_get_opt(keys, &self.ro);
let res = res.into_iter().collect::<Result<_, _>>()?;
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,
{
self.ensure_write(version)?;
let key = key.encode_owned()?;
let val = val.into();
self.inner.as_ref().unwrap().put(key, val)?;
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,
{
self.ensure_write(version)?;
let key = key.encode_owned()?;
let val = val.into();
match self.inner.as_ref().unwrap().get_pinned_opt(&key, &self.ro)? {
None => self.inner.as_ref().unwrap().put(key, val)?,
_ => return Err(Error::TxKeyAlreadyExists),
};
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,
{
self.ensure_write(None)?;
let key = key.encode_owned()?;
let val = val.into();
let chk = chk.map(Into::into);
match (self.inner.as_ref().unwrap().get_pinned_opt(&key, &self.ro)?, chk) {
(Some(v), Some(w)) if v.eq(&w) => self.inner.as_ref().unwrap().put(key, val)?,
(None, None) => self.inner.as_ref().unwrap().put(key, val)?,
_ => return Err(Error::TxConditionNotMet),
};
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,
{
self.ensure_deletion(None)?;
let key = key.encode_owned()?;
self.inner.as_ref().unwrap().delete(key)?;
Ok(())
}
#[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,
{
self.ensure_deletion(None)?;
let key = key.encode_owned()?;
let chk = chk.map(Into::into);
match (self.inner.as_ref().unwrap().get_pinned_opt(&key, &self.ro)?, chk) {
(Some(v), Some(w)) if v.eq(&w) => self.inner.as_ref().unwrap().delete(key)?,
(None, None) => self.inner.as_ref().unwrap().delete(key)?,
_ => return Err(Error::TxConditionNotMet),
};
Ok(())
}
#[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,
{
self.ensure_read(version)?;
let rng: Range<Key> = Range {
start: rng.start.encode_owned()?,
end: rng.end.encode_owned()?,
};
let res = affinitypool::spawn_local(move || {
let mut res = vec![];
let beg = rng.start.as_slice();
let end = rng.end.as_slice();
let mut ro = ReadOptions::default();
ro.set_snapshot(&self.inner.as_ref().unwrap().snapshot());
ro.set_iterate_lower_bound(beg);
ro.set_iterate_upper_bound(end);
ro.set_async_io(true);
ro.fill_cache(true);
let mut iter = self.inner.as_ref().unwrap().raw_iterator_opt(ro);
iter.seek(&rng.start);
while res.len() < limit as usize {
if let Some(k) = iter.key() {
if k >= beg && k < end {
res.push(k.to_vec());
iter.next();
continue;
}
}
break;
}
drop(iter);
res
})
.await;
Ok(res)
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::api", skip(self), fields(rng = rng.sprint()))]
async fn keysr<K>(
&mut self,
rng: Range<K>,
limit: u32,
version: Option<u64>,
) -> Result<Vec<Key>, Error>
where
K: KeyEncode + Sprintable + Debug,
{
self.ensure_read(version)?;
let rng: Range<Key> = Range {
start: rng.start.encode_owned()?,
end: rng.end.encode_owned()?,
};
let inner = self.inner.as_ref().unwrap();
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());
ro.set_iterate_lower_bound(beg);
ro.set_iterate_upper_bound(end);
ro.set_async_io(true);
ro.fill_cache(true);
let mut iter = inner.raw_iterator_opt(ro);
iter.seek_for_prev(&rng.end);
while res.len() < limit as usize {
if let Some(k) = iter.key() {
if k >= beg && k < end {
res.push(k.to_vec());
iter.prev();
continue;
}
}
break;
}
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,
{
self.ensure_read(version)?;
let rng: Range<Key> = Range {
start: rng.start.encode_owned()?,
end: rng.end.encode_owned()?,
};
let res = affinitypool::spawn_local(move || {
let mut res = vec![];
let beg = rng.start.as_slice();
let end = rng.end.as_slice();
let mut ro = ReadOptions::default();
ro.set_snapshot(&self.inner.as_ref().unwrap().snapshot());
ro.set_iterate_lower_bound(beg);
ro.set_iterate_upper_bound(end);
ro.set_async_io(true);
ro.fill_cache(true);
let mut iter = self.inner.as_ref().unwrap().raw_iterator_opt(ro);
iter.seek(&rng.start);
while res.len() < limit as usize {
if let Some((k, v)) = iter.item() {
if k >= beg && k < end {
res.push((k.to_vec(), v.to_vec()));
iter.next();
continue;
}
}
break;
}
drop(iter);
res
})
.await;
Ok(res)
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::api", skip(self), fields(rng = rng.sprint()))]
async fn scanr<K>(
&mut self,
rng: Range<K>,
limit: u32,
version: Option<u64>,
) -> Result<Vec<(Key, Val)>, Error>
where
K: KeyEncode + Sprintable + Debug,
{
self.ensure_read(version)?;
let rng: Range<Key> = Range {
start: rng.start.encode_owned()?,
end: rng.end.encode_owned()?,
};
let inner = self.inner.as_ref().unwrap();
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());
ro.set_iterate_lower_bound(beg);
ro.set_iterate_upper_bound(end);
ro.set_async_io(true);
ro.fill_cache(true);
let mut iter = inner.raw_iterator_opt(ro);
iter.seek_for_prev(&rng.end);
while res.len() < limit as usize {
if let Some((k, v)) = iter.item() {
if k >= beg && k < end {
res.push((k.to_vec(), v.to_vec()));
iter.prev();
continue;
}
}
break;
}
Ok(res)
}
}
impl Transaction {
pub(crate) fn new_save_point(&mut self) {
let inner = self.inner.as_ref().unwrap();
inner.set_savepoint();
}
pub(crate) async fn rollback_to_save_point(&mut self) -> Result<(), Error> {
let inner = self.inner.as_ref().unwrap();
inner.rollback_to_savepoint()?;
Ok(())
}
pub(crate) fn release_last_save_point(&mut self) -> Result<(), Error> {
Ok(())
}
}