Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 5 additions & 3 deletions src/backend/broadcast.rs
Original file line number Diff line number Diff line change
Expand Up @@ -186,16 +186,17 @@ pub(crate) use tail_only::refuse_start_other_than_tail;
// settles through `router::route_reject` / `nack_requeue` because an AMQP
// delivery has to be nacked on the channel it arrived on — the router already
// records the terminal metric there, so going through this helper too would
// count it twice. InMemory *is* in the cfg list, but only for
// `BROADCAST_DEFER_DELAY`: its in-place broadcast `Defer` paces on the same
// constant, so the test substrate redelivers on the schedule production does.
// count it twice. InMemory and RabbitMQ *are* in the cfg list, but only for
// `BROADCAST_DEFER_DELAY`: their in-place broadcast `Defer` paces on the same
// constant, so every backend redelivers on one schedule.
// The re-exports below are split to match — a name re-exported under a feature
// that never uses it is an unused import, which `-D warnings` rejects and
// CI's per-feature clippy legs would catch.
#[cfg(any(
feature = "inmemory",
feature = "kafka",
feature = "nats",
feature = "rabbitmq",
feature = "redis-streams"
))]
mod settling {
Expand Down Expand Up @@ -352,6 +353,7 @@ mod settling {
feature = "inmemory",
feature = "kafka",
feature = "nats",
feature = "rabbitmq",
feature = "redis-streams"
))]
pub(crate) use settling::BROADCAST_DEFER_DELAY;
Expand Down
26 changes: 20 additions & 6 deletions src/backends/rabbitmq/consumer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ use crate::backend::batch_consumer::settling::{
use crate::backend::batch_consumer::{
BatchConsumerOptionsInner, BatchSettlement, settle_batch_outcome,
};
use crate::backend::broadcast::BROADCAST_DEFER_DELAY;
use crate::backends::rabbitmq::client::RabbitMqClient;
use crate::backends::rabbitmq::headers::{
extract_dead_metadata, extract_message_metadata, get_retry_count,
Expand Down Expand Up @@ -1247,6 +1248,7 @@ impl RabbitMqConsumer {
&publisher,
retry_count,
group.as_deref(),
&options.shutdown,
)
.await?;
}
Expand All @@ -1264,6 +1266,7 @@ impl RabbitMqConsumer {
&publisher,
retry_count,
group.as_deref(),
&options.shutdown,
)
.await?;
}
Expand Down Expand Up @@ -1314,6 +1317,7 @@ impl RabbitMqConsumer {
&publisher,
retry_count,
group.as_deref(),
&options.shutdown,
)
.await
.ok();
Expand Down Expand Up @@ -1521,13 +1525,14 @@ async fn route_outcome_for(
publisher: &ChannelPublisher,
retry_count: u32,
group: Option<&str>,
shutdown: &CancellationToken,
) -> Result<()> {
match attachment {
Attachment::Shared(_) => {
route_outcome(received, outcome, topology, publisher, retry_count, group).await
}
Attachment::Broadcast { .. } => {
route_broadcast_outcome(received, outcome, topology, publisher, group).await
route_broadcast_outcome(received, outcome, topology, publisher, group, shutdown).await
}
}
}
Expand All @@ -1542,16 +1547,19 @@ async fn route_outcome_for(
/// — matching what `decide_retry` yields on the InMemory path at
/// `max_retries = 0`, so one dashboard reads both backends.
///
/// `Defer` nack-requeues. That is redelivery to *this subscriber only*, which
/// is the contract rather than an approximation of it: the queue is exclusive,
/// so it has exactly one consumer and a requeued message can reach no other
/// instance's copy of the fan-out.
/// `Defer` waits [`BROADCAST_DEFER_DELAY`], then nack-requeues. That is
/// redelivery to *this subscriber only*, which is the contract rather than an
/// approximation of it: the queue is exclusive, so it has exactly one consumer
/// and a requeued message can reach no other instance's copy of the fan-out.
/// The wait holds the single delivery slot, so the requeued message is handed
/// back before anything queued behind it.
async fn route_broadcast_outcome(
received: &ReceivedDelivery,
outcome: Outcome,
topology: &'static QueueTopology,
publisher: &ChannelPublisher,
group: Option<&str>,
shutdown: &CancellationToken,
) -> Result<()> {
let delivery = &received.delivery;
match outcome {
Expand Down Expand Up @@ -1581,7 +1589,13 @@ async fn route_broadcast_outcome(
)
.await
}
Outcome::Defer => router::nack_requeue(delivery, publisher).await,
Outcome::Defer => {
tokio::select! {
_ = tokio::time::sleep(BROADCAST_DEFER_DELAY) => {}
_ = shutdown.cancelled() => {}
}
router::nack_requeue(delivery, publisher).await
}
}
}

Expand Down
59 changes: 59 additions & 0 deletions tests/inmemory_broadcast.rs
Original file line number Diff line number Diff line change
Expand Up @@ -707,6 +707,65 @@ async fn a_deferred_message_is_retried_in_place_before_those_behind_it() {
);
}

/// Defers the first sighting of key `1` and acks everything else, recording
/// when each call started.
#[derive(Clone, Default)]
struct DeferFirstTimed {
calls: Arc<Mutex<Vec<(u64, tokio::time::Instant)>>>,
}

impl MessageHandler<CacheInvalidations> for DeferFirstTimed {
type Context = ();
async fn handle(&self, msg: Invalidate, _meta: MessageMetadata, _ctx: &()) -> Outcome {
let mut calls = self.calls.lock().expect("calls lock");
let first = msg.key == 1 && calls.iter().all(|(k, _)| *k != 1);
calls.push((msg.key, tokio::time::Instant::now()));
if first { Outcome::Defer } else { Outcome::Ack }
}
}

