use crate::error::Error;
use crate::topic::backend::{InMemoryBackend, SubscriptionBackend};
use crate::topic::subscription::{DeliveryBackend, DeliveryStorageKind};
use crate::topic::topic::{BroadcastReadResult, FastBroadcastRing};
use crate::topic::types::{
AckFromResult, BroadcastSubscriberId, Envelope, NackFromResult, PublishOutcome, RecvItem,
SubscriberOptions, SubscriptionMode, TopicOptions, TopicPublishOutcomeConfig,
TrackedPublishOutcome, TrackedPublishPermit, TrackedPublishTracker, TrackedTryPublishOutcome,
};
use crate::topic::{Delivery, RecvDelivery, Subscription, TopicBroker, TopicSet};
use otel_arrow_dfe_config::topic::{TopicBroadcastAckMode, TopicBroadcastOnLagPolicy};
use otel_arrow_dfe_config::{SubscriptionGroupName, TopicName};
use std::collections::HashSet;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::task::{Context, Poll};
use std::time::Duration;
use tokio::sync::Semaphore;
#[tokio::test]
async fn balanced_single_group_all_messages_delivered_exactly_once() {
let broker = TopicBroker::new();
let topic = broker
.create_in_memory_topic("test", TopicOptions::default())
.unwrap();
let mut sub1 = topic
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
)
.unwrap();
let mut sub2 = topic
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
)
.unwrap();
let n = 100u64;
for i in 0..n {
topic.publish(Arc::new(i)).await.unwrap();
}
topic.close();
let mut received = HashSet::new();
let mut order1 = Vec::new();
let mut order2 = Vec::new();
loop {
tokio::select! {
r = sub1.recv() => {
match r {
Ok(RecvItem::Message(env)) => {
assert!(received.insert(env.id), "duplicate id={}", env.id);
order1.push(*env.payload);
}
Err(_) => break,
_ => {}
}
}
r = sub2.recv() => {
match r {
Ok(RecvItem::Message(env)) => {
assert!(received.insert(env.id), "duplicate id={}", env.id);
order2.push(*env.payload);
}
Err(_) => break,
_ => {}
}
}
}
}
while let Ok(RecvItem::Message(env)) = sub1.recv().await {
assert!(received.insert(env.id), "duplicate id={}", env.id);
order1.push(*env.payload);
}
while let Ok(RecvItem::Message(env)) = sub2.recv().await {
assert!(received.insert(env.id), "duplicate id={}", env.id);
order2.push(*env.payload);
}
let all_payloads: HashSet<u64> = order1.iter().chain(order2.iter()).copied().collect();
assert_eq!(all_payloads.len(), n as usize);
for i in 0..n {
assert!(all_payloads.contains(&i), "missing payload {i}");
}
}
#[tokio::test]
async fn balanced_single_group_preserves_order_per_subscriber() {
let broker = TopicBroker::new();
let topic = broker
.create_in_memory_topic("test-order", TopicOptions::default())
.unwrap();
let mut sub1 = topic
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
)
.unwrap();
let mut sub2 = topic
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
)
.unwrap();
let n = 200u64;
for i in 0..n {
topic.publish(Arc::new(i)).await.unwrap();
}
topic.close();
let mut ids1 = Vec::new();
let mut ids2 = Vec::new();
while let Ok(RecvItem::Message(env)) = sub1.recv().await {
ids1.push(env.id);
}
while let Ok(RecvItem::Message(env)) = sub2.recv().await {
ids2.push(env.id);
}
for w in ids1.windows(2) {
assert!(w[0] < w[1], "out of order in sub1: {} >= {}", w[0], w[1]);
}
for w in ids2.windows(2) {
assert!(w[0] < w[1], "out of order in sub2: {} >= {}", w[0], w[1]);
}
}
#[tokio::test]
async fn balanced_no_duplicates_within_group() {
let broker = TopicBroker::new();
let topic = broker
.create_in_memory_topic("dup-test", TopicOptions::default())
.unwrap();
let num_subs = 4;
let mut subs: Vec<_> = (0..num_subs)
.map(|_| {
topic
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
)
.unwrap()
})
.collect();
let n = 500u64;
for i in 0..n {
topic.publish(Arc::new(i)).await.unwrap();
}
topic.close();
let mut all_ids = HashSet::new();
for sub in subs.iter_mut() {
while let Ok(RecvItem::Message(env)) = sub.recv().await {
assert!(all_ids.insert(env.id), "duplicate id={}", env.id);
}
}
assert_eq!(all_ids.len(), n as usize);
}
#[tokio::test]
async fn balanced_multiple_groups_independent() {
let broker = TopicBroker::new();
let topic = broker
.create_in_memory_topic("multi-group", TopicOptions::default())
.unwrap();
let mut sub_g1_a = topic
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("group-A"),
},
SubscriberOptions::default(),
)
.unwrap();
let mut sub_g1_b = topic
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("group-A"),
},
SubscriberOptions::default(),
)
.unwrap();
let mut sub_g2 = topic
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("group-B"),
},
SubscriberOptions::default(),
)
.unwrap();
let n = 100u64;
for i in 0..n {
topic.publish(Arc::new(i)).await.unwrap();
}
topic.close();
let mut ga_ids = HashSet::new();
let mut ga_a_ids = Vec::new();
let mut ga_b_ids = Vec::new();
while let Ok(RecvItem::Message(env)) = sub_g1_a.recv().await {
ga_a_ids.push(env.id);
_ = ga_ids.insert(env.id);
}
while let Ok(RecvItem::Message(env)) = sub_g1_b.recv().await {
ga_b_ids.push(env.id);
assert!(ga_ids.insert(env.id), "duplicate in group-A");
}
assert_eq!(ga_ids.len(), n as usize);
let mut gb_ids = HashSet::new();
let mut gb_order = Vec::new();
while let Ok(RecvItem::Message(env)) = sub_g2.recv().await {
_ = gb_ids.insert(env.id);
gb_order.push(env.id);
}
assert_eq!(gb_ids.len(), n as usize);
for w in gb_order.windows(2) {
assert!(w[0] < w[1]);
}
for w in ga_a_ids.windows(2) {
assert!(w[0] < w[1]);
}
for w in ga_b_ids.windows(2) {
assert!(w[0] < w[1]);
}
}
#[tokio::test]
async fn broadcast_all_subscribers_see_all_messages_in_order() {
let broker = TopicBroker::new();
let topic = broker
.create_in_memory_topic(
"broadcast-test",
TopicOptions::Mixed {
balanced_capacity: TopicOptions::DEFAULT_BALANCED_CAPACITY,
broadcast_capacity: 1024,
on_lag: TopicBroadcastOnLagPolicy::DropOldest,
ack_mode: TopicBroadcastAckMode::First,
},
)
.unwrap();
let mut sub1 = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
let mut sub2 = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
let mut sub3 = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
let n = 100u64;
for i in 0..n {
topic.publish(Arc::new(i)).await.unwrap();
}
topic.close();
for sub in [&mut sub1, &mut sub2, &mut sub3] {
let mut received = Vec::new();
while let Ok(RecvItem::Message(env)) = sub.recv().await {
received.push(*env.payload);
}
assert_eq!(received.len(), n as usize);
for (i, &val) in received.iter().enumerate() {
assert_eq!(val, i as u64);
}
}
}
#[tokio::test]
async fn broadcast_lag_reported_on_slow_subscriber() {
let broker = TopicBroker::new();
let topic = broker
.create_in_memory_topic(
"broadcast-lag",
TopicOptions::Mixed {
balanced_capacity: TopicOptions::DEFAULT_BALANCED_CAPACITY,
broadcast_capacity: 8,
on_lag: TopicBroadcastOnLagPolicy::DropOldest,
ack_mode: TopicBroadcastAckMode::First,
},
)
.unwrap();
let mut sub = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
let n = 50u64;
for i in 0..n {
topic.publish(Arc::new(i)).await.unwrap();
}
topic.close();
let mut messages = Vec::new();
let mut total_lagged = 0u64;
loop {
match sub.recv().await {
Ok(RecvItem::Message(env)) => messages.push(*env.payload),
Ok(RecvItem::Lagged { missed }) => total_lagged += missed,
Err(_) => break,
}
}
assert!(
messages.len() < n as usize,
"expected lag but got all {} messages",
n
);
assert!(total_lagged > 0, "expected lag notification");
for w in messages.windows(2) {
assert!(w[0] < w[1], "out of order: {} >= {}", w[0], w[1]);
}
}
#[tokio::test]
async fn broadcast_slow_subscriber_does_not_block_fast_subscriber() {
let broker = TopicBroker::new();
let topic = broker
.create_in_memory_topic(
"broadcast-no-block",
TopicOptions::Mixed {
balanced_capacity: TopicOptions::DEFAULT_BALANCED_CAPACITY,
broadcast_capacity: 4,
on_lag: TopicBroadcastOnLagPolicy::DropOldest,
ack_mode: TopicBroadcastAckMode::First,
},
)
.unwrap();
let _slow_sub = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
let mut fast_sub = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
for i in 0..20u64 {
topic.publish(Arc::new(i)).await.unwrap();
match fast_sub.recv().await {
Ok(RecvItem::Message(env)) => assert_eq!(*env.payload, i),
other => panic!("unexpected: {:?}", other),
}
}
}
#[tokio::test]
async fn broadcast_disconnects_slow_subscriber_on_lag() {
let broker = TopicBroker::new();
let topic = broker
.create_topic(
"broadcast-disconnect",
TopicOptions::BroadcastOnly {
capacity: 8,
on_lag: TopicBroadcastOnLagPolicy::Disconnect,
ack_mode: TopicBroadcastAckMode::First,
},
InMemoryBackend,
)
.unwrap();
let mut sub = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
for i in 0..50u64 {
topic.publish(Arc::new(i)).await.unwrap();
}
match sub.recv().await {
Ok(RecvItem::Lagged { missed }) => assert!(missed > 0, "expected missed count > 0"),
other => panic!("expected lag notification before disconnect, got {other:?}"),
}
assert!(matches!(sub.recv().await, Err(Error::SubscriptionClosed)));
}
#[tokio::test]
async fn mixed_async_publish_waits_for_balanced_admission_before_broadcast() {
let broker = TopicBroker::new();
let topic = broker
.create_topic(
"mixed-bcast-priority",
TopicOptions::Mixed {
balanced_capacity: 1,
broadcast_capacity: 16,
on_lag: TopicBroadcastOnLagPolicy::DropOldest,
ack_mode: TopicBroadcastAckMode::First,
},
InMemoryBackend,
)
.unwrap();
let mut balanced = topic
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
)
.unwrap();
let mut broadcast = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
topic.publish(Arc::new(1u64)).await.unwrap();
let topic_clone = topic.clone();
let second_publish = tokio::spawn(async move {
topic_clone.publish(Arc::new(2u64)).await.unwrap();
});
match broadcast.recv().await {
Ok(RecvItem::Message(env)) => assert_eq!(*env.payload, 1),
other => panic!("unexpected: {:?}", other),
}
assert!(
tokio::time::timeout(Duration::from_millis(200), async {
loop {
match broadcast.recv().await {
Ok(RecvItem::Message(env)) => break *env.payload,
Ok(RecvItem::Lagged { .. }) => continue,
Err(e) => panic!("unexpected recv error: {e:?}"),
}
}
})
.await
.is_err(),
"broadcast delivery should wait for balanced admission on async publish"
);
assert!(
!second_publish.is_finished(),
"second publish should still be blocked by balanced backpressure"
);
_ = balanced.recv().await.unwrap();
second_publish.await.unwrap();
let second = tokio::time::timeout(Duration::from_millis(200), async {
loop {
match broadcast.recv().await {
Ok(RecvItem::Message(env)) => break *env.payload,
Ok(RecvItem::Lagged { .. }) => continue,
Err(e) => panic!("unexpected recv error: {e:?}"),
}
}
})
.await
.expect("broadcast delivery should resume after balanced admission succeeds");
assert_eq!(second, 2);
_ = balanced.recv().await.unwrap();
topic.close();
}
#[tokio::test]
async fn mixed_async_publish_does_not_reserve_free_groups_while_waiting() {
let broker = TopicBroker::<u64>::new();
let topic = broker
.create_topic(
"mixed-no-convoy",
TopicOptions::Mixed {
balanced_capacity: 1,
broadcast_capacity: 16,
on_lag: TopicBroadcastOnLagPolicy::DropOldest,
ack_mode: TopicBroadcastAckMode::First,
},
InMemoryBackend,
)
.unwrap();
let mut fast = topic
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("fast"),
},
SubscriberOptions::default(),
)
.unwrap();
let mut slow = topic
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("slow"),
},
SubscriberOptions::default(),
)
.unwrap();
topic.publish(Arc::new(1u64)).await.unwrap();
match fast.recv().await.unwrap() {
RecvItem::Message(env) => assert_eq!(*env.payload, 1),
other => panic!("unexpected first fast item: {other:?}"),
}
let topic_clone = topic.clone();
let blocked_publish = tokio::spawn(async move {
topic_clone.publish(Arc::new(2u64)).await.unwrap();
});
tokio::time::sleep(Duration::from_millis(50)).await;
assert!(
!blocked_publish.is_finished(),
"publish should still be blocked by the slow group"
);
let mut permits = topic.debug_balanced_available_permits();
permits.sort_by(|(left, _), (right, _)| left.as_ref().cmp(right.as_ref()));
assert_eq!(
permits,
vec![
(SubscriptionGroupName::from("fast"), 1),
(SubscriptionGroupName::from("slow"), 0),
],
"blocked mixed publish should not hold the fast group's permit"
);
match slow.recv().await.unwrap() {
RecvItem::Message(env) => assert_eq!(*env.payload, 1),
other => panic!("unexpected first slow item: {other:?}"),
}
blocked_publish.await.unwrap();
match fast.recv().await.unwrap() {
RecvItem::Message(env) => assert_eq!(*env.payload, 2),
other => panic!("unexpected second fast item: {other:?}"),
}
match slow.recv().await.unwrap() {
RecvItem::Message(env) => assert_eq!(*env.payload, 2),
other => panic!("unexpected second slow item: {other:?}"),
}
topic.close();
}
#[tokio::test]
async fn mixed_async_publish_drop_does_not_publish_to_broadcast() {
let broker = TopicBroker::<u64>::new();
let topic = broker
.create_topic(
"mixed-cancel-no-broadcast",
TopicOptions::Mixed {
balanced_capacity: 1,
broadcast_capacity: 16,
on_lag: TopicBroadcastOnLagPolicy::DropOldest,
ack_mode: TopicBroadcastAckMode::First,
},
InMemoryBackend,
)
.unwrap();
let mut balanced = topic
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
)
.unwrap();
let mut broadcast = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
topic.publish(Arc::new(1u64)).await.unwrap();
let mut publish = Box::pin(topic.publish(Arc::new(2u64)));
assert!(
tokio::time::timeout(Duration::from_millis(200), publish.as_mut())
.await
.is_err(),
"second publish should block while balanced capacity is full"
);
match broadcast.recv().await.unwrap() {
RecvItem::Message(env) => assert_eq!(*env.payload, 1),
other => panic!("unexpected first broadcast item: {other:?}"),
}
drop(publish);
assert!(
tokio::time::timeout(Duration::from_millis(200), async {
loop {
match broadcast.recv().await {
Ok(RecvItem::Message(env)) => break *env.payload,
Ok(RecvItem::Lagged { .. }) => continue,
Err(e) => panic!("unexpected recv error: {e:?}"),
}
}
})
.await
.is_err(),
"dropping a blocked publish must not leak a broadcast delivery"
);
match balanced.recv().await.unwrap() {
RecvItem::Message(env) => assert_eq!(*env.payload, 1),
other => panic!("unexpected balanced item: {other:?}"),
}
}
#[tokio::test]
async fn blocked_tracked_publish_drop_releases_in_flight_capacity() {
let broker = TopicBroker::<u64>::new();
let topic = broker
.create_topic(
"tracked-publish-drop",
TopicOptions::BalancedOnly { capacity: 1 },
InMemoryBackend,
)
.unwrap();
let _balanced = topic
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
)
.unwrap();
topic.publish(Arc::new(1u64)).await.unwrap();
let tracked = topic.tracked_publisher_with_config(TopicPublishOutcomeConfig {
max_in_flight: 1,
timeout: Duration::from_secs(30),
});
let mut blocked = Box::pin(tracked.publish(Arc::new(2u64)));
assert!(
tokio::time::timeout(Duration::from_millis(200), blocked.as_mut())
.await
.is_err(),
"tracked publish should block while balanced capacity is full"
);
drop(blocked);
assert!(
matches!(
tracked.try_publish(Arc::new(3u64)).unwrap(),
TrackedTryPublishOutcome::DroppedOnFull
),
"dropping the blocked tracked publish should release the in-flight slot"
);
}
#[tokio::test]
async fn broadcast_delivery_commit_advances_to_next_message() {
let broker = TopicBroker::<u64>::new();
let topic = broker
.create_topic(
"broadcast-delivery-permit",
TopicOptions::BroadcastOnly {
capacity: 16,
on_lag: TopicBroadcastOnLagPolicy::DropOldest,
ack_mode: TopicBroadcastAckMode::First,
},
InMemoryBackend,
)
.unwrap();
let mut sub = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
topic.publish(Arc::new(1u64)).await.unwrap();
topic.publish(Arc::new(2u64)).await.unwrap();
let first = match sub.recv_delivery().await.unwrap() {
RecvDelivery::Message(delivery) => delivery,
RecvDelivery::Lagged { .. } => panic!("unexpected lag notification"),
};
assert_eq!(first.message_id(), 1);
assert_eq!(*first.envelope().payload, 1);
assert!(
tokio::time::timeout(Duration::from_millis(200), sub.recv_delivery())
.await
.is_err(),
"second delivery should stay blocked while the first delivery permit is held"
);
first.commit();
let second = match sub.recv_delivery().await.unwrap() {
RecvDelivery::Message(delivery) => delivery,
RecvDelivery::Lagged { .. } => panic!("unexpected lag notification"),
};
assert_eq!(second.message_id(), 2);
assert_eq!(*second.envelope().payload, 2);
second.commit();
}
#[tokio::test]
async fn balanced_delivery_uses_specialized_inline_storage() {
let broker = TopicBroker::<u64>::new();
let topic = broker
.create_topic(
"balanced-inline-delivery",
TopicOptions::default(),
InMemoryBackend,
)
.unwrap();
let mut sub = topic
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
)
.unwrap();
topic.publish(Arc::new(7u64)).await.unwrap();
let delivery = match sub.recv_delivery().await.unwrap() {
RecvDelivery::Message(delivery) => delivery,
RecvDelivery::Lagged { .. } => panic!("unexpected lag notification"),
};
assert_eq!(delivery.storage_kind(), DeliveryStorageKind::Balanced);
delivery.commit();
}
#[tokio::test]
async fn broadcast_delivery_uses_specialized_inline_storage() {
let broker = TopicBroker::<u64>::new();
let topic = broker
.create_topic(
"broadcast-inline-delivery",
TopicOptions::BroadcastOnly {
capacity: 16,
on_lag: TopicBroadcastOnLagPolicy::DropOldest,
ack_mode: TopicBroadcastAckMode::First,
},
InMemoryBackend,
)
.unwrap();
let mut sub = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
topic.publish(Arc::new(11u64)).await.unwrap();
let delivery = match sub.recv_delivery().await.unwrap() {
RecvDelivery::Message(delivery) => delivery,
RecvDelivery::Lagged { .. } => panic!("unexpected lag notification"),
};
assert_eq!(delivery.storage_kind(), DeliveryStorageKind::Broadcast);
delivery.commit();
}
#[tokio::test]
async fn opaque_delivery_fallback_still_works() {
struct OpaqueDelivery {
envelope: Envelope<u64>,
aborted: Arc<AtomicBool>,
}
impl DeliveryBackend<u64> for OpaqueDelivery {
fn envelope(&self) -> &Envelope<u64> {
&self.envelope
}
fn commit(&mut self) {}
fn abort(&mut self, _reason: Arc<str>) -> Result<(), Error> {
self.aborted.store(true, Ordering::SeqCst);
Ok(())
}
fn abandon(&mut self) {}
}
struct OpaqueSubscription {
yielded: bool,
aborted: Arc<AtomicBool>,
}
impl SubscriptionBackend<u64> for OpaqueSubscription {
fn poll_recv_delivery(
&mut self,
_cx: &mut Context<'_>,
) -> Poll<Result<RecvDelivery<u64>, Error>> {
if self.yielded {
Poll::Ready(Err(Error::SubscriptionClosed))
} else {
self.yielded = true;
Poll::Ready(Ok(RecvDelivery::Message(Delivery::new_opaque(Box::new(
OpaqueDelivery {
envelope: Envelope {
id: 41,
payload: Arc::new(99u64),
tracked: false,
},
aborted: Arc::clone(&self.aborted),
},
)))))
}
}
fn ack(&self, _id: u64) -> Result<(), Error> {
Ok(())
}
fn nack(&self, _id: u64, _reason: Arc<str>) -> Result<(), Error> {
Ok(())
}
}
let aborted = Arc::new(AtomicBool::new(false));
let mut subscription = Subscription::new(Box::new(OpaqueSubscription {
yielded: false,
aborted: Arc::clone(&aborted),
}));
let delivery = match subscription.recv_delivery().await.unwrap() {
RecvDelivery::Message(delivery) => delivery,
RecvDelivery::Lagged { .. } => panic!("unexpected lag notification"),
};
assert_eq!(delivery.storage_kind(), DeliveryStorageKind::Opaque);
delivery.abort("opaque fallback abort").unwrap();
assert!(aborted.load(Ordering::SeqCst));
}
#[tokio::test]
async fn tracked_delivery_abort_resolves_publish_outcome() {
let broker = TopicBroker::<u64>::new();
let topic = broker
.create_topic(
"tracked-delivery-abort",
TopicOptions::default(),
InMemoryBackend,
)
.unwrap();
let tracked = topic.tracked_publisher();
let mut sub = topic
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
)
.unwrap();
let receipt = tracked.publish(Arc::new(42u64)).await.unwrap();
let delivery = match sub.recv_delivery().await.unwrap() {
RecvDelivery::Message(delivery) => delivery,
RecvDelivery::Lagged { .. } => panic!("unexpected lag notification"),
};
delivery
.abort("downstream rejected before forward")
.unwrap();
assert_eq!(
receipt.wait_for_outcome().await,
TrackedPublishOutcome::Nack {
reason: Arc::from("downstream rejected before forward"),
}
);
}
#[tokio::test]
async fn mixed_disconnects_only_lagging_broadcast_subscriber() {
let broker = TopicBroker::new();
let topic = broker
.create_in_memory_topic(
"mixed-broadcast-disconnect",
TopicOptions::Mixed {
balanced_capacity: 32,
broadcast_capacity: 4,
on_lag: TopicBroadcastOnLagPolicy::Disconnect,
ack_mode: TopicBroadcastAckMode::First,
},
)
.unwrap();
let mut balanced = topic
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("workers"),
},
SubscriberOptions::default(),
)
.unwrap();
let mut slow_broadcast = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
let mut fast_broadcast = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
for i in 0..20u64 {
topic.publish(Arc::new(i)).await.unwrap();
match fast_broadcast.recv().await {
Ok(RecvItem::Message(env)) => assert_eq!(*env.payload, i),
other => panic!("fast broadcast subscriber should stay current: {other:?}"),
}
}
match slow_broadcast.recv().await {
Ok(RecvItem::Lagged { missed }) => assert!(missed > 0, "expected lagged slow subscriber"),
other => panic!("expected lag notification for slow broadcast subscriber, got {other:?}"),
}
assert!(matches!(
slow_broadcast.recv().await,
Err(Error::SubscriptionClosed)
));
match balanced.recv().await {
Ok(RecvItem::Message(env)) => assert_eq!(*env.payload, 0),
other => panic!("balanced subscriber should still receive messages: {other:?}"),
}
}
#[tokio::test]
async fn publisher_api_identical_for_all_subscriber_modes() {
let broker = TopicBroker::new();
let topic = broker
.create_topic("choice-test", TopicOptions::default(), InMemoryBackend)
.unwrap();
let mut balanced = topic
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
)
.unwrap();
let mut broadcast = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
topic.publish(Arc::new("hello".to_string())).await.unwrap();
topic.close();
match balanced.recv().await {
Ok(RecvItem::Message(env)) => assert_eq!(*env.payload, "hello"),
other => panic!("unexpected: {:?}", other),
}
match broadcast.recv().await {
Ok(RecvItem::Message(env)) => assert_eq!(*env.payload, "hello"),
other => panic!("unexpected: {:?}", other),
}
}
#[tokio::test]
async fn tracked_publish_outcomes_received_when_enabled() {
let broker = TopicBroker::new();
let base = broker
.create_topic("ack-test", TopicOptions::default(), InMemoryBackend)
.unwrap();
let topic = base.tracked_publisher();
let mut sub = base
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
)
.unwrap();
let receipt1 = topic.publish(Arc::new(42u64)).await.unwrap();
let receipt2 = topic.publish(Arc::new(43u64)).await.unwrap();
let env1 = match sub.recv().await.unwrap() {
RecvItem::Message(e) => e,
_ => panic!(),
};
let env2 = match sub.recv().await.unwrap() {
RecvItem::Message(e) => e,
_ => panic!(),
};
sub.ack(env1.id).unwrap();
sub.nack(env2.id, "bad data").unwrap();
assert_eq!(
receipt1.wait_for_outcome().await,
TrackedPublishOutcome::Ack
);
assert_eq!(
receipt2.wait_for_outcome().await,
TrackedPublishOutcome::Nack {
reason: Arc::from("bad data"),
}
);
}
#[tokio::test]
async fn ack_nack_fail_when_message_is_not_tracked() {
let broker = TopicBroker::new();
let topic = broker
.create_topic("no-ack", TopicOptions::default(), InMemoryBackend)
.unwrap();
let mut sub = topic
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
)
.unwrap();
topic.publish(Arc::new(1)).await.unwrap();
let env = match sub.recv().await.unwrap() {
RecvItem::Message(e) => e,
_ => panic!(),
};
match sub.ack(env.id) {
Err(Error::MessageNotTracked) => {}
other => panic!("expected MessageNotTracked, got {:?}", other),
}
}
#[tokio::test]
async fn tracked_try_publish_respects_max_in_flight() {
let broker = TopicBroker::new();
let base = broker
.create_topic("ack-full", TopicOptions::default(), InMemoryBackend)
.unwrap();
let topic = base.tracked_publisher_with_config(TopicPublishOutcomeConfig {
max_in_flight: 1,
timeout: Duration::from_secs(30),
});
let mut sub = base
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
)
.unwrap();
let receipt = topic.publish(Arc::new(1u64)).await.unwrap();
assert!(matches!(
topic.try_publish(Arc::new(2u64)).unwrap(),
TrackedTryPublishOutcome::MaxInFlightReached
));
let env = match sub.recv().await.unwrap() {
RecvItem::Message(e) => e,
_ => panic!(),
};
sub.ack(env.id).unwrap();
assert_eq!(receipt.wait_for_outcome().await, TrackedPublishOutcome::Ack);
}
#[tokio::test]
async fn balanced_backpressure_blocks_publisher() {
let broker = TopicBroker::new();
let topic = broker
.create_topic(
"backpressure",
TopicOptions::Mixed {
balanced_capacity: 2,
broadcast_capacity: TopicOptions::DEFAULT_BROADCAST_CAPACITY,
on_lag: TopicBroadcastOnLagPolicy::DropOldest,
ack_mode: TopicBroadcastAckMode::First,
},
InMemoryBackend,
)
.unwrap();
let mut sub = topic
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
)
.unwrap();
topic.publish(Arc::new(1)).await.unwrap();
topic.publish(Arc::new(2)).await.unwrap();
let topic_clone = topic.clone();
let handle = tokio::spawn(async move {
topic_clone.publish(Arc::new(3)).await.unwrap();
});
tokio::time::sleep(Duration::from_millis(50)).await;
assert!(!handle.is_finished(), "publish should be blocked");
let _ = sub.recv().await.unwrap();
tokio::time::sleep(Duration::from_millis(50)).await;
assert!(handle.is_finished(), "publish should have completed");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn balanced_multi_threaded_no_duplicates() {
let broker = TopicBroker::new();
let topic = broker
.create_topic("mt-balanced", TopicOptions::default(), InMemoryBackend)
.unwrap();
let num_publishers = 4;
let msgs_per_publisher = 250;
let total = num_publishers * msgs_per_publisher;
let mut sub1 = topic
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
)
.unwrap();
let mut sub2 = topic
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
)
.unwrap();
let mut pub_handles = Vec::new();
for p in 0..num_publishers {
let t = topic.clone();
pub_handles.push(tokio::spawn(async move {
for i in 0..msgs_per_publisher {
let val = (p * msgs_per_publisher + i) as u64;
t.publish(Arc::new(val)).await.unwrap();
}
}));
}
for h in pub_handles {
h.await.unwrap();
}
topic.close();
let mut all_ids = HashSet::new();
while let Ok(RecvItem::Message(env)) = sub1.recv().await {
assert!(all_ids.insert(env.id));
}
while let Ok(RecvItem::Message(env)) = sub2.recv().await {
assert!(all_ids.insert(env.id));
}
assert_eq!(all_ids.len(), total);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn broadcast_multi_threaded_all_receive() {
let broker = TopicBroker::new();
let topic = broker
.create_topic(
"mt-broadcast",
TopicOptions::Mixed {
balanced_capacity: TopicOptions::DEFAULT_BALANCED_CAPACITY,
broadcast_capacity: 2048,
on_lag: TopicBroadcastOnLagPolicy::DropOldest,
ack_mode: TopicBroadcastAckMode::First,
},
InMemoryBackend,
)
.unwrap();
let n = 500u64;
let mut sub1 = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
let mut sub2 = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
let t = topic.clone();
let pub_handle = tokio::spawn(async move {
for i in 0..n {
t.publish(Arc::new(i)).await.unwrap();
}
t.close();
});
let h1 = tokio::spawn(async move {
let mut received = Vec::new();
loop {
match sub1.recv().await {
Ok(RecvItem::Message(env)) => received.push(*env.payload),
Ok(RecvItem::Lagged { .. }) => {}
Err(_) => break,
}
}
received
});
let h2 = tokio::spawn(async move {
let mut received = Vec::new();
loop {
match sub2.recv().await {
Ok(RecvItem::Message(env)) => received.push(*env.payload),
Ok(RecvItem::Lagged { .. }) => {}
Err(_) => break,
}
}
received
});
pub_handle.await.unwrap();
let r1 = h1.await.unwrap();
let r2 = h2.await.unwrap();
assert_eq!(r1.len(), n as usize);
assert_eq!(r2.len(), n as usize);
for w in r1.windows(2) {
assert!(w[0] < w[1]);
}
for w in r2.windows(2) {
assert!(w[0] < w[1]);
}
}
#[tokio::test]
async fn publish_with_no_subscribers_succeeds() {
let broker = TopicBroker::new();
let topic = broker
.create_topic("empty", TopicOptions::default(), InMemoryBackend)
.unwrap();
topic.publish(Arc::new(42)).await.unwrap();
}
#[tokio::test]
async fn subscriber_only_receives_messages_after_subscribe() {
let broker = TopicBroker::new();
let topic = broker
.create_topic("late-sub", TopicOptions::default(), InMemoryBackend)
.unwrap();
topic.publish(Arc::new(1)).await.unwrap();
topic.publish(Arc::new(2)).await.unwrap();
let mut sub = topic
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
)
.unwrap();
topic.publish(Arc::new(3)).await.unwrap();
topic.close();
let mut received = Vec::new();
while let Ok(RecvItem::Message(env)) = sub.recv().await {
received.push(*env.payload);
}
assert_eq!(received, vec![3]);
}
#[tokio::test]
async fn balanced_only_basic_delivery() {
let broker = TopicBroker::new();
let topic = broker
.create_topic(
"bo-basic",
TopicOptions::BalancedOnly {
capacity: TopicOptions::DEFAULT_BALANCED_CAPACITY,
},
InMemoryBackend,
)
.unwrap();
let mut sub = topic
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
)
.unwrap();
let n = 50u64;
for i in 0..n {
topic.publish(Arc::new(i)).await.unwrap();
}
topic.close();
let mut received = Vec::new();
while let Ok(RecvItem::Message(env)) = sub.recv().await {
received.push(*env.payload);
}
assert_eq!(received.len(), n as usize);
for (i, &val) in received.iter().enumerate() {
assert_eq!(val, i as u64);
}
}
#[tokio::test]
async fn balanced_only_rejects_broadcast() {
let broker: TopicBroker<u64> = TopicBroker::new();
let topic = broker
.create_topic(
"bo-no-bcast",
TopicOptions::BalancedOnly {
capacity: TopicOptions::DEFAULT_BALANCED_CAPACITY,
},
InMemoryBackend,
)
.unwrap();
match topic.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default()) {
Err(Error::SubscribeBroadcastNotSupported) => {}
Ok(_) => panic!("expected BroadcastNotSupported, got Ok"),
Err(e) => panic!("expected BroadcastNotSupported, got Err({e:?})"),
}
}
#[tokio::test]
async fn balanced_only_rejects_second_group() {
let broker: TopicBroker<u64> = TopicBroker::new();
let topic = broker
.create_topic(
"bo-single-group",
TopicOptions::BalancedOnly {
capacity: TopicOptions::DEFAULT_BALANCED_CAPACITY,
},
InMemoryBackend,
)
.unwrap();
let _sub1 = topic
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
)
.unwrap();
let _sub2 = topic
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
)
.unwrap();
match topic.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g2"),
},
SubscriberOptions::default(),
) {
Err(Error::SubscribeSingleGroupViolation) => {}
Ok(_) => panic!("expected SingleGroupViolation, got Ok"),
Err(e) => panic!("expected SingleGroupViolation, got Err({e:?})"),
}
}
#[tokio::test]
async fn balanced_only_rejects_balanced_subscribe_after_close() {
let broker: TopicBroker<u64> = TopicBroker::new();
let topic = broker
.create_topic(
"bo-closed-subscribe",
TopicOptions::BalancedOnly {
capacity: TopicOptions::DEFAULT_BALANCED_CAPACITY,
},
InMemoryBackend,
)
.unwrap();
topic.close();
match topic.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
) {
Err(Error::TopicClosed) => {}
Ok(_) => panic!("expected TopicClosed, got Ok"),
Err(e) => panic!("expected TopicClosed, got Err({e:?})"),
}
}
#[tokio::test]
async fn mixed_rejects_balanced_subscribe_after_close() {
let broker: TopicBroker<u64> = TopicBroker::new();
let topic = broker
.create_topic(
"mixed-closed-subscribe",
TopicOptions::Mixed {
balanced_capacity: TopicOptions::DEFAULT_BALANCED_CAPACITY,
broadcast_capacity: TopicOptions::DEFAULT_BROADCAST_CAPACITY,
on_lag: TopicBroadcastOnLagPolicy::DropOldest,
ack_mode: TopicBroadcastAckMode::First,
},
InMemoryBackend,
)
.unwrap();
topic.close();
match topic.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
) {
Err(Error::TopicClosed) => {}
Ok(_) => panic!("expected TopicClosed, got Ok"),
Err(e) => panic!("expected TopicClosed, got Err({e:?})"),
}
}
#[tokio::test]
async fn balanced_only_no_messages_lost() {
let broker = TopicBroker::new();
let topic = broker
.create_topic(
"bo-no-loss",
TopicOptions::BalancedOnly {
capacity: TopicOptions::DEFAULT_BALANCED_CAPACITY,
},
InMemoryBackend,
)
.unwrap();
let mut sub1 = topic
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
)
.unwrap();
let mut sub2 = topic
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
)
.unwrap();
let n = 200u64;
for i in 0..n {
topic.publish(Arc::new(i)).await.unwrap();
}
topic.close();
let mut all_ids = HashSet::new();
while let Ok(RecvItem::Message(env)) = sub1.recv().await {
assert!(all_ids.insert(env.id), "duplicate id={}", env.id);
}
while let Ok(RecvItem::Message(env)) = sub2.recv().await {
assert!(all_ids.insert(env.id), "duplicate id={}", env.id);
}
assert_eq!(all_ids.len(), n as usize);
}
#[tokio::test]
async fn broadcast_only_basic_delivery() {
let broker = TopicBroker::new();
let topic = broker
.create_topic(
"bro-basic",
TopicOptions::BroadcastOnly {
capacity: 1024,
on_lag: TopicBroadcastOnLagPolicy::DropOldest,
ack_mode: TopicBroadcastAckMode::First,
},
InMemoryBackend,
)
.unwrap();
let mut sub1 = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
let mut sub2 = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
let n = 50u64;
for i in 0..n {
topic.publish(Arc::new(i)).await.unwrap();
}
topic.close();
for sub in [&mut sub1, &mut sub2] {
let mut received = Vec::new();
while let Ok(RecvItem::Message(env)) = sub.recv().await {
received.push(*env.payload);
}
assert_eq!(received.len(), n as usize);
for (i, &val) in received.iter().enumerate() {
assert_eq!(val, i as u64);
}
}
}
#[tokio::test]
async fn broadcast_only_rejects_balanced() {
let broker: TopicBroker<u64> = TopicBroker::new();
let topic = broker
.create_topic(
"bro-no-balanced",
TopicOptions::BroadcastOnly {
capacity: TopicOptions::DEFAULT_BROADCAST_CAPACITY,
on_lag: TopicBroadcastOnLagPolicy::DropOldest,
ack_mode: TopicBroadcastAckMode::First,
},
InMemoryBackend,
)
.unwrap();
match topic.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
) {
Err(Error::SubscribeBalancedNotSupported) => {}
Ok(_) => panic!("expected BalancedNotSupported, got Ok"),
Err(e) => panic!("expected BalancedNotSupported, got Err({e:?})"),
}
}
#[tokio::test]
async fn broadcast_only_lag_reported() {
let broker = TopicBroker::new();
let topic = broker
.create_topic(
"bro-lag",
TopicOptions::BroadcastOnly {
capacity: 8,
on_lag: TopicBroadcastOnLagPolicy::DropOldest,
ack_mode: TopicBroadcastAckMode::First,
},
InMemoryBackend,
)
.unwrap();
let mut sub = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
let n = 50u64;
for i in 0..n {
topic.publish(Arc::new(i)).await.unwrap();
}
topic.close();
let mut messages = Vec::new();
let mut total_lagged = 0u64;
loop {
match sub.recv().await {
Ok(RecvItem::Message(env)) => messages.push(*env.payload),
Ok(RecvItem::Lagged { missed }) => total_lagged += missed,
Err(_) => break,
}
}
assert!(
messages.len() < n as usize,
"expected lag but got all {} messages",
n
);
assert!(total_lagged > 0, "expected lag notification");
for w in messages.windows(2) {
assert!(w[0] < w[1], "out of order: {} >= {}", w[0], w[1]);
}
}
#[tokio::test]
async fn tracked_publishers_resolve_independently() {
let broker = TopicBroker::new();
let base = broker
.create_topic("per-pub-ack", TopicOptions::default(), InMemoryBackend)
.unwrap();
let handle_a = base.tracked_publisher();
let handle_b = base.tracked_publisher();
let mut sub = base
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
)
.unwrap();
let receipt_a = handle_a.publish(Arc::new(100u64)).await.unwrap();
let receipt_b = handle_b.publish(Arc::new(200u64)).await.unwrap();
let env1 = match sub.recv().await.unwrap() {
RecvItem::Message(e) => e,
_ => panic!(),
};
let env2 = match sub.recv().await.unwrap() {
RecvItem::Message(e) => e,
_ => panic!(),
};
sub.ack(env1.id).unwrap();
sub.ack(env2.id).unwrap();
assert_eq!(
receipt_a.wait_for_outcome().await,
TrackedPublishOutcome::Ack
);
assert_eq!(
receipt_b.wait_for_outcome().await,
TrackedPublishOutcome::Ack
);
}
#[tokio::test]
async fn tracked_publish_broadcast_mode() {
let broker = TopicBroker::new();
let base = broker
.create_topic(
"per-pub-bcast",
TopicOptions::BroadcastOnly {
capacity: 1024,
on_lag: TopicBroadcastOnLagPolicy::DropOldest,
ack_mode: TopicBroadcastAckMode::First,
},
InMemoryBackend,
)
.unwrap();
let handle = base.tracked_publisher();
let mut sub = base
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
let receipt = handle.publish(Arc::new(42u64)).await.unwrap();
let env = match sub.recv().await.unwrap() {
RecvItem::Message(e) => e,
_ => panic!(),
};
sub.ack(env.id).unwrap();
assert_eq!(receipt.wait_for_outcome().await, TrackedPublishOutcome::Ack);
}
fn all_mode_topic(
broker: &TopicBroker<u64>,
name: &'static str,
capacity: usize,
on_lag: TopicBroadcastOnLagPolicy,
) -> crate::topic::TopicHandle<u64> {
broker
.create_topic(
name,
TopicOptions::BroadcastOnly {
capacity,
on_lag,
ack_mode: TopicBroadcastAckMode::All,
},
InMemoryBackend,
)
.unwrap()
}
fn recv_message_id(item: Result<RecvItem<u64>, Error>) -> u64 {
match item.unwrap() {
RecvItem::Message(env) => env.id,
RecvItem::Lagged { missed } => panic!("unexpected lag of {missed} messages"),
}
}
#[tokio::test]
async fn broadcast_all_mode_acks_only_after_all_subscribers_ack() {
let broker = TopicBroker::<u64>::new();
let topic = all_mode_topic(
&broker,
"all-ack",
1024,
TopicBroadcastOnLagPolicy::DropOldest,
);
let handle = topic.tracked_publisher();
let mut sub1 = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
let mut sub2 = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
let receipt = handle.publish(Arc::new(1)).await.unwrap();
let id1 = recv_message_id(sub1.recv().await);
let id2 = recv_message_id(sub2.recv().await);
assert_eq!(id1, id2);
let outcome = tokio::spawn(receipt.wait_for_outcome());
sub1.ack(id1).unwrap();
tokio::task::yield_now().await;
assert!(
!outcome.is_finished(),
"upstream resolved before all subscribers acked"
);
sub2.ack(id2).unwrap();
assert_eq!(outcome.await.unwrap(), TrackedPublishOutcome::Ack);
}
#[tokio::test]
async fn broadcast_all_mode_single_nack_resolves_nack() {
let broker = TopicBroker::<u64>::new();
let topic = all_mode_topic(
&broker,
"all-ack",
1024,
TopicBroadcastOnLagPolicy::DropOldest,
);
let handle = topic.tracked_publisher();
let mut sub1 = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
let mut sub2 = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
let receipt = handle.publish(Arc::new(1)).await.unwrap();
let id1 = recv_message_id(sub1.recv().await);
let id2 = recv_message_id(sub2.recv().await);
sub1.ack(id1).unwrap();
sub2.nack(id2, Arc::from("downstream rejected")).unwrap();
match receipt.wait_for_outcome().await {
TrackedPublishOutcome::Nack { reason } => assert_eq!(&*reason, "downstream rejected"),
other => panic!("expected Nack, got {other:?}"),
}
}
#[tokio::test]
async fn broadcast_all_mode_lag_disconnect_before_ack_nacks() {
let broker = TopicBroker::<u64>::new();
let topic = all_mode_topic(&broker, "all-lag", 4, TopicBroadcastOnLagPolicy::Disconnect);
let handle = topic.tracked_publisher();
let mut sub = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
let receipt = handle.publish(Arc::new(1)).await.unwrap();
for i in 0..8u64 {
topic.publish(Arc::new(i)).await.unwrap();
}
match sub.recv().await.unwrap() {
RecvItem::Lagged { .. } => {}
RecvItem::Message(env) => panic!("expected lag, got message {}", env.id),
}
assert!(matches!(
receipt.wait_for_outcome().await,
TrackedPublishOutcome::Nack { .. }
));
}
#[tokio::test]
async fn broadcast_all_mode_lag_drop_oldest_before_ack_nacks() {
let broker = TopicBroker::<u64>::new();
let topic = all_mode_topic(
&broker,
"all-lag-drop",
4,
TopicBroadcastOnLagPolicy::DropOldest,
);
let handle = topic.tracked_publisher();
let mut sub = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
let receipt = handle.publish(Arc::new(1)).await.unwrap();
for i in 0..8u64 {
topic.publish(Arc::new(i)).await.unwrap();
}
match sub.recv().await.unwrap() {
RecvItem::Lagged { .. } => {}
RecvItem::Message(env) => panic!("expected lag, got message {}", env.id),
}
assert!(matches!(
receipt.wait_for_outcome().await,
TrackedPublishOutcome::Nack { .. }
));
}
#[tokio::test]
async fn broadcast_all_mode_drop_before_ack_nacks() {
let broker = TopicBroker::<u64>::new();
let topic = all_mode_topic(
&broker,
"all-drop",
1024,
TopicBroadcastOnLagPolicy::DropOldest,
);
let handle = topic.tracked_publisher();
let mut sub = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
let receipt = handle.publish(Arc::new(1)).await.unwrap();
let _id = recv_message_id(sub.recv().await);
drop(sub);
assert!(matches!(
receipt.wait_for_outcome().await,
TrackedPublishOutcome::Nack { .. }
));
}
#[tokio::test]
async fn broadcast_all_mode_zero_subscribers_nacks_immediately() {
let broker = TopicBroker::<u64>::new();
let topic = all_mode_topic(
&broker,
"all-empty",
1024,
TopicBroadcastOnLagPolicy::DropOldest,
);
let handle = topic.tracked_publisher();
let receipt = handle.publish(Arc::new(1)).await.unwrap();
assert!(matches!(
receipt.wait_for_outcome().await,
TrackedPublishOutcome::Nack { reason }
if reason.as_ref() == "broadcast publish had no eligible subscribers"
));
}
#[tokio::test]
async fn broadcast_all_mode_late_subscriber_not_required() {
let broker = TopicBroker::<u64>::new();
let topic = all_mode_topic(
&broker,
"all-late",
1024,
TopicBroadcastOnLagPolicy::DropOldest,
);
let handle = topic.tracked_publisher();
let mut early = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
let receipt = handle.publish(Arc::new(1)).await.unwrap();
let _late = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
let id = recv_message_id(early.recv().await);
early.ack(id).unwrap();
assert_eq!(receipt.wait_for_outcome().await, TrackedPublishOutcome::Ack);
}
#[tokio::test]
async fn broadcast_all_mode_multi_message_nacks_only_unacked() {
let broker = TopicBroker::<u64>::new();
let topic = all_mode_topic(
&broker,
"all-multi",
1024,
TopicBroadcastOnLagPolicy::DropOldest,
);
let handle = topic.tracked_publisher();
let mut sub = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
let receipt1 = handle.publish(Arc::new(1)).await.unwrap();
let receipt2 = handle.publish(Arc::new(2)).await.unwrap();
let id1 = recv_message_id(sub.recv().await);
sub.ack(id1).unwrap();
assert_eq!(
receipt1.wait_for_outcome().await,
TrackedPublishOutcome::Ack
);
let _id2 = recv_message_id(sub.recv().await);
drop(sub);
assert!(matches!(
receipt2.wait_for_outcome().await,
TrackedPublishOutcome::Nack { .. }
));
}
#[tokio::test]
async fn broadcast_all_mode_ack_then_disconnect_keeps_ack() {
let broker = TopicBroker::<u64>::new();
let topic = all_mode_topic(
&broker,
"all-ack-then-drop",
1024,
TopicBroadcastOnLagPolicy::DropOldest,
);
let handle = topic.tracked_publisher();
let mut sub = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
let receipt = handle.publish(Arc::new(1)).await.unwrap();
let id = recv_message_id(sub.recv().await);
sub.ack(id).unwrap();
drop(sub);
assert_eq!(receipt.wait_for_outcome().await, TrackedPublishOutcome::Ack);
}
#[tokio::test]
async fn broadcast_all_mode_disconnect_then_ack_keeps_nack() {
let broker = TopicBroker::<u64>::new();
let topic = all_mode_topic(
&broker,
"all-drop-then-ack",
1024,
TopicBroadcastOnLagPolicy::DropOldest,
);
let handle = topic.tracked_publisher();
let mut held = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
let mut dropped = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
let receipt = handle.publish(Arc::new(1)).await.unwrap();
let held_id = recv_message_id(held.recv().await);
let _dropped_id = recv_message_id(dropped.recv().await);
drop(dropped);
assert!(matches!(held.ack(held_id), Err(Error::MessageNotTracked)));
assert!(matches!(
receipt.wait_for_outcome().await,
TrackedPublishOutcome::Nack { .. }
));
}
#[tokio::test]
async fn broadcast_all_mode_ack_then_nack_is_ignored() {
let broker = TopicBroker::<u64>::new();
let topic = all_mode_topic(
&broker,
"all-ack-then-nack",
1024,
TopicBroadcastOnLagPolicy::DropOldest,
);
let handle = topic.tracked_publisher();
let mut sub1 = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
let mut sub2 = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
let receipt = handle.publish(Arc::new(1)).await.unwrap();
let id1 = recv_message_id(sub1.recv().await);
let id2 = recv_message_id(sub2.recv().await);
sub1.ack(id1).unwrap();
sub1.nack(id1, Arc::from("late retraction")).unwrap();
sub2.ack(id2).unwrap();
assert_eq!(receipt.wait_for_outcome().await, TrackedPublishOutcome::Ack);
}
#[tokio::test]
async fn broadcast_all_mode_topic_close_resolves_topic_closed() {
let broker = TopicBroker::<u64>::new();
let topic = all_mode_topic(
&broker,
"all-close",
1024,
TopicBroadcastOnLagPolicy::DropOldest,
);
let handle = topic.tracked_publisher();
let mut sub = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
let receipt = handle.publish(Arc::new(1)).await.unwrap();
let _id = recv_message_id(sub.recv().await);
topic.close();
assert_eq!(
receipt.wait_for_outcome().await,
TrackedPublishOutcome::TopicClosed
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn broadcast_all_mode_concurrent_publish_subscribe_resolves() {
let broker = TopicBroker::<u64>::new();
let topic = all_mode_topic(
&broker,
"all-concurrent",
4096,
TopicBroadcastOnLagPolicy::DropOldest,
);
let handle = topic.tracked_publisher();
const MESSAGES: u64 = 64;
let mut ackers = Vec::new();
for _ in 0..3 {
let mut sub = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
ackers.push(tokio::spawn(async move {
let mut acked = 0u64;
while acked < MESSAGES {
match sub.recv().await {
Ok(RecvItem::Message(env)) => {
let _ = sub.ack(env.id);
acked += 1;
}
Ok(RecvItem::Lagged { .. }) => {}
Err(_) => break,
}
}
sub
}));
}
let topic_for_joins = topic.clone();
let joiner = tokio::spawn(async move {
let mut subs = Vec::new();
for _ in 0..MESSAGES {
let mut sub = topic_for_joins
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
while let Ok(Ok(item)) =
tokio::time::timeout(Duration::from_millis(1), sub.recv()).await
{
if let RecvItem::Message(env) = item {
let _ = sub.ack(env.id);
}
}
subs.push(sub);
tokio::task::yield_now().await;
}
subs
});
let mut receipts = Vec::new();
for i in 0..MESSAGES {
receipts.push(handle.publish(Arc::new(i)).await.unwrap());
}
let outcomes = tokio::time::timeout(Duration::from_secs(5), async {
let mut outcomes = Vec::new();
for receipt in receipts {
outcomes.push(receipt.wait_for_outcome().await);
}
outcomes
})
.await
.expect("receipts did not resolve before timeout");
for outcome in &outcomes {
assert!(
matches!(
outcome,
TrackedPublishOutcome::Ack | TrackedPublishOutcome::Nack { .. }
),
"unexpected non-terminal outcome: {outcome:?}"
);
}
let _ackers = ackers;
let _joined = joiner.await.unwrap();
}
#[tokio::test]
async fn try_publish_balanced_only_reports_drop_on_full() {
let broker = TopicBroker::<u64>::new();
let topic = broker
.create_topic(
"nb-balanced",
TopicOptions::BalancedOnly { capacity: 1 },
InMemoryBackend,
)
.unwrap();
let _sub = topic
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
)
.unwrap();
assert_eq!(
topic.try_publish(Arc::new(1)).unwrap(),
PublishOutcome::Published
);
assert_eq!(
topic.try_publish(Arc::new(2)).unwrap(),
PublishOutcome::DroppedOnFull
);
}
#[tokio::test]
async fn try_publish_mixed_rejects_broadcast_when_balanced_is_full() {
let broker = TopicBroker::<u64>::new();
let topic = broker
.create_topic(
"nb-mixed",
TopicOptions::Mixed {
balanced_capacity: 1,
broadcast_capacity: 8,
on_lag: TopicBroadcastOnLagPolicy::DropOldest,
ack_mode: TopicBroadcastAckMode::First,
},
InMemoryBackend,
)
.unwrap();
let _balanced_sub = topic
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
)
.unwrap();
let mut broadcast_sub = topic
.subscribe(SubscriptionMode::Broadcast, SubscriberOptions::default())
.unwrap();
assert_eq!(
topic.try_publish(Arc::new(10)).unwrap(),
PublishOutcome::Published
);
assert_eq!(
topic.try_publish(Arc::new(20)).unwrap(),
PublishOutcome::DroppedOnFull
);
topic.close();
let mut received = Vec::new();
while let Ok(item) = broadcast_sub.recv().await {
if let RecvItem::Message(env) = item {
received.push(*env.payload);
}
}
assert_eq!(received, vec![10]);
}
#[tokio::test]
async fn get_topic_returns_existing() {
let broker = TopicBroker::<u64>::new();
let handle = broker
.create_topic("my-topic", TopicOptions::default(), InMemoryBackend)
.unwrap();
let retrieved = broker.get_topic("my-topic");
assert!(retrieved.is_some());
assert_eq!(retrieved.unwrap().name().as_ref(), handle.name().as_ref());
}
#[tokio::test]
async fn get_topic_returns_none_for_missing() {
let broker = TopicBroker::<u64>::new();
assert!(broker.get_topic("nonexistent").is_none());
}
#[tokio::test]
async fn get_topic_required_returns_error_for_missing() {
let broker = TopicBroker::<u64>::new();
let missing = "missing";
match broker.get_topic_required(missing) {
Err(Error::UnknownTopic { topic }) => assert_eq!(topic, missing),
_ => panic!("expected UnknownTopic"),
}
}
#[tokio::test]
async fn has_topic_reflects_state() {
let broker = TopicBroker::<u64>::new();
assert!(!broker.has_topic("t1"));
_ = broker
.create_topic("t1", TopicOptions::default(), InMemoryBackend)
.unwrap();
assert!(broker.has_topic("t1"));
}
#[tokio::test]
async fn create_topics_batch_success() {
let broker = TopicBroker::<u64>::new();
let declarations = vec![
(TopicName::parse("t-a").unwrap(), TopicOptions::default()),
(TopicName::parse("t-b").unwrap(), TopicOptions::default()),
];
let handles = broker.create_topics(declarations, InMemoryBackend).unwrap();
assert_eq!(handles.len(), 2);
assert!(broker.has_topic("t-a"));
assert!(broker.has_topic("t-b"));
}
#[tokio::test]
async fn create_topics_batch_duplicate_is_rejected_atomically() {
let broker = TopicBroker::<u64>::new();
let dup = TopicName::parse("dup").unwrap();
let result = broker.create_topics(
vec![
(dup.clone(), TopicOptions::default()),
(dup.clone(), TopicOptions::default()),
],
InMemoryBackend,
);
match result {
Err(Error::TopicAlreadyExists { topic }) => assert_eq!(topic, dup),
_ => panic!("expected TopicAlreadyExists"),
}
assert!(!broker.has_topic("dup"));
}
#[tokio::test]
async fn remove_topic_closes_and_removes() {
let broker = TopicBroker::<u64>::new();
let handle = broker
.create_topic("removable", TopicOptions::default(), InMemoryBackend)
.unwrap();
let mut sub = handle
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
)
.unwrap();
handle.publish(Arc::new(1)).await.unwrap();
assert!(broker.remove_topic("removable"));
assert!(!broker.has_topic("removable"));
match handle.publish(Arc::new(2)).await {
Err(Error::TopicClosed) => {}
other => panic!("expected PublishError::Closed, got {:?}", other),
}
loop {
match sub.recv().await {
Ok(RecvItem::Message(_)) => continue,
Err(Error::SubscriptionClosed) => break,
other => panic!("expected RecvError::Closed, got {:?}", other),
}
}
}
#[tokio::test]
async fn remove_topic_returns_false_for_missing() {
let broker = TopicBroker::<u64>::new();
assert!(!broker.remove_topic("nonexistent"));
}
#[tokio::test]
async fn remove_topic_allows_recreation_with_same_name() {
let broker = TopicBroker::<u64>::new();
let old_handle = broker
.create_topic("reuse", TopicOptions::default(), InMemoryBackend)
.unwrap();
old_handle.publish(Arc::new(1)).await.unwrap();
_ = broker.remove_topic("reuse");
let new_handle = broker
.create_topic("reuse", TopicOptions::default(), InMemoryBackend)
.unwrap();
let mut sub = new_handle
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
)
.unwrap();
new_handle.publish(Arc::new(42)).await.unwrap();
new_handle.close();
let mut received = Vec::new();
while let Ok(RecvItem::Message(env)) = sub.recv().await {
received.push(*env.payload);
}
assert_eq!(received, vec![42]);
}
#[tokio::test]
async fn topic_names_returns_all_names() {
let broker = TopicBroker::<u64>::new();
_ = broker
.create_topic("alpha", TopicOptions::default(), InMemoryBackend)
.unwrap();
_ = broker
.create_topic("beta", TopicOptions::default(), InMemoryBackend)
.unwrap();
_ = broker
.create_topic("gamma", TopicOptions::default(), InMemoryBackend)
.unwrap();
let mut names: Vec<String> = broker
.topic_names()
.into_iter()
.map(|n| n.to_string())
.collect();
names.sort();
assert_eq!(names, vec!["alpha", "beta", "gamma"]);
}
#[tokio::test]
async fn close_all_closes_every_topic() {
let broker = TopicBroker::<u64>::new();
let h1 = broker
.create_topic("t1", TopicOptions::default(), InMemoryBackend)
.unwrap();
let h2 = broker
.create_topic("t2", TopicOptions::default(), InMemoryBackend)
.unwrap();
broker.close_all();
match h1.publish(Arc::new(1)).await {
Err(Error::TopicClosed) => {}
other => panic!("expected Closed, got {:?}", other),
}
match h2.publish(Arc::new(2)).await {
Err(Error::TopicClosed) => {}
other => panic!("expected Closed, got {:?}", other),
}
assert!(broker.topic_names().is_empty());
}
#[tokio::test]
async fn tracked_publish_tracker_register_after_close_resolves_immediately() {
let tracker = TrackedPublishTracker::new();
tracker.close_all();
let permit = TrackedPublishPermit::from_tokio_owned(
Arc::new(Semaphore::new(1))
.acquire_owned()
.await
.expect("semaphore should not be closed"),
);
let receipt = tracker.register(1, Duration::from_secs(30), permit);
assert_eq!(receipt.message_id(), 1);
assert_eq!(
receipt.wait_for_outcome().await,
TrackedPublishOutcome::TopicClosed
);
}
fn consensus_permit_from(sem: &Arc<Semaphore>) -> TrackedPublishPermit {
TrackedPublishPermit::from_tokio_owned(
sem.clone()
.try_acquire_owned()
.expect("semaphore permit available"),
)
}
fn consensus_permit() -> TrackedPublishPermit {
consensus_permit_from(&Arc::new(Semaphore::new(1)))
}
fn subscriber_set(ids: impl IntoIterator<Item = u64>) -> HashSet<BroadcastSubscriberId> {
ids.into_iter().map(BroadcastSubscriberId).collect()
}
#[tokio::test]
async fn consensus_resolves_ack_only_after_all_subscribers_ack() {
let tracker = TrackedPublishTracker::new();
let receipt = tracker.register_consensus(
7,
Duration::from_secs(30),
consensus_permit(),
subscriber_set([1, 2]),
7,
);
assert_eq!(
tracker.resolve_ack_from(7, BroadcastSubscriberId(1)),
AckFromResult::StillPending
);
assert_eq!(
tracker.resolve_ack_from(7, BroadcastSubscriberId(2)),
AckFromResult::Resolved
);
assert_eq!(receipt.wait_for_outcome().await, TrackedPublishOutcome::Ack);
assert_eq!(
tracker.resolve_ack_from(7, BroadcastSubscriberId(1)),
AckFromResult::NotTracked
);
}
#[tokio::test]
async fn consensus_duplicate_ack_does_not_complete() {
let tracker = TrackedPublishTracker::new();
let _receipt = tracker.register_consensus(
1,
Duration::from_secs(30),
consensus_permit(),
subscriber_set([1, 2]),
1,
);
assert_eq!(
tracker.resolve_ack_from(1, BroadcastSubscriberId(1)),
AckFromResult::StillPending
);
assert_eq!(
tracker.resolve_ack_from(1, BroadcastSubscriberId(1)),
AckFromResult::StillPending
);
assert_eq!(
tracker.resolve_ack_from(1, BroadcastSubscriberId(2)),
AckFromResult::Resolved
);
}
#[tokio::test]
async fn consensus_ack_from_non_member_does_not_advance() {
let tracker = TrackedPublishTracker::new();
let _receipt = tracker.register_consensus(
1,
Duration::from_secs(30),
consensus_permit(),
subscriber_set([1, 2]),
1,
);
assert_eq!(
tracker.resolve_ack_from(1, BroadcastSubscriberId(99)),
AckFromResult::StillPending
);
assert_eq!(
tracker.resolve_ack_from(1, BroadcastSubscriberId(1)),
AckFromResult::StillPending
);
assert_eq!(
tracker.resolve_ack_from(1, BroadcastSubscriberId(2)),
AckFromResult::Resolved
);
}
#[tokio::test]
async fn consensus_nack_from_respects_pending_membership() {
let tracker = TrackedPublishTracker::new();
let _receipt = tracker.register_consensus(
1,
Duration::from_secs(30),
consensus_permit(),
subscriber_set([1, 2]),
1,
);
assert_eq!(
tracker.resolve_ack_from(1, BroadcastSubscriberId(1)),
AckFromResult::StillPending
);
assert_eq!(
tracker.resolve_nack_from(1, BroadcastSubscriberId(1), Arc::from("retraction")),
NackFromResult::NotRequired
);
assert_eq!(
tracker.resolve_nack_from(1, BroadcastSubscriberId(2), Arc::from("rejected")),
NackFromResult::Resolved
);
assert_eq!(
tracker.resolve_nack_from(1, BroadcastSubscriberId(2), Arc::from("again")),
NackFromResult::NotTracked
);
assert_eq!(
tracker.resolve_nack_from(999, BroadcastSubscriberId(1), Arc::from("unknown")),
NackFromResult::NotTracked
);
let first_receipt = tracker.register(2, Duration::from_secs(30), consensus_permit());
assert_eq!(
tracker.resolve_nack_from(2, BroadcastSubscriberId(1), Arc::from("wrong mode")),
NackFromResult::NotTracked
);
assert!(tracker.resolve(2, TrackedPublishOutcome::Ack));
assert_eq!(
first_receipt.wait_for_outcome().await,
TrackedPublishOutcome::Ack
);
}
#[tokio::test]
async fn consensus_empty_set_resolves_nack_immediately() {
let tracker = TrackedPublishTracker::new();
let receipt = tracker.register_consensus(
1,
Duration::from_secs(30),
consensus_permit(),
HashSet::new(),
1,
);
assert!(matches!(
receipt.wait_for_outcome().await,
TrackedPublishOutcome::Nack { reason }
if reason.as_ref() == "broadcast publish had no eligible subscribers"
));
assert_eq!(
tracker.resolve_ack_from(1, BroadcastSubscriberId(1)),
AckFromResult::NotTracked
);
}
#[tokio::test]
async fn consensus_resolve_ack_from_unknown_message_not_tracked() {
let tracker = TrackedPublishTracker::new();
assert_eq!(
tracker.resolve_ack_from(99, BroadcastSubscriberId(1)),
AckFromResult::NotTracked
);
}
#[tokio::test]
async fn resolve_ack_from_first_kind_entry_is_not_tracked() {
let tracker = TrackedPublishTracker::new();
let receipt = tracker.register(1, Duration::from_secs(30), consensus_permit());
assert_eq!(
tracker.resolve_ack_from(1, BroadcastSubscriberId(1)),
AckFromResult::NotTracked
);
assert!(tracker.resolve(1, TrackedPublishOutcome::Ack));
assert_eq!(receipt.wait_for_outcome().await, TrackedPublishOutcome::Ack);
}
#[tokio::test]
async fn consensus_subscriber_disappearance_nacks_outstanding() {
let tracker = TrackedPublishTracker::new();
let receipt = tracker.register_consensus(
5,
Duration::from_secs(30),
consensus_permit(),
subscriber_set([1, 2]),
5,
);
tracker.nack_pending_for_subscriber(
BroadcastSubscriberId(2),
Arc::from("subscriber disconnected"),
);
match receipt.wait_for_outcome().await {
TrackedPublishOutcome::Nack { reason } => assert_eq!(&*reason, "subscriber disconnected"),
other => panic!("expected Nack, got {other:?}"),
}
assert_eq!(
tracker.resolve_ack_from(5, BroadcastSubscriberId(1)),
AckFromResult::NotTracked
);
}
#[tokio::test]
async fn consensus_disappearance_only_nacks_entries_requiring_subscriber() {
let tracker = TrackedPublishTracker::new();
let r1 = tracker.register_consensus(
1,
Duration::from_secs(30),
consensus_permit(),
subscriber_set([1]),
1,
);
let r2 = tracker.register_consensus(
2,
Duration::from_secs(30),
consensus_permit(),
subscriber_set([2]),
2,
);
tracker.nack_pending_for_subscriber(BroadcastSubscriberId(1), Arc::from("gone"));
assert!(matches!(
r1.wait_for_outcome().await,
TrackedPublishOutcome::Nack { .. }
));
assert_eq!(
tracker.resolve_ack_from(2, BroadcastSubscriberId(2)),
AckFromResult::Resolved
);
assert_eq!(r2.wait_for_outcome().await, TrackedPublishOutcome::Ack);
}
#[tokio::test]
async fn consensus_nack_owed_before_respects_seq_threshold() {
let tracker = TrackedPublishTracker::new();
let skipped = tracker.register_consensus(
10,
Duration::from_secs(30),
consensus_permit(),
subscriber_set([1]),
3,
);
let readable = tracker.register_consensus(
11,
Duration::from_secs(30),
consensus_permit(),
subscriber_set([1]),
5,
);
tracker.nack_owed_before(
BroadcastSubscriberId(1),
5,
Arc::from("lagged past message"),
);
match skipped.wait_for_outcome().await {
TrackedPublishOutcome::Nack { reason } => assert_eq!(&*reason, "lagged past message"),
other => panic!("expected Nack for skipped message, got {other:?}"),
}
assert_eq!(
tracker.resolve_ack_from(11, BroadcastSubscriberId(1)),
AckFromResult::Resolved
);
assert_eq!(
readable.wait_for_outcome().await,
TrackedPublishOutcome::Ack
);
}
#[tokio::test]
async fn consensus_single_resolve_nack_overrides_consensus() {
let tracker = TrackedPublishTracker::new();
let receipt = tracker.register_consensus(
1,
Duration::from_secs(30),
consensus_permit(),
subscriber_set([1, 2]),
1,
);
assert_eq!(
tracker.resolve_ack_from(1, BroadcastSubscriberId(1)),
AckFromResult::StillPending
);
assert!(tracker.resolve(
1,
TrackedPublishOutcome::Nack {
reason: Arc::from("boom")
}
));
assert!(matches!(
receipt.wait_for_outcome().await,
TrackedPublishOutcome::Nack { .. }
));
assert_eq!(
tracker.resolve_ack_from(1, BroadcastSubscriberId(2)),
AckFromResult::NotTracked
);
}
#[tokio::test]
async fn consensus_ack_completion_beats_late_disconnect() {
let tracker = TrackedPublishTracker::new();
let receipt = tracker.register_consensus(
1,
Duration::from_secs(30),
consensus_permit(),
subscriber_set([1]),
1,
);
assert_eq!(
tracker.resolve_ack_from(1, BroadcastSubscriberId(1)),
AckFromResult::Resolved
);
tracker.nack_pending_for_subscriber(BroadcastSubscriberId(1), Arc::from("late"));
assert_eq!(receipt.wait_for_outcome().await, TrackedPublishOutcome::Ack);
}
#[tokio::test]
async fn consensus_register_after_close_resolves_topic_closed() {
let tracker = TrackedPublishTracker::new();
tracker.close_all();
let receipt = tracker.register_consensus(
1,
Duration::from_secs(30),
consensus_permit(),
subscriber_set([1]),
1,
);
assert_eq!(
receipt.wait_for_outcome().await,
TrackedPublishOutcome::TopicClosed
);
}
#[tokio::test]
async fn consensus_releases_permit_on_resolution() {
let sem = Arc::new(Semaphore::new(1));
let tracker = TrackedPublishTracker::new();
let receipt = tracker.register_consensus(
1,
Duration::from_secs(30),
consensus_permit_from(&sem),
subscriber_set([1]),
1,
);
assert_eq!(sem.available_permits(), 0);
assert_eq!(
tracker.resolve_ack_from(1, BroadcastSubscriberId(1)),
AckFromResult::Resolved
);
let _ = receipt.wait_for_outcome().await;
assert_eq!(sem.available_permits(), 1);
}
#[tokio::test]
async fn topic_set_insert_and_get() {
let broker = TopicBroker::<u64>::new();
let handle = broker
.create_topic("t1", TopicOptions::default(), InMemoryBackend)
.unwrap();
let set = TopicSet::new("pipeline-1");
_ = set.insert("output", handle);
let retrieved = set.get("output");
assert!(retrieved.is_some());
assert_eq!(retrieved.unwrap().name().as_ref(), "t1");
}
#[tokio::test]
async fn topic_set_get_missing_returns_none() {
let set = TopicSet::<u64>::new("empty-set");
assert!(set.get("nonexistent").is_none());
}
#[tokio::test]
async fn topic_set_get_required_returns_error_for_missing() {
let set = TopicSet::<u64>::new("empty-set");
let missing = "missing-local";
match set.get_required(missing) {
Err(Error::UnknownTopic { topic }) => assert_eq!(topic, missing),
_ => panic!("expected UnknownTopic"),
}
}
#[tokio::test]
async fn topic_set_remove() {
let broker = TopicBroker::<u64>::new();
let handle = broker
.create_topic("t1", TopicOptions::default(), InMemoryBackend)
.unwrap();
let set = TopicSet::new("p1");
_ = set.insert("output", handle);
assert!(set.contains("output"));
let removed = set.remove("output");
assert!(removed.is_some());
assert!(!set.contains("output"));
}
#[tokio::test]
async fn topic_set_remove_does_not_close_topic() {
let broker = TopicBroker::<u64>::new();
let handle = broker
.create_topic("t1", TopicOptions::default(), InMemoryBackend)
.unwrap();
let set = TopicSet::new("p1");
_ = set.insert("output", handle.clone());
_ = set.remove("output");
handle.publish(Arc::new(42)).await.unwrap();
}
#[tokio::test]
async fn topic_set_insert_overwrites() {
let broker = TopicBroker::<u64>::new();
let h1 = broker
.create_topic("t1", TopicOptions::default(), InMemoryBackend)
.unwrap();
let h2 = broker
.create_topic("t2", TopicOptions::default(), InMemoryBackend)
.unwrap();
let set = TopicSet::new("p1");
assert!(set.insert("output", h1).is_none());
let prev = set.insert("output", h2);
assert!(prev.is_some());
assert_eq!(prev.unwrap().name().as_ref(), "t1");
assert_eq!(set.get("output").unwrap().name().as_ref(), "t2");
}
#[tokio::test]
async fn topic_set_topic_names() {
let broker = TopicBroker::<u64>::new();
let h1 = broker
.create_topic("t1", TopicOptions::default(), InMemoryBackend)
.unwrap();
let h2 = broker
.create_topic("t2", TopicOptions::default(), InMemoryBackend)
.unwrap();
let set = TopicSet::new("p1");
_ = set.insert("alpha", h1);
_ = set.insert("beta", h2);
let mut names: Vec<String> = set
.topic_names()
.into_iter()
.map(|n| n.to_string())
.collect();
names.sort();
assert_eq!(names, vec!["alpha", "beta"]);
}
#[tokio::test]
async fn topic_set_tracked_publisher() {
let broker = TopicBroker::<u64>::new();
let handle = broker
.create_topic("t1", TopicOptions::default(), InMemoryBackend)
.unwrap();
let set = TopicSet::new("p1");
_ = set.insert("output", handle);
let handle = set.get("output").unwrap();
let pub_handle = handle.tracked_publisher();
let mut sub = handle
.subscribe(
SubscriptionMode::Balanced {
group: SubscriptionGroupName::from("g1"),
},
SubscriberOptions::default(),
)
.unwrap();
let receipt = pub_handle.publish(Arc::new(99u64)).await.unwrap();
let env = match sub.recv().await.unwrap() {
RecvItem::Message(e) => e,
_ => panic!(),
};
sub.ack(env.id).unwrap();
assert_eq!(receipt.wait_for_outcome().await, TrackedPublishOutcome::Ack);
}
#[tokio::test]
async fn topic_set_name() {
let set = TopicSet::<u64>::new("my-pipeline");
assert_eq!(set.name(), "my-pipeline");
}
#[tokio::test]
async fn topic_set_clone_shares_state() {
let broker = TopicBroker::<u64>::new();
let handle = broker
.create_topic("t1", TopicOptions::default(), InMemoryBackend)
.unwrap();
let set1 = TopicSet::new("p1");
let set2 = set1.clone();
_ = set1.insert("output", handle);
assert!(set2.contains("output"));
assert_eq!(set2.len(), 1);
}
#[tokio::test]
async fn topic_set_len_and_is_empty() {
let broker = TopicBroker::<u64>::new();
let handle = broker
.create_topic("t1", TopicOptions::default(), InMemoryBackend)
.unwrap();
let set = TopicSet::new("p1");
assert!(set.is_empty());
assert_eq!(set.len(), 0);
_ = set.insert("output", handle);
assert!(!set.is_empty());
assert_eq!(set.len(), 1);
}
#[tokio::test]
async fn broadcast_ring_stale_commit_does_not_overwrite_newer_wrapped_sequence() {
let ring = FastBroadcastRing::new(2);
let stale_seq = ring.reserve_seq();
let tracker = TrackedPublishTracker::new();
let receipt = tracker.register_consensus(
stale_seq,
Duration::from_secs(30),
consensus_permit(),
subscriber_set([1]),
stale_seq,
);
let next_seq = ring.reserve_seq();
assert!(ring.commit_slot(
next_seq,
Envelope {
id: next_seq,
tracked: false,
payload: Arc::new(next_seq),
},
));
let wrapped_seq = ring.reserve_seq();
assert!(ring.commit_slot(
wrapped_seq,
Envelope {
id: wrapped_seq,
tracked: false,
payload: Arc::new(wrapped_seq),
},
));
let committed = ring.commit_slot(
stale_seq,
Envelope {
id: stale_seq,
tracked: true,
payload: Arc::new(stale_seq),
},
);
assert!(!committed);
if !committed {
let _ = tracker.resolve(
stale_seq,
TrackedPublishOutcome::Nack {
reason: Arc::from("broadcast message was overtaken before ring commit"),
},
);
}
match ring.try_read(wrapped_seq) {
BroadcastReadResult::Ready(envelope) => {
assert_eq!(envelope.id, wrapped_seq);
assert_eq!(*envelope.payload, wrapped_seq);
}
_ => panic!("wrapped sequence should remain readable"),
}
assert!(matches!(
receipt.wait_for_outcome().await,
TrackedPublishOutcome::Nack { .. }
));
}
#[tokio::test]
async fn consensus_disconnect_does_not_nack_completed_ack() {
let tracker = TrackedPublishTracker::new();
let receipt = tracker.register_consensus(
1,
Duration::from_secs(30),
consensus_permit(),
subscriber_set([1, 2]),
1,
);
assert_eq!(
tracker.resolve_ack_from(1, BroadcastSubscriberId(1)),
AckFromResult::StillPending
);
tracker.nack_pending_for_subscriber(BroadcastSubscriberId(1), Arc::from("disconnected"));
assert_eq!(
tracker.resolve_ack_from(1, BroadcastSubscriberId(2)),
AckFromResult::Resolved
);
assert_eq!(receipt.wait_for_outcome().await, TrackedPublishOutcome::Ack);
}