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