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