Skip to main content

lora_database/database/
stream.rs

1//! Streaming entry points on [`Database<InMemoryGraph>`].
2//!
3//! The actual [`QueryStream`] type and its three cursor variants live
4//! in [`crate::stream`]; this module owns the *opening* of those
5//! streams from a `Database` — including the read/mutating shape
6//! split, the hidden auto-commit transaction wrapping for mutating
7//! queries, and the `unsafe` `'static`-lifetime escape hatch used by
8//! language bindings.
9//!
10//! [`Database::begin_transaction`] also lives here because the
11//! mutating-stream path uses it as the hidden auto-commit transaction
12//! origin.
13
14use std::collections::BTreeMap;
15use std::sync::Arc;
16
17use anyhow::Result;
18use lora_executor::{classify_stream, compiled_result_columns, LoraValue, StreamShape};
19use lora_store::InMemoryGraph;
20
21use crate::database::Database;
22use crate::error::LoraError;
23use crate::stream::{AutoCommitGuard, LiveCursor, QueryStream};
24use crate::transaction::{LiveStoreGuard, Transaction, TransactionMode, WriteLease};
25
26impl Database<InMemoryGraph> {
27    /// Start an explicit transaction.
28    ///
29    /// Read-only transactions hold a shared read lock for their
30    /// lifetime; read-write transactions hold the write lock. The
31    /// staging clone is **lazy** — it only happens when a
32    /// [`TransactionMode::ReadWrite`] transaction sees its first
33    /// mutating statement. Materialized read-only statements run
34    /// straight against the live graph; tx-bound streams may still
35    /// clone so their cursors can own a stable view. ReadWrite
36    /// transactions that perform only materialized reads (or commit
37    /// empty) pay nothing for staging.
38    pub fn begin_transaction(&self, mode: TransactionMode) -> Result<Transaction<'_>, LoraError> {
39        let live = match mode {
40            TransactionMode::ReadOnly => LiveStoreGuard::Read(self.store.load_full()),
41            TransactionMode::ReadWrite => {
42                // Acquire the writer lock — writers serialize, but we
43                // do NOT clone the graph yet. Staging is lazy: the
44                // working copy is built only when the first mutating
45                // statement runs. This keeps a `begin_transaction →
46                // commit` round trip with no mutations cheap (matches
47                // the previous RwLock-based behavior).
48                let lock = self
49                    .writer
50                    .lock()
51                    .unwrap_or_else(|poisoned| poisoned.into_inner());
52                let snapshot = self.store.load_full();
53                LiveStoreGuard::Write(WriteLease {
54                    _writer_lock: lock,
55                    store: self.store.clone(),
56                    snapshot,
57                })
58            }
59        };
60        Ok(Transaction::new(
61            live,
62            self.wal.clone(),
63            self.snapshots.clone(),
64            self.changes.clone(),
65            mode,
66        ))
67    }
68
69    /// Execute a query and return an owning row stream.
70    pub fn stream(&self, query: &str) -> Result<QueryStream<'_>, LoraError> {
71        self.stream_with_params(query, BTreeMap::new())
72    }
73
74    /// Execute a parameterised query and return an owning row stream.
75    ///
76    /// The compiled plan is classified at open time. Read-only
77    /// queries run directly off the live store and yield a
78    /// buffered cursor with plan-derived columns. Mutating queries
79    /// are routed through a hidden read-write [`Transaction`]:
80    /// full cursor exhaustion calls `tx.commit` (publishing staged
81    /// changes and replaying the tx-local WAL buffer); a premature
82    /// drop or any error from `next_row` calls `tx.rollback` so
83    /// the live store and the WAL stay untouched.
84    pub fn stream_with_params(
85        &self,
86        query: &str,
87        params: BTreeMap<String, LoraValue>,
88    ) -> Result<QueryStream<'_>, LoraError> {
89        // Classify by fetching (or compiling once into) the plan cache. The
90        // mutating branch hands the same `Arc<CompiledQuery>` straight to
91        // the hidden transaction, so we no longer recompile against the
92        // staged graph — and the read-only branch reuses the cached plan
93        // for every subsequent stream.
94        let (store_guard, store_epoch) = self.read_store_with_epoch_deadline(None)?;
95        let compiled_arc = self.compile_query_cached(query, &*store_guard, store_epoch)?;
96        let columns = compiled_result_columns(&compiled_arc);
97        let shape = classify_stream(&compiled_arc);
98        // Release the analyzer's lock before either branch
99        // re-acquires (read-only path keeps it; mutating path
100        // delegates to begin_transaction which takes its own).
101        drop(store_guard);
102
103        match shape {
104            StreamShape::ReadOnly => {
105                // True pull-shaped streaming. `LiveCursor` holds
106                // the live store lock and the cursor that
107                // borrows from it; its `Drop` releases them in
108                // the right order so the caller observes pure
109                // pull semantics with no intermediate
110                // materialization.
111                let compiled = (*compiled_arc).clone();
112                let live = LiveCursor::open(self.store.clone(), compiled, params)?;
113                Ok(QueryStream::live(live, columns))
114            }
115            StreamShape::Mutating => {
116                // Hidden auto-commit transaction. The transaction
117                // owns staging, the buffering recorder, savepoint
118                // management, and the WAL replay-on-commit logic;
119                // we just pick commit-on-exhaustion vs
120                // rollback-on-drop based on cursor state.
121                //
122                // The cursor returned by `open_streaming_compiled_autocommit`
123                // may be a real per-row `StreamingWriteCursor`,
124                // a mutable UNION cursor, or a buffered leaf for
125                // operators that still need full materialization.
126                // Either way the AutoCommit guard's
127                // drop/exhaustion semantics are identical. The
128                // compiled plan is already wrapped in an `Arc` so
129                // the cursor's `'static` borrows into it remain
130                // valid for the cursor's lifetime.
131                let mut tx = self.begin_transaction(TransactionMode::ReadWrite)?;
132                let cursor =
133                    match tx.open_streaming_compiled_autocommit(compiled_arc.clone(), params) {
134                        Ok(c) => c,
135                        Err(err) => {
136                            // Tx rolls back implicitly on drop here.
137                            return Err(err.into());
138                        }
139                    };
140                let guard = AutoCommitGuard {
141                    tx: Some(tx),
142                    finalized: false,
143                };
144                Ok(QueryStream::auto_commit(cursor, columns, guard))
145            }
146        }
147    }
148
149    /// Begin a transaction that does not borrow the database handle.
150    ///
151    /// # Safety
152    ///
153    /// The transaction holds lock guards that borrow from the database's
154    /// internal locks. The caller must keep this exact `Arc` alive until
155    /// the transaction is committed, rolled back or dropped, and must not
156    /// move it to another thread while it holds the writer lock (the guard
157    /// is not `Send`). Intended for language bindings that run one
158    /// interactive transaction on a dedicated thread next to its `Arc`.
159    pub unsafe fn begin_transaction_owned(
160        self: &Arc<Self>,
161        mode: TransactionMode,
162    ) -> Result<Transaction<'static>, LoraError> {
163        let tx = self.begin_transaction(mode)?;
164        Ok(std::mem::transmute::<Transaction<'_>, Transaction<'static>>(tx))
165    }
166
167    /// Open a stream whose lifetime can be carried by an outer owner that
168    /// also retains an `Arc<Database>`.
169    ///
170    /// # Safety
171    ///
172    /// The returned stream may contain lock guards that borrow from the
173    /// database's internal `RwLock`. The caller must keep this exact `Arc`
174    /// alive until the stream is dropped. This is intended for language
175    /// bindings that store both the `Arc<Database>` and the `QueryStream` in
176    /// the same opaque stream handle.
177    pub unsafe fn stream_with_params_owned(
178        self: &Arc<Self>,
179        query: &str,
180        params: BTreeMap<String, LoraValue>,
181    ) -> Result<QueryStream<'static>, LoraError> {
182        let stream = self.stream_with_params(query, params)?;
183        Ok(std::mem::transmute::<QueryStream<'_>, QueryStream<'static>>(stream))
184    }
185}