Skip to main content

lora_database/
transaction.rs

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