- Add QuicSyncBackend + quic_scheduler; fix push stream/close race by waiting for peer close instead of calling close_and_drain() on the client side (SyncComplete stream was racing CONNECTION_CLOSE) - Add PipeWriteHalf::shutdown_push() for client-sends-last QUIC paths; use it in cmd_push so the server can process SyncComplete before the connection tears down - Add SshSyncBackend + ssh_scheduler with subprocess integration tests - Add clawsync diff command with --porcelain flag and subprocess tests - Add QUIC subprocess push/pull integration tests - Fix clippy --tests violations across 5 crates - Fix broken intra-doc links (reader.rs, lib.rs, scheduler.rs) - Rewrite crate READMEs; update BENCHMARKS.md with clawsync-fs CDC row 653 tests, 0 failures. Co-Authored-By: Claude Sonnet 4.6 <[email protected]>
302 lines
12 KiB
Rust
302 lines
12 KiB
Rust
//! Integration tests for `QuicSyncBackend` push and pull.
|
|
//!
|
|
//! Spins up a real QUIC server (`QuicConfig::self_signed()`) on loopback, then
|
|
//! exercises `QuicSyncBackend` (which uses `QuicConfig::insecure()` on the
|
|
//! client side so cert verification is skipped) end-to-end.
|
|
|
|
use std::net::{IpAddr, Ipv4Addr, SocketAddr};
|
|
use std::path::PathBuf;
|
|
use std::sync::OnceLock;
|
|
|
|
use clawhdf5_onion::format::NO_PARENT;
|
|
use clawhdf5_onion::writer::OnionFile;
|
|
use clawsync_agent::backend::{QuicSyncBackend, SyncBackend};
|
|
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::quic::{QuicConfig, QuicServer};
|
|
use tempfile::TempDir;
|
|
|
|
const H5_MAGIC: &[u8] = b"\x89HDF\r\n\x1a\n";
|
|
|
|
static CRYPTO: OnceLock<()> = OnceLock::new();
|
|
|
|
fn init_crypto() {
|
|
CRYPTO.get_or_init(|| {
|
|
rustls::crypto::ring::default_provider()
|
|
.install_default()
|
|
.ok();
|
|
});
|
|
}
|
|
|
|
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 QUIC server helper (mirrors backend_integration::run_server)
|
|
// ─────────────────────────────────────────────────────────────────────────────
|
|
|
|
async fn run_quic_server(server: QuicServer, mut dst: OnionFile) -> OnionFile {
|
|
let 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();
|
|
|
|
let diff = client_iblt.diff_against(&server_revs).unwrap();
|
|
|
|
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(),
|
|
})
|
|
.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();
|
|
// Keep connection alive until client closes it, so SyncComplete
|
|
// stream arrives before any CONNECTION_CLOSE races it.
|
|
conn.wait_for_peer_close().await;
|
|
} 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
|
|
// ─────────────────────────────────────────────────────────────────────────────
|
|
|
|
/// `QuicSyncBackend::push` sends all revisions to an empty server.
|
|
#[tokio::test]
|
|
async fn quic_backend_push_cold_copy() {
|
|
init_crypto();
|
|
let (_src_dir, src_h5, src_onion) = make_onion(5);
|
|
let (_dst_dir, _dst_h5, dst_onion) = make_onion(0);
|
|
|
|
let server = QuicServer::bind(any_addr(), QuicConfig::self_signed().unwrap())
|
|
.await
|
|
.unwrap();
|
|
let addr = server.local_addr;
|
|
let server_task = tokio::spawn(run_quic_server(server, dst_onion));
|
|
|
|
let backend = QuicSyncBackend::new(addr, "test-agent");
|
|
let stats = backend.push(&src_h5, &SyncSelector::All).await.unwrap();
|
|
|
|
assert_eq!(stats.revisions_transferred, 5);
|
|
|
|
let dst_final = server_task.await.unwrap();
|
|
assert_eq!(dst_final.revision_count(), 5);
|
|
let _ = src_onion;
|
|
}
|
|
|
|
/// `QuicSyncBackend::push` sends only the missing delta (N-K).
|
|
#[tokio::test]
|
|
async fn quic_backend_push_partial_delta() {
|
|
init_crypto();
|
|
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 = QuicServer::bind(any_addr(), QuicConfig::self_signed().unwrap())
|
|
.await
|
|
.unwrap();
|
|
let addr = server.local_addr;
|
|
let server_task = tokio::spawn(run_quic_server(server, dst_onion));
|
|
|
|
let backend = QuicSyncBackend::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;
|
|
}
|
|
|
|
/// `QuicSyncBackend::push` when already in sync transfers zero revisions.
|
|
#[tokio::test]
|
|
async fn quic_backend_push_already_in_sync() {
|
|
init_crypto();
|
|
let (_src_dir, src_h5, src_onion) = make_onion(4);
|
|
let (_dst_dir, _dst_h5, dst_onion) = make_onion(4);
|
|
|
|
let server = QuicServer::bind(any_addr(), QuicConfig::self_signed().unwrap())
|
|
.await
|
|
.unwrap();
|
|
let addr = server.local_addr;
|
|
let server_task = tokio::spawn(run_quic_server(server, dst_onion));
|
|
|
|
let backend = QuicSyncBackend::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
|
|
// ─────────────────────────────────────────────────────────────────────────────
|
|
|
|
/// `QuicSyncBackend::pull` receives all revisions from a full server.
|
|
#[tokio::test]
|
|
async fn quic_backend_pull_cold_copy() {
|
|
init_crypto();
|
|
let (_srv_dir, _srv_h5, srv_onion) = make_onion(5);
|
|
let (cli_dir, cli_h5, cli_onion) = make_onion(0);
|
|
|
|
let server = QuicServer::bind(any_addr(), QuicConfig::self_signed().unwrap())
|
|
.await
|
|
.unwrap();
|
|
let addr = server.local_addr;
|
|
let server_task = tokio::spawn(run_quic_server(server, srv_onion));
|
|
|
|
let backend = QuicSyncBackend::new(addr, "agent");
|
|
let stats = backend.pull(&cli_h5).await.unwrap();
|
|
|
|
assert_eq!(stats.revisions_transferred, 5);
|
|
|
|
server_task.await.unwrap();
|
|
let cli_after = OnionFile::open(&cli_h5).unwrap();
|
|
assert_eq!(cli_after.revision_count(), 5);
|
|
let _ = (cli_dir, cli_onion);
|
|
}
|
|
|
|
/// `QuicSyncBackend::pull` fetches only the missing delta.
|
|
#[tokio::test]
|
|
async fn quic_backend_pull_partial_delta() {
|
|
init_crypto();
|
|
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 = QuicServer::bind(any_addr(), QuicConfig::self_signed().unwrap())
|
|
.await
|
|
.unwrap();
|
|
let addr = server.local_addr;
|
|
let server_task = tokio::spawn(run_quic_server(server, srv_onion));
|
|
|
|
let backend = QuicSyncBackend::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;
|
|
}
|