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
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
//! Query-execute family on [`Database<S>`].
//!
//! These methods cover the `execute_*` and `execute_rows_*` entry points
//! plus the cooperative-deadline variants and the plan-cache-backed
//! compile path. The actual write routing (optimistic OCC vs. pessimistic
//! `with_logged_write_guard`) is dispatched from
//! [`Database::execute_rows_with_params_deadline`]; the OCC body lives
//! in [`super::occ`] and the WAL-bracketed write closure lives in
//! [`super::write_guard`].
use std::any::Any;
use std::collections::BTreeMap;
use std::time::{Duration, Instant};
use anyhow::Result;
use lora_ast::Statement;
use lora_executor::{
classify_stream, collect_compiled, project_rows, ExecuteOptions, LoraValue,
MutableExecutionContext, MutableExecutor, QueryResult, Row, StreamShape,
};
use lora_parser::parse_query;
use lora_store::{GraphStorage, GraphStorageMut, InMemoryGraph};
use crate::database::{Database, QUERY_FAILURE_POISON};
use crate::error::LoraError;
use crate::wal::write_scope::{ensure_wal_query_can_start, WalAbortPolicy};
use super::pull_mode::should_collect_read_via_pull;
impl<S> Database<S>
where
S: GraphStorage + GraphStorageMut + Any + Clone + Send + Sync + 'static,
{
/// Execute a query and return its result.
pub fn execute(
&self,
query: &str,
options: Option<ExecuteOptions>,
) -> Result<QueryResult, LoraError> {
self.execute_with_params(query, options, BTreeMap::new())
}
/// Execute a query with a cooperative deadline. The timeout is checked at
/// executor operator boundaries and hot scan loops; if it fires, the query
/// returns an error and any WAL-backed mutating query is aborted through
/// the existing failure path.
pub fn execute_with_timeout(
&self,
query: &str,
options: Option<ExecuteOptions>,
timeout: Duration,
) -> Result<QueryResult, LoraError> {
let deadline = Instant::now()
.checked_add(timeout)
.unwrap_or_else(Instant::now);
let rows =
self.execute_rows_with_params_deadline(query, BTreeMap::new(), Some(deadline))?;
Ok(project_rows(rows, options.unwrap_or_default()))
}
/// Execute a query with bound parameters.
///
/// When a WAL is attached the call is bracketed by a transaction:
///
/// 1. `recorder.arm()` after analyze + compile (so a parse /
/// semantic / compile error never opens a tx that has to be
/// immediately aborted). Arming is *cheap*: no record is
/// appended to the WAL yet, so a pure read query that
/// completes here pays nothing for the WAL hot path.
/// 2. The executor runs; every primitive mutation fires
/// `MutationRecorder::record`, which buffers events in memory.
/// 3. On Ok, `recorder.commit()` writes `TxBegin`, one batched
/// mutation record, and `TxCommit` only when mutations occurred,
/// then applies the WAL's configured single-thread flush policy.
/// 4. On Err, `recorder.abort()` clears the pending batch. The
/// engine has no rollback, so the in-memory state may already
/// be partially mutated; the live handle is quarantined while
/// durable recovery stays atomic because no committed batch was
/// written.
/// 5. The recorder's poisoned flag is polled once. If set, the query
/// fails loudly with the durability error so the caller can act on it;
/// the WAL refuses further appends until the operator restarts the
/// database, which recovers from the last consistent snapshot + WAL.
pub fn execute_with_params(
&self,
query: &str,
options: Option<ExecuteOptions>,
params: BTreeMap<String, LoraValue>,
) -> Result<QueryResult, LoraError> {
let rows = self.execute_rows_with_params_deadline(query, params, None)?;
Ok(project_rows(rows, options.unwrap_or_default()))
}
/// Execute a parameterised query with a cooperative deadline.
pub fn execute_with_params_timeout(
&self,
query: &str,
options: Option<ExecuteOptions>,
params: BTreeMap<String, LoraValue>,
timeout: Duration,
) -> Result<QueryResult, LoraError> {
let deadline = Instant::now()
.checked_add(timeout)
.unwrap_or_else(Instant::now);
let rows = self.execute_rows_with_params_deadline(query, params, Some(deadline))?;
Ok(project_rows(rows, options.unwrap_or_default()))
}
/// Execute a query and return hydrated rows before final result-format
/// projection.
pub fn execute_rows(&self, query: &str) -> Result<Vec<Row>, LoraError> {
self.execute_rows_with_params(query, BTreeMap::new())
}
/// Execute a query with parameters and return hydrated rows before final
/// result-format projection.
pub fn execute_rows_with_params(
&self,
query: &str,
params: BTreeMap<String, LoraValue>,
) -> Result<Vec<Row>, LoraError> {
Ok(self.execute_rows_with_params_deadline(query, params, None)?)
}
pub(crate) fn execute_rows_with_params_deadline(
&self,
query: &str,
params: BTreeMap<String, LoraValue>,
deadline: Option<Instant>,
) -> Result<Vec<Row>> {
// Schema commands (CREATE INDEX, SHOW INDEXES) are catalog-level
// mutations and don't share the row-producing pipeline. A cheap
// textual prefix check routes them out before the parser cache and
// analyzer get involved.
if Self::is_schema_command_text(query) {
let document = parse_query(query)?;
if let Statement::Schema(cmd) = &document.statement {
return self.execute_schema_command(cmd, params, deadline);
}
// Fall through if the prefix matched but the parser disagreed —
// shouldn't happen given the grammar, but the regular path
// below will surface a parse error coherently.
}
// Procedure calls we route ourselves (vector kNN today; full-text
// queries land here too once they ship). These bypass the
// analyzer/compiler, which doesn't yet model standalone CALL.
if super::procedures::is_procedure_call_text(query) {
let document = parse_query(query)?;
if let lora_ast::Statement::Query(lora_ast::Query::StandaloneCall(call)) =
&document.statement
{
return self.execute_procedure_call(call, params);
}
}
// Compile (or fetch from the plan cache) under the read lock. The
// read lock is also what the read-only fast path runs under, so we
// can reuse it without a release/reacquire when the plan turns out
// to be a pure read.
let (store, store_epoch) = self.read_store_with_epoch_deadline(deadline)?;
let compiled = self.compile_query_cached(query, &*store, store_epoch)?;
let shape = classify_stream(&compiled);
if matches!(shape, StreamShape::ReadOnly) {
if let Some(rec) = &self.wal {
ensure_wal_query_can_start(rec)?;
}
if deadline.is_none() && should_collect_read_via_pull(&compiled) {
return collect_compiled(&*store, params, &compiled).map_err(anyhow::Error::from);
}
let executor = lora_executor::Executor::with_deadline(
lora_executor::ExecutionContext {
storage: &*store,
params,
},
deadline,
);
return executor
.execute_compiled_rows_parallel_safe(&compiled)
.map_err(anyhow::Error::from);
}
// Mutating path. Drop the snapshot used for compilation and route
// through the single-writer auto-commit path. The writer mutex is held
// for mutate + WAL commit so live state and WAL order stay identical.
// The method name still says "optimistic" because this is the seam a
// future write-set/CAS implementation can reuse.
drop(store);
debug_assert!(shape.is_mutating());
if std::any::TypeId::of::<S>() == std::any::TypeId::of::<InMemoryGraph>() {
return self.execute_mutating_optimistic(params, deadline, &compiled);
}
// Backends that aren't `InMemoryGraph` don't support the
// recorder install hook, so fall back to the pessimistic path.
let store = self.write_store_deadline(deadline)?;
self.with_logged_write_guard(
store,
WalAbortPolicy::PoisonIfMutated(QUERY_FAILURE_POISON),
|store| {
let mut executor = MutableExecutor::with_deadline(
MutableExecutionContext {
storage: store,
params,
},
deadline,
);
Ok(executor.execute_compiled_rows(&compiled)?)
},
)
}
}