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