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 mode,
65 ))
66 }
67
68 /// Execute a query and return an owning row stream.
69 pub fn stream(&self, query: &str) -> Result<QueryStream<'_>, LoraError> {
70 self.stream_with_params(query, BTreeMap::new())
71 }
72
73 /// Execute a parameterised query and return an owning row stream.
74 ///
75 /// The compiled plan is classified at open time. Read-only
76 /// queries run directly off the live store and yield a
77 /// buffered cursor with plan-derived columns. Mutating queries
78 /// are routed through a hidden read-write [`Transaction`]:
79 /// full cursor exhaustion calls `tx.commit` (publishing staged
80 /// changes and replaying the tx-local WAL buffer); a premature
81 /// drop or any error from `next_row` calls `tx.rollback` so
82 /// the live store and the WAL stay untouched.
83 pub fn stream_with_params(
84 &self,
85 query: &str,
86 params: BTreeMap<String, LoraValue>,
87 ) -> Result<QueryStream<'_>, LoraError> {
88 // Classify by fetching (or compiling once into) the plan cache. The
89 // mutating branch hands the same `Arc<CompiledQuery>` straight to
90 // the hidden transaction, so we no longer recompile against the
91 // staged graph — and the read-only branch reuses the cached plan
92 // for every subsequent stream.
93 let (store_guard, store_epoch) = self.read_store_with_epoch_deadline(None)?;
94 let compiled_arc = self.compile_query_cached(query, &*store_guard, store_epoch)?;
95 let columns = compiled_result_columns(&compiled_arc);
96 let shape = classify_stream(&compiled_arc);
97 // Release the analyzer's lock before either branch
98 // re-acquires (read-only path keeps it; mutating path
99 // delegates to begin_transaction which takes its own).
100 drop(store_guard);
101
102 match shape {
103 StreamShape::ReadOnly => {
104 // True pull-shaped streaming. `LiveCursor` holds
105 // the live store lock and the cursor that
106 // borrows from it; its `Drop` releases them in
107 // the right order so the caller observes pure
108 // pull semantics with no intermediate
109 // materialization.
110 let compiled = (*compiled_arc).clone();
111 let live = LiveCursor::open(self.store.clone(), compiled, params)?;
112 Ok(QueryStream::live(live, columns))
113 }
114 StreamShape::Mutating => {
115 // Hidden auto-commit transaction. The transaction
116 // owns staging, the buffering recorder, savepoint
117 // management, and the WAL replay-on-commit logic;
118 // we just pick commit-on-exhaustion vs
119 // rollback-on-drop based on cursor state.
120 //
121 // The cursor returned by `open_streaming_compiled_autocommit`
122 // may be a real per-row `StreamingWriteCursor`,
123 // a mutable UNION cursor, or a buffered leaf for
124 // operators that still need full materialization.
125 // Either way the AutoCommit guard's
126 // drop/exhaustion semantics are identical. The
127 // compiled plan is already wrapped in an `Arc` so
128 // the cursor's `'static` borrows into it remain
129 // valid for the cursor's lifetime.
130 let mut tx = self.begin_transaction(TransactionMode::ReadWrite)?;
131 let cursor =
132 match tx.open_streaming_compiled_autocommit(compiled_arc.clone(), params) {
133 Ok(c) => c,
134 Err(err) => {
135 // Tx rolls back implicitly on drop here.
136 return Err(err.into());
137 }
138 };
139 let guard = AutoCommitGuard {
140 tx: Some(tx),
141 finalized: false,
142 };
143 Ok(QueryStream::auto_commit(cursor, columns, guard))
144 }
145 }
146 }
147
148 /// Open a stream whose lifetime can be carried by an outer owner that
149 /// also retains an `Arc<Database>`.
150 ///
151 /// # Safety
152 ///
153 /// The returned stream may contain lock guards that borrow from the
154 /// database's internal `RwLock`. The caller must keep this exact `Arc`
155 /// alive until the stream is dropped. This is intended for language
156 /// bindings that store both the `Arc<Database>` and the `QueryStream` in
157 /// the same opaque stream handle.
158 pub unsafe fn stream_with_params_owned(
159 self: &Arc<Self>,
160 query: &str,
161 params: BTreeMap<String, LoraValue>,
162 ) -> Result<QueryStream<'static>, LoraError> {
163 let stream = self.stream_with_params(query, params)?;
164 Ok(std::mem::transmute::<QueryStream<'_>, QueryStream<'static>>(stream))
165 }
166}