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]
