MAILBOX-376 Plug DeadLetters into RabbitMQEventBus

Project: http://git-wip-us.apache.org/repos/asf/james-project/repo
Commit: http://git-wip-us.apache.org/repos/asf/james-project/commit/98cc6604
Tree: http://git-wip-us.apache.org/repos/asf/james-project/tree/98cc6604
Diff: http://git-wip-us.apache.org/repos/asf/james-project/diff/98cc6604

Branch: refs/heads/master
Commit: 98cc66045851c7575236862bde64ed7e86cb4609
Parents: 00973f9
Author: tran tien duc <[email protected]>
Authored: Mon Jan 21 11:50:57 2019 +0700
Committer: Benoit Tellier <[email protected]>
Committed: Tue Jan 22 09:20:12 2019 +0700

----------------------------------------------------------------------
 .../mailbox/events/ErrorHandlingContract.java   | 63 ++++++++++++++++++++
 .../james/mailbox/events/InVMEventBusTest.java  | 25 ++++++++
 mailbox/event/event-rabbitmq/pom.xml            |  5 ++
 .../mailbox/events/GroupConsumerRetry.java      | 32 ++++++----
 .../james/mailbox/events/GroupRegistration.java |  3 +-
 .../events/GroupRegistrationHandler.java        |  7 ++-
 .../james/mailbox/events/RabbitMQEventBus.java  |  9 ++-
 .../mailbox/events/RabbitMQEventBusTest.java    | 10 ++--
 8 files changed, 134 insertions(+), 20 deletions(-)
----------------------------------------------------------------------


http://git-wip-us.apache.org/repos/asf/james-project/blob/98cc6604/mailbox/api/src/test/java/org/apache/james/mailbox/events/ErrorHandlingContract.java
----------------------------------------------------------------------
diff --git 
a/mailbox/api/src/test/java/org/apache/james/mailbox/events/ErrorHandlingContract.java
 
b/mailbox/api/src/test/java/org/apache/james/mailbox/events/ErrorHandlingContract.java
index 069798d..be4b7ef 100644
--- 
a/mailbox/api/src/test/java/org/apache/james/mailbox/events/ErrorHandlingContract.java
+++ 
b/mailbox/api/src/test/java/org/apache/james/mailbox/events/ErrorHandlingContract.java
@@ -60,6 +60,8 @@ interface ErrorHandlingContract extends EventBusContract {
         }
     }
 
+    EventDeadLetters deadLetter();
+
     default EventCollector eventCollector() {
         return spy(new EventCollector());
     }
@@ -196,4 +198,65 @@ interface ErrorHandlingContract extends EventBusContract {
 
         WAIT_CONDITION.until(successfulRetry::get);
     }
+
+    @Test
+    default void deadLettersIsNotAppliedForKeyRegistrations() throws Exception 
{
+        EventCollector eventCollector = eventCollector();
+
+        doThrow(new RuntimeException())
+            .doThrow(new RuntimeException())
+            .doThrow(new RuntimeException())
+            .doThrow(new RuntimeException())
+            .doCallRealMethod()
+            .when(eventCollector).event(EVENT);
+
+        eventBus().register(eventCollector, KEY_1);
+        eventBus().dispatch(EVENT, ImmutableSet.of(KEY_1)).block();
+
+        TimeUnit.SECONDS.sleep(1);
+        SoftAssertions.assertSoftly(softly -> {
+            softly.assertThat(eventCollector.getEvents()).isEmpty();
+            
softly.assertThat(deadLetter().groupsWithFailedEvents().toIterable())
+                .isEmpty();
+        });
+    }
+
+    @Test
+    default void deadLetterShouldNotStoreWhenFailsLessThanMaxRetries() {
+        EventCollector eventCollector = eventCollector();
+
+        doThrow(new RuntimeException())
+            .doThrow(new RuntimeException())
+            .doCallRealMethod()
+            .when(eventCollector).event(EVENT);
+
+        eventBus().register(eventCollector, new EventBusTestFixture.GroupA());
+        eventBus().dispatch(EVENT, NO_KEYS).block();
+
+        WAIT_CONDITION
+            .until(() -> assertThat(eventCollector.getEvents()).hasSize(1));
+
+        assertThat(deadLetter().groupsWithFailedEvents().toIterable())
+            .isEmpty();
+    }
+
+    @Test
+    default void deadLetterShouldStoreWhenFailsGreaterThanMaxRetries() throws 
Exception {
+        EventCollector eventCollector = eventCollector();
+
+        doThrow(new RuntimeException())
+            .doThrow(new RuntimeException())
+            .doThrow(new RuntimeException())
+            .doThrow(new RuntimeException())
+            .doCallRealMethod()
+            .when(eventCollector).event(EVENT);
+
+        eventBus().register(eventCollector, GROUP_A);
+        eventBus().dispatch(EVENT, NO_KEYS).block();
+
+        WAIT_CONDITION.until(() -> 
assertThat(deadLetter().failedEventIds(GROUP_A).toIterable())
+            .containsOnly(EVENT.getEventId()));
+        assertThat(eventCollector.getEvents())
+            .isEmpty();
+    }
 }

