1use std::borrow::Cow;
2use std::collections::{HashMap, HashSet};
3use std::io::Write;
4use std::path::{Path, PathBuf};
5use std::sync::atomic::{AtomicU64, Ordering};
6use std::sync::{Arc, LazyLock, Mutex, RwLock};
7
8use rusqlite::Connection;
9
10use crate::db::backups::BackupRow;
11use crate::error::AftError;
12use sha2::{Digest, Sha256};
13
14pub const DEFAULT_MAX_UNDO_DEPTH: usize = 20;
15#[cfg(test)]
16const MAX_UNDO_DEPTH: usize = DEFAULT_MAX_UNDO_DEPTH;
17const V2_FORMAT_VERSION: &str = "v2";
18const DB_RESTORE_META_VERSION: u32 = 1;
19const MAX_RESTORE_OPERATION_LOCK_RETRIES: usize = 32;
20
21#[cfg(test)]
22type RestoreBeforeLockHook = Box<dyn FnMut(usize) -> bool + Send>;
23
24#[cfg(test)]
25static RESTORE_BEFORE_LOCK_HOOKS: LazyLock<Mutex<HashMap<String, RestoreBeforeLockHook>>> =
26 LazyLock::new(|| Mutex::new(HashMap::new()));
27
28static BACKUP_MAINTENANCE_KEYS: LazyLock<Mutex<HashSet<(PathBuf, Option<String>)>>> =
29 LazyLock::new(|| Mutex::new(HashSet::new()));
30
31#[cfg(test)]
32fn set_restore_before_lock_hook_for_tests(
33 session: &str,
34 hook: impl FnMut(usize) -> bool + Send + 'static,
35) {
36 RESTORE_BEFORE_LOCK_HOOKS
37 .lock()
38 .unwrap()
39 .insert(session.to_string(), Box::new(hook));
40}
41
42#[cfg(test)]
43fn run_restore_before_lock_hook_for_tests(session: &str, attempt: usize) {
44 let mut hooks = RESTORE_BEFORE_LOCK_HOOKS.lock().unwrap();
45 let Some(mut hook) = hooks.remove(session) else {
46 return;
47 };
48 drop(hooks);
49 let keep_hook = hook(attempt);
50 if keep_hook {
51 RESTORE_BEFORE_LOCK_HOOKS
52 .lock()
53 .unwrap()
54 .insert(session.to_string(), hook);
55 }
56}
57
58#[cfg(not(test))]
59fn run_restore_before_lock_hook_for_tests(_session: &str, _attempt: usize) {}
60
61const SCHEMA_VERSION: u32 = 4;
66
67#[derive(Debug, Clone)]
69pub struct BackupEntry {
70 pub backup_id: String,
71 pub content: String,
74 pub content_bytes: Arc<[u8]>,
75 pub timestamp: u64,
76 pub order: u128,
77 pub description: String,
78 pub op_id: Option<String>,
79 pub kind: BackupEntryKind,
80 pub mode: Option<u32>,
81 pub link_target: Option<PathBuf>,
82 pub created_dirs: Vec<PathBuf>,
83}
84
85#[derive(Debug, Clone, Copy, PartialEq, Eq)]
86pub enum BackupEntryKind {
87 Content,
88 Symlink,
89 Tombstone,
90}
91
92#[derive(Debug, Clone)]
100pub(crate) struct CapturedRegularFile {
101 metadata: std::fs::Metadata,
102 bytes: Arc<[u8]>,
103}
104
105impl CapturedRegularFile {
106 pub(crate) fn read(path: &Path) -> std::io::Result<Option<Self>> {
107 let before = std::fs::symlink_metadata(path)?;
108 if !before.is_file() {
109 return Ok(None);
110 }
111
112 let mut bytes: Arc<[u8]> = read_captured_content(path)?.into();
113 let mut metadata = std::fs::symlink_metadata(path)?;
114 if !same_capture_stat(&before, &metadata) {
115 if !metadata.is_file() {
116 return Ok(None);
117 }
118 bytes = read_captured_content(path)?.into();
119 metadata = std::fs::symlink_metadata(path)?;
120 }
121
122 Ok(Some(Self { metadata, bytes }))
123 }
124
125 pub(crate) fn read_text(path: &Path) -> std::io::Result<Option<(Self, String)>> {
126 let Some(capture) = Self::read(path)? else {
127 return Ok(None);
128 };
129 let text = std::str::from_utf8(capture.bytes())
130 .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidData, error))?
131 .to_owned();
132 Ok(Some((capture, text)))
133 }
134
135 pub(crate) fn refresh_if_stale(&mut self, path: &Path) -> std::io::Result<bool> {
136 let current = std::fs::symlink_metadata(path)?;
137 if current.is_file() && same_capture_stat(&self.metadata, ¤t) {
138 return Ok(false);
139 }
140
141 let Some(fresh) = Self::read(path)? else {
142 return Err(std::io::Error::new(
143 std::io::ErrorKind::InvalidInput,
144 "captured regular file is no longer a regular file",
145 ));
146 };
147 *self = fresh;
148 Ok(true)
149 }
150
151 pub(crate) fn bytes(&self) -> &[u8] {
152 &self.bytes
153 }
154
155 pub(crate) fn shared_bytes(&self) -> Arc<[u8]> {
156 Arc::clone(&self.bytes)
157 }
158
159 pub(crate) fn metadata(&self) -> &std::fs::Metadata {
160 &self.metadata
161 }
162}
163
164fn same_capture_stat(left: &std::fs::Metadata, right: &std::fs::Metadata) -> bool {
165 left.len() == right.len() && left.modified().ok() == right.modified().ok()
166}
167
168fn read_captured_content(path: &Path) -> std::io::Result<Vec<u8>> {
169 #[cfg(test)]
170 {
171 let key = std::fs::canonicalize(path).unwrap_or_else(|_| path.to_path_buf());
172 *CAPTURE_READ_COUNTS.lock().unwrap().entry(key).or_default() += 1;
173 }
174 std::fs::read(path)
175}
176
177#[cfg(test)]
178static CAPTURE_READ_COUNTS: LazyLock<Mutex<HashMap<PathBuf, usize>>> =
179 LazyLock::new(|| Mutex::new(HashMap::new()));
180
181#[cfg(test)]
182pub(crate) fn reset_capture_read_count(path: &Path) {
183 let key = std::fs::canonicalize(path).unwrap_or_else(|_| path.to_path_buf());
184 CAPTURE_READ_COUNTS.lock().unwrap().remove(&key);
185}
186
187#[cfg(test)]
188pub(crate) fn capture_read_count(path: &Path) -> usize {
189 let key = std::fs::canonicalize(path).unwrap_or_else(|_| path.to_path_buf());
190 CAPTURE_READ_COUNTS
191 .lock()
192 .unwrap()
193 .get(&key)
194 .copied()
195 .unwrap_or_default()
196}
197
198#[derive(Debug, Clone)]
199struct BackupEntryHead {
200 order: u128,
201 op_id: Option<String>,
202}
203
204impl BackupEntryHead {
205 fn from_entry(entry: &BackupEntry) -> Self {
206 Self {
207 order: entry.order,
208 op_id: entry.op_id.clone(),
209 }
210 }
211
212 fn from_row(row: &BackupRow) -> Self {
213 Self {
214 order: row.order,
215 op_id: row.op_id.clone(),
216 }
217 }
218}
219
220impl BackupEntry {
221 fn to_backup_row(
222 &self,
223 harness: &str,
224 session_id: &str,
225 project_key: &str,
226 file_path: &str,
227 path_hash: &str,
228 backup_path: Option<&str>,
229 ) -> BackupRow {
230 BackupRow {
231 backup_id: self.backup_id.clone(),
232 harness: harness.to_string(),
233 session_id: session_id.to_string(),
234 project_key: project_key.to_string(),
235 op_id: self.op_id.clone(),
236 order: self.order,
237 file_path: file_path.to_string(),
238 path_hash: path_hash.to_string(),
239 backup_path: backup_path.map(str::to_string),
240 kind: match self.kind {
241 BackupEntryKind::Content => "content".to_string(),
242 BackupEntryKind::Symlink => "symlink".to_string(),
243 BackupEntryKind::Tombstone => "tombstone".to_string(),
244 },
245 description: self.description.clone(),
246 created_at: i64::try_from(self.timestamp).unwrap_or(i64::MAX),
247 is_tombstone: matches!(self.kind, BackupEntryKind::Tombstone),
248 restore_meta: Some(restore_metadata_json(self)),
249 }
250 }
251}
252
253impl TryFrom<BackupRow> for BackupEntry {
254 type Error = std::io::Error;
255
256 fn try_from(row: BackupRow) -> Result<Self, Self::Error> {
257 let kind = if row.is_tombstone || row.kind == "tombstone" {
258 BackupEntryKind::Tombstone
259 } else if row.kind == "symlink" {
260 BackupEntryKind::Symlink
261 } else {
262 BackupEntryKind::Content
263 };
264 let backup_path = row.backup_path.clone();
265 let persisted_metadata = row
266 .restore_meta
267 .as_deref()
268 .and_then(restore_metadata_from_json);
269 let restore_metadata = persisted_metadata.or_else(|| {
270 backup_path
271 .as_deref()
272 .and_then(|path| read_entry_disk_metadata(Path::new(path), &row.backup_id))
273 });
274 let content_bytes = match kind {
275 BackupEntryKind::Content | BackupEntryKind::Symlink => {
276 let backup_path = backup_path.ok_or_else(|| {
277 std::io::Error::new(
278 std::io::ErrorKind::NotFound,
279 format!("backup DB row {} has no backup_path", row.backup_id),
280 )
281 })?;
282 std::fs::read(backup_path)?
283 }
284 BackupEntryKind::Tombstone => Vec::new(),
285 };
286 let link_target = if kind == BackupEntryKind::Symlink {
287 restore_metadata
288 .as_ref()
289 .and_then(|metadata| metadata.link_target.clone())
290 .or_else(|| {
291 Some(PathBuf::from(
292 String::from_utf8_lossy(&content_bytes).into_owned(),
293 ))
294 })
295 } else {
296 None
297 };
298 let content = match kind {
299 BackupEntryKind::Content => String::from_utf8_lossy(&content_bytes).into_owned(),
300 BackupEntryKind::Symlink => link_target
301 .as_ref()
302 .map(|target| target.display().to_string())
303 .unwrap_or_default(),
304 BackupEntryKind::Tombstone => String::new(),
305 };
306
307 Ok(BackupEntry {
308 backup_id: row.backup_id,
309 content,
310 content_bytes: content_bytes.into(),
311 timestamp: u64::try_from(row.created_at).unwrap_or_default(),
312 order: row.order,
313 description: row.description,
314 op_id: row.op_id,
315 kind,
316 mode: restore_metadata.as_ref().and_then(|metadata| metadata.mode),
317 link_target,
318 created_dirs: restore_metadata
319 .map(|metadata| metadata.created_dirs)
320 .unwrap_or_default(),
321 })
322 }
323}
324
325#[derive(Debug, Clone)]
326pub struct RestoredOperation {
327 pub op_id: String,
328 pub restored: Vec<RestoredFile>,
329 pub warnings: Vec<String>,
330}
331
332#[derive(Debug, Clone)]
333pub struct RestoredFile {
334 pub path: PathBuf,
335 pub backup_id: String,
336}
337
338#[derive(Debug, Clone, Copy, PartialEq, Eq)]
339pub struct BackupPolicy {
340 pub enabled: bool,
341 pub max_depth: usize,
342 pub max_file_size: Option<u64>,
343}
344
345impl Default for BackupPolicy {
346 fn default() -> Self {
347 Self {
348 enabled: true,
349 max_depth: DEFAULT_MAX_UNDO_DEPTH,
350 max_file_size: None,
351 }
352 }
353}
354
355#[derive(Debug)]
374pub struct BackupStore {
375 entries: HashMap<String, HashMap<PathBuf, Vec<BackupEntry>>>,
377 disk_index: HashMap<String, HashMap<PathBuf, DiskMeta>>,
379 session_meta: HashMap<String, SessionMeta>,
381 counter: AtomicU64,
382 storage_dir: Option<PathBuf>,
383 storage_harness: Option<String>,
384 maintenance_ttl_hours: u32,
385 db_pool: RwLock<Option<Arc<Mutex<Connection>>>>,
386 db_harness: RwLock<Option<String>>,
387 db_project_key: RwLock<Option<String>>,
388 db_mirrored_stacks: RwLock<HashMap<String, HashSet<PathBuf>>>,
391 policy: BackupPolicy,
392 #[cfg(test)]
393 disk_io_count: AtomicU64,
394 #[cfg(test)]
395 fail_next_disk_write: bool,
396}
397
398#[derive(Debug, Clone)]
399struct DiskMeta {
400 dir: PathBuf,
401 count: usize,
402}
403
404enum DbMirrorPlan<'a> {
405 Full,
406 Append {
407 evicted_orders: &'a [u128],
408 new_entry: Option<&'a BackupEntry>,
409 },
410}
411
412struct DbMirrorContext<'a> {
413 harness: &'a str,
414 session: &'a str,
415 project_key: &'a str,
416 file_path: &'a str,
417 path_hash: &'a str,
418}
419
420#[derive(Debug, Clone, Default)]
421struct SessionMeta {
422 last_accessed: u64,
425}
426
427impl BackupStore {
428 pub fn new() -> Self {
429 BackupStore {
430 entries: HashMap::new(),
431 disk_index: HashMap::new(),
432 session_meta: HashMap::new(),
433 counter: AtomicU64::new(0),
434 storage_dir: None,
435 storage_harness: None,
436 maintenance_ttl_hours: 0,
437 db_pool: RwLock::new(None),
438 db_harness: RwLock::new(None),
439 db_project_key: RwLock::new(None),
440 db_mirrored_stacks: RwLock::new(HashMap::new()),
441 policy: BackupPolicy::default(),
442 #[cfg(test)]
443 disk_io_count: AtomicU64::new(0),
444 #[cfg(test)]
445 fail_next_disk_write: false,
446 }
447 }
448
449 pub fn set_policy(&mut self, policy: BackupPolicy) {
450 let old_policy = self.policy;
451 self.policy = policy;
452
453 let failed_disk_prunes = if policy.max_depth < old_policy.max_depth {
454 self.prune_disk_stacks_to_depth(policy.max_depth)
455 } else {
456 HashSet::new()
457 };
458
459 for (session, files) in &mut self.entries {
460 for (key, stack) in files {
461 if failed_disk_prunes.contains(&(session.clone(), key.clone())) {
462 continue;
463 }
464 trim_stack_to_depth(stack, self.policy.max_depth);
465 }
466 }
467 self.entries.retain(|_, files| {
468 files.retain(|_, stack| !stack.is_empty());
469 !files.is_empty()
470 });
471 }
472
473 pub fn policy(&self) -> BackupPolicy {
474 self.policy
475 }
476
477 #[cfg(test)]
478 fn fail_next_disk_write_for_tests(&mut self) {
479 self.fail_next_disk_write = true;
480 }
481
482 pub fn set_db_pool(&self, conn: Arc<Mutex<Connection>>) {
483 if let Ok(mut slot) = self.db_pool.write() {
484 *slot = Some(conn);
485 }
486 self.clear_db_mirror_sync();
487 }
488
489 pub fn clear_db_pool(&self) {
490 if let Ok(mut slot) = self.db_pool.write() {
491 *slot = None;
492 }
493 self.clear_db_mirror_sync();
494 }
495
496 pub fn set_db_harness(&self, harness: crate::harness::Harness) {
497 if let Ok(mut slot) = self.db_harness.write() {
498 *slot = Some(harness.storage_segment());
499 }
500 self.clear_db_mirror_sync();
501 }
502
503 pub fn set_db_project_key(&self, project_key: String) {
504 if let Ok(mut slot) = self.db_project_key.write() {
505 *slot = Some(project_key);
506 }
507 self.clear_db_mirror_sync();
508 }
509
510 pub fn set_storage_dir(&mut self, dir: PathBuf, ttl_hours: u32) {
515 self.set_storage_dir_inner(dir, None, ttl_hours);
516 }
517
518 pub fn set_storage_dir_for_harness(
519 &mut self,
520 dir: PathBuf,
521 harness: crate::harness::Harness,
522 ttl_hours: u32,
523 ) {
524 self.set_storage_dir_inner(dir, Some(harness.storage_segment()), ttl_hours);
525 }
526
527 fn set_storage_dir_inner(&mut self, dir: PathBuf, harness: Option<String>, ttl_hours: u32) {
528 if self.storage_dir.as_ref() == Some(&dir) && self.storage_harness == harness {
529 return;
530 }
531
532 self.storage_dir = Some(dir);
533 self.storage_harness = harness;
534 self.maintenance_ttl_hours = ttl_hours;
535 self.entries.clear();
536 self.disk_index.clear();
537 self.session_meta.clear();
538 self.clear_db_mirror_sync();
539 }
540
541 pub fn run_process_maintenance_once(&mut self) {
548 let Some(storage_dir) = self.storage_dir.clone() else {
549 return;
550 };
551 let key = (storage_dir, self.storage_harness.clone());
552 if !BACKUP_MAINTENANCE_KEYS.lock().unwrap().insert(key) {
553 return;
554 }
555
556 self.record_disk_io_for_tests();
557 self.repair_root_backups_if_needed();
558 self.gc_stale_sessions(self.maintenance_ttl_hours);
559 self.migrate_legacy_layout_if_needed();
560 }
561
562 #[cfg(test)]
563 pub(crate) fn disk_io_count_for_tests(&self) -> u64 {
564 self.disk_io_count.load(Ordering::SeqCst)
565 }
566
567 #[cfg(test)]
568 fn record_disk_io_for_tests(&self) {
569 self.disk_io_count.fetch_add(1, Ordering::SeqCst);
570 }
571
572 #[cfg(not(test))]
573 fn record_disk_io_for_tests(&self) {}
574
575 pub fn snapshot(
577 &mut self,
578 session: &str,
579 path: &Path,
580 description: &str,
581 ) -> Result<Option<String>, AftError> {
582 self.snapshot_with_op(session, path, description, None)
583 }
584
585 pub fn snapshot_with_op(
589 &mut self,
590 session: &str,
591 path: &Path,
592 description: &str,
593 op_id: Option<&str>,
594 ) -> Result<Option<String>, AftError> {
595 if !self.should_snapshot_path(path)? {
596 return Ok(None);
597 }
598 self.run_process_maintenance_once();
599 let key = canonicalize_key(path);
600 let _disk_lock = self.acquire_stack_disk_lock(session, &key)?;
601 self.ensure_stack_hydrated_locked(session, &key)?;
606 let (id, order) = self.next_id_and_order();
607 let entry = backup_entry_from_path(path, id.clone(), order, description, op_id)?;
608
609 self.persist_new_entry_locked(session, &key, entry)?;
610 self.touch_session(session);
611
612 Ok(Some(id))
613 }
614
615 pub(crate) fn snapshot_with_op_from_capture(
618 &mut self,
619 session: &str,
620 path: &Path,
621 description: &str,
622 op_id: Option<&str>,
623 capture: &CapturedRegularFile,
624 ) -> Result<Option<String>, AftError> {
625 if !self.should_snapshot_path(path)? {
626 return Ok(None);
627 }
628 self.run_process_maintenance_once();
629 let key = canonicalize_key(path);
630 let _disk_lock = self.acquire_stack_disk_lock(session, &key)?;
631 self.ensure_stack_hydrated_locked(session, &key)?;
632 let (id, order) = self.next_id_and_order();
633 let entry = backup_entry_from_capture(capture, id.clone(), order, description, op_id);
634
635 self.persist_new_entry_locked(session, &key, entry)?;
636 self.touch_session(session);
637
638 Ok(Some(id))
639 }
640
641 pub fn snapshot_op_tombstone(
644 &mut self,
645 session: &str,
646 op_id: &str,
647 path: &Path,
648 description: &str,
649 ) -> Result<Option<String>, AftError> {
650 if !self.policy.enabled {
651 return Ok(None);
652 }
653 self.run_process_maintenance_once();
654 let key = canonicalize_key(path);
655 let _disk_lock = self.acquire_stack_disk_lock(session, &key)?;
656 self.ensure_stack_hydrated_locked(session, &key)?;
657 let created_dirs = path.parent().map(missing_parent_dirs).unwrap_or_default();
658 let (id, order) = self.next_id_and_order();
659 let entry = BackupEntry {
660 backup_id: id.clone(),
661 content: String::new(),
662 content_bytes: Arc::from([]),
663 timestamp: current_timestamp(),
664 order,
665 description: description.to_string(),
666 op_id: Some(op_id.to_string()),
667 kind: BackupEntryKind::Tombstone,
668 mode: None,
669 link_target: None,
670 created_dirs,
671 };
672
673 self.persist_new_entry_locked(session, &key, entry)?;
674 self.touch_session(session);
675
676 Ok(Some(id))
677 }
678
679 pub fn restore_last_operation(&mut self, session: &str) -> Result<RestoredOperation, AftError> {
682 self.run_process_maintenance_once();
683 let mut candidate_keys = self.restore_operation_candidate_keys(session)?;
684 if candidate_keys.is_empty() {
685 self.load_latest_operation_from_db_or_log(session);
686 candidate_keys = self.restore_operation_candidate_keys(session)?;
687 }
688
689 for attempt in 0..MAX_RESTORE_OPERATION_LOCK_RETRIES {
690 if candidate_keys.is_empty() {
691 return Err(AftError::NoUndoHistory {
692 path: "operation".to_string(),
693 });
694 }
695
696 run_restore_before_lock_hook_for_tests(session, attempt);
697
698 let disk_locks = self.acquire_stack_disk_locks(session, &candidate_keys)?;
699 let locked_keys: HashSet<PathBuf> = candidate_keys.iter().cloned().collect();
700 let current_keys = self.restore_operation_candidate_keys(session)?;
701 let current_key_set: HashSet<PathBuf> = current_keys.iter().cloned().collect();
702 if !current_key_set.is_subset(&locked_keys) {
703 drop(disk_locks);
704 candidate_keys.extend(current_key_set);
705 candidate_keys.sort();
706 candidate_keys.dedup();
707 continue;
708 }
709
710 for key in ¤t_keys {
711 self.load_from_disk_if_needed_locked(session, key)?;
712 }
713
714 if !self.has_in_memory_entries(session) {
715 self.load_latest_operation_from_db_or_log(session);
716 }
717
718 let Some(op_id) = self.latest_operation_id_from_memory(session) else {
719 return Err(AftError::NoUndoHistory {
720 path: "operation".to_string(),
721 });
722 };
723
724 let keys_to_restore = self.operation_keys_for_top_op(session, &op_id);
725 if keys_to_restore.is_empty() {
726 return Err(AftError::NoUndoHistory {
727 path: "operation".to_string(),
728 });
729 }
730 if !keys_to_restore.iter().all(|key| locked_keys.contains(key)) {
731 drop(disk_locks);
732 candidate_keys.extend(keys_to_restore);
733 candidate_keys.sort();
734 candidate_keys.dedup();
735 continue;
736 }
737
738 let mut content_targets = Vec::new();
739 let mut tombstone_targets = Vec::new();
740 for key in &keys_to_restore {
741 let entry = self
742 .entries
743 .get(session)
744 .and_then(|files| files.get(key))
745 .and_then(|stack| stack.last())
746 .cloned()
747 .ok_or_else(|| AftError::NoUndoHistory {
748 path: key.display().to_string(),
749 })?;
750 match entry.kind {
751 BackupEntryKind::Content | BackupEntryKind::Symlink => {
752 let existing_state = capture_path_state(key)?;
753 let warning = self.check_external_modification(session, key, key);
754 content_targets.push((key.clone(), entry, warning, existing_state));
755 }
756 BackupEntryKind::Tombstone => {
757 let existing_state = capture_path_state(key)?;
758 tombstone_targets.push((key.clone(), entry, existing_state));
759 }
760 }
761 }
762
763 let mut created_dirs = Vec::new();
764 for (key, _, _, _) in &content_targets {
765 if let Some(parent) = key.parent() {
766 if !parent.as_os_str().is_empty() {
767 let missing_dirs = missing_parent_dirs(parent);
768 if let Err(e) = std::fs::create_dir_all(parent) {
769 let mut dirs_to_remove = created_dirs;
770 dirs_to_remove.extend(missing_dirs);
771 let rollback_ok = rollback_created_dirs(&dirs_to_remove);
772 return Err(AftError::IoError {
773 path: parent.display().to_string(),
774 message: format!(
775 "{}; restore_last_operation aborted; partial_rollback: {}; rollback_succeeded: {}",
776 e,
777 !rollback_ok,
778 rollback_ok
779 ),
780 });
781 }
782 created_dirs.extend(missing_dirs);
783 }
784 }
785 }
786
787 let mut written = Vec::new();
788 for (key, entry, _, existing_state) in &content_targets {
789 if let Err(e) = restore_entry_to_path(key, entry) {
790 let files_rollback_ok =
791 rollback_transactional_restore(&written, Some((key, existing_state)));
792 let dirs_rollback_ok = rollback_created_dirs(&created_dirs);
793 let rollback_ok = files_rollback_ok && dirs_rollback_ok;
794 return Err(AftError::IoError {
795 path: key.display().to_string(),
796 message: format!(
797 "{}; restore_last_operation aborted; partial_rollback: {}; rollback_succeeded: {}",
798 e,
799 !rollback_ok,
800 rollback_ok
801 ),
802 });
803 }
804 written.push((key.clone(), existing_state.clone()));
805 }
806
807 let mut deleted_tombstones = Vec::new();
808 for (key, _, existing_state) in &tombstone_targets {
809 match remove_tombstone_path(key) {
810 Ok(()) => deleted_tombstones.push((key.clone(), existing_state.clone())),
811 Err(e) => {
812 let files_rollback_ok = rollback_transactional_restore(&written, None);
813 let tombstone_rollback_ok =
814 rollback_deleted_tombstones(&deleted_tombstones);
815 let dirs_rollback_ok = rollback_created_dirs(&created_dirs);
816 let rollback_ok =
817 files_rollback_ok && tombstone_rollback_ok && dirs_rollback_ok;
818 return Err(AftError::IoError {
819 path: key.display().to_string(),
820 message: format!(
821 "{}; restore_last_operation aborted; partial_rollback: {}; rollback_succeeded: {}",
822 e,
823 !rollback_ok,
824 rollback_ok
825 ),
826 });
827 }
828 }
829 }
830 let tombstone_created_dirs = tombstone_targets
831 .iter()
832 .flat_map(|(_, entry, _)| entry.created_dirs.iter().cloned())
833 .collect::<Vec<_>>();
834 remove_created_dirs_best_effort(&tombstone_created_dirs);
835
836 let mut restored = Vec::new();
837 let mut warnings = Vec::new();
838 for (key, entry, warning, _) in content_targets {
839 self.commit_restored_backup_locked(session, &key)?;
840 if let Some(warning) = warning {
841 warnings.push(format!("{}: {}", key.display(), warning));
842 }
843 restored.push(RestoredFile {
844 path: key,
845 backup_id: entry.backup_id,
846 });
847 }
848 for (key, entry, _) in tombstone_targets {
849 self.commit_restored_backup_locked(session, &key)?;
850 restored.push(RestoredFile {
851 path: key,
852 backup_id: entry.backup_id,
853 });
854 }
855 self.touch_session(session);
856 drop(disk_locks);
857
858 return Ok(RestoredOperation {
859 op_id,
860 restored,
861 warnings,
862 });
863 }
864
865 Err(AftError::IoError {
866 path: "operation".to_string(),
867 message: "backup stack changing under concurrent activity; retry".to_string(),
868 })
869 }
870
871 pub fn restore_latest(
874 &mut self,
875 session: &str,
876 path: &Path,
877 ) -> Result<(BackupEntry, Option<String>), AftError> {
878 self.run_process_maintenance_once();
879 let key = canonicalize_key(path);
880 let _disk_lock = self.acquire_stack_disk_lock(session, &key)?;
881
882 match self.read_stack_from_disk_unlocked(session, &key) {
883 Ok(Some(entries)) if !entries.is_empty() => {
884 self.update_counter_from_entries(&entries);
885 self.entries
886 .entry(session.to_string())
887 .or_default()
888 .insert(key.to_path_buf(), entries);
889 }
890 Ok(_) => {
891 if self.session_dir(session).is_some() {
892 self.restore_in_memory_stack(session, &key, None);
893 }
894 }
895 Err(error) => {
896 return Err(AftError::IoError {
897 path: key.display().to_string(),
898 message: error,
899 });
900 }
901 }
902
903 if self
904 .entries
905 .get(session)
906 .and_then(|s| s.get(&key))
907 .is_none_or(|s| s.is_empty())
908 {
909 match self.load_from_db_if_present(session, &key) {
910 Some(Ok(true)) => {}
911 Some(Ok(false)) => {
912 crate::slog_info!(
913 "backup DB miss for session {} path {}; disk meta is authoritative",
914 session,
915 key.display()
916 );
917 }
918 Some(Err(error)) => {
919 crate::slog_warn!(
920 "backup DB lookup failed for session {} path {}: {}",
921 session,
922 key.display(),
923 error
924 );
925 }
926 None => {
927 crate::slog_info!(
928 "backup DB unavailable for session {} path {}",
929 session,
930 key.display()
931 );
932 }
933 }
934 }
935
936 let in_memory = self
938 .entries
939 .get(session)
940 .and_then(|s| s.get(&key))
941 .map_or(false, |s| !s.is_empty());
942 if in_memory {
943 let warning = self.check_external_modification(session, &key, path);
944 let result = self
945 .do_restore_locked(session, &key, path)
946 .map(|(entry, _)| (entry, warning));
947 if result.is_ok() {
948 self.touch_session(session);
949 }
950 return result;
951 }
952
953 Err(AftError::NoUndoHistory {
954 path: path.display().to_string(),
955 })
956 }
957
958 pub fn history(&self, session: &str, path: &Path) -> Vec<BackupEntry> {
960 let key = canonicalize_key(path);
961 let _disk_lock = match self.acquire_stack_disk_lock(session, &key) {
962 Ok(lock) => lock,
963 Err(error) => {
964 crate::slog_warn!(
965 "backup disk read lock failed for {}: {}",
966 key.display(),
967 error
968 );
969 return Vec::new();
970 }
971 };
972
973 match self.read_stack_from_disk_unlocked(session, &key) {
974 Ok(Some(stack)) if !stack.is_empty() => return stack,
975 Ok(_) => {}
976 Err(error) => {
977 crate::slog_warn!("backup disk read failed for {}: {}", key.display(), error);
978 return Vec::new();
979 }
980 }
981
982 if let Some(stack) = self.entries.get(session).and_then(|s| s.get(&key)).cloned() {
983 if !stack.is_empty() {
984 return stack;
985 }
986 }
987
988 match self.read_stack_from_db(session, &key) {
989 Some(Ok(stack)) if !stack.is_empty() => stack,
990 Some(Ok(_)) => Vec::new(),
991 Some(Err(error)) => {
992 crate::slog_warn!(
993 "backup history DB lookup failed for session {} path {}: {}",
994 session,
995 key.display(),
996 error
997 );
998 Vec::new()
999 }
1000 None => Vec::new(),
1001 }
1002 }
1003
1004 pub fn disk_history_count(&self, session: &str, path: &Path) -> usize {
1006 let key = canonicalize_key(path);
1007 self.disk_index
1008 .get(session)
1009 .and_then(|s| s.get(&key))
1010 .map(|m| m.count)
1011 .unwrap_or(0)
1012 }
1013
1014 pub fn tracked_files(&self, session: &str) -> Vec<PathBuf> {
1017 let mut files: std::collections::HashSet<PathBuf> = self
1018 .entries
1019 .get(session)
1020 .map(|s| s.keys().cloned().collect())
1021 .unwrap_or_default();
1022 if let Some(disk) = self.disk_index.get(session) {
1023 for key in disk.keys() {
1024 files.insert(key.clone());
1025 }
1026 }
1027 files.into_iter().collect()
1028 }
1029
1030 pub fn preview_latest_path(&self, session: &str, path: &Path) -> Result<PathBuf, AftError> {
1035 let key = canonicalize_key(path);
1036 if self.latest_head_for_key(session, &key).is_some() {
1037 Ok(key)
1038 } else {
1039 Err(AftError::NoUndoHistory {
1040 path: path.display().to_string(),
1041 })
1042 }
1043 }
1044
1045 pub fn preview_last_operation_paths(&self, session: &str) -> Result<Vec<PathBuf>, AftError> {
1051 let mut heads_by_path: HashMap<PathBuf, BackupEntryHead> = self
1052 .entries
1053 .get(session)
1054 .map(|files| {
1055 files
1056 .iter()
1057 .filter_map(|(key, stack)| {
1058 stack
1059 .last()
1060 .map(|entry| (key.clone(), BackupEntryHead::from_entry(entry)))
1061 })
1062 .collect()
1063 })
1064 .unwrap_or_default();
1065
1066 match self.read_latest_operation_heads_from_db(session) {
1067 Some(Ok(db_heads)) if !db_heads.is_empty() => {
1068 for (key, head) in db_heads {
1069 heads_by_path.insert(key, head);
1070 }
1071 self.merge_disk_stack_heads(session, &mut heads_by_path);
1072 }
1073 Some(Ok(_)) => {
1074 crate::slog_info!(
1075 "backup latest operation preview DB miss for session {}; falling back to disk",
1076 session
1077 );
1078 self.merge_disk_stack_heads(session, &mut heads_by_path);
1079 }
1080 Some(Err(error)) => {
1081 crate::slog_warn!(
1082 "backup latest operation preview DB lookup failed for session {}; falling back to disk: {}",
1083 session,
1084 error
1085 );
1086 self.merge_disk_stack_heads(session, &mut heads_by_path);
1087 }
1088 None => {
1089 crate::slog_info!(
1090 "backup latest operation preview DB unavailable for session {}; falling back to disk",
1091 session
1092 );
1093 self.merge_disk_stack_heads(session, &mut heads_by_path);
1094 }
1095 }
1096
1097 let mut latest: Option<(u128, String)> = None;
1098 for head in heads_by_path.values() {
1099 if let Some(op_id) = &head.op_id {
1100 if latest
1101 .as_ref()
1102 .map_or(true, |(latest_order, _)| head.order > *latest_order)
1103 {
1104 latest = Some((head.order, op_id.clone()));
1105 }
1106 }
1107 }
1108
1109 let Some((_, op_id)) = latest else {
1110 return Err(AftError::NoUndoHistory {
1111 path: "operation".to_string(),
1112 });
1113 };
1114
1115 let mut paths: Vec<PathBuf> = heads_by_path
1116 .into_iter()
1117 .filter_map(|(key, head)| {
1118 (head.op_id.as_deref() == Some(op_id.as_str())).then_some(key)
1119 })
1120 .collect();
1121 paths.sort();
1122
1123 if paths.is_empty() {
1124 Err(AftError::NoUndoHistory {
1125 path: "operation".to_string(),
1126 })
1127 } else {
1128 Ok(paths)
1129 }
1130 }
1131
1132 pub fn sessions_with_backups(&self) -> Vec<String> {
1135 let mut sessions: std::collections::HashSet<String> =
1136 self.entries.keys().cloned().collect();
1137 for s in self.disk_index.keys() {
1138 sessions.insert(s.clone());
1139 }
1140 sessions.into_iter().collect()
1141 }
1142
1143 pub fn total_disk_bytes(&self) -> u64 {
1146 let mut total = 0u64;
1147 for session_dirs in self.disk_index.values() {
1148 for meta in session_dirs.values() {
1149 if let Ok(read_dir) = std::fs::read_dir(&meta.dir) {
1150 for entry in read_dir.flatten() {
1151 if let Ok(m) = entry.metadata() {
1152 if m.is_file() {
1153 total += m.len();
1154 }
1155 }
1156 }
1157 }
1158 }
1159 }
1160 total
1161 }
1162
1163 fn next_id_and_order(&self) -> (String, u128) {
1164 let n = self.counter.fetch_add(1, Ordering::Relaxed);
1165 let order = ((current_timestamp_nanos() as u128) << 32) | u128::from(n);
1166 (format!("backup-{}", n), order)
1167 }
1168
1169 fn db_pool_and_harness(&self) -> Option<(Arc<Mutex<Connection>>, String)> {
1170 let pool = self.db_pool.read().ok().and_then(|slot| slot.clone())?;
1171 let harness = self.db_harness.read().ok().and_then(|slot| slot.clone())?;
1172 Some((pool, harness))
1173 }
1174
1175 fn clear_db_mirror_sync(&self) {
1176 if let Ok(mut synced) = self.db_mirrored_stacks.write() {
1177 synced.clear();
1178 }
1179 }
1180
1181 fn db_mirror_is_synced(&self, session: &str, key: &Path) -> bool {
1182 self.db_mirrored_stacks
1183 .read()
1184 .is_ok_and(|synced| synced.get(session).is_some_and(|keys| keys.contains(key)))
1185 }
1186
1187 fn set_db_mirror_synced(&self, session: &str, key: &Path, is_synced: bool) {
1188 if let Ok(mut synced) = self.db_mirrored_stacks.write() {
1189 if is_synced {
1190 synced
1191 .entry(session.to_string())
1192 .or_default()
1193 .insert(key.to_path_buf());
1194 } else if let Some(keys) = synced.get_mut(session) {
1195 keys.remove(key);
1196 if keys.is_empty() {
1197 synced.remove(session);
1198 }
1199 }
1200 }
1201 }
1202
1203 fn latest_head_for_key(&self, session: &str, key: &Path) -> Option<BackupEntryHead> {
1204 self.entries
1205 .get(session)
1206 .and_then(|files| files.get(key))
1207 .and_then(|stack| stack.last())
1208 .map(BackupEntryHead::from_entry)
1209 .or_else(|| {
1210 self.read_stack_heads_from_disk(session, key)
1211 .and_then(|stack| stack.last().cloned())
1212 })
1213 .or_else(|| match self.read_stack_heads_from_db(session, key) {
1214 Some(Ok(stack)) if !stack.is_empty() => stack.last().cloned(),
1215 Some(Err(error)) => {
1216 crate::slog_warn!(
1217 "backup preview DB lookup failed for session {} path {}: {}",
1218 session,
1219 key.display(),
1220 error
1221 );
1222 None
1223 }
1224 _ => None,
1225 })
1226 }
1227
1228 fn merge_disk_stack_heads(
1229 &self,
1230 session: &str,
1231 heads_by_path: &mut HashMap<PathBuf, BackupEntryHead>,
1232 ) {
1233 let disk_keys: Vec<PathBuf> = self
1234 .disk_index
1235 .get(session)
1236 .map(|files| files.keys().cloned().collect())
1237 .unwrap_or_default();
1238 for key in disk_keys {
1239 if let Some(head) = self
1240 .read_stack_heads_from_disk(session, &key)
1241 .and_then(|stack| stack.last().cloned())
1242 {
1243 heads_by_path.insert(key, head);
1244 }
1245 }
1246 }
1247
1248 fn read_stack_heads_from_db(
1249 &self,
1250 session: &str,
1251 key: &Path,
1252 ) -> Option<Result<Vec<BackupEntryHead>, String>> {
1253 let (pool, harness) = self.db_pool_and_harness()?;
1254 let conn = match pool.lock() {
1255 Ok(conn) => conn,
1256 Err(_) => return Some(Err("db mutex poisoned".to_string())),
1257 };
1258 let path_hash = Self::path_hash(key);
1259 Some(
1260 crate::db::backups::list_backups(&conn, &harness, session, &path_hash)
1261 .map_err(|error| error.to_string())
1262 .map(|rows| {
1263 rows.iter()
1264 .map(BackupEntryHead::from_row)
1265 .collect::<Vec<_>>()
1266 }),
1267 )
1268 }
1269
1270 fn read_latest_operation_heads_from_db(
1271 &self,
1272 session: &str,
1273 ) -> Option<Result<HashMap<PathBuf, BackupEntryHead>, String>> {
1274 let (pool, harness) = self.db_pool_and_harness()?;
1275 let conn = match pool.lock() {
1276 Ok(conn) => conn,
1277 Err(_) => return Some(Err("db mutex poisoned".to_string())),
1278 };
1279 let latest = match crate::db::backups::get_latest_operation_backup(&conn, &harness, session)
1280 {
1281 Ok(Some(row)) => row,
1282 Ok(None) => return Some(Ok(HashMap::new())),
1283 Err(error) => return Some(Err(error.to_string())),
1284 };
1285 let Some(op_id) = latest.op_id else {
1286 return Some(Ok(HashMap::new()));
1287 };
1288 let rows = match crate::db::backups::list_backups_by_op(&conn, &harness, session, &op_id) {
1289 Ok(rows) => rows,
1290 Err(error) => return Some(Err(error.to_string())),
1291 };
1292 if rows.is_empty() {
1293 return Some(Ok(HashMap::new()));
1294 }
1295 let path_hashes: std::collections::HashSet<String> =
1296 rows.into_iter().map(|row| row.path_hash).collect();
1297 drop(conn);
1298
1299 let mut heads = HashMap::new();
1300 for path_hash in path_hashes {
1301 let conn = match pool.lock() {
1302 Ok(conn) => conn,
1303 Err(_) => return Some(Err("db mutex poisoned".to_string())),
1304 };
1305 let rows = match crate::db::backups::list_backups(&conn, &harness, session, &path_hash)
1306 {
1307 Ok(rows) => rows,
1308 Err(error) => return Some(Err(error.to_string())),
1309 };
1310 drop(conn);
1311
1312 let Some(file_path) = rows.first().map(|row| row.file_path.clone()) else {
1313 continue;
1314 };
1315 let Some(head) = rows.last().map(BackupEntryHead::from_row) else {
1316 continue;
1317 };
1318 heads.insert(PathBuf::from(file_path), head);
1319 }
1320
1321 Some(Ok(heads))
1322 }
1323
1324 fn read_stack_from_db(
1325 &self,
1326 session: &str,
1327 key: &Path,
1328 ) -> Option<Result<Vec<BackupEntry>, String>> {
1329 let (pool, harness) = self.db_pool_and_harness()?;
1330 let conn = match pool.lock() {
1331 Ok(conn) => conn,
1332 Err(_) => return Some(Err("db mutex poisoned".to_string())),
1333 };
1334 let path_hash = Self::path_hash(key);
1335 Some(
1336 crate::db::backups::list_backups(&conn, &harness, session, &path_hash)
1337 .map_err(|error| error.to_string())
1338 .and_then(|rows| {
1339 rows.into_iter()
1340 .map(|row| self.backup_entry_from_db_row(row))
1341 .collect::<Result<Vec<_>, _>>()
1342 .map_err(|error| error.to_string())
1343 }),
1344 )
1345 }
1346
1347 fn load_from_db_if_present(
1348 &mut self,
1349 session: &str,
1350 key: &Path,
1351 ) -> Option<Result<bool, String>> {
1352 match self.read_stack_from_db(session, key) {
1353 Some(Ok(stack)) if !stack.is_empty() => {
1354 self.update_counter_from_entries(&stack);
1355 self.entries
1356 .entry(session.to_string())
1357 .or_default()
1358 .insert(key.to_path_buf(), stack);
1359 Some(Ok(true))
1360 }
1361 Some(Ok(_)) => Some(Ok(false)),
1362 Some(Err(error)) => Some(Err(error)),
1363 None => None,
1364 }
1365 }
1366
1367 fn load_latest_operation_from_db(&mut self, session: &str) -> Option<Result<bool, String>> {
1368 let (pool, harness) = self.db_pool_and_harness()?;
1369 let conn = match pool.lock() {
1370 Ok(conn) => conn,
1371 Err(_) => return Some(Err("db mutex poisoned".to_string())),
1372 };
1373 let latest = match crate::db::backups::get_latest_operation_backup(&conn, &harness, session)
1374 {
1375 Ok(Some(row)) => row,
1376 Ok(None) => return Some(Ok(false)),
1377 Err(error) => return Some(Err(error.to_string())),
1378 };
1379 let Some(op_id) = latest.op_id else {
1380 return Some(Ok(false));
1381 };
1382 let rows = match crate::db::backups::list_backups_by_op(&conn, &harness, session, &op_id) {
1383 Ok(rows) => rows,
1384 Err(error) => return Some(Err(error.to_string())),
1385 };
1386 if rows.is_empty() {
1387 return Some(Ok(false));
1388 }
1389 let path_hashes: std::collections::HashSet<String> =
1390 rows.into_iter().map(|row| row.path_hash).collect();
1391 drop(conn);
1392
1393 let mut loaded_any = false;
1394 for path_hash in path_hashes {
1395 let conn = match pool.lock() {
1396 Ok(conn) => conn,
1397 Err(_) => return Some(Err("db mutex poisoned".to_string())),
1398 };
1399 let loaded =
1400 match crate::db::backups::list_backups(&conn, &harness, session, &path_hash) {
1401 Ok(rows) => {
1402 let file_path = rows.first().map(|row| row.file_path.clone());
1403 rows.into_iter()
1404 .map(|row| self.backup_entry_from_db_row(row))
1405 .collect::<Result<Vec<_>, _>>()
1406 .map(|stack| (file_path, stack))
1407 .map_err(|error| error.to_string())
1408 }
1409 Err(error) => Err(error.to_string()),
1410 };
1411 drop(conn);
1412 let (file_path, stack) = match loaded {
1413 Ok((file_path, stack)) if !stack.is_empty() => (file_path, stack),
1414 Ok(_) => continue,
1415 Err(error) => return Some(Err(error)),
1416 };
1417 let Some(file_path) = file_path else {
1418 return Some(Err(format!(
1419 "backup DB rows for path hash {path_hash} have no file path"
1420 )));
1421 };
1422 let key = PathBuf::from(file_path);
1423 self.update_counter_from_entries(&stack);
1424 self.entries
1425 .entry(session.to_string())
1426 .or_default()
1427 .insert(key, stack);
1428 loaded_any = true;
1429 }
1430
1431 Some(Ok(loaded_any))
1432 }
1433
1434 fn update_counter_from_entries(&self, entries: &[BackupEntry]) {
1435 if let Some(next_counter) = entries
1436 .iter()
1437 .filter_map(|entry| backup_sequence(&entry.backup_id))
1438 .max()
1439 .and_then(|max| max.checked_add(1))
1440 {
1441 let _ = self
1442 .counter
1443 .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |current| {
1444 (current < next_counter).then_some(next_counter)
1445 });
1446 }
1447 }
1448
1449 fn persist_new_entry_locked(
1450 &mut self,
1451 session: &str,
1452 key: &Path,
1453 entry: BackupEntry,
1454 ) -> Result<(), AftError> {
1455 let max_depth = self.policy.max_depth;
1456 let mut stack = self
1460 .entries
1461 .get_mut(session)
1462 .and_then(|files| files.remove(key))
1463 .unwrap_or_default();
1464 let mut evicted = drain_stack_to_depth(&mut stack, max_depth.saturating_sub(1));
1465 if max_depth > 0 {
1466 stack.push(entry);
1467 }
1468
1469 let evicted_orders = evicted.iter().map(|entry| entry.order).collect::<Vec<_>>();
1473 let new_entry = (max_depth > 0).then(|| stack.last().expect("new backup entry"));
1474 if let Err(error) = self.write_appended_snapshot_to_disk_locked(
1475 session,
1476 key,
1477 &stack,
1478 &evicted_orders,
1479 new_entry,
1480 ) {
1481 if max_depth > 0 {
1482 stack.pop();
1483 }
1484 evicted.append(&mut stack);
1485 self.restore_in_memory_stack(session, key, Some(evicted));
1486 return Err(error);
1487 }
1488
1489 self.entries
1490 .entry(session.to_string())
1491 .or_default()
1492 .insert(key.to_path_buf(), stack);
1493 Ok(())
1494 }
1495
1496 fn restore_in_memory_stack(
1497 &mut self,
1498 session: &str,
1499 key: &Path,
1500 stack: Option<Vec<BackupEntry>>,
1501 ) {
1502 match stack {
1503 Some(stack) if !stack.is_empty() => {
1504 self.entries
1505 .entry(session.to_string())
1506 .or_default()
1507 .insert(key.to_path_buf(), stack);
1508 }
1509 _ => {
1510 if let Some(files) = self.entries.get_mut(session) {
1511 files.remove(key);
1512 if files.is_empty() {
1513 self.entries.remove(session);
1514 }
1515 }
1516 }
1517 }
1518 }
1519
1520 fn has_in_memory_entries(&self, session: &str) -> bool {
1521 self.entries
1522 .get(session)
1523 .is_some_and(|files| files.values().any(|stack| !stack.is_empty()))
1524 }
1525
1526 fn latest_operation_id_from_memory(&self, session: &str) -> Option<String> {
1527 let mut latest: Option<(u128, String)> = None;
1528 if let Some(files) = self.entries.get(session) {
1529 for stack in files.values() {
1530 if let Some(entry) = stack.last() {
1531 if let Some(op_id) = &entry.op_id {
1532 if latest
1533 .as_ref()
1534 .is_none_or(|(latest_order, _)| entry.order > *latest_order)
1535 {
1536 latest = Some((entry.order, op_id.clone()));
1537 }
1538 }
1539 }
1540 }
1541 }
1542 latest.map(|(_, op_id)| op_id)
1543 }
1544
1545 fn operation_keys_for_top_op(&self, session: &str, op_id: &str) -> Vec<PathBuf> {
1546 let mut keys: Vec<PathBuf> = self
1547 .entries
1548 .get(session)
1549 .map(|files| {
1550 files
1551 .iter()
1552 .filter_map(|(key, stack)| {
1553 stack.last().and_then(|entry| {
1554 (entry.op_id.as_deref() == Some(op_id)).then(|| key.clone())
1555 })
1556 })
1557 .collect()
1558 })
1559 .unwrap_or_default();
1560 keys.sort();
1561 keys
1562 }
1563
1564 fn load_latest_operation_from_db_or_log(&mut self, session: &str) {
1565 match self.load_latest_operation_from_db(session) {
1566 Some(Ok(true)) => {}
1567 Some(Ok(false)) => {
1568 crate::slog_info!(
1569 "backup latest operation DB miss for session {}; disk meta is authoritative",
1570 session
1571 );
1572 }
1573 Some(Err(error)) => {
1574 crate::slog_warn!(
1575 "backup latest operation DB lookup failed for session {}: {}",
1576 session,
1577 error
1578 );
1579 }
1580 None => {
1581 crate::slog_info!(
1582 "backup latest operation DB unavailable for session {}",
1583 session
1584 );
1585 }
1586 }
1587 }
1588
1589 fn resolve_db_backup_row_path(&self, mut row: BackupRow) -> BackupRow {
1590 if let Some(backup_path) = row.backup_path.clone() {
1591 let path = PathBuf::from(&backup_path);
1592 if path.is_relative() {
1593 if let Some(session_dir) = self.session_dir(&row.session_id) {
1594 row.backup_path = Some(
1595 session_dir
1596 .join(&row.path_hash)
1597 .join(path)
1598 .display()
1599 .to_string(),
1600 );
1601 }
1602 }
1603 }
1604 row
1605 }
1606
1607 fn backup_entry_from_db_row(&self, row: BackupRow) -> Result<BackupEntry, std::io::Error> {
1608 BackupEntry::try_from(self.resolve_db_backup_row_path(row))
1609 }
1610
1611 pub fn discard_operation_entries(&mut self, session: &str, op_id: &str) {
1612 let keys: Vec<PathBuf> = self
1613 .entries
1614 .get(session)
1615 .map(|files| files.keys().cloned().collect())
1616 .unwrap_or_default();
1617
1618 for key in keys {
1619 let mut remove_key = false;
1620 let mut remaining_stack = None;
1621 if let Some(session_entries) = self.entries.get_mut(session) {
1622 if let Some(stack) = session_entries.get_mut(&key) {
1623 while stack
1624 .last()
1625 .is_some_and(|entry| entry.op_id.as_deref() == Some(op_id))
1626 {
1627 stack.pop();
1628 }
1629 if stack.is_empty() {
1630 remove_key = true;
1631 } else {
1632 remaining_stack = Some(stack.clone());
1633 }
1634 }
1635 if remove_key {
1636 session_entries.remove(&key);
1637 }
1638 }
1639
1640 if remove_key {
1641 if let Err(error) = self.remove_disk_backups(session, &key) {
1642 crate::slog_warn!(
1643 "failed to remove backup stack for {} during operation discard: {}",
1644 key.display(),
1645 error
1646 );
1647 }
1648 } else if let Some(stack) = remaining_stack {
1649 if let Err(error) = self.write_snapshot_to_disk(session, &key, &stack) {
1650 crate::slog_warn!(
1651 "failed to persist backup stack for {} during operation discard: {}",
1652 key.display(),
1653 error
1654 );
1655 }
1656 }
1657 }
1658
1659 if self
1660 .entries
1661 .get(session)
1662 .is_some_and(|session_entries| session_entries.is_empty())
1663 {
1664 self.entries.remove(session);
1665 }
1666 }
1667
1668 pub(crate) fn discard_latest_operation_entry_for_path(
1669 &mut self,
1670 session: &str,
1671 op_id: &str,
1672 path: &Path,
1673 ) {
1674 let key = canonicalize_key(path);
1675 let mut remove_key = false;
1676 let mut remaining_stack = None;
1677
1678 if let Some(session_entries) = self.entries.get_mut(session) {
1679 if let Some(stack) = session_entries.get_mut(&key) {
1680 if stack
1681 .last()
1682 .is_some_and(|entry| entry.op_id.as_deref() == Some(op_id))
1683 {
1684 stack.pop();
1685 if stack.is_empty() {
1686 remove_key = true;
1687 } else {
1688 remaining_stack = Some(stack.clone());
1689 }
1690 }
1691 }
1692 if remove_key {
1693 session_entries.remove(&key);
1694 }
1695 }
1696
1697 if remove_key {
1698 if let Err(error) = self.remove_disk_backups(session, &key) {
1699 crate::slog_warn!(
1700 "failed to remove backup stack for {} during single-entry discard: {}",
1701 key.display(),
1702 error
1703 );
1704 }
1705 } else if let Some(stack) = remaining_stack {
1706 if let Err(error) = self.write_snapshot_to_disk(session, &key, &stack) {
1707 crate::slog_warn!(
1708 "failed to persist backup stack for {} during single-entry discard: {}",
1709 key.display(),
1710 error
1711 );
1712 }
1713 }
1714
1715 if self
1716 .entries
1717 .get(session)
1718 .is_some_and(|session_entries| session_entries.is_empty())
1719 {
1720 self.entries.remove(session);
1721 }
1722 }
1723
1724 fn touch_session(&mut self, session: &str) {
1725 let now = current_timestamp();
1726 self.session_meta
1727 .entry(session.to_string())
1728 .or_default()
1729 .last_accessed = now;
1730 self.write_session_marker(session, now);
1731 }
1732
1733 fn do_restore_locked(
1736 &mut self,
1737 session: &str,
1738 key: &Path,
1739 path: &Path,
1740 ) -> Result<(BackupEntry, Option<String>), AftError> {
1741 let session_entries =
1742 self.entries
1743 .get_mut(session)
1744 .ok_or_else(|| AftError::NoUndoHistory {
1745 path: path.display().to_string(),
1746 })?;
1747 let stack = session_entries
1748 .get_mut(key)
1749 .ok_or_else(|| AftError::NoUndoHistory {
1750 path: path.display().to_string(),
1751 })?;
1752
1753 let entry = stack
1754 .last()
1755 .cloned()
1756 .ok_or_else(|| AftError::NoUndoHistory {
1757 path: path.display().to_string(),
1758 })?;
1759
1760 match entry.kind {
1761 BackupEntryKind::Content | BackupEntryKind::Symlink => {
1762 restore_entry_to_path(path, &entry).map_err(|e| AftError::IoError {
1763 path: path.display().to_string(),
1764 message: e.to_string(),
1765 })?;
1766 }
1767 BackupEntryKind::Tombstone => {
1768 remove_tombstone_path(path).map_err(|e| AftError::IoError {
1769 path: path.display().to_string(),
1770 message: e.to_string(),
1771 })?;
1772 remove_created_dirs_best_effort(&entry.created_dirs);
1773 }
1774 }
1775
1776 stack.pop();
1777 if stack.is_empty() {
1778 session_entries.remove(key);
1779 if session_entries.is_empty() {
1781 self.entries.remove(session);
1782 }
1783 self.remove_disk_backups_locked(session, key)?;
1784 } else {
1785 let stack_clone = self
1786 .entries
1787 .get(session)
1788 .and_then(|s| s.get(key))
1789 .cloned()
1790 .unwrap_or_default();
1791 self.write_snapshot_to_disk_locked(session, key, &stack_clone)?;
1792 }
1793
1794 Ok((entry, None))
1795 }
1796
1797 fn commit_restored_backup_locked(&mut self, session: &str, key: &Path) -> Result<(), AftError> {
1798 let mut remove_key = false;
1799 let mut remove_session = false;
1800 let mut remaining_stack = None;
1801
1802 if let Some(session_entries) = self.entries.get_mut(session) {
1803 if let Some(stack) = session_entries.get_mut(key) {
1804 stack.pop();
1805 if stack.is_empty() {
1806 remove_key = true;
1807 } else {
1808 remaining_stack = Some(stack.clone());
1809 }
1810 }
1811
1812 if remove_key {
1813 session_entries.remove(key);
1814 remove_session = session_entries.is_empty();
1815 }
1816 }
1817
1818 if remove_session {
1819 self.entries.remove(session);
1820 }
1821
1822 if remove_key {
1823 self.remove_disk_backups_locked(session, key)?;
1824 } else if let Some(stack) = remaining_stack {
1825 self.write_snapshot_to_disk_locked(session, key, &stack)?;
1826 }
1827
1828 Ok(())
1829 }
1830
1831 fn check_external_modification(
1832 &self,
1833 session: &str,
1834 key: &Path,
1835 path: &Path,
1836 ) -> Option<String> {
1837 let stack = self.entries.get(session).and_then(|s| s.get(key))?;
1838 let latest = stack.last()?;
1839 let modified = match latest.kind {
1840 BackupEntryKind::Content => std::fs::read(path)
1841 .map(|current| current.as_slice() != latest.content_bytes.as_ref())
1842 .unwrap_or(true),
1843 BackupEntryKind::Symlink => std::fs::read_link(path)
1844 .map(|target| latest.link_target.as_ref() != Some(&target))
1845 .unwrap_or(true),
1846 BackupEntryKind::Tombstone => false,
1847 };
1848 modified.then(|| "file was modified externally since last backup".to_string())
1849 }
1850
1851 fn backups_dir(&self) -> Option<PathBuf> {
1854 self.storage_dir
1855 .as_ref()
1856 .map(|dir| match &self.storage_harness {
1857 Some(harness) => dir.join(harness).join("backups"),
1858 None => dir.join("backups"),
1859 })
1860 }
1861
1862 fn session_dir(&self, session: &str) -> Option<PathBuf> {
1863 self.backups_dir()
1864 .map(|d| d.join(Self::session_hash(session)))
1865 }
1866
1867 fn session_hash(session: &str) -> String {
1868 hash_session(session)
1869 }
1870
1871 fn path_hash(key: &Path) -> String {
1872 stable_hash_16(key.to_string_lossy().as_bytes())
1877 }
1878
1879 fn write_session_marker(&self, session: &str, last_accessed: u64) {
1880 let Some(session_dir) = self.session_dir(session) else {
1881 return;
1882 };
1883 if let Err(e) = std::fs::create_dir_all(&session_dir) {
1884 crate::slog_warn!("failed to create session dir: {}", e);
1885 return;
1886 }
1887 let marker = session_dir.join("session.json");
1888 let json = serde_json::json!({
1889 "schema_version": SCHEMA_VERSION,
1890 "session_id": session,
1891 "last_accessed": last_accessed,
1892 });
1893 if let Ok(s) = serde_json::to_string_pretty(&json) {
1894 let tmp = session_dir.join("session.json.tmp");
1895 if std::fs::write(&tmp, s).is_ok() {
1896 let _ = std::fs::rename(&tmp, marker);
1897 }
1898 }
1899 }
1900
1901 fn repair_root_backups_if_needed(&self) {
1902 let (Some(storage_dir), Some(harness)) = (&self.storage_dir, &self.storage_harness) else {
1903 return;
1904 };
1905 let root_backups = storage_dir.join("backups");
1906 if !dir_has_entries(&root_backups) {
1907 return;
1908 }
1909 let harness_backups = storage_dir.join(harness).join("backups");
1910 if dir_has_entries(&harness_backups) {
1911 return;
1912 }
1913 if let Some(parent) = harness_backups.parent() {
1914 if let Err(error) = std::fs::create_dir_all(parent) {
1915 crate::slog_warn!(
1916 "failed to create harness backup dir {}: {}",
1917 parent.display(),
1918 error
1919 );
1920 return;
1921 }
1922 }
1923 if harness_backups.exists() {
1924 let _ = std::fs::remove_dir(&harness_backups);
1925 }
1926 match std::fs::rename(&root_backups, &harness_backups) {
1927 Ok(()) => {
1928 crate::slog_info!(
1929 "moved legacy root backups into harness namespace: {}",
1930 harness_backups.display()
1931 );
1932 }
1933 Err(error) => {
1934 crate::slog_warn!(
1935 "failed to move legacy root backups into {}: {}; trying child merge",
1936 harness_backups.display(),
1937 error
1938 );
1939 if std::fs::create_dir_all(&harness_backups).is_err() {
1940 return;
1941 }
1942 if let Ok(entries) = std::fs::read_dir(&root_backups) {
1943 for entry in entries.flatten() {
1944 let source = entry.path();
1945 let target = harness_backups.join(entry.file_name());
1946 if !target.exists() {
1947 let _ = std::fs::rename(source, target);
1948 }
1949 }
1950 }
1951 let _ = std::fs::remove_dir(&root_backups);
1952 }
1953 }
1954 }
1955
1956 fn gc_stale_sessions(&mut self, ttl_hours: u32) {
1957 let backups_dir = match self.backups_dir() {
1958 Some(d) if d.exists() => d,
1959 _ => return,
1960 };
1961 let ttl_secs = u64::from(if ttl_hours == 0 { 72 } else { ttl_hours }) * 60 * 60;
1962 let cutoff = current_timestamp().saturating_sub(ttl_secs);
1963 let entries = match std::fs::read_dir(&backups_dir) {
1964 Ok(entries) => entries,
1965 Err(_) => return,
1966 };
1967
1968 for entry in entries.flatten() {
1969 let session_dir = entry.path();
1970 if !session_dir.is_dir() || session_dir.join("meta.json").exists() {
1971 continue;
1972 }
1973 let Some(last_accessed) = Self::read_session_last_accessed(&session_dir) else {
1974 continue;
1975 };
1976 if last_accessed >= cutoff {
1977 continue;
1978 }
1979 if let Err(e) = std::fs::remove_dir_all(&session_dir) {
1980 crate::slog_warn!(
1981 "failed to remove stale backup session {}: {}",
1982 session_dir.display(),
1983 e
1984 );
1985 } else {
1986 crate::slog_warn!(
1987 "removed stale backup session {} (last_accessed={})",
1988 session_dir.display(),
1989 last_accessed
1990 );
1991 }
1992 }
1993 }
1994
1995 fn migrate_legacy_layout_if_needed(&mut self) {
2003 let backups_dir = match self.backups_dir() {
2004 Some(d) if d.exists() => d,
2005 _ => return,
2006 };
2007 let default_session_dir =
2008 backups_dir.join(Self::session_hash(crate::protocol::DEFAULT_SESSION_ID));
2009
2010 let entries = match std::fs::read_dir(&backups_dir) {
2011 Ok(e) => e,
2012 Err(_) => return,
2013 };
2014 let mut migrated = 0usize;
2015 for entry in entries.flatten() {
2016 let entry_path = entry.path();
2017 if !entry_path.is_dir() {
2019 continue;
2020 }
2021 if entry_path == default_session_dir {
2022 continue;
2023 }
2024 let meta_path = entry_path.join("meta.json");
2025 if !meta_path.exists() {
2026 continue; }
2028 if let Err(e) = std::fs::create_dir_all(&default_session_dir) {
2031 crate::slog_warn!("failed to create default session dir: {}", e);
2032 return;
2033 }
2034 let leaf = match entry_path.file_name() {
2035 Some(n) => n,
2036 None => continue,
2037 };
2038 let target = default_session_dir.join(leaf);
2039 if target.exists() {
2040 continue;
2043 }
2044 match std::fs::rename(&entry_path, &target) {
2045 Ok(()) => {
2046 Self::upgrade_meta_file(
2048 &target.join("meta.json"),
2049 crate::protocol::DEFAULT_SESSION_ID,
2050 );
2051 migrated += 1;
2052 }
2053 Err(e) => {
2054 crate::slog_warn!(
2055 "failed to migrate legacy backup {}: {}",
2056 entry_path.display(),
2057 e
2058 );
2059 }
2060 }
2061 }
2062 if migrated > 0 {
2063 crate::slog_info!(
2064 "migrated {} legacy backup entries into default session namespace",
2065 migrated
2066 );
2067 let marker = default_session_dir.join("session.json");
2069 let json = serde_json::json!({
2070 "schema_version": SCHEMA_VERSION,
2071 "session_id": crate::protocol::DEFAULT_SESSION_ID,
2072 "last_accessed": current_timestamp(),
2073 });
2074 if let Ok(s) = serde_json::to_string_pretty(&json) {
2075 let _ = std::fs::write(&marker, s);
2076 }
2077 }
2078 }
2079
2080 fn upgrade_meta_file(meta_path: &Path, session_id: &str) {
2081 let content = match std::fs::read_to_string(meta_path) {
2082 Ok(c) => c,
2083 Err(_) => return,
2084 };
2085 let mut parsed: serde_json::Value = match serde_json::from_str(&content) {
2086 Ok(v) => v,
2087 Err(_) => return,
2088 };
2089 if let Some(obj) = parsed.as_object_mut() {
2090 let count = obj.get("count").and_then(|v| v.as_u64()).unwrap_or(0);
2091 obj.insert(
2092 "schema_version".to_string(),
2093 serde_json::json!(SCHEMA_VERSION),
2094 );
2095 obj.insert("session_id".to_string(), serde_json::json!(session_id));
2096 obj.entry("entries").or_insert_with(|| {
2097 serde_json::Value::Array(
2098 (0..count)
2099 .map(|i| {
2100 serde_json::json!({
2101 "backup_id": format!("disk-{}", i),
2102 "timestamp": 0,
2103 "description": "restored from disk",
2104 "op_id": null,
2105 })
2106 })
2107 .collect(),
2108 )
2109 });
2110 }
2111 if let Ok(s) = serde_json::to_string_pretty(&parsed) {
2112 let tmp = meta_path.with_extension("json.tmp");
2113 if std::fs::write(&tmp, &s).is_ok() {
2114 let _ = std::fs::rename(&tmp, meta_path);
2115 }
2116 }
2117 }
2118
2119 fn read_session_last_accessed(session_dir: &Path) -> Option<u64> {
2120 let marker = session_dir.join("session.json");
2121 let content = std::fs::read_to_string(&marker).ok()?;
2122 let parsed: serde_json::Value = serde_json::from_str(&content).ok()?;
2123 parsed.get("last_accessed").and_then(|v| v.as_u64())
2124 }
2125
2126 fn should_snapshot_path(&self, path: &Path) -> Result<bool, AftError> {
2127 if !self.policy.enabled {
2128 return Ok(false);
2129 }
2130 let Some(max_file_size) = self.policy.max_file_size else {
2131 return Ok(true);
2132 };
2133 match std::fs::symlink_metadata(path) {
2134 Ok(metadata) if metadata.is_file() && metadata.len() > max_file_size => Ok(false),
2135 Ok(_) => Ok(true),
2136 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
2137 Err(AftError::FileNotFound {
2138 path: path.display().to_string(),
2139 })
2140 }
2141 Err(error) => Err(AftError::IoError {
2142 path: path.display().to_string(),
2143 message: error.to_string(),
2144 }),
2145 }
2146 }
2147
2148 fn ensure_session_marker(&self, session_dir: &Path, session: &str) -> Result<(), AftError> {
2149 let marker = session_dir.join("session.json");
2150 if marker.exists() {
2151 return Ok(());
2152 }
2153 let json = serde_json::json!({
2154 "schema_version": SCHEMA_VERSION,
2155 "session_id": session,
2156 "last_accessed": current_timestamp(),
2157 });
2158 let content = serde_json::to_string_pretty(&json).map_err(|error| AftError::IoError {
2159 path: marker.display().to_string(),
2160 message: error.to_string(),
2161 })?;
2162 write_temp_fsync_rename(session_dir, "session.json", content.as_bytes()).map_err(
2163 |error| AftError::IoError {
2164 path: marker.display().to_string(),
2165 message: error.to_string(),
2166 },
2167 )?;
2168 let _ = fsync_dir(session_dir);
2169 Ok(())
2170 }
2171
2172 fn acquire_stack_disk_lock(
2173 &self,
2174 session: &str,
2175 key: &Path,
2176 ) -> Result<Option<crate::fs_lock::LockGuard>, AftError> {
2177 let Some(session_dir) = self.session_dir(session) else {
2178 return Ok(None);
2179 };
2180 self.record_disk_io_for_tests();
2181 let lock_dir = session_dir.join(".locks");
2182 std::fs::create_dir_all(&lock_dir).map_err(|error| AftError::IoError {
2183 path: lock_dir.display().to_string(),
2184 message: error.to_string(),
2185 })?;
2186 let lock_path = lock_dir.join(format!("{}.lock", Self::path_hash(key)));
2187 crate::fs_lock::acquire(&lock_path)
2188 .map(Some)
2189 .map_err(|error| AftError::IoError {
2190 path: lock_path.display().to_string(),
2191 message: error.to_string(),
2192 })
2193 }
2194
2195 fn acquire_stack_disk_locks(
2196 &self,
2197 session: &str,
2198 keys: &[PathBuf],
2199 ) -> Result<Vec<crate::fs_lock::LockGuard>, AftError> {
2200 let mut keys = keys.to_vec();
2201 keys.sort();
2202 keys.dedup();
2203 let mut guards = Vec::with_capacity(keys.len());
2204 for key in keys {
2205 if let Some(guard) = self.acquire_stack_disk_lock(session, &key)? {
2206 guards.push(guard);
2207 }
2208 }
2209 Ok(guards)
2210 }
2211
2212 #[cfg(test)]
2213 fn load_from_disk_if_needed(&mut self, session: &str, key: &Path) -> Result<bool, AftError> {
2214 let _disk_lock = self.acquire_stack_disk_lock(session, key)?;
2215 self.load_from_disk_if_needed_locked(session, key)
2216 }
2217
2218 fn load_from_disk_if_needed_locked(
2219 &mut self,
2220 session: &str,
2221 key: &Path,
2222 ) -> Result<bool, AftError> {
2223 let entries = match self.read_stack_from_disk_unlocked(session, key) {
2224 Ok(Some(entries)) => entries,
2225 Ok(None) => {
2226 if self.session_dir(session).is_some() {
2227 self.restore_in_memory_stack(session, key, None);
2228 }
2229 if let Some(files) = self.disk_index.get_mut(session) {
2230 files.remove(key);
2231 if files.is_empty() {
2232 self.disk_index.remove(session);
2233 }
2234 }
2235 return Ok(false);
2236 }
2237 Err(error) => {
2238 return Err(AftError::IoError {
2239 path: key.display().to_string(),
2240 message: error,
2241 });
2242 }
2243 };
2244
2245 self.update_counter_from_entries(&entries);
2246 if let Ok(Some((disk_meta, _))) = self.read_disk_meta_value(session, key) {
2247 self.disk_index
2248 .entry(session.to_string())
2249 .or_default()
2250 .insert(key.to_path_buf(), disk_meta);
2251 }
2252 self.entries
2253 .entry(session.to_string())
2254 .or_default()
2255 .insert(key.to_path_buf(), entries);
2256 Ok(true)
2257 }
2258
2259 fn ensure_stack_hydrated_locked(&mut self, session: &str, key: &Path) -> Result<(), AftError> {
2266 self.load_from_disk_if_needed_locked(session, key)?;
2267 Ok(())
2268 }
2269
2270 fn refresh_disk_index_for_session(&mut self, session: &str) -> Result<Vec<PathBuf>, AftError> {
2271 let Some(session_dir) = self.session_dir(session) else {
2272 self.disk_index.remove(session);
2273 return Ok(Vec::new());
2274 };
2275 if !session_dir.exists() {
2276 self.disk_index.remove(session);
2277 return Ok(Vec::new());
2278 }
2279
2280 let path_dirs = std::fs::read_dir(&session_dir).map_err(|error| AftError::IoError {
2281 path: session_dir.display().to_string(),
2282 message: error.to_string(),
2283 })?;
2284 let mut per_session = HashMap::new();
2285 for path_entry in path_dirs {
2286 let path_entry = path_entry.map_err(|error| AftError::IoError {
2287 path: session_dir.display().to_string(),
2288 message: error.to_string(),
2289 })?;
2290 let path_dir = path_entry.path();
2291 if !path_dir.is_dir() {
2292 continue;
2293 }
2294 let meta_path = path_dir.join("meta.json");
2295 if !meta_path.exists() {
2296 continue;
2297 }
2298 let content =
2299 std::fs::read_to_string(&meta_path).map_err(|error| AftError::IoError {
2300 path: meta_path.display().to_string(),
2301 message: error.to_string(),
2302 })?;
2303 let meta = serde_json::from_str::<serde_json::Value>(&content).map_err(|error| {
2304 AftError::IoError {
2305 path: meta_path.display().to_string(),
2306 message: error.to_string(),
2307 }
2308 })?;
2309 let path_str = meta
2310 .get("path")
2311 .and_then(|value| value.as_str())
2312 .ok_or_else(|| AftError::IoError {
2313 path: meta_path.display().to_string(),
2314 message: "backup meta missing path".to_string(),
2315 })?;
2316 let key = PathBuf::from(path_str);
2317 if !is_loadable_backup_path(&key, &path_dir) {
2318 continue;
2319 }
2320 let count = meta_entry_count(&meta).ok_or_else(|| AftError::IoError {
2321 path: meta_path.display().to_string(),
2322 message: "backup meta missing entry count".to_string(),
2323 })?;
2324 if count > 0 {
2325 per_session.insert(
2326 key,
2327 DiskMeta {
2328 dir: path_dir,
2329 count,
2330 },
2331 );
2332 }
2333 }
2334
2335 let keys = per_session.keys().cloned().collect::<Vec<_>>();
2336 if per_session.is_empty() {
2337 self.disk_index.remove(session);
2338 } else {
2339 self.disk_index.insert(session.to_string(), per_session);
2340 }
2341 Ok(keys)
2342 }
2343
2344 fn restore_operation_candidate_keys(
2345 &mut self,
2346 session: &str,
2347 ) -> Result<Vec<PathBuf>, AftError> {
2348 let mut keys: HashSet<PathBuf> = self
2349 .refresh_disk_index_for_session(session)?
2350 .into_iter()
2351 .collect();
2352 if let Some(files) = self.entries.get(session) {
2353 keys.extend(files.keys().cloned());
2354 }
2355 let mut keys = keys.into_iter().collect::<Vec<_>>();
2356 keys.sort();
2357 Ok(keys)
2358 }
2359
2360 fn read_stack_heads_from_disk(
2361 &self,
2362 session: &str,
2363 key: &Path,
2364 ) -> Option<Vec<BackupEntryHead>> {
2365 let _disk_lock = match self.acquire_stack_disk_lock(session, key) {
2366 Ok(lock) => lock,
2367 Err(error) => {
2368 crate::slog_warn!(
2369 "backup disk head read lock failed for {}: {}",
2370 key.display(),
2371 error
2372 );
2373 return None;
2374 }
2375 };
2376 match self.read_stack_heads_from_disk_unlocked(session, key) {
2377 Ok(heads) => heads,
2378 Err(error) => {
2379 crate::slog_warn!(
2380 "backup disk head read failed for {}: {}",
2381 key.display(),
2382 error
2383 );
2384 None
2385 }
2386 }
2387 }
2388
2389 fn read_stack_heads_from_disk_unlocked(
2390 &self,
2391 session: &str,
2392 key: &Path,
2393 ) -> Result<Option<Vec<BackupEntryHead>>, String> {
2394 let Some((disk_meta, meta)) = self.read_disk_meta_value(session, key)? else {
2395 return Ok(None);
2396 };
2397 if disk_meta.count == 0 {
2398 return Ok(None);
2399 }
2400
2401 let heads = if is_v2_meta(&meta) {
2402 let entries = meta_entries(&meta)?;
2403 for entry in entries {
2404 self.validate_v2_content_reference(&disk_meta.dir, entry)?;
2405 }
2406 entries
2407 .iter()
2408 .enumerate()
2409 .map(|(i, entry)| backup_head_from_meta(Some(entry), i))
2410 .collect::<Vec<_>>()
2411 } else {
2412 let entries = meta.get("entries").and_then(|value| value.as_array());
2413 (0..disk_meta.count)
2414 .map(|i| backup_head_from_meta(entries.and_then(|entries| entries.get(i)), i))
2415 .collect::<Vec<_>>()
2416 };
2417
2418 Ok((!heads.is_empty()).then_some(heads))
2419 }
2420
2421 fn read_stack_from_disk_unlocked(
2422 &self,
2423 session: &str,
2424 key: &Path,
2425 ) -> Result<Option<Vec<BackupEntry>>, String> {
2426 let Some((disk_meta, meta)) = self.read_disk_meta_value(session, key)? else {
2427 return Ok(None);
2428 };
2429 if disk_meta.count == 0 {
2430 return Ok(None);
2431 }
2432
2433 let entries = if is_v2_meta(&meta) {
2434 meta_entries(&meta)?
2435 .iter()
2436 .enumerate()
2437 .map(|(i, entry_meta)| self.entry_from_v2_meta(&disk_meta.dir, entry_meta, i))
2438 .collect::<Result<Vec<_>, _>>()?
2439 } else {
2440 let entries = meta.get("entries").and_then(|value| value.as_array());
2441 let mut loaded = Vec::new();
2442 for i in 0..disk_meta.count {
2443 let entry_meta = entries.and_then(|entries| entries.get(i));
2444 if let Some(entry) = legacy_entry_from_meta(&disk_meta.dir, entry_meta, i) {
2445 loaded.push(entry);
2446 }
2447 }
2448 loaded
2449 };
2450
2451 Ok((!entries.is_empty()).then_some(entries))
2452 }
2453
2454 fn read_disk_meta_value(
2455 &self,
2456 session: &str,
2457 key: &Path,
2458 ) -> Result<Option<(DiskMeta, serde_json::Value)>, String> {
2459 let Some(session_dir) = self.session_dir(session) else {
2460 return Ok(None);
2461 };
2462 let dir = session_dir.join(Self::path_hash(key));
2463 let meta_path = dir.join("meta.json");
2464 if !meta_path.exists() {
2465 return Ok(None);
2466 }
2467 let content = std::fs::read_to_string(&meta_path)
2468 .map_err(|error| format!("failed to read {}: {}", meta_path.display(), error))?;
2469 let meta = serde_json::from_str::<serde_json::Value>(&content)
2470 .map_err(|error| format!("failed to parse {}: {}", meta_path.display(), error))?;
2471 let path_str = meta
2472 .get("path")
2473 .and_then(|value| value.as_str())
2474 .ok_or_else(|| format!("backup meta {} missing path", meta_path.display()))?;
2475 let stored_key = PathBuf::from(path_str);
2476 if stored_key != key || !is_loadable_backup_path(&stored_key, &dir) {
2477 return Ok(None);
2478 }
2479 let count = meta_entry_count(&meta)
2480 .ok_or_else(|| format!("backup meta {} missing entry count", meta_path.display()))?;
2481 Ok(Some((DiskMeta { dir, count }, meta)))
2482 }
2483
2484 fn validate_v2_content_reference(
2485 &self,
2486 dir: &Path,
2487 entry_meta: &serde_json::Value,
2488 ) -> Result<(), String> {
2489 let kind = entry_kind_from_meta(Some(entry_meta));
2490 if matches!(kind, BackupEntryKind::Tombstone) {
2491 return Ok(());
2492 }
2493 let content_path = content_path_from_meta(entry_meta)?;
2494 let path = dir.join(content_path);
2495 if !path.is_file() {
2496 return Err(format!(
2497 "v2 backup meta references missing content file {}",
2498 path.display()
2499 ));
2500 }
2501 Ok(())
2502 }
2503
2504 fn entry_from_v2_meta(
2505 &self,
2506 dir: &Path,
2507 entry_meta: &serde_json::Value,
2508 index: usize,
2509 ) -> Result<BackupEntry, String> {
2510 let kind = entry_kind_from_meta(Some(entry_meta));
2511 let content_bytes = match kind {
2512 BackupEntryKind::Content | BackupEntryKind::Symlink => {
2513 let content_path = content_path_from_meta(entry_meta)?;
2514 let path = dir.join(content_path);
2515 std::fs::read(&path).map_err(|error| {
2516 format!(
2517 "failed to read v2 backup content {}: {}",
2518 path.display(),
2519 error
2520 )
2521 })?
2522 }
2523 BackupEntryKind::Tombstone => Vec::new(),
2524 };
2525 Ok(entry_from_meta(
2526 Some(entry_meta),
2527 index,
2528 kind,
2529 content_bytes,
2530 ))
2531 }
2532
2533 fn write_snapshot_to_disk(
2534 &mut self,
2535 session: &str,
2536 key: &Path,
2537 stack: &[BackupEntry],
2538 ) -> Result<(), AftError> {
2539 let _disk_lock = self.acquire_stack_disk_lock(session, key)?;
2540 self.write_snapshot_to_disk_locked(session, key, stack)
2541 }
2542
2543 fn write_snapshot_to_disk_locked(
2544 &mut self,
2545 session: &str,
2546 key: &Path,
2547 stack: &[BackupEntry],
2548 ) -> Result<(), AftError> {
2549 self.write_snapshot_to_disk_locked_with_db_plan(session, key, stack, DbMirrorPlan::Full)
2550 }
2551
2552 fn write_appended_snapshot_to_disk_locked(
2553 &mut self,
2554 session: &str,
2555 key: &Path,
2556 stack: &[BackupEntry],
2557 evicted_orders: &[u128],
2558 new_entry: Option<&BackupEntry>,
2559 ) -> Result<(), AftError> {
2560 self.write_snapshot_to_disk_locked_with_db_plan(
2561 session,
2562 key,
2563 stack,
2564 DbMirrorPlan::Append {
2565 evicted_orders,
2566 new_entry,
2567 },
2568 )
2569 }
2570
2571 fn write_snapshot_to_disk_locked_with_db_plan(
2572 &mut self,
2573 session: &str,
2574 key: &Path,
2575 stack: &[BackupEntry],
2576 db_plan: DbMirrorPlan<'_>,
2577 ) -> Result<(), AftError> {
2578 #[cfg(test)]
2579 if self.fail_next_disk_write {
2580 self.fail_next_disk_write = false;
2581 return Err(AftError::IoError {
2582 path: key.display().to_string(),
2583 message: "injected backup disk write failure".to_string(),
2584 });
2585 }
2586
2587 let Some(session_dir) = self.session_dir(session) else {
2588 return Ok(());
2589 };
2590
2591 std::fs::create_dir_all(&session_dir).map_err(|error| AftError::IoError {
2592 path: session_dir.display().to_string(),
2593 message: error.to_string(),
2594 })?;
2595 self.ensure_session_marker(&session_dir, session)?;
2596
2597 let hash = Self::path_hash(key);
2598 let dir = session_dir.join(&hash);
2599 std::fs::create_dir_all(&dir).map_err(|error| AftError::IoError {
2600 path: dir.display().to_string(),
2601 message: error.to_string(),
2602 })?;
2603
2604 let max_depth = self.policy.max_depth;
2605 let retained_start = stack.len().saturating_sub(max_depth);
2606 let retained = &stack[retained_start..];
2607 let mut referenced_content = HashSet::new();
2608 let mut wrote_content = false;
2609
2610 for entry in retained {
2611 if let Some(content_path) = content_filename_for_entry(entry) {
2612 referenced_content.insert(content_path.clone());
2613 let final_path = dir.join(&content_path);
2614 if final_path.exists() {
2615 continue;
2616 }
2617 let bytes = content_bytes_for_disk(entry);
2618 write_temp_fsync_rename(&dir, &content_path, &bytes).map_err(|error| {
2619 AftError::IoError {
2620 path: final_path.display().to_string(),
2621 message: error.to_string(),
2622 }
2623 })?;
2624 wrote_content = true;
2625 }
2626 }
2627 if wrote_content {
2628 fsync_dir(&dir).map_err(|error| AftError::IoError {
2629 path: dir.display().to_string(),
2630 message: error.to_string(),
2631 })?;
2632 }
2633
2634 let entries: Vec<serde_json::Value> = retained.iter().map(entry_meta_json).collect();
2635 let meta = serde_json::json!({
2636 "schema_version": SCHEMA_VERSION,
2637 "format_version": V2_FORMAT_VERSION,
2638 "session_id": session,
2639 "path": key.display().to_string(),
2640 "count": retained.len(),
2641 "entries": entries,
2642 });
2643 let meta_content =
2644 serde_json::to_string_pretty(&meta).map_err(|error| AftError::IoError {
2645 path: dir.join("meta.json").display().to_string(),
2646 message: error.to_string(),
2647 })?;
2648 write_temp_fsync_rename(&dir, "meta.json", meta_content.as_bytes()).map_err(|error| {
2649 AftError::IoError {
2650 path: dir.join("meta.json").display().to_string(),
2651 message: error.to_string(),
2652 }
2653 })?;
2654 fsync_dir(&dir).map_err(|error| AftError::IoError {
2655 path: dir.display().to_string(),
2656 message: error.to_string(),
2657 })?;
2658
2659 prune_unreferenced_backup_files(&dir, &referenced_content).map_err(|error| {
2660 AftError::IoError {
2661 path: dir.display().to_string(),
2662 message: error.to_string(),
2663 }
2664 })?;
2665 let _ = fsync_dir(&dir);
2666
2667 self.disk_index
2670 .entry(session.to_string())
2671 .or_default()
2672 .insert(
2673 key.to_path_buf(),
2674 DiskMeta {
2675 dir: dir.clone(),
2676 count: retained.len(),
2677 },
2678 );
2679 self.mirror_stack_to_db(session, key, retained, db_plan);
2680 Ok(())
2681 }
2682
2683 fn mirror_stack_to_db(
2684 &self,
2685 session: &str,
2686 key: &Path,
2687 stack: &[BackupEntry],
2688 plan: DbMirrorPlan<'_>,
2689 ) {
2690 let pool = self.db_pool.read().ok().and_then(|slot| slot.clone());
2691 let Some(pool) = pool else {
2692 return;
2693 };
2694 let harness = self.db_harness.read().ok().and_then(|slot| slot.clone());
2695 let Some(harness) = harness else {
2696 crate::slog_warn!(
2697 "dual-write backup to DB skipped for {}: harness not configured",
2698 key.display()
2699 );
2700 return;
2701 };
2702 let project_key = self
2703 .db_project_key
2704 .read()
2705 .ok()
2706 .and_then(|slot| slot.clone());
2707 let Some(project_key) = project_key else {
2708 crate::slog_warn!(
2709 "dual-write backup to DB skipped for {}: project key not configured",
2710 key.display()
2711 );
2712 return;
2713 };
2714
2715 let conn = match pool.lock() {
2716 Ok(conn) => conn,
2717 Err(_) => {
2718 self.set_db_mirror_synced(session, key, false);
2719 crate::slog_warn!(
2720 "dual-write backup to DB failed for {}: db mutex poisoned",
2721 key.display()
2722 );
2723 return;
2724 }
2725 };
2726 let path_hash = Self::path_hash(key);
2727 let file_path = key.display().to_string();
2728
2729 let context = DbMirrorContext {
2730 harness: &harness,
2731 session,
2732 project_key: &project_key,
2733 file_path: &file_path,
2734 path_hash: &path_hash,
2735 };
2736 let (write_result, operation) = match plan {
2737 DbMirrorPlan::Append {
2738 evicted_orders,
2739 new_entry,
2740 } if self.db_mirror_is_synced(session, key) => (
2741 apply_backup_append_delta_in_db(&conn, &context, evicted_orders, new_entry),
2742 "append delta",
2743 ),
2744 DbMirrorPlan::Append { .. } | DbMirrorPlan::Full => (
2748 replace_backup_stack_in_db(&conn, &context, stack),
2749 "full stack",
2750 ),
2751 };
2752 self.set_db_mirror_synced(session, key, write_result.is_ok());
2753 if let Err(error) = write_result {
2754 crate::slog_warn!(
2755 "dual-write backup {} to DB failed for {} (rolled back, prior stack kept): {}",
2756 operation,
2757 key.display(),
2758 error
2759 );
2760 }
2761 }
2762
2763 fn prune_disk_stacks_to_depth(&mut self, max_depth: usize) -> HashSet<(String, PathBuf)> {
2764 let disk_keys = self
2768 .disk_index
2769 .iter()
2770 .flat_map(|(session, files)| {
2771 files
2772 .keys()
2773 .cloned()
2774 .map(|key| (session.clone(), key))
2775 .collect::<Vec<_>>()
2776 })
2777 .collect::<Vec<_>>();
2778 let mut failed = HashSet::new();
2779
2780 for (session, key) in disk_keys {
2781 let disk_lock = match self.acquire_stack_disk_lock(&session, &key) {
2782 Ok(lock) => lock,
2783 Err(error) => {
2784 crate::slog_warn!(
2785 "failed to lock backup stack for {} while applying max_depth: {}",
2786 key.display(),
2787 error
2788 );
2789 failed.insert((session, key));
2790 continue;
2791 }
2792 };
2793
2794 let mut stack = match self.read_stack_from_disk_unlocked(&session, &key) {
2795 Ok(Some(stack)) => stack,
2796 Ok(None) => Vec::new(),
2797 Err(error) => {
2798 crate::slog_warn!(
2799 "failed to read backup stack for {} while applying max_depth: {}",
2800 key.display(),
2801 error
2802 );
2803 failed.insert((session, key));
2804 drop(disk_lock);
2805 continue;
2806 }
2807 };
2808 trim_stack_to_depth(&mut stack, max_depth);
2809 if let Err(error) = self.write_snapshot_to_disk_locked(&session, &key, &stack) {
2810 crate::slog_warn!(
2811 "failed to prune backup stack for {} while applying max_depth: {}",
2812 key.display(),
2813 error
2814 );
2815 failed.insert((session, key));
2816 drop(disk_lock);
2817 continue;
2818 }
2819 if stack.is_empty() {
2820 if let Some(files) = self.entries.get_mut(&session) {
2821 files.remove(&key);
2822 if files.is_empty() {
2823 self.entries.remove(&session);
2824 }
2825 }
2826 } else {
2827 self.entries
2828 .entry(session.clone())
2829 .or_default()
2830 .insert(key.clone(), stack);
2831 }
2832 drop(disk_lock);
2833 }
2834
2835 failed
2836 }
2837
2838 fn remove_disk_backups(&mut self, session: &str, key: &Path) -> Result<(), AftError> {
2839 let _disk_lock = self.acquire_stack_disk_lock(session, key)?;
2840 self.remove_disk_backups_locked(session, key)
2841 }
2842
2843 fn remove_disk_backups_locked(&mut self, session: &str, key: &Path) -> Result<(), AftError> {
2844 self.remove_db_backups(session, key);
2845 let removed = self.disk_index.get_mut(session).and_then(|s| s.remove(key));
2846 if let Some(meta) = removed {
2847 if let Err(error) = std::fs::remove_dir_all(&meta.dir) {
2848 return Err(AftError::IoError {
2849 path: meta.dir.display().to_string(),
2850 message: error.to_string(),
2851 });
2852 }
2853 } else if let Some(session_dir) = self.session_dir(session) {
2854 let hash = Self::path_hash(key);
2855 let dir = session_dir.join(&hash);
2856 if dir.exists() {
2857 if let Err(error) = std::fs::remove_dir_all(&dir) {
2858 return Err(AftError::IoError {
2859 path: dir.display().to_string(),
2860 message: error.to_string(),
2861 });
2862 }
2863 }
2864 }
2865
2866 let empty = self
2869 .disk_index
2870 .get(session)
2871 .map(|s| s.is_empty())
2872 .unwrap_or(false);
2873 if empty {
2874 self.disk_index.remove(session);
2875 }
2876 Ok(())
2877 }
2878
2879 fn remove_db_backups(&self, session: &str, key: &Path) {
2880 let Some((pool, harness)) = self.db_pool_and_harness() else {
2881 return;
2882 };
2883 let conn = match pool.lock() {
2884 Ok(conn) => conn,
2885 Err(_) => {
2886 self.set_db_mirror_synced(session, key, false);
2887 crate::slog_warn!(
2888 "delete backup DB rows failed for {}: db mutex poisoned",
2889 key.display()
2890 );
2891 return;
2892 }
2893 };
2894 let path_hash = Self::path_hash(key);
2895 match crate::db::backups::delete_backups_for_path(&conn, &harness, session, &path_hash) {
2896 Ok(_) => self.set_db_mirror_synced(session, key, false),
2899 Err(error) => {
2900 self.set_db_mirror_synced(session, key, false);
2901 crate::slog_warn!(
2902 "delete backup DB rows failed for {}: {}",
2903 key.display(),
2904 error
2905 );
2906 }
2907 }
2908 }
2909}
2910
2911fn backup_row_for_db(entry: &BackupEntry, context: &DbMirrorContext<'_>) -> BackupRow {
2912 let backup_path = content_filename_for_entry(entry);
2913 entry.to_backup_row(
2914 context.harness,
2915 context.session,
2916 context.project_key,
2917 context.file_path,
2918 context.path_hash,
2919 backup_path.as_deref(),
2920 )
2921}
2922
2923fn replace_backup_stack_in_db(
2924 conn: &Connection,
2925 context: &DbMirrorContext<'_>,
2926 stack: &[BackupEntry],
2927) -> rusqlite::Result<()> {
2928 let tx = conn.unchecked_transaction()?;
2933 crate::db::backups::delete_backups_for_path(
2934 &tx,
2935 context.harness,
2936 context.session,
2937 context.path_hash,
2938 )?;
2939 for entry in stack {
2940 crate::db::backups::insert_backup(&tx, &backup_row_for_db(entry, context))?;
2941 }
2942 tx.commit()
2943}
2944
2945fn apply_backup_append_delta_in_db(
2946 conn: &Connection,
2947 context: &DbMirrorContext<'_>,
2948 evicted_orders: &[u128],
2949 new_entry: Option<&BackupEntry>,
2950) -> rusqlite::Result<()> {
2951 let tx = conn.unchecked_transaction()?;
2955 for &order in evicted_orders {
2956 crate::db::backups::delete_backup_for_order(
2957 &tx,
2958 context.harness,
2959 context.session,
2960 context.path_hash,
2961 order,
2962 )?;
2963 }
2964 if let Some(entry) = new_entry {
2965 crate::db::backups::insert_backup(&tx, &backup_row_for_db(entry, context))?;
2966 }
2967 tx.commit()
2968}
2969
2970pub fn hash_session(session: &str) -> String {
2971 stable_hash_16(session.as_bytes())
2972}
2973
2974pub fn new_op_id() -> String {
2975 let mut bytes = [0u8; 4];
2976 if getrandom::fill(&mut bytes).is_err() {
2977 bytes = current_timestamp().to_le_bytes()[..4]
2978 .try_into()
2979 .unwrap_or([0; 4]);
2980 }
2981 let rand = u32::from_le_bytes(bytes);
2982 format!("op-{}-{:08x}", current_timestamp() * 1000, rand)
2983}
2984
2985#[derive(Debug, Clone, PartialEq, Eq)]
2986struct BackupEntryDiskMetadata {
2987 mode: Option<u32>,
2988 link_target: Option<PathBuf>,
2989 created_dirs: Vec<PathBuf>,
2990}
2991
2992fn restore_metadata_json(entry: &BackupEntry) -> String {
2993 serde_json::json!({
2994 "version": DB_RESTORE_META_VERSION,
2995 "mode": entry.mode,
2996 "link_target": entry.link_target.as_ref().map(|target| target.display().to_string()),
2997 "created_dirs": entry
2998 .created_dirs
2999 .iter()
3000 .map(|dir| dir.display().to_string())
3001 .collect::<Vec<_>>(),
3002 })
3003 .to_string()
3004}
3005
3006fn restore_metadata_from_json(value: &str) -> Option<BackupEntryDiskMetadata> {
3007 let value: serde_json::Value = serde_json::from_str(value).ok()?;
3008 if value.get("version")?.as_u64()? != u64::from(DB_RESTORE_META_VERSION) {
3009 return None;
3010 }
3011
3012 let mode = value.get("mode")?;
3013 if !mode.is_null()
3014 && mode
3015 .as_u64()
3016 .and_then(|mode| u32::try_from(mode).ok())
3017 .is_none()
3018 {
3019 return None;
3020 }
3021 let link_target = value.get("link_target")?;
3022 if !link_target.is_null() && !link_target.is_string() {
3023 return None;
3024 }
3025 if !value
3026 .get("created_dirs")?
3027 .as_array()?
3028 .iter()
3029 .all(serde_json::Value::is_string)
3030 {
3031 return None;
3032 }
3033
3034 Some(restore_metadata_fields(&value))
3035}
3036
3037fn restore_metadata_fields(value: &serde_json::Value) -> BackupEntryDiskMetadata {
3038 BackupEntryDiskMetadata {
3039 mode: value
3040 .get("mode")
3041 .and_then(|value| value.as_u64())
3042 .and_then(|mode| u32::try_from(mode).ok()),
3043 link_target: value
3044 .get("link_target")
3045 .and_then(|value| value.as_str())
3046 .map(PathBuf::from),
3047 created_dirs: value
3048 .get("created_dirs")
3049 .and_then(|value| value.as_array())
3050 .map(|dirs| {
3051 dirs.iter()
3052 .filter_map(|dir| dir.as_str())
3053 .map(PathBuf::from)
3054 .collect()
3055 })
3056 .unwrap_or_default(),
3057 }
3058}
3059
3060#[derive(Debug, Clone)]
3061enum RestorePathState {
3062 Missing,
3063 Regular {
3064 content_bytes: Vec<u8>,
3065 mode: Option<u32>,
3066 },
3067 Symlink {
3068 target: PathBuf,
3069 },
3070 Directory,
3071}
3072
3073fn backup_entry_from_path(
3074 path: &Path,
3075 backup_id: String,
3076 order: u128,
3077 description: &str,
3078 op_id: Option<&str>,
3079) -> Result<BackupEntry, AftError> {
3080 let metadata = std::fs::symlink_metadata(path).map_err(|error| match error.kind() {
3081 std::io::ErrorKind::NotFound => AftError::FileNotFound {
3082 path: path.display().to_string(),
3083 },
3084 _ => AftError::IoError {
3085 path: path.display().to_string(),
3086 message: error.to_string(),
3087 },
3088 })?;
3089 let mode = file_mode(&metadata);
3090
3091 let (kind, content, content_bytes, link_target) = if metadata.file_type().is_symlink() {
3092 let target = std::fs::read_link(path).map_err(|error| AftError::IoError {
3093 path: path.display().to_string(),
3094 message: error.to_string(),
3095 })?;
3096 (
3097 BackupEntryKind::Symlink,
3098 target.display().to_string(),
3099 Arc::from([]),
3100 Some(target),
3101 )
3102 } else if metadata.is_file() {
3103 let bytes: Arc<[u8]> = read_captured_content(path)
3104 .map_err(|error| AftError::IoError {
3105 path: path.display().to_string(),
3106 message: error.to_string(),
3107 })?
3108 .into();
3109 (
3110 BackupEntryKind::Content,
3111 String::from_utf8_lossy(&bytes).into_owned(),
3112 bytes,
3113 None,
3114 )
3115 } else {
3116 return Err(AftError::InvalidRequest {
3117 message: format!(
3118 "backup: '{}' is not a regular file or symlink",
3119 path.display()
3120 ),
3121 });
3122 };
3123
3124 Ok(BackupEntry {
3125 backup_id,
3126 content,
3127 content_bytes,
3128 timestamp: current_timestamp(),
3129 order,
3130 description: description.to_string(),
3131 op_id: op_id.map(str::to_string),
3132 kind,
3133 mode,
3134 link_target,
3135 created_dirs: Vec::new(),
3136 })
3137}
3138
3139fn backup_entry_from_capture(
3140 capture: &CapturedRegularFile,
3141 backup_id: String,
3142 order: u128,
3143 description: &str,
3144 op_id: Option<&str>,
3145) -> BackupEntry {
3146 BackupEntry {
3147 backup_id,
3148 content: String::from_utf8_lossy(capture.bytes()).into_owned(),
3149 content_bytes: capture.shared_bytes(),
3150 timestamp: current_timestamp(),
3151 order,
3152 description: description.to_string(),
3153 op_id: op_id.map(str::to_string),
3154 kind: BackupEntryKind::Content,
3155 mode: file_mode(capture.metadata()),
3156 link_target: None,
3157 created_dirs: Vec::new(),
3158 }
3159}
3160
3161fn canonicalize_key(path: &Path) -> PathBuf {
3162 let absolute = if path.is_absolute() {
3163 path.to_path_buf()
3164 } else {
3165 std::env::current_dir()
3166 .unwrap_or_else(|_| PathBuf::from("."))
3167 .join(path)
3168 };
3169
3170 match std::fs::symlink_metadata(&absolute) {
3171 Ok(metadata) if metadata.file_type().is_symlink() => {
3172 canonicalize_parent_join_leaf(&absolute)
3173 }
3174 Ok(_) => std::fs::canonicalize(&absolute)
3175 .map(|path| normalize_absolute_key(&path))
3176 .unwrap_or_else(|_| canonicalize_existing_ancestor(&absolute)),
3177 Err(_) => canonicalize_existing_ancestor(&absolute),
3178 }
3179}
3180
3181fn canonicalize_parent_join_leaf(path: &Path) -> PathBuf {
3182 let Some(parent) = path.parent() else {
3183 return normalize_absolute_key(path);
3184 };
3185 let mut key = canonicalize_existing_ancestor(parent);
3186 if let Some(file_name) = path.file_name() {
3187 key.push(file_name);
3188 }
3189 key
3190}
3191
3192fn canonicalize_existing_ancestor(path: &Path) -> PathBuf {
3193 let mut suffix = Vec::new();
3194 let mut current = path;
3195
3196 loop {
3197 if let Ok(mut base) = std::fs::canonicalize(current) {
3198 for component in suffix.iter().rev() {
3199 base.push(Path::new(component));
3200 }
3201 return normalize_absolute_key(&base);
3202 }
3203 let Some(parent) = current.parent() else {
3204 return normalize_absolute_key(path);
3205 };
3206 if let Some(file_name) = current.file_name() {
3207 suffix.push(file_name.to_os_string());
3208 }
3209 current = parent;
3210 }
3211}
3212
3213fn normalize_absolute_key(path: &Path) -> PathBuf {
3214 let mut normalized = PathBuf::new();
3215
3216 for component in path.components() {
3217 match component {
3218 std::path::Component::CurDir => {}
3219 std::path::Component::ParentDir => {
3220 if !normalized.pop() {
3221 normalized.push(component.as_os_str());
3222 }
3223 }
3224 other => normalized.push(other.as_os_str()),
3225 }
3226 }
3227
3228 normalized
3229}
3230
3231fn file_mode(metadata: &std::fs::Metadata) -> Option<u32> {
3232 #[cfg(unix)]
3233 {
3234 use std::os::unix::fs::PermissionsExt;
3235 Some(metadata.permissions().mode())
3236 }
3237 #[cfg(not(unix))]
3238 {
3239 let _ = metadata;
3240 None
3241 }
3242}
3243
3244fn set_file_mode(path: &Path, mode: Option<u32>) -> std::io::Result<()> {
3245 #[cfg(unix)]
3246 {
3247 use std::os::unix::fs::PermissionsExt;
3248 if let Some(mode) = mode {
3249 std::fs::set_permissions(path, std::fs::Permissions::from_mode(mode))?;
3250 }
3251 }
3252 #[cfg(not(unix))]
3253 {
3254 let _ = (path, mode);
3255 }
3256 Ok(())
3257}
3258
3259fn capture_path_state(path: &Path) -> Result<RestorePathState, AftError> {
3260 let metadata = match std::fs::symlink_metadata(path) {
3261 Ok(metadata) => metadata,
3262 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
3263 return Ok(RestorePathState::Missing);
3264 }
3265 Err(error) => {
3266 return Err(AftError::IoError {
3267 path: path.display().to_string(),
3268 message: error.to_string(),
3269 });
3270 }
3271 };
3272
3273 if metadata.file_type().is_symlink() {
3274 let target = std::fs::read_link(path).map_err(|error| AftError::IoError {
3275 path: path.display().to_string(),
3276 message: error.to_string(),
3277 })?;
3278 Ok(RestorePathState::Symlink { target })
3279 } else if metadata.is_file() {
3280 let content_bytes = std::fs::read(path).map_err(|error| AftError::IoError {
3281 path: path.display().to_string(),
3282 message: error.to_string(),
3283 })?;
3284 Ok(RestorePathState::Regular {
3285 content_bytes,
3286 mode: file_mode(&metadata),
3287 })
3288 } else {
3289 Ok(RestorePathState::Directory)
3290 }
3291}
3292
3293fn restore_entry_to_path(path: &Path, entry: &BackupEntry) -> std::io::Result<()> {
3294 match entry.kind {
3295 BackupEntryKind::Content => restore_regular_file(path, &entry.content_bytes, entry.mode),
3296 BackupEntryKind::Symlink => {
3297 let target = entry.link_target.as_ref().ok_or_else(|| {
3298 std::io::Error::new(
3299 std::io::ErrorKind::InvalidData,
3300 "symlink backup entry missing target",
3301 )
3302 })?;
3303 restore_symlink(path, target)
3304 }
3305 BackupEntryKind::Tombstone => remove_tombstone_path(path),
3306 }
3307}
3308
3309fn restore_path_state(path: &Path, state: &RestorePathState) -> bool {
3310 match state {
3311 RestorePathState::Missing => remove_file_or_symlink_if_present(path).is_ok(),
3312 RestorePathState::Regular {
3313 content_bytes,
3314 mode,
3315 } => restore_regular_file(path, content_bytes, *mode).is_ok(),
3316 RestorePathState::Symlink { target } => restore_symlink(path, target).is_ok(),
3317 RestorePathState::Directory => true,
3318 }
3319}
3320
3321fn restore_regular_file(
3322 path: &Path,
3323 content_bytes: &[u8],
3324 mode: Option<u32>,
3325) -> std::io::Result<()> {
3326 if let Some(parent) = path.parent() {
3327 if !parent.as_os_str().is_empty() {
3328 std::fs::create_dir_all(parent)?;
3329 }
3330 }
3331 if std::fs::symlink_metadata(path)
3332 .map(|metadata| metadata.file_type().is_symlink())
3333 .unwrap_or(false)
3334 {
3335 std::fs::remove_file(path)?;
3336 }
3337 std::fs::write(path, content_bytes)?;
3338 set_file_mode(path, mode)
3339}
3340
3341fn restore_symlink(path: &Path, target: &Path) -> std::io::Result<()> {
3342 if let Some(parent) = path.parent() {
3343 if !parent.as_os_str().is_empty() {
3344 std::fs::create_dir_all(parent)?;
3345 }
3346 }
3347 remove_file_or_symlink_if_present(path)?;
3348 create_symlink(target, path)
3349}
3350
3351#[cfg(unix)]
3352fn create_symlink(target: &Path, link: &Path) -> std::io::Result<()> {
3353 std::os::unix::fs::symlink(target, link)
3354}
3355
3356#[cfg(windows)]
3357fn create_symlink(target: &Path, link: &Path) -> std::io::Result<()> {
3358 if target.is_dir() {
3359 std::os::windows::fs::symlink_dir(target, link)
3360 } else {
3361 std::os::windows::fs::symlink_file(target, link)
3362 }
3363}
3364
3365fn remove_tombstone_path(path: &Path) -> std::io::Result<()> {
3366 match std::fs::symlink_metadata(path) {
3367 Ok(metadata) if metadata.file_type().is_symlink() || metadata.is_file() => {
3368 std::fs::remove_file(path)
3369 }
3370 Ok(_) => Err(std::io::Error::new(
3371 std::io::ErrorKind::IsADirectory,
3372 "tombstone target is a directory",
3373 )),
3374 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
3375 Err(error) => Err(error),
3376 }
3377}
3378
3379fn remove_file_or_symlink_if_present(path: &Path) -> std::io::Result<()> {
3380 match std::fs::symlink_metadata(path) {
3381 Ok(metadata) if metadata.file_type().is_symlink() || metadata.is_file() => {
3382 std::fs::remove_file(path)
3383 }
3384 Ok(_) => Err(std::io::Error::new(
3385 std::io::ErrorKind::IsADirectory,
3386 "path is a directory",
3387 )),
3388 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
3389 Err(error) => Err(error),
3390 }
3391}
3392
3393fn read_entry_disk_metadata(
3394 backup_path: &Path,
3395 backup_id: &str,
3396) -> Option<BackupEntryDiskMetadata> {
3397 let meta_path = if backup_path.file_name().and_then(|name| name.to_str()) == Some("meta.json") {
3398 backup_path.to_path_buf()
3399 } else {
3400 backup_path.parent()?.join("meta.json")
3401 };
3402 let content = std::fs::read_to_string(meta_path).ok()?;
3403 let meta: serde_json::Value = serde_json::from_str(&content).ok()?;
3404 let entries = meta.get("entries")?.as_array()?;
3405 let entry = entries
3406 .iter()
3407 .find(|entry| entry.get("backup_id").and_then(|value| value.as_str()) == Some(backup_id))?;
3408 Some(restore_metadata_fields(entry))
3409}
3410
3411fn rollback_transactional_restore(
3412 written: &[(PathBuf, RestorePathState)],
3413 attempted: Option<(&PathBuf, &RestorePathState)>,
3414) -> bool {
3415 let mut ok = true;
3416
3417 if let Some((path, state)) = attempted {
3418 ok &= restore_path_state(path, state);
3419 }
3420
3421 for (path, state) in written.iter().rev() {
3422 ok &= restore_path_state(path, state);
3423 }
3424
3425 ok
3426}
3427
3428fn rollback_deleted_tombstones(deleted: &[(PathBuf, RestorePathState)]) -> bool {
3429 let mut ok = true;
3430 for (path, state) in deleted.iter().rev() {
3431 ok &= restore_path_state(path, state);
3432 }
3433 ok
3434}
3435
3436fn missing_parent_dirs(parent: &Path) -> Vec<PathBuf> {
3437 let mut dirs = Vec::new();
3438 let mut current = Some(parent);
3439
3440 while let Some(dir) = current {
3441 if dir.as_os_str().is_empty() || dir.exists() {
3442 break;
3443 }
3444 dirs.push(dir.to_path_buf());
3445 current = dir.parent();
3446 }
3447
3448 dirs
3449}
3450
3451fn rollback_created_dirs(dirs: &[PathBuf]) -> bool {
3452 let mut dirs = dirs.to_vec();
3453 dirs.sort_by_key(|dir| std::cmp::Reverse(dir.components().count()));
3454 dirs.dedup();
3455
3456 let mut ok = true;
3457 for dir in dirs {
3458 match std::fs::remove_dir(&dir) {
3459 Ok(()) => {}
3460 Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
3461 Err(_) => ok = false,
3462 }
3463 }
3464
3465 ok
3466}
3467
3468fn remove_created_dirs_best_effort(dirs: &[PathBuf]) {
3469 let mut dirs = dirs.to_vec();
3470 dirs.sort_by_key(|dir| std::cmp::Reverse(dir.components().count()));
3471 dirs.dedup();
3472
3473 for dir in dirs {
3474 match std::fs::remove_dir(&dir) {
3475 Ok(()) => {}
3476 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
3477 Err(_) => {}
3478 }
3479 }
3480}
3481
3482fn dir_has_entries(path: &Path) -> bool {
3483 std::fs::read_dir(path)
3484 .map(|mut entries| entries.next().is_some())
3485 .unwrap_or(false)
3486}
3487
3488fn current_timestamp() -> u64 {
3489 std::time::SystemTime::now()
3490 .duration_since(std::time::UNIX_EPOCH)
3491 .unwrap_or_default()
3492 .as_secs()
3493}
3494
3495fn current_timestamp_nanos() -> u64 {
3496 let nanos = std::time::SystemTime::now()
3497 .duration_since(std::time::UNIX_EPOCH)
3498 .unwrap_or_default()
3499 .as_nanos();
3500 nanos.min(u128::from(u64::MAX)) as u64
3501}
3502
3503fn legacy_entry_order(timestamp_secs: u64, backup_id: &str) -> u128 {
3504 let nanos = timestamp_secs.saturating_mul(1_000_000_000);
3505 ((nanos as u128) << 32) | u128::from(backup_sequence(backup_id).unwrap_or(0))
3506}
3507
3508fn parse_order_value(value: &serde_json::Value) -> Option<u128> {
3509 value
3510 .as_str()
3511 .and_then(|s| s.parse::<u128>().ok())
3512 .or_else(|| value.as_u64().map(u128::from))
3513}
3514
3515fn is_v2_meta(meta: &serde_json::Value) -> bool {
3516 meta.get("format_version").and_then(|value| value.as_str()) == Some(V2_FORMAT_VERSION)
3517}
3518
3519fn meta_entries(meta: &serde_json::Value) -> Result<&Vec<serde_json::Value>, String> {
3520 meta.get("entries")
3521 .and_then(|value| value.as_array())
3522 .ok_or_else(|| "backup meta missing entries array".to_string())
3523}
3524
3525fn meta_entry_count(meta: &serde_json::Value) -> Option<usize> {
3526 if is_v2_meta(meta) {
3527 return meta
3528 .get("entries")
3529 .and_then(|value| value.as_array())
3530 .map(Vec::len);
3531 }
3532 meta.get("count")
3533 .and_then(|value| value.as_u64())
3534 .and_then(|count| usize::try_from(count).ok())
3535 .or_else(|| {
3536 meta.get("entries")
3537 .and_then(|value| value.as_array())
3538 .map(Vec::len)
3539 })
3540}
3541
3542fn entry_kind_from_meta(entry_meta: Option<&serde_json::Value>) -> BackupEntryKind {
3543 match entry_meta
3544 .and_then(|meta| meta.get("kind"))
3545 .and_then(|value| value.as_str())
3546 {
3547 Some("tombstone") => BackupEntryKind::Tombstone,
3548 Some("symlink") => BackupEntryKind::Symlink,
3549 _ => BackupEntryKind::Content,
3550 }
3551}
3552
3553fn backup_head_from_meta(entry_meta: Option<&serde_json::Value>, index: usize) -> BackupEntryHead {
3554 let backup_id = entry_backup_id(entry_meta, index);
3555 let timestamp = entry_meta
3556 .and_then(|meta| meta.get("timestamp"))
3557 .and_then(|value| value.as_u64())
3558 .unwrap_or(0);
3559 let order = entry_meta
3560 .and_then(|meta| meta.get("order"))
3561 .and_then(parse_order_value)
3562 .unwrap_or_else(|| legacy_entry_order(timestamp, &backup_id));
3563 BackupEntryHead {
3564 order,
3565 op_id: entry_meta
3566 .and_then(|meta| meta.get("op_id"))
3567 .and_then(|value| value.as_str())
3568 .map(str::to_string),
3569 }
3570}
3571
3572fn entry_backup_id(entry_meta: Option<&serde_json::Value>, index: usize) -> String {
3573 entry_meta
3574 .and_then(|meta| meta.get("backup_id"))
3575 .and_then(|value| value.as_str())
3576 .map(str::to_string)
3577 .unwrap_or_else(|| format!("disk-{}", index))
3578}
3579
3580fn entry_from_meta(
3581 entry_meta: Option<&serde_json::Value>,
3582 index: usize,
3583 kind: BackupEntryKind,
3584 content_bytes: Vec<u8>,
3585) -> BackupEntry {
3586 let backup_id = entry_backup_id(entry_meta, index);
3587 let timestamp = entry_meta
3588 .and_then(|meta| meta.get("timestamp"))
3589 .and_then(|value| value.as_u64())
3590 .unwrap_or(0);
3591 let order = entry_meta
3592 .and_then(|meta| meta.get("order"))
3593 .and_then(parse_order_value)
3594 .unwrap_or_else(|| legacy_entry_order(timestamp, &backup_id));
3595 let link_target = if kind == BackupEntryKind::Symlink {
3596 entry_meta
3597 .and_then(|meta| meta.get("link_target"))
3598 .and_then(|value| value.as_str())
3599 .map(PathBuf::from)
3600 .or_else(|| {
3601 Some(PathBuf::from(
3602 String::from_utf8_lossy(&content_bytes).into_owned(),
3603 ))
3604 })
3605 } else {
3606 None
3607 };
3608 let content = match kind {
3609 BackupEntryKind::Content => String::from_utf8_lossy(&content_bytes).into_owned(),
3610 BackupEntryKind::Symlink => link_target
3611 .as_ref()
3612 .map(|target| target.display().to_string())
3613 .unwrap_or_default(),
3614 BackupEntryKind::Tombstone => String::new(),
3615 };
3616 BackupEntry {
3617 backup_id,
3618 content,
3619 content_bytes: content_bytes.into(),
3620 timestamp,
3621 order,
3622 description: entry_meta
3623 .and_then(|meta| meta.get("description"))
3624 .and_then(|value| value.as_str())
3625 .unwrap_or("restored from disk")
3626 .to_string(),
3627 op_id: entry_meta
3628 .and_then(|meta| meta.get("op_id"))
3629 .and_then(|value| value.as_str())
3630 .map(str::to_string),
3631 kind,
3632 mode: entry_meta
3633 .and_then(|meta| meta.get("mode"))
3634 .and_then(|value| value.as_u64())
3635 .and_then(|mode| u32::try_from(mode).ok()),
3636 link_target,
3637 created_dirs: entry_meta
3638 .and_then(|meta| meta.get("created_dirs"))
3639 .and_then(|value| value.as_array())
3640 .map(|dirs| {
3641 dirs.iter()
3642 .filter_map(|dir| dir.as_str())
3643 .map(PathBuf::from)
3644 .collect()
3645 })
3646 .unwrap_or_default(),
3647 }
3648}
3649
3650fn legacy_entry_from_meta(
3651 dir: &Path,
3652 entry_meta: Option<&serde_json::Value>,
3653 index: usize,
3654) -> Option<BackupEntry> {
3655 let kind = entry_kind_from_meta(entry_meta);
3656 let content_bytes = match kind {
3657 BackupEntryKind::Content | BackupEntryKind::Symlink => {
3658 std::fs::read(dir.join(format!("{}.bak", index))).ok()?
3659 }
3660 BackupEntryKind::Tombstone => Vec::new(),
3661 };
3662 Some(entry_from_meta(entry_meta, index, kind, content_bytes))
3663}
3664
3665fn content_path_from_meta(entry_meta: &serde_json::Value) -> Result<&str, String> {
3666 let value = entry_meta
3667 .get("content_path")
3668 .and_then(|value| value.as_str())
3669 .ok_or_else(|| "v2 backup entry missing content_path".to_string())?;
3670 let path = Path::new(value);
3671 let mut components = path.components();
3672 match (components.next(), components.next()) {
3673 (Some(std::path::Component::Normal(_)), None) => Ok(value),
3674 _ => Err(format!("invalid backup content_path '{value}'")),
3675 }
3676}
3677
3678fn sanitize_backup_id(value: &str) -> String {
3679 value
3680 .chars()
3681 .map(|ch| {
3682 if ch.is_ascii_alphanumeric() || ch == '-' || ch == '_' {
3683 ch
3684 } else {
3685 '_'
3686 }
3687 })
3688 .collect()
3689}
3690
3691fn content_filename_for_entry(entry: &BackupEntry) -> Option<String> {
3692 match entry.kind {
3693 BackupEntryKind::Content | BackupEntryKind::Symlink => Some(format!(
3694 "bak_{}_{}.bak",
3695 entry.order,
3696 sanitize_backup_id(&entry.backup_id)
3697 )),
3698 BackupEntryKind::Tombstone => None,
3699 }
3700}
3701
3702fn content_bytes_for_disk(entry: &BackupEntry) -> Cow<'_, [u8]> {
3703 match entry.kind {
3704 BackupEntryKind::Content => Cow::Borrowed(&entry.content_bytes),
3705 BackupEntryKind::Symlink => Cow::Owned(
3706 entry
3707 .link_target
3708 .as_ref()
3709 .map(|target| target.as_os_str().to_string_lossy().as_bytes().to_vec())
3710 .unwrap_or_default(),
3711 ),
3712 BackupEntryKind::Tombstone => Cow::Borrowed(&[]),
3713 }
3714}
3715
3716fn entry_meta_json(entry: &BackupEntry) -> serde_json::Value {
3717 serde_json::json!({
3718 "backup_id": entry.backup_id,
3719 "timestamp": entry.timestamp,
3720 "order": entry.order.to_string(),
3721 "description": entry.description,
3722 "op_id": entry.op_id,
3723 "kind": match entry.kind {
3724 BackupEntryKind::Content => "content",
3725 BackupEntryKind::Symlink => "symlink",
3726 BackupEntryKind::Tombstone => "tombstone",
3727 },
3728 "content_path": content_filename_for_entry(entry),
3729 "mode": entry.mode,
3730 "link_target": entry.link_target.as_ref().map(|target| target.display().to_string()),
3731 "created_dirs": entry
3732 .created_dirs
3733 .iter()
3734 .map(|dir| dir.display().to_string())
3735 .collect::<Vec<_>>(),
3736 })
3737}
3738
3739fn drain_stack_to_depth(stack: &mut Vec<BackupEntry>, max_depth: usize) -> Vec<BackupEntry> {
3740 let overflow = stack.len().saturating_sub(max_depth);
3741 stack.drain(..overflow).collect()
3742}
3743
3744fn trim_stack_to_depth(stack: &mut Vec<BackupEntry>, max_depth: usize) {
3745 let overflow = stack.len().saturating_sub(max_depth);
3746 drop(stack.drain(..overflow));
3747}
3748
3749fn write_temp_fsync_rename(dir: &Path, final_name: &str, content: &[u8]) -> std::io::Result<()> {
3750 let tmp_name = format!(
3751 ".{}.{}.{}.tmp",
3752 final_name,
3753 std::process::id(),
3754 current_timestamp_nanos()
3755 );
3756 let tmp_path = dir.join(tmp_name);
3757 let final_path = dir.join(final_name);
3758 {
3759 let mut file = std::fs::OpenOptions::new()
3760 .write(true)
3761 .create_new(true)
3762 .open(&tmp_path)?;
3763 file.write_all(content)?;
3764 file.sync_all()?;
3765 }
3766 replace_file(&tmp_path, &final_path)
3767}
3768
3769fn replace_file(from: &Path, to: &Path) -> std::io::Result<()> {
3770 std::fs::rename(from, to)
3773}
3774
3775#[cfg(unix)]
3776fn fsync_dir(path: &Path) -> std::io::Result<()> {
3777 std::fs::File::open(path)?.sync_all()
3778}
3779
3780#[cfg(not(unix))]
3781fn fsync_dir(_path: &Path) -> std::io::Result<()> {
3782 Ok(())
3790}
3791
3792fn prune_unreferenced_backup_files(
3793 dir: &Path,
3794 referenced: &HashSet<String>,
3795) -> std::io::Result<()> {
3796 for entry in std::fs::read_dir(dir)? {
3797 let entry = entry?;
3798 let path = entry.path();
3799 if !path.is_file() {
3800 continue;
3801 }
3802 let Some(name) = path.file_name().and_then(|name| name.to_str()) else {
3803 continue;
3804 };
3805 let is_backup_content = (name.starts_with("bak_") && name.ends_with(".bak"))
3806 || legacy_numeric_backup_name(name);
3807 let is_temp = name.ends_with(".tmp") || name.contains(".tmp.");
3808 if is_temp || (is_backup_content && !referenced.contains(name)) {
3809 let _ = std::fs::remove_file(path);
3810 }
3811 }
3812 Ok(())
3813}
3814
3815fn legacy_numeric_backup_name(name: &str) -> bool {
3816 name.strip_suffix(".bak")
3817 .is_some_and(|stem| !stem.is_empty() && stem.chars().all(|ch| ch.is_ascii_digit()))
3818}
3819
3820fn is_loadable_backup_path(key: &Path, path_dir: &Path) -> bool {
3821 if !key.is_absolute()
3822 || key
3823 .components()
3824 .any(|c| matches!(c, std::path::Component::ParentDir))
3825 {
3826 return false;
3827 }
3828 let Some(dir_name) = path_dir.file_name().and_then(|name| name.to_str()) else {
3829 return false;
3830 };
3831 BackupStore::path_hash(key) == dir_name
3832}
3833
3834fn stable_hash_16(bytes: &[u8]) -> String {
3835 let digest = Sha256::digest(bytes);
3836 digest[..8]
3837 .iter()
3838 .map(|byte| format!("{:02x}", byte))
3839 .collect()
3840}
3841
3842fn backup_sequence(backup_id: &str) -> Option<u64> {
3843 backup_id
3844 .strip_prefix("backup-")
3845 .or_else(|| backup_id.strip_prefix("disk-"))
3846 .and_then(|s| s.parse().ok())
3847}
3848
3849#[cfg(test)]
3850mod tests {
3851 use super::*;
3852 use crate::harness::Harness;
3853 use crate::protocol::DEFAULT_SESSION_ID;
3854 use std::fs;
3855 #[cfg(unix)]
3856 use std::os::unix::fs::PermissionsExt;
3857 use std::sync::{Arc, LazyLock, Mutex};
3858
3859 const DB_MIRROR_HARNESS: &str = "opencode";
3860 const DB_MIRROR_SESSION: &str = "db-mirror-session";
3861 const DB_MIRROR_PROJECT: &str = "db-mirror-project";
3862 const DB_MIRROR_FILE: &str = "/project/src/file.rs";
3863 const DB_MIRROR_PATH_HASH: &str = "db-mirror-path";
3864
3865 static MIRROR_SQL_TRACE: LazyLock<Mutex<Vec<String>>> =
3866 LazyLock::new(|| Mutex::new(Vec::new()));
3867
3868 fn capture_mirror_sql(sql: &str) {
3869 MIRROR_SQL_TRACE.lock().unwrap().push(sql.to_string());
3870 }
3871
3872 fn temp_file(name: &str, content: &str) -> PathBuf {
3873 let dir = std::env::temp_dir().join("aft_backup_tests");
3874 fs::create_dir_all(&dir).unwrap();
3875 let path = dir.join(name);
3876 fs::write(&path, content).unwrap();
3877 path
3878 }
3879
3880 fn db_mirror_context() -> DbMirrorContext<'static> {
3881 DbMirrorContext {
3882 harness: DB_MIRROR_HARNESS,
3883 session: DB_MIRROR_SESSION,
3884 project_key: DB_MIRROR_PROJECT,
3885 file_path: DB_MIRROR_FILE,
3886 path_hash: DB_MIRROR_PATH_HASH,
3887 }
3888 }
3889
3890 fn db_mirror_entry(index: u128, kind: BackupEntryKind) -> BackupEntry {
3891 BackupEntry {
3892 backup_id: format!("backup-{index}"),
3893 content: format!("content-{index}"),
3894 content_bytes: format!("content-{index}").into_bytes().into(),
3895 timestamp: u64::try_from(index + 1).unwrap(),
3896 order: index + 1,
3897 description: format!("generated-{index}"),
3898 op_id: index.is_multiple_of(3).then(|| format!("op-{}", index % 5)),
3899 kind,
3900 mode: None,
3901 link_target: None,
3902 created_dirs: Vec::new(),
3903 }
3904 }
3905
3906 fn db_mirror_rows(conn: &Connection) -> Vec<BackupRow> {
3907 crate::db::backups::list_backups(
3908 conn,
3909 DB_MIRROR_HARNESS,
3910 DB_MIRROR_SESSION,
3911 DB_MIRROR_PATH_HASH,
3912 )
3913 .unwrap()
3914 }
3915
3916 #[test]
3917 fn backup_db_restore_metadata_round_trips_all_fields() {
3918 let dir = tempfile::tempdir().unwrap();
3919 let conn = crate::db::open(&dir.path().join("aft.db")).unwrap();
3920 let context = db_mirror_context();
3921 let mut entry = db_mirror_entry(0, BackupEntryKind::Symlink);
3922 entry.mode = Some(0o100755);
3923 entry.link_target = Some(PathBuf::from("../target"));
3924 entry.created_dirs = vec![PathBuf::from("/project/new"), PathBuf::from("/project")];
3925
3926 crate::db::backups::insert_backup(&conn, &backup_row_for_db(&entry, &context)).unwrap();
3927 let row = db_mirror_rows(&conn).pop().unwrap();
3928 let metadata = row
3929 .restore_meta
3930 .as_deref()
3931 .and_then(restore_metadata_from_json)
3932 .unwrap();
3933
3934 assert_eq!(
3935 metadata,
3936 BackupEntryDiskMetadata {
3937 mode: entry.mode,
3938 link_target: entry.link_target,
3939 created_dirs: entry.created_dirs,
3940 }
3941 );
3942 }
3943
3944 #[test]
3945 fn backup_db_append_delta_matches_full_rebuild_for_generated_sequences() {
3946 let context = db_mirror_context();
3947 let mut saw_eviction = false;
3948 let mut saw_tombstone = false;
3949
3950 for max_depth in [1, 2, 5, 8] {
3951 let delta_dir = tempfile::tempdir().unwrap();
3952 let full_dir = tempfile::tempdir().unwrap();
3953 let delta_conn = crate::db::open(&delta_dir.path().join("aft.db")).unwrap();
3954 let full_conn = crate::db::open(&full_dir.path().join("aft.db")).unwrap();
3955 let mut stack = Vec::new();
3956 let mut generated = 0x9e37_79b9_u64 ^ u64::try_from(max_depth).unwrap();
3957
3958 for step in 0..64_u128 {
3959 generated = generated
3960 .wrapping_mul(6_364_136_223_846_793_005)
3961 .wrapping_add(1_442_695_040_888_963_407);
3962 let kind = if generated.is_multiple_of(4) {
3963 saw_tombstone = true;
3964 BackupEntryKind::Tombstone
3965 } else {
3966 BackupEntryKind::Content
3967 };
3968 let entry = db_mirror_entry(step, kind);
3969 let evicted = drain_stack_to_depth(&mut stack, max_depth - 1);
3970 saw_eviction |= !evicted.is_empty();
3971 let evicted_orders = evicted.iter().map(|entry| entry.order).collect::<Vec<_>>();
3972 stack.push(entry);
3973
3974 apply_backup_append_delta_in_db(
3975 &delta_conn,
3976 &context,
3977 &evicted_orders,
3978 stack.last(),
3979 )
3980 .unwrap();
3981 replace_backup_stack_in_db(&full_conn, &context, &stack).unwrap();
3982
3983 assert_eq!(
3984 db_mirror_rows(&delta_conn),
3985 db_mirror_rows(&full_conn),
3986 "delta mirror drifted at depth {max_depth}, step {step}"
3987 );
3988 }
3989 }
3990
3991 assert!(saw_eviction, "generated sequence must exercise eviction");
3992 assert!(saw_tombstone, "generated sequence must exercise tombstones");
3993 }
3994
3995 #[test]
3996 fn backup_db_append_delta_executes_constant_dml_statement_count() {
3997 let project = tempfile::tempdir().unwrap();
3998 let storage = tempfile::tempdir().unwrap();
3999 let file = project.path().join("statement-count.txt");
4000 let mut store = BackupStore::new();
4001 store.set_storage_dir(storage.path().to_path_buf(), 72);
4002 store.set_db_harness(Harness::Opencode);
4003 store.set_db_project_key(DB_MIRROR_PROJECT.to_string());
4004 store.set_policy(BackupPolicy {
4005 enabled: true,
4006 max_depth: 4,
4007 max_file_size: None,
4008 });
4009 let shared = Arc::new(Mutex::new(
4010 crate::db::open(&storage.path().join("aft.db")).unwrap(),
4011 ));
4012 store.set_db_pool(shared.clone());
4013 for version in 0..4 {
4014 fs::write(&file, format!("version-{version}")).unwrap();
4015 store
4016 .snapshot(DB_MIRROR_SESSION, &file, "fill retained stack")
4017 .unwrap();
4018 }
4019
4020 MIRROR_SQL_TRACE.lock().unwrap().clear();
4021 shared.lock().unwrap().trace(Some(capture_mirror_sql));
4022 fs::write(&file, "version-4").unwrap();
4023 store
4024 .snapshot(DB_MIRROR_SESSION, &file, "measured append")
4025 .unwrap();
4026 shared.lock().unwrap().trace(None);
4027
4028 let traced = MIRROR_SQL_TRACE.lock().unwrap().clone();
4029 let dml = traced
4030 .iter()
4031 .filter(|sql| {
4032 let sql = sql.trim_start();
4033 sql.starts_with("DELETE FROM backups") || sql.starts_with("INSERT INTO backups")
4034 })
4035 .collect::<Vec<_>>();
4036 assert_eq!(
4037 dml.len(),
4038 2,
4039 "one eviction plus one append must execute exactly two DML statements: {traced:?}"
4040 );
4041 assert_eq!(
4042 dml.iter()
4043 .filter(|sql| sql.trim_start().starts_with("DELETE FROM backups"))
4044 .count(),
4045 1
4046 );
4047 assert_eq!(
4048 dml.iter()
4049 .filter(|sql| sql.trim_start().starts_with("INSERT INTO backups"))
4050 .count(),
4051 1
4052 );
4053 }
4054
4055 #[test]
4056 fn backup_db_append_delta_rolls_back_eviction_when_insert_fails() {
4057 let context = db_mirror_context();
4058 let dir = tempfile::tempdir().unwrap();
4059 let conn = crate::db::open(&dir.path().join("aft.db")).unwrap();
4060 let prior_stack = vec![
4061 db_mirror_entry(0, BackupEntryKind::Content),
4062 db_mirror_entry(1, BackupEntryKind::Content),
4063 ];
4064 replace_backup_stack_in_db(&conn, &context, &prior_stack).unwrap();
4065 let prior_rows = db_mirror_rows(&conn);
4066 conn.execute_batch(
4067 "CREATE TRIGGER fail_delta_insert
4068 BEFORE INSERT ON backups
4069 WHEN NEW.backup_id = 'backup-2'
4070 BEGIN
4071 SELECT RAISE(ABORT, 'forced delta insert failure');
4072 END;",
4073 )
4074 .unwrap();
4075
4076 let new_entry = db_mirror_entry(2, BackupEntryKind::Tombstone);
4077 let error = apply_backup_append_delta_in_db(
4078 &conn,
4079 &context,
4080 &[prior_stack[0].order],
4081 Some(&new_entry),
4082 )
4083 .unwrap_err();
4084
4085 assert!(error.to_string().contains("forced delta insert failure"));
4086 assert_eq!(
4087 db_mirror_rows(&conn),
4088 prior_rows,
4089 "failed append must roll back the preceding eviction"
4090 );
4091 }
4092
4093 #[test]
4094 fn backup_db_unknown_mirror_repairs_disk_history_before_append_deltas() {
4095 let project = tempfile::tempdir().unwrap();
4096 let storage = tempfile::tempdir().unwrap();
4097 let file = project.path().join("repair-before-delta.txt");
4098 let mut disk_only = BackupStore::new();
4099 disk_only.set_storage_dir(storage.path().to_path_buf(), 72);
4100 fs::write(&file, "v1").unwrap();
4101 disk_only
4102 .snapshot(DB_MIRROR_SESSION, &file, "disk first")
4103 .unwrap();
4104 fs::write(&file, "v2").unwrap();
4105 disk_only
4106 .snapshot(DB_MIRROR_SESSION, &file, "disk second")
4107 .unwrap();
4108
4109 let mut mirrored = BackupStore::new();
4110 mirrored.set_storage_dir(storage.path().to_path_buf(), 72);
4111 mirrored.set_db_harness(Harness::Opencode);
4112 mirrored.set_db_project_key(DB_MIRROR_PROJECT.to_string());
4113 let conn = Arc::new(Mutex::new(
4114 crate::db::open(&storage.path().join("aft.db")).unwrap(),
4115 ));
4116 mirrored.set_db_pool(conn.clone());
4117 fs::write(&file, "v3").unwrap();
4118 mirrored
4119 .snapshot(DB_MIRROR_SESSION, &file, "first mirrored append")
4120 .unwrap();
4121
4122 let key = canonicalize_key(&file);
4123 let rows = crate::db::backups::list_backups(
4124 &conn.lock().unwrap(),
4125 DB_MIRROR_HARNESS,
4126 DB_MIRROR_SESSION,
4127 &BackupStore::path_hash(&key),
4128 )
4129 .unwrap();
4130 assert_eq!(
4131 rows.iter()
4132 .map(|row| row.description.as_str())
4133 .collect::<Vec<_>>(),
4134 vec!["disk first", "disk second", "first mirrored append"]
4135 );
4136 }
4137
4138 #[test]
4139 fn snapshot_and_restore_round_trip() {
4140 let path = temp_file("round_trip.txt", "original");
4141 let mut store = BackupStore::new();
4142
4143 let id = store
4144 .snapshot(DEFAULT_SESSION_ID, &path, "before edit")
4145 .unwrap()
4146 .unwrap();
4147 assert!(id.starts_with("backup-"));
4148
4149 fs::write(&path, "modified").unwrap();
4150 assert_eq!(fs::read_to_string(&path).unwrap(), "modified");
4151
4152 let (entry, _) = store.restore_latest(DEFAULT_SESSION_ID, &path).unwrap();
4153 assert_eq!(entry.content, "original");
4154 assert_eq!(fs::read_to_string(&path).unwrap(), "original");
4155 }
4156
4157 #[test]
4158 fn multiple_snapshots_preserve_order() {
4159 let path = temp_file("order.txt", "v1");
4160 let mut store = BackupStore::new();
4161
4162 store.snapshot(DEFAULT_SESSION_ID, &path, "first").unwrap();
4163 fs::write(&path, "v2").unwrap();
4164 store.snapshot(DEFAULT_SESSION_ID, &path, "second").unwrap();
4165 fs::write(&path, "v3").unwrap();
4166 store.snapshot(DEFAULT_SESSION_ID, &path, "third").unwrap();
4167
4168 let history = store.history(DEFAULT_SESSION_ID, &path);
4169 assert_eq!(history.len(), 3);
4170 assert_eq!(history[0].content, "v1");
4171 assert_eq!(history[1].content, "v2");
4172 assert_eq!(history[2].content, "v3");
4173 }
4174
4175 #[test]
4176 fn restore_pops_from_stack() {
4177 let path = temp_file("pop.txt", "v1");
4178 let mut store = BackupStore::new();
4179
4180 store.snapshot(DEFAULT_SESSION_ID, &path, "first").unwrap();
4181 fs::write(&path, "v2").unwrap();
4182 store.snapshot(DEFAULT_SESSION_ID, &path, "second").unwrap();
4183
4184 let (entry, _) = store.restore_latest(DEFAULT_SESSION_ID, &path).unwrap();
4185 assert_eq!(entry.description, "second");
4186 assert_eq!(entry.content, "v2");
4187
4188 let history = store.history(DEFAULT_SESSION_ID, &path);
4189 assert_eq!(history.len(), 1);
4190 }
4191
4192 #[test]
4193 fn empty_history_returns_empty_vec() {
4194 let store = BackupStore::new();
4195 let path = Path::new("/tmp/aft_backup_tests/nonexistent_history.txt");
4196 assert!(store.history(DEFAULT_SESSION_ID, path).is_empty());
4197 }
4198
4199 #[test]
4200 fn snapshot_nonexistent_file_returns_error() {
4201 let mut store = BackupStore::new();
4202 let path = Path::new("/tmp/aft_backup_tests/absolutely_does_not_exist.txt");
4203 assert!(store.snapshot(DEFAULT_SESSION_ID, path, "test").is_err());
4204 }
4205
4206 #[test]
4207 fn tracked_files_lists_snapshotted_paths() {
4208 let path1 = temp_file("tracked1.txt", "a");
4209 let path2 = temp_file("tracked2.txt", "b");
4210 let mut store = BackupStore::new();
4211
4212 store.snapshot(DEFAULT_SESSION_ID, &path1, "snap1").unwrap();
4213 store.snapshot(DEFAULT_SESSION_ID, &path2, "snap2").unwrap();
4214 assert_eq!(store.tracked_files(DEFAULT_SESSION_ID).len(), 2);
4215 }
4216
4217 #[test]
4218 fn sessions_are_isolated() {
4219 let path = temp_file("isolated.txt", "original");
4220 let mut store = BackupStore::new();
4221
4222 store.snapshot("session_a", &path, "a's snapshot").unwrap();
4223
4224 assert!(store.history("session_b", &path).is_empty());
4226 assert_eq!(store.tracked_files("session_b").len(), 0);
4227
4228 let err = store.restore_latest("session_b", &path);
4230 assert!(matches!(err, Err(AftError::NoUndoHistory { .. })));
4231
4232 assert_eq!(store.history("session_a", &path).len(), 1);
4234 assert_eq!(store.tracked_files("session_a").len(), 1);
4235 }
4236
4237 #[test]
4238 fn per_session_per_file_cap_is_independent() {
4239 let path = temp_file("cap_indep.txt", "v0");
4242 let mut store = BackupStore::new();
4243
4244 for i in 0..(MAX_UNDO_DEPTH + 5) {
4245 fs::write(&path, format!("a{}", i)).unwrap();
4246 store.snapshot("session_a", &path, "a").unwrap();
4247 }
4248 fs::write(&path, "b_initial").unwrap();
4249 store.snapshot("session_b", &path, "b").unwrap();
4250
4251 assert_eq!(store.history("session_a", &path).len(), MAX_UNDO_DEPTH);
4253 assert_eq!(store.history("session_b", &path).len(), 1);
4255 }
4256
4257 #[test]
4258 fn sessions_with_backups_lists_all_namespaces() {
4259 let path_a = temp_file("sessions_list_a.txt", "a");
4260 let path_b = temp_file("sessions_list_b.txt", "b");
4261 let mut store = BackupStore::new();
4262
4263 store.snapshot("alice", &path_a, "from alice").unwrap();
4264 store.snapshot("bob", &path_b, "from bob").unwrap();
4265
4266 let sessions = store.sessions_with_backups();
4267 assert_eq!(sessions.len(), 2);
4268 assert!(sessions.iter().any(|s| s == "alice"));
4269 assert!(sessions.iter().any(|s| s == "bob"));
4270 }
4271
4272 #[test]
4273 fn disk_persistence_survives_reload() {
4274 let dir = std::env::temp_dir().join("aft_backup_disk_test");
4275 let _ = fs::remove_dir_all(&dir);
4276 fs::create_dir_all(&dir).unwrap();
4277
4278 let file_path = temp_file("disk_persist.txt", "original");
4279
4280 {
4282 let mut store = BackupStore::new();
4283 store.set_storage_dir(dir.clone(), 72);
4284 store
4285 .snapshot(DEFAULT_SESSION_ID, &file_path, "before edit")
4286 .unwrap();
4287 }
4288
4289 fs::write(&file_path, "externally modified").unwrap();
4291
4292 let mut store2 = BackupStore::new();
4294 store2.set_storage_dir(dir.clone(), 72);
4295
4296 let (entry, warning) = store2
4297 .restore_latest(DEFAULT_SESSION_ID, &file_path)
4298 .unwrap();
4299 assert_eq!(entry.content, "original");
4300 assert!(warning.is_some()); assert_eq!(fs::read_to_string(&file_path).unwrap(), "original");
4302
4303 let _ = fs::remove_dir_all(&dir);
4304 }
4305
4306 #[test]
4307 fn snapshot_after_restart_preserves_history_and_unique_ids() {
4308 let dir = std::env::temp_dir().join("aft_backup_restart_history_test");
4314 let _ = fs::remove_dir_all(&dir);
4315 fs::create_dir_all(&dir).unwrap();
4316 let file_path = temp_file("restart_history.txt", "v0");
4317
4318 let first_id = {
4320 let mut store = BackupStore::new();
4321 store.set_storage_dir(dir.clone(), 72);
4322 let id = store
4323 .snapshot(DEFAULT_SESSION_ID, &file_path, "edit 1")
4324 .unwrap()
4325 .unwrap();
4326 fs::write(&file_path, "v1").unwrap();
4327 id
4328 };
4329
4330 let second_id = {
4333 let mut store = BackupStore::new();
4334 store.set_storage_dir(dir.clone(), 72);
4335 let id = store
4336 .snapshot(DEFAULT_SESSION_ID, &file_path, "edit 2")
4337 .unwrap()
4338 .unwrap();
4339 fs::write(&file_path, "v2").unwrap();
4340 id
4341 };
4342
4343 assert_ne!(
4346 first_id, second_id,
4347 "post-restart snapshot reused backup id {first_id}"
4348 );
4349
4350 let mut store = BackupStore::new();
4353 store.set_storage_dir(dir.clone(), 72);
4354 assert_eq!(
4355 store.history(DEFAULT_SESSION_ID, &file_path).len(),
4356 2,
4357 "prior history was overwritten by the post-restart snapshot"
4358 );
4359
4360 let (entry1, _) = store
4361 .restore_latest(DEFAULT_SESSION_ID, &file_path)
4362 .unwrap();
4363 assert_eq!(entry1.content, "v1", "first undo should restore v1");
4364 let (entry0, _) = store
4365 .restore_latest(DEFAULT_SESSION_ID, &file_path)
4366 .unwrap();
4367 assert_eq!(entry0.content, "v0", "second undo should restore v0");
4368
4369 let _ = fs::remove_dir_all(&dir);
4370 }
4371
4372 #[test]
4373 fn fresh_store_defers_backup_io_until_first_snapshot_and_preserves_undo() {
4374 let project = tempfile::tempdir().unwrap();
4375 let storage = tempfile::tempdir().unwrap();
4376 let path = project.path().join("lazy-history.txt");
4377 fs::write(&path, "v0").unwrap();
4378
4379 {
4380 let mut store = BackupStore::new();
4381 store.set_storage_dir(storage.path().to_path_buf(), 72);
4382 assert_eq!(store.disk_io_count_for_tests(), 0);
4383 store.snapshot("session-a", &path, "captures v0").unwrap();
4384 fs::write(&path, "v1").unwrap();
4385 }
4386
4387 let mut fresh = BackupStore::new();
4388 fresh.set_storage_dir(storage.path().to_path_buf(), 72);
4389 assert_eq!(
4390 fresh.disk_io_count_for_tests(),
4391 0,
4392 "binding a fresh store must not inspect backup directories"
4393 );
4394
4395 fresh.snapshot("session-a", &path, "captures v1").unwrap();
4396 assert!(fresh.disk_io_count_for_tests() > 0);
4397 fs::write(&path, "v2").unwrap();
4398
4399 fresh.restore_latest("session-a", &path).unwrap();
4400 assert_eq!(fs::read_to_string(&path).unwrap(), "v1");
4401 fresh.restore_latest("session-a", &path).unwrap();
4402 assert_eq!(fs::read_to_string(&path).unwrap(), "v0");
4403 }
4404
4405 #[test]
4406 fn same_namespace_bind_is_idempotent_and_session_history_survives() {
4407 let project = tempfile::tempdir().unwrap();
4408 let storage = tempfile::tempdir().unwrap();
4409 let path = project.path().join("session-isolation.txt");
4410 fs::write(&path, "v0").unwrap();
4411
4412 let mut store = BackupStore::new();
4413 store.set_storage_dir_for_harness(storage.path().to_path_buf(), Harness::Opencode, 72);
4414 assert_eq!(store.disk_io_count_for_tests(), 0);
4415 store.snapshot("session-a", &path, "session A").unwrap();
4416 fs::write(&path, "v1").unwrap();
4417
4418 let io_before_rebind = store.disk_io_count_for_tests();
4419 store.set_storage_dir_for_harness(storage.path().to_path_buf(), Harness::Opencode, 72);
4420 assert_eq!(store.disk_io_count_for_tests(), io_before_rebind);
4421 assert_eq!(store.disk_history_count("session-a", &path), 1);
4422
4423 store.snapshot("session-b", &path, "session B").unwrap();
4424 fs::write(&path, "v2").unwrap();
4425 store.restore_latest("session-a", &path).unwrap();
4426 assert_eq!(fs::read_to_string(&path).unwrap(), "v0");
4427 assert_eq!(store.disk_history_count("session-b", &path), 1);
4428 }
4429
4430 #[test]
4431 fn legacy_flat_layout_migrates_to_default_session() {
4432 let dir = std::env::temp_dir().join("aft_backup_migration_test");
4435 let _ = fs::remove_dir_all(&dir);
4436 fs::create_dir_all(&dir).unwrap();
4437 let backups = dir.join("backups");
4438 fs::create_dir_all(&backups).unwrap();
4439
4440 let legacy_hash = "deadbeefcafebabe";
4442 let legacy_dir = backups.join(legacy_hash);
4443 fs::create_dir_all(&legacy_dir).unwrap();
4444 fs::write(legacy_dir.join("0.bak"), "original content").unwrap();
4445 let legacy_meta = serde_json::json!({
4446 "path": "/tmp/migrated_file.txt",
4447 "count": 1,
4448 });
4449 fs::write(
4450 legacy_dir.join("meta.json"),
4451 serde_json::to_string_pretty(&legacy_meta).unwrap(),
4452 )
4453 .unwrap();
4454
4455 let mut store = BackupStore::new();
4457 store.set_storage_dir(dir.clone(), 72);
4458 assert_eq!(store.disk_io_count_for_tests(), 0);
4459 store.run_process_maintenance_once();
4460
4461 let default_session_dir = backups.join(BackupStore::session_hash(DEFAULT_SESSION_ID));
4464 assert!(default_session_dir.exists());
4465 assert!(default_session_dir.join(legacy_hash).exists());
4466 assert!(!backups.join(legacy_hash).exists());
4467
4468 let meta_content =
4470 fs::read_to_string(default_session_dir.join(legacy_hash).join("meta.json")).unwrap();
4471 let meta: serde_json::Value = serde_json::from_str(&meta_content).unwrap();
4472 assert_eq!(meta["session_id"], DEFAULT_SESSION_ID);
4473 assert_eq!(meta["schema_version"], SCHEMA_VERSION);
4474
4475 let _ = fs::remove_dir_all(&dir);
4476 }
4477
4478 #[test]
4479 fn process_maintenance_removes_stale_backup_sessions() {
4480 let dir = std::env::temp_dir().join("aft_backup_gc_test");
4481 let _ = fs::remove_dir_all(&dir);
4482 let backups = dir.join("backups");
4483 fs::create_dir_all(&backups).unwrap();
4484
4485 let stale_session_dir = backups.join("stale-session");
4486 fs::create_dir_all(&stale_session_dir).unwrap();
4487 let stale_marker = serde_json::json!({
4488 "schema_version": SCHEMA_VERSION,
4489 "session_id": "stale",
4490 "last_accessed": 1,
4491 });
4492 fs::write(
4493 stale_session_dir.join("session.json"),
4494 serde_json::to_string_pretty(&stale_marker).unwrap(),
4495 )
4496 .unwrap();
4497
4498 let mut store = BackupStore::new();
4499 store.set_storage_dir(dir.clone(), 1);
4500 assert_eq!(store.disk_io_count_for_tests(), 0);
4501 store.run_process_maintenance_once();
4502
4503 assert!(!stale_session_dir.exists());
4504 let _ = fs::remove_dir_all(&dir);
4505 }
4506
4507 #[test]
4508 fn markerless_session_dir_is_skipped_not_mapped_to_default() {
4509 let dir = std::env::temp_dir().join("aft_backup_markerless_skip_test");
4510 let _ = fs::remove_dir_all(&dir);
4511 let file_path = temp_file("markerless.txt", "original");
4512 let key = canonicalize_key(&file_path);
4513 let path_dir = dir
4514 .join("backups")
4515 .join("corrupt-session")
4516 .join("path-entry");
4517 fs::create_dir_all(&path_dir).unwrap();
4518 fs::write(path_dir.join("0.bak"), "original").unwrap();
4519 fs::write(
4520 path_dir.join("meta.json"),
4521 serde_json::to_string_pretty(&serde_json::json!({
4522 "schema_version": SCHEMA_VERSION,
4523 "session_id": "lost-session",
4524 "path": key.display().to_string(),
4525 "count": 1,
4526 "entries": [{
4527 "backup_id": "disk-0",
4528 "timestamp": 0,
4529 "description": "corrupt marker test",
4530 "op_id": null,
4531 "kind": "content",
4532 }]
4533 }))
4534 .unwrap(),
4535 )
4536 .unwrap();
4537
4538 let mut store = BackupStore::new();
4539 store.set_storage_dir(dir.clone(), 72);
4540
4541 assert_eq!(store.disk_history_count(DEFAULT_SESSION_ID, &file_path), 0);
4542 assert!(store.sessions_with_backups().is_empty());
4543 let _ = fs::remove_dir_all(&dir);
4544 }
4545
4546 #[test]
4547 fn set_storage_dir_reconfiguration_drops_previous_disk_index() {
4548 let dir_a = std::env::temp_dir().join("aft_backup_storage_a_test");
4549 let dir_b = std::env::temp_dir().join("aft_backup_storage_b_test");
4550 let _ = fs::remove_dir_all(&dir_a);
4551 let _ = fs::remove_dir_all(&dir_b);
4552 fs::create_dir_all(&dir_a).unwrap();
4553 fs::create_dir_all(&dir_b).unwrap();
4554 let file_path = temp_file("storage_reconfigure.txt", "original");
4555
4556 let mut store = BackupStore::new();
4557 store.set_storage_dir(dir_a.clone(), 72);
4558 store
4559 .snapshot(DEFAULT_SESSION_ID, &file_path, "stored in a")
4560 .unwrap();
4561 assert_eq!(store.disk_history_count(DEFAULT_SESSION_ID, &file_path), 1);
4562
4563 store.set_storage_dir(dir_b.clone(), 72);
4564
4565 assert_eq!(store.disk_history_count(DEFAULT_SESSION_ID, &file_path), 0);
4566 assert!(store.tracked_files(DEFAULT_SESSION_ID).is_empty());
4567 let _ = fs::remove_dir_all(&dir_a);
4568 let _ = fs::remove_dir_all(&dir_b);
4569 }
4570
4571 #[test]
4572 fn restore_last_operation_restores_all_top_entries_for_same_op() {
4573 let path_a = temp_file("op_restore_a.txt", "a1");
4574 let path_b = temp_file("op_restore_b.txt", "b1");
4575 let mut store = BackupStore::new();
4576 let op_id = "op-test-00000001";
4577
4578 store
4579 .snapshot_with_op(DEFAULT_SESSION_ID, &path_a, "a", Some(op_id))
4580 .unwrap();
4581 store
4582 .snapshot_with_op(DEFAULT_SESSION_ID, &path_b, "b", Some(op_id))
4583 .unwrap();
4584 fs::write(&path_a, "a2").unwrap();
4585 fs::write(&path_b, "b2").unwrap();
4586
4587 let restored = store.restore_last_operation(DEFAULT_SESSION_ID).unwrap();
4588 assert_eq!(restored.op_id, op_id);
4589 assert_eq!(restored.restored.len(), 2);
4590 assert_eq!(fs::read_to_string(&path_a).unwrap(), "a1");
4591 assert_eq!(fs::read_to_string(&path_b).unwrap(), "b1");
4592 }
4593
4594 #[test]
4595 fn restore_last_operation_deletes_tombstone_destination() {
4596 let dir = std::env::temp_dir().join("aft_backup_tombstone_delete_test");
4597 let _ = fs::remove_dir_all(&dir);
4598 fs::create_dir_all(&dir).unwrap();
4599 let source = dir.join("source.txt");
4600 let destination = dir.join("destination.txt");
4601 let canonical_dir = fs::canonicalize(&dir).unwrap();
4602 fs::write(&source, "original").unwrap();
4603
4604 let mut store = BackupStore::new();
4605 let op_id = "op-tombstone-delete";
4606 store
4607 .snapshot_with_op(DEFAULT_SESSION_ID, &source, "move source", Some(op_id))
4608 .unwrap();
4609 fs::rename(&source, &destination).unwrap();
4610 store
4611 .snapshot_op_tombstone(DEFAULT_SESSION_ID, op_id, &destination, "created dest")
4612 .unwrap();
4613
4614 let restored = store.restore_last_operation(DEFAULT_SESSION_ID).unwrap();
4615 assert_eq!(restored.op_id, op_id);
4616 assert_eq!(restored.restored.len(), 2);
4617 let restored_paths = restored
4618 .restored
4619 .iter()
4620 .map(|file| file.path.clone())
4621 .collect::<HashSet<_>>();
4622 assert_eq!(
4623 restored_paths,
4624 HashSet::from([
4625 canonical_dir.join("source.txt"),
4626 canonical_dir.join("destination.txt")
4627 ])
4628 );
4629 assert!(restored
4630 .restored
4631 .iter()
4632 .all(|file| !file.backup_id.is_empty()));
4633 assert_eq!(fs::read_to_string(&source).unwrap(), "original");
4634 assert!(!destination.exists());
4635 let _ = fs::remove_dir_all(&dir);
4636 }
4637
4638 #[test]
4639 fn restore_last_operation_rolls_back_source_when_tombstone_delete_fails() {
4640 let dir = std::env::temp_dir().join("aft_backup_tombstone_atomic_test");
4641 let _ = fs::remove_dir_all(&dir);
4642 fs::create_dir_all(&dir).unwrap();
4643 let source = dir.join("source.txt");
4644 let destination = dir.join("destination.txt");
4645 fs::write(&source, "original").unwrap();
4646
4647 let mut store = BackupStore::new();
4648 let op_id = "op-tombstone-atomic";
4649 store
4650 .snapshot_with_op(DEFAULT_SESSION_ID, &source, "move source", Some(op_id))
4651 .unwrap();
4652 fs::rename(&source, &destination).unwrap();
4653 store
4654 .snapshot_op_tombstone(DEFAULT_SESSION_ID, op_id, &destination, "created dest")
4655 .unwrap();
4656
4657 fs::remove_file(&destination).unwrap();
4658 fs::create_dir(&destination).unwrap();
4659 let result = store.restore_last_operation(DEFAULT_SESSION_ID);
4660
4661 assert!(result.is_err(), "directory tombstone target should fail");
4662 assert!(
4663 !source.exists(),
4664 "source restore must roll back when destination deletion fails"
4665 );
4666 assert!(
4667 destination.is_dir(),
4668 "failed tombstone target should remain"
4669 );
4670 let _ = fs::remove_dir_all(&dir);
4671 }
4672
4673 #[cfg(unix)]
4679 #[test]
4680 fn restore_last_operation_is_atomic_when_a_write_fails() {
4681 let dir = std::env::temp_dir().join("aft_backup_tests_atomic_restore");
4682 let _ = fs::remove_dir_all(&dir);
4683 fs::create_dir_all(&dir).unwrap();
4684 let path_a = dir.join("a.txt");
4685 let path_b = dir.join("b.txt");
4686 let path_c = dir.join("c.txt");
4687 fs::write(&path_a, "a-original").unwrap();
4688 fs::write(&path_b, "b-original").unwrap();
4689 fs::write(&path_c, "c-original").unwrap();
4690
4691 let mut store = BackupStore::new();
4692 let op_id = "op-atomic-restore-01";
4693 let id_a = store
4694 .snapshot_with_op(DEFAULT_SESSION_ID, &path_a, "a", Some(op_id))
4695 .unwrap()
4696 .unwrap();
4697 let id_b = store
4698 .snapshot_with_op(DEFAULT_SESSION_ID, &path_b, "b", Some(op_id))
4699 .unwrap()
4700 .unwrap();
4701 let id_c = store
4702 .snapshot_with_op(DEFAULT_SESSION_ID, &path_c, "c", Some(op_id))
4703 .unwrap()
4704 .unwrap();
4705 fs::write(&path_a, "a-modified").unwrap();
4706 fs::write(&path_b, "b-modified").unwrap();
4707 fs::write(&path_c, "c-modified").unwrap();
4708
4709 let original_permissions = fs::metadata(&path_b).unwrap().permissions();
4710 let mut readonly_permissions = original_permissions.clone();
4711 readonly_permissions.set_mode(0o444);
4712 fs::set_permissions(&path_b, readonly_permissions).unwrap();
4713
4714 let result = store.restore_last_operation(DEFAULT_SESSION_ID);
4715 fs::set_permissions(&path_b, original_permissions).unwrap();
4716
4717 assert!(result.is_err());
4718 assert_eq!(fs::read_to_string(&path_a).unwrap(), "a-modified");
4719 assert_eq!(fs::read_to_string(&path_b).unwrap(), "b-modified");
4720 assert_eq!(fs::read_to_string(&path_c).unwrap(), "c-modified");
4721
4722 let history_a = store.history(DEFAULT_SESSION_ID, &path_a);
4723 let history_b = store.history(DEFAULT_SESSION_ID, &path_b);
4724 let history_c = store.history(DEFAULT_SESSION_ID, &path_c);
4725 assert_eq!(history_a.len(), 1);
4726 assert_eq!(history_b.len(), 1);
4727 assert_eq!(history_c.len(), 1);
4728 assert_eq!(history_a[0].backup_id, id_a);
4729 assert_eq!(history_b[0].backup_id, id_b);
4730 assert_eq!(history_c[0].backup_id, id_c);
4731 assert_eq!(history_a[0].op_id.as_deref(), Some(op_id));
4732 assert_eq!(history_b[0].op_id.as_deref(), Some(op_id));
4733 assert_eq!(history_c[0].op_id.as_deref(), Some(op_id));
4734
4735 let restored = store.restore_last_operation(DEFAULT_SESSION_ID).unwrap();
4736 assert_eq!(restored.op_id, op_id);
4737 assert_eq!(restored.restored.len(), 3);
4738 assert_eq!(fs::read_to_string(&path_a).unwrap(), "a-original");
4739 assert_eq!(fs::read_to_string(&path_b).unwrap(), "b-original");
4740 assert_eq!(fs::read_to_string(&path_c).unwrap(), "c-original");
4741
4742 let _ = fs::remove_dir_all(&dir);
4743 }
4744
4745 #[test]
4746 fn restore_last_operation_restores_only_most_recent_op() {
4747 let path_a = temp_file("op_recent_a.txt", "a1");
4748 let path_b = temp_file("op_recent_b.txt", "b1");
4749 let mut store = BackupStore::new();
4750
4751 store
4752 .snapshot_with_op(DEFAULT_SESSION_ID, &path_a, "older", Some("op-older"))
4753 .unwrap();
4754 store
4755 .snapshot_with_op(DEFAULT_SESSION_ID, &path_b, "newer", Some("op-newer"))
4756 .unwrap();
4757 fs::write(&path_a, "a2").unwrap();
4758 fs::write(&path_b, "b2").unwrap();
4759
4760 let restored = store.restore_last_operation(DEFAULT_SESSION_ID).unwrap();
4761 assert_eq!(restored.op_id, "op-newer");
4762 assert_eq!(restored.restored.len(), 1);
4763 assert_eq!(fs::read_to_string(&path_a).unwrap(), "a2");
4764 assert_eq!(fs::read_to_string(&path_b).unwrap(), "b1");
4765 }
4766
4767 #[test]
4768 fn restore_recreates_missing_parent_directories() {
4769 let dir = std::env::temp_dir().join("aft_backup_tests_recreate_parents");
4772 let _ = fs::remove_dir_all(&dir);
4773 let nested = dir.join("nested");
4774 fs::create_dir_all(&nested).unwrap();
4775 let path = nested.join("inner.txt");
4776 fs::write(&path, "original").unwrap();
4777
4778 let mut store = BackupStore::new();
4779 let op_id = "op-recreate-parents-01";
4780 store
4781 .snapshot_with_op(DEFAULT_SESSION_ID, &path, "original", Some(op_id))
4782 .unwrap();
4783
4784 fs::remove_dir_all(&dir).unwrap();
4786 assert!(!path.exists());
4787 assert!(!nested.exists());
4788 assert!(!dir.exists());
4789
4790 let restored = store.restore_last_operation(DEFAULT_SESSION_ID).unwrap();
4791 assert_eq!(restored.op_id, op_id);
4792 assert_eq!(restored.restored.len(), 1);
4793 assert!(
4794 path.exists(),
4795 "file should be restored even though both nested/ and dir/ were missing"
4796 );
4797 assert_eq!(fs::read_to_string(&path).unwrap(), "original");
4798
4799 let _ = fs::remove_dir_all(&dir);
4800 }
4801
4802 #[test]
4803 fn restore_last_operation_ignores_legacy_entries_without_op_id() {
4804 let path = temp_file("op_legacy_none.txt", "v1");
4805 let mut store = BackupStore::new();
4806
4807 store.snapshot(DEFAULT_SESSION_ID, &path, "legacy").unwrap();
4808 fs::write(&path, "v2").unwrap();
4809
4810 let err = store.restore_last_operation(DEFAULT_SESSION_ID);
4811 assert!(matches!(err, Err(AftError::NoUndoHistory { .. })));
4812 assert_eq!(fs::read_to_string(&path).unwrap(), "v2");
4813 }
4814
4815 #[test]
4816 fn schema_v2_meta_loads_with_none_op_id_and_persists_as_v3() {
4817 let dir = std::env::temp_dir().join("aft_backup_v2_to_v3_test");
4818 let _ = fs::remove_dir_all(&dir);
4819 fs::create_dir_all(&dir).unwrap();
4820 let file_path = temp_file("v2_to_v3.txt", "original");
4821 let key = canonicalize_key(&file_path);
4822 let session_dir = dir
4823 .join("backups")
4824 .join(BackupStore::session_hash(DEFAULT_SESSION_ID));
4825 let path_dir = session_dir.join(BackupStore::path_hash(&key));
4826 fs::create_dir_all(&path_dir).unwrap();
4827 fs::write(path_dir.join("0.bak"), "original").unwrap();
4828 fs::write(
4829 session_dir.join("session.json"),
4830 serde_json::to_string_pretty(&serde_json::json!({
4831 "schema_version": 2,
4832 "session_id": DEFAULT_SESSION_ID,
4833 "last_accessed": current_timestamp(),
4834 }))
4835 .unwrap(),
4836 )
4837 .unwrap();
4838 fs::write(
4839 path_dir.join("meta.json"),
4840 serde_json::to_string_pretty(&serde_json::json!({
4841 "schema_version": 2,
4842 "session_id": DEFAULT_SESSION_ID,
4843 "path": key.display().to_string(),
4844 "count": 1,
4845 }))
4846 .unwrap(),
4847 )
4848 .unwrap();
4849
4850 let mut store = BackupStore::new();
4851 store.set_storage_dir(dir.clone(), 72);
4852 assert!(store
4853 .load_from_disk_if_needed(DEFAULT_SESSION_ID, &key)
4854 .unwrap());
4855 let history = store.history(DEFAULT_SESSION_ID, &file_path);
4856 assert_eq!(history.len(), 1);
4857 assert_eq!(history[0].op_id, None);
4858
4859 fs::write(&file_path, "second").unwrap();
4860 store
4861 .snapshot_with_op(DEFAULT_SESSION_ID, &file_path, "second", Some("op-v3"))
4862 .unwrap();
4863 let written: serde_json::Value =
4864 serde_json::from_str(&fs::read_to_string(path_dir.join("meta.json")).unwrap()).unwrap();
4865 assert_eq!(written["schema_version"], SCHEMA_VERSION);
4866 assert_eq!(written["entries"][0]["op_id"], serde_json::Value::Null);
4867 assert_eq!(written["entries"][1]["op_id"], "op-v3");
4868 let _ = fs::remove_dir_all(&dir);
4869 }
4870
4871 #[test]
4872 fn per_file_restore_latest_still_works_with_op_ids() {
4873 let path = temp_file("op_per_file.txt", "v1");
4874 let mut store = BackupStore::new();
4875
4876 store
4877 .snapshot_with_op(DEFAULT_SESSION_ID, &path, "op", Some("op-file"))
4878 .unwrap();
4879 fs::write(&path, "v2").unwrap();
4880
4881 let (entry, _) = store.restore_latest(DEFAULT_SESSION_ID, &path).unwrap();
4882 assert_eq!(entry.op_id.as_deref(), Some("op-file"));
4883 assert_eq!(fs::read_to_string(&path).unwrap(), "v1");
4884 }
4885
4886 #[test]
4887 fn per_file_restore_latest_deletes_tombstone() {
4888 let dir = std::env::temp_dir().join("aft_backup_per_file_tombstone_test");
4889 let _ = fs::remove_dir_all(&dir);
4890 fs::create_dir_all(&dir).unwrap();
4891 let path = dir.join("created.txt");
4892 fs::write(&path, "created").unwrap();
4893
4894 let mut store = BackupStore::new();
4895 let id = store
4896 .snapshot_op_tombstone(DEFAULT_SESSION_ID, "op-create", &path, "created")
4897 .unwrap()
4898 .unwrap();
4899
4900 let (entry, _) = store.restore_latest(DEFAULT_SESSION_ID, &path).unwrap();
4901 assert_eq!(entry.backup_id, id);
4902 assert!(!path.exists(), "tombstone undo should delete the file");
4903 let _ = fs::remove_dir_all(&dir);
4904 }
4905
4906 #[test]
4907 fn lazy_stack_read_skips_tampered_meta_path_hash_mismatch() {
4908 let dir = std::env::temp_dir().join("aft_backup_tampered_meta_skip_test");
4909 let _ = fs::remove_dir_all(&dir);
4910 let backups = dir.join("backups");
4911 let session_dir = backups.join(BackupStore::session_hash(DEFAULT_SESSION_ID));
4912 let path_dir = session_dir.join("not-the-path-hash");
4913 fs::create_dir_all(&path_dir).unwrap();
4914 fs::write(
4915 session_dir.join("session.json"),
4916 serde_json::to_string_pretty(&serde_json::json!({
4917 "schema_version": SCHEMA_VERSION,
4918 "session_id": DEFAULT_SESSION_ID,
4919 "last_accessed": current_timestamp(),
4920 }))
4921 .unwrap(),
4922 )
4923 .unwrap();
4924 fs::write(path_dir.join("0.bak"), "outside").unwrap();
4925 fs::write(
4926 path_dir.join("meta.json"),
4927 serde_json::to_string_pretty(&serde_json::json!({
4928 "schema_version": SCHEMA_VERSION,
4929 "session_id": DEFAULT_SESSION_ID,
4930 "path": "/tmp/aft-malicious-overwrite-target.txt",
4931 "count": 1,
4932 "entries": [{
4933 "backup_id": "backup-0",
4934 "timestamp": current_timestamp(),
4935 "order": "1",
4936 "description": "tampered",
4937 "op_id": "op-tampered",
4938 "kind": "content",
4939 }]
4940 }))
4941 .unwrap(),
4942 )
4943 .unwrap();
4944
4945 let mut store = BackupStore::new();
4946 store.set_storage_dir(dir.clone(), 72);
4947
4948 assert!(store
4949 .history(
4950 DEFAULT_SESSION_ID,
4951 Path::new("/tmp/aft-malicious-overwrite-target.txt")
4952 )
4953 .is_empty());
4954 assert!(store.sessions_with_backups().is_empty());
4955 let _ = fs::remove_dir_all(&dir);
4956 }
4957
4958 #[test]
4959 fn restore_last_operation_uses_only_top_entries_and_persisted_order() {
4960 let path_a = temp_file("op_order_a.txt", "a1");
4961 let path_b = temp_file("op_order_b.txt", "b1");
4962 let mut store = BackupStore::new();
4963
4964 store
4965 .snapshot_with_op(DEFAULT_SESSION_ID, &path_a, "buried", Some("op-buried"))
4966 .unwrap();
4967 store
4968 .snapshot(DEFAULT_SESSION_ID, &path_a, "top without op")
4969 .unwrap();
4970 store
4971 .snapshot_with_op(DEFAULT_SESSION_ID, &path_b, "top", Some("op-top"))
4972 .unwrap();
4973
4974 let key_a = canonicalize_key(&path_a);
4975 let key_b = canonicalize_key(&path_b);
4976 let files = store.entries.get_mut(DEFAULT_SESSION_ID).unwrap();
4977 files.get_mut(&key_a).unwrap()[0].order = u128::MAX;
4978 files.get_mut(&key_a).unwrap()[1].order = 1;
4979 files.get_mut(&key_b).unwrap()[0].order = 2;
4980
4981 fs::write(&path_a, "a2").unwrap();
4982 fs::write(&path_b, "b2").unwrap();
4983
4984 let restored = store.restore_last_operation(DEFAULT_SESSION_ID).unwrap();
4985 assert_eq!(restored.op_id, "op-top");
4986 assert_eq!(restored.restored.len(), 1);
4987 assert_eq!(fs::read_to_string(&path_a).unwrap(), "a2");
4988 assert_eq!(fs::read_to_string(&path_b).unwrap(), "b1");
4989 }
4990
4991 #[test]
4992 fn append_only_v2_adds_one_content_file_at_steady_depth() {
4993 let dir = tempfile::tempdir().unwrap();
4994 let path = dir.path().join("append_only.txt");
4995 fs::write(&path, "v0").unwrap();
4996 let mut store = BackupStore::new();
4997 store.set_storage_dir(dir.path().to_path_buf(), 72);
4998
4999 for i in 0..MAX_UNDO_DEPTH {
5000 store
5001 .snapshot(DEFAULT_SESSION_ID, &path, "push")
5002 .unwrap()
5003 .unwrap();
5004 fs::write(&path, format!("v{}", i + 1)).unwrap();
5005 }
5006
5007 let key = canonicalize_key(&path);
5008 let stack_dir = store
5009 .session_dir(DEFAULT_SESSION_ID)
5010 .unwrap()
5011 .join(BackupStore::path_hash(&key));
5012 let before = backup_content_names(&stack_dir);
5013 assert_eq!(before.len(), MAX_UNDO_DEPTH);
5014
5015 store
5016 .snapshot(DEFAULT_SESSION_ID, &path, "steady push")
5017 .unwrap()
5018 .unwrap();
5019 let after = backup_content_names(&stack_dir);
5020 assert_eq!(after.len(), MAX_UNDO_DEPTH);
5021 assert_eq!(after.difference(&before).count(), 1);
5022 assert_eq!(before.difference(&after).count(), 1);
5023
5024 let meta: serde_json::Value =
5025 serde_json::from_str(&fs::read_to_string(stack_dir.join("meta.json")).unwrap())
5026 .unwrap();
5027 assert_eq!(
5028 meta.get("format_version").and_then(|v| v.as_str()),
5029 Some("v2")
5030 );
5031 assert!(meta_entries(&meta)
5032 .unwrap()
5033 .iter()
5034 .all(|entry| entry.get("content_path").and_then(|v| v.as_str()).is_some()));
5035 }
5036
5037 #[test]
5038 fn legacy_stack_migrates_to_v2_on_next_write() {
5039 let dir = tempfile::tempdir().unwrap();
5040 let path = dir.path().join("legacy.txt");
5041 fs::write(&path, "current").unwrap();
5042 let key = canonicalize_key(&path);
5043 let session_dir = dir
5044 .path()
5045 .join("backups")
5046 .join(BackupStore::session_hash(DEFAULT_SESSION_ID));
5047 let stack_dir = session_dir.join(BackupStore::path_hash(&key));
5048 fs::create_dir_all(&stack_dir).unwrap();
5049 fs::write(
5050 session_dir.join("session.json"),
5051 serde_json::to_string_pretty(&serde_json::json!({
5052 "schema_version": SCHEMA_VERSION,
5053 "session_id": DEFAULT_SESSION_ID,
5054 "last_accessed": current_timestamp(),
5055 }))
5056 .unwrap(),
5057 )
5058 .unwrap();
5059 fs::write(stack_dir.join("0.bak"), "legacy").unwrap();
5060 fs::write(
5061 stack_dir.join("meta.json"),
5062 serde_json::to_string_pretty(&serde_json::json!({
5063 "schema_version": SCHEMA_VERSION,
5064 "session_id": DEFAULT_SESSION_ID,
5065 "path": key.display().to_string(),
5066 "count": 1,
5067 "entries": [{
5068 "backup_id": "backup-0",
5069 "timestamp": current_timestamp(),
5070 "order": "1",
5071 "description": "legacy",
5072 "kind": "content",
5073 }]
5074 }))
5075 .unwrap(),
5076 )
5077 .unwrap();
5078
5079 let mut store = BackupStore::new();
5080 store.set_storage_dir(dir.path().to_path_buf(), 72);
5081 assert_eq!(
5082 store.history(DEFAULT_SESSION_ID, &path)[0].content,
5083 "legacy"
5084 );
5085
5086 store
5087 .snapshot(DEFAULT_SESSION_ID, &path, "migrate")
5088 .unwrap()
5089 .unwrap();
5090 let meta: serde_json::Value =
5091 serde_json::from_str(&fs::read_to_string(stack_dir.join("meta.json")).unwrap())
5092 .unwrap();
5093 assert_eq!(
5094 meta.get("format_version").and_then(|v| v.as_str()),
5095 Some("v2")
5096 );
5097 assert!(!stack_dir.join("0.bak").exists());
5098 assert_eq!(backup_content_names(&stack_dir).len(), 2);
5099 }
5100
5101 #[test]
5102 fn snapshot_reloads_non_empty_stale_stack_before_append() {
5103 let project = tempfile::tempdir().unwrap();
5104 let storage = tempfile::tempdir().unwrap();
5105 let path = project.path().join("stale-memory.txt");
5106 fs::write(&path, "v0").unwrap();
5107 let policy = BackupPolicy {
5108 enabled: true,
5109 max_depth: 2,
5110 max_file_size: None,
5111 };
5112
5113 let mut store_a = BackupStore::new();
5114 store_a.set_storage_dir(storage.path().to_path_buf(), 72);
5115 store_a.set_policy(policy);
5116 store_a
5117 .snapshot(DEFAULT_SESSION_ID, &path, "a captures v0")
5118 .unwrap();
5119 fs::write(&path, "v1").unwrap();
5120
5121 let mut store_b = BackupStore::new();
5122 store_b.set_storage_dir(storage.path().to_path_buf(), 72);
5123 store_b.set_policy(policy);
5124 store_b
5125 .snapshot(DEFAULT_SESSION_ID, &path, "b captures v1")
5126 .unwrap();
5127 fs::write(&path, "v2").unwrap();
5128
5129 store_a
5130 .snapshot(DEFAULT_SESSION_ID, &path, "a captures v2")
5131 .unwrap();
5132
5133 let mut fresh = BackupStore::new();
5134 fresh.set_storage_dir(storage.path().to_path_buf(), 72);
5135 let contents = fresh
5136 .history(DEFAULT_SESSION_ID, &path)
5137 .into_iter()
5138 .map(|entry| entry.content)
5139 .collect::<Vec<_>>();
5140 assert_eq!(contents, vec!["v1".to_string(), "v2".to_string()]);
5141 }
5142
5143 #[test]
5144 fn restore_latest_clears_stale_memory_when_disk_stack_disappears() {
5145 let project = tempfile::tempdir().unwrap();
5146 let storage = tempfile::tempdir().unwrap();
5147 let session = "stale-resurrection-session";
5148 let path = project.path().join("stale-resurrection.txt");
5149 fs::write(&path, "v0").unwrap();
5150
5151 let mut store_a = BackupStore::new();
5152 store_a.set_storage_dir(storage.path().to_path_buf(), 72);
5153 store_a.snapshot(session, &path, "a captures v0").unwrap();
5154 fs::write(&path, "v1").unwrap();
5155
5156 let mut store_b = BackupStore::new();
5157 store_b.set_storage_dir(storage.path().to_path_buf(), 72);
5158 let (restored, _) = store_b.restore_latest(session, &path).unwrap();
5159 assert_eq!(restored.content, "v0");
5160
5161 fs::write(&path, "current after other restore").unwrap();
5162 let error = store_a.restore_latest(session, &path).unwrap_err();
5163
5164 assert_eq!(error.code(), "no_undo_history");
5165 assert_eq!(
5166 fs::read_to_string(&path).unwrap(),
5167 "current after other restore"
5168 );
5169 let key = canonicalize_key(&path);
5170 assert!(store_a
5171 .entries
5172 .get(session)
5173 .and_then(|files| files.get(&key))
5174 .is_none());
5175
5176 let snapshot_path = project.path().join("stale-snapshot.txt");
5177 fs::write(&snapshot_path, "snapshot v0").unwrap();
5178 let mut store_c = BackupStore::new();
5179 store_c.set_storage_dir(storage.path().to_path_buf(), 72);
5180 store_c
5181 .snapshot(session, &snapshot_path, "c captures v0")
5182 .unwrap();
5183 fs::write(&snapshot_path, "snapshot v1").unwrap();
5184 let mut store_d = BackupStore::new();
5185 store_d.set_storage_dir(storage.path().to_path_buf(), 72);
5186 store_d.restore_latest(session, &snapshot_path).unwrap();
5187
5188 fs::write(&snapshot_path, "snapshot current").unwrap();
5189 store_c
5190 .snapshot(session, &snapshot_path, "c captures current")
5191 .unwrap();
5192 let mut fresh = BackupStore::new();
5193 fresh.set_storage_dir(storage.path().to_path_buf(), 72);
5194 let contents = fresh
5195 .history(session, &snapshot_path)
5196 .into_iter()
5197 .map(|entry| entry.content)
5198 .collect::<Vec<_>>();
5199 assert_eq!(contents, vec!["snapshot current".to_string()]);
5200 }
5201
5202 #[test]
5203 fn restore_last_operation_returns_retry_error_under_unbounded_key_churn() {
5204 let project = tempfile::tempdir().unwrap();
5205 let storage = tempfile::tempdir().unwrap();
5206 let session = "restore-churn-session";
5207 let base_path = project.path().join("base.txt");
5208 fs::write(&base_path, "base before").unwrap();
5209 let mut base_store = BackupStore::new();
5210 base_store.set_storage_dir(storage.path().to_path_buf(), 72);
5211 base_store
5212 .snapshot_with_op(session, &base_path, "base op", Some("op-base"))
5213 .unwrap();
5214 fs::write(&base_path, "base after").unwrap();
5215
5216 let churn_count = Arc::new(Mutex::new(0usize));
5217 let hook_count = churn_count.clone();
5218 let hook_project = project.path().to_path_buf();
5219 let hook_storage = storage.path().to_path_buf();
5220 set_restore_before_lock_hook_for_tests(session, move |_| {
5221 let mut count = hook_count.lock().unwrap();
5222 let churn_path = hook_project.join(format!("churn-{}.txt", *count));
5223 fs::write(&churn_path, format!("churn before {}", *count)).unwrap();
5224 let mut churn_store = BackupStore::new();
5225 churn_store.set_storage_dir(hook_storage.clone(), 72);
5226 let op_id = format!("op-churn-{}", *count);
5227 churn_store
5228 .snapshot_with_op(session, &churn_path, "churn op", Some(&op_id))
5229 .unwrap();
5230 fs::write(&churn_path, format!("churn after {}", *count)).unwrap();
5231 *count += 1;
5232 *count < MAX_RESTORE_OPERATION_LOCK_RETRIES
5233 });
5234
5235 let mut restore_store = BackupStore::new();
5236 restore_store.set_storage_dir(storage.path().to_path_buf(), 72);
5237 let error = restore_store.restore_last_operation(session).unwrap_err();
5238
5239 assert_eq!(error.code(), "io_error");
5240 assert!(error
5241 .to_string()
5242 .contains("backup stack changing under concurrent activity; retry"));
5243 assert_eq!(
5244 *churn_count.lock().unwrap(),
5245 MAX_RESTORE_OPERATION_LOCK_RETRIES
5246 );
5247 }
5248
5249 #[test]
5250 fn restore_last_operation_test_hooks_are_isolated_per_session() {
5251 let project = tempfile::tempdir().unwrap();
5252 let storage = tempfile::tempdir().unwrap();
5253 let session_a = "restore-hook-session-a";
5254 let session_b = "restore-hook-session-b";
5255 let path_a = project.path().join("restore-hook-a.txt");
5256 fs::write(&path_a, "v0").unwrap();
5257
5258 let mut store_a = BackupStore::new();
5259 store_a.set_storage_dir(storage.path().to_path_buf(), 72);
5260 store_a
5261 .snapshot_with_op(session_a, &path_a, "old op", Some("op-old-a"))
5262 .unwrap();
5263 fs::write(&path_a, "v1").unwrap();
5264
5265 let hook_storage = storage.path().to_path_buf();
5266 let hook_path_a = path_a.clone();
5267 set_restore_before_lock_hook_for_tests(session_a, move |_| {
5268 let mut hook_store = BackupStore::new();
5269 hook_store.set_storage_dir(hook_storage.clone(), 72);
5270 hook_store
5271 .snapshot_with_op(session_a, &hook_path_a, "new op", Some("op-new-a"))
5272 .unwrap();
5273 fs::write(&hook_path_a, "v2").unwrap();
5274 false
5275 });
5276 set_restore_before_lock_hook_for_tests(session_b, |_| false);
5277
5278 let restored = store_a.restore_last_operation(session_a).unwrap();
5279
5280 assert_eq!(restored.op_id, "op-new-a");
5281 assert_eq!(fs::read_to_string(&path_a).unwrap(), "v1");
5282 run_restore_before_lock_hook_for_tests(session_b, 0);
5283 }
5284
5285 #[test]
5286 fn restore_last_operation_rescans_stack_after_locking() {
5287 let project = tempfile::tempdir().unwrap();
5288 let storage = tempfile::tempdir().unwrap();
5289 let session = "restore-toctou-session";
5290 let path = project.path().join("restore-toctou.txt");
5291 fs::write(&path, "v0").unwrap();
5292
5293 let mut store_a = BackupStore::new();
5294 store_a.set_storage_dir(storage.path().to_path_buf(), 72);
5295 store_a
5296 .snapshot_with_op(session, &path, "old op", Some("op-old"))
5297 .unwrap();
5298 fs::write(&path, "v1").unwrap();
5299
5300 let hook_storage = storage.path().to_path_buf();
5301 let hook_path = path.clone();
5302 set_restore_before_lock_hook_for_tests(session, move |_| {
5303 let mut store_b = BackupStore::new();
5304 store_b.set_storage_dir(hook_storage.clone(), 72);
5305 store_b
5306 .snapshot_with_op(session, &hook_path, "new op", Some("op-new"))
5307 .unwrap();
5308 fs::write(&hook_path, "v2").unwrap();
5309 false
5310 });
5311
5312 let restored = store_a.restore_last_operation(session).unwrap();
5313
5314 assert_eq!(restored.op_id, "op-new");
5315 assert_eq!(fs::read_to_string(&path).unwrap(), "v1");
5316 }
5317
5318 #[test]
5319 fn corrupt_v2_meta_fails_closed_for_operation_and_single_restore() {
5320 let project = tempfile::tempdir().unwrap();
5321 let storage = tempfile::tempdir().unwrap();
5322 let session = "corrupt-v2-session";
5323 let path = project.path().join("corrupt-v2.txt");
5324 fs::write(&path, "current").unwrap();
5325 let key = canonicalize_key(&path);
5326 let session_dir = storage
5327 .path()
5328 .join("backups")
5329 .join(BackupStore::session_hash(session));
5330 let stack_dir = session_dir.join(BackupStore::path_hash(&key));
5331 fs::create_dir_all(&stack_dir).unwrap();
5332 fs::write(
5333 session_dir.join("session.json"),
5334 serde_json::to_string_pretty(&serde_json::json!({
5335 "schema_version": SCHEMA_VERSION,
5336 "session_id": session,
5337 "last_accessed": current_timestamp(),
5338 }))
5339 .unwrap(),
5340 )
5341 .unwrap();
5342 fs::write(
5343 stack_dir.join("meta.json"),
5344 serde_json::to_string_pretty(&serde_json::json!({
5345 "schema_version": SCHEMA_VERSION,
5346 "format_version": "v2",
5347 "session_id": session,
5348 "path": key.display().to_string(),
5349 "count": 1,
5350 "entries": [{
5351 "backup_id": "backup-corrupt",
5352 "timestamp": current_timestamp(),
5353 "order": "9",
5354 "description": "corrupt disk should win over DB fallback",
5355 "op_id": "op-corrupt",
5356 "kind": "content",
5357 "content_path": "bak_9_backup-corrupt.bak",
5358 }]
5359 }))
5360 .unwrap(),
5361 )
5362 .unwrap();
5363
5364 let conn = crate::db::open(&storage.path().join("aft.db")).unwrap();
5365 let fallback_path = stack_dir.join("db-fallback.bak");
5366 fs::write(&fallback_path, "db fallback").unwrap();
5367 crate::db::backups::upsert_backup(
5368 &conn,
5369 &BackupRow {
5370 backup_id: "backup-db".to_string(),
5371 harness: "opencode".to_string(),
5372 session_id: session.to_string(),
5373 project_key: "project".to_string(),
5374 op_id: Some("op-corrupt".to_string()),
5375 order: 9,
5376 file_path: key.display().to_string(),
5377 path_hash: BackupStore::path_hash(&key),
5378 backup_path: Some(fallback_path.display().to_string()),
5379 kind: "content".to_string(),
5380 description: "db fallback".to_string(),
5381 created_at: i64::try_from(current_timestamp()).unwrap(),
5382 is_tombstone: false,
5383 restore_meta: None,
5384 },
5385 )
5386 .unwrap();
5387 let shared = Arc::new(Mutex::new(conn));
5388
5389 let mut single = BackupStore::new();
5390 single.set_storage_dir(storage.path().to_path_buf(), 72);
5391 single.set_db_harness(Harness::Opencode);
5392 single.set_db_project_key("project".to_string());
5393 single.set_db_pool(shared.clone());
5394 let single_error = single.restore_latest(session, &path).unwrap_err();
5395 assert_eq!(single_error.code(), "io_error");
5396 assert_eq!(fs::read_to_string(&path).unwrap(), "current");
5397
5398 let mut operation = BackupStore::new();
5399 operation.set_storage_dir(storage.path().to_path_buf(), 72);
5400 operation.set_db_harness(Harness::Opencode);
5401 operation.set_db_project_key("project".to_string());
5402 operation.set_db_pool(shared);
5403 let operation_error = operation.restore_last_operation(session).unwrap_err();
5404 assert_eq!(operation_error.code(), "io_error");
5405 assert_eq!(fs::read_to_string(&path).unwrap(), "current");
5406 }
5407
5408 #[test]
5409 fn replace_file_replaces_existing_meta_with_single_rename_path() {
5410 let dir = tempfile::tempdir().unwrap();
5411 let meta_path = dir.path().join("meta.json");
5412 let temp_path = dir.path().join("meta.tmp");
5413 fs::write(&meta_path, "old").unwrap();
5414 fs::write(&temp_path, "new").unwrap();
5415
5416 replace_file(&temp_path, &meta_path).unwrap();
5417
5418 assert_eq!(fs::read_to_string(&meta_path).unwrap(), "new");
5419 assert!(!temp_path.exists());
5420 }
5421
5422 #[test]
5423 fn snapshot_write_failure_restores_full_pre_trim_stack() {
5424 let project = tempfile::tempdir().unwrap();
5425 let storage = tempfile::tempdir().unwrap();
5426 let session = "rollback-pretrim-session";
5427 let path = project.path().join("rollback.txt");
5428 fs::write(&path, "v0").unwrap();
5429 let mut store = BackupStore::new();
5430 store.set_storage_dir(storage.path().to_path_buf(), 72);
5431 store.set_policy(BackupPolicy {
5432 enabled: true,
5433 max_depth: 2,
5434 max_file_size: None,
5435 });
5436
5437 store.snapshot(session, &path, "first").unwrap();
5438 fs::write(&path, "v1").unwrap();
5439 store.snapshot(session, &path, "second").unwrap();
5440 fs::write(&path, "v2").unwrap();
5441 let key = canonicalize_key(&path);
5442 let before_file_stack = store.entries.get(session).unwrap().get(&key).unwrap();
5443 let before_file_identity = before_file_stack
5444 .iter()
5445 .map(|entry| (entry.backup_id.clone(), entry.order))
5446 .collect::<Vec<_>>();
5447
5448 store.fail_next_disk_write_for_tests();
5449 let error = store.snapshot(session, &path, "third").unwrap_err();
5450 assert_eq!(error.code(), "io_error");
5451 let after_file_stack = store.entries.get(session).unwrap().get(&key).unwrap();
5452 assert_eq!(after_file_stack.len(), before_file_identity.len());
5453 assert_eq!(
5454 after_file_stack
5455 .iter()
5456 .map(|entry| (entry.backup_id.clone(), entry.order))
5457 .collect::<Vec<_>>(),
5458 before_file_identity
5459 );
5460
5461 let successful_id = store
5462 .snapshot(session, &path, "after failure")
5463 .unwrap()
5464 .unwrap();
5465 let successful_stack = store.entries.get(session).unwrap().get(&key).unwrap();
5466 assert_eq!(successful_stack.len(), 2);
5467 assert_eq!(successful_stack[0].description, "second");
5468 assert_eq!(successful_stack[1].description, "after failure");
5469 assert_eq!(successful_stack[1].backup_id, successful_id);
5470
5471 let tombstone = project.path().join("created-by-op.txt");
5472 store
5473 .snapshot_op_tombstone(session, "op-one", &tombstone, "created one")
5474 .unwrap();
5475 store
5476 .snapshot_op_tombstone(session, "op-two", &tombstone, "created two")
5477 .unwrap();
5478 let tombstone_key = canonicalize_key(&tombstone);
5479 let before_tombstone_stack = store
5480 .entries
5481 .get(session)
5482 .unwrap()
5483 .get(&tombstone_key)
5484 .unwrap();
5485 let before_tombstone_identity = before_tombstone_stack
5486 .iter()
5487 .map(|entry| (entry.backup_id.clone(), entry.order, entry.op_id.clone()))
5488 .collect::<Vec<_>>();
5489
5490 store.fail_next_disk_write_for_tests();
5491 let error = store
5492 .snapshot_op_tombstone(session, "op-three", &tombstone, "created three")
5493 .unwrap_err();
5494 assert_eq!(error.code(), "io_error");
5495 let after_tombstone_stack = store
5496 .entries
5497 .get(session)
5498 .unwrap()
5499 .get(&tombstone_key)
5500 .unwrap();
5501 assert_eq!(
5502 after_tombstone_stack
5503 .iter()
5504 .map(|entry| (entry.backup_id.clone(), entry.order, entry.op_id.clone()))
5505 .collect::<Vec<_>>(),
5506 before_tombstone_identity
5507 );
5508
5509 let successful_id = store
5510 .snapshot_op_tombstone(session, "op-four", &tombstone, "created four")
5511 .unwrap()
5512 .unwrap();
5513 let successful_stack = store
5514 .entries
5515 .get(session)
5516 .unwrap()
5517 .get(&tombstone_key)
5518 .unwrap();
5519 assert_eq!(successful_stack.len(), 2);
5520 assert_eq!(successful_stack[0].op_id.as_deref(), Some("op-two"));
5521 assert_eq!(successful_stack[1].op_id.as_deref(), Some("op-four"));
5522 assert_eq!(successful_stack[1].backup_id, successful_id);
5523 }
5524
5525 #[test]
5526 fn snapshot_at_max_depth_keeps_newest_window() {
5527 let project = tempfile::tempdir().unwrap();
5528 let storage = tempfile::tempdir().unwrap();
5529 let session = "depth-window-session";
5530 let path = project.path().join("window.txt");
5531 fs::write(&path, "v0").unwrap();
5532 let mut store = BackupStore::new();
5533 store.set_storage_dir(storage.path().to_path_buf(), 72);
5534 store.set_policy(BackupPolicy {
5535 enabled: true,
5536 max_depth: 2,
5537 max_file_size: None,
5538 });
5539
5540 store.snapshot(session, &path, "first").unwrap();
5541 fs::write(&path, "v1").unwrap();
5542 let second_id = store.snapshot(session, &path, "second").unwrap().unwrap();
5543 fs::write(&path, "v2").unwrap();
5544 let third_id = store.snapshot(session, &path, "third").unwrap().unwrap();
5545
5546 let history = store.history(session, &path);
5547 assert_eq!(history.len(), 2);
5548 assert_eq!(
5549 history
5550 .iter()
5551 .map(|entry| entry.backup_id.as_str())
5552 .collect::<Vec<_>>(),
5553 vec![second_id.as_str(), third_id.as_str()]
5554 );
5555 assert_eq!(history[0].content_bytes.as_ref(), b"v1");
5556 assert_eq!(history[1].content_bytes.as_ref(), b"v2");
5557 }
5558
5559 #[test]
5560 fn lowering_max_depth_prunes_disk_content_immediately() {
5561 let project = tempfile::tempdir().unwrap();
5562 let storage = tempfile::tempdir().unwrap();
5563 let path = project.path().join("policy-prune.txt");
5564 fs::write(&path, "v0").unwrap();
5565 let mut store = BackupStore::new();
5566 store.set_storage_dir(storage.path().to_path_buf(), 72);
5567
5568 for i in 0..3 {
5569 store
5570 .snapshot(DEFAULT_SESSION_ID, &path, &format!("snapshot {i}"))
5571 .unwrap();
5572 fs::write(&path, format!("v{}", i + 1)).unwrap();
5573 }
5574
5575 let key = canonicalize_key(&path);
5576 let stack_dir = store
5577 .session_dir(DEFAULT_SESSION_ID)
5578 .unwrap()
5579 .join(BackupStore::path_hash(&key));
5580 assert_eq!(backup_content_names(&stack_dir).len(), 3);
5581
5582 store.set_policy(BackupPolicy {
5583 enabled: true,
5584 max_depth: 1,
5585 max_file_size: None,
5586 });
5587
5588 assert_eq!(backup_content_names(&stack_dir).len(), 1);
5589 let meta: serde_json::Value =
5590 serde_json::from_str(&fs::read_to_string(stack_dir.join("meta.json")).unwrap())
5591 .unwrap();
5592 assert_eq!(meta_entry_count(&meta), Some(1));
5593 let mut fresh = BackupStore::new();
5594 fresh.set_storage_dir(storage.path().to_path_buf(), 72);
5595 assert_eq!(fresh.history(DEFAULT_SESSION_ID, &path).len(), 1);
5596 }
5597
5598 #[test]
5599 fn v2_missing_content_fails_closed() {
5600 let dir = tempfile::tempdir().unwrap();
5601 let path = dir.path().join("missing-content.txt");
5602 fs::write(&path, "current").unwrap();
5603 let key = canonicalize_key(&path);
5604 let session_dir = dir
5605 .path()
5606 .join("backups")
5607 .join(BackupStore::session_hash(DEFAULT_SESSION_ID));
5608 let stack_dir = session_dir.join(BackupStore::path_hash(&key));
5609 fs::create_dir_all(&stack_dir).unwrap();
5610 fs::write(
5611 session_dir.join("session.json"),
5612 serde_json::to_string_pretty(&serde_json::json!({
5613 "schema_version": SCHEMA_VERSION,
5614 "session_id": DEFAULT_SESSION_ID,
5615 "last_accessed": current_timestamp(),
5616 }))
5617 .unwrap(),
5618 )
5619 .unwrap();
5620 fs::write(
5621 stack_dir.join("meta.json"),
5622 serde_json::to_string_pretty(&serde_json::json!({
5623 "schema_version": SCHEMA_VERSION,
5624 "format_version": "v2",
5625 "session_id": DEFAULT_SESSION_ID,
5626 "path": key.display().to_string(),
5627 "count": 1,
5628 "entries": [{
5629 "backup_id": "backup-0",
5630 "timestamp": current_timestamp(),
5631 "order": "1",
5632 "description": "missing",
5633 "kind": "content",
5634 "content_path": "bak_1_backup-0.bak",
5635 }]
5636 }))
5637 .unwrap(),
5638 )
5639 .unwrap();
5640
5641 let mut store = BackupStore::new();
5642 store.set_storage_dir(dir.path().to_path_buf(), 72);
5643 let error = store.restore_latest(DEFAULT_SESSION_ID, &path).unwrap_err();
5644 assert_eq!(error.code(), "io_error");
5645 }
5646
5647 #[test]
5648 fn v2_orphan_files_are_ignored_then_pruned() {
5649 let dir = tempfile::tempdir().unwrap();
5650 let path = dir.path().join("orphan.txt");
5651 fs::write(&path, "v0").unwrap();
5652 let mut store = BackupStore::new();
5653 store.set_storage_dir(dir.path().to_path_buf(), 72);
5654 store
5655 .snapshot(DEFAULT_SESSION_ID, &path, "first")
5656 .unwrap()
5657 .unwrap();
5658 let key = canonicalize_key(&path);
5659 let stack_dir = store
5660 .session_dir(DEFAULT_SESSION_ID)
5661 .unwrap()
5662 .join(BackupStore::path_hash(&key));
5663 fs::write(stack_dir.join("bak_999_orphan.bak"), "orphan").unwrap();
5664
5665 assert_eq!(store.history(DEFAULT_SESSION_ID, &path).len(), 1);
5666 fs::write(&path, "v1").unwrap();
5667 store
5668 .snapshot(DEFAULT_SESSION_ID, &path, "second")
5669 .unwrap()
5670 .unwrap();
5671 assert!(!stack_dir.join("bak_999_orphan.bak").exists());
5672 }
5673
5674 fn backup_content_names(dir: &Path) -> HashSet<String> {
5675 fs::read_dir(dir)
5676 .unwrap()
5677 .filter_map(|entry| entry.ok())
5678 .filter_map(|entry| entry.file_name().to_str().map(str::to_string))
5679 .filter(|name| name.starts_with("bak_") && name.ends_with(".bak"))
5680 .collect()
5681 }
5682}