Skip to main content

lora_database/
transaction.rs

1use std::collections::BTreeMap;
2use std::sync::atomic::{AtomicBool, Ordering};
3use std::sync::{Arc, Mutex, MutexGuard};
4use std::time::Duration;
5
6use web_time::Instant;
7
8use anyhow::Result;
9use thiserror::Error;
10
11use lora_analyzer::Analyzer;
12use lora_compiler::{CompiledQuery, Compiler};
13use lora_executor::{
14    classify_stream, compiled_result_columns, project_rows, ExecuteOptions, ExecutionContext,
15    Executor, LoraValue, MutableExecutionContext, MutableExecutor, MutablePullExecutor,
16    PullExecutor, QueryResult, Row, RowSource,
17};
18use lora_parser::parse_query;
19use lora_store::{InMemoryGraph, MutationEvent, MutationRecorder};
20use lora_wal::WalRecorder;
21
22use crate::changes::ChangeHub;
23use crate::error::LoraError;
24use crate::explain::{OperatorMetrics, ProfileMetrics, QueryPlan, QueryProfile};
25use crate::live_store::LiveStore;
26use crate::snapshot::ManagedSnapshotStore;
27use crate::stream::QueryStream;
28use crate::wal::write_scope::ensure_wal_not_poisoned;
29use lora_compiler::plan_tree_from_compiled;
30use lora_executor::{plan_result_columns, CollectorGuard, MetricsCollector};
31
32/// Transaction-lifecycle invariant violations.
33///
34/// All variants used to be raised as `anyhow!("...")` strings. Surfacing
35/// them as a typed enum lets [`crate::LoraError`] route them onto stable
36/// `LoraErrorCode`s without phrase-matching the `Display` text.
37#[derive(Debug, Clone, PartialEq, Eq, Error)]
38pub enum TransactionError {
39    #[error("transaction is already closed")]
40    AlreadyClosed,
41
42    #[error("transaction has no live graph guard")]
43    NoGraphGuard,
44
45    #[error("transaction has no staged graph")]
46    NoStagedGraph,
47
48    #[error("cannot commit transaction while a streaming cursor is still active")]
49    CursorActiveCommit,
50
51    #[error("cannot start a new statement while a streaming cursor is still active")]
52    CursorActiveStatement,
53
54    #[error("cannot execute mutating query in read-only transaction")]
55    ReadOnlyMutation,
56
57    #[error("streaming write cursor requires a ReadWrite transaction")]
58    StreamingRequiresReadWrite,
59
60    #[error("read-only transaction cannot publish staged graph")]
61    ReadOnlyCommit,
62
63    #[error("transaction state lock is poisoned")]
64    Poisoned,
65}
66
67/// Transaction execution mode.
68#[derive(Debug, Clone, Copy, PartialEq, Eq)]
69pub enum TransactionMode {
70    /// Use the read-only executor. Write operators return read-only errors.
71    ReadOnly,
72    /// Execute reads and writes against a staged graph, then publish on commit.
73    ReadWrite,
74}
75
76/// What a transaction holds onto for the duration of its statements.
77///
78/// `Read` simply pins an `Arc<InMemoryGraph>` snapshot — readers don't
79/// take any lock, so nothing observable changes when a writer commits
80/// a new version mid-transaction. `Write` holds the writer Mutex
81/// (serializing commit ordering) and a snapshot of the live graph at
82/// the point the transaction began; the working copy is built lazily
83/// in `TxInner::staged` on first mutation, mirroring the previous
84/// "clone on first mutation" behavior.
85pub(crate) enum LiveStoreGuard<'db> {
86    Read(Arc<InMemoryGraph>),
87    Write(WriteLease<'db>),
88}
89
90/// Writer lease for a `ReadWrite` transaction. Holds the per-database
91/// writer Mutex plus a read snapshot of the graph at lease open time.
92/// The mutating working copy lives in `TxInner::staged` and is
93/// cloned from `snapshot` lazily.
94pub(crate) struct WriteLease<'db> {
95    /// Held for the tx lifetime so concurrent ReadWrite txns serialize.
96    pub(crate) _writer_lock: MutexGuard<'db, ()>,
97    /// Pointer back to the live store so commit can publish.
98    pub(crate) store: Arc<LiveStore<InMemoryGraph>>,
99    /// Read-only view of the graph at lease open time. The first
100    /// mutating statement clones from this into `TxInner::staged`.
101    pub(crate) snapshot: Arc<InMemoryGraph>,
102}
103
104impl LiveStoreGuard<'_> {
105    fn as_graph(&self) -> &InMemoryGraph {
106        match self {
107            Self::Read(arc) => arc,
108            Self::Write(lease) => &lease.snapshot,
109        }
110    }
111}
112
113/// Captures the staged graph and tx-local mutation buffer at the point
114/// a statement is opened, so a failed/dropped statement can be rolled
115/// back to that point without affecting earlier work in the same
116/// transaction.
117pub(crate) struct Savepoint {
118    staged: Option<InMemoryGraph>,
119    buffer_len: usize,
120}
121
122/// How a transaction-bound stream finished. Exhaustion commits that
123/// statement's staged changes into the transaction; interruption means
124/// drop or runtime error before all rows were observed.
125#[derive(Debug, Clone, Copy, PartialEq, Eq)]
126pub(crate) enum TxStreamOutcome {
127    Exhausted,
128    Interrupted,
129}
130
131impl TxStreamOutcome {
132    fn should_restore_savepoint(self, rollback_on_drop: bool) -> bool {
133        matches!(self, Self::Interrupted) && rollback_on_drop
134    }
135}
136
137/// Owns the active-cursor token for a transaction-bound stream.
138///
139/// The lease is created only after a cursor has opened successfully. If stream
140/// code forgets to finalize it explicitly, `Drop` treats the cursor as
141/// interrupted so a mutating statement cannot accidentally commit partial work
142/// into the transaction.
143pub(crate) struct TxCursorLease {
144    handle: Arc<Mutex<TxInner>>,
145    rollback_on_drop: bool,
146    finalized: bool,
147}
148
149impl TxCursorLease {
150    pub(crate) fn new(handle: Arc<Mutex<TxInner>>, rollback_on_drop: bool) -> Self {
151        Self {
152            handle,
153            rollback_on_drop,
154            finalized: false,
155        }
156    }
157
158    pub(crate) fn finalize(&mut self, outcome: TxStreamOutcome) {
159        if self.finalized {
160            return;
161        }
162        finalize_tx_stream(&self.handle, outcome, self.rollback_on_drop);
163        self.finalized = true;
164    }
165}
166
167impl Drop for TxCursorLease {
168    fn drop(&mut self) {
169        self.finalize(TxStreamOutcome::Interrupted);
170    }
171}
172
173/// Buffers `MutationEvent`s emitted by the staged graph while a
174/// transaction is in progress. The buffer replaces direct WAL writes
175/// during the transaction body; on commit the host replays the
176/// buffer into the real `WalRecorder` as a single durable
177/// transaction. Statement rollback truncates the buffer back to its
178/// pre-statement length; transaction rollback drops it entirely.
179pub(crate) struct BufferingRecorder {
180    buffer: Arc<Mutex<Vec<MutationEvent>>>,
181    failed: Arc<AtomicBool>,
182}
183
184impl BufferingRecorder {
185    pub(crate) fn new(buffer: Arc<Mutex<Vec<MutationEvent>>>, failed: Arc<AtomicBool>) -> Self {
186        Self { buffer, failed }
187    }
188}
189
190impl MutationRecorder for BufferingRecorder {
191    fn record(&self, event: MutationEvent) {
192        match self.buffer.lock() {
193            Ok(mut buf) => buf.push(event),
194            Err(_) => self.failed.store(true, Ordering::Release),
195        }
196    }
197}
198
199/// Shared transaction state. Wrapped in `Arc<Mutex<>>` so a
200/// `QueryStream` opened against the transaction can release its
201/// cursor token and signal savepoint-rollback intent on drop without
202/// borrowing the [`Transaction`] handle.
203pub(crate) struct TxInner {
204    /// The cloned staging graph. Mutated by write statements through
205    /// the [`MutableExecutor`]; read by read-only statements through
206    /// [`PullExecutor`]. `None` once the transaction has been closed.
207    pub(crate) staged: Option<InMemoryGraph>,
208    /// Tx-local mutation log, populated by the [`BufferingRecorder`]
209    /// installed on `staged`. Replayed into the real WAL exactly once
210    /// at commit time.
211    pub(crate) buffer: Arc<Mutex<Vec<MutationEvent>>>,
212    /// Latched if a mutation event could not be appended to `buffer`.
213    /// Commit must fail rather than publishing a graph whose WAL event stream
214    /// is incomplete.
215    pub(crate) buffer_failed: Arc<AtomicBool>,
216    /// Per-statement savepoint snapshot. Set when a statement opens,
217    /// cleared on successful completion, restored on
218    /// failure/premature drop.
219    pub(crate) pending_savepoint: Option<Savepoint>,
220    /// True while a `QueryStream` opened against this transaction is
221    /// alive. Blocks new statements and prevents commit until the
222    /// cursor is released.
223    pub(crate) cursor_active: bool,
224    /// True after `commit` or `rollback` has run, regardless of
225    /// outcome. Subsequent operations fail loudly instead of silently
226    /// running on stale state.
227    pub(crate) closed: bool,
228    /// Transaction execution mode chosen at `begin_transaction` time.
229    pub(crate) mode: TransactionMode,
230    /// Whether this transaction needs a mutation buffer for durable WAL
231    /// replay. Databases without a WAL can skip recorder installation and
232    /// avoid cloning mutation payloads into an unused buffer.
233    pub(crate) buffer_mutations: bool,
234}
235
236/// Explicit transaction over the in-memory graph.
237///
238/// The implementation is conservative: read-only transactions hold a
239/// database read lock, and read-write transactions hold the database
240/// write lock. Read-write transactions lazily create a cloned staging
241/// graph on the first mutating statement, then either swap that graph
242/// into place on commit or drop it on rollback. Explicit mutating
243/// statements capture a graph +
244/// WAL-buffer savepoint so a failed or dropped streaming statement
245/// only rolls back its own effects, not the transaction as a whole.
246///
247/// When a WAL is attached, mutation events fire into a tx-local
248/// buffer rather than the durable log. The buffer is replayed into
249/// the WAL exactly once at commit, so recovery never observes
250/// partial / aborted / dropped statements.
251pub struct Transaction<'db> {
252    pub(crate) live: Option<LiveStoreGuard<'db>>,
253    pub(crate) inner: Arc<Mutex<TxInner>>,
254    pub(crate) wal: Option<Arc<WalRecorder>>,
255    pub(crate) snapshots: Option<Arc<ManagedSnapshotStore>>,
256    pub(crate) changes: Arc<ChangeHub>,
257    mode: TransactionMode,
258}
259
260impl<'db> Transaction<'db> {
261    /// Build a fresh transaction. Used by `Database::begin_transaction`.
262    pub(crate) fn new(
263        live: LiveStoreGuard<'db>,
264        wal: Option<Arc<WalRecorder>>,
265        snapshots: Option<Arc<ManagedSnapshotStore>>,
266        changes: Arc<ChangeHub>,
267        mode: TransactionMode,
268    ) -> Self {
269        // The buffer feeds the WAL at commit and, while a change feed is
270        // capturing, the feed's batch.
271        let buffer_mutations = wal.is_some() || changes.is_active();
272        let inner = TxInner {
273            staged: None,
274            buffer: Arc::new(Mutex::new(Vec::new())),
275            buffer_failed: Arc::new(AtomicBool::new(false)),
276            pending_savepoint: None,
277            cursor_active: false,
278            closed: false,
279            mode,
280            buffer_mutations,
281        };
282        Self {
283            live: Some(live),
284            inner: Arc::new(Mutex::new(inner)),
285            wal,
286            snapshots,
287            changes,
288            mode,
289        }
290    }
291
292    /// Transaction mode chosen at begin time.
293    pub fn mode(&self) -> TransactionMode {
294        self.mode
295    }
296
297    /// Execute a query inside the transaction and return a materialized
298    /// `QueryResult`.
299    pub fn execute(
300        &mut self,
301        query: &str,
302        options: Option<ExecuteOptions>,
303    ) -> Result<QueryResult, LoraError> {
304        self.execute_with_params(query, options, BTreeMap::new())
305    }
306
307    /// Execute a query inside the transaction with a cooperative deadline.
308    pub fn execute_with_timeout(
309        &mut self,
310        query: &str,
311        options: Option<ExecuteOptions>,
312        timeout: Duration,
313    ) -> Result<QueryResult, LoraError> {
314        let rows = self.execute_rows_with_params_deadline(
315            query,
316            BTreeMap::new(),
317            Some(deadline_after(timeout)),
318        )?;
319        Ok(project_rows(rows, options.unwrap_or_default()))
320    }
321
322    /// Execute a parameterised query inside the transaction.
323    pub fn execute_with_params(
324        &mut self,
325        query: &str,
326        options: Option<ExecuteOptions>,
327        params: BTreeMap<String, LoraValue>,
328    ) -> Result<QueryResult, LoraError> {
329        let rows = self.execute_rows_with_params_deadline(query, params, None)?;
330        Ok(project_rows(rows, options.unwrap_or_default()))
331    }
332
333    /// Execute a parameterised query inside the transaction with a cooperative
334    /// deadline.
335    pub fn execute_with_params_timeout(
336        &mut self,
337        query: &str,
338        options: Option<ExecuteOptions>,
339        params: BTreeMap<String, LoraValue>,
340        timeout: Duration,
341    ) -> Result<QueryResult, LoraError> {
342        let rows =
343            self.execute_rows_with_params_deadline(query, params, Some(deadline_after(timeout)))?;
344        Ok(project_rows(rows, options.unwrap_or_default()))
345    }
346
347    /// Execute a parameterised query inside the transaction that stops at an
348    /// absolute `deadline` (see [`crate::Database::execute_with_params_deadline`]).
349    /// Several statements can share one deadline to bound the whole
350    /// transaction.
351    pub fn execute_with_params_deadline(
352        &mut self,
353        query: &str,
354        options: Option<ExecuteOptions>,
355        params: BTreeMap<String, LoraValue>,
356        deadline: Instant,
357    ) -> Result<QueryResult, LoraError> {
358        let rows = self.execute_rows_with_params_deadline(query, params, Some(deadline))?;
359        Ok(project_rows(rows, options.unwrap_or_default()))
360    }
361
362    /// Execute a query inside the transaction and return hydrated rows before
363    /// final result-format projection.
364    pub fn execute_rows(&mut self, query: &str) -> Result<Vec<Row>, LoraError> {
365        self.execute_rows_with_params(query, BTreeMap::new())
366    }
367
368    /// Execute a parameterised query inside the transaction and return hydrated
369    /// rows before final result-format projection.
370    pub fn execute_rows_with_params(
371        &mut self,
372        query: &str,
373        params: BTreeMap<String, LoraValue>,
374    ) -> Result<Vec<Row>, LoraError> {
375        Ok(self.execute_rows_with_params_deadline(query, params, None)?)
376    }
377
378    fn execute_rows_with_params_deadline(
379        &mut self,
380        query: &str,
381        params: BTreeMap<String, LoraValue>,
382        deadline: Option<Instant>,
383    ) -> Result<Vec<Row>> {
384        if crate::Database::<InMemoryGraph>::is_schema_command_text(query) {
385            let document = parse_query(query)?;
386            if let lora_ast::Statement::Schema(command) = &document.statement {
387                return self.execute_schema_in_tx(command, &params);
388            }
389        }
390        let compiled = self.compile_in_tx(query)?;
391        self.execute_rows_compiled_deadline(&compiled, params, deadline)
392    }
393
394    fn execute_rows_compiled_deadline(
395        &mut self,
396        compiled: &CompiledQuery,
397        params: BTreeMap<String, LoraValue>,
398        deadline: Option<Instant>,
399    ) -> Result<Vec<Row>> {
400        // ReadOnly tx: never clones, runs straight against live.
401        if self.is_read_only_unchecked() {
402            self.precheck_open_no_savepoint()?;
403            return self.execute_live_compiled(compiled, params, deadline);
404        }
405
406        // ReadWrite tx, lazy-clone aware.
407        let mut inner = self.begin_statement()?;
408        let is_mutating = classify_stream(compiled).is_mutating();
409
410        if !is_mutating {
411            // Read-only statement in a ReadWrite tx. Run against
412            // staged if it has been materialized (so the read
413            // sees prior in-tx writes), otherwise straight off
414            // the live graph — which equals staged-as-it-would-be
415            // because no writes have happened yet.
416            return self.execute_read_statement(inner, compiled, params, deadline);
417        }
418
419        // Mutating statement: lazy-clone the live graph if this
420        // is the first write in the tx, then capture a savepoint
421        // and run the mutable executor.
422        let savepoint = self.prepare_mutating_statement(&mut inner)?;
423
424        let exec_result: ExecResultRows = {
425            let staged = inner.staged_mut()?;
426            execute_mutable_compiled(staged, compiled, params, deadline)
427        };
428
429        match exec_result {
430            Ok(rows) => Ok(rows),
431            Err(err) => {
432                restore_savepoint(&mut inner, savepoint);
433                Err(err)
434            }
435        }
436    }
437
438    /// Open a streaming write cursor over the staged graph for a
439    /// pre-compiled mutating plan, used by the hidden auto-commit
440    /// stream path in `Database::stream_with_params`.
441    ///
442    /// The returned `Box<dyn RowSource + 'static>` may be either a
443    /// real per-row [`StreamingWriteCursor`][lora_executor::StreamingWriteCursor],
444    /// a mutable UNION cursor, or a [`BufferedRowSource`][lora_executor::BufferedRowSource]
445    /// for the remaining materialized leaves. Either way:
446    ///
447    /// * The cursor mutates the *staged* graph, never the live store.
448    /// * Mutations fire the [`BufferingRecorder`] installed on staged
449    ///   by [`Self::ensure_staged_locked`], which accumulates into
450    ///   `inner.buffer` and is replayed into the WAL on commit.
451    /// * `cursor_active` is set to `true` here. The caller MUST clear
452    ///   it before invoking [`Self::commit`] or [`Self::rollback`] —
453    ///   the cursor itself does not.
454    ///
455    /// # Safety
456    ///
457    /// The cursor is `'static` because it owns its compiled query (via
458    /// the supplied `Arc`) and aliases the staged graph through a raw
459    /// pointer. Soundness depends on the invariant that
460    /// `inner.staged` remains `Some(_)` at a stable address for the
461    /// cursor's lifetime. That invariant holds while
462    /// `cursor_active = true` blocks every other path that could
463    /// move or drop staged: explicit statements (`begin_statement`
464    /// rejects), `commit` and `rollback` (rejected until the caller
465    /// clears `cursor_active`).
466    pub(crate) fn open_streaming_compiled_autocommit(
467        &mut self,
468        compiled: Arc<CompiledQuery>,
469        params: BTreeMap<String, LoraValue>,
470    ) -> Result<Box<dyn RowSource + 'static>> {
471        if self.is_read_only_unchecked() {
472            return Err(TransactionError::StreamingRequiresReadWrite.into());
473        }
474
475        let mut inner = self.begin_statement()?;
476        self.ensure_staged_locked(&mut inner)?;
477        inner.activate_cursor();
478
479        // SAFETY: `inner.staged` is `Some` after `ensure_staged_locked`,
480        // and stays at the same address while `cursor_active = true`
481        // (see method-level safety note).
482        let staged_ptr: *mut InMemoryGraph = inner
483            .staged
484            .as_mut()
485            .ok_or(TransactionError::NoStagedGraph)?
486            as *mut _;
487        drop(inner);
488
489        // SAFETY: `compiled` (Arc held by the caller / AutoCommit guard)
490        // keeps the plan alive; `staged_ptr` is valid for the cursor's
491        // lifetime per the invariant above. We extend both lifetimes
492        // to `'static` so the resulting cursor can sit inside the
493        // `'static`-shaped AutoCommit variant of `QueryStream`.
494        let storage_static: &'static mut InMemoryGraph = unsafe { &mut *staged_ptr };
495        let compiled_static: &'static CompiledQuery =
496            unsafe { std::mem::transmute::<&CompiledQuery, _>(compiled.as_ref()) };
497
498        // `MutablePullExecutor::open_compiled` picks the narrowest
499        // cursor shape it can: per-row write cursor, branch-wise
500        // mutable UNION cursor, or a buffered materialized leaf.
501        let cursor = MutablePullExecutor::new(storage_static, params)
502            .open_compiled(compiled_static)
503            .map_err(|e| {
504                // Roll back: the cursor build never happened, so the
505                // tx is in a clean-but-poisoned state. Discard
506                // everything and let the caller bubble the error.
507                if let Ok(mut inner) = self.inner.lock() {
508                    discard_transaction_state(&mut inner);
509                }
510                self.live.take();
511                anyhow::Error::from(e)
512            })?;
513
514        // The Arc<CompiledQuery> is the safety anchor for the
515        // `'static` plan reference. Keep it alive for the cursor's
516        // lifetime by leaking a clone into the cursor's owned data.
517        // We can't store it on the cursor itself (it's a Box<dyn>),
518        // so we wrap the cursor in a guard that owns the Arc.
519        Ok(Box::new(StreamingCursorWithArc {
520            cursor,
521            _compiled: compiled,
522        }))
523    }
524
525    /// Compile a query in this transaction's view of the world:
526    /// against `staged` if it has been materialized, otherwise
527    /// straight against `live`. The two are equivalent before the
528    /// first mutating statement, so the resulting plan is valid
529    /// either way.
530    fn compile_in_tx(&self, query: &str) -> Result<CompiledQuery> {
531        let document = parse_query(query)?;
532        let (resolved, stats) = {
533            let inner = self.lock_inner()?;
534            if let Some(staged) = &inner.staged {
535                let mut analyzer = Analyzer::new(staged);
536                let resolved = analyzer.analyze(&document)?;
537                let stats = staged.graph_stats();
538                (resolved, stats)
539            } else {
540                drop(inner);
541                let live = self.live.as_ref().ok_or(TransactionError::NoGraphGuard)?;
542                let graph = live.as_graph();
543                let mut analyzer = Analyzer::new(graph);
544                let resolved = analyzer.analyze(&document)?;
545                let stats = graph.graph_stats();
546                (resolved, stats)
547            }
548        };
549        Ok(Compiler::compile(&resolved, &stats))
550    }
551
552    /// Materialize `inner.staged` if it doesn't exist yet —
553    /// ReadWrite transactions defer this clone until the first
554    /// mutating statement.
555    fn ensure_staged_locked(&self, inner: &mut MutexGuard<'_, TxInner>) -> Result<()> {
556        if inner.staged.is_some() {
557            return Ok(());
558        }
559        let live = self.live.as_ref().ok_or(TransactionError::NoGraphGuard)?;
560        let mut staged: InMemoryGraph = live.as_graph().clone();
561        if matches!(inner.mode, TransactionMode::ReadWrite) && inner.buffer_mutations {
562            staged.set_mutation_recorder(Some(Arc::new(BufferingRecorder::new(
563                inner.buffer.clone(),
564                inner.buffer_failed.clone(),
565            )) as Arc<dyn MutationRecorder>));
566        }
567        inner.staged = Some(staged);
568        Ok(())
569    }
570
571    fn execute_live_compiled(
572        &self,
573        compiled: &CompiledQuery,
574        params: BTreeMap<String, LoraValue>,
575        deadline: Option<Instant>,
576    ) -> Result<Vec<Row>> {
577        let live = self.live.as_ref().ok_or(TransactionError::NoGraphGuard)?;
578        execute_read_compiled(live.as_graph(), compiled, params, deadline)
579    }
580
581    fn execute_read_statement(
582        &self,
583        inner: MutexGuard<'_, TxInner>,
584        compiled: &CompiledQuery,
585        params: BTreeMap<String, LoraValue>,
586        deadline: Option<Instant>,
587    ) -> Result<Vec<Row>> {
588        match inner.staged.as_ref() {
589            Some(staged) => execute_read_compiled(staged, compiled, params, deadline),
590            None => {
591                drop(inner);
592                self.execute_live_compiled(compiled, params, deadline)
593            }
594        }
595    }
596
597    /// Schema DDL inside a transaction. `SHOW` reads the transaction's
598    /// view (including DDL applied earlier in it). `CREATE` / `DROP`
599    /// apply to the staged graph like a data write: they emit catalog
600    /// events into the transaction's WAL buffer, become visible to later
601    /// statements, and are published by `commit` or discarded by
602    /// `rollback`, all-or-nothing with the data statements.
603    fn execute_schema_in_tx(
604        &mut self,
605        command: &lora_ast::SchemaCommand,
606        params: &BTreeMap<String, LoraValue>,
607    ) -> Result<Vec<Row>> {
608        use crate::database::schema::{apply_schema_mutation, schema_command_is_read, show_schema};
609
610        if schema_command_is_read(command) {
611            if self.is_read_only_unchecked() {
612                self.precheck_open_no_savepoint()?;
613                let live = self.live.as_ref().ok_or(TransactionError::NoGraphGuard)?;
614                return show_schema(live.as_graph(), command, params);
615            }
616            let inner = self.begin_statement()?;
617            if let Some(staged) = inner.staged.as_ref() {
618                return show_schema(staged, command, params);
619            }
620            drop(inner);
621            let live = self.live.as_ref().ok_or(TransactionError::NoGraphGuard)?;
622            return show_schema(live.as_graph(), command, params);
623        }
624
625        if self.is_read_only_unchecked() {
626            return Err(TransactionError::ReadOnlyMutation.into());
627        }
628        let mut inner = self.begin_statement()?;
629        let savepoint = self.prepare_mutating_statement(&mut inner)?;
630        let result = {
631            let staged = inner.staged_mut()?;
632            apply_schema_mutation(staged, command, params)
633        };
634        if result.is_err() {
635            restore_savepoint(&mut inner, savepoint);
636        }
637        result
638    }
639
640    fn prepare_mutating_statement(
641        &self,
642        inner: &mut MutexGuard<'_, TxInner>,
643    ) -> Result<Option<Savepoint>> {
644        let clone_savepoint_graph = inner.staged.is_some();
645        self.ensure_staged_locked(inner)?;
646        Ok(Some(take_savepoint(inner, clone_savepoint_graph)))
647    }
648
649    /// Compile `query` against the transaction's view of the graph and
650    /// return the plan that *would* run. Never executes the query, so
651    /// running `explain` on a mutating statement leaves the transaction's
652    /// staged graph untouched.
653    pub fn explain(
654        &self,
655        query: &str,
656        _params: Option<BTreeMap<String, LoraValue>>,
657    ) -> Result<QueryPlan, LoraError> {
658        let compiled = self.compile_in_tx(query).map_err(LoraError::from_anyhow)?;
659        Ok(query_plan_for(query, &compiled))
660    }
661
662    /// Execute `query` inside the transaction and return the plan plus
663    /// runtime metrics.
664    ///
665    /// **PROFILE executes the query for real.** Mutating statements
666    /// affect the transaction's staged graph and are persisted at the
667    /// usual `commit` point. Use [`Transaction::explain`] to inspect a
668    /// plan without running it.
669    pub fn profile(
670        &mut self,
671        query: &str,
672        params: Option<BTreeMap<String, LoraValue>>,
673    ) -> Result<QueryProfile, LoraError> {
674        let params = params.unwrap_or_default();
675        let compiled = self.compile_in_tx(query).map_err(LoraError::from_anyhow)?;
676        let plan = query_plan_for(query, &compiled);
677        let shape = plan.shape;
678
679        let collector = Arc::new(MetricsCollector::new());
680        let _guard = CollectorGuard::install(collector.clone());
681
682        let started = Instant::now();
683        let rows = self
684            .execute_rows_compiled_deadline(&compiled, params, None)
685            .map_err(LoraError::from_anyhow)?;
686        let total_elapsed_ns = started.elapsed().as_nanos() as u64;
687
688        drop(_guard);
689        let per_operator = collector
690            .snapshot()
691            .into_iter()
692            .map(|(id, op)| {
693                (
694                    id,
695                    OperatorMetrics {
696                        rows: op.rows,
697                        elapsed_ns: op.elapsed_ns,
698                        next_calls: op.next_calls,
699                        db_hits: 0,
700                    },
701                )
702            })
703            .collect();
704
705        let metrics = ProfileMetrics {
706            total_elapsed_ns,
707            total_rows: rows.len() as u64,
708            mutated: shape.is_mutating(),
709            per_operator,
710        };
711
712        Ok(QueryProfile { plan, metrics })
713    }
714
715    /// Execute a query inside the transaction and return an owning row stream.
716    pub fn stream(&mut self, query: &str) -> Result<QueryStream<'static>, LoraError> {
717        self.stream_with_params(query, BTreeMap::new())
718    }
719
720    /// Execute a parameterised query inside the transaction and return an
721    /// owning row stream.
722    pub fn stream_with_params(
723        &mut self,
724        query: &str,
725        params: BTreeMap<String, LoraValue>,
726    ) -> Result<QueryStream<'static>, LoraError> {
727        let compiled = Arc::new(self.compile_in_tx(query)?);
728        let columns = compiled_result_columns(&compiled);
729        Ok(self.stream_compiled(compiled, columns, params)?)
730    }
731
732    /// Open a tx-bound stream for an already-compiled plan. Lets
733    /// `Database::stream_with_params` reuse the plan it built for
734    /// classification.
735    pub(crate) fn stream_compiled(
736        &mut self,
737        compiled: Arc<CompiledQuery>,
738        columns: Vec<String>,
739        params: BTreeMap<String, LoraValue>,
740    ) -> Result<QueryStream<'static>> {
741        let mut inner = self.begin_statement()?;
742        let is_mutating = classify_stream(&compiled).is_mutating();
743        ensure_stream_allowed(&inner, is_mutating)?;
744
745        let rollback_on_drop = stream_rolls_back_on_drop(is_mutating);
746        let staged_ptr = self.prepare_stream_staging(&mut inner, rollback_on_drop)?;
747        drop(inner);
748
749        let cursor = open_tx_stream_cursor(staged_ptr, compiled, params, is_mutating);
750
751        match cursor {
752            Ok(cursor) => Ok(QueryStream::for_tx_cursor(
753                cursor,
754                columns,
755                TxCursorLease::new(self.inner.clone(), rollback_on_drop),
756            )),
757            Err(err) => {
758                finalize_tx_stream(&self.inner, TxStreamOutcome::Interrupted, rollback_on_drop);
759                Err(err)
760            }
761        }
762    }
763
764    fn prepare_stream_staging(
765        &self,
766        inner: &mut MutexGuard<'_, TxInner>,
767        rollback_on_drop: bool,
768    ) -> Result<*mut InMemoryGraph> {
769        // Transaction streams borrow from the staged graph. Even read-only
770        // streams materialize staging when needed so the cursor can outlive
771        // the `&mut Transaction` borrow without borrowing from the
772        // transaction-owned live write guard.
773        let clone_savepoint_graph = inner.staged.is_some();
774        self.ensure_staged_locked(inner)?;
775        inner.activate_cursor();
776
777        if rollback_on_drop {
778            inner.pending_savepoint = Some(take_savepoint(inner, clone_savepoint_graph));
779        } else {
780            inner.pending_savepoint = None;
781        }
782
783        Ok(inner
784            .staged
785            .as_mut()
786            .ok_or(TransactionError::NoStagedGraph)? as *mut _)
787    }
788
789    /// Commit the transaction and publish staged changes.
790    ///
791    /// When WAL is attached the buffered tx-local mutation log is
792    /// replayed into the durable WAL as a single committed
793    /// transaction; recovery therefore observes either every write
794    /// in this transaction or none.
795    pub fn commit(mut self) -> Result<(), LoraError> {
796        let CommitState {
797            staged,
798            buffer_events,
799            mode,
800        } = self.take_commit_state()?;
801
802        let capture = self.changes.is_active()
803            && matches!(mode, TransactionMode::ReadWrite)
804            && staged.is_some()
805            && !buffer_events.is_empty();
806        let (wrote_wal_commit, captured) = if capture && self.wal.is_some() {
807            let events = buffer_events.clone();
808            let lsn = self.replay_commit_wal_lsn(mode, buffer_events)?;
809            (lsn.is_some(), lsn.map(|lsn| (Some(lsn), events)))
810        } else if capture {
811            (false, Some((None, buffer_events)))
812        } else {
813            (self.replay_commit_wal(mode, buffer_events)?, None)
814        };
815        self.publish_staged_graph(mode, staged, wrote_wal_commit)?;
816
817        if let Some((lsn, events)) = captured {
818            // Still holding the writer lease, so the batch is ordered
819            // against the next commit.
820            if let Some(LiveStoreGuard::Write(lease)) = &self.live {
821                let post = lease.store.load_full();
822                crate::changes::publish_committed(
823                    &self.changes,
824                    lsn.map(|lsn| lsn.raw()),
825                    &events,
826                    &crate::changes::PreImages::for_events(&events, Some(&lease.snapshot)),
827                    &post,
828                );
829            }
830        }
831
832        self.live.take();
833        Ok(())
834    }
835
836    fn take_commit_state(&self) -> Result<CommitState> {
837        let mut inner = self.lock_inner()?;
838        if inner.cursor_active {
839            return Err(TransactionError::CursorActiveCommit.into());
840        }
841        if inner.closed {
842            return Err(TransactionError::AlreadyClosed.into());
843        }
844
845        let mode = inner.mode;
846        if inner.buffer_failed.load(Ordering::Acquire) {
847            inner.closed = true;
848            return Err(TransactionError::Poisoned.into());
849        }
850        let buffer_events = {
851            let mut buffer = inner
852                .buffer
853                .lock()
854                .map_err(|_| TransactionError::Poisoned)?;
855            std::mem::take(&mut *buffer)
856        };
857        // Both modes can have `staged = None`: ReadOnly never clones,
858        // and ReadWrite transactions that performed no writes leave
859        // staging unmaterialized too.
860        let staged = inner.staged.take();
861        inner.closed = true;
862
863        Ok(CommitState {
864            staged,
865            buffer_events,
866            mode,
867        })
868    }
869
870    fn replay_commit_wal(
871        &self,
872        mode: TransactionMode,
873        buffer_events: Vec<MutationEvent>,
874    ) -> Result<bool> {
875        let Some(rec) = &self.wal else {
876            return Ok(false);
877        };
878
879        if !matches!(mode, TransactionMode::ReadWrite) {
880            ensure_wal_not_poisoned(rec)?;
881            return Ok(false);
882        }
883
884        Ok(rec.commit_events(buffer_events)?.wrote())
885    }
886
887    /// [`Self::replay_commit_wal`] for a capturing change feed: returns the
888    /// commit LSN instead of a flag.
889    fn replay_commit_wal_lsn(
890        &self,
891        mode: TransactionMode,
892        buffer_events: Vec<MutationEvent>,
893    ) -> Result<Option<lora_wal::Lsn>> {
894        let Some(rec) = &self.wal else {
895            return Ok(None);
896        };
897        if !matches!(mode, TransactionMode::ReadWrite) {
898            ensure_wal_not_poisoned(rec)?;
899            return Ok(None);
900        }
901        Ok(rec.commit_events_lsn(buffer_events)?)
902    }
903
904    fn publish_staged_graph(
905        &mut self,
906        mode: TransactionMode,
907        staged: Option<InMemoryGraph>,
908        wrote_wal_commit: bool,
909    ) -> Result<()> {
910        if !matches!(mode, TransactionMode::ReadWrite) {
911            return Ok(());
912        }
913
914        let Some(mut staged) = staged else {
915            return Ok(());
916        };
917
918        // Strip the buffering recorder from the staged graph before
919        // publishing it as the live store; the live store either has
920        // the durable WAL recorder reinstalled below or no recorder at
921        // all (for non-WAL databases).
922        staged.set_mutation_recorder(None);
923        let wal = self.wal.clone();
924        if let Some(rec) = &wal {
925            staged.set_mutation_recorder(Some(rec.clone() as Arc<dyn MutationRecorder>));
926        }
927
928        let live = self.live.as_mut().ok_or(TransactionError::NoGraphGuard)?;
929        let lease = match live {
930            LiveStoreGuard::Write(lease) => lease,
931            LiveStoreGuard::Read(_) => {
932                return Err(TransactionError::ReadOnlyCommit.into());
933            }
934        };
935
936        if wrote_wal_commit {
937            if let (Some(snapshots), Some(rec)) = (&self.snapshots, wal.as_ref()) {
938                snapshots.observe_commit(&staged, rec)?;
939            }
940        }
941
942        // Atomic publish — concurrent readers will see the new state on
943        // their next `load_full()`, while in-flight readers keep their
944        // existing `Arc<InMemoryGraph>` snapshot until they drop it.
945        lease.store.store(Arc::new(staged));
946
947        Ok(())
948    }
949
950    /// Roll back the transaction. Staged graph changes and buffered
951    /// mutations are discarded; the WAL is never armed.
952    pub fn rollback(mut self) -> Result<(), LoraError> {
953        let mut inner = self.lock_inner()?;
954        if inner.closed {
955            return Err(TransactionError::AlreadyClosed.into());
956        }
957        discard_transaction_state(&mut inner);
958        drop(inner);
959        self.live.take();
960        Ok(())
961    }
962
963    /// Acquire the inner state for a new statement. Validates that
964    /// the transaction is still open and no cursor is active. The
965    /// staged graph is *not* required: ReadWrite transactions
966    /// defer the staging clone until the first mutating statement
967    /// (see [`Transaction::ensure_staged_locked`]).
968    fn begin_statement(&self) -> Result<MutexGuard<'_, TxInner>> {
969        let inner = self.lock_inner()?;
970        if inner.closed {
971            return Err(TransactionError::AlreadyClosed.into());
972        }
973        if inner.cursor_active {
974            return Err(TransactionError::CursorActiveStatement.into());
975        }
976        Ok(inner)
977    }
978
979    /// Cheap state check for the ReadOnly fast path: closed +
980    /// cursor_active. No staged-graph check — ReadOnly tx has no
981    /// staged graph by construction.
982    fn precheck_open_no_savepoint(&self) -> Result<()> {
983        let inner = self.lock_inner()?;
984        if inner.closed {
985            return Err(TransactionError::AlreadyClosed.into());
986        }
987        if inner.cursor_active {
988            return Err(TransactionError::CursorActiveStatement.into());
989        }
990        Ok(())
991    }
992
993    /// True if the transaction was begun in `ReadOnly` mode. Cheap
994    /// — `mode` doesn't change after `begin_transaction`, so we
995    /// pay one small state-lock acquisition.
996    fn is_read_only_unchecked(&self) -> bool {
997        matches!(self.mode, TransactionMode::ReadOnly)
998    }
999
1000    fn lock_inner(&self) -> Result<MutexGuard<'_, TxInner>> {
1001        self.inner
1002            .lock()
1003            .map_err(|_| TransactionError::Poisoned.into())
1004    }
1005
1006    pub(crate) fn release_streaming_cursor(&self) {
1007        if let Ok(mut inner) = self.inner.lock() {
1008            inner.release_cursor();
1009        }
1010    }
1011}
1012
1013type ExecResultRows = Result<Vec<Row>>;
1014
1015fn deadline_after(timeout: Duration) -> Instant {
1016    Instant::now()
1017        .checked_add(timeout)
1018        .unwrap_or_else(Instant::now)
1019}
1020
1021fn execute_read_compiled(
1022    storage: &InMemoryGraph,
1023    compiled: &CompiledQuery,
1024    params: BTreeMap<String, LoraValue>,
1025    deadline: Option<Instant>,
1026) -> Result<Vec<Row>> {
1027    // Same fast path as auto-commit reads: a plan with an early LIMIT runs
1028    // on the pull pipeline so it stops scanning once the LIMIT is met.
1029    if crate::database::pull_mode::should_collect_read_via_pull(compiled) {
1030        return lora_executor::collect_compiled_with_deadline(storage, params, compiled, deadline)
1031            .map_err(anyhow::Error::from);
1032    }
1033    let executor = Executor::with_deadline(ExecutionContext { storage, params }, deadline);
1034    executor
1035        .execute_compiled_rows(compiled)
1036        .map_err(anyhow::Error::from)
1037}
1038
1039fn execute_mutable_compiled(
1040    storage: &mut InMemoryGraph,
1041    compiled: &CompiledQuery,
1042    params: BTreeMap<String, LoraValue>,
1043    deadline: Option<Instant>,
1044) -> Result<Vec<Row>> {
1045    let mut executor =
1046        MutableExecutor::with_deadline(MutableExecutionContext { storage, params }, deadline);
1047    executor
1048        .execute_compiled_rows(compiled)
1049        .map_err(anyhow::Error::from)
1050}
1051
1052fn query_plan_for(query: &str, compiled: &CompiledQuery) -> QueryPlan {
1053    QueryPlan {
1054        query: query.to_string(),
1055        tree: plan_tree_from_compiled(compiled),
1056        shape: classify_stream(compiled).into(),
1057        result_columns: plan_result_columns(&compiled.physical),
1058    }
1059}
1060
1061fn ensure_stream_allowed(inner: &TxInner, is_mutating: bool) -> Result<()> {
1062    if matches!(inner.mode, TransactionMode::ReadOnly) && is_mutating {
1063        Err(TransactionError::ReadOnlyMutation.into())
1064    } else {
1065        Ok(())
1066    }
1067}
1068
1069fn stream_rolls_back_on_drop(is_mutating: bool) -> bool {
1070    is_mutating
1071}
1072
1073fn open_tx_stream_cursor(
1074    staged_ptr: *mut InMemoryGraph,
1075    compiled: Arc<CompiledQuery>,
1076    params: BTreeMap<String, LoraValue>,
1077    is_mutating: bool,
1078) -> Result<Box<dyn RowSource + 'static>> {
1079    // SAFETY: `stream_compiled` calls this only after `prepare_stream_staging`
1080    // has materialized `inner.staged` and activated the cursor guard. While the
1081    // guard is active, statements/commit/rollback cannot move or drop staged.
1082    let compiled_static: &'static CompiledQuery =
1083        unsafe { std::mem::transmute::<&CompiledQuery, _>(compiled.as_ref()) };
1084
1085    if is_mutating {
1086        // SAFETY: see function-level note above. Mutating streams need exclusive
1087        // access, and the transaction cursor token prevents any second borrow.
1088        let storage_static: &'static mut InMemoryGraph = unsafe { &mut *staged_ptr };
1089        MutablePullExecutor::new(storage_static, params)
1090            .open_compiled(compiled_static)
1091            .map(|cursor| boxed_streaming_cursor(cursor, compiled))
1092            .map_err(anyhow::Error::from)
1093    } else {
1094        // SAFETY: see function-level note above. Read streams borrow staged
1095        // immutably until the cursor lease is finalized.
1096        let storage_static: &'static InMemoryGraph = unsafe { &*staged_ptr };
1097        PullExecutor::new(storage_static, params)
1098            .open_compiled(compiled_static)
1099            .map(|cursor| boxed_streaming_cursor(cursor, compiled))
1100            .map_err(anyhow::Error::from)
1101    }
1102}
1103
1104fn boxed_streaming_cursor(
1105    cursor: Box<dyn RowSource + 'static>,
1106    compiled: Arc<CompiledQuery>,
1107) -> Box<dyn RowSource + 'static> {
1108    Box::new(StreamingCursorWithArc {
1109        cursor,
1110        _compiled: compiled,
1111    })
1112}
1113
1114struct CommitState {
1115    staged: Option<InMemoryGraph>,
1116    buffer_events: Vec<MutationEvent>,
1117    mode: TransactionMode,
1118}
1119
1120impl TxInner {
1121    fn staged_mut(&mut self) -> Result<&mut InMemoryGraph> {
1122        self.staged
1123            .as_mut()
1124            .ok_or(TransactionError::NoStagedGraph.into())
1125    }
1126
1127    fn activate_cursor(&mut self) {
1128        self.cursor_active = true;
1129    }
1130
1131    fn release_cursor(&mut self) {
1132        self.cursor_active = false;
1133    }
1134
1135    fn clear_pending_savepoint(&mut self) {
1136        self.pending_savepoint = None;
1137    }
1138
1139    fn restore_pending_savepoint(&mut self) {
1140        if let Some(sp) = self.pending_savepoint.take() {
1141            apply_savepoint(self, sp);
1142        }
1143    }
1144
1145    fn finalize_stream(&mut self, outcome: TxStreamOutcome, rollback_on_drop: bool) {
1146        self.release_cursor();
1147
1148        if self.closed {
1149            discard_transaction_state(self);
1150            return;
1151        }
1152
1153        if outcome.should_restore_savepoint(rollback_on_drop) {
1154            self.restore_pending_savepoint();
1155        } else {
1156            self.clear_pending_savepoint();
1157        }
1158    }
1159}
1160
1161/// `RowSource` adapter that owns an `Arc<CompiledQuery>` so the
1162/// inner cursor's `'static` borrows into the plan stay valid for the
1163/// life of the wrapper. The inner cursor is stored first so it drops
1164/// before the Arc, releasing any borrows back into the plan.
1165struct StreamingCursorWithArc {
1166    cursor: Box<dyn RowSource + 'static>,
1167    _compiled: Arc<CompiledQuery>,
1168}
1169
1170impl RowSource for StreamingCursorWithArc {
1171    fn next_row(&mut self) -> lora_executor::ExecResult<Option<Row>> {
1172        self.cursor.next_row()
1173    }
1174}
1175
1176fn finalize_tx_stream(
1177    handle: &Arc<Mutex<TxInner>>,
1178    outcome: TxStreamOutcome,
1179    rollback_on_drop: bool,
1180) {
1181    if let Ok(mut inner) = handle.lock() {
1182        inner.finalize_stream(outcome, rollback_on_drop);
1183    }
1184}
1185
1186fn discard_transaction_state(inner: &mut TxInner) {
1187    // A full transaction rollback supersedes any pending cursor savepoint.
1188    inner.clear_pending_savepoint();
1189    inner.release_cursor();
1190    inner.staged = None;
1191    if let Ok(mut buf) = inner.buffer.lock() {
1192        buf.clear();
1193    } else {
1194        inner.buffer_failed.store(true, Ordering::Release);
1195    }
1196    inner.closed = true;
1197}
1198
1199fn take_savepoint(inner: &TxInner, clone_staged: bool) -> Savepoint {
1200    let buffer_len = match inner.buffer.lock() {
1201        Ok(buffer) => buffer.len(),
1202        Err(_) => {
1203            inner.buffer_failed.store(true, Ordering::Release);
1204            0
1205        }
1206    };
1207    Savepoint {
1208        staged: if clone_staged {
1209            inner.staged.as_ref().cloned()
1210        } else {
1211            None
1212        },
1213        buffer_len,
1214    }
1215}
1216
1217fn restore_savepoint(inner: &mut TxInner, savepoint: Option<Savepoint>) {
1218    if let Some(sp) = savepoint {
1219        apply_savepoint(inner, sp);
1220    }
1221}
1222
1223fn apply_savepoint(inner: &mut TxInner, sp: Savepoint) {
1224    if let Ok(mut buf) = inner.buffer.lock() {
1225        buf.truncate(sp.buffer_len);
1226    } else {
1227        inner.buffer_failed.store(true, Ordering::Release);
1228    }
1229
1230    let Some(mut graph) = sp.staged else {
1231        inner.staged = None;
1232        return;
1233    };
1234
1235    // Rebuild the staged graph from the snapshot and re-install the
1236    // buffering recorder. `InMemoryGraph::clone` deliberately drops
1237    // recorders, so the snapshot has none until we put it back.
1238    if matches!(inner.mode, TransactionMode::ReadWrite) && inner.buffer_mutations {
1239        graph.set_mutation_recorder(Some(Arc::new(BufferingRecorder::new(
1240            inner.buffer.clone(),
1241            inner.buffer_failed.clone(),
1242        )) as Arc<dyn MutationRecorder>));
1243    }
1244    inner.staged = Some(graph);
1245}
1246
1247impl Drop for Transaction<'_> {
1248    fn drop(&mut self) {
1249        // If the user never called commit/rollback, treat it as a
1250        // rollback: drop staged changes and the buffered mutations.
1251        // The live RwLock guard is released as part of dropping `self.live`.
1252        if let Ok(mut inner) = self.inner.lock() {
1253            if !inner.closed {
1254                if inner.cursor_active {
1255                    // A tx-bound stream may still be borrowing the
1256                    // staged graph through `inner`. Leave that graph
1257                    // in place until the stream drops, but mark the
1258                    // transaction closed so finalization discards it
1259                    // instead of making it commit-eligible.
1260                    inner.closed = true;
1261                } else {
1262                    discard_transaction_state(&mut inner);
1263                }
1264            }
1265        }
1266    }
1267}
1268
1269#[cfg(test)]
1270mod tests {
1271    use std::thread;
1272
1273    use super::*;
1274
1275    #[test]
1276    fn buffering_recorder_latches_poisoned_buffer() {
1277        let buffer = Arc::new(Mutex::new(Vec::new()));
1278        let failed = Arc::new(AtomicBool::new(false));
1279
1280        let poisoned_buffer = buffer.clone();
1281        let _ = thread::spawn(move || {
1282            let _guard = poisoned_buffer.lock().unwrap();
1283            panic!("poison mutation buffer");
1284        })
1285        .join();
1286
1287        let recorder = BufferingRecorder::new(buffer, failed.clone());
1288        recorder.record(MutationEvent::Clear);
1289
1290        assert!(failed.load(Ordering::Acquire));
1291    }
1292}