Skip to main content

lora_database/database/
mod.rs

1use std::any::Any;
2use std::collections::BTreeMap;
3use std::sync::{Arc, Mutex};
4
5use anyhow::{anyhow, Result};
6use lora_ast::{Direction, Document};
7use lora_executor::{lora_value_to_property, ExecuteOptions, LoraValue, QueryResult};
8use lora_parser::parse_query;
9use lora_store::{GraphStorage, GraphStorageMut, InMemoryGraph, Properties};
10use lora_wal::WalRecorder;
11
12mod builder;
13mod compile;
14mod execute;
15mod explain;
16mod graph_api;
17mod occ;
18mod profile;
19mod pull_mode;
20mod replay;
21mod stream;
22mod write_guard;
23
24use crate::error::LoraError;
25use crate::explain::{QueryPlan, QueryProfile};
26use crate::live_store::LiveStore;
27use crate::plan_cache::PlanCache;
28use crate::snapshot::ManagedSnapshotStore;
29use crate::wal::write_scope::WalAbortPolicy;
30
31/// Minimal abstraction any transport can depend on to run Lora queries.
32///
33/// `execute` runs a query and returns rows. `explain` and `profile` are
34/// deliberately separate methods: `explain` never invokes the executor
35/// (so it can be called on mutating queries without side effects) and
36/// `profile` runs the executor and reports runtime metrics. Transports
37/// MUST NOT route plan / profile requests through `execute` — exposing
38/// the plan-only and profile-with-metrics behaviours as separate
39/// methods is part of the public contract.
40pub trait QueryRunner: Send + Sync + 'static {
41    fn execute(
42        &self,
43        query: &str,
44        options: Option<ExecuteOptions>,
45    ) -> Result<QueryResult, LoraError>;
46
47    /// Compile a query and return its plan without executing it.
48    fn explain(
49        &self,
50        query: &str,
51        params: Option<BTreeMap<String, LoraValue>>,
52    ) -> Result<QueryPlan, LoraError>;
53
54    /// Execute a query and return its plan plus runtime metrics.
55    /// Mutating queries are persisted exactly as in `execute`.
56    fn profile(
57        &self,
58        query: &str,
59        params: Option<BTreeMap<String, LoraValue>>,
60    ) -> Result<QueryProfile, LoraError>;
61}
62
63/// Owns the graph store and orchestrates parse → analyze → compile → execute.
64///
65/// Optionally drives a write-ahead log: when constructed via
66/// [`Database::open_with_wal`] or [`Database::recover`] the database
67/// holds an [`Arc<WalRecorder>`] that brackets every query with
68/// `begin → mutations → commit/abort → flush` while the store write
69/// lock is held, so the WAL order is exactly the in-memory commit order.
70/// When constructed via [`Database::in_memory`] / [`Database::from_graph`]
71/// the WAL handle is `None` and the engine pays only the existing
72/// `MutationRecorder::record` null-pointer check per mutation.
73pub struct Database<S> {
74    /// The current authoritative store. Reads call `store.load_full()`
75    /// to obtain an independent `Arc<S>` snapshot; writers take the
76    /// `writer` Mutex and then a brief write-lock on the inner
77    /// `RwLock<Arc<S>>`, mutating in-place via `Arc::make_mut`. When
78    /// no in-flight reader holds a snapshot Arc, `make_mut` returns
79    /// `&mut S` without cloning the graph — that's the single-writer
80    /// fast path that "CREATE one node" / "SET property" depend on
81    /// for graph-size-independent throughput. When concurrent readers
82    /// are alive, `make_mut` clones once and the readers keep
83    /// observing the pre-mutation state via their old Arc.
84    pub(crate) store: Arc<LiveStore<S>>,
85    /// Serializes commit ordering. Held across `mutate-WAL-publish` so
86    /// WAL records are appended in the same order live state advances.
87    /// Readers never touch this Mutex; only writers contend.
88    pub(crate) writer: Arc<Mutex<()>>,
89    /// Per-record write locks. Plumbed for a future phase that allows
90    /// concurrent commits across disjoint write sets; today the writer
91    /// Mutex provides single-writer-at-a-time semantics so this table
92    /// is idle.
93    #[allow(dead_code)]
94    pub(crate) lock_table: Arc<lora_store::LockTable>,
95    pub(crate) wal: Option<Arc<WalRecorder>>,
96    pub(crate) snapshots: Option<Arc<ManagedSnapshotStore>>,
97    /// Cache of compiled query plans, content-keyed by raw query text. Shared
98    /// across the read- and write-lock phases of a single execute (so a
99    /// mutating query compiles at most once instead of twice) and across
100    /// every subsequent call that uses the same query string.
101    pub(crate) plan_cache: Arc<PlanCache>,
102}
103
104#[derive(Debug, Clone, Copy, PartialEq, Eq)]
105pub enum GraphDirection {
106    Outgoing,
107    Incoming,
108    Both,
109}
110
111impl GraphDirection {
112    pub(crate) fn as_store_direction(self) -> Direction {
113        match self {
114            Self::Outgoing => Direction::Right,
115            Self::Incoming => Direction::Left,
116            Self::Both => Direction::Undirected,
117        }
118    }
119}
120
121pub(crate) fn values_to_properties(values: BTreeMap<String, LoraValue>) -> Result<Properties> {
122    values
123        .into_iter()
124        .map(|(key, value)| {
125            let value = lora_value_to_property(value).map_err(|e| anyhow!(e))?;
126            Ok((key, value))
127        })
128        .collect()
129}
130
131pub(crate) const QUERY_FAILURE_POISON: &str =
132    "query mutated the live graph before failing; restart from snapshot + WAL required";
133
134impl Database<InMemoryGraph> {
135    /// Force any pending WAL bytes to durable storage and, for archive-backed
136    /// databases, refresh the portable `.loradb` file before returning.
137    ///
138    /// Managed snapshot checkpoints are explicit via
139    /// [`Self::checkpoint_managed`] or threshold-driven via
140    /// [`SnapshotConfig::checkpoint_every_commits`]; `sync()` remains a
141    /// durability operation rather than an O(graph) checkpoint.
142    pub fn sync(&self) -> Result<(), LoraError> {
143        if let Some(wal) = &self.wal {
144            wal.force_fsync()?;
145        }
146        Ok(())
147    }
148}
149
150impl<S> Database<S>
151where
152    S: GraphStorage + GraphStorageMut + Any + Clone + Send + Sync + 'static,
153{
154    /// Handle to the installed WAL recorder, if any. Exposed for
155    /// admin paths (checkpoint, truncate, observability) that need
156    /// to drive the WAL outside the standard query lifecycle.
157    pub fn wal(&self) -> Option<&Arc<WalRecorder>> {
158        self.wal.as_ref()
159    }
160
161    /// Snapshot the current authoritative graph. Equivalent to the
162    /// historical `database.store().load_full()` pattern, exposed here
163    /// so external callers don't need to name the internal storage
164    /// wrapper.
165    pub fn snapshot(&self) -> Arc<S> {
166        self.store.load_full()
167    }
168
169    /// Parse a query string into an AST without executing it.
170    pub fn parse(&self, query: &str) -> Result<Document, LoraError> {
171        Ok(parse_query(query)?)
172    }
173
174    /// Read the current authoritative snapshot. Lock-free: returns an
175    /// `Arc<S>` whose lifetime is independent of any writer.
176    pub(crate) fn read_store(&self) -> Arc<S> {
177        self.store.load_full()
178    }
179}
180
181impl<S> Database<S>
182where
183    S: GraphStorage + GraphStorageMut + Any + Clone + Send + Sync + 'static,
184{
185    // ---------- Storage-agnostic utility helpers ----------
186    //
187    // Bindings previously reached into the shared store lock to answer
188    // stat / admin calls; these helpers let them depend on `Database<S>`
189    // instead, so swapping in a new backend only requires changing one type
190    // parameter.
191
192    /// Drop every node and relationship, returning WAL/archive errors to the
193    /// caller.
194    ///
195    /// When a WAL is attached, the clear is wrapped in `arm`/`commit` so the
196    /// `MutationEvent::Clear` fired by the store reaches the log inside a
197    /// transaction. If a failure happens after the in-memory graph has been
198    /// cleared, the recorder is poisoned by the failing WAL path and future
199    /// writes fail until the database is reopened from durable state.
200    pub fn try_clear(&self) -> Result<(), LoraError> {
201        let guard = self.write_store();
202        self.with_logged_write_guard(guard, WalAbortPolicy::AbortOnly, |store| {
203            store.clear();
204            Ok(())
205        })
206        .map_err(LoraError::from_anyhow)
207    }
208
209    /// Drop every node and relationship.
210    ///
211    /// This compatibility helper keeps the historical infallible Rust API.
212    /// Bindings that can report errors should call [`Self::try_clear`].
213    pub fn clear(&self) {
214        let _ = self.try_clear();
215    }
216
217    /// Number of nodes currently in the graph.
218    pub fn node_count(&self) -> usize {
219        let snapshot = self.read_store();
220        snapshot.node_count()
221    }
222
223    /// Number of relationships currently in the graph.
224    pub fn relationship_count(&self) -> usize {
225        let snapshot = self.read_store();
226        snapshot.relationship_count()
227    }
228
229    /// Run a closure with a shared borrow of the underlying store.
230    /// Lock-free: callers see a consistent snapshot for the duration of
231    /// the closure even while writers commit new versions.
232    pub fn with_store<R>(&self, f: impl FnOnce(&S) -> R) -> R {
233        let snapshot = self.read_store();
234        f(&*snapshot)
235    }
236
237    /// Run a closure with an exclusive borrow of the underlying store. Reserved
238    /// for admin paths (restore, bulk load); regular mutation goes through
239    /// `execute_with_params`. The closure mutates a staged copy that is
240    /// published atomically when the closure returns.
241    pub fn with_store_mut<R>(&self, f: impl FnOnce(&mut S) -> R) -> R {
242        let mut guard = self.write_store();
243        let result = f(&mut *guard);
244        guard.publish();
245        result
246    }
247}
248
249impl<S> QueryRunner for Database<S>
250where
251    S: GraphStorage + GraphStorageMut + Any + Clone + Send + Sync + 'static,
252{
253    fn execute(
254        &self,
255        query: &str,
256        options: Option<ExecuteOptions>,
257    ) -> Result<QueryResult, LoraError> {
258        Database::execute(self, query, options)
259    }
260
261    fn explain(
262        &self,
263        query: &str,
264        params: Option<BTreeMap<String, LoraValue>>,
265    ) -> Result<QueryPlan, LoraError> {
266        Database::explain(self, query, params)
267    }
268
269    fn profile(
270        &self,
271        query: &str,
272        params: Option<BTreeMap<String, LoraValue>>,
273    ) -> Result<QueryProfile, LoraError> {
274        Database::profile(self, query, params)
275    }
276}