lora_database/database/mod.rs
1use std::any::Any;
2use std::collections::BTreeMap;
3use std::sync::{Arc, Mutex};
4
5use anyhow::{Error, 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 changes;
14mod compile;
15mod execute;
16mod explain;
17mod graph_api;
18mod occ;
19mod procedures;
20mod profile;
21pub(crate) mod pull_mode;
22pub(crate) mod replay;
23mod row_projection;
24pub(crate) mod schema;
25mod show_pipeline;
26mod stream;
27mod write_guard;
28
29use crate::error::LoraError;
30use crate::explain::{QueryPlan, QueryProfile};
31use crate::live_store::LiveStore;
32use crate::plan_cache::PlanCache;
33use crate::snapshot::ManagedSnapshotStore;
34use crate::wal::archive::WalArchive;
35
36/// Minimal abstraction any transport can depend on to run Lora queries.
37///
38/// `execute` runs a query and returns rows. `explain` and `profile` are
39/// deliberately separate methods: `explain` never invokes the executor
40/// (so it can be called on mutating queries without side effects) and
41/// `profile` runs the executor and reports runtime metrics. Transports
42/// MUST NOT route plan / profile requests through `execute` — exposing
43/// the plan-only and profile-with-metrics behaviours as separate
44/// methods is part of the public contract.
45pub trait QueryRunner: Send + Sync + 'static {
46 fn execute(
47 &self,
48 query: &str,
49 options: Option<ExecuteOptions>,
50 ) -> Result<QueryResult, LoraError>;
51
52 /// Execute a query with bound parameters. Implementors that cannot bind
53 /// parameters should reject non-empty maps instead of silently ignoring
54 /// them.
55 fn execute_with_params(
56 &self,
57 query: &str,
58 options: Option<ExecuteOptions>,
59 params: BTreeMap<String, LoraValue>,
60 ) -> Result<QueryResult, LoraError> {
61 if params.is_empty() {
62 return self.execute(query, options);
63 }
64 Err(LoraError::new(
65 crate::error::LoraErrorCode::InvalidParams,
66 "query parameters are not supported by this database runner",
67 ))
68 }
69
70 /// Compile a query and return its plan without executing it.
71 fn explain(
72 &self,
73 query: &str,
74 params: Option<BTreeMap<String, LoraValue>>,
75 ) -> Result<QueryPlan, LoraError>;
76
77 /// Execute a query and return its plan plus runtime metrics.
78 /// Mutating queries are persisted exactly as in `execute`.
79 fn profile(
80 &self,
81 query: &str,
82 params: Option<BTreeMap<String, LoraValue>>,
83 ) -> Result<QueryProfile, LoraError>;
84}
85
86/// Owns the graph store and orchestrates parse → analyze → compile → execute.
87///
88/// Optionally drives a write-ahead log: when constructed via
89/// [`Database::open_with_wal`] or [`Database::recover`] the database
90/// holds an [`Arc<WalRecorder>`] that brackets every query with
91/// `begin → mutations → commit/abort → flush` while the store write
92/// lock is held, so the WAL order is exactly the in-memory commit order.
93/// When constructed via [`Database::in_memory`] / [`Database::from_graph`]
94/// the WAL handle is `None` and the engine pays only the existing
95/// `MutationRecorder::record` null-pointer check per mutation.
96pub struct Database<S> {
97 /// The current authoritative store. Reads call `store.load_full()`
98 /// to obtain an independent `Arc<S>` snapshot; writers take the
99 /// `writer` Mutex and then a brief write-lock on the inner
100 /// `RwLock<Arc<S>>`, mutating in-place via `Arc::make_mut`. When
101 /// no in-flight reader holds a snapshot Arc, `make_mut` returns
102 /// `&mut S` without cloning the graph — that's the single-writer
103 /// fast path that "CREATE one node" / "SET property" depend on
104 /// for graph-size-independent throughput. When concurrent readers
105 /// are alive, `make_mut` clones once and the readers keep
106 /// observing the pre-mutation state via their old Arc.
107 pub(crate) store: Arc<LiveStore<S>>,
108 /// Serializes commit ordering. Held across `mutate-WAL-publish` so
109 /// WAL records are appended in the same order live state advances.
110 /// Readers never touch this Mutex; only writers contend.
111 pub(crate) writer: Arc<Mutex<()>>,
112 /// Per-record write locks. Plumbed for a future phase that allows
113 /// concurrent commits across disjoint write sets; today the writer
114 /// Mutex provides single-writer-at-a-time semantics so this table
115 /// is idle.
116 #[allow(dead_code)]
117 pub(crate) lock_table: Arc<lora_store::LockTable>,
118 pub(crate) wal: Option<Arc<WalRecorder>>,
119 pub(crate) snapshots: Option<Arc<ManagedSnapshotStore>>,
120 /// Present for named `.loradb` databases. Runtime state is still the
121 /// in-memory graph; this handle lets `sync()` refresh the portable
122 /// checkpointed container with a base snapshot frame plus WAL delta frames.
123 pub(crate) named_archive: Option<Arc<WalArchive>>,
124 /// Cache of compiled query plans, content-keyed by raw query text. Shared
125 /// across the read- and write-lock phases of a single execute (so a
126 /// mutating query compiles at most once instead of twice) and across
127 /// every subsequent call that uses the same query string.
128 pub(crate) plan_cache: Arc<PlanCache>,
129 /// Committed-change fan-out for [`Database::changes`]. Idle (one atomic
130 /// load per write) until the first feed opens.
131 pub(crate) changes: Arc<crate::changes::ChangeHub>,
132}
133
134impl<S> Drop for Database<S> {
135 fn drop(&mut self) {
136 // End every open change feed so consumers see the feed close
137 // instead of waiting forever.
138 self.changes.close();
139 }
140}
141
142#[derive(Debug, Clone, Copy, PartialEq, Eq)]
143pub enum GraphDirection {
144 Outgoing,
145 Incoming,
146 Both,
147}
148
149impl GraphDirection {
150 pub(crate) fn as_store_direction(self) -> Direction {
151 match self {
152 Self::Outgoing => Direction::Right,
153 Self::Incoming => Direction::Left,
154 Self::Both => Direction::Undirected,
155 }
156 }
157}
158
159pub(crate) fn values_to_properties(values: BTreeMap<String, LoraValue>) -> Result<Properties> {
160 values
161 .into_iter()
162 .map(|(key, value)| {
163 let value = lora_value_to_property(value).map_err(Error::from)?;
164 Ok((lora_store::intern_owned(key), value))
165 })
166 .collect()
167}
168
169pub(crate) const QUERY_FAILURE_POISON: &str =
170 "query mutated the live graph before failing; restart from snapshot + WAL required";
171
172impl Database<InMemoryGraph> {
173 /// Force any pending WAL bytes to durable storage and, for container-backed
174 /// named databases, refresh the portable `.loradb` file before returning.
175 ///
176 /// Managed snapshot checkpoints are explicit via
177 /// [`Self::checkpoint_managed`] or threshold-driven via
178 /// [`SnapshotConfig::checkpoint_every_commits`]; `sync()` remains a
179 /// durability operation rather than an O(graph) checkpoint.
180 pub fn sync(&self) -> Result<(), LoraError> {
181 if let Some(wal) = &self.wal {
182 if let Some(archive) = &self.named_archive {
183 let _commit_lock = self
184 .writer
185 .lock()
186 .unwrap_or_else(|poisoned| poisoned.into_inner());
187 wal.force_fsync_wal_only()?;
188 let snapshot_lsn = wal.wal().durable_lsn();
189 let graph = self.store.load_full();
190 let payload = graph.snapshot_payload();
191 let mut bytes = Vec::new();
192 let options = lora_snapshot::SnapshotOptions::default();
193 crate::snapshot::encode_snapshot_to(
194 &mut bytes,
195 &payload,
196 Some(snapshot_lsn.raw()),
197 &options,
198 )?;
199 archive.persist_snapshot_bytes(bytes)?;
200 } else {
201 wal.force_fsync()?;
202 }
203 }
204 Ok(())
205 }
206}
207
208impl<S> Database<S>
209where
210 S: GraphStorage + GraphStorageMut + Any + Clone + Send + Sync + 'static,
211{
212 /// Handle to the installed WAL recorder, if any. Exposed for
213 /// admin paths (checkpoint, truncate, observability) that need
214 /// to drive the WAL outside the standard query lifecycle.
215 pub fn wal(&self) -> Option<&Arc<WalRecorder>> {
216 self.wal.as_ref()
217 }
218
219 /// Snapshot the current authoritative graph. Equivalent to the
220 /// historical `database.store().load_full()` pattern, exposed here
221 /// so external callers don't need to name the internal storage
222 /// wrapper.
223 pub fn snapshot(&self) -> Arc<S> {
224 self.store.load_full()
225 }
226
227 /// Parse a query string into an AST without executing it.
228 pub fn parse(&self, query: &str) -> Result<Document, LoraError> {
229 Ok(parse_query(query)?)
230 }
231
232 /// Read the current authoritative snapshot. Lock-free: returns an
233 /// `Arc<S>` whose lifetime is independent of any writer.
234 pub(crate) fn read_store(&self) -> Arc<S> {
235 self.store.load_full()
236 }
237}
238
239impl<S> Database<S>
240where
241 S: GraphStorage + GraphStorageMut + Any + Clone + Send + Sync + 'static,
242{
243 // ---------- Storage-agnostic utility helpers ----------
244 //
245 // Bindings previously reached into the shared store lock to answer
246 // stat / admin calls; these helpers let them depend on `Database<S>`
247 // instead, so swapping in a new backend only requires changing one type
248 // parameter.
249
250 /// Drop every node and relationship, returning WAL/archive errors to the
251 /// caller.
252 ///
253 /// When a WAL is attached, the buffered `MutationEvent::Clear` is appended
254 /// to the log on success. The clear runs in place against the live
255 /// graph via the same fast path as `with_logged_store_mut`, so a
256 /// large graph clear no longer pays an O(N+E) snapshot clone.
257 pub fn try_clear(&self) -> Result<(), LoraError> {
258 self.with_logged_store_mut(|store| {
259 store.clear();
260 Ok(())
261 })
262 .map_err(LoraError::from_anyhow)
263 }
264
265 /// Drop every node and relationship.
266 ///
267 /// This compatibility helper keeps the historical infallible Rust API.
268 /// Bindings that can report errors should call [`Self::try_clear`].
269 pub fn clear(&self) {
270 let _ = self.try_clear();
271 }
272
273 /// Number of nodes currently in the graph.
274 pub fn node_count(&self) -> usize {
275 let snapshot = self.read_store();
276 snapshot.node_count()
277 }
278
279 /// Number of relationships currently in the graph.
280 pub fn relationship_count(&self) -> usize {
281 let snapshot = self.read_store();
282 snapshot.relationship_count()
283 }
284
285 /// Run a closure with a shared borrow of the underlying store.
286 /// Lock-free: callers see a consistent snapshot for the duration of
287 /// the closure even while writers commit new versions.
288 pub fn with_store<R>(&self, f: impl FnOnce(&S) -> R) -> R {
289 let snapshot = self.read_store();
290 f(&*snapshot)
291 }
292
293 /// Run a closure with an exclusive borrow of a staged store, then publish
294 /// it atomically after the closure returns. Reserved for admin paths;
295 /// regular mutation goes through `execute_with_params`.
296 ///
297 /// If the closure panics, the staged copy is dropped and the live store is
298 /// left untouched.
299 pub fn with_store_mut<R>(&self, f: impl FnOnce(&mut S) -> R) -> R {
300 let _writer_lock = self
301 .writer
302 .lock()
303 .unwrap_or_else(|poisoned| poisoned.into_inner());
304 let snapshot = self.store.load_full();
305 let mut staged: S = (*snapshot).clone();
306 let result = f(&mut staged);
307 self.store.store(Arc::new(staged));
308 // The closure bypasses the recorder, so a change feed can only be
309 // told that the graph was replaced.
310 self.publish_reset();
311 result
312 }
313}
314
315impl<S> QueryRunner for Database<S>
316where
317 S: GraphStorage + GraphStorageMut + Any + Clone + Send + Sync + 'static,
318{
319 fn execute(
320 &self,
321 query: &str,
322 options: Option<ExecuteOptions>,
323 ) -> Result<QueryResult, LoraError> {
324 Database::execute(self, query, options)
325 }
326
327 fn execute_with_params(
328 &self,
329 query: &str,
330 options: Option<ExecuteOptions>,
331 params: BTreeMap<String, LoraValue>,
332 ) -> Result<QueryResult, LoraError> {
333 Database::execute_with_params(self, query, options, params)
334 }
335
336 fn explain(
337 &self,
338 query: &str,
339 params: Option<BTreeMap<String, LoraValue>>,
340 ) -> Result<QueryPlan, LoraError> {
341 Database::explain(self, query, params)
342 }
343
344 fn profile(
345 &self,
346 query: &str,
347 params: Option<BTreeMap<String, LoraValue>>,
348 ) -> Result<QueryProfile, LoraError> {
349 Database::profile(self, query, params)
350 }
351}