//! Integration tests for `TcpSyncBackend` push and pull. //! //! Spins up a real TCP server (same direction-aware logic as the CLI), then //! exercises `TcpSyncBackend::push` and `TcpSyncBackend::pull` end-to-end. use std::net::{IpAddr, Ipv4Addr, SocketAddr}; use std::path::PathBuf; use clawhdf5_onion::format::NO_PARENT; use clawhdf5_onion::writer::OnionFile; use clawsync_agent::backend::{SyncBackend, TcpSyncBackend}; use clawsync_onion::differ::diff_revisions; use clawsync_onion::iblt::IbltSketch; use clawsync_onion::manifest::{ClawSyncManifest, IbltManifest}; use clawsync_onion::merger::merge_packets; use clawsync_onion::selector::SyncSelector; use clawsync_transport::protocol::SyncMessage; use clawsync_transport::tcp::TcpServer; use tempfile::TempDir; const H5_MAGIC: &[u8] = b"\x89HDF\r\n\x1a\n"; fn any_addr() -> SocketAddr { SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0) } /// Create a temp HDF5 stub with `n` committed onion revisions, flushed. fn make_onion(n: u8) -> (TempDir, PathBuf, OnionFile) { let dir = TempDir::new().unwrap(); let h5 = dir.path().join("data.h5"); std::fs::write(&h5, H5_MAGIC).unwrap(); let mut onion = OnionFile::create(&h5, 4096).unwrap(); for i in 0..n { let mut s = onion.begin_session(None).unwrap(); s.record_page(0, &vec![i; 4096]); onion.commit_session(s, Some(&format!("rev {i}"))).unwrap(); } onion.flush().unwrap(); (dir, h5, onion) } // ───────────────────────────────────────────────────────────────────────────── // Direction-aware server helper (mirrors CLI handle_client) // ───────────────────────────────────────────────────────────────────────────── async fn run_server(server: TcpServer, mut dst: OnionFile) -> OnionFile { let (mut conn, _) = server.accept().await.unwrap(); match conn.recv().await.unwrap() { // ── IBLT pre-flight (push direction) ───────────────────────────── SyncMessage::IbltRequest { sketch: client_iblt, } => { let server_revs: Vec = dst.list_revisions().iter().map(|s| s.revision).collect(); // Decode diff: what client has that server lacks let diff = client_iblt.diff_against(&server_revs).unwrap(); // Build server's sketch with the same embedded seed let client_sketch = IbltSketch::from_bytes(&client_iblt.sketch).unwrap(); let seed = client_sketch.seed(); let srv_sketch = IbltSketch::from_keys(&server_revs, seed); let server_iblt = IbltManifest { agent_id: "test-server".to_string(), file_blake3: [0u8; 32], revision_count: server_revs.len() as u64, head_revision: server_revs.last().copied().unwrap_or(0), head_blake3: [0u8; 32], last_write: 0.0, sketch_cells: srv_sketch.cell_count() as u32, sketch: srv_sketch.to_bytes(), }; conn.send(&SyncMessage::IbltResponse { sketch: server_iblt, missing_from_remote: diff.only_in_b.clone(), // client has, server lacks }) .await .unwrap(); let mut packets = Vec::new(); loop { match conn.recv().await.unwrap() { SyncMessage::LayerPacket { packet } => { let rev = packet.revision; conn.send(&SyncMessage::Ack { revision: rev }) .await .unwrap(); packets.push(packet); } SyncMessage::SyncComplete { .. } => break, other => panic!("server: unexpected {other:?}"), } } if !packets.is_empty() { merge_packets(&mut dst, packets, true).unwrap(); } } // ── Flat manifest (pull direction) ──────────────────────────────── SyncMessage::ManifestRequest { revision_count: remote_rev_count, .. } => { let local_rev_count = dst.revision_count(); let manifest = ClawSyncManifest::from_onion("server", &dst, H5_MAGIC); conn.send(&SyncMessage::ManifestResponse { manifest }) .await .unwrap(); if local_rev_count > remote_rev_count { let remote_head = if remote_rev_count == 0 { NO_PARENT } else { remote_rev_count - 1 }; let packets = diff_revisions(&dst, remote_head).unwrap(); let total = packets.len() as u64; let mut bytes_sent = 0u64; for packet in packets { bytes_sent += packet.page_data_size() as u64; conn.send(&SyncMessage::LayerPacket { packet }) .await .unwrap(); match conn.recv().await.unwrap() { SyncMessage::Ack { .. } => {} other => panic!("expected Ack, got {other:?}"), } } conn.send(&SyncMessage::SyncComplete { revisions_transferred: total, bytes_transferred: bytes_sent, }) .await .unwrap(); } else { let mut packets = Vec::new(); loop { match conn.recv().await.unwrap() { SyncMessage::LayerPacket { packet } => { let rev = packet.revision; conn.send(&SyncMessage::Ack { revision: rev }) .await .unwrap(); packets.push(packet); } SyncMessage::SyncComplete { .. } => break, other => panic!("server: unexpected {other:?}"), } } if !packets.is_empty() { merge_packets(&mut dst, packets, true).unwrap(); } } } other => panic!("server: expected IbltRequest or ManifestRequest, got {other:?}"), } dst } // ───────────────────────────────────────────────────────────────────────────── // Push tests // ───────────────────────────────────────────────────────────────────────────── /// `TcpSyncBackend::push` sends all revisions to an empty server. #[tokio::test] async fn backend_push_all_to_empty() { let (_src_dir, src_h5, src_onion) = make_onion(4); let (_dst_dir, _dst_h5, dst_onion) = make_onion(0); let server = TcpServer::bind(any_addr()).await.unwrap(); let addr = server.local_addr; let server_task = tokio::spawn(run_server(server, dst_onion)); let backend = TcpSyncBackend::new(addr, "test-agent"); let stats = backend.push(&src_h5, &SyncSelector::All).await.unwrap(); assert_eq!(stats.revisions_transferred, 4); let dst_final = server_task.await.unwrap(); assert_eq!(dst_final.revision_count(), 4); let _ = src_onion; } /// `TcpSyncBackend::push` sends only the missing delta (N-K). #[tokio::test] async fn backend_push_partial_delta() { const K: u8 = 3; const N: u8 = 7; let (_src_dir, src_h5, src_onion) = make_onion(N); let (_dst_dir, _dst_h5, dst_onion) = make_onion(K); let server = TcpServer::bind(any_addr()).await.unwrap(); let addr = server.local_addr; let server_task = tokio::spawn(run_server(server, dst_onion)); let backend = TcpSyncBackend::new(addr, "agent"); let stats = backend.push(&src_h5, &SyncSelector::All).await.unwrap(); assert_eq!(stats.revisions_transferred, (N - K) as u64); let dst_final = server_task.await.unwrap(); assert_eq!(dst_final.revision_count(), N as u64); let _ = src_onion; } /// `TcpSyncBackend::push` when already in sync transfers zero revisions. #[tokio::test] async fn backend_push_already_in_sync() { let (_src_dir, src_h5, src_onion) = make_onion(3); let (_dst_dir, _dst_h5, dst_onion) = make_onion(3); let server = TcpServer::bind(any_addr()).await.unwrap(); let addr = server.local_addr; let server_task = tokio::spawn(run_server(server, dst_onion)); let backend = TcpSyncBackend::new(addr, "agent"); let stats = backend.push(&src_h5, &SyncSelector::All).await.unwrap(); assert_eq!(stats.revisions_transferred, 0); let _ = (src_onion, server_task.await.unwrap()); } // ───────────────────────────────────────────────────────────────────────────── // Pull tests // ───────────────────────────────────────────────────────────────────────────── /// `TcpSyncBackend::pull` receives all revisions from a full server. #[tokio::test] async fn backend_pull_all_from_server() { let (_srv_dir, _srv_h5, srv_onion) = make_onion(5); let (cli_dir, cli_h5, cli_onion) = make_onion(0); let server = TcpServer::bind(any_addr()).await.unwrap(); let addr = server.local_addr; let server_task = tokio::spawn(run_server(server, srv_onion)); let backend = TcpSyncBackend::new(addr, "agent"); let stats = backend.pull(&cli_h5).await.unwrap(); assert_eq!(stats.revisions_transferred, 5); server_task.await.unwrap(); // Re-open from disk to verify the sidecar was updated. let cli_after = OnionFile::open(&cli_h5).unwrap(); assert_eq!(cli_after.revision_count(), 5); let _ = (cli_dir, cli_onion); } /// `TcpSyncBackend::pull` when already in sync returns zero transfers. #[tokio::test] async fn backend_pull_already_in_sync() { let (_srv_dir, _srv_h5, srv_onion) = make_onion(3); let (_cli_dir, cli_h5, cli_onion) = make_onion(3); let server = TcpServer::bind(any_addr()).await.unwrap(); let addr = server.local_addr; let server_task = tokio::spawn(run_server(server, srv_onion)); let backend = TcpSyncBackend::new(addr, "agent"); let stats = backend.pull(&cli_h5).await.unwrap(); assert_eq!( stats.revisions_transferred, 0, "nothing to pull when in sync" ); let _ = (cli_onion, server_task.await.unwrap()); } /// `TcpSyncBackend::pull` fetches only the missing delta. #[tokio::test] async fn backend_pull_partial_delta() { const K: u8 = 2; const N: u8 = 6; let (_srv_dir, _srv_h5, srv_onion) = make_onion(N); let (_cli_dir, cli_h5, cli_onion) = make_onion(K); let server = TcpServer::bind(any_addr()).await.unwrap(); let addr = server.local_addr; let server_task = tokio::spawn(run_server(server, srv_onion)); let backend = TcpSyncBackend::new(addr, "agent"); let stats = backend.pull(&cli_h5).await.unwrap(); assert_eq!(stats.revisions_transferred, (N - K) as u64); server_task.await.unwrap(); let cli_after = OnionFile::open(&cli_h5).unwrap(); assert_eq!(cli_after.revision_count(), N as u64); let _ = cli_onion; } // ───────────────────────────────────────────────────────────────────────────── // local_manifest // ───────────────────────────────────────────────────────────────────────────── /// `local_manifest` reflects the correct revision count. #[test] fn backend_local_manifest_revision_count() { let (_dir, h5, _onion) = make_onion(4); let backend = TcpSyncBackend::new(any_addr(), "agent"); let mf = backend.local_manifest(&h5).unwrap(); assert_eq!(mf.revision_count, 4); assert_eq!(mf.agent_id, "agent"); }