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_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    ///    then applies the WAL's configured single-thread flush policy.
73    /// 4. On Err, `recorder.abort()` clears the pending batch. The
74    ///    engine has no rollback, so the in-memory state may already
75    ///    be partially mutated; the live handle is quarantined while
76    ///    durable recovery stays atomic because no committed batch was
77    ///    written.
78    /// 5. The recorder's poisoned flag is polled once. If set, the query
79    ///    fails loudly with the durability error so the caller can act on it;
80    ///    the WAL refuses further appends until the operator restarts the
81    ///    database, which recovers from the last consistent snapshot + WAL.
82    pub fn execute_with_params(
83        &self,
84        query: &str,
85        options: Option<ExecuteOptions>,
86        params: BTreeMap<String, LoraValue>,
87    ) -> Result<QueryResult, LoraError> {
88        let rows = self.execute_rows_with_params_deadline(query, params, None)?;
89        Ok(project_rows(rows, options.unwrap_or_default()))
90    }
91
92    /// Execute a parameterised query with a cooperative deadline.
93    pub fn execute_with_params_timeout(
94        &self,
95        query: &str,
96        options: Option<ExecuteOptions>,
97        params: BTreeMap<String, LoraValue>,
98        timeout: Duration,
99    ) -> Result<QueryResult, LoraError> {
100        let deadline = Instant::now()
101            .checked_add(timeout)
102            .unwrap_or_else(Instant::now);
103        let rows = self.execute_rows_with_params_deadline(query, params, Some(deadline))?;
104        Ok(project_rows(rows, options.unwrap_or_default()))
105    }
106
107    /// Execute a query and return hydrated rows before final result-format
108    /// projection.
109    pub fn execute_rows(&self, query: &str) -> Result<Vec<Row>, LoraError> {
110        self.execute_rows_with_params(query, BTreeMap::new())
111    }
112
113    /// Execute a query with parameters and return hydrated rows before final
114    /// result-format projection.
115    pub fn execute_rows_with_params(
116        &self,
117        query: &str,
118        params: BTreeMap<String, LoraValue>,
119    ) -> Result<Vec<Row>, LoraError> {
120        Ok(self.execute_rows_with_params_deadline(query, params, None)?)
121    }
122
123    pub(crate) fn execute_rows_with_params_deadline(
124        &self,
125        query: &str,
126        params: BTreeMap<String, LoraValue>,
127        deadline: Option<Instant>,
128    ) -> Result<Vec<Row>> {
129        // Compile (or fetch from the plan cache) under the read lock. The
130        // read lock is also what the read-only fast path runs under, so we
131        // can reuse it without a release/reacquire when the plan turns out
132        // to be a pure read.
133        let store = self.read_store_deadline(deadline)?;
134        let compiled = self.compile_query_cached(query, &*store)?;
135        let shape = classify_stream(&compiled);
136
137        if matches!(shape, StreamShape::ReadOnly) {
138            if let Some(rec) = &self.wal {
139                ensure_wal_query_can_start(rec)?;
140            }
141            if deadline.is_none() && should_collect_read_via_pull(&compiled) {
142                return collect_compiled(&*store, params, &compiled).map_err(anyhow::Error::from);
143            }
144            let executor = lora_executor::Executor::with_deadline(
145                lora_executor::ExecutionContext {
146                    storage: &*store,
147                    params,
148                },
149                deadline,
150            );
151            return executor
152                .execute_compiled_rows(&compiled)
153                .map_err(anyhow::Error::from);
154        }
155
156        // Mutating path. Drop the snapshot used for compilation and route
157        // through the single-writer auto-commit path. The writer mutex is held
158        // for mutate + WAL commit so live state and WAL order stay identical.
159        // The method name still says "optimistic" because this is the seam a
160        // future write-set/CAS implementation can reuse.
161        drop(store);
162        debug_assert!(shape.is_mutating());
163
164        if std::any::TypeId::of::<S>() == std::any::TypeId::of::<InMemoryGraph>() {
165            return self.execute_mutating_optimistic(params, deadline, &compiled);
166        }
167
168        // Backends that aren't `InMemoryGraph` don't support the
169        // recorder install hook, so fall back to the pessimistic path.
170        let store = self.write_store_deadline(deadline)?;
171        self.with_logged_write_guard(
172            store,
173            WalAbortPolicy::PoisonIfMutated(QUERY_FAILURE_POISON),
174            |store| {
175                let mut executor = MutableExecutor::with_deadline(
176                    MutableExecutionContext {
177                        storage: store,
178                        params,
179                    },
180                    deadline,
181                );
182                Ok(executor.execute_compiled_rows(&compiled)?)
183            },
184        )
185    }
186}