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


The following commit(s) were added to refs/heads/master by this push:
     new e49adae86b [FIX] Memory: findMetadata ignored the messageId, Email/set 
update corrupted data and stalled (#3233)
e49adae86b is described below

commit e49adae86b5ee26b5cb0d222fc4e1a6cc8d1fd6c
Author: Benoit TELLIER <[email protected]>
AuthorDate: Mon Oct 5 11:45:38 2026 +0200

    [FIX] Memory: findMetadata ignored the messageId, Email/set update 
corrupted data and stalled (#3233)
    
    It returned the metadata of every message of every mailbox. JMAP Email/set
    update then applied range updates (flags, moves) to unrequested messages.
    
    groupBy + filterWhen left the groups of unreadable mailboxes unsubscribed:
    their elements stayed buffered and, once the groupBy prefetch (256) was
    reached, the Flux never completed.
    
    Filter each element against a per call cache of the Read right instead.
    
    Defensive: restrict the store metadata to the requested ids so that range
    updates and the `updated` field can not span unrequested messages.
---
 .../inmemory/mail/InMemoryMessageIdMapper.java     |  3 +-
 .../james/mailbox/store/StoreMessageIdManager.java | 18 ++++++++----
 .../store/AbstractMessageIdManagerStorageTest.java | 34 ++++++++++++++++++++++
 .../store/mail/model/MessageIdMapperTest.java      | 14 +++++++++
 .../jmap/method/EmailSetUpdatePerformer.scala      |  5 +++-
 5 files changed, 65 insertions(+), 9 deletions(-)

diff --git 
a/mailbox/memory/src/main/java/org/apache/james/mailbox/inmemory/mail/InMemoryMessageIdMapper.java
 
b/mailbox/memory/src/main/java/org/apache/james/mailbox/inmemory/mail/InMemoryMessageIdMapper.java
index b90eb181d0..431b08d0e3 100644
--- 
a/mailbox/memory/src/main/java/org/apache/james/mailbox/inmemory/mail/InMemoryMessageIdMapper.java
+++ 
b/mailbox/memory/src/main/java/org/apache/james/mailbox/inmemory/mail/InMemoryMessageIdMapper.java
@@ -72,8 +72,7 @@ public class InMemoryMessageIdMapper implements 
MessageIdMapper {
 
     @Override
     public Publisher<ComposedMessageIdWithMetaData> findMetadata(MessageId 
messageId) {
-        return mailboxMapper.list()
-            .flatMap(mailbox -> messageMapper.findInMailboxReactive(mailbox, 
MessageRange.all(), MessageMapper.FetchType.FULL, UNLIMITED), 
DEFAULT_CONCURRENCY)
+        return findReactive(ImmutableList.of(messageId), 
MessageMapper.FetchType.METADATA)
             .map(message -> new ComposedMessageIdWithMetaData(
                 new ComposedMessageId(
                     message.getMailboxId(),
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 ac30c59119..d0ffdf0492 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
@@ -29,6 +29,7 @@ import java.util.List;
 import java.util.Map;
 import java.util.Optional;
 import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
 import java.util.function.Function;
 import java.util.function.Predicate;
 import java.util.stream.Stream;
@@ -184,10 +185,9 @@ public class StoreMessageIdManager implements 
MessageIdManager {
         MessageMapper.FetchType fetchType = 
FetchGroupConverter.getFetchType(fetchGroup);
         boolean delayError = false;
         int prefetch = 1;
+        Function<MailboxId, Mono<Boolean>> canRead = 
cachedReadRight(mailboxSession);
         return messageIdMapper.findReactive(messageIds, fetchType)
-            .groupBy(MailboxMessage::getMailboxId)
-            .filterWhen(groupedFlux -> 
hasRightsOnMailboxReactive(mailboxSession, 
Right.Read).apply(groupedFlux.key()), DEFAULT_CONCURRENCY)
-            .flatMap(Function.identity(), DEFAULT_CONCURRENCY)
+            .filterWhen(message -> canRead.apply(message.getMailboxId()), 
DEFAULT_CONCURRENCY)
             .publishOn(forFetchType(fetchType), delayError, prefetch)
             
.map(Throwing.function(messageResultConverter(fetchGroup)).sneakyThrow());
     }
@@ -203,11 +203,17 @@ public class StoreMessageIdManager implements 
MessageIdManager {
     public Publisher<ComposedMessageIdWithMetaData> 
messagesMetadata(Collection<MessageId> ids, MailboxSession session) {
         MessageIdMapper messageIdMapper = 
mailboxSessionMapperFactory.getMessageIdMapper(session);
         int concurrency = 4;
+        Function<MailboxId, Mono<Boolean>> canRead = cachedReadRight(session);
         return Flux.fromIterable(ids)
             .flatMap(messageIdMapper::findMetadata, concurrency)
-            .groupBy(metaData -> 
metaData.getComposedMessageId().getMailboxId())
-            .filterWhen(groupedFlux -> hasRightsOnMailboxReactive(session, 
Right.Read).apply(groupedFlux.key()), DEFAULT_CONCURRENCY)
-            .flatMap(Function.identity(), DEFAULT_CONCURRENCY);
+            .filterWhen(metaData -> 
canRead.apply(metaData.getComposedMessageId().getMailboxId()), 
DEFAULT_CONCURRENCY);
+    }
+
+    // Unlike groupBy + filterWhen, never leaves rejected elements buffered: 
those would stall the Flux once the groupBy prefetch is reached
+    private Function<MailboxId, Mono<Boolean>> cachedReadRight(MailboxSession 
session) {
+        Function<MailboxId, Mono<Boolean>> hasReadRight = 
hasRightsOnMailboxReactive(session, Right.Read);
+        Map<MailboxId, Mono<Boolean>> rightsByMailbox = new 
ConcurrentHashMap<>();
+        return mailboxId -> rightsByMailbox.computeIfAbsent(mailboxId, id -> 
hasReadRight.apply(id).cache());
     }
 
     private Mono<ImmutableSet<MailboxId>> getAllowedMailboxIds(MailboxSession 
mailboxSession, Stream<MailboxId> idList, Right... rights) {
diff --git 
a/mailbox/store/src/test/java/org/apache/james/mailbox/store/AbstractMessageIdManagerStorageTest.java
 
b/mailbox/store/src/test/java/org/apache/james/mailbox/store/AbstractMessageIdManagerStorageTest.java
index c3d6aa6eda..8eda7442b5 100644
--- 
a/mailbox/store/src/test/java/org/apache/james/mailbox/store/AbstractMessageIdManagerStorageTest.java
+++ 
b/mailbox/store/src/test/java/org/apache/james/mailbox/store/AbstractMessageIdManagerStorageTest.java
@@ -23,10 +23,12 @@ import static 
org.apache.james.mailbox.fixture.MailboxFixture.BOB;
 import static org.assertj.core.api.Assertions.assertThat;
 import static org.assertj.core.api.Assertions.assertThatThrownBy;
 
+import java.time.Duration;
 import java.util.List;
 import java.util.Map;
 import java.util.Set;
 import java.util.function.Predicate;
+import java.util.stream.IntStream;
 
 import jakarta.mail.Flags;
 
@@ -40,6 +42,7 @@ import org.apache.james.mailbox.ModSeq;
 import org.apache.james.mailbox.exception.MailboxException;
 import org.apache.james.mailbox.exception.MailboxNotFoundException;
 import org.apache.james.mailbox.fixture.MailboxFixture;
+import org.apache.james.mailbox.model.ComposedMessageIdWithMetaData;
 import org.apache.james.mailbox.model.DeleteResult;
 import org.apache.james.mailbox.model.FetchGroup;
 import org.apache.james.mailbox.model.Mailbox;
@@ -55,11 +58,14 @@ import org.junit.jupiter.api.Test;
 import com.google.common.collect.ImmutableList;
 import com.google.common.collect.ImmutableMap;
 
+import reactor.core.publisher.Flux;
+
 public abstract class AbstractMessageIdManagerStorageTest {
     public static final Flags FLAGS = new Flags();
 
     private static final MessageUid messageUid1 = MessageUid.of(111);
     private static final MessageUid messageUid2 = MessageUid.of(222);
+    private static final int MANY_MESSAGES = 300;
 
     private MessageIdManagerTestSystem testingData;
     private MessageIdManager messageIdManager;
@@ -119,6 +125,34 @@ public abstract class AbstractMessageIdManagerStorageTest {
         messageIdManager.setInMailboxes(messageId, 
ImmutableList.of(aliceMailbox1.getMailboxId()), aliceSession);
     }
 
+    @Test
+    void messagesMetadataShouldCompleteWhenManyMessagesAreNotReadable() {
+        List<MessageId> aliceMessageIds = persistAliceMessages(MANY_MESSAGES);
+
+        List<ComposedMessageIdWithMetaData> bobMetadata = 
Flux.from(messageIdManager.messagesMetadata(aliceMessageIds, bobSession))
+            .collectList()
+            .block(Duration.ofSeconds(10));
+
+        assertThat(bobMetadata).isEmpty();
+    }
+
+    @Test
+    void getMessagesReactiveShouldCompleteWhenManyMessagesAreNotReadable() {
+        List<MessageId> aliceMessageIds = persistAliceMessages(MANY_MESSAGES);
+
+        List<MessageResult> bobMessages = 
Flux.from(messageIdManager.getMessagesReactive(aliceMessageIds, 
FetchGroup.MINIMAL, bobSession))
+            .collectList()
+            .block(Duration.ofSeconds(10));
+
+        assertThat(bobMessages).isEmpty();
+    }
+
+    private List<MessageId> persistAliceMessages(int count) {
+        return IntStream.rangeClosed(1, count)
+            .mapToObj(uid -> testingData.persist(aliceMailbox1.getMailboxId(), 
MessageUid.of(uid), FLAGS, aliceSession))
+            .collect(ImmutableList.toImmutableList());
+    }
+
     @Test
     void getMessagesShouldReturnStoredResults() throws Exception {
         MessageId messageId = 
testingData.persist(aliceMailbox1.getMailboxId(), messageUid1, FLAGS, 
aliceSession);
diff --git 
a/mailbox/store/src/test/java/org/apache/james/mailbox/store/mail/model/MessageIdMapperTest.java
 
b/mailbox/store/src/test/java/org/apache/james/mailbox/store/mail/model/MessageIdMapperTest.java
index f5b2856754..899a24c85b 100644
--- 
a/mailbox/store/src/test/java/org/apache/james/mailbox/store/mail/model/MessageIdMapperTest.java
+++ 
b/mailbox/store/src/test/java/org/apache/james/mailbox/store/mail/model/MessageIdMapperTest.java
@@ -63,6 +63,8 @@ import com.google.common.collect.ImmutableList;
 import com.google.common.collect.ImmutableMultimap;
 import com.google.common.collect.Multimap;
 
+import reactor.core.publisher.Flux;
+
 public abstract class MessageIdMapperTest {
     private static final Username BENWA = Username.of("benwa");
 
@@ -132,6 +134,18 @@ public abstract class MessageIdMapperTest {
         assertMessages(messages).containOnly(message1, message4, message3);
     }
 
+    @Test
+    void findMetadataShouldReturnOnlyTheGivenMessage() throws MailboxException 
{
+        saveMessages();
+
+        List<MessageId> messageIds = 
Flux.from(sut.findMetadata(message1.getMessageId()))
+            .map(metaData -> metaData.getComposedMessageId().getMessageId())
+            .collectList()
+            .block();
+
+        assertThat(messageIds).containsOnly(message1.getMessageId());
+    }
+
     @Test
     void findMailboxesShouldReturnEmptyWhenMessageDoesntExist() {
         
assertThat(sut.findMailboxes(mapperProvider.generateMessageId())).isEmpty();
diff --git 
a/server/protocols/jmap-rfc-8621/src/main/scala/org/apache/james/jmap/method/EmailSetUpdatePerformer.scala
 
b/server/protocols/jmap-rfc-8621/src/main/scala/org/apache/james/jmap/method/EmailSetUpdatePerformer.scala
index 4aea71c270..d5bfaf6760 100644
--- 
a/server/protocols/jmap-rfc-8621/src/main/scala/org/apache/james/jmap/method/EmailSetUpdatePerformer.scala
+++ 
b/server/protocols/jmap-rfc-8621/src/main/scala/org/apache/james/jmap/method/EmailSetUpdatePerformer.scala
@@ -114,8 +114,11 @@ class EmailSetUpdatePerformer @Inject() (serializer: 
EmailSetSerializer,
       case _ => None
     })
 
+    val requestedIds: Set[MessageId] = validUpdates.map(_._1).toSet
+
     for {
-      updates <- 
SFlux.fromPublisher(messageIdManager.messagesMetadata(validUpdates.map(_._1).asJavaCollection,
 session))
+      updates <- 
SFlux.fromPublisher(messageIdManager.messagesMetadata(requestedIds.asJavaCollection,
 session))
+        .filter(metaData => 
requestedIds.contains(metaData.getComposedMessageId.getMessageId))
         .collectMultimap(metaData => 
metaData.getComposedMessageId.getMessageId)
         .flatMap(doUpdate(validUpdates, _, session))
     } yield {


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

Reply via email to