Skip to main content

alopex_server/ops/
restore.rs

1use std::collections::HashMap;
2use std::fs;
3use std::io::Read;
4use std::path::{Path, PathBuf};
5use std::sync::{Arc, Mutex};
6use std::time::{SystemTime, UNIX_EPOCH};
7
8use serde::{Deserialize, Serialize};
9use sha2::{Digest, Sha256};
10use tokio::task;
11use uuid::Uuid;
12
13use crate::error::{Result, ServerError};
14use crate::ops::backup::{copy_dir_filtered, verify_snapshot_integrity};
15use crate::ops::state::{LifecycleStateManager, Mode, OperationState, Progress, RestoreMetadata};
16
17#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
18pub struct RestoreHandle {
19    pub id: Uuid,
20}
21
22#[derive(Debug, Clone)]
23pub struct RestoreSource {
24    pub path: PathBuf,
25}
26
27#[derive(Debug, Default)]
28struct RestoreRuntime {
29    active: Option<RestoreHandle>,
30    history: HashMap<RestoreHandle, RestoreRecord>,
31}
32
33#[derive(Debug, Clone)]
34pub struct RestoreCoordinator {
35    data_dir: PathBuf,
36    state: Arc<LifecycleStateManager>,
37    runtime: Arc<Mutex<RestoreRuntime>>,
38}
39
40impl RestoreCoordinator {
41    pub fn new(data_dir: PathBuf, state: Arc<LifecycleStateManager>) -> Self {
42        Self {
43            data_dir,
44            state,
45            runtime: Arc::new(Mutex::new(RestoreRuntime::default())),
46        }
47    }
48
49    pub async fn start_restore(&self, source: RestoreSource) -> Result<RestoreHandle> {
50        preflight_restore_source(&source.path)?;
51        let mut runtime = self.runtime.lock().expect("restore runtime lock poisoned");
52        if runtime.active.is_some() {
53            return Err(ServerError::Conflict("restore already running".to_string()));
54        }
55
56        let handle = RestoreHandle { id: Uuid::new_v4() };
57        let mut running = OperationState::running();
58        running.set_progress(Progress::percent(0))?;
59        runtime.active = Some(handle.clone());
60        runtime.history.insert(
61            handle.clone(),
62            RestoreRecord {
63                state: running.clone(),
64                metadata: None,
65            },
66        );
67        self.state.set_restore_state(running);
68        self.state.set_mode(Mode::Maintenance);
69
70        let state = self.state.clone();
71        let data_dir = self.data_dir.clone();
72        let runtime = self.runtime.clone();
73        let handle_for_task = handle.clone();
74        task::spawn(async move {
75            let result = task::spawn_blocking(move || run_restore(&data_dir, &source.path))
76                .await
77                .map_err(|err| ServerError::Internal(err.to_string()))
78                .and_then(|res| res);
79
80            let mut runtime = runtime.lock().expect("restore runtime lock poisoned");
81            runtime.active = None;
82
83            match result {
84                Ok(metadata) => {
85                    state.set_mode(Mode::Normal);
86                    state.set_restore_metadata(Some(metadata.clone()));
87                    let completed = OperationState::completed(Some(Progress::percent(100)))
88                        .unwrap_or_else(|err| OperationState::failed(err.to_string()));
89                    if let Some(record) = runtime.history.get_mut(&handle_for_task) {
90                        record.state = completed.clone();
91                        record.metadata = Some(metadata);
92                    }
93                    state.set_restore_state(completed);
94                }
95                Err(err) => {
96                    state.set_mode(Mode::ReadOnly);
97                    let failed = OperationState::failed(err.to_string());
98                    if let Some(record) = runtime.history.get_mut(&handle_for_task) {
99                        record.state = failed.clone();
100                    }
101                    state.set_restore_state(failed);
102                }
103            }
104        });
105
106        Ok(handle)
107    }
108
109    pub fn status(&self, handle: &RestoreHandle) -> Result<OperationState> {
110        let runtime = self.runtime.lock().expect("restore runtime lock poisoned");
111        runtime
112            .history
113            .get(handle)
114            .map(|record| record.state.clone())
115            .ok_or_else(|| ServerError::NotFound("restore handle not found".to_string()))
116    }
117
118    pub fn metadata(&self, handle: &RestoreHandle) -> Result<Option<RestoreMetadata>> {
119        let runtime = self.runtime.lock().expect("restore runtime lock poisoned");
120        runtime
121            .history
122            .get(handle)
123            .map(|record| record.metadata.clone())
124            .ok_or_else(|| ServerError::NotFound("restore handle not found".to_string()))
125    }
126}
127
128fn run_restore(data_dir: &Path, source: &Path) -> Result<RestoreMetadata> {
129    let backup_dir = restore_backup_dir(data_dir);
130    fs::create_dir_all(&backup_dir)?;
131    copy_dir_filtered(data_dir, &backup_dir)?;
132    clear_data_dir(data_dir)?;
133    copy_dir_filtered(source, data_dir)?;
134    Ok(RestoreMetadata {
135        backup_id: backup_id_from_path(source),
136        location: source.display().to_string(),
137        restored_at_ms: now_ms(),
138        size_bytes: dir_size_bytes(source)?,
139    })
140}
141
142pub fn resolve_default_source(data_dir: &Path) -> Result<PathBuf> {
143    let lifecycle_root = data_dir.join(".lifecycle");
144    let backup_root = lifecycle_root.join("backup");
145    let archive_root = lifecycle_root.join("archive");
146    if let Ok(path) = read_latest_marker(&backup_root) {
147        return Ok(path);
148    }
149    read_latest_marker(&archive_root)
150}
151
152fn validate_source(source: &Path) -> Result<()> {
153    if !source.exists() {
154        return Err(ServerError::NotFound(format!(
155            "restore source does not exist: {}",
156            source.display()
157        )));
158    }
159    if !source.is_dir() {
160        return Err(ServerError::BadRequest(format!(
161            "restore source is not a directory: {}",
162            source.display()
163        )));
164    }
165    Ok(())
166}
167
168fn preflight_restore_source(source: &Path) -> Result<()> {
169    validate_source(source)?;
170    let manifest_path = source.join("snapshot.manifest");
171    if !manifest_path.exists() {
172        return Ok(());
173    }
174    verify_snapshot_integrity(source).map_err(|err| {
175        ServerError::RestoreIntegrityMismatch(format!(
176            "source snapshot failed integrity checks: {err}"
177        ))
178    })
179}
180
181/// Computes a deterministic SHA-256 identity for a restore/upgrade source.
182/// Every directory entry path, type, and file byte is included in sorted
183/// order, while `.lifecycle` remains excluded because it is server-local
184/// operation history rather than source database content.
185pub fn restore_source_fingerprint(source: &Path) -> Result<String> {
186    validate_source(source)?;
187    let mut digest = Sha256::new();
188    hash_restore_tree(source, source, &mut digest)?;
189    Ok(format!("{:x}", digest.finalize()))
190}
191
192fn hash_restore_tree(root: &Path, path: &Path, digest: &mut Sha256) -> Result<()> {
193    let mut entries = fs::read_dir(path)?.collect::<std::result::Result<Vec<_>, _>>()?;
194    entries.sort_by_key(|entry| entry.file_name());
195    for entry in entries {
196        let name = entry.file_name();
197        if path == root && name == ".lifecycle" {
198            continue;
199        }
200        let child = entry.path();
201        let relative = child
202            .strip_prefix(root)
203            .map_err(|error| ServerError::Internal(error.to_string()))?;
204        digest.update(relative.to_string_lossy().as_bytes());
205        let metadata = entry.metadata()?;
206        if metadata.is_dir() {
207            digest.update(b"directory");
208            hash_restore_tree(root, &child, digest)?;
209        } else if metadata.is_file() {
210            digest.update(b"file");
211            let mut file = fs::File::open(&child)?;
212            let mut buffer = [0u8; 8192];
213            loop {
214                let read = file.read(&mut buffer)?;
215                if read == 0 {
216                    break;
217                }
218                digest.update(&buffer[..read]);
219            }
220        }
221    }
222    Ok(())
223}
224
225fn restore_backup_dir(data_dir: &Path) -> PathBuf {
226    data_dir
227        .join(".lifecycle")
228        .join("restore-backup")
229        .join(timestamp_dir())
230}
231
232fn timestamp_dir() -> String {
233    let seconds = SystemTime::now()
234        .duration_since(UNIX_EPOCH)
235        .unwrap_or_default()
236        .as_secs();
237    format!("ts-{seconds}")
238}
239
240fn clear_data_dir(dir: &Path) -> Result<()> {
241    for entry in fs::read_dir(dir)? {
242        let entry = entry?;
243        let name = entry.file_name();
244        // 裁定 D15: restore runs *inside* the server that holds this directory's
245        // lock. Deleting the lock file here would unlink the inode we hold, so
246        // a second process could create a fresh one at the same path and lock
247        // it — two live writers, which is issue #181 all over again.
248        if name == ".lifecycle" || alopex_core::lsm::is_lock_file(Path::new(&name)) {
249            continue;
250        }
251        let path = entry.path();
252        if path.is_dir() {
253            fs::remove_dir_all(&path)?;
254        } else {
255            fs::remove_file(&path)?;
256        }
257    }
258    Ok(())
259}
260
261fn backup_id_from_path(source: &Path) -> String {
262    source
263        .file_name()
264        .map(|name| name.to_string_lossy().into_owned())
265        .unwrap_or_else(|| source.display().to_string())
266}
267
268fn read_latest_marker(root: &Path) -> Result<PathBuf> {
269    let marker = root.join("latest");
270    if !marker.exists() {
271        return Err(ServerError::NotFound(format!(
272            "no snapshots found in {}",
273            root.display()
274        )));
275    }
276    let path = fs::read_to_string(&marker)?;
277    let path = PathBuf::from(path.trim());
278    if !path.exists() {
279        return Err(ServerError::NotFound(format!(
280            "latest snapshot path does not exist: {}",
281            path.display()
282        )));
283    }
284    Ok(path)
285}
286
287fn dir_size_bytes(dir: &Path) -> Result<u64> {
288    let mut size = 0u64;
289    for entry in fs::read_dir(dir)? {
290        let entry = entry?;
291        let name = entry.file_name();
292        if name == ".lifecycle" || alopex_core::lsm::is_lock_file(Path::new(&name)) {
293            continue;
294        }
295        let path = entry.path();
296        let metadata = entry.metadata()?;
297        if metadata.is_dir() {
298            size = size.saturating_add(dir_size_bytes(&path)?);
299        } else {
300            size = size.saturating_add(metadata.len());
301        }
302    }
303    Ok(size)
304}
305
306fn now_ms() -> u64 {
307    SystemTime::now()
308        .duration_since(UNIX_EPOCH)
309        .unwrap_or_default()
310        .as_millis() as u64
311}
312
313#[derive(Debug, Clone)]
314struct RestoreRecord {
315    state: OperationState,
316    metadata: Option<RestoreMetadata>,
317}
318
319#[cfg(test)]
320mod tests {
321    use super::*;
322
323    #[test]
324    fn clear_data_dir_preserves_lifecycle_and_lock_files() {
325        let data_dir = tempfile::tempdir().expect("data tempdir");
326        fs::create_dir_all(data_dir.path().join(".lifecycle/restore")).expect("lifecycle");
327        fs::create_dir_all(data_dir.path().join("sst")).expect("sst dir");
328        fs::write(data_dir.path().join(".alopex.lock"), b"pid=1").expect("plain lock");
329        fs::write(data_dir.path().join("mydb.alopex.lock"), b"pid=1").expect("sidecar lock");
330        fs::write(data_dir.path().join("lsm.wal"), b"wal").expect("wal");
331        fs::write(data_dir.path().join("sst/1.sst"), b"sst").expect("sst");
332
333        assert_eq!(
334            dir_size_bytes(data_dir.path()).expect("size data directory"),
335            6,
336            "restore metadata must count canonical data, not local lock artifacts"
337        );
338        clear_data_dir(data_dir.path()).expect("clear data directory");
339
340        assert!(data_dir.path().join(".lifecycle").exists());
341        assert!(data_dir.path().join(".alopex.lock").exists());
342        assert!(data_dir.path().join("mydb.alopex.lock").exists());
343        assert!(!data_dir.path().join("lsm.wal").exists());
344        assert!(!data_dir.path().join("sst").exists());
345    }
346}