This is an automated email from the ASF dual-hosted git repository. chibenwa pushed a commit to branch master in repository https://gitbox.apache.org/repos/asf/james-project.git
commit a1a3581b37eb0b59ba86eeeb9511920d228894e0 Author: Benoit TELLIER <[email protected]> AuthorDate: Wed Sep 30 09:18:59 2026 +0200 [FIX] Re-enable SolveMessageInconsistenciesServiceTest The test was commented out by JAMES-3576 (message denormalization) and never restored. Adapt it to the current DAO APIs and query strings. Doing so reveals that outdated messageIdTable entries were not fully fixed: flags were diffed against an empty set, so flags only present in the stale record (eg \Seen) were never removed. --- .../task/SolveMessageInconsistenciesService.java | 4 +- .../SolveMessageInconsistenciesServiceTest.java | 149 ++++++++++++++------- 2 files changed, 101 insertions(+), 52 deletions(-) diff --git a/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/task/SolveMessageInconsistenciesService.java b/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/task/SolveMessageInconsistenciesService.java index c3713aadb2..611b2771a3 100644 --- a/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/task/SolveMessageInconsistenciesService.java +++ b/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/task/SolveMessageInconsistenciesService.java @@ -32,7 +32,6 @@ import java.util.concurrent.atomic.AtomicLong; import java.util.function.Predicate; import jakarta.inject.Inject; -import jakarta.mail.Flags; import org.apache.james.backends.cassandra.init.configuration.CassandraConfiguration; import org.apache.james.backends.cassandra.init.configuration.JamesExecutionProfiles.ConsistencyChoice; @@ -125,7 +124,8 @@ public class SolveMessageInconsistenciesService { ComposedMessageIdWithMetaData id = messageFromImapUid.getComposedMessageId(); return messageIdDAO.updateMetadata(id.getComposedMessageId(), UpdatedFlags.builder() - .oldFlags(new Flags()) + // The update only writes flag changes: diff against the stale record so that extra flags get removed + .oldFlags(messageFromMessageId.getComposedMessageId().getFlags()) .newFlags(id.getFlags()) .modSeq(id.getModSeq()) .messageId(id.getComposedMessageId().getMessageId()) diff --git a/mailbox/cassandra/src/test/java/org/apache/james/mailbox/cassandra/mail/task/SolveMessageInconsistenciesServiceTest.java b/mailbox/cassandra/src/test/java/org/apache/james/mailbox/cassandra/mail/task/SolveMessageInconsistenciesServiceTest.java index decabe8663..36072ad236 100644 --- a/mailbox/cassandra/src/test/java/org/apache/james/mailbox/cassandra/mail/task/SolveMessageInconsistenciesServiceTest.java +++ b/mailbox/cassandra/src/test/java/org/apache/james/mailbox/cassandra/mail/task/SolveMessageInconsistenciesServiceTest.java @@ -19,7 +19,47 @@ package org.apache.james.mailbox.cassandra.mail.task; -/* +import static org.apache.james.backends.cassandra.Scenario.Builder.awaitOn; +import static org.apache.james.backends.cassandra.Scenario.Builder.fail; +import static org.assertj.core.api.Assertions.assertThat; + +import java.util.Date; +import java.util.Optional; + +import jakarta.mail.Flags; + +import org.apache.james.backends.cassandra.CassandraCluster; +import org.apache.james.backends.cassandra.CassandraClusterExtension; +import org.apache.james.backends.cassandra.Scenario; +import org.apache.james.backends.cassandra.components.CassandraDataDefinition; +import org.apache.james.backends.cassandra.init.configuration.CassandraConfiguration; +import org.apache.james.backends.cassandra.versions.CassandraSchemaVersionDataDefinition; +import org.apache.james.blob.api.PlainBlobId; +import org.apache.james.junit.categories.Unstable; +import org.apache.james.mailbox.MessageUid; +import org.apache.james.mailbox.ModSeq; +import org.apache.james.mailbox.cassandra.ids.CassandraId; +import org.apache.james.mailbox.cassandra.ids.CassandraMessageId; +import org.apache.james.mailbox.cassandra.mail.CassandraMessageIdDAO; +import org.apache.james.mailbox.cassandra.mail.CassandraMessageIdToImapUidDAO; +import org.apache.james.mailbox.cassandra.mail.CassandraMessageMetadata; +import org.apache.james.mailbox.cassandra.mail.task.SolveMessageInconsistenciesService.Context; +import org.apache.james.mailbox.cassandra.mail.task.SolveMessageInconsistenciesService.RunningOptions; +import org.apache.james.mailbox.cassandra.modules.CassandraMessageDataDefinition; +import org.apache.james.mailbox.model.ComposedMessageId; +import org.apache.james.mailbox.model.ComposedMessageIdWithMetaData; +import org.apache.james.mailbox.model.ThreadId; +import org.apache.james.task.Task; +import org.assertj.core.api.SoftAssertions; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Nested; +import org.junit.jupiter.api.Tag; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.RegisterExtension; + +import reactor.core.publisher.Mono; +import reactor.core.scheduler.Schedulers; + public class SolveMessageInconsistenciesServiceTest { private static final CassandraId MAILBOX_ID = CassandraId.timeBased(); @@ -29,40 +69,54 @@ public class SolveMessageInconsistenciesServiceTest { private static final MessageUid MESSAGE_UID_2 = MessageUid.of(2L); private static final ModSeq MOD_SEQ_1 = ModSeq.of(1L); private static final ModSeq MOD_SEQ_2 = ModSeq.of(2L); + private static final PlainBlobId HEADER_BLOB_ID = new PlainBlobId.Factory().of("header"); + private static final Date INTERNAL_DATE = new Date(1586000000000L); + private static final long SIZE = 36L; + private static final int BODY_START_OCTET = 18; - private static final ComposedMessageIdWithMetaData MESSAGE_1 = ComposedMessageIdWithMetaData.builder() + private static final CassandraMessageMetadata MESSAGE_1 = metadata(ComposedMessageIdWithMetaData.builder() .composedMessageId(new ComposedMessageId(MAILBOX_ID, MESSAGE_ID_1, MESSAGE_UID_1)) .modSeq(MOD_SEQ_1) .flags(new Flags()) .threadId(ThreadId.fromBaseMessageId(MESSAGE_ID_1)) - .build(); + .build()); - private static final ComposedMessageIdWithMetaData MESSAGE_1_WITH_SEEN_FLAG = ComposedMessageIdWithMetaData.builder() + private static final CassandraMessageMetadata MESSAGE_1_WITH_SEEN_FLAG = metadata(ComposedMessageIdWithMetaData.builder() .composedMessageId(new ComposedMessageId(MAILBOX_ID, MESSAGE_ID_1, MESSAGE_UID_1)) .modSeq(MOD_SEQ_1) .flags(new Flags(Flags.Flag.SEEN)) .threadId(ThreadId.fromBaseMessageId(MESSAGE_ID_1)) - .build(); + .build()); - private static final ComposedMessageIdWithMetaData MESSAGE_1_WITH_MOD_SEQ_2 = ComposedMessageIdWithMetaData.builder() + private static final CassandraMessageMetadata MESSAGE_1_WITH_MOD_SEQ_2 = metadata(ComposedMessageIdWithMetaData.builder() .composedMessageId(new ComposedMessageId(MAILBOX_ID, MESSAGE_ID_1, MESSAGE_UID_1)) .modSeq(MOD_SEQ_2) .flags(new Flags(Flags.Flag.SEEN)) .threadId(ThreadId.fromBaseMessageId(MESSAGE_ID_1)) - .build(); + .build()); - private static final ComposedMessageIdWithMetaData MESSAGE_2 = ComposedMessageIdWithMetaData.builder() + private static final CassandraMessageMetadata MESSAGE_2 = metadata(ComposedMessageIdWithMetaData.builder() .composedMessageId(new ComposedMessageId(MAILBOX_ID, MESSAGE_ID_2, MESSAGE_UID_2)) .modSeq(MOD_SEQ_2) .flags(new Flags()) .threadId(ThreadId.fromBaseMessageId(MESSAGE_ID_2)) - .build(); + .build()); + + private static CassandraMessageMetadata metadata(ComposedMessageIdWithMetaData ids) { + return CassandraMessageMetadata.builder() + .ids(ids) + .internalDate(INTERNAL_DATE) + .bodyStartOctet(BODY_START_OCTET) + .size(SIZE) + .headerContent(Optional.of(HEADER_BLOB_ID)) + .build(); + } @RegisterExtension static CassandraClusterExtension cassandraCluster = new CassandraClusterExtension( - CassandraModule.aggregateModules( - CassandraSchemaVersionModule.MODULE, - CassandraMessageModule.MODULE)); + CassandraDataDefinition.aggregateModules( + CassandraSchemaVersionDataDefinition.MODULE, + CassandraMessageDataDefinition.MODULE)); CassandraMessageIdToImapUidDAO imapUidDAO; CassandraMessageIdDAO messageIdDAO; @@ -70,12 +124,9 @@ public class SolveMessageInconsistenciesServiceTest { @BeforeEach void setUp(CassandraCluster cassandra) { - imapUidDAO = new CassandraMessageIdToImapUidDAO( - cassandra.getConf(), - cassandraCluster.getCassandraConsistenciesConfiguration(), - new CassandraMessageId.Factory(), - CassandraConfiguration.DEFAULT_CONFIGURATION); - messageIdDAO = new CassandraMessageIdDAO(cassandra.getConf(), blobIdFactory, new CassandraMessageId.Factory()); + PlainBlobId.Factory blobIdFactory = new PlainBlobId.Factory(); + imapUidDAO = new CassandraMessageIdToImapUidDAO(cassandra.getConf(), blobIdFactory, CassandraConfiguration.DEFAULT_CONFIGURATION); + messageIdDAO = new CassandraMessageIdDAO(cassandra.getConf(), blobIdFactory); testee = new SolveMessageInconsistenciesService(imapUidDAO, messageIdDAO, CassandraConfiguration.DEFAULT_CONFIGURATION); } @@ -140,10 +191,10 @@ public class SolveMessageInconsistenciesServiceTest { .registerScenario(awaitOn(barrier) .thenExecuteNormally() .times(1) - .whenQueryStartsWith("SELECT messageId,mailboxId,uid,threadId,modSeq,flagAnswered,flagDeleted,flagDraft,flagFlagged,flagRecent,flagSeen,flagUser,userFlags FROM messageIdTable WHERE mailboxId=:mailboxId AND uid=:uid;")); + .whenQueryStartsWith("SELECT * FROM messageidtable WHERE mailboxid=:mailboxid AND uid=:uid")); Context context = new Context(); - Mono<Task.Result> task = testee.fixMessageInconsistencies(context, RunningOptions.DEFAULT).subscribeOn(Schedulers.elastic()).cache(); + Mono<Task.Result> task = testee.fixMessageInconsistencies(context, RunningOptions.DEFAULT).subscribeOn(Schedulers.boundedElastic()).cache(); task.subscribe(); barrier.awaitCaller(); @@ -169,10 +220,10 @@ public class SolveMessageInconsistenciesServiceTest { .registerScenario(awaitOn(barrier) .thenExecuteNormally() .times(1) - .whenQueryStartsWith("SELECT messageId,mailboxId,uid,threadId,modSeq,flagAnswered,flagDeleted,flagDraft,flagFlagged,flagRecent,flagSeen,flagUser,userFlags FROM messageIdTable WHERE mailboxId=:mailboxId AND uid=:uid;")); + .whenQueryStartsWith("SELECT * FROM messageidtable WHERE mailboxid=:mailboxid AND uid=:uid")); Context context = new Context(); - Mono<Task.Result> task = testee.fixMessageInconsistencies(context, RunningOptions.DEFAULT).subscribeOn(Schedulers.elastic()).cache(); + Mono<Task.Result> task = testee.fixMessageInconsistencies(context, RunningOptions.DEFAULT).subscribeOn(Schedulers.boundedElastic()).cache(); task.subscribe(); barrier.awaitCaller(); @@ -260,7 +311,7 @@ public class SolveMessageInconsistenciesServiceTest { cassandra.getConf() .registerScenario(fail() .forever() - .whenQueryStartsWith("INSERT INTO messageIdTable (mailboxId,uid,threadId,modSeq,messageId,flagAnswered,flagDeleted,flagDraft,flagFlagged,flagRecent,flagSeen,flagUser,userFlags) VALUES (:mailboxId,:uid,:threadId,:modSeq,:messageId,:flagAnswered,:flagDeleted,:flagDraft,:flagFlagged,:flagRecent,:flagSeen,:flagUser,:userFlags)")); + .whenQueryStartsWith("UPDATE messageidtable SET threadid=:threadid, messageid=:messageid")); assertThat(testee.fixMessageInconsistencies(new Context(), RunningOptions.DEFAULT).block()) .isEqualTo(Task.Result.PARTIAL); @@ -274,7 +325,7 @@ public class SolveMessageInconsistenciesServiceTest { cassandra.getConf() .registerScenario(fail() .times(1) - .whenQueryStartsWith("INSERT INTO messageIdTable (mailboxId,uid,threadId,modSeq,messageId,flagAnswered,flagDeleted,flagDraft,flagFlagged,flagRecent,flagSeen,flagUser,userFlags) VALUES (:mailboxId,:uid,:threadId,:modSeq,:messageId,:flagAnswered,:flagDeleted,:flagDraft,:flagFlagged,:flagRecent,:flagSeen,:flagUser,:userFlags)")); + .whenQueryStartsWith("UPDATE messageidtable SET threadid=:threadid, messageid=:messageid")); assertThat(testee.fixMessageInconsistencies(new Context(), RunningOptions.DEFAULT).block()) .isEqualTo(Task.Result.PARTIAL); @@ -288,7 +339,7 @@ public class SolveMessageInconsistenciesServiceTest { cassandra.getConf() .registerScenario(fail() .times(1) - .whenQueryStartsWith("INSERT INTO messageIdTable (mailboxId,uid,modSeq,messageId,flagAnswered,flagDeleted,flagDraft,flagFlagged,flagRecent,flagSeen,flagUser,userFlags) VALUES (:mailboxId,:uid,:modSeq,:messageId,:flagAnswered,:flagDeleted,:flagDraft,:flagFlagged,:flagRecent,:flagSeen,:flagUser,:userFlags)")); + .whenQueryStartsWith("UPDATE messageidtable SET threadid=:threadid, messageid=:messageid")); testee.fixMessageInconsistencies(new Context(), new RunningOptions(1)).block(); @@ -309,14 +360,14 @@ public class SolveMessageInconsistenciesServiceTest { cassandra.getConf() .registerScenario(fail() .times(1) - .whenQueryStartsWith("SELECT messageId,mailboxId,uid,threadId,modSeq,flagAnswered,flagDeleted,flagDraft,flagFlagged,flagRecent,flagSeen,flagUser,userFlags FROM messageIdTable WHERE mailboxId=:mailboxId AND uid=:uid;")); + .whenQueryStartsWith("SELECT * FROM messageidtable WHERE mailboxid=:mailboxid AND uid=:uid")); testee.fixMessageInconsistencies(context, RunningOptions.DEFAULT).block(); assertThat(context.snapshot()) .isEqualTo(Context.Snapshot.builder() .processedImapUidEntries(1) - .errors(MESSAGE_1.getComposedMessageId()) + .errors(MESSAGE_1.getComposedMessageId().getComposedMessageId()) .build()); } @@ -329,14 +380,14 @@ public class SolveMessageInconsistenciesServiceTest { cassandra.getConf() .registerScenario(fail() .times(1) - .whenQueryStartsWith("SELECT messageId,mailboxId,uid,threadId,modSeq,flagAnswered,flagDeleted,flagDraft,flagFlagged,flagRecent,flagSeen,flagUser,userFlags FROM messageIdTable WHERE mailboxId=:mailboxId AND uid=:uid")); + .whenQueryStartsWith("SELECT * FROM messageidtable WHERE mailboxid=:mailboxid AND uid=:uid")); testee.fixMessageInconsistencies(context, RunningOptions.DEFAULT).block(); assertThat(context.snapshot()) .isEqualTo(Context.Snapshot.builder() .processedImapUidEntries(1) - .errors(MESSAGE_1.getComposedMessageId()) + .errors(MESSAGE_1.getComposedMessageId().getComposedMessageId()) .build()); } } @@ -362,12 +413,10 @@ public class SolveMessageInconsistenciesServiceTest { .registerScenario(awaitOn(barrier) .thenExecuteNormally() .times(1) - .whenQueryStartsWith("SELECT messageId,mailboxId,uid,threadId,modSeq,flagAnswered,flagDeleted," + - "flagDraft,flagFlagged,flagRecent,flagSeen,flagUser,userFlags FROM messageIdTable " + - "WHERE mailboxId=:mailboxId AND uid=:uid;")); + .whenQueryStartsWith("SELECT * FROM messageidtable WHERE mailboxid=:mailboxid AND uid=:uid")); Context context = new Context(); - Mono<Task.Result> task = testee.fixMessageInconsistencies(context, RunningOptions.DEFAULT).subscribeOn(Schedulers.elastic()).cache(); + Mono<Task.Result> task = testee.fixMessageInconsistencies(context, RunningOptions.DEFAULT).subscribeOn(Schedulers.boundedElastic()).cache(); task.subscribe(); barrier.awaitCaller(); @@ -435,7 +484,7 @@ public class SolveMessageInconsistenciesServiceTest { cassandra.getConf() .registerScenario(fail() .forever() - .whenQueryStartsWith("DELETE FROM messageIdTable WHERE mailboxId=:mailboxId AND uid=:uid")); + .whenQueryStartsWith("DELETE FROM messageidtable WHERE mailboxid=:mailboxid AND uid=:uid")); assertThat(testee.fixMessageInconsistencies(new Context(), RunningOptions.DEFAULT).block()) .isEqualTo(Task.Result.PARTIAL); @@ -449,7 +498,7 @@ public class SolveMessageInconsistenciesServiceTest { cassandra.getConf() .registerScenario(fail() .times(1) - .whenQueryStartsWith("DELETE FROM messageIdTable WHERE mailboxId=:mailboxId AND uid=:uid")); + .whenQueryStartsWith("DELETE FROM messageidtable WHERE mailboxid=:mailboxid AND uid=:uid")); assertThat(testee.fixMessageInconsistencies(new Context(), RunningOptions.DEFAULT).block()) .isEqualTo(Task.Result.PARTIAL); @@ -463,7 +512,7 @@ public class SolveMessageInconsistenciesServiceTest { cassandra.getConf() .registerScenario(fail() .times(1) - .whenQueryStartsWith("DELETE FROM messageIdTable WHERE mailboxId=:mailboxId AND uid=:uid;")); + .whenQueryStartsWith("DELETE FROM messageidtable WHERE mailboxid=:mailboxid AND uid=:uid")); testee.fixMessageInconsistencies(new Context(), new RunningOptions(1)).block(); @@ -484,14 +533,14 @@ public class SolveMessageInconsistenciesServiceTest { cassandra.getConf() .registerScenario(fail() .times(1) - .whenQueryStartsWith("SELECT messageId,mailboxId,uid,threadId,modSeq,flagAnswered,flagDeleted,flagDraft,flagFlagged,flagRecent,flagSeen,flagUser,userFlags FROM messageIdTable WHERE mailboxId=:mailboxId AND uid=:uid")); + .whenQueryStartsWith("SELECT * FROM messageidtable WHERE mailboxid=:mailboxid AND uid=:uid")); testee.fixMessageInconsistencies(context, RunningOptions.DEFAULT).block(); assertThat(context.snapshot()) .isEqualTo(Context.Snapshot.builder() .processedMessageIdEntries(1) - .errors(MESSAGE_1.getComposedMessageId()) + .errors(MESSAGE_1.getComposedMessageId().getComposedMessageId()) .build()); } @@ -504,14 +553,14 @@ public class SolveMessageInconsistenciesServiceTest { cassandra.getConf() .registerScenario(fail() .times(1) - .whenQueryStartsWith("SELECT messageId,mailboxId,uid,threadId,modSeq,flagAnswered,flagDeleted,flagDraft,flagFlagged,flagRecent,flagSeen,flagUser,userFlags FROM imapUidTable WHERE messageid=:messageid AND mailboxid=:mailboxid")); + .whenQueryStartsWith("SELECT * FROM imapuidtable WHERE messageid=:messageid AND mailboxid=:mailboxid")); testee.fixMessageInconsistencies(context, RunningOptions.DEFAULT).block(); assertThat(context.snapshot()) .isEqualTo(Context.Snapshot.builder() .processedMessageIdEntries(1) - .errors(MESSAGE_1.getComposedMessageId()) + .errors(MESSAGE_1.getComposedMessageId().getComposedMessageId()) .build()); } } @@ -523,7 +572,7 @@ public class SolveMessageInconsistenciesServiceTest { testee.fixMessageInconsistencies(context, RunningOptions.DEFAULT).block(); - assertThat(context.snapshot()).isEqualToComparingFieldByFieldRecursively(new Context().snapshot()); + assertThat(context.snapshot()).isEqualTo(new Context().snapshot()); } @Test @@ -554,7 +603,7 @@ public class SolveMessageInconsistenciesServiceTest { .isEqualTo(Context.Snapshot.builder() .processedImapUidEntries(1) .addedMessageIdEntries(1) - .addFixedInconsistencies(MESSAGE_1.getComposedMessageId()) + .addFixedInconsistencies(MESSAGE_1.getComposedMessageId().getComposedMessageId()) .build()); } @@ -572,7 +621,7 @@ public class SolveMessageInconsistenciesServiceTest { .processedImapUidEntries(1) .processedMessageIdEntries(1) .updatedMessageIdEntries(1) - .addFixedInconsistencies(MESSAGE_1.getComposedMessageId()) + .addFixedInconsistencies(MESSAGE_1.getComposedMessageId().getComposedMessageId()) .build()); } @@ -590,7 +639,7 @@ public class SolveMessageInconsistenciesServiceTest { .processedImapUidEntries(1) .processedMessageIdEntries(1) .updatedMessageIdEntries(1) - .addFixedInconsistencies(MESSAGE_1.getComposedMessageId()) + .addFixedInconsistencies(MESSAGE_1.getComposedMessageId().getComposedMessageId()) .build()); } @@ -606,9 +655,9 @@ public class SolveMessageInconsistenciesServiceTest { .isEqualTo(Context.Snapshot.builder() .processedMessageIdEntries(1) .removedMessageIdEntries(1) - .addFixedInconsistencies(MESSAGE_1.getComposedMessageId()) + .addFixedInconsistencies(MESSAGE_1.getComposedMessageId().getComposedMessageId()) .build()); - }*/ + } /* Error @@ -618,7 +667,7 @@ public class SolveMessageInconsistenciesServiceTest { Caused by: com.datastax.driver.core.exceptions.WriteTimeoutException: Cassandra timeout during SIMPLE write query at consistency QUORUM (1 replica were required but only 0 acknowledged the write) https://builds.apache.org/blue/organizations/jenkins/james%2FApacheJames/detail/PR-268/39/tests */ - /* @Test + @Test @Tag(Unstable.TAG) void fixMailboxInconsistenciesShouldUpdateContextWhenDeleteError(CassandraCluster cassandra) { Context context = new Context(); @@ -628,14 +677,14 @@ public class SolveMessageInconsistenciesServiceTest { cassandra.getConf() .registerScenario(fail() .times(1) - .whenQueryStartsWith("DELETE FROM messageIdTable WHERE mailboxId=:mailboxId AND uid=:uid;")); + .whenQueryStartsWith("DELETE FROM messageidtable WHERE mailboxid=:mailboxid AND uid=:uid")); testee.fixMessageInconsistencies(context, RunningOptions.DEFAULT).block(); assertThat(context.snapshot()) .isEqualTo(Context.Snapshot.builder() .processedMessageIdEntries(1) - .errors(MESSAGE_1.getComposedMessageId()) + .errors(MESSAGE_1.getComposedMessageId().getComposedMessageId()) .build()); } -}*/ +} --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
