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/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)); + } +} 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()); + } +} 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/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(), 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) } 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 + ); + } +}