Skip to main content

radixdb_executor/dispatch/
transaction.rs

1//! SQL transaction-control routing.
2
3use 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
14/// Parse the SQL spelling of a supported isolation level.
15pub 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
25/// Transaction-control phase owned by the executor crate.
26pub 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        // An explicit transaction may legitimately discover that its pinned
63        // catalog is stale. It must not, however, publish in the middle of an
64        // auto-commit writer's pin-to-commit interval and make that statement
65        // fail spuriously.
66        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 {}