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]
