#![allow(clippy::clone_on_ref_ptr)]
#![allow(private_bounds, private_interfaces)]
use std::any::Any;
use std::collections::{HashMap, HashSet};
use std::fmt::Debug;
use std::ops::{Deref, Range};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, OnceLock};
use std::time::Duration;
use anyhow::Result;
use chrono::Utc;
use futures::future::try_join_all;
use tokio::sync::{Mutex, Notify};
use tokio::time::sleep;
use tracing::Instrument;
use uuid::Uuid;
use web_time::Instant;
use super::api::{
KeyVisitor, KeysBatch, ScanChunkStats, ScanCursorKeys, ScanCursorVals, ValVisitor, ValsBatch,
};
use super::batch::Batch;
use super::{Key, LockType, TransactionFactory, TransactionType, Val, util};
use crate::catalog::providers::{
ApiProvider, AuthorisationProvider, BoxProviderFut, BucketProvider, CatalogProvider,
DatabaseProvider, NamespaceProvider, NodeProvider, RootProvider, TableProvider, UserProvider,
};
use crate::catalog::{
self, ApiDefinition, ConfigDefinition, DatabaseDefinition, DatabaseId, DefaultConfig, IndexId,
NamespaceDefinition, NamespaceId, Record, TableDefinition, TableId,
};
use crate::cf::Changefeed;
use crate::cnf::CommonConfig;
use crate::ctx::Context;
use crate::dbs::node::Node;
use crate::doc::CursorRecord;
use crate::err::Error;
use crate::idx::IndexKeyBase;
use crate::idx::planner::ScanDirection;
use crate::key::database::sq::Sq;
use crate::key::index::all as index_all;
use crate::key::table::bg::Bg;
use crate::key::table::br::Br;
use crate::key::table::bs::Bs;
use crate::key::table::ix as table_ix;
use crate::kvs::cache::tx::TransactionCache;
use crate::kvs::index::{
BuildGeneration, BuildTicket, BuildTicketMutationSeq, IndexBuildPhase, IndexBuildReportStatus,
IndexBuildState, IndexBuilder,
};
use crate::kvs::sequences::Sequences;
#[cfg(test)]
use crate::kvs::testing::{
NonRetryableErrorSite, RetryableConflictSite, maybe_inject_non_retryable_error,
maybe_inject_retryable_conflict,
};
use crate::kvs::{
BoxTimeStamp, BoxTimeStampImpl, Direction, Error as KvsError, KVKey, KVValue, Transactor,
cache, is_retryable_transaction_conflict,
};
use crate::lq::writer::LiveEventBuffer;
use crate::observe::{
ExecutionObserver, Outcome, TenantIdentity, TransactionEvent, TransactionEventSafe,
TransactionMetrics,
};
use crate::val::{RecordId, RecordIdKey, TableName};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum CachePolicy {
ReadWrite,
ReadOnly,
}
pub struct Transaction {
local: bool,
started_at: Instant,
observer: Arc<dyn ExecutionObserver>,
metrics: TransactionMetrics,
tenant_identity: OnceLock<Arc<TenantIdentity>>,
tr: Transactor,
cache: TransactionCache,
sequences: Sequences,
changefeed: OnceLock<Changefeed>,
live_events: OnceLock<LiveEventBuffer>,
async_event_trigger: Arc<Notify>,
trigger_async_event: AtomicBool,
pending_index_build_reservations: Mutex<Vec<IndexBuildReservationRelease>>,
cached_index_build_reservations:
Mutex<HashMap<CachedIndexBuildReservationKey, CachedIndexBuildReservation>>,
pending_index_builder_aborts: Mutex<Vec<PendingIndexBuilderAbort>>,
pending_uncommitted_index_builds: Mutex<Vec<PendingUncommittedIndexBuild>>,
}
const INDEX_BUILD_RESERVATION_RELEASE_RETRY_SLEEP: Duration = Duration::from_millis(100);
#[derive(Clone, Debug, Eq, Hash, PartialEq)]
pub(crate) struct CachedIndexBuildReservationKey {
pub(crate) ns: NamespaceId,
pub(crate) db: DatabaseId,
pub(crate) tb: TableName,
pub(crate) ix: IndexId,
}
pub(crate) struct CachedIndexBuildReservation {
pub(crate) generation: BuildGeneration,
pub(crate) ticket: BuildTicket,
pub(crate) initial_complete: bool,
pub(crate) next_mutation_seq: BuildTicketMutationSeq,
}
#[derive(Clone, Copy, Debug)]
pub(crate) enum CachedIndexBuildReservationLookup {
FirstUse {
generation: BuildGeneration,
ticket: BuildTicket,
mutation_seq: BuildTicketMutationSeq,
initial_complete: bool,
},
Reused {
generation: BuildGeneration,
ticket: BuildTicket,
mutation_seq: BuildTicketMutationSeq,
initial_complete: bool,
},
}
struct PendingIndexBuilderAbort {
builder: IndexBuilder,
ns: NamespaceId,
db: DatabaseId,
tb: TableName,
ix: IndexId,
}
impl PendingIndexBuilderAbort {
async fn abort(self) {
if let Err(err) = self.builder.remove_index(self.ns, self.db, &self.tb, self.ix).await {
tracing::warn!(
target: "surrealdb::core::kvs::tx",
"failed to abort local index builder after committed schema retirement: {err}"
);
}
}
}
struct PendingUncommittedIndexBuild {
builder: IndexBuilder,
tf: TransactionFactory,
sequences: Sequences,
ns: NamespaceId,
db: DatabaseId,
tb: TableName,
ix: IndexId,
}
impl PendingUncommittedIndexBuild {
async fn cleanup_once(&self) -> Result<()> {
if let Err(err) = self.builder.remove_index(self.ns, self.db, &self.tb, self.ix).await {
tracing::warn!(
target: "surrealdb::core::kvs::tx",
"failed to abort uncommitted local index builder during rollback cleanup: {err}"
);
}
let tx = self
.tf
.transaction(TransactionType::Write, LockType::Optimistic, self.sequences.clone())
.await?;
let ikb = IndexKeyBase::new(self.ns, self.db, self.tb.clone(), self.ix);
let index_prefix = index_all::new(self.ns, self.db, &self.tb, self.ix).encode_key()?;
let result: Result<()> = async {
tx.tr.del(ikb.new_bs_key().encode_key()?).await.map_err(Error::from)?;
tx.tr.delr(ikb.new_bg_all_generations_range()?).await.map_err(Error::from)?;
tx.tr.delr(ikb.new_bp_all_generations_range()?).await.map_err(Error::from)?;
tx.tr.delr(ikb.new_br_all_generations_range()?).await.map_err(Error::from)?;
tx.tr.delp(index_prefix).await.map_err(Error::from)?;
tx.tr.commit().await.map_err(Error::from)?;
Ok(())
}
.await;
if let Err(err) = result {
let _ = tx.tr.cancel().await;
return Err(err);
}
Ok(())
}
async fn cleanup(self) -> Result<()> {
loop {
match self.cleanup_once().await {
Ok(()) => return Ok(()),
Err(err) if is_retryable_transaction_conflict(&err) => {
tracing::debug!(
target: "surrealdb::core::kvs::tx",
error = %err,
"retryable conflict while cleaning uncommitted index build, retrying"
);
sleep(INDEX_BUILD_RESERVATION_RELEASE_RETRY_SLEEP).await;
}
Err(err) => return Err(err),
}
}
}
}
#[derive(Clone)]
pub(crate) struct IndexBuildReservationRelease {
tf: TransactionFactory,
sequences: Sequences,
node: Uuid,
key: Key,
val: Val,
}
impl IndexBuildReservationRelease {
pub(crate) fn new(
tf: TransactionFactory,
sequences: Sequences,
node: Uuid,
key: Key,
val: Val,
) -> Self {
Self {
tf,
sequences,
node,
key,
val,
}
}
async fn release_once(&self) -> Result<()> {
let tx = self
.tf
.transaction(TransactionType::Write, LockType::Optimistic, self.sequences.clone())
.await?;
#[cfg(test)]
if let Err(err) = maybe_inject_non_retryable_error(
NonRetryableErrorSite::ConcurrentIndexReservationRelease,
self.node,
) {
let _ = tx.tr.cancel().await;
return Err(err);
}
match tx.tr.delc(self.key.clone(), Some(self.val.clone())).await {
Ok(()) => {}
Err(KvsError::TransactionConditionNotMet) => {
let _ = tx.tr.cancel().await;
return Ok(());
}
Err(err) => {
let _ = tx.tr.cancel().await;
return Err(err.into());
}
}
#[cfg(test)]
if let Err(err) = maybe_inject_retryable_conflict(
RetryableConflictSite::ConcurrentIndexReservationRelease,
self.node,
) {
let _ = tx.tr.cancel().await;
return Err(err);
}
if let Err(err) = tx.tr.commit().await {
let _ = tx.tr.cancel().await;
return Err(err.into());
}
Ok(())
}
async fn mark_build_error_if_uncommitted(&self, release_err: &anyhow::Error) -> Result<()> {
let br = Br::decode_key(&self.key)?;
let bg_range_start = Bg::new(
br.ns,
br.db,
br.tb.as_ref(),
br.ix,
br.generation,
br.ticket,
BuildTicketMutationSeq::MIN,
)
.encode_key()?;
let bg_range_end = Bg::new(
br.ns,
br.db,
br.tb.as_ref(),
br.ix,
br.generation,
br.ticket,
BuildTicketMutationSeq::MAX,
)
.encode_key()?;
let bs = Bs::new(br.ns, br.db, br.tb.as_ref(), br.ix).encode_key()?;
let reason = format!(
"Failed to release durable index-build reservation for generation {} ticket {} after transaction close: {release_err}",
br.generation, br.ticket
);
loop {
let tx = self
.tf
.transaction(TransactionType::Write, LockType::Optimistic, self.sequences.clone())
.await?;
let current_reservation = match tx.tr.get(self.key.clone(), None).await {
Ok(current) => current,
Err(err) => {
let _ = tx.tr.cancel().await;
return Err(err.into());
}
};
if current_reservation.as_deref() != Some(self.val.as_slice()) {
let _ = tx.tr.cancel().await;
return Ok(());
}
match tx.tr.keys(bg_range_start.clone()..bg_range_end.clone(), 1, 0, None).await {
Ok(res) if !res.keys.is_empty() => {
let _ = tx.tr.cancel().await;
return Ok(());
}
Ok(_) => {}
Err(err) => {
let _ = tx.tr.cancel().await;
return Err(err.into());
}
}
let current_state = match tx.tr.get(bs.clone(), None).await {
Ok(Some(current_state)) => current_state,
Ok(None) => {
let _ = tx.tr.cancel().await;
return Ok(());
}
Err(err) => {
let _ = tx.tr.cancel().await;
return Err(err.into());
}
};
let current = IndexBuildState::kv_decode_value(¤t_state, ())?;
if current.generation != br.generation
|| !matches!(current.phase, IndexBuildPhase::Building | IndexBuildPhase::Closing)
{
let _ = tx.tr.cancel().await;
return Ok(());
}
let mut next = current.clone();
next.phase = IndexBuildPhase::Error;
next.owner = None;
next.owner_heartbeat_at = None;
next.updated_at = Utc::now();
next.error = Some(reason.clone());
next.report_status = Some(IndexBuildReportStatus::Error);
let next_state = next.kv_encode_value()?;
match tx.tr.putc(bs.clone(), next_state, Some(current_state)).await {
Ok(()) => {}
Err(KvsError::TransactionConditionNotMet) => {
let _ = tx.tr.cancel().await;
continue;
}
Err(err) => {
let _ = tx.tr.cancel().await;
return Err(err.into());
}
}
match tx.tr.commit().await {
Ok(()) => return Ok(()),
Err(err) if err.is_retryable() => {
let _ = tx.tr.cancel().await;
sleep(INDEX_BUILD_RESERVATION_RELEASE_RETRY_SLEEP).await;
}
Err(err) => {
let _ = tx.tr.cancel().await;
return Err(err.into());
}
}
}
}
pub(crate) async fn release(self) -> Result<()> {
loop {
match self.release_once().await {
Ok(()) => return Ok(()),
Err(err) if is_retryable_transaction_conflict(&err) => {
tracing::debug!(
target: "surrealdb::core::kvs::tx",
node = %self.node,
error = %err,
"retryable conflict while releasing durable index-build reservation, retrying"
);
sleep(INDEX_BUILD_RESERVATION_RELEASE_RETRY_SLEEP).await;
}
Err(err) => {
tracing::warn!(
target: "surrealdb::core::kvs::tx",
node = %self.node,
"failed to release durable index-build reservation: {err}"
);
if let Err(mark_err) = self.mark_build_error_if_uncommitted(&err).await {
tracing::warn!(
target: "surrealdb::core::kvs::tx",
node = %self.node,
"failed to mark durable index build error after reservation release failure: {mark_err}"
);
}
return Err(err);
}
}
}
}
async fn release_batch(reservations: Vec<Self>) -> Result<(), Vec<Self>> {
if reservations.len() <= 1 {
return Err(reservations);
}
let Some(first) = reservations.first() else {
return Ok(());
};
let tf = first.tf.clone();
let sequences = first.sequences.clone();
let tx = match tf.transaction(TransactionType::Write, LockType::Optimistic, sequences).await
{
Ok(tx) => tx,
Err(_) => return Err(reservations),
};
for reservation in &reservations {
match tx.tr.delc(reservation.key.clone(), Some(reservation.val.clone())).await {
Ok(()) => {}
Err(KvsError::TransactionConditionNotMet) => {
}
Err(_) => {
let _ = tx.tr.cancel().await;
return Err(reservations);
}
}
}
match tx.tr.commit().await {
Ok(()) => Ok(()),
Err(_) => {
let _ = tx.tr.cancel().await;
Err(reservations)
}
}
}
}
impl Deref for Transaction {
type Target = Transactor;
fn deref(&self) -> &Self::Target {
&self.tr
}
}
pub struct MeteredKeysCursor<'a> {
inner: Box<dyn ScanCursorKeys + 'a>,
metrics: &'a TransactionMetrics,
}
impl<'a> MeteredKeysCursor<'a> {
pub async fn next_batch<'s>(&'s mut self, limit: u32) -> Result<KeysBatch<'s>> {
let batch = self.inner.next_batch(limit).await.map_err(Error::from)?;
self.metrics.record_scan(batch.len() as u64, batch.key_bytes, 0);
Ok(batch)
}
pub async fn for_each(&mut self, limit: u32, f: &mut dyn KeyVisitor) -> Result<ScanChunkStats> {
let stats = self.inner.for_each(limit, f).await.map_err(Error::from)?;
self.metrics.record_scan(stats.rows, stats.key_bytes, 0);
Ok(stats)
}
}
pub struct MeteredValsCursor<'a> {
inner: Box<dyn ScanCursorVals + 'a>,
metrics: &'a TransactionMetrics,
}
impl<'a> MeteredValsCursor<'a> {
pub async fn next_batch<'s>(&'s mut self, limit: u32) -> Result<ValsBatch<'s>> {
let batch = self.inner.next_batch(limit).await.map_err(Error::from)?;
self.metrics.record_scan(batch.len() as u64, batch.key_bytes, batch.value_bytes);
Ok(batch)
}
pub async fn for_each(&mut self, limit: u32, f: &mut dyn ValVisitor) -> Result<ScanChunkStats> {
let stats = self.inner.for_each(limit, f).await.map_err(Error::from)?;
self.metrics.record_scan(stats.rows, stats.key_bytes, stats.value_bytes);
Ok(stats)
}
}
struct ReferenceTargets {
any: bool,
tables: HashSet<TableName>,
}
impl ReferenceTargets {
fn can_target(&self, table: &TableName) -> bool {
self.any || self.tables.contains(table)
}
}
impl Transaction {
pub(crate) async fn table_may_have_incoming_references(
&self,
ns: NamespaceId,
db: DatabaseId,
table: &TableName,
) -> Result<bool> {
Ok(self.database_reference_targets(ns, db).await?.can_target(table))
}
async fn database_reference_targets(
&self,
ns: NamespaceId,
db: DatabaseId,
) -> Result<Arc<ReferenceTargets>> {
let qey = cache::tx::Lookup::DbReferenceTargets(ns, db);
if let Some(entry) = self.cache.get(&qey) {
return entry.try_into_type::<ReferenceTargets>();
}
let mut any = false;
let mut tables = HashSet::new();
for tb in self.all_tb(ns, db, None).await?.iter() {
for fd in self.all_tb_fields(ns, db, &tb.name, None).await?.iter() {
if fd.reference.is_none() {
continue;
}
match &fd.field_kind {
Some(kind) => {
if kind.collect_reference_target_tables(&mut tables) {
any = true;
}
}
None => any = true,
}
}
}
let targets = Arc::new(ReferenceTargets {
any,
tables,
});
self.cache.insert(qey, cache::tx::Entry::Any(targets.clone()));
Ok(targets)
}
pub fn new(
local: bool,
sequences: Sequences,
async_event_trigger: Arc<Notify>,
observer: Arc<dyn ExecutionObserver>,
tr: Transactor,
config: &CommonConfig,
) -> Transaction {
Transaction {
local,
started_at: Instant::now(),
observer,
metrics: TransactionMetrics::new(),
tenant_identity: OnceLock::new(),
tr,
cache: TransactionCache::new(config.transaction_cache_size),
sequences,
changefeed: OnceLock::new(),
live_events: OnceLock::new(),
async_event_trigger,
trigger_async_event: AtomicBool::new(false),
pending_index_build_reservations: Mutex::new(Vec::new()),
cached_index_build_reservations: Mutex::new(HashMap::new()),
pending_index_builder_aborts: Mutex::new(Vec::new()),
pending_uncommitted_index_builds: Mutex::new(Vec::new()),
}
}
pub fn with_tenant_identity(self, identity: Option<Arc<TenantIdentity>>) -> Self {
if let Some(id) = identity {
let _ = self.tenant_identity.set(id);
}
self
}
pub fn set_tenant_identity(&self, identity: Arc<TenantIdentity>) {
let _ = self.tenant_identity.set(identity);
}
#[cfg(test)]
pub(crate) fn metrics_snapshot_for_test(&self) -> crate::observe::TransactionMetricsSnapshot {
self.metrics.snapshot()
}
fn emit_transaction_event(&self, outcome: Outcome) {
self.emit_transaction_event_with_class(outcome, None);
}
fn emit_transaction_event_with_class(
&self,
outcome: Outcome,
error_class: Option<&'static str>,
) {
if self.observer.is_noop() {
return;
}
self.observer.on_transaction_complete(&TransactionEvent {
safe: TransactionEventSafe {
outcome,
write: self.tr.writeable(),
duration: self.started_at.elapsed(),
metrics: self.metrics.snapshot(),
error_class,
},
ctx: self.tenant_identity.get().map(|t| t.to_transaction_ctx()).unwrap_or_default(),
});
}
pub(crate) async fn register_index_build_reservation_release(
&self,
release: IndexBuildReservationRelease,
) {
self.pending_index_build_reservations.lock().await.push(release);
}
pub(crate) async fn lookup_cached_index_build_reservation(
&self,
key: &CachedIndexBuildReservationKey,
) -> Result<Option<CachedIndexBuildReservationLookup>> {
let mut cache = self.cached_index_build_reservations.lock().await;
let Some(entry) = cache.get_mut(key) else {
return Ok(None);
};
let mutation_seq = entry.next_mutation_seq;
let next_seq =
mutation_seq.checked_add(1).ok_or_else(|| Error::IndexingBuildingCancelled {
reason: "Per-user-transaction index build mutation sequence overflowed u32::MAX"
.to_string(),
})?;
entry.next_mutation_seq = next_seq;
Ok(Some(CachedIndexBuildReservationLookup::Reused {
generation: entry.generation,
ticket: entry.ticket,
mutation_seq,
initial_complete: entry.initial_complete,
}))
}
#[cfg(test)]
#[cfg_attr(not(feature = "kv-mem"), allow(dead_code))]
pub(crate) async fn seed_cached_index_build_reservation_for_test(
&self,
key: CachedIndexBuildReservationKey,
generation: BuildGeneration,
ticket: BuildTicket,
initial_complete: bool,
next_mutation_seq: BuildTicketMutationSeq,
) {
self.cached_index_build_reservations.lock().await.insert(
key,
CachedIndexBuildReservation {
generation,
ticket,
initial_complete,
next_mutation_seq,
},
);
}
pub(crate) async fn remove_cached_index_build_reservation(
&self,
key: &CachedIndexBuildReservationKey,
) {
self.cached_index_build_reservations.lock().await.remove(key);
}
pub(crate) async fn insert_cached_index_build_reservation(
&self,
key: CachedIndexBuildReservationKey,
generation: BuildGeneration,
ticket: BuildTicket,
initial_complete: bool,
) -> CachedIndexBuildReservationLookup {
let mut cache = self.cached_index_build_reservations.lock().await;
cache.insert(
key,
CachedIndexBuildReservation {
generation,
ticket,
initial_complete,
next_mutation_seq: 1,
},
);
CachedIndexBuildReservationLookup::FirstUse {
generation,
ticket,
mutation_seq: 0,
initial_complete,
}
}
pub(crate) async fn register_index_builder_abort_after_commit(
&self,
builder: IndexBuilder,
ns: NamespaceId,
db: DatabaseId,
tb: TableName,
ix: IndexId,
) {
self.pending_index_builder_aborts.lock().await.push(PendingIndexBuilderAbort {
builder,
ns,
db,
tb,
ix,
});
}
pub(crate) async fn register_uncommitted_index_build_cleanup(
&self,
builder: IndexBuilder,
tf: TransactionFactory,
ns: NamespaceId,
db: DatabaseId,
tb: TableName,
ix: IndexId,
) {
self.pending_uncommitted_index_builds.lock().await.push(PendingUncommittedIndexBuild {
builder,
tf,
sequences: self.sequences.clone(),
ns,
db,
tb,
ix,
});
}
pub fn is_local(&self) -> bool {
self.local
}
pub fn enclose(self) -> Arc<Transaction> {
Arc::new(self)
}
pub fn closed(&self) -> bool {
self.tr.closed()
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip_all)]
pub async fn cancel(&self) -> Result<()> {
if let Some(changefeed) = self.changefeed.get() {
changefeed.clear();
}
if let Some(live_events) = self.live_events.get() {
live_events.clear();
}
let result = self.tr.cancel().await.map_err(Error::from);
let cleanup_result = self.cleanup_uncommitted_index_builds().await;
let release_result = self.release_index_build_reservations().await;
self.discard_index_builder_aborts().await;
self.emit_transaction_event(Outcome::from(&result));
result?;
cleanup_result?;
release_result?;
Ok(())
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip_all)]
pub async fn commit(&self) -> Result<()> {
if let Err(e) = self.store_changes().await {
if let Err(err) = self.cancel().await {
tracing::warn!(
target: "surrealdb::core::kvs::tx",
"transaction cleanup failed after changefeed storage failed; preserving original store_changes error {e}: {err}"
);
}
return Err(e);
}
if let Err(e) = self.tr.commit().await {
let cleanup_result = self.cleanup_uncommitted_index_builds().await;
let release_result = self.release_index_build_reservations().await;
self.discard_index_builder_aborts().await;
let class = if e.is_retryable() {
crate::observe::error_class::TXN_CONFLICT
} else {
crate::observe::error_class::STORAGE
};
self.emit_transaction_event_with_class(Outcome::Error, Some(class));
if let Err(err) = release_result {
tracing::warn!(
target: "surrealdb::core::kvs::tx",
"durable index-build reservation cleanup failed after transaction commit failed; preserving original commit error {e}: {err}"
);
}
if let Err(err) = cleanup_result {
tracing::warn!(
target: "surrealdb::core::kvs::tx",
"uncommitted index-build cleanup failed after transaction commit failed; preserving original commit error {e}: {err}"
);
}
anyhow::bail!(e);
}
if let Err(err) = self.release_index_build_reservations().await {
tracing::warn!(
target: "surrealdb::core::kvs::tx",
"durable index-build reservation cleanup failed after transaction commit; committed appendings remain recoverable: {err}"
);
}
self.discard_uncommitted_index_builds().await;
self.run_index_builder_aborts().await;
if self.trigger_async_event.load(Ordering::Relaxed) {
self.async_event_trigger.notify_one();
}
self.emit_transaction_event(Outcome::Success);
Ok(())
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip_all)]
pub async fn exists<K>(&self, key: &K, version: Option<u64>) -> Result<bool>
where
K: KVKey + Debug,
{
let key = key.encode_key()?;
let key_bytes = key.len() as u64;
let found = self.tr.exists(key, version).await.map_err(Error::from)?;
self.metrics.record_get(u64::from(found), key_bytes, 0);
Ok(found)
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip_all)]
pub async fn get<K>(&self, key: &K, version: Option<u64>) -> Result<Option<K::ValueType>>
where
K: KVKey + Debug,
{
let encoded = key.encode_key()?;
let key_bytes = encoded.len() as u64;
let val = self.tr.get(encoded, version).await.map_err(Error::from)?;
let (keys_found, value_bytes) = match &val {
Some(v) => (1, v.len() as u64),
None => (0, 0),
};
self.metrics.record_get(keys_found, key_bytes, value_bytes);
val.map(|v| K::ValueType::kv_decode_value(&v, key.value_context())).transpose()
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip_all)]
pub async fn getm<K>(
&self,
keys: Vec<K>,
version: Option<u64>,
) -> Result<Vec<Option<K::ValueType>>>
where
K: KVKey + Debug,
{
let encoded_keys: Vec<_> = keys.iter().map(|k| k.encode_key()).collect::<Result<_>>()?;
let key_bytes: u64 = encoded_keys.iter().map(|k| k.len() as u64).sum();
let res = self.tr.getm(encoded_keys, version).await.map_err(Error::from)?;
self.metrics.record_get(res.records, key_bytes, res.value_bytes);
res.values
.into_iter()
.zip(keys)
.map(|(v, k)| match v {
Some(v) => K::ValueType::kv_decode_value(&v, k.value_context()).map(Some),
None => Ok(None),
})
.collect()
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip_all)]
pub async fn getp<K>(&self, key: &K, version: Option<u64>) -> Result<Vec<(Key, K::ValueType)>>
where
K: KVKey + Debug,
K::ValueType: KVValue<KeyContext = ()>,
{
let key = key.encode_key()?;
let res = self.tr.getp(key, version).await.map_err(Error::from)?;
self.metrics.record_scan(res.values.len() as u64, res.key_bytes, res.value_bytes);
res.values
.into_iter()
.map(|(k, v)| Ok((k, K::ValueType::kv_decode_value(&v, ())?)))
.collect()
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip_all)]
pub async fn getr<K>(
&self,
rng: Range<K>,
version: Option<u64>,
) -> Result<Vec<(Key, K::ValueType)>>
where
K: KVKey + Debug,
K::ValueType: KVValue<KeyContext = ()>,
{
let beg = rng.start.encode_key()?;
let end = rng.end.encode_key()?;
let res = self.tr.getr(beg..end, version).await.map_err(Error::from)?;
self.metrics.record_scan(res.values.len() as u64, res.key_bytes, res.value_bytes);
res.values
.into_iter()
.map(|(k, v)| Ok((k, K::ValueType::kv_decode_value(&v, ())?)))
.collect()
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip_all)]
pub async fn del<K>(&self, key: &K) -> Result<()>
where
K: KVKey + Debug,
{
let key = key.encode_key()?;
let key_bytes = key.len() as u64;
self.tr.del(key).await.map_err(Error::from)?;
self.metrics.record_del(1, key_bytes);
Ok(())
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip_all)]
pub async fn delc<K>(&self, key: &K, chk: Option<&K::ValueType>) -> Result<()>
where
K: KVKey + Debug,
{
let key = key.encode_key()?;
let key_bytes = key.len() as u64;
let chk = chk.map(|v| v.kv_encode_value()).transpose()?;
self.tr.delc(key, chk).await.map_err(Error::from)?;
self.metrics.record_del(1, key_bytes);
Ok(())
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip_all)]
pub async fn delr<K>(&self, rng: Range<K>) -> Result<()>
where
K: KVKey + Debug,
{
let beg = rng.start.encode_key()?;
let end = rng.end.encode_key()?;
self.tr.delr(beg..end).await.map_err(Error::from)?;
self.metrics.record_del(0, 0);
Ok(())
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip_all)]
pub async fn delp<K>(&self, key: &K) -> Result<()>
where
K: KVKey + Debug,
{
let key = key.encode_key()?;
self.tr.delp(key).await.map_err(Error::from)?;
self.metrics.record_del(0, 0);
Ok(())
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip_all)]
pub async fn clr<K>(&self, key: &K) -> Result<()>
where
K: KVKey + Debug,
{
let key = key.encode_key()?;
let key_bytes = key.len() as u64;
self.tr.clr(key).await.map_err(Error::from)?;
self.metrics.record_del(1, key_bytes);
Ok(())
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip_all)]
pub async fn clrc<K>(&self, key: &K, chk: Option<&K::ValueType>) -> Result<()>
where
K: KVKey + Debug,
{
let key = key.encode_key()?;
let key_bytes = key.len() as u64;
let chk = chk.map(|v| v.kv_encode_value()).transpose()?;
self.tr.clrc(key, chk).await.map_err(Error::from)?;
self.metrics.record_del(1, key_bytes);
Ok(())
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip_all)]
pub async fn clrr<K>(&self, rng: Range<K>) -> Result<()>
where
K: KVKey + Debug,
{
let beg = rng.start.encode_key()?;
let end = rng.end.encode_key()?;
self.tr.clrr(beg..end).await.map_err(Error::from)?;
self.metrics.record_del(0, 0);
Ok(())
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip_all)]
pub async fn clrp<K>(&self, key: &K) -> Result<()>
where
K: KVKey + Debug,
{
let key = key.encode_key()?;
self.tr.clrp(key).await.map_err(Error::from)?;
self.metrics.record_del(0, 0);
Ok(())
}
pub(crate) async fn del_ns_deferred(
&self,
ns: &str,
expunge: bool,
) -> Result<Option<NamespaceId>> {
let Some(ns_def) = self.get_ns_by_name(ns, None).await? else {
return Ok(None);
};
let key = crate::key::root::ns::new(&ns_def.name);
if expunge {
self.clr(&key).await?;
} else {
self.del(&key).await?;
}
let rc = crate::key::root::rc::ReclaimKey::namespace(
ns_def.namespace_id,
expunge,
Uuid::now_v7(),
);
self.set(
&rc,
&crate::key::root::rc::ReclaimState {
observed_ms: 0,
},
)
.await?;
self.cache.remove(&cache::tx::Lookup::Nss);
self.cache.remove(&cache::tx::Lookup::NsByName(&ns_def.name));
Ok(Some(ns_def.namespace_id))
}
pub(crate) async fn del_db_deferred(
&self,
ns: &str,
db: &str,
expunge: bool,
) -> Result<Option<DatabaseId>> {
let Some(db_def) = self.get_db_by_name(ns, db, None).await? else {
return Ok(None);
};
let key = crate::key::namespace::db::new(db_def.namespace_id, &db_def.name);
if expunge {
self.clr(&key).await?;
} else {
self.del(&key).await?;
}
let rc = crate::key::root::rc::ReclaimKey::database(
db_def.namespace_id,
db_def.database_id,
expunge,
Uuid::now_v7(),
);
self.set(
&rc,
&crate::key::root::rc::ReclaimState {
observed_ms: 0,
},
)
.await?;
self.cache.remove(&cache::tx::Lookup::Dbs(db_def.namespace_id));
self.cache.remove(&cache::tx::Lookup::DbByName(ns, &db_def.name));
Ok(Some(db_def.database_id))
}
pub(crate) async fn del_tb_index_deferred(
&self,
ns: NamespaceId,
db: DatabaseId,
tb: &TableName,
ix: &str,
) -> Result<()> {
let Some(ix_def) = self.get_tb_index(ns, db, tb, ix, None).await? else {
return Ok(());
};
let key = crate::key::table::ix::new(ns, db, tb, &ix_def.name);
self.del(&key).await?;
let name_lookup_key =
crate::key::table::ix::IndexNameLookupKey::new(ns, db, tb, ix_def.index_id);
self.del(&name_lookup_key).await?;
let rc = crate::key::root::rc::ReclaimKey::index(
ns,
db,
std::borrow::Cow::Borrowed(tb),
ix_def.index_id,
false,
Uuid::now_v7(),
);
self.set(
&rc,
&crate::key::root::rc::ReclaimState {
observed_ms: 0,
},
)
.await?;
self.cache.remove(&cache::tx::Lookup::Ixs(ns, db, tb.as_ref()));
self.cache.remove(&cache::tx::Lookup::Ix(ns, db, tb.as_ref(), &ix_def.name));
Ok(())
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip_all)]
pub async fn set<K>(&self, key: &K, val: &K::ValueType) -> Result<()>
where
K: KVKey + Debug,
{
let key = key.encode_key()?;
let val = val.kv_encode_value()?;
let key_bytes = key.len() as u64;
let value_bytes = val.len() as u64;
self.tr.set(key, val).await.map_err(Error::from)?;
self.metrics.record_set(key_bytes, value_bytes);
Ok(())
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip_all)]
pub async fn put<K>(&self, key: &K, val: &K::ValueType) -> Result<()>
where
K: KVKey + Debug,
{
let key = key.encode_key()?;
let val = val.kv_encode_value()?;
let key_bytes = key.len() as u64;
let value_bytes = val.len() as u64;
self.tr.put(key, val).await.map_err(Error::from)?;
self.metrics.record_put(key_bytes, value_bytes);
Ok(())
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip_all)]
pub async fn putc<K>(
&self,
key: &K,
val: &K::ValueType,
chk: Option<&K::ValueType>,
) -> Result<()>
where
K: KVKey + Debug,
{
let key = key.encode_key()?;
let val = val.kv_encode_value()?;
let chk = chk.map(|v| v.kv_encode_value()).transpose()?;
let key_bytes = key.len() as u64;
let value_bytes = val.len() as u64;
self.tr.putc(key, val, chk).await.map_err(Error::from)?;
self.metrics.record_put(key_bytes, value_bytes);
Ok(())
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip_all)]
pub async fn replace<K>(&self, key: &K, val: &K::ValueType) -> Result<()>
where
K: KVKey + Debug,
{
let key = key.encode_key()?;
let val = val.kv_encode_value()?;
let key_bytes = key.len() as u64;
let value_bytes = val.len() as u64;
self.tr.replace(key, val).await.map_err(Error::from)?;
self.metrics.record_put(key_bytes, value_bytes);
Ok(())
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip_all)]
pub async fn get_raw<K>(&self, key: &K, version: Option<u64>) -> Result<Option<Val>>
where
K: KVKey + Debug,
{
let key = key.encode_key()?;
let key_bytes = key.len() as u64;
let val = self.tr.get(key, version).await.map_err(Error::from)?;
let (keys_found, value_bytes) = match &val {
Some(v) => (1, v.len() as u64),
None => (0, 0),
};
self.metrics.record_get(keys_found, key_bytes, value_bytes);
Ok(val)
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip_all)]
pub async fn getm_raw<K>(&self, keys: Vec<K>, version: Option<u64>) -> Result<Vec<Option<Val>>>
where
K: KVKey + Debug,
{
let keys = keys.iter().map(|k| k.encode_key()).collect::<Result<Vec<_>>>()?;
let key_bytes: u64 = keys.iter().map(|k| k.len() as u64).sum();
let res = self.tr.getm(keys, version).await.map_err(Error::from)?;
self.metrics.record_get(res.records, key_bytes, res.value_bytes);
Ok(res.values)
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip_all)]
pub async fn keys<K>(
&self,
rng: Range<K>,
limit: u32,
skip: u32,
version: Option<u64>,
) -> Result<Vec<Key>>
where
K: KVKey + Debug,
{
let beg = rng.start.encode_key()?;
let end = rng.end.encode_key()?;
let res = self.tr.keys(beg..end, limit, skip, version).await.map_err(Error::from)?;
self.metrics.record_scan(res.keys.len() as u64, res.key_bytes, 0);
Ok(res.keys)
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip_all)]
pub async fn keysr<K>(
&self,
rng: Range<K>,
limit: u32,
skip: u32,
version: Option<u64>,
) -> Result<Vec<Key>>
where
K: KVKey + Debug,
{
let beg = rng.start.encode_key()?;
let end = rng.end.encode_key()?;
let res = self.tr.keysr(beg..end, limit, skip, version).await.map_err(Error::from)?;
self.metrics.record_scan(res.keys.len() as u64, res.key_bytes, 0);
Ok(res.keys)
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip_all)]
pub async fn scan<K>(
&self,
rng: Range<K>,
limit: u32,
skip: u32,
version: Option<u64>,
) -> Result<Vec<(Key, Val)>>
where
K: KVKey + Debug,
{
let beg = rng.start.encode_key()?;
let end = rng.end.encode_key()?;
let res = self.tr.scan(beg..end, limit, skip, version).await.map_err(Error::from)?;
self.metrics.record_scan(res.values.len() as u64, res.key_bytes, res.value_bytes);
Ok(res.values)
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip_all)]
pub async fn scanr<K>(
&self,
rng: Range<K>,
limit: u32,
skip: u32,
version: Option<u64>,
) -> Result<Vec<(Key, Val)>>
where
K: KVKey + Debug,
{
let beg = rng.start.encode_key()?;
let end = rng.end.encode_key()?;
let res = self.tr.scanr(beg..end, limit, skip, version).await.map_err(Error::from)?;
self.metrics.record_scan(res.values.len() as u64, res.key_bytes, res.value_bytes);
Ok(res.values)
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip_all)]
pub async fn count<K>(&self, rng: Range<K>, version: Option<u64>) -> Result<usize>
where
K: KVKey + Debug,
{
let beg = rng.start.encode_key()?;
let end = rng.end.encode_key()?;
let n = self.tr.count(beg..end, version).await.map_err(Error::from)?;
self.metrics.record_scan(n as u64, 0, 0);
Ok(n)
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip_all)]
pub async fn open_keys_cursor(
&self,
rng: Range<Key>,
dir: ScanDirection,
skip: u32,
version: Option<u64>,
) -> Result<MeteredKeysCursor<'_>> {
let inner = self
.tr
.open_keys_cursor(
rng,
match dir {
ScanDirection::Forward => Direction::Forward,
ScanDirection::Backward => Direction::Backward,
},
skip,
version,
)
.await
.map_err(Error::from)?;
Ok(MeteredKeysCursor {
inner,
metrics: &self.metrics,
})
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip_all)]
pub async fn open_vals_cursor(
&self,
rng: Range<Key>,
dir: ScanDirection,
skip: u32,
version: Option<u64>,
) -> Result<MeteredValsCursor<'_>> {
let inner = self
.tr
.open_vals_cursor(
rng,
match dir {
ScanDirection::Forward => Direction::Forward,
ScanDirection::Backward => Direction::Backward,
},
skip,
version,
)
.await
.map_err(Error::from)?;
Ok(MeteredValsCursor {
inner,
metrics: &self.metrics,
})
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip_all)]
pub async fn batch_keys<K>(
&self,
rng: Range<K>,
batch: u32,
version: Option<u64>,
) -> Result<Batch<Key>>
where
K: KVKey + Debug,
{
let beg = rng.start.encode_key()?;
let end = rng.end.encode_key()?;
Ok(self.tr.batch_keys(beg..end, batch, version).await.map_err(Error::from)?)
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip_all)]
pub async fn batch_keys_vals<K>(
&self,
rng: Range<K>,
batch: u32,
version: Option<u64>,
) -> Result<Batch<(Key, Val)>>
where
K: KVKey + Debug,
{
let beg = rng.start.encode_key()?;
let end = rng.end.encode_key()?;
Ok(self.tr.batch_keys_vals(beg..end, batch, version).await.map_err(Error::from)?)
}
pub async fn new_save_point(&self) -> Result<()> {
Ok(self.inner.new_save_point().await.map_err(Error::from)?)
}
pub async fn release_last_save_point(&self) -> Result<()> {
Ok(self.inner.release_last_save_point().await.map_err(Error::from)?)
}
pub async fn rollback_to_save_point(&self) -> Result<()> {
Ok(self.inner.rollback_to_save_point().await.map_err(Error::from)?)
}
pub async fn timestamp(&self) -> Result<BoxTimeStamp> {
Ok(self.tr.timestamp().await.map_err(Error::from)?)
}
pub async fn safe_timestamp(&self) -> Result<BoxTimeStamp> {
Ok(self.tr.safe_timestamp().await.map_err(Error::from)?)
}
pub(crate) async fn table_has_live_query(
&self,
ns: NamespaceId,
db: DatabaseId,
tb: &TableName,
) -> Result<bool> {
let beg = crate::key::table::lq::prefix(ns, db, tb)?;
let end = crate::key::table::lq::suffix(ns, db, tb)?;
Ok(!self.keys(beg..end, 1, 0, None).await?.is_empty())
}
pub fn timestamp_impl(&self) -> BoxTimeStampImpl {
self.tr.timestamp_impl()
}
pub(crate) fn changefeed_buffer_table_change(
&self,
ns: NamespaceId,
db: DatabaseId,
tb: &TableName,
dt: &TableDefinition,
) {
self.changefeed.get_or_init(Changefeed::new).buffer_table_change(ns, db, tb, dt)
}
#[expect(clippy::too_many_arguments)]
pub(crate) fn changefeed_buffer_record_change(
&self,
ns: NamespaceId,
db: DatabaseId,
tb: &TableName,
id: &RecordId,
previous: CursorRecord,
current: CursorRecord,
store_difference: bool,
) {
self.changefeed.get_or_init(Changefeed::new).buffer_record_change(
ns,
db,
tb,
id.clone(),
previous,
current,
store_difference,
)
}
pub(crate) fn live_event_buffer_record_change(
&self,
ns: NamespaceId,
db: DatabaseId,
tb: &TableName,
id: &RecordId,
previous: CursorRecord,
current: CursorRecord,
) {
self.live_events.get_or_init(LiveEventBuffer::new).buffer_record_change(
ns,
db,
tb,
id.clone(),
previous.into_owned(),
current.into_owned(),
)
}
pub(crate) async fn store_changes(&self) -> Result<()> {
let cf_changes = match self.changefeed.get() {
Some(changefeed) => changefeed.changes()?,
None => Vec::new(),
};
let lqe_changes = match self.live_events.get() {
Some(live_events) => live_events.changes()?,
None => Vec::new(),
};
if cf_changes.is_empty() && lqe_changes.is_empty() {
return Ok(());
}
let buf = &mut [0u8; _];
let ts = self.timestamp().await?.encode(buf);
let cf_futures = cf_changes.into_iter().map(|(ns, db, tb, value)| async move {
let key = crate::key::change::new(ns, db, ts, &tb).encode_key()?;
self.tr.set(key, value).await.map_err(Error::from)?;
Ok::<(), anyhow::Error>(())
});
try_join_all(cf_futures).await?;
let lqe_futures = lqe_changes.into_iter().map(|(ns, db, tb, value)| async move {
let key = crate::key::lqe::new(ns, db, ts, &tb).encode_key()?;
self.tr.set(key, value).await.map_err(Error::from)?;
Ok::<(), anyhow::Error>(())
});
try_join_all(lqe_futures).await?;
Ok(())
}
async fn release_index_build_reservations(&self) -> Result<()> {
let reservations = {
let mut pending = self.pending_index_build_reservations.lock().await;
std::mem::take(&mut *pending)
};
if reservations.is_empty() {
return Ok(());
}
let reservations = match IndexBuildReservationRelease::release_batch(reservations).await {
Ok(()) => return Ok(()),
Err(remaining) => remaining,
};
let mut first_error = None;
for reservation in reservations {
if let Err(err) = reservation.release().await
&& first_error.is_none()
{
first_error = Some(err);
}
}
if let Some(err) = first_error {
Err(err)
} else {
Ok(())
}
}
async fn cleanup_uncommitted_index_builds(&self) -> Result<()> {
let builds = {
let mut pending = self.pending_uncommitted_index_builds.lock().await;
std::mem::take(&mut *pending)
};
let mut first_error = None;
for build in builds {
if let Err(err) = build.cleanup().await
&& first_error.is_none()
{
first_error = Some(err);
}
}
if let Some(err) = first_error {
Err(err)
} else {
Ok(())
}
}
async fn discard_uncommitted_index_builds(&self) {
self.pending_uncommitted_index_builds.lock().await.clear();
}
async fn run_index_builder_aborts(&self) {
let aborts = {
let mut pending = self.pending_index_builder_aborts.lock().await;
std::mem::take(&mut *pending)
};
for abort in aborts {
abort.abort().await;
}
}
async fn discard_index_builder_aborts(&self) {
self.pending_index_builder_aborts.lock().await.clear();
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip(self))]
pub fn clear_cache(&self) {
self.cache.clear()
}
#[instrument(level = "trace", target = "surrealdb::core::kvs::tx", skip_all)]
pub async fn compact<K>(&self, key: Option<K>) -> Result<()>
where
K: KVKey + Debug,
{
let rng = match key {
Some(key) => Some(util::to_prefix_range(&key)?),
None => None,
};
self.tr.inner.compact(rng).await
}
pub(crate) fn trigger_async_event(&self) {
self.trigger_async_event.store(true, Ordering::Relaxed);
}
}
impl NodeProvider for Transaction {
fn all_nodes(&self) -> BoxProviderFut<'_, Result<Arc<[Node]>>> {
Box::pin(
async move {
let qey = cache::tx::Lookup::Nds;
match self.cache.get(&qey) {
Some(val) => val.try_into_nds(),
None => {
let beg = crate::key::root::nd::prefix();
let end = crate::key::root::nd::suffix();
let val = self.getr(beg..end, None).await?;
let val = util::deserialize_cache(val.iter().map(|x| x.1.as_slice()))?;
let entry = cache::tx::Entry::Nds(Arc::clone(&val));
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "all_nodes")),
)
}
fn get_node(&self, id: Uuid) -> BoxProviderFut<'_, Result<Arc<Node>>> {
Box::pin(
async move {
let qey = cache::tx::Lookup::Nd(id);
match self.cache.get(&qey) {
Some(val) => val,
None => {
let key = crate::key::root::nd::new(id);
let val = self.get(&key, None).await?.ok_or_else(|| Error::NdNotFound {
uuid: id.to_string(),
})?;
let val = cache::tx::Entry::Any(Arc::new(val));
self.cache.insert(qey, val.clone());
val
}
}
.try_into_type()
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "get_node")),
)
}
}
impl RootProvider for Transaction {
fn get_default_config(&self) -> BoxProviderFut<'_, Result<Option<Arc<DefaultConfig>>>> {
Box::pin(async move {
let qey = cache::tx::Lookup::Rcg("default");
match self.cache.get(&qey) {
Some(val) => val,
None => {
let key = crate::key::root::root_config::new("default");
let Some(val) = self.get(&key, None).await? else {
return Ok(None);
};
let ConfigDefinition::Default(val) = val else {
fail!("Expected a default config but found {val:?} instead");
};
let val = cache::tx::Entry::Any(Arc::new(val));
self.cache.insert(qey, val.clone());
val
}
}
.try_into_type()
.map(Option::Some)
})
}
fn get_root_config<'a>(
&'a self,
cg: &'a str,
) -> BoxProviderFut<'a, Result<Option<Arc<ConfigDefinition>>>> {
Box::pin(
async move {
let qey = cache::tx::Lookup::Rcg(cg);
match self.cache.get(&qey) {
Some(val) => val.try_into_type().map(Option::Some),
None => {
let key = crate::key::root::root_config::new(cg);
if let Some(val) = self.get(&key, None).await? {
let val = Arc::new(val);
let entr = cache::tx::Entry::Any(val.clone());
self.cache.insert(qey, entr);
Ok(Some(val))
} else {
Ok(None)
}
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "get_root_config")),
)
}
}
impl NamespaceProvider for Transaction {
fn all_ns(
&self,
version: Option<u64>,
) -> BoxProviderFut<'_, Result<Arc<[NamespaceDefinition]>>> {
Box::pin(
async move {
if version.is_some() {
let beg = crate::key::root::ns::prefix();
let end = crate::key::root::ns::suffix();
let val = self.getr(beg..end, version).await?;
return util::deserialize_cache(val.iter().map(|x| x.1.as_slice()));
}
let qey = cache::tx::Lookup::Nss;
match self.cache.get(&qey) {
Some(val) => val.try_into_nss(),
None => {
let beg = crate::key::root::ns::prefix();
let end = crate::key::root::ns::suffix();
let val = self.getr(beg..end, None).await?;
let val = util::deserialize_cache(val.iter().map(|x| x.1.as_slice()))?;
let entry = cache::tx::Entry::Nss(Arc::clone(&val));
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "all_ns")),
)
}
fn get_ns_by_name<'a>(
&'a self,
ns: &'a str,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Option<Arc<NamespaceDefinition>>>> {
Box::pin(async move {
if version.is_some() {
let key = crate::key::root::ns::new(ns);
let Some(ns) = self.get(&key, version).await? else {
return Ok(None);
};
return Ok(Some(Arc::new(ns)));
}
let qey = cache::tx::Lookup::NsByName(ns);
match self.cache.get(&qey) {
Some(val) => val.try_into_type().map(Some),
None => {
let key = crate::key::root::ns::new(ns);
let Some(ns) = self.get(&key, None).await? else {
return Ok(None);
};
let ns = Arc::new(ns);
let entr = cache::tx::Entry::Any(ns.clone());
self.cache.insert(qey, entr);
Ok(Some(ns))
}
}
})
}
fn expect_ns_by_name<'a>(
&'a self,
ns: &'a str,
) -> BoxProviderFut<'a, Result<Arc<NamespaceDefinition>>> {
Box::pin(async move {
match self.get_ns_by_name(ns, None).await? {
Some(val) => Ok(val),
None => anyhow::bail!(Error::NsNotFound {
name: ns.to_owned(),
}),
}
})
}
fn put_ns(
&self,
ns: NamespaceDefinition,
) -> BoxProviderFut<'_, Result<Arc<NamespaceDefinition>>> {
Box::pin(async move {
let key = crate::key::root::ns::new(&ns.name);
self.set(&key, &ns).await?;
let list_key = cache::tx::Lookup::Nss;
self.cache.remove(&list_key);
let cached_ns = Arc::new(ns.clone());
let entry = cache::tx::Entry::Any(Arc::clone(&cached_ns) as Arc<dyn Any + Send + Sync>);
let qey = cache::tx::Lookup::NsByName(&ns.name);
self.cache.insert(qey, entry);
Ok(cached_ns)
})
}
fn get_next_ns_id<'a>(
&'a self,
ctx: Option<&'a Context>,
) -> BoxProviderFut<'a, Result<NamespaceId>> {
Box::pin(async move { self.sequences.next_namespace_id(ctx).await })
}
}
impl DatabaseProvider for Transaction {
fn all_db(
&self,
ns: NamespaceId,
version: Option<u64>,
) -> BoxProviderFut<'_, Result<Arc<[DatabaseDefinition]>>> {
Box::pin(
async move {
if version.is_some() {
let beg = crate::key::namespace::db::prefix(ns)?;
let end = crate::key::namespace::db::suffix(ns)?;
let val = self.getr(beg..end, version).await?;
return util::deserialize_cache(val.iter().map(|x| x.1.as_slice()));
}
let qey = cache::tx::Lookup::Dbs(ns);
match self.cache.get(&qey) {
Some(val) => val.try_into_dbs(),
None => {
let beg = crate::key::namespace::db::prefix(ns)?;
let end = crate::key::namespace::db::suffix(ns)?;
let val = self.getr(beg..end, None).await?;
let val = util::deserialize_cache(val.iter().map(|x| x.1.as_slice()))?;
let entry = cache::tx::Entry::Dbs(Arc::clone(&val));
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "all_db")),
)
}
fn get_db_by_name<'a>(
&'a self,
ns: &'a str,
db: &'a str,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Option<Arc<DatabaseDefinition>>>> {
Box::pin(
async move {
if version.is_some() {
let Some(ns) = self.get_ns_by_name(ns, version).await? else {
return Ok(None);
};
let key = crate::key::namespace::db::new(ns.namespace_id, db);
let Some(db_def) = self.get(&key, version).await? else {
return Ok(None);
};
return Ok(Some(Arc::new(db_def)));
}
let qey = cache::tx::Lookup::DbByName(ns, db);
match self.cache.get(&qey) {
Some(val) => val.try_into_type().map(Some),
None => {
let Some(ns) = self.get_ns_by_name(ns, None).await? else {
return Ok(None);
};
let key = crate::key::namespace::db::new(ns.namespace_id, db);
let Some(db_def) = self.get(&key, None).await? else {
return Ok(None);
};
let val = Arc::new(db_def);
let entry = cache::tx::Entry::Any(val.clone());
self.cache.insert(qey, entry);
Ok(Some(val))
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "get_db_by_name")),
)
}
fn get_or_add_db_upwards<'a>(
&'a self,
ctx: Option<&'a Context>,
ns: &'a str,
db: &'a str,
upwards: bool,
) -> BoxProviderFut<'a, Result<Arc<DatabaseDefinition>>> {
Box::pin(
async move {
let qey = cache::tx::Lookup::DbByName(ns, db);
match self.cache.get(&qey) {
Some(val) => {
let t = val.try_into_type()?;
Ok(t)
}
None => {
let db_def = self.get_db_by_name(ns, db, None).await?;
if let Some(db_def) = db_def {
return Ok(db_def);
}
let ns_def = if upwards {
self.get_or_add_ns(ctx, ns).await?
} else {
match self.get_ns_by_name(ns, None).await? {
Some(ns_def) => ns_def,
None => {
return Err(Error::NsNotFound {
name: ns.to_owned(),
}
.into());
}
}
};
let db_def = DatabaseDefinition {
namespace_id: ns_def.namespace_id,
database_id: self.get_next_db_id(ctx, ns_def.namespace_id).await?,
name: db.into(),
comment: None,
changefeed: None,
strict: false,
};
return self.put_db(ns_def.name.as_str(), db_def).await;
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "get_or_add_db_upwards")),
)
}
fn get_next_db_id<'a>(
&'a self,
ctx: Option<&'a Context>,
ns: NamespaceId,
) -> BoxProviderFut<'a, Result<DatabaseId>> {
Box::pin(async move { self.sequences.next_database_id(ctx, ns).await })
}
fn put_db<'a>(
&'a self,
ns: &'a str,
db: DatabaseDefinition,
) -> BoxProviderFut<'a, Result<Arc<DatabaseDefinition>>> {
Box::pin(async move {
let key = crate::key::namespace::db::new(db.namespace_id, &db.name);
self.set(&key, &db).await?;
let list_key = cache::tx::Lookup::Dbs(db.namespace_id);
self.cache.remove(&list_key);
let cached_db = Arc::new(db.clone());
let entry = cache::tx::Entry::Any(Arc::clone(&cached_db) as Arc<dyn Any + Send + Sync>);
let qey = cache::tx::Lookup::DbByName(ns, &db.name);
self.cache.insert(qey, entry);
Ok(cached_db)
})
}
fn all_db_analyzers(
&self,
ns: NamespaceId,
db: DatabaseId,
version: Option<u64>,
) -> BoxProviderFut<'_, Result<Arc<[catalog::AnalyzerDefinition]>>> {
Box::pin(
async move {
if version.is_some() {
let beg = crate::key::database::az::prefix(ns, db)?;
let end = crate::key::database::az::suffix(ns, db)?;
let val = self.getr(beg..end, version).await?;
return util::deserialize_cache(val.iter().map(|x| x.1.as_slice()));
}
let qey = cache::tx::Lookup::Azs(ns, db);
match self.cache.get(&qey) {
Some(val) => val.try_into_azs(),
None => {
let beg = crate::key::database::az::prefix(ns, db)?;
let end = crate::key::database::az::suffix(ns, db)?;
let val = self.getr(beg..end, None).await?;
let val = util::deserialize_cache(val.iter().map(|x| x.1.as_slice()))?;
let entry = cache::tx::Entry::Azs(Arc::clone(&val));
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "all_db_analyzers")),
)
}
fn all_db_sequences(
&self,
ns: NamespaceId,
db: DatabaseId,
version: Option<u64>,
) -> BoxProviderFut<'_, Result<Arc<[catalog::SequenceDefinition]>>> {
Box::pin(
async move {
if version.is_some() {
let beg = crate::key::database::sq::prefix(ns, db)?;
let end = crate::key::database::sq::suffix(ns, db)?;
let val = self.getr(beg..end, version).await?;
return util::deserialize_cache(val.iter().map(|x| x.1.as_slice()));
}
let qey = cache::tx::Lookup::Sqs(ns, db);
match self.cache.get(&qey) {
Some(val) => val.try_into_sqs(),
None => {
let beg = crate::key::database::sq::prefix(ns, db)?;
let end = crate::key::database::sq::suffix(ns, db)?;
let val = self.getr(beg..end, None).await?;
let val = util::deserialize_cache(val.iter().map(|x| x.1.as_slice()))?;
let entry = cache::tx::Entry::Sqs(Arc::clone(&val));
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "all_db_sequences")),
)
}
fn all_db_functions(
&self,
ns: NamespaceId,
db: DatabaseId,
version: Option<u64>,
) -> BoxProviderFut<'_, Result<Arc<[catalog::FunctionDefinition]>>> {
Box::pin(
async move {
if version.is_some() {
let beg = crate::key::database::fc::prefix(ns, db)?;
let end = crate::key::database::fc::suffix(ns, db)?;
let val = self.getr(beg..end, version).await?;
return util::deserialize_cache(val.iter().map(|x| x.1.as_slice()));
}
let qey = cache::tx::Lookup::Fcs(ns, db);
match self.cache.get(&qey) {
Some(val) => val.try_into_fcs(),
None => {
let beg = crate::key::database::fc::prefix(ns, db)?;
let end = crate::key::database::fc::suffix(ns, db)?;
let val = self.getr(beg..end, None).await?;
let val = util::deserialize_cache(val.iter().map(|x| x.1.as_slice()))?;
let entry = cache::tx::Entry::Fcs(Arc::clone(&val));
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "all_db_functions")),
)
}
fn all_db_modules(
&self,
ns: NamespaceId,
db: DatabaseId,
version: Option<u64>,
) -> BoxProviderFut<'_, Result<Arc<[catalog::ModuleDefinition]>>> {
Box::pin(
async move {
if version.is_some() {
let beg = crate::key::database::md::prefix(ns, db)?;
let end = crate::key::database::md::suffix(ns, db)?;
let val = self.getr(beg..end, version).await?;
return util::deserialize_cache(val.iter().map(|x| x.1.as_slice()));
}
let qey = cache::tx::Lookup::Mds(ns, db);
match self.cache.get(&qey) {
Some(val) => val.try_into_mds(),
None => {
let beg = crate::key::database::md::prefix(ns, db)?;
let end = crate::key::database::md::suffix(ns, db)?;
let val = self.getr(beg..end, None).await?;
let val = util::deserialize_cache(val.iter().map(|x| x.1.as_slice()))?;
let entry = cache::tx::Entry::Mds(Arc::clone(&val));
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "all_db_modules")),
)
}
fn all_db_params(
&self,
ns: NamespaceId,
db: DatabaseId,
version: Option<u64>,
) -> BoxProviderFut<'_, Result<Arc<[catalog::ParamDefinition]>>> {
Box::pin(
async move {
if version.is_some() {
let beg = crate::key::database::pa::prefix(ns, db)?;
let end = crate::key::database::pa::suffix(ns, db)?;
let val = self.getr(beg..end, version).await?;
return util::deserialize_cache(val.iter().map(|x| x.1.as_slice()));
}
let qey = cache::tx::Lookup::Pas(ns, db);
match self.cache.get(&qey) {
Some(val) => val.try_into_pas(),
None => {
let beg = crate::key::database::pa::prefix(ns, db)?;
let end = crate::key::database::pa::suffix(ns, db)?;
let val = self.getr(beg..end, None).await?;
let val = util::deserialize_cache(val.iter().map(|x| x.1.as_slice()))?;
let entry = cache::tx::Entry::Pas(Arc::clone(&val));
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "all_db_params")),
)
}
fn all_db_models(
&self,
ns: NamespaceId,
db: DatabaseId,
version: Option<u64>,
) -> BoxProviderFut<'_, Result<Arc<[catalog::MlModelDefinition]>>> {
Box::pin(
async move {
if version.is_some() {
let beg = crate::key::database::ml::prefix(ns, db)?;
let end = crate::key::database::ml::suffix(ns, db)?;
let val = self.getr(beg..end, version).await?;
return util::deserialize_cache(val.iter().map(|x| x.1.as_slice()));
}
let qey = cache::tx::Lookup::Mls(ns, db);
match self.cache.get(&qey) {
Some(val) => val.try_into_mls(),
None => {
let beg = crate::key::database::ml::prefix(ns, db)?;
let end = crate::key::database::ml::suffix(ns, db)?;
let val = self.getr(beg..end, None).await?;
let val = util::deserialize_cache(val.iter().map(|x| x.1.as_slice()))?;
let entry = cache::tx::Entry::Mls(Arc::clone(&val));
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "all_db_models")),
)
}
fn all_db_configs(
&self,
ns: NamespaceId,
db: DatabaseId,
version: Option<u64>,
) -> BoxProviderFut<'_, Result<Arc<[ConfigDefinition]>>> {
Box::pin(
async move {
if version.is_some() {
let beg = crate::key::database::cg::prefix(ns, db)?;
let end = crate::key::database::cg::suffix(ns, db)?;
let val = self.getr(beg..end, version).await?;
return util::deserialize_cache(val.iter().map(|x| x.1.as_slice()));
}
let qey = cache::tx::Lookup::Cgs(ns, db);
match self.cache.get(&qey) {
Some(val) => val.try_into_cgs(),
None => {
let beg = crate::key::database::cg::prefix(ns, db)?;
let end = crate::key::database::cg::suffix(ns, db)?;
let val = self.getr(beg..end, None).await?;
let val = util::deserialize_cache(val.iter().map(|x| x.1.as_slice()))?;
let entry = cache::tx::Entry::Cgs(Arc::clone(&val));
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "all_db_configs")),
)
}
fn get_db_model<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
ml: &'a str,
vn: &'a str,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Option<Arc<catalog::MlModelDefinition>>>> {
Box::pin(
async move {
if version.is_some() {
let key = crate::key::database::ml::new(ns, db, ml, vn);
let Some(val) = self.get(&key, version).await? else {
return Ok(None);
};
return Ok(Some(Arc::new(val)));
}
let qey = cache::tx::Lookup::Ml(ns, db, ml, vn);
match self.cache.get(&qey) {
Some(val) => val.try_into_type().map(Some),
None => {
let key = crate::key::database::ml::new(ns, db, ml, vn);
let Some(val) = self.get(&key, None).await? else {
return Ok(None);
};
let val = Arc::new(val);
let entry = cache::tx::Entry::Any(val.clone());
self.cache.insert(qey, entry);
Ok(Some(val))
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "get_db_model")),
)
}
fn get_db_analyzer<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
az: &'a str,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Arc<catalog::AnalyzerDefinition>>> {
Box::pin(
async move {
if version.is_some() {
let key = crate::key::database::az::new(ns, db, az);
let val = self.get(&key, version).await?.ok_or_else(|| Error::AzNotFound {
name: az.to_owned(),
})?;
return Ok(Arc::new(val));
}
let qey = cache::tx::Lookup::Az(ns, db, az);
match self.cache.get(&qey) {
Some(val) => val.try_into_type(),
None => {
let key = crate::key::database::az::new(ns, db, az);
let val = self.get(&key, None).await?.ok_or_else(|| Error::AzNotFound {
name: az.to_owned(),
})?;
let val = Arc::new(val);
let entry = cache::tx::Entry::Any(val.clone());
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "get_db_analyzer")),
)
}
fn get_db_sequence<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
sq: &'a str,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Arc<catalog::SequenceDefinition>>> {
Box::pin(
async move {
if version.is_some() {
let key = Sq::new(ns, db, sq);
let val = self.get(&key, version).await?.ok_or_else(|| Error::SeqNotFound {
name: sq.to_owned(),
})?;
return Ok(Arc::new(val));
}
let qey = cache::tx::Lookup::Sq(ns, db, sq);
match self.cache.get(&qey) {
Some(val) => val.try_into_type(),
None => {
let key = Sq::new(ns, db, sq);
let val =
self.get(&key, None).await?.ok_or_else(|| Error::SeqNotFound {
name: sq.to_owned(),
})?;
let val = Arc::new(val);
let entry = cache::tx::Entry::Any(val.clone());
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "get_db_sequence")),
)
}
fn get_db_function<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
fc: &'a str,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Arc<catalog::FunctionDefinition>>> {
Box::pin(
async move {
if version.is_some() {
let key = crate::key::database::fc::new(ns, db, fc);
let val = self.get(&key, version).await?.ok_or_else(|| Error::FcNotFound {
name: format!("fn::{fc}"),
})?;
return Ok(Arc::new(val));
}
let qey = cache::tx::Lookup::Fc(ns, db, fc);
match self.cache.get(&qey) {
Some(val) => val.try_into_type(),
None => {
let key = crate::key::database::fc::new(ns, db, fc);
let val = self.get(&key, None).await?.ok_or_else(|| Error::FcNotFound {
name: format!("fn::{fc}"),
})?;
let val = Arc::new(val);
let entry = cache::tx::Entry::Any(val.clone());
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "get_db_function")),
)
}
fn get_db_module<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
md: &'a str,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Arc<catalog::ModuleDefinition>>> {
Box::pin(
async move {
if version.is_some() {
let key = crate::key::database::md::new(ns, db, md);
let val = self.get(&key, version).await?.ok_or_else(|| Error::MdNotFound {
name: md.to_owned(),
})?;
return Ok(Arc::new(val));
}
let qey = cache::tx::Lookup::Md(ns, db, md);
match self.cache.get(&qey) {
Some(val) => val.try_into_type(),
None => {
let key = crate::key::database::md::new(ns, db, md);
let val = self.get(&key, None).await?.ok_or_else(|| Error::MdNotFound {
name: md.to_owned(),
})?;
let val = Arc::new(val);
let entry = cache::tx::Entry::Any(val.clone());
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "get_db_module")),
)
}
fn get_db_param<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
pa: &'a str,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Arc<catalog::ParamDefinition>>> {
Box::pin(
async move {
if version.is_some() {
let key = crate::key::database::pa::new(ns, db, pa);
let val = self.get(&key, version).await?.ok_or_else(|| Error::PaNotFound {
name: pa.to_owned(),
})?;
return Ok(Arc::new(val));
}
let qey = cache::tx::Lookup::Pa(ns, db, pa);
match self.cache.get(&qey) {
Some(val) => val.try_into_type(),
None => {
let key = crate::key::database::pa::new(ns, db, pa);
let val = self.get(&key, None).await?.ok_or_else(|| Error::PaNotFound {
name: pa.to_owned(),
})?;
let val = Arc::new(val);
let entry = cache::tx::Entry::Any(val.clone());
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "get_db_param")),
)
}
fn get_db_config<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
cg: &'a str,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Option<Arc<ConfigDefinition>>>> {
Box::pin(
async move {
if version.is_some() {
let key = crate::key::database::cg::new(ns, db, cg);
if let Some(val) = self.get(&key, version).await? {
return Ok(Some(Arc::new(val)));
} else {
return Ok(None);
}
}
let qey = cache::tx::Lookup::Cg(ns, db, cg);
match self.cache.get(&qey) {
Some(val) => val.try_into_type().map(Option::Some),
None => {
let key = crate::key::database::cg::new(ns, db, cg);
if let Some(val) = self.get(&key, None).await? {
let val = Arc::new(val);
let entr = cache::tx::Entry::Any(val.clone());
self.cache.insert(qey, entr);
Ok(Some(val))
} else {
Ok(None)
}
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "get_db_config")),
)
}
fn put_db_function<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
fc: &'a catalog::FunctionDefinition,
) -> BoxProviderFut<'a, Result<()>> {
Box::pin(async move {
let key = crate::key::database::fc::new(ns, db, &fc.name);
self.set(&key, fc).await?;
let list_key = cache::tx::Lookup::Fcs(ns, db);
self.cache.remove(&list_key);
let qey = cache::tx::Lookup::Fc(ns, db, &fc.name);
let entry = cache::tx::Entry::Any(Arc::new(fc.clone()));
self.cache.insert(qey, entry);
Ok(())
})
}
fn put_db_module<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
md: &'a catalog::ModuleDefinition,
) -> BoxProviderFut<'a, Result<()>> {
Box::pin(async move {
let name = md.get_storage_name()?;
let key = crate::key::database::md::new(ns, db, &name);
self.set(&key, md).await?;
let list_key = cache::tx::Lookup::Mds(ns, db);
self.cache.remove(&list_key);
let qey = cache::tx::Lookup::Md(ns, db, &name);
let entry = cache::tx::Entry::Any(Arc::new(md.clone()));
self.cache.insert(qey, entry);
Ok(())
})
}
fn put_db_param<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
pa: &'a catalog::ParamDefinition,
) -> BoxProviderFut<'a, Result<()>> {
Box::pin(async move {
let key = crate::key::database::pa::new(ns, db, &pa.name);
self.set(&key, pa).await?;
let list_key = cache::tx::Lookup::Pas(ns, db);
self.cache.remove(&list_key);
let qey = cache::tx::Lookup::Pa(ns, db, &pa.name);
let entry = cache::tx::Entry::Any(Arc::new(pa.clone()));
self.cache.insert(qey, entry);
Ok(())
})
}
}
impl TableProvider for Transaction {
fn all_tb(
&self,
ns: NamespaceId,
db: DatabaseId,
version: Option<u64>,
) -> BoxProviderFut<'_, Result<Arc<[TableDefinition]>>> {
Box::pin(
async move {
if version.is_some() {
let beg = crate::key::database::tb::prefix(ns, db)?;
let end = crate::key::database::tb::suffix(ns, db)?;
let val = self.getr(beg..end, version).await?;
return util::deserialize_cache(val.iter().map(|x| x.1.as_slice()));
}
let qey = cache::tx::Lookup::Tbs(ns, db);
match self.cache.get(&qey) {
Some(val) => val.try_into_tbs(),
None => {
let beg = crate::key::database::tb::prefix(ns, db)?;
let end = crate::key::database::tb::suffix(ns, db)?;
let val = self.getr(beg..end, None).await?;
let val = util::deserialize_cache(val.iter().map(|x| x.1.as_slice()))?;
let entry = cache::tx::Entry::Tbs(Arc::clone(&val));
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "all_tb")),
)
}
fn all_tb_views<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
tb: &'a TableName,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Arc<[catalog::TableDefinition]>>> {
Box::pin(
async move {
if version.is_some() {
let beg = crate::key::table::ft::prefix(ns, db, tb)?;
let end = crate::key::table::ft::suffix(ns, db, tb)?;
let val = self.getr(beg..end, version).await?;
return util::deserialize_cache(val.iter().map(|x| x.1.as_slice()));
}
let qey = cache::tx::Lookup::Fts(ns, db, tb);
match self.cache.get(&qey) {
Some(val) => val.try_into_fts(),
None => {
let beg = crate::key::table::ft::prefix(ns, db, tb)?;
let end = crate::key::table::ft::suffix(ns, db, tb)?;
let val = self.getr(beg..end, None).await?;
let val = util::deserialize_cache(val.iter().map(|x| x.1.as_slice()))?;
let entry = cache::tx::Entry::Fts(Arc::clone(&val));
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "all_tb_views")),
)
}
fn get_or_add_tb<'a>(
&'a self,
ctx: Option<&'a Context>,
ns: &'a str,
db: &'a str,
tb: &'a TableName,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Arc<TableDefinition>>> {
Box::pin(
async move {
if version.is_some() {
let Some(db_def) = self.get_db_by_name(ns, db, version).await? else {
return Err(anyhow::anyhow!(Error::DbNotFound {
name: db.to_owned(),
}));
};
let table_key =
crate::key::database::tb::new(db_def.namespace_id, db_def.database_id, tb);
if let Some(tb_def) = self.get(&table_key, version).await? {
return Ok(Arc::new(tb_def));
}
return Err(Error::TbNotFound {
name: tb.to_owned(),
}
.into());
}
let qey = cache::tx::Lookup::TbByName(ns, db, tb);
match self.cache.get(&qey) {
Some(val) => val.try_into_type(),
None => {
let Some(db_def) = self.get_db_by_name(ns, db, None).await? else {
return Err(anyhow::anyhow!(Error::DbNotFound {
name: db.to_owned(),
}));
};
let table_key = crate::key::database::tb::new(
db_def.namespace_id,
db_def.database_id,
tb,
);
if let Some(tb_def) = self.get(&table_key, None).await? {
let cached_tb = Arc::new(tb_def);
let cached_entry = cache::tx::Entry::Any(
Arc::clone(&cached_tb) as Arc<dyn Any + Send + Sync>
);
self.cache.insert(qey, cached_entry);
return Ok(cached_tb);
}
if db_def.strict {
return Err(Error::TbNotFound {
name: tb.to_owned(),
}
.into());
}
let tb_def = TableDefinition::new(
db_def.namespace_id,
db_def.database_id,
self.get_next_tb_id(ctx, db_def.namespace_id, db_def.database_id)
.await?,
tb.clone(),
);
self.put_tb(ns, db, &tb_def).await
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "get_or_add_tb")),
)
}
fn get_tb_by_name<'a>(
&'a self,
ns: &'a str,
db: &'a str,
tb: &'a TableName,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Option<Arc<TableDefinition>>>> {
Box::pin(async move {
if version.is_some() {
let Some(db) = self.get_db_by_name(ns, db, version).await? else {
return Ok(None);
};
let key = crate::key::database::tb::new(db.namespace_id, db.database_id, tb);
let Some(tb) = self.get(&key, version).await? else {
return Ok(None);
};
return Ok(Some(Arc::new(tb)));
}
let qey = cache::tx::Lookup::TbByName(ns, db, tb);
match self.cache.get(&qey) {
Some(val) => val.try_into_type().map(Some),
None => {
let Some(db) = self.get_db_by_name(ns, db, None).await? else {
return Ok(None);
};
let key = crate::key::database::tb::new(db.namespace_id, db.database_id, tb);
let Some(tb) = self.get(&key, None).await? else {
return Ok(None);
};
let tb = Arc::new(tb);
let entr = cache::tx::Entry::Any(tb.clone());
self.cache.insert(qey, entr);
Ok(Some(tb))
}
}
})
}
fn put_tb<'a>(
&'a self,
ns: &'a str,
db: &'a str,
tb: &'a TableDefinition,
) -> BoxProviderFut<'a, Result<Arc<TableDefinition>>> {
Box::pin(async move {
let key = crate::key::database::tb::new(tb.namespace_id, tb.database_id, &tb.name);
match self.set(&key, tb).await {
Ok(_) => {}
Err(e) => {
if matches!(
e.downcast_ref(),
Some(Error::Kvs(crate::kvs::Error::TransactionReadonly))
) {
return Err(Error::TbNotFound {
name: tb.name.clone(),
}
.into());
}
return Err(e);
}
}
let list_key = cache::tx::Lookup::Tbs(tb.namespace_id, tb.database_id);
self.cache.remove(&list_key);
let cached_tb = Arc::new(tb.clone());
let cached_entry =
cache::tx::Entry::Any(Arc::clone(&cached_tb) as Arc<dyn Any + Send + Sync>);
let qey = cache::tx::Lookup::Tb(tb.namespace_id, tb.database_id, &tb.name);
self.cache.insert(qey, cached_entry.clone());
let qey = cache::tx::Lookup::TbByName(ns, db, &tb.name);
self.cache.insert(qey, cached_entry);
Ok(cached_tb)
})
}
fn del_tb<'a>(
&'a self,
ns: &'a str,
db: &'a str,
tb: &'a TableName,
) -> BoxProviderFut<'a, Result<()>> {
Box::pin(async move {
let Some(tb) = self.get_tb_by_name(ns, db, tb, None).await? else {
return Err(Error::TbNotFound {
name: tb.clone(),
}
.into());
};
let key = crate::key::database::tb::new(tb.namespace_id, tb.database_id, &tb.name);
self.del(&key).await?;
let list_key = cache::tx::Lookup::Tbs(tb.namespace_id, tb.database_id);
self.cache.remove(&list_key);
let qey = cache::tx::Lookup::Tb(tb.namespace_id, tb.database_id, &tb.name);
self.cache.remove(&qey);
let qey = cache::tx::Lookup::TbByName(ns, db, &tb.name);
self.cache.remove(&qey);
Ok(())
})
}
fn clr_tb<'a>(
&'a self,
ns: &'a str,
db: &'a str,
tb: &'a TableName,
) -> BoxProviderFut<'a, Result<()>> {
Box::pin(async move {
let Some(tb) = self.get_tb_by_name(ns, db, tb, None).await? else {
return Err(Error::TbNotFound {
name: tb.clone(),
}
.into());
};
let key = crate::key::database::tb::new(tb.namespace_id, tb.database_id, &tb.name);
self.clr(&key).await?;
let list_key = cache::tx::Lookup::Tbs(tb.namespace_id, tb.database_id);
self.cache.remove(&list_key);
let qey = cache::tx::Lookup::Tb(tb.namespace_id, tb.database_id, &tb.name);
self.cache.remove(&qey);
let qey = cache::tx::Lookup::TbByName(ns, db, &tb.name);
self.cache.remove(&qey);
Ok(())
})
}
fn all_tb_events<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
tb: &'a TableName,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Arc<[catalog::EventDefinition]>>> {
Box::pin(
async move {
if version.is_some() {
let beg = crate::key::table::ev::prefix(ns, db, tb)?;
let end = crate::key::table::ev::suffix(ns, db, tb)?;
let val = self.getr(beg..end, version).await?;
return util::deserialize_cache(val.iter().map(|x| x.1.as_slice()));
}
let qey = cache::tx::Lookup::Evs(ns, db, tb);
match self.cache.get(&qey) {
Some(val) => val.try_into_evs(),
None => {
let beg = crate::key::table::ev::prefix(ns, db, tb)?;
let end = crate::key::table::ev::suffix(ns, db, tb)?;
let val = self.getr(beg..end, None).await?;
let val = util::deserialize_cache(val.iter().map(|x| x.1.as_slice()))?;
let entry = cache::tx::Entry::Evs(Arc::clone(&val));
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "all_tb_events")),
)
}
fn all_tb_fields<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
tb: &'a TableName,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Arc<[catalog::FieldDefinition]>>> {
Box::pin(
async move {
if version.is_some() {
let beg = crate::key::table::fd::prefix(ns, db, tb)?;
let end = crate::key::table::fd::suffix(ns, db, tb)?;
let val = self.getr(beg..end, version).await?;
return util::deserialize_cache(val.iter().map(|x| x.1.as_slice()));
}
let qey = cache::tx::Lookup::Fds(ns, db, tb);
match self.cache.get(&qey) {
Some(val) => val.try_into_fds(),
None => {
let beg = crate::key::table::fd::prefix(ns, db, tb)?;
let end = crate::key::table::fd::suffix(ns, db, tb)?;
let val = self.getr(beg..end, None).await?;
let val = util::deserialize_cache(val.iter().map(|x| x.1.as_slice()))?;
let entry = cache::tx::Entry::Fds(Arc::clone(&val));
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "all_tb_fields")),
)
}
fn all_tb_indexes<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
tb: &'a TableName,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Arc<[catalog::IndexDefinition]>>> {
Box::pin(
async move {
if version.is_some() {
let beg = table_ix::prefix(ns, db, tb)?;
let end = table_ix::suffix(ns, db, tb)?;
let val = self.getr(beg..end, version).await?;
return util::deserialize_cache(val.iter().map(|x| x.1.as_slice()));
}
let qey = cache::tx::Lookup::Ixs(ns, db, tb);
match self.cache.get(&qey) {
Some(val) => val.try_into_ixs(),
None => {
let beg = table_ix::prefix(ns, db, tb)?;
let end = table_ix::suffix(ns, db, tb)?;
let val = self.getr(beg..end, None).await?;
let val = util::deserialize_cache(val.iter().map(|x| x.1.as_slice()))?;
let entry = cache::tx::Entry::Ixs(Arc::clone(&val));
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "all_tb_indexes")),
)
}
fn all_tb_lives<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
tb: &'a TableName,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Arc<[catalog::SubscriptionDefinition]>>> {
Box::pin(
async move {
if version.is_some() {
let beg = crate::key::table::lq::prefix(ns, db, tb)?;
let end = crate::key::table::lq::suffix(ns, db, tb)?;
let val = self.getr(beg..end, version).await?;
return util::deserialize_cache(val.iter().map(|x| x.1.as_slice()));
}
let qey = cache::tx::Lookup::Lvs(ns, db, tb);
match self.cache.get(&qey) {
Some(val) => val.try_into_lvs(),
None => {
let beg = crate::key::table::lq::prefix(ns, db, tb)?;
let end = crate::key::table::lq::suffix(ns, db, tb)?;
let val = self.getr(beg..end, None).await?;
let val = util::deserialize_cache(val.iter().map(|x| x.1.as_slice()))?;
let entry = cache::tx::Entry::Lvs(Arc::clone(&val));
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "all_tb_lives")),
)
}
fn get_tb<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
tb: &'a TableName,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Option<Arc<TableDefinition>>>> {
Box::pin(
async move {
if version.is_some() {
let key = crate::key::database::tb::new(ns, db, tb);
let Some(val) = self.get(&key, version).await? else {
return Ok(None);
};
return Ok(Some(Arc::new(val)));
}
let qey = cache::tx::Lookup::Tb(ns, db, tb);
match self.cache.get(&qey) {
Some(val) => val.try_into_type().map(Some),
None => {
let key = crate::key::database::tb::new(ns, db, tb);
let Some(val) = self.get(&key, None).await? else {
return Ok(None);
};
let val = Arc::new(val);
let entry = cache::tx::Entry::Any(val.clone());
self.cache.insert(qey, entry);
Ok(Some(val))
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "get_tb")),
)
}
fn get_tb_event<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
tb: &'a TableName,
ev: &'a str,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Arc<catalog::EventDefinition>>> {
Box::pin(
async move {
if version.is_some() {
let key = crate::key::table::ev::new(ns, db, tb, ev);
let val = self.get(&key, version).await?.ok_or_else(|| Error::EvNotFound {
name: ev.to_owned(),
})?;
return Ok(Arc::new(val));
}
let qey = cache::tx::Lookup::Ev(ns, db, tb, ev);
match self.cache.get(&qey) {
Some(val) => val.try_into_type(),
None => {
let key = crate::key::table::ev::new(ns, db, tb, ev);
let val = self.get(&key, None).await?.ok_or_else(|| Error::EvNotFound {
name: ev.to_owned(),
})?;
let val = Arc::new(val);
let entry = cache::tx::Entry::Any(val.clone());
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "get_tb_event")),
)
}
fn get_tb_field<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
tb: &'a TableName,
fd: &'a str,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Option<Arc<catalog::FieldDefinition>>>> {
Box::pin(
async move {
if version.is_some() {
let key = crate::key::table::fd::new(ns, db, tb, fd);
let Some(val) = self.get(&key, version).await? else {
return Ok(None);
};
return Ok(Some(Arc::new(val)));
}
let qey = cache::tx::Lookup::Fd(ns, db, tb, fd);
match self.cache.get(&qey) {
Some(val) => val.try_into_type().map(Some),
None => {
let key = crate::key::table::fd::new(ns, db, tb, fd);
let Some(val) = self.get(&key, None).await? else {
return Ok(None);
};
let val = Arc::new(val);
let entry = cache::tx::Entry::Any(val.clone());
self.cache.insert(qey, entry);
Ok(Some(val))
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "get_tb_field")),
)
}
fn put_tb_field<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
tb: &'a TableName,
fd: &'a catalog::FieldDefinition,
) -> BoxProviderFut<'a, Result<()>> {
Box::pin(async move {
let name = fd.name.to_raw_string();
let key = crate::key::table::fd::new(ns, db, tb, &name);
self.set(&key, fd).await?;
let list_key = cache::tx::Lookup::Fds(ns, db, tb.as_ref());
self.cache.remove(&list_key);
self.cache.remove(&cache::tx::Lookup::DbReferenceTargets(ns, db));
let qey = cache::tx::Lookup::Fd(ns, db, tb, &name);
let entry = cache::tx::Entry::Any(Arc::new(fd.clone()));
self.cache.insert(qey, entry);
Ok(())
})
}
fn get_tb_index<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
tb: &'a TableName,
ix: &'a str,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Option<Arc<catalog::IndexDefinition>>>> {
Box::pin(
async move {
if version.is_some() {
let key = table_ix::new(ns, db, tb, ix);
let Some(val) = self.get(&key, version).await? else {
return Ok(None);
};
return Ok(Some(Arc::new(val)));
}
let qey = cache::tx::Lookup::Ix(ns, db, tb, ix);
match self.cache.get(&qey) {
Some(val) => val.try_into_type().map(Some),
None => {
let key = table_ix::new(ns, db, tb, ix);
let Some(val) = self.get(&key, None).await? else {
return Ok(None);
};
let val = Arc::new(val);
let entry = cache::tx::Entry::Any(val.clone());
self.cache.insert(qey, entry);
Ok(Some(val))
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "get_tb_index")),
)
}
fn get_tb_index_by_id<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
tb: &'a TableName,
ix: IndexId,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Option<Arc<catalog::IndexDefinition>>>> {
Box::pin(async move {
let key = table_ix::IndexNameLookupKey::new(ns, db, tb, ix);
let Some(index_name) = self.get(&key, version).await? else {
return Ok(None);
};
self.get_tb_index(ns, db, tb, &index_name, version).await
})
}
fn put_tb_index<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
tb: &'a TableName,
ix: &'a catalog::IndexDefinition,
) -> BoxProviderFut<'a, Result<()>> {
Box::pin(async move {
let key = table_ix::new(ns, db, tb, &ix.name);
self.set(&key, ix).await?;
let name_lookup_key = table_ix::IndexNameLookupKey::new(ns, db, tb, ix.index_id);
self.set(&name_lookup_key, &ix.name.to_string()).await?;
let list_key = cache::tx::Lookup::Ixs(ns, db, tb.as_ref());
self.cache.remove(&list_key);
let qey = cache::tx::Lookup::Ix(ns, db, tb, &ix.name);
let entry = cache::tx::Entry::Any(Arc::new(ix.clone()));
self.cache.insert(qey, entry);
Ok(())
})
}
fn del_tb_index<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
tb: &'a TableName,
ix: &'a str,
) -> BoxProviderFut<'a, Result<()>> {
Box::pin(async move {
let Some(ix) = self.get_tb_index(ns, db, tb, ix, None).await? else {
return Ok(());
};
let key = index_all::new(ns, db, tb, ix.index_id);
self.delp(&key).await?;
let key = table_ix::new(ns, db, tb, &ix.name);
self.del(&key).await?;
let name_lookup_key = table_ix::IndexNameLookupKey::new(ns, db, tb, ix.index_id);
self.del(&name_lookup_key).await?;
let list_key = cache::tx::Lookup::Ixs(ns, db, tb.as_ref());
self.cache.remove(&list_key);
let index_key = cache::tx::Lookup::Ix(ns, db, tb.as_ref(), &ix.name);
self.cache.remove(&index_key);
Ok(())
})
}
fn get_record<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
tb: &'a TableName,
id: &'a RecordIdKey,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Arc<Record>>> {
Box::pin(
async move {
if version.is_some() {
let key = crate::key::record::new(ns, db, tb, id);
match self.get(&key, version).await? {
Some(record) => Ok(record.into_read_only()),
None => Ok(Arc::new(Default::default())),
}
} else {
let qey = cache::tx::Lookup::Record(ns, db, tb, id);
match self.cache.get(&qey) {
Some(val) => val.try_into_record(),
None => {
let key = crate::key::record::new(ns, db, tb, id);
match self.get(&key, None).await? {
Some(record) => {
let record = record.into_read_only();
let entry = cache::tx::Entry::Val(Arc::clone(&record));
self.cache.insert(qey, entry);
Ok(record)
}
None => Ok(Arc::new(Default::default())),
}
}
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "get_record")),
)
}
fn get_records<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
rids: &'a [RecordId],
version: Option<u64>,
cache_policy: CachePolicy,
) -> BoxProviderFut<'a, Result<Vec<Arc<Record>>>> {
Box::pin(
async move {
if rids.is_empty() {
return Ok(Vec::new());
}
if version.is_some() {
let keys: Vec<crate::key::record::RecordKey<'_>> = rids
.iter()
.map(|rid| crate::key::record::new(ns, db, &rid.table, &rid.key))
.collect();
let values = self.getm(keys, version).await?;
let out: Vec<Arc<Record>> = values
.into_iter()
.map(|opt| match opt {
Some(record) => record.into_read_only(),
None => Arc::new(Default::default()),
})
.collect();
return Ok(out);
}
let mut out: Vec<Option<Arc<Record>>> = vec![None; rids.len()];
let mut uncached_rids: Vec<(usize, &RecordId)> = Vec::new();
for (i, rid) in rids.iter().enumerate() {
let qey = cache::tx::Lookup::Record(ns, db, rid.table.as_str(), &rid.key);
match self.cache.get(&qey) {
Some(entry) => out[i] = Some(entry.try_into_record()?),
None => uncached_rids.push((i, rid)),
}
}
if !uncached_rids.is_empty() {
let keys: Vec<crate::key::record::RecordKey<'_>> = uncached_rids
.iter()
.map(|(_, rid)| crate::key::record::new(ns, db, &rid.table, &rid.key))
.collect();
let values = self.getm(keys, None).await?;
for ((i, rid), opt) in uncached_rids.into_iter().zip(values) {
let record = match opt {
Some(record) => {
let record = record.into_read_only();
if matches!(cache_policy, CachePolicy::ReadWrite) {
let qey = cache::tx::Lookup::Record(
ns,
db,
rid.table.as_str(),
&rid.key,
);
let entry = cache::tx::Entry::Val(Arc::clone(&record));
self.cache.insert(qey, entry);
}
record
}
None => Arc::new(Default::default()),
};
out[i] = Some(record);
}
}
out.into_iter()
.map(|o| {
o.ok_or_else(|| {
Error::Internal("missing record in multi-get batch".into()).into()
})
})
.collect()
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "get_records")),
)
}
fn record_exists<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
tb: &'a TableName,
id: &'a RecordIdKey,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<bool>> {
Box::pin(async move {
let key = crate::key::record::new(ns, db, tb, id);
self.exists(&key, version).await
})
}
fn put_record<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
tb: &'a TableName,
id: &'a RecordIdKey,
record: Arc<Record>,
) -> BoxProviderFut<'a, Result<()>> {
Box::pin(
async move {
let key = crate::key::record::new(ns, db, tb, id);
self.put(&key, record.as_ref()).await?;
let qey = cache::tx::Lookup::Record(ns, db, tb, id);
self.cache.insert(qey, cache::tx::Entry::Val(record));
Ok(())
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "put_record")),
)
}
fn set_record<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
tb: &'a TableName,
id: &'a RecordIdKey,
record: Arc<Record>,
) -> BoxProviderFut<'a, Result<()>> {
Box::pin(
async move {
let key = crate::key::record::new(ns, db, tb, id);
self.set(&key, record.as_ref()).await?;
let qey = cache::tx::Lookup::Record(ns, db, tb, id);
self.cache.remove(&qey);
Ok(())
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "set_record")),
)
}
fn del_record<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
tb: &'a TableName,
id: &'a RecordIdKey,
) -> BoxProviderFut<'a, Result<()>> {
Box::pin(
async move {
let key = crate::key::record::new(ns, db, tb, id);
self.del(&key).await?;
let qey = cache::tx::Lookup::Record(ns, db, tb, id);
self.cache.remove(&qey);
Ok(())
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "del_record")),
)
}
fn get_next_tb_id<'a>(
&'a self,
ctx: Option<&'a Context>,
ns: NamespaceId,
db: DatabaseId,
) -> BoxProviderFut<'a, Result<TableId>> {
Box::pin(async move { self.sequences.next_table_id(ctx, ns, db).await })
}
}
impl UserProvider for Transaction {
fn all_root_users(
&self,
version: Option<u64>,
) -> BoxProviderFut<'_, Result<Arc<[catalog::UserDefinition]>>> {
Box::pin(
async move {
if version.is_some() {
let beg = crate::key::root::us::prefix();
let end = crate::key::root::us::suffix();
let val = self.getr(beg..end, version).await?;
return util::deserialize_cache(val.iter().map(|x| x.1.as_slice()));
}
let qey = cache::tx::Lookup::Rus;
match self.cache.get(&qey) {
Some(val) => val.try_into_rus(),
None => {
let beg = crate::key::root::us::prefix();
let end = crate::key::root::us::suffix();
let val = self.getr(beg..end, None).await?;
let val = util::deserialize_cache(val.iter().map(|x| x.1.as_slice()))?;
let entry = cache::tx::Entry::Rus(Arc::clone(&val));
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "all_root_users")),
)
}
fn all_ns_users(
&self,
ns: NamespaceId,
version: Option<u64>,
) -> BoxProviderFut<'_, Result<Arc<[catalog::UserDefinition]>>> {
Box::pin(
async move {
if version.is_some() {
let beg = crate::key::namespace::us::prefix(ns)?;
let end = crate::key::namespace::us::suffix(ns)?;
let val = self.getr(beg..end, version).await?;
return util::deserialize_cache(val.iter().map(|x| x.1.as_slice()));
}
let qey = cache::tx::Lookup::Nus(ns);
match self.cache.get(&qey) {
Some(val) => val.try_into_nus(),
None => {
let beg = crate::key::namespace::us::prefix(ns)?;
let end = crate::key::namespace::us::suffix(ns)?;
let val = self.getr(beg..end, None).await?;
let val = util::deserialize_cache(val.iter().map(|x| x.1.as_slice()))?;
let entry = cache::tx::Entry::Nus(Arc::clone(&val));
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "all_ns_users")),
)
}
fn all_db_users(
&self,
ns: NamespaceId,
db: DatabaseId,
version: Option<u64>,
) -> BoxProviderFut<'_, Result<Arc<[catalog::UserDefinition]>>> {
Box::pin(
async move {
if version.is_some() {
let beg = crate::key::database::us::prefix(ns, db)?;
let end = crate::key::database::us::suffix(ns, db)?;
let val = self.getr(beg..end, version).await?;
return util::deserialize_cache(val.iter().map(|x| x.1.as_slice()));
}
let qey = cache::tx::Lookup::Dus(ns, db);
match self.cache.get(&qey) {
Some(val) => val.try_into_dus(),
None => {
let beg = crate::key::database::us::prefix(ns, db)?;
let end = crate::key::database::us::suffix(ns, db)?;
let val = self.getr(beg..end, None).await?;
let val = util::deserialize_cache(val.iter().map(|x| x.1.as_slice()))?;
let entry = cache::tx::Entry::Dus(Arc::clone(&val));
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "all_db_users")),
)
}
fn get_root_user<'a>(
&'a self,
us: &'a str,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Option<Arc<catalog::UserDefinition>>>> {
Box::pin(
async move {
if version.is_some() {
let key = crate::key::root::us::new(us);
let Some(val) = self.get(&key, version).await? else {
return Ok(None);
};
return Ok(Some(Arc::new(val)));
}
let qey = cache::tx::Lookup::Ru(us);
match self.cache.get(&qey) {
Some(val) => val.try_into_type().map(Some),
None => {
let key = crate::key::root::us::new(us);
let Some(val) = self.get(&key, None).await? else {
return Ok(None);
};
let val = Arc::new(val);
let entry = cache::tx::Entry::Any(val.clone());
self.cache.insert(qey, entry);
Ok(Some(val))
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "get_root_user")),
)
}
fn get_ns_user<'a>(
&'a self,
ns: NamespaceId,
us: &'a str,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Option<Arc<catalog::UserDefinition>>>> {
Box::pin(
async move {
if version.is_some() {
let key = crate::key::namespace::us::new(ns, us);
let Some(val) = self.get(&key, version).await? else {
return Ok(None);
};
return Ok(Some(Arc::new(val)));
}
let qey = cache::tx::Lookup::Nu(ns, us);
match self.cache.get(&qey) {
Some(val) => val.try_into_type().map(Some),
None => {
let key = crate::key::namespace::us::new(ns, us);
let Some(val) = self.get(&key, None).await? else {
return Ok(None);
};
let val = Arc::new(val);
let entry = cache::tx::Entry::Any(val.clone());
self.cache.insert(qey, entry);
Ok(Some(val))
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "get_ns_user")),
)
}
fn get_db_user<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
us: &'a str,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Option<Arc<catalog::UserDefinition>>>> {
Box::pin(
async move {
if version.is_some() {
let key = crate::key::database::us::new(ns, db, us);
let Some(val) = self.get(&key, version).await? else {
return Ok(None);
};
return Ok(Some(Arc::new(val)));
}
let qey = cache::tx::Lookup::Du(ns, db, us);
match self.cache.get(&qey) {
Some(val) => val.try_into_type().map(Some),
None => {
let key = crate::key::database::us::new(ns, db, us);
let Some(val) = self.get(&key, None).await? else {
return Ok(None);
};
let val = Arc::new(val);
let entry = cache::tx::Entry::Any(val.clone());
self.cache.insert(qey, entry);
Ok(Some(val))
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "get_db_user")),
)
}
fn put_root_user<'a>(
&'a self,
us: &'a catalog::UserDefinition,
) -> BoxProviderFut<'a, Result<()>> {
Box::pin(async move {
let key = crate::key::root::us::new(&us.name);
self.set(&key, us).await?;
let list_key = cache::tx::Lookup::Rus;
self.cache.remove(&list_key);
let qey = cache::tx::Lookup::Ru(&us.name);
let entry = cache::tx::Entry::Any(Arc::new(us.clone()));
self.cache.insert(qey, entry);
Ok(())
})
}
fn put_ns_user<'a>(
&'a self,
ns: NamespaceId,
us: &'a catalog::UserDefinition,
) -> BoxProviderFut<'a, Result<()>> {
Box::pin(async move {
let key = crate::key::namespace::us::new(ns, &us.name);
self.set(&key, us).await?;
let list_key = cache::tx::Lookup::Nus(ns);
self.cache.remove(&list_key);
let qey = cache::tx::Lookup::Nu(ns, &us.name);
let entry = cache::tx::Entry::Any(Arc::new(us.clone()));
self.cache.insert(qey, entry);
Ok(())
})
}
fn put_db_user<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
us: &'a catalog::UserDefinition,
) -> BoxProviderFut<'a, Result<()>> {
Box::pin(async move {
let key = crate::key::database::us::new(ns, db, &us.name);
self.set(&key, us).await?;
let list_key = cache::tx::Lookup::Dus(ns, db);
self.cache.remove(&list_key);
let qey = cache::tx::Lookup::Du(ns, db, &us.name);
let entry = cache::tx::Entry::Any(Arc::new(us.clone()));
self.cache.insert(qey, entry);
Ok(())
})
}
}
impl AuthorisationProvider for Transaction {
fn all_root_accesses(
&self,
version: Option<u64>,
) -> BoxProviderFut<'_, Result<Arc<[catalog::AccessDefinition]>>> {
Box::pin(
async move {
if version.is_some() {
let beg = crate::key::root::ac::prefix();
let end = crate::key::root::ac::suffix();
let val = self.getr(beg..end, version).await?;
return util::deserialize_cache(val.iter().map(|x| x.1.as_slice()));
}
let qey = cache::tx::Lookup::Ras;
match self.cache.get(&qey) {
Some(val) => val.try_into_ras(),
None => {
let beg = crate::key::root::ac::prefix();
let end = crate::key::root::ac::suffix();
let val = self.getr(beg..end, None).await?;
let val = util::deserialize_cache(val.iter().map(|x| x.1.as_slice()))?;
let entry = cache::tx::Entry::Ras(Arc::clone(&val));
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "all_root_accesses")),
)
}
fn all_root_access_grants<'a>(
&'a self,
ra: &'a str,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Arc<[catalog::AccessGrant]>>> {
Box::pin(
async move {
if version.is_some() {
let beg = crate::key::root::access::gr::prefix(ra)?;
let end = crate::key::root::access::gr::suffix(ra)?;
let val = self.getr(beg..end, version).await?;
return util::deserialize_cache(val.iter().map(|x| x.1.as_slice()));
}
let qey = cache::tx::Lookup::Rgs(ra);
match self.cache.get(&qey) {
Some(val) => val.try_into_rag(),
None => {
let beg = crate::key::root::access::gr::prefix(ra)?;
let end = crate::key::root::access::gr::suffix(ra)?;
let val = self.getr(beg..end, None).await?;
let val = util::deserialize_cache(val.iter().map(|x| x.1.as_slice()))?;
let entry = cache::tx::Entry::Rag(Arc::clone(&val));
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "all_root_access_grants")),
)
}
fn all_ns_accesses(
&self,
ns: NamespaceId,
version: Option<u64>,
) -> BoxProviderFut<'_, Result<Arc<[catalog::AccessDefinition]>>> {
Box::pin(
async move {
if version.is_some() {
let beg = crate::key::namespace::ac::prefix(ns)?;
let end = crate::key::namespace::ac::suffix(ns)?;
let val = self.getr(beg..end, version).await?;
return util::deserialize_cache(val.iter().map(|x| x.1.as_slice()));
}
let qey = cache::tx::Lookup::Nas(ns);
match self.cache.get(&qey) {
Some(val) => val.try_into_nas(),
None => {
let beg = crate::key::namespace::ac::prefix(ns)?;
let end = crate::key::namespace::ac::suffix(ns)?;
let val = self.getr(beg..end, None).await?;
let val = util::deserialize_cache(val.iter().map(|x| x.1.as_slice()))?;
let entry = cache::tx::Entry::Nas(Arc::clone(&val));
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "all_ns_accesses")),
)
}
fn all_ns_access_grants<'a>(
&'a self,
ns: NamespaceId,
na: &'a str,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Arc<[catalog::AccessGrant]>>> {
Box::pin(
async move {
if version.is_some() {
let beg = crate::key::namespace::access::gr::prefix(ns, na)?;
let end = crate::key::namespace::access::gr::suffix(ns, na)?;
let val = self.getr(beg..end, version).await?;
return util::deserialize_cache(val.iter().map(|x| x.1.as_slice()));
}
let qey = cache::tx::Lookup::Ngs(ns, na);
match self.cache.get(&qey) {
Some(val) => val.try_into_nag(),
None => {
let beg = crate::key::namespace::access::gr::prefix(ns, na)?;
let end = crate::key::namespace::access::gr::suffix(ns, na)?;
let val = self.getr(beg..end, None).await?;
let val = util::deserialize_cache(val.iter().map(|x| x.1.as_slice()))?;
let entry = cache::tx::Entry::Nag(Arc::clone(&val));
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "all_ns_access_grants")),
)
}
fn all_db_accesses(
&self,
ns: NamespaceId,
db: DatabaseId,
version: Option<u64>,
) -> BoxProviderFut<'_, Result<Arc<[catalog::AccessDefinition]>>> {
Box::pin(
async move {
if version.is_some() {
let beg = crate::key::database::ac::prefix(ns, db)?;
let end = crate::key::database::ac::suffix(ns, db)?;
let val = self.getr(beg..end, version).await?;
return util::deserialize_cache(val.iter().map(|x| x.1.as_slice()));
}
let qey = cache::tx::Lookup::Das(ns, db);
match self.cache.get(&qey) {
Some(val) => val.try_into_das(),
None => {
let beg = crate::key::database::ac::prefix(ns, db)?;
let end = crate::key::database::ac::suffix(ns, db)?;
let val = self.getr(beg..end, None).await?;
let val = util::deserialize_cache(val.iter().map(|x| x.1.as_slice()))?;
let entry = cache::tx::Entry::Das(Arc::clone(&val));
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "all_db_accesses")),
)
}
fn all_db_access_grants<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
da: &'a str,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Arc<[catalog::AccessGrant]>>> {
Box::pin(
async move {
if version.is_some() {
let beg = crate::key::database::access::gr::prefix(ns, db, da)?;
let end = crate::key::database::access::gr::suffix(ns, db, da)?;
let val = self.getr(beg..end, version).await?;
return util::deserialize_cache(val.iter().map(|x| x.1.as_slice()));
}
let qey = cache::tx::Lookup::Dgs(ns, db, da);
match self.cache.get(&qey) {
Some(val) => val.try_into_dag(),
None => {
let beg = crate::key::database::access::gr::prefix(ns, db, da)?;
let end = crate::key::database::access::gr::suffix(ns, db, da)?;
let val = self.getr(beg..end, None).await?;
let val = util::deserialize_cache(val.iter().map(|x| x.1.as_slice()))?;
let entry = cache::tx::Entry::Dag(Arc::clone(&val));
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "all_db_access_grants")),
)
}
fn get_root_access<'a>(
&'a self,
ra: &'a str,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Option<Arc<catalog::AccessDefinition>>>> {
Box::pin(
async move {
if version.is_some() {
let key = crate::key::root::ac::new(ra);
let Some(val) = self.get(&key, version).await? else {
return Ok(None);
};
return Ok(Some(Arc::new(val)));
}
let qey = cache::tx::Lookup::Ra(ra);
match self.cache.get(&qey) {
Some(val) => val.try_into_type().map(Some),
None => {
let key = crate::key::root::ac::new(ra);
let Some(val) = self.get(&key, None).await? else {
return Ok(None);
};
let val = Arc::new(val);
let entry = cache::tx::Entry::Any(val.clone());
self.cache.insert(qey, entry);
Ok(Some(val))
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "get_root_access")),
)
}
fn get_root_access_grant<'a>(
&'a self,
ac: &'a str,
gr: &'a str,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Option<Arc<catalog::AccessGrant>>>> {
Box::pin(
async move {
if version.is_some() {
let key = crate::key::root::access::gr::new(ac, gr);
let Some(val) = self.get(&key, version).await? else {
return Ok(None);
};
return Ok(Some(Arc::new(val)));
}
let qey = cache::tx::Lookup::Rg(ac, gr);
match self.cache.get(&qey) {
Some(val) => val.try_into_type().map(Some),
None => {
let key = crate::key::root::access::gr::new(ac, gr);
let Some(val) = self.get(&key, None).await? else {
return Ok(None);
};
let val = Arc::new(val);
let entry = cache::tx::Entry::Any(val.clone());
self.cache.insert(qey, entry);
Ok(Some(val))
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "get_root_access_grant")),
)
}
fn get_ns_access<'a>(
&'a self,
ns: NamespaceId,
na: &'a str,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Option<Arc<catalog::AccessDefinition>>>> {
Box::pin(
async move {
if version.is_some() {
let key = crate::key::namespace::ac::new(ns, na);
let Some(val) = self.get(&key, version).await? else {
return Ok(None);
};
return Ok(Some(Arc::new(val)));
}
let qey = cache::tx::Lookup::Na(ns, na);
match self.cache.get(&qey) {
Some(val) => val.try_into_type().map(Some),
None => {
let key = crate::key::namespace::ac::new(ns, na);
let Some(val) = self.get(&key, None).await? else {
return Ok(None);
};
let val = Arc::new(val);
let entry = cache::tx::Entry::Any(val.clone());
self.cache.insert(qey, entry);
Ok(Some(val))
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "get_ns_access")),
)
}
fn get_ns_access_grant<'a>(
&'a self,
ns: NamespaceId,
ac: &'a str,
gr: &'a str,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Option<Arc<catalog::AccessGrant>>>> {
Box::pin(
async move {
if version.is_some() {
let key = crate::key::namespace::access::gr::new(ns, ac, gr);
let Some(val) = self.get(&key, version).await? else {
return Ok(None);
};
return Ok(Some(Arc::new(val)));
}
let qey = cache::tx::Lookup::Ng(ns, ac, gr);
match self.cache.get(&qey) {
Some(val) => val.try_into_type().map(Some),
None => {
let key = crate::key::namespace::access::gr::new(ns, ac, gr);
let Some(val) = self.get(&key, None).await? else {
return Ok(None);
};
let val = Arc::new(val);
let entry = cache::tx::Entry::Any(val.clone());
self.cache.insert(qey, entry);
Ok(Some(val))
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "get_ns_access_grant")),
)
}
fn get_db_access<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
da: &'a str,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Option<Arc<catalog::AccessDefinition>>>> {
Box::pin(
async move {
if version.is_some() {
let key = crate::key::database::ac::new(ns, db, da);
let Some(val) = self.get(&key, version).await? else {
return Ok(None);
};
return Ok(Some(Arc::new(val)));
}
let qey = cache::tx::Lookup::Da(ns, db, da);
match self.cache.get(&qey) {
Some(val) => val.try_into_type().map(Some),
None => {
let key = crate::key::database::ac::new(ns, db, da);
let Some(val) = self.get(&key, None).await? else {
return Ok(None);
};
let val = Arc::new(val);
let entry = cache::tx::Entry::Any(val.clone());
self.cache.insert(qey, entry);
Ok(Some(val))
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "get_db_access")),
)
}
fn get_db_access_grant<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
ac: &'a str,
gr: &'a str,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Option<Arc<catalog::AccessGrant>>>> {
Box::pin(
async move {
if version.is_some() {
let key = crate::key::database::access::gr::new(ns, db, ac, gr);
let Some(val) = self.get(&key, version).await? else {
return Ok(None);
};
return Ok(Some(Arc::new(val)));
}
let qey = cache::tx::Lookup::Dg(ns, db, ac, gr);
match self.cache.get(&qey) {
Some(val) => val.try_into_type().map(Some),
None => {
let key = crate::key::database::access::gr::new(ns, db, ac, gr);
let Some(val) = self.get(&key, None).await? else {
return Ok(None);
};
let val = Arc::new(val);
let entry = cache::tx::Entry::Any(val.clone());
self.cache.insert(qey, entry);
Ok(Some(val))
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "get_db_access_grant")),
)
}
fn del_root_access<'a>(&'a self, ra: &'a str) -> BoxProviderFut<'a, Result<()>> {
Box::pin(async move {
let key = crate::key::root::ac::new(ra);
self.del(&key).await?;
let key = crate::key::root::access::all::new(ra);
self.delp(&key).await?;
let list_key = cache::tx::Lookup::Ras;
self.cache.remove(&list_key);
let access_key = cache::tx::Lookup::Ra(ra);
self.cache.remove(&access_key);
let grants_key = cache::tx::Lookup::Rgs(ra);
self.cache.remove(&grants_key);
Ok(())
})
}
fn del_ns_access<'a>(&'a self, ns: NamespaceId, na: &'a str) -> BoxProviderFut<'a, Result<()>> {
Box::pin(async move {
let key = crate::key::namespace::ac::new(ns, na);
self.del(&key).await?;
let key = crate::key::namespace::access::all::new(ns, na);
self.delp(&key).await?;
let list_key = cache::tx::Lookup::Nas(ns);
self.cache.remove(&list_key);
let access_key = cache::tx::Lookup::Na(ns, na);
self.cache.remove(&access_key);
let grants_key = cache::tx::Lookup::Ngs(ns, na);
self.cache.remove(&grants_key);
Ok(())
})
}
fn del_db_access<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
da: &'a str,
) -> BoxProviderFut<'a, Result<()>> {
Box::pin(async move {
let key = crate::key::database::ac::new(ns, db, da);
self.del(&key).await?;
let key = crate::key::database::access::all::new(ns, db, da);
self.delp(&key).await?;
let list_key = cache::tx::Lookup::Das(ns, db);
self.cache.remove(&list_key);
let access_key = cache::tx::Lookup::Da(ns, db, da);
self.cache.remove(&access_key);
let grants_key = cache::tx::Lookup::Dgs(ns, db, da);
self.cache.remove(&grants_key);
Ok(())
})
}
}
impl ApiProvider for Transaction {
fn all_db_apis(
&self,
ns: NamespaceId,
db: DatabaseId,
version: Option<u64>,
) -> BoxProviderFut<'_, Result<Arc<[ApiDefinition]>>> {
Box::pin(
async move {
if version.is_some() {
let beg = crate::key::database::ap::prefix(ns, db)?;
let end = crate::key::database::ap::suffix(ns, db)?;
let val = self.getr(beg..end, version).await?;
return util::deserialize_cache(val.iter().map(|x| x.1.as_slice()));
}
let qey = cache::tx::Lookup::Aps(ns, db);
match self.cache.get(&qey) {
Some(val) => val,
None => {
let beg = crate::key::database::ap::prefix(ns, db)?;
let end = crate::key::database::ap::suffix(ns, db)?;
let val = self.getr(beg..end, None).await?;
let val = util::deserialize_cache(val.iter().map(|x| x.1.as_slice()))?;
let val = cache::tx::Entry::Aps(Arc::clone(&val));
self.cache.insert(qey, val.clone());
val
}
}
.try_into_aps()
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "all_db_apis")),
)
}
fn get_db_api<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
ap: &'a str,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Option<Arc<ApiDefinition>>>> {
Box::pin(
async move {
if version.is_some() {
let key = crate::key::database::ap::new(ns, db, ap);
let Some(val) = self.get(&key, version).await? else {
return Ok(None);
};
return Ok(Some(Arc::new(val)));
}
let qey = cache::tx::Lookup::Ap(ns, db, ap);
match self.cache.get(&qey) {
Some(val) => val.try_into_type().map(Some),
None => {
let key = crate::key::database::ap::new(ns, db, ap);
let Some(val) = self.get(&key, None).await? else {
return Ok(None);
};
let val = Arc::new(val);
let entry = cache::tx::Entry::Any(val.clone());
self.cache.insert(qey, entry);
Ok(Some(val))
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "get_db_api")),
)
}
fn put_db_api<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
ap: &'a catalog::ApiDefinition,
) -> BoxProviderFut<'a, Result<()>> {
Box::pin(async move {
let name = ap.path.to_string();
let key = crate::key::database::ap::new(ns, db, &name);
self.set(&key, ap).await?;
let list_key = cache::tx::Lookup::Aps(ns, db);
self.cache.remove(&list_key);
let qey = cache::tx::Lookup::Ap(ns, db, &name);
let entry = cache::tx::Entry::Any(Arc::new(ap.clone()));
self.cache.insert(qey, entry);
Ok(())
})
}
}
impl BucketProvider for Transaction {
fn all_db_buckets(
&self,
ns: NamespaceId,
db: DatabaseId,
version: Option<u64>,
) -> BoxProviderFut<'_, Result<Arc<[catalog::BucketDefinition]>>> {
Box::pin(
async move {
if version.is_some() {
let beg = crate::key::database::bu::prefix(ns, db)?;
let end = crate::key::database::bu::suffix(ns, db)?;
let val = self.getr(beg..end, version).await?;
return util::deserialize_cache(val.iter().map(|x| x.1.as_slice()));
}
let qey = cache::tx::Lookup::Bus(ns, db);
match self.cache.get(&qey) {
Some(val) => val.try_into_bus(),
None => {
let beg = crate::key::database::bu::prefix(ns, db)?;
let end = crate::key::database::bu::suffix(ns, db)?;
let val = self.getr(beg..end, None).await?;
let val = util::deserialize_cache(val.iter().map(|x| x.1.as_slice()))?;
let entry = cache::tx::Entry::Bus(Arc::clone(&val));
self.cache.insert(qey, entry);
Ok(val)
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "all_db_buckets")),
)
}
fn get_db_bucket<'a>(
&'a self,
ns: NamespaceId,
db: DatabaseId,
bu: &'a str,
version: Option<u64>,
) -> BoxProviderFut<'a, Result<Option<Arc<catalog::BucketDefinition>>>> {
Box::pin(
async move {
if version.is_some() {
let key = crate::key::database::bu::new(ns, db, bu);
let Some(val) = self.get(&key, version).await? else {
return Ok(None);
};
return Ok(Some(Arc::new(val)));
}
let qey = cache::tx::Lookup::Bu(ns, db, bu);
match self.cache.get(&qey) {
Some(val) => val.try_into_type().map(Some),
None => {
let key = crate::key::database::bu::new(ns, db, bu);
let Some(val) = self.get(&key, None).await? else {
return Ok(None);
};
let bucket_def = Arc::new(val);
let entr = cache::tx::Entry::Any(bucket_def.clone());
self.cache.insert(qey, entr);
Ok(Some(bucket_def))
}
}
}
.instrument(trace_span!(target: "surrealdb::core::kvs::tx", "get_db_bucket")),
)
}
}
impl CatalogProvider for Transaction {}