mod common;
use std::time::Duration;
use common::Server;
use weida::{Error, GuaranteeSet, Limits, OrderingMode, Publisher, RuntimeConfig, Subscriber};
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")
}
async fn await_filters(publisher: &Publisher, count: usize) {
within(async {
while publisher.filter_count() != count {
tokio::time::sleep(Duration::from_millis(1)).await;
}
})
.await;
}
async fn recv_one(sub: &Subscriber) -> (String, Vec<u8>) {
let transfer = within(sub.recv()).await.expect("recv");
let topic = transfer
.meta()
.topic
.clone()
.expect("a published message carries its topic");
let body = within(transfer.collect(1024 * 1024))
.await
.expect("collect");
(topic, body)
}
#[tokio::test]
async fn subscribe_filters_topics_by_segment() {
let server = Server::start().await;
let publisher = server.listener.publisher("/md").expect("publisher");
let client = server.client_runtime();
let sub = client.subscriber(server.trust());
within(sub.connect(&server.url("/md")))
.await
.expect("connect");
within(sub.subscribe("px.*")).await.expect("subscribe px.*");
within(sub.subscribe("ctl.#"))
.await
.expect("subscribe ctl.#");
within(sub.subscribe("sensors.temp"))
.await
.expect("subscribe sensors.temp");
await_filters(&publisher, 3).await;
assert_eq!(
publisher.publish("px.eur", &b"one"[..]).expect("publish"),
1
);
assert_eq!(
publisher
.publish("px.eur.spot", &b"too deep"[..])
.expect("publish"),
0
);
assert_eq!(
publisher.publish("px", &b"too short"[..]).expect("publish"),
0
);
assert_eq!(
publisher
.publish("fx.usd", &b"nobody"[..])
.expect("publish"),
0
);
assert_eq!(
publisher
.publish("sensors.temperature", &b"not mine"[..])
.expect("publish"),
0
);
assert_eq!(
publisher
.publish("sensors.temp", &b"mine"[..])
.expect("publish"),
1
);
assert_eq!(
publisher.publish("ctl", &b"parent"[..]).expect("publish"),
1
);
assert_eq!(
publisher
.publish("ctl.end.now", &b"deep"[..])
.expect("publish"),
1
);
assert_eq!(recv_one(&sub).await, ("px.eur".to_owned(), b"one".to_vec()));
assert_eq!(
recv_one(&sub).await,
("sensors.temp".to_owned(), b"mine".to_vec())
);
assert_eq!(recv_one(&sub).await, ("ctl".to_owned(), b"parent".to_vec()));
assert_eq!(
recv_one(&sub).await,
("ctl.end.now".to_owned(), b"deep".to_vec())
);
client.shutdown().await;
}
#[tokio::test]
async fn a_topic_containing_a_wildcard_byte_is_literal() {
let server = Server::start().await;
let publisher = server.listener.publisher("/md").expect("publisher");
let client = server.client_runtime();
let sub = client.subscriber(server.trust());
within(sub.connect(&server.url("/md")))
.await
.expect("connect");
within(sub.subscribe("px.*")).await.expect("subscribe");
await_filters(&publisher, 1).await;
assert_eq!(publisher.publish("px.*", &b"star"[..]).expect("publish"), 1);
assert_eq!(recv_one(&sub).await, ("px.*".to_owned(), b"star".to_vec()));
client.shutdown().await;
}
#[tokio::test]
async fn an_illegal_filter_is_refused_before_it_reaches_the_wire() {
let server = Server::start().await;
let _publisher = server.listener.publisher("/md").expect("publisher");
let client = server.client_runtime();
let sub = client.subscriber(server.trust());
within(sub.connect(&server.url("/md")))
.await
.expect("connect");
for bad in ["px*", "px.#.eur", "p*x"] {
let err = within(sub.subscribe(bad))
.await
.expect_err("the grammar must refuse it");
assert!(matches!(err, Error::Protocol(_)), "{bad}: {err:?}");
}
assert_eq!(sub.filter_count(), 0);
client.shutdown().await;
}
#[tokio::test]
async fn empty_filter_receives_all() {
let server = Server::start().await;
let publisher = server.listener.publisher("/md").expect("publisher");
let client = server.client_runtime();
let sub = client.subscriber(server.trust());
within(sub.connect(&server.url("/md")))
.await
.expect("connect");
within(sub.subscribe("")).await.expect("subscribe all");
await_filters(&publisher, 1).await;
for topic in ["px.eur", "fx.usd", "anything"] {
assert_eq!(publisher.publish(topic, &b"x"[..]).expect("publish"), 1);
}
for topic in ["px.eur", "fx.usd", "anything"] {
assert_eq!(recv_one(&sub).await, (topic.to_owned(), b"x".to_vec()));
}
client.shutdown().await;
}
#[tokio::test]
async fn two_subscribers_both_receive() {
let server = Server::start().await;
let publisher = server.listener.publisher("/md").expect("publisher");
let client_a = server.client_runtime();
let client_b = server.client_runtime();
let a = client_a.subscriber(server.trust());
let b = client_b.subscriber(server.trust());
within(a.connect(&server.url("/md")))
.await
.expect("connect a");
within(b.connect(&server.url("/md")))
.await
.expect("connect b");
within(a.subscribe("px.#")).await.expect("subscribe a");
within(b.subscribe("px.#")).await.expect("subscribe b");
await_filters(&publisher, 2).await;
assert_eq!(publisher.subscriber_count(), 2);
assert_eq!(
publisher.publish("px.eur", &b"tick"[..]).expect("publish"),
2
);
assert_eq!(recv_one(&a).await, ("px.eur".to_owned(), b"tick".to_vec()));
assert_eq!(recv_one(&b).await, ("px.eur".to_owned(), b"tick".to_vec()));
client_a.shutdown().await;
client_b.shutdown().await;
}
#[tokio::test]
async fn unsubscribe_stops_delivery() {
let server = Server::start().await;
let publisher = server.listener.publisher("/md").expect("publisher");
let client = server.client_runtime();
let sub = client.subscriber(server.trust());
within(sub.connect(&server.url("/md")))
.await
.expect("connect");
within(sub.subscribe("px.#")).await.expect("subscribe px.#");
within(sub.subscribe("ctl.#"))
.await
.expect("subscribe ctl.#");
await_filters(&publisher, 2).await;
within(sub.unsubscribe("px.#")).await.expect("unsubscribe");
await_filters(&publisher, 1).await;
assert_eq!(publisher.publish("px.x", &b"gone"[..]).expect("publish"), 0);
assert_eq!(
publisher.publish("ctl.end", &b"kept"[..]).expect("publish"),
1
);
assert_eq!(
recv_one(&sub).await,
("ctl.end".to_owned(), b"kept".to_vec())
);
client.shutdown().await;
}
#[tokio::test]
async fn late_publisher_receives_early_subscription() {
let server = Server::start().await;
let client = server.client_runtime();
let sub = client.subscriber(server.trust());
within(sub.connect(&server.url("/md")))
.await
.expect("connect");
within(sub.subscribe("px.#")).await.expect("subscribe");
let publisher = server.listener.publisher("/md").expect("publisher");
await_filters(&publisher, 1).await;
assert_eq!(
publisher.publish("px.eur", &b"late"[..]).expect("publish"),
1
);
assert_eq!(
recv_one(&sub).await,
("px.eur".to_owned(), b"late".to_vec())
);
client.shutdown().await;
}
#[tokio::test]
async fn slow_subscriber_drops_not_blocks() {
const MSG: usize = 32 * 1024;
const COUNT: usize = 100;
let server = Server::start_with(Limits {
subscriber_buffer_bytes: 64 * 1024,
..Limits::default()
})
.await;
let publisher = server.listener.publisher("/md").expect("publisher");
let slow_rt = server.client_runtime_with(Limits {
connection_receive_window: 128 * 1024,
stream_receive_window: 64 * 1024,
..Limits::default()
});
let fast_rt = server.client_runtime();
let slow = slow_rt.subscriber(server.trust());
let fast = fast_rt.subscriber(server.trust());
within(slow.connect(&server.url("/md")))
.await
.expect("connect slow");
within(fast.connect(&server.url("/md")))
.await
.expect("connect fast");
within(slow.subscribe("")).await.expect("subscribe slow");
within(fast.subscribe("")).await.expect("subscribe fast");
await_filters(&publisher, 2).await;
let payload = vec![0xa5u8; MSG];
let mut fast_received = 0usize;
within(async {
publisher.publish("fx.usd", &b"tiny"[..]).expect("publish");
let (topic, _) = recv_one(&fast).await;
assert_eq!(topic, "fx.usd");
for _ in 0..COUNT {
publisher
.publish("px.eur", payload.clone())
.expect("publish");
let (topic, body) = recv_one(&fast).await;
assert_eq!(topic, "px.eur");
assert_eq!(body.len(), MSG);
fast_received += 1;
}
})
.await;
assert_eq!(fast_received, COUNT, "the fast subscriber lost messages");
assert!(
publisher.dropped() > 0,
"the slow subscriber should have lost messages"
);
let starved = publisher
.dropped_on("px.eur")
.expect("the dropped topic is counted");
assert_eq!(starved.total(), publisher.dropped());
assert!(starved.subscriber_budget > 0, "{starved:?}");
assert_eq!(starved.subscriber_queue, 0, "{starved:?}");
assert_eq!(starved.no_parked_connection, 0, "{starved:?}");
assert_eq!(
publisher.dropped_on("fx.usd"),
None,
"the other topic's count is untouched"
);
assert_eq!(publisher.drops().len(), 1);
drop(fast);
drop(slow);
slow_rt.shutdown().await;
fast_rt.shutdown().await;
}
#[tokio::test]
async fn sub_meta_carries_topic_and_the_trace_the_publisher_propagated() {
let server = Server::start().await;
let publisher = server.listener.publisher("/md").expect("publisher");
let client = server.client_runtime();
let sub = client.subscriber(server.trust());
within(sub.connect(&server.url("/md")))
.await
.expect("connect");
within(sub.subscribe("")).await.expect("subscribe");
await_filters(&publisher, 1).await;
publisher.publish("px.eur", &b"body"[..]).expect("publish");
let transfer = within(sub.recv()).await.expect("recv");
let meta = transfer.meta().clone();
assert_eq!(meta.topic.as_deref(), Some("px.eur"));
assert_eq!(meta.endpoint.as_deref(), Some("/md"));
assert_eq!(meta.content_len, Some(4));
assert!(
meta.trace.is_none(),
"a fan-out mints no trace context: it cost 60 bytes on every copy"
);
let trace = weida::new_trace();
publisher
.publish_with_trace("px.eur", &b"body"[..], trace)
.expect("publish with a trace");
let transfer = within(sub.recv()).await.expect("recv the traced copy");
assert_eq!(
transfer.meta().trace,
Some(trace),
"a propagated context reaches every subscriber verbatim"
);
client.shutdown().await;
}
#[tokio::test]
async fn a_payload_larger_than_the_budget_is_refused() {
let server = Server::start().await;
let publisher = server.listener.publisher("/md").expect("publisher");
let oversized = vec![0u8; Limits::default().subscriber_buffer_bytes + 1];
let err = publisher.publish("px.eur", oversized).unwrap_err();
assert!(matches!(err, Error::LimitExceeded), "{err:?}");
}
#[tokio::test]
async fn a_streamed_publish_carries_a_payload_no_publish_could_take() {
const BUDGET: usize = 64 * 1024;
const CHUNK: usize = 16 * 1024;
const CHUNKS: usize = 64;
let server = Server::start_with(Limits {
subscriber_buffer_bytes: BUDGET,
..Limits::default()
})
.await;
let publisher = server.listener.publisher("/md").expect("publisher");
let whole = vec![0xa5u8; CHUNK * CHUNKS];
assert!(
matches!(
publisher.publish("px.eur", whole.clone()).unwrap_err(),
Error::LimitExceeded
),
"the whole-payload publish must still refuse what cannot be enqueued"
);
let first_rt = server.client_runtime();
let second_rt = server.client_runtime();
let first = first_rt.subscriber(server.trust());
let second = second_rt.subscriber(server.trust());
within(first.connect(&server.url("/md")))
.await
.expect("connect first");
within(second.connect(&server.url("/md")))
.await
.expect("connect second");
within(first.subscribe("px.#")).await.expect("subscribe");
within(second.subscribe("px.#")).await.expect("subscribe");
await_filters(&publisher, 2).await;
let readers = tokio::spawn(async move {
let one = recv_one(&first).await;
let two = recv_one(&second).await;
(one, two)
});
let mut fan = publisher.open("px.eur");
assert_eq!(fan.subscribers(), 2);
assert_eq!(fan.topic(), "px.eur");
let chunk = vec![0xa5u8; CHUNK];
for _ in 0..CHUNKS {
let still = within(fan.write_within(chunk.clone(), DEADLINE))
.await
.expect("write");
assert_eq!(still, 2, "a reading subscriber must not lose the transfer");
}
assert_eq!(fan.finish(), 2);
let ((topic_one, body_one), (topic_two, body_two)) =
within(readers).await.expect("both subscribers");
assert_eq!(topic_one, "px.eur");
assert_eq!(topic_two, "px.eur");
assert_eq!(body_one.len(), CHUNK * CHUNKS);
assert_eq!(body_two, body_one);
assert_eq!(body_one, whole);
assert_eq!(publisher.dropped(), 0, "nobody was behind");
}
#[tokio::test]
async fn a_streamed_publish_drops_the_subscriber_that_stalls_and_keeps_the_other() {
const BUDGET: usize = 64 * 1024;
const CHUNK: usize = 16 * 1024;
const CHUNKS: usize = 64;
let server = Server::start_with(Limits {
subscriber_buffer_bytes: BUDGET,
..Limits::default()
})
.await;
let publisher = server.listener.publisher("/md").expect("publisher");
let slow_rt = server.client_runtime_with(Limits {
connection_receive_window: 128 * 1024,
stream_receive_window: 64 * 1024,
..Limits::default()
});
let fast_rt = server.client_runtime();
let slow = slow_rt.subscriber(server.trust());
let fast = fast_rt.subscriber(server.trust());
within(slow.connect(&server.url("/md")))
.await
.expect("connect slow");
within(fast.connect(&server.url("/md")))
.await
.expect("connect fast");
within(slow.subscribe("")).await.expect("subscribe slow");
within(fast.subscribe("")).await.expect("subscribe fast");
await_filters(&publisher, 2).await;
let reader = tokio::spawn(async move { recv_one(&fast).await });
let squeeze = Duration::from_millis(100);
let (remaining, delivered) = within(async {
let mut fan = publisher.open("px.eur");
assert_eq!(fan.subscribers(), 2);
let chunk = vec![0x5au8; CHUNK];
let mut remaining = 2;
for _ in 0..CHUNKS {
remaining = fan
.write_within(chunk.clone(), squeeze)
.await
.expect("write");
}
let delivered = fan.finish();
(remaining, delivered)
})
.await;
assert_eq!(
remaining, 1,
"the subscriber that never read must lose this transfer and the other must keep it"
);
assert_eq!(delivered, 1);
let (topic, body) = within(reader).await.expect("the reading subscriber");
assert_eq!(topic, "px.eur");
assert_eq!(body.len(), CHUNK * CHUNKS);
assert!(
publisher.dropped() >= 1,
"the abandoned copy is counted like any other fan-out drop"
);
let drops = publisher
.dropped_on("px.eur")
.expect("the topic that lost a copy");
assert_eq!(drops.total(), 1, "one copy, counted once: {drops:?}");
}
#[tokio::test]
async fn a_streamed_publish_that_never_waits_drops_at_the_budget() {
const BUDGET: usize = 64 * 1024;
const CHUNK: usize = 16 * 1024;
let server = Server::start_with(Limits {
subscriber_buffer_bytes: BUDGET,
..Limits::default()
})
.await;
let publisher = server.listener.publisher("/md").expect("publisher");
let client = server.client_runtime_with(Limits {
connection_receive_window: 128 * 1024,
stream_receive_window: 64 * 1024,
..Limits::default()
});
let sub = client.subscriber(server.trust());
within(sub.connect(&server.url("/md")))
.await
.expect("connect");
within(sub.subscribe("")).await.expect("subscribe");
await_filters(&publisher, 1).await;
let mut fan = publisher.open("px.eur");
assert_eq!(fan.subscribers(), 1);
let chunk = vec![0x11u8; CHUNK];
let mut written = 0usize;
within(async {
while fan.write_now(chunk.clone()).expect("write") == 1 {
written += 1;
}
})
.await;
assert!(written >= 1, "the first chunks fit the budget");
assert_eq!(fan.subscribers(), 0);
assert_eq!(fan.finish(), 0);
let drops = publisher
.dropped_on("px.eur")
.expect("the topic that lost the copy");
assert_eq!(drops.total(), 1, "{drops:?}");
assert_eq!(drops.subscriber_budget, 1);
}
#[tokio::test]
async fn publishing_to_nobody_is_not_an_error() {
let server = Server::start().await;
let publisher = server.listener.publisher("/md").expect("publisher");
assert_eq!(publisher.subscriber_count(), 0);
assert_eq!(publisher.publish("px.eur", &b"x"[..]).expect("publish"), 0);
}
#[tokio::test]
async fn a_publisher_path_refuses_inbound_transfers() {
let server = Server::start().await;
let _publisher = server.listener.publisher("/md").expect("publisher");
let client = server.client_runtime();
let pusher = client.pusher(server.trust());
within(pusher.connect(&server.url("/md")))
.await
.expect("connect");
let mut transfer = within(pusher.open(weida::TransferMeta::default()))
.await
.expect("open");
let payload = vec![0u8; 2 * 1024 * 1024];
let refused = async {
transfer.write_all(&payload).await?;
transfer.finish()?.delivered().await
};
let err = within(refused)
.await
.expect_err("a push to a publisher path must be refused");
assert!(matches!(err, Error::Unsupported), "{err:?}");
let requester = client.requester(server.trust());
within(requester.connect(&server.url("/md")))
.await
.expect("connect");
let err = within(requester.request(b"nope"))
.await
.expect_err("an exchange with a publisher path must be refused");
assert!(matches!(err, Error::Unsupported), "{err:?}");
client.shutdown().await;
}
#[tokio::test]
async fn a_second_subscriber_on_one_connection_collides() {
let server = Server::start().await;
let _publisher = server.listener.publisher("/md").expect("publisher");
let client = server.client_runtime();
let first = client.subscriber(server.trust());
within(first.connect(&server.url("/md")))
.await
.expect("connect first");
let second = client.subscriber(server.trust());
let err = within(second.connect(&server.url("/md")))
.await
.expect_err("the path is already claimed on this connection");
assert!(matches!(err, Error::AlreadyRegistered), "{err:?}");
client.shutdown().await;
}
#[tokio::test]
async fn a_dropped_fan_out_copy_shows_up_as_a_gap() {
const MSG: usize = 32 * 1024;
let detect = GuaranteeSet {
ordering: OrderingMode::PerProducerDetect,
..GuaranteeSet::CORE
};
let server = Server::start_with_config(RuntimeConfig {
limits: Limits {
subscriber_buffer_bytes: 64 * 1024,
..Limits::default()
},
guarantees: detect,
..RuntimeConfig::default()
})
.await;
let publisher = server.listener.publisher("/md").expect("publisher");
let client = server.client_runtime_with_config(RuntimeConfig {
limits: Limits {
connection_receive_window: 128 * 1024,
stream_receive_window: 64 * 1024,
..Limits::default()
},
guarantees: detect,
..RuntimeConfig::default()
});
let sub = client.subscriber(server.trust());
within(sub.connect(&server.url("/md")))
.await
.expect("connect");
within(sub.subscribe("px.#")).await.expect("subscribe");
await_filters(&publisher, 1).await;
let payload = vec![0u8; MSG];
let mut published = 0usize;
within(async {
while publisher.dropped() == 0 {
publisher
.publish("px.eur", payload.clone())
.expect("publish");
published += 1;
tokio::task::yield_now().await;
}
})
.await;
assert!(published > 1, "the first copy must have been deliverable");
let lost = publisher.dropped();
let mut gaps = Vec::new();
let mut received = 0usize;
while received < published - lost as usize {
let transfer = within(sub.recv()).await.expect("recv");
let meta = transfer.meta().clone();
assert!(
meta.sequence.is_some(),
"detect mode numbers every fan-out copy"
);
if let Some(gap) = meta.gap {
gaps.push(gap);
}
within(transfer.collect(MSG)).await.expect("collect");
received += 1;
}
assert!(gaps.is_empty(), "nothing was missing yet: {gaps:?}");
publisher
.publish("px.eur", payload.clone())
.expect("publish the sentinel");
let sentinel = within(sub.recv()).await.expect("recv the sentinel");
let gap = sentinel
.meta()
.gap
.expect("the sentinel must carry the gap");
assert_eq!(
gap.missed(),
lost,
"the gap must name exactly the copies the publisher dropped: {gap:?}"
);
assert!(gap.seen > gap.expected, "{gap:?}");
assert_eq!(
publisher.dropped(),
lost,
"the sentinel must not be dropped"
);
client.shutdown().await;
}
#[tokio::test]
async fn a_full_hold_reports_the_pub_sub_drop_it_was_waiting_for() {
const MSG: usize = 32 * 1024;
const HOLD: usize = 2;
let reassemble = GuaranteeSet {
ordering: OrderingMode::PerProducerReassemble,
..GuaranteeSet::CORE
};
let server = Server::start_with_config(RuntimeConfig {
limits: Limits {
subscriber_buffer_bytes: 64 * 1024,
..Limits::default()
},
guarantees: reassemble,
..RuntimeConfig::default()
})
.await;
let publisher = server.listener.publisher("/md").expect("publisher");
let client = server.client_runtime_with_config(RuntimeConfig {
limits: Limits {
connection_receive_window: 512 * 1024,
stream_receive_window: 64 * 1024,
max_reorder_hold: HOLD,
..Limits::default()
},
guarantees: reassemble,
..RuntimeConfig::default()
});
let sub = client.subscriber(server.trust());
within(sub.connect(&server.url("/md")))
.await
.expect("connect");
within(sub.subscribe("px.#")).await.expect("subscribe");
await_filters(&publisher, 1).await;
let payload = vec![0u8; MSG];
let mut published = 0usize;
within(async {
while publisher.dropped() == 0 {
publisher
.publish("px.eur", payload.clone())
.expect("publish");
published += 1;
tokio::task::yield_now().await;
}
})
.await;
let lost = publisher.dropped();
for _ in 0..published - lost as usize {
let transfer = within(sub.recv()).await.expect("recv");
assert_eq!(transfer.meta().gap, None, "nothing is missing yet");
within(transfer.collect(MSG)).await.expect("collect");
}
let mut sequences = Vec::new();
let gap = within(async {
loop {
publisher
.publish("px.eur", payload.clone())
.expect("publish");
tokio::task::yield_now().await;
let Ok(transfer) = tokio::time::timeout(Duration::from_millis(20), sub.recv()).await
else {
continue;
};
let transfer = transfer.expect("recv");
let meta = transfer.meta().clone();
sequences.push(meta.sequence.expect("fan-out copies are numbered"));
transfer.collect(MSG).await.expect("collect");
if let Some(gap) = meta.gap {
break gap;
}
}
})
.await;
assert!(
gap.missed() >= lost,
"the gap must cover the copies the publisher dropped: {gap:?}, dropped {}",
publisher.dropped()
);
assert!(
gap.missed() <= publisher.dropped(),
"the gap must not invent losses: {gap:?}, dropped {}",
publisher.dropped()
);
assert!(
sequences.windows(2).all(|pair| pair[0] < pair[1]),
"reassemble mode never delivers backwards: {sequences:?}"
);
client.shutdown().await;
}