Skip to main content

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