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]