1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
//! Streaming entry points on [`Database<InMemoryGraph>`].
//!
//! The actual [`QueryStream`] type and its three cursor variants live
//! in [`crate::stream`]; this module owns the *opening* of those
//! streams from a `Database` — including the read/mutating shape
//! split, the hidden auto-commit transaction wrapping for mutating
//! queries, and the `unsafe` `'static`-lifetime escape hatch used by
//! language bindings.
//!
//! [`Database::begin_transaction`] also lives here because the
//! mutating-stream path uses it as the hidden auto-commit transaction
//! origin.
use std::collections::BTreeMap;
use std::sync::Arc;
use anyhow::Result;
use lora_executor::{classify_stream, compiled_result_columns, LoraValue, StreamShape};
use lora_store::InMemoryGraph;
use crate::database::Database;
use crate::error::LoraError;
use crate::stream::{AutoCommitGuard, LiveCursor, QueryStream};
use crate::transaction::{LiveStoreGuard, Transaction, TransactionMode, WriteLease};
impl Database<InMemoryGraph> {
/// Start an explicit transaction.
///
/// Read-only transactions hold a shared read lock for their
/// lifetime; read-write transactions hold the write lock. The
/// staging clone is **lazy** — it only happens when a
/// [`TransactionMode::ReadWrite`] transaction sees its first
/// mutating statement. Materialized read-only statements run
/// straight against the live graph; tx-bound streams may still
/// clone so their cursors can own a stable view. ReadWrite
/// transactions that perform only materialized reads (or commit
/// empty) pay nothing for staging.
pub fn begin_transaction(&self, mode: TransactionMode) -> Result<Transaction<'_>, LoraError> {
let live = match mode {
TransactionMode::ReadOnly => LiveStoreGuard::Read(self.store.load_full()),
TransactionMode::ReadWrite => {
// Acquire the writer lock — writers serialize, but we
// do NOT clone the graph yet. Staging is lazy: the
// working copy is built only when the first mutating
// statement runs. This keeps a `begin_transaction →
// commit` round trip with no mutations cheap (matches
// the previous RwLock-based behavior).
let lock = self
.writer
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
let snapshot = self.store.load_full();
LiveStoreGuard::Write(WriteLease {
_writer_lock: lock,
store: self.store.clone(),
snapshot,
})
}
};
Ok(Transaction::new(
live,
self.wal.clone(),
self.snapshots.clone(),
self.changes.clone(),
mode,
))
}
/// Execute a query and return an owning row stream.
pub fn stream(&self, query: &str) -> Result<QueryStream<'_>, LoraError> {
self.stream_with_params(query, BTreeMap::new())
}
/// Execute a parameterised query and return an owning row stream.
///
/// The compiled plan is classified at open time. Read-only
/// queries run directly off the live store and yield a
/// buffered cursor with plan-derived columns. Mutating queries
/// are routed through a hidden read-write [`Transaction`]:
/// full cursor exhaustion calls `tx.commit` (publishing staged
/// changes and replaying the tx-local WAL buffer); a premature
/// drop or any error from `next_row` calls `tx.rollback` so
/// the live store and the WAL stay untouched.
pub fn stream_with_params(
&self,
query: &str,
params: BTreeMap<String, LoraValue>,
) -> Result<QueryStream<'_>, LoraError> {
// Classify by fetching (or compiling once into) the plan cache. The
// mutating branch hands the same `Arc<CompiledQuery>` straight to
// the hidden transaction, so we no longer recompile against the
// staged graph — and the read-only branch reuses the cached plan
// for every subsequent stream.
let (store_guard, store_epoch) = self.read_store_with_epoch_deadline(None)?;
let compiled_arc = self.compile_query_cached(query, &*store_guard, store_epoch)?;
let columns = compiled_result_columns(&compiled_arc);
let shape = classify_stream(&compiled_arc);
// Release the analyzer's lock before either branch
// re-acquires (read-only path keeps it; mutating path
// delegates to begin_transaction which takes its own).
drop(store_guard);
match shape {
StreamShape::ReadOnly => {
// True pull-shaped streaming. `LiveCursor` holds
// the live store lock and the cursor that
// borrows from it; its `Drop` releases them in
// the right order so the caller observes pure
// pull semantics with no intermediate
// materialization.
let compiled = (*compiled_arc).clone();
let live = LiveCursor::open(self.store.clone(), compiled, params)?;
Ok(QueryStream::live(live, columns))
}
StreamShape::Mutating => {
// Hidden auto-commit transaction. The transaction
// owns staging, the buffering recorder, savepoint
// management, and the WAL replay-on-commit logic;
// we just pick commit-on-exhaustion vs
// rollback-on-drop based on cursor state.
//
// The cursor returned by `open_streaming_compiled_autocommit`
// may be a real per-row `StreamingWriteCursor`,
// a mutable UNION cursor, or a buffered leaf for
// operators that still need full materialization.
// Either way the AutoCommit guard's
// drop/exhaustion semantics are identical. The
// compiled plan is already wrapped in an `Arc` so
// the cursor's `'static` borrows into it remain
// valid for the cursor's lifetime.
let mut tx = self.begin_transaction(TransactionMode::ReadWrite)?;
let cursor =
match tx.open_streaming_compiled_autocommit(compiled_arc.clone(), params) {
Ok(c) => c,
Err(err) => {
// Tx rolls back implicitly on drop here.
return Err(err.into());
}
};
let guard = AutoCommitGuard {
tx: Some(tx),
finalized: false,
};
Ok(QueryStream::auto_commit(cursor, columns, guard))
}
}
}
/// Begin a transaction that does not borrow the database handle.
///
/// # Safety
///
/// The transaction holds lock guards that borrow from the database's
/// internal locks. The caller must keep this exact `Arc` alive until
/// the transaction is committed, rolled back or dropped, and must not
/// move it to another thread while it holds the writer lock (the guard
/// is not `Send`). Intended for language bindings that run one
/// interactive transaction on a dedicated thread next to its `Arc`.
pub unsafe fn begin_transaction_owned(
self: &Arc<Self>,
mode: TransactionMode,
) -> Result<Transaction<'static>, LoraError> {
let tx = self.begin_transaction(mode)?;
Ok(std::mem::transmute::<Transaction<'_>, Transaction<'static>>(tx))
}
/// Open a stream whose lifetime can be carried by an outer owner that
/// also retains an `Arc<Database>`.
///
/// # Safety
///
/// The returned stream may contain lock guards that borrow from the
/// database's internal `RwLock`. The caller must keep this exact `Arc`
/// alive until the stream is dropped. This is intended for language
/// bindings that store both the `Arc<Database>` and the `QueryStream` in
/// the same opaque stream handle.
pub unsafe fn stream_with_params_owned(
self: &Arc<Self>,
query: &str,
params: BTreeMap<String, LoraValue>,
) -> Result<QueryStream<'static>, LoraError> {
let stream = self.stream_with_params(query, params)?;
Ok(std::mem::transmute::<QueryStream<'_>, QueryStream<'static>>(stream))
}
}