radixdb_executor/dispatch/
transaction.rs1use crate::catalog::DdlTransaction;
4use crate::context::ExecutionContext;
5use crate::mutation::host::{ActiveTransaction, MutationHost};
6use crate::result::ExecResult;
7use radixdb_core::{Error, IsolationLevel, Result};
8use radixdb_sql::ast::{
9 BeginStatement, CommitStatement, ReleaseSavepointStatement, RollbackStatement,
10 SavepointStatement,
11};
12use radixdb_storage::traits::{Engine, QueryResult};
13
14pub fn parse_isolation_level(level: &str) -> Result<IsolationLevel> {
16 match level {
17 "READ COMMITTED" => Ok(IsolationLevel::ReadCommitted),
18 "SNAPSHOT" => Ok(IsolationLevel::SnapshotIsolation),
19 _ => Err(Error::internal(format!(
20 "unsupported isolation level: '{level}'. Supported: READ COMMITTED, SNAPSHOT"
21 ))),
22 }
23}
24
25pub trait TransactionControlExt: MutationHost {
27 fn execute_begin(
28 &self,
29 statement: &BeginStatement,
30 _context: &ExecutionContext,
31 ) -> Result<Box<dyn QueryResult>> {
32 let mut active = self.mutation_active_transaction().lock().unwrap();
33 if active.is_some() {
34 return Err(Error::TransactionAlreadyStarted);
35 }
36
37 let transaction = if let Some(level) = &statement.isolation_level {
38 self.mutation_engine()
39 .begin_transaction_with_level(parse_isolation_level(level)?)?
40 } else {
41 self.mutation_engine()
42 .begin_transaction_with_level(self.mutation_default_isolation_level())?
43 };
44 let catalog = self.mutation_engine().pin_catalog()?;
45 let catalog = DdlTransaction::begin_shared_with_plugin_registry(
46 catalog,
47 std::sync::Arc::clone(self.mutation_plugin_registry()),
48 );
49 *active = Some(ActiveTransaction::new(transaction, catalog));
50 Ok(Box::new(ExecResult::empty()))
51 }
52
53 fn execute_commit_stmt(
54 &self,
55 _statement: &CommitStatement,
56 _context: &ExecutionContext,
57 ) -> Result<Box<dyn QueryResult>> {
58 let mut active = self.mutation_active_transaction().lock().unwrap();
59 let Some(mut state) = active.take() else {
60 return Err(Error::TransactionNotStarted);
61 };
62 let _catalog_write_fence = state
67 .has_pending_catalog_changes()
68 .then(|| self.mutation_engine().acquire_catalog_write_fence());
69 if let Err(error) = state.stage_catalog_for_commit() {
70 *active = Some(state);
71 return Err(error);
72 }
73 match state.transaction.commit() {
74 Ok(()) => Ok(Box::new(ExecResult::empty())),
75 Err(error) => {
76 if state.transaction.is_active() {
77 *active = Some(state);
78 }
79 Err(error)
80 }
81 }
82 }
83
84 fn execute_rollback_stmt(
85 &self,
86 statement: &RollbackStatement,
87 _context: &ExecutionContext,
88 ) -> Result<Box<dyn QueryResult>> {
89 let mut active = self.mutation_active_transaction().lock().unwrap();
90 if let Some(savepoint) = &statement.savepoint_name {
91 let state = active.as_mut().ok_or_else(|| {
92 Error::internal("ROLLBACK TO SAVEPOINT can only be used within a transaction")
93 })?;
94 let name = if savepoint.token.quoted {
95 savepoint.value.as_str()
96 } else {
97 savepoint.value_lower.as_str()
98 };
99 state.rollback_to_savepoint(name)?;
100 return Ok(Box::new(ExecResult::empty()));
101 }
102
103 let Some(mut state) = active.take() else {
104 return Err(Error::TransactionNotStarted);
105 };
106 state.rollback()?;
107 Ok(Box::new(ExecResult::empty()))
108 }
109
110 fn execute_savepoint(
111 &self,
112 statement: &SavepointStatement,
113 _context: &ExecutionContext,
114 ) -> Result<Box<dyn QueryResult>> {
115 let mut active = self.mutation_active_transaction().lock().unwrap();
116 let state = active.as_mut().ok_or_else(|| {
117 Error::internal("SAVEPOINT can only be used within a transaction (after BEGIN)")
118 })?;
119 let name = if statement.savepoint_name.token.quoted {
120 statement.savepoint_name.value.as_str()
121 } else {
122 statement.savepoint_name.value_lower.as_str()
123 };
124 state.create_savepoint(name)?;
125 Ok(Box::new(ExecResult::empty()))
126 }
127
128 fn execute_release_savepoint(
129 &self,
130 statement: &ReleaseSavepointStatement,
131 _context: &ExecutionContext,
132 ) -> Result<Box<dyn QueryResult>> {
133 let mut active = self.mutation_active_transaction().lock().unwrap();
134 let state = active.as_mut().ok_or_else(|| {
135 Error::internal("RELEASE SAVEPOINT can only be used within a transaction (after BEGIN)")
136 })?;
137 let name = if statement.savepoint_name.token.quoted {
138 statement.savepoint_name.value.as_str()
139 } else {
140 statement.savepoint_name.value_lower.as_str()
141 };
142 state.release_savepoint(name)?;
143 Ok(Box::new(ExecResult::empty()))
144 }
145}
146
147impl<T: MutationHost + ?Sized> TransactionControlExt for T {}