Skip to main content

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