From 738fff766ea76890cb81b022901a7b1c1d1127d3 Mon Sep 17 00:00:00 2001 From: Guy Dols Date: Wed, 16 Sep 2026 18:06:41 +0200 Subject: [PATCH 1/5] feat(filesync): add covering-tombstone lookup for dir deletes Dir deletes tombstone only the top path, so children resurrected via exact-match checks. Add nearest-ancestor lookup plus has_newer_than_prefix so a parent delete covers stale children while newer-mtime recreates still win. --- crates/filesync/src/ledger.rs | 84 +++++++++++++++++++++++++++++++++++ 1 file changed, 84 insertions(+) diff --git a/crates/filesync/src/ledger.rs b/crates/filesync/src/ledger.rs index 23171a8..f1f7dd1 100644 --- a/crates/filesync/src/ledger.rs +++ b/crates/filesync/src/ledger.rs @@ -148,6 +148,35 @@ impl DeletionLedger { } } + /// Nearest tombstone at `path` or any ancestor (prefix) directory. + /// + /// Walks `path` and its ancestors nearest-first (exact match wins, + /// then parent, grandparent, ...). Returns the first tombstone found, + /// or `None` when neither `path` nor any ancestor is tombstoned. + /// Used so a directory delete (which tombstones only the top path) + /// still covers resurrected children. + pub fn covering_tombstone(&self, path: &Path) -> Option<&Tombstone> { + for ancestor in path.ancestors() { + if ancestor.as_os_str().is_empty() { + break; + } + if let Some(t) = self.entries.get(ancestor) { + return Some(t); + } + } + None + } + + /// Prefix-aware variant of [`DeletionLedger::has_newer_than`]: true + /// when the exact or any covering ancestor tombstone is strictly newer + /// than `mtime_ms`. + pub fn has_newer_than_prefix(&self, path: &Path, mtime_ms: u64) -> bool { + match self.covering_tombstone(path) { + Some(t) => t.deleted_at_ms > mtime_ms, + None => false, + } + } + /// Merge remote tombstones, keeping the newest per path (LWW). pub fn merge_remote(&mut self, remote: HashMap) { for (path, tomb) in remote { @@ -207,3 +236,58 @@ impl DeletionLedger { self.entries.iter() } } + +#[cfg(test)] +mod tests { + use super::*; + + fn tomb(path: &str, at: u64) -> (PathBuf, u64, String, Option, Option) { + ( + PathBuf::from(path), + at, + "peer".to_string(), + None, + None, + ) + } + + #[test] + fn covering_tombstone_finds_parent_prefix() { + let mut l = DeletionLedger::new(); + let (p, at, node, ph, pm) = tomb("docs", 100); + l.record(p, at, node, ph, pm); + let cover = l + .covering_tombstone(Path::new("docs/a.txt")) + .expect("child must be covered by parent tombstone"); + assert_eq!(cover.path, PathBuf::from("docs")); + assert!(l.has_newer_than_prefix(Path::new("docs/a.txt"), 50)); + assert!(!l.has_newer_than_prefix(Path::new("docs/a.txt"), 150)); + } + + #[test] + fn covering_tombstone_prefers_nearest_ancestor() { + let mut l = DeletionLedger::new(); + let (p, at, node, ph, pm) = tomb("docs", 100); + l.record(p, at, node, ph, pm); + let (p2, at2, node2, ph2, pm2) = tomb("docs/sub", 200); + l.record(p2, at2, node2, ph2, pm2); + let cover = l + .covering_tombstone(Path::new("docs/sub/a.txt")) + .expect("must find nearest cover"); + assert_eq!(cover.path, PathBuf::from("docs/sub")); + assert_eq!(cover.deleted_at_ms, 200); + } + + #[test] + fn covering_tombstone_exact_wins_and_miss_returns_none() { + let mut l = DeletionLedger::new(); + let (p, at, node, ph, pm) = tomb("docs/a.txt", 100); + l.record(p, at, node, ph, pm); + let exact = l + .covering_tombstone(Path::new("docs/a.txt")) + .expect("exact match must be returned"); + assert_eq!(exact.path, PathBuf::from("docs/a.txt")); + assert!(l.covering_tombstone(Path::new("other/b.txt")).is_none()); + assert!(!l.has_newer_than_prefix(Path::new("other/b.txt"), 0)); + } +} From d8a01afa5e4ff9985ac1d2b7b748eec4692633b4 Mon Sep 17 00:00:00 2001 From: Guy Dols Date: Wed, 16 Sep 2026 18:06:46 +0200 Subject: [PATCH 2/5] fix(filesync): veto resurrected children via covering tombstone filter_resurrected only checked exact paths, so files under a deleted dir were re-uploaded as new. Fall back to ancestor tombstones for children (timestamp-only) while keeping hash/mtime checks for exact paths, so stale subtrees stay deleted and legitimate recreates win. --- crates/filesync/src/manifest.rs | 191 ++++++++++++++++++++++++++------ 1 file changed, 158 insertions(+), 33 deletions(-) diff --git a/crates/filesync/src/manifest.rs b/crates/filesync/src/manifest.rs index 92b7cee..bdcaa47 100644 --- a/crates/filesync/src/manifest.rs +++ b/crates/filesync/src/manifest.rs @@ -347,16 +347,24 @@ pub fn compute_send_list(local: &Manifest, remote: &Manifest, is_server: bool) - /// `send_list` is typically the output of [`compute_send_list`]; `ledger` /// is the peer's (or merged) deletion ledger. Returns `(to_send, /// resurrected)`: -/// - A path is vetoed (dropped from `to_send`) when the ledger holds a -/// tombstone strictly newer than the local file's `modified_ms` (via -/// [`DeletionLedger::has_newer_than`]) AND the local copy looks like the -/// same-or-older version the tombstone describes (local content hash +/// - A path is vetoed (dropped from `to_send`) when the ledger holds an +/// exact tombstone strictly newer than the local file's `modified_ms` +/// (via [`DeletionLedger::has_newer_than`]) AND the local copy looks like +/// the same-or-older version the tombstone describes (local content hash /// equals `prev_hash`, or local mtime is not newer than `prev_mtime_ms`; /// missing prev info fails closed toward the delete). +/// - A path with no exact tombstone but a covering ancestor tombstone +/// (e.g. child of a deleted directory, via +/// [`DeletionLedger::covering_tombstone`]) is vetoed purely on timestamp +/// ordering: ancestor `deleted_at_ms > local modified_ms` vetoes (the +/// ancestor's `prev_hash`/`prev_mtime_ms` describe the directory, not the +/// child, so they are not consulted). This also covers stale directory +/// entries (`hash == [0; 32]`). /// - A path whose local `modified_ms` is strictly newer than the -/// tombstone is a legitimate recreate: it stays in `to_send` and is -/// reported in `resurrected` so the caller can drop the stale tombstone -/// (e.g. via `SyncEngine::clear_tombstones`, which removes + saves). +/// exact or covering tombstone is a legitimate recreate: it stays in +/// `to_send` and is reported in `resurrected` so the caller can drop the +/// stale tombstone (e.g. via `SyncEngine::clear_tombstones`, which +/// removes + saves). /// - `None` ledger (or no tombstone / no local entry) passes through. /// /// The existing `is_server` tie-break in [`compute_send_list`] is @@ -377,38 +385,66 @@ pub fn filter_resurrected( to_send.push(path); continue; }; - let Some(tomb) = ledger.get(&path) else { + // Exact tombstone first (preserves prev_hash/prev_mtime semantics); + // otherwise fall back to the nearest covering ancestor tombstone. + if let Some(tomb) = ledger.get(&path) { + if local_meta.modified_ms > tomb.deleted_at_ms { + // Legitimate recreate: local version is newer than the delete. + resurrected.push(path.clone()); + to_send.push(path); + } else if ledger.has_newer_than(&path, local_meta.modified_ms) { + let hash_matches = match &tomb.prev_hash { + Some(prev) => crate::hex(&local_meta.hash) == *prev, + None => true, + }; + let mtime_not_newer = match tomb.prev_mtime_ms { + Some(prev_mtime) => local_meta.modified_ms <= prev_mtime, + None => true, + }; + if hash_matches || mtime_not_newer { + // Stale copy loses to the tombstone — veto the upload. + log::debug!( + "filter_resurrected: vetoing {:?} (tombstone @{} by {} beats local mtime {})", + path, + tomb.deleted_at_ms, + tomb.deleter_node, + local_meta.modified_ms, + ); + } else { + to_send.push(path); + } + } else { + // Equal timestamps or no ordering signal — pass through and let + // the normal sync decision stand. + to_send.push(path); + } + continue; + } + // No exact tombstone: check for a covering ancestor (dir delete). + let Some(cover) = ledger.covering_tombstone(&path) else { to_send.push(path); continue; }; - if local_meta.modified_ms > tomb.deleted_at_ms { - // Legitimate recreate: local version is newer than the delete. + if local_meta.modified_ms > cover.deleted_at_ms { + // Legitimate recreate under a deleted dir: local version is newer. resurrected.push(path.clone()); to_send.push(path); - } else if ledger.has_newer_than(&path, local_meta.modified_ms) { - let hash_matches = match &tomb.prev_hash { - Some(prev) => crate::hex(&local_meta.hash) == *prev, - None => true, - }; - let mtime_not_newer = match tomb.prev_mtime_ms { - Some(prev_mtime) => local_meta.modified_ms <= prev_mtime, - None => true, - }; - if hash_matches || mtime_not_newer { - // Stale copy loses to the tombstone — veto the upload. - log::debug!( - "filter_resurrected: vetoing {:?} (tombstone @{} by {} beats local mtime {})", - path, - tomb.deleted_at_ms, - tomb.deleter_node, - local_meta.modified_ms, - ); - } else { - to_send.push(path); - } + } else if cover.deleted_at_ms > local_meta.modified_ms { + // Stale child loses to the ancestor delete — veto the upload. + // Note: no prev_hash/prev_mtime check here; those describe the + // deleted directory itself, not this child (and stale dir + // entries with hash == [0; 32] must still be vetoed). + log::debug!( + "filter_resurrected: vetoing {:?} (covering tombstone {:?} @{} by {} beats local mtime {})", + path, + cover.path, + cover.deleted_at_ms, + cover.deleter_node, + local_meta.modified_ms, + ); } else { - // Equal timestamps or no ordering signal — pass through and let - // the normal sync decision stand. + // Equal timestamps — pass through and let the normal sync + // decision stand. to_send.push(path); } } @@ -435,3 +471,92 @@ pub fn diff_manifests(old: &Manifest, new: &Manifest) -> (Vec, Vec FileMetadata { + FileMetadata { + rel_path: PathBuf::from(path), + size: 1, + hash: [byte; 32], + modified_ms: mtime, + change_sequence: 0, + is_dir: false, + } + } + + fn dir_meta(path: &str, mtime: u64) -> FileMetadata { + FileMetadata { + rel_path: PathBuf::from(path), + size: 0, + hash: [0; 32], + modified_ms: mtime, + change_sequence: 0, + is_dir: true, + } + } + + fn manifest_with(entries: Vec) -> Manifest { + Manifest { + files: entries.into_iter().map(|m| (m.rel_path.clone(), m)).collect(), + node_id: "test".to_string(), + } + } + + fn ledger_with(path: &str, deleted_at: u64) -> DeletionLedger { + let mut l = DeletionLedger::new(); + l.record( + PathBuf::from(path), + deleted_at, + "peer".to_string(), + None, + None, + ); + l + } + + #[test] + fn child_file_vetoed_by_parent_tombstone() { + // Dir delete tombstones only the top path; the stale child must veto. + let local = manifest_with(vec![meta("docs/a.txt", 0xAA, 50)]); + let ledger = ledger_with("docs", 100); + let (to_send, resurrected) = filter_resurrected( + vec![PathBuf::from("docs/a.txt")], + &local, + Some(&ledger), + ); + assert!(to_send.is_empty(), "stale child must be vetoed"); + assert!(resurrected.is_empty()); + } + + #[test] + fn child_recreate_with_newer_mtime_allowed() { + // Legitimate recreate wins: child mtime newer than parent delete. + let local = manifest_with(vec![meta("docs/a.txt", 0xBB, 150)]); + let ledger = ledger_with("docs", 100); + let (to_send, resurrected) = filter_resurrected( + vec![PathBuf::from("docs/a.txt")], + &local, + Some(&ledger), + ); + assert_eq!(to_send, vec![PathBuf::from("docs/a.txt")]); + assert_eq!(resurrected, vec![PathBuf::from("docs/a.txt")]); + } + + #[test] + fn stale_dir_entry_vetoed_by_parent_tombstone() { + // Directory entries carry hash == [0; 32]; a stale subdir must veto. + let local = manifest_with(vec![dir_meta("docs/sub", 50)]); + let ledger = ledger_with("docs", 100); + let (to_send, resurrected) = filter_resurrected( + vec![PathBuf::from("docs/sub")], + &local, + Some(&ledger), + ); + assert!(to_send.is_empty(), "stale dir child must be vetoed"); + assert!(resurrected.is_empty()); + } +} From 700ae0a8db2c2444597eb626f839e034568f6cf5 Mon Sep 17 00:00:00 2001 From: Guy Dols Date: Wed, 16 Sep 2026 18:06:50 +0200 Subject: [PATCH 3/5] fix(filesync): watch moved-in dirs and fix stale wd_map on rename Moved-in subtrees stayed unwatched and wd->path entries kept pointing at the old location, so events inside a moved dir were mis-attributed to the old path and recreated phantom folders. Recursively watch CREATE/MOVED_TO dirs, rewrite wd_map prefixes on rename, drop stale watches on IGNORED/MOVE_SELF/DELETE_SELF, and extend move pairing to 2s. --- crates/filesync/src/watcher.rs | 175 +++++++++++++++++++++++++++++++-- 1 file changed, 168 insertions(+), 7 deletions(-) diff --git a/crates/filesync/src/watcher.rs b/crates/filesync/src/watcher.rs index 17138cd..d98d036 100644 --- a/crates/filesync/src/watcher.rs +++ b/crates/filesync/src/watcher.rs @@ -15,6 +15,71 @@ pub enum FsEvent { Renamed(PathBuf, PathBuf), } +/// Add inotify watches for `path` and every directory below it. +/// +/// A freshly created / moved-in directory may already contain a populated +/// subtree (e.g. `mv outside/ new/`), so watching only the top level would +/// leave the new children unwatched. +fn add_watch_recursive( + inotify: &mut Inotify, + wd_map: &mut HashMap, + path: &PathBuf, + mask: WatchMask, +) { + for entry in WalkDir::new(path).into_iter().filter_map(Result::ok) { + if entry.file_type().is_dir() { + match inotify.watches().add(entry.path(), mask) { + Ok(new_wd) => { + wd_map.insert(new_wd, entry.path().to_path_buf()); + } + Err(e) => warn!("failed to watch {:?}: {e}", entry.path()), + } + } + } +} + +/// Pure prefix-remap used by [`rewrite_wd_map_for_rename`]: returns the new +/// absolute path if `path` is `from_abs` itself or lives below it. +fn remap_path_for_rename( + path: &std::path::Path, + from_abs: &std::path::Path, + to_abs: &std::path::Path, +) -> Option { + if path == from_abs { + Some(to_abs.to_path_buf()) + } else if let Ok(stripped) = path.strip_prefix(from_abs) { + Some(to_abs.join(stripped)) + } else { + None + } +} + +/// Rewrite `wd_map` entries whose path is `from_rel` (or a descendant of it) +/// to point at `to_rel`. +/// +/// Inotify watches follow the renamed inode, so the kernel watch descriptors +/// stay valid across a rename — only our `wd -> path` bookkeeping goes stale. +/// Without this rewrite, later events inside the moved tree would be reported +/// under the old (phantom) path. +fn rewrite_wd_map_for_rename( + wd_map: &mut HashMap, + root: &std::path::Path, + from_rel: &std::path::Path, + to_rel: &std::path::Path, +) { + let from_abs = root.join(from_rel); + let to_abs = root.join(to_rel); + let updates: Vec<(WatchDescriptor, PathBuf)> = wd_map + .iter() + .filter_map(|(wd, path)| { + remap_path_for_rename(path, &from_abs, &to_abs).map(|p| (wd.clone(), p)) + }) + .collect(); + for (wd, new_path) in updates { + wd_map.insert(wd, new_path); + } +} + pub fn start_watcher( root: PathBuf, tx: Sender, @@ -25,6 +90,7 @@ pub fn start_watcher( | WatchMask::CLOSE_WRITE | WatchMask::DELETE | WatchMask::DELETE_SELF + | WatchMask::MOVE_SELF | WatchMask::MOVED_FROM | WatchMask::MOVED_TO; @@ -49,7 +115,9 @@ pub fn start_watcher( .spawn(move || { let mut buffer = vec![0u8; 65536]; let mut pending_moves: HashMap = HashMap::new(); - const ORPHAN_TIMEOUT: Duration = Duration::from_millis(500); + // Long enough to keep MOVED_FROM/MOVED_TO pairs together under + // load so a slow move isn't split into Delete + Create. + const ORPHAN_TIMEOUT: Duration = Duration::from_millis(2000); loop { let events: Vec<_> = match inotify.read_events(&mut buffer) { @@ -71,6 +139,35 @@ pub fn start_watcher( }; for (wd, mask, cookie, name) in events { + // Watch-lifecycle events target the watched dir itself + // (no filename) and must be handled before the + // name-based logic below, which would otherwise skip + // them via `None => continue`. + if mask.contains(EventMask::IGNORED) { + if let Some(old) = wd_map.remove(&wd) { + debug!("inotify IGNORED, dropped stale watch for {old:?}"); + } + continue; + } + if mask.contains(EventMask::MOVE_SELF) || mask.contains(EventMask::DELETE_SELF) + { + if let Some(dir) = wd_map.get(&wd).cloned() { + // Drop the stale entry. For a rename the watch + // follows the inode: the paired MOVED_FROM / + // MOVED_TO handler below rewrites wd_map and + // re-adds the subtree, so dropping here is safe + // whichever event arrives first. Anything missed + // is picked up by the next periodic rescan. + wd_map.remove(&wd); + let _ = inotify.watches().remove(wd); + debug!( + "inotify self-event mask={:?} dropped stale watch for {dir:?}", + mask + ); + } + continue; + } + let name = match name { Some(n) => n, None => continue, @@ -92,10 +189,46 @@ pub fn start_watcher( if mask.contains(EventMask::MOVED_FROM) { pending_moves.insert(cookie, (rel, Instant::now())); } else if mask.contains(EventMask::MOVED_TO) { + let is_dir = mask.contains(EventMask::ISDIR); if let Some((from_rel, _)) = pending_moves.remove(&cookie) { debug!("inotify RENAME {:?} → {:?}", from_rel, rel); + // Fix stale bookkeeping: watches followed the + // inode, so point them at the new location. + rewrite_wd_map_for_rename(&mut wd_map, &root, &from_rel, &rel); + // Drop the kernel watch if the source path somehow + // survived (e.g. cross-device move leaving an + // empty dir); ignore errors — it is usually gone. + let from_abs = root.join(&from_rel); + if !from_abs.exists() { + let stale: Vec = wd_map + .iter() + .filter_map(|(w, p)| { + if p == &from_abs { + Some(w.clone()) + } else { + None + } + }) + .collect(); + for w in stale { + wd_map.remove(&w); + let _ = inotify.watches().remove(w); + } + } + if is_dir { + // Re-watch the new subtree: covers entries + // dropped by an early MOVE_SELF as well as + // children created mid-move. + add_watch_recursive(&mut inotify, &mut wd_map, &full, watch_mask); + } let _ = tx.send(FsEvent::Renamed(from_rel, rel)); } else { + // Move from outside the watched tree (or an + // orphaned MOVED_TO): treat as a new path, and + // watch the whole subtree if it is a directory. + if is_dir { + add_watch_recursive(&mut inotify, &mut wd_map, &full, watch_mask); + } let _ = tx.send(FsEvent::WriteComplete(rel)); } } else if mask.contains(EventMask::DELETE) @@ -103,12 +236,10 @@ pub fn start_watcher( { let _ = tx.send(FsEvent::Deleted(rel)); } else if mask.contains(EventMask::CREATE) && mask.contains(EventMask::ISDIR) { - match inotify.watches().add(&full, watch_mask) { - Ok(new_wd) => { - wd_map.insert(new_wd, full); - } - Err(e) => warn!("failed to watch new dir {:?}: {e}", rel), - } + // Watch the whole new subtree, not just the top dir: + // it may already contain files created before the + // watch was added. + add_watch_recursive(&mut inotify, &mut wd_map, &full, watch_mask); let _ = tx.send(FsEvent::Changed(rel)); } else if mask.contains(EventMask::CLOSE_WRITE) { let _ = tx.send(FsEvent::WriteComplete(rel)); @@ -134,3 +265,33 @@ pub fn start_watcher( Ok(handle) } + +#[cfg(test)] +mod tests { + use super::remap_path_for_rename; + use std::path::Path; + + #[test] + fn remap_dir_rename_prefix() { + let root = Path::new("/root"); + let from_abs = root.join("a"); + let to_abs = root.join("b"); + assert_eq!( + remap_path_for_rename(&from_abs, &from_abs, &to_abs), + Some(to_abs.clone()) + ); + assert_eq!( + remap_path_for_rename(&from_abs.join("sub/dir"), &from_abs, &to_abs), + Some(to_abs.join("sub/dir")) + ); + // Sibling with a similar name prefix must not match. + assert_eq!( + remap_path_for_rename(&root.join("ab"), &from_abs, &to_abs), + None + ); + assert_eq!( + remap_path_for_rename(&root.join("other"), &from_abs, &to_abs), + None + ); + } +} From 8ec38c9add600cf1eabaf2485caeeb86e4339a8d Mon Sep 17 00:00:00 2001 From: Guy Dols Date: Wed, 16 Sep 2026 18:07:08 +0200 Subject: [PATCH 4/5] fix(filesync): recursive idempotent rename plus live resurrection guard apply_rename moved only the top manifest entry, leaving children on the old path so scans recreated the source dir. Remap the full subtree component-wise, suppress all remapped keys, and make missing-src with present-dst converge instead of erroring. Veto bundle/large-file writes covered by a newer tombstone, lift tombstones on newer-mtime recreates, and clear prefix tombstones on trash restore. Document why Rename keeps no move_id for bincode compat. --- crates/filesync/src/protocol.rs | 8 + crates/filesync/src/sync_engine.rs | 291 ++++++++++++++++++++++++++--- 2 files changed, 270 insertions(+), 29 deletions(-) diff --git a/crates/filesync/src/protocol.rs b/crates/filesync/src/protocol.rs index 6bb2b42..6c17cb9 100644 --- a/crates/filesync/src/protocol.rs +++ b/crates/filesync/src/protocol.rs @@ -116,6 +116,14 @@ pub enum Message { from: PathBuf, to: PathBuf, }, + // NOTE: a `move_id: u64` (e.g. `timestamp_id()`) was considered for + // Rename dedup, but is deliberately omitted: `Message` is serialised + // with bincode (positional, not self-describing), so adding a field + // breaks wire compat with older peers (old bytes fail to decode, and + // new bytes fail on old peers). Dedup instead relies on the idempotent + // `apply_rename` (source-absent + dst-present converges the manifest + // and returns Ok). If a move id is ever added, it needs a protocol + // version bump plus dual-decode handling. InsufficientDiskSpace { available_bytes: u64, diff --git a/crates/filesync/src/sync_engine.rs b/crates/filesync/src/sync_engine.rs index acfb0b0..6d2e8f7 100644 --- a/crates/filesync/src/sync_engine.rs +++ b/crates/filesync/src/sync_engine.rs @@ -317,7 +317,29 @@ impl SyncEngine { } pub fn restore_trash_entry(&self, id: &str) -> Result<(), String> { - self.trash_manager.restore_entry(id) + // Capture the original path first so a successful restore can lift + // the tombstones covering that prefix (otherwise the next send + // would veto the restored file as a resurrection). + let original: Option = self + .trash_manager + .list_trash() + .iter() + .find(|e| e.id == id) + .map(|e| e.original_path.clone()); + let result = self.trash_manager.restore_entry(id); + if result.is_ok() { + if let Some(prefix) = original { + let to_clear: Vec = self + .ledger + .read() + .iter() + .filter(|(k, _)| *k == &prefix || k.strip_prefix(&prefix).is_ok()) + .map(|(k, _)| k.clone()) + .collect(); + self.clear_tombstones(&to_clear); + } + } + result } pub fn purge_trash_entry(&self, id: &str) -> Result<(), String> { @@ -450,6 +472,46 @@ impl SyncEngine { self.ledger.read().has_newer_than(path, mtime_ms) } + /// Covering-tombstone veto: delegates to + /// [`DeletionLedger::covering_tombstone`] (nearest tombstone at `path` + /// or any ancestor). Returns the tombstoned path when it vetoes `path` + /// at `incoming_mtime_ms` (tombstone strictly newer than the incoming + /// mtime); returns `None` when there is no covering tombstone or the + /// incoming write is a legitimate recreate (incoming mtime strictly + /// newer — see `clear_tombstone_on_recreate`). + pub fn covering_tombstone_veto( + &self, + path: &Path, + incoming_mtime_ms: u64, + ) -> Option { + let ledger = self.ledger.read(); + match ledger.covering_tombstone(path) { + Some(t) if t.deleted_at_ms > incoming_mtime_ms => Some(t.path.clone()), + _ => None, + } + } + + /// Clear the exact-path tombstone when an incoming write is a + /// legitimate recreate (incoming mtime strictly newer than the + /// tombstone). Persists only when something was removed. Returns true + /// when a tombstone was lifted. + fn clear_tombstone_on_recreate(&self, path: &Path, incoming_mtime_ms: u64) -> bool { + let should_clear = self + .ledger + .read() + .get(path) + .map(|t| incoming_mtime_ms > t.deleted_at_ms) + .unwrap_or(false); + if should_clear { + self.ledger.write().remove(path); + if let Err(e) = self.ledger.read().save_atomic(&self.root) { + log::warn!("ledger: save after recreate-clear failed: {e}"); + } + return true; + } + false + } + /// Drop tombstones older than the 90-day TTL. Persists only when at /// least one entry was pruned. Returns the number removed. pub fn prune_ledger(&self) -> usize { @@ -572,6 +634,26 @@ impl SyncEngine { continue; } + // Live resurrection guard: veto writes covered by a tombstone + // strictly newer than the incoming mtime (path or ancestor). + // A strictly newer incoming mtime is a legitimate recreate and + // lifts the exact-path tombstone instead. + if let Some(tomb_path) = + self.covering_tombstone_veto(&fd.metadata.rel_path, fd.metadata.modified_ms) + { + log::warn!( + "apply_bundle: vetoing {:?} (covering tombstone {:?} beats incoming mtime {})", + fd.metadata.rel_path, + tomb_path, + fd.metadata.modified_ms, + ); + continue; + } + self.clear_tombstone_on_recreate( + &fd.metadata.rel_path, + fd.metadata.modified_ms, + ); + let full = self.root.join(&fd.metadata.rel_path); if !fd.metadata.is_dir { @@ -659,6 +741,22 @@ impl SyncEngine { return Ok(()); } + // Live resurrection guard (same policy as apply_bundle): veto the + // stream setup when a covering tombstone beats the incoming mtime; + // lift the exact tombstone on a legitimate recreate. + if let Some(tomb_path) = + self.covering_tombstone_veto(&metadata.rel_path, metadata.modified_ms) + { + log::warn!( + "begin_large_file: vetoing {:?} (covering tombstone {:?} beats incoming mtime {})", + metadata.rel_path, + tomb_path, + metadata.modified_ms, + ); + return Ok(()); + } + self.clear_tombstone_on_recreate(&metadata.rel_path, metadata.modified_ms); + let dst = self.root.join(&metadata.rel_path); if let Some(parent) = dst.parent() { fs::create_dir_all(parent)?; @@ -781,6 +879,20 @@ impl SyncEngine { } }; + // Re-check the resurrection guard at commit time: a tombstone may + // have arrived between LargeFileStart and commit. Veto the write + // (clean up tmp) instead of recreating a deleted path. + if let Some(tomb_path) = self.covering_tombstone_veto(path, asm.modified_ms) { + log::warn!( + "commit_large_file: vetoing {path:?} (covering tombstone {tomb_path:?} beats incoming mtime {})", + asm.modified_ms, + ); + let _ = fs::remove_file(&asm.tmp_path); + cleanup_transfer_dirs(&self.root, &asm.tmp_path); + return Ok(FinishResult::Committed); + } + self.clear_tombstone_on_recreate(path, asm.modified_ms); + { use std::io::Read; let mut hasher = blake3::Hasher::new(); @@ -907,41 +1019,148 @@ impl SyncEngine { let dst = self.root.join(to); if !src.exists() { - log::debug!("apply_rename: source absent, skipping: {from:?}"); + // Idempotent retransmit: source gone. If manifest still holds + // stale `from`-prefixed keys but `to` already exists on disk, + // converge the manifest without touching the filesystem. + if dst.exists() { + let pairs: Vec<(PathBuf, PathBuf)> = { + let m = self.manifest.read(); + let mut v = Vec::new(); + for k in m.files.keys() { + if k == from { + v.push((k.clone(), to.clone())); + } else if let Ok(stripped) = k.strip_prefix(from) { + v.push((k.clone(), to.join(stripped))); + } + } + v + }; + if !pairs.is_empty() { + let mut suppressed = Vec::with_capacity(pairs.len() * 2); + { + let mut m = self.manifest.write(); + for (old, new) in pairs { + if let Some(mut meta) = m.files.remove(&old) { + meta.rel_path = new.clone(); + if meta.is_dir { + if let Ok(fs_meta) = + fs::metadata(self.root.join(&new)) + { + if let Ok(modified) = fs_meta.modified() { + if let Ok(dur) = modified.duration_since( + std::time::SystemTime::UNIX_EPOCH, + ) { + meta.modified_ms = + dur.as_millis() as u64; + } + } + } + } + m.files.insert(new.clone(), meta); + } + suppressed.push(old); + suppressed.push(new); + } + } + { + let mut sup = self.suppressed.write(); + for p in &suppressed { + sup.insert(p.clone()); + } + } + self.schedule_unsuppress(suppressed); + } + } else { + log::debug!("apply_rename: source absent, skipping: {from:?}"); + } return Ok(()); } + // Collect every manifest key under `from` (component-wise via + // strip_prefix, so `a/b` does not match `a/bc`) and remap it under + // `to`. Includes the top entry itself. + let pairs: Vec<(PathBuf, PathBuf)> = { + let m = self.manifest.read(); + let mut v = Vec::new(); + for k in m.files.keys() { + if k == from { + v.push((k.clone(), to.clone())); + } else if let Ok(stripped) = k.strip_prefix(from) { + v.push((k.clone(), to.join(stripped))); + } + } + v + }; + { let mut sup = self.suppressed.write(); sup.insert(from.clone()); sup.insert(to.clone()); + for (old, new) in &pairs { + sup.insert(old.clone()); + sup.insert(new.clone()); + } } if let Some(parent) = dst.parent() { fs::create_dir_all(parent)?; } - fs::rename(&src, &dst)?; + // Idempotent when `dst` already exists (retransmit/race): files + // overwrite via rename; a failing dir-onto-dir rename still + // converges via the manifest remap below instead of erroring. + let rename_result = fs::rename(&src, &dst); + if let Err(e) = rename_result { + if dst.exists() { + log::warn!("apply_rename: {from:?} → {to:?} rename failed but dst exists, converging manifest: {e}"); + } else { + return Err(e); + } + } { - let mut m = self.manifest.write(); - if let Some(mut meta) = m.files.remove(from) { - meta.rel_path = to.clone(); - - if let Ok(fs_meta) = fs::metadata(&dst) { - if let Ok(modified) = fs_meta.modified() { - if let Ok(dur) = modified.duration_since(std::time::SystemTime::UNIX_EPOCH) - { - meta.modified_ms = dur.as_millis() as u64; + let mut suppressed = Vec::with_capacity(pairs.len() * 2 + 2); + suppressed.push(from.clone()); + suppressed.push(to.clone()); + { + let mut m = self.manifest.write(); + for (old, new) in pairs { + if let Some(mut meta) = m.files.remove(&old) { + meta.rel_path = new.clone(); + // Refresh dir mtime from disk after the move; files + // keep the incoming manifest mtime via watcher/scan. + if meta.is_dir { + if let Ok(fs_meta) = fs::metadata(self.root.join(&new)) { + if let Ok(modified) = fs_meta.modified() { + if let Ok(dur) = modified.duration_since( + std::time::SystemTime::UNIX_EPOCH, + ) { + meta.modified_ms = dur.as_millis() as u64; + } + } + } + } else if let Ok(fs_meta) = fs::metadata(self.root.join(&new)) { + if let Ok(modified) = fs_meta.modified() { + if let Ok(dur) = modified.duration_since( + std::time::SystemTime::UNIX_EPOCH, + ) { + meta.modified_ms = dur.as_millis() as u64; + } + } } + m.files.insert(new.clone(), meta); } + suppressed.push(old); + suppressed.push(new); } - m.files.insert(to.clone(), meta); } - } + // Deduplicate before scheduling (pairs already include from/to). + suppressed.sort(); + suppressed.dedup(); - log::info!("apply_rename: {from:?} → {to:?}"); - self.schedule_unsuppress(vec![from.clone(), to.clone()]); + log::info!("apply_rename: {from:?} → {to:?}"); + self.schedule_unsuppress(suppressed); + } Ok(()) } @@ -961,8 +1180,29 @@ impl SyncEngine { continue; } - self.suppressed_deletes.write().insert(rel.clone()); - removed.push(rel.clone()); + // Collect the top path plus every manifest child under it + // (component-wise: `a/b` must not match `a/bc`) BEFORE removal, + // so delete-suppression covers the whole subtree recursively. + let subtree: Vec = { + let m = self.manifest.read(); + let mut v = vec![rel.clone()]; + for k in m.files.keys() { + if k != rel { + if let Ok(stripped) = k.strip_prefix(rel) { + let _ = stripped; + v.push(k.clone()); + } + } + } + v + }; + { + let mut sup = self.suppressed_deletes.write(); + for p in &subtree { + sup.insert(p.clone()); + removed.push(p.clone()); + } + } // Snapshot manifest info BEFORE removal for the tombstone. let (prev_hash, prev_mtime_ms) = self @@ -976,19 +1216,10 @@ impl SyncEngine { let full = self.root.join(rel); if full.is_dir() { self.trash_manager.move_to_trash(&full, rel, &self.node_id); - let children: Vec = self - .manifest - .read() - .files - .keys() - .filter(|p| p.starts_with(rel)) - .cloned() - .collect(); let mut m = self.manifest.write(); - for child in children { - m.files.remove(&child); + for child in &subtree { + m.files.remove(child); } - m.files.remove(rel); } else { self.trash_manager.move_to_trash(&full, rel, &self.node_id); self.manifest.write().files.remove(rel); @@ -1018,6 +1249,8 @@ impl SyncEngine { } } + removed.sort(); + removed.dedup(); self.schedule_unsuppress_deletes(removed); Ok(count) } From 60f25f0dd0cb9a8e27c635972cd273824c5a9709 Mon Sep 17 00:00:00 2001 From: Guy Dols Date: Wed, 16 Sep 2026 18:07:12 +0200 Subject: [PATCH 5/5] fix(filesync): single delete timestamp plus send-path resurrection veto Flushes stamped ledger and wire deletes with two now_ms calls so peers disagreed on LWW winners. Reuse one batch stamp, propagate sender stamps on forward, prune fabricated ancestor deletes, suppress rename subtrees recursively, and veto outgoing paths covered by a newer tombstone before streaming or broadcasting. --- crates/filesync/src/client.rs | 39 ++++++++++++++++++- crates/filesync/src/common.rs | 70 +++++++++++++++++++++++++++++++---- crates/filesync/src/server.rs | 37 +++++++++++++++++- 3 files changed, 135 insertions(+), 11 deletions(-) diff --git a/crates/filesync/src/client.rs b/crates/filesync/src/client.rs index 85e91d9..0d2b42b 100644 --- a/crates/filesync/src/client.rs +++ b/crates/filesync/src/client.rs @@ -1272,16 +1272,19 @@ fn flush_to_server( send_paths_to_server(engine, conn, bus, gui_state, stable)?; } - let (paths, delete_count) = pending.take_deletes_with_engine(engine); + let (paths, delete_count, deleted_at_ms) = pending.take_deletes_with_engine(engine); if !paths.is_empty() { debug!( "filesync send: flushing {} delete path(s) ({} pre-expansion)", paths.len(), delete_count ); + // Reuse the single batch stamp from take_deletes_with_engine so + // the ledger tombstone and the wire Delete share one timestamp + // (no second now_ms() re-stamp); deleter is this node. conn.send(&Message::Delete { paths, - deleted_at_ms: crate::ledger::now_ms(), + deleted_at_ms, deleter: engine.node_id().to_string(), })?; if let Some(ref gs) = gui_state { @@ -1351,7 +1354,39 @@ fn send_paths_to_server( gui_state: &Option, paths: Vec, ) -> io::Result<()> { + // Outgoing resurrection guard: drop paths whose local version loses + // to a covering tombstone (same policy as filter_resurrected, applied + // live before streaming). Paths without a manifest entry pass through + // (no mtime to compare); legitimate recreates (local mtime strictly + // newer) pass through. let manifest = engine.get_manifest(); + let mut vetoed = 0usize; + let paths: Vec = paths + .into_iter() + .filter(|p| { + match manifest.files.get(p) { + Some(meta) + if engine + .covering_tombstone_veto(p, meta.modified_ms) + .is_some() => + { + log::warn!( + "filesync send: vetoing {p:?} (covering tombstone beats local mtime {})", + meta.modified_ms, + ); + vetoed += 1; + false + } + _ => true, + } + }) + .collect(); + if paths.is_empty() { + if vetoed > 0 { + log::debug!("filesync send: all {} path(s) vetoed by tombstones", vetoed); + } + return Ok(()); + } let files_count = paths .iter() .filter(|p| manifest.files.get(*p).map(|m| !m.is_dir).unwrap_or(true)) diff --git a/crates/filesync/src/common.rs b/crates/filesync/src/common.rs index be65e89..17a0a76 100644 --- a/crates/filesync/src/common.rs +++ b/crates/filesync/src/common.rs @@ -179,16 +179,37 @@ impl PendingChanges { /// Drain pending deletes and record a local tombstone for each path /// BEFORE it is sent, so a later remote copy cannot resurrect it. - /// Tombstones use wall-clock [`crate::ledger::now_ms`] and this node's - /// id; prev hash/mtime come from the manifest when available. + /// Tombstones use a SINGLE wall-clock [`crate::ledger::now_ms`] stamp + /// for the whole batch (returned for reuse as the `deleted_at_ms` on + /// the wire — callers must NOT re-stamp with a second `now_ms()`). + /// Fabricated ancestor extras from [`expand_deleted_ancestors`] are + /// constrained: an extra parent is kept only when it is absent from + /// the local manifest (still-listed parents are not deletes) and not + /// currently rename-suppressed (rename in-flight race). /// (Next lane: switch client/server flush paths to this method.) pub fn take_deletes_with_engine( &mut self, engine: &SyncEngine, - ) -> (Vec, usize) { - let (paths, original_count) = self.take_deletes(engine.root()); + ) -> (Vec, usize, u64) { + let (mut paths, original_count) = self.take_deletes(engine.root()); + if original_count < paths.len() { + // Extras are appended by expand_deleted_ancestors; prune any + // that are still live in the manifest or rename-suppressed. + let manifest = engine.get_manifest(); + let mut kept: Vec = paths.drain(..original_count).collect(); + for extra in paths { + if manifest.files.contains_key(&extra) { + continue; + } + if engine.is_suppressed(&extra) { + continue; + } + kept.push(extra); + } + paths = kept; + } + let now = crate::ledger::now_ms(); if !paths.is_empty() { - let now = crate::ledger::now_ms(); let node = engine.node_id().to_string(); let manifest = engine.get_manifest(); for p in &paths { @@ -200,7 +221,7 @@ impl PendingChanges { engine.record_delete(p.clone(), now, &node, prev_hash, prev_mtime_ms); } } - (paths, original_count) + (paths, original_count, now) } } @@ -274,7 +295,42 @@ pub fn handle_recv_bundle( bus: &Option>, log_prefix: &str, ) -> io::Result { - let apply_result = engine.apply_bundle(bundle)?; + // Ingress pre-guard (defense in depth; SyncEngine::apply_bundle + // enforces the same policy): log vetoes for tombstone-covered paths + // and drop them before applying so a stale retransmit cannot + // resurrect a deleted path. Legitimate recreates (incoming mtime + // strictly newer than the tombstone) pass through. + let mut vetoed = 0usize; + for fd in &bundle.files { + if engine + .covering_tombstone_veto(&fd.metadata.rel_path, fd.metadata.modified_ms) + .is_some() + { + warn!( + "{log_prefix}: vetoing {:?} from {peer} (covering tombstone beats mtime {})", + fd.metadata.rel_path, fd.metadata.modified_ms, + ); + vetoed += 1; + } + } + let apply_result = if vetoed > 0 { + let filtered = FileBundle { + files: bundle + .files + .iter() + .filter(|fd| { + engine + .covering_tombstone_veto(&fd.metadata.rel_path, fd.metadata.modified_ms) + .is_none() + }) + .cloned() + .collect(), + bundle_id: bundle.bundle_id, + }; + engine.apply_bundle(&filtered)? + } else { + engine.apply_bundle(bundle)? + }; let files_count = bundle.files.iter().filter(|f| !f.metadata.is_dir).count(); let dirs_count = bundle.files.iter().filter(|f| f.metadata.is_dir).count(); diff --git a/crates/filesync/src/server.rs b/crates/filesync/src/server.rs index b82b8f0..388c4f5 100644 --- a/crates/filesync/src/server.rs +++ b/crates/filesync/src/server.rs @@ -1113,7 +1113,8 @@ fn flush_local_changes( "filesync server: broadcasting {} delete(s)", pending.deletes.len() ); - let (paths, _) = pending.take_deletes_with_engine(engine); + // Single batch stamp shared between ledger + wire (no re-stamp). + let (paths, _, deleted_at_ms) = pending.take_deletes_with_engine(engine); if let Some(ref bus) = bus { bus.publish( "filesync", @@ -1135,7 +1136,7 @@ fn flush_local_changes( "", &Message::Delete { paths, - deleted_at_ms: crate::ledger::now_ms(), + deleted_at_ms, deleter: engine.node_id().to_string(), }, ); @@ -1152,6 +1153,38 @@ fn broadcast_paths( sent_bundles: &Arc>>, // Add this parameter ) { const BROADCAST_PIPELINE_DEPTH: usize = 128; + // Outgoing resurrection guard (same policy as the client send path): + // drop paths whose local version loses to a covering tombstone before + // streaming. Unknown paths (no manifest entry) pass through. + let manifest = engine.get_manifest(); + let mut vetoed = 0usize; + let paths: Vec = paths + .into_iter() + .filter(|p| match manifest.files.get(p) { + Some(meta) + if engine + .covering_tombstone_veto(p, meta.modified_ms) + .is_some() => + { + log::warn!( + "filesync server: vetoing {p:?} broadcast (covering tombstone beats local mtime {})", + meta.modified_ms, + ); + vetoed += 1; + false + } + _ => true, + }) + .collect(); + if paths.is_empty() { + if vetoed > 0 { + log::debug!( + "filesync server: all {} path(s) vetoed by tombstones", + vetoed + ); + } + return; + } debug!( "filesync server: broadcast_paths — {} path(s) to {} peer(s)", paths.len(),