use std::{mem::take, ops::Bound, sync::Arc};
use reifydb_codec::{key::encoded::EncodedKey, row::bytes::EncodedBytes};
use reifydb_core::{
common::{CommitVersion, SourceVersion},
event::EventBus,
execution::ExecutionResult,
interface::{
WithEventBus,
catalog::{object::ObjectId, storage::StorageId},
change::{Change, ChangeOrigin, Diff},
store::{MultiVersionBatch, MultiVersionRow},
},
key::{
any::TaggedKey,
bound::TaggedKeyBoundRange,
row::{StoragePartitionedRowKey, StorageRowKey},
},
};
use reifydb_runtime::context::clock::Clock;
use reifydb_value::{Result, error::Diagnostic, params::Params, reifydb_assertions, value::identity::IdentityId};
use tracing::instrument;
use crate::{
TransactionId,
accumulator::ChangeAccumulator,
change::{RowChange, TransactionalCatalogChanges},
dictionary::DictionaryAllocatorRegistry,
error::TransactionError,
interceptor::{
WithInterceptors,
authentication::{AuthenticationPostCreateInterceptor, AuthenticationPreDeleteInterceptor},
chain::InterceptorChain as Chain,
dictionary::{
DictionaryPostCreateInterceptor, DictionaryPostUpdateInterceptor,
DictionaryPreDeleteInterceptor, DictionaryPreUpdateInterceptor,
},
dictionary_row::{
DictionaryRowPostDeleteInterceptor, DictionaryRowPostInsertInterceptor,
DictionaryRowPostUpdateInterceptor, DictionaryRowPreDeleteInterceptor,
DictionaryRowPreInsertInterceptor, DictionaryRowPreUpdateInterceptor,
},
granted_role::{GrantedRolePostCreateInterceptor, GrantedRolePreDeleteInterceptor},
identity::{IdentityPostCreateInterceptor, IdentityPreDeleteInterceptor},
identity_attribute::{IdentityAttributePostCreateInterceptor, IdentityAttributePreDeleteInterceptor},
identity_attribute_value::{
IdentityAttributeValuePostCreateInterceptor, IdentityAttributeValuePreDeleteInterceptor,
},
interceptors::Interceptors,
namespace::{
NamespacePostCreateInterceptor, NamespacePostUpdateInterceptor, NamespacePreDeleteInterceptor,
NamespacePreUpdateInterceptor,
},
ringbuffer::{
RingBufferPostCreateInterceptor, RingBufferPostUpdateInterceptor,
RingBufferPreDeleteInterceptor, RingBufferPreUpdateInterceptor,
},
ringbuffer_row::{
RingBufferRowPostDeleteInterceptor, RingBufferRowPostInsertInterceptor,
RingBufferRowPostUpdateInterceptor, RingBufferRowPreDeleteInterceptor,
RingBufferRowPreInsertInterceptor, RingBufferRowPreUpdateInterceptor,
},
role::{RolePostCreateInterceptor, RolePreDeleteInterceptor},
series::{
SeriesPostCreateInterceptor, SeriesPostUpdateInterceptor, SeriesPreDeleteInterceptor,
SeriesPreUpdateInterceptor,
},
series_row::{
SeriesRowPostDeleteInterceptor, SeriesRowPostInsertInterceptor, SeriesRowPostUpdateInterceptor,
SeriesRowPreDeleteInterceptor, SeriesRowPreInsertInterceptor, SeriesRowPreUpdateInterceptor,
},
table::{
TablePostCreateInterceptor, TablePostUpdateInterceptor, TablePreDeleteInterceptor,
TablePreUpdateInterceptor,
},
table_row::{
TableRowPostDeleteInterceptor, TableRowPostInsertInterceptor, TableRowPostUpdateInterceptor,
TableRowPreDeleteInterceptor, TableRowPreInsertInterceptor, TableRowPreUpdateInterceptor,
},
transaction::{PostCommitContext, PostCommitInterceptor, PreCommitContext, PreCommitInterceptor},
view::{
ViewPostCreateInterceptor, ViewPostUpdateInterceptor, ViewPreDeleteInterceptor,
ViewPreUpdateInterceptor,
},
},
multi::{
RangeScope,
pending::PendingWrites,
transaction::{MultiTransaction, write::MultiWriteTransaction},
},
single::{SingleTransaction, read::SingleReadTransaction, write::SingleWriteTransaction},
transaction::{RqlExecutor, Transaction, apply_pre_commit_writes, collect_transaction_writes, write::Write},
};
pub struct CommandTransaction {
pub multi: MultiTransaction,
pub single: SingleTransaction,
state: TransactionState,
pub cmd: Option<MultiWriteTransaction>,
pub event_bus: EventBus,
pub(crate) row_changes: Vec<RowChange>,
pub(crate) interceptors: Interceptors,
pub(crate) accumulator: ChangeAccumulator,
pub identity: IdentityId,
pub(crate) executor: Option<Arc<dyn RqlExecutor>>,
pub(crate) dictionary_allocators: Option<DictionaryAllocatorRegistry>,
pub(crate) clock: Clock,
poison_cause: Option<Diagnostic>,
}
#[derive(Clone, Copy, PartialEq)]
enum TransactionState {
Active,
Committed,
RolledBack,
Poisoned,
}
impl CommandTransaction {
#[instrument(name = "transaction::command::new", level = "debug", skip_all)]
pub fn new(
multi: MultiTransaction,
single: SingleTransaction,
event_bus: EventBus,
interceptors: Interceptors,
identity: IdentityId,
clock: Clock,
) -> Result<Self> {
let cmd = multi.begin_command()?;
Ok(Self {
cmd: Some(cmd),
multi,
single,
state: TransactionState::Active,
event_bus,
interceptors,
row_changes: Vec::new(),
accumulator: ChangeAccumulator::new(),
identity,
executor: None,
dictionary_allocators: None,
clock,
poison_cause: None,
})
}
pub fn set_executor(&mut self, executor: Arc<dyn RqlExecutor>) {
self.executor = Some(executor);
}
pub fn set_dictionary_allocators(&mut self, registry: DictionaryAllocatorRegistry) {
self.dictionary_allocators = Some(registry);
}
pub fn dictionary_allocators(&self) -> Option<DictionaryAllocatorRegistry> {
self.dictionary_allocators.clone()
}
pub fn rql(&mut self, rql: &str, params: Params) -> ExecutionResult {
if let Err(e) = self.check_active() {
return ExecutionResult {
frames: vec![],
error: Some(e),
metrics: Default::default(),
};
}
let executor = self.executor.clone().expect("RqlExecutor not set");
let result = executor.rql(&mut Transaction::Command(self), rql, params);
if let Some(ref e) = result.error {
self.poison(*e.0.clone());
}
result
}
#[instrument(name = "transaction::command::event_bus", level = "trace", skip(self))]
pub fn event_bus(&self) -> &EventBus {
&self.event_bus
}
fn check_active(&self) -> Result<()> {
match self.state {
TransactionState::Active => Ok(()),
TransactionState::Committed => Err(TransactionError::AlreadyCommitted.into()),
TransactionState::RolledBack => Err(TransactionError::AlreadyRolledBack.into()),
TransactionState::Poisoned => Err(TransactionError::Poisoned {
cause: Box::new(self.poison_cause.clone().unwrap()),
}
.into()),
}
}
pub(crate) fn poison(&mut self, cause: Diagnostic) {
self.state = TransactionState::Poisoned;
self.poison_cause = Some(cause);
}
#[instrument(name = "transaction::command::commit", level = "debug", skip(self))]
pub fn commit(&mut self) -> Result<CommitVersion> {
self.check_active()?;
let mut ctx = self.build_pre_commit_context()?;
self.interceptors.pre_commit.execute(&mut ctx)?;
self.finalize_commit(ctx, false)
}
#[inline]
fn build_pre_commit_context(&mut self) -> Result<PreCommitContext> {
let transaction_writes = collect_transaction_writes(self.pending_writes());
Ok(PreCommitContext {
flow_changes: self.accumulator.take_changes(CommitVersion(0), self.clock.now())?,
pending_writes: Vec::new(),
transaction_writes,
view_entries: Vec::new(),
})
}
fn finalize_commit(&mut self, ctx: PreCommitContext, unchecked: bool) -> Result<CommitVersion> {
let Some(mut multi) = self.cmd.take() else {
unreachable!("Transaction state inconsistency")
};
reifydb_assertions! {
assert!(
self.state == TransactionState::Active,
"finalize_commit entered in non-Active state; commit()/commit_unchecked() must \
pass check_active() first, otherwise this double-commits or commits a \
rolled-back/poisoned transaction"
);
}
let id = self.apply_writes_and_mark_committed(&mut multi, &ctx)?;
let row_changes = take(&mut self.row_changes);
let flow_changes = self.merge_view_entries(ctx.flow_changes, ctx.view_entries)?;
let version = self.commit_and_post(multi, id, flow_changes, row_changes, unchecked)?;
Ok(version)
}
#[inline]
fn apply_writes_and_mark_committed(
&mut self,
multi: &mut MultiWriteTransaction,
ctx: &PreCommitContext,
) -> Result<TransactionId> {
apply_pre_commit_writes(multi, &ctx.pending_writes)?;
let id = multi.id();
self.state = TransactionState::Committed;
Ok(id)
}
#[inline]
fn merge_view_entries(
&self,
mut flow_changes: Vec<Change>,
view_entries: Vec<(ObjectId, Diff)>,
) -> Result<Vec<Change>> {
if !view_entries.is_empty() {
let mut accumulator = ChangeAccumulator::new();
for (object, diff) in view_entries {
accumulator.track(object, diff);
}
let changed_at = self.clock.now();
flow_changes.extend(accumulator.take_changes(CommitVersion(0), changed_at)?);
}
Ok(flow_changes)
}
#[inline]
fn commit_and_post(
&self,
mut multi: MultiWriteTransaction,
id: TransactionId,
flow_changes: Vec<Change>,
row_changes: Vec<RowChange>,
unchecked: bool,
) -> Result<CommitVersion> {
let changes = TransactionalCatalogChanges::default();
let version = if unchecked {
multi.commit_unchecked(flow_changes)?
} else {
multi.commit(flow_changes)?
};
let _self_lease = multi.take_self_lease();
self.interceptors.post_commit.execute(PostCommitContext::new(id, version, changes, row_changes))?;
Ok(version)
}
pub fn execute_bulk_unchecked<F, R>(&mut self, body: F) -> Result<R>
where
F: FnOnce(&mut CommandTransaction) -> Result<R>,
{
self.disable_conflict_tracking()?;
let r = match body(self) {
Ok(r) => r,
Err(e) => {
let _ = self.rollback();
return Err(e);
}
};
self.commit_unchecked()?;
Ok(r)
}
#[instrument(name = "transaction::command::commit_unchecked", level = "debug", skip(self))]
pub fn commit_unchecked(&mut self) -> Result<CommitVersion> {
self.check_active()?;
let mut ctx = self.build_pre_commit_context()?;
self.interceptors.pre_commit.execute(&mut ctx)?;
self.finalize_commit(ctx, true)
}
#[instrument(name = "transaction::command::rollback", level = "debug", skip(self))]
pub fn rollback(&mut self) -> Result<()> {
self.check_active()?;
if let Some(mut multi) = self.cmd.take() {
self.state = TransactionState::RolledBack;
multi.rollback()
} else {
unreachable!("Transaction state inconsistency")
}
}
#[instrument(name = "transaction::command::pending_writes", level = "trace", skip(self))]
pub fn pending_writes(&self) -> &PendingWrites {
self.cmd.as_ref().unwrap().pending_writes()
}
#[instrument(name = "transaction::command::with_single_command", level = "trace", skip(self, keys, f))]
pub fn with_single_command<'a, I, F, R>(&self, keys: I, f: F) -> Result<R>
where
I: IntoIterator<Item = &'a EncodedKey> + Send,
F: FnOnce(&mut SingleWriteTransaction<'_>) -> Result<R> + Send,
R: Send,
{
self.check_active()?;
self.single.with_command(keys, f)
}
#[instrument(name = "transaction::command::begin_single_query", level = "trace", skip(self, keys))]
pub fn begin_single_query<'a, I>(&self, keys: I) -> Result<SingleReadTransaction<'_>>
where
I: IntoIterator<Item = &'a EncodedKey>,
{
self.check_active()?;
self.single.begin_query(keys)
}
#[instrument(name = "transaction::command::begin_single_command", level = "trace", skip(self, keys))]
pub fn begin_single_command<'a, I>(&self, keys: I) -> Result<SingleWriteTransaction<'_>>
where
I: IntoIterator<Item = &'a EncodedKey>,
{
self.check_active()?;
self.single.begin_command(keys)
}
pub fn track_row_change(&mut self, changes: &[RowChange]) {
self.row_changes.extend_from_slice(changes);
}
pub fn track_flow_change(&mut self, change: Change) {
if let ChangeOrigin::Object(id) = change.origin {
for diff in change.diffs {
self.accumulator.track(id, diff);
}
}
}
#[inline]
pub fn version(&self) -> CommitVersion {
self.cmd.as_ref().unwrap().version()
}
#[inline]
pub fn id(&self) -> TransactionId {
self.cmd.as_ref().unwrap().id()
}
#[inline]
pub fn get<K: Into<TaggedKey> + Clone>(&mut self, key: &K) -> Result<Option<MultiVersionRow<TaggedKey>>> {
self.check_active()?;
Ok(self.cmd.as_mut().unwrap().get(key)?.map(|v| v.into_multi_version_row()))
}
#[inline]
pub fn get_committed<K: Into<TaggedKey> + Clone>(
&mut self,
key: &K,
) -> Result<Option<MultiVersionRow<TaggedKey>>> {
self.check_active()?;
Ok(self.cmd.as_mut().unwrap().get_committed(key)?.map(|v| v.into_multi_version_row()))
}
#[inline]
pub fn contains<K: Into<TaggedKey> + Clone>(&mut self, key: &K) -> Result<bool> {
self.check_active()?;
self.cmd.as_mut().unwrap().contains(key)
}
#[inline]
pub fn prefix(&mut self, prefix: &EncodedKey) -> Result<MultiVersionBatch<TaggedKey>> {
self.check_active()?;
self.cmd.as_mut().unwrap().prefix(prefix)
}
#[inline]
pub fn prefix_rev(&mut self, prefix: &EncodedKey) -> Result<MultiVersionBatch<TaggedKey>> {
self.check_active()?;
self.cmd.as_mut().unwrap().prefix_rev(prefix)
}
#[inline]
pub fn read_as_of_version_exclusive(&mut self, version: CommitVersion) -> Result<()> {
self.check_active()?;
self.cmd.as_mut().unwrap().read_as_of_version_exclusive(version);
Ok(())
}
pub fn stamp_source(&mut self, source: SourceVersion) -> Result<()> {
self.check_active()?;
self.cmd.as_mut().unwrap().stamp_source(source);
Ok(())
}
#[inline]
pub fn set<K: Into<TaggedKey> + Clone>(&mut self, key: &K, bytes: impl Into<EncodedBytes>) -> Result<()> {
self.check_active()?;
self.cmd.as_mut().unwrap().set(key, bytes.into())
}
#[inline]
pub fn reserve_writes(&mut self, additional: usize) -> Result<()> {
self.check_active()?;
self.cmd.as_mut().unwrap().reserve_writes(additional);
Ok(())
}
#[inline]
pub fn disable_conflict_tracking(&mut self) -> Result<()> {
self.check_active()?;
self.cmd.as_mut().unwrap().disable_conflict_tracking();
Ok(())
}
#[inline]
pub fn remove_with_pre<K: Into<TaggedKey> + Clone>(&mut self, key: &K, pre: EncodedBytes) -> Result<()> {
self.check_active()?;
self.cmd.as_mut().unwrap().remove_with_pre(key, pre)
}
#[inline]
pub fn remove<K: Into<TaggedKey> + Clone>(&mut self, key: &K) -> Result<()> {
self.check_active()?;
self.cmd.as_mut().unwrap().remove(key)
}
#[inline]
pub fn remove_unobserved<K: Into<TaggedKey> + Clone>(&mut self, key: &K) -> Result<()> {
self.check_active()?;
self.cmd.as_mut().unwrap().remove_unobserved(key)
}
#[inline]
pub fn remove_unobserved_with_pre<K: Into<TaggedKey> + Clone>(
&mut self,
key: &K,
pre: EncodedBytes,
) -> Result<()> {
self.check_active()?;
self.cmd.as_mut().unwrap().remove_unobserved_with_pre(key, pre)
}
#[inline]
pub fn remove_silent<K: Into<TaggedKey> + Clone>(&mut self, key: &K) -> Result<()> {
self.check_active()?;
self.cmd.as_mut().unwrap().remove_silent(key)
}
#[inline]
pub fn mark_preexisting<K: Into<TaggedKey> + Clone>(&mut self, key: &K) -> Result<()> {
self.check_active()?;
self.cmd.as_mut().unwrap().mark_preexisting(key);
Ok(())
}
#[inline]
pub fn range(
&mut self,
range: TaggedKeyBoundRange,
scope: RangeScope,
batch_size: usize,
) -> Result<Box<dyn Iterator<Item = Result<MultiVersionRow<TaggedKey>>> + Send + '_>> {
self.check_active()?;
Ok(self.cmd.as_mut().unwrap().range(range, scope, batch_size))
}
pub fn range_row(
&mut self,
storage: StorageId,
start: Bound<StorageRowKey>,
end: Bound<StorageRowKey>,
scope: RangeScope,
batch_size: usize,
) -> Result<Box<dyn Iterator<Item = Result<MultiVersionRow<StorageRowKey>>> + Send + '_>> {
self.check_active()?;
Ok(self.cmd.as_mut().unwrap().range_row(storage, start, end, scope, batch_size))
}
#[inline]
pub fn range_partitioned_row(
&mut self,
storage: StorageId,
start: Bound<StoragePartitionedRowKey>,
end: Bound<StoragePartitionedRowKey>,
scope: RangeScope,
batch_size: usize,
) -> Result<Box<dyn Iterator<Item = Result<MultiVersionRow<StoragePartitionedRowKey>>> + Send + '_>> {
self.check_active()?;
Ok(self.cmd.as_mut().unwrap().range_partitioned_row(storage, start, end, scope, batch_size))
}
#[inline]
pub fn range_persistence(
&mut self,
range: TaggedKeyBoundRange,
scope: RangeScope,
batch_size: usize,
) -> Result<Box<dyn Iterator<Item = Result<MultiVersionRow<TaggedKey>>> + Send + '_>> {
self.check_active()?;
Ok(self.cmd.as_mut().unwrap().range_persistence(range, scope, batch_size))
}
#[inline]
pub fn range_rev(
&mut self,
range: TaggedKeyBoundRange,
scope: RangeScope,
batch_size: usize,
) -> Result<Box<dyn Iterator<Item = Result<MultiVersionRow<TaggedKey>>> + Send + '_>> {
self.check_active()?;
Ok(self.cmd.as_mut().unwrap().range_rev(range, scope, batch_size))
}
#[inline]
pub fn range_rev_persistence(
&mut self,
range: TaggedKeyBoundRange,
scope: RangeScope,
batch_size: usize,
) -> Result<Box<dyn Iterator<Item = Result<MultiVersionRow<TaggedKey>>> + Send + '_>> {
self.check_active()?;
Ok(self.cmd.as_mut().unwrap().range_rev_persistence(range, scope, batch_size))
}
}
impl WithEventBus for CommandTransaction {
fn event_bus(&self) -> &EventBus {
&self.event_bus
}
}
impl Write for CommandTransaction {
#[inline]
fn set(&mut self, key: &TaggedKey, bytes: EncodedBytes) -> Result<()> {
CommandTransaction::set(self, key, bytes)
}
#[inline]
fn remove_with_pre(&mut self, key: &TaggedKey, pre: EncodedBytes) -> Result<()> {
CommandTransaction::remove_with_pre(self, key, pre)
}
#[inline]
fn remove(&mut self, key: &TaggedKey) -> Result<()> {
CommandTransaction::remove(self, key)
}
#[inline]
fn mark_preexisting(&mut self, key: &TaggedKey) -> Result<()> {
CommandTransaction::mark_preexisting(self, key)
}
#[inline]
fn track_row_change(&mut self, changes: &[RowChange]) {
CommandTransaction::track_row_change(self, changes)
}
#[inline]
fn track_flow_change(&mut self, change: Change) {
CommandTransaction::track_flow_change(self, change)
}
}
impl WithInterceptors for CommandTransaction {
fn table_row_pre_insert_interceptors(&mut self) -> &mut Chain<dyn TableRowPreInsertInterceptor + Send + Sync> {
&mut self.interceptors.table_row_pre_insert
}
fn table_row_post_insert_interceptors(
&mut self,
) -> &mut Chain<dyn TableRowPostInsertInterceptor + Send + Sync> {
&mut self.interceptors.table_row_post_insert
}
fn table_row_pre_update_interceptors(&mut self) -> &mut Chain<dyn TableRowPreUpdateInterceptor + Send + Sync> {
&mut self.interceptors.table_row_pre_update
}
fn table_row_post_update_interceptors(
&mut self,
) -> &mut Chain<dyn TableRowPostUpdateInterceptor + Send + Sync> {
&mut self.interceptors.table_row_post_update
}
fn table_row_pre_delete_interceptors(&mut self) -> &mut Chain<dyn TableRowPreDeleteInterceptor + Send + Sync> {
&mut self.interceptors.table_row_pre_delete
}
fn table_row_post_delete_interceptors(
&mut self,
) -> &mut Chain<dyn TableRowPostDeleteInterceptor + Send + Sync> {
&mut self.interceptors.table_row_post_delete
}
fn ringbuffer_row_pre_insert_interceptors(
&mut self,
) -> &mut Chain<dyn RingBufferRowPreInsertInterceptor + Send + Sync> {
&mut self.interceptors.ringbuffer_row_pre_insert
}
fn ringbuffer_row_post_insert_interceptors(
&mut self,
) -> &mut Chain<dyn RingBufferRowPostInsertInterceptor + Send + Sync> {
&mut self.interceptors.ringbuffer_row_post_insert
}
fn ringbuffer_row_pre_update_interceptors(
&mut self,
) -> &mut Chain<dyn RingBufferRowPreUpdateInterceptor + Send + Sync> {
&mut self.interceptors.ringbuffer_row_pre_update
}
fn ringbuffer_row_post_update_interceptors(
&mut self,
) -> &mut Chain<dyn RingBufferRowPostUpdateInterceptor + Send + Sync> {
&mut self.interceptors.ringbuffer_row_post_update
}
fn ringbuffer_row_pre_delete_interceptors(
&mut self,
) -> &mut Chain<dyn RingBufferRowPreDeleteInterceptor + Send + Sync> {
&mut self.interceptors.ringbuffer_row_pre_delete
}
fn ringbuffer_row_post_delete_interceptors(
&mut self,
) -> &mut Chain<dyn RingBufferRowPostDeleteInterceptor + Send + Sync> {
&mut self.interceptors.ringbuffer_row_post_delete
}
fn pre_commit_interceptors(&mut self) -> &mut Chain<dyn PreCommitInterceptor + Send + Sync> {
&mut self.interceptors.pre_commit
}
fn post_commit_interceptors(&mut self) -> &mut Chain<dyn PostCommitInterceptor + Send + Sync> {
&mut self.interceptors.post_commit
}
fn namespace_post_create_interceptors(
&mut self,
) -> &mut Chain<dyn NamespacePostCreateInterceptor + Send + Sync> {
&mut self.interceptors.namespace_post_create
}
fn namespace_pre_update_interceptors(&mut self) -> &mut Chain<dyn NamespacePreUpdateInterceptor + Send + Sync> {
&mut self.interceptors.namespace_pre_update
}
fn namespace_post_update_interceptors(
&mut self,
) -> &mut Chain<dyn NamespacePostUpdateInterceptor + Send + Sync> {
&mut self.interceptors.namespace_post_update
}
fn namespace_pre_delete_interceptors(&mut self) -> &mut Chain<dyn NamespacePreDeleteInterceptor + Send + Sync> {
&mut self.interceptors.namespace_pre_delete
}
fn table_post_create_interceptors(&mut self) -> &mut Chain<dyn TablePostCreateInterceptor + Send + Sync> {
&mut self.interceptors.table_post_create
}
fn table_pre_update_interceptors(&mut self) -> &mut Chain<dyn TablePreUpdateInterceptor + Send + Sync> {
&mut self.interceptors.table_pre_update
}
fn table_post_update_interceptors(&mut self) -> &mut Chain<dyn TablePostUpdateInterceptor + Send + Sync> {
&mut self.interceptors.table_post_update
}
fn table_pre_delete_interceptors(&mut self) -> &mut Chain<dyn TablePreDeleteInterceptor + Send + Sync> {
&mut self.interceptors.table_pre_delete
}
fn view_post_create_interceptors(&mut self) -> &mut Chain<dyn ViewPostCreateInterceptor + Send + Sync> {
&mut self.interceptors.view_post_create
}
fn view_pre_update_interceptors(&mut self) -> &mut Chain<dyn ViewPreUpdateInterceptor + Send + Sync> {
&mut self.interceptors.view_pre_update
}
fn view_post_update_interceptors(&mut self) -> &mut Chain<dyn ViewPostUpdateInterceptor + Send + Sync> {
&mut self.interceptors.view_post_update
}
fn view_pre_delete_interceptors(&mut self) -> &mut Chain<dyn ViewPreDeleteInterceptor + Send + Sync> {
&mut self.interceptors.view_pre_delete
}
fn ringbuffer_post_create_interceptors(
&mut self,
) -> &mut Chain<dyn RingBufferPostCreateInterceptor + Send + Sync> {
&mut self.interceptors.ringbuffer_post_create
}
fn ringbuffer_pre_update_interceptors(
&mut self,
) -> &mut Chain<dyn RingBufferPreUpdateInterceptor + Send + Sync> {
&mut self.interceptors.ringbuffer_pre_update
}
fn ringbuffer_post_update_interceptors(
&mut self,
) -> &mut Chain<dyn RingBufferPostUpdateInterceptor + Send + Sync> {
&mut self.interceptors.ringbuffer_post_update
}
fn ringbuffer_pre_delete_interceptors(
&mut self,
) -> &mut Chain<dyn RingBufferPreDeleteInterceptor + Send + Sync> {
&mut self.interceptors.ringbuffer_pre_delete
}
fn dictionary_row_pre_insert_interceptors(
&mut self,
) -> &mut Chain<dyn DictionaryRowPreInsertInterceptor + Send + Sync> {
&mut self.interceptors.dictionary_row_pre_insert
}
fn dictionary_row_post_insert_interceptors(
&mut self,
) -> &mut Chain<dyn DictionaryRowPostInsertInterceptor + Send + Sync> {
&mut self.interceptors.dictionary_row_post_insert
}
fn dictionary_row_pre_update_interceptors(
&mut self,
) -> &mut Chain<dyn DictionaryRowPreUpdateInterceptor + Send + Sync> {
&mut self.interceptors.dictionary_row_pre_update
}
fn dictionary_row_post_update_interceptors(
&mut self,
) -> &mut Chain<dyn DictionaryRowPostUpdateInterceptor + Send + Sync> {
&mut self.interceptors.dictionary_row_post_update
}
fn dictionary_row_pre_delete_interceptors(
&mut self,
) -> &mut Chain<dyn DictionaryRowPreDeleteInterceptor + Send + Sync> {
&mut self.interceptors.dictionary_row_pre_delete
}
fn dictionary_row_post_delete_interceptors(
&mut self,
) -> &mut Chain<dyn DictionaryRowPostDeleteInterceptor + Send + Sync> {
&mut self.interceptors.dictionary_row_post_delete
}
fn dictionary_post_create_interceptors(
&mut self,
) -> &mut Chain<dyn DictionaryPostCreateInterceptor + Send + Sync> {
&mut self.interceptors.dictionary_post_create
}
fn dictionary_pre_update_interceptors(
&mut self,
) -> &mut Chain<dyn DictionaryPreUpdateInterceptor + Send + Sync> {
&mut self.interceptors.dictionary_pre_update
}
fn dictionary_post_update_interceptors(
&mut self,
) -> &mut Chain<dyn DictionaryPostUpdateInterceptor + Send + Sync> {
&mut self.interceptors.dictionary_post_update
}
fn dictionary_pre_delete_interceptors(
&mut self,
) -> &mut Chain<dyn DictionaryPreDeleteInterceptor + Send + Sync> {
&mut self.interceptors.dictionary_pre_delete
}
fn series_row_pre_insert_interceptors(
&mut self,
) -> &mut Chain<dyn SeriesRowPreInsertInterceptor + Send + Sync> {
&mut self.interceptors.series_row_pre_insert
}
fn series_row_post_insert_interceptors(
&mut self,
) -> &mut Chain<dyn SeriesRowPostInsertInterceptor + Send + Sync> {
&mut self.interceptors.series_row_post_insert
}
fn series_row_pre_update_interceptors(
&mut self,
) -> &mut Chain<dyn SeriesRowPreUpdateInterceptor + Send + Sync> {
&mut self.interceptors.series_row_pre_update
}
fn series_row_post_update_interceptors(
&mut self,
) -> &mut Chain<dyn SeriesRowPostUpdateInterceptor + Send + Sync> {
&mut self.interceptors.series_row_post_update
}
fn series_row_pre_delete_interceptors(
&mut self,
) -> &mut Chain<dyn SeriesRowPreDeleteInterceptor + Send + Sync> {
&mut self.interceptors.series_row_pre_delete
}
fn series_row_post_delete_interceptors(
&mut self,
) -> &mut Chain<dyn SeriesRowPostDeleteInterceptor + Send + Sync> {
&mut self.interceptors.series_row_post_delete
}
fn series_post_create_interceptors(&mut self) -> &mut Chain<dyn SeriesPostCreateInterceptor + Send + Sync> {
&mut self.interceptors.series_post_create
}
fn series_pre_update_interceptors(&mut self) -> &mut Chain<dyn SeriesPreUpdateInterceptor + Send + Sync> {
&mut self.interceptors.series_pre_update
}
fn series_post_update_interceptors(&mut self) -> &mut Chain<dyn SeriesPostUpdateInterceptor + Send + Sync> {
&mut self.interceptors.series_post_update
}
fn series_pre_delete_interceptors(&mut self) -> &mut Chain<dyn SeriesPreDeleteInterceptor + Send + Sync> {
&mut self.interceptors.series_pre_delete
}
fn identity_post_create_interceptors(&mut self) -> &mut Chain<dyn IdentityPostCreateInterceptor + Send + Sync> {
&mut self.interceptors.identity_post_create
}
fn identity_pre_delete_interceptors(&mut self) -> &mut Chain<dyn IdentityPreDeleteInterceptor + Send + Sync> {
&mut self.interceptors.identity_pre_delete
}
fn role_post_create_interceptors(&mut self) -> &mut Chain<dyn RolePostCreateInterceptor + Send + Sync> {
&mut self.interceptors.role_post_create
}
fn role_pre_delete_interceptors(&mut self) -> &mut Chain<dyn RolePreDeleteInterceptor + Send + Sync> {
&mut self.interceptors.role_pre_delete
}
fn granted_role_post_create_interceptors(
&mut self,
) -> &mut Chain<dyn GrantedRolePostCreateInterceptor + Send + Sync> {
&mut self.interceptors.granted_role_post_create
}
fn granted_role_pre_delete_interceptors(
&mut self,
) -> &mut Chain<dyn GrantedRolePreDeleteInterceptor + Send + Sync> {
&mut self.interceptors.granted_role_pre_delete
}
fn identity_attribute_post_create_interceptors(
&mut self,
) -> &mut Chain<dyn IdentityAttributePostCreateInterceptor + Send + Sync> {
&mut self.interceptors.identity_attribute_post_create
}
fn identity_attribute_pre_delete_interceptors(
&mut self,
) -> &mut Chain<dyn IdentityAttributePreDeleteInterceptor + Send + Sync> {
&mut self.interceptors.identity_attribute_pre_delete
}
fn identity_attribute_value_post_create_interceptors(
&mut self,
) -> &mut Chain<dyn IdentityAttributeValuePostCreateInterceptor + Send + Sync> {
&mut self.interceptors.identity_attribute_value_post_create
}
fn identity_attribute_value_pre_delete_interceptors(
&mut self,
) -> &mut Chain<dyn IdentityAttributeValuePreDeleteInterceptor + Send + Sync> {
&mut self.interceptors.identity_attribute_value_pre_delete
}
fn authentication_post_create_interceptors(
&mut self,
) -> &mut Chain<dyn AuthenticationPostCreateInterceptor + Send + Sync> {
&mut self.interceptors.authentication_post_create
}
fn authentication_pre_delete_interceptors(
&mut self,
) -> &mut Chain<dyn AuthenticationPreDeleteInterceptor + Send + Sync> {
&mut self.interceptors.authentication_pre_delete
}
}
impl Drop for CommandTransaction {
fn drop(&mut self) {
if let Some(mut multi) = self.cmd.take()
&& (self.state == TransactionState::Active || self.state == TransactionState::Poisoned)
{
let _ = multi.rollback();
}
}
}