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 arc_swap::ArcSwap;
25use lora_store::{GraphStorage, GraphStorageMut, InMemoryGraph, MutationRecorder};
26use lora_wal::{replay_dir, Lsn, Wal, WalConfig, WalMirror, WalRecorder};
27
28use crate::database::Database;
29use crate::error::{LoraError, LoraErrorCode};
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))
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 ))
105 }
106 }
107 }
108
109 /// Open or create a named portable database rooted under
110 /// `options.database_dir`.
111 ///
112 /// The database name may be either a portable basename (`app` or
113 /// `app.loradb`) or a safe relative path (`tenant/app`). It is resolved
114 /// under `options.database_dir` before the WAL archive backend opens.
115 pub fn open_named(
116 database_name: impl AsRef<str>,
117 options: DatabaseOpenOptions,
118 ) -> Result<Self, LoraError> {
119 let name = DatabaseName::parse(database_name.as_ref())?;
120 let archive = Arc::new(WalArchive::open(
121 options.database_path_for(&name),
122 options.max_database_bytes,
123 )?);
124 let mut graph = InMemoryGraph::new();
125 let (wal, events) = Wal::open(
126 archive.work_dir(),
127 options.sync_mode,
128 options.segment_target_bytes,
129 Lsn::ZERO,
130 )?;
131 replay_into(&mut graph, events).map_err(LoraError::from_anyhow)?;
132 let mirror: Arc<dyn WalMirror> = archive;
133 let recorder = Arc::new(WalRecorder::new_with_mirror(wal, Some(mirror)));
134 // Mark the archive dirty so a fresh named database is materialized as
135 // a portable ZIP. The archive writer coalesces this with any immediate
136 // follow-up writes and flushes it in the background, with a final flush
137 // on database drop.
138 recorder.flush()?;
139 Ok(Self::from_graph_with_wal(graph, recorder, None))
140 }
141
142 /// Restore from a snapshot file then replay any WAL records past
143 /// it.
144 ///
145 /// The snapshot's `wal_lsn` (when set) becomes the replay fence —
146 /// events at or below that LSN are already represented in the
147 /// loaded snapshot and are skipped. A missing snapshot file is
148 /// treated as "fresh start" so operators can pass the same path
149 /// on every boot.
150 ///
151 /// If the WAL contains a checkpoint marker newer than the
152 /// snapshot's `wal_lsn`, a one-line warning is printed to stderr
153 /// — the snapshot is stale relative to a more recent checkpoint
154 /// the operator is presumably aware of. Recovery still proceeds
155 /// from the snapshot's fence (replay re-applies every record
156 /// above it, which is conservative-correct); a tighter contract
157 /// is deferred to v2 because verifying that the marker's
158 /// snapshot file actually exists and is loadable is a separate
159 /// observability concern.
160 pub fn recover(
161 snapshot_path: impl AsRef<Path>,
162 wal_config: WalConfig,
163 ) -> Result<Self, LoraError> {
164 let snapshot_path = snapshot_path.as_ref();
165 let mut graph = InMemoryGraph::new();
166 let snapshot_lsn = match File::open(snapshot_path) {
167 Ok(f) => {
168 let reader = BufReader::new(f);
169 let (payload, info) = crate::snapshot::read_snapshot_from(reader, None)?;
170 graph.load_snapshot_payload(payload)?;
171 info.wal_lsn.map(Lsn::new).unwrap_or(Lsn::ZERO)
172 }
173 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Lsn::ZERO,
174 Err(e) => return Err(e.into()),
175 };
176
177 match wal_config {
178 WalConfig::Disabled => Ok(Self::from_graph(graph)),
179 WalConfig::Enabled {
180 dir,
181 sync_mode,
182 segment_target_bytes,
183 } => {
184 // Diagnostic peek at the WAL's newest checkpoint
185 // marker so we can warn the operator about a stale
186 // snapshot before we start replaying. Treat any error
187 // as "no marker" — the subsequent `Wal::open` will
188 // surface the real failure if there is one.
189 if dir.exists() {
190 if let Ok(outcome) = replay_dir(&dir, Lsn::ZERO) {
191 if let Some(marker) = outcome.checkpoint_lsn_observed {
192 if marker > snapshot_lsn {
193 eprintln!(
194 "lora-wal: snapshot at LSN {} is older than the newest \
195 checkpoint marker on disk (LSN {}). Replaying every WAL \
196 record above LSN {}; consider passing the more recent \
197 snapshot to --restore-from.",
198 snapshot_lsn.raw(),
199 marker.raw(),
200 snapshot_lsn.raw()
201 );
202 }
203 }
204 }
205 }
206
207 let (wal, events) = Wal::open(dir, sync_mode, segment_target_bytes, snapshot_lsn)?;
208 replay_into(&mut graph, events).map_err(LoraError::from_anyhow)?;
209 let recorder = Arc::new(WalRecorder::new(wal));
210 Ok(Self::from_graph_with_wal(graph, recorder, None))
211 }
212 }
213 }
214
215 /// Install the durable recorder on `graph` and assemble the
216 /// `Database` envelope. Shared by every WAL-backed constructor.
217 fn from_graph_with_wal(
218 mut graph: InMemoryGraph,
219 recorder: Arc<WalRecorder>,
220 snapshots: Option<Arc<ManagedSnapshotStore>>,
221 ) -> Self {
222 graph.set_mutation_recorder(Some(recorder.clone() as Arc<dyn MutationRecorder>));
223 Self {
224 store: Arc::new(ArcSwap::from(Arc::new(graph))),
225 writer: Arc::new(Mutex::new(())),
226 lock_table: Arc::new(lora_store::LockTable::new()),
227 wal: Some(recorder),
228 snapshots,
229 plan_cache: Arc::new(PlanCache::new()),
230 }
231 }
232}
233
234impl<S> Database<S>
235where
236 S: GraphStorage + GraphStorageMut + Any + Clone + Send + Sync + 'static,
237{
238 /// Build a database from a pre-wrapped, shared store.
239 pub fn new(store: Arc<ArcSwap<S>>) -> Self {
240 Self {
241 store,
242 writer: Arc::new(Mutex::new(())),
243 lock_table: Arc::new(lora_store::LockTable::new()),
244 wal: None,
245 snapshots: None,
246 plan_cache: Arc::new(PlanCache::new()),
247 }
248 }
249
250 /// Build a database by taking ownership of a bare graph store.
251 pub fn from_graph(graph: S) -> Self {
252 Self::new(Arc::new(ArcSwap::from(Arc::new(graph))))
253 }
254}