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 849b1514c0 [FIX] CassandraMailRepository should be more resilient to 
extra large mail repository (#3115) (#3121)
849b1514c0 is described below

commit 849b1514c074b5404c4837df55c15fc21e01e659
Author: Benoit TELLIER <[email protected]>
AuthorDate: Wed Aug 19 07:54:35 2026 +0700

    [FIX] CassandraMailRepository should be more resilient to extra large mail 
repository (#3115) (#3121)
---
 .../CassandraEventDeadLettersDataDefinition.java   |  6 +-
 .../cassandra/CassandraMailRepository.java         | 37 +++++++++---
 .../CassandraMailRepositoryDataDefinition.java     |  6 ++
 .../cassandra/CassandraMailRepositoryTest.java     | 67 ++++++++++++++++++++++
 4 files changed, 107 insertions(+), 9 deletions(-)

diff --git 
a/event-bus/cassandra/src/main/java/org/apache/james/events/CassandraEventDeadLettersDataDefinition.java
 
b/event-bus/cassandra/src/main/java/org/apache/james/events/CassandraEventDeadLettersDataDefinition.java
index 13efd02367..6edc611038 100644
--- 
a/event-bus/cassandra/src/main/java/org/apache/james/events/CassandraEventDeadLettersDataDefinition.java
+++ 
b/event-bus/cassandra/src/main/java/org/apache/james/events/CassandraEventDeadLettersDataDefinition.java
@@ -33,7 +33,11 @@ public interface CassandraEventDeadLettersDataDefinition {
         .table(CassandraEventDeadLettersTable.TABLE_NAME)
         .comment("Holds event dead letter")
         .options(options -> options
-            .withCaching(true, 
rows(CassandraConstants.DEFAULT_CACHED_ROW_PER_PARTITION)))
+            .withCaching(true, 
rows(CassandraConstants.DEFAULT_CACHED_ROW_PER_PARTITION))
+            // All the dead lettered events of a group are stored within a 
single partition, thus deletions
+            // (upon redelivery) quickly pile up tombstones there, up to the 
point reads start failing.
+            // A resurrected event would only lead to an extra redelivery, 
which is harmless.
+            .withGcGraceSeconds(0))
         .statement(statement -> types -> statement
             .withPartitionKey(CassandraEventDeadLettersTable.GROUP, 
DataTypes.TEXT)
             .withClusteringColumn(CassandraEventDeadLettersTable.INSERTION_ID, 
DataTypes.UUID)
diff --git 
a/server/mailrepository/mailrepository-cassandra/src/main/java/org/apache/james/mailrepository/cassandra/CassandraMailRepository.java
 
b/server/mailrepository/mailrepository-cassandra/src/main/java/org/apache/james/mailrepository/cassandra/CassandraMailRepository.java
index 15dbe63d64..6f21eab31d 100644
--- 
a/server/mailrepository/mailrepository-cassandra/src/main/java/org/apache/james/mailrepository/cassandra/CassandraMailRepository.java
+++ 
b/server/mailrepository/mailrepository-cassandra/src/main/java/org/apache/james/mailrepository/cassandra/CassandraMailRepository.java
@@ -44,6 +44,8 @@ import org.apache.james.util.AuditTrail;
 import org.apache.james.util.streams.Iterators;
 import org.apache.mailet.Mail;
 import org.reactivestreams.Publisher;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
 
 import com.github.fge.lambdas.Throwing;
 import com.google.common.collect.ImmutableMap;
@@ -51,6 +53,8 @@ import com.google.common.collect.ImmutableMap;
 import reactor.core.publisher.Mono;
 
 public class CassandraMailRepository implements MailRepository {
+    private static final Logger LOGGER = 
LoggerFactory.getLogger(CassandraMailRepository.class);
+
     private final MailRepositoryUrl url;
     private final CassandraMailRepositoryKeysDAO keysDAO;
     private final CassandraMailRepositoryMailDaoV2 mailDAO;
@@ -102,28 +106,45 @@ public class CassandraMailRepository implements 
MailRepository {
     @Override
     public Iterator<MailKey> list(Condition condition) {
         return Iterators.toStream(list())
-            .filter(key -> {
-                Mail mail = retrieveMetadata(key);
-                return condition.test(mail);
-            }).iterator();
+            .filter(key -> Optional.ofNullable(retrieveMetadata(key))
+                .map(condition::test)
+                .orElse(false))
+            .iterator();
     }
 
     private Mail retrieveMetadata(MailKey key) {
-        return mailDAO.read(url, key)
-            .handle(publishIfPresent())
+        return readMail(key)
             .map(mailDTO -> mailDTO.getMailBuilder().build())
             .block();
     }
 
     @Override
     public Mail retrieve(MailKey key) {
-        return mailDAO.read(url, key)
-            .handle(publishIfPresent())
+        return readMail(key)
             .flatMap(this::toMail)
             .blockOptional()
             .orElse(null);
     }
 
+    /**
+     * Reads the content of a mail, auto-healing the keys that are no longer 
backed by any content.
+     *
+     * Such 'lonely' keys can be encountered when a key deletion gets 
resurrected - the keys table relies on a
+     * zero gc_grace_seconds in order not to accumulate tombstones. Leaving 
them behind would make the repository
+     * report mails that can not be read anymore.
+     */
+    private Mono<MailDTO> readMail(MailKey key) {
+        return mailDAO.read(url, key)
+            .handle(publishIfPresent())
+            .switchIfEmpty(Mono.defer(() -> removeLonelyKey(key)));
+    }
+
+    private Mono<MailDTO> removeLonelyKey(MailKey key) {
+        return keysDAO.remove(url, key)
+            .doOnNext(any -> LOGGER.info("Removed key {} of mail repository {} 
as it was not backed by any content", key.asString(), url.asString()))
+            .then(Mono.empty());
+    }
+
     private Mono<Mail> toMail(MailDTO mailDTO) {
         MimeMessagePartsId parts = blobIds(mailDTO);
 
diff --git 
a/server/mailrepository/mailrepository-cassandra/src/main/java/org/apache/james/mailrepository/cassandra/CassandraMailRepositoryDataDefinition.java
 
b/server/mailrepository/mailrepository-cassandra/src/main/java/org/apache/james/mailrepository/cassandra/CassandraMailRepositoryDataDefinition.java
index 602829666f..78cbbda1f4 100644
--- 
a/server/mailrepository/mailrepository-cassandra/src/main/java/org/apache/james/mailrepository/cassandra/CassandraMailRepositoryDataDefinition.java
+++ 
b/server/mailrepository/mailrepository-cassandra/src/main/java/org/apache/james/mailrepository/cassandra/CassandraMailRepositoryDataDefinition.java
@@ -37,6 +37,12 @@ public interface CassandraMailRepositoryDataDefinition {
 
         .table(MailRepositoryTable.KEYS_TABLE_NAME)
         .comment("Per-mailRepository mail key list")
+        .options(options -> options
+            // A mail repository holds all its keys within a single partition, 
thus deletions quickly pile up
+            // tombstones there, up to the point reads (list, count) start 
failing. As CassandraMailRepository
+            // auto-heals the keys that would be resurrected by an 
unreplicated deletion, we can afford to
+            // collect those tombstones right away.
+            .withGcGraceSeconds(0))
         .statement(statement -> types -> statement
             .withPartitionKey(MailRepositoryTable.REPOSITORY_NAME, TEXT)
             .withClusteringColumn(MailRepositoryTable.MAIL_KEY, TEXT))
diff --git 
a/server/mailrepository/mailrepository-cassandra/src/test/java/org/apache/james/mailrepository/cassandra/CassandraMailRepositoryTest.java
 
b/server/mailrepository/mailrepository-cassandra/src/test/java/org/apache/james/mailrepository/cassandra/CassandraMailRepositoryTest.java
index aaf1473929..fa33eb4be4 100644
--- 
a/server/mailrepository/mailrepository-cassandra/src/test/java/org/apache/james/mailrepository/cassandra/CassandraMailRepositoryTest.java
+++ 
b/server/mailrepository/mailrepository-cassandra/src/test/java/org/apache/james/mailrepository/cassandra/CassandraMailRepositoryTest.java
@@ -212,4 +212,71 @@ class CassandraMailRepositoryTest {
         }
     }
 
+    @Nested
+    class LonelyKeysTest {
+        static final MailKey LONELY_KEY = new MailKey("lonely");
+
+        CassandraMailRepositoryKeysDAO keysDAO;
+        MailRepository testee;
+
+        @BeforeEach
+        void setUp(CassandraCluster cassandra) {
+            CassandraMailRepositoryMailDaoV2 v2 = new 
CassandraMailRepositoryMailDaoV2(cassandra.getConf(), BLOB_ID_FACTORY);
+            keysDAO = new CassandraMailRepositoryKeysDAO(cassandra.getConf(), 
CassandraConfiguration.DEFAULT_CONFIGURATION);
+            BlobStore blobStore = 
CassandraBlobStoreFactory.forTesting(cassandra.getConf(), new 
RecordingMetricFactory())
+                .passthrough();
+
+            testee = new CassandraMailRepository(URL, keysDAO, v2, 
MimeMessageStore.factory(blobStore));
+        }
+
+        @Test
+        void retrieveShouldReturnNullWhenKeyIsNotBackedByAnyContent() throws 
Exception {
+            keysDAO.store(URL, LONELY_KEY).block();
+
+            assertThat(testee.retrieve(LONELY_KEY)).isNull();
+        }
+
+        @Test
+        void retrieveShouldRemoveKeysThatAreNotBackedByAnyContent() throws 
Exception {
+            keysDAO.store(URL, LONELY_KEY).block();
+
+            testee.retrieve(LONELY_KEY);
+
+            assertThat(testee.list()).toIterable().isEmpty();
+        }
+
+        @Test
+        void listWithConditionShouldIgnoreKeysThatAreNotBackedByAnyContent() 
throws Exception {
+            keysDAO.store(URL, LONELY_KEY).block();
+            MailKey key = testee.store(MailImpl.builder()
+                .name("mail1")
+                .sender("[email protected]")
+                .addRecipient("[email protected]")
+                .mimeMessage(MimeMessageBuilder.mimeMessageBuilder().build())
+                .build());
+
+            assertThat(testee.list(new 
MailRepository.SenderCondition("[email protected]")))
+                .toIterable()
+                .containsOnly(key);
+        }
+
+        @Test
+        void listWithConditionShouldRemoveKeysThatAreNotBackedByAnyContent() 
throws Exception {
+            keysDAO.store(URL, LONELY_KEY).block();
+
+            testee.list(new 
MailRepository.SenderCondition("[email protected]")).forEachRemaining(key -> { 
});
+
+            assertThat(testee.list()).toIterable().isEmpty();
+        }
+
+        @Test
+        void removeAllShouldRemoveKeysThatAreNotBackedByAnyContent() throws 
Exception {
+            keysDAO.store(URL, LONELY_KEY).block();
+
+            testee.removeAll();
+
+            assertThat(testee.list()).toIterable().isEmpty();
+        }
+    }
+
 }
\ No newline at end of file


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

Reply via email to