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