lora_database/database/execute.rs
1//! Query-execute family on [`Database<S>`].
2//!
3//! These methods cover the `execute_*` and `execute_rows_*` entry points
4//! plus the cooperative-deadline variants and the plan-cache-backed
5//! compile path. The actual write routing (optimistic OCC vs. pessimistic
6//! `with_logged_write_guard`) is dispatched from
7//! [`Database::execute_rows_with_params_deadline`]; the OCC body lives
8//! in [`super::occ`] and the WAL-bracketed write closure lives in
9//! [`super::write_guard`].
10
11use std::any::Any;
12use std::collections::BTreeMap;
13use std::time::Duration;
14
15use web_time::Instant;
16
17use anyhow::Result;
18use lora_ast::Statement;
19use lora_executor::{
20 classify_stream, collect_compiled_with_deadline, project_rows, ExecuteOptions, LoraValue,
21 MutableExecutionContext, MutableExecutor, QueryResult, Row, StreamShape,
22};
23use lora_parser::parse_query;
24use lora_store::{GraphStorage, GraphStorageMut, InMemoryGraph};
25
26use crate::database::{Database, QUERY_FAILURE_POISON};
27use crate::error::LoraError;
28use crate::wal::write_scope::{ensure_wal_query_can_start, WalAbortPolicy};
29
30use super::pull_mode::should_collect_read_via_pull;
31
32impl<S> Database<S>
33where
34 S: GraphStorage + GraphStorageMut + Any + Clone + Send + Sync + 'static,
35{
36 /// Execute a query and return its result.
37 pub fn execute(
38 &self,
39 query: &str,
40 options: Option<ExecuteOptions>,
41 ) -> Result<QueryResult, LoraError> {
42 self.execute_with_params(query, options, BTreeMap::new())
43 }
44
45 /// Execute a query with a cooperative deadline. The timeout is checked at
46 /// executor operator boundaries and hot scan loops; if it fires, the query
47 /// returns an error and any WAL-backed mutating query is aborted through
48 /// the existing failure path.
49 pub fn execute_with_timeout(
50 &self,
51 query: &str,
52 options: Option<ExecuteOptions>,
53 timeout: Duration,
54 ) -> Result<QueryResult, LoraError> {
55 let deadline = Instant::now()
56 .checked_add(timeout)
57 .unwrap_or_else(Instant::now);
58 let rows =
59 self.execute_rows_with_params_deadline(query, BTreeMap::new(), Some(deadline))?;
60 Ok(project_rows(rows, options.unwrap_or_default()))
61 }
62
63 /// Execute a query with bound parameters.
64 ///
65 /// When a WAL is attached the call is bracketed by a transaction:
66 ///
67 /// 1. `recorder.arm()` after analyze + compile (so a parse /
68 /// semantic / compile error never opens a tx that has to be
69 /// immediately aborted). Arming is *cheap*: no record is
70 /// appended to the WAL yet, so a pure read query that
71 /// completes here pays nothing for the WAL hot path.
72 /// 2. The executor runs; every primitive mutation fires
73 /// `MutationRecorder::record`, which buffers events in memory.
74 /// 3. On Ok, `recorder.commit()` writes `TxBegin`, one batched
75 /// mutation record, and `TxCommit` only when mutations occurred,
76 /// then applies the WAL's configured single-thread flush policy.
77 /// 4. On Err, `recorder.abort()` clears the pending batch. The
78 /// staged graph is discarded without publishing, so the live
79 /// in-memory state and durable WAL both remain at the last
80 /// successful commit.
81 /// 5. The recorder's poisoned flag is polled once. If set, the query
82 /// fails loudly with the durability error so the caller can act on it;
83 /// the WAL refuses further appends until the operator restarts the
84 /// database, which recovers from the last consistent snapshot + WAL.
85 pub fn execute_with_params(
86 &self,
87 query: &str,
88 options: Option<ExecuteOptions>,
89 params: BTreeMap<String, LoraValue>,
90 ) -> Result<QueryResult, LoraError> {
91 let rows = self.execute_rows_with_params_deadline(query, params, None)?;
92 Ok(project_rows(rows, options.unwrap_or_default()))
93 }
94
95 /// Execute a parameterised query with a cooperative deadline.
96 pub fn execute_with_params_timeout(
97 &self,
98 query: &str,
99 options: Option<ExecuteOptions>,
100 params: BTreeMap<String, LoraValue>,
101 timeout: Duration,
102 ) -> Result<QueryResult, LoraError> {
103 let deadline = Instant::now()
104 .checked_add(timeout)
105 .unwrap_or_else(Instant::now);
106 let rows = self.execute_rows_with_params_deadline(query, params, Some(deadline))?;
107 Ok(project_rows(rows, options.unwrap_or_default()))
108 }
109
110 /// Execute a parameterised query that stops at an absolute `deadline`.
111 ///
112 /// Pair with [`lora_executor::CancellableDeadline`] to also cancel the
113 /// query from another thread before the deadline: a cancelled query
114 /// fails with a timeout error, aborts its WAL transaction, and releases
115 /// its locks, exactly like one that ran out of time.
116 pub fn execute_with_params_deadline(
117 &self,
118 query: &str,
119 options: Option<ExecuteOptions>,
120 params: BTreeMap<String, LoraValue>,
121 deadline: Instant,
122 ) -> Result<QueryResult, LoraError> {
123 let rows = self.execute_rows_with_params_deadline(query, params, Some(deadline))?;
124 Ok(project_rows(rows, options.unwrap_or_default()))
125 }
126
127 /// Execute a query and return hydrated rows before final result-format
128 /// projection.
129 pub fn execute_rows(&self, query: &str) -> Result<Vec<Row>, LoraError> {
130 self.execute_rows_with_params(query, BTreeMap::new())
131 }
132
133 /// Execute a query with parameters and return hydrated rows before final
134 /// result-format projection.
135 pub fn execute_rows_with_params(
136 &self,
137 query: &str,
138 params: BTreeMap<String, LoraValue>,
139 ) -> Result<Vec<Row>, LoraError> {
140 Ok(self.execute_rows_with_params_deadline(query, params, None)?)
141 }
142
143 pub(crate) fn execute_rows_with_params_deadline(
144 &self,
145 query: &str,
146 params: BTreeMap<String, LoraValue>,
147 deadline: Option<Instant>,
148 ) -> Result<Vec<Row>> {
149 // Schema commands (CREATE INDEX, SHOW INDEXES) are catalog-level
150 // mutations and don't share the row-producing pipeline. A cheap
151 // textual prefix check routes them out before the parser cache and
152 // analyzer get involved.
153 if Self::is_schema_command_text(query) {
154 let document = parse_query(query)?;
155 if let Statement::Schema(cmd) = &document.statement {
156 return self.execute_schema_command(cmd, params, deadline);
157 }
158 // Fall through if the prefix matched but the parser disagreed —
159 // shouldn't happen given the grammar, but the regular path
160 // below will surface a parse error coherently.
161 }
162
163 // Procedure calls we route ourselves (vector kNN today; full-text
164 // queries land here too once they ship). These bypass the
165 // analyzer/compiler, which doesn't yet model standalone CALL.
166 if super::procedures::is_procedure_call_text(query) {
167 // `CALL ... YIELD ... RETURN` also starts with a procedure name
168 // but runs through the regular pipeline; reuse the cached parse
169 // so telling the two apart costs nothing after the first call.
170 let document = match self.plan_cache.document(query) {
171 Some(document) => document,
172 None => std::sync::Arc::new(parse_query(query)?),
173 };
174 if let lora_ast::Statement::Query(lora_ast::Query::StandaloneCall(call)) =
175 &document.statement
176 {
177 return self.execute_procedure_call(call, params);
178 }
179 }
180
181 // Compile (or fetch from the plan cache) under the read lock. The
182 // read lock is also what the read-only fast path runs under, so we
183 // can reuse it without a release/reacquire when the plan turns out
184 // to be a pure read.
185 let (store, store_epoch) = self.read_store_with_epoch_deadline(deadline)?;
186 let compiled = self.compile_query_cached(query, &*store, store_epoch)?;
187 let shape = classify_stream(&compiled);
188
189 if matches!(shape, StreamShape::ReadOnly) {
190 if let Some(rec) = &self.wal {
191 ensure_wal_query_can_start(rec)?;
192 }
193 if should_collect_read_via_pull(&compiled) {
194 // The pull pipeline checks the deadline as rows are pulled,
195 // so a bounded read still stops early at its LIMIT.
196 return collect_compiled_with_deadline(&*store, params, &compiled, deadline)
197 .map_err(anyhow::Error::from);
198 }
199 let executor = lora_executor::Executor::with_deadline(
200 lora_executor::ExecutionContext {
201 storage: &*store,
202 params,
203 },
204 deadline,
205 );
206 return executor
207 .execute_compiled_rows_parallel_safe(&compiled)
208 .map_err(anyhow::Error::from);
209 }
210
211 // Mutating path. Drop the snapshot used for compilation and route
212 // through the single-writer auto-commit path. The writer mutex is held
213 // for mutate + WAL commit so live state and WAL order stay identical.
214 // The method name still says "optimistic" because this is the seam a
215 // future write-set/CAS implementation can reuse.
216 drop(store);
217 debug_assert!(shape.is_mutating());
218
219 if std::any::TypeId::of::<S>() == std::any::TypeId::of::<InMemoryGraph>() {
220 return self.execute_mutating_optimistic(params, deadline, &compiled);
221 }
222
223 // Backends that aren't `InMemoryGraph` don't support the
224 // recorder install hook, so fall back to the pessimistic path.
225 let store = self.write_store_deadline(deadline)?;
226 self.with_logged_write_guard(
227 store,
228 WalAbortPolicy::PoisonIfMutated(QUERY_FAILURE_POISON),
229 |store| {
230 let mut executor = MutableExecutor::with_deadline(
231 MutableExecutionContext {
232 storage: store,
233 params,
234 },
235 deadline,
236 );
237 Ok(executor.execute_compiled_rows(&compiled)?)
238 },
239 )
240 }
241}