Skip to main content

alopex_server/ops/
backup.rs

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        // 裁定 D15: the data-directory lock names *this* process. Copying it
226        // into a backup would make a restore later stomp on a running server's
227        // live lock file (or, on Windows, fail to delete it at all).
228        if name == ".lifecycle" || alopex_core::lsm::is_lock_file(Path::new(&name)) {
229            continue;
230        }
231        let dest_path = dest.join(name);
232        if file_type.is_dir() {
233            fs::create_dir_all(&dest_path)?;
234            copy_dir_filtered(&entry.path(), &dest_path)?;
235        } else {
236            copy_file_preserving_sparse_zeros(&entry.path(), &dest_path)?;
237        }
238    }
239    Ok(())
240}
241
242fn copy_file_preserving_sparse_zeros(src: &Path, dest: &Path) -> Result<()> {
243    let metadata = fs::metadata(src)?;
244    let mut input = fs::File::open(src)?;
245    let mut output = fs::File::create(dest)?;
246    let mut buffer = vec![0u8; SPARSE_COPY_BUFFER_SIZE];
247
248    loop {
249        let read = input.read(&mut buffer)?;
250        if read == 0 {
251            break;
252        }
253        if buffer[..read].iter().all(|byte| *byte == 0) {
254            output.seek(SeekFrom::Current(read as i64))?;
255        } else {
256            output.write_all(&buffer[..read])?;
257        }
258    }
259
260    output.set_len(metadata.len())?;
261    fs::set_permissions(dest, metadata.permissions())?;
262    Ok(())
263}
264
265pub(crate) fn export_snapshot(source: &Path, dest: &Path) -> Result<()> {
266    let manifest = build_snapshot_manifest(source)?;
267    copy_dir_filtered(source, dest)?;
268    write_snapshot_manifest(dest, &manifest)?;
269    verify_snapshot(dest)
270}
271
272fn verify_snapshot(dest: &Path) -> Result<()> {
273    let manifest = read_snapshot_manifest(dest)?;
274    validate_manifest(dest, &manifest)?;
275
276    let checkpoint_path = dest.join("checkpoint.meta");
277    let meta = load_checkpoint_meta(&checkpoint_path)?;
278    if meta.is_none() {
279        return Err(ServerError::Internal(
280            "checkpoint metadata missing or corrupted".to_string(),
281        ));
282    }
283    let wal_path = dest.join("lsm.wal");
284    if !wal_path.exists() {
285        return Err(ServerError::Internal(
286            "snapshot missing lsm.wal".to_string(),
287        ));
288    }
289    let sst_dir = dest.join("sst");
290    if !sst_dir.exists() {
291        return Err(ServerError::Internal(
292            "snapshot missing sst directory".to_string(),
293        ));
294    }
295
296    let wal_config = LsmKVConfig::default().wal;
297    let mut reader = WalReader::open(&wal_path, wal_config)?;
298    let _ = reader.replay()?;
299
300    for entry in fs::read_dir(&sst_dir)? {
301        let entry = entry?;
302        let path = entry.path();
303        if path.is_file()
304            && path
305                .extension()
306                .and_then(|ext| ext.to_str())
307                .is_some_and(|ext| ext.eq_ignore_ascii_case("sst"))
308        {
309            let _ = SSTableReader::open(&path)?;
310        }
311    }
312    Ok(())
313}
314
315pub(crate) fn verify_snapshot_integrity(dest: &Path) -> Result<()> {
316    verify_snapshot(dest)
317}
318
319fn build_snapshot_manifest(source: &Path) -> Result<SnapshotManifest> {
320    let mut entries = Vec::new();
321    collect_manifest_entries(source, source, &mut entries, true)?;
322    Ok(SnapshotManifest {
323        version: SNAPSHOT_MANIFEST_VERSION,
324        entries,
325    })
326}
327
328fn write_snapshot_manifest(dest: &Path, manifest: &SnapshotManifest) -> Result<()> {
329    let manifest_path = dest.join(SNAPSHOT_MANIFEST_NAME);
330    let payload = serde_json::to_vec_pretty(&manifest)
331        .map_err(|err| ServerError::Internal(format!("manifest encode failed: {err}")))?;
332    fs::write(&manifest_path, payload)?;
333    Ok(())
334}
335
336fn read_snapshot_manifest(dest: &Path) -> Result<SnapshotManifest> {
337    let manifest_path = dest.join(SNAPSHOT_MANIFEST_NAME);
338    let payload = fs::read(&manifest_path)?;
339    let manifest: SnapshotManifest = serde_json::from_slice(&payload)
340        .map_err(|err| ServerError::Internal(format!("manifest decode failed: {err}")))?;
341    if manifest.version != SNAPSHOT_MANIFEST_VERSION {
342        return Err(ServerError::Internal(format!(
343            "unsupported manifest version: {}",
344            manifest.version
345        )));
346    }
347    Ok(manifest)
348}
349
350fn validate_manifest(dest: &Path, manifest: &SnapshotManifest) -> Result<()> {
351    for entry in &manifest.entries {
352        let path = dest.join(&entry.path);
353        let metadata = path.metadata().map_err(|err| {
354            ServerError::Internal(format!("snapshot entry missing {}: {err}", entry.path))
355        })?;
356        if metadata.len() != entry.size {
357            return Err(ServerError::Internal(format!(
358                "snapshot entry size mismatch {}",
359                entry.path
360            )));
361        }
362        let crc = crc32_file(&path)?;
363        if crc != entry.crc32 {
364            return Err(ServerError::Internal(format!(
365                "snapshot entry crc mismatch {}",
366                entry.path
367            )));
368        }
369    }
370    Ok(())
371}
372
373fn collect_manifest_entries(
374    root: &Path,
375    current: &Path,
376    entries: &mut Vec<SnapshotEntry>,
377    skip_lifecycle: bool,
378) -> Result<()> {
379    for entry in fs::read_dir(current)? {
380        let entry = entry?;
381        let path = entry.path();
382        let name = entry.file_name();
383        if skip_lifecycle && name == ".lifecycle" {
384            continue;
385        }
386        if name == SNAPSHOT_MANIFEST_NAME || alopex_core::lsm::is_lock_file(&path) {
387            continue;
388        }
389        let metadata = entry.metadata()?;
390        if metadata.is_dir() {
391            collect_manifest_entries(root, &path, entries, skip_lifecycle)?;
392        } else if metadata.is_file() {
393            let relative = path
394                .strip_prefix(root)
395                .map_err(|err| ServerError::Internal(format!("manifest path error: {err}")))?;
396            let crc32 = crc32_file(&path)?;
397            entries.push(SnapshotEntry {
398                path: relative.to_string_lossy().replace('\\', "/"),
399                size: metadata.len(),
400                crc32,
401            });
402        }
403    }
404    Ok(())
405}
406
407fn crc32_file(path: &Path) -> Result<u32> {
408    let mut file = fs::File::open(path)?;
409    let mut buf = [0u8; 8192];
410    let mut hasher = Hasher::new();
411    loop {
412        let read = std::io::Read::read(&mut file, &mut buf)?;
413        if read == 0 {
414            break;
415        }
416        hasher.update(&buf[..read]);
417    }
418    Ok(hasher.finalize())
419}
420
421#[cfg(test)]
422mod tests {
423    use super::*;
424
425    /// 裁定 D15: the data-directory lock is host-local state. A backup that
426    /// captured it would, on restore, drop another process's pid into a live
427    /// directory — and `clear_data_dir` would have to delete the very file the
428    /// running server holds.
429    #[test]
430    fn backup_snapshot_skips_the_data_directory_lock() {
431        let source = tempfile::tempdir().expect("source tempdir");
432        let dest = tempfile::tempdir().expect("dest tempdir");
433        fs::create_dir_all(source.path().join("sst")).expect("sst dir");
434        fs::write(source.path().join("lsm.wal"), b"wal").expect("wal");
435        fs::write(source.path().join("sst/1.sst"), b"sst").expect("sst");
436        fs::write(source.path().join(".alopex.lock"), b"pid=1").expect("lock");
437        fs::write(source.path().join("mydb.alopex.lock"), b"pid=1").expect("sidecar lock");
438
439        let manifest = build_snapshot_manifest(source.path()).expect("manifest");
440        assert!(
441            manifest
442                .entries
443                .iter()
444                .all(|entry| !entry.path.ends_with(".alopex.lock")),
445            "host-local lock files must not be promised by the snapshot manifest"
446        );
447
448        copy_dir_filtered(source.path(), dest.path()).expect("copy");
449
450        assert!(dest.path().join("lsm.wal").exists());
451        assert!(dest.path().join("sst/1.sst").exists());
452        assert!(
453            !dest.path().join(".alopex.lock").exists(),
454            "a plain-directory lock must not land in a backup"
455        );
456        assert!(
457            !dest.path().join("mydb.alopex.lock").exists(),
458            "a sidecar-shape lock must not land in a backup"
459        );
460    }
461
462    #[test]
463    fn lifecycle_copy_preserves_sparse_file_without_materializing_holes() {
464        // APFS may reserve an 8 MiB allocation extent for the first write.
465        // Use a larger logical hole so the assertion measures materialization
466        // rather than filesystem extent granularity.
467        const FILE_LEN: u64 = 64 * 1024 * 1024;
468
469        let source = tempfile::tempdir().expect("source tempdir");
470        let destination = tempfile::tempdir().expect("destination tempdir");
471        let source_path = source.path().join("sparse.bin");
472        let destination_path = destination.path().join("sparse.bin");
473
474        let mut file = fs::File::create(&source_path).expect("create sparse source");
475        file.write_all(b"head").expect("write sparse head");
476        file.seek(SeekFrom::Start(FILE_LEN - 4))
477            .expect("seek sparse tail");
478        file.write_all(b"tail").expect("write sparse tail");
479        drop(file);
480
481        copy_dir_filtered(source.path(), destination.path()).expect("copy sparse source");
482
483        let mut copied = fs::File::open(&destination_path).expect("open copied file");
484        let mut head = [0u8; 4];
485        copied.read_exact(&mut head).expect("read copied head");
486        assert_eq!(&head, b"head");
487        copied
488            .seek(SeekFrom::Start(FILE_LEN - 4))
489            .expect("seek copied tail");
490        let mut tail = [0u8; 4];
491        copied.read_exact(&mut tail).expect("read copied tail");
492        assert_eq!(&tail, b"tail");
493        assert_eq!(copied.metadata().expect("copied metadata").len(), FILE_LEN);
494
495        #[cfg(unix)]
496        {
497            use std::os::unix::fs::MetadataExt;
498
499            let allocated = copied.metadata().expect("copied metadata").blocks() * 512;
500            assert!(
501                allocated < FILE_LEN / 4,
502                "sparse copy allocated {allocated} bytes for a {FILE_LEN}-byte file"
503            );
504        }
505    }
506}