Skip to main content

lora_database/
stream.rs

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