Skip to main content

lora_database/snapshot/
store.rs

1use std::fs::{self, File, OpenOptions};
2use std::io::{BufReader, BufWriter, Write};
3use std::path::{Path, PathBuf};
4use std::sync::atomic::{AtomicU64, Ordering};
5
6use anyhow::{anyhow, Context, Result};
7use lora_snapshot::{read_snapshot, write_snapshot, SnapshotOptions};
8use lora_store::{InMemoryGraph, SnapshotMeta};
9use lora_wal::{Lsn, WalRecorder};
10
11use crate::durable_io::{sync_dir, sync_file};
12
13const CURRENT_FILE: &str = "CURRENT";
14const SNAPSHOT_PREFIX: &str = "snapshot-";
15const SNAPSHOT_SUFFIX: &str = ".lsnap";
16
17#[derive(Debug, Clone)]
18pub struct SnapshotConfig {
19    pub dir: PathBuf,
20    /// When set, create a managed checkpoint after this many committed WAL
21    /// transactions. `None` keeps checkpointing manual via `sync()` /
22    /// `checkpoint_managed()`.
23    pub checkpoint_every_commits: Option<u64>,
24    /// Number of older checkpoint files to retain in addition to `CURRENT`.
25    pub keep_old: usize,
26    /// Columnar snapshot codec options. Defaults to fast gzip compression and
27    /// no encryption.
28    pub codec: SnapshotOptions,
29}
30
31impl SnapshotConfig {
32    pub fn enabled(dir: impl Into<PathBuf>) -> Self {
33        Self {
34            dir: dir.into(),
35            checkpoint_every_commits: None,
36            keep_old: 1,
37            codec: SnapshotOptions::default(),
38        }
39    }
40
41    pub fn every_commits(mut self, commits: u64) -> Self {
42        self.checkpoint_every_commits = Some(commits.max(1));
43        self
44    }
45
46    pub fn keep_old(mut self, keep_old: usize) -> Self {
47        self.keep_old = keep_old;
48        self
49    }
50
51    pub fn codec(mut self, codec: SnapshotOptions) -> Self {
52        self.codec = codec;
53        self
54    }
55}
56
57pub(crate) struct ManagedSnapshotStore {
58    config: SnapshotConfig,
59    commits_since_checkpoint: AtomicU64,
60}
61
62impl ManagedSnapshotStore {
63    pub(crate) fn open(config: SnapshotConfig) -> Result<Self> {
64        fs::create_dir_all(&config.dir)
65            .with_context(|| format!("create snapshot dir {}", config.dir.display()))?;
66        Ok(Self {
67            config,
68            commits_since_checkpoint: AtomicU64::new(0),
69        })
70    }
71
72    pub(crate) fn load_latest(&self, graph: &mut InMemoryGraph) -> Result<Lsn> {
73        let Some(path) = self.latest_snapshot_path()? else {
74            return Ok(Lsn::ZERO);
75        };
76        let file =
77            File::open(&path).with_context(|| format!("open snapshot {}", path.display()))?;
78        let (payload, info) =
79            read_snapshot(BufReader::new(file), self.config.codec.encryption.as_ref())
80                .with_context(|| format!("load snapshot {}", path.display()))?;
81        graph.load_snapshot_payload(payload)?;
82        Ok(info.wal_lsn.map(Lsn::new).unwrap_or(Lsn::ZERO))
83    }
84
85    /// Managed snapshot files whose LSN fence is at or below `lsn`,
86    /// newest first. Used by change-feed history replay to find a base.
87    pub(crate) fn snapshot_files_at_or_below(&self, lsn: Lsn) -> Result<Vec<(Lsn, PathBuf)>> {
88        let mut files = snapshot_files(&self.config.dir)?;
89        files.retain(|(file_lsn, _)| *file_lsn <= lsn);
90        files.sort_by_key(|(file_lsn, _)| std::cmp::Reverse(*file_lsn));
91        Ok(files)
92    }
93
94    /// Load one managed snapshot file into `graph`, returning its LSN
95    /// fence.
96    pub(crate) fn load_file_into(&self, path: &Path, graph: &mut InMemoryGraph) -> Result<Lsn> {
97        let file = File::open(path).with_context(|| format!("open snapshot {}", path.display()))?;
98        let (payload, info) =
99            read_snapshot(BufReader::new(file), self.config.codec.encryption.as_ref())
100                .with_context(|| format!("load snapshot {}", path.display()))?;
101        graph.load_snapshot_payload(payload)?;
102        Ok(info.wal_lsn.map(Lsn::new).unwrap_or(Lsn::ZERO))
103    }
104
105    pub(crate) fn checkpoint(
106        &self,
107        graph: &InMemoryGraph,
108        recorder: &WalRecorder,
109    ) -> Result<SnapshotMeta> {
110        recorder
111            .force_fsync()
112            .map_err(|e| anyhow!("WAL fsync before managed snapshot failed: {e}"))?;
113        let snapshot_lsn = recorder.wal().durable_lsn();
114        let meta = self.write_snapshot(graph, snapshot_lsn)?;
115
116        recorder
117            .checkpoint_marker(snapshot_lsn)
118            .map_err(|e| anyhow!("WAL checkpoint marker failed: {e}"))?;
119        recorder
120            .force_fsync()
121            .map_err(|e| anyhow!("WAL fsync after checkpoint marker failed: {e}"))?;
122        if let Err(err) = recorder.truncate_up_to(snapshot_lsn) {
123            tracing::warn!(
124                lsn = snapshot_lsn.raw(),
125                error = %err,
126                "WAL truncation after managed checkpoint failed; will retry later"
127            );
128        }
129
130        self.commits_since_checkpoint.store(0, Ordering::Relaxed);
131        self.prune_old_snapshots(snapshot_lsn)?;
132        Ok(meta)
133    }
134
135    pub(crate) fn observe_commit(
136        &self,
137        graph: &InMemoryGraph,
138        recorder: &WalRecorder,
139    ) -> Result<()> {
140        let Some(every) = self.config.checkpoint_every_commits else {
141            return Ok(());
142        };
143        let commits = self
144            .commits_since_checkpoint
145            .fetch_add(1, Ordering::Relaxed)
146            + 1;
147        if commits >= every {
148            self.checkpoint(graph, recorder)?;
149        }
150        Ok(())
151    }
152
153    fn write_snapshot(&self, graph: &InMemoryGraph, snapshot_lsn: Lsn) -> Result<SnapshotMeta> {
154        let target = snapshot_path(&self.config.dir, snapshot_lsn);
155        let tmp = tmp_path(&target);
156        let file = OpenOptions::new()
157            .write(true)
158            .create(true)
159            .truncate(true)
160            .open(&tmp)
161            .with_context(|| format!("open temp snapshot {}", tmp.display()))?;
162        let mut writer = BufWriter::new(file);
163        let payload = graph.snapshot_payload();
164        let info = write_snapshot(
165            &mut writer,
166            &payload,
167            Some(snapshot_lsn.raw()),
168            &self.config.codec,
169        )
170        .map_err(|e| anyhow!("encode managed snapshot failed: {e}"))?;
171        let meta = SnapshotMeta {
172            format_version: info.format_version,
173            node_count: info.node_count,
174            relationship_count: info.relationship_count,
175            wal_lsn: info.wal_lsn,
176        };
177        writer.flush()?;
178        let file = writer.into_inner().map_err(|e| e.into_error())?;
179        sync_file(&file)?;
180        drop(file);
181
182        fs::rename(&tmp, &target)
183            .with_context(|| format!("rename {} to {}", tmp.display(), target.display()))?;
184        sync_dir(&self.config.dir)
185            .with_context(|| format!("sync snapshot dir {}", self.config.dir.display()))?;
186        write_current(&self.config.dir, &target)?;
187        Ok(meta)
188    }
189
190    fn latest_snapshot_path(&self) -> Result<Option<PathBuf>> {
191        let current = self.config.dir.join(CURRENT_FILE);
192        match fs::read_to_string(&current) {
193            Ok(name) => {
194                let name = name.trim();
195                if name.is_empty() {
196                    return Ok(None);
197                }
198                let path = self.config.dir.join(name);
199                if path.exists() {
200                    return Ok(Some(path));
201                }
202            }
203            Err(err) if err.kind() == std::io::ErrorKind::NotFound => {}
204            Err(err) => return Err(err).with_context(|| format!("read {}", current.display())),
205        }
206
207        let latest = snapshot_files(&self.config.dir)?
208            .into_iter()
209            .max_by_key(|(lsn, _)| *lsn)
210            .map(|(_, path)| path);
211        Ok(latest)
212    }
213
214    fn prune_old_snapshots(&self, current_lsn: Lsn) -> Result<()> {
215        let mut snapshots = snapshot_files(&self.config.dir)?;
216        snapshots.retain(|(lsn, _)| *lsn != current_lsn);
217        snapshots.sort_by_key(|(lsn, _)| *lsn);
218        let retain = self.config.keep_old;
219        let remove_count = snapshots.len().saturating_sub(retain);
220        for (_, path) in snapshots.into_iter().take(remove_count) {
221            fs::remove_file(&path)
222                .with_context(|| format!("remove old snapshot {}", path.display()))?;
223        }
224        sync_dir(&self.config.dir)
225            .with_context(|| format!("sync snapshot dir {}", self.config.dir.display()))?;
226        Ok(())
227    }
228}
229
230fn snapshot_path(dir: &Path, lsn: Lsn) -> PathBuf {
231    dir.join(format!(
232        "{SNAPSHOT_PREFIX}{:020}{SNAPSHOT_SUFFIX}",
233        lsn.raw()
234    ))
235}
236
237fn snapshot_files(dir: &Path) -> Result<Vec<(Lsn, PathBuf)>> {
238    let mut out = Vec::new();
239    for entry in
240        fs::read_dir(dir).with_context(|| format!("read snapshot dir {}", dir.display()))?
241    {
242        let entry = entry?;
243        let path = entry.path();
244        let Some(name) = path.file_name().and_then(|name| name.to_str()) else {
245            continue;
246        };
247        let Some(raw) = name
248            .strip_prefix(SNAPSHOT_PREFIX)
249            .and_then(|name| name.strip_suffix(SNAPSHOT_SUFFIX))
250        else {
251            continue;
252        };
253        if let Ok(lsn) = raw.parse::<u64>() {
254            out.push((Lsn::new(lsn), path));
255        }
256    }
257    Ok(out)
258}
259
260fn write_current(dir: &Path, target: &Path) -> Result<()> {
261    let name = target
262        .file_name()
263        .and_then(|name| name.to_str())
264        .ok_or_else(|| {
265            anyhow!(
266                "snapshot path has no portable filename: {}",
267                target.display()
268            )
269        })?;
270    let current = dir.join(CURRENT_FILE);
271    let tmp = tmp_path(&current);
272    let mut file = OpenOptions::new()
273        .write(true)
274        .create(true)
275        .truncate(true)
276        .open(&tmp)
277        .with_context(|| format!("open temp CURRENT {}", tmp.display()))?;
278    writeln!(file, "{name}")?;
279    sync_file(&file)?;
280    drop(file);
281    fs::rename(&tmp, &current)
282        .with_context(|| format!("rename {} to {}", tmp.display(), current.display()))?;
283    sync_dir(dir).with_context(|| format!("sync snapshot dir {}", dir.display()))?;
284    Ok(())
285}
286
287fn tmp_path(path: &Path) -> PathBuf {
288    let mut tmp = path.as_os_str().to_owned();
289    tmp.push(".tmp");
290    PathBuf::from(tmp)
291}