lora_database/stream.rs
1use std::collections::BTreeMap;
2use std::mem::ManuallyDrop;
3use std::sync::Arc;
4
5use anyhow::{anyhow, Result};
6use lora_compiler::CompiledQuery;
7use lora_executor::{ExecResult, LoraValue, PullExecutor, Row, RowSource};
8use lora_store::InMemoryGraph;
9
10use crate::live_store::LiveStore;
11use crate::transaction::{Transaction, TxCursorLease, TxStreamOutcome};
12
13/// Owning row stream returned by [`crate::Database::stream`] and transaction
14/// streaming methods.
15///
16/// The cursor is fallible (`next_row()` surfaces execution errors)
17/// and exposes plan-derived column names populated even for empty
18/// results. The lifetime parameter `'a` is bound to the source the
19/// cursor borrows from — typically the database for auto-commit
20/// write streams that hold the live write guard until exhaustion or
21/// drop. Read-only and transaction-bound streams need no live
22/// borrow and use `'static` (the buffered variant).
23pub struct QueryStream<'a> {
24 columns: Vec<String>,
25 inner: StreamInner<'a>,
26}
27
28enum StreamInner<'a> {
29 /// Transaction-bound streaming cursor. The cursor borrows from
30 /// the transaction's staged graph, which is kept alive by
31 /// `lease`; finalization releases the cursor token and either
32 /// clears or restores the pending statement savepoint.
33 Tx {
34 cursor: Option<Box<dyn RowSource + 'static>>,
35 state: StreamState,
36 /// Lease releases the transaction cursor token and either
37 /// clears or restores the pending statement savepoint.
38 lease: TxCursorLease,
39 },
40 /// True pull-based read-only stream. Holds a live store read
41 /// lock through the cursor's lifetime and emits rows as the
42 /// caller pulls them, without any intermediate
43 /// materialization. Backed by a [`LiveCursor`] which uses
44 /// `self_cell` to safely co-own the lock guard and the
45 /// borrowing cursor.
46 Live {
47 cursor: LiveCursor,
48 state: StreamState,
49 // The 'a parameter is unused for this variant — the
50 // self-cell hides the borrow. We carry a phantom to keep
51 // the enum's lifetime parameter consistent with the
52 // other variants.
53 _phantom: std::marker::PhantomData<&'a ()>,
54 },
55 /// Auto-commit write stream backed by a hidden staged
56 /// transaction. The graph is mutated on a clone held in
57 /// `guard.tx.inner.staged`; the live store write lock stays locked
58 /// through the tx's `live` guard so no other writer races. On
59 /// full exhaustion the staged graph is published and the WAL
60 /// replays the buffered events; on premature drop or error the
61 /// staged graph and buffer are discarded and the live store is
62 /// untouched.
63 ///
64 /// `cursor` is a streaming `RowSource` that may apply mutations
65 /// row-by-row (via `StreamingWriteCursor`) or yield from a
66 /// pre-materialized buffer (via `BufferedRowSource`); see
67 /// `Transaction::open_streaming_compiled_autocommit`. It is
68 /// taken and dropped before the guard's commit/rollback so any
69 /// borrows back into the staged graph are released first.
70 AutoCommit {
71 cursor: Option<Box<dyn RowSource + 'static>>,
72 state: StreamState,
73 guard: AutoCommitGuard<'a>,
74 },
75}
76
77/// Self-referential cursor that pulls rows directly from a snapshot
78/// of the live store. We hold an `Arc<InMemoryGraph>` (loaded once
79/// from the database's `LiveStore` at open time) and the boxed
80/// `RowSource` borrows from `&*snapshot`. Drop order — `cursor`
81/// first, then `_snapshot` — guarantees the cursor never sees a
82/// freed graph.
83///
84/// Snapshot isolation is automatic: even if a writer commits a new
85/// version while this cursor is live, the cursor keeps observing the
86/// graph it was opened against until it drops the Arc.
87pub(crate) struct LiveCursor {
88 /// SAFETY invariant: borrows from `&*_snapshot` and `&*_compiled`.
89 /// Must drop before either.
90 cursor: ManuallyDrop<Box<dyn RowSource + 'static>>,
91 /// Pinned snapshot the cursor borrows from. Dropped after `cursor`.
92 _snapshot: Arc<InMemoryGraph>,
93 /// Pinned live-store owner so the database storage outlives the borrowed
94 /// cursor shape, even though the cursor itself only reads from `_snapshot`.
95 _store: Arc<LiveStore<InMemoryGraph>>,
96 /// Keeps the compiled plan alive — operator sources hold
97 /// references into it (e.g. predicate `ResolvedExpr`s). Boxed
98 /// so the plan address is stable across the move into the
99 /// struct.
100 _compiled: Box<CompiledQuery>,
101}
102
103impl LiveCursor {
104 /// Snapshot the live store and open a streaming cursor against
105 /// the given compiled query. Internal helper for
106 /// `Database::stream_with_params` — never expose the
107 /// constructed `LiveCursor` to callers without the
108 /// surrounding `QueryStream`, which makes the `'static`
109 /// transmutes invisible.
110 pub(crate) fn open(
111 store: Arc<LiveStore<InMemoryGraph>>,
112 compiled: CompiledQuery,
113 params: BTreeMap<String, LoraValue>,
114 ) -> Result<Self> {
115 let compiled = Box::new(compiled);
116 let snapshot = store.load_full();
117
118 // SAFETY: We extend the lifetime of borrows into `&*snapshot`
119 // and `&*compiled` to `'static`. This is sound because the
120 // surrounding `LiveCursor` keeps:
121 // (a) the `Arc<InMemoryGraph>` alive while the cursor is
122 // alive — the graph behind it is never freed; and
123 // (b) the `Box<CompiledQuery>` alive while the cursor is
124 // alive.
125 // The `Drop` impl below releases `cursor` before `_snapshot`,
126 // so neither borrow can outlive its backing storage.
127 let storage_ref: &'static InMemoryGraph =
128 unsafe { std::mem::transmute::<&InMemoryGraph, _>(&*snapshot) };
129 let compiled_ref: &'static CompiledQuery =
130 unsafe { std::mem::transmute::<&CompiledQuery, _>(&*compiled) };
131
132 let cursor = PullExecutor::new(storage_ref, params)
133 .open_compiled(compiled_ref)
134 .map_err(|e| anyhow!(e))?;
135
136 Ok(Self {
137 cursor: ManuallyDrop::new(cursor),
138 _snapshot: snapshot,
139 _store: store,
140 _compiled: compiled,
141 })
142 }
143
144 fn next_row(&mut self) -> ExecResult<Option<Row>> {
145 self.cursor.next_row()
146 }
147}
148
149impl Drop for LiveCursor {
150 fn drop(&mut self) {
151 // SAFETY: drop in the documented order — cursor first
152 // (releases its borrow into `*_snapshot` / `*_compiled`),
153 // then the rest drops via field-drop ordering. After this
154 // call we never touch `cursor` again.
155 unsafe {
156 ManuallyDrop::drop(&mut self.cursor);
157 }
158 }
159}
160
161/// Per-stream state held only by auto-commit write streams.
162///
163/// The auto-commit guard is a thin wrapper around an explicit
164/// [`Transaction`]: full cursor exhaustion calls `commit`,
165/// premature drop or error calls `rollback`. All staged-graph,
166/// savepoint, and WAL replay logic lives on `Transaction` itself,
167/// so the guard contributes no behavior of its own beyond the
168/// commit-vs-rollback decision.
169pub(crate) struct AutoCommitGuard<'a> {
170 /// The hidden transaction. `None` once the guard has finalized
171 /// (commit consumes the tx; rollback consumes it; both leave
172 /// `None` behind).
173 pub(crate) tx: Option<Transaction<'a>>,
174 /// Set once a finalization (commit or rollback) has run so
175 /// duplicate calls — including the `Drop` path after a
176 /// successful `next_row`-driven commit — are no-ops.
177 pub(crate) finalized: bool,
178}
179
180#[derive(Debug, Clone, Copy, PartialEq, Eq)]
181enum StreamState {
182 Active,
183 Exhausted,
184 Errored,
185}
186
187impl<'a> std::fmt::Debug for QueryStream<'a> {
188 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
189 let state = match &self.inner {
190 StreamInner::Tx { state, .. }
191 | StreamInner::AutoCommit { state, .. }
192 | StreamInner::Live { state, .. } => *state,
193 };
194 f.debug_struct("QueryStream")
195 .field("columns", &self.columns)
196 .field("state", &state)
197 .finish()
198 }
199}
200
201impl<'a> QueryStream<'a> {
202 pub(crate) fn for_tx_cursor(
203 cursor: Box<dyn RowSource + 'static>,
204 columns: Vec<String>,
205 lease: TxCursorLease,
206 ) -> Self {
207 Self {
208 columns,
209 inner: StreamInner::Tx {
210 cursor: Some(cursor),
211 state: StreamState::Active,
212 lease,
213 },
214 }
215 }
216
217 pub(crate) fn auto_commit(
218 cursor: Box<dyn RowSource + 'static>,
219 columns: Vec<String>,
220 guard: AutoCommitGuard<'a>,
221 ) -> Self {
222 Self {
223 columns,
224 inner: StreamInner::AutoCommit {
225 cursor: Some(cursor),
226 state: StreamState::Active,
227 guard,
228 },
229 }
230 }
231
232 pub(crate) fn live(cursor: LiveCursor, columns: Vec<String>) -> Self {
233 Self {
234 columns,
235 inner: StreamInner::Live {
236 cursor,
237 state: StreamState::Active,
238 _phantom: std::marker::PhantomData,
239 },
240 }
241 }
242
243 /// Plan-derived column names. Populated even when the result is
244 /// empty so callers can drive a row-arrays format off this list
245 /// without first peeking at a materialized row.
246 pub fn columns(&self) -> &[String] {
247 &self.columns
248 }
249
250 /// Pull the next row. Returns `Ok(None)` once the cursor is
251 /// exhausted, `Ok(Some(row))` for the next hydrated row, or an
252 /// error if the underlying execution failed. Once an error has
253 /// been observed, subsequent calls keep returning that terminal
254 /// state — the cursor never tries to recover or re-execute.
255 pub fn next_row(&mut self) -> Result<Option<Row>> {
256 match &mut self.inner {
257 StreamInner::Live { state, cursor, .. } => match *state {
258 StreamState::Errored => Err(anyhow!("query stream errored")),
259 StreamState::Exhausted => Ok(None),
260 StreamState::Active => match cursor.next_row() {
261 Ok(Some(row)) => Ok(Some(row)),
262 Ok(None) => {
263 *state = StreamState::Exhausted;
264 Ok(None)
265 }
266 Err(e) => {
267 *state = StreamState::Errored;
268 Err(anyhow!(e))
269 }
270 },
271 },
272 StreamInner::Tx {
273 state,
274 cursor,
275 lease,
276 } => match *state {
277 StreamState::Errored => Err(anyhow!("query stream errored")),
278 StreamState::Exhausted => Ok(None),
279 StreamState::Active => {
280 let pull = match cursor.as_mut() {
281 Some(c) => c.next_row(),
282 None => {
283 *state = StreamState::Errored;
284 return Err(anyhow!("transaction cursor missing"));
285 }
286 };
287 match pull {
288 Ok(Some(row)) => Ok(Some(row)),
289 Ok(None) => {
290 cursor.take();
291 lease.finalize(TxStreamOutcome::Exhausted);
292 *state = StreamState::Exhausted;
293 Ok(None)
294 }
295 Err(e) => {
296 cursor.take();
297 lease.finalize(TxStreamOutcome::Interrupted);
298 *state = StreamState::Errored;
299 Err(anyhow!(e))
300 }
301 }
302 }
303 },
304 StreamInner::AutoCommit {
305 state,
306 cursor,
307 guard,
308 } => match *state {
309 StreamState::Errored => Err(anyhow!("query stream errored")),
310 StreamState::Exhausted => Ok(None),
311 StreamState::Active => {
312 let pull = match cursor.as_mut() {
313 Some(c) => c.next_row(),
314 None => {
315 *state = StreamState::Errored;
316 return Err(anyhow!("auto-commit cursor missing"));
317 }
318 };
319 match pull {
320 Ok(Some(row)) => Ok(Some(row)),
321 Ok(None) => {
322 // Drop the cursor first so its borrows
323 // into the staged graph release before
324 // commit moves staged out of inner.
325 cursor.take();
326 match guard.commit() {
327 Ok(()) => {
328 *state = StreamState::Exhausted;
329 Ok(None)
330 }
331 Err(e) => {
332 *state = StreamState::Errored;
333 Err(e)
334 }
335 }
336 }
337 Err(e) => {
338 cursor.take();
339 guard.rollback();
340 *state = StreamState::Errored;
341 Err(anyhow!(e))
342 }
343 }
344 }
345 },
346 }
347 }
348
349 /// True once the stream has produced its last row.
350 fn is_exhausted(&self) -> bool {
351 match &self.inner {
352 StreamInner::Tx { state, .. }
353 | StreamInner::AutoCommit { state, .. }
354 | StreamInner::Live { state, .. } => matches!(state, StreamState::Exhausted),
355 }
356 }
357}
358
359impl<'a> Iterator for QueryStream<'a> {
360 type Item = Row;
361
362 fn next(&mut self) -> Option<Self::Item> {
363 match self.next_row() {
364 Ok(Some(row)) => Some(row),
365 Ok(None) => None,
366 Err(_) => None,
367 }
368 }
369
370 fn size_hint(&self) -> (usize, Option<usize>) {
371 match &self.inner {
372 // Live and AutoCommit (now backed by a streaming cursor)
373 // don't know their length until drained.
374 StreamInner::Live { .. } | StreamInner::Tx { .. } | StreamInner::AutoCommit { .. } => {
375 (0, None)
376 }
377 }
378 }
379}
380
381// Note: `ExactSizeIterator` intentionally not implemented. The
382// `Live` variant produces rows lazily and can't report an exact
383// remaining count.
384
385impl<'a> Drop for QueryStream<'a> {
386 fn drop(&mut self) {
387 let exhausted = self.is_exhausted();
388 match &mut self.inner {
389 StreamInner::Tx { cursor, lease, .. } => {
390 cursor.take();
391 let outcome = if exhausted {
392 TxStreamOutcome::Exhausted
393 } else {
394 TxStreamOutcome::Interrupted
395 };
396 lease.finalize(outcome);
397 }
398 StreamInner::Live { .. } => {
399 // Drop releases the cursor, then the read guard,
400 // which releases the live store read lock. No
401 // additional cleanup needed — live streams never
402 // mutate, so there is nothing to commit or roll back.
403 }
404 StreamInner::AutoCommit {
405 state,
406 cursor,
407 guard,
408 } => {
409 // Drop the cursor first so its borrows into the
410 // staged graph release before the guard rolls back
411 // (which moves staged to None).
412 cursor.take();
413 // Premature drop = rollback. Successful exhaustion
414 // already finalized the guard via `commit()` in
415 // `next_row`, so this path is a no-op for the
416 // exhausted case.
417 if !guard.finalized && !matches!(state, StreamState::Exhausted) {
418 guard.rollback();
419 }
420 }
421 }
422 }
423}
424
425impl<'a> AutoCommitGuard<'a> {
426 /// Publish the staged graph as the live store. Delegates to
427 /// [`Transaction::commit`] which owns the WAL replay + swap
428 /// logic. Idempotent — subsequent calls are no-ops once
429 /// finalized, regardless of whether the previous attempt
430 /// succeeded or failed.
431 fn commit(&mut self) -> Result<()> {
432 if self.finalized {
433 return Ok(());
434 }
435 // Mark finalized before consuming the tx so a commit
436 // failure still prevents Drop from later trying to roll
437 // back a tx that no longer exists.
438 self.finalized = true;
439 match self.tx.take() {
440 Some(tx) => {
441 // The streaming auto-commit cursor sets
442 // `cursor_active = true` at construction; it must
443 // be cleared before `tx.commit` (which rejects on
444 // an active cursor). The cursor itself was already
445 // dropped by the caller in `next_row` — its
446 // borrows back into staged are gone, so we can
447 // safely flip the flag here. For the buffered
448 // fallback path the flag was never set, so this
449 // assignment is a no-op.
450 tx.release_streaming_cursor();
451 Ok(tx.commit()?)
452 }
453 None => Ok(()),
454 }
455 }
456
457 /// Discard the staged graph. Delegates to
458 /// [`Transaction::rollback`]; failures are swallowed because
459 /// the rollback path runs from `Drop` and has nowhere to
460 /// surface an error.
461 fn rollback(&mut self) {
462 if self.finalized {
463 return;
464 }
465 self.finalized = true;
466 if let Some(tx) = self.tx.take() {
467 // Clear the streaming-cursor flag before delegating to
468 // tx.rollback so the rollback can finalize without
469 // stumbling over a stale `cursor_active = true`.
470 tx.release_streaming_cursor();
471 let _ = tx.rollback();
472 }
473 }
474}