This is an automated email from the ASF dual-hosted git repository.

quantranhong1999 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/james-project.git

commit 12b00468d12005e7070b7d9b3ddd78b9d09808d5
Author: Quan Tran <[email protected]>
AuthorDate: Wed Oct 7 15:10:30 2026 +0700

    [FIX] KeyRegistrationHandler: bound the event bus stop
    
    Stopping a RabbitMQ event bus deleted its key registration queue first,
    with a one minute timeout retried up to 8 times, and only then disposed the
    registration consumer and its scheduler. With an unresponsive broker one
    bus could block stop for about ten minutes, then skip the local cleanup.
    Each bus is stopped in turn, and these stops now run on every server
    shutdown since provider-built singletons get their @PreDestroy called.
    
    The consumer and scheduler are now disposed first. Deleting the queue is a
    single best effort attempt bounded to 10 seconds; a failure is logged
    instead of thrown. The registration queue is auto-delete by default with
    classic queues; a queue left behind otherwise is the same leftover a crash
    already leaves.
    
    New test: stopping the event bus while RabbitMQ is paused returns within
    30 seconds (it timed out before). RabbitMQEventBusTest: 81 tests green.
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 .../org/apache/james/events/KeyRegistrationHandler.java     | 13 +++++++++----
 .../java/org/apache/james/events/RabbitMQEventBusTest.java  |  9 +++++++++
 2 files changed, 18 insertions(+), 4 deletions(-)

diff --git 
a/event-bus/distributed/src/main/java/org/apache/james/events/KeyRegistrationHandler.java
 
b/event-bus/distributed/src/main/java/org/apache/james/events/KeyRegistrationHandler.java
index b1a53a4444..2dc2f2d5ab 100644
--- 
a/event-bus/distributed/src/main/java/org/apache/james/events/KeyRegistrationHandler.java
+++ 
b/event-bus/distributed/src/main/java/org/apache/james/events/KeyRegistrationHandler.java
@@ -59,6 +59,7 @@ class KeyRegistrationHandler {
     private static final Logger LOGGER = 
LoggerFactory.getLogger(KeyRegistrationHandler.class);
 
     private static final Duration TOPOLOGY_CHANGES_TIMEOUT = 
Duration.ofMinutes(1);
+    private static final Duration STOP_QUEUE_DELETION_TIMEOUT = 
Duration.ofSeconds(10);
 
     private final EventBusId eventBusId;
     private final LocalListenerRegistry localListenerRegistry;
@@ -135,13 +136,17 @@ class KeyRegistrationHandler {
     }
 
     void stop() {
-        sender.delete(QueueSpecification.queue(registrationQueue.asString()))
-            .timeout(TOPOLOGY_CHANGES_TIMEOUT)
-            
.retryWhen(configurations.retryBackoff().asReactorRetry().scheduler(Schedulers.parallel()))
-            .block();
         receiverSubscriber.filter(Predicate.not(Disposable::isDisposed))
                 .ifPresent(Disposable::dispose);
         Optional.ofNullable(scheduler).ifPresent(Scheduler::dispose);
+        // Best effort and bounded: stopping must not wait minutes on an 
unavailable broker
+        sender.delete(QueueSpecification.queue(registrationQueue.asString()))
+            .timeout(STOP_QUEUE_DELETION_TIMEOUT)
+            .onErrorResume(e -> {
+                LOGGER.warn("Could not delete key registration queue {} while 
stopping the event bus", registrationQueue.asString(), e);
+                return Mono.empty();
+            })
+            .block();
     }
 
     Mono<Registration> register(EventListener.ReactiveEventListener listener, 
RegistrationKey key) {
diff --git 
a/event-bus/distributed/src/test/java/org/apache/james/events/RabbitMQEventBusTest.java
 
b/event-bus/distributed/src/test/java/org/apache/james/events/RabbitMQEventBusTest.java
index ce12c5d628..52bcc47b18 100644
--- 
a/event-bus/distributed/src/test/java/org/apache/james/events/RabbitMQEventBusTest.java
+++ 
b/event-bus/distributed/src/test/java/org/apache/james/events/RabbitMQEventBusTest.java
@@ -42,6 +42,7 @@ import static org.awaitility.Awaitility.await;
 import static org.awaitility.Durations.FIVE_SECONDS;
 import static org.awaitility.Durations.TEN_MINUTES;
 import static org.awaitility.Durations.TEN_SECONDS;
+import static org.junit.jupiter.api.Assertions.assertTimeoutPreemptively;
 import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.Mockito.doAnswer;
 import static org.mockito.Mockito.mock;
@@ -766,6 +767,14 @@ class RabbitMQEventBusTest implements 
GroupContract.SingleEventBusGroupContract,
                     .anySatisfy(queue -> 
assertThat(queue.getName()).contains(GroupA.class.getName()));
             }
 
+            @Test
+            void stopShouldNotWaitLongWhenRabbitMQIsUnavailable() {
+                eventBus.start();
+                rabbitMQExtension.getRabbitMQ().pause();
+
+                assertTimeoutPreemptively(Duration.ofSeconds(30), () -> 
eventBus.stop());
+            }
+
             @Test
             void eventBusShouldNotThrowWhenContinuouslyStartAndStop() {
                 assertThatCode(() -> {


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to