Skip to main content

lora_database/database/
builder.rs

1//! Constructors for [`Database<InMemoryGraph>`] and [`Database<S>`].
2//!
3//! Five public entry points live here:
4//!
5//! * [`Database::in_memory`] — fresh empty in-memory database.
6//! * [`Database::open_with_wal`] — open or create a WAL-backed database.
7//! * [`Database::open_with_wal_snapshots`] — same, with managed snapshots.
8//! * [`Database::open_named`] — open a portable `.loradb` archive.
9//! * [`Database::recover`] — restore from a snapshot then replay the WAL.
10//! * [`Database::new`] / [`Database::from_graph`] — build from a generic store.
11//!
12//! All four WAL paths share the same "install recorder, assemble
13//! Database" tail; that is captured in the private
14//! [`Database::from_graph_with_wal`] helper to keep the public
15//! constructors focused on the recovery decisions specific to each
16//! entry point.
17
18use std::any::Any;
19use std::fs::File;
20use std::io::BufReader;
21use std::path::Path;
22use std::sync::{Arc, Mutex};
23
24use lora_store::{GraphStorage, GraphStorageMut, InMemoryGraph, MutationRecorder};
25use lora_wal::{replay_dir, Lsn, Wal, WalConfig, WalMirror, WalRecorder};
26
27use crate::database::Database;
28use crate::error::{LoraError, LoraErrorCode};
29use crate::live_store::LiveStore;
30use crate::named::{DatabaseName, DatabaseOpenOptions};
31use crate::plan_cache::PlanCache;
32use crate::snapshot::{ManagedSnapshotStore, SnapshotConfig};
33use crate::wal::archive::WalArchive;
34
35use super::replay::replay_into;
36
37impl Database<InMemoryGraph> {
38    /// Convenience constructor: a fresh, empty in-memory graph database.
39    pub fn in_memory() -> Self {
40        Self::from_graph(InMemoryGraph::new())
41    }
42
43    /// Open or create a WAL-enabled in-memory database from a fresh
44    /// graph.
45    ///
46    /// `WalConfig::Disabled` falls back to [`Database::in_memory`].
47    /// Otherwise, opens the WAL directory, replays any committed
48    /// events into a fresh graph, installs a [`WalRecorder`] on the
49    /// graph, and returns a database ready to serve queries.
50    ///
51    /// To restore from a snapshot in addition to the WAL, use
52    /// [`Database::recover`] instead.
53    pub fn open_with_wal(wal_config: WalConfig) -> Result<Self, LoraError> {
54        match wal_config {
55            WalConfig::Disabled => Ok(Self::in_memory()),
56            WalConfig::Enabled {
57                dir,
58                sync_mode,
59                segment_target_bytes,
60            } => {
61                let mut graph = InMemoryGraph::new();
62                let (wal, events) = Wal::open(dir, sync_mode, segment_target_bytes, Lsn::ZERO)?;
63                replay_into(&mut graph, events).map_err(LoraError::from_anyhow)?;
64                let recorder = Arc::new(WalRecorder::new(wal));
65                Ok(Self::from_graph_with_wal(graph, recorder, None, None))
66            }
67        }
68    }
69
70    /// Open or create a WAL-backed database with managed snapshots beside it.
71    ///
72    /// Recovery loads the newest managed snapshot first, then replays WAL
73    /// records above the snapshot's LSN fence. Checkpoints are written through
74    /// [`Self::checkpoint_managed`] / [`Self::sync`], or automatically when
75    /// `snapshot_config.checkpoint_every_commits` is set.
76    pub fn open_with_wal_snapshots(
77        wal_config: WalConfig,
78        snapshot_config: SnapshotConfig,
79    ) -> Result<Self, LoraError> {
80        let snapshot_store =
81            Arc::new(ManagedSnapshotStore::open(snapshot_config).map_err(LoraError::from_anyhow)?);
82        let mut graph = InMemoryGraph::new();
83
84        match wal_config {
85            WalConfig::Disabled => Err(LoraError::new(
86                LoraErrorCode::Config,
87                "managed snapshots require WAL enabled",
88            )),
89            WalConfig::Enabled {
90                dir,
91                sync_mode,
92                segment_target_bytes,
93            } => {
94                let snapshot_lsn = snapshot_store
95                    .load_latest(&mut graph)
96                    .map_err(LoraError::from_anyhow)?;
97                let (wal, events) = Wal::open(dir, sync_mode, segment_target_bytes, snapshot_lsn)?;
98                replay_into(&mut graph, events).map_err(LoraError::from_anyhow)?;
99                let recorder = Arc::new(WalRecorder::new(wal));
100                Ok(Self::from_graph_with_wal(
101                    graph,
102                    recorder,
103                    Some(snapshot_store),
104                    None,
105                ))
106            }
107        }
108    }
109
110    /// Open or create a named portable database rooted under
111    /// `options.database_dir`.
112    ///
113    /// The database name may be either a portable basename (`app` or
114    /// `app.loradb`) or a safe relative path (`tenant/app`). It is resolved
115    /// under `options.database_dir` before the WAL archive backend opens.
116    pub fn open_named(
117        database_name: impl AsRef<str>,
118        options: DatabaseOpenOptions,
119    ) -> Result<Self, LoraError> {
120        let name = DatabaseName::parse(database_name.as_ref())?;
121        let archive = Arc::new(WalArchive::open(
122            options.database_path_for(&name),
123            options.max_database_bytes,
124        )?);
125        let mut graph = InMemoryGraph::new();
126        let snapshot_lsn = if let Some(bytes) = archive.snapshot_bytes()? {
127            let (payload, info) = crate::snapshot::decode_snapshot_bytes(&bytes, None)?;
128            graph.load_snapshot_payload(payload)?;
129            info.wal_lsn.map(Lsn::new).unwrap_or(Lsn::ZERO)
130        } else {
131            Lsn::ZERO
132        };
133        let (wal, events) = Wal::open(
134            archive.work_dir(),
135            options.sync_mode,
136            options.segment_target_bytes,
137            snapshot_lsn,
138        )?;
139        replay_into(&mut graph, events).map_err(LoraError::from_anyhow)?;
140        let mirror: Arc<dyn WalMirror> = archive.clone();
141        let recorder = Arc::new(WalRecorder::new_with_mirror(wal, Some(mirror)));
142        // Mark the archive dirty so a fresh named database is materialized as
143        // a portable Lora container. Follow-up writes refresh the container on explicit
144        // `sync()` / checkpoint and on clean database drop.
145        recorder.flush()?;
146        Ok(Self::from_graph_with_wal(
147            graph,
148            recorder,
149            None,
150            Some(archive),
151        ))
152    }
153
154    /// Restore from a snapshot file then replay any WAL records past
155    /// it.
156    ///
157    /// The snapshot's `wal_lsn` (when set) becomes the replay fence —
158    /// events at or below that LSN are already represented in the
159    /// loaded snapshot and are skipped. A missing snapshot file is
160    /// treated as "fresh start" so operators can pass the same path
161    /// on every boot.
162    ///
163    /// If the WAL contains a checkpoint marker newer than the
164    /// snapshot's `wal_lsn`, a one-line warning is printed to stderr
165    /// — the snapshot is stale relative to a more recent checkpoint
166    /// the operator is presumably aware of. Recovery still proceeds
167    /// from the snapshot's fence (replay re-applies every record
168    /// above it, which is conservative-correct); a tighter contract
169    /// is deferred to v2 because verifying that the marker's
170    /// snapshot file actually exists and is loadable is a separate
171    /// observability concern.
172    pub fn recover(
173        snapshot_path: impl AsRef<Path>,
174        wal_config: WalConfig,
175    ) -> Result<Self, LoraError> {
176        let snapshot_path = snapshot_path.as_ref();
177        let mut graph = InMemoryGraph::new();
178        let snapshot_lsn = match File::open(snapshot_path) {
179            Ok(f) => {
180                let reader = BufReader::new(f);
181                let (payload, info) = crate::snapshot::read_snapshot_from(reader, None)?;
182                graph.load_snapshot_payload(payload)?;
183                info.wal_lsn.map(Lsn::new).unwrap_or(Lsn::ZERO)
184            }
185            Err(e) if e.kind() == std::io::ErrorKind::NotFound => Lsn::ZERO,
186            Err(e) => return Err(e.into()),
187        };
188
189        match wal_config {
190            WalConfig::Disabled => Ok(Self::from_graph(graph)),
191            WalConfig::Enabled {
192                dir,
193                sync_mode,
194                segment_target_bytes,
195            } => {
196                // Diagnostic peek at the WAL's newest checkpoint
197                // marker so we can warn the operator about a stale
198                // snapshot before we start replaying. Treat any error
199                // as "no marker" — the subsequent `Wal::open` will
200                // surface the real failure if there is one.
201                if dir.exists() {
202                    if let Ok(outcome) = replay_dir(&dir, Lsn::ZERO) {
203                        if let Some(marker) = outcome.checkpoint_lsn_observed {
204                            if marker > snapshot_lsn {
205                                eprintln!(
206                                    "lora-wal: snapshot at LSN {} is older than the newest \
207                                     checkpoint marker on disk (LSN {}). Replaying every WAL \
208                                     record above LSN {}; consider passing the more recent \
209                                     snapshot to --restore-from.",
210                                    snapshot_lsn.raw(),
211                                    marker.raw(),
212                                    snapshot_lsn.raw()
213                                );
214                            }
215                        }
216                    }
217                }
218
219                let (wal, events) = Wal::open(dir, sync_mode, segment_target_bytes, snapshot_lsn)?;
220                replay_into(&mut graph, events).map_err(LoraError::from_anyhow)?;
221                let recorder = Arc::new(WalRecorder::new(wal));
222                Ok(Self::from_graph_with_wal(graph, recorder, None, None))
223            }
224        }
225    }
226
227    /// Install the durable recorder on `graph` and assemble the
228    /// `Database` envelope. Shared by every WAL-backed constructor.
229    fn from_graph_with_wal(
230        mut graph: InMemoryGraph,
231        recorder: Arc<WalRecorder>,
232        snapshots: Option<Arc<ManagedSnapshotStore>>,
233        named_archive: Option<Arc<WalArchive>>,
234    ) -> Self {
235        graph.set_mutation_recorder(Some(recorder.clone() as Arc<dyn MutationRecorder>));
236        Self {
237            store: Arc::new(LiveStore::new(Arc::new(graph))),
238            writer: Arc::new(Mutex::new(())),
239            lock_table: Arc::new(lora_store::LockTable::new()),
240            wal: Some(recorder),
241            snapshots,
242            named_archive,
243            plan_cache: Arc::new(PlanCache::new()),
244            changes: Arc::new(crate::changes::ChangeHub::default()),
245        }
246    }
247}
248
249impl<S> Database<S>
250where
251    S: GraphStorage + GraphStorageMut + Any + Clone + Send + Sync + 'static,
252{
253    /// Build a database from a pre-wrapped, shared store.
254    pub(crate) fn new(store: Arc<LiveStore<S>>) -> Self {
255        Self {
256            store,
257            writer: Arc::new(Mutex::new(())),
258            lock_table: Arc::new(lora_store::LockTable::new()),
259            wal: None,
260            snapshots: None,
261            named_archive: None,
262            plan_cache: Arc::new(PlanCache::new()),
263            changes: Arc::new(crate::changes::ChangeHub::default()),
264        }
265    }
266
267    /// Build a database by taking ownership of a bare graph store.
268    pub fn from_graph(graph: S) -> Self {
269        Self::new(Arc::new(LiveStore::new(Arc::new(graph))))
270    }
271}