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 pub checkpoint_every_commits: Option<u64>,
24 pub keep_old: usize,
26 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 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 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(¤t) {
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(¤t);
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, ¤t)
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}