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;
14
15use web_time::Instant;
16
17use anyhow::Result;
18use lora_ast::Statement;
19use lora_executor::{
20    classify_stream, collect_compiled, 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 query and return hydrated rows before final result-format
111    /// projection.
112    pub fn execute_rows(&self, query: &str) -> Result<Vec<Row>, LoraError> {
113        self.execute_rows_with_params(query, BTreeMap::new())
114    }
115
116    /// Execute a query with parameters and return hydrated rows before final
117    /// result-format projection.
118    pub fn execute_rows_with_params(
119        &self,
120        query: &str,
121        params: BTreeMap<String, LoraValue>,
122    ) -> Result<Vec<Row>, LoraError> {
123        Ok(self.execute_rows_with_params_deadline(query, params, None)?)
124    }
125
126    pub(crate) fn execute_rows_with_params_deadline(
127        &self,
128        query: &str,
129        params: BTreeMap<String, LoraValue>,
130        deadline: Option<Instant>,
131    ) -> Result<Vec<Row>> {
132        // Schema commands (CREATE INDEX, SHOW INDEXES) are catalog-level
133        // mutations and don't share the row-producing pipeline. A cheap
134        // textual prefix check routes them out before the parser cache and
135        // analyzer get involved.
136        if Self::is_schema_command_text(query) {
137            let document = parse_query(query)?;
138            if let Statement::Schema(cmd) = &document.statement {
139                return self.execute_schema_command(cmd, params, deadline);
140            }
141            // Fall through if the prefix matched but the parser disagreed —
142            // shouldn't happen given the grammar, but the regular path
143            // below will surface a parse error coherently.
144        }
145
146        // Procedure calls we route ourselves (vector kNN today; full-text
147        // queries land here too once they ship). These bypass the
148        // analyzer/compiler, which doesn't yet model standalone CALL.
149        if super::procedures::is_procedure_call_text(query) {
150            let document = parse_query(query)?;
151            if let lora_ast::Statement::Query(lora_ast::Query::StandaloneCall(call)) =
152                &document.statement
153            {
154                return self.execute_procedure_call(call, params);
155            }
156        }
157
158        // Compile (or fetch from the plan cache) under the read lock. The
159        // read lock is also what the read-only fast path runs under, so we
160        // can reuse it without a release/reacquire when the plan turns out
161        // to be a pure read.
162        let (store, store_epoch) = self.read_store_with_epoch_deadline(deadline)?;
163        let compiled = self.compile_query_cached(query, &*store, store_epoch)?;
164        let shape = classify_stream(&compiled);
165
166        if matches!(shape, StreamShape::ReadOnly) {
167            if let Some(rec) = &self.wal {
168                ensure_wal_query_can_start(rec)?;
169            }
170            if deadline.is_none() && should_collect_read_via_pull(&compiled) {
171                return collect_compiled(&*store, params, &compiled).map_err(anyhow::Error::from);
172            }
173            let executor = lora_executor::Executor::with_deadline(
174                lora_executor::ExecutionContext {
175                    storage: &*store,
176                    params,
177                },
178                deadline,
179            );
180            return executor
181                .execute_compiled_rows_parallel_safe(&compiled)
182                .map_err(anyhow::Error::from);
183        }
184
185        // Mutating path. Drop the snapshot used for compilation and route
186        // through the single-writer auto-commit path. The writer mutex is held
187        // for mutate + WAL commit so live state and WAL order stay identical.
188        // The method name still says "optimistic" because this is the seam a
189        // future write-set/CAS implementation can reuse.
190        drop(store);
191        debug_assert!(shape.is_mutating());
192
193        if std::any::TypeId::of::<S>() == std::any::TypeId::of::<InMemoryGraph>() {
194            return self.execute_mutating_optimistic(params, deadline, &compiled);
195        }
196
197        // Backends that aren't `InMemoryGraph` don't support the
198        // recorder install hook, so fall back to the pessimistic path.
199        let store = self.write_store_deadline(deadline)?;
200        self.with_logged_write_guard(
201            store,
202            WalAbortPolicy::PoisonIfMutated(QUERY_FAILURE_POISON),
203            |store| {
204                let mut executor = MutableExecutor::with_deadline(
205                    MutableExecutionContext {
206                        storage: store,
207                        params,
208                    },
209                    deadline,
210                );
211                Ok(executor.execute_compiled_rows(&compiled)?)
212            },
213        )
214    }
215}