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]

Reply via email to