From 551a40777ff1f2cfab95d5597744311670dcda2c Mon Sep 17 00:00:00 2001 From: Guy Dols Date: Mon, 6 Jul 2026 00:03:29 +0200 Subject: [PATCH 1/6] Cleanup --- crates/filesync/src/client.rs | 140 +---------- crates/filesync/src/known_hosts.rs | 284 ---------------------- crates/filesync/src/protocol.rs | 2 +- crates/filesync/src/server.rs | 21 +- crates/filesync/src/suspend_detector.rs | 6 +- crates/filesync/src/sync_engine.rs | 56 +---- crates/filesync/tests/test_client.rs | 20 +- crates/filesync/tests/test_known_hosts.rs | 7 +- 8 files changed, 27 insertions(+), 509 deletions(-) diff --git a/crates/filesync/src/client.rs b/crates/filesync/src/client.rs index 5c61913..a114178 100644 --- a/crates/filesync/src/client.rs +++ b/crates/filesync/src/client.rs @@ -34,15 +34,9 @@ pub struct Client { tls_config: Arc, gui_state: Option, awaiting_approval: Arc, - /// The connection currently in use by an in-flight `session()` call, if - /// any. This lets `shutdown()` force-close a live connection (e.g. after - /// the system resumes from suspend) instead of only preventing *future* - /// connection attempts. active_conn: Arc>>>, } -/// RAII guard that clears the client's `active_conn` slot when a `session()` -/// call finishes, regardless of which return path is taken. struct ActiveConnGuard<'a> { slot: &'a Mutex>>, } @@ -142,10 +136,6 @@ impl Client { pub fn shutdown(&self) { debug!("filesync client: shutdown requested"); self.stopped.store(true, Ordering::SeqCst); - // Force-close any connection that a blocked `session()` call is - // currently using. This unblocks the recv/send loops the same way a - // normal network drop would, so the caller's reconnect logic kicks in - // immediately instead of waiting on a possibly-stale socket. if let Some(conn) = self.active_conn.lock().clone() { debug!("filesync client: forcing active connection closed"); conn.shutdown(); @@ -265,9 +255,6 @@ impl Client { ); debug!("filesync session: TLS 1.3 handshake complete"); - // Make this connection reachable from `shutdown()` for the rest of the - // session so a forced shutdown (e.g. suspend/resume) can close it even - // while we're blocked in the receive/send loops below. *self.active_conn.lock() = Some(conn.clone()); let _active_conn_guard = ActiveConnGuard { slot: &self.active_conn, @@ -437,8 +424,6 @@ impl Client { l_files, l_dirs, l_bytes ); - // Update GUI state with local stats right after scan so the Stats - // panel shows real values while the sync is still in progress. if let Some(ref gs) = self.gui_state { let mut s = gs.write(); s.file_count = l_files; @@ -446,11 +431,6 @@ impl Client { s.total_bytes = l_bytes; } - // ── Preemptive disk-space check (client) ────────────────────────────── - // Simulate the server's send-list computation (is_server = true) to - // predict exactly which bytes will arrive. We do this *before* sending - // our ManifestExchange so that, if we are out of space, the server is - // notified before it starts transmitting anything. let bytes_incoming: u64 = manifest::compute_send_list(&remote, &local, true) .iter() .filter_map(|p| remote.files.get(p)) @@ -472,8 +452,7 @@ impl Client { {} B available, {} B required; aborting sync", avail, bytes_incoming ); - // Notify the server before closing so it can surface the - // real reason rather than a generic connection error. + let _ = conn.send(&Message::InsufficientDiskSpace { available_bytes: avail, required_bytes: bytes_incoming, @@ -945,9 +924,6 @@ fn recv_loop( s.begin_sync_activity(); s.files_received += 1; } - - // Note: For large files, we don't have the metadata here to get the sequence number - // The acknowledgment would need to be handled differently for large files } Ok(LargeFileEndOutcome::CommittedWithConflict(ci)) => { debug!("{prefix}: LargeFileEnd committed with conflict {path:?}"); @@ -1013,7 +989,6 @@ fn recv_loop( "{prefix}: ChangeAcknowledgment bundle_id={} sequences={:?}", bundle_id, sequence_numbers ); - // Handle acknowledgment - could be used to track which changes were received } Ok(other) => { warn!("{prefix}: unexpected message in live sync phase — possible protocol issue"); @@ -1289,7 +1264,7 @@ fn format_modified(path: &Path) -> String { } } -fn format_unix_secs(secs: u64) -> String { +pub fn format_unix_secs(secs: u64) -> String { let days = (secs / 86_400) as i64; let rem = secs % 86_400; let hour = rem / 3600; @@ -1308,114 +1283,3 @@ fn format_unix_secs(secs: u64) -> String { format!("{y:04}-{m:02}-{d:02} {hour:02}:{minute:02}") } - -#[cfg(test)] -mod format_tests { - use super::format_unix_secs; - - #[test] - fn epoch_formats_correctly() { - assert_eq!(format_unix_secs(0), "1970-01-01 00:00"); - } - - #[test] - fn known_date_formats_correctly() { - assert_eq!(format_unix_secs(1_717_236_000), "2024-06-01 10:00"); - } -} - -#[cfg(test)] -mod push_gui_conflicts_tests { - use super::push_gui_conflicts; - use crate::exclusions::{ExclusionConfig, Exclusions}; - use crate::gui::state::new_shared_state; - use crate::protocol::{FileBundle, FileData, FileMetadata}; - use crate::sync_engine::SyncEngine; - use std::fs; - use std::path::PathBuf; - use std::sync::Arc; - - fn tmp_dir(label: &str) -> PathBuf { - let d = std::env::temp_dir().join(format!( - "filesync_push_gui_conflicts_{label}_{:x}", - crate::timestamp_id() - )); - fs::create_dir_all(&d).unwrap(); - d - } - - fn make_engine(root: PathBuf) -> SyncEngine { - let ex = Arc::new(Exclusions::compile(&ExclusionConfig::default())); - SyncEngine::new(root, "test-node".to_string(), ex) - } - - fn file_bundle(rel: &str, content: &[u8]) -> FileBundle { - let hash: [u8; 32] = blake3::hash(content).into(); - FileBundle { - files: vec![FileData { - metadata: FileMetadata { - change_sequence: 0, - rel_path: PathBuf::from(rel), - size: content.len() as u64, - hash, - modified_ms: 1_000, - is_dir: false, - }, - content: content.to_vec(), - }], - bundle_id: 1, - } - } - - /// Reproduces the exact root cause of the "GUI can't show conflicts" bug: - /// a `ConflictInfo` produced by `apply_bundle` must end up visible in - /// `SyncSnapshot.conflicts` once routed through `push_gui_conflicts`. - #[test] - fn conflict_from_apply_bundle_is_visible_in_gui_snapshot() { - let dir = tmp_dir("visible"); - let engine = make_engine(dir.clone()); - - // Establish the ancestor state, then simulate an offline local edit - // (bypassing apply_bundle, exactly like a real filesystem edit while - // disconnected), then apply a genuinely different incoming version. - engine - .apply_bundle(&file_bundle("shared.txt", b"version A")) - .unwrap(); - fs::write(dir.join("shared.txt"), b"version B (local)").unwrap(); - let result = engine - .apply_bundle(&file_bundle("shared.txt", b"version C (remote)")) - .unwrap(); - assert_eq!(result.conflicts.len(), 1, "expected exactly one conflict"); - - let gui_state = new_shared_state(); - push_gui_conflicts(&Some(gui_state.clone()), &engine, &result.conflicts); - - let snap = gui_state.read().clone(); - assert_eq!(snap.conflicts.len(), 1); - assert_eq!(snap.conflicts[0].filename, "shared.txt"); - - fs::remove_dir_all(&dir).ok(); - } - - #[test] - fn no_gui_state_does_not_panic() { - let dir = tmp_dir("no_gui"); - let engine = make_engine(dir.clone()); - engine.apply_bundle(&file_bundle("f.txt", b"A")).unwrap(); - fs::write(dir.join("f.txt"), b"B").unwrap(); - let result = engine.apply_bundle(&file_bundle("f.txt", b"C")).unwrap(); - // Must be a no-op, not a panic, when gui_state is None. - push_gui_conflicts(&None, &engine, &result.conflicts); - fs::remove_dir_all(&dir).ok(); - } - - #[test] - fn empty_conflicts_does_not_touch_gui_state() { - let dir = tmp_dir("empty"); - let engine = make_engine(dir.clone()); - let gui_state = new_shared_state(); - push_gui_conflicts(&Some(gui_state.clone()), &engine, &[]); - assert!(gui_state.read().conflicts.is_empty()); - fs::remove_dir_all(&dir).ok(); - } -} diff --git a/crates/filesync/src/known_hosts.rs b/crates/filesync/src/known_hosts.rs index 420f7f8..c4677f5 100644 --- a/crates/filesync/src/known_hosts.rs +++ b/crates/filesync/src/known_hosts.rs @@ -225,7 +225,6 @@ impl KnownClients { } fn splice_known_clients(original: &str, clients: &[KnownClient]) -> String { - // Build the replacement section text up front. let new_section = if clients.is_empty() { String::new() } else { @@ -235,7 +234,6 @@ fn splice_known_clients(original: &str, clients: &[KnownClient]) -> String { Ok(s) => s, Err(e) => { log::error!("filesync: failed to serialize known_clients: {e}"); - // Return original unchanged rather than corrupting the file. return original.to_string(); } } @@ -245,9 +243,6 @@ fn splice_known_clients(original: &str, clients: &[KnownClient]) -> String { let mut after: Vec<&str> = Vec::new(); let mut found_first = false; let mut in_section = false; - - // Blank / comment lines that trail an auth block are "pending": we keep - // them only if the next non-blank line belongs to a non-auth section. let mut pending: Vec<&str> = Vec::new(); for line in original.lines() { @@ -262,8 +257,6 @@ fn splice_known_clients(original: &str, clients: &[KnownClient]) -> String { let inner = &trimmed[2..trimmed.len() - 2]; let name = inner.trim().to_ascii_lowercase(); if name == "filesync_known_clients" { - // Entering our managed section — discard any trailing - // whitespace that followed the previous entry. pending.clear(); in_section = true; found_first = true; @@ -286,8 +279,6 @@ fn splice_known_clients(original: &str, clients: &[KnownClient]) -> String { } if in_section { - // Inside a managed block: buffer blank/comment lines; drop value - // lines (they belong to the old entry we're replacing). if trimmed.is_empty() || trimmed.starts_with('#') { pending.push(line); } else { @@ -430,278 +421,3 @@ impl KnownServers { } } } - -#[cfg(test)] -mod tests { - use super::*; - use std::sync::atomic::{AtomicU64, Ordering}; - - static COUNTER: AtomicU64 = AtomicU64::new(0); - - /// Return a unique path in the system temp dir for each test invocation. - fn tmp_path(name: &str) -> std::path::PathBuf { - let n = COUNTER.fetch_add(1, Ordering::Relaxed); - std::env::temp_dir().join(format!("bh_kh_unit_{n}_{name}")) - } - - // ── splice_known_clients unit tests ────────────────────────────────────── - - #[test] - fn splice_into_empty_string_produces_section() { - let out = splice_known_clients( - "", - &[KnownClient { - node_id: "n1".into(), - fingerprint: "fp1".into(), - label: String::new(), - status: ClientStatus::Pending, - addr: "1.2.3.4:1".into(), - first_seen_ms: 1, - last_seen_ms: 1, - }], - ); - assert!(out.contains("[[filesync_known_clients]]")); - assert!(out.contains("fp1")); - } - - #[test] - fn splice_empty_clients_removes_section() { - let original = "[framework]\nname = \"test\"\n\n[[filesync_known_clients]]\nfingerprint = \"fp1\"\nnode_id = \"n1\"\nstatus = \"pending\"\nfirst_seen_ms = 1\nlast_seen_ms = 1\n"; - let out = splice_known_clients(original, &[]); - assert!(!out.contains("filesync_known_clients")); - assert!(out.contains("[framework]")); - } - - #[test] - fn splice_preserves_surrounding_config() { - let original = "[framework]\nname = \"test\"\n\n[[filesync_known_clients]]\nfingerprint = \"old\"\nnode_id = \"n\"\nstatus = \"pending\"\nfirst_seen_ms = 1\nlast_seen_ms = 1\n\n[apps.filesync]\nroot = \"/data\"\n"; - let replacement = vec![KnownClient { - node_id: "n".into(), - fingerprint: "new".into(), - label: String::new(), - status: ClientStatus::Allowed, - addr: String::new(), - first_seen_ms: 1, - last_seen_ms: 2, - }]; - let out = splice_known_clients(original, &replacement); - assert!(out.contains("[framework]")); - assert!(out.contains("[apps.filesync]")); - assert!(out.contains("new")); - assert!(!out.contains("\"old\"")); - } - - #[test] - fn splice_roundtrip_multiple_entries() { - let clients = vec![ - KnownClient { - node_id: "a".into(), - fingerprint: "fp-a".into(), - label: String::new(), - status: ClientStatus::Allowed, - addr: "1.1.1.1:1".into(), - first_seen_ms: 10, - last_seen_ms: 20, - }, - KnownClient { - node_id: "b".into(), - fingerprint: "fp-b".into(), - label: "my label".into(), - status: ClientStatus::Rejected, - addr: "2.2.2.2:2".into(), - first_seen_ms: 30, - last_seen_ms: 40, - }, - ]; - let spliced = splice_known_clients("", &clients); - let parsed: KnownClientsRaw = toml::from_str(&spliced).expect("valid toml"); - assert_eq!(parsed.filesync_known_clients.len(), 2); - assert_eq!(parsed.filesync_known_clients[0].fingerprint, "fp-a"); - assert_eq!(parsed.filesync_known_clients[1].fingerprint, "fp-b"); - assert_eq!(parsed.filesync_known_clients[1].label, "my label"); - } - - // ── KnownClients integration tests ─────────────────────────────────────── - - #[test] - fn new_client_upsert_returns_true() { - let p = tmp_path("kc1.toml"); - let mut kc = KnownClients::load_from_config(&p); - assert!(kc.upsert_pending("node-1", "fp-aaa", "127.0.0.1:1234")); - let _ = std::fs::remove_file(&p); - } - - #[test] - fn repeat_upsert_returns_false() { - let p = tmp_path("kc2.toml"); - let mut kc = KnownClients::load_from_config(&p); - kc.upsert_pending("node-1", "fp-bbb", "127.0.0.1:1"); - assert!(!kc.upsert_pending("node-1", "fp-bbb", "127.0.0.1:2")); - let _ = std::fs::remove_file(&p); - } - - #[test] - fn new_client_has_pending_status() { - let p = tmp_path("kc3.toml"); - let mut kc = KnownClients::load_from_config(&p); - kc.upsert_pending("node-1", "fp-ccc", "127.0.0.1:1"); - assert_eq!(kc.status("fp-ccc"), Some(ClientStatus::Pending)); - let _ = std::fs::remove_file(&p); - } - - #[test] - fn approve_changes_status_to_allowed() { - let p = tmp_path("kc4.toml"); - let mut kc = KnownClients::load_from_config(&p); - kc.upsert_pending("node-1", "fp-ddd", "127.0.0.1:1"); - assert!(kc.set_status("fp-ddd", ClientStatus::Allowed)); - assert_eq!(kc.status("fp-ddd"), Some(ClientStatus::Allowed)); - let _ = std::fs::remove_file(&p); - } - - #[test] - fn set_status_unknown_fingerprint_returns_false() { - let p = tmp_path("kc5.toml"); - let mut kc = KnownClients::load_from_config(&p); - assert!(!kc.set_status("no-such-fp", ClientStatus::Allowed)); - let _ = std::fs::remove_file(&p); - } - - #[test] - fn remove_deletes_entry() { - let p = tmp_path("kc6.toml"); - let mut kc = KnownClients::load_from_config(&p); - kc.upsert_pending("node-1", "fp-eee", "127.0.0.1:1"); - assert!(kc.remove("fp-eee")); - assert_eq!(kc.status("fp-eee"), None); - assert!(!kc.remove("fp-eee")); // second remove returns false - let _ = std::fs::remove_file(&p); - } - - #[test] - fn pending_count_is_accurate() { - let p = tmp_path("kc7.toml"); - let mut kc = KnownClients::load_from_config(&p); - kc.upsert_pending("n1", "fp-f1", ""); - kc.upsert_pending("n2", "fp-f2", ""); - kc.upsert_pending("n3", "fp-f3", ""); - assert_eq!(kc.pending_count(), 3); - kc.set_status("fp-f1", ClientStatus::Allowed); - assert_eq!(kc.pending_count(), 2); - kc.set_status("fp-f2", ClientStatus::Rejected); - assert_eq!(kc.pending_count(), 1); - let _ = std::fs::remove_file(&p); - } - - #[test] - fn persists_and_reloads() { - let p = tmp_path("kc8.toml"); - { - let mut kc = KnownClients::load_from_config(&p); - kc.upsert_pending("node-1", "fp-ppp", "10.0.0.1:7878"); - kc.set_status("fp-ppp", ClientStatus::Allowed); - } - // reload from disk - let kc2 = KnownClients::load_from_config(&p); - assert_eq!(kc2.status("fp-ppp"), Some(ClientStatus::Allowed)); - assert_eq!(kc2.list()[0].addr, "10.0.0.1:7878"); - let _ = std::fs::remove_file(&p); - } - - #[test] - fn load_nonexistent_file_gives_empty_store() { - let p = tmp_path("kc9_nonexistent.toml"); - let kc = KnownClients::load_from_config(&p); - assert_eq!(kc.list().len(), 0); - } - - #[test] - fn persists_alongside_existing_config_content() { - let p = tmp_path("kc10.toml"); - // Seed the "config file" with some pre-existing content. - std::fs::write( - &p, - "[framework]\nname = \"test\"\n\n[apps.filesync]\nroot = \"/tmp\"\n", - ) - .unwrap(); - - let mut kc = KnownClients::load_from_config(&p); - kc.upsert_pending("node-x", "fp-xyz", "9.9.9.9:1"); - - // Reload and verify the client was saved. - let kc2 = KnownClients::load_from_config(&p); - assert_eq!(kc2.status("fp-xyz"), Some(ClientStatus::Pending)); - - // The rest of the config must still be intact. - let raw = std::fs::read_to_string(&p).unwrap(); - assert!(raw.contains("[framework]")); - assert!(raw.contains("[apps.filesync]")); - assert!(raw.contains("[[filesync_known_clients]]")); - - let _ = std::fs::remove_file(&p); - } - - #[test] - fn remove_last_entry_strips_section_from_config() { - let p = tmp_path("kc11.toml"); - std::fs::write(&p, "[framework]\nname = \"test\"\n").unwrap(); - - let mut kc = KnownClients::load_from_config(&p); - kc.upsert_pending("n", "fp-del", ""); - kc.remove("fp-del"); - - let raw = std::fs::read_to_string(&p).unwrap(); - assert!(!raw.contains("filesync_known_clients")); - assert!(raw.contains("[framework]")); - - let _ = std::fs::remove_file(&p); - } - - // ── KnownServers tests (unchanged) ──────────────────────────────────────── - - #[test] - fn first_pin_is_tofu() { - let p = tmp_path("ks1.toml"); - let mut ks = KnownServers::load_or_create(&p); - assert!(ks.get_fingerprint("server:7878").is_none()); - ks.pin("server:7878", "abcdef"); - assert_eq!(ks.get_fingerprint("server:7878"), Some("abcdef")); - let _ = std::fs::remove_file(&p); - } - - #[test] - fn re_pin_updates_fingerprint() { - let p = tmp_path("ks2.toml"); - let mut ks = KnownServers::load_or_create(&p); - ks.pin("server:7878", "old-fp"); - ks.pin("server:7878", "new-fp"); - assert_eq!(ks.get_fingerprint("server:7878"), Some("new-fp")); - let _ = std::fs::remove_file(&p); - } - - #[test] - fn remove_server_clears_pin() { - let p = tmp_path("ks3.toml"); - let mut ks = KnownServers::load_or_create(&p); - ks.pin("srv:1", "fp-x"); - assert!(ks.remove("srv:1")); - assert!(ks.get_fingerprint("srv:1").is_none()); - assert!(!ks.remove("srv:1")); - let _ = std::fs::remove_file(&p); - } - - #[test] - fn server_persists_and_reloads() { - let p = tmp_path("ks4.toml"); - { - let mut ks = KnownServers::load_or_create(&p); - ks.pin("192.168.1.10:7878", "server-fingerprint-hex"); - } - let ks2 = KnownServers::load_or_create(&p); - assert_eq!( - ks2.get_fingerprint("192.168.1.10:7878"), - Some("server-fingerprint-hex") - ); - let _ = std::fs::remove_file(&p); - } -} diff --git a/crates/filesync/src/protocol.rs b/crates/filesync/src/protocol.rs index 4434b1f..6a4ccb8 100644 --- a/crates/filesync/src/protocol.rs +++ b/crates/filesync/src/protocol.rs @@ -21,7 +21,7 @@ pub const BH_DIR: &str = ".bh_filesync"; pub const TMP_DIR: &str = ".bh_filesync/transfers"; pub const TRASH_DIR: &str = ".bh_filesync/trash"; pub const TRASH_INDEX_FILE: &str = "index.json"; -pub const FULL_SCAN_INTERVAL_SECS: u64 = 900; // 15 minutes +pub const FULL_SCAN_INTERVAL_SECS: u64 = 3600; // 1 hour #[derive(Debug, Clone, Serialize, Deserialize)] pub struct FileMetadata { diff --git a/crates/filesync/src/server.rs b/crates/filesync/src/server.rs index c9269d7..e6c6abd 100644 --- a/crates/filesync/src/server.rs +++ b/crates/filesync/src/server.rs @@ -35,18 +35,15 @@ pub struct Server { peers: Arc>>, active_conns: Arc>>>, tls_config: Arc, - - // For tracking sent bundles and sending acknowledgments sent_bundles: Arc>>, } -/// Tracks information about sent bundles for acknowledgment purposes #[derive(Debug, Clone)] struct BundleTracking { bundle_id: u64, sequence_numbers: Vec, sent_time: Instant, - acknowledged_by: HashSet, // client IDs that have acknowledged + acknowledged_by: HashSet, } impl BundleTracking { @@ -209,15 +206,9 @@ fn handle_client( let conn = Arc::new(Connection::new_server(stream, tls_config)?); debug!("filesync: TLS handshake complete with new client"); - // ── Known-host check ──────────────────────────────────────────────────── - // The client's certificate was exchanged and key-ownership was proven - // during the mutual-TLS handshake. We now derive its fingerprint and - // look it up in the known_clients table. let client_fp = match &conn.peer_cert { Some(der) => cert_fingerprint(der), None => { - // Should never happen with client_auth_mandatory = true, but - // guard defensively. warn!("filesync: client at {addr} did not present a certificate — rejecting"); let _ = conn.send(&Message::Rejected { reason: "no client certificate presented".into(), @@ -253,7 +244,6 @@ fn handle_client( } debug!("filesync: received Hello from {node_id} proto={protocol_version}"); - // ── Authorization check ────────────────────────────────────── let status = { let mut kc = known_clients.lock(); let is_new = kc.upsert_pending(&node_id, &client_fp, &addr); @@ -331,10 +321,6 @@ fn handle_client( })?; debug!("filesync: sent Hello response to {client_id}"); - // Use the cached manifest from the initial (or most recent periodic) scan - // when available, rather than blocking on a full re-walk + BLAKE3 hash of - // every file. A fresh scan of a large tree (tens of thousands of dirs) - // can take minutes and starves the client waiting for ManifestExchange. let cached = engine.get_manifest(); let local = if cached.files.is_empty() { debug!("filesync: no cached manifest yet, running full scan for {client_id}"); @@ -388,10 +374,6 @@ fn handle_client( r_files, r_dirs ); - // ── Preemptive disk-space check (server) ────────────────────────────── - // Simulate the client's send-list computation (is_server = false) to - // predict the bytes the client will upload. Abort before sending - // anything if the server filesystem cannot accommodate them. let bytes_from_client: u64 = manifest::compute_send_list(&remote, &local, false) .iter() .filter_map(|p| remote.files.get(p)) @@ -1108,7 +1090,6 @@ fn broadcast_paths( } } - // Track this bundle for acknowledgment purposes let sequence_numbers: Vec = bundle .files .iter() diff --git a/crates/filesync/src/suspend_detector.rs b/crates/filesync/src/suspend_detector.rs index 1edaaf1..a6901db 100644 --- a/crates/filesync/src/suspend_detector.rs +++ b/crates/filesync/src/suspend_detector.rs @@ -1,8 +1,8 @@ use std::fs; use std::time::{Duration, Instant}; -/// Detects system suspend/resume by comparing wall-clock time against system uptime. -/// When wall-clock advances significantly more than uptime, the system was suspended. +const SUSPEND_THRESHOLD_SECS: u64 = 5; + pub struct SuspendDetector { last_check: Instant, last_uptime: Duration, @@ -30,8 +30,6 @@ impl SuspendDetector { }; let uptime_elapsed = current_uptime.saturating_sub(self.last_uptime); - - const SUSPEND_THRESHOLD_SECS: u64 = 5; let suspended = wall_elapsed.as_secs() > uptime_elapsed.as_secs() + SUSPEND_THRESHOLD_SECS; self.last_check = now; diff --git a/crates/filesync/src/sync_engine.rs b/crates/filesync/src/sync_engine.rs index 8391736..bba3994 100644 --- a/crates/filesync/src/sync_engine.rs +++ b/crates/filesync/src/sync_engine.rs @@ -17,7 +17,6 @@ use std::time::{Duration, Instant, SystemTime}; const PIPELINE_DEPTH: usize = 128; -/// Tracks recent changes to a file for coalescing rapid modifications #[derive(Debug, Clone)] struct FileChangeHistory { last_change_time: Instant, @@ -42,7 +41,6 @@ impl FileChangeHistory { self.change_count += 1; self.last_sequence_number = sequence_number; - // Return true if this change is part of a rapid sequence elapsed < FILE_CHANGE_COALESCE_MS } @@ -70,17 +68,13 @@ pub enum FinishResult { #[derive(Debug, Clone)] pub struct ConflictInfo { - /// Original file path (incoming file applied here). pub original_path: PathBuf, - /// Path where the diverged local copy was saved. pub conflict_copy_path: PathBuf, } #[derive(Debug, Default)] pub struct ApplyResult { - /// Number of entries actually written (files + dirs). pub written: usize, - /// Conflict copies created during this apply. pub conflicts: Vec, } @@ -92,9 +86,6 @@ struct LargeFileAssembly { expected_hash: [u8; 32], dst: PathBuf, file_size: u64, - /// Original modification timestamp (ms since Unix epoch) from the sender. - /// Restored onto the destination file after the transfer is committed so - /// that the receiver's filesystem mtime matches the sender's. modified_ms: u64, final_hash_pending: Option<[u8; 32]>, @@ -102,11 +93,7 @@ struct LargeFileAssembly { #[derive(Debug, Default, Clone)] pub struct SyncEngineConfig { - /// Automatically purge trash entries older than this many days. - /// `None` or `0` disables automatic purging. pub trash_expiry_days: Option, - /// Override for the periodic full-rescan interval in seconds. - /// Defaults to `FULL_SCAN_INTERVAL_SECS` (900) when `None`. pub full_scan_interval_secs: Option, } @@ -114,18 +101,12 @@ pub struct SyncEngine { root: PathBuf, node_id: String, manifest: RwLock, - suppressed: Arc>>, - suppressed_deletes: Arc>>, - in_progress: RwLock>, - exclusions: Arc, trash_manager: Arc, full_scan_interval_secs: u64, - - // For tracking rapid changes and coalescing change_sequences: RwLock>, last_sequence_number: RwLock, } @@ -146,8 +127,6 @@ pub fn conflict_copy_name(rel_path: &Path, node_id: &str, unix_secs: u64) -> Pat } } -/// Compute the BLAKE3 hash of a file on disk. Returns `None` if the file -/// cannot be read (e.g. does not exist yet). fn hash_file(path: &Path) -> Option<[u8; 32]> { use std::io::Read; let file = fs::File::open(path).ok()?; @@ -164,10 +143,6 @@ fn hash_file(path: &Path) -> Option<[u8; 32]> { Some(hasher.finalize().into()) } -/// Remove empty intermediate directories left behind inside the transfers -/// folder after a large-file tmp file has been consumed or deleted. -/// Walks upward from `tmp_path`'s parent, removing directories that are -/// empty, stopping when it reaches the transfers root or a non-empty dir. fn cleanup_transfer_dirs(root: &Path, tmp_path: &Path) { let transfers_dir = root.join(crate::protocol::TMP_DIR); let mut current = tmp_path.parent(); @@ -176,7 +151,6 @@ fn cleanup_transfer_dirs(root: &Path, tmp_path: &Path) { break; } if fs::remove_dir(dir).is_err() { - // Non-empty or permission error — stop climbing break; } current = dir.parent(); @@ -280,14 +254,9 @@ impl SyncEngine { let age_ms = now_ms.saturating_sub(modified_ms); - // Improved stability detection: consider file stable if: - // 1. It hasn't been modified for FILE_STABILITY_MS, OR - // 2. It's been modified recently but we've seen multiple rapid changes (coalescing) if age_ms >= FILE_STABILITY_MS { return true; } - - // Check if this file has recent rapid changes that should be coalesced self.has_recent_rapid_changes(rel, modified_ms) } @@ -298,7 +267,6 @@ impl SyncEngine { pub fn has_recent_rapid_changes(&self, rel: &Path, _current_modified_ms: u64) -> bool { let history = self.change_sequences.read(); if let Some(file_history) = history.get(rel) { - // If we have recent rapid changes, consider the file stable for coalescing file_history.should_coalesce() } else { false @@ -328,7 +296,6 @@ impl SyncEngine { histories.remove(rel); } - /// Create a new SyncEngine with the same configuration but fresh state pub fn clone_with_fresh_state(&self) -> Self { Self { root: self.root.clone(), @@ -431,25 +398,19 @@ impl SyncEngine { .collect() } - /// Returns `Some(manifest_hash)` when both the on-disk file and the incoming - /// file have diverged from the recorded manifest hash (i.e. a genuine - /// two-sided conflict). Returns `None` when there is no conflict. fn detect_conflict( &self, rel_path: &Path, full_path: &Path, incoming_hash: &[u8; 32], ) -> Option<[u8; 32]> { - // Must be a known file (present in our manifest) let manifest_hash = self.manifest.read().files.get(rel_path).map(|m| m.hash)?; - // The incoming file must differ from the last-synced state if incoming_hash == &manifest_hash { return None; } - // The on-disk file must also differ from the last-synced state let on_disk_hash = hash_file(full_path)?; if on_disk_hash == manifest_hash { - return None; // local never changed — no conflict + return None; } if &on_disk_hash == incoming_hash { return None; @@ -460,9 +421,6 @@ impl SyncEngine { pub fn apply_bundle(&self, bundle: &FileBundle) -> std::io::Result { let mut written = Vec::new(); let mut result = ApplyResult::default(); - // Collect (absolute_path, modified_ms) for directories so we can stamp - // them *after* all their contents have been written. Writing files into - // a directory updates the directory's mtime, so we must restore it last. let mut dir_mtimes: Vec<(PathBuf, u64)> = Vec::new(); for fd in &bundle.files { @@ -473,7 +431,6 @@ impl SyncEngine { let full = self.root.join(&fd.metadata.rel_path); - // --- Conflict detection (files only) --- if !fd.metadata.is_dir { if let Some(_ancestor_hash) = self.detect_conflict(&fd.metadata.rel_path, &full, &fd.metadata.hash) @@ -507,21 +464,17 @@ impl SyncEngine { } } - // --- Apply the incoming entry --- self.suppressed.write().insert(fd.metadata.rel_path.clone()); written.push(fd.metadata.rel_path.clone()); if fd.metadata.is_dir { fs::create_dir_all(&full)?; - // Defer mtime restoration until after all files in this bundle - // have been written, otherwise child-file writes will clobber it. dir_mtimes.push((full.clone(), fd.metadata.modified_ms)); } else { if let Some(parent) = full.parent() { fs::create_dir_all(parent)?; } fs::write(&full, &fd.content)?; - // Restore the sender's modification timestamp on the file. let mtime = filetime::FileTime::from_unix_time( (fd.metadata.modified_ms / 1000) as i64, ((fd.metadata.modified_ms % 1000) * 1_000_000) as u32, @@ -538,9 +491,6 @@ impl SyncEngine { result.written += 1; } - // Restore directory mtimes deepest-first so that stamping a child - // directory does not update its parent's mtime before we stamp the - // parent too. dir_mtimes.sort_by(|a, b| b.0.components().count().cmp(&a.0.components().count())); for (path, modified_ms) in dir_mtimes { let mtime = filetime::FileTime::from_unix_time( @@ -711,7 +661,6 @@ impl SyncEngine { } } - // --- Conflict detection --- let conflict = { let manifest_hash = self.manifest.read().files.get(path).map(|m| m.hash); if let Some(manifest_hash) = manifest_hash { @@ -762,9 +711,6 @@ impl SyncEngine { fs::rename(&asm.tmp_path, &asm.dst)?; cleanup_transfer_dirs(&self.root, &asm.tmp_path); - // Restore the sender's original modification timestamp. The rename - // above (and the OS itself) would otherwise stamp the file with the - // current time. let mtime = filetime::FileTime::from_unix_time( (asm.modified_ms / 1000) as i64, ((asm.modified_ms % 1000) * 1_000_000) as u32, diff --git a/crates/filesync/tests/test_client.rs b/crates/filesync/tests/test_client.rs index af8c0c7..5914eb6 100644 --- a/crates/filesync/tests/test_client.rs +++ b/crates/filesync/tests/test_client.rs @@ -1,9 +1,15 @@ use std::collections::HashMap; +use std::fs; use std::path::PathBuf; +use std::sync::Arc; use bytehive_filesync::{ - client::count_manifest, - protocol::{FileMetadata, Manifest}, + client::{count_manifest, format_unix_secs}, + exclusions::{ExclusionConfig, Exclusions}, + gui::state::new_shared_state, + protocol::{FileBundle, FileData, FileMetadata, Manifest}, + sync_engine::SyncEngine, + timestamp_id, }; fn make_manifest(entries: &[(&str, u64, bool)]) -> Manifest { @@ -128,3 +134,13 @@ fn count_many_files_and_dirs() { assert_eq!(d, 5); assert_eq!(b, 500); } + +#[test] +fn epoch_formats_correctly() { + assert_eq!(format_unix_secs(0), "1970-01-01 00:00"); +} + +#[test] +fn known_date_formats_correctly() { + assert_eq!(format_unix_secs(1_717_236_000), "2024-06-01 10:00"); +} diff --git a/crates/filesync/tests/test_known_hosts.rs b/crates/filesync/tests/test_known_hosts.rs index f6dfddd..2ea2508 100644 --- a/crates/filesync/tests/test_known_hosts.rs +++ b/crates/filesync/tests/test_known_hosts.rs @@ -13,8 +13,6 @@ mod tests { std::env::temp_dir().join(format!("bh_kh_test_{n}_{name}")) } - // ── KnownClients ───────────────────────────────────────────────────────── - #[test] fn new_client_upsert_returns_true() { let p = tmp_path("kc1.toml"); @@ -54,6 +52,7 @@ mod tests { #[test] fn set_status_unknown_fingerprint_returns_false() { let p = tmp_path("kc5.toml"); + let mut kc = KnownClients::load_from_config(&p); assert!(!kc.set_status("no-such-fp", ClientStatus::Allowed)); let _ = std::fs::remove_file(&p); @@ -66,7 +65,7 @@ mod tests { kc.upsert_pending("node-1", "fp-eee", "127.0.0.1:1"); assert!(kc.remove("fp-eee")); assert_eq!(kc.status("fp-eee"), None); - assert!(!kc.remove("fp-eee")); // second remove returns false + assert!(!kc.remove("fp-eee")); let _ = std::fs::remove_file(&p); } @@ -107,8 +106,6 @@ mod tests { assert_eq!(kc.list().len(), 0); } - // ── KnownServers ───────────────────────────────────────────────────────── - #[test] fn first_pin_is_tofu() { let p = tmp_path("ks1.toml"); From eb04bc62f6265506ffd126aafcae24d21defeb05 Mon Sep 17 00:00:00 2001 From: Guy Dols Date: Thu, 10 Sep 2026 22:05:22 +0200 Subject: [PATCH 2/6] feat(filesync): add persistent deletion ledger with 90d TTL Offline clients re-uploaded files deleted while they were away because deletes only lived in memory. The ledger in .bh_filesync/deletion-ledger.json keeps tombstones across restarts so deletes survive reconnects. --- crates/filesync/src/ledger.rs | 209 ++++++++++++++++++++++++++++++++++ crates/filesync/src/lib.rs | 2 + 2 files changed, 211 insertions(+) create mode 100644 crates/filesync/src/ledger.rs diff --git a/crates/filesync/src/ledger.rs b/crates/filesync/src/ledger.rs new file mode 100644 index 0000000..23171a8 --- /dev/null +++ b/crates/filesync/src/ledger.rs @@ -0,0 +1,209 @@ +//! Persistent deletion ledger (tombstones) for filesync. +//! +//! Foundation for resurrection-safe deletes (TODO 1): +//! - Tombstones persist in `/.bh_filesync/deletion-ledger.json` +//! - Last-writer-wins per path on `deleted_at_ms` +//! - 90-day TTL via [`DeletionLedger::prune`] +//! - Never panics on corrupt / missing files: corrupt JSON is backed up +//! to `deletion-ledger.corrupt-.json` and an empty ledger is returned. + +use serde::{Deserialize, Serialize}; +use std::collections::HashMap; +use std::io; +use std::path::{Path, PathBuf}; +use std::time::{SystemTime, UNIX_EPOCH}; + +/// Directory (under the sync root) holding filesync state. +pub const LEDGER_DIR_NAME: &str = ".bh_filesync"; +/// File name of the persistent deletion ledger. +pub const LEDGER_FILE_NAME: &str = "deletion-ledger.json"; +/// Default time-to-live for tombstones: 90 days in milliseconds. +pub const DELETION_LEDGER_TTL_MS: u64 = 90 * 24 * 60 * 60 * 1000; + +/// A single deletion tombstone. +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct Tombstone { + pub path: PathBuf, + pub deleted_at_ms: u64, + pub deleter_node: String, + pub prev_hash: Option, + pub prev_mtime_ms: Option, +} + +impl Tombstone { + pub fn new( + path: PathBuf, + deleted_at_ms: u64, + deleter_node: String, + prev_hash: Option, + prev_mtime_ms: Option, + ) -> Self { + Self { + path, + deleted_at_ms, + deleter_node, + prev_hash, + prev_mtime_ms, + } + } +} + +/// Persistent set of deletion tombstones, keyed by relative path. +#[derive(Debug, Clone, Default, Serialize, Deserialize)] +pub struct DeletionLedger { + #[serde(default)] + pub entries: HashMap, +} + +fn ledger_path(root: &Path) -> PathBuf { + root.join(LEDGER_DIR_NAME).join(LEDGER_FILE_NAME) +} + +/// Wall-clock milliseconds since the Unix epoch. Never panics: on clock +/// error (time before epoch) returns 0. +pub fn now_ms() -> u64 { + SystemTime::now() + .duration_since(UNIX_EPOCH) + .map(|d| d.as_millis() as u64) + .unwrap_or(0) +} + +impl DeletionLedger { + pub fn new() -> Self { + Self { + entries: HashMap::new(), + } + } + + /// Load the ledger from `/.bh_filesync/deletion-ledger.json`. + /// Missing file => empty ledger. Corrupt JSON => back the file up to + /// `deletion-ledger.corrupt-.json` (best effort) and return empty. + /// Never panics. + pub fn load(root: &Path) -> Self { + let path = ledger_path(root); + let raw = match std::fs::read_to_string(&path) { + Ok(s) => s, + Err(e) if e.kind() == io::ErrorKind::NotFound => return Self::new(), + Err(_) => return Self::new(), + }; + // Primary format: {"entries": {...}}. Fall back to a bare map for + // forward-compat, then give up (backup + empty). + if let Ok(ledger) = serde_json::from_str::(&raw) { + return ledger; + } + if let Ok(map) = serde_json::from_str::>(&raw) { + return Self { entries: map }; + } + // Corrupt: back up best-effort, then start empty. + let backup = path.with_extension(format!("corrupt-{}.json", now_ms())); + let _ = std::fs::copy(&path, &backup); + Self::new() + } + + /// Persist atomically: write `.tmp` then rename. Creates + /// `/.bh_filesync` if needed. + pub fn save_atomic(&self, root: &Path) -> io::Result<()> { + let dir = root.join(LEDGER_DIR_NAME); + std::fs::create_dir_all(&dir)?; + let dest = dir.join(LEDGER_FILE_NAME); + let tmp = dest.with_extension("json.tmp"); + let raw = serde_json::to_string_pretty(self) + .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?; + std::fs::write(&tmp, raw)?; + std::fs::rename(&tmp, &dest)?; + Ok(()) + } + + /// Record a tombstone. Last-writer-wins: keeps whichever entry has the + /// newest `deleted_at_ms`. + pub fn record( + &mut self, + path: PathBuf, + deleted_at_ms: u64, + deleter_node: String, + prev_hash: Option, + prev_mtime_ms: Option, + ) { + let entry = Tombstone::new( + path.clone(), + deleted_at_ms, + deleter_node, + prev_hash, + prev_mtime_ms, + ); + match self.entries.get(&path) { + Some(existing) if existing.deleted_at_ms >= deleted_at_ms => {} + _ => { + self.entries.insert(path, entry); + } + } + } + + /// Returns true if a tombstone for `path` exists and is strictly newer + /// than `mtime_ms` (i.e. the delete wins over that file version). + pub fn has_newer_than(&self, path: &Path, mtime_ms: u64) -> bool { + match self.entries.get(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 { + match self.entries.get(&path) { + Some(existing) if existing.deleted_at_ms >= tomb.deleted_at_ms => {} + _ => { + self.entries.insert(path, tomb); + } + } + } + } + + /// Drop tombstones older than `ttl_ms` relative to `now_ms`. + /// Entries from the future (`deleted_at_ms > now_ms`) are kept. + /// Returns the number of entries removed. + pub fn prune(&mut self, now_ms: u64, ttl_ms: u64) -> usize { + let before = self.entries.len(); + self.entries.retain(|_, t| { + if t.deleted_at_ms > now_ms { + return true; + } + now_ms.saturating_sub(t.deleted_at_ms) <= ttl_ms + }); + before - self.entries.len() + } + + /// Drop tombstones older than the default 90-day TTL. Returns the + /// number of entries removed. + pub fn prune_expired(&mut self, now_ms: u64) -> usize { + self.prune(now_ms, DELETION_LEDGER_TTL_MS) + } + + /// Remove the tombstone for `path` (e.g. on recreate / upload win). + /// Returns the removed tombstone, if any. + pub fn remove(&mut self, path: &Path) -> Option { + self.entries.remove(path) + } + + /// Borrow the tombstone for `path`, if any. + pub fn get(&self, path: &Path) -> Option<&Tombstone> { + self.entries.get(path) + } + + pub fn contains(&self, path: &Path) -> bool { + self.entries.contains_key(path) + } + + pub fn len(&self) -> usize { + self.entries.len() + } + + pub fn is_empty(&self) -> bool { + self.entries.is_empty() + } + + pub fn iter(&self) -> impl Iterator { + self.entries.iter() + } +} diff --git a/crates/filesync/src/lib.rs b/crates/filesync/src/lib.rs index 253ff4f..46defb8 100644 --- a/crates/filesync/src/lib.rs +++ b/crates/filesync/src/lib.rs @@ -5,6 +5,7 @@ pub mod common; pub mod exclusions; pub mod gui; pub mod known_hosts; +pub mod ledger; pub mod manifest; pub mod protocol; pub mod server; @@ -15,6 +16,7 @@ pub mod watcher; pub use app::FileSyncApp; pub use known_hosts::{ClientStatus, KnownClient, KnownClients, KnownServer, KnownServers}; +pub use ledger::{DeletionLedger, Tombstone, DELETION_LEDGER_TTL_MS}; pub fn timestamp_id() -> u64 { use std::time::{SystemTime, UNIX_EPOCH}; From 60c52b41a0eb2b17990d8ae73549ae1d3d88d242 Mon Sep 17 00:00:00 2001 From: Guy Dols Date: Thu, 10 Sep 2026 22:05:22 +0200 Subject: [PATCH 3/6] feat(filesync): bump protocol to v8 with timestamped deletes Message::Delete now carries deleted_at_ms and deleter plus LedgerExchange so peers can compare delete time versus file mtime. Clean break, no compat shim. --- crates/filesync/src/protocol.rs | 19 ++++++++++++++++++- crates/filesync/tests/test_protocol.rs | 14 +++++++++++++- 2 files changed, 31 insertions(+), 2 deletions(-) diff --git a/crates/filesync/src/protocol.rs b/crates/filesync/src/protocol.rs index 6a4ccb8..6bb2b42 100644 --- a/crates/filesync/src/protocol.rs +++ b/crates/filesync/src/protocol.rs @@ -8,7 +8,7 @@ pub const BUNDLE_MAX_FILES: usize = 500; pub const LARGE_FILE_THRESHOLD: u64 = 8 * 1024 * 1024; pub const FILE_CHUNK_SIZE: usize = 8 * 1024 * 1024; pub const MAX_FRAME_BYTES: usize = 32 * 1024 * 1024; -pub const PROTOCOL_VERSION: u32 = 7; +pub const PROTOCOL_VERSION: u32 = 8; pub const DEBOUNCE_MS: u64 = 200; pub const FILE_STABILITY_MS: u64 = 500; pub const FILE_CHANGE_COALESCE_MS: u64 = 100; @@ -51,6 +51,15 @@ pub struct FileBundle { pub bundle_id: u64, } +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct TombstoneDto { + pub path: PathBuf, + pub deleted_at_ms: u64, + pub deleter_node: String, + pub prev_hash: Option, + pub prev_mtime_ms: Option, +} + #[derive(Debug, Clone, Serialize, Deserialize)] pub enum Message { Hello { @@ -73,6 +82,14 @@ pub enum Message { }, Delete { paths: Vec, + deleted_at_ms: u64, + deleter: String, + }, + LedgerExchange { + entries: Vec, + }, + LedgerAck { + received: usize, }, SyncComplete, LargeFileStart { diff --git a/crates/filesync/tests/test_protocol.rs b/crates/filesync/tests/test_protocol.rs index 7a0ef9f..b47ca05 100644 --- a/crates/filesync/tests/test_protocol.rs +++ b/crates/filesync/tests/test_protocol.rs @@ -79,14 +79,22 @@ fn roundtrip_bundle() { fn roundtrip_delete() { let msg = Message::Delete { paths: vec![PathBuf::from("a.txt"), PathBuf::from("sub/b.txt")], + deleted_at_ms: 1_700_000_000_000, + deleter: "test-node".to_string(), }; let frame = serialise_message(&msg).unwrap(); let mut cur = Cursor::new(frame); match read_message(&mut cur).unwrap() { - Message::Delete { paths } => { + Message::Delete { + paths, + deleted_at_ms, + deleter, + } => { assert_eq!(paths.len(), 2); assert_eq!(paths[0], PathBuf::from("a.txt")); assert_eq!(paths[1], PathBuf::from("sub/b.txt")); + assert_eq!(deleted_at_ms, 1_700_000_000_000); + assert_eq!(deleter, "test-node"); } _ => panic!("expected Delete"), } @@ -331,9 +339,13 @@ fn roundtrip_hello_no_credential() { fn delete_message_length_prefix_is_consistent() { let msg1 = Message::Delete { paths: vec![PathBuf::from("a.txt")], + deleted_at_ms: 1, + deleter: String::new(), }; let msg2 = Message::Delete { paths: vec![PathBuf::from("a.txt"), PathBuf::from("b.txt")], + deleted_at_ms: 1, + deleter: String::new(), }; let frame1 = serialise_message(&msg1).unwrap(); let frame2 = serialise_message(&msg2).unwrap(); From cac11aba0b7881adb53f95d1dcd5588703471e9b Mon Sep 17 00:00:00 2001 From: Guy Dols Date: Thu, 10 Sep 2026 22:05:27 +0200 Subject: [PATCH 4/6] fix(filesync): record tombstones and veto resurrection uploads apply_deletes now persists a tombstone even when the file is already absent, and filter_resurrected blocks re-upload when the remote tombstone is newer than local mtime. Legitimate recreates with newer mtime still win and clear the tombstone. --- crates/filesync/src/common.rs | 63 ++++++++- crates/filesync/src/manifest.rs | 74 ++++++++++ crates/filesync/src/sync_engine.rs | 161 +++++++++++++++++++++- crates/filesync/tests/test_manifest.rs | 94 ++++++++++++- crates/filesync/tests/test_sync_engine.rs | 144 ++++++++++++++++++- 5 files changed, 525 insertions(+), 11 deletions(-) diff --git a/crates/filesync/src/common.rs b/crates/filesync/src/common.rs index 82aa81c..be65e89 100644 --- a/crates/filesync/src/common.rs +++ b/crates/filesync/src/common.rs @@ -176,6 +176,32 @@ impl PendingChanges { expand_deleted_ancestors(root, &mut paths); (paths, original_count) } + + /// 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. + /// (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()); + 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 { + let (prev_hash, prev_mtime_ms) = manifest + .files + .get(p) + .map(|m| (Some(crate::hex(&m.hash)), Some(m.modified_ms))) + .unwrap_or((None, None)); + engine.record_delete(p.clone(), now, &node, prev_hash, prev_mtime_ms); + } + } + (paths, original_count) + } } pub fn expand_deleted_ancestors(root: &Path, paths: &mut Vec) { @@ -421,20 +447,29 @@ pub fn handle_recv_large_file_end( } } -pub fn handle_recv_delete( +/// Apply a remote delete, preserving the SENDER's timestamp/deleter for +/// last-writer-wins (never re-stamp locally). +pub fn handle_recv_delete_with_meta( engine: &SyncEngine, paths: &[PathBuf], + deleted_at_ms: u64, + deleter: &str, peer: &str, bus: &Option>, log_prefix: &str, ) -> io::Result { - let n = engine.apply_deletes(paths)?; + let n = engine.apply_deletes(paths, deleted_at_ms, deleter)?; info!("{log_prefix}: -{n} path(s) from {peer}"); if let Some(ref bus) = bus { bus.publish( "filesync", "filesync.file_deleted", - serde_json::json!({ "paths": paths, "node": peer }), + serde_json::json!({ + "paths": paths, + "node": peer, + "deleted_at_ms": deleted_at_ms, + "deleter": deleter, + }), ); bus.publish( "filesync", @@ -449,6 +484,28 @@ pub fn handle_recv_delete( Ok(n) } +/// Legacy 5-arg entry point kept so the current client/server recv loops +/// compile untouched in this lane. Stamps with local wall-clock time; +/// prefer [`handle_recv_delete_with_meta`] (next lane switches callers to +/// it so sender timestamps survive for LWW). +pub fn handle_recv_delete( + engine: &SyncEngine, + paths: &[PathBuf], + peer: &str, + bus: &Option>, + log_prefix: &str, +) -> io::Result { + handle_recv_delete_with_meta( + engine, + paths, + crate::ledger::now_ms(), + peer, + peer, + bus, + log_prefix, + ) +} + pub fn handle_recv_rename( engine: &SyncEngine, from: &PathBuf, diff --git a/crates/filesync/src/manifest.rs b/crates/filesync/src/manifest.rs index a53d694..7df40d1 100644 --- a/crates/filesync/src/manifest.rs +++ b/crates/filesync/src/manifest.rs @@ -1,4 +1,5 @@ use crate::exclusions::Exclusions; +use crate::ledger::DeletionLedger; use crate::protocol::{FileMetadata, Manifest, HASH_THREADS}; use rayon::prelude::*; use std::io::Read; @@ -120,6 +121,79 @@ pub fn compute_send_list(local: &Manifest, remote: &Manifest, is_server: bool) - out } +/// Veto uploads that would resurrect a peer-deleted file. +/// +/// `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 +/// equals `prev_hash`, or local mtime is not newer than `prev_mtime_ms`; +/// missing prev info fails closed toward the delete). +/// - 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). +/// - `None` ledger (or no tombstone / no local entry) passes through. +/// +/// The existing `is_server` tie-break in [`compute_send_list`] is +/// untouched; this only filters its output. +pub fn filter_resurrected( + send_list: Vec, + local: &Manifest, + ledger: Option<&DeletionLedger>, +) -> (Vec, Vec) { + let Some(ledger) = ledger else { + return (send_list, Vec::new()); + }; + let mut to_send = Vec::with_capacity(send_list.len()); + let mut resurrected = Vec::new(); + for path in send_list { + let Some(local_meta) = local.files.get(&path) else { + // No local entry to evaluate — pass through. + to_send.push(path); + continue; + }; + let Some(tomb) = ledger.get(&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. + 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); + } + } + (to_send, resurrected) +} + pub fn diff_manifests(old: &Manifest, new: &Manifest) -> (Vec, Vec) { let mut changed = Vec::new(); let mut deleted = Vec::new(); diff --git a/crates/filesync/src/sync_engine.rs b/crates/filesync/src/sync_engine.rs index bba3994..fbbca69 100644 --- a/crates/filesync/src/sync_engine.rs +++ b/crates/filesync/src/sync_engine.rs @@ -1,5 +1,6 @@ use crate::bundler; use crate::exclusions::Exclusions; +use crate::ledger::{now_ms, DeletionLedger, Tombstone, DELETION_LEDGER_TTL_MS}; use crate::manifest; use crate::protocol::FILE_STABILITY_MS; use crate::protocol::*; @@ -101,6 +102,7 @@ pub struct SyncEngine { root: PathBuf, node_id: String, manifest: RwLock, + ledger: RwLock, suppressed: Arc>>, suppressed_deletes: Arc>>, in_progress: RwLock>, @@ -190,11 +192,13 @@ impl SyncEngine { } ); let trash_manager = TrashManager::new(root.clone(), config.trash_expiry_days); + let ledger = DeletionLedger::load(&root); Self { manifest: RwLock::new(Manifest { files: HashMap::new(), node_id: node_id.clone(), }), + ledger: RwLock::new(ledger), root, node_id, suppressed: Arc::new(RwLock::new(HashSet::new())), @@ -301,6 +305,7 @@ impl SyncEngine { root: self.root.clone(), node_id: self.node_id.clone(), manifest: RwLock::new(self.manifest.read().clone()), + ledger: RwLock::new(self.ledger.read().clone()), suppressed: self.suppressed.clone(), suppressed_deletes: self.suppressed_deletes.clone(), in_progress: RwLock::new(HashMap::new()), @@ -339,6 +344,8 @@ impl SyncEngine { pub fn scan(&self) -> std::io::Result { let m = manifest::build_manifest(&self.root, &self.node_id, &self.exclusions)?; *self.manifest.write() = m.clone(); + // Cheap upkeep (persisted only when entries actually expire). + self.prune_ledger(); Ok(m) } @@ -346,6 +353,120 @@ impl SyncEngine { self.manifest.read().clone() } + // ---- Deletion ledger (tombstones) ---- + + /// Record a deletion tombstone and persist the ledger (best effort: + /// save errors are logged, not propagated). + pub fn record_delete( + &self, + path: PathBuf, + deleted_at_ms: u64, + deleter: &str, + prev_hash: Option, + prev_mtime_ms: Option, + ) { + self.ledger.write().record( + path, + deleted_at_ms, + deleter.to_string(), + prev_hash, + prev_mtime_ms, + ); + if let Err(e) = self.ledger.read().save_atomic(&self.root) { + log::warn!("ledger: save after record failed: {e}"); + } + } + + /// Merge remote tombstones (LWW per path) and persist. Returns the + /// number of DTOs received. + pub fn merge_ledger(&self, dtos: Vec) -> usize { + let n = dtos.len(); + if n == 0 { + return 0; + } + let remote: HashMap = dtos + .into_iter() + .map(|d| { + let path = d.path.clone(); + ( + path, + Tombstone::new( + d.path, + d.deleted_at_ms, + d.deleter_node, + d.prev_hash, + d.prev_mtime_ms, + ), + ) + }) + .collect(); + self.ledger.write().merge_remote(remote); + if let Err(e) = self.ledger.read().save_atomic(&self.root) { + log::warn!("ledger: save after merge failed: {e}"); + } + n + } + + /// Snapshot local tombstones as DTOs for a LedgerExchange. + pub fn ledger_entries(&self) -> Vec { + self.ledger + .read() + .iter() + .map(|(_, t)| TombstoneDto { + path: t.path.clone(), + deleted_at_ms: t.deleted_at_ms, + deleter_node: t.deleter_node.clone(), + prev_hash: t.prev_hash.clone(), + prev_mtime_ms: t.prev_mtime_ms, + }) + .collect() + } + + /// Returns true if a tombstone for `path` exists and is strictly newer + /// than `mtime_ms`. + pub fn has_tombstone_newer_than(&self, path: &Path, mtime_ms: u64) -> bool { + self.ledger.read().has_newer_than(path, mtime_ms) + } + + /// 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 { + let pruned = self + .ledger + .write() + .prune(now_ms(), DELETION_LEDGER_TTL_MS); + if pruned > 0 { + if let Err(e) = self.ledger.read().save_atomic(&self.root) { + log::warn!("ledger: save after prune failed: {e}"); + } else { + log::info!("ledger: pruned {pruned} expired tombstone(s)"); + } + } + pruned + } + + /// Clear tombstones for legitimately recreated paths (recreate win) + /// and persist once. Used with [`crate::manifest::filter_resurrected`]. + pub fn clear_tombstones(&self, paths: &[PathBuf]) { + if paths.is_empty() { + return; + } + let mut removed_any = false; + { + let mut ledger = self.ledger.write(); + for p in paths { + if ledger.remove(p).is_some() { + removed_any = true; + } + } + } + if removed_any { + if let Err(e) = self.ledger.read().save_atomic(&self.root) { + log::warn!("ledger: save after clear failed: {e}"); + } + } + } + pub fn send_paths(&self, paths: &[PathBuf], conn: &Arc) -> std::io::Result<()> { if paths.is_empty() { return Ok(()); @@ -802,9 +923,16 @@ impl SyncEngine { Ok(()) } - pub fn apply_deletes(&self, paths: &[PathBuf]) -> std::io::Result { + pub fn apply_deletes( + &self, + paths: &[PathBuf], + deleted_at_ms: u64, + deleter: &str, + ) -> std::io::Result { let mut removed = Vec::new(); let mut count = 0usize; + // Stage tombstone data in memory; persist once at the end. + let mut tombstones: Vec<(PathBuf, Option, Option)> = Vec::new(); for rel in paths { if !safe_relative(rel) { @@ -814,6 +942,15 @@ impl SyncEngine { self.suppressed_deletes.write().insert(rel.clone()); removed.push(rel.clone()); + // Snapshot manifest info BEFORE removal for the tombstone. + let (prev_hash, prev_mtime_ms) = self + .manifest + .read() + .files + .get(rel) + .map(|m| (Some(crate::hex(&m.hash)), Some(m.modified_ms))) + .unwrap_or((None, None)); + let full = self.root.join(rel); if full.is_dir() { self.trash_manager.move_to_trash(&full, rel, &self.node_id); @@ -834,9 +971,31 @@ impl SyncEngine { self.trash_manager.move_to_trash(&full, rel, &self.node_id); self.manifest.write().files.remove(rel); } + // Record ALWAYS, even when the file was absent on disk and no + // manifest entry existed (fixes the silent-absent gap: the peer + // may still hold the file and would otherwise resurrect it). + tombstones.push((rel.clone(), prev_hash, prev_mtime_ms)); count += 1; } + if !tombstones.is_empty() { + { + let mut ledger = self.ledger.write(); + for (path, prev_hash, prev_mtime_ms) in tombstones { + ledger.record( + path, + deleted_at_ms, + deleter.to_string(), + prev_hash, + prev_mtime_ms, + ); + } + } + if let Err(e) = self.ledger.read().save_atomic(&self.root) { + log::warn!("ledger: save after apply_deletes failed: {e}"); + } + } + self.schedule_unsuppress_deletes(removed); Ok(count) } diff --git a/crates/filesync/tests/test_manifest.rs b/crates/filesync/tests/test_manifest.rs index 26ae5ee..4e2d0c3 100644 --- a/crates/filesync/tests/test_manifest.rs +++ b/crates/filesync/tests/test_manifest.rs @@ -2,7 +2,9 @@ use std::{collections::HashMap, path::PathBuf}; use bytehive_filesync::{ exclusions::{ExclusionConfig, Exclusions}, - manifest::{build_manifest, compute_send_list}, + hex, + ledger::DeletionLedger, + manifest::{build_manifest, compute_send_list, filter_resurrected}, protocol::{FileMetadata, Manifest}, }; @@ -281,3 +283,93 @@ fn send_list_multiple_files_some_sent_some_not() { "identical hash file must not be sent" ); } + +// ---- Resurrection-filter tests, relocated from src/manifest.rs ---- +// (repo policy: no inline #[cfg(test)] in implementation files). + +fn filter_meta(rel: &str, hash_byte: u8, mtime: u64) -> FileMetadata { + FileMetadata { + rel_path: PathBuf::from(rel), + size: 4, + hash: [hash_byte; 32], + modified_ms: mtime, + change_sequence: 0, + is_dir: false, + } +} + +fn filter_manifest_with(entries: Vec) -> Manifest { + Manifest { + files: entries + .into_iter() + .map(|m| (m.rel_path.clone(), m)) + .collect(), + node_id: "test".to_string(), + } +} + +fn filter_ledger_with(path: &str, deleted_at: u64, prev_hash: Option) -> DeletionLedger { + let mut l = DeletionLedger::new(); + l.record( + PathBuf::from(path), + deleted_at, + "peer".to_string(), + prev_hash, + Some(50), + ); + l +} + +#[test] +fn filter_resurrected_vetoes_stale_copy() { + let local = filter_manifest_with(vec![filter_meta("gone.txt", 0xAA, 50)]); + let prev_hash = hex(&[0xAA; 32]); + let ledger = filter_ledger_with("gone.txt", 100, Some(prev_hash)); + let (to_send, resurrected) = + filter_resurrected(vec![PathBuf::from("gone.txt")], &local, Some(&ledger)); + assert!(to_send.is_empty(), "stale copy must be vetoed"); + assert!(resurrected.is_empty()); +} + +#[test] +fn filter_resurrected_allows_legitimate_recreate() { + let local = filter_manifest_with(vec![filter_meta("back.txt", 0xBB, 150)]); + let ledger = filter_ledger_with("back.txt", 100, Some(hex(&[0xAA; 32]))); + let (to_send, resurrected) = + filter_resurrected(vec![PathBuf::from("back.txt")], &local, Some(&ledger)); + assert_eq!(to_send, vec![PathBuf::from("back.txt")]); + assert_eq!(resurrected, vec![PathBuf::from("back.txt")]); +} + +#[test] +fn filter_resurrected_passes_through_without_ledger() { + let local = filter_manifest_with(vec![filter_meta("a.txt", 1, 10)]); + let (to_send, resurrected) = + filter_resurrected(vec![PathBuf::from("a.txt")], &local, None); + assert_eq!(to_send, vec![PathBuf::from("a.txt")]); + assert!(resurrected.is_empty()); +} + +#[test] +fn filter_resurrected_diverged_content_is_not_vetoed() { + // Different hash AND newer mtime than prev => not the deleted version. + let local = filter_manifest_with(vec![filter_meta("c.txt", 0xCC, 60)]); + let ledger = filter_ledger_with("c.txt", 100, Some(hex(&[0xAA; 32]))); + let (to_send, _) = + filter_resurrected(vec![PathBuf::from("c.txt")], &local, Some(&ledger)); + // mtime 60 <= prev_mtime 50? No (60 > 50) and hash differs, so no veto. + assert_eq!(to_send, vec![PathBuf::from("c.txt")]); +} + +#[test] +fn compute_send_list_unchanged_without_ledger() { + let local = filter_manifest_with(vec![filter_meta("x.txt", 9, 200)]); + let remote = Manifest { + files: std::collections::HashMap::new(), + node_id: "r".to_string(), + }; + assert_eq!( + compute_send_list(&local, &remote, false), + vec![PathBuf::from("x.txt")] + ); +} diff --git a/crates/filesync/tests/test_sync_engine.rs b/crates/filesync/tests/test_sync_engine.rs index 150ff94..8f805f3 100644 --- a/crates/filesync/tests/test_sync_engine.rs +++ b/crates/filesync/tests/test_sync_engine.rs @@ -6,7 +6,9 @@ use std::{ use bytehive_filesync::{ exclusions::{ExclusionConfig, Exclusions}, - protocol::{FileBundle, FileData, FileMetadata}, + hex, + ledger::{now_ms, DeletionLedger}, + protocol::{FileBundle, FileData, FileMetadata, TombstoneDto}, sync_engine::{safe_relative, SyncEngine}, }; @@ -239,7 +241,7 @@ fn apply_deletes_removes_file() { assert!(dir.join("victim.txt").exists()); let n = engine - .apply_deletes(&[PathBuf::from("victim.txt")]) + .apply_deletes(&[PathBuf::from("victim.txt")], now_ms(), "test-node") .unwrap(); assert_eq!(n, 1); assert!(!dir.join("victim.txt").exists()); @@ -253,7 +255,9 @@ fn apply_deletes_nonexistent_file_still_counts() { let dir = tmp_dir("delete_nonexistent"); let engine = make_engine(dir.clone()); - let n = engine.apply_deletes(&[PathBuf::from("ghost.txt")]).unwrap(); + let n = engine + .apply_deletes(&[PathBuf::from("ghost.txt")], now_ms(), "test-node") + .unwrap(); assert_eq!(n, 1); fs::remove_dir_all(&dir).unwrap(); } @@ -263,7 +267,7 @@ fn apply_deletes_rejects_path_traversal() { let dir = tmp_dir("delete_traversal"); let engine = make_engine(dir.clone()); let n = engine - .apply_deletes(&[PathBuf::from("../not_here")]) + .apply_deletes(&[PathBuf::from("../not_here")], now_ms(), "test-node") .unwrap(); assert_eq!(n, 0, "path traversal in delete must be silently skipped"); fs::remove_dir_all(&dir).unwrap(); @@ -275,7 +279,9 @@ fn apply_deletes_removes_directory_recursively() { let engine = make_engine(dir.clone()); engine.apply_bundle(&dir_bundle("mydir")).unwrap(); fs::write(dir.join("mydir/file.txt"), b"data").unwrap(); - engine.apply_deletes(&[PathBuf::from("mydir")]).unwrap(); + engine + .apply_deletes(&[PathBuf::from("mydir")], now_ms(), "test-node") + .unwrap(); assert!(!dir.join("mydir").exists()); fs::remove_dir_all(&dir).unwrap(); } @@ -591,7 +597,9 @@ fn apply_deletes_directory_clears_manifest_entries() { let before = engine.get_manifest(); assert!(before.files.contains_key(&PathBuf::from("mydir"))); assert!(before.files.contains_key(&PathBuf::from("mydir/child.txt"))); - engine.apply_deletes(&[PathBuf::from("mydir")]).unwrap(); + engine + .apply_deletes(&[PathBuf::from("mydir")], now_ms(), "test-node") + .unwrap(); let after = engine.get_manifest(); assert!(!after.files.contains_key(&PathBuf::from("mydir"))); assert!(!after.files.contains_key(&PathBuf::from("mydir/child.txt"))); @@ -695,3 +703,127 @@ fn clear_in_progress_removes_assembly() { assert!(matches!(result, FinishResult::Committed)); fs::remove_dir_all(&dir).unwrap(); } + +// ---- Deletion-ledger engine tests, relocated from src/sync_engine.rs ---- +// (repo policy: no inline #[cfg(test)] in implementation files; private +// `ledger` field access was rewritten through public API + disk reload). + +#[test] +fn apply_deletes_records_tombstone_with_prev_info() { + let dir = tmp_dir("ledger_prev_info"); + let engine = make_engine(dir.clone()); + fs::write(dir.join("victim.txt"), b"data").unwrap(); + engine.scan().unwrap(); + let before = engine + .get_manifest() + .files + .get(&PathBuf::from("victim.txt")) + .cloned(); + assert!(before.is_some()); + let before = before.unwrap(); + + let n = engine + .apply_deletes(&[PathBuf::from("victim.txt")], 1_700_000_000_000, "node-peer") + .unwrap(); + assert_eq!(n, 1); + assert!(!dir.join("victim.txt").exists()); + + // Read back through the persisted ledger (public API only). + let tomb = DeletionLedger::load(&dir) + .get(&PathBuf::from("victim.txt")) + .cloned() + .expect("tombstone must be recorded"); + assert_eq!(tomb.deleted_at_ms, 1_700_000_000_000); + assert_eq!(tomb.deleter_node, "node-peer"); + assert_eq!(tomb.prev_hash.as_deref(), Some(hex(&before.hash).as_str())); + assert_eq!(tomb.prev_mtime_ms, Some(before.modified_ms)); + fs::remove_dir_all(&dir).unwrap(); +} + +#[test] +fn apply_deletes_absent_path_still_records_tombstone() { + let dir = tmp_dir("ledger_absent"); + let engine = make_engine(dir.clone()); + let n = engine + .apply_deletes(&[PathBuf::from("ghost.txt")], 42, "node-peer") + .unwrap(); + assert_eq!(n, 1); + let tomb = DeletionLedger::load(&dir) + .get(&PathBuf::from("ghost.txt")) + .cloned() + .expect("absent path must still get a tombstone"); + assert_eq!(tomb.deleted_at_ms, 42); + assert!(tomb.prev_hash.is_none()); + fs::remove_dir_all(&dir).unwrap(); +} + +#[test] +fn apply_deletes_preserves_sender_timestamp_lww() { + let dir = tmp_dir("ledger_lww"); + let engine = make_engine(dir.clone()); + engine + .apply_deletes(&[PathBuf::from("a.txt")], 200, "n1") + .unwrap(); + // Older sender timestamp must not clobber the newer tombstone. + engine + .apply_deletes(&[PathBuf::from("a.txt")], 100, "n2") + .unwrap(); + assert_eq!( + DeletionLedger::load(&dir) + .get(&PathBuf::from("a.txt")) + .unwrap() + .deleter_node, + "n1" + ); + fs::remove_dir_all(&dir).unwrap(); +} + +#[test] +fn merge_and_export_ledger_entries_roundtrip() { + let dir = tmp_dir("ledger_merge"); + let engine = make_engine(dir.clone()); + let dtos = vec![TombstoneDto { + path: PathBuf::from("r.txt"), + deleted_at_ms: 777, + deleter_node: "node-remote".to_string(), + prev_hash: None, + prev_mtime_ms: None, + }]; + assert_eq!(engine.merge_ledger(dtos), 1); + let entries = engine.ledger_entries(); + assert_eq!(entries.len(), 1); + assert_eq!(entries[0].path, PathBuf::from("r.txt")); + assert_eq!(entries[0].deleted_at_ms, 777); + assert!(engine.has_tombstone_newer_than(&PathBuf::from("r.txt"), 776)); + assert!(!engine.has_tombstone_newer_than(&PathBuf::from("r.txt"), 777)); + fs::remove_dir_all(&dir).unwrap(); +} + +#[test] +fn prune_ledger_drops_expired_only() { + let dir = tmp_dir("ledger_prune"); + let engine = make_engine(dir.clone()); + engine.merge_ledger(vec![TombstoneDto { + path: PathBuf::from("old.txt"), + deleted_at_ms: 1, // ancient + deleter_node: "n".to_string(), + prev_hash: None, + prev_mtime_ms: None, + }]); + engine.record_delete(PathBuf::from("fresh.txt"), now_ms(), "n", None, None); + assert_eq!(engine.prune_ledger(), 1); + let loaded = DeletionLedger::load(&dir); + assert!(!loaded.contains(&PathBuf::from("old.txt"))); + assert!(loaded.contains(&PathBuf::from("fresh.txt"))); + fs::remove_dir_all(&dir).unwrap(); +} + +#[test] +fn clear_tombstones_removes_on_recreate_win() { + let dir = tmp_dir("ledger_clear"); + let engine = make_engine(dir.clone()); + engine.record_delete(PathBuf::from("back.txt"), 100, "n", None, None); + engine.clear_tombstones(&[PathBuf::from("back.txt")]); + assert!(!DeletionLedger::load(&dir).contains(&PathBuf::from("back.txt"))); + fs::remove_dir_all(&dir).unwrap(); +} From 31ede327a74f823b9b930f40139e5292a1fa981b Mon Sep 17 00:00:00 2001 From: Guy Dols Date: Thu, 10 Sep 2026 22:05:27 +0200 Subject: [PATCH 5/6] feat(filesync): exchange ledgers on reconnect to push missed deletes After ManifestExchange both peers swap LedgerExchange, merge with last-writer-wins, delete locally-resurrected files with the original stamp, and push deletes the offline peer missed. Live deletes preserve the sender timestamp on forward. --- crates/filesync/src/client.rs | 159 +++++++++++++++++++++++++++++-- crates/filesync/src/server.rs | 171 ++++++++++++++++++++++++++++++++-- 2 files changed, 316 insertions(+), 14 deletions(-) diff --git a/crates/filesync/src/client.rs b/crates/filesync/src/client.rs index a114178..85e91d9 100644 --- a/crates/filesync/src/client.rs +++ b/crates/filesync/src/client.rs @@ -3,6 +3,7 @@ use crate::common::{self, LargeFileEndOutcome, PendingChanges}; use crate::gui::state::ConnectionStatus; use crate::gui::state::{ConflictKind, SharedState}; use crate::known_hosts::KnownServers; +use crate::ledger::{DeletionLedger, Tombstone}; use crate::manifest; use crate::protocol::*; use crate::sync_engine::{ConflictInfo, SyncEngine}; @@ -16,6 +17,7 @@ use log::{debug, error, info, warn}; use parking_lot::Mutex; use std::io; +use std::collections::HashSet; use std::net::TcpStream; use std::path::{Path, PathBuf}; use std::sync::atomic::{AtomicBool, Ordering}; @@ -472,6 +474,30 @@ impl Client { debug!("filesync session: sending local ManifestExchange"); conn.send(&Message::ManifestExchange(local.clone()))?; + // ---- Deletion-ledger exchange (symmetric order: server sent first, + // we reply — avoids both-wait deadlock) ---- + debug!("filesync session: waiting for server LedgerExchange"); + let peer_ledger_entries = match conn.recv()? { + Message::LedgerExchange { entries } => entries, + other => { + error!("filesync session: expected LedgerExchange from server"); + debug!( + "filesync session: unexpected message instead of LedgerExchange: {other:?}" + ); + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "filesync: expected LedgerExchange", + )); + } + }; + let merged = self.engine.merge_ledger(peer_ledger_entries); + debug!("filesync session: merged {merged} tombstone(s) from server"); + conn.send(&Message::LedgerExchange { + entries: self.engine.ledger_entries(), + })?; + debug!("filesync session: sent LedgerExchange to server"); + // NOTE: LedgerAck intentionally not sent — no side reads it. + let expected_rx: u64 = remote .files .iter() @@ -632,6 +658,30 @@ impl Client { } } } + Message::Delete { + paths, + deleted_at_ms, + deleter, + } => { + debug!( + "filesync session: recv Delete {} path(s) during initial sync", + paths.len() + ); + if let Err(e) = common::handle_recv_delete_with_meta( + &self.engine, + &paths, + deleted_at_ms, + &deleter, + "server", + &self.bus, + "filesync session", + ) { + error!("filesync session: apply_deletes during initial sync: {e}"); + } + if let Some(ref gs) = self.gui_state { + gs.write().begin_sync_activity(); + } + } Message::SyncComplete => { debug!( "filesync session: received SyncComplete — files_received={} dirs_received={} bytes_received={} B elapsed={}ms", @@ -665,7 +715,38 @@ impl Client { } } - let to_send = manifest::compute_send_list(&local, &remote, false); + let raw_to_send = manifest::compute_send_list(&local, &remote, false); + let snap = local_ledger_snapshot(&self.engine); + let (to_send, recreated) = + manifest::filter_resurrected(raw_to_send.clone(), &local, Some(&snap)); + // Tombstone wins over our stale copy: converge by deleting locally + // while preserving the original tombstone stamp (never re-stamp). + { + let send_set: HashSet<&PathBuf> = to_send.iter().collect(); + for path in raw_to_send.iter().filter(|p| !send_set.contains(p)) { + if let Some(tomb) = snap.get(path) { + debug!( + "filesync session: tombstone wins for {path:?} — deleting local stale copy" + ); + self.engine.apply_deletes( + &[path.clone()], + tomb.deleted_at_ms, + &tomb.deleter_node, + )?; + } + } + } + // Refresh: veto step changed disk + manifest. + let local = self.engine.get_manifest(); + // Tell the server about deletes it missed while we were offline. + for (path, deleted_at_ms, deleter) in compute_deletes_to_push(&local, &remote, &snap) { + debug!("filesync session: pushing missed delete {path:?} to server"); + conn.send(&Message::Delete { + paths: vec![path], + deleted_at_ms, + deleter, + })?; + } let files_sent = to_send .iter() .filter(|p| local.files.get(*p).map(|m| !m.is_dir).unwrap_or(false)) @@ -688,6 +769,8 @@ impl Client { self.engine.send_paths(&to_send, &conn)?; debug!("filesync session: send_paths complete"); } + // Legitimate recreates were just uploaded — lift their tombstones. + self.engine.clear_tombstones(&recreated); debug!("filesync session: sending SyncComplete to server"); conn.send(&Message::SyncComplete)?; @@ -937,10 +1020,21 @@ fn recv_loop( Err(e) => error!("{prefix}: large_file_end: {e}"), } } - Ok(Message::Delete { paths }) => { + Ok(Message::Delete { + paths, + deleted_at_ms, + deleter, + }) => { debug!("{prefix}: Delete {} path(s)", paths.len()); - if let Err(e) = common::handle_recv_delete(&engine, &paths, "server", &bus, prefix) - { + if let Err(e) = common::handle_recv_delete_with_meta( + &engine, + &paths, + deleted_at_ms, + &deleter, + "server", + &bus, + prefix, + ) { error!("{prefix}: apply_deletes: {e}"); } if let Some(ref gs) = gui_state { @@ -1080,6 +1174,55 @@ fn send_loop( debug!("filesync send: loop exited after {} flush(es)", flush_count); } +/// In-memory snapshot of the engine's deletion ledger for sync decisions +/// (rebuilt from `ledger_entries`; avoids new sync_engine accessors). +fn local_ledger_snapshot(engine: &SyncEngine) -> DeletionLedger { + let mut snap = DeletionLedger::new(); + snap.merge_remote( + engine + .ledger_entries() + .into_iter() + .map(|d| { + let path = d.path.clone(); + ( + path, + Tombstone::new( + d.path, + d.deleted_at_ms, + d.deleter_node, + d.prev_hash, + d.prev_mtime_ms, + ), + ) + }) + .collect(), + ); + snap +} + +/// Paths the peer still lists but our merged ledger marks deleted (tombstone +/// newer than the peer's copy): the peer missed the delete while offline. +/// Returns `(path, deleted_at_ms, deleter)` so the original stamp survives. +fn compute_deletes_to_push( + local: &Manifest, + remote: &Manifest, + ledger: &DeletionLedger, +) -> Vec<(PathBuf, u64, String)> { + let mut out: Vec<(PathBuf, u64, String)> = Vec::new(); + for (path, tomb) in ledger.iter() { + if local.files.contains_key(path) { + continue; + } + if let Some(peer_meta) = remote.files.get(path) { + if tomb.deleted_at_ms > peer_meta.modified_ms { + out.push((path.clone(), tomb.deleted_at_ms, tomb.deleter_node.clone())); + } + } + } + out.sort_by(|a, b| a.0.cmp(&b.0)); + out +} + fn flush_to_server( engine: &Arc, conn: &Arc, @@ -1129,14 +1272,18 @@ fn flush_to_server( send_paths_to_server(engine, conn, bus, gui_state, stable)?; } - let (paths, delete_count) = pending.take_deletes(engine.root()); + let (paths, delete_count) = pending.take_deletes_with_engine(engine); if !paths.is_empty() { debug!( "filesync send: flushing {} delete path(s) ({} pre-expansion)", paths.len(), delete_count ); - conn.send(&Message::Delete { paths })?; + conn.send(&Message::Delete { + paths, + deleted_at_ms: crate::ledger::now_ms(), + deleter: engine.node_id().to_string(), + })?; if let Some(ref gs) = gui_state { gs.write().begin_sync_activity(); } diff --git a/crates/filesync/src/server.rs b/crates/filesync/src/server.rs index e6c6abd..b82b8f0 100644 --- a/crates/filesync/src/server.rs +++ b/crates/filesync/src/server.rs @@ -2,6 +2,7 @@ use crate::bundler; use crate::cert_fingerprint; use crate::common::{self, LargeFileEndOutcome, PendingChanges}; use crate::known_hosts::{ClientStatus, KnownClients}; +use crate::ledger::{DeletionLedger, Tombstone}; use crate::manifest; use crate::protocol::*; use crate::sync_engine::SyncEngine; @@ -374,6 +375,25 @@ fn handle_client( r_files, r_dirs ); + // ---- Deletion-ledger exchange (symmetric order: server sends first, + // client replies — avoids both-wait deadlock) ---- + conn.send(&Message::LedgerExchange { + entries: engine.ledger_entries(), + })?; + debug!("filesync: sent LedgerExchange to {client_id}"); + let peer_ledger_entries = match conn.recv()? { + Message::LedgerExchange { entries } => entries, + _ => { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "filesync: expected LedgerExchange", + )) + } + }; + let merged = engine.merge_ledger(peer_ledger_entries); + debug!("filesync: merged {merged} tombstone(s) from {client_id}"); + // NOTE: LedgerAck intentionally not sent — no side reads it. + let bytes_from_client: u64 = manifest::compute_send_list(&remote, &local, false) .iter() .filter_map(|p| remote.files.get(p)) @@ -405,7 +425,32 @@ fn handle_client( ); } - let to_client = manifest::compute_send_list(&local, &remote, true); + let raw_to_client = manifest::compute_send_list(&local, &remote, true); + let snap = local_ledger_snapshot(&engine); + let (to_client, recreated) = + manifest::filter_resurrected(raw_to_client.clone(), &local, Some(&snap)); + // Tombstone wins over our stale copy: converge by deleting locally while + // preserving the original tombstone stamp (never re-stamp here). + { + let send_set: HashSet<&PathBuf> = to_client.iter().collect(); + for path in raw_to_client.iter().filter(|p| !send_set.contains(p)) { + if let Some(tomb) = snap.get(path) { + debug!("filesync: tombstone wins for {path:?} — deleting local stale copy"); + engine.apply_deletes(&[path.clone()], tomb.deleted_at_ms, &tomb.deleter_node)?; + } + } + } + // Refresh: veto step changed disk + manifest. + let local = engine.get_manifest(); + // Tell the peer about deletes it missed while offline. + for (path, deleted_at_ms, deleter) in compute_deletes_to_push(&local, &remote, &snap) { + debug!("filesync: pushing missed delete {path:?} to {client_id}"); + conn.send(&Message::Delete { + paths: vec![path], + deleted_at_ms, + deleter, + })?; + } let files_sent_to_client = to_client .iter() .filter(|p| local.files.get(*p).map(|m| !m.is_dir).unwrap_or(false)) @@ -430,6 +475,8 @@ fn handle_client( engine.send_paths(&to_client, &conn)?; debug!("filesync: send_paths to {client_id} complete"); } + // Legitimate recreates were just uploaded — lift their stale tombstones. + engine.clear_tombstones(&recreated); conn.send(&Message::SyncComplete)?; debug!("filesync: sent SyncComplete to {client_id}, waiting for client's initial sync"); @@ -580,6 +627,38 @@ fn handle_client( Err(e) => error!("filesync: initial sync large_file_end from {client_id}: {e}"), } } + Message::Delete { + paths, + deleted_at_ms, + deleter, + } => { + debug!( + "filesync: initial recv Delete from {client_id}: {} path(s)", + paths.len() + ); + match common::handle_recv_delete_with_meta( + &engine, + &paths, + deleted_at_ms, + &deleter, + &client_id, + &bus, + "filesync", + ) { + Ok(_) => broadcast_to_others( + &peers, + &client_id, + &Message::Delete { + paths, + deleted_at_ms, + deleter, + }, + ), + Err(e) => { + error!("filesync: initial sync apply_deletes from {client_id}: {e}") + } + } + } Message::SyncComplete => { debug!( "filesync: received SyncComplete from {client_id} — initial sync complete in {}ms", @@ -793,13 +872,33 @@ fn client_recv_loop( Err(e) => error!("filesync: large_file_end from {client_id}: {e}"), } } - Ok(Message::Delete { paths }) => { + Ok(Message::Delete { + paths, + deleted_at_ms, + deleter, + }) => { debug!( "filesync: recv Delete from {client_id}: {} path(s)", paths.len() ); - match common::handle_recv_delete(&engine, &paths, &client_id, &bus, "filesync") { - Ok(_) => broadcast_to_others(&peers, &client_id, &Message::Delete { paths }), + match common::handle_recv_delete_with_meta( + &engine, + &paths, + deleted_at_ms, + &deleter, + &client_id, + &bus, + "filesync", + ) { + Ok(_) => broadcast_to_others( + &peers, + &client_id, + &Message::Delete { + paths, + deleted_at_ms, + deleter, + }, + ), Err(e) => error!("filesync: apply_deletes from {client_id}: {e}"), } } @@ -1014,7 +1113,7 @@ fn flush_local_changes( "filesync server: broadcasting {} delete(s)", pending.deletes.len() ); - let (paths, _) = pending.take_deletes(engine.root()); + let (paths, _) = pending.take_deletes_with_engine(engine); if let Some(ref bus) = bus { bus.publish( "filesync", @@ -1031,7 +1130,15 @@ fn flush_local_changes( }), ); } - broadcast_to_others(peers, "", &Message::Delete { paths }); + broadcast_to_others( + peers, + "", + &Message::Delete { + paths, + deleted_at_ms: crate::ledger::now_ms(), + deleter: engine.node_id().to_string(), + }, + ); } pending.reset_timer(); @@ -1145,8 +1252,56 @@ fn broadcast_paths( } } -fn broadcast_to_others(peers: &Arc>>, sender_id: &str, msg: &Message) { - let peers = peers.read(); +/// In-memory snapshot of the engine's deletion ledger for sync decisions +/// (rebuilt from `ledger_entries`; avoids new sync_engine accessors). +fn local_ledger_snapshot(engine: &SyncEngine) -> DeletionLedger { + let mut snap = DeletionLedger::new(); + snap.merge_remote( + engine + .ledger_entries() + .into_iter() + .map(|d| { + let path = d.path.clone(); + ( + path, + Tombstone::new( + d.path, + d.deleted_at_ms, + d.deleter_node, + d.prev_hash, + d.prev_mtime_ms, + ), + ) + }) + .collect(), + ); + snap +} + +/// Paths the peer still lists but our merged ledger marks deleted (tombstone +/// newer than the peer's copy): the peer missed the delete while offline. +/// Returns `(path, deleted_at_ms, deleter)` so the original stamp survives. +fn compute_deletes_to_push( + local: &Manifest, + remote: &Manifest, + ledger: &DeletionLedger, +) -> Vec<(PathBuf, u64, String)> { + let mut out: Vec<(PathBuf, u64, String)> = Vec::new(); + for (path, tomb) in ledger.iter() { + if local.files.contains_key(path) { + continue; + } + if let Some(peer_meta) = remote.files.get(path) { + if tomb.deleted_at_ms > peer_meta.modified_ms { + out.push((path.clone(), tomb.deleted_at_ms, tomb.deleter_node.clone())); + } + } + } + out.sort_by(|a, b| a.0.cmp(&b.0)); + out +} + +fn broadcast_to_others(peers: &Arc>>, sender_id: &str, msg: &Message) { let peers = peers.read(); if peers.is_empty() { return; } From 558bd9c7423f0a69706a8d4f4d87bed541348b61 Mon Sep 17 00:00:00 2001 From: Guy Dols Date: Thu, 10 Sep 2026 22:05:27 +0200 Subject: [PATCH 6/6] test(filesync): add offline-delete regression and ledger coverage Covers the offline-delete-stays-deleted repro, recreate-wins, ledger persist/prune/TTL, corrupt-backup handling, and v8 wire roundtrips. --- crates/filesync/tests/test_deletion_ledger.rs | 318 ++++++++++++++++++ crates/filesync/tests/test_ledger_unit.rs | 124 +++++++ 2 files changed, 442 insertions(+) create mode 100644 crates/filesync/tests/test_deletion_ledger.rs create mode 100644 crates/filesync/tests/test_ledger_unit.rs diff --git a/crates/filesync/tests/test_deletion_ledger.rs b/crates/filesync/tests/test_deletion_ledger.rs new file mode 100644 index 0000000..3489c6c --- /dev/null +++ b/crates/filesync/tests/test_deletion_ledger.rs @@ -0,0 +1,318 @@ +//! Regression tests for the filesync deletion ledger (TODO 4): +//! persistent tombstones, 90-day TTL, protocol v8, no compat shim. +//! +//! Fast, no network: two `SyncEngine`s over `tempfile` dirs simulate the +//! offline-delete / recreate races; wire roundtrips cover the v8 messages. + +use std::io::Cursor; +use std::path::PathBuf; +use std::sync::Arc; + +use bytehive_filesync::{ + exclusions::{ExclusionConfig, Exclusions}, + ledger::{now_ms, DeletionLedger, DELETION_LEDGER_TTL_MS}, + manifest, + protocol::{ + read_message, serialise_message, Message, TombstoneDto, PROTOCOL_VERSION, + }, + sync_engine::SyncEngine, +}; + +const DAY_MS: u64 = 24 * 60 * 60 * 1000; + +fn make_engine(root: PathBuf, node: &str) -> SyncEngine { + let ex = Arc::new(Exclusions::compile(&ExclusionConfig::default())); + SyncEngine::new(root, node.to_string(), ex) +} + +/// `(path, deleted_at_ms, deleter, prev_hash, prev_mtime)` — order-stable +/// projection for ledger-equality asserts (`TombstoneDto` has no `PartialEq`). +fn dto_key(d: &TombstoneDto) -> (PathBuf, u64, String, Option, Option) { + ( + d.path.clone(), + d.deleted_at_ms, + d.deleter_node.clone(), + d.prev_hash.clone(), + d.prev_mtime_ms, + ) +} + +fn sorted_entries(engine: &SyncEngine) -> Vec<(PathBuf, u64, String, Option, Option)> { + let mut v: Vec<_> = engine.ledger_entries().iter().map(dto_key).collect(); + v.sort(); + v +} + +#[test] +fn ttl_constant_is_90_days() { + assert_eq!(DELETION_LEDGER_TTL_MS, 90 * DAY_MS); + assert_eq!(PROTOCOL_VERSION, 8, "ledger work rode the clean break to v8"); +} + +/// Core resurrection repro: A deletes `f` while B is offline holding a stale +/// copy. After ledger merge, B must veto the upload, converge locally, and +/// both ledgers must agree. +#[test] +fn offline_delete_stays_deleted() { + let dir_a = tempfile::tempdir().unwrap(); + let dir_b = tempfile::tempdir().unwrap(); + let engine_a = make_engine(dir_a.path().to_path_buf(), "node-a"); + let engine_b = make_engine(dir_b.path().to_path_buf(), "node-b"); + + // Same content on both sides; B's copy goes stale (B is "offline"). + std::fs::write(dir_a.path().join("f.txt"), b"shared content").unwrap(); + std::fs::write(dir_b.path().join("f.txt"), b"shared content").unwrap(); + engine_a.scan().unwrap(); + engine_b.scan().unwrap(); + + // A deletes f at T0 (after B's mtime, so the tombstone is strictly newer). + let t0 = now_ms() + 5_000; + let n = engine_a + .apply_deletes(&[PathBuf::from("f.txt")], t0, "node-a") + .unwrap(); + assert_eq!(n, 1); + assert!(!dir_a.path().join("f.txt").exists()); + + // B comes back online: merge A's ledger, then plan its upload. + assert_eq!(engine_b.merge_ledger(engine_a.ledger_entries()), 1); + let local_b = engine_b.get_manifest(); + let remote_a = engine_a.get_manifest(); + let send = manifest::compute_send_list(&local_b, &remote_a, false); + assert!( + send.contains(&PathBuf::from("f.txt")), + "without the ledger, B would resurrect f.txt by uploading it" + ); + + let ledger_b = DeletionLedger::load(dir_b.path()); + let (to_send, recreated) = + manifest::filter_resurrected(send, &local_b, Some(&ledger_b)); + assert!( + !to_send.contains(&PathBuf::from("f.txt")), + "stale copy must be vetoed by the newer tombstone" + ); + assert!(recreated.is_empty()); + + // Veto convergence (mirrors the client/server lane): delete locally with + // the ORIGINAL stamp so the tombstone survives. + let tomb = ledger_b + .get(&PathBuf::from("f.txt")) + .expect("tombstone must exist after merge"); + assert_eq!(tomb.deleted_at_ms, t0); + engine_b + .apply_deletes( + &[PathBuf::from("f.txt")], + tomb.deleted_at_ms, + &tomb.deleter_node.clone(), + ) + .unwrap(); + assert!(!dir_b.path().join("f.txt").exists()); + assert!(!engine_b + .get_manifest() + .files + .contains_key(&PathBuf::from("f.txt"))); + + // Both sides agree on the tombstone. + assert_eq!(sorted_entries(&engine_a), sorted_entries(&engine_b)); + assert_eq!(sorted_entries(&engine_b)[0].1, t0); + assert_eq!(sorted_entries(&engine_b)[0].2, "node-a"); +} + +/// Recreate wins: B's copy is newer than the tombstone (and different +/// content) → upload proceeds and the stale tombstone lifts. +#[test] +fn recreate_after_delete_wins() { + let dir_a = tempfile::tempdir().unwrap(); + let dir_b = tempfile::tempdir().unwrap(); + let engine_a = make_engine(dir_a.path().to_path_buf(), "node-a"); + let engine_b = make_engine(dir_b.path().to_path_buf(), "node-b"); + + std::fs::write(dir_a.path().join("f.txt"), b"old content").unwrap(); + engine_a.scan().unwrap(); + std::fs::write(dir_b.path().join("f.txt"), b"new content").unwrap(); + engine_b.scan().unwrap(); + let mtime_b = engine_b + .get_manifest() + .files + .get(&PathBuf::from("f.txt")) + .unwrap() + .modified_ms; + + // Delete predates B's current copy. + let t0 = mtime_b.saturating_sub(10_000); + engine_a + .apply_deletes(&[PathBuf::from("f.txt")], t0, "node-a") + .unwrap(); + engine_b.merge_ledger(engine_a.ledger_entries()); + + let local_b = engine_b.get_manifest(); + let send = vec![PathBuf::from("f.txt")]; + let ledger_b = DeletionLedger::load(dir_b.path()); + let (to_send, recreated) = + manifest::filter_resurrected(send, &local_b, Some(&ledger_b)); + assert_eq!(to_send, vec![PathBuf::from("f.txt")]); + assert_eq!(recreated, vec![PathBuf::from("f.txt")]); + + engine_b.clear_tombstones(&recreated); + assert!(engine_b.ledger_entries().is_empty()); + assert!(DeletionLedger::load(dir_b.path()).is_empty()); + // The recreated file itself is untouched. + assert_eq!( + std::fs::read(dir_b.path().join("f.txt")).unwrap(), + b"new content" + ); +} + +/// Persistence roundtrip, 90-day TTL boundaries, corrupt-file backup. +#[test] +fn ledger_persists_and_prunes() { + let dir = tempfile::tempdir().unwrap(); + let engine = make_engine(dir.path().to_path_buf(), "node-a"); + + let t = now_ms(); + engine.record_delete( + PathBuf::from("gone.txt"), + t, + "node-a", + Some("ab".repeat(32)), + Some(t - 1_000), + ); + + // On-disk artifact exists. + let ledger_file = dir.path().join(".bh_filesync").join("deletion-ledger.json"); + assert!(ledger_file.is_file(), "ledger must persist under .bh_filesync"); + + // Save→load roundtrip preserves the tombstone. + let loaded = DeletionLedger::load(dir.path()); + let tomb = loaded.get(&PathBuf::from("gone.txt")).unwrap(); + assert_eq!(tomb.deleted_at_ms, t); + assert_eq!(tomb.deleter_node, "node-a"); + assert_eq!(tomb.prev_mtime_ms, Some(t - 1_000)); + + // TTL boundaries: exactly 90d keeps, 90d+1ms drops, 89d keeps. + let mut at_90d = loaded.clone(); + assert_eq!(at_90d.prune(t + 90 * DAY_MS, DELETION_LEDGER_TTL_MS), 0); + let mut past = loaded.clone(); + assert_eq!( + past.prune(t + 90 * DAY_MS + 1, DELETION_LEDGER_TTL_MS), + 1 + ); + assert!(past.is_empty()); + let mut before = loaded.clone(); + assert_eq!(before.prune(t + 89 * DAY_MS, DELETION_LEDGER_TTL_MS), 0); + assert!(!before.is_empty()); + + // Corrupt JSON backs up and loads empty instead of panicking. + std::fs::write(&ledger_file, "{ this is not json").unwrap(); + let recovered = DeletionLedger::load(dir.path()); + assert!(recovered.is_empty()); + let backups: Vec<_> = std::fs::read_dir(dir.path().join(".bh_filesync")) + .unwrap() + .filter_map(|e| e.ok()) + .filter(|e| { + e.file_name() + .to_string_lossy() + .starts_with("deletion-ledger.corrupt-") + }) + .collect(); + assert_eq!(backups.len(), 1); +} + +/// Protocol v8: LedgerExchange encode/decode (Delete shape already covered +/// in test_protocol.rs; re-asserted here for the no-compat story). +#[test] +fn ledger_exchange_roundtrip() { + let msg = Message::LedgerExchange { + entries: vec![ + TombstoneDto { + path: PathBuf::from("gone.txt"), + deleted_at_ms: 1_700_000_000_000, + deleter_node: "node-a".to_string(), + prev_hash: Some("ab".repeat(32)), + prev_mtime_ms: Some(1_699_999_999_000), + }, + TombstoneDto { + path: PathBuf::from("old/dir"), + deleted_at_ms: 42, + deleter_node: "node-b".to_string(), + prev_hash: None, + prev_mtime_ms: None, + }, + ], + }; + let frame = serialise_message(&msg).unwrap(); + let mut cur = Cursor::new(frame); + let expected_prev = "ab".repeat(32); + match read_message(&mut cur).unwrap() { + Message::LedgerExchange { entries } => { + assert_eq!(entries.len(), 2); + assert_eq!(entries[0].path, PathBuf::from("gone.txt")); + assert_eq!(entries[0].deleted_at_ms, 1_700_000_000_000); + assert_eq!(entries[0].deleter_node, "node-a"); + assert_eq!(entries[0].prev_hash.as_deref(), Some(expected_prev.as_str())); + assert_eq!(entries[0].prev_mtime_ms, Some(1_699_999_999_000)); + assert_eq!(entries[1].path, PathBuf::from("old/dir")); + assert_eq!(entries[1].prev_hash, None); + } + _ => panic!("expected LedgerExchange"), + } +} + +#[test] +fn v8_delete_shape_roundtrip() { + let msg = Message::Delete { + paths: vec![PathBuf::from("f.txt")], + deleted_at_ms: 1_700_000_000_001, + deleter: "node-a".to_string(), + }; + let frame = serialise_message(&msg).unwrap(); + let mut cur = Cursor::new(frame); + match read_message(&mut cur).unwrap() { + Message::Delete { + paths, + deleted_at_ms, + deleter, + } => { + assert_eq!(paths, vec![PathBuf::from("f.txt")]); + assert_eq!(deleted_at_ms, 1_700_000_000_001); + assert_eq!(deleter, "node-a"); + } + _ => panic!("expected Delete"), + } +} + +/// Tombstone payload survives the engine DTO export used by the exchange. +#[test] +fn ledger_dto_export_preserves_prev_info() { + let dir = tempfile::tempdir().unwrap(); + let engine = make_engine(dir.path().to_path_buf(), "node-a"); + std::fs::write(dir.path().join("f.txt"), b"content").unwrap(); + engine.scan().unwrap(); + let meta = engine + .get_manifest() + .files + .get(&PathBuf::from("f.txt")) + .unwrap() + .clone(); + + let t0 = now_ms() + 5_000; + engine + .apply_deletes(&[PathBuf::from("f.txt")], t0, "node-a") + .unwrap(); + let entries = engine.ledger_entries(); + assert_eq!(entries.len(), 1); + let expected_hash = bytehive_filesync::hex(&meta.hash); + assert_eq!(entries[0].prev_hash.as_deref(), Some(expected_hash.as_str())); + assert_eq!(entries[0].prev_mtime_ms, Some(meta.modified_ms)); + + // And the DTOs feed back through merge on a second engine unchanged. + let dir_b = tempfile::tempdir().unwrap(); + let engine_b = make_engine(dir_b.path().to_path_buf(), "node-b"); + assert_eq!(engine_b.merge_ledger(entries), 1); + let got = sorted_entries(&engine_b); + assert_eq!(got.len(), 1); + assert_eq!(got[0].0, PathBuf::from("f.txt")); + assert_eq!(got[0].1, t0); + assert_eq!(got[0].2, "node-a"); + assert_eq!(got[0].3.as_deref(), Some(expected_hash.as_str())); + assert_eq!(got[0].4, Some(meta.modified_ms)); +} diff --git a/crates/filesync/tests/test_ledger_unit.rs b/crates/filesync/tests/test_ledger_unit.rs new file mode 100644 index 0000000..ec10134 --- /dev/null +++ b/crates/filesync/tests/test_ledger_unit.rs @@ -0,0 +1,124 @@ +//! Unit tests for [`bytehive_filesync::ledger`] relocated from `src/ledger.rs` +//! (repo policy: no inline `#[cfg(test)]` in implementation files). + +use std::collections::HashMap; +use std::path::{Path, PathBuf}; + +use bytehive_filesync::ledger::{ + now_ms, DeletionLedger, Tombstone, LEDGER_DIR_NAME, LEDGER_FILE_NAME, +}; + +fn tomb(path: &str, at: u64) -> Tombstone { + Tombstone::new(PathBuf::from(path), at, "node-a".to_string(), None, None) +} + +#[test] +fn record_keeps_newest() { + let mut l = DeletionLedger::new(); + l.record(PathBuf::from("a.txt"), 100, "n1".to_string(), None, None); + l.record(PathBuf::from("a.txt"), 50, "n2".to_string(), None, None); + assert_eq!(l.get(Path::new("a.txt")).unwrap().deleted_at_ms, 100); + l.record(PathBuf::from("a.txt"), 200, "n2".to_string(), None, None); + assert_eq!(l.get(Path::new("a.txt")).unwrap().deleted_at_ms, 200); +} + +#[test] +fn has_newer_than_strict() { + let mut l = DeletionLedger::new(); + assert!(!l.has_newer_than(Path::new("a.txt"), 10)); + l.record(PathBuf::from("a.txt"), 100, "n1".to_string(), None, None); + assert!(l.has_newer_than(Path::new("a.txt"), 99)); + assert!(!l.has_newer_than(Path::new("a.txt"), 100)); + assert!(!l.has_newer_than(Path::new("a.txt"), 101)); +} + +#[test] +fn merge_remote_lww() { + let mut local = DeletionLedger::new(); + local.record(PathBuf::from("a.txt"), 100, "n1".to_string(), None, None); + let mut remote = HashMap::new(); + remote.insert(PathBuf::from("a.txt"), tomb("a.txt", 50)); + remote.insert(PathBuf::from("b.txt"), tomb("b.txt", 300)); + local.merge_remote(remote); + assert_eq!(local.get(Path::new("a.txt")).unwrap().deleted_at_ms, 100); + assert_eq!(local.get(Path::new("b.txt")).unwrap().deleted_at_ms, 300); +} + +#[test] +fn prune_ttl() { + let mut l = DeletionLedger::new(); + l.record(PathBuf::from("old.txt"), 0, "n".to_string(), None, None); + l.record( + PathBuf::from("new.txt"), + 1_000, + "n".to_string(), + None, + None, + ); + let removed = l.prune(1_000 + 10, 100); + assert_eq!(removed, 1); + assert!(l.get(Path::new("old.txt")).is_none()); + assert!(l.get(Path::new("new.txt")).is_some()); +} + +#[test] +fn remove_clears_tombstone() { + let mut l = DeletionLedger::new(); + l.record(PathBuf::from("a.txt"), 5, "n".to_string(), None, None); + assert!(l.remove(Path::new("a.txt")).is_some()); + assert!(!l.has_newer_than(Path::new("a.txt"), 0)); +} + +#[test] +fn save_load_roundtrip() { + let dir = tempfile::tempdir().unwrap(); + let mut l = DeletionLedger::new(); + l.record( + PathBuf::from("a.txt"), + 123, + "node-x".to_string(), + Some("abc".to_string()), + Some(100), + ); + l.save_atomic(dir.path()).unwrap(); + let loaded = DeletionLedger::load(dir.path()); + let t = loaded.get(Path::new("a.txt")).unwrap(); + assert_eq!(t.deleted_at_ms, 123); + assert_eq!(t.deleter_node, "node-x"); + assert_eq!(t.prev_hash.as_deref(), Some("abc")); + assert_eq!(t.prev_mtime_ms, Some(100)); +} + +#[test] +fn load_missing_is_empty() { + let dir = tempfile::tempdir().unwrap(); + let l = DeletionLedger::load(dir.path()); + assert!(l.is_empty()); +} + +#[test] +fn load_corrupt_backs_up_and_returns_empty() { + let dir = tempfile::tempdir().unwrap(); + let bh = dir.path().join(LEDGER_DIR_NAME); + std::fs::create_dir_all(&bh).unwrap(); + std::fs::write(bh.join(LEDGER_FILE_NAME), "{ not json").unwrap(); + let l = DeletionLedger::load(dir.path()); + assert!(l.is_empty()); + let backups: Vec<_> = std::fs::read_dir(&bh) + .unwrap() + .filter_map(|e| e.ok()) + .filter(|e| { + e.file_name() + .to_string_lossy() + .starts_with("deletion-ledger.corrupt-") + }) + .collect(); + assert_eq!(backups.len(), 1); +} + +#[test] +fn now_ms_is_sane() { + // Smoke-cover the clock helper (also guards against future removal). + let t = now_ms(); + assert!(t > 1_700_000_000_000, "wall clock must be post-2023"); +}