Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions crates/rds-bench/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ transport-noq = ["rds-net/transport-noq", "dep:noq"]
[dependencies]
rds-observe.workspace = true
anyhow.workspace = true
blake3.workspace = true
clap.workspace = true
ed25519-dalek.workspace = true
iroh-relay = { workspace = true, features = ["server"] }
Expand Down
1 change: 1 addition & 0 deletions crates/rds-bench/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,4 +13,5 @@ pub mod capacity;
pub mod impair;
pub mod report;
pub mod scenario;
mod transfer;
pub mod world;
126 changes: 96 additions & 30 deletions crates/rds-bench/src/scenario.rs
Original file line number Diff line number Diff line change
Expand Up @@ -63,7 +63,7 @@ pub enum Scenario {
Handshake,
/// Established connection, N ping probes, RTT percentiles.
Ping,
/// Bulk stream throughput to a TCP discard target.
/// Bulk stream goodput through receiver byte/digest acknowledgement.
Transfer,
/// Handshake under 5% loss (interop `multiconnect` analogue).
Multiconnect,
Expand All @@ -83,7 +83,7 @@ impl Scenario {
match self {
Scenario::Handshake => "handshake",
Scenario::Ping => "ping",
Scenario::Transfer => "transfer",
Scenario::Transfer => "transfer-receiver-ack-v1",
Scenario::Multiconnect => "multiconnect",
Scenario::RelayFallback => "relay-fallback",
Scenario::Impaired => "impaired",
Expand Down Expand Up @@ -262,43 +262,56 @@ async fn ping(
})
}

/// `transfer_mib` MiB over one forwarded stream to a discard sink.
/// `transfer_mib` MiB over one forwarded stream, ending at a verified receipt.
async fn transfer(p: &Params) -> anyhow::Result<BenchReport> {
let world = World::spawn(Path::Direct, p.transport_backend()?)
.await
.context("spawn world")?;
let conn = tokio::time::timeout(
let total = p
.transfer_mib
.checked_mul(1024 * 1024)
.context("transfer size overflow")?;
anyhow::ensure!(total > 0, "transfer size must be positive");
anyhow::ensure!(!p.timeout.is_zero(), "transfer timeout must be positive");
let world = tokio::time::timeout(
p.timeout,
rds_cli::connect(&world.client, world.target.clone()),
World::spawn(Path::Direct, p.transport_backend()?),
)
.await
.context("connect timed out")?
.context("connect failed")?;
let (host, port) = world.discard_target();
let (mut send, _recv) = rds_cli::open_tcp(&conn, &host, port).await?;
let chunk = vec![0xABu8; 256 * 1024];
let mut written = 0u64;
let total = p.transfer_mib * 1024 * 1024;
let t0 = Instant::now();
while written < total {
let n = (total - written).min(chunk.len() as u64) as usize;
send.write_all(&chunk[..n]).await?;
written += n as u64;
}
send.finish()?;
// Wait until the receiver has drained: finish() returns when our
// side is done sending; add a grace read timeout on recv to bound it.
let elapsed = t0.elapsed();
let mib_s = written as f64 / (1024.0 * 1024.0) / elapsed.as_secs_f64();
let metrics = world.metrics_snapshot(Some(&conn));
world.close().await;
.context("world startup timed out")?
.context("spawn world")?;
// A single deadline covers connect, OpenTcp, upload, receipt and EOF.
let outcome = tokio::time::timeout(p.timeout, async {
let conn = rds_cli::connect(&world.client, world.target.clone())
.await
.context("connect failed")?;
let (host, port) = world.transfer_target();
let (mut send, mut recv) = rds_cli::open_tcp(&conn, &host, port).await?;
let measurement = crate::transfer::send_verified(&mut send, &mut recv, total).await?;
anyhow::Ok((measurement, world.metrics_snapshot(Some(&conn))))
})
.await
.context("transfer operation timed out")
.and_then(|result| result);
let cleanup = tokio::time::timeout(Duration::from_secs(5), world.close()).await;
let (measurement, mut metrics) = outcome?;
cleanup.context("world shutdown timed out")?;
let mib_s = measurement.bytes as f64 / (1024.0 * 1024.0) / measurement.elapsed.as_secs_f64();
metrics.insert("transfer_verified_bytes".into(), measurement.bytes);
metrics.insert(
"transfer_completion_ns".into(),
measurement
.elapsed
.as_nanos()
.try_into()
.unwrap_or(u64::MAX),
);
let mut notes = proxy_note(&world);
notes.push("receiver-ack-v1: byte count + BLAKE3 digest + EOF; payload generation/hash, upload and receipt are timed; connect/OpenTcp excluded; not comparable to historical sender-finish results".into());
Ok(BenchReport {
meta: meta("transfer", p, world.path_label(), None),
meta: meta(Scenario::Transfer.name(), p, world.path_label(), None),
rtt: None,
throughput_mib_s: Some(mib_s),
attempts: None,
metrics,
notes: proxy_note(&world),
notes,
})
}

Expand Down Expand Up @@ -492,3 +505,56 @@ fn impairment_of(path: &Path) -> Option<Impairment> {
_ => None,
}
}

#[cfg(test)]
mod tests {
use super::*;

async fn verified_transfer(backend: &str) {
let p = Params {
transfer_mib: 1,
backend: backend.into(),
..Params::default()
};
let report = tokio::time::timeout(Duration::from_secs(40), transfer(&p))
.await
.unwrap()
.unwrap();
assert_eq!(report.meta.scenario, "transfer-receiver-ack-v1");
assert_eq!(report.metrics["transfer_verified_bytes"], 1024 * 1024);
assert!(report.metrics["transfer_completion_ns"] > 0);
let rate = report.throughput_mib_s.unwrap();
assert!(rate.is_finite() && rate > 0.0);
}

#[tokio::test]
async fn transfer_receipt_crosses_real_iroh_and_forwarded_tcp() {
verified_transfer("iroh").await;
}

#[cfg(feature = "transport-noq")]
#[tokio::test]
async fn transfer_receipt_crosses_real_noq_and_forwarded_tcp() {
verified_transfer("noq").await;
}

#[tokio::test]
async fn invalid_transfer_parameters_fail_before_startup() {
for p in [
Params {
transfer_mib: 0,
..Params::default()
},
Params {
transfer_mib: u64::MAX,
..Params::default()
},
Params {
timeout: Duration::ZERO,
..Params::default()
},
] {
assert!(transfer(&p).await.is_err());
}
}
}
197 changes: 197 additions & 0 deletions crates/rds-bench/src/transfer.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,197 @@
//! Benchmark-only receipt: big-endian u64 byte count, BLAKE3 digest, EOF.
//! This is the local TCP target's protocol, not an RDS wire extension.

use std::time::{Duration, Instant};

use anyhow::{Context, ensure};
use tokio::io::{AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt};

const CHUNK_BYTES: usize = 64 * 1024;

pub(crate) struct Measurement {
pub bytes: u64,
pub elapsed: Duration,
}

/// Time payload generation, hashing, upload and receiver completion together.
/// The caller owns the absolute operation deadline and both streams.
pub(crate) async fn send_verified(
send: &mut (impl AsyncWrite + Unpin),
recv: &mut (impl AsyncRead + Unpin),
total: u64,
) -> anyhow::Result<Measurement> {
ensure!(total > 0, "transfer payload must be nonempty");
let mut generator = blake3::Hasher::new();
generator.update(b"rds-bench-transfer/v1");
let mut generator = generator.finalize_xof();
let mut digest = blake3::Hasher::new();
let mut chunk = vec![0; CHUNK_BYTES];
let started = Instant::now();
let mut written = 0;
while written < total {
let n = (total - written).min(CHUNK_BYTES as u64) as usize;
// Position-dependent data also exposes duplicated/reordered chunks.
generator.fill(&mut chunk[..n]);
digest.update(&chunk[..n]);
send.write_all(&chunk[..n]).await.context("write payload")?;
written += n as u64;
}
send.shutdown().await.context("finish payload")?;
let mut receipt = [0; 40];
recv.read_exact(&mut receipt)
.await
.context("read receiver receipt")?;
ensure!(
receipt[..8] == total.to_be_bytes(),
"receiver byte count mismatch"
);
ensure!(
digest.finalize().eq(&receipt[8..]),
"receiver digest mismatch"
);
ensure!(
recv.read(&mut [0]).await.context("read receipt EOF")? == 0,
"trailing receiver receipt data"
);
Ok(Measurement {
bytes: total,
elapsed: started.elapsed(),
})
}

/// Acknowledge only after all payload bytes and FIN have been consumed.
pub(crate) async fn receive(
socket: &mut (impl AsyncRead + AsyncWrite + Unpin),
) -> anyhow::Result<()> {
let mut buffer = vec![0; CHUNK_BYTES];
let mut received = 0u64;
let mut digest = blake3::Hasher::new();
loop {
let n = socket.read(&mut buffer).await?;
if n == 0 {
break;
}
received = received
.checked_add(n as u64)
.context("receiver count overflow")?;
digest.update(&buffer[..n]);
}
socket.write_all(&received.to_be_bytes()).await?;
socket.write_all(digest.finalize().as_bytes()).await?;
socket.shutdown().await?;
Ok(())
}

#[cfg(test)]
mod tests {
use super::*;

#[tokio::test]
async fn large_send_window_cannot_complete_before_receiver_ack() {
let total = CHUNK_BYTES as u64 * 2 + 17;
let (sender, mut receiver) = tokio::io::duplex(total as usize + 1);
let (mut read, mut write) = tokio::io::split(sender);
let work = send_verified(&mut write, &mut read, total);
tokio::pin!(work);
// The entire body fits the transport buffer. A sender-only timer
// would already finish, while the receiver has consumed nothing.
tokio::select! {
biased;
result = &mut work => panic!("completed without receipt: {:?}", result.err()),
() = tokio::time::sleep(Duration::from_millis(20)) => {},
}
let (sent, received) = tokio::time::timeout(Duration::from_secs(2), async {
tokio::join!(work, receive(&mut receiver))
})
.await
.unwrap();
received.unwrap();
let sent = sent.unwrap();
assert_eq!(sent.bytes, total);
assert!(sent.elapsed >= Duration::from_millis(20));
}

#[tokio::test]
async fn invalid_or_missing_receipts_never_produce_throughput() {
for fault in [
"count",
"digest",
"truncated",
"trailing",
"missing",
"corrupt_body",
"short_body",
"reordered_body",
] {
let (sender, mut receiver) = tokio::io::duplex(1024);
let (mut read, mut write) = tokio::io::split(sender);
let (result, ()) = tokio::time::timeout(Duration::from_secs(2), async {
tokio::join!(send_verified(&mut write, &mut read, 257), async {
let mut body = Vec::new();
receiver.read_to_end(&mut body).await.unwrap();
assert_eq!(body.len(), 257);
match fault {
"corrupt_body" => body[0] ^= 1,
"short_body" => {
body.pop();
}
"reordered_body" => body.rotate_left(1),
_ => {}
}
let mut receipt = (body.len() as u64).to_be_bytes().to_vec();
receipt.extend_from_slice(blake3::hash(&body).as_bytes());
match fault {
"count" => receipt[7] ^= 1,
"digest" => receipt[8] ^= 1,
"truncated" => {
receipt.pop();
}
"trailing" => receipt.push(0),
"missing" => receipt.clear(),
"corrupt_body" | "short_body" | "reordered_body" => {}
_ => unreachable!(),
}
receiver.write_all(&receipt).await.unwrap();
receiver.shutdown().await.unwrap();
})
})
.await
.unwrap();
assert!(result.is_err(), "accepted {fault} receipt");
}
}

#[tokio::test]
async fn missing_receipt_or_eof_remains_subject_to_the_operation_deadline() {
for ack_before_stall in [false, true] {
let (sender, mut receiver) = tokio::io::duplex(1024);
let (mut read, mut write) = tokio::io::split(sender);
let (result, ()) = tokio::time::timeout(Duration::from_secs(2), async {
tokio::join!(
tokio::time::timeout(
Duration::from_millis(20),
send_verified(&mut write, &mut read, 257)
),
async {
let mut body = Vec::new();
receiver.read_to_end(&mut body).await.unwrap();
if ack_before_stall {
receiver
.write_all(&(body.len() as u64).to_be_bytes())
.await
.unwrap();
receiver
.write_all(blake3::hash(&body).as_bytes())
.await
.unwrap();
}
// Keep the stream alive without EOF until the deadline.
}
)
})
.await
.unwrap();
assert!(result.is_err());
}
}
}
Loading
Loading