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
274            .staged_mut_or_error()?
275            .load_snapshot_payload(payload)?;
276        // Publish the staged graph atomically into the live store; dropping
277        // the guard without `publish` would discard the restore (rollback
278        // semantics on the writer lease).
279        guard.publish_in_place();
280        // Still under the writer lock: change feeds see the restore in
281        // commit order.
282        self.publish_reset();
283        drop(guard);
284        Ok(meta)
285    }
286
287    /// Replace the current graph state from a database snapshot file,
288    /// supplying credentials when the snapshot is encrypted.
289    pub fn load_snapshot_from_with_credentials(
290        &self,
291        path: impl AsRef<Path>,
292        credentials: Option<&SnapshotCredentials>,
293    ) -> Result<SnapshotMeta, LoraError> {
294        let bytes = std::fs::read(path.as_ref())?;
295        self.load_snapshot_from_bytes_with_credentials(&bytes, credentials)
296    }
297
298    /// Take a checkpoint: snapshot the current state with the WAL's
299    /// `durable_lsn` stamped into the header, append a `Checkpoint`
300    /// marker to the WAL, then drop sealed segments at or below the
301    /// fence.
302    ///
303    /// Errors with "checkpoint requires WAL enabled" when called on a
304    /// database constructed without a WAL — operators that just want
305    /// a fence-less dump should use [`Self::save_snapshot_to`] instead.
306    ///
307    /// The write-lock-held window covers snapshot serialization plus the
308    /// checkpoint marker append. Truncation runs after the rename
309    /// but still under the write lock; making it concurrent with queries
310    /// is a v2 concern (see `docs/decisions/0004-wal.md`).
311    pub fn checkpoint_to(&self, path: impl AsRef<Path>) -> Result<SnapshotMeta, LoraError> {
312        let recorder = self.wal.as_ref().ok_or_else(|| {
313            LoraError::new(LoraErrorCode::Internal, "checkpoint requires WAL enabled")
314        })?;
315        let path = path.as_ref();
316        let tmp = snapshot_tmp_path(path);
317
318        let guard = self.write_store();
319
320        // Make every record appended so far durable, then capture
321        // the LSN that becomes the snapshot fence.
322        recorder.force_fsync()?;
323        let snapshot_lsn = recorder.wal().durable_lsn();
324
325        let file = OpenOptions::new()
326            .write(true)
327            .create(true)
328            .truncate(true)
329            .open(&tmp)?;
330        let tmp_guard = TempFileGuard::new(tmp.clone());
331        let mut writer = BufWriter::new(file);
332        let payload = guard.staged_or_error()?.snapshot_payload();
333        let options = SnapshotOptions {
334            compression: Compression::None,
335            encryption: None,
336        };
337        let info = encode_snapshot_to(&mut writer, &payload, Some(snapshot_lsn.raw()), &options)?;
338        let meta = snapshot_info_to_meta(info);
339
340        writer.flush()?;
341        let file = writer.into_inner().map_err(|e| e.into_error())?;
342        sync_file(&file)?;
343        drop(file);
344
345        std::fs::rename(&tmp, path)?;
346        tmp_guard.commit();
347
348        sync_parent_dir(path).map_err(|e| LoraError::new(LoraErrorCode::Io, e.to_string()))?;
349
350        // Append the checkpoint marker AFTER the rename succeeds —
351        // this preserves the invariant that a `Checkpoint` record
352        // in the WAL implies the snapshot it points at exists.
353        recorder.checkpoint_marker(snapshot_lsn)?;
354        recorder.force_fsync()?;
355
356        // Best-effort segment truncation. Failure here doesn't undo
357        // the checkpoint — the next call will retry.
358        if let Err(err) = recorder.truncate_up_to(snapshot_lsn) {
359            tracing::warn!(
360                lsn = snapshot_lsn.raw(),
361                error = %err,
362                "WAL truncation after checkpoint failed; will retry later"
363            );
364        }
365
366        Ok(meta)
367    }
368
369    /// Take a checkpoint into the managed snapshot directory configured by
370    /// [`Self::open_with_wal_snapshots`].
371    pub fn checkpoint_managed(&self) -> Result<SnapshotMeta, LoraError> {
372        let recorder = self.wal.as_ref().ok_or_else(|| {
373            LoraError::new(
374                LoraErrorCode::Internal,
375                "managed checkpoint requires WAL enabled",
376            )
377        })?;
378        let snapshots = self.snapshots.as_ref().ok_or_else(|| {
379            LoraError::new(
380                LoraErrorCode::Internal,
381                "managed checkpoint requires snapshots enabled",
382            )
383        })?;
384        let guard = self.write_store();
385        snapshots
386            .checkpoint(guard.staged_or_error()?, recorder)
387            .map_err(Into::into)
388    }
389}
390
391// ---------------------------------------------------------------------------
392// SnapshotAdmin — type-erased admin entry for transports.
393// ---------------------------------------------------------------------------
394
395/// Storage-agnostic admin surface for HTTP / binding callers that want to
396/// drive snapshot operations without naming the backend type parameter.
397///
398/// Implemented on `Database<InMemoryGraph>` since the in-memory backend
399/// is currently the only one that bridges to the columnar
400/// `lora-snapshot` codec. Transports (e.g. `lora-server`) type-erase on
401/// `Arc<dyn SnapshotAdmin>`.
402pub trait SnapshotAdmin: Send + Sync + 'static {
403    fn save_snapshot(&self, path: &Path) -> Result<SnapshotMeta, LoraError>;
404    fn load_snapshot(&self, path: &Path) -> Result<SnapshotMeta, LoraError>;
405}
406
407impl SnapshotAdmin for Database<InMemoryGraph> {
408    fn save_snapshot(&self, path: &Path) -> Result<SnapshotMeta, LoraError> {
409        self.save_snapshot_to(path)
410    }
411
412    fn load_snapshot(&self, path: &Path) -> Result<SnapshotMeta, LoraError> {
413        self.load_snapshot_from(path)
414    }
415}