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}