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}