1use std::collections::HashMap;
2use std::fs;
3use std::io::{Read, Seek, SeekFrom, Write};
4use std::path::{Path, PathBuf};
5use std::sync::{Arc, Mutex};
6use std::time::{SystemTime, UNIX_EPOCH};
7
8use alopex_core::lsm::checkpoint::load_checkpoint_meta;
9use alopex_core::lsm::sstable::SSTableReader;
10use alopex_core::lsm::wal::WalReader;
11use alopex_core::lsm::LsmKVConfig;
12use crc32fast::Hasher;
13use serde::{Deserialize, Serialize};
14use tokio::task;
15use uuid::Uuid;
16
17use crate::error::{Result, ServerError};
18use crate::ops::state::{LifecycleStateManager, OperationState, Progress};
19
20const SNAPSHOT_MANIFEST_NAME: &str = "snapshot.manifest";
21const SNAPSHOT_MANIFEST_VERSION: u32 = 1;
22const SPARSE_COPY_BUFFER_SIZE: usize = 64 * 1024;
23
24#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
25pub struct BackupHandle {
26 pub id: Uuid,
27}
28
29#[derive(Debug, Clone)]
30pub struct BackupMetadata {
31 pub handle: BackupHandle,
32 pub location: PathBuf,
33}
34
35#[derive(Debug, Clone)]
36struct BackupRecord {
37 metadata: BackupMetadata,
38 state: OperationState,
39}
40
41#[derive(Debug, Default)]
42struct BackupRuntime {
43 active: Option<BackupHandle>,
44 history: HashMap<BackupHandle, BackupRecord>,
45 last_location: Option<PathBuf>,
46}
47
48#[derive(Debug, Clone, Serialize, Deserialize)]
49struct SnapshotManifest {
50 version: u32,
51 entries: Vec<SnapshotEntry>,
52}
53
54#[derive(Debug, Clone, Serialize, Deserialize)]
55struct SnapshotEntry {
56 path: String,
57 size: u64,
58 crc32: u32,
59}
60
61#[derive(Clone)]
62pub struct BackupCoordinator {
63 data_dir: PathBuf,
64 state: Arc<LifecycleStateManager>,
65 checkpoint: Arc<dyn Fn() -> Result<()> + Send + Sync>,
66 runtime: Arc<Mutex<BackupRuntime>>,
67}
68
69impl BackupCoordinator {
70 pub fn new(
71 data_dir: PathBuf,
72 state: Arc<LifecycleStateManager>,
73 checkpoint: Arc<dyn Fn() -> Result<()> + Send + Sync>,
74 ) -> Self {
75 Self {
76 data_dir,
77 state,
78 checkpoint,
79 runtime: Arc::new(Mutex::new(BackupRuntime::default())),
80 }
81 }
82
83 pub async fn start_backup(&self) -> Result<BackupHandle> {
84 let mut runtime = self.runtime.lock().expect("backup runtime lock poisoned");
85 if runtime.active.is_some() {
86 return Err(ServerError::Conflict("backup already running".to_string()));
87 }
88
89 let handle = BackupHandle { id: Uuid::new_v4() };
90 let dest = backup_destination(&self.data_dir);
91 fs::create_dir_all(&dest)?;
92 let metadata = BackupMetadata {
93 handle: handle.clone(),
94 location: dest.clone(),
95 };
96 let mut running = OperationState::running();
97 running.set_progress(Progress::percent(0))?;
98 runtime.active = Some(handle.clone());
99 runtime.last_location = Some(dest.clone());
100 runtime.history.insert(
101 handle.clone(),
102 BackupRecord {
103 metadata: metadata.clone(),
104 state: running.clone(),
105 },
106 );
107 self.state.set_backup_state(running);
108
109 let state = self.state.clone();
110 let data_dir = self.data_dir.clone();
111 let runtime = self.runtime.clone();
112 let checkpoint = self.checkpoint.clone();
113 let handle_for_task = handle.clone();
114 task::spawn(async move {
115 let result = task::spawn_blocking(move || run_backup(&data_dir, &dest, checkpoint))
116 .await
117 .map_err(|err| ServerError::Internal(err.to_string()))
118 .and_then(|res| res);
119
120 let mut runtime = runtime.lock().expect("backup runtime lock poisoned");
121 runtime.active = None;
122
123 match result {
124 Ok(()) => {
125 let completed = OperationState::completed(Some(Progress::percent(100)))
126 .unwrap_or_else(|err| OperationState::failed(err.to_string()));
127 if let Some(record) = runtime.history.get_mut(&handle_for_task) {
128 record.state = completed.clone();
129 }
130 state.set_backup_state(completed);
131 }
132 Err(err) => {
133 let failed = OperationState::failed(err.to_string());
134 if let Some(record) = runtime.history.get_mut(&handle_for_task) {
135 record.state = failed.clone();
136 }
137 state.set_backup_state(failed);
138 }
139 }
140 });
141
142 Ok(handle)
143 }
144
145 pub fn status(&self, handle: &BackupHandle) -> Result<OperationState> {
146 let runtime = self.runtime.lock().expect("backup runtime lock poisoned");
147 runtime
148 .history
149 .get(handle)
150 .map(|record| record.state.clone())
151 .ok_or_else(|| ServerError::NotFound("backup handle not found".to_string()))
152 }
153
154 pub fn location(&self, handle: &BackupHandle) -> Result<PathBuf> {
155 let runtime = self.runtime.lock().expect("backup runtime lock poisoned");
156 runtime
157 .history
158 .get(handle)
159 .map(|record| record.metadata.location.clone())
160 .ok_or_else(|| ServerError::NotFound("backup handle not found".to_string()))
161 }
162
163 pub fn latest_location(&self) -> Option<PathBuf> {
164 let runtime = self.runtime.lock().expect("backup runtime lock poisoned");
165 runtime.last_location.clone()
166 }
167}
168
169fn run_backup(
170 data_dir: &Path,
171 dest: &Path,
172 checkpoint: Arc<dyn Fn() -> Result<()> + Send + Sync>,
173) -> Result<()> {
174 if !data_dir.exists() {
175 return Err(ServerError::NotFound(format!(
176 "data directory does not exist: {}",
177 data_dir.display()
178 )));
179 }
180 if !data_dir.is_dir() {
181 return Err(ServerError::BadRequest(format!(
182 "data directory is not a directory: {}",
183 data_dir.display()
184 )));
185 }
186
187 checkpoint().map_err(|err| ServerError::Internal(format!("checkpoint failed: {err}")))?;
188 fs::create_dir_all(dest)?;
189 let manifest = build_snapshot_manifest(data_dir)?;
190 copy_dir_filtered(data_dir, dest)?;
191 write_snapshot_manifest(dest, &manifest)?;
192 verify_snapshot(dest)?;
193 write_latest_marker(&backup_root(data_dir), dest)?;
194 Ok(())
195}
196
197fn backup_destination(data_dir: &Path) -> PathBuf {
198 backup_root(data_dir).join(timestamp_dir())
199}
200
201fn backup_root(data_dir: &Path) -> PathBuf {
202 data_dir.join(".lifecycle").join("backup")
203}
204
205fn timestamp_dir() -> String {
206 let seconds = SystemTime::now()
207 .duration_since(UNIX_EPOCH)
208 .unwrap_or_default()
209 .as_secs();
210 format!("ts-{seconds}")
211}
212
213fn write_latest_marker(root: &Path, latest: &Path) -> Result<()> {
214 fs::create_dir_all(root)?;
215 let marker = root.join("latest");
216 fs::write(marker, latest.display().to_string().as_bytes())?;
217 Ok(())
218}
219
220pub(crate) fn copy_dir_filtered(src: &Path, dest: &Path) -> Result<()> {
221 for entry in fs::read_dir(src)? {
222 let entry = entry?;
223 let file_type = entry.file_type()?;
224 let name = entry.file_name();
225 if name == ".lifecycle" {
226 continue;
227 }
228 let dest_path = dest.join(name);
229 if file_type.is_dir() {
230 fs::create_dir_all(&dest_path)?;
231 copy_dir_filtered(&entry.path(), &dest_path)?;
232 } else {
233 copy_file_preserving_sparse_zeros(&entry.path(), &dest_path)?;
234 }
235 }
236 Ok(())
237}
238
239fn copy_file_preserving_sparse_zeros(src: &Path, dest: &Path) -> Result<()> {
240 let metadata = fs::metadata(src)?;
241 let mut input = fs::File::open(src)?;
242 let mut output = fs::File::create(dest)?;
243 let mut buffer = vec![0u8; SPARSE_COPY_BUFFER_SIZE];
244
245 loop {
246 let read = input.read(&mut buffer)?;
247 if read == 0 {
248 break;
249 }
250 if buffer[..read].iter().all(|byte| *byte == 0) {
251 output.seek(SeekFrom::Current(read as i64))?;
252 } else {
253 output.write_all(&buffer[..read])?;
254 }
255 }
256
257 output.set_len(metadata.len())?;
258 fs::set_permissions(dest, metadata.permissions())?;
259 Ok(())
260}
261
262pub(crate) fn export_snapshot(source: &Path, dest: &Path) -> Result<()> {
263 let manifest = build_snapshot_manifest(source)?;
264 copy_dir_filtered(source, dest)?;
265 write_snapshot_manifest(dest, &manifest)?;
266 verify_snapshot(dest)
267}
268
269fn verify_snapshot(dest: &Path) -> Result<()> {
270 let manifest = read_snapshot_manifest(dest)?;
271 validate_manifest(dest, &manifest)?;
272
273 let checkpoint_path = dest.join("checkpoint.meta");
274 let meta = load_checkpoint_meta(&checkpoint_path)?;
275 if meta.is_none() {
276 return Err(ServerError::Internal(
277 "checkpoint metadata missing or corrupted".to_string(),
278 ));
279 }
280 let wal_path = dest.join("lsm.wal");
281 if !wal_path.exists() {
282 return Err(ServerError::Internal(
283 "snapshot missing lsm.wal".to_string(),
284 ));
285 }
286 let sst_dir = dest.join("sst");
287 if !sst_dir.exists() {
288 return Err(ServerError::Internal(
289 "snapshot missing sst directory".to_string(),
290 ));
291 }
292
293 let wal_config = LsmKVConfig::default().wal;
294 let mut reader = WalReader::open(&wal_path, wal_config)?;
295 let _ = reader.replay()?;
296
297 for entry in fs::read_dir(&sst_dir)? {
298 let entry = entry?;
299 let path = entry.path();
300 if path.is_file()
301 && path
302 .extension()
303 .and_then(|ext| ext.to_str())
304 .is_some_and(|ext| ext.eq_ignore_ascii_case("sst"))
305 {
306 let _ = SSTableReader::open(&path)?;
307 }
308 }
309 Ok(())
310}
311
312pub(crate) fn verify_snapshot_integrity(dest: &Path) -> Result<()> {
313 verify_snapshot(dest)
314}
315
316fn build_snapshot_manifest(source: &Path) -> Result<SnapshotManifest> {
317 let mut entries = Vec::new();
318 collect_manifest_entries(source, source, &mut entries, true)?;
319 Ok(SnapshotManifest {
320 version: SNAPSHOT_MANIFEST_VERSION,
321 entries,
322 })
323}
324
325fn write_snapshot_manifest(dest: &Path, manifest: &SnapshotManifest) -> Result<()> {
326 let manifest_path = dest.join(SNAPSHOT_MANIFEST_NAME);
327 let payload = serde_json::to_vec_pretty(&manifest)
328 .map_err(|err| ServerError::Internal(format!("manifest encode failed: {err}")))?;
329 fs::write(&manifest_path, payload)?;
330 Ok(())
331}
332
333fn read_snapshot_manifest(dest: &Path) -> Result<SnapshotManifest> {
334 let manifest_path = dest.join(SNAPSHOT_MANIFEST_NAME);
335 let payload = fs::read(&manifest_path)?;
336 let manifest: SnapshotManifest = serde_json::from_slice(&payload)
337 .map_err(|err| ServerError::Internal(format!("manifest decode failed: {err}")))?;
338 if manifest.version != SNAPSHOT_MANIFEST_VERSION {
339 return Err(ServerError::Internal(format!(
340 "unsupported manifest version: {}",
341 manifest.version
342 )));
343 }
344 Ok(manifest)
345}
346
347fn validate_manifest(dest: &Path, manifest: &SnapshotManifest) -> Result<()> {
348 for entry in &manifest.entries {
349 let path = dest.join(&entry.path);
350 let metadata = path.metadata().map_err(|err| {
351 ServerError::Internal(format!("snapshot entry missing {}: {err}", entry.path))
352 })?;
353 if metadata.len() != entry.size {
354 return Err(ServerError::Internal(format!(
355 "snapshot entry size mismatch {}",
356 entry.path
357 )));
358 }
359 let crc = crc32_file(&path)?;
360 if crc != entry.crc32 {
361 return Err(ServerError::Internal(format!(
362 "snapshot entry crc mismatch {}",
363 entry.path
364 )));
365 }
366 }
367 Ok(())
368}
369
370fn collect_manifest_entries(
371 root: &Path,
372 current: &Path,
373 entries: &mut Vec<SnapshotEntry>,
374 skip_lifecycle: bool,
375) -> Result<()> {
376 for entry in fs::read_dir(current)? {
377 let entry = entry?;
378 let path = entry.path();
379 let name = entry.file_name();
380 if skip_lifecycle && name == ".lifecycle" {
381 continue;
382 }
383 if name == SNAPSHOT_MANIFEST_NAME {
384 continue;
385 }
386 let metadata = entry.metadata()?;
387 if metadata.is_dir() {
388 collect_manifest_entries(root, &path, entries, skip_lifecycle)?;
389 } else if metadata.is_file() {
390 let relative = path
391 .strip_prefix(root)
392 .map_err(|err| ServerError::Internal(format!("manifest path error: {err}")))?;
393 let crc32 = crc32_file(&path)?;
394 entries.push(SnapshotEntry {
395 path: relative.to_string_lossy().replace('\\', "/"),
396 size: metadata.len(),
397 crc32,
398 });
399 }
400 }
401 Ok(())
402}
403
404fn crc32_file(path: &Path) -> Result<u32> {
405 let mut file = fs::File::open(path)?;
406 let mut buf = [0u8; 8192];
407 let mut hasher = Hasher::new();
408 loop {
409 let read = std::io::Read::read(&mut file, &mut buf)?;
410 if read == 0 {
411 break;
412 }
413 hasher.update(&buf[..read]);
414 }
415 Ok(hasher.finalize())
416}
417
418#[cfg(test)]
419mod tests {
420 use super::*;
421
422 #[test]
423 fn lifecycle_copy_preserves_sparse_file_without_materializing_holes() {
424 const FILE_LEN: u64 = 64 * 1024 * 1024;
428
429 let source = tempfile::tempdir().expect("source tempdir");
430 let destination = tempfile::tempdir().expect("destination tempdir");
431 let source_path = source.path().join("sparse.bin");
432 let destination_path = destination.path().join("sparse.bin");
433
434 let mut file = fs::File::create(&source_path).expect("create sparse source");
435 file.write_all(b"head").expect("write sparse head");
436 file.seek(SeekFrom::Start(FILE_LEN - 4))
437 .expect("seek sparse tail");
438 file.write_all(b"tail").expect("write sparse tail");
439 drop(file);
440
441 copy_dir_filtered(source.path(), destination.path()).expect("copy sparse source");
442
443 let mut copied = fs::File::open(&destination_path).expect("open copied file");
444 let mut head = [0u8; 4];
445 copied.read_exact(&mut head).expect("read copied head");
446 assert_eq!(&head, b"head");
447 copied
448 .seek(SeekFrom::Start(FILE_LEN - 4))
449 .expect("seek copied tail");
450 let mut tail = [0u8; 4];
451 copied.read_exact(&mut tail).expect("read copied tail");
452 assert_eq!(&tail, b"tail");
453 assert_eq!(copied.metadata().expect("copied metadata").len(), FILE_LEN);
454
455 #[cfg(unix)]
456 {
457 use std::os::unix::fs::MetadataExt;
458
459 let allocated = copied.metadata().expect("copied metadata").blocks() * 512;
460 assert!(
461 allocated < FILE_LEN / 4,
462 "sparse copy allocated {allocated} bytes for a {FILE_LEN}-byte file"
463 );
464 }
465 }
466}