MAILBOX-376 QueueName is only required when initializing group retry
Project: http://git-wip-us.apache.org/repos/asf/james-project/repo Commit: http://git-wip-us.apache.org/repos/asf/james-project/commit/301329fd Tree: http://git-wip-us.apache.org/repos/asf/james-project/tree/301329fd Diff: http://git-wip-us.apache.org/repos/asf/james-project/diff/301329fd Branch: refs/heads/master Commit: 301329fd36d3ce3c73f4cd10c3c5678de1620a0b Parents: 0c64eec Author: Benoit Tellier <[email protected]> Authored: Wed Jan 23 10:42:43 2019 +0700 Committer: Benoit Tellier <[email protected]> Committed: Wed Jan 23 10:44:33 2019 +0700 ---------------------------------------------------------------------- .../apache/james/mailbox/events/GroupConsumerRetry.java | 11 +++++------ .../apache/james/mailbox/events/GroupRegistration.java | 8 ++------ 2 files changed, 7 insertions(+), 12 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/james-project/blob/301329fd/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 ff10d23..e789465 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 @@ -64,21 +64,20 @@ class GroupConsumerRetry { private static final Logger LOGGER = LoggerFactory.getLogger(GroupConsumerRetry.class); private final Sender sender; - private final GroupRegistration.WorkQueueName queueName; private final RetryExchangeName retryExchangeName; private final RetryBackoffConfiguration retryBackoff; private final EventDeadLetters eventDeadLetters; + private final Group group; - GroupConsumerRetry(Sender sender, GroupRegistration.WorkQueueName queueName, Group group, - RetryBackoffConfiguration retryBackoff, EventDeadLetters eventDeadLetters) { + GroupConsumerRetry(Sender sender, Group group, RetryBackoffConfiguration retryBackoff, EventDeadLetters eventDeadLetters) { this.sender = sender; - this.queueName = queueName; this.retryExchangeName = RetryExchangeName.of(group); this.retryBackoff = retryBackoff; this.eventDeadLetters = eventDeadLetters; + this.group = group; } - Mono<Void> createRetryExchange() { + Mono<Void> createRetryExchange(GroupRegistration.WorkQueueName queueName) { return Flux.concat( sender.declareExchange(ExchangeSpecification.exchange(retryExchangeName.asString()) .durable(DURABLE) @@ -99,7 +98,7 @@ class GroupConsumerRetry { private Mono<Void> retryOrStoreToDeadLetter(Event event, byte[] eventAsByte, int currentRetryCount) { if (currentRetryCount >= retryBackoff.getMaxRetries()) { - return eventDeadLetters.store(queueName.getGroup(), event); + return eventDeadLetters.store(group, event); } return sendRetryMessage(event, eventAsByte, currentRetryCount); } http://git-wip-us.apache.org/repos/asf/james-project/blob/301329fd/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 d7c4c91..2b44ab8 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 @@ -68,10 +68,6 @@ class GroupRegistration implements Registration { this.group = group; } - Group getGroup() { - return group; - } - String asString() { return MAILBOX_EVENT_WORK_QUEUE_PREFIX + group.asString(); } @@ -101,13 +97,13 @@ 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, eventDeadLetters); + this.retryHandler = new GroupConsumerRetry(sender, group, retryBackoff, eventDeadLetters); this.delayGenerator = WaitDelayGenerator.of(retryBackoff); } GroupRegistration start() { createGroupWorkQueue() - .then(retryHandler.createRetryExchange()) + .then(retryHandler.createRetryExchange(queueName)) .doOnSuccess(any -> this.subscribeWorkQueue()) .block(); return this; --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
