Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
39 changes: 37 additions & 2 deletions crates/filesync/src/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -1351,7 +1354,39 @@ fn send_paths_to_server(
gui_state: &Option<SharedState>,
paths: Vec<PathBuf>,
) -> 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<PathBuf> = 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))
Expand Down
70 changes: 63 additions & 7 deletions crates/filesync/src/common.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<PathBuf>, usize) {
let (paths, original_count) = self.take_deletes(engine.root());
) -> (Vec<PathBuf>, 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<PathBuf> = 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 {
Expand All @@ -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)
}
}

Expand Down Expand Up @@ -274,7 +295,42 @@ pub fn handle_recv_bundle(
bus: &Option<Arc<MessageBus>>,
log_prefix: &str,
) -> io::Result<BundleApplied> {
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();
Expand Down
84 changes: 84 additions & 0 deletions crates/filesync/src/ledger.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<PathBuf, Tombstone>) {
for (path, tomb) in remote {
Expand Down Expand Up @@ -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<String>, Option<u64>) {
(
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));
}
}
Loading
Loading