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