diff --git a/Cargo.lock b/Cargo.lock index 8d98d7c..cdcd56c 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -515,6 +515,8 @@ dependencies = [ "cfg-if", "constant_time_eq 0.4.2", "cpufeatures 0.3.0", + "memmap2", + "rayon-core", ] [[package]] @@ -638,8 +640,10 @@ dependencies = [ "iced", "image", "inotify", + "jwalk", "log", "lz4_flex", + "num_cpus", "parking_lot", "rayon", "rcgen", @@ -1117,6 +1121,19 @@ dependencies = [ "cfg-if", ] +[[package]] +name = "crossbeam" +version = "0.8.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e71406cd8807725f7ac2f999a4cdd32e98f829fdf65f528343cebf945e41df1e" +dependencies = [ + "crossbeam-channel", + "crossbeam-deque", + "crossbeam-epoch", + "crossbeam-queue", + "crossbeam-utils", +] + [[package]] name = "crossbeam-channel" version = "0.5.15" @@ -1145,6 +1162,15 @@ dependencies = [ "crossbeam-utils", ] +[[package]] +name = "crossbeam-queue" +version = "0.3.14" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "03e8bd762f7479489c70ed6c768ddca99d7296857de437a68dcb2a94365b3fae" +dependencies = [ + "crossbeam-utils", +] + [[package]] name = "crossbeam-utils" version = "0.8.21" @@ -2901,6 +2927,16 @@ dependencies = [ "wasm-bindgen", ] +[[package]] +name = "jwalk" +version = "0.8.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2735847566356cd2179a2a38264839308f7079fa96e6bd5a42d740460e003c56" +dependencies = [ + "crossbeam", + "rayon", +] + [[package]] name = "keyboard-types" version = "0.7.0" @@ -3359,6 +3395,16 @@ dependencies = [ "libm", ] +[[package]] +name = "num_cpus" +version = "1.17.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "91df4bbde75afed763b708b7eee1e8e7651e02d97f6d5dd763e89367e957b23b" +dependencies = [ + "hermit-abi", + "libc", +] + [[package]] name = "num_enum" version = "0.7.6" diff --git a/crates/filesync/Cargo.toml b/crates/filesync/Cargo.toml index 09ad272..08cf95a 100644 --- a/crates/filesync/Cargo.toml +++ b/crates/filesync/Cargo.toml @@ -19,7 +19,9 @@ crossbeam-channel = { workspace = true } parking_lot = { workspace = true } log = { workspace = true } bincode = "1" -blake3 = "1" +blake3 = { version = "1", features = ["mmap", "rayon"] } +jwalk = "0.8" +num_cpus = "1" lz4_flex = "0.11" rayon = "1" walkdir = "2" diff --git a/crates/filesync/src/bundler.rs b/crates/filesync/src/bundler.rs index 02d28db..2d45ca5 100644 --- a/crates/filesync/src/bundler.rs +++ b/crates/filesync/src/bundler.rs @@ -153,18 +153,11 @@ fn stream_large_file( ); let hash_start = std::time::Instant::now(); + // Unified tiered digest (mmap for big files, streaming otherwise); + // the chunk protocol below is unchanged. let final_hash: [u8; 32] = { - let mut hasher = blake3::Hasher::new(); - let mut f = std::fs::File::open(&full)?; - let mut buf = vec![0u8; 64 * 1024]; - loop { - let n = f.read(&mut buf)?; - if n == 0 { - break; - } - hasher.update(&buf[..n]); - } - hasher.finalize().into() + let (_, h) = crate::manifest::hash_file_tiered(&full)?; + h }; debug!( "bundler: large file {:?} hash computed in {}ms", diff --git a/crates/filesync/src/lib.rs b/crates/filesync/src/lib.rs index 46defb8..4622f14 100644 --- a/crates/filesync/src/lib.rs +++ b/crates/filesync/src/lib.rs @@ -7,6 +7,7 @@ pub mod gui; pub mod known_hosts; pub mod ledger; pub mod manifest; +pub mod manifest_cache; pub mod protocol; pub mod server; pub mod suspend_detector; @@ -17,6 +18,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 use manifest_cache::{CachedEntry, ManifestCache, CACHE_MIN_AGE_MS}; pub fn timestamp_id() -> u64 { use std::time::{SystemTime, UNIX_EPOCH}; diff --git a/crates/filesync/src/manifest.rs b/crates/filesync/src/manifest.rs index 7df40d1..92b7cee 100644 --- a/crates/filesync/src/manifest.rs +++ b/crates/filesync/src/manifest.rs @@ -1,106 +1,327 @@ use crate::exclusions::Exclusions; use crate::ledger::DeletionLedger; +use crate::manifest_cache::{inode_of, mtime_parts, CachedEntry, ManifestCache}; use crate::protocol::{FileMetadata, Manifest, HASH_THREADS}; -use rayon::prelude::*; -use std::io::Read; +use crossbeam_channel::bounded; +use std::collections::{HashMap, HashSet}; +use std::io; use std::path::{Path, PathBuf}; -use std::time::SystemTime; -use walkdir::WalkDir; +use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering}; +use std::time::{Instant, SystemTime}; -fn hash_file_streaming(path: &Path) -> io::Result<(u64, [u8; 32])> { - let file = std::fs::File::open(path)?; - let mut reader = std::io::BufReader::with_capacity(64 * 1024, file); - let mut hasher = blake3::Hasher::new(); - let mut buf = [0u8; 64 * 1024]; - let mut size: u64 = 0; - loop { - let n = reader.read(&mut buf)?; - if n == 0 { - break; +/// Files at or above this size try `update_mmap_rayon` first (with streaming +/// fallback); everything else streams with a 256 KiB buffer. +const TIER_MMAP_MIN_BYTES: u64 = 4 * 1024 * 1024; +/// Streaming read buffer (also used for files under 128 KiB — a single read). +const HASH_BUF_BYTES: usize = 256 * 1024; +/// In-flight paths between the walker and the hash pool (backpressure). +const WALK_CHANNEL_DEPTH: usize = 1024; +/// Traversal threads for jwalk's isolated pool (see below). +const TRAVERSAL_THREADS: usize = 4; + +/// Tiered blake3 digest shared by scan / conflict-check / large-file paths. +/// Returns `(size, hash)`. Never switches algorithms (blake3 only). +pub fn hash_file_tiered(path: &Path) -> io::Result<(u64, [u8; 32])> { + let size = std::fs::metadata(path)?.len(); + if size >= TIER_MMAP_MIN_BYTES { + let mut hasher = blake3::Hasher::new(); + match hasher.update_mmap_rayon(path) { + Ok(_) => return Ok((size, hasher.finalize().into())), + Err(e) => { + log::debug!( + "hash_file_tiered: mmap failed for {path:?}, streaming fallback: {e}" + ); + } } - hasher.update(&buf[..n]); - size += n as u64; } + let file = std::fs::File::open(path)?; + let mut hasher = blake3::Hasher::new(); + hasher.update_reader(std::io::BufReader::with_capacity(HASH_BUF_BYTES, file))?; Ok((size, hasher.finalize().into())) } -use std::io; +fn parse_hex32(s: &str) -> Result<[u8; 32], ()> { + if s.len() != 64 { + return Err(()); + } + let mut out = [0u8; 32]; + for (i, chunk) in s.as_bytes().chunks(2).enumerate() { + let hi = (chunk[0] as char).to_digit(16).ok_or(())?; + let lo = (chunk[1] as char).to_digit(16).ok_or(())?; + out[i] = ((hi << 4) | lo) as u8; + } + Ok(out) +} + +/// Per-scan metrics. `hash_ms` is summed cpu time across hash threads; +/// `walk_ms` is walker wall time (the two overlap by design). +#[derive(Debug, Clone, Default)] +pub struct ScanStats { + pub walk_ms: u64, + pub hash_ms: u64, + pub files: usize, + pub dirs: usize, + pub bytes: u64, + pub cache_hits: usize, + pub cache_misses: usize, + pub hash_mb_s: f64, +} + +struct BuildOutcome { + rel: PathBuf, + meta: FileMetadata, + cache_upsert: Option<(PathBuf, CachedEntry)>, +} -pub fn build_manifest(root: &Path, node_id: &str, exclusions: &Exclusions) -> io::Result { - if HASH_THREADS > 0 { - let _ = rayon::ThreadPoolBuilder::new() - .num_threads(HASH_THREADS) - .build_global(); +#[allow(clippy::too_many_arguments)] +fn hash_one( + root: &Path, + rel: &PathBuf, + cache: Option<&ManifestCache>, + dirty: &HashSet, + hash_nanos: &AtomicU64, + hash_bytes: &AtomicU64, + hits: &AtomicUsize, + misses: &AtomicUsize, +) -> Option { + let full = root.join(rel); + let meta = std::fs::metadata(&full).ok()?; + let modified_ms = meta + .modified() + .ok()? + .duration_since(SystemTime::UNIX_EPOCH) + .ok()? + .as_millis() as u64; + + if meta.is_dir() { + return Some(BuildOutcome { + rel: rel.clone(), + meta: FileMetadata { + rel_path: rel.clone(), + size: 0, + hash: [0u8; 32], + modified_ms, + change_sequence: 0, + is_dir: true, + }, + cache_upsert: None, + }); } - let paths: Vec = WalkDir::new(root) - .into_iter() - .filter_map(Result::ok) - .filter_map(|e| { - let p = e.into_path(); - let rel = p.strip_prefix(root).ok()?; - if rel.as_os_str().is_empty() { - return None; + // Fast path: trusted cache hit (stable + old + identity match). + if !dirty.contains(rel) { + if let Some(cache) = cache { + if let Some((mtime_s, mtime_ns)) = mtime_parts(&meta) { + let now = crate::ledger::now_ms(); + if let Some(hex) = cache.lookup( + rel, + meta.len(), + mtime_s, + mtime_ns, + inode_of(&meta), + now, + ) { + if let Ok(hash) = parse_hex32(&hex) { + hits.fetch_add(1, Ordering::Relaxed); + log::trace!("manifest: cache hit for {rel:?}"); + return Some(BuildOutcome { + rel: rel.clone(), + meta: FileMetadata { + rel_path: rel.clone(), + size: meta.len(), + hash, + modified_ms, + change_sequence: 0, + is_dir: false, + }, + cache_upsert: None, + }); + } + // Corrupt cache payload: fall through and re-hash + // (the upsert below self-heals the entry). + } } + } + } - if exclusions.is_excluded(rel) { + misses.fetch_add(1, Ordering::Relaxed); + let t = Instant::now(); + let (size, hash) = hash_file_tiered(&full).ok()?; + hash_nanos.fetch_add(t.elapsed().as_nanos() as u64, Ordering::Relaxed); + hash_bytes.fetch_add(size, Ordering::Relaxed); + + let cache_upsert = mtime_parts(&meta).map(|(mtime_s, mtime_ns)| { + ( + rel.clone(), + CachedEntry { + size, + mtime_s, + mtime_ns, + inode: inode_of(&meta), + hash: crate::hex(&hash), + }, + ) + }); + + Some(BuildOutcome { + rel: rel.clone(), + meta: FileMetadata { + rel_path: rel.clone(), + size, + hash, + modified_ms, + change_sequence: 0, + is_dir: false, + }, + cache_upsert, + }) +} + +/// Build a manifest with a pipelined parallel scan: jwalk traverses on the +/// calling thread and streams relative paths through a bounded channel to a +/// dedicated rayon hash pool (`max(num_cpus, HASH_THREADS)` threads) — +/// hashing overlaps walking, no full path vec is collected first. +/// +/// `cache` enables the hash cache (see [`ManifestCache`]); `dirty` lists +/// watcher-pending paths that bypass the cache (still re-cached after +/// hashing). Returns `(manifest, updated_cache, stats)`; the caller owns +/// pruning/saving the cache. +pub fn build_manifest( + root: &Path, + node_id: &str, + exclusions: &Exclusions, + cache: Option<&ManifestCache>, + dirty: &HashSet, +) -> io::Result<(Manifest, ManifestCache, ScanStats)> { + let threads = num_cpus::get().max(HASH_THREADS); + let pool = rayon::ThreadPoolBuilder::new() + .num_threads(threads) + .thread_name(|i| format!("hash-{i}")) + .build() + .map_err(io::Error::other)?; + + let (work_tx, work_rx) = bounded::(WALK_CHANNEL_DEPTH); + let (res_tx, res_rx) = crossbeam_channel::unbounded::(); + let walk_ms = AtomicU64::new(0); + let hash_nanos = AtomicU64::new(0); + let hash_bytes = AtomicU64::new(0); + let hits = AtomicUsize::new(0); + let misses = AtomicUsize::new(0); + + let root_p = root.to_path_buf(); + // Shared metric refs (Copy) for the worker closures below. + let hash_nanos_r = &hash_nanos; + let hash_bytes_r = &hash_bytes; + let hits_r = &hits; + let misses_r = &misses; + pool.scope(|s| { + for _ in 0..threads { + let rx = work_rx.clone(); + let tx = res_tx.clone(); + let root_c = root_p.clone(); + s.spawn(move |_| { + for rel in rx { + if let Some(outcome) = hash_one( + &root_c, + &rel, + cache, + dirty, + hash_nanos_r, + hash_bytes_r, + hits_r, + misses_r, + ) { + if tx.send(outcome).is_err() { + break; + } + } + } + }); + } + // The scope thread walks: traversal overlaps hashing via backpressure. + // NOTE: skip_hidden(false) is required — walkdir parity (hidden files + // and .bh_filesync/ are visited, then filtered by exclusions below). + // NOTE: isolated RayonNewPool, NOT the rayon global pool — jwalk's + // default global-pool handshake (1s busy-timeout, silent abort) can + // misfire when many scans share a process, silently degrading to an + // EMPTY walk (empty manifest, data looks deleted). Sharing our hash + // pool would deadlock instead (all hash threads parked on recv while + // traversal waits for a thread). Verified by failing→passing repro; + // do not change this without a multi-scan-in-one-process test. + let wstart = Instant::now(); + let walk_opts = jwalk::WalkDir::new(&root_p) + .skip_hidden(false) + .parallelism(jwalk::Parallelism::RayonNewPool(TRAVERSAL_THREADS)); + for entry in walk_opts { + let entry = match entry { + Ok(e) => e, + Err(_) => continue, + }; + let full = entry.path(); + let rel = match full.strip_prefix(&root_p) { + Ok(r) if !r.as_os_str().is_empty() => r.to_path_buf(), + _ => continue, + }; + if exclusions.is_excluded(&rel) { log::debug!( "manifest: skipping {:?} (rule: {:?})", rel, - exclusions.matching_rule(rel) + exclusions.matching_rule(&rel) ); - return None; + continue; } - Some(p) - }) - .collect(); - - let entries: Vec<(PathBuf, FileMetadata)> = paths - .par_iter() - .filter_map(|full| { - let rel = full.strip_prefix(root).ok()?.to_path_buf(); - let meta = std::fs::metadata(full).ok()?; - let modified_ms = meta - .modified() - .ok()? - .duration_since(SystemTime::UNIX_EPOCH) - .ok()? - .as_millis() as u64; - - if meta.is_dir() { - return Some(( - rel.clone(), - FileMetadata { - rel_path: rel, - size: 0, - hash: [0u8; 32], - modified_ms, - change_sequence: 0, - is_dir: true, - }, - )); + if work_tx.send(rel).is_err() { + break; } + } + walk_ms.store(wstart.elapsed().as_millis() as u64, Ordering::Relaxed); + drop(work_tx); + }); + drop(work_rx); + drop(res_tx); - let (size, hash) = hash_file_streaming(full).ok()?; - - Some(( - rel.clone(), - FileMetadata { - rel_path: rel, - size, - hash, - modified_ms, - change_sequence: 0, - is_dir: false, - }, - )) - }) - .collect(); - - Ok(Manifest { - files: entries.into_iter().collect(), - node_id: node_id.to_string(), - }) + let mut files_map: HashMap = HashMap::new(); + let mut updated = cache.cloned().unwrap_or_default(); + let mut dirs = 0usize; + for outcome in res_rx { + if outcome.meta.is_dir { + dirs += 1; + } + if let Some((p, e)) = outcome.cache_upsert { + updated.upsert(p, e); + } + files_map.insert(outcome.rel, outcome.meta); + } + let files = files_map.values().filter(|m| !m.is_dir).count(); + let bytes: u64 = files_map.values().map(|m| m.size).sum(); + let existing: HashSet = files_map.keys().cloned().collect(); + updated.prune_missing(&existing); + + let nanos = hash_nanos.load(Ordering::Relaxed); + let hbytes = hash_bytes.load(Ordering::Relaxed); + let stats = ScanStats { + walk_ms: walk_ms.load(Ordering::Relaxed), + hash_ms: (nanos / 1_000_000) as u64, + files, + dirs, + bytes, + cache_hits: hits.load(Ordering::Relaxed), + cache_misses: misses.load(Ordering::Relaxed), + hash_mb_s: if nanos > 0 { + hbytes as f64 / (nanos as f64 / 1e9) / 1e6 + } else { + 0.0 + }, + }; + + Ok(( + Manifest { + files: files_map, + node_id: node_id.to_string(), + }, + updated, + stats, + )) } pub fn compute_send_list(local: &Manifest, remote: &Manifest, is_server: bool) -> Vec { diff --git a/crates/filesync/src/manifest_cache.rs b/crates/filesync/src/manifest_cache.rs new file mode 100644 index 0000000..9cf184e --- /dev/null +++ b/crates/filesync/src/manifest_cache.rs @@ -0,0 +1,174 @@ +//! Persistent manifest hash cache (initial-scan speedup). +//! +//! Template: [`crate::ledger`] — state lives in +//! `/.bh_filesync/manifest-cache.json`, loads tolerate missing/corrupt +//! files (corrupt JSON is backed up, never panics), saves are atomic +//! (`.tmp` + rename) and skipped when nothing changed. +//! +//! Trust rule: a cached hash is reused only when size, mtime (secs + nanos) +//! and inode ALL match the current stat AND the file is older than +//! [`CACHE_MIN_AGE_MS`] (stability gate — a recently modified file may still +//! be mid-write, so it is always re-hashed). + +use serde::{Deserialize, Serialize}; +use std::collections::{HashMap, HashSet}; +use std::io; +use std::path::{Path, PathBuf}; +use std::time::UNIX_EPOCH; + +/// Directory (under the sync root) holding filesync state. +pub const CACHE_DIR_NAME: &str = ".bh_filesync"; +/// File name of the persistent manifest hash cache. +pub const CACHE_FILE_NAME: &str = "manifest-cache.json"; +/// Minimum file age before a cache hit is trusted (stability gate). +pub const CACHE_MIN_AGE_MS: u64 = 5_000; + +/// Hash + identity snapshot for one relative path. `hash` is lowercase hex. +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct CachedEntry { + pub size: u64, + pub mtime_s: i64, + pub mtime_ns: u32, + pub inode: u64, + pub hash: String, +} + +/// Persistent path → [`CachedEntry`] map with dirty tracking so warm scans +/// skip disk writes entirely. +#[derive(Debug, Clone, Default, Serialize, Deserialize)] +pub struct ManifestCache { + #[serde(default)] + pub entries: HashMap, + #[serde(skip)] + dirty: bool, +} + +fn cache_path(root: &Path) -> PathBuf { + root.join(CACHE_DIR_NAME).join(CACHE_FILE_NAME) +} + +/// Wall-clock milliseconds since the Unix epoch (re-export of the ledger +/// clock so all timestamps share one definition). +pub fn now_ms() -> u64 { + crate::ledger::now_ms() +} + +/// Split `Metadata::modified()` into `(unix secs, subsec nanos)`. +/// Returns `None` for unavailable/pre-epoch mtimes (caller re-hashes). +pub fn mtime_parts(meta: &std::fs::Metadata) -> Option<(i64, u32)> { + let d = meta.modified().ok()?.duration_since(UNIX_EPOCH).ok()?; + let secs: i64 = d.as_secs().try_into().ok()?; + Some((secs, d.subsec_nanos())) +} + +/// Inode for cache identity (unix; the crate already depends on inotify). +pub fn inode_of(meta: &std::fs::Metadata) -> u64 { + use std::os::unix::fs::MetadataExt; + meta.ino() +} + +impl ManifestCache { + pub fn new() -> Self { + Self { + entries: HashMap::new(), + dirty: false, + } + } + + /// Load the cache from `/.bh_filesync/manifest-cache.json`. + /// Missing file => empty cache. Corrupt JSON => back the file up to + /// `manifest-cache.corrupt-.json` (best effort) and return empty. + /// Never panics. + pub fn load(root: &Path) -> Self { + let path = cache_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(), + }; + if let Ok(cache) = serde_json::from_str::(&raw) { + return cache; + } + let backup = path.with_extension(format!("corrupt-{}.json", now_ms())); + let _ = std::fs::copy(&path, &backup); + Self::new() + } + + /// Persist atomically (write `.tmp` then rename), creating + /// `/.bh_filesync` if needed. Skips the write entirely when + /// nothing changed since load/last save. + pub fn save_atomic(&mut self, root: &Path) -> io::Result<()> { + if !self.dirty { + return Ok(()); + } + let dir = root.join(CACHE_DIR_NAME); + std::fs::create_dir_all(&dir)?; + let dest = dir.join(CACHE_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)?; + self.dirty = false; + Ok(()) + } + + /// Whether the cache changed since load/last save. + pub fn is_dirty(&self) -> bool { + self.dirty + } + + /// Look up a cached hash. Returns `Some(hex)` only when size, mtime + /// (secs + nanos) and inode all match AND the file is older than + /// [`CACHE_MIN_AGE_MS`] as of `now_ms`. + pub fn lookup( + &self, + path: &Path, + size: u64, + mtime_s: i64, + mtime_ns: u32, + inode: u64, + now_ms: u64, + ) -> Option { + let e = self.entries.get(path)?; + if e.size != size || e.mtime_s != mtime_s || e.mtime_ns != mtime_ns || e.inode != inode { + return None; + } + let mtime_ms = (mtime_s.max(0) as u64) + .saturating_mul(1_000) + .saturating_add(mtime_ns as u64 / 1_000_000); + if now_ms.saturating_sub(mtime_ms) <= CACHE_MIN_AGE_MS { + return None; // too fresh — stability gate forces a re-hash + } + Some(e.hash.clone()) + } + + /// Insert/replace an entry. Marks dirty only on actual change so + /// fully-warm scans skip the save. + pub fn upsert(&mut self, path: PathBuf, entry: CachedEntry) { + if self.entries.get(&path) != Some(&entry) { + self.entries.insert(path, entry); + self.dirty = true; + } + } + + /// Drop entries for paths no longer present. Returns the number removed + /// (marks dirty only when something was removed). + pub fn prune_missing(&mut self, existing: &HashSet) -> usize { + let before = self.entries.len(); + self.entries.retain(|p, _| existing.contains(p)); + let removed = before - self.entries.len(); + if removed > 0 { + self.dirty = true; + } + removed + } + + pub fn len(&self) -> usize { + self.entries.len() + } + + pub fn is_empty(&self) -> bool { + self.entries.is_empty() + } +} diff --git a/crates/filesync/src/sync_engine.rs b/crates/filesync/src/sync_engine.rs index fbbca69..acfb0b0 100644 --- a/crates/filesync/src/sync_engine.rs +++ b/crates/filesync/src/sync_engine.rs @@ -2,6 +2,7 @@ use crate::bundler; use crate::exclusions::Exclusions; use crate::ledger::{now_ms, DeletionLedger, Tombstone, DELETION_LEDGER_TTL_MS}; use crate::manifest; +use crate::manifest_cache::ManifestCache; use crate::protocol::FILE_STABILITY_MS; use crate::protocol::*; use crate::transport::Connection; @@ -130,19 +131,9 @@ pub fn conflict_copy_name(rel_path: &Path, node_id: &str, unix_secs: u64) -> Pat } fn hash_file(path: &Path) -> Option<[u8; 32]> { - use std::io::Read; - let file = fs::File::open(path).ok()?; - let mut reader = std::io::BufReader::with_capacity(64 * 1024, file); - let mut hasher = blake3::Hasher::new(); - let mut buf = [0u8; 64 * 1024]; - loop { - let n = reader.read(&mut buf).ok()?; - if n == 0 { - break; - } - hasher.update(&buf[..n]); - } - Some(hasher.finalize().into()) + crate::manifest::hash_file_tiered(path) + .ok() + .map(|(_, hash)| hash) } fn cleanup_transfer_dirs(root: &Path, tmp_path: &Path) { @@ -342,7 +333,38 @@ impl SyncEngine { } pub fn scan(&self) -> std::io::Result { - let m = manifest::build_manifest(&self.root, &self.node_id, &self.exclusions)?; + self.scan_with_dirty(&HashSet::new()) + } + + /// Full scan with a dirty-path set that bypasses the manifest hash cache + /// (e.g. watcher-pending paths). Drains nothing; defaults to empty via + /// [`SyncEngine::scan`]. + pub fn scan_with_dirty(&self, dirty: &HashSet) -> std::io::Result { + let total = std::time::Instant::now(); + let cache = ManifestCache::load(&self.root); + let (m, mut updated, stats) = manifest::build_manifest( + &self.root, + &self.node_id, + &self.exclusions, + Some(&cache), + dirty, + )?; + if let Err(e) = updated.save_atomic(&self.root) { + log::warn!("filesync: manifest cache save failed: {e}"); + } + log::info!( + "filesync: scan complete — {} file(s), {} dir(s), {} byte(s) in {}ms \ + (walk {}ms, hash {}ms cpu, {:.1} MB/s, cache {} hit/{} miss)", + stats.files, + stats.dirs, + stats.bytes, + total.elapsed().as_millis(), + stats.walk_ms, + stats.hash_ms, + stats.hash_mb_s, + stats.cache_hits, + stats.cache_misses, + ); *self.manifest.write() = m.clone(); // Cheap upkeep (persisted only when entries actually expire). self.prune_ledger(); diff --git a/crates/filesync/tests/test_manifest.rs b/crates/filesync/tests/test_manifest.rs index 4e2d0c3..c8500ad 100644 --- a/crates/filesync/tests/test_manifest.rs +++ b/crates/filesync/tests/test_manifest.rs @@ -7,6 +7,12 @@ use bytehive_filesync::{ manifest::{build_manifest, compute_send_list, filter_resurrected}, protocol::{FileMetadata, Manifest}, }; +use std::collections::HashSet; + +/// Uncached `build_manifest` wrapper returning just the manifest. +fn build(dir: &PathBuf, node: &str, excl: &Exclusions) -> std::io::Result { + build_manifest(dir, node, excl, None, &HashSet::new()).map(|t| t.0) +} fn no_exclusions() -> Exclusions { Exclusions::compile(&ExclusionConfig::default()) @@ -49,7 +55,7 @@ fn make_manifest(entries: &[(&str, u64, [u8; 32], bool, u64)], node: &str) -> Ma #[test] fn build_manifest_empty_directory() { let dir = tmp_dir("empty"); - let m = build_manifest(&dir, "node-1", &no_exclusions()).unwrap(); + let m = build(&dir, "node-1", &no_exclusions()).unwrap(); assert_eq!(m.node_id, "node-1"); assert!(m.files.is_empty(), "empty root must produce empty manifest"); std::fs::remove_dir_all(&dir).unwrap(); @@ -61,7 +67,7 @@ fn build_manifest_single_file_hash_and_size() { let content = b"hello filesync"; std::fs::write(dir.join("hello.txt"), content).unwrap(); - let m = build_manifest(&dir, "n", &no_exclusions()).unwrap(); + let m = build(&dir, "n", &no_exclusions()).unwrap(); assert_eq!(m.files.len(), 1); let meta = m.files.get(&PathBuf::from("hello.txt")).unwrap(); let expected: [u8; 32] = blake3::hash(content).into(); @@ -76,7 +82,7 @@ fn build_manifest_directory_entry_has_zero_hash() { let dir = tmp_dir("direntry"); std::fs::create_dir_all(dir.join("subdir")).unwrap(); - let m = build_manifest(&dir, "n", &no_exclusions()).unwrap(); + let m = build(&dir, "n", &no_exclusions()).unwrap(); let meta = m.files.get(&PathBuf::from("subdir")).unwrap(); assert!(meta.is_dir); assert_eq!(meta.size, 0); @@ -91,7 +97,7 @@ fn build_manifest_nested_files_and_dirs() { std::fs::write(dir.join("a/b/deep.txt"), b"deep").unwrap(); std::fs::write(dir.join("root.txt"), b"root").unwrap(); - let m = build_manifest(&dir, "n", &no_exclusions()).unwrap(); + let m = build(&dir, "n", &no_exclusions()).unwrap(); assert!(m.files.contains_key(&PathBuf::from("a"))); assert!(m.files.contains_key(&PathBuf::from("a/b"))); assert!(m.files.contains_key(&PathBuf::from("a/b/deep.txt"))); @@ -109,7 +115,7 @@ fn build_manifest_respects_glob_exclusion() { exclude_patterns: vec!["*.log".to_string()], exclude_regex: vec![], }); - let m = build_manifest(&dir, "n", &excl).unwrap(); + let m = build(&dir, "n", &excl).unwrap(); assert!(m.files.contains_key(&PathBuf::from("keep.txt"))); assert!(!m.files.contains_key(&PathBuf::from("skip.log"))); std::fs::remove_dir_all(&dir).unwrap(); @@ -125,7 +131,7 @@ fn build_manifest_respects_regex_exclusion() { exclude_patterns: vec![], exclude_regex: vec![r".*\.tmp$".to_string()], }); - let m = build_manifest(&dir, "n", &excl).unwrap(); + let m = build(&dir, "n", &excl).unwrap(); assert!(!m.files.contains_key(&PathBuf::from("file.tmp"))); assert!(m.files.contains_key(&PathBuf::from("file.rs"))); std::fs::remove_dir_all(&dir).unwrap(); @@ -134,7 +140,7 @@ fn build_manifest_respects_regex_exclusion() { #[test] fn build_manifest_node_id_is_preserved() { let dir = tmp_dir("nodeid"); - let m = build_manifest(&dir, "my-special-node", &no_exclusions()).unwrap(); + let m = build(&dir, "my-special-node", &no_exclusions()).unwrap(); assert_eq!(m.node_id, "my-special-node"); std::fs::remove_dir_all(&dir).unwrap(); } @@ -219,7 +225,7 @@ fn build_manifest_excludes_filesync_tmp_dir() { std::fs::create_dir_all(dir.join(".bh_filesync/transfers")).unwrap(); std::fs::write(dir.join(".bh_filesync/transfers/partial.tmp"), b"temp").unwrap(); std::fs::write(dir.join("real.txt"), b"real").unwrap(); - let m = build_manifest(&dir, "n", &no_exclusions()).unwrap(); + let m = build(&dir, "n", &no_exclusions()).unwrap(); assert!(m.files.contains_key(&PathBuf::from("real.txt"))); assert!( !m.files.contains_key(&PathBuf::from(".bh_filesync")), @@ -238,7 +244,7 @@ fn different_content_produces_different_hashes() { let dir = tmp_dir("diff_hash"); std::fs::write(dir.join("a.bin"), b"content A").unwrap(); std::fs::write(dir.join("b.bin"), b"content B").unwrap(); - let m = build_manifest(&dir, "n", &no_exclusions()).unwrap(); + let m = build(&dir, "n", &no_exclusions()).unwrap(); let ha = m.files.get(&PathBuf::from("a.bin")).unwrap().hash; let hb = m.files.get(&PathBuf::from("b.bin")).unwrap().hash; assert_ne!( diff --git a/crates/filesync/tests/test_manifest_cache.rs b/crates/filesync/tests/test_manifest_cache.rs new file mode 100644 index 0000000..d4d0f84 --- /dev/null +++ b/crates/filesync/tests/test_manifest_cache.rs @@ -0,0 +1,341 @@ +//! Regression tests for the manifest hash cache (initial-scan speedup). +//! Fast, no network: tempdirs + `SyncEngine::scan` + direct cache API. +//! Policy: zero inline tests in `src/` — all coverage lives here. + +use std::collections::HashSet; +use std::os::unix::fs::MetadataExt; +use std::path::{Path, PathBuf}; +use std::sync::Arc; + +use bytehive_filesync::{ + exclusions::{ExclusionConfig, Exclusions}, + hex, + manifest_cache::{ + inode_of, mtime_parts, now_ms, CachedEntry, ManifestCache, CACHE_MIN_AGE_MS, + }, + sync_engine::SyncEngine, +}; + +fn make_engine(root: PathBuf, node: &str) -> SyncEngine { + let ex = Arc::new(Exclusions::compile(&ExclusionConfig::default())); + SyncEngine::new(root, node.to_string(), ex) +} + +fn tmp_engine(label: &str, node: &str) -> (tempfile::TempDir, SyncEngine) { + let dir = tempfile::tempdir().unwrap(); + let _ = label; + let engine = make_engine(dir.path().to_path_buf(), node); + (dir, engine) +} + +/// Set mtime to an exact instant (controls the 5s stability gate). +fn set_mtime(path: &Path, ms: u64) { + let ft = filetime::FileTime::from_unix_time( + (ms / 1000) as i64, + ((ms % 1000) * 1_000_000) as u32, + ); + filetime::set_file_mtime(path, ft).unwrap(); +} + +fn manifest_hash(engine: &SyncEngine, rel: &str) -> [u8; 32] { + engine + .get_manifest() + .files + .get(&PathBuf::from(rel)) + .unwrap_or_else(|| panic!("{rel} missing from manifest")) + .hash +} + +#[test] +fn cache_roundtrip() { + let dir = tempfile::tempdir().unwrap(); + let mut c = ManifestCache::load(dir.path()); + assert!(c.is_empty()); + c.upsert( + PathBuf::from("a.txt"), + CachedEntry { + size: 12, + mtime_s: 1_700_000_000, + mtime_ns: 0, + inode: 4242, + hash: "ab".repeat(32), + }, + ); + assert!(c.is_dirty()); + c.save_atomic(dir.path()).unwrap(); + assert!(!c.is_dirty()); + assert!(dir + .path() + .join(".bh_filesync") + .join("manifest-cache.json") + .is_file()); + let loaded = ManifestCache::load(dir.path()); + assert_eq!(loaded.len(), 1); + assert_eq!(loaded.entries, c.entries); +} + +#[test] +fn stable_file_skips_rehash() { + // Same size + same mtime + same inode + older than 5s: the second scan + // must reuse the cached hash WITHOUT reading the file. Proved by + // rewriting different content under an identical stat fingerprint. + let (dir, engine) = tmp_engine("skip", "n"); + let f = dir.path().join("f.txt"); + std::fs::write(&f, b"content-one!").unwrap(); + let old_ms = now_ms() - 10_000; + set_mtime(&f, old_ms); + engine.scan().unwrap(); + let h1 = manifest_hash(&engine, "f.txt"); + assert_eq!(blake3::hash(b"content-one!").to_hex().as_str(), hex(&h1)); + + std::fs::write(&f, b"content-two!").unwrap(); // same 12 B, new inode? no — same file + set_mtime(&f, old_ms); // restore identical stat fingerprint + engine.scan().unwrap(); + let h2 = manifest_hash(&engine, "f.txt"); + assert_eq!(h1, h2, "stable stat fingerprint must reuse the cached hash"); + assert_ne!( + blake3::hash(b"content-two!").to_hex().as_str(), + hex(&h2), + "test is meaningful only if content really changed" + ); +} + +#[test] +fn mtime_bump_rehashes() { + let (dir, engine) = tmp_engine("mtime", "n"); + let f = dir.path().join("f.txt"); + std::fs::write(&f, b"content-one!").unwrap(); + set_mtime(&f, now_ms() - 30_000); + engine.scan().unwrap(); + let h1 = manifest_hash(&engine, "f.txt"); + + // Same size, different content, DIFFERENT (but still old) mtime. + std::fs::write(&f, b"content-two!").unwrap(); + set_mtime(&f, now_ms() - 20_000); + engine.scan().unwrap(); + let h2 = manifest_hash(&engine, "f.txt"); + assert_ne!(h1, h2); + assert_eq!(blake3::hash(b"content-two!").to_hex().as_str(), hex(&h2)); +} + +#[test] +fn recreate_new_inode_rehashes() { + let (dir, engine) = tmp_engine("inode", "n"); + let f = dir.path().join("f.txt"); + let old_ms = now_ms() - 30_000; + std::fs::write(&f, b"content-one!").unwrap(); + set_mtime(&f, old_ms); + engine.scan().unwrap(); + let h1 = manifest_hash(&engine, "f.txt"); + + // Delete + recreate: new inode, same size, same mtime, new content. + std::fs::remove_file(&f).unwrap(); + std::fs::write(&f, b"content-two!").unwrap(); + set_mtime(&f, old_ms); + engine.scan().unwrap(); + let h2 = manifest_hash(&engine, "f.txt"); + assert_ne!(h1, h2, "inode change must force a re-hash"); + assert_eq!(blake3::hash(b"content-two!").to_hex().as_str(), hex(&h2)); +} + +#[test] +fn fresh_file_age_gate_forces_rehash() { + // A prefabricated cache entry matching the live stat EXACTLY must still + // be ignored when the file is younger than 5s. + let (dir, engine) = tmp_engine("gate", "n"); + let f = dir.path().join("f.txt"); + std::fs::write(&f, b"fresh content here").unwrap(); + let meta = std::fs::metadata(&f).unwrap(); + let (ms, mns) = mtime_parts(&meta).unwrap(); + let age = now_ms().saturating_sub(ms as u64 * 1000 + mns as u64 / 1_000_000); + assert!(age <= CACHE_MIN_AGE_MS, "test file must be fresh"); + + let mut c = ManifestCache::load(dir.path()); + c.upsert( + PathBuf::from("f.txt"), + CachedEntry { + size: meta.len(), + mtime_s: ms, + mtime_ns: mns, + inode: inode_of(&meta), + hash: "ff".repeat(32), // bogus — must NOT surface + }, + ); + c.save_atomic(dir.path()).unwrap(); + + engine.scan().unwrap(); + let h = manifest_hash(&engine, "f.txt"); + assert_eq!( + blake3::hash(b"fresh content here").to_hex().as_str(), + hex(&h) + ); + assert_ne!(hex(&h), "ff".repeat(32)); +} + +#[test] +fn prune_missing_drops_gone_paths() { + let (dir, engine) = tmp_engine("prune", "n"); + std::fs::write(dir.path().join("a.txt"), b"a").unwrap(); + std::fs::write(dir.path().join("b.txt"), b"b").unwrap(); + // Old mtimes so both entries actually cache. + set_mtime(&dir.path().join("a.txt"), now_ms() - 30_000); + set_mtime(&dir.path().join("b.txt"), now_ms() - 30_000); + engine.scan().unwrap(); + assert_eq!(ManifestCache::load(dir.path()).len(), 2); + + std::fs::remove_file(dir.path().join("b.txt")).unwrap(); + engine.scan().unwrap(); + let loaded = ManifestCache::load(dir.path()); + assert!(loaded.entries.contains_key(&PathBuf::from("a.txt"))); + assert!(!loaded.entries.contains_key(&PathBuf::from("b.txt"))); +} + +#[test] +fn corrupt_cache_backs_up_no_panic() { + let dir = tempfile::tempdir().unwrap(); + let bh = dir.path().join(".bh_filesync"); + std::fs::create_dir_all(&bh).unwrap(); + std::fs::write(bh.join("manifest-cache.json"), "{ not json").unwrap(); + let c = ManifestCache::load(dir.path()); + assert!(c.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("manifest-cache.corrupt-") + }) + .collect(); + assert_eq!(backups.len(), 1); +} + +#[test] +fn warm_scan_returns_identical_manifest() { + let (dir, engine) = tmp_engine("warm", "n"); + std::fs::write(dir.path().join("a.txt"), b"aaa").unwrap(); + std::fs::create_dir_all(dir.path().join("sub")).unwrap(); + std::fs::write(dir.path().join("sub/b.txt"), b"bbb").unwrap(); + engine.scan().unwrap(); + let first = engine.get_manifest(); + engine.scan().unwrap(); + let second = engine.get_manifest(); + fn proj(m: &bytehive_filesync::protocol::Manifest) -> Vec<(PathBuf, u64, [u8; 32], u64, bool)> { + let mut v: Vec<_> = m + .files + .iter() + .map(|(p, f)| (p.clone(), f.size, f.hash, f.modified_ms, f.is_dir)) + .collect(); + v.sort(); + v + } + assert_eq!(proj(&first), proj(&second)); + assert_eq!(first.node_id, second.node_id); +} + +#[test] +fn rapid_rewrite_picks_up_new_hash() { + // Rewrite within 5s (fresh mtime): the stability gate must re-hash. + let (dir, engine) = tmp_engine("rapid", "n"); + let f = dir.path().join("f.txt"); + std::fs::write(&f, b"version one").unwrap(); + engine.scan().unwrap(); + let h1 = manifest_hash(&engine, "f.txt"); + + std::fs::write(&f, b"version two!!").unwrap(); + engine.scan().unwrap(); + let h2 = manifest_hash(&engine, "f.txt"); + assert_ne!(h1, h2); + assert_eq!(blake3::hash(b"version two!!").to_hex().as_str(), hex(&h2)); +} + +#[test] +fn unchanged_scan_skips_cache_save() { + let (dir, engine) = tmp_engine("skip-save", "n"); + std::fs::write(dir.path().join("a.txt"), b"aaa").unwrap(); + engine.scan().unwrap(); + let cache_file = dir + .path() + .join(".bh_filesync") + .join("manifest-cache.json"); + assert!(cache_file.is_file()); + let m1 = std::fs::metadata(&cache_file).unwrap().ino(); + let t1 = std::fs::metadata(&cache_file).unwrap().modified().unwrap(); + assert!(m1 > 0); + engine.scan().unwrap(); + let t2 = std::fs::metadata(&cache_file).unwrap().modified().unwrap(); + assert_eq!(t1, t2, "warm scan with no changes must skip the save"); +} + +#[test] +fn dirty_set_bypasses_cache() { + // A watcher-dirty path is re-hashed even with a matching stable entry. + let (dir, engine) = tmp_engine("dirty", "n"); + let f = dir.path().join("f.txt"); + std::fs::write(&f, b"content-one!").unwrap(); + let old_ms = now_ms() - 30_000; + set_mtime(&f, old_ms); + engine.scan().unwrap(); + + std::fs::write(&f, b"content-two!").unwrap(); + set_mtime(&f, old_ms); // identical fingerprint — would hit… + let mut dirty = HashSet::new(); + dirty.insert(PathBuf::from("f.txt")); + engine.scan_with_dirty(&dirty).unwrap(); // …but dirty bypasses + let h = manifest_hash(&engine, "f.txt"); + assert_eq!(blake3::hash(b"content-two!").to_hex().as_str(), hex(&h)); +} + +#[test] +fn cache_lookup_helpers_agree_with_fs() { + let dir = tempfile::tempdir().unwrap(); + let f = dir.path().join("f.txt"); + std::fs::write(&f, b"abc").unwrap(); + let meta = std::fs::metadata(&f).unwrap(); + let (ms, mns) = mtime_parts(&meta).expect("mtime must parse"); + assert!(ms > 0 && mns < 1_000_000_000); + assert_eq!(inode_of(&meta), meta.ino()); + + let mut c = ManifestCache::new(); + c.upsert( + PathBuf::from("f.txt"), + CachedEntry { + size: 3, + mtime_s: ms, + mtime_ns: mns, + inode: meta.ino(), + hash: "00".repeat(32), + }, + ); + // Fresh file → gate denies even a perfect match. + assert!(c + .lookup( + Path::new("f.txt"), + 3, + ms, + mns, + meta.ino(), + now_ms() + ) + .is_none()); + // Same entry, far-future clock → hit. + assert_eq!( + c.lookup( + Path::new("f.txt"), + 3, + ms, + mns, + meta.ino(), + (ms as u64) * 1000 + 60_000 + ) + .as_deref(), + Some(&"00".repeat(32) as &str) + ); + // Wrong size/inode → miss. + assert!(c + .lookup(Path::new("f.txt"), 4, ms, mns, meta.ino(), u64::MAX) + .is_none()); + assert!(c + .lookup(Path::new("f.txt"), 3, ms, mns, meta.ino() + 1, u64::MAX) + .is_none()); +}