diff --git a/src/backend/broadcast.rs b/src/backend/broadcast.rs index 8c5717b2..34b86745 100644 --- a/src/backend/broadcast.rs +++ b/src/backend/broadcast.rs @@ -186,9 +186,9 @@ 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. @@ -196,6 +196,7 @@ pub(crate) use tail_only::refuse_start_other_than_tail; feature = "inmemory", feature = "kafka", feature = "nats", + feature = "rabbitmq", feature = "redis-streams" ))] mod settling { @@ -352,6 +353,7 @@ mod settling { feature = "inmemory", feature = "kafka", feature = "nats", + feature = "rabbitmq", feature = "redis-streams" ))] pub(crate) use settling::BROADCAST_DEFER_DELAY; diff --git a/src/backends/rabbitmq/consumer.rs b/src/backends/rabbitmq/consumer.rs index b9f6abb7..4243879d 100644 --- a/src/backends/rabbitmq/consumer.rs +++ b/src/backends/rabbitmq/consumer.rs @@ -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, @@ -1247,6 +1248,7 @@ impl RabbitMqConsumer { &publisher, retry_count, group.as_deref(), + &options.shutdown, ) .await?; } @@ -1264,6 +1266,7 @@ impl RabbitMqConsumer { &publisher, retry_count, group.as_deref(), + &options.shutdown, ) .await?; } @@ -1314,6 +1317,7 @@ impl RabbitMqConsumer { &publisher, retry_count, group.as_deref(), + &options.shutdown, ) .await .ok(); @@ -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 } } } @@ -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 { @@ -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 + } } } diff --git a/tests/inmemory_broadcast.rs b/tests/inmemory_broadcast.rs index 85fd79fb..49aa1240 100644 --- a/tests/inmemory_broadcast.rs +++ b/tests/inmemory_broadcast.rs @@ -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>>, +} + +impl MessageHandler 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::::from_client(client.clone()); + let publisher = broker.publisher().await.expect("publisher"); + + let handler = DeferFirstTimed::default(); + let mut subscriber = broker.broadcast_subscriber(); + subscriber + .subscribe::(handler.clone(), ConsumerOptions::new()) + .expect("subscribe"); + wait_for_subscribers(&client, "cache-invalidations-bcast", 1).await; + + publisher + .publish::(&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] diff --git a/tests/rabbitmq_broadcast_integration.rs b/tests/rabbitmq_broadcast_integration.rs index f61113af..65229065 100644 --- a/tests/rabbitmq_broadcast_integration.rs +++ b/tests/rabbitmq_broadcast_integration.rs @@ -244,6 +244,14 @@ define_topic!( .build() ); +define_topic!( + DeferInPlaceTopic, + Invalidate, + TopologyBuilder::new("rmq-broadcast-defer-in-place") + .broadcast() + .build() +); + define_topic!( RetryTopic, Invalidate, @@ -328,6 +336,27 @@ impl MessageHandler for DeferOnce { } } +/// Defers the first sighting of `"a"` and acks everything else, recording when +/// each call started. +#[derive(Clone, Default)] +struct DeferFirstTimed { + calls: Arc>>, +} + +impl MessageHandler 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 // --------------------------------------------------------------------------- @@ -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::() + .await + .expect("failed to declare broadcast topology"); + + let handler = DeferFirstTimed::default(); + let mut sub = broker.broadcast_subscriber(); + sub.subscribe::(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::(&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