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, Instant};
14
15use anyhow::Result;
16use lora_executor::{
17 classify_stream, collect_compiled, project_rows, ExecuteOptions, LoraValue,
18 MutableExecutionContext, MutableExecutor, QueryResult, Row, StreamShape,
19};
20use lora_store::{GraphStorage, GraphStorageMut, InMemoryGraph};
21
22use crate::database::{Database, QUERY_FAILURE_POISON};
23use crate::error::LoraError;
24use crate::wal::write_scope::{ensure_wal_query_can_start, WalAbortPolicy};
25
26use super::pull_mode::should_collect_read_via_pull;
27
28impl<S> Database<S>
29where
30 S: GraphStorage + GraphStorageMut + Any + Clone + Send + Sync + 'static,
31{
32 /// Execute a query and return its result.
33 pub fn execute(
34 &self,
35 query: &str,
36 options: Option<ExecuteOptions>,
37 ) -> Result<QueryResult, LoraError> {
38 self.execute_with_params(query, options, BTreeMap::new())
39 }
40
41 /// Execute a query with a cooperative deadline. The timeout is checked at
42 /// executor operator boundaries and hot scan loops; if it fires, the query
43 /// returns an error and any WAL-backed mutating query is aborted through
44 /// the existing failure path.
45 pub fn execute_with_timeout(
46 &self,
47 query: &str,
48 options: Option<ExecuteOptions>,
49 timeout: Duration,
50 ) -> Result<QueryResult, LoraError> {
51 let deadline = Instant::now()
52 .checked_add(timeout)
53 .unwrap_or_else(Instant::now);
54 let rows =
55 self.execute_rows_with_params_deadline(query, BTreeMap::new(), Some(deadline))?;
56 Ok(project_rows(rows, options.unwrap_or_default()))
57 }
58
59 /// Execute a query with bound parameters.
60 ///
61 /// When a WAL is attached the call is bracketed by a transaction:
62 ///
63 /// 1. `recorder.arm()` after analyze + compile (so a parse /
64 /// semantic / compile error never opens a tx that has to be
65 /// immediately aborted). Arming is *cheap*: no record is
66 /// appended to the WAL yet, so a pure read query that
67 /// completes here pays nothing for the WAL hot path.
68 /// 2. The executor runs; every primitive mutation fires
69 /// `MutationRecorder::record`, which buffers events in memory.
70 /// 3. On Ok, `recorder.commit()` writes `TxBegin`, one batched
71 /// mutation record, and `TxCommit` only when mutations occurred;
72 /// the surrounding `recorder.flush()` runs only in that case so
73 /// a read-only query never pays an `fsync`.
74 /// 4. On Err, `recorder.abort()` clears the pending batch. The
75 /// engine has no rollback, so the in-memory state may already
76 /// be partially mutated; the live handle is quarantined while
77 /// durable recovery stays atomic because no committed batch was
78 /// written.
79 /// 5. The recorder's poisoned flag is polled once (it also
80 /// surfaces background-flusher fsync failures from
81 /// `SyncMode::Group`). If set, the query fails loudly with the
82 /// durability error so the caller can act on it; the WAL
83 /// refuses further appends until the operator restarts the
84 /// database, which recovers from the last consistent
85 /// snapshot + WAL.
86 pub fn execute_with_params(
87 &self,
88 query: &str,
89 options: Option<ExecuteOptions>,
90 params: BTreeMap<String, LoraValue>,
91 ) -> Result<QueryResult, LoraError> {
92 let rows = self.execute_rows_with_params_deadline(query, params, None)?;
93 Ok(project_rows(rows, options.unwrap_or_default()))
94 }
95
96 /// Execute a parameterised query with a cooperative deadline.
97 pub fn execute_with_params_timeout(
98 &self,
99 query: &str,
100 options: Option<ExecuteOptions>,
101 params: BTreeMap<String, LoraValue>,
102 timeout: Duration,
103 ) -> Result<QueryResult, LoraError> {
104 let deadline = Instant::now()
105 .checked_add(timeout)
106 .unwrap_or_else(Instant::now);
107 let rows = self.execute_rows_with_params_deadline(query, params, Some(deadline))?;
108 Ok(project_rows(rows, options.unwrap_or_default()))
109 }
110
111 /// Execute a query and return hydrated rows before final result-format
112 /// projection.
113 pub fn execute_rows(&self, query: &str) -> Result<Vec<Row>, LoraError> {
114 self.execute_rows_with_params(query, BTreeMap::new())
115 }
116
117 /// Execute a query with parameters and return hydrated rows before final
118 /// result-format projection.
119 pub fn execute_rows_with_params(
120 &self,
121 query: &str,
122 params: BTreeMap<String, LoraValue>,
123 ) -> Result<Vec<Row>, LoraError> {
124 Ok(self.execute_rows_with_params_deadline(query, params, None)?)
125 }
126
127 pub(crate) fn execute_rows_with_params_deadline(
128 &self,
129 query: &str,
130 params: BTreeMap<String, LoraValue>,
131 deadline: Option<Instant>,
132 ) -> Result<Vec<Row>> {
133 // Compile (or fetch from the plan cache) under the read lock. The
134 // read lock is also what the read-only fast path runs under, so we
135 // can reuse it without a release/reacquire when the plan turns out
136 // to be a pure read.
137 let store = self.read_store_deadline(deadline)?;
138 let compiled = self.compile_query_cached(query, &*store)?;
139 let shape = classify_stream(&compiled);
140
141 if matches!(shape, StreamShape::ReadOnly) {
142 if let Some(rec) = &self.wal {
143 ensure_wal_query_can_start(rec)?;
144 }
145 if deadline.is_none() && should_collect_read_via_pull(&compiled) {
146 return collect_compiled(&*store, params, &compiled).map_err(anyhow::Error::from);
147 }
148 let executor = lora_executor::Executor::with_deadline(
149 lora_executor::ExecutionContext {
150 storage: &*store,
151 params,
152 },
153 deadline,
154 );
155 return executor
156 .execute_compiled_rows(&compiled)
157 .map_err(anyhow::Error::from);
158 }
159
160 // Mutating path. Drop the snapshot we used for compilation and
161 // route through the optimistic auto-commit path: build the
162 // working copy + mutate without holding any lock, then take
163 // the writer Mutex only briefly to do a CAS publish (replays
164 // buffered mutation events to the durable WAL inside the
165 // critical section). Multiple auto-commit writers can run
166 // their prep work in parallel; they only serialize at the
167 // commit point. On conflict (another writer published since
168 // we took our snapshot), retry from a fresh snapshot.
169 drop(store);
170 debug_assert!(shape.is_mutating());
171
172 if std::any::TypeId::of::<S>() == std::any::TypeId::of::<InMemoryGraph>() {
173 return self.execute_mutating_optimistic(params, deadline, &compiled);
174 }
175
176 // Backends that aren't `InMemoryGraph` don't support the
177 // recorder install hook, so fall back to the pessimistic path.
178 let store = self.write_store_deadline(deadline)?;
179 self.with_logged_write_guard(
180 store,
181 WalAbortPolicy::PoisonIfMutated(QUERY_FAILURE_POISON),
182 |store| {
183 let mut executor = MutableExecutor::with_deadline(
184 MutableExecutionContext {
185 storage: store,
186 params,
187 },
188 deadline,
189 );
190 Ok(executor.execute_compiled_rows(&compiled)?)
191 },
192 )
193 }
194}