#![cfg(all(feature = "blocking", feature = "generate"))]
use std::time::Duration;
use weida::blocking::Runtime;
use weida::{
Acknowledgement, CursorLevel, Error, Identity, ReportMode, Reported, RuntimeConfig,
TransferMeta, Trust,
};
const CAP: usize = 1024 * 1024;
fn served(path: &str) -> (Runtime, weida::blocking::Binding, String) {
let runtime = Runtime::new(RuntimeConfig::default()).expect("an owned reactor");
let identity = Identity::generate().expect("identity");
let fingerprint = identity.fingerprint().expect("fingerprint");
let binding = runtime
.bind_quic("127.0.0.1:0".parse().expect("addr"), identity)
.expect("bind");
let url = format!("weida://{fingerprint}@{}{path}", binding.local_addr());
(runtime, binding, url)
}
#[test]
fn a_request_and_its_reply_cross_between_two_threads() {
let (server, binding, url) = served("/echo");
let replier = binding.replier("/echo").expect("replier");
let answering = std::thread::spawn(move || {
let request = replier.accept(CAP).expect("accept");
assert_eq!(request.message().payload, b"ping");
assert!(request.message().meta.peer.is_none());
assert_eq!(request.message().meta.endpoint.as_deref(), Some("/echo"));
request.reply(b"pong").expect("reply");
});
let client = Runtime::new(RuntimeConfig::default()).expect("client runtime");
let requester = client.requester(Trust::by_address());
requester.connect(&url).expect("connect");
assert_eq!(requester.request(b"ping", CAP).expect("request"), b"pong");
answering.join().expect("the replier thread");
client.shutdown().expect("client shutdown");
server.shutdown().expect("server shutdown");
}
#[test]
fn a_push_reaches_a_puller_and_the_receipt_comes_back() {
let (server, binding, url) = served("/ingest");
let puller = binding.puller("/ingest").expect("puller");
let draining = std::thread::spawn(move || {
let first = puller.recv(CAP).expect("recv");
assert_eq!(first.payload, b"one");
let second = puller.recv(CAP).expect("recv");
assert_eq!(second.payload, b"two");
});
let client = Runtime::new(RuntimeConfig::default()).expect("client runtime");
let pusher = client.pusher(Trust::by_address());
pusher.connect(&url).expect("connect");
pusher.send(b"one").expect("send");
pusher.send(b"two").expect("send");
draining.join().expect("the puller thread");
client.shutdown().expect("client shutdown");
server.shutdown().expect("server shutdown");
}
#[test]
fn a_published_message_reaches_a_blocking_subscriber() {
let (server, binding, url) = served("/md");
let publisher = binding.publisher("/md").expect("publisher");
let client = Runtime::new(RuntimeConfig::default()).expect("client runtime");
let subscriber = client.subscriber(Trust::by_address());
subscriber.connect(&url).expect("connect");
subscriber.subscribe("px.#").expect("subscribe");
let deadline = std::time::Instant::now() + Duration::from_secs(15);
loop {
assert!(
std::time::Instant::now() < deadline,
"the subscription never reached the publisher"
);
if publisher
.publish("px.eur", &b"1.0812"[..])
.expect("publish")
> 0
{
break;
}
std::thread::sleep(Duration::from_millis(5));
}
let message = subscriber.recv(CAP).expect("recv");
assert_eq!(message.payload, b"1.0812");
assert_eq!(message.meta.topic.as_deref(), Some("px.eur"));
assert_eq!(publisher.publish("fx.chf", &b"x"[..]).expect("publish"), 0);
client.shutdown().expect("client shutdown");
server.shutdown().expect("server shutdown");
}
#[test]
fn a_pair_carries_both_directions_and_refuses_a_second_peer() {
let (server, binding, url) = served("/link");
let bound = binding.pair("/link").expect("bound pair");
let answering = std::thread::spawn(move || {
let first = bound.recv(CAP).expect("recv");
assert_eq!(first.payload, b"ping");
bound.send(b"pong").expect("send back");
bound
});
let client = Runtime::new(RuntimeConfig::default()).expect("client runtime");
let paired = client.pair(Trust::by_address());
paired.connect(&url).expect("connect");
paired.send(b"ping").expect("send");
assert_eq!(paired.recv(CAP).expect("recv").payload, b"pong");
let bound = answering.join().expect("the bound thread");
let newcomer = Runtime::new(RuntimeConfig::default()).expect("second client runtime");
let second = newcomer.pair(Trust::by_address());
second
.connect(&url)
.expect("the refusal is per stream, not per connection");
let refused = second
.send(&vec![0x7au8; 2 * 1024 * 1024])
.expect_err("a second peer is refused");
assert!(
matches!(refused, Error::LimitExceeded),
"a capacity decision said out loud, got {refused:?}"
);
paired.send(b"still mine").expect("the first peer is kept");
assert_eq!(bound.recv(CAP).expect("recv").payload, b"still mine");
newcomer.shutdown().expect("second client shutdown");
client.shutdown().expect("client shutdown");
server.shutdown().expect("server shutdown");
}
#[test]
fn a_survey_collects_what_answers_and_counts_what_does_not() {
let (server, binding, url) = served("/poll");
let quiet_url = url.replace("/poll", "/quiet");
let answering = binding.respondent("/poll").expect("respondent");
let silent = binding.respondent("/quiet").expect("a second respondent");
let answers = std::thread::spawn(move || {
let question = answering.accept(CAP).expect("accept");
assert_eq!(question.message().payload, b"who is there");
question.reply(b"me").expect("reply");
});
let holding = std::thread::spawn(move || {
let question = silent.accept(CAP).expect("accept");
std::thread::sleep(Duration::from_secs(3));
drop(question);
});
let client = Runtime::new(RuntimeConfig::default()).expect("client runtime");
let surveyor = client.surveyor(Trust::by_address());
surveyor.connect(&url).expect("connect the answering one");
surveyor
.connect(&quiet_url)
.expect("connect the silent one");
let survey = surveyor
.survey(b"who is there", Duration::from_millis(500), CAP)
.expect("the survey ran");
assert_eq!(survey.asked, 2, "both respondents were asked: {survey:?}");
assert_eq!(
survey.replies,
vec![b"me".to_vec()],
"one answered: {survey:?}"
);
assert_eq!(
survey.silent(),
1,
"and the other's silence is a number, not an error: {survey:?}"
);
answers.join().expect("the answering thread");
client.shutdown().expect("client shutdown");
server.shutdown().expect("server shutdown");
holding.join().expect("the silent thread");
}
#[test]
fn a_bus_message_reaches_every_other_member_and_never_the_sender() {
let (first_rt, first_binding, first_url) = served("/bus1");
let (second_rt, second_binding, second_url) = served("/bus2");
let one = first_binding
.bus("/bus1", Trust::by_address())
.expect("first member");
let two = second_binding
.bus("/bus2", Trust::by_address())
.expect("second member");
one.connect(&second_url).expect("one dials two");
two.connect(&first_url).expect("two dials one");
let listening = std::thread::spawn(move || {
let heard = two.recv(CAP).expect("recv");
assert_eq!(heard.payload, b"hello all");
two.send(b"and back").expect("send back");
two
});
assert_eq!(one.send(b"hello all").expect("send"), 1, "one other member");
assert_eq!(one.recv(CAP).expect("recv").payload, b"and back");
let two = listening.join().expect("the second member's thread");
one.send(b"mine alone").expect("send");
assert_eq!(two.recv(CAP).expect("recv").payload, b"mine alone");
drop(two);
first_rt.shutdown().expect("first shutdown");
second_rt.shutdown().expect("second shutdown");
}
#[test]
fn the_facade_refuses_to_block_a_reactor_thread() {
let reactor = tokio::runtime::Builder::new_current_thread()
.build()
.expect("a runtime to be inside of");
let refused = reactor.block_on(async { Runtime::new(RuntimeConfig::default()) });
match refused {
Err(Error::Runtime(message)) => {
assert!(
message.contains("deadlock"),
"the refusal must say why: {message}"
);
}
Ok(_) => panic!("a blocking runtime built inside a reactor would deadlock on first use"),
Err(other) => panic!("{other:?}"),
}
let outside = Runtime::new(RuntimeConfig::default()).expect("outside a reactor");
let requester = outside.requester(Trust::by_address());
let refused = reactor.block_on(async { requester.connect("weida://127.0.0.1:1/x") });
assert!(
matches!(refused, Err(Error::Runtime(_))),
"{refused:?} must be refused rather than deadlock"
);
outside.shutdown().expect("shutdown");
}
#[test]
fn a_push_gets_a_verdict_from_a_cursor_with_no_exchange() {
const ACCEPTED: CursorLevel = CursorLevel::Known(Acknowledgement::Accepted);
const PROCESSED: CursorLevel = CursorLevel::Known(Acknowledgement::Processed);
let (server, binding, url) = served("/work");
let puller = binding.puller("/work").expect("puller");
let receiving = std::thread::spawn(move || {
let (message, reporter) = puller.recv_reporting(CAP).expect("recv");
assert_eq!(message.payload, b"a unit of work");
assert_eq!(message.meta.report, vec![ACCEPTED, PROCESSED]);
assert_eq!(message.meta.report_mode, ReportMode::Progress);
assert!(message.meta.report_id.is_some());
let mut reporter = reporter.expect("the sender ordered a report");
assert_eq!(reporter.levels(), [ACCEPTED, PROCESSED]);
reporter
.report(ACCEPTED, message.payload.len() as u64)
.expect("accepted");
reporter
.report(CursorLevel::Known(Acknowledgement::Stored), 1)
.expect("ignored");
reporter
.report(PROCESSED, message.payload.len() as u64)
.expect("processed");
reporter.finish().expect("finish the report");
});
let client = Runtime::new(RuntimeConfig::default()).expect("client runtime");
let pusher = client.pusher(Trust::by_address());
pusher.connect(&url).expect("connect");
let body = b"a unit of work";
let mut cursors = pusher
.send_reporting(
TransferMeta::default()
.with_content_len(body.len() as u64)
.with_report([PROCESSED, ACCEPTED]),
body,
)
.expect("send")
.expect("the metadata ordered a report");
let mut processed = None;
while processed.is_none() {
match cursors.changed(Duration::from_secs(15)).expect("a cursor") {
Reported::Changed => {
processed = cursors.snapshot().offset(PROCESSED);
}
Reported::Waiting => panic!("no cursor within the deadline"),
Reported::Ended => panic!("the report ended without a verdict"),
}
}
assert_eq!(processed, Some(body.len() as u64));
assert_eq!(
cursors.snapshot().offset(ACCEPTED),
Some(body.len() as u64),
"a cursor is absolute, so the earlier level is still readable"
);
assert_eq!(
cursors
.snapshot()
.offset(CursorLevel::Known(Acknowledgement::Stored)),
None
);
receiving.join().expect("the receiving thread");
client.shutdown().expect("client shutdown");
server.shutdown().expect("server shutdown");
}
#[test]
fn a_transfer_that_orders_no_report_hands_out_no_handles() {
let (server, binding, url) = served("/plain");
let puller = binding.puller("/plain").expect("puller");
let receiving = std::thread::spawn(move || {
let (message, reporter) = puller.recv_reporting(CAP).expect("recv");
assert_eq!(message.payload, b"no report");
assert!(message.meta.report.is_empty());
assert!(message.meta.report_id.is_none());
assert!(
reporter.is_none(),
"a reporter with nothing ordered would be a handle that writes to nobody"
);
});
let client = Runtime::new(RuntimeConfig::default()).expect("client runtime");
let pusher = client.pusher(Trust::by_address());
pusher.connect(&url).expect("connect");
assert!(
pusher
.send_reporting(TransferMeta::default(), b"no report")
.expect("send")
.is_none(),
"nothing was ordered, so there is nothing to read"
);
receiving.join().expect("the receiving thread");
client.shutdown().expect("client shutdown");
server.shutdown().expect("server shutdown");
}