http://git-wip-us.apache.org/repos/asf/james-project/blob/98cc6604/mailbox/event/event-memory/src/test/java/org/apache/james/mailbox/events/InVMEventBusTest.java
----------------------------------------------------------------------
diff --git 
a/mailbox/event/event-memory/src/test/java/org/apache/james/mailbox/events/InVMEventBusTest.java
 
b/mailbox/event/event-memory/src/test/java/org/apache/james/mailbox/events/InVMEventBusTest.java
index 2ee725b..c7e7edd 100644
--- 
a/mailbox/event/event-memory/src/test/java/org/apache/james/mailbox/events/InVMEventBusTest.java
+++ 
b/mailbox/event/event-memory/src/test/java/org/apache/james/mailbox/events/InVMEventBusTest.java
@@ -22,6 +22,8 @@ package org.apache.james.mailbox.events;
 import org.apache.james.mailbox.events.delivery.InVmEventDelivery;
 import org.apache.james.metrics.api.NoopMetricFactory;
 import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Disabled;
+import org.junit.jupiter.api.Test;
 
 public class InVMEventBusTest implements 
KeyContract.SingleEventBusKeyContract, 
GroupContract.SingleEventBusGroupContract,
     ErrorHandlingContract{
@@ -39,4 +41,27 @@ public class InVMEventBusTest implements 
KeyContract.SingleEventBusKeyContract,
     public EventBus eventBus() {
         return eventBus;
     }
+
+    @Override
+    public EventDeadLetters deadLetter() {
+        throw new RuntimeException("this method is not a part of this task 
contents, will be handled in another pull request");
+    }
+
+    @Test
+    @Disabled("this test is not a part of this task contents, will be handled 
in another pull request")
+    @Override
+    public void deadLettersIsNotAppliedForKeyRegistrations() {
+    }
+
+    @Test
+    @Disabled("this test is not a part of this task contents, will be handled 
in another pull request")
+    @Override
+    public void deadLetterShouldNotStoreWhenFailsLessThanMaxRetries() {
+    }
+
+    @Test
+    @Disabled("this test is not a part of this task contents, will be handled 
in another pull request")
+    @Override
+    public void deadLetterShouldStoreWhenFailsGreaterThanMaxRetries() {
+    }
 }
\ No newline at end of file

http://git-wip-us.apache.org/repos/asf/james-project/blob/98cc6604/mailbox/event/event-rabbitmq/pom.xml
----------------------------------------------------------------------
diff --git a/mailbox/event/event-rabbitmq/pom.xml 
b/mailbox/event/event-rabbitmq/pom.xml
index 383172c..605fe6d 100644
--- a/mailbox/event/event-rabbitmq/pom.xml
+++ b/mailbox/event/event-rabbitmq/pom.xml
@@ -58,6 +58,11 @@
         </dependency>
         <dependency>
             <groupId>${james.groupId}</groupId>
+            <artifactId>apache-james-mailbox-event-memory</artifactId>
+            <scope>test</scope>
+        </dependency>
+        <dependency>
+            <groupId>${james.groupId}</groupId>
             <artifactId>james-server-testing</artifactId>
             <scope>test</scope>
         </dependency>

http://git-wip-us.apache.org/repos/asf/james-project/blob/98cc6604/mailbox/event/event-rabbitmq/src/main/java/org/apache/james/mailbox/events/GroupConsumerRetry.java
----------------------------------------------------------------------
diff --git 
a/mailbox/event/event-rabbitmq/src/main/java/org/apache/james/mailbox/events/GroupConsumerRetry.java
 
b/mailbox/event/event-rabbitmq/src/main/java/org/apache/james/mailbox/events/GroupConsumerRetry.java
index 04a513a..2b99c49 100644
--- 
a/mailbox/event/event-rabbitmq/src/main/java/org/apache/james/mailbox/events/GroupConsumerRetry.java
+++ 
b/mailbox/event/event-rabbitmq/src/main/java/org/apache/james/mailbox/events/GroupConsumerRetry.java
@@ -68,32 +68,44 @@ class GroupConsumerRetry {
         private final Sender sender;
         private final RetryExchangeName retryExchangeName;
         private final RetryBackoffConfiguration retryBackoff;
+        private final EventDeadLetters eventDeadLetters;
+        private final GroupRegistration.WorkQueueName queueName;
 
-        RetryPublisher(Sender sender, RetryExchangeName retryExchangeName, 
RetryBackoffConfiguration retryBackoff) {
+        RetryPublisher(Sender sender, RetryExchangeName retryExchangeName, 
RetryBackoffConfiguration retryBackoff,
+                       EventDeadLetters eventDeadLetters, 
GroupRegistration.WorkQueueName queueName) {
             this.sender = sender;
             this.retryExchangeName = retryExchangeName;
             this.retryBackoff = retryBackoff;
+            this.eventDeadLetters = eventDeadLetters;
+            this.queueName = queueName;
         }
 
         Mono<Void> publish(Event event, byte[] eventAsByte, int 
currentRetryCount) {
-            return sender.send(createRetryMessage(eventAsByte, 
currentRetryCount))
-                .doOnError(throwable -> LOGGER.error("Exception happens when 
publishing event of user {} to retry exchange," +
-                        "this event will be lost forever",
-                    event.getUser().asString(), throwable));
+            return Mono.just(currentRetryCount)
+                .flatMap(retryCount -> retryOrStoreToDeadLetter(event, 
eventAsByte, retryCount));
         }
 
-        private Mono<OutboundMessage> createRetryMessage(byte[] eventAsByte, 
int currentRetryCount) {
+        private Mono<Void> retryOrStoreToDeadLetter(Event event, byte[] 
eventAsByte, int currentRetryCount) {
             if (currentRetryCount >= retryBackoff.getMaxRetries()) {
-                return Mono.empty(); // will store event to deadletter latter
+                return eventDeadLetters.store(queueName.getGroup(), event);
             }
 
-            return Mono.just(new OutboundMessage(
+            return sendRetryMessage(event, eventAsByte, currentRetryCount);
+        }
+
+        private Mono<Void> sendRetryMessage(Event event, byte[] eventAsByte, 
int currentRetryCount) {
+            Mono<OutboundMessage> retryMessage = Mono.just(new OutboundMessage(
                 retryExchangeName.asString(),
                 EMPTY_ROUTING_KEY,
                 new AMQP.BasicProperties.Builder()
                     .headers(ImmutableMap.of(RETRY_COUNT, currentRetryCount + 
1))
                     .build(),
                 eventAsByte));
+
+            return sender.send(retryMessage)
+                .doOnError(throwable -> LOGGER.error("Exception happens when 
publishing event of user {} to retry exchange," +
+                    "this event will be lost forever",
+                    event.getUser().asString(), throwable));
         }
     }
 
@@ -105,11 +117,11 @@ class GroupConsumerRetry {
     private final RetryPublisher retryPublisher;
 
     GroupConsumerRetry(Sender sender, GroupRegistration.WorkQueueName 
queueName, Group group,
-                       RetryBackoffConfiguration retryBackoff) {
+                       RetryBackoffConfiguration retryBackoff, 
EventDeadLetters eventDeadLetters) {
         this.sender = sender;
         this.queueName = queueName;
         this.retryExchangeName = RetryExchangeName.of(group);
-        this.retryPublisher = new RetryPublisher(sender, retryExchangeName, 
retryBackoff);
+        this.retryPublisher = new RetryPublisher(sender, retryExchangeName, 
retryBackoff, eventDeadLetters, queueName);
     }
 
     Mono<Void> createRetryExchange() {

http://git-wip-us.apache.org/repos/asf/james-project/blob/98cc6604/mailbox/event/event-rabbitmq/src/main/java/org/apache/james/mailbox/events/GroupRegistration.java
----------------------------------------------------------------------
diff --git 
a/mailbox/event/event-rabbitmq/src/main/java/org/apache/james/mailbox/events/GroupRegistration.java
 
b/mailbox/event/event-rabbitmq/src/main/java/org/apache/james/mailbox/events/GroupRegistration.java
index 36a6d1a..7e4fe6d 100644
--- 
a/mailbox/event/event-rabbitmq/src/main/java/org/apache/james/mailbox/events/GroupRegistration.java
+++ 
b/mailbox/event/event-rabbitmq/src/main/java/org/apache/james/mailbox/events/GroupRegistration.java
@@ -101,6 +101,7 @@ class GroupRegistration implements Registration {
 
     GroupRegistration(Mono<Connection> connectionSupplier, Sender sender, 
EventSerializer eventSerializer,
                       MailboxListener mailboxListener, Group group, 
RetryBackoffConfiguration retryBackoff,
+                      EventDeadLetters eventDeadLetters,
                       Runnable unregisterGroup) {
         this.eventSerializer = eventSerializer;
         this.mailboxListener = mailboxListener;
@@ -109,7 +110,7 @@ class GroupRegistration implements Registration {
         this.receiver = RabbitFlux.createReceiver(new 
ReceiverOptions().connectionMono(connectionSupplier));
         this.receiverSubscriber = Optional.empty();
         this.unregisterGroup = unregisterGroup;
-        this.retryHandler = new GroupConsumerRetry(sender, queueName, group, 
retryBackoff);
+        this.retryHandler = new GroupConsumerRetry(sender, queueName, group, 
retryBackoff, eventDeadLetters);
         this.delayGenerator = WaitDelayGenerator.of(retryBackoff);
     }
 

http://git-wip-us.apache.org/repos/asf/james-project/blob/98cc6604/mailbox/event/event-rabbitmq/src/main/java/org/apache/james/mailbox/events/GroupRegistrationHandler.java
----------------------------------------------------------------------
diff --git 
a/mailbox/event/event-rabbitmq/src/main/java/org/apache/james/mailbox/events/GroupRegistrationHandler.java
 
b/mailbox/event/event-rabbitmq/src/main/java/org/apache/james/mailbox/events/GroupRegistrationHandler.java
index 33ba13b..9e3e062 100644
--- 
a/mailbox/event/event-rabbitmq/src/main/java/org/apache/james/mailbox/events/GroupRegistrationHandler.java
+++ 
b/mailbox/event/event-rabbitmq/src/main/java/org/apache/james/mailbox/events/GroupRegistrationHandler.java
@@ -36,12 +36,16 @@ class GroupRegistrationHandler {
     private final Sender sender;
     private final Mono<Connection> connectionMono;
     private final RetryBackoffConfiguration retryBackoff;
+    private final EventDeadLetters eventDeadLetters;
 
-    GroupRegistrationHandler(EventSerializer eventSerializer, Sender sender, 
Mono<Connection> connectionMono, RetryBackoffConfiguration retryBackoff) {
+    GroupRegistrationHandler(EventSerializer eventSerializer, Sender sender, 
Mono<Connection> connectionMono,
+                             RetryBackoffConfiguration retryBackoff,
+                             EventDeadLetters eventDeadLetters) {
         this.eventSerializer = eventSerializer;
         this.sender = sender;
         this.connectionMono = connectionMono;
         this.retryBackoff = retryBackoff;
+        this.eventDeadLetters = eventDeadLetters;
         this.groupRegistrations = new ConcurrentHashMap<>();
     }
 
@@ -68,6 +72,7 @@ class GroupRegistrationHandler {
             listener,
             group,
             retryBackoff,
+            eventDeadLetters,
             () -> groupRegistrations.remove(group));
     }
 }
\ No newline at end of file

http://git-wip-us.apache.org/repos/asf/james-project/blob/98cc6604/mailbox/event/event-rabbitmq/src/main/java/org/apache/james/mailbox/events/RabbitMQEventBus.java
----------------------------------------------------------------------
diff --git 
a/mailbox/event/event-rabbitmq/src/main/java/org/apache/james/mailbox/events/RabbitMQEventBus.java
 
b/mailbox/event/event-rabbitmq/src/main/java/org/apache/james/mailbox/events/RabbitMQEventBus.java
index b2aaa8d..0ab082c 100644
--- 
a/mailbox/event/event-rabbitmq/src/main/java/org/apache/james/mailbox/events/RabbitMQEventBus.java
+++ 
b/mailbox/event/event-rabbitmq/src/main/java/org/apache/james/mailbox/events/RabbitMQEventBus.java
@@ -47,6 +47,7 @@ class RabbitMQEventBus implements EventBus {
     private final RoutingKeyConverter routingKeyConverter;
     private final RetryBackoffConfiguration retryBackoff;
     private final EventBusId eventBusId;
+    private final EventDeadLetters eventDeadLetters;
 
     private MailboxListenerRegistry mailboxListenerRegistry;
     private GroupRegistrationHandler groupRegistrationHandler;
@@ -56,21 +57,23 @@ class RabbitMQEventBus implements EventBus {
 
     RabbitMQEventBus(RabbitMQConnectionFactory rabbitMQConnectionFactory, 
EventSerializer eventSerializer,
                      RetryBackoffConfiguration retryBackoff,
-                     RoutingKeyConverter routingKeyConverter) {
-        this.eventBusId = EventBusId.random();
+                     RoutingKeyConverter routingKeyConverter,
+                     EventDeadLetters eventDeadLetters) {
+            this.eventBusId = EventBusId.random();
         this.connectionMono = 
Mono.fromSupplier(rabbitMQConnectionFactory::create).cache();
         this.eventSerializer = eventSerializer;
         this.routingKeyConverter = routingKeyConverter;
         this.retryBackoff = retryBackoff;
+        this.eventDeadLetters = eventDeadLetters;
         this.isRunning = new AtomicBoolean(false);
     }
 
     public void start() {
         if (!isRunning.get()) {
             sender = RabbitFlux.createSender(new 
SenderOptions().connectionMono(connectionMono));
-            groupRegistrationHandler = new 
GroupRegistrationHandler(eventSerializer, sender, connectionMono, retryBackoff);
             mailboxListenerRegistry = new MailboxListenerRegistry();
             keyRegistrationHandler = new KeyRegistrationHandler(eventBusId, 
eventSerializer, sender, connectionMono, routingKeyConverter, 
mailboxListenerRegistry);
+            groupRegistrationHandler = new 
GroupRegistrationHandler(eventSerializer, sender, connectionMono, retryBackoff, 
eventDeadLetters);
             eventDispatcher = new EventDispatcher(eventBusId, eventSerializer, 
sender, mailboxListenerRegistry);
 
             eventDispatcher.start();

http://git-wip-us.apache.org/repos/asf/james-project/blob/98cc6604/mailbox/event/event-rabbitmq/src/test/java/org/apache/james/mailbox/events/RabbitMQEventBusTest.java
----------------------------------------------------------------------
diff --git 
a/mailbox/event/event-rabbitmq/src/test/java/org/apache/james/mailbox/events/RabbitMQEventBusTest.java
 
b/mailbox/event/event-rabbitmq/src/test/java/org/apache/james/mailbox/events/RabbitMQEventBusTest.java
index 910b83e..0d55f30 100644
--- 
a/mailbox/event/event-rabbitmq/src/test/java/org/apache/james/mailbox/events/RabbitMQEventBusTest.java
+++ 
b/mailbox/event/event-rabbitmq/src/test/java/org/apache/james/mailbox/events/RabbitMQEventBusTest.java
@@ -90,10 +90,12 @@ class RabbitMQEventBusTest implements 
GroupContract.SingleEventBusGroupContract,
     private RabbitMQConnectionFactory connectionFactory;
     private EventSerializer eventSerializer;
     private RoutingKeyConverter routingKeyConverter;
+    private MemoryEventDeadLetters memoryEventDeadLetters;
 
     @BeforeEach
     void setUp() {
         connectionFactory = rabbitMQExtension.getConnectionFactory();
+        memoryEventDeadLetters = new MemoryEventDeadLetters();
         Mono<Connection> connectionMono = 
Mono.fromSupplier(connectionFactory::create).cache();
 
         TestId.Factory mailboxIdFactory = new TestId.Factory();
@@ -123,7 +125,7 @@ class RabbitMQEventBusTest implements 
GroupContract.SingleEventBusGroupContract,
     }
 
     private RabbitMQEventBus newEventBus() {
-        return new RabbitMQEventBus(connectionFactory, eventSerializer, 
RetryBackoffConfiguration.DEFAULT, routingKeyConverter);
+        return new RabbitMQEventBus(connectionFactory, eventSerializer, 
RetryBackoffConfiguration.DEFAULT, routingKeyConverter, memoryEventDeadLetters);
     }
 
     @Override
@@ -137,10 +139,8 @@ class RabbitMQEventBusTest implements 
GroupContract.SingleEventBusGroupContract,
     }
 
     @Override
-    @Test
-    @Disabled("This test is failing by RabbitMQEventBus exponential backoff is 
not implemented at this time")
-    public void failingRegisteredListenersShouldNotAbortRegisteredDelivery() {
-
+    public EventDeadLetters deadLetter() {
+        return memoryEventDeadLetters;
     }
 
     @Override


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

Reply via email to