Skip to main content

lora_database/snapshot/
mod.rs

1//! Snapshot integration for [`Database`]. Owns:
2//!
3//! * the byte-level entry points (`save_snapshot_to_*`, `load_snapshot_*`,
4//!   `checkpoint_to`, `checkpoint_managed`, `recover`-side helpers),
5//! * the JSON option/credential adapters used by the language bindings,
6//! * the [`SnapshotByteFormat`] sniff and the [`SnapshotAdmin`] trait,
7//! * filesystem hygiene helpers (`TempFileGuard`, `snapshot_tmp_path`,
8//!   `sync_parent_dir`).
9//!
10//! Every byte-level path here goes through the columnar `lora-snapshot`
11//! codec exclusively; the legacy in-store `LORASNAP` format was retired.
12//! `lora-store` now owns only the [`SnapshotPayload`] vocabulary, this
13//! module owns the encode/decode integration.
14
15mod json;
16pub(crate) mod store;
17
18pub use json::{snapshot_credentials_from_json, snapshot_options_from_json};
19pub(crate) use store::ManagedSnapshotStore;
20pub use store::SnapshotConfig;
21
22use std::fs::OpenOptions;
23use std::io::{BufWriter, Read, Write};
24use std::path::{Path, PathBuf};
25
26use anyhow::Result;
27
28use lora_snapshot::{
29    decode_snapshot as decode_database_snapshot, read_snapshot as read_database_snapshot,
30    write_snapshot as write_database_snapshot, Compression, SnapshotCodecError,
31    SnapshotCredentials, SnapshotInfo, SnapshotOptions, DATABASE_SNAPSHOT_MAGIC,
32};
33use lora_store::{InMemoryGraph, SnapshotMeta, SnapshotPayload};
34
35use crate::durable_io::{sync_dir, sync_file};
36use crate::error::{LoraError, LoraErrorCode};
37use crate::Database;
38
39/// Magic-byte sniff for snapshot bytes. The legacy in-store `LORASNAP`
40/// codec was removed in favour of the columnar `lora-snapshot` format,
41/// so this collapses to a single recognized variant; kept as a typed
42/// detect for forward compatibility if a future format is introduced.
43#[derive(Debug, Clone, Copy, PartialEq, Eq)]
44pub enum SnapshotByteFormat {
45    Database,
46}
47
48impl SnapshotByteFormat {
49    pub fn detect(bytes: &[u8]) -> Option<Self> {
50        if bytes.starts_with(DATABASE_SNAPSHOT_MAGIC) {
51            Some(Self::Database)
52        } else {
53            None
54        }
55    }
56}
57
58pub(crate) fn snapshot_info_to_meta(info: SnapshotInfo) -> SnapshotMeta {
59    SnapshotMeta {
60        format_version: info.format_version,
61        node_count: info.node_count,
62        relationship_count: info.relationship_count,
63        wal_lsn: info.wal_lsn,
64    }
65}
66
67// ---------------------------------------------------------------------------
68// Filesystem hygiene helpers
69// ---------------------------------------------------------------------------
70
71pub(crate) fn snapshot_tmp_path(target: &Path) -> PathBuf {
72    let mut tmp = target.as_os_str().to_owned();
73    tmp.push(".tmp");
74    PathBuf::from(tmp)
75}
76
77pub(crate) fn sync_parent_dir(path: &Path) -> Result<()> {
78    let Some(parent) = path.parent() else {
79        return Ok(());
80    };
81    Ok(sync_dir(parent)?)
82}
83
84/// RAII handle that deletes its path on drop unless [`commit`] is called.
85///
86/// The snapshot save path creates `<target>.tmp` before the payload is
87/// written; if any step between then and the final rename fails (or the
88/// thread unwinds), the guard's `Drop` removes the scratch file so a crashed
89/// save never leaves leftovers on disk.
90///
91/// [`commit`]: Self::commit
92pub(crate) struct TempFileGuard {
93    path: Option<PathBuf>,
94}
95
96impl TempFileGuard {
97    pub(crate) fn new(path: PathBuf) -> Self {
98        Self { path: Some(path) }
99    }
100
101    /// Disarm the guard. Call this once the tmp file's contents have been
102    /// handed off (e.g. renamed to their final destination) so the `Drop`
103    /// impl does not try to remove them.
104    pub(crate) fn commit(mut self) {
105        self.path.take();
106    }
107}
108
109impl Drop for TempFileGuard {
110    fn drop(&mut self) {
111        if let Some(path) = self.path.take() {
112            // Best-effort: cleanup failure is not worth surfacing — the
113            // worst case is a leaked scratch file that the next save
114            // overwrites via `OpenOptions::truncate(true)`.
115            let _ = std::fs::remove_file(path);
116        }
117    }
118}
119
120/// Decode columnar snapshot bytes into a payload + info. Returns the
121/// underlying typed [`SnapshotCodecError`] so callers can `?` it
122/// straight into a [`LoraError`] via the `From` impl.
123pub(crate) fn decode_snapshot_bytes(
124    bytes: &[u8],
125    credentials: Option<&SnapshotCredentials>,
126) -> Result<(SnapshotPayload, SnapshotInfo), SnapshotCodecError> {
127    decode_database_snapshot(bytes, credentials)
128}
129
130/// Decode columnar snapshot bytes streamed from a reader. Used by
131/// `Database::recover` at startup.
132pub(crate) fn read_snapshot_from<R: Read>(
133    reader: R,
134    credentials: Option<&SnapshotCredentials>,
135) -> Result<(SnapshotPayload, SnapshotInfo), SnapshotCodecError> {
136    read_database_snapshot(reader, credentials)
137}
138
139/// Encode a payload through the columnar codec.
140pub(crate) fn encode_snapshot_to<W: Write>(
141    writer: W,
142    payload: &SnapshotPayload,
143    wal_lsn: Option<u64>,
144    options: &SnapshotOptions,
145) -> Result<SnapshotInfo, SnapshotCodecError> {
146    write_database_snapshot(writer, payload, wal_lsn, options)
147}
148
149// ---------------------------------------------------------------------------
150// Database<InMemoryGraph> snapshot surface
151// ---------------------------------------------------------------------------
152
153impl Database<InMemoryGraph> {
154    /// Serialize the current graph state to `path` using the default
155    /// columnar codec options (uncompressed, unencrypted). Writes are
156    /// atomic via `<path>.tmp` + rename + parent-dir fsync.
157    ///
158    /// Callers that need compression or encryption should reach for
159    /// [`Self::save_snapshot_to_with_options`] directly.
160    pub fn save_snapshot_to(&self, path: impl AsRef<Path>) -> Result<SnapshotMeta, LoraError> {
161        let options = SnapshotOptions {
162            compression: Compression::None,
163            encryption: None,
164        };
165        self.save_snapshot_to_with_options(path, &options)
166    }
167
168    /// Replace the current graph state with a snapshot loaded from `path`.
169    /// Decodes via the columnar codec; encrypted snapshots require
170    /// [`Self::load_snapshot_from_with_credentials`] instead.
171    pub fn load_snapshot_from(&self, path: impl AsRef<Path>) -> Result<SnapshotMeta, LoraError> {
172        self.load_snapshot_from_with_credentials(path, None)
173    }
174
175    /// Convenience constructor: open (or create) an empty in-memory database
176    /// and immediately restore it from `path`. Errors if the file cannot be
177    /// opened or the snapshot is malformed.
178    pub fn in_memory_from_snapshot(path: impl AsRef<Path>) -> Result<Self, LoraError> {
179        let db = Self::in_memory();
180        db.load_snapshot_from_with_credentials(path, None)?;
181        Ok(db)
182    }
183
184    /// Serialize the current graph state into the database snapshot byte
185    /// format.
186    ///
187    /// This uses the column-oriented `lora-snapshot` codec — the same one
188    /// driven by `save_snapshot_to_with_options`, but without a WAL fence.
189    /// The default is uncompressed so bytes stay portable across native
190    /// and WASM builds; callers that want a specific codec can use
191    /// [`Self::save_snapshot_to_bytes_with_options`].
192    pub fn save_snapshot_to_bytes(&self) -> Result<Vec<u8>, LoraError> {
193        let options = SnapshotOptions {
194            compression: Compression::None,
195            encryption: None,
196        };
197        let (bytes, _) = self.save_snapshot_to_bytes_with_options(&options)?;
198        Ok(bytes)
199    }
200
201    /// Serialize the current graph state into database snapshot bytes with
202    /// explicit codec options.
203    pub fn save_snapshot_to_bytes_with_options(
204        &self,
205        options: &SnapshotOptions,
206    ) -> Result<(Vec<u8>, SnapshotInfo), LoraError> {
207        let guard = self.read_store();
208        let payload = guard.snapshot_payload();
209        let mut bytes = Vec::new();
210        let info = encode_snapshot_to(&mut bytes, &payload, None, options)?;
211        Ok((bytes, info))
212    }
213
214    /// Serialize the current graph state to a database snapshot file with
215    /// explicit codec options. This is the path form of
216    /// [`Self::save_snapshot_to_bytes_with_options`] and supports the same
217    /// compression and encryption options.
218    pub fn save_snapshot_to_with_options(
219        &self,
220        path: impl AsRef<Path>,
221        options: &SnapshotOptions,
222    ) -> Result<SnapshotMeta, LoraError> {
223        let path = path.as_ref();
224        let tmp = snapshot_tmp_path(path);
225        let guard = self.read_store();
226
227        let file = OpenOptions::new()
228            .write(true)
229            .create(true)
230            .truncate(true)
231            .open(&tmp)?;
232        let tmp_guard = TempFileGuard::new(tmp.clone());
233        let mut writer = BufWriter::new(file);
234
235        let payload = guard.snapshot_payload();
236        let info = encode_snapshot_to(&mut writer, &payload, None, options)?;
237
238        writer.flush()?;
239        let file = writer.into_inner().map_err(|e| e.into_error())?;
240        sync_file(&file)?;
241        drop(file);
242
243        std::fs::rename(&tmp, path)?;
244        tmp_guard.commit();
245
246        sync_parent_dir(path).map_err(|e| LoraError::new(LoraErrorCode::Io, e.to_string()))?;
247
248        Ok(snapshot_info_to_meta(info))
249    }
250
251    /// Replace the current graph state from snapshot bytes (columnar
252    /// `lora-snapshot` format).
253    pub fn load_snapshot_from_bytes(&self, bytes: &[u8]) -> Result<SnapshotMeta, LoraError> {
254        self.load_snapshot_from_bytes_with_credentials(bytes, None)
255    }
256
257    /// Replace the current graph state from snapshot bytes, supplying
258    /// credentials when loading an encrypted database snapshot.
259    pub fn load_snapshot_from_bytes_with_credentials(
260        &self,
261        bytes: &[u8],
262        credentials: Option<&SnapshotCredentials>,
263    ) -> Result<SnapshotMeta, LoraError> {
264        if SnapshotByteFormat::detect(bytes).is_none() {
265            return Err(LoraError::new(
266                LoraErrorCode::SnapshotCodec,
267                "snapshot bytes have unrecognized magic",
268            ));
269        }
270        let mut guard = self.write_store();
271        let (payload, info) = decode_snapshot_bytes(bytes, credentials)?;
272        let meta = snapshot_info_to_meta(info);
273        guard.load_snapshot_payload(payload)?;
274        // Publish the staged graph atomically into the live store; dropping
275        // the guard without `publish` would discard the restore (rollback
276        // semantics on the writer lease).
277        guard.publish();
278        Ok(meta)
279    }
280
281    /// Replace the current graph state from a database snapshot file,
282    /// supplying credentials when the snapshot is encrypted.
283    pub fn load_snapshot_from_with_credentials(
284        &self,
285        path: impl AsRef<Path>,
286        credentials: Option<&SnapshotCredentials>,
287    ) -> Result<SnapshotMeta, LoraError> {
288        let bytes = std::fs::read(path.as_ref())?;
289        self.load_snapshot_from_bytes_with_credentials(&bytes, credentials)
290    }
291
292    /// Take a checkpoint: snapshot the current state with the WAL's
293    /// `durable_lsn` stamped into the header, append a `Checkpoint`
294    /// marker to the WAL, then drop sealed segments at or below the
295    /// fence.
296    ///
297    /// Errors with "checkpoint requires WAL enabled" when called on a
298    /// database constructed without a WAL — operators that just want
299    /// a fence-less dump should use [`Self::save_snapshot_to`] instead.
300    ///
301    /// The write-lock-held window covers snapshot serialization plus the
302    /// checkpoint marker append. Truncation runs after the rename
303    /// but still under the write lock; making it concurrent with queries
304    /// is a v2 concern (see `docs/decisions/0004-wal.md`).
305    pub fn checkpoint_to(&self, path: impl AsRef<Path>) -> Result<SnapshotMeta, LoraError> {
306        let recorder = self.wal.as_ref().ok_or_else(|| {
307            LoraError::new(LoraErrorCode::Internal, "checkpoint requires WAL enabled")
308        })?;
309        let path = path.as_ref();
310        let tmp = snapshot_tmp_path(path);
311
312        let guard = self.write_store();
313
314        // Make every record appended so far durable, then capture
315        // the LSN that becomes the snapshot fence.
316        recorder.force_fsync()?;
317        let snapshot_lsn = recorder.wal().durable_lsn();
318
319        let file = OpenOptions::new()
320            .write(true)
321            .create(true)
322            .truncate(true)
323            .open(&tmp)?;
324        let tmp_guard = TempFileGuard::new(tmp.clone());
325        let mut writer = BufWriter::new(file);
326        let payload = guard.snapshot_payload();
327        let options = SnapshotOptions {
328            compression: Compression::None,
329            encryption: None,
330        };
331        let info = encode_snapshot_to(&mut writer, &payload, Some(snapshot_lsn.raw()), &options)?;
332        let meta = snapshot_info_to_meta(info);
333
334        writer.flush()?;
335        let file = writer.into_inner().map_err(|e| e.into_error())?;
336        sync_file(&file)?;
337        drop(file);
338
339        std::fs::rename(&tmp, path)?;
340        tmp_guard.commit();
341
342        sync_parent_dir(path).map_err(|e| LoraError::new(LoraErrorCode::Io, e.to_string()))?;
343
344        // Append the checkpoint marker AFTER the rename succeeds —
345        // this preserves the invariant that a `Checkpoint` record
346        // in the WAL implies the snapshot it points at exists.
347        recorder.checkpoint_marker(snapshot_lsn)?;
348        recorder.force_fsync()?;
349
350        // Best-effort segment truncation. Failure here doesn't undo
351        // the checkpoint — the next call will retry.
352        if let Err(err) = recorder.truncate_up_to(snapshot_lsn) {
353            tracing::warn!(
354                lsn = snapshot_lsn.raw(),
355                error = %err,
356                "WAL truncation after checkpoint failed; will retry later"
357            );
358        }
359
360        Ok(meta)
361    }
362
363    /// Take a checkpoint into the managed snapshot directory configured by
364    /// [`Self::open_with_wal_snapshots`].
365    pub fn checkpoint_managed(&self) -> Result<SnapshotMeta, LoraError> {
366        let recorder = self.wal.as_ref().ok_or_else(|| {
367            LoraError::new(
368                LoraErrorCode::Internal,
369                "managed checkpoint requires WAL enabled",
370            )
371        })?;
372        let snapshots = self.snapshots.as_ref().ok_or_else(|| {
373            LoraError::new(
374                LoraErrorCode::Internal,
375                "managed checkpoint requires snapshots enabled",
376            )
377        })?;
378        let guard = self.write_store();
379        snapshots.checkpoint(&guard, recorder).map_err(Into::into)
380    }
381}
382
383// ---------------------------------------------------------------------------
384// SnapshotAdmin — type-erased admin entry for transports.
385// ---------------------------------------------------------------------------
386
387/// Storage-agnostic admin surface for HTTP / binding callers that want to
388/// drive snapshot operations without naming the backend type parameter.
389///
390/// Implemented on `Database<InMemoryGraph>` since the in-memory backend
391/// is currently the only one that bridges to the columnar
392/// `lora-snapshot` codec. Transports (e.g. `lora-server`) type-erase on
393/// `Arc<dyn SnapshotAdmin>`.
394pub trait SnapshotAdmin: Send + Sync + 'static {
395    fn save_snapshot(&self, path: &Path) -> Result<SnapshotMeta, LoraError>;
396    fn load_snapshot(&self, path: &Path) -> Result<SnapshotMeta, LoraError>;
397}
398
399impl SnapshotAdmin for Database<InMemoryGraph> {
400    fn save_snapshot(&self, path: &Path) -> Result<SnapshotMeta, LoraError> {
401        self.save_snapshot_to(path)
402    }
403
404    fn load_snapshot(&self, path: &Path) -> Result<SnapshotMeta, LoraError> {
405        self.load_snapshot_from(path)
406    }
407}