Runs cargo fmt --all; all 573 tests still passing, clippy still clean. No logic changes — formatting only. Co-Authored-By: Claude Sonnet 4.6 <[email protected]>
309 lines
13 KiB
Rust
309 lines
13 KiB
Rust
//! 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<u64> = 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");
|
|
}
|