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