use runsync_transfer::transport::quic;
use runsync_transfer::{
human_bytes, receive, send, Config, Progress, QuicTransport, Secrecy, Source, Transport,
};
use std::collections::BTreeMap;
use std::io::Write;
use std::net::SocketAddr;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::time::Instant;
type Any = Box<dyn std::error::Error + Send + Sync>;
#[tokio::main]
async fn main() -> Result<(), Any> {
tracing_subscriber::fmt()
.with_env_filter(
tracing_subscriber::EnvFilter::try_from_default_env()
.unwrap_or_else(|_| "runsync_transfer=info".into()),
)
.with_writer(std::io::stderr)
.init();
let args: Vec<String> = std::env::args().skip(1).collect();
let cmd = args.first().map(String::as_str).unwrap_or("help");
let opts = Opts::parse(&args[1.min(args.len())..]);
match cmd {
"gen" => gen(&opts),
"serve" => serve(&opts).await,
"send" => send_cmd(&opts).await,
"selftest" => selftest(&opts).await,
"verify" => verify(&opts),
_ => {
eprintln!("{}", include_str!("rst_usage.txt"));
Ok(())
}
}
}
struct Opts {
flags: BTreeMap<String, String>,
positional: Vec<String>,
}
const BOOL_FLAGS: &[&str] = &[
"no-compress",
"no-verify",
"no-resume",
"no-sparse",
"no-audio",
"no-prealloc",
"lz4",
"once",
"keep",
];
impl Opts {
fn parse(args: &[String]) -> Self {
let mut flags = BTreeMap::new();
let mut positional = Vec::new();
let mut i = 0;
while i < args.len() {
let a = &args[i];
if let Some(name) = a.strip_prefix("--") {
if BOOL_FLAGS.contains(&name) {
flags.insert(name.to_string(), "1".into());
i += 1;
continue;
}
match args.get(i + 1) {
Some(v) if !v.starts_with("--") => {
flags.insert(name.to_string(), v.clone());
i += 2;
}
_ => {
flags.insert(name.to_string(), "1".into());
i += 1;
}
}
} else {
positional.push(a.clone());
i += 1;
}
}
Self { flags, positional }
}
fn get(&self, k: &str) -> Option<&str> {
self.flags.get(k).map(String::as_str)
}
fn req(&self, k: &str) -> Result<&str, Any> {
self.get(k)
.ok_or_else(|| format!("missing required --{k}").into())
}
fn has(&self, k: &str) -> bool {
self.flags.contains_key(k)
}
fn num(&self, k: &str, default: u64) -> u64 {
self.get(k).and_then(parse_size).unwrap_or(default)
}
}
fn parse_size(s: &str) -> Option<u64> {
let s = s.trim();
let (num, mult) = match s.chars().last()? {
'k' | 'K' => (&s[..s.len() - 1], 1u64 << 10),
'm' | 'M' => (&s[..s.len() - 1], 1 << 20),
'g' | 'G' => (&s[..s.len() - 1], 1 << 30),
't' | 'T' => (&s[..s.len() - 1], 1u64 << 40),
_ => (s, 1),
};
num.trim().parse::<u64>().ok().map(|n| n * mult)
}
fn psk_from(opts: &Opts) -> Result<Secrecy, Any> {
match opts.get("psk") {
None => Ok(Secrecy::TransportOnly),
Some(hex) => {
let bytes = decode_hex(hex).ok_or("--psk must be 64 hex characters")?;
if bytes.len() != 32 {
return Err("--psk must decode to 32 bytes".into());
}
let mut k = [0u8; 32];
k.copy_from_slice(&bytes);
Ok(Secrecy::Psk(k))
}
}
}
fn decode_hex(s: &str) -> Option<Vec<u8>> {
if s.len() % 2 != 0 {
return None;
}
(0..s.len())
.step_by(2)
.map(|i| u8::from_str_radix(&s[i..i + 2], 16).ok())
.collect()
}
fn encode_hex(b: &[u8]) -> String {
b.iter().map(|x| format!("{x:02x}")).collect()
}
fn config_from(opts: &Opts) -> Result<Config, Any> {
let mut cfg = match opts.get("profile") {
Some("throughput") => Config::throughput(),
Some("bandwidth") => Config::bandwidth_saving(),
_ => Config::default(),
};
if let Some(v) = opts.get("chunk").and_then(parse_size) {
cfg.chunk_size = v as usize;
}
if let Some(v) = opts.get("streams").and_then(|s| s.parse().ok()) {
cfg.streams = v;
}
if let Some(v) = opts.get("workers").and_then(|s| s.parse().ok()) {
cfg.workers = v;
}
if let Some(v) = opts.get("queue").and_then(|s| s.parse().ok()) {
cfg.queue_depth = v;
}
if opts.has("no-compress") {
cfg.compression.mode = runsync_transfer::CompressionMode::Off;
}
if opts.has("lz4") {
cfg.compression.algorithm = runsync_transfer::Algorithm::Lz4;
}
if opts.has("no-verify") {
cfg.verify_hashes = false;
}
if opts.has("no-resume") {
cfg.resume = false;
}
if opts.has("no-audio") {
cfg.compression.audio_codec = false;
}
if opts.has("no-sparse") {
cfg.sparse = false;
}
if opts.has("no-prealloc") {
cfg.preallocate = false;
}
cfg.secrecy = psk_from(opts)?;
Ok(cfg)
}
fn progress_printer(label: &'static str) -> runsync_transfer::ProgressFn {
Arc::new(move |p: Progress| {
eprintln!("[{label}] {p}");
})
}
fn fill_prng(seed: u64, buf: &mut [u8]) {
let mut s = seed | 1;
for c in buf.chunks_mut(8) {
s ^= s << 13;
s ^= s >> 7;
s ^= s << 17;
let b = s.to_le_bytes();
c.copy_from_slice(&b[..c.len()]);
}
}
fn write_random(path: &Path, size: u64, seed: u64) -> Result<(), Any> {
if let Some(p) = path.parent() {
std::fs::create_dir_all(p)?;
}
let mut f = std::io::BufWriter::with_capacity(4 << 20, std::fs::File::create(path)?);
let mut block = vec![0u8; 4 << 20];
let mut written = 0u64;
let mut s = seed;
while written < size {
let n = ((size - written) as usize).min(block.len());
s = s.wrapping_mul(6364136223846793005).wrapping_add(1);
fill_prng(s, &mut block[..n]);
f.write_all(&block[..n])?;
written += n as u64;
}
f.flush()?;
Ok(())
}
fn write_textish(path: &Path, size: u64) -> Result<(), Any> {
if let Some(p) = path.parent() {
std::fs::create_dir_all(p)?;
}
let mut f = std::io::BufWriter::with_capacity(4 << 20, std::fs::File::create(path)?);
let line = b"2026-08-01T00:00:00Z INFO transfer chunk=000000 offset=00000000 status=ok\n";
let mut written = 0u64;
while written < size {
let n = ((size - written) as usize).min(line.len());
f.write_all(&line[..n])?;
written += n as u64;
}
f.flush()?;
Ok(())
}
fn write_pcm(path: &Path, size: u64) -> Result<(), Any> {
if let Some(p) = path.parent() {
std::fs::create_dir_all(p)?;
}
let mut f = std::io::BufWriter::with_capacity(4 << 20, std::fs::File::create(path)?);
let mut block = Vec::with_capacity(4 << 20);
let mut i: u64 = 0;
let mut written = 0u64;
while written < size {
block.clear();
while block.len() < (4 << 20) && written + block.len() as u64 <= size {
let t = i as f64 / 44_100.0;
let l = ((t * 440.0 * std::f64::consts::TAU).sin() * 12_000.0) as i16;
let r = ((t * 587.33 * std::f64::consts::TAU).sin() * 9_500.0) as i16;
block.extend_from_slice(&l.to_le_bytes());
block.extend_from_slice(&r.to_le_bytes());
i += 1;
}
let n = ((size - written) as usize).min(block.len());
f.write_all(&block[..n])?;
written += n as u64;
}
f.flush()?;
Ok(())
}
fn write_sparse(path: &Path, size: u64, data_bytes: u64) -> Result<(), Any> {
use std::io::{Seek, SeekFrom};
if let Some(p) = path.parent() {
std::fs::create_dir_all(p)?;
}
let mut f = std::fs::File::create(path)?;
f.set_len(size)?;
if data_bytes == 0 {
return Ok(());
}
let islands = 16u64;
let per = data_bytes / islands;
let stride = size / islands;
let mut block = vec![0u8; per.min(64 << 20) as usize];
for k in 0..islands {
fill_prng(k + 1, &mut block);
f.seek(SeekFrom::Start(k * stride))?;
let mut left = per;
while left > 0 {
let n = (left as usize).min(block.len());
f.write_all(&block[..n])?;
left -= n as u64;
}
}
f.flush()?;
Ok(())
}
fn gen(opts: &Opts) -> Result<(), Any> {
let out = PathBuf::from(opts.req("out")?);
let size = opts.num("size", 1 << 30);
let profile = opts.get("profile").unwrap_or("mixed");
std::fs::create_dir_all(&out)?;
let t0 = Instant::now();
match profile {
"large" => write_random(&out.join("large.bin"), size, 42)?,
"sparse" => {
let data = opts.num("data", size / 100);
write_sparse(&out.join("disk.img"), size, data)?;
}
"manysmall" => {
let count = opts.num("count", 20_000);
let each = (size / count).max(256);
for i in 0..count {
let d = out.join(format!("d{:03}", i % 200));
write_random(&d.join(format!("f{i:06}.bin")), each, i + 1)?;
}
}
_ => {
let q = size / 8;
write_textish(&out.join("logs/service.log"), q)?;
write_textish(&out.join("logs/access.log"), q / 2)?;
write_pcm(&out.join("audio/master.wav"), q)?;
write_pcm(&out.join("audio/session.wav"), q / 2)?;
write_random(&out.join("audio/album/track01.flac"), q / 2, 11)?;
write_random(&out.join("audio/album/track02.flac"), q / 2, 12)?;
write_random(&out.join("video/capture.mp4"), q, 21)?;
write_random(&out.join("images/scan.jpg"), q / 4, 31)?;
write_random(&out.join("blobs/opaque.dat"), q, 41)?;
write_random(&out.join("blobs/odd.bin"), q / 2 + 12_345, 51)?;
std::fs::write(out.join("empty.bin"), b"")?;
std::fs::write(out.join("tiny.txt"), b"x")?;
std::fs::create_dir_all(out.join("emptydir"))?;
for i in 0..500 {
std::fs::write(
out.join(format!("small/f{i:04}.txt")),
format!("file number {i}\n").repeat(20),
)
.or_else(|_| {
std::fs::create_dir_all(out.join("small"))?;
std::fs::write(
out.join(format!("small/f{i:04}.txt")),
format!("file number {i}\n").repeat(20),
)
})?;
}
}
}
let (files, bytes) = tree_stats(&out);
println!(
"generated {profile}: {files} files, {} in {:.1}s at {}",
human_bytes(bytes),
t0.elapsed().as_secs_f64(),
out.display()
);
Ok(())
}
fn tree_stats(root: &Path) -> (u64, u64) {
let mut files = 0;
let mut bytes = 0;
let mut stack = vec![root.to_path_buf()];
while let Some(d) = stack.pop() {
let Ok(rd) = std::fs::read_dir(&d) else {
continue;
};
for e in rd.flatten() {
let p = e.path();
match std::fs::symlink_metadata(&p) {
Ok(m) if m.is_dir() => stack.push(p),
Ok(m) if m.is_file() => {
files += 1;
bytes += m.len();
}
_ => {}
}
}
}
(files, bytes)
}
async fn serve(opts: &Opts) -> Result<(), Any> {
let addr: SocketAddr = opts.get("addr").unwrap_or("0.0.0.0:5555").parse()?;
let dest = PathBuf::from(opts.req("dest")?);
let cfg = config_from(opts)?;
let (cert, key) = quic::self_signed(vec!["localhost".into(), "rst".into()])?;
if let Some(p) = opts.get("cert-out") {
std::fs::write(p, cert.as_ref())?;
eprintln!("wrote certificate to {p}");
}
let ep = quic::server_endpoint(addr, vec![cert.clone()], key)?;
println!("CERT_HEX {}", encode_hex(cert.as_ref()));
println!("LISTENING {}", ep.local_addr()?);
std::io::stdout().flush()?;
let once = opts.has("once");
loop {
let Some(incoming) = ep.accept().await else {
break;
};
let conn = match incoming.await {
Ok(c) => c,
Err(e) => {
eprintln!("handshake failed: {e}");
continue;
}
};
let peer = conn.remote_address();
eprintln!("connection from {peer}");
let transport: Arc<dyn Transport> = Arc::new(QuicTransport::from_connection(conn));
let t0 = Instant::now();
let r = receive(
transport.clone(),
&dest,
&cfg,
Some(progress_printer("recv")),
)
.await;
match r {
Ok(p) => {
println!(
"RESULT ok files={} logical={} wire={} secs={:.2} rate={}/s ratio={:.2}",
p.files_completed,
p.logical_bytes,
p.wire_bytes,
t0.elapsed().as_secs_f64(),
human_bytes(p.throughput() as u64),
p.compression_ratio()
);
}
Err(e) => println!("RESULT error {e}"),
}
std::io::stdout().flush()?;
transport.close(0, b"done");
if once {
break;
}
}
Ok(())
}
async fn send_cmd(opts: &Opts) -> Result<(), Any> {
let addr: SocketAddr = opts.req("addr")?.parse()?;
let cfg = config_from(opts)?;
let cert = match (opts.get("cert"), opts.get("cert-hex")) {
(Some(p), _) => std::fs::read(p)?,
(None, Some(h)) => decode_hex(h).ok_or("--cert-hex is not valid hex")?,
_ => return Err("need --cert FILE or --cert-hex HEX".into()),
};
let cert = rustls::pki_types::CertificateDer::from(cert);
let sources: Vec<Source> = opts.positional.iter().map(Source::new).collect();
if sources.is_empty() {
return Err("no source paths given".into());
}
let bind: SocketAddr = if addr.is_ipv6() {
"[::]:0".parse()?
} else {
"0.0.0.0:0".parse()?
};
let ep = quic::client_endpoint(bind, cert)?;
let conn = ep.connect(addr, "localhost")?.await?;
let transport: Arc<dyn Transport> = Arc::new(QuicTransport::from_connection(conn));
let t0 = Instant::now();
let p = send(
transport.clone(),
&sources,
&cfg,
Some(progress_printer("send")),
)
.await?;
println!(
"RESULT ok files={} logical={} wire={} secs={:.2} rate={}/s ratio={:.2} udp={}",
p.files_completed,
p.logical_bytes,
p.wire_bytes,
t0.elapsed().as_secs_f64(),
human_bytes(p.throughput() as u64),
p.compression_ratio(),
transport.bytes_sent().unwrap_or(0)
);
transport.close(0, b"done");
tokio::time::sleep(std::time::Duration::from_millis(200)).await;
Ok(())
}
async fn selftest(opts: &Opts) -> Result<(), Any> {
let size = opts.num("size", 1 << 30);
let profile = opts.get("profile").unwrap_or("mixed").to_string();
let workdir = match opts.get("workdir") {
Some(d) => PathBuf::from(d),
None => std::env::temp_dir().join("rst-selftest"),
};
let _ = std::fs::remove_dir_all(&workdir);
let src = workdir.join("src");
let dest = workdir.join("dest");
std::fs::create_dir_all(&src)?;
std::fs::create_dir_all(&dest)?;
let mut gen_opts = Opts {
flags: opts.flags.clone(),
positional: vec![],
};
gen_opts
.flags
.insert("out".into(), src.display().to_string());
gen_opts.flags.insert("size".into(), size.to_string());
gen_opts.flags.insert("profile".into(), profile.clone());
gen(&gen_opts)?;
let cfg = config_from(opts)?;
let cfg2 = cfg.clone();
let dest2 = dest.clone();
if opts.get("transport") == Some("mem") {
let (ta, tb) = runsync_transfer::transport::mem::pair(4 << 20);
let ta: Arc<dyn Transport> = Arc::new(ta);
let tb: Arc<dyn Transport> = Arc::new(tb);
let src3 = src.clone();
let t0 = Instant::now();
let rh = tokio::spawn(async move {
receive(tb, &dest2, &cfg2, Some(progress_printer("recv"))).await
});
let sp = send(
ta,
&[Source::new(&src3)],
&cfg,
Some(progress_printer("send")),
)
.await?;
let rp = rh.await??;
let elapsed = t0.elapsed();
println!("--- selftest {profile} (mem transport) ---");
println!(" sent : {sp}");
println!(" received : {rp}");
println!(
" wall : {:.2}s effective {}/s on-wire {}/s ratio {:.2}x",
elapsed.as_secs_f64(),
human_bytes((sp.logical_bytes as f64 / elapsed.as_secs_f64()) as u64),
human_bytes((sp.wire_bytes as f64 / elapsed.as_secs_f64()) as u64),
sp.compression_ratio()
);
let mut vopts = Opts {
flags: BTreeMap::new(),
positional: vec![],
};
vopts.flags.insert("a".into(), src.display().to_string());
vopts
.flags
.insert("b".into(), dest.join("src").display().to_string());
verify(&vopts)?;
if !opts.has("keep") {
let _ = std::fs::remove_dir_all(&workdir);
}
return Ok(());
}
let (cert, key) = quic::self_signed(vec!["localhost".into()])?;
let server = quic::server_endpoint("127.0.0.1:0".parse()?, vec![cert.clone()], key)?;
let server_addr = server.local_addr()?;
let accept = tokio::spawn(async move {
let conn = server.accept().await.unwrap().await.unwrap();
let t: Arc<dyn Transport> = Arc::new(QuicTransport::from_connection(conn));
let r = receive(t.clone(), &dest2, &cfg2, Some(progress_printer("recv"))).await;
drop(server);
(r, t)
});
let client_ep = quic::client_endpoint("127.0.0.1:0".parse()?, cert)?;
let conn = client_ep.connect(server_addr, "localhost")?.await?;
let transport: Arc<dyn Transport> = Arc::new(QuicTransport::from_connection(conn));
let t0 = Instant::now();
let sp = send(
transport.clone(),
&[Source::new(&src)],
&cfg,
Some(progress_printer("send")),
)
.await?;
let (rp, _keep) = accept.await?;
let rp = rp?;
let elapsed = t0.elapsed();
println!("--- selftest {profile} ---");
println!(" sent : {}", sp);
println!(" received : {}", rp);
println!(
" wall : {:.2}s effective {}/s on-wire {}/s ratio {:.2}x",
elapsed.as_secs_f64(),
human_bytes((sp.logical_bytes as f64 / elapsed.as_secs_f64()) as u64),
human_bytes((sp.wire_bytes as f64 / elapsed.as_secs_f64()) as u64),
sp.compression_ratio()
);
let mut vopts = Opts {
flags: BTreeMap::new(),
positional: vec![],
};
vopts.flags.insert("a".into(), src.display().to_string());
vopts
.flags
.insert("b".into(), dest.join("src").display().to_string());
verify(&vopts)?;
if !opts.has("keep") {
let _ = std::fs::remove_dir_all(&workdir);
}
Ok(())
}
fn verify(opts: &Opts) -> Result<(), Any> {
let a = PathBuf::from(opts.req("a")?);
let b = PathBuf::from(opts.req("b")?);
let ha = hash_tree(&a);
let hb = hash_tree(&b);
let mut problems = Vec::new();
for (rel, h) in &ha {
match hb.get(rel) {
None => problems.push(format!("missing in dest: {rel}")),
Some(x) if x != h => problems.push(format!("content differs: {rel}")),
_ => {}
}
}
for rel in hb.keys() {
if !ha.contains_key(rel) {
problems.push(format!("extra in dest: {rel}"));
}
}
if problems.is_empty() {
println!("VERIFY ok {} files match", ha.len());
Ok(())
} else {
for p in problems.iter().take(20) {
println!("VERIFY {p}");
}
Err(format!(
"VERIFY failed: {} problems across {} files",
problems.len(),
ha.len()
)
.into())
}
}
fn hash_tree(root: &Path) -> BTreeMap<String, blake3::Hash> {
use rayon::prelude::*;
let mut paths = Vec::new();
let mut stack = vec![root.to_path_buf()];
while let Some(d) = stack.pop() {
let Ok(rd) = std::fs::read_dir(&d) else {
continue;
};
for e in rd.flatten() {
let p = e.path();
match std::fs::symlink_metadata(&p) {
Ok(m) if m.is_dir() => stack.push(p),
Ok(m) if m.is_file() => paths.push(p),
_ => {}
}
}
}
paths
.par_iter()
.filter_map(|p| {
let rel = p.strip_prefix(root).ok()?.to_string_lossy().into_owned();
if rel.contains(".rst-part") || rel.contains(".rst-state") || rel.contains(".rst-index")
{
return None;
}
Some((rel, hash_file(p)))
})
.collect::<Vec<_>>()
.into_iter()
.collect()
}
fn hash_file(p: &Path) -> blake3::Hash {
use std::io::Read;
let Ok(mut f) = std::fs::File::open(p) else {
return blake3::Hash::from([0u8; 32]);
};
let mut h = blake3::Hasher::new();
let mut buf = vec![0u8; 4 << 20];
loop {
match f.read(&mut buf) {
Ok(0) | Err(_) => break,
Ok(n) => h.update(&buf[..n]),
};
}
h.finalize()
}