runsync-transfer 2026.1.0

High-throughput P2P file transfer engine: adaptive compression, end-to-end AEAD, parallel chunked pipeline over QUIC or any async transport.
Documentation
//! Transfers over a real QUIC connection.
//!
//! The in-memory transport proves the protocol; this proves the QUIC adapter —
//! stream limits, flow control, and the shutdown handshake behave the same way
//! over a socket as they do over a duplex pipe.

#![cfg(feature = "test-certs")]

use runsync_transfer::transport::quic;
use runsync_transfer::{receive, send, Config, QuicTransport, Secrecy, Source, Transport};
use std::net::{Ipv4Addr, SocketAddr};
use std::path::Path;
use std::sync::Arc;

fn prng(seed: u64, n: usize) -> Vec<u8> {
    let mut s = seed | 1;
    let mut out = Vec::with_capacity(n);
    while out.len() < n {
        s ^= s << 13;
        s ^= s >> 7;
        s ^= s << 17;
        out.extend_from_slice(&s.to_le_bytes());
    }
    out.truncate(n);
    out
}

fn textish(n: usize) -> Vec<u8> {
    b"INFO chunk complete offset=0 len=1048576 file=archive.tar\n"
        .iter()
        .copied()
        .cycle()
        .take(n)
        .collect()
}

fn hash_file(p: &Path) -> blake3::Hash {
    use std::io::Read;
    let mut f = std::fs::File::open(p).unwrap();
    let mut h = blake3::Hasher::new();
    let mut buf = vec![0u8; 1 << 20];
    loop {
        let n = f.read(&mut buf).unwrap();
        if n == 0 {
            break;
        }
        h.update(&buf[..n]);
    }
    h.finalize()
}

/// Bring up a loopback QUIC pair and hand both ends back as transports.
async fn quic_pair() -> (Arc<dyn Transport>, Arc<dyn Transport>) {
    let (cert, key) = quic::self_signed(vec!["localhost".into()]).unwrap();
    let server = quic::server_endpoint(
        SocketAddr::from((Ipv4Addr::LOCALHOST, 0)),
        vec![cert.clone()],
        key,
    )
    .unwrap();
    let server_addr = server.local_addr().unwrap();

    let accept = tokio::spawn(async move {
        let incoming = server.accept().await.expect("no inbound connection");
        let conn = incoming.await.expect("handshake failed");
        // Hold the endpoint for the lifetime of the connection; dropping it
        // would close the socket underneath us.
        (conn, server)
    });

    let client_ep =
        quic::client_endpoint(SocketAddr::from((Ipv4Addr::LOCALHOST, 0)), cert).unwrap();
    let client_conn = client_ep
        .connect(server_addr, "localhost")
        .unwrap()
        .await
        .unwrap();

    let (server_conn, server_ep) = accept.await.unwrap();

    // Leak the endpoints for the duration of the test process; they must
    // outlive the connections and nothing else owns them.
    Box::leak(Box::new((client_ep, server_ep)));

    (
        Arc::new(QuicTransport::from_connection(client_conn)),
        Arc::new(QuicTransport::from_connection(server_conn)),
    )
}

#[tokio::test]
async fn quic_roundtrip_mixed_files() {
    let tmp = tempfile::tempdir().unwrap();
    let src = tmp.path().join("src");
    std::fs::create_dir_all(src.join("audio")).unwrap();
    std::fs::create_dir_all(src.join("logs")).unwrap();

    let files: Vec<(&str, Vec<u8>)> = vec![
        ("logs/server.log", textish(3_000_000)),
        ("audio/track.flac", prng(1, 2_500_000)),
        ("audio/master.wav", textish(1_800_000)),
        ("blob.bin", prng(2, 4_000_000)),
        ("tiny.txt", b"hello".to_vec()),
        ("empty.dat", Vec::new()),
    ];
    for (rel, data) in &files {
        std::fs::write(src.join(rel), data).unwrap();
    }
    let dest = tmp.path().join("dest");

    let (client, server) = quic_pair().await;
    // Hold both connections for the duration: the last handle to drop closes
    // the QUIC connection, and the engine deliberately never closes one it
    // did not open.
    let (_keep_c, _keep_s) = (client.clone(), server.clone());
    let cfg = Config::default();
    let cfg2 = cfg.clone();
    let src2 = src.clone();
    let dest2 = dest.clone();

    let sh = tokio::spawn(async move { send(client, &[Source::new(&src2)], &cfg, None).await });
    let rh = tokio::spawn(async move { receive(server, &dest2, &cfg2, None).await });
    let s = sh.await.unwrap().unwrap();
    let r = rh.await.unwrap().unwrap();

    for (rel, data) in &files {
        let got = std::fs::read(dest.join("src").join(rel)).unwrap();
        assert_eq!(&got, data, "mismatch at {rel}");
    }
    assert_eq!(s.files_completed, files.len() as u64);
    assert_eq!(r.files_completed, files.len() as u64);
    assert_eq!(s.logical_bytes, r.logical_bytes);
}

