Skip to main content

lora_database/
transaction.rs

1use std::collections::BTreeMap;
2use std::sync::{Arc, Mutex, MutexGuard};
3use std::time::{Duration, Instant};
4
5use arc_swap::ArcSwap;
6
7use anyhow::Result;
8use thiserror::Error;
9
10use lora_analyzer::Analyzer;
11use lora_compiler::{CompiledQuery, Compiler};
12use lora_executor::{
13    classify_stream, compiled_result_columns, project_rows, ExecuteOptions, ExecutionContext,
14    Executor, LoraValue, MutableExecutionContext, MutableExecutor, MutablePullExecutor,
15    PullExecutor, QueryResult, Row, RowSource,
16};
17use lora_parser::parse_query;
18use lora_store::{InMemoryGraph, MutationEvent, MutationRecorder};
19use lora_wal::WalRecorder;
20
21use crate::error::LoraError;
22use crate::snapshot::ManagedSnapshotStore;
23use crate::stream::QueryStream;
24use crate::wal::write_scope::ensure_wal_not_poisoned;
25
26/// Transaction-lifecycle invariant violations.
27///
28/// All variants used to be raised as `anyhow!("...")` strings. Surfacing
29/// them as a typed enum lets [`crate::LoraError`] route them onto stable
30/// `LoraErrorCode`s without phrase-matching the `Display` text.
31#[derive(Debug, Clone, PartialEq, Eq, Error)]
32pub enum TransactionError {
33    #[error("transaction is already closed")]
34    AlreadyClosed,
35
36    #[error("transaction has no live graph guard")]
37    NoGraphGuard,
38
39    #[error("transaction has no staged graph")]
40    NoStagedGraph,
41
42    #[error("cannot commit transaction while a streaming cursor is still active")]
43    CursorActiveCommit,
44
45    #[error("cannot start a new statement while a streaming cursor is still active")]
46    CursorActiveStatement,
47
48    #[error("cannot execute mutating query in read-only transaction")]
49    ReadOnlyMutation,
50
51    #[error("streaming write cursor requires a ReadWrite transaction")]
52    StreamingRequiresReadWrite,
53
54    #[error("read-only transaction cannot publish staged graph")]
55    ReadOnlyCommit,
56}
57
58/// Transaction execution mode.
59#[derive(Debug, Clone, Copy, PartialEq, Eq)]
60pub enum TransactionMode {
61    /// Use the read-only executor. Write operators return read-only errors.
62    ReadOnly,
63    /// Execute reads and writes against a staged graph, then publish on commit.
64    ReadWrite,
65}
66
67/// What a transaction holds onto for the duration of its statements.
68///
69/// `Read` simply pins an `Arc<InMemoryGraph>` snapshot — readers don't
70/// take any lock, so nothing observable changes when a writer commits
71/// a new version mid-transaction. `Write` holds the writer Mutex
72/// (serializing commit ordering) and a snapshot of the live graph at
73/// the point the transaction began; the working copy is built lazily
74/// in `TxInner::staged` on first mutation, mirroring the previous
75/// "clone on first mutation" behavior.
76pub(crate) enum LiveStoreGuard<'db> {
77    Read(Arc<InMemoryGraph>),
78    Write(WriteLease<'db>),
79}
80
81/// Writer lease for a `ReadWrite` transaction. Holds the per-database
82/// writer Mutex plus a read snapshot of the graph at lease open time.
83/// The mutating working copy lives in `TxInner::staged` and is
84/// cloned from `snapshot` lazily.
85pub(crate) struct WriteLease<'db> {
86    /// Held for the tx lifetime so concurrent ReadWrite txns serialize.
87    pub(crate) _writer_lock: MutexGuard<'db, ()>,
88    /// Pointer back to the live `ArcSwap` so commit can publish.
89    pub(crate) store: Arc<ArcSwap<InMemoryGraph>>,
90    /// Read-only view of the graph at lease open time. The first
91    /// mutating statement clones from this into `TxInner::staged`.
92    pub(crate) snapshot: Arc<InMemoryGraph>,
93}
94
95impl LiveStoreGuard<'_> {
96    fn as_graph(&self) -> &InMemoryGraph {
97        match self {
98            Self::Read(arc) => arc,
99            Self::Write(lease) => &lease.snapshot,
100        }
101    }
102}
103
104/// Captures the staged graph and tx-local mutation buffer at the point
105/// a statement is opened, so a failed/dropped statement can be rolled
106/// back to that point without affecting earlier work in the same
107/// transaction.
108pub(crate) struct Savepoint {
109    staged: Option<InMemoryGraph>,
110    buffer_len: usize,
111}
112
113/// How a transaction-bound stream finished. Exhaustion commits that
114/// statement's staged changes into the transaction; interruption means
115/// drop or runtime error before all rows were observed.
116#[derive(Debug, Clone, Copy, PartialEq, Eq)]
117pub(crate) enum TxStreamOutcome {
118    Exhausted,
119    Interrupted,
120}
121
122impl TxStreamOutcome {
123    fn should_restore_savepoint(self, rollback_on_drop: bool) -> bool {
124        matches!(self, Self::Interrupted) && rollback_on_drop
125    }
126}
127
128/// Owns the active-cursor token for a transaction-bound stream.
129///
130/// The lease is created only after a cursor has opened successfully. If stream
131/// code forgets to finalize it explicitly, `Drop` treats the cursor as
132/// interrupted so a mutating statement cannot accidentally commit partial work
133/// into the transaction.
134pub(crate) struct TxCursorLease {
135    handle: Arc<Mutex<TxInner>>,
136    rollback_on_drop: bool,
137    finalized: bool,
138}
139
140impl TxCursorLease {
141    pub(crate) fn new(handle: Arc<Mutex<TxInner>>, rollback_on_drop: bool) -> Self {
142        Self {
143            handle,
144            rollback_on_drop,
145            finalized: false,
146        }
147    }
148
149    pub(crate) fn finalize(&mut self, outcome: TxStreamOutcome) {
150        if self.finalized {
151            return;
152        }
153        finalize_tx_stream(&self.handle, outcome, self.rollback_on_drop);
154        self.finalized = true;
155    }
156}
157
158impl Drop for TxCursorLease {
159    fn drop(&mut self) {
160        self.finalize(TxStreamOutcome::Interrupted);
161    }
162}
163
164/// Buffers `MutationEvent`s emitted by the staged graph while a
165/// transaction is in progress. The buffer replaces direct WAL writes
166/// during the transaction body; on commit the host replays the
167/// buffer into the real `WalRecorder` as a single durable
168/// transaction. Statement rollback truncates the buffer back to its
169/// pre-statement length; transaction rollback drops it entirely.
170///
171/// Also reused by the optimistic auto-commit write path
172/// ([`Database::execute_rows_with_params_deadline`]) so multiple
173/// concurrent writers can buffer mutations off-thread and only
174/// serialize at the brief WAL append + ArcSwap publish step.
175pub(crate) struct BufferingRecorder {
176    buffer: Arc<Mutex<Vec<MutationEvent>>>,
177}
178
179impl BufferingRecorder {
180    pub(crate) fn new(buffer: Arc<Mutex<Vec<MutationEvent>>>) -> Self {
181        Self { buffer }
182    }
183}
184
185impl MutationRecorder for BufferingRecorder {
186    fn record(&self, event: MutationEvent) {
187        if let Ok(mut buf) = self.buffer.lock() {
188            buf.push(event);
189        }
190    }
191}
192
193/// Shared transaction state. Wrapped in `Arc<Mutex<>>` so a
194/// `QueryStream` opened against the transaction can release its
195/// cursor token and signal savepoint-rollback intent on drop without
196/// borrowing the [`Transaction`] handle.
197pub(crate) struct TxInner {
198    /// The cloned staging graph. Mutated by write statements through
199    /// the [`MutableExecutor`]; read by read-only statements through
200    /// [`PullExecutor`]. `None` once the transaction has been closed.
201    pub(crate) staged: Option<InMemoryGraph>,
202    /// Tx-local mutation log, populated by the [`BufferingRecorder`]
203    /// installed on `staged`. Replayed into the real WAL exactly once
204    /// at commit time.
205    pub(crate) buffer: Arc<Mutex<Vec<MutationEvent>>>,
206    /// Per-statement savepoint snapshot. Set when a statement opens,
207    /// cleared on successful completion, restored on
208    /// failure/premature drop.
209    pub(crate) pending_savepoint: Option<Savepoint>,
210    /// True while a `QueryStream` opened against this transaction is
211    /// alive. Blocks new statements and prevents commit until the
212    /// cursor is released.
213    pub(crate) cursor_active: bool,
214    /// True after `commit` or `rollback` has run, regardless of
215    /// outcome. Subsequent operations fail loudly instead of silently
216    /// running on stale state.
217    pub(crate) closed: bool,
218    /// Transaction execution mode chosen at `begin_transaction` time.
219    pub(crate) mode: TransactionMode,
220    /// Whether this transaction needs a mutation buffer for durable WAL
221    /// replay. Databases without a WAL can skip recorder installation and
222    /// avoid cloning mutation payloads into an unused buffer.
223    pub(crate) buffer_mutations: bool,
224}
225
226/// Explicit transaction over the in-memory graph.
227///
228/// The implementation is conservative: read-only transactions hold a
229/// database read lock, and read-write transactions hold the database
230/// write lock. Read-write transactions lazily create a cloned staging
231/// graph on the first mutating statement, then either swap that graph
232/// into place on commit or drop it on rollback. Explicit mutating
233/// statements capture a graph +
234/// WAL-buffer savepoint so a failed or dropped streaming statement
235/// only rolls back its own effects, not the transaction as a whole.
236///
237/// When a WAL is attached, mutation events fire into a tx-local
238/// buffer rather than the durable log. The buffer is replayed into
239/// the WAL exactly once at commit, so recovery never observes
240/// partial / aborted / dropped statements.
241pub struct Transaction<'db> {
242    pub(crate) live: Option<LiveStoreGuard<'db>>,
243    pub(crate) inner: Arc<Mutex<TxInner>>,
244    pub(crate) wal: Option<Arc<WalRecorder>>,
245    pub(crate) snapshots: Option<Arc<ManagedSnapshotStore>>,
246    mode: TransactionMode,
247}
248
249impl<'db> Transaction<'db> {
250    /// Build a fresh transaction. Used by `Database::begin_transaction`.
251    pub(crate) fn new(
252        live: LiveStoreGuard<'db>,
253        wal: Option<Arc<WalRecorder>>,
254        snapshots: Option<Arc<ManagedSnapshotStore>>,
255        mode: TransactionMode,
256    ) -> Self {
257        let buffer_mutations = wal.is_some();
258        let inner = TxInner {
259            staged: None,
260            buffer: Arc::new(Mutex::new(Vec::new())),
261            pending_savepoint: None,
262            cursor_active: false,
263            closed: false,
264            mode,
265            buffer_mutations,
266        };
267        Self {
268            live: Some(live),
269            inner: Arc::new(Mutex::new(inner)),
270            wal,
271            snapshots,
272            mode,
273        }
274    }
275
276    /// Transaction mode chosen at begin time.
277    pub fn mode(&self) -> TransactionMode {
278        self.mode
279    }
280
281    /// Execute a query inside the transaction and return a materialized
282    /// `QueryResult`.
283    pub fn execute(
284        &mut self,
285        query: &str,
286        options: Option<ExecuteOptions>,
287    ) -> Result<QueryResult, LoraError> {
288        self.execute_with_params(query, options, BTreeMap::new())
289    }
290
291    /// Execute a query inside the transaction with a cooperative deadline.
292    pub fn execute_with_timeout(
293        &mut self,
294        query: &str,
295        options: Option<ExecuteOptions>,
296        timeout: Duration,
297    ) -> Result<QueryResult, LoraError> {
298        let deadline = Instant::now()
299            .checked_add(timeout)
300            .unwrap_or_else(Instant::now);
301        let rows =
302            self.execute_rows_with_params_deadline(query, BTreeMap::new(), Some(deadline))?;
303        Ok(project_rows(rows, options.unwrap_or_default()))
304    }
305
306    /// Execute a parameterised query inside the transaction.
307    pub fn execute_with_params(
308        &mut self,
309        query: &str,
310        options: Option<ExecuteOptions>,
311        params: BTreeMap<String, LoraValue>,
312    ) -> Result<QueryResult, LoraError> {
313        let rows = self.execute_rows_with_params_deadline(query, params, None)?;
314        Ok(project_rows(rows, options.unwrap_or_default()))
315    }
316
317    /// Execute a parameterised query inside the transaction with a cooperative
318    /// deadline.
319    pub fn execute_with_params_timeout(
320        &mut self,
321        query: &str,
322        options: Option<ExecuteOptions>,
323        params: BTreeMap<String, LoraValue>,
324        timeout: Duration,
325    ) -> Result<QueryResult, LoraError> {
326        let deadline = Instant::now()
327            .checked_add(timeout)
328            .unwrap_or_else(Instant::now);
329        let rows = self.execute_rows_with_params_deadline(query, params, Some(deadline))?;
330        Ok(project_rows(rows, options.unwrap_or_default()))
331    }
332
333    /// Execute a query inside the transaction and return hydrated rows before
334    /// final result-format projection.
335    pub fn execute_rows(&mut self, query: &str) -> Result<Vec<Row>, LoraError> {
336        self.execute_rows_with_params(query, BTreeMap::new())
337    }
338
339    /// Execute a parameterised query inside the transaction and return hydrated
340    /// rows before final result-format projection.
341    pub fn execute_rows_with_params(
342        &mut self,
343        query: &str,
344        params: BTreeMap<String, LoraValue>,
345    ) -> Result<Vec<Row>, LoraError> {
346        Ok(self.execute_rows_with_params_deadline(query, params, None)?)
347    }
348
349    fn execute_rows_with_params_deadline(
350        &mut self,
351        query: &str,
352        params: BTreeMap<String, LoraValue>,
353        deadline: Option<Instant>,
354    ) -> Result<Vec<Row>> {
355        let compiled = self.compile_in_tx(query)?;
356        self.execute_rows_compiled_deadline(&compiled, params, deadline)
357    }
358
359    fn execute_rows_compiled_deadline(
360        &mut self,
361        compiled: &CompiledQuery,
362        params: BTreeMap<String, LoraValue>,
363        deadline: Option<Instant>,
364    ) -> Result<Vec<Row>> {
365        // ReadOnly tx: never clones, runs straight against live.
366        if self.is_read_only_unchecked() {
367            self.precheck_open_no_savepoint()?;
368            let live = self.live.as_ref().ok_or(TransactionError::NoGraphGuard)?;
369            let storage = live.as_graph();
370            let executor = Executor::with_deadline(ExecutionContext { storage, params }, deadline);
371            return executor
372                .execute_compiled_rows(compiled)
373                .map_err(anyhow::Error::from);
374        }
375
376        // ReadWrite tx, lazy-clone aware.
377        let mut inner = self.begin_statement()?;
378        let is_mutating = classify_stream(compiled).is_mutating();
379
380        if !is_mutating {
381            // Read-only statement in a ReadWrite tx. Run against
382            // staged if it has been materialized (so the read
383            // sees prior in-tx writes), otherwise straight off
384            // the live graph — which equals staged-as-it-would-be
385            // because no writes have happened yet.
386            return match inner.staged.as_ref() {
387                Some(staged) => {
388                    let executor = Executor::with_deadline(
389                        ExecutionContext {
390                            storage: staged,
391                            params,
392                        },
393                        deadline,
394                    );
395                    executor
396                        .execute_compiled_rows(compiled)
397                        .map_err(anyhow::Error::from)
398                }
399                None => {
400                    drop(inner);
401                    let live = self.live.as_ref().ok_or(TransactionError::NoGraphGuard)?;
402                    let storage = live.as_graph();
403                    let executor =
404                        Executor::with_deadline(ExecutionContext { storage, params }, deadline);
405                    executor
406                        .execute_compiled_rows(compiled)
407                        .map_err(anyhow::Error::from)
408                }
409            };
410        }
411
412        // Mutating statement: lazy-clone the live graph if this
413        // is the first write in the tx, then capture a savepoint
414        // and run the mutable executor.
415        let clone_savepoint_graph = inner.staged.is_some();
416        self.ensure_staged_locked(&mut inner)?;
417        let savepoint = Some(take_savepoint(&inner, clone_savepoint_graph));
418
419        let exec_result: ExecResultRows = {
420            let staged = inner.staged_mut()?;
421            let mut executor = MutableExecutor::with_deadline(
422                MutableExecutionContext {
423                    storage: staged,
424                    params,
425                },
426                deadline,
427            );
428            executor
429                .execute_compiled_rows(compiled)
430                .map_err(anyhow::Error::from)
431        };
432
433        match exec_result {
434            Ok(rows) => Ok(rows),
435            Err(err) => {
436                restore_savepoint(&mut inner, savepoint);
437                Err(err)
438            }
439        }
440    }
441
442    /// Open a streaming write cursor over the staged graph for a
443    /// pre-compiled mutating plan, used by the hidden auto-commit
444    /// stream path in `Database::stream_with_params`.
445    ///
446    /// The returned `Box<dyn RowSource + 'static>` may be either a
447    /// real per-row [`StreamingWriteCursor`][lora_executor::StreamingWriteCursor],
448    /// a mutable UNION cursor, or a [`BufferedRowSource`][lora_executor::BufferedRowSource]
449    /// for the remaining materialized leaves. Either way:
450    ///
451    /// * The cursor mutates the *staged* graph, never the live store.
452    /// * Mutations fire the [`BufferingRecorder`] installed on staged
453    ///   by [`Self::ensure_staged_locked`], which accumulates into
454    ///   `inner.buffer` and is replayed into the WAL on commit.
455    /// * `cursor_active` is set to `true` here. The caller MUST clear
456    ///   it before invoking [`Self::commit`] or [`Self::rollback`] —
457    ///   the cursor itself does not.
458    ///
459    /// # Safety
460    ///
461    /// The cursor is `'static` because it owns its compiled query (via
462    /// the supplied `Arc`) and aliases the staged graph through a raw
463    /// pointer. Soundness depends on the invariant that
464    /// `inner.staged` remains `Some(_)` at a stable address for the
465    /// cursor's lifetime. That invariant holds while
466    /// `cursor_active = true` blocks every other path that could
467    /// move or drop staged: explicit statements (`begin_statement`
468    /// rejects), `commit` and `rollback` (rejected until the caller
469    /// clears `cursor_active`).
470    pub(crate) fn open_streaming_compiled_autocommit(
471        &mut self,
472        compiled: Arc<CompiledQuery>,
473        params: BTreeMap<String, LoraValue>,
474    ) -> Result<Box<dyn RowSource + 'static>> {
475        if self.is_read_only_unchecked() {
476            return Err(TransactionError::StreamingRequiresReadWrite.into());
477        }
478
479        let mut inner = self.begin_statement()?;
480        self.ensure_staged_locked(&mut inner)?;
481        inner.activate_cursor();
482
483        // SAFETY: `inner.staged` is `Some` after `ensure_staged_locked`,
484        // and stays at the same address while `cursor_active = true`
485        // (see method-level safety note).
486        let staged_ptr: *mut InMemoryGraph = inner
487            .staged
488            .as_mut()
489            .expect("ensure_staged_locked guarantees Some")
490            as *mut _;
491        drop(inner);
492
493        // SAFETY: `compiled` (Arc held by the caller / AutoCommit guard)
494        // keeps the plan alive; `staged_ptr` is valid for the cursor's
495        // lifetime per the invariant above. We extend both lifetimes
496        // to `'static` so the resulting cursor can sit inside the
497        // `'static`-shaped AutoCommit variant of `QueryStream`.
498        let storage_static: &'static mut InMemoryGraph = unsafe { &mut *staged_ptr };
499        let compiled_static: &'static CompiledQuery =
500            unsafe { std::mem::transmute::<&CompiledQuery, _>(compiled.as_ref()) };
501
502        // `MutablePullExecutor::open_compiled` picks the narrowest
503        // cursor shape it can: per-row write cursor, branch-wise
504        // mutable UNION cursor, or a buffered materialized leaf.
505        let cursor = MutablePullExecutor::new(storage_static, params)
506            .open_compiled(compiled_static)
507            .map_err(|e| {
508                // Roll back: the cursor build never happened, so the
509                // tx is in a clean-but-poisoned state. Discard
510                // everything and let the caller bubble the error.
511                if let Ok(mut inner) = self.inner.lock() {
512                    discard_transaction_state(&mut inner);
513                }
514                self.live.take();
515                anyhow::Error::from(e)
516            })?;
517
518        // The Arc<CompiledQuery> is the safety anchor for the
519        // `'static` plan reference. Keep it alive for the cursor's
520        // lifetime by leaking a clone into the cursor's owned data.
521        // We can't store it on the cursor itself (it's a Box<dyn>),
522        // so we wrap the cursor in a guard that owns the Arc.
523        Ok(Box::new(StreamingCursorWithArc {
524            cursor,
525            _compiled: compiled,
526        }))
527    }
528
529    /// Compile a query in this transaction's view of the world:
530    /// against `staged` if it has been materialized, otherwise
531    /// straight against `live`. The two are equivalent before the
532    /// first mutating statement, so the resulting plan is valid
533    /// either way.
534    fn compile_in_tx(&self, query: &str) -> Result<CompiledQuery> {
535        let document = parse_query(query)?;
536        let resolved = {
537            let inner = self.lock_inner_unchecked();
538            if let Some(staged) = &inner.staged {
539                let mut analyzer = Analyzer::new(staged);
540                analyzer.analyze(&document)?
541            } else {
542                drop(inner);
543                let live = self.live.as_ref().ok_or(TransactionError::NoGraphGuard)?;
544                let mut analyzer = Analyzer::new(live.as_graph());
545                analyzer.analyze(&document)?
546            }
547        };
548        Ok(Compiler::compile(&resolved))
549    }
550
551    /// Materialize `inner.staged` if it doesn't exist yet —
552    /// ReadWrite transactions defer this clone until the first
553    /// mutating statement.
554    fn ensure_staged_locked(&self, inner: &mut MutexGuard<'_, TxInner>) -> Result<()> {
555        if inner.staged.is_some() {
556            return Ok(());
557        }
558        let live = self.live.as_ref().ok_or(TransactionError::NoGraphGuard)?;
559        let mut staged: InMemoryGraph = live.as_graph().clone();
560        if matches!(inner.mode, TransactionMode::ReadWrite) && inner.buffer_mutations {
561            staged.set_mutation_recorder(Some(
562                Arc::new(BufferingRecorder::new(inner.buffer.clone())) as Arc<dyn MutationRecorder>,
563            ));
564        }
565        inner.staged = Some(staged);
566        Ok(())
567    }
568
569    /// Execute a query inside the transaction and return an owning row stream.
570    pub fn stream(&mut self, query: &str) -> Result<QueryStream<'static>, LoraError> {
571        self.stream_with_params(query, BTreeMap::new())
572    }
573
574    /// Execute a parameterised query inside the transaction and return an
575    /// owning row stream.
576    pub fn stream_with_params(
577        &mut self,
578        query: &str,
579        params: BTreeMap<String, LoraValue>,
580    ) -> Result<QueryStream<'static>, LoraError> {
581        let compiled = Arc::new(self.compile_in_tx(query)?);
582        let columns = compiled_result_columns(&compiled);
583        Ok(self.stream_compiled(compiled, columns, params)?)
584    }
585
586    /// Open a tx-bound stream for an already-compiled plan. Lets
587    /// `Database::stream_with_params` reuse the plan it built for
588    /// classification.
589    pub(crate) fn stream_compiled(
590        &mut self,
591        compiled: Arc<CompiledQuery>,
592        columns: Vec<String>,
593        params: BTreeMap<String, LoraValue>,
594    ) -> Result<QueryStream<'static>> {
595        let mut inner = self.begin_statement()?;
596        let is_mutating = classify_stream(&compiled).is_mutating();
597        if matches!(inner.mode, TransactionMode::ReadOnly) && is_mutating {
598            return Err(TransactionError::ReadOnlyMutation.into());
599        }
600
601        // Transaction streams borrow from the staged graph. Even
602        // read-only streams materialize staging when needed so the
603        // cursor can outlive the `&mut Transaction` borrow without
604        // borrowing from the transaction-owned live write guard.
605        let clone_savepoint_graph = inner.staged.is_some();
606        self.ensure_staged_locked(&mut inner)?;
607        inner.activate_cursor();
608
609        let rollback_on_drop = is_mutating;
610        if rollback_on_drop {
611            inner.pending_savepoint = Some(take_savepoint(&inner, clone_savepoint_graph));
612        } else {
613            inner.pending_savepoint = None;
614        }
615
616        let staged_ptr: *mut InMemoryGraph = inner
617            .staged
618            .as_mut()
619            .expect("ensure_staged_locked guarantees Some")
620            as *mut _;
621        drop(inner);
622
623        let compiled_static: &'static CompiledQuery =
624            unsafe { std::mem::transmute::<&CompiledQuery, _>(compiled.as_ref()) };
625        let cursor: Result<Box<dyn RowSource + 'static>> = if is_mutating {
626            let storage_static: &'static mut InMemoryGraph = unsafe { &mut *staged_ptr };
627            MutablePullExecutor::new(storage_static, params)
628                .open_compiled(compiled_static)
629                .map(|cursor| {
630                    Box::new(StreamingCursorWithArc {
631                        cursor,
632                        _compiled: compiled.clone(),
633                    }) as Box<dyn RowSource + 'static>
634                })
635                .map_err(anyhow::Error::from)
636        } else {
637            let storage_static: &'static InMemoryGraph = unsafe { &*staged_ptr };
638            PullExecutor::new(storage_static, params)
639                .open_compiled(compiled_static)
640                .map(|cursor| {
641                    Box::new(StreamingCursorWithArc {
642                        cursor,
643                        _compiled: compiled.clone(),
644                    }) as Box<dyn RowSource + 'static>
645                })
646                .map_err(anyhow::Error::from)
647        };
648
649        match cursor {
650            Ok(cursor) => Ok(QueryStream::for_tx_cursor(
651                cursor,
652                columns,
653                TxCursorLease::new(self.inner.clone(), rollback_on_drop),
654            )),
655            Err(err) => {
656                finalize_tx_stream(&self.inner, TxStreamOutcome::Interrupted, rollback_on_drop);
657                Err(err)
658            }
659        }
660    }
661
662    /// Commit the transaction and publish staged changes.
663    ///
664    /// When WAL is attached the buffered tx-local mutation log is
665    /// replayed into the durable WAL as a single committed
666    /// transaction; recovery therefore observes either every write
667    /// in this transaction or none.
668    pub fn commit(mut self) -> Result<(), LoraError> {
669        let CommitState {
670            staged,
671            buffer_events,
672            mode,
673        } = self.take_commit_state()?;
674
675        let wrote_wal_commit = self.replay_commit_wal(mode, buffer_events)?;
676        self.publish_staged_graph(mode, staged, wrote_wal_commit)?;
677
678        self.live.take();
679        Ok(())
680    }
681
682    fn take_commit_state(&self) -> Result<CommitState> {
683        let mut inner = self.inner.lock().unwrap();
684        if inner.cursor_active {
685            return Err(TransactionError::CursorActiveCommit.into());
686        }
687        if inner.closed {
688            return Err(TransactionError::AlreadyClosed.into());
689        }
690
691        let mode = inner.mode;
692        // Both modes can have `staged = None`: ReadOnly never clones,
693        // and ReadWrite transactions that performed no writes leave
694        // staging unmaterialized too.
695        let staged = inner.staged.take();
696        let buffer_events = std::mem::take(&mut *inner.buffer.lock().unwrap());
697        inner.closed = true;
698
699        Ok(CommitState {
700            staged,
701            buffer_events,
702            mode,
703        })
704    }
705
706    fn replay_commit_wal(
707        &self,
708        mode: TransactionMode,
709        buffer_events: Vec<MutationEvent>,
710    ) -> Result<bool> {
711        let Some(rec) = &self.wal else {
712            return Ok(false);
713        };
714
715        if !matches!(mode, TransactionMode::ReadWrite) {
716            ensure_wal_not_poisoned(rec)?;
717            return Ok(false);
718        }
719
720        Ok(rec.commit_events(buffer_events)?.wrote())
721    }
722
723    fn publish_staged_graph(
724        &mut self,
725        mode: TransactionMode,
726        staged: Option<InMemoryGraph>,
727        wrote_wal_commit: bool,
728    ) -> Result<()> {
729        if !matches!(mode, TransactionMode::ReadWrite) {
730            return Ok(());
731        }
732
733        let Some(mut staged) = staged else {
734            return Ok(());
735        };
736
737        // Strip the buffering recorder from the staged graph before
738        // publishing it as the live store; the live store either has
739        // the durable WAL recorder reinstalled below or no recorder at
740        // all (for non-WAL databases).
741        staged.set_mutation_recorder(None);
742        let wal = self.wal.clone();
743        if let Some(rec) = &wal {
744            staged.set_mutation_recorder(Some(rec.clone() as Arc<dyn MutationRecorder>));
745        }
746
747        let live = self.live.as_mut().ok_or(TransactionError::NoGraphGuard)?;
748        let lease = match live {
749            LiveStoreGuard::Write(lease) => lease,
750            LiveStoreGuard::Read(_) => {
751                return Err(TransactionError::ReadOnlyCommit.into());
752            }
753        };
754
755        if wrote_wal_commit {
756            if let (Some(snapshots), Some(rec)) = (&self.snapshots, wal.as_ref()) {
757                snapshots.observe_commit(&staged, rec)?;
758            }
759        }
760
761        // Atomic publish — concurrent readers will see the new state on
762        // their next `load_full()`, while in-flight readers keep their
763        // existing `Arc<InMemoryGraph>` snapshot until they drop it.
764        lease.store.store(Arc::new(staged));
765
766        Ok(())
767    }
768
769    /// Roll back the transaction. Staged graph changes and buffered
770    /// mutations are discarded; the WAL is never armed.
771    pub fn rollback(mut self) -> Result<(), LoraError> {
772        let mut inner = self.inner.lock().unwrap();
773        if inner.closed {
774            return Err(TransactionError::AlreadyClosed.into());
775        }
776        discard_transaction_state(&mut inner);
777        drop(inner);
778        self.live.take();
779        Ok(())
780    }
781
782    /// Acquire the inner state for a new statement. Validates that
783    /// the transaction is still open and no cursor is active. The
784    /// staged graph is *not* required: ReadWrite transactions
785    /// defer the staging clone until the first mutating statement
786    /// (see [`Transaction::ensure_staged_locked`]).
787    fn begin_statement(&self) -> Result<MutexGuard<'_, TxInner>> {
788        let inner = self.inner.lock().unwrap();
789        if inner.closed {
790            return Err(TransactionError::AlreadyClosed.into());
791        }
792        if inner.cursor_active {
793            return Err(TransactionError::CursorActiveStatement.into());
794        }
795        Ok(inner)
796    }
797
798    /// Cheap state check for the ReadOnly fast path: closed +
799    /// cursor_active. No staged-graph check — ReadOnly tx has no
800    /// staged graph by construction.
801    fn precheck_open_no_savepoint(&self) -> Result<()> {
802        let inner = self.inner.lock().unwrap();
803        if inner.closed {
804            return Err(TransactionError::AlreadyClosed.into());
805        }
806        if inner.cursor_active {
807            return Err(TransactionError::CursorActiveStatement.into());
808        }
809        Ok(())
810    }
811
812    /// True if the transaction was begun in `ReadOnly` mode. Cheap
813    /// — `mode` doesn't change after `begin_transaction`, so we
814    /// pay one small state-lock acquisition.
815    fn is_read_only_unchecked(&self) -> bool {
816        matches!(self.mode, TransactionMode::ReadOnly)
817    }
818
819    fn lock_inner_unchecked(&self) -> MutexGuard<'_, TxInner> {
820        self.inner
821            .lock()
822            .unwrap_or_else(|poisoned| poisoned.into_inner())
823    }
824
825    pub(crate) fn release_streaming_cursor(&self) {
826        if let Ok(mut inner) = self.inner.lock() {
827            inner.release_cursor();
828        }
829    }
830}
831
832type ExecResultRows = Result<Vec<Row>>;
833
834struct CommitState {
835    staged: Option<InMemoryGraph>,
836    buffer_events: Vec<MutationEvent>,
837    mode: TransactionMode,
838}
839
840impl TxInner {
841    fn staged_mut(&mut self) -> Result<&mut InMemoryGraph> {
842        self.staged
843            .as_mut()
844            .ok_or(TransactionError::NoStagedGraph.into())
845    }
846
847    fn activate_cursor(&mut self) {
848        self.cursor_active = true;
849    }
850
851    fn release_cursor(&mut self) {
852        self.cursor_active = false;
853    }
854
855    fn clear_pending_savepoint(&mut self) {
856        self.pending_savepoint = None;
857    }
858
859    fn restore_pending_savepoint(&mut self) {
860        if let Some(sp) = self.pending_savepoint.take() {
861            apply_savepoint(self, sp);
862        }
863    }
864
865    fn finalize_stream(&mut self, outcome: TxStreamOutcome, rollback_on_drop: bool) {
866        self.release_cursor();
867
868        if self.closed {
869            discard_transaction_state(self);
870            return;
871        }
872
873        if outcome.should_restore_savepoint(rollback_on_drop) {
874            self.restore_pending_savepoint();
875        } else {
876            self.clear_pending_savepoint();
877        }
878    }
879}
880
881/// `RowSource` adapter that owns an `Arc<CompiledQuery>` so the
882/// inner cursor's `'static` borrows into the plan stay valid for the
883/// life of the wrapper. The inner cursor is stored first so it drops
884/// before the Arc, releasing any borrows back into the plan.
885struct StreamingCursorWithArc {
886    cursor: Box<dyn RowSource + 'static>,
887    _compiled: Arc<CompiledQuery>,
888}
889
890impl RowSource for StreamingCursorWithArc {
891    fn next_row(&mut self) -> lora_executor::ExecResult<Option<Row>> {
892        self.cursor.next_row()
893    }
894}
895
896fn finalize_tx_stream(
897    handle: &Arc<Mutex<TxInner>>,
898    outcome: TxStreamOutcome,
899    rollback_on_drop: bool,
900) {
901    if let Ok(mut inner) = handle.lock() {
902        inner.finalize_stream(outcome, rollback_on_drop);
903    }
904}
905
906fn discard_transaction_state(inner: &mut TxInner) {
907    // A full transaction rollback supersedes any pending cursor savepoint.
908    inner.clear_pending_savepoint();
909    inner.release_cursor();
910    inner.staged = None;
911    if let Ok(mut buf) = inner.buffer.lock() {
912        buf.clear();
913    }
914    inner.closed = true;
915}
916
917fn take_savepoint(inner: &TxInner, clone_staged: bool) -> Savepoint {
918    let buffer_len = inner.buffer.lock().ok().map(|b| b.len()).unwrap_or(0);
919    Savepoint {
920        staged: if clone_staged {
921            inner.staged.as_ref().cloned()
922        } else {
923            None
924        },
925        buffer_len,
926    }
927}
928
929fn restore_savepoint(inner: &mut TxInner, savepoint: Option<Savepoint>) {
930    if let Some(sp) = savepoint {
931        apply_savepoint(inner, sp);
932    }
933}
934
935fn apply_savepoint(inner: &mut TxInner, sp: Savepoint) {
936    if let Ok(mut buf) = inner.buffer.lock() {
937        buf.truncate(sp.buffer_len);
938    }
939
940    let Some(mut graph) = sp.staged else {
941        inner.staged = None;
942        return;
943    };
944
945    // Rebuild the staged graph from the snapshot and re-install the
946    // buffering recorder. `InMemoryGraph::clone` deliberately drops
947    // recorders, so the snapshot has none until we put it back.
948    if matches!(inner.mode, TransactionMode::ReadWrite) && inner.buffer_mutations {
949        graph.set_mutation_recorder(Some(
950            Arc::new(BufferingRecorder::new(inner.buffer.clone())) as Arc<dyn MutationRecorder>
951        ));
952    }
953    inner.staged = Some(graph);
954}
955
956impl Drop for Transaction<'_> {
957    fn drop(&mut self) {
958        // If the user never called commit/rollback, treat it as a
959        // rollback: drop staged changes and the buffered mutations.
960        // The live RwLock guard is released as part of dropping `self.live`.
961        if let Ok(mut inner) = self.inner.lock() {
962            if !inner.closed {
963                if inner.cursor_active {
964                    // A tx-bound stream may still be borrowing the
965                    // staged graph through `inner`. Leave that graph
966                    // in place until the stream drops, but mark the
967                    // transaction closed so finalization discards it
968                    // instead of making it commit-eligible.
969                    inner.closed = true;
970                } else {
971                    discard_transaction_state(&mut inner);
972                }
973            }
974        }
975    }
976}