This is an automated email from the ASF dual-hosted git repository. quantranhong1999 pushed a commit to branch 3.9.x in repository https://gitbox.apache.org/repos/asf/james-project.git
commit 5158569c9a6dc6f77825a9d80638da8008fad944 Author: Benoit TELLIER <[email protected]> AuthorDate: Wed Sep 30 09:19:26 2026 +0200 [ENHANCEMENT] SolveMessageInconsistencies: do not propagate imapUidTable entries without messageV3 An orphan imapUidTable entry was blindly re-inserted into messageIdTable. When the message content (messageV3) is missing, typically because deleted data got resurrected after tombstones were purged without repair, this exposes a message that can not be read. Check messageV3 (READ profile, QUORUM by default) before propagating. When absent, skip the entry, log it and report it as an error. --- .../operate/webadmin/admin-mailboxes-extend.adoc | 8 +++ .../operate/webadmin/admin-messages-extend.adoc | 8 +++ .../task/SolveMessageInconsistenciesService.java | 42 +++++++++++++- .../SolveMessageInconsistenciesServiceTest.java | 66 +++++++++++++++++++++- 4 files changed, 121 insertions(+), 3 deletions(-) diff --git a/docs/modules/servers/pages/distributed/operate/webadmin/admin-mailboxes-extend.adoc b/docs/modules/servers/pages/distributed/operate/webadmin/admin-mailboxes-extend.adoc index a1dabd7290..df7194bd54 100644 --- a/docs/modules/servers/pages/distributed/operate/webadmin/admin-mailboxes-extend.adoc +++ b/docs/modules/servers/pages/distributed/operate/webadmin/admin-mailboxes-extend.adoc @@ -167,6 +167,14 @@ The message is SEEN via JMAP The message is UNSEEN via IMAP .... +Entries of `imapUidTable` missing in `messageIdTable` are propagated to +`messageIdTable` only if the message content (`messageV3` table) still +exists. Otherwise the entry is left untouched, logged and reported in +the `errors` of the task, as exposing it would lead to a message that +can not be read. Such entries typically result from deleted data being +resurrected, for instance when Cassandra repairs are not run within +`gc_grace_seconds`. + link:#_endpoints_returning_a_task[More details about endpoints returning a task]. diff --git a/docs/modules/servers/pages/distributed/operate/webadmin/admin-messages-extend.adoc b/docs/modules/servers/pages/distributed/operate/webadmin/admin-messages-extend.adoc index 1f77c27658..0d5df6519c 100644 --- a/docs/modules/servers/pages/distributed/operate/webadmin/admin-messages-extend.adoc +++ b/docs/modules/servers/pages/distributed/operate/webadmin/admin-messages-extend.adoc @@ -28,6 +28,14 @@ The message is SEEN via JMAP The message is UNSEEN via IMAP .... +Entries of `imapUidTable` missing in `messageIdTable` are propagated to +`messageIdTable` only if the message content (`messageV3` table) still +exists. Otherwise the entry is left untouched, logged and reported in +the `errors` of the task, as exposing it would lead to a message that +can not be read. Such entries typically result from deleted data being +resurrected, for instance when Cassandra repairs are not run within +`gc_grace_seconds`. + link:#_endpoints_returning_a_task[More details about endpoints returning a task]. 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 d4c98bfe7e..0b2fd7cfe2 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 @@ -38,12 +38,14 @@ import org.apache.james.backends.cassandra.init.configuration.JamesExecutionProf import org.apache.james.mailbox.MessageUid; import org.apache.james.mailbox.cassandra.ids.CassandraId; import org.apache.james.mailbox.cassandra.ids.CassandraMessageId; +import org.apache.james.mailbox.cassandra.mail.CassandraMessageDAOV3; 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.model.ComposedMessageId; import org.apache.james.mailbox.model.ComposedMessageIdWithMetaData; import org.apache.james.mailbox.model.UpdatedFlags; +import org.apache.james.mailbox.store.mail.MessageMapper.FetchType; import org.apache.james.task.Task; import org.apache.james.util.ReactorUtils; import org.slf4j.Logger; @@ -110,6 +112,28 @@ public class SolveMessageInconsistenciesService { } } + /** + * The message is referenced in ImapUid but its content is missing in MessageV3. + * + * This happens when deleted index entries are resurrected (eg: tombstones purged before being repaired) + * while the message content is not. The entry is not propagated to MessageId: doing so would expose + * a message that can not be read. It is reported instead. + */ + private static class ImapUidEntryWithoutContent implements Inconsistency { + private final CassandraMessageMetadata message; + + private ImapUidEntryWithoutContent(CassandraMessageMetadata message) { + this.message = message; + } + + @Override + public Mono<Task.Result> fix(Context context, CassandraMessageIdToImapUidDAO imapUidDAO, CassandraMessageIdDAO messageIdDAO) { + context.addErrors(message.getComposedMessageId().getComposedMessageId()); + LOGGER.warn("Skipping orphan message in ImapUid as its content is missing in MessageV3: {}", message.getComposedMessageId()); + return Mono.just(Task.Result.PARTIAL); + } + } + private static class OutdatedMessageIdEntry implements Inconsistency { private final CassandraMessageMetadata messageFromMessageId; private final CassandraMessageMetadata messageFromImapUid; @@ -424,13 +448,15 @@ public class SolveMessageInconsistenciesService { private final CassandraMessageIdToImapUidDAO messageIdToImapUidDAO; private final CassandraMessageIdDAO messageIdDAO; + private final CassandraMessageDAOV3 messageDAOV3; private final CassandraConfiguration cassandraConfiguration; @Inject SolveMessageInconsistenciesService(CassandraMessageIdToImapUidDAO messageIdToImapUidDAO, CassandraMessageIdDAO messageIdDAO, - CassandraConfiguration cassandraConfiguration) { + CassandraMessageDAOV3 messageDAOV3, CassandraConfiguration cassandraConfiguration) { this.messageIdToImapUidDAO = messageIdToImapUidDAO; this.messageIdDAO = messageIdDAO; + this.messageDAOV3 = messageDAOV3; this.cassandraConfiguration = cassandraConfiguration; } @@ -492,10 +518,22 @@ public class SolveMessageInconsistenciesService { private Mono<Inconsistency> detectOrphanImapUidEntry(CassandraId mailboxId, CassandraMessageId messageId) { return messageIdToImapUidDAO.retrieve(messageId, Optional.of(mailboxId), chooseReadConsistency()) .next() - .<Inconsistency>map(OrphanImapUidEntry::new) + .flatMap(orphanEntry -> hasContent(messageId) + .map(hasContent -> { + if (hasContent) { + return new OrphanImapUidEntry(orphanEntry); + } + return new ImapUidEntryWithoutContent(orphanEntry); + })) .switchIfEmpty(Mono.just(NO_INCONSISTENCY)); } + // Upon optimistic consistency, an empty read is retried with the READ execution profile (QUORUM by default) + private Mono<Boolean> hasContent(CassandraMessageId messageId) { + return messageDAOV3.retrieveMessage(messageId, FetchType.METADATA) + .hasElement(); + } + private Flux<Task.Result> fixInconsistenciesInMessageId(Context context, RunningOptions runningOptions) { return messageIdDAO.retrieveAllMessages() .transform(ReactorUtils.<CassandraMessageMetadata, Task.Result>throttle() 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 36072ad236..2311f55d98 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 @@ -22,6 +22,7 @@ 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 static org.mockito.Mockito.mock; import java.util.Date; import java.util.Optional; @@ -34,18 +35,24 @@ 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.BlobStore; +import org.apache.james.blob.api.BlobStoreCacheCallback; +import org.apache.james.blob.api.BlobStoreDAO; 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.CassandraMessageDAOV3; 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.MessageRepresentation; 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.ByteContent; import org.apache.james.mailbox.model.ComposedMessageId; import org.apache.james.mailbox.model.ComposedMessageIdWithMetaData; import org.apache.james.mailbox.model.ThreadId; @@ -57,6 +64,8 @@ import org.junit.jupiter.api.Tag; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.RegisterExtension; +import com.google.common.collect.ImmutableList; + import reactor.core.publisher.Mono; import reactor.core.scheduler.Schedulers; @@ -65,11 +74,14 @@ public class SolveMessageInconsistenciesServiceTest { private static final CassandraId MAILBOX_ID = CassandraId.timeBased(); private static final CassandraMessageId MESSAGE_ID_1 = new CassandraMessageId.Factory().fromString("d2bee791-7e63-11ea-883c-95b84008f979"); private static final CassandraMessageId MESSAGE_ID_2 = new CassandraMessageId.Factory().fromString("eeeeeeee-7e63-11ea-883c-95b84008f979"); + private static final CassandraMessageId MESSAGE_ID_3 = new CassandraMessageId.Factory().fromString("ffffffff-7e63-11ea-883c-95b84008f979"); private static final MessageUid MESSAGE_UID_1 = MessageUid.of(1L); private static final MessageUid MESSAGE_UID_2 = MessageUid.of(2L); + private static final MessageUid MESSAGE_UID_3 = MessageUid.of(3L); 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 PlainBlobId BODY_BLOB_ID = new PlainBlobId.Factory().of("body"); private static final Date INTERNAL_DATE = new Date(1586000000000L); private static final long SIZE = 36L; private static final int BODY_START_OCTET = 18; @@ -102,6 +114,14 @@ public class SolveMessageInconsistenciesServiceTest { .threadId(ThreadId.fromBaseMessageId(MESSAGE_ID_2)) .build()); + // No content is stored for this message within MessageV3 + private static final CassandraMessageMetadata MESSAGE_3 = metadata(ComposedMessageIdWithMetaData.builder() + .composedMessageId(new ComposedMessageId(MAILBOX_ID, MESSAGE_ID_3, MESSAGE_UID_3)) + .modSeq(MOD_SEQ_1) + .flags(new Flags()) + .threadId(ThreadId.fromBaseMessageId(MESSAGE_ID_3)) + .build()); + private static CassandraMessageMetadata metadata(ComposedMessageIdWithMetaData ids) { return CassandraMessageMetadata.builder() .ids(ids) @@ -120,6 +140,7 @@ public class SolveMessageInconsistenciesServiceTest { CassandraMessageIdToImapUidDAO imapUidDAO; CassandraMessageIdDAO messageIdDAO; + CassandraMessageDAOV3 messageDAOV3; SolveMessageInconsistenciesService testee; @BeforeEach @@ -127,7 +148,19 @@ public class SolveMessageInconsistenciesServiceTest { 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); + // Only MessageV3 metadata is read: blobs are never accessed + messageDAOV3 = new CassandraMessageDAOV3(cassandra.getConf(), cassandra.getTypesProvider(), mock(BlobStore.class), + mock(BlobStoreDAO.class), blobIdFactory, CassandraConfiguration.DEFAULT_CONFIGURATION, BlobStoreCacheCallback.NOOP); + testee = new SolveMessageInconsistenciesService(imapUidDAO, messageIdDAO, messageDAOV3, CassandraConfiguration.DEFAULT_CONFIGURATION); + + saveContent(MESSAGE_ID_1); + saveContent(MESSAGE_ID_2); + } + + private void saveContent(CassandraMessageId messageId) { + messageDAOV3.save(new MessageRepresentation(messageId, INTERNAL_DATE, SIZE, BODY_START_OCTET, + new ByteContent(new byte[0]), ImmutableList.of(), HEADER_BLOB_ID, BODY_BLOB_ID)) + .block(); } @Test @@ -607,6 +640,37 @@ public class SolveMessageInconsistenciesServiceTest { .build()); } + @Test + void orphanImapUidEntryWithoutContentShouldNotBePropagated() { + imapUidDAO.insert(MESSAGE_3).block(); + + testee.fixMessageInconsistencies(new Context(), RunningOptions.DEFAULT).block(); + + SoftAssertions.assertSoftly(softly -> { + softly.assertThat(imapUidDAO.retrieve(MESSAGE_ID_3, Optional.of(MAILBOX_ID)).collectList().block()) + .containsExactly(MESSAGE_3); + softly.assertThat(messageIdDAO.retrieve(MAILBOX_ID, MESSAGE_UID_3).block()) + .isEmpty(); + }); + } + + @Test + void orphanImapUidEntryWithoutContentShouldBeReportedAsError() { + Context context = new Context(); + imapUidDAO.insert(MESSAGE_3).block(); + + Task.Result result = testee.fixMessageInconsistencies(context, RunningOptions.DEFAULT).block(); + + SoftAssertions.assertSoftly(softly -> { + softly.assertThat(result).isEqualTo(Task.Result.PARTIAL); + softly.assertThat(context.snapshot()) + .isEqualTo(Context.Snapshot.builder() + .processedImapUidEntries(1) + .errors(MESSAGE_3.getComposedMessageId().getComposedMessageId()) + .build()); + }); + } + @Test void fixMailboxInconsistenciesShouldUpdateContextWhenInconsistentModSeq() { Context context = new Context(); --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
