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