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