diff --git a/crates/filesync/src/bundler.rs b/crates/filesync/src/bundler.rs index 9d248fc..02d28db 100644 --- a/crates/filesync/src/bundler.rs +++ b/crates/filesync/src/bundler.rs @@ -1,9 +1,11 @@ use crate::protocol::*; +use crate::sync_engine::SyncEngine; use crossbeam_channel::Sender; use log::{debug, warn}; use std::io::Read; use std::path::{Path, PathBuf}; use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::Arc; use std::time::SystemTime; static NEXT_ID: AtomicU64 = AtomicU64::new(1); @@ -20,7 +22,12 @@ fn modified_ms(meta: &std::fs::Metadata) -> u64 { .as_millis() as u64 } -pub fn stream_messages(root: &Path, rel_paths: &[PathBuf], tx: &Sender) { +pub fn stream_messages( + root: &Path, + rel_paths: &[PathBuf], + tx: &Sender, + engine: Option>, +) { debug!( "bundler: stream_messages starting — {} path(s) from {:?}", rel_paths.len(), @@ -47,6 +54,10 @@ pub fn stream_messages(root: &Path, rel_paths: &[PathBuf], tx: &Sender) size: 0, hash: [0u8; 32], modified_ms: modified_ms(&meta), + change_sequence: engine + .as_ref() + .map(|e| e.record_file_change(rel)) + .unwrap_or(0), is_dir: true, }, content: Vec::new(), @@ -63,7 +74,7 @@ pub fn stream_messages(root: &Path, rel_paths: &[PathBuf], tx: &Sender) ); flush_bundle(&mut cur, &mut cur_bytes, tx); - if let Err(e) = stream_large_file(root, rel, &meta, tx) { + if let Err(e) = stream_large_file(root, rel, &meta, tx, engine.clone()) { warn!("large-file stream({rel:?}): {e}"); } continue; @@ -91,6 +102,10 @@ pub fn stream_messages(root: &Path, rel_paths: &[PathBuf], tx: &Sender) size: size as u64, hash, modified_ms: modified_ms(&meta), + change_sequence: engine + .as_ref() + .map(|e| e.record_file_change(rel)) + .unwrap_or(0), is_dir: false, }, content, @@ -125,6 +140,7 @@ fn stream_large_file( rel: &PathBuf, meta: &std::fs::Metadata, tx: &Sender, + engine: Option>, ) -> std::io::Result<()> { let full = root.join(rel); let file_size = meta.len(); @@ -163,6 +179,10 @@ fn stream_large_file( size: file_size, hash: final_hash, modified_ms: mms, + change_sequence: engine + .as_ref() + .map(|e| e.record_file_change(rel)) + .unwrap_or(0), is_dir: false, }, total_chunks, diff --git a/crates/filesync/src/client.rs b/crates/filesync/src/client.rs index eb6e4af..31188d5 100644 --- a/crates/filesync/src/client.rs +++ b/crates/filesync/src/client.rs @@ -789,6 +789,31 @@ fn recv_loop(engine: Arc, conn: Arc, bus: Option = b + .files + .iter() + .map(|fd| fd.metadata.change_sequence) + .filter(|&seq| seq > 0) + .collect(); + + if !sequence_numbers.is_empty() { + if let Err(e) = conn.send(&Message::ChangeAcknowledgment { + bundle_id: b.bundle_id, + sequence_numbers: sequence_numbers.clone(), + }) { + warn!( + "{prefix}: failed to send acknowledgment for bundle {}: {e}", + b.bundle_id + ); + } else { + debug!( + "{prefix}: sent acknowledgment for bundle {} (sequences: {:?})", + b.bundle_id, sequence_numbers + ); + } + } } Ok(Message::LargeFileStart { ref metadata, @@ -846,6 +871,9 @@ fn recv_loop(engine: Arc, conn: Arc, bus: Option { debug!("{prefix}: LargeFileEnd committed {path:?}"); + + // 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 } Err(e) => error!("{prefix}: large_file_end: {e}"), } @@ -888,6 +916,16 @@ fn recv_loop(engine: Arc, conn: Arc, bus: Option { + debug!( + "{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"); debug!("{prefix}: unexpected message variant: {other:?}"); diff --git a/crates/filesync/src/common.rs b/crates/filesync/src/common.rs index 265cc49..e326717 100644 --- a/crates/filesync/src/common.rs +++ b/crates/filesync/src/common.rs @@ -160,7 +160,9 @@ impl PendingChanges { let mut stable = Vec::new(); for path in paths { if engine.is_file_stable(&path) { - stable.push(path); + stable.push(path.clone()); + // Clear change history for this file since we're processing it + engine.clear_change_history(&path); } else { self.changes.insert(path); } diff --git a/crates/filesync/src/manifest.rs b/crates/filesync/src/manifest.rs index 9e018b1..a53d694 100644 --- a/crates/filesync/src/manifest.rs +++ b/crates/filesync/src/manifest.rs @@ -74,6 +74,7 @@ pub fn build_manifest(root: &Path, node_id: &str, exclusions: &Exclusions) -> io size: 0, hash: [0u8; 32], modified_ms, + change_sequence: 0, is_dir: true, }, )); @@ -88,6 +89,7 @@ pub fn build_manifest(root: &Path, node_id: &str, exclusions: &Exclusions) -> io size, hash, modified_ms, + change_sequence: 0, is_dir: false, }, )) diff --git a/crates/filesync/src/protocol.rs b/crates/filesync/src/protocol.rs index 0837460..4434b1f 100644 --- a/crates/filesync/src/protocol.rs +++ b/crates/filesync/src/protocol.rs @@ -8,9 +8,10 @@ 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 = 6; +pub const PROTOCOL_VERSION: u32 = 7; pub const DEBOUNCE_MS: u64 = 200; pub const FILE_STABILITY_MS: u64 = 500; +pub const FILE_CHANGE_COALESCE_MS: u64 = 100; pub const SUPPRESSION_SECS: u64 = 2; pub const SEND_QUEUE_DEPTH: usize = 512; pub const CLIENT_BROADCAST_DEPTH: usize = 512; @@ -28,6 +29,7 @@ pub struct FileMetadata { pub size: u64, pub hash: [u8; 32], pub modified_ms: u64, + pub change_sequence: u64, pub is_dir: bool, } @@ -65,6 +67,10 @@ pub enum Message { }, ManifestExchange(Manifest), Bundle(FileBundle), + ChangeAcknowledgment { + bundle_id: u64, + sequence_numbers: Vec, + }, Delete { paths: Vec, }, diff --git a/crates/filesync/src/server.rs b/crates/filesync/src/server.rs index 169285a..6e1cf07 100644 --- a/crates/filesync/src/server.rs +++ b/crates/filesync/src/server.rs @@ -12,7 +12,7 @@ use bytehive_core::MessageBus; use crossbeam_channel::{bounded, Receiver, RecvTimeoutError, Sender, TrySendError}; use log::{debug, error, info, warn}; use parking_lot::{Mutex, RwLock}; -use std::collections::HashMap; +use std::collections::{HashMap, HashSet}; use std::io; use std::net::{TcpListener, TcpStream}; use std::path::PathBuf; @@ -33,9 +33,35 @@ pub struct Server { known_clients: Arc>, stopped: Arc, 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 +} + +impl BundleTracking { + fn new(bundle_id: u64, sequence_numbers: Vec) -> Self { + Self { + bundle_id, + sequence_numbers, + sent_time: Instant::now(), + acknowledged_by: HashSet::new(), + } + } + + fn record_acknowledgment(&mut self, client_id: &str) { + self.acknowledged_by.insert(client_id.to_string()); + } } struct ConnGuard { @@ -68,6 +94,7 @@ impl Server { peers: Arc::new(RwLock::new(HashMap::new())), active_conns: Arc::new(RwLock::new(Vec::new())), tls_config, + sent_bundles: Arc::new(RwLock::new(HashMap::new())), } } @@ -107,9 +134,12 @@ impl Server { let peers = self.peers.clone(); let bus = self.bus.clone(); let stopped = self.stopped.clone(); + let sent_bundles = self.sent_bundles.clone(); thread::Builder::new() .name("srv-watcher-broadcast".into()) - .spawn(move || local_change_broadcaster(eng, peers, fs_rx, bus, stopped))?; + .spawn(move || { + local_change_broadcaster(eng, peers, fs_rx, bus, stopped, sent_bundles) + })?; info!("filesync server listening on {} (TLS 1.3)", self.bind_addr); debug!("filesync server: accept loop started"); @@ -130,6 +160,7 @@ impl Server { let active_conns = self.active_conns.clone(); let known_clients = self.known_clients.clone(); let tls = self.tls_config.clone(); + let sent_bundles = self.sent_bundles.clone(); thread::Builder::new() .name(format!("srv-client-{addr}")) .spawn(move || { @@ -142,6 +173,7 @@ impl Server { bus, known_clients, tls, + sent_bundles, ) { error!("filesync: client {addr} session error: {e}"); } @@ -172,6 +204,7 @@ fn handle_client( bus: Option>, known_clients: Arc>, tls_config: Arc, + sent_bundles: Arc>>, // Add this parameter ) -> io::Result<()> { let conn = Arc::new(Connection::new_server(stream, tls_config)?); debug!("filesync: TLS handshake complete with new client"); @@ -654,7 +687,7 @@ fn handle_client( let recv_handle = thread::Builder::new() .name(format!("srv-recv-{client_id}")) .spawn(move || { - client_recv_loop(cid, eng_r, conn_r, peers_r, bus_r); + client_recv_loop(cid, eng_r, conn_r, peers_r, bus_r, sent_bundles); let _ = done_tx.send(()); })?; @@ -680,6 +713,7 @@ fn client_recv_loop( conn: Arc, peers: Arc>>, bus: Option>, + sent_bundles: Arc>>, // Add this parameter ) { debug!("filesync: recv loop started for {client_id}"); loop { @@ -796,6 +830,32 @@ fn client_recv_loop( Err(e) => error!("filesync: apply_rename from {client_id}: {e}"), } } + Ok(Message::ChangeAcknowledgment { + bundle_id, + sequence_numbers, + }) => { + debug!( + "filesync: recv ChangeAcknowledgment from {client_id}: bundle_id={} sequences={:?}", + bundle_id, sequence_numbers + ); + + // Record the acknowledgment + let mut sent_bundles_lock = sent_bundles.write(); + if let Some(bundle_tracking) = sent_bundles_lock.get_mut(&bundle_id) { + bundle_tracking.record_acknowledgment(&client_id); + debug!( + "filesync: bundle {} acknowledged by {} (total: {})", + bundle_id, + client_id, + bundle_tracking.acknowledged_by.len() + ); + } else { + debug!( + "filesync: received acknowledgment for unknown bundle {} from {}", + bundle_id, client_id + ); + } + } Ok(Message::RequestChunks { ref path, ref chunk_indices, @@ -896,6 +956,7 @@ fn local_change_broadcaster( fs_rx: Receiver, bus: Option>, stopped: Arc, + sent_bundles: Arc>>, // Add this parameter ) { let mut pending = PendingChanges::new(); debug!("filesync server: local change broadcaster started"); @@ -915,7 +976,7 @@ fn local_change_broadcaster( if pending.should_flush() { // debug!("filesync server: flushing local changes"); - flush_local_changes(&engine, &peers, &bus, &mut pending); + flush_local_changes(&engine, &peers, &bus, &mut pending, &sent_bundles); } } } @@ -925,6 +986,7 @@ fn flush_local_changes( peers: &Arc>>, bus: &Option>, pending: &mut PendingChanges, + sent_bundles: &Arc>>, // Add this parameter ) { let rn = pending.renames.len(); if rn > 0 { @@ -952,7 +1014,7 @@ fn flush_local_changes( "filesync server: broadcasting {} ready path(s)", ready_paths.len() ); - broadcast_paths(engine, peers, bus, ready_paths); + broadcast_paths(engine, peers, bus, ready_paths, sent_bundles); } let stable_paths = pending.take_stable_changes(engine); @@ -961,7 +1023,7 @@ fn flush_local_changes( "filesync server: broadcasting {} stable-change path(s)", stable_paths.len() ); - broadcast_paths(engine, peers, bus, stable_paths); + broadcast_paths(engine, peers, bus, stable_paths, sent_bundles); } if !pending.deletes.is_empty() { @@ -997,6 +1059,7 @@ fn broadcast_paths( peers: &Arc>>, bus: &Option>, paths: Vec, + sent_bundles: &Arc>>, // Add this parameter ) { const BROADCAST_PIPELINE_DEPTH: usize = 128; debug!( @@ -1008,9 +1071,17 @@ fn broadcast_paths( let (msg_tx, msg_rx) = bounded::(BROADCAST_PIPELINE_DEPTH); let root = engine.root().to_path_buf(); let paths_for_thread = paths.clone(); + let engine_for_thread = engine.clone(); std::thread::Builder::new() .name("broadcast-reader".into()) - .spawn(move || bundler::stream_messages(&root, &paths_for_thread, &msg_tx)) + .spawn(move || { + bundler::stream_messages( + &root, + &paths_for_thread, + &msg_tx, + Some(engine_for_thread.clone()), + ) + }) .expect("spawn broadcast-reader"); let mut total_files = 0usize; @@ -1035,6 +1106,22 @@ fn broadcast_paths( ); } } + + // Track this bundle for acknowledgment purposes + let sequence_numbers: Vec = bundle + .files + .iter() + .map(|fd| fd.metadata.change_sequence) + .filter(|&seq| seq > 0) + .collect(); + + if !sequence_numbers.is_empty() { + let mut sent_bundles_lock = sent_bundles.write(); + sent_bundles_lock.insert( + bundle.bundle_id, + BundleTracking::new(bundle.bundle_id, sequence_numbers), + ); + } } Message::LargeFileStart { metadata, .. } => { diff --git a/crates/filesync/src/sync_engine.rs b/crates/filesync/src/sync_engine.rs index 1412592..8391736 100644 --- a/crates/filesync/src/sync_engine.rs +++ b/crates/filesync/src/sync_engine.rs @@ -13,10 +13,44 @@ use std::fs; use std::path::{Component, Path, PathBuf}; use std::sync::Arc; use std::thread; -use std::time::{Duration, SystemTime}; +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, + change_count: u32, + last_sequence_number: u64, +} + +impl FileChangeHistory { + fn new(sequence_number: u64) -> Self { + Self { + last_change_time: Instant::now(), + change_count: 1, + last_sequence_number: sequence_number, + } + } + + fn record_change(&mut self, sequence_number: u64) -> bool { + let now = Instant::now(); + let elapsed = now.duration_since(self.last_change_time).as_millis() as u64; + + self.last_change_time = now; + 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 + } + + fn should_coalesce(&self) -> bool { + self.change_count > 1 + } +} + pub fn safe_relative(p: &Path) -> bool { p.components().all(|c| matches!(c, Component::Normal(_))) } @@ -90,6 +124,10 @@ pub struct SyncEngine { exclusions: Arc, trash_manager: Arc, full_scan_interval_secs: u64, + + // For tracking rapid changes and coalescing + change_sequences: RwLock>, + last_sequence_number: RwLock, } pub fn conflict_copy_name(rel_path: &Path, node_id: &str, unix_secs: u64) -> PathBuf { @@ -191,6 +229,8 @@ impl SyncEngine { exclusions, trash_manager, full_scan_interval_secs, + change_sequences: RwLock::new(HashMap::new()), + last_sequence_number: RwLock::new(0), } } @@ -239,13 +279,72 @@ impl SyncEngine { .as_millis() as u64; let age_ms = now_ms.saturating_sub(modified_ms); - age_ms >= FILE_STABILITY_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) } pub fn full_scan_interval_secs(&self) -> u64 { self.full_scan_interval_secs } + 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 + } + } + + pub fn record_file_change(&self, rel: &Path) -> u64 { + let mut seq_lock = self.last_sequence_number.write(); + *seq_lock += 1; + let sequence_number = *seq_lock; + + let mut histories = self.change_sequences.write(); + let entry = histories + .entry(rel.to_path_buf()) + .or_insert_with(|| FileChangeHistory::new(sequence_number)); + + entry.record_change(sequence_number); + sequence_number + } + + pub fn get_current_sequence_number(&self) -> u64 { + *self.last_sequence_number.read() + } + + pub fn clear_change_history(&self, rel: &Path) { + let mut histories = self.change_sequences.write(); + 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(), + node_id: self.node_id.clone(), + manifest: RwLock::new(self.manifest.read().clone()), + suppressed: self.suppressed.clone(), + suppressed_deletes: self.suppressed_deletes.clone(), + in_progress: RwLock::new(HashMap::new()), + exclusions: self.exclusions.clone(), + trash_manager: self.trash_manager.clone(), + full_scan_interval_secs: self.full_scan_interval_secs, + change_sequences: RwLock::new(HashMap::new()), + last_sequence_number: RwLock::new(0), + } + } + pub fn trash_manager(&self) -> &Arc { &self.trash_manager } @@ -297,7 +396,7 @@ impl SyncEngine { thread::Builder::new() .name("file-reader".into()) .spawn(move || { - bundler::stream_messages(&root, &slice, &tx); + bundler::stream_messages(&root, &slice, &tx, None); }) .expect("spawn file-reader") }) @@ -319,7 +418,8 @@ impl SyncEngine { let (tx, rx) = bounded::(PIPELINE_DEPTH); let root = self.root.clone(); let paths = paths.to_vec(); - thread::spawn(move || bundler::stream_messages(&root, &paths, &tx)); + let engine = self.clone_with_fresh_state(); + thread::spawn(move || bundler::stream_messages(&root, &paths, &tx, Some(Arc::new(engine)))); rx.into_iter() .filter_map(|m| { if let Message::Bundle(b) = m { @@ -678,6 +778,7 @@ impl SyncEngine { size: asm.file_size, hash: final_hash, modified_ms: asm.modified_ms, + change_sequence: 0, is_dir: false, }; self.manifest.write().files.insert(path.clone(), file_meta); diff --git a/crates/filesync/tests/test_bundler.rs b/crates/filesync/tests/test_bundler.rs index 6692460..da4b8f3 100644 --- a/crates/filesync/tests/test_bundler.rs +++ b/crates/filesync/tests/test_bundler.rs @@ -19,7 +19,7 @@ fn tmp_dir(suffix: &str) -> std::path::PathBuf { fn collect_messages(root: &std::path::Path, paths: &[PathBuf]) -> Vec { let (tx, rx) = bounded(256); - stream_messages(root, paths, &tx); + stream_messages(root, paths, &tx, None); drop(tx); rx.iter().collect() } diff --git a/crates/filesync/tests/test_client.rs b/crates/filesync/tests/test_client.rs index 98c1192..af8c0c7 100644 --- a/crates/filesync/tests/test_client.rs +++ b/crates/filesync/tests/test_client.rs @@ -13,6 +13,7 @@ fn make_manifest(entries: &[(&str, u64, bool)]) -> Manifest { files.insert( p.clone(), FileMetadata { + change_sequence: 0, rel_path: p, size: *size, hash: [0u8; 32], diff --git a/crates/filesync/tests/test_conflict.rs b/crates/filesync/tests/test_conflict.rs index 2b6668e..f0387fe 100644 --- a/crates/filesync/tests/test_conflict.rs +++ b/crates/filesync/tests/test_conflict.rs @@ -55,6 +55,7 @@ fn file_bundle(rel: &str, content: &[u8]) -> FileBundle { FileBundle { files: vec![FileData { metadata: FileMetadata { + change_sequence: 0, rel_path: PathBuf::from(rel), size: content.len() as u64, hash, @@ -93,6 +94,7 @@ fn large_file_with_prior( let incoming_hash: [u8; 32] = blake3::hash(incoming_content).into(); let rel_path = PathBuf::from(rel); let meta = FileMetadata { + change_sequence: 0, rel_path: rel_path.clone(), size: incoming_content.len() as u64, hash: incoming_hash, @@ -291,6 +293,7 @@ fn no_conflict_for_directory_entries_in_bundle() { let dir_bundle = FileBundle { files: vec![FileData { metadata: FileMetadata { + change_sequence: 0, rel_path: PathBuf::from("mydir"), size: 0, hash: [0u8; 32], @@ -528,6 +531,7 @@ fn only_conflicted_files_get_copies_in_multi_file_bundle() { files: vec![ FileData { metadata: FileMetadata { + change_sequence: 0, rel_path: PathBuf::from("will_conflict.txt"), size: 5, hash: blake3::hash(b"base1").into(), @@ -538,6 +542,7 @@ fn only_conflicted_files_get_copies_in_multi_file_bundle() { }, FileData { metadata: FileMetadata { + change_sequence: 0, rel_path: PathBuf::from("no_conflict.txt"), size: 5, hash: blake3::hash(b"base2").into(), @@ -560,6 +565,7 @@ fn only_conflicted_files_get_copies_in_multi_file_bundle() { files: vec![ FileData { metadata: FileMetadata { + change_sequence: 0, rel_path: PathBuf::from("will_conflict.txt"), size: 13, hash: blake3::hash(b"remote update1").into(), @@ -570,6 +576,7 @@ fn only_conflicted_files_get_copies_in_multi_file_bundle() { }, FileData { metadata: FileMetadata { + change_sequence: 0, rel_path: PathBuf::from("no_conflict.txt"), size: 13, hash: blake3::hash(b"remote update2").into(), @@ -611,6 +618,7 @@ fn unsafe_path_still_rejected_even_with_conflict_logic_active() { let bad_bundle = FileBundle { files: vec![FileData { metadata: FileMetadata { + change_sequence: 0, rel_path: PathBuf::from("../escape.txt"), size: 4, hash: blake3::hash(b"evil").into(), @@ -739,6 +747,7 @@ fn large_file_no_conflict_for_brand_new_file() { let hash: [u8; 32] = blake3::hash(content).into(); let rel = PathBuf::from("new_large.bin"); let meta = FileMetadata { + change_sequence: 0, rel_path: rel.clone(), size: content.len() as u64, hash, @@ -792,6 +801,7 @@ fn large_file_no_conflict_when_incoming_equals_manifest() { let incoming_hash: [u8; 32] = blake3::hash(content).into(); let rel = PathBuf::from("stable.bin"); let meta = FileMetadata { + change_sequence: 0, rel_path: rel.clone(), size: content.len() as u64, hash: incoming_hash, @@ -1021,6 +1031,7 @@ fn regression_large_file_happy_path_yields_committed() { let hash: [u8; 32] = blake3::hash(content).into(); let rel = PathBuf::from("lf.bin"); let meta = FileMetadata { + change_sequence: 0, rel_path: rel.clone(), size: content.len() as u64, hash, @@ -1049,6 +1060,7 @@ fn regression_large_file_hash_mismatch_still_errors() { let wrong_hash = [0xFFu8; 32]; let rel = PathBuf::from("corrupt.bin"); let meta = FileMetadata { + change_sequence: 0, rel_path: rel.clone(), size: content.len() as u64, hash: wrong_hash, @@ -1071,6 +1083,7 @@ fn regression_apply_bundle_multiple_files_written_correctly() { files: vec![ FileData { metadata: FileMetadata { + change_sequence: 0, rel_path: PathBuf::from("one.txt"), size: 3, hash: blake3::hash(b"aaa").into(), @@ -1081,6 +1094,7 @@ fn regression_apply_bundle_multiple_files_written_correctly() { }, FileData { metadata: FileMetadata { + change_sequence: 0, rel_path: PathBuf::from("two.txt"), size: 3, hash: blake3::hash(b"bbb").into(), diff --git a/crates/filesync/tests/test_manifest.rs b/crates/filesync/tests/test_manifest.rs index 909b595..26ae5ee 100644 --- a/crates/filesync/tests/test_manifest.rs +++ b/crates/filesync/tests/test_manifest.rs @@ -29,6 +29,7 @@ fn make_manifest(entries: &[(&str, u64, [u8; 32], bool, u64)], node: &str) -> Ma files.insert( p.clone(), FileMetadata { + change_sequence: 0, rel_path: p, size: *size, hash: *hash, diff --git a/crates/filesync/tests/test_protocol.rs b/crates/filesync/tests/test_protocol.rs index af4d8a9..7a0ef9f 100644 --- a/crates/filesync/tests/test_protocol.rs +++ b/crates/filesync/tests/test_protocol.rs @@ -24,6 +24,7 @@ fn bundle_msg(filename: &str, content: &[u8]) -> Message { size: content.len() as u64, hash, modified_ms: 1_000_000, + change_sequence: 0, is_dir: false, }, content: content.to_vec(), @@ -114,6 +115,7 @@ fn roundtrip_manifest_exchange() { size: 42, hash, modified_ms: 9999, + change_sequence: 0, is_dir: false, }, ); @@ -144,6 +146,7 @@ fn roundtrip_large_file_start() { size: 32 * 1024 * 1024, hash, modified_ms: 12345, + change_sequence: 0, is_dir: false, }; let msg = Message::LargeFileStart { diff --git a/crates/filesync/tests/test_sync_engine.rs b/crates/filesync/tests/test_sync_engine.rs index a05d3fd..150ff94 100644 --- a/crates/filesync/tests/test_sync_engine.rs +++ b/crates/filesync/tests/test_sync_engine.rs @@ -32,6 +32,7 @@ fn file_bundle(rel: &str, content: &[u8]) -> FileBundle { FileBundle { files: vec![FileData { metadata: FileMetadata { + change_sequence: 0, rel_path: PathBuf::from(rel), size: content.len() as u64, hash, @@ -48,6 +49,7 @@ fn dir_bundle(rel: &str) -> FileBundle { FileBundle { files: vec![FileData { metadata: FileMetadata { + change_sequence: 0, rel_path: PathBuf::from(rel), size: 0, hash: [0u8; 32], @@ -136,6 +138,7 @@ fn apply_bundle_rejects_path_traversal() { let evil_bundle = FileBundle { files: vec![FileData { metadata: FileMetadata { + change_sequence: 0, rel_path: PathBuf::from("../escaped.txt"), size: 4, hash: [0u8; 32], @@ -159,6 +162,7 @@ fn apply_bundle_rejects_absolute_path() { let evil_bundle = FileBundle { files: vec![FileData { metadata: FileMetadata { + change_sequence: 0, rel_path: PathBuf::from("/tmp/evil.txt"), size: 4, hash: [0u8; 32], @@ -194,6 +198,7 @@ fn apply_bundle_handles_multiple_files_in_one_bundle() { files: vec![ FileData { metadata: FileMetadata { + change_sequence: 0, rel_path: PathBuf::from("a.txt"), size: 1, hash: blake3::hash(b"a").into(), @@ -204,6 +209,7 @@ fn apply_bundle_handles_multiple_files_in_one_bundle() { }, FileData { metadata: FileMetadata { + change_sequence: 0, rel_path: PathBuf::from("b.txt"), size: 1, hash: blake3::hash(b"b").into(), @@ -303,6 +309,7 @@ fn large_file_flow_happy_path() { let hash: [u8; 32] = blake3::hash(content).into(); let rel = PathBuf::from("large.bin"); let meta = FileMetadata { + change_sequence: 0, rel_path: rel.clone(), size: content.len() as u64, hash, @@ -329,6 +336,7 @@ fn large_file_finish_rejects_hash_mismatch() { let wrong_hash = [0xFFu8; 32]; let rel = PathBuf::from("bad.bin"); let meta = FileMetadata { + change_sequence: 0, rel_path: rel.clone(), size: content.len() as u64, hash: wrong_hash, @@ -351,6 +359,7 @@ fn large_file_rejects_unsafe_path() { let dir = tmp_dir("lf_unsafe"); let engine = make_engine(dir.clone()); let meta = FileMetadata { + change_sequence: 0, rel_path: PathBuf::from("../outside.bin"), size: 0, hash: [0u8; 32], @@ -513,6 +522,7 @@ fn large_file_finish_detects_missing_chunks() { let hash: [u8; 32] = blake3::hash(&all_content).into(); let rel = PathBuf::from("partial.bin"); let meta = FileMetadata { + change_sequence: 0, rel_path: rel.clone(), size: all_content.len() as u64, hash, @@ -621,6 +631,7 @@ fn begin_large_file_rejects_path_traversal() { let dir = tmp_dir("lf_unsafe_path"); let engine = make_engine(dir.clone()); let meta = FileMetadata { + change_sequence: 0, rel_path: PathBuf::from("../escape.bin"), size: 100, hash: [0u8; 32], @@ -666,6 +677,7 @@ fn clear_in_progress_removes_assembly() { let engine = make_engine(dir.clone()); // Begin a large file assembly let meta = FileMetadata { + change_sequence: 0, rel_path: PathBuf::from("bigfile.bin"), size: 1024, hash: [0u8; 32],