/// A deferred broadcast message waits the broadcast defer delay before it is
/// handed back, the same pacing the production backends apply.
#[tokio::test]
async fn a_deferred_message_waits_the_defer_delay_before_redelivery() {
// Mirrors the crate-private `BROADCAST_DEFER_DELAY`.
const DEFER_DELAY: Duration = Duration::from_secs(1);

let client = InMemoryBroker::new();
let broker = Broker::<InMemory>::from_client(client.clone());
let publisher = broker.publisher().await.expect("publisher");

let handler = DeferFirstTimed::default();
let mut subscriber = broker.broadcast_subscriber();
subscriber
.subscribe::<CacheInvalidations, _>(handler.clone(), ConsumerOptions::new())
.expect("subscribe");
wait_for_subscribers(&client, "cache-invalidations-bcast", 1).await;

publisher
.publish::<CacheInvalidations>(&Invalidate { key: 1 })
.await
.expect("publish");

let calls = Arc::clone(&handler.calls);
wait_until("the deferred message is redelivered", || {
calls.lock().expect("calls lock").len() == 2
})
.await;

subscriber.cancellation_token().cancel();
let _ = subscriber
.run_until_timeout(std::future::pending::<()>(), Duration::from_secs(2))
.await;

let calls = handler.calls.lock().expect("calls lock").clone();
let gap = calls[1].1.duration_since(calls[0].1);
assert!(
gap >= DEFER_DELAY,
"a deferred broadcast message was redelivered after {gap:?}, before the {DEFER_DELAY:?} defer delay"
);
}

/// AC8 at the wire level rather than the name level: an ordinary topology still
/// behaves exactly as it did, on the same broker, alongside broadcast traffic.
#[tokio::test]
Expand Down
91 changes: 91 additions & 0 deletions tests/rabbitmq_broadcast_integration.rs
Original file line number Diff line number Diff line change
Expand Up @@ -244,6 +244,14 @@ define_topic!(
.build()
);

define_topic!(
DeferInPlaceTopic,
Invalidate,
TopologyBuilder::new("rmq-broadcast-defer-in-place")
.broadcast()
.build()
);

define_topic!(
RetryTopic,
Invalidate,
Expand Down Expand Up @@ -328,6 +336,27 @@ impl MessageHandler<DeferTopic> for DeferOnce {
}
}

/// Defers the first sighting of `"a"` and acks everything else, recording when
/// each call started.
#[derive(Clone, Default)]
struct DeferFirstTimed {
calls: Arc<Mutex<Vec<(String, Instant)>>>,
}

impl MessageHandler<DeferInPlaceTopic> for DeferFirstTimed {
type Context = ();
async fn handle(&self, msg: Invalidate, _meta: MessageMetadata, _: &()) -> Outcome {
let mut calls = self.calls.lock().await;
let first_a = msg.key == "a" && calls.iter().all(|(k, _)| k != "a");
calls.push((msg.key, Instant::now()));
if first_a {
Outcome::Defer
} else {
Outcome::Ack
}
}
}

// ---------------------------------------------------------------------------
// Tests
// ---------------------------------------------------------------------------
Expand Down Expand Up @@ -596,6 +625,68 @@ async fn defer_redelivers_only_to_the_subscriber_that_deferred() {
ctx.cleanup().await;
}

/// A deferred broadcast message comes back to the same handler after the
/// broadcast defer delay, ahead of the message queued behind it — the pacing
/// and order the Kafka, NATS, Redis and InMemory broadcast loops share.
#[tokio::test]
async fn defer_waits_the_defer_delay_and_redelivers_in_place() {
// Mirrors the crate-private `BROADCAST_DEFER_DELAY`.
const DEFER_DELAY: Duration = Duration::from_secs(1);

let ctx = TestContext::new().await;
let broker = ctx.broker().await;
broker
.topology()
.declare::<DeferInPlaceTopic>()
.await
.expect("failed to declare broadcast topology");

let handler = DeferFirstTimed::default();
let mut sub = broker.broadcast_subscriber();
sub.subscribe::<DeferInPlaceTopic, _>(handler.clone(), ConsumerOptions::new())
.expect("failed to subscribe");

let deadline = Instant::now() + Duration::from_secs(20);
while ctx.queue_names().await.is_empty() && Instant::now() < deadline {
tokio::time::sleep(Duration::from_millis(100)).await;
}
assert_eq!(ctx.queue_names().await.len(), 1);

let publisher = broker.publisher().await.expect("failed to build publisher");
for key in ["a", "b"] {
publisher
.publish::<DeferInPlaceTopic>(&Invalidate { key: key.into() })
.await
.expect("broadcast publish failed");
}

let deadline = Instant::now() + Duration::from_secs(20);
while handler.calls.lock().await.len() < 3 && Instant::now() < deadline {
tokio::time::sleep(Duration::from_millis(50)).await;
}

sub.cancellation_token().cancel();
let _ = sub
.run_until_timeout(std::future::pending(), Duration::from_secs(10))
.await;
broker.close().await;

let calls = handler.calls.lock().await.clone();
let keys: Vec<&str> = calls.iter().map(|(k, _)| k.as_str()).collect();
assert_eq!(
keys,
vec!["a", "a", "b"],
"a deferred broadcast message is retried in place, before the message behind it"
);
let gap = calls[1].1.duration_since(calls[0].1);
assert!(
gap >= DEFER_DELAY,
"a deferred broadcast message was redelivered after {gap:?}, before the {DEFER_DELAY:?} defer delay"
);

ctx.cleanup().await;
}

/// `Retry` on a broadcast subscription discards, and does not loop.
///
/// RabbitMQ's shared `route_retry` falls back to a nack-**requeue** when a
Expand Down
Loading