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 dab5a9fd3c705d6f65ebb67308f338704dea9c88
Author: Quan Tran <[email protected]>
AuthorDate: Wed Dec 24 11:49:31 2025 +0700

    [IMPROVEMENT] MailboxMergingTask: better reactive
---
 .../mail/task/MailboxMergingTaskRunner.java        | 60 ++++++++++++++--------
 .../james/mailbox/store/StoreMessageIdManager.java | 30 +++++------
 2 files changed, 52 insertions(+), 38 deletions(-)

diff --git 
a/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/task/MailboxMergingTaskRunner.java
 
b/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/task/MailboxMergingTaskRunner.java
index d184b4672e..20d8e28054 100644
--- 
a/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/task/MailboxMergingTaskRunner.java
+++ 
b/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/task/MailboxMergingTaskRunner.java
@@ -19,6 +19,8 @@
 
 package org.apache.james.mailbox.cassandra.mail.task;
 
+import java.util.function.Function;
+
 import jakarta.inject.Inject;
 
 import org.apache.james.core.Username;
@@ -67,35 +69,49 @@ public class MailboxMergingTaskRunner {
 
     public Task.Result run(CassandraId oldMailboxId, CassandraId newMailboxId, 
MailboxMergingTask.Context context) {
         return moveMessages(oldMailboxId, newMailboxId, mailboxSession, 
context)
-            .onComplete(
-                () -> mergeRights(oldMailboxId, newMailboxId).block(),
-                () -> mailboxDAO.delete(oldMailboxId).block());
+            .flatMap(onMoveCompleteOperations(oldMailboxId, newMailboxId))
+            .block();
     }
 
-    private Task.Result moveMessages(CassandraId oldMailboxId, CassandraId 
newMailboxId, MailboxSession session, MailboxMergingTask.Context context) {
+    private Function<Task.Result, Mono<Task.Result>> 
onMoveCompleteOperations(CassandraId oldMailboxId, CassandraId newMailboxId) {
+        return result -> {
+            if (result == Task.Result.COMPLETED) {
+                return mergeRights(oldMailboxId, newMailboxId)
+                    .then(mailboxDAO.delete(oldMailboxId))
+                    .thenReturn(result)
+                    .onErrorResume(e -> {
+                        LOGGER.error("Error while executing move completion 
operation", e);
+                        return Mono.just(Task.Result.PARTIAL);
+                    });
+            }
+            return Mono.just(result);
+        };
+    }
+
+    private Mono<Task.Result> moveMessages(CassandraId oldMailboxId, 
CassandraId newMailboxId, MailboxSession session, MailboxMergingTask.Context 
context) {
         return cassandraMessageIdDAO.retrieveMessages(oldMailboxId, 
MessageRange.all(), Limit.unlimited())
             .map(CassandraMessageMetadata::getComposedMessageId)
             .map(ComposedMessageIdWithMetaData::getComposedMessageId)
-            .flatMap(messageId -> Mono.fromCallable(() -> 
moveMessage(newMailboxId, messageId, session, context))
-                .subscribeOn(ReactorUtils.BLOCKING_CALL_WRAPPER), 
ReactorUtils.DEFAULT_CONCURRENCY)
-            .reduce(Task.Result.COMPLETED, Task::combine)
-            .block();
+            .flatMap(messageId -> moveMessage(newMailboxId, messageId, 
session, context), ReactorUtils.DEFAULT_CONCURRENCY)
+            .reduce(Task.Result.COMPLETED, Task::combine);
     }
 
-    private Task.Result moveMessage(CassandraId newMailboxId, 
ComposedMessageId composedMessageId, MailboxSession session, 
MailboxMergingTask.Context context) {
-        try {
-            
messageIdManager.setInMailboxesNoCheck(composedMessageId.getMessageId(), 
newMailboxId, session);
-            context.incrementMovedCount();
-            return Task.Result.COMPLETED;
-        } catch (OverQuotaException e) {
-            LOGGER.warn("Failed moving message {} due to quota error", 
composedMessageId.getMessageId(), e);
-            context.incrementFailedCount();
-            return Task.Result.PARTIAL;
-        } catch (MailboxException e) {
-            LOGGER.warn("Failed moving message {}", 
composedMessageId.getMessageId(), e);
-            context.incrementFailedCount();
-            return Task.Result.PARTIAL;
-        }
+    private Mono<Task.Result> moveMessage(CassandraId newMailboxId, 
ComposedMessageId composedMessageId, MailboxSession session, 
MailboxMergingTask.Context context) {
+        return 
messageIdManager.setInMailboxesNoCheck(composedMessageId.getMessageId(), 
newMailboxId, session)
+            .then(Mono.fromCallable(() -> {
+                context.incrementMovedCount();
+                return Task.Result.COMPLETED;
+            }))
+            .onErrorResume(OverQuotaException.class, e -> {
+                LOGGER.warn("Failed moving message {} due to quota error", 
composedMessageId.getMessageId(), e);
+                context.incrementFailedCount();
+                return Mono.just(Task.Result.PARTIAL);
+            })
+            .onErrorResume(MailboxException.class, e -> {
+                LOGGER.warn("Failed moving message {}", 
composedMessageId.getMessageId(), e);
+                context.incrementFailedCount();
+                return Mono.just(Task.Result.PARTIAL);
+            });
     }
 
     private Mono<Void> mergeRights(CassandraId oldMailboxId, CassandraId 
newMailboxId) {
diff --git 
a/mailbox/store/src/main/java/org/apache/james/mailbox/store/StoreMessageIdManager.java
 
b/mailbox/store/src/main/java/org/apache/james/mailbox/store/StoreMessageIdManager.java
index 3c5d4ef76e..ac30c59119 100644
--- 
a/mailbox/store/src/main/java/org/apache/james/mailbox/store/StoreMessageIdManager.java
+++ 
b/mailbox/store/src/main/java/org/apache/james/mailbox/store/StoreMessageIdManager.java
@@ -363,22 +363,20 @@ public class StoreMessageIdManager implements 
MessageIdManager {
             .flatMap(eventBus::dispatch);
     }
 
-    public void setInMailboxesNoCheck(MessageId messageId, MailboxId 
targetMailboxId, MailboxSession mailboxSession) throws MailboxException {
-        MessageIdMapper messageIdMapper = 
mailboxSessionMapperFactory.getMessageIdMapper(mailboxSession);
-        List<MailboxMessage> currentMailboxMessages = 
messageIdMapper.find(ImmutableList.of(messageId), 
MessageMapper.FetchType.METADATA);
-
-        
MailboxReactorUtils.block(messageMovesWithMailbox(MessageMoves.builder()
-            .targetMailboxIds(targetMailboxId)
-            .previousMailboxIds(toMailboxIds(currentMailboxMessages))
-            .build(), mailboxSession)
-            .flatMapMany(messageMove -> {
-                if (messageMove.isChange()) {
-                    return applyMessageMoveNoMailboxChecks(mailboxSession, 
currentMailboxMessages, messageMove);
-                }
-                return Flux.empty();
-            })
-            .collectList()
-            .flatMap(eventBus::dispatch));
+    public Mono<Void> setInMailboxesNoCheck(MessageId messageId, MailboxId 
targetMailboxId, MailboxSession mailboxSession) {
+        return findRelatedMailboxMessages(messageId, mailboxSession)
+            .flatMap(currentMailboxMessages -> 
messageMovesWithMailbox(MessageMoves.builder()
+                .targetMailboxIds(targetMailboxId)
+                .previousMailboxIds(toMailboxIds(currentMailboxMessages))
+                .build(), mailboxSession)
+                .flatMapMany(messageMove -> {
+                    if (messageMove.isChange()) {
+                        return applyMessageMoveNoMailboxChecks(mailboxSession, 
currentMailboxMessages, messageMove);
+                    }
+                    return Flux.empty();
+                })
+                .collectList()
+                .flatMap(eventBus::dispatch));
     }
 
     private Mono<List<MailboxMessage>> findRelatedMailboxMessages(MessageId 
messageId, MailboxSession mailboxSession) {


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

Reply via email to