#[tokio::test]
async fn quic_roundtrip_encrypted_end_to_end() {
    let tmp = tempfile::tempdir().unwrap();
    let src = tmp.path().join("src");
    std::fs::create_dir_all(&src).unwrap();
    let data = prng(9, 6_000_000);
    std::fs::write(src.join("confidential.bin"), &data).unwrap();
    let dest = tmp.path().join("dest");

    let psk = runsync_transfer::crypto::random_key();
    let (client, server) = quic_pair().await;
    let (_keep_c, _keep_s) = (client.clone(), server.clone());
    let cfg = Config::default().with_secrecy(Secrecy::Psk(psk));
    let cfg2 = cfg.clone();
    let src2 = src.clone();
    let dest2 = dest.clone();

    let sh = tokio::spawn(async move { send(client, &[Source::new(&src2)], &cfg, None).await });
    let rh = tokio::spawn(async move { receive(server, &dest2, &cfg2, None).await });
    sh.await.unwrap().unwrap();
    rh.await.unwrap().unwrap();

    assert_eq!(
        std::fs::read(dest.join("src/confidential.bin")).unwrap(),
        data
    );
}

#[tokio::test]
async fn quic_handles_a_single_large_file() {
    // 512 MiB over loopback QUIC. Large enough to exercise flow control,
    // stream scheduling, and the buffer pool under sustained load.
    let tmp = tempfile::tempdir().unwrap();
    let src = tmp.path().join("src");
    std::fs::create_dir_all(&src).unwrap();
    let path = src.join("large.bin");
    {
        use std::io::Write;
        let mut f = std::io::BufWriter::new(std::fs::File::create(&path).unwrap());
        let block = prng(4242, 4 * 1024 * 1024);
        for _ in 0..128 {
            f.write_all(&block).unwrap();
        }
        f.flush().unwrap();
    }
    let size = std::fs::metadata(&path).unwrap().len();
    assert_eq!(size, 512 * 1024 * 1024);
    let dest = tmp.path().join("dest");

    let (client, server) = quic_pair().await;
    let (_keep_c, _keep_s) = (client.clone(), server.clone());
    let cfg = Config::throughput();
    let cfg2 = cfg.clone();
    let src2 = src.clone();
    let dest2 = dest.clone();

    let sh = tokio::spawn(async move { send(client, &[Source::new(&src2)], &cfg, None).await });
    let rh = tokio::spawn(async move { receive(server, &dest2, &cfg2, None).await });
    let s = sh.await.unwrap().unwrap();
    rh.await.unwrap().unwrap();

    let out = dest.join("src/large.bin");
    assert_eq!(std::fs::metadata(&out).unwrap().len(), size);
    assert_eq!(hash_file(&path), hash_file(&out));
    assert_eq!(s.logical_bytes, size);
    eprintln!(
        "512 MiB over loopback QUIC: {:.0} MiB/s",
        s.throughput() / (1024.0 * 1024.0)
    );
}

#[tokio::test]
async fn quic_transport_reports_bytes_sent() {
    let tmp = tempfile::tempdir().unwrap();
    let src = tmp.path().join("src");
    std::fs::create_dir_all(&src).unwrap();
    std::fs::write(src.join("a.bin"), prng(3, 1_000_000)).unwrap();
    let dest = tmp.path().join("dest");

    let (client, server) = quic_pair().await;
    let (probe, _keep_s) = (client.clone(), server.clone());
    assert_eq!(probe.bytes_sent().map(|_| true), Some(true));

    let cfg = Config::default();
    let cfg2 = cfg.clone();
    let src2 = src.clone();
    let dest2 = dest.clone();
    let sh = tokio::spawn(async move { send(client, &[Source::new(&src2)], &cfg, None).await });
    let rh = tokio::spawn(async move { receive(server, &dest2, &cfg2, None).await });
    sh.await.unwrap().unwrap();
    rh.await.unwrap().unwrap();

    // UDP bytes must exceed the payload: headers and acks are not free.
    let udp = probe.bytes_sent().unwrap();
    assert!(udp > 1_000_000, "only {udp} UDP bytes for a 1 MB payload");
}