This is an automated email from the ASF dual-hosted git repository.
chibenwa pushed a commit to branch 3.9.x
in repository https://gitbox.apache.org/repos/asf/james-project.git
The following commit(s) were added to refs/heads/3.9.x by this push:
new 201aca3e79 [FIX] Plug MailboxMerging error stream in parent task
(#3220)
201aca3e79 is described below
commit 201aca3e79978c0dad767c0c083f35bc55d55f40
Author: Benoit TELLIER <[email protected]>
AuthorDate: Wed Sep 30 15:12:11 2026 +0200
[FIX] Plug MailboxMerging error stream in parent task (#3220)
---
.../task/SolveMailboxInconsistenciesService.java | 97 ++++++++++++++++------
.../SolveMailboxInconsistenciesServiceTest.java | 72 +++++++++++++++-
2 files changed, 143 insertions(+), 26 deletions(-)
diff --git
a/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/task/SolveMailboxInconsistenciesService.java
b/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/task/SolveMailboxInconsistenciesService.java
index 654f631469..a21a5b564c 100644
---
a/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/task/SolveMailboxInconsistenciesService.java
+++
b/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/task/SolveMailboxInconsistenciesService.java
@@ -235,25 +235,9 @@ public class SolveMailboxInconsistenciesService {
// for the loser *nowhere*. If the loser owns another path, it is a
genuine mailbox with its
// own (same-id) conflict to resolve first, so we report rather than
destroy it.
private Mono<Result> autoMergeConflict(Context context,
CassandraMailboxDAO mailboxDAO, CassandraMailboxPathV3DAO pathV3DAO) {
- CassandraId winnerId = (CassandraId)
mailboxPathEntry.getMailboxId();
- CassandraId loserId = (CassandraId) mailboxDaoEntry.getMailboxId();
- MailboxPath conflictingPath =
mailboxPathEntry.generateAssociatedPath();
-
- Mono<Boolean> pathStillOwnedByWinner =
pathV3DAO.retrieve(conflictingPath, STRONG)
- .map(entry -> entry.getMailboxId().equals(winnerId))
- .defaultIfEmpty(false);
- Mono<Boolean> loserStillResolvesToPath =
mailboxDAO.retrieveMailbox(loserId)
- .map(projection ->
projection.generateAssociatedPath().equals(conflictingPath))
- .defaultIfEmpty(false);
- Mono<Boolean> loserIsUnregistered =
pathV3DAO.listUserMailboxes(mailboxDaoEntry.getNamespace(),
mailboxDaoEntry.getUser(), STRONG)
- .filter(entry -> entry.getMailboxId().equals(loserId))
- .hasElements()
- .map(referenced -> !referenced);
-
- return Mono.zip(pathStillOwnedByWinner, loserStillResolvesToPath,
loserIsUnregistered)
- .flatMap(state -> {
- boolean cleanGhost = state.getT1() && state.getT2() &&
state.getT3();
- if (!cleanGhost) {
+ return validateMergeNeeded(mailboxDAO, pathV3DAO)
+ .flatMap(mergeNeeded -> {
+ if (!mergeNeeded) {
// State no longer matches the clean-ghost picture:
either already reconciled,
// or the loser owns another path. In the latter case
the loser is a genuine
// mailbox whose stale projection gets realigned onto
its registered path by the
@@ -261,15 +245,78 @@ public class SolveMailboxInconsistenciesService {
// We therefore leave it to that resolution rather
than destroying or reporting it.
return Mono.just(Result.COMPLETED);
}
- return mergingRunner.runReactive(loserId, winnerId, new
MailboxMergingTask.Context(0))
- .doOnNext(result -> {
- LOGGER.info("Auto-merged ghost mailbox {} into {}
at path {}",
- loserId.serialize(), winnerId.serialize(),
conflictingPath.asString());
- context.addFixedInconsistency(winnerId);
- });
+ return mergeBothMailboxes(context);
+ })
+ .onErrorResume(e -> {
+ LOGGER.error("Failed auto-merging ghost mailbox {} into {}
at path {}",
+ loserId().serialize(), winnerId().serialize(),
conflictingPath().asString(), e);
+ context.incrementErrors();
+ return Mono.just(Result.PARTIAL);
});
}
+ private Mono<Boolean> validateMergeNeeded(CassandraMailboxDAO
mailboxDAO, CassandraMailboxPathV3DAO pathV3DAO) {
+ return Mono.zip(pathStillOwnedByWinner(pathV3DAO),
loserStillResolvesToPath(mailboxDAO), loserIsUnregistered(pathV3DAO))
+ .map(state -> state.getT1() && state.getT2() && state.getT3());
+ }
+
+ private Mono<Boolean> pathStillOwnedByWinner(CassandraMailboxPathV3DAO
pathV3DAO) {
+ return pathV3DAO.retrieve(conflictingPath(), STRONG)
+ .map(entry -> entry.getMailboxId().equals(winnerId()))
+ .defaultIfEmpty(false);
+ }
+
+ private Mono<Boolean> loserStillResolvesToPath(CassandraMailboxDAO
mailboxDAO) {
+ return mailboxDAO.retrieveMailbox(loserId())
+ .map(projection ->
projection.generateAssociatedPath().equals(conflictingPath()))
+ .defaultIfEmpty(false);
+ }
+
+ private Mono<Boolean> loserIsUnregistered(CassandraMailboxPathV3DAO
pathV3DAO) {
+ return pathV3DAO.listUserMailboxes(mailboxDaoEntry.getNamespace(),
mailboxDaoEntry.getUser(), STRONG)
+ .filter(entry -> entry.getMailboxId().equals(loserId()))
+ .hasElements()
+ .map(referenced -> !referenced);
+ }
+
+ private Mono<Result> mergeBothMailboxes(Context context) {
+ MailboxMergingTask.Context mergingContext = new
MailboxMergingTask.Context(0);
+ return mergingRunner.runReactive(loserId(), winnerId(),
mergingContext)
+ .doOnNext(result -> {
+ if (result == Result.COMPLETED) {
+ notifyMergeSuccess(context);
+ } else {
+ notifyMergeFailure(context, mergingContext);
+ }
+ });
+ }
+
+ private void notifyMergeSuccess(Context context) {
+ LOGGER.info("Auto-merged ghost mailbox {} into {} at path {}",
+ loserId().serialize(), winnerId().serialize(),
conflictingPath().asString());
+ context.addFixedInconsistency(winnerId());
+ }
+
+ // The ghost is kept when the merge is partial: this is not a fix.
+ private void notifyMergeFailure(Context context,
MailboxMergingTask.Context mergingContext) {
+ LOGGER.error("Failed auto-merging ghost mailbox {} into {} at path
{}: {} message(s) moved, {} message(s) failed",
+ loserId().serialize(), winnerId().serialize(),
conflictingPath().asString(),
+ mergingContext.getMessageMovedCount(),
mergingContext.getMessageFailedCount());
+ context.incrementErrors();
+ }
+
+ private CassandraId winnerId() {
+ return (CassandraId) mailboxPathEntry.getMailboxId();
+ }
+
+ private CassandraId loserId() {
+ return (CassandraId) mailboxDaoEntry.getMailboxId();
+ }
+
+ private MailboxPath conflictingPath() {
+ return mailboxPathEntry.generateAssociatedPath();
+ }
+
// Auto-resolution is restricted to the case where BOTH the
conflicting path entry and the
// projection's path are registered to the same mailbox id: only then
are they two genuine
// aliases of a single mailbox, hence safe to deduplicate without data
loss. If the
diff --git
a/mailbox/cassandra/src/test/java/org/apache/james/mailbox/cassandra/mail/task/SolveMailboxInconsistenciesServiceTest.java
b/mailbox/cassandra/src/test/java/org/apache/james/mailbox/cassandra/mail/task/SolveMailboxInconsistenciesServiceTest.java
index a2cd31802e..644d3c92bd 100644
---
a/mailbox/cassandra/src/test/java/org/apache/james/mailbox/cassandra/mail/task/SolveMailboxInconsistenciesServiceTest.java
+++
b/mailbox/cassandra/src/test/java/org/apache/james/mailbox/cassandra/mail/task/SolveMailboxInconsistenciesServiceTest.java
@@ -20,6 +20,7 @@
package org.apache.james.mailbox.cassandra.mail.task;
import static
org.apache.james.JsonSerializationVerifier.recursiveComparisonConfiguration;
+import static org.apache.james.backends.cassandra.Scenario.Builder.fail;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatCode;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
@@ -38,21 +39,29 @@ import
org.apache.james.backends.cassandra.versions.CassandraSchemaVersionDataDe
import
org.apache.james.backends.cassandra.versions.CassandraSchemaVersionManager;
import org.apache.james.backends.cassandra.versions.SchemaVersion;
import org.apache.james.core.Username;
+import org.apache.james.mailbox.MailboxManager;
import org.apache.james.mailbox.cassandra.ids.CassandraId;
+import org.apache.james.mailbox.cassandra.mail.ACLMapper;
import org.apache.james.mailbox.cassandra.mail.CassandraMailboxDAO;
import org.apache.james.mailbox.cassandra.mail.CassandraMailboxPathV3DAO;
+import org.apache.james.mailbox.cassandra.mail.CassandraMessageIdDAO;
import
org.apache.james.mailbox.cassandra.mail.task.SolveMailboxInconsistenciesService.Context;
import org.apache.james.mailbox.cassandra.modules.CassandraAclDataDefinition;
import
org.apache.james.mailbox.cassandra.modules.CassandraMailboxDataDefinition;
import org.apache.james.mailbox.model.Mailbox;
+import org.apache.james.mailbox.model.MailboxACL;
import org.apache.james.mailbox.model.MailboxPath;
import org.apache.james.mailbox.model.UidValidity;
+import org.apache.james.mailbox.store.StoreMessageIdManager;
import org.apache.james.task.Task.Result;
import org.assertj.core.api.SoftAssertions;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.RegisterExtension;
+import reactor.core.publisher.Flux;
+import reactor.core.publisher.Mono;
+
class SolveMailboxInconsistenciesServiceTest {
private static final UidValidity UID_VALIDITY_1 = UidValidity.of(145);
private static final UidValidity UID_VALIDITY_2 = UidValidity.of(147);
@@ -491,4 +500,65 @@ class SolveMailboxInconsistenciesServiceTest {
softly.assertThat(context.snapshot().getConflictingEntries()).isEmpty();
});
}
-}
\ No newline at end of file
+
+ @Test
+ void autoMergeShouldReportAnErrorWhenMergeIsPartial(CassandraCluster
cassandra) {
+ // Ghost CASSANDRA_ID_1 squats path "abc" owned by CASSANDRA_ID_2.
+ Mailbox pathOwnerProjection = new Mailbox(MAILBOX_PATH,
UID_VALIDITY_2, CASSANDRA_ID_2);
+ mailboxDAO.save(MAILBOX).block();
+ mailboxDAO.save(pathOwnerProjection).block();
+ mailboxPathV3DAO.save(MAILBOX_2).block();
+
+ // Real merging runner: failing to drop the ghost projection makes the
merge partial.
+ testee = new SolveMailboxInconsistenciesService(mailboxDAO,
mailboxPathV3DAO,
+ new CassandraSchemaVersionManager(versionDAO),
realMergingRunner());
+ cassandra.getConf().registerScenario(fail()
+ .forever()
+ .whenQueryStartsWith("DELETE FROM mailbox WHERE"));
+
+ Context context = new Context();
+ Result result = testee.fixMailboxInconsistencies(context, new
SolveMailboxInconsistenciesService.RunningOptions(5, true)).block();
+
+ SoftAssertions.assertSoftly(softly -> {
+ softly.assertThat(result).isEqualTo(Result.PARTIAL);
+ softly.assertThat(context.snapshot().getErrors()).isEqualTo(1);
+
softly.assertThat(context.snapshot().getFixedInconsistencies()).isEmpty();
+
softly.assertThat(mailboxDAO.retrieveAllMailboxes().collectList().block())
+ .containsExactlyInAnyOrder(MAILBOX, pathOwnerProjection);
+ });
+ }
+
+ @Test
+ void autoMergeShouldReportAnErrorWhenCleanGhostCheckFails(CassandraCluster
cassandra) {
+ // Ghost CASSANDRA_ID_1 squats path "abc" owned by CASSANDRA_ID_2.
+ Mailbox pathOwnerProjection = new Mailbox(MAILBOX_PATH,
UID_VALIDITY_2, CASSANDRA_ID_2);
+ mailboxDAO.save(MAILBOX).block();
+ mailboxDAO.save(pathOwnerProjection).block();
+ mailboxPathV3DAO.save(MAILBOX_2).block();
+
+ // The first read by id is the one checking the ghost still resolves
to the conflicting path.
+ cassandra.getConf().registerScenario(fail()
+ .times(1)
+ .whenQueryStartsWith("SELECT id,mailboxbase,uidvalidity,name FROM
mailbox WHERE"));
+
+ Context context = new Context();
+ Result result = testee.fixMailboxInconsistencies(context, new
SolveMailboxInconsistenciesService.RunningOptions(1, true)).block();
+
+ SoftAssertions.assertSoftly(softly -> {
+ verify(mergingRunner, never()).runReactive(any(), any(), any());
+ softly.assertThat(result).isEqualTo(Result.PARTIAL);
+ softly.assertThat(context.snapshot().getErrors()).isEqualTo(1);
+
softly.assertThat(context.snapshot().getFixedInconsistencies()).isEmpty();
+ });
+ }
+
+ private MailboxMergingTaskRunner realMergingRunner() {
+ CassandraMessageIdDAO messageIdDAO = mock(CassandraMessageIdDAO.class);
+ when(messageIdDAO.retrieveMessages(any(), any(),
any())).thenReturn(Flux.empty());
+ ACLMapper aclMapper = mock(ACLMapper.class);
+ when(aclMapper.getACL(any())).thenReturn(Mono.just(MailboxACL.EMPTY));
+ when(aclMapper.setACL(any(), any())).thenReturn(Mono.empty());
+ return new MailboxMergingTaskRunner(mock(MailboxManager.class),
mock(StoreMessageIdManager.class),
+ messageIdDAO, mailboxDAO, aclMapper);
+ }
+}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]