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]
