mod common;
use std::net::SocketAddr;
use std::sync::Arc;
use std::time::Duration;
use common::{Certs, Server};
use tokio::io::AsyncReadExt;
use tokio::sync::mpsc;
use tokio::time::timeout;
use weida::{Binding, Error, Limits, Listener, LossCause, Runtime, RuntimeConfig, TransferMeta};
const DEADLINE: Duration = Duration::from_secs(15);
const STALL: Duration = Duration::from_millis(300);
async fn within<F: Future>(f: F) -> F::Output {
timeout(DEADLINE, f).await.expect("operation timed out")
}
#[derive(Debug, PartialEq, Eq)]
enum Stage {
Opened,
Wrote,
Delivered,
}
struct ManualServer {
_runtime: Runtime,
listener: Listener,
binding: Binding,
addr: SocketAddr,
}
impl ManualServer {
async fn start(certs: &Certs, config: RuntimeConfig) -> ManualServer {
let runtime = Runtime::new(config).expect("runtime");
let listener = runtime.listener();
let binding = listener
.bind_quic(
"127.0.0.1:0".parse().expect("loopback address"),
certs.server_tls(),
)
.await
.expect("bind");
let addr = binding.local_addr();
ManualServer {
_runtime: runtime,
listener,
binding,
addr,
}
}
fn url(&self, path: &str) -> String {
format!("weida://127.0.0.1:{}{}", self.addr.port(), path)
}
}
#[tokio::test]
async fn a_finished_transfer_needs_no_local_handle_to_arrive() {
let server = Server::start().await;
let puller = server.listener.puller("/jobs").expect("puller");
let client = server.client_runtime();
let payload = vec![0xa5u8; 64 * 1024];
{
let pusher = client.pusher(server.trust());
within(pusher.connect(&server.url("/jobs")))
.await
.expect("connect");
let mut transfer = within(pusher.open(TransferMeta::default()))
.await
.expect("open");
within(transfer.write_all(&payload)).await.expect("write");
drop(transfer.finish().expect("finish"));
}
let inbound = within(puller.recv()).await.expect("recv");
let body = within(inbound.collect(128 * 1024)).await.expect("collect");
assert_eq!(body.len(), payload.len());
assert_eq!(body, payload);
client.shutdown().await;
}
#[tokio::test]
async fn a_receipt_beyond_the_window_implies_the_reader_consumed() {
const WINDOW: usize = 64 * 1024;
const TOTAL: usize = WINDOW * 5 / 2;
const CHUNK: usize = 8 * 1024;
let server = Server::start_with(Limits {
stream_receive_window: WINDOW as u64,
connection_receive_window: 1024 * 1024,
..Limits::default()
})
.await;
let puller = server.listener.puller("/jobs").expect("puller");
let client = server.client_runtime();
let pusher = Arc::new(client.pusher(server.trust()));
within(pusher.connect(&server.url("/jobs")))
.await
.expect("connect");
let (tx, mut rx) = mpsc::unbounded_channel();
let sender = {
let pusher = Arc::clone(&pusher);
tokio::spawn(async move {
let mut transfer = pusher.open(TransferMeta::default()).await.expect("open");
let _ = tx.send(Stage::Opened);
transfer
.write_all(&vec![0x7eu8; TOTAL])
.await
.expect("write");
let _ = tx.send(Stage::Wrote);
transfer
.finish()
.expect("finish")
.delivered()
.await
.expect("delivered");
let _ = tx.send(Stage::Delivered);
})
};
assert_eq!(within(rx.recv()).await, Some(Stage::Opened));
assert!(
timeout(STALL, rx.recv()).await.is_err(),
"a payload past the window must stall while the reader reads nothing"
);
let mut inbound = within(puller.recv()).await.expect("recv");
let mut chunk = vec![0u8; CHUNK];
let mut consumed = 0usize;
let mut wrote_at = None;
let mut delivered_at = None;
while consumed < TOTAL {
within(inbound.read_exact(&mut chunk)).await.expect("read");
consumed += CHUNK;
while let Ok(stage) = rx.try_recv() {
match stage {
Stage::Wrote => wrote_at = Some(consumed),
Stage::Delivered => delivered_at = Some(consumed),
Stage::Opened => unreachable!("already observed"),
}
}
}
while delivered_at.is_none() {
match within(rx.recv()).await.expect("sender reported") {
Stage::Wrote => wrote_at = Some(consumed),
Stage::Delivered => delivered_at = Some(consumed),
Stage::Opened => unreachable!("already observed"),
}
}
let wrote_at = wrote_at.expect("the write completed");
let delivered_at = delivered_at.expect("the receipt resolved");
assert!(
wrote_at >= TOTAL - WINDOW,
"the write completed after {wrote_at} bytes were consumed, less than the {} the window allows",
TOTAL - WINDOW
);
assert!(
delivered_at >= TOTAL - WINDOW,
"the receipt resolved after {delivered_at} bytes were consumed, less than the {} the window allows",
TOTAL - WINDOW
);
within(sender).await.expect("sender task");
client.shutdown().await;
}
#[tokio::test]
async fn a_payload_the_size_of_the_stream_window_waits_for_the_reader() {
const WINDOW: usize = 64 * 1024;
let server = Server::start_with(Limits {
stream_receive_window: WINDOW as u64,
connection_receive_window: 1024 * 1024,
..Limits::default()
})
.await;
let puller = server.listener.puller("/jobs").expect("puller");
let client = server.client_runtime();
let pusher = Arc::new(client.pusher(server.trust()));
within(pusher.connect(&server.url("/jobs")))
.await
.expect("connect");
let (tx, mut rx) = mpsc::unbounded_channel();
let sender = {
let pusher = Arc::clone(&pusher);
tokio::spawn(async move {
let mut transfer = pusher.open(TransferMeta::default()).await.expect("open");
let _ = tx.send(Stage::Opened);
transfer
.write_all(&vec![0x11u8; WINDOW])
.await
.expect("write");
let _ = tx.send(Stage::Wrote);
})
};
assert_eq!(within(rx.recv()).await, Some(Stage::Opened));
assert!(
timeout(STALL, rx.recv()).await.is_err(),
"the header shares the window, so a full-window payload cannot complete unread"
);
let mut inbound = within(puller.recv()).await.expect("recv");
let mut eighth = vec![0u8; WINDOW / 8];
within(inbound.read_exact(&mut eighth)).await.expect("read");
assert_eq!(
within(rx.recv()).await,
Some(Stage::Wrote),
"one eighth of the window is the update threshold, so the write resumes"
);
within(sender).await.expect("sender task");
client.shutdown().await;
}
#[tokio::test]
async fn a_stalled_stream_does_not_block_its_siblings() {
const STREAM_WINDOW: usize = 64 * 1024;
const PAYLOAD: usize = 32 * 1024;
const CONN_WINDOW: usize = 256 * 1024;
const CAPACITY: usize = CONN_WINDOW / PAYLOAD;
let server = Server::start_with(Limits {
stream_receive_window: STREAM_WINDOW as u64,
connection_receive_window: CONN_WINDOW as u64,
..Limits::default()
})
.await;
let puller = server.listener.puller("/jobs").expect("puller");
let client = server.client_runtime();
let pusher = Arc::new(client.pusher(server.trust()));
within(pusher.connect(&server.url("/jobs")))
.await
.expect("connect");
let mut a = within(pusher.open(TransferMeta::default()))
.await
.expect("open A");
within(a.write_all(&vec![0xaau8; PAYLOAD]))
.await
.expect("write A");
let mut inbound_a = within(puller.recv()).await.expect("recv A");
within(pusher.send(b"B")).await.expect("send B");
within(pusher.send(b"C")).await.expect("send C");
let b = within(puller.recv()).await.expect("recv B");
assert_eq!(within(b.collect(64)).await.expect("collect B"), b"B");
let c = within(puller.recv()).await.expect("recv C");
assert_eq!(within(c.collect(64)).await.expect("collect C"), b"C");
let (tx, mut rx) = mpsc::unbounded_channel();
let mut writers = Vec::new();
let mut blocked_at = None;
for i in 0..CAPACITY + 4 {
let mut transfer = within(pusher.open(TransferMeta::default()))
.await
.expect("open filler");
let tx = tx.clone();
writers.push(tokio::spawn(async move {
transfer
.write_all(&vec![0xbbu8; PAYLOAD])
.await
.expect("write filler");
let _ = tx.send(i);
std::future::pending::<()>().await;
}));
match timeout(STALL, rx.recv()).await {
Ok(Some(done)) => assert_eq!(done, i, "writers report in order"),
Ok(None) => unreachable!("a writer holds the sender"),
Err(_) => {
blocked_at = Some(i);
break;
}
}
}
let blocked_at = blocked_at.expect("the connection window must run out");
let unread = 1 + blocked_at;
assert!(
unread <= CAPACITY,
"{unread} unread {PAYLOAD}-byte streams fit in a {CONN_WINDOW}-byte connection window"
);
let mut drained = vec![0u8; PAYLOAD];
within(inbound_a.read_exact(&mut drained))
.await
.expect("read A");
assert!(
drained.iter().all(|&b| b == 0xaa),
"a stream parked behind flow control loses nothing"
);
assert_eq!(
within(rx.recv()).await,
Some(blocked_at),
"the stalled write must complete once the slow reader consumes"
);
client.shutdown().await;
}
#[tokio::test]
async fn a_stalled_path_does_not_stall_another_path() {
const STREAM_WINDOW: usize = 64 * 1024;
const PAYLOAD: usize = 32 * 1024;
const CONN_WINDOW: usize = 256 * 1024;
const CAPACITY: usize = CONN_WINDOW / PAYLOAD;
let server = Server::start_with(Limits {
stream_receive_window: STREAM_WINDOW as u64,
connection_receive_window: CONN_WINDOW as u64,
..Limits::default()
})
.await;
let slow = server.listener.puller("/slow").expect("slow puller");
let fast = server.listener.puller("/fast").expect("fast puller");
let client = server.client_runtime();
let slow_only = client.pusher(server.trust());
within(slow_only.connect(&server.url("/slow")))
.await
.expect("connect slow");
let fast_only = client.pusher(server.trust());
within(fast_only.connect(&server.url("/fast")))
.await
.expect("connect fast");
let (tx, mut rx) = mpsc::unbounded_channel();
let mut writers = Vec::new();
let mut blocked_at = None;
for i in 0..CAPACITY + 4 {
let mut transfer = within(slow_only.open(TransferMeta::default()))
.await
.expect("open filler");
let tx = tx.clone();
writers.push(tokio::spawn(async move {
transfer
.write_all(&vec![0xccu8; PAYLOAD])
.await
.expect("write filler");
let _ = tx.send(i);
std::future::pending::<()>().await;
}));
match timeout(STALL, rx.recv()).await {
Ok(Some(done)) => assert_eq!(done, i, "writers report in order"),
Ok(None) => unreachable!("a writer holds the sender"),
Err(_) => {
blocked_at = Some(i);
break;
}
}
}
let blocked_at = blocked_at.expect("the slow path's connection window must run out");
assert!(blocked_at <= CAPACITY, "{blocked_at} unread streams");
within(fast_only.send(b"fast"))
.await
.expect("send on /fast");
let arrived = within(fast.recv()).await.expect("recv on /fast");
assert_eq!(
within(arrived.collect(64)).await.expect("collect"),
b"fast",
"a path with its own connection is unaffected by another path's stall"
);
let mut inbound = within(slow.recv()).await.expect("recv on /slow");
let mut drained = vec![0u8; PAYLOAD];
within(inbound.read_exact(&mut drained))
.await
.expect("read the slow path");
assert_eq!(
within(rx.recv()).await,
Some(blocked_at),
"the stalled write completes once its own path's reader consumes"
);
client.shutdown().await;
}
#[tokio::test]
async fn the_stream_budget_is_backpressure_not_an_error() {
let server = Server::start_with_config(RuntimeConfig {
limits: Limits {
max_concurrent_uni_streams: 2,
..Limits::default()
},
endpoint_queue: 1,
..RuntimeConfig::default()
})
.await;
let puller = server.listener.puller("/jobs").expect("puller");
let client = server.client_runtime();
let pusher = Arc::new(client.pusher(server.trust()));
within(pusher.connect(&server.url("/jobs")))
.await
.expect("connect");
within(pusher.send(b"one")).await.expect("send one");
within(pusher.send(b"two")).await.expect("send two");
let (tx, mut rx) = mpsc::unbounded_channel();
let third = {
let pusher = Arc::clone(&pusher);
tokio::spawn(async move {
let mut transfer = pusher.open(TransferMeta::default()).await.expect("open");
let _ = tx.send(Stage::Opened);
transfer.write_all(b"three").await.expect("write");
let _ = tx.send(Stage::Wrote);
transfer
.finish()
.expect("finish")
.delivered()
.await
.expect("delivered");
let _ = tx.send(Stage::Delivered);
})
};
assert!(
timeout(STALL, rx.recv()).await.is_err(),
"a spent stream budget must block the third transfer, not fail it"
);
let first = within(puller.recv()).await.expect("recv one");
assert_eq!(within(first.collect(64)).await.expect("collect"), b"one");
assert_eq!(within(rx.recv()).await, Some(Stage::Opened));
assert_eq!(within(rx.recv()).await, Some(Stage::Wrote));
assert_eq!(within(rx.recv()).await, Some(Stage::Delivered));
within(third).await.expect("third task");
client.shutdown().await;
}
#[tokio::test]
async fn a_deeper_endpoint_queue_does_not_raise_the_stream_budget() {
let server = Server::start_with_config(RuntimeConfig {
limits: Limits {
max_concurrent_uni_streams: 2,
..Limits::default()
},
endpoint_queue: 8,
..RuntimeConfig::default()
})
.await;
let puller = server.listener.puller("/jobs").expect("puller");
let client = server.client_runtime();
let pusher = Arc::new(client.pusher(server.trust()));
within(pusher.connect(&server.url("/jobs")))
.await
.expect("connect");
within(pusher.send(b"one")).await.expect("send one");
within(pusher.send(b"two")).await.expect("send two");
let (tx, mut rx) = mpsc::unbounded_channel();
{
let pusher = Arc::clone(&pusher);
tokio::spawn(async move {
let mut transfer = pusher.open(TransferMeta::default()).await.expect("open");
let _ = tx.send(Stage::Opened);
transfer.write_all(b"three").await.expect("write");
let _ = tx.send(Stage::Wrote);
});
}
assert!(
timeout(STALL, rx.recv()).await.is_err(),
"queue depth is not stream credit: the third transfer must still block"
);
let first = within(puller.recv()).await.expect("recv one");
assert_eq!(within(first.collect(64)).await.expect("collect"), b"one");
assert_eq!(within(rx.recv()).await, Some(Stage::Opened));
assert_eq!(within(rx.recv()).await, Some(Stage::Wrote));
client.shutdown().await;
}
#[tokio::test]
async fn cancel_discards_unread_bytes_and_keeps_read_ones() {
let server = Server::start().await;
let puller = server.listener.puller("/jobs").expect("puller");
let client = server.client_runtime();
let pusher = client.pusher(server.trust());
within(pusher.connect(&server.url("/jobs")))
.await
.expect("connect");
let payload: Vec<u8> = (0..8 * 1024u32).map(|i| (i % 251) as u8).collect();
let mut transfer = within(pusher.open(TransferMeta::default()))
.await
.expect("open");
within(transfer.write_all(&payload)).await.expect("write");
let mut inbound = within(puller.recv()).await.expect("recv");
let mut head = vec![0u8; 4 * 1024];
within(inbound.read_exact(&mut head)).await.expect("read");
assert_eq!(head, payload[..4 * 1024]);
transfer.cancel();
let mut scratch = vec![0u8; 1024];
let mut after_cancel = 0usize;
let err = loop {
match within(inbound.read(&mut scratch)).await {
Ok(0) => panic!("a canceled transfer must never surface as EOF"),
Ok(n) => after_cancel += n,
Err(e) => break e,
}
};
assert_eq!(
err.kind(),
std::io::ErrorKind::ConnectionReset,
"quinn maps RESET_STREAM to ConnectionReset: {err:?}"
);
assert!(
after_cancel <= 4 * 1024,
"the reset cannot deliver more than the payload's remainder"
);
assert_eq!(head, payload[..4 * 1024]);
let mut transfer = within(pusher.open(TransferMeta::default()))
.await
.expect("open again");
within(transfer.write_all(&payload)).await.expect("write");
let mut inbound = within(puller.recv()).await.expect("recv again");
let mut head = vec![0u8; 4 * 1024];
within(inbound.read_exact(&mut head)).await.expect("read");
transfer.cancel();
let err = within(inbound.read_capped(64 * 1024))
.await
.expect_err("a reset stream must not read as complete");
assert!(matches!(err, Error::Canceled), "{err:?}");
within(pusher.send(b"next")).await.expect("send next");
let next = within(puller.recv()).await.expect("recv next");
assert_eq!(within(next.collect(64)).await.expect("collect"), b"next");
client.shutdown().await;
}
#[tokio::test]
async fn a_server_that_goes_away_is_reported_as_connection_lost() {
let certs = Certs::generate();
let first = ManualServer::start(&certs, RuntimeConfig::default()).await;
let puller = first.listener.puller("/jobs").expect("puller");
let client = Runtime::new(RuntimeConfig {
reconnect: weida::ReconnectPolicy::never(),
..RuntimeConfig::default()
})
.expect("client runtime");
let pusher = client.pusher(certs.client_tls());
within(pusher.connect(&first.url("/jobs")))
.await
.expect("connect");
within(pusher.send(b"before")).await.expect("send before");
let inbound = within(puller.recv()).await.expect("recv before");
assert_eq!(
within(inbound.collect(64)).await.expect("collect"),
b"before"
);
first.binding.close().await;
drop(puller);
drop(first);
let mut attempts = 0usize;
let err = within(async {
loop {
attempts += 1;
match pusher.send(b"orphan").await {
Err(e) => return e,
Ok(()) => tokio::task::yield_now().await,
}
}
})
.await;
assert!(
matches!(err, Error::ConnectionLost(_)),
"a closed peer must be reported as ConnectionLost, got {err:?}"
);
assert_eq!(attempts, 1, "the first send after the close already fails");
assert_eq!(pusher.peer_count(), 0, "a dead peer does not count");
let second = ManualServer::start(&certs, RuntimeConfig::default()).await;
let puller = second.listener.puller("/jobs").expect("puller");
within(pusher.connect(&second.url("/jobs")))
.await
.expect("connect again");
within(pusher.send(b"after")).await.expect("send after");
let inbound = within(puller.recv()).await.expect("recv after");
assert_eq!(
within(inbound.collect(64)).await.expect("collect"),
b"after"
);
assert_eq!(pusher.peer_count(), 1, "only the live peer counts");
client.shutdown().await;
}
#[tokio::test]
async fn idle_timeout_reports_loss_within_the_window() {
let certs = Certs::generate();
let server = ManualServer::start(
&certs,
RuntimeConfig {
limits: Limits {
idle_timeout: Duration::from_millis(500),
..Limits::default()
},
..RuntimeConfig::default()
},
)
.await;
let replier = server.listener.replier("/rpc").expect("replier");
let replier = std::sync::Arc::new(replier);
let handler = tokio::spawn({
let replier = std::sync::Arc::clone(&replier);
async move {
let request = replier.accept().await.expect("accept");
let mut out = request
.reply(TransferMeta::default())
.await
.expect("open reply");
out.write_all(b"pong").await.expect("write reply");
out.finish().expect("finish reply");
}
});
let client = Runtime::new(RuntimeConfig::default()).expect("client runtime");
assert!(
client.config().limits.keep_alive > Duration::from_millis(500),
"the client's keep-alive must be too slow to save this connection"
);
let requester = client.requester(certs.client_tls());
let mut events = requester.events();
within(requester.connect(&server.url("/rpc")))
.await
.expect("connect");
assert!(matches!(
within(events.recv()).await,
Some(weida::PeerEvent::Connected { .. })
));
let reply = within(requester.request(b"ping")).await.expect("request");
assert_eq!(within(reply.collect(64)).await.expect("collect"), b"pong");
within(handler).await.expect("handler");
tokio::time::sleep(Duration::from_millis(1500)).await;
let lost = within(events.recv()).await;
assert!(
matches!(
lost,
Some(weida::PeerEvent::Lost {
cause: LossCause::IdleTimeout,
..
})
),
"a timed-out peer must surface as an idle timeout, got {lost:?}"
);
let handler = tokio::spawn(async move {
let request = replier.accept().await.expect("accept");
let mut out = request
.reply(TransferMeta::default())
.await
.expect("open reply");
out.write_all(b"pong").await.expect("write reply");
out.finish().expect("finish reply");
});
let reply = within(requester.request(b"ping again"))
.await
.expect("the redialled connection serves the request");
assert_eq!(within(reply.collect(64)).await.expect("collect"), b"pong");
within(handler).await.expect("handler");
client.shutdown().await;
}
#[tokio::test]
async fn a_replier_that_stops_accepting_stalls_requesters_after_the_queue_fills() {
const BUDGET: usize = 2;
let server = Server::start_with_config(RuntimeConfig {
limits: Limits {
max_concurrent_bidi_streams: BUDGET as u32,
..Limits::default()
},
endpoint_queue: 1,
..RuntimeConfig::default()
})
.await;
let replier = server.listener.replier("/rpc").expect("replier");
let client = server.client_runtime();
let requester = Arc::new(client.requester(server.trust()));
within(requester.connect(&server.url("/rpc")))
.await
.expect("connect");
let mut replies = Vec::new();
let (tx, mut rx) = mpsc::unbounded_channel();
let mut blocked_at = None;
for i in 0..BUDGET + 2 {
let requester = Arc::clone(&requester);
let tx = tx.clone();
let started = tokio::spawn(async move {
let (mut transfer, reply) =
requester.open(TransferMeta::default()).await.expect("open");
transfer.write_all(b"req").await.expect("write");
transfer.finish().expect("finish");
let _ = tx.send(i);
reply
});
match timeout(STALL, rx.recv()).await {
Ok(Some(done)) => {
assert_eq!(done, i, "exchanges report in order");
replies.push(within(started).await.expect("exchange task"));
}
Ok(None) => unreachable!("a task holds the sender"),
Err(_) => {
blocked_at = Some((i, started));
break;
}
}
}
let (blocked_at, blocked_task) = blocked_at.expect("the stream budget must run out");
assert_eq!(
blocked_at, BUDGET,
"exactly the bidi budget worth of exchanges gets through; the one-deep \
accept queue does not add a slot, because a queued request still owns \
its stream"
);
assert_eq!(replies.len(), BUDGET);
let accepted = within(replier.accept()).await.expect("accept");
assert_eq!(accepted.meta().endpoint.as_deref(), Some("/rpc"));
drop(accepted);
assert_eq!(
within(rx.recv()).await,
Some(blocked_at),
"the blocked exchange must proceed once a request is accepted"
);
replies.push(within(blocked_task).await.expect("blocked task"));
client.shutdown().await;
}
struct Reorder {
arrivals: Vec<u64>,
out_of_order: usize,
peak_held: usize,
first_divergence: Option<usize>,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum Finish {
Batched,
Sequenced,
}
async fn reorder_probe(n: usize, finish: Finish) -> Reorder {
let server = Server::start_with_config(RuntimeConfig {
endpoint_queue: n + 1,
..RuntimeConfig::default()
})
.await;
let puller = server.listener.puller("/reorder").expect("puller");
let client = server.client_runtime();
let pusher = client.pusher(server.trust());
within(pusher.connect(&server.url("/reorder")))
.await
.expect("connect");
let (tx, mut rx) = mpsc::unbounded_channel();
let reader = tokio::spawn(async move {
for _ in 0..n {
let transfer = puller.recv().await.expect("recv");
let tx = tx.clone();
tokio::spawn(async move {
let body = transfer.collect(64).await.expect("collect");
let bytes: [u8; 8] = body[..].try_into().expect("8-byte sequence");
let seq = u64::from_le_bytes(bytes);
let _ = tx.send(seq);
});
}
puller
});
let mut transfers = Vec::with_capacity(n);
for seq in 0..n as u64 {
let mut transfer = within(pusher.open(TransferMeta::default()))
.await
.expect("open");
within(transfer.write_all(&seq.to_le_bytes()))
.await
.expect("write");
transfers.push(transfer);
}
for transfer in transfers.into_iter().rev() {
let delivery = transfer.finish().expect("finish");
if finish == Finish::Sequenced {
within(delivery.delivered()).await.expect("delivered");
}
}
let mut arrivals = Vec::with_capacity(n);
for _ in 0..n {
arrivals.push(within(rx.recv()).await.expect("arrival"));
}
let _puller = within(reader).await.expect("reader task");
let mut held = std::collections::BTreeSet::new();
let mut next_expected = 0u64;
let mut peak_held = 0usize;
for &seq in &arrivals {
held.insert(seq);
while held.remove(&next_expected) {
next_expected += 1;
}
peak_held = peak_held.max(held.len());
}
assert!(
held.is_empty() && next_expected == n as u64,
"the reorder buffer must drain: {} held, next_expected {next_expected}",
held.len()
);
let diverged: Vec<usize> = arrivals
.iter()
.enumerate()
.filter(|(position, seq)| **seq != *position as u64)
.map(|(position, _)| position)
.collect();
client.shutdown().await;
Reorder {
arrivals,
out_of_order: diverged.len(),
peak_held,
first_divergence: diverged.first().copied(),
}
}
#[tokio::test]
async fn reverse_order_completion_measures_the_reorder_buffer() {
let cases = [
(16usize, Finish::Batched),
(256, Finish::Batched),
(16, Finish::Sequenced),
];
for (n, finish) in cases {
let probe = within(reorder_probe(n, finish)).await;
assert_eq!(probe.arrivals.len(), n, "every transfer must arrive");
let mut seen = probe.arrivals.clone();
seen.sort_unstable();
assert!(
seen.iter().copied().eq(0..n as u64),
"each dispatched transfer must arrive exactly once"
);
assert!(probe.peak_held < n, "the buffer cannot hold every transfer");
eprintln!(
"reorder n={n} {finish:?}: {} of {n} arrived out of dispatch position, peak reorder \
buffer {} transfers, first divergence at position {:?}, first five arrivals {:?}",
probe.out_of_order,
probe.peak_held,
probe.first_divergence,
&probe.arrivals[..probe.arrivals.len().min(5)],
);
}
}