Skip to main content

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 execute;
16mod graph_api;
17mod occ;
18mod pull_mode;
19mod replay;
20mod stream;
21mod write_guard;
22
23use crate::error::LoraError;
24use crate::plan_cache::PlanCache;
25use crate::snapshot::ManagedSnapshotStore;
26use crate::wal::write_scope::WalAbortPolicy;
27
28/// Minimal abstraction any transport can depend on to run Lora queries.
29pub trait QueryRunner: Send + Sync + 'static {
30    fn execute(
31        &self,
32        query: &str,
33        options: Option<ExecuteOptions>,
34    ) -> Result<QueryResult, LoraError>;
35}
36
37/// Owns the graph store and orchestrates parse → analyze → compile → execute.
38///
39/// Optionally drives a write-ahead log: when constructed via
40/// [`Database::open_with_wal`] or [`Database::recover`] the database
41/// holds an [`Arc<WalRecorder>`] that brackets every query with
42/// `begin → mutations → commit/abort → flush` while the store write
43/// lock is held, so the WAL order is exactly the in-memory commit order.
44/// When constructed via [`Database::in_memory`] / [`Database::from_graph`]
45/// the WAL handle is `None` and the engine pays only the existing
46/// `MutationRecorder::record` null-pointer check per mutation.
47pub struct Database<S> {
48    /// The current authoritative store, atomically swappable. Reads call
49    /// `store.load_full()` to obtain an `Arc<S>` snapshot — no lock, no
50    /// blocking — and run their executor against `&*snapshot`. Writes
51    /// take the `writer` Mutex (for commit-order serialization), clone
52    /// the current snapshot into a working copy, mutate that copy, append
53    /// to the WAL, then `store.store(Arc::new(staged))` to publish.
54    /// Concurrent reads keep their old `Arc<S>` alive until they drop it,
55    /// which gives natural snapshot isolation.
56    pub(crate) store: Arc<ArcSwap<S>>,
57    /// Serializes commit ordering. Held only across `clone-mutate-WAL-publish`,
58    /// not around any read. Multiple readers proceed concurrently with a
59    /// writer; only writers contend with each other on this Mutex.
60    pub(crate) writer: Arc<Mutex<()>>,
61    /// Per-record write locks, plumbed in Phase 4.1. Phase 4.2 keeps
62    /// the writer Mutex serialization model so the lock table sits
63    /// idle for now; it becomes load-bearing in a future phase that
64    /// drops the writer Mutex in favour of ArcSwap CAS for true
65    /// concurrent commits, where two writers with overlapping write
66    /// sets need a real serialization point.
67    #[allow(dead_code)]
68    pub(crate) lock_table: Arc<lora_store::LockTable>,
69    pub(crate) wal: Option<Arc<WalRecorder>>,
70    pub(crate) snapshots: Option<Arc<ManagedSnapshotStore>>,
71    /// Cache of compiled query plans, content-keyed by raw query text. Shared
72    /// across the read- and write-lock phases of a single execute (so a
73    /// mutating query compiles at most once instead of twice) and across
74    /// every subsequent call that uses the same query string.
75    pub(crate) plan_cache: Arc<PlanCache>,
76}
77
78#[derive(Debug, Clone, Copy, PartialEq, Eq)]
79pub enum GraphDirection {
80    Outgoing,
81    Incoming,
82    Both,
83}
84
85impl GraphDirection {
86    pub(crate) fn as_store_direction(self) -> Direction {
87        match self {
88            Self::Outgoing => Direction::Right,
89            Self::Incoming => Direction::Left,
90            Self::Both => Direction::Undirected,
91        }
92    }
93}
94
95pub(crate) fn values_to_properties(values: BTreeMap<String, LoraValue>) -> Result<Properties> {
96    values
97        .into_iter()
98        .map(|(key, value)| {
99            let value = lora_value_to_property(value).map_err(|e| anyhow!(e))?;
100            Ok((key, value))
101        })
102        .collect()
103}
104
105pub(crate) const QUERY_FAILURE_POISON: &str =
106    "query mutated the live graph before failing; restart from snapshot + WAL required";
107
108impl Database<InMemoryGraph> {
109    /// Force any pending WAL bytes to durable storage and, for archive-backed
110    /// databases, refresh the portable `.loradb` file before returning.
111    ///
112    /// Managed snapshot checkpoints are explicit via
113    /// [`Self::checkpoint_managed`] or threshold-driven via
114    /// [`SnapshotConfig::checkpoint_every_commits`]; `sync()` remains a
115    /// durability operation rather than an O(graph) checkpoint.
116    pub fn sync(&self) -> Result<(), LoraError> {
117        if let Some(wal) = &self.wal {
118            wal.force_fsync()?;
119        }
120        Ok(())
121    }
122}
123
124impl<S> Database<S>
125where
126    S: GraphStorage + GraphStorageMut + Any + Clone + Send + Sync + 'static,
127{
128    /// Handle to the installed WAL recorder, if any. Exposed for
129    /// admin paths (checkpoint, truncate, observability) that need
130    /// to drive the WAL outside the standard query lifecycle.
131    pub fn wal(&self) -> Option<&Arc<WalRecorder>> {
132        self.wal.as_ref()
133    }
134
135    /// Handle to the underlying shared store — useful for callers that need
136    /// to snapshot or share the graph across multiple databases.
137    pub fn store(&self) -> &Arc<ArcSwap<S>> {
138        &self.store
139    }
140
141    /// Parse a query string into an AST without executing it.
142    pub fn parse(&self, query: &str) -> Result<Document, LoraError> {
143        Ok(parse_query(query)?)
144    }
145
146    /// Read the current authoritative snapshot. Lock-free: returns an
147    /// `Arc<S>` whose lifetime is independent of any writer.
148    pub(crate) fn read_store(&self) -> Arc<S> {
149        self.store.load_full()
150    }
151}
152
153impl<S> Database<S>
154where
155    S: GraphStorage + GraphStorageMut + Any + Clone + Send + Sync + 'static,
156{
157    // ---------- Storage-agnostic utility helpers ----------
158    //
159    // Bindings previously reached into the shared store lock to answer
160    // stat / admin calls; these helpers let them depend on `Database<S>`
161    // instead, so swapping in a new backend only requires changing one type
162    // parameter.
163
164    /// Drop every node and relationship, returning WAL/archive errors to the
165    /// caller.
166    ///
167    /// When a WAL is attached, the clear is wrapped in `arm`/`commit` so the
168    /// `MutationEvent::Clear` fired by the store reaches the log inside a
169    /// transaction. If a failure happens after the in-memory graph has been
170    /// cleared, the recorder is poisoned by the failing WAL path and future
171    /// writes fail until the database is reopened from durable state.
172    pub fn try_clear(&self) -> Result<(), LoraError> {
173        let guard = self.write_store();
174        self.with_logged_write_guard(guard, WalAbortPolicy::AbortOnly, |store| {
175            store.clear();
176            Ok(())
177        })
178        .map_err(LoraError::from_anyhow)
179    }
180
181    /// Drop every node and relationship.
182    ///
183    /// This compatibility helper keeps the historical infallible Rust API.
184    /// Bindings that can report errors should call [`Self::try_clear`].
185    pub fn clear(&self) {
186        let _ = self.try_clear();
187    }
188
189    /// Number of nodes currently in the graph.
190    pub fn node_count(&self) -> usize {
191        let snapshot = self.read_store();
192        snapshot.node_count()
193    }
194
195    /// Number of relationships currently in the graph.
196    pub fn relationship_count(&self) -> usize {
197        let snapshot = self.read_store();
198        snapshot.relationship_count()
199    }
200
201    /// Run a closure with a shared borrow of the underlying store.
202    /// Lock-free: callers see a consistent snapshot for the duration of
203    /// the closure even while writers commit new versions.
204    pub fn with_store<R>(&self, f: impl FnOnce(&S) -> R) -> R {
205        let snapshot = self.read_store();
206        f(&*snapshot)
207    }
208
209    /// Run a closure with an exclusive borrow of the underlying store. Reserved
210    /// for admin paths (restore, bulk load); regular mutation goes through
211    /// `execute_with_params`. The closure mutates a staged copy that is
212    /// published atomically when the closure returns.
213    pub fn with_store_mut<R>(&self, f: impl FnOnce(&mut S) -> R) -> R {
214        let mut guard = self.write_store();
215        let result = f(&mut *guard);
216        guard.publish();
217        result
218    }
219}
220
221impl<S> QueryRunner for Database<S>
222where
223    S: GraphStorage + GraphStorageMut + Any + Clone + Send + Sync + 'static,
224{
225    fn execute(
226        &self,
227        query: &str,
228        options: Option<ExecuteOptions>,
229    ) -> Result<QueryResult, LoraError> {
230        Database::execute(self, query, options)
231    }
232}