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