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