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

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

commit 439aa21cf6340ec347488eb4cbc7901f33f67547
Author: RĂ©mi KOWALSKI <[email protected]>
AuthorDate: Tue Mar 3 14:15:33 2020 +0100

    [Refactoring] do ack in mono when consuming workqueue
---
 .../task/eventsourcing/distributed/RabbitMQWorkQueue.java    | 12 ++++++------
 1 file changed, 6 insertions(+), 6 deletions(-)

diff --git 
a/server/task/task-distributed/src/main/java/org/apache/james/task/eventsourcing/distributed/RabbitMQWorkQueue.java
 
b/server/task/task-distributed/src/main/java/org/apache/james/task/eventsourcing/distributed/RabbitMQWorkQueue.java
index 005b86d..93208fe 100644
--- 
a/server/task/task-distributed/src/main/java/org/apache/james/task/eventsourcing/distributed/RabbitMQWorkQueue.java
+++ 
b/server/task/task-distributed/src/main/java/org/apache/james/task/eventsourcing/distributed/RabbitMQWorkQueue.java
@@ -129,13 +129,13 @@ public class RabbitMQWorkQueue implements WorkQueue {
     }
 
     private Mono<Task.Result> executeTask(AcknowledgableDelivery delivery) {
-        delivery.ack();
-        String json = new String(delivery.getBody(), StandardCharsets.UTF_8);
-
         TaskId taskId = 
TaskId.fromString(delivery.getProperties().getHeaders().get(TASK_ID).toString());
-
-        return deserialize(json, taskId)
-            .flatMap(task -> executeOnWorker(taskId, task));
+        return Mono.fromCallable(() -> {
+            delivery.ack();
+            return new String(delivery.getBody(), StandardCharsets.UTF_8);
+        }).flatMap(json ->
+            deserialize(json, taskId)
+                .flatMap(task -> executeOnWorker(taskId, task)));
     }
 
     private Mono<Task> deserialize(String json, TaskId taskId) {


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

Reply via email to