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 `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}