mod common;
use std::time::Duration;
use common::Server;
use tokio::io::AsyncReadExt;
use weida::{
Acknowledgement, CursorLevel, Error, Limits, RuntimeConfig, TransferMeta, TransferMeta as Meta,
};
const DEADLINE: Duration = Duration::from_secs(15);
async fn within<F: Future>(f: F) -> F::Output {
tokio::time::timeout(DEADLINE, f)
.await
.expect("operation timed out")
}
const STORED: CursorLevel = CursorLevel::Known(Acknowledgement::Stored);
#[tokio::test]
async fn a_transfer_interrupted_before_fin_is_connection_lost_for_the_sender() {
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 mut transfer = within(pusher.open(Meta::default())).await.expect("open");
within(transfer.write_all(b"the first half"))
.await
.expect("write");
let inbound = within(puller.recv()).await.expect("recv");
server.runtime.shutdown().await;
let err = loop {
match within(transfer.write_all(&[0u8; 64 * 1024])).await {
Ok(()) => continue,
Err(e) => break e,
}
};
assert!(
matches!(err, Error::ConnectionLost(_)),
"a write before FIN reports a definite loss: {err:?}"
);
let err = within(inbound.collect(1024))
.await
.expect_err("an unfinished stream is never complete");
assert!(
!matches!(err, Error::Indeterminate),
"the receiver's side of a break is definite too: {err:?}"
);
client.shutdown().await;
}
#[tokio::test]
async fn a_reader_never_sees_an_interrupted_stream_as_complete() {
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 mut transfer = within(pusher.open(Meta::default())).await.expect("open");
within(transfer.write_all(b"half a message"))
.await
.expect("write");
let mut inbound = within(puller.recv()).await.expect("recv");
transfer.cancel();
let err = within(inbound.read_capped(1024))
.await
.expect_err("collect never reports a short payload as whole");
assert!(matches!(err, Error::Canceled), "{err:?}");
let mut transfer = within(pusher.open(Meta::default())).await.expect("open");
within(transfer.write_all(b"half again"))
.await
.expect("write");
let mut inbound = within(puller.recv()).await.expect("recv");
let mut head = [0u8; 4];
within(inbound.read_exact(&mut head))
.await
.expect("the prefix arrived");
assert_eq!(&head, b"half");
transfer.cancel();
let mut rest = [0u8; 64];
loop {
match within(inbound.read(&mut rest)).await {
Ok(n) if n > 0 => continue,
Ok(0) => panic!("an interrupted stream must never read as EOF"),
Ok(_) => unreachable!("n is either zero or positive"),
Err(e) => {
assert_eq!(
e.kind(),
std::io::ErrorKind::ConnectionReset,
"AsyncRead reports a reset, not an end: {e:?}"
);
break;
}
}
}
client.shutdown().await;
}
#[tokio::test]
async fn a_reqrep_requester_can_reissue_after_an_indeterminate_outcome() {
let server = Server::start().await;
let replier = server.listener.replier("/t").expect("replier");
let (taken_tx, taken_rx) = tokio::sync::oneshot::channel();
let taking = tokio::spawn(async move {
let mut request = replier.accept().await.expect("accept");
let body = request.body().read_capped(64).await.expect("body");
assert_eq!(body, b"do the work");
taken_tx.send(()).expect("the requester is waiting");
std::mem::forget(request);
std::future::pending::<()>().await;
});
let client = server.client_runtime();
let requester = client.requester(server.trust());
within(requester.connect(&server.url("/t")))
.await
.expect("connect");
let pending = tokio::spawn(async move { requester.request(b"do the work").await });
taken_rx.await.expect("the replier took the request");
server.runtime.shutdown().await;
let outcome = within(pending)
.await
.expect("the request task")
.expect_err("a request whose FIN landed and whose reply never came is not a success");
taking.abort();
assert!(
matches!(outcome, Error::Indeterminate),
"the replier may have acted, so the outcome is indeterminate: {outcome:?}"
);
let second = Server::start().await;
let replier = second.listener.replier("/t").expect("replier");
let answering = tokio::spawn(async move {
let mut request = replier.accept().await.expect("accept");
let body = request.body().read_capped(64).await.expect("body");
assert_eq!(body, b"do the work");
let mut reply = request.reply(Meta::default()).await.expect("reply");
reply.write_all(b"done").await.expect("write");
reply.finish().expect("finish");
});
let reissued = client.requester(second.trust());
within(reissued.connect(&second.url("/t")))
.await
.expect("connect");
let answer = within(reissued.request(b"do the work"))
.await
.expect("the re-issued request is answered");
assert_eq!(within(answer.collect(64)).await.expect("collect"), b"done");
answering.await.expect("replier");
client.shutdown().await;
}
#[tokio::test]
async fn a_push_producer_is_the_only_side_that_can_reschedule() {
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 mut transfer = within(pusher.open(Meta::default())).await.expect("open");
within(transfer.write_all(b"work")).await.expect("write");
let inbound = within(puller.recv()).await.expect("recv");
transfer.cancel();
let err = within(inbound.collect(64))
.await
.expect_err("the puller sees a failure");
assert!(matches!(err, Error::Canceled), "{err:?}");
within(pusher.send(b"work")).await.expect("resend");
let inbound = within(puller.recv()).await.expect("recv");
assert_eq!(within(inbound.collect(64)).await.expect("collect"), b"work");
client.shutdown().await;
}
#[tokio::test]
async fn a_pubsub_copy_lost_to_a_dead_subscriber_is_counted_not_retried() {
let server = Server::start_with(Limits {
subscriber_buffer_bytes: 4096,
..Limits::default()
})
.await;
let publisher = server.listener.publisher("/md").expect("publisher");
let client = server.client_runtime();
let subscriber = client.subscriber(server.trust());
within(subscriber.connect(&server.url("/md")))
.await
.expect("connect");
within(subscriber.subscribe("")).await.expect("subscribe");
let deadline = std::time::Instant::now() + Duration::from_secs(5);
while publisher.subscriber_count() == 0 && std::time::Instant::now() < deadline {
tokio::time::sleep(Duration::from_millis(10)).await;
}
assert_eq!(publisher.subscriber_count(), 1);
let payload = vec![0x33u8; 4000];
let mut reached = 0usize;
for _ in 0..32u32 {
reached += publisher
.publish("burst", payload.clone())
.expect("publish never fails for a slow subscriber");
}
let drops: u64 = publisher.drops().iter().map(|d| d.total()).sum();
assert!(
drops > 0,
"a copy the subscriber's budget cannot take is counted: {:?}",
publisher.drops()
);
assert_eq!(
reached as u64 + drops,
32,
"every copy is either sent or counted"
);
publisher.publish("end", &b"end"[..]).expect("publish");
let mut burst_seen = 0usize;
loop {
let arrived = within(subscriber.recv()).await.expect("recv");
let topic = arrived.meta().topic.clone().unwrap_or_default();
let body = within(arrived.collect(8192)).await.expect("collect");
if topic == "end" {
assert_eq!(body, b"end");
break;
}
assert_eq!(topic, "burst");
burst_seen += 1;
}
assert_eq!(
burst_seen, reached,
"a dropped copy is never re-sent: {burst_seen} arrived of {reached} sent"
);
client.shutdown().await;
}
#[tokio::test]
async fn a_cursor_reported_before_the_break_survives_the_break() {
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 mut transfer = within(pusher.open(TransferMeta::default().with_report([STORED])))
.await
.expect("open");
let mut cursors = transfer.cursors().expect("cursors");
within(transfer.write_all(&[0x5au8; 900]))
.await
.expect("write");
let mut inbound = within(puller.recv()).await.expect("recv");
let mut reporter = inbound.reporter().expect("a reporter");
let mut prefix = [0u8; 900];
within(inbound.read_exact(&mut prefix))
.await
.expect("the prefix arrived");
within(reporter.report(STORED, 900)).await.expect("report");
let set = within(cursors.changed()).await.expect("a cursor arrived");
assert_eq!(set.offset(STORED), Some(900));
server.runtime.shutdown().await;
let _ = transfer.write_all(&[0u8; 64 * 1024]).await;
assert_eq!(cursors.snapshot().offset(STORED), Some(900));
assert_eq!(within(cursors.changed()).await, None);
client.shutdown().await;
}
#[tokio::test]
async fn a_bound_side_observes_a_dead_dialler_within_the_idle_timeout() {
let limits = Limits {
keep_alive: Duration::from_millis(50),
idle_timeout: Duration::from_millis(300),
..Limits::default()
};
let server = Server::start_with(limits).await;
let paired = server.listener.pair("/link").expect("bound pair");
let url = server.url("/link");
let client = weida::Runtime::owned(RuntimeConfig {
limits,
worker_threads: 1,
..RuntimeConfig::default()
})
.expect("owned runtime");
let dialling = client.pair(server.trust());
within(dialling.connect(&url)).await.expect("connect");
within(dialling.send(b"alive")).await.expect("send");
let inbound = within(paired.recv()).await.expect("recv");
assert_eq!(
within(inbound.collect(64)).await.expect("collect"),
b"alive"
);
drop(dialling);
drop(client);
let err = tokio::time::timeout(limits.idle_timeout * 10, async {
loop {
match paired.send(b"still there?").await {
Ok(()) => tokio::time::sleep(Duration::from_millis(20)).await,
Err(e) => return e,
}
}
})
.await
.expect("the bound side observes the death inside ten idle timeouts");
assert!(
matches!(err, Error::ConnectionLost(_) | Error::NotConnected),
"a dead dialler surfaces as a lost connection: {err:?}"
);
}