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}