//! End-to-end sync benchmark — regression guard for the full push/pull cycle. //! //! Groups: //! - `end_to_end/push/N` — cold push: N revisions client → empty server //! - `end_to_end/push_delta/N/D` — delta push: only D new revisions transferred //! - `end_to_end/pull/N` — cold pull: N revisions server → empty client //! //! Each iteration spins up a real TCP loopback server and performs a complete //! sync handshake. Setup (file creation) is excluded from timing via //! `iter_batched`. //! //! Run: //! cargo bench -p clawsync-agent -- end_to_end_bench use std::net::{IpAddr, Ipv4Addr, SocketAddr}; use std::path::{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::manifest::ClawSyncManifest; use clawsync_onion::merger::merge_packets; use clawsync_onion::selector::SyncSelector; use clawsync_transport::protocol::SyncMessage; use clawsync_transport::tcp::TcpServer; use criterion::{BatchSize, BenchmarkId, Criterion, Throughput, criterion_group, criterion_main}; use tempfile::TempDir; // ───────────────────────────────────────────────────────────────────────────── const H5_MAGIC: &[u8] = b"\x89HDF\r\n\x1a\n"; fn any_addr() -> SocketAddr { SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0) } /// Build an OnionFile with `n` committed revisions in `dir` and return (path, onion). fn make_onion_n(dir: &TempDir, name: &str, n: usize) -> (PathBuf, OnionFile) { let h5 = dir.path().join(name); 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 as u8; 4096]); onion.commit_session(s, Some(&format!("rev {i}"))).unwrap(); } onion.flush().unwrap(); (h5, onion) } /// Create an empty `.onion` sidecar (no revisions) in `dir`. fn make_empty_onion(dir: &TempDir, name: &str) -> PathBuf { let h5 = dir.path().join(name); std::fs::write(&h5, H5_MAGIC).unwrap(); let mut onion = OnionFile::create(&h5, 4096).unwrap(); onion.flush().unwrap(); h5 } // ───────────────────────────────────────────────────────────────────────────── // Direction-aware one-shot TCP server // // Receives one connection and handles it as either: // - a push receiver (server has fewer revisions → receives packets from client) // - a pull sender (server has more revisions → sends packets to client) // ───────────────────────────────────────────────────────────────────────────── async fn serve_one(server: TcpServer, server_onion_path: PathBuf) { let (mut conn, _) = server.accept().await.unwrap(); // Read the client's ManifestRequest let client_rev_count = match conn.recv().await.unwrap() { SyncMessage::ManifestRequest { revision_count, .. } => revision_count, other => panic!("server: expected ManifestRequest, got {other:?}"), }; let server_onion = OnionFile::open(&server_onion_path).unwrap(); let server_rev_count = server_onion.revision_count(); let manifest = ClawSyncManifest::from_onion("bench-server", &server_onion, H5_MAGIC); conn.send(&SyncMessage::ManifestResponse { manifest }) .await .unwrap(); if server_rev_count > client_rev_count { // Server is ahead — push its extra revisions to client let client_head = if client_rev_count == 0 { NO_PARENT } else { client_rev_count - 1 }; let packets = diff_revisions(&server_onion, client_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!("server: expected Ack, got {other:?}"), } } conn.send(&SyncMessage::SyncComplete { revisions_transferred: total, bytes_transferred: bytes_sent, }) .await .unwrap(); } else { // Server is behind or equal — receive packets from client let mut server_onion_rw = OnionFile::open(&server_onion_path).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 { .. } => { merge_packets(&mut server_onion_rw, packets, true).unwrap(); return; } other => panic!("server: unexpected {other:?}"), } } } } // ───────────────────────────────────────────────────────────────────────────── // Benchmarks // ───────────────────────────────────────────────────────────────────────────── fn bench_e2e(c: &mut Criterion) { let rt = tokio::runtime::Builder::new_multi_thread() .worker_threads(2) .enable_all() .build() .unwrap(); // ── Cold push: N revisions client → empty server ────────────────────────── { let mut group = c.benchmark_group("end_to_end/push"); for &n in &[10usize, 100, 1_000] { group.sample_size(if n >= 1_000 { 10 } else { 20 }); group.throughput(Throughput::Elements(n as u64)); group.bench_with_input(BenchmarkId::from_parameter(n), &n, |b, &n| { b.iter_batched( || { // SETUP (not timed) let src_dir = TempDir::new().unwrap(); let (src_h5, _) = make_onion_n(&src_dir, "src.h5", n); let dst_dir = TempDir::new().unwrap(); let dst_h5 = make_empty_onion(&dst_dir, "dst.h5"); let (server, addr) = rt.block_on(async { let s = TcpServer::bind(any_addr()).await.unwrap(); let a = s.local_addr; (s, a) }); (src_dir, src_h5, dst_dir, dst_h5, server, addr) }, |(src_dir, src_h5, _dst_dir, dst_h5, server, addr)| { // TIMED rt.block_on(async { let server_task = tokio::spawn(serve_one(server, dst_h5)); let backend = TcpSyncBackend::new(addr, "bench-agent"); backend .push(Path::new(&src_h5), &SyncSelector::All) .await .unwrap(); server_task.await.unwrap(); }); drop(src_dir); }, BatchSize::SmallInput, ); }); } group.finish(); } // ── Delta push: only D new revisions transferred ────────────────────────── { let mut group = c.benchmark_group("end_to_end/push_delta"); let n = 100usize; for &d in &[1usize, 5, 10] { group.sample_size(20); group.throughput(Throughput::Elements(d as u64)); group.bench_with_input( BenchmarkId::new(format!("n{n}"), format!("d{d}")), &(n, d), |b, &(n, d)| { b.iter_batched( || { // SETUP: client has n, server has n-d let src_dir = TempDir::new().unwrap(); let (src_h5, _) = make_onion_n(&src_dir, "src.h5", n); let dst_dir = TempDir::new().unwrap(); let (dst_h5, _) = make_onion_n(&dst_dir, "dst.h5", n - d); let (server, addr) = rt.block_on(async { let s = TcpServer::bind(any_addr()).await.unwrap(); let a = s.local_addr; (s, a) }); (src_dir, src_h5, dst_dir, dst_h5, server, addr) }, |(_src_dir, src_h5, _dst_dir, dst_h5, server, addr)| { // TIMED: should only transfer d packets rt.block_on(async { let server_task = tokio::spawn(serve_one(server, dst_h5)); let backend = TcpSyncBackend::new(addr, "bench-agent"); let stats = backend .push(Path::new(&src_h5), &SyncSelector::All) .await .unwrap(); assert_eq!( stats.revisions_transferred, d as u64, "delta push should transfer exactly d={d} revisions" ); server_task.await.unwrap(); }); }, BatchSize::SmallInput, ); }, ); } group.finish(); } // ── Cold pull: N revisions server → empty client ────────────────────────── { let mut group = c.benchmark_group("end_to_end/pull"); for &n in &[10usize, 100, 1_000] { group.sample_size(if n >= 1_000 { 10 } else { 20 }); group.throughput(Throughput::Elements(n as u64)); group.bench_with_input(BenchmarkId::from_parameter(n), &n, |b, &n| { b.iter_batched( || { // SETUP: server has n revisions, client is empty let srv_dir = TempDir::new().unwrap(); let (srv_h5, _) = make_onion_n(&srv_dir, "srv.h5", n); let cli_dir = TempDir::new().unwrap(); let cli_h5 = make_empty_onion(&cli_dir, "cli.h5"); let (server, addr) = rt.block_on(async { let s = TcpServer::bind(any_addr()).await.unwrap(); let a = s.local_addr; (s, a) }); (srv_dir, srv_h5, cli_dir, cli_h5, server, addr) }, |(_srv_dir, srv_h5, _cli_dir, cli_h5, server, addr)| { // TIMED rt.block_on(async { let server_task = tokio::spawn(serve_one(server, srv_h5)); let backend = TcpSyncBackend::new(addr, "bench-agent"); backend.pull(Path::new(&cli_h5)).await.unwrap(); server_task.await.unwrap(); }); }, BatchSize::SmallInput, ); }); } group.finish(); } // ── Wire size / throughput summary ──────────────────────────────────────── println!("\n=== End-to-end: revisions transferred per push (delta guard) ==="); println!("{:>8} {:>8} {:>10}", "N", "D (new)", "should xfer"); let rt2 = tokio::runtime::Builder::new_current_thread() .enable_all() .build() .unwrap(); for &(n, d) in &[(100usize, 1usize), (100, 10), (1000, 10), (1000, 50)] { let src_dir = TempDir::new().unwrap(); let (src_h5, _) = make_onion_n(&src_dir, "src.h5", n); let dst_dir = TempDir::new().unwrap(); let (dst_h5, _) = make_onion_n(&dst_dir, "dst.h5", n - d); let actual = rt2.block_on(async { let server = TcpServer::bind(any_addr()).await.unwrap(); let addr = server.local_addr; let server_task = tokio::spawn(serve_one(server, dst_h5)); let backend = TcpSyncBackend::new(addr, "bench-agent"); let stats = backend .push(Path::new(&src_h5), &SyncSelector::All) .await .unwrap(); server_task.await.unwrap(); stats.revisions_transferred }); println!("{:>8} {:>8} {:>9} ✓", n, d, actual); } } criterion_group!(benches, bench_e2e); criterion_main!(benches);