From 8b6308434ca10a28956187cbc032c4e13619de15 Mon Sep 17 00:00:00 2001 From: Benoit Tellier Date: Sat, 8 May 2021 14:46:48 +0700 Subject: [PATCH 01/15] [REFACTORING] DefaultMailboxesProvisioner: Avoid re-opening a session --- .../james/jmap/http/DefaultMailboxesProvisioner.java | 10 ++-------- 1 file changed, 2 insertions(+), 8 deletions(-) diff --git a/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/http/DefaultMailboxesProvisioner.java b/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/http/DefaultMailboxesProvisioner.java index a5766115174..8e4ff9336c7 100644 --- a/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/http/DefaultMailboxesProvisioner.java +++ b/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/http/DefaultMailboxesProvisioner.java @@ -26,7 +26,6 @@ import javax.inject.Inject; -import org.apache.james.core.Username; import org.apache.james.mailbox.DefaultMailboxes; import org.apache.james.mailbox.MailboxManager; import org.apache.james.mailbox.MailboxSession; @@ -63,15 +62,10 @@ public DefaultMailboxesProvisioner(MailboxManager mailboxManager, public Mono createMailboxesIfNeeded(MailboxSession session) { return metricFactory.decorateSupplierWithTimerMetric("JMAP-mailboxes-provisioning", - () -> { - Username username = session.getUser(); - return createDefaultMailboxes(username); - }); + () -> createDefaultMailboxes(session)); } - private Mono createDefaultMailboxes(Username username) { - MailboxSession session = mailboxManager.createSystemSession(username); - + private Mono createDefaultMailboxes(MailboxSession session) { return Flux.fromIterable(DefaultMailboxes.DEFAULT_MAILBOXES) .map(toMailboxPath(session)) .filterWhen(mailboxPath -> mailboxDoesntExist(mailboxPath, session), DEFAULT_CONCURRENCY) From 82d434fd73a5367410804b082aa0342c92bd0f9c Mon Sep 17 00:00:00 2001 From: Benoit Tellier Date: Sat, 8 May 2021 09:48:30 +0700 Subject: [PATCH 02/15] [REFACTORING] Avoid a blocking modseq call upon message flags update --- .../cassandra/mail/CassandraMessageIdMapper.java | 11 +++++------ 1 file changed, 5 insertions(+), 6 deletions(-) diff --git a/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/CassandraMessageIdMapper.java b/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/CassandraMessageIdMapper.java index 59ea9873d5a..e5db7eafb1b 100644 --- a/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/CassandraMessageIdMapper.java +++ b/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/CassandraMessageIdMapper.java @@ -49,7 +49,6 @@ import org.apache.james.mailbox.store.mail.MailboxMapper; import org.apache.james.mailbox.store.mail.MessageIdMapper; import org.apache.james.mailbox.store.mail.MessageMapper.FetchType; -import org.apache.james.mailbox.store.mail.ModSeqProvider; import org.apache.james.mailbox.store.mail.model.MailboxMessage; import org.apache.james.util.FunctionalUtils; import org.apache.james.util.ReactorUtils; @@ -80,14 +79,14 @@ public class CassandraMessageIdMapper implements MessageIdMapper { private final CassandraMessageDAO messageDAO; private final CassandraMessageDAOV3 messageDAOV3; private final CassandraIndexTableHandler indexTableHandler; - private final ModSeqProvider modSeqProvider; + private final CassandraModSeqProvider modSeqProvider; private final AttachmentLoader attachmentLoader; private final CassandraConfiguration cassandraConfiguration; public CassandraMessageIdMapper(MailboxMapper mailboxMapper, CassandraMailboxDAO mailboxDAO, CassandraAttachmentMapper attachmentMapper, CassandraMessageIdToImapUidDAO imapUidDAO, CassandraMessageIdDAO messageIdDAO, CassandraMessageDAO messageDAO, CassandraMessageDAOV3 messageDAOV3, CassandraIndexTableHandler indexTableHandler, - ModSeqProvider modSeqProvider, CassandraConfiguration cassandraConfiguration) { + CassandraModSeqProvider modSeqProvider, CassandraConfiguration cassandraConfiguration) { this.mailboxMapper = mailboxMapper; this.mailboxDAO = mailboxDAO; @@ -305,11 +304,11 @@ private Mono> updateFlags(Flags newSt if (identicalFlags(oldComposedId, newFlags)) { return Mono.just(Pair.of(oldComposedId.getFlags(), oldComposedId)); } else { - return Mono - .fromCallable(() -> new ComposedMessageIdWithMetaData( + return modSeqProvider.nextModSeq(cassandraId) + .map(modSeq -> new ComposedMessageIdWithMetaData( oldComposedId.getComposedMessageId(), newFlags, - modSeqProvider.nextModSeq(cassandraId))) + modSeq)) .flatMap(newComposedId -> updateFlags(oldComposedId, newComposedId)); } } From c605f5f2fbf4b675ef4e5bf0429ec9e92ecbd236 Mon Sep 17 00:00:00 2001 From: Benoit Tellier Date: Fri, 7 May 2021 18:13:13 +0700 Subject: [PATCH 03/15] [REFACTORING] Fasten JMAP Mailbox/get Reactive quotas --- .../james/mailbox/quota/MaxQuotaManager.java | 22 +++++++ .../james/mailbox/quota/QuotaManager.java | 3 + .../mailbox/quota/QuotaRootResolver.java | 4 ++ .../CassandraPerUserMaxQuotaManager.java | 44 +++++++------ .../quota/DefaultUserQuotaRootResolver.java | 15 ++++- .../quota/ListeningCurrentQuotaUpdater.java | 4 +- .../mailbox/store/quota/NoQuotaManager.java | 8 +++ .../store/quota/StoreQuotaManager.java | 27 ++++++-- .../utils/quotas/DefaultQuotaLoader.java | 25 ++++---- .../jmap/draft/utils/quotas/QuotaLoader.java | 5 +- .../QuotaLoaderWithDefaultPreloaded.java | 61 ++++++++----------- .../james/jmap/utils/quotas/QuotaLoader.scala | 2 - .../QuotaLoaderWithPreloadedDefault.scala | 13 ++-- .../james/jmap/utils/quotas/QuotaReader.scala | 16 +++-- 14 files changed, 155 insertions(+), 94 deletions(-) diff --git a/mailbox/api/src/main/java/org/apache/james/mailbox/quota/MaxQuotaManager.java b/mailbox/api/src/main/java/org/apache/james/mailbox/quota/MaxQuotaManager.java index b105b42c185..3d2cb5c655b 100644 --- a/mailbox/api/src/main/java/org/apache/james/mailbox/quota/MaxQuotaManager.java +++ b/mailbox/api/src/main/java/org/apache/james/mailbox/quota/MaxQuotaManager.java @@ -32,9 +32,13 @@ import org.apache.james.mailbox.model.Quota; import org.apache.james.mailbox.model.Quota.Scope; import org.apache.james.mailbox.model.QuotaRoot; +import org.reactivestreams.Publisher; import com.github.fge.lambdas.Throwing; +import reactor.core.publisher.Mono; +import reactor.core.scheduler.Schedulers; + /** * This interface describe how to set the max quotas for users * Part of RFC 2087 implementation @@ -150,12 +154,30 @@ default Optional getMaxMessage(Map listMaxMessagesDetails(QuotaRoot quotaRoot); + default Publisher> listMaxMessagesDetailsReactive(QuotaRoot quotaRoot) { + return Mono.fromCallable(() -> listMaxMessagesDetails(quotaRoot)) + .subscribeOn(Schedulers.elastic()); + } + Map listMaxStorageDetails(QuotaRoot quotaRoot); + default Publisher> listMaxStorageDetailsReactive(QuotaRoot quotaRoot) { + return Mono.fromCallable(() -> listMaxStorageDetails(quotaRoot)) + .subscribeOn(Schedulers.elastic()); + } + + default QuotaDetails quotaDetails(QuotaRoot quotaRoot) { return new QuotaDetails(listMaxMessagesDetails(quotaRoot), listMaxStorageDetails(quotaRoot)); } + default Publisher quotaDetailsReactive(QuotaRoot quotaRoot) { + return Mono.zip( + Mono.from(listMaxMessagesDetailsReactive(quotaRoot)), + Mono.from(listMaxStorageDetailsReactive(quotaRoot))) + .map(tuple -> new QuotaDetails(tuple.getT1(), tuple.getT2())); + } + Optional getDomainMaxMessage(Domain domain); void setDomainMaxMessage(Domain domain, QuotaCountLimit count) throws MailboxException; diff --git a/mailbox/api/src/main/java/org/apache/james/mailbox/quota/QuotaManager.java b/mailbox/api/src/main/java/org/apache/james/mailbox/quota/QuotaManager.java index a08c0bcd1e6..4d34ecf5f72 100644 --- a/mailbox/api/src/main/java/org/apache/james/mailbox/quota/QuotaManager.java +++ b/mailbox/api/src/main/java/org/apache/james/mailbox/quota/QuotaManager.java @@ -25,6 +25,7 @@ import org.apache.james.mailbox.exception.MailboxException; import org.apache.james.mailbox.model.Quota; import org.apache.james.mailbox.model.QuotaRoot; +import org.reactivestreams.Publisher; /** @@ -68,4 +69,6 @@ public Quota getStorageQuota() { Quota getStorageQuota(QuotaRoot quotaRoot) throws MailboxException; Quotas getQuotas(QuotaRoot quotaRoot) throws MailboxException; + + Publisher getQuotasReactive(QuotaRoot quotaRoot); } diff --git a/mailbox/api/src/main/java/org/apache/james/mailbox/quota/QuotaRootResolver.java b/mailbox/api/src/main/java/org/apache/james/mailbox/quota/QuotaRootResolver.java index 60510f60e8e..efbc244ab3b 100644 --- a/mailbox/api/src/main/java/org/apache/james/mailbox/quota/QuotaRootResolver.java +++ b/mailbox/api/src/main/java/org/apache/james/mailbox/quota/QuotaRootResolver.java @@ -37,10 +37,14 @@ public interface QuotaRootResolver extends QuotaRootDeserializer { */ QuotaRoot getQuotaRoot(MailboxPath mailboxPath) throws MailboxException; + Publisher getQuotaRootReactive(MailboxPath mailboxPath); + QuotaRoot getQuotaRoot(MailboxId mailboxId) throws MailboxException; QuotaRoot getQuotaRoot(Mailbox mailbox) throws MailboxException; + Publisher getQuotaRootReactive(Mailbox mailbox); + Publisher getQuotaRootReactive(MailboxId mailboxId); Publisher retrieveAssociatedMailboxes(QuotaRoot quotaRoot, MailboxSession mailboxSession); diff --git a/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/quota/CassandraPerUserMaxQuotaManager.java b/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/quota/CassandraPerUserMaxQuotaManager.java index 5fe0095e5fd..dd6884b0bea 100644 --- a/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/quota/CassandraPerUserMaxQuotaManager.java +++ b/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/quota/CassandraPerUserMaxQuotaManager.java @@ -136,34 +136,42 @@ public Optional getGlobalMaxMessage() { @Override public Map listMaxMessagesDetails(QuotaRoot quotaRoot) { + return listMaxMessagesDetailsReactive(quotaRoot).block(); + } + + @Override + public Mono> listMaxMessagesDetailsReactive(QuotaRoot quotaRoot) { return Flux.merge( - perUserQuota.getMaxMessage(quotaRoot) - .map(limit -> Pair.of(Quota.Scope.User, limit)), - Mono.justOrEmpty(quotaRoot.getDomain()) - .flatMap(perDomainQuota::getMaxMessage) - .map(limit -> Pair.of(Quota.Scope.Domain, limit)), - globalQuota.getGlobalMaxMessage() - .map(limit -> Pair.of(Quota.Scope.Global, limit))) + perUserQuota.getMaxMessage(quotaRoot) + .map(limit -> Pair.of(Quota.Scope.User, limit)), + Mono.justOrEmpty(quotaRoot.getDomain()) + .flatMap(perDomainQuota::getMaxMessage) + .map(limit -> Pair.of(Quota.Scope.Domain, limit)), + globalQuota.getGlobalMaxMessage() + .map(limit -> Pair.of(Quota.Scope.Global, limit))) .collect(Guavate.toImmutableMap( Pair::getKey, - Pair::getValue)) - .block(); + Pair::getValue)); } @Override public Map listMaxStorageDetails(QuotaRoot quotaRoot) { + return listMaxStorageDetailsReactive(quotaRoot).block(); + } + + @Override + public Mono> listMaxStorageDetailsReactive(QuotaRoot quotaRoot) { return Flux.merge( - perUserQuota.getMaxStorage(quotaRoot) - .map(limit -> Pair.of(Quota.Scope.User, limit)), - Mono.justOrEmpty(quotaRoot.getDomain()) - .flatMap(perDomainQuota::getMaxStorage) - .map(limit -> Pair.of(Quota.Scope.Domain, limit)), - globalQuota.getGlobalMaxStorage() - .map(limit -> Pair.of(Quota.Scope.Global, limit))) + perUserQuota.getMaxStorage(quotaRoot) + .map(limit -> Pair.of(Quota.Scope.User, limit)), + Mono.justOrEmpty(quotaRoot.getDomain()) + .flatMap(perDomainQuota::getMaxStorage) + .map(limit -> Pair.of(Quota.Scope.Domain, limit)), + globalQuota.getGlobalMaxStorage() + .map(limit -> Pair.of(Quota.Scope.Global, limit))) .collect(Guavate.toImmutableMap( Pair::getKey, - Pair::getValue)) - .block(); + Pair::getValue)); } @Override diff --git a/mailbox/store/src/main/java/org/apache/james/mailbox/store/quota/DefaultUserQuotaRootResolver.java b/mailbox/store/src/main/java/org/apache/james/mailbox/store/quota/DefaultUserQuotaRootResolver.java index b6f206939e0..629d6aa056e 100644 --- a/mailbox/store/src/main/java/org/apache/james/mailbox/store/quota/DefaultUserQuotaRootResolver.java +++ b/mailbox/store/src/main/java/org/apache/james/mailbox/store/quota/DefaultUserQuotaRootResolver.java @@ -37,6 +37,7 @@ import org.apache.james.mailbox.quota.QuotaRootDeserializer; import org.apache.james.mailbox.quota.UserQuotaRootResolver; import org.apache.james.mailbox.store.MailboxSessionMapperFactory; +import org.reactivestreams.Publisher; import com.google.common.base.Preconditions; import com.google.common.base.Splitter; @@ -95,12 +96,12 @@ public QuotaRoot getQuotaRoot(MailboxPath mailboxPath) { } @Override - public QuotaRoot getQuotaRoot(Mailbox mailbox) throws MailboxException { + public QuotaRoot getQuotaRoot(Mailbox mailbox) { return getQuotaRoot(mailbox.generateAssociatedPath()); } @Override - public QuotaRoot getQuotaRoot(MailboxId mailboxId) throws MailboxException { + public QuotaRoot getQuotaRoot(MailboxId mailboxId) { return getQuotaRootReactive(mailboxId).block(); } @@ -115,6 +116,16 @@ public Mono getQuotaRootReactive(MailboxId mailboxId) { .map(this::forUser); } + @Override + public Publisher getQuotaRootReactive(MailboxPath mailboxPath) { + return Mono.just(getQuotaRoot(mailboxPath)); + } + + @Override + public Publisher getQuotaRootReactive(Mailbox mailbox) { + return Mono.just(getQuotaRoot(mailbox)); + } + @Override public QuotaRoot fromString(String serializedQuotaRoot) throws MailboxException { return QUOTA_ROOT_DESERIALIZER.fromString(serializedQuotaRoot); diff --git a/mailbox/store/src/main/java/org/apache/james/mailbox/store/quota/ListeningCurrentQuotaUpdater.java b/mailbox/store/src/main/java/org/apache/james/mailbox/store/quota/ListeningCurrentQuotaUpdater.java index 2bafe6462b8..cfe2d801582 100644 --- a/mailbox/store/src/main/java/org/apache/james/mailbox/store/quota/ListeningCurrentQuotaUpdater.java +++ b/mailbox/store/src/main/java/org/apache/james/mailbox/store/quota/ListeningCurrentQuotaUpdater.java @@ -110,9 +110,7 @@ private Mono handleAddedEvent(Added added, QuotaRoot quotaRoot) { } private Mono dispatchNewQuota(QuotaRoot quotaRoot, Username username) { - Mono quotasMono = Mono.fromCallable(() -> quotaManager.getQuotas(quotaRoot)); - - return quotasMono.subscribeOn(Schedulers.elastic()) + return Mono.from(quotaManager.getQuotasReactive(quotaRoot)) .flatMap(quotas -> eventBus.dispatch( EventFactory.quotaUpdated() .randomEventId() diff --git a/mailbox/store/src/main/java/org/apache/james/mailbox/store/quota/NoQuotaManager.java b/mailbox/store/src/main/java/org/apache/james/mailbox/store/quota/NoQuotaManager.java index e40b084e52b..89e3d639b1d 100644 --- a/mailbox/store/src/main/java/org/apache/james/mailbox/store/quota/NoQuotaManager.java +++ b/mailbox/store/src/main/java/org/apache/james/mailbox/store/quota/NoQuotaManager.java @@ -26,6 +26,9 @@ import org.apache.james.mailbox.model.Quota; import org.apache.james.mailbox.model.QuotaRoot; import org.apache.james.mailbox.quota.QuotaManager; +import org.reactivestreams.Publisher; + +import reactor.core.publisher.Mono; /** * This quota manager is intended to be used when you want to deactivate the Quota feature @@ -52,4 +55,9 @@ public Quota getStorageQuota(QuotaRoot quotaRoot public Quotas getQuotas(QuotaRoot quotaRoot) { return new Quotas(getMessageQuota(quotaRoot), getStorageQuota(quotaRoot)); } + + @Override + public Publisher getQuotasReactive(QuotaRoot quotaRoot) { + return Mono.just(getQuotas(quotaRoot)); + } } diff --git a/mailbox/store/src/main/java/org/apache/james/mailbox/store/quota/StoreQuotaManager.java b/mailbox/store/src/main/java/org/apache/james/mailbox/store/quota/StoreQuotaManager.java index 17d80e7ac0c..aa010bce42a 100644 --- a/mailbox/store/src/main/java/org/apache/james/mailbox/store/quota/StoreQuotaManager.java +++ b/mailbox/store/src/main/java/org/apache/james/mailbox/store/quota/StoreQuotaManager.java @@ -27,7 +27,6 @@ import org.apache.james.core.quota.QuotaCountUsage; import org.apache.james.core.quota.QuotaSizeLimit; import org.apache.james.core.quota.QuotaSizeUsage; -import org.apache.james.mailbox.exception.MailboxException; import org.apache.james.mailbox.model.CurrentQuotas; import org.apache.james.mailbox.model.Quota; import org.apache.james.mailbox.model.Quota.Scope; @@ -35,8 +34,10 @@ import org.apache.james.mailbox.quota.CurrentQuotaManager; import org.apache.james.mailbox.quota.MaxQuotaManager; import org.apache.james.mailbox.quota.QuotaManager; +import org.reactivestreams.Publisher; import reactor.core.publisher.Mono; +import reactor.core.scheduler.Schedulers; /** * Default implementation for the Quota Manager. @@ -54,7 +55,7 @@ public StoreQuotaManager(CurrentQuotaManager currentQuotaManager, MaxQuotaManage } @Override - public Quota getMessageQuota(QuotaRoot quotaRoot) throws MailboxException { + public Quota getMessageQuota(QuotaRoot quotaRoot) { Map maxMessageDetails = maxQuotaManager.listMaxMessagesDetails(quotaRoot); return Quota.builder() .used(Mono.from(currentQuotaManager.getCurrentMessageCount(quotaRoot)).block()) @@ -65,7 +66,7 @@ public Quota getMessageQuota(QuotaRoot quotaRo @Override - public Quota getStorageQuota(QuotaRoot quotaRoot) throws MailboxException { + public Quota getStorageQuota(QuotaRoot quotaRoot) { Map maxStorageDetails = maxQuotaManager.listMaxStorageDetails(quotaRoot); return Quota.builder() .used(Mono.from(currentQuotaManager.getCurrentStorage(quotaRoot)).block()) @@ -75,7 +76,7 @@ public Quota getStorageQuota(QuotaRoot quotaRoot } @Override - public Quotas getQuotas(QuotaRoot quotaRoot) throws MailboxException { + public Quotas getQuotas(QuotaRoot quotaRoot) { MaxQuotaManager.QuotaDetails quotaDetails = maxQuotaManager.quotaDetails(quotaRoot); CurrentQuotas currentQuotas = Mono.from(currentQuotaManager.getCurrentQuotas(quotaRoot)).block(); return new Quotas( @@ -90,4 +91,22 @@ public Quotas getQuotas(QuotaRoot quotaRoot) throws MailboxException { .limitsByScope(quotaDetails.getMaxStorageDetails()) .build()); } + + @Override + public Publisher getQuotasReactive(QuotaRoot quotaRoot) { + return Mono.zip( + Mono.from(maxQuotaManager.quotaDetailsReactive(quotaRoot)), + Mono.from(currentQuotaManager.getCurrentQuotas(quotaRoot))) + .map(tuple -> new Quotas( + Quota.builder() + .used(tuple.getT2().count()) + .computedLimit(maxQuotaManager.getMaxMessage(tuple.getT1().getMaxMessageDetails()).orElse(QuotaCountLimit.unlimited())) + .limitsByScope(tuple.getT1().getMaxMessageDetails()) + .build(), + Quota.builder() + .used(tuple.getT2().size()) + .computedLimit(maxQuotaManager.getMaxStorage(tuple.getT1().getMaxStorageDetails()).orElse(QuotaSizeLimit.unlimited())) + .limitsByScope(tuple.getT1().getMaxStorageDetails()) + .build())); + } } \ No newline at end of file diff --git a/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/utils/quotas/DefaultQuotaLoader.java b/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/utils/quotas/DefaultQuotaLoader.java index 25a851f854a..e1aa41504e2 100644 --- a/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/utils/quotas/DefaultQuotaLoader.java +++ b/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/utils/quotas/DefaultQuotaLoader.java @@ -21,12 +21,12 @@ import javax.inject.Inject; import org.apache.james.jmap.draft.model.mailbox.Quotas; -import org.apache.james.mailbox.exception.MailboxException; import org.apache.james.mailbox.model.MailboxPath; -import org.apache.james.mailbox.model.QuotaRoot; import org.apache.james.mailbox.quota.QuotaManager; import org.apache.james.mailbox.quota.QuotaRootResolver; +import reactor.core.publisher.Mono; + public class DefaultQuotaLoader extends QuotaLoader { private final QuotaRootResolver quotaRootResolver; @@ -38,15 +38,18 @@ public DefaultQuotaLoader(QuotaRootResolver quotaRootResolver, QuotaManager quot this.quotaManager = quotaManager; } - public Quotas getQuotas(MailboxPath mailboxPath) throws MailboxException { - QuotaRoot quotaRoot = quotaRootResolver.getQuotaRoot(mailboxPath); - Quotas.QuotaId quotaId = Quotas.QuotaId.fromQuotaRoot(quotaRoot); - QuotaManager.Quotas quotas = quotaManager.getQuotas(quotaRoot); - return Quotas.from( - quotaId, - Quotas.Quota.from( - quotaToValue(quotas.getStorageQuota()), - quotaToValue(quotas.getMessageQuota()))); + public Mono getQuotas(MailboxPath mailboxPath) { + return Mono.from(quotaRootResolver.getQuotaRootReactive(mailboxPath)) + .flatMap(quotaRoot -> Mono.from(quotaManager.getQuotasReactive(quotaRoot)) + .map(quotas -> { + Quotas.QuotaId quotaId = Quotas.QuotaId.fromQuotaRoot(quotaRoot); + + return Quotas.from( + quotaId, + Quotas.Quota.from( + quotaToValue(quotas.getStorageQuota()), + quotaToValue(quotas.getMessageQuota()))); + })); } } diff --git a/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/utils/quotas/QuotaLoader.java b/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/utils/quotas/QuotaLoader.java index 76bd3710f63..77a6aea13cc 100644 --- a/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/utils/quotas/QuotaLoader.java +++ b/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/utils/quotas/QuotaLoader.java @@ -24,13 +24,14 @@ import org.apache.james.core.quota.QuotaUsageValue; import org.apache.james.jmap.draft.model.Number; import org.apache.james.jmap.draft.model.mailbox.Quotas; -import org.apache.james.mailbox.exception.MailboxException; import org.apache.james.mailbox.model.MailboxPath; import org.apache.james.mailbox.model.Quota; +import reactor.core.publisher.Mono; + public abstract class QuotaLoader { - public abstract Quotas getQuotas(MailboxPath mailboxPath) throws MailboxException; + public abstract Mono getQuotas(MailboxPath mailboxPath); protected , U extends QuotaUsageValue> Quotas.Value quotaToValue(Quota quota) { return new Quotas.Value<>( diff --git a/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/utils/quotas/QuotaLoaderWithDefaultPreloaded.java b/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/utils/quotas/QuotaLoaderWithDefaultPreloaded.java index 326f0d27993..876dfbdeafb 100644 --- a/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/utils/quotas/QuotaLoaderWithDefaultPreloaded.java +++ b/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/utils/quotas/QuotaLoaderWithDefaultPreloaded.java @@ -22,42 +22,44 @@ import org.apache.james.jmap.draft.model.mailbox.Quotas; import org.apache.james.mailbox.MailboxSession; -import org.apache.james.mailbox.exception.MailboxException; import org.apache.james.mailbox.model.MailboxPath; -import org.apache.james.mailbox.model.QuotaRoot; import org.apache.james.mailbox.quota.QuotaManager; import org.apache.james.mailbox.quota.QuotaRootResolver; +import reactor.core.publisher.Mono; + public class QuotaLoaderWithDefaultPreloaded extends QuotaLoader { + public static Mono preLoad(QuotaRootResolver quotaRootResolver, + QuotaManager quotaManager, + MailboxSession session) { + DefaultQuotaLoader defaultQuotaLoader = new DefaultQuotaLoader(quotaRootResolver, quotaManager); + + return defaultQuotaLoader.getQuotas(MailboxPath.inbox(session)) + .map(Optional::of) + .switchIfEmpty(Mono.just(Optional.empty())) + .map(quotas -> new QuotaLoaderWithDefaultPreloaded(quotaRootResolver, defaultQuotaLoader, quotas)); + } + private final QuotaRootResolver quotaRootResolver; - private final QuotaManager quotaManager; + private final DefaultQuotaLoader defaultQuotaLoader; private final Optional preloadedUserDefaultQuotas; - private final MailboxSession session; - public QuotaLoaderWithDefaultPreloaded(QuotaRootResolver quotaRootResolver, - QuotaManager quotaManager, - MailboxSession session) throws MailboxException { + private QuotaLoaderWithDefaultPreloaded(QuotaRootResolver quotaRootResolver, DefaultQuotaLoader defaultQuotaLoader, Optional preloadedUserDefaultQuotas) { this.quotaRootResolver = quotaRootResolver; - this.quotaManager = quotaManager; - this.session = session; - preloadedUserDefaultQuotas = Optional.of(getUserDefaultQuotas()); - + this.defaultQuotaLoader = defaultQuotaLoader; + this.preloadedUserDefaultQuotas = preloadedUserDefaultQuotas; } - public Quotas getQuotas(MailboxPath mailboxPath) throws MailboxException { - QuotaRoot quotaRoot = quotaRootResolver.getQuotaRoot(mailboxPath); - Quotas.QuotaId quotaId = Quotas.QuotaId.fromQuotaRoot(quotaRoot); - - if (containsQuotaId(preloadedUserDefaultQuotas, quotaId)) { - return preloadedUserDefaultQuotas.get(); - } - QuotaManager.Quotas quotas = quotaManager.getQuotas(quotaRoot); - return Quotas.from( - quotaId, - Quotas.Quota.from( - quotaToValue(quotas.getStorageQuota()), - quotaToValue(quotas.getMessageQuota()))); + public Mono getQuotas(MailboxPath mailboxPath) { + return Mono.from(quotaRootResolver.getQuotaRootReactive(mailboxPath)) + .flatMap(quotaRoot -> { + Quotas.QuotaId quotaId = Quotas.QuotaId.fromQuotaRoot(quotaRoot); + if (containsQuotaId(preloadedUserDefaultQuotas, quotaId)) { + return Mono.just(preloadedUserDefaultQuotas.get()); + } + return defaultQuotaLoader.getQuotas(mailboxPath); + }); } private boolean containsQuotaId(Optional preloadedUserDefaultQuotas, Quotas.QuotaId quotaId) { @@ -67,15 +69,4 @@ private boolean containsQuotaId(Optional preloadedUserDefaultQuotas, Quo .orElse(false); } - private Quotas getUserDefaultQuotas() throws MailboxException { - QuotaRoot quotaRoot = quotaRootResolver.getQuotaRoot(MailboxPath.inbox(session)); - Quotas.QuotaId quotaId = Quotas.QuotaId.fromQuotaRoot(quotaRoot); - QuotaManager.Quotas quotas = quotaManager.getQuotas(quotaRoot); - return Quotas.from( - quotaId, - Quotas.Quota.from( - quotaToValue(quotas.getStorageQuota()), - quotaToValue(quotas.getMessageQuota()))); - } - } diff --git a/server/protocols/jmap-rfc-8621/src/main/scala/org/apache/james/jmap/utils/quotas/QuotaLoader.scala b/server/protocols/jmap-rfc-8621/src/main/scala/org/apache/james/jmap/utils/quotas/QuotaLoader.scala index 0c770c75471..447ca25fef2 100644 --- a/server/protocols/jmap-rfc-8621/src/main/scala/org/apache/james/jmap/utils/quotas/QuotaLoader.scala +++ b/server/protocols/jmap-rfc-8621/src/main/scala/org/apache/james/jmap/utils/quotas/QuotaLoader.scala @@ -20,11 +20,9 @@ package org.apache.james.jmap.utils.quotas import org.apache.james.jmap.mail.Quotas -import org.apache.james.mailbox.exception.MailboxException import org.apache.james.mailbox.model.MailboxPath import reactor.core.scala.publisher.SMono trait QuotaLoader { - @throws[MailboxException] def getQuotas(mailboxPath: MailboxPath): SMono[Quotas] } diff --git a/server/protocols/jmap-rfc-8621/src/main/scala/org/apache/james/jmap/utils/quotas/QuotaLoaderWithPreloadedDefault.scala b/server/protocols/jmap-rfc-8621/src/main/scala/org/apache/james/jmap/utils/quotas/QuotaLoaderWithPreloadedDefault.scala index 8a4c45889d1..6e56928c89d 100644 --- a/server/protocols/jmap-rfc-8621/src/main/scala/org/apache/james/jmap/utils/quotas/QuotaLoaderWithPreloadedDefault.scala +++ b/server/protocols/jmap-rfc-8621/src/main/scala/org/apache/james/jmap/utils/quotas/QuotaLoaderWithPreloadedDefault.scala @@ -22,7 +22,6 @@ package org.apache.james.jmap.utils.quotas import javax.inject.Inject import org.apache.james.jmap.mail.{QuotaRoot, Quotas} import org.apache.james.mailbox.MailboxSession -import org.apache.james.mailbox.exception.MailboxException import org.apache.james.mailbox.model.{MailboxPath, QuotaRoot => ModelQuotaRoot} import org.apache.james.mailbox.quota.UserQuotaRootResolver import reactor.core.scala.publisher.SMono @@ -31,14 +30,13 @@ import reactor.core.scala.publisher.SMono class QuotaLoaderWithPreloadedDefaultFactory @Inject()(quotaRootResolver: UserQuotaRootResolver, quotaReader: QuotaReader) { def loadFor(session: MailboxSession): SMono[QuotaLoaderWithPreloadedDefault] = - SMono.fromCallable(() => new QuotaLoaderWithPreloadedDefault( + getUserDefaultQuotas(session) + .map(qotas => new QuotaLoaderWithPreloadedDefault( quotaRootResolver, quotaReader, session, - getUserDefaultQuotas(session))) + qotas)) - - @throws[MailboxException] private def getUserDefaultQuotas(session:MailboxSession): SMono[Quotas] = { val quotaRoot: ModelQuotaRoot = quotaRootResolver.forUser(session.getUser) quotaReader.retrieveQuotas(QuotaRoot.toJmap(quotaRoot)) @@ -48,11 +46,10 @@ class QuotaLoaderWithPreloadedDefaultFactory @Inject()(quotaRootResolver: UserQu class QuotaLoaderWithPreloadedDefault(quotaRootResolver: UserQuotaRootResolver, quotaReader: QuotaReader, session: MailboxSession, - preloadedUserDefaultQuotas: SMono[Quotas]) extends QuotaLoader { - @throws[MailboxException] + preloadedUserDefaultQuotas: Quotas) extends QuotaLoader { override def getQuotas(mailboxPath: MailboxPath): SMono[Quotas] = if (mailboxPath.belongsTo(session)) { - preloadedUserDefaultQuotas + SMono.just(preloadedUserDefaultQuotas) } else { val quotaRoot: ModelQuotaRoot = quotaRootResolver.getQuotaRoot(mailboxPath) quotaReader.retrieveQuotas(QuotaRoot.toJmap(quotaRoot)) diff --git a/server/protocols/jmap-rfc-8621/src/main/scala/org/apache/james/jmap/utils/quotas/QuotaReader.scala b/server/protocols/jmap-rfc-8621/src/main/scala/org/apache/james/jmap/utils/quotas/QuotaReader.scala index de772456aa5..de646c0478e 100644 --- a/server/protocols/jmap-rfc-8621/src/main/scala/org/apache/james/jmap/utils/quotas/QuotaReader.scala +++ b/server/protocols/jmap-rfc-8621/src/main/scala/org/apache/james/jmap/utils/quotas/QuotaReader.scala @@ -31,15 +31,13 @@ import reactor.core.scala.publisher.SMono class QuotaReader @Inject() (quotaManager: QuotaManager) { @throws[MailboxException] - def retrieveQuotas(quotaRoot: QuotaRoot): SMono[Quotas] = { - val quotaId = QuotaId.fromQuotaRoot(quotaRoot) - val quotas = quotaManager.getQuotas(quotaRoot.toModel) - SMono.just(Quotas.from( - quotaId, - Quota.from(Map( - Quotas.Storage -> quotaToValue(quotas.getStorageQuota), - Quotas.Message -> quotaToValue(quotas.getMessageQuota))))) - } + def retrieveQuotas(quotaRoot: QuotaRoot): SMono[Quotas] = + SMono(quotaManager.getQuotasReactive(quotaRoot.toModel)) + .map(quotas => Quotas.from( + QuotaId.fromQuotaRoot(quotaRoot), + Quota.from(Map( + Quotas.Storage -> quotaToValue(quotas.getStorageQuota), + Quotas.Message -> quotaToValue(quotas.getMessageQuota))))) private def quotaToValue[T <: QuotaLimitValue[T], U <: QuotaUsageValue[U, T]](quota: ModelQuota[T, U]): Value = Value( From 07a5fddf8e800f511ec5adfa68c015fd1e581476 Mon Sep 17 00:00:00 2001 From: Benoit Tellier Date: Fri, 7 May 2021 18:12:28 +0700 Subject: [PATCH 04/15] [REFACTORING] Fasten Mailbox/get & Email/query with reactive MailboxManager::getMailbox --- .../apache/james/mailbox/MailboxManager.java | 4 + .../quota/InMemoryCurrentQuotaManager.java | 14 +- .../mailbox/store/StoreMailboxManager.java | 54 ++++--- .../draft/methods/GetMailboxesMethod.java | 18 +-- .../draft/methods/GetMessageListMethod.java | 9 +- .../SetMailboxesCreationProcessor.java | 3 +- .../SetMailboxesDestructionProcessor.java | 3 +- .../methods/SetMailboxesUpdateProcessor.java | 1 + .../jmap/draft/model/MailboxFactory.java | 138 ++++++++++-------- .../SetMailboxesUpdateProcessorTest.java | 4 +- .../jmap/draft/model/MailboxFactoryTest.java | 31 ++-- .../james/jmap/method/EmailQueryMethod.scala | 6 +- .../data/jmap/EmailQueryViewPopulator.java | 3 +- .../MessageFastViewProjectionCorrector.java | 3 +- 14 files changed, 154 insertions(+), 137 deletions(-) diff --git a/mailbox/api/src/main/java/org/apache/james/mailbox/MailboxManager.java b/mailbox/api/src/main/java/org/apache/james/mailbox/MailboxManager.java index a175127ff2a..6ee804bb263 100644 --- a/mailbox/api/src/main/java/org/apache/james/mailbox/MailboxManager.java +++ b/mailbox/api/src/main/java/org/apache/james/mailbox/MailboxManager.java @@ -139,6 +139,10 @@ enum SearchCapabilities { */ MessageManager getMailbox(MailboxId mailboxId, MailboxSession session) throws MailboxException; + Publisher getMailboxReactive(MailboxId mailboxId, MailboxSession session); + + Publisher getMailboxReactive(MailboxPath mailboxPath, MailboxSession session); + /** * Creates a new mailbox. Any intermediary mailboxes missing from the * hierarchy should be created. diff --git a/mailbox/memory/src/main/java/org/apache/james/mailbox/inmemory/quota/InMemoryCurrentQuotaManager.java b/mailbox/memory/src/main/java/org/apache/james/mailbox/inmemory/quota/InMemoryCurrentQuotaManager.java index 120df80cd2e..26b35dc1b2d 100644 --- a/mailbox/memory/src/main/java/org/apache/james/mailbox/inmemory/quota/InMemoryCurrentQuotaManager.java +++ b/mailbox/memory/src/main/java/org/apache/james/mailbox/inmemory/quota/InMemoryCurrentQuotaManager.java @@ -60,7 +60,6 @@ public AtomicReference load(QuotaRoot quotaRoot) { public CurrentQuotas loadQuotas(QuotaRoot quotaRoot, CurrentQuotaCalculator quotaCalculator, SessionProvider sessionProvider) { return quotaCalculator.recalculateCurrentQuotas(quotaRoot, sessionProvider.createSystemSession(Username.of(quotaRoot.getValue()))) - .subscribeOn(Schedulers.elastic()) .block(); } @@ -77,22 +76,20 @@ public Mono decrease(QuotaOperation quotaOperation) { @Override public Mono getCurrentMessageCount(QuotaRoot quotaRoot) { return Mono.fromCallable(() -> quotaCache.get(quotaRoot).get().count()) - .onErrorMap(this::wrapAsMailboxException) - .subscribeOn(Schedulers.elastic()); + .onErrorMap(this::wrapAsMailboxException); } @Override public Mono getCurrentStorage(QuotaRoot quotaRoot) { return Mono.fromCallable(() -> quotaCache.get(quotaRoot).get().size()) - .onErrorMap(this::wrapAsMailboxException) - .subscribeOn(Schedulers.elastic()); + .onErrorMap(this::wrapAsMailboxException); } @Override public Mono getCurrentQuotas(QuotaRoot quotaRoot) { return Mono.fromCallable(() -> quotaCache.get(quotaRoot).get()) - .onErrorMap(this::wrapAsMailboxException) - .subscribeOn(Schedulers.elastic()); + .subscribeOn(Schedulers.elastic()) + .onErrorMap(this::wrapAsMailboxException); } @Override @@ -100,8 +97,7 @@ public Mono setCurrentQuotas(QuotaOperation quotaOperation) { return getCurrentQuotas(quotaOperation.quotaRoot()) .filter(Predicate.not(Predicate.isEqual(CurrentQuotas.from(quotaOperation)))) .flatMap(storedQuotas -> decrease(new QuotaOperation(quotaOperation.quotaRoot(), storedQuotas.count(), storedQuotas.size())) - .then(increase(quotaOperation))) - .subscribeOn(Schedulers.elastic()); + .then(increase(quotaOperation))); } private Mono updateQuota(QuotaRoot quotaRoot, UnaryOperator quotaFunction) { diff --git a/mailbox/store/src/main/java/org/apache/james/mailbox/store/StoreMailboxManager.java b/mailbox/store/src/main/java/org/apache/james/mailbox/store/StoreMailboxManager.java index ceabd4d44e6..683727214da 100644 --- a/mailbox/store/src/main/java/org/apache/james/mailbox/store/StoreMailboxManager.java +++ b/mailbox/store/src/main/java/org/apache/james/mailbox/store/StoreMailboxManager.java @@ -87,6 +87,7 @@ import org.apache.james.mailbox.store.user.SubscriptionMapper; import org.apache.james.mailbox.store.user.model.Subscription; import org.apache.james.util.FunctionalUtils; +import org.reactivestreams.Publisher; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -265,37 +266,50 @@ protected StoreMessageManager createMessageManager(Mailbox mailbox, MailboxSessi @Override public MessageManager getMailbox(MailboxPath mailboxPath, MailboxSession session) throws MailboxException { - final MailboxMapper mapper = mailboxSessionMapperFactory.getMailboxMapper(session); - Mailbox mailboxRow = mapper.findMailboxByPath(mailboxPath) - .blockOptional() - .orElseThrow(() -> { - LOGGER.info("Mailbox '{}' not found.", mailboxPath); - return new MailboxNotFoundException(mailboxPath); - }); + return MailboxReactorUtils.block(getMailboxReactive(mailboxPath, session)); + } - if (!assertUserHasAccessTo(mailboxRow, session)) { - LOGGER.info("Mailbox '{}' does not belong to user '{}' but to '{}'", mailboxPath, session.getUser(), mailboxRow.getUser()); - throw new MailboxNotFoundException(mailboxPath); - } + @Override + public Mono getMailboxReactive(MailboxPath mailboxPath, MailboxSession session) { + MailboxMapper mapper = mailboxSessionMapperFactory.getMailboxMapper(session); + + return mapper.findMailboxByPath(mailboxPath) + .map(Throwing.function(mailboxRow -> { + if (!assertUserHasAccessTo(mailboxRow, session)) { + LOGGER.info("Mailbox '{}' does not belong to user '{}' but to '{}'", mailboxPath, session.getUser(), mailboxRow.getUser()); + throw new MailboxNotFoundException(mailboxPath); + } - LOGGER.debug("Loaded mailbox {}", mailboxPath); + LOGGER.debug("Loaded mailbox {}", mailboxPath); - return createMessageManager(mailboxRow, session); + return createMessageManager(mailboxRow, session); + }).sneakyThrow()) + .switchIfEmpty(Mono.fromCallable(() -> { + LOGGER.info("Mailbox '{}' not found.", mailboxPath); + throw new MailboxNotFoundException(mailboxPath); + })); } @Override public MessageManager getMailbox(MailboxId mailboxId, MailboxSession session) throws MailboxException { + return block(getMailboxReactive(mailboxId, session)); + } + + @Override + public Publisher getMailboxReactive(MailboxId mailboxId, MailboxSession session) { MailboxMapper mapper = mailboxSessionMapperFactory.getMailboxMapper(session); - Mailbox mailboxRow = block(mapper.findMailboxById(mailboxId)); - if (!assertUserHasAccessTo(mailboxRow, session)) { - LOGGER.info("Mailbox '{}' does not belong to user '{}' but to '{}'", mailboxId.serialize(), session.getUser(), mailboxRow.getUser()); - throw new MailboxNotFoundException(mailboxId); - } + return mapper.findMailboxById(mailboxId) + .map(Throwing.function(mailboxRow -> { + if (!assertUserHasAccessTo(mailboxRow, session)) { + LOGGER.info("Mailbox '{} {}' does not belong to user '{}' but to '{}'", mailboxRow.getMailboxId().serialize(), mailboxRow.generateAssociatedPath(), session.getUser(), mailboxRow.getUser()); + throw new MailboxNotFoundException(mailboxId); + } - LOGGER.debug("Loaded mailbox {}", mailboxId.serialize()); + LOGGER.debug("Loaded mailbox {} {}", mailboxRow.getMailboxId().serialize(), mailboxRow.generateAssociatedPath()); - return createMessageManager(mailboxRow, session); + return createMessageManager(mailboxRow, session); + }).sneakyThrow()); } private boolean assertUserHasAccessTo(Mailbox mailbox, MailboxSession session) { diff --git a/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/methods/GetMailboxesMethod.java b/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/methods/GetMailboxesMethod.java index 610924ff938..8381fdc8b3a 100644 --- a/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/methods/GetMailboxesMethod.java +++ b/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/methods/GetMailboxesMethod.java @@ -21,7 +21,6 @@ import static org.apache.james.util.ReactorUtils.DEFAULT_CONCURRENCY; import static org.apache.james.util.ReactorUtils.context; -import static org.apache.james.util.ReactorUtils.publishIfPresent; import java.util.Comparator; import java.util.List; @@ -57,7 +56,6 @@ import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; -import reactor.core.scheduler.Schedulers; public class GetMailboxesMethod implements Method { @@ -143,32 +141,26 @@ private Flux retrieveMailboxes(Optional> mailb private Flux retrieveSpecificMailboxes(MailboxSession mailboxSession, ImmutableList mailboxIds) { return Flux.fromIterable(mailboxIds) - .flatMap(mailboxId -> Mono.fromCallable(() -> - mailboxFactory.builder() + .flatMap(mailboxId -> mailboxFactory.builder() .id(mailboxId) .session(mailboxSession) .usingPreloadedMailboxesMetadata(NO_PRELOADED_METADATA) - .build()) - .subscribeOn(Schedulers.elastic()), DEFAULT_CONCURRENCY) - .handle(publishIfPresent()); + .build(), DEFAULT_CONCURRENCY); } private Flux retrieveAllMailboxes(MailboxSession mailboxSession) { Mono> userMailboxesMono = getAllMailboxesMetaData(mailboxSession).collectList(); - Mono quotaLoaderMono = Mono.fromCallable(() -> - new QuotaLoaderWithDefaultPreloaded(quotaRootResolver, quotaManager, mailboxSession)) - .subscribeOn(Schedulers.elastic()); + Mono quotaLoaderMono = QuotaLoaderWithDefaultPreloaded.preLoad(quotaRootResolver, quotaManager, mailboxSession); return userMailboxesMono.zipWith(quotaLoaderMono) .flatMapMany( tuple -> Flux.fromIterable(tuple.getT1()) - .map(mailboxMetaData -> mailboxFactory.builder() + .flatMap(mailboxMetaData -> mailboxFactory.builder() .mailboxMetadata(mailboxMetaData) .session(mailboxSession) .usingPreloadedMailboxesMetadata(Optional.of(tuple.getT1())) .quotaLoader(tuple.getT2()) - .build()) - .handle(publishIfPresent())); + .build())); } private Flux getAllMailboxesMetaData(MailboxSession mailboxSession) { diff --git a/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/methods/GetMessageListMethod.java b/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/methods/GetMessageListMethod.java index 69d64d5eabb..caad3960780 100644 --- a/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/methods/GetMessageListMethod.java +++ b/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/methods/GetMessageListMethod.java @@ -115,8 +115,7 @@ public Flux process(JmapRequest request, MethodCallId methodCallId return metricFactory.decorateSupplierWithTimerMetricLogP99(JMAP_PREFIX + METHOD_NAME.getName(), () -> process(methodCallId, mailboxSession, messageListRequest) - .subscriberContext(context("GET_MESSAGE_LIST", mdc(messageListRequest)))) - .subscribeOn(Schedulers.elastic()); + .subscriberContext(context("GET_MESSAGE_LIST", mdc(messageListRequest)))); } private MDCBuilder mdc(GetMessageListRequest messageListRequest) { @@ -171,8 +170,7 @@ private Mono getMessageListResponse(GetMessageListReques MailboxId mailboxId = mailboxIdFactory.fromString(mailboxIdAsString); Limit aLimit = Limit.from(Math.toIntExact(limit)); - return Mono.fromCallable(() -> mailboxManager.getMailbox(mailboxId, mailboxSession)) - .subscribeOn(Schedulers.elastic()) + return Mono.from(mailboxManager.getMailboxReactive(mailboxId, mailboxSession)) .then(emailQueryView.listMailboxContent(mailboxId, aLimit) .skip(position) .take(limit) @@ -188,8 +186,7 @@ private Mono getMessageListResponse(GetMessageListReques ZonedDateTime after = condition.getAfter().get(); Limit aLimit = Limit.from(Math.toIntExact(limit)); - return Mono.fromCallable(() -> mailboxManager.getMailbox(mailboxId, mailboxSession)) - .subscribeOn(Schedulers.elastic()) + return Mono.from(mailboxManager.getMailboxReactive(mailboxId, mailboxSession)) .then(emailQueryView.listMailboxContentSinceReceivedAt(mailboxId, after, aLimit) .skip(position) .take(limit) diff --git a/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/methods/SetMailboxesCreationProcessor.java b/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/methods/SetMailboxesCreationProcessor.java index c75b462a202..25a1a69a00d 100644 --- a/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/methods/SetMailboxesCreationProcessor.java +++ b/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/methods/SetMailboxesCreationProcessor.java @@ -121,7 +121,8 @@ private void createMailbox(MailboxCreationId mailboxCreationId, MailboxCreateReq Optional mailbox = mailboxId.flatMap(id -> mailboxFactory.builder() .id(id) .session(mailboxSession) - .build()); + .build() + .blockOptional()); if (mailbox.isPresent()) { subscriptionManager.subscribe(mailboxSession, mailboxPath.getName()); builder.created(mailboxCreationId, mailbox.get()); diff --git a/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/methods/SetMailboxesDestructionProcessor.java b/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/methods/SetMailboxesDestructionProcessor.java index 1238147e54d..dbd2e1e4b84 100644 --- a/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/methods/SetMailboxesDestructionProcessor.java +++ b/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/methods/SetMailboxesDestructionProcessor.java @@ -98,7 +98,8 @@ private ImmutableMap mapDestroyRequests(SetMailboxesRequest .map(id -> mailboxFactory.builder() .id(id) .session(mailboxSession) - .build()) + .build() + .blockOptional()) .flatMap(Optional::stream) .forEach(mailbox -> idToMailboxBuilder.put(mailbox.getId(), mailbox)); return idToMailboxBuilder.build(); diff --git a/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/methods/SetMailboxesUpdateProcessor.java b/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/methods/SetMailboxesUpdateProcessor.java index acf274752a1..33aa98ea68b 100644 --- a/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/methods/SetMailboxesUpdateProcessor.java +++ b/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/methods/SetMailboxesUpdateProcessor.java @@ -181,6 +181,7 @@ private Mailbox getMailbox(MailboxId mailboxId, MailboxSession mailboxSession) t .id(mailboxId) .session(mailboxSession) .build() + .blockOptional() .orElseThrow(() -> new MailboxNotFoundException(mailboxId)); } diff --git a/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/model/MailboxFactory.java b/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/model/MailboxFactory.java index d4597b7e6ce..7f3f723d7dd 100644 --- a/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/model/MailboxFactory.java +++ b/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/model/MailboxFactory.java @@ -26,7 +26,6 @@ import org.apache.james.core.Username; import org.apache.james.jmap.draft.model.mailbox.Mailbox; import org.apache.james.jmap.draft.model.mailbox.MailboxNamespace; -import org.apache.james.jmap.draft.model.mailbox.Quotas; import org.apache.james.jmap.draft.model.mailbox.Rights; import org.apache.james.jmap.draft.model.mailbox.SortOrder; import org.apache.james.jmap.draft.utils.quotas.DefaultQuotaLoader; @@ -35,7 +34,6 @@ import org.apache.james.mailbox.MailboxSession; import org.apache.james.mailbox.MessageManager; import org.apache.james.mailbox.Role; -import org.apache.james.mailbox.exception.MailboxException; import org.apache.james.mailbox.exception.MailboxNotFoundException; import org.apache.james.mailbox.model.MailboxACL; import org.apache.james.mailbox.model.MailboxCounters; @@ -52,8 +50,31 @@ import com.google.common.primitives.Booleans; import reactor.core.publisher.Mono; +import reactor.core.scheduler.Schedulers; public class MailboxFactory { + private static class MailboxTuple { + public static MailboxTuple from(MailboxMetaData metaData) { + return new MailboxTuple(metaData.getPath(), metaData.getCounters().sanitize(), metaData.getResolvedAcls()); + } + + public static Mono from(MessageManager messageManager, MailboxSession session) { + return Mono.fromCallable(() -> + new MailboxTuple(messageManager.getMailboxPath(), messageManager.getMailboxCounters(session).sanitize(), messageManager.getResolvedAcl(session))) + .subscribeOn(Schedulers.elastic()); + } + + private final MailboxPath mailboxPath; + private final MailboxCounters.Sanitized mailboxCounters; + private final MailboxACL acl; + + private MailboxTuple(MailboxPath mailboxPath, MailboxCounters.Sanitized mailboxCounters, MailboxACL acl) { + this.mailboxPath = mailboxPath; + this.mailboxCounters = mailboxCounters; + this.acl = acl; + } + } + private final MailboxManager mailboxManager; private final QuotaManager quotaManager; private final QuotaRootResolver quotaRootResolver; @@ -96,36 +117,25 @@ public MailboxBuilder usingPreloadedMailboxesMetadata(Optional build() { + public Mono build() { Preconditions.checkNotNull(session); - try { - MailboxId mailboxId = computeMailboxId(); - Mono mailbox = mailbox(mailboxId).cache(); - - MailboxACL mailboxACL = mailboxMetaData.map(MailboxMetaData::getResolvedAcls) - .orElseGet(Throwing.supplier(() -> retrieveCachedMailbox(mailboxId, mailbox).getResolvedAcl(session)).sneakyThrow()); - - MailboxPath mailboxPath = mailboxMetaData.map(MailboxMetaData::getPath) - .orElseGet(Throwing.supplier(() -> retrieveCachedMailbox(mailboxId, mailbox).getMailboxPath()).sneakyThrow()); - - MailboxCounters.Sanitized mailboxCounters = mailboxMetaData.map(MailboxMetaData::getCounters) - .orElseGet(Throwing.supplier(() -> retrieveCachedMailbox(mailboxId, mailbox).getMailboxCounters(session)).sneakyThrow()) - .sanitize(); - - return Optional.of(mailboxFactory.from( - mailboxId, - mailboxPath, - mailboxCounters, - mailboxACL, - userMailboxesMetadata, - quotaLoader, - session)); - } catch (MailboxNotFoundException e) { - return Optional.empty(); - } catch (MailboxException e) { - throw new RuntimeException(e); - } + MailboxId mailboxId = computeMailboxId(); + + Mono mailbox = mailboxMetaData.map(MailboxTuple::from) + .map(Mono::just) + .orElse(Mono.from(mailboxFactory.mailboxManager.getMailboxReactive(mailboxId, session)) + .flatMap(messageManager -> MailboxTuple.from(messageManager, session))); + + return mailbox.flatMap(tuple -> mailboxFactory.from( + mailboxId, + tuple.mailboxPath, + tuple.mailboxCounters, + tuple.acl, + userMailboxesMetadata, + quotaLoader, + session)) + .onErrorResume(MailboxNotFoundException.class, e -> Mono.empty()); } private MailboxId computeMailboxId() { @@ -137,14 +147,13 @@ private MailboxId computeMailboxId() { } private Mono mailbox(MailboxId mailboxId) { - return Mono.fromCallable(() -> mailboxFactory.mailboxManager.getMailbox(mailboxId, session)); + return Mono.from(mailboxFactory.mailboxManager.getMailboxReactive(mailboxId, session)); } - private MessageManager retrieveCachedMailbox(MailboxId mailboxId, Mono mailbox) throws MailboxNotFoundException { + private Mono retrieveCachedMailbox(MailboxId mailboxId, Mono mailbox) throws MailboxNotFoundException { return mailbox .onErrorResume(MailboxNotFoundException.class, any -> Mono.empty()) - .blockOptional() - .orElseThrow(() -> new MailboxNotFoundException(mailboxId)); + .switchIfEmpty(Mono.error(() -> new MailboxNotFoundException(mailboxId))); } } @@ -160,13 +169,13 @@ public MailboxBuilder builder() { return new MailboxBuilder(this, defaultQuotaLoader); } - private Mailbox from(MailboxId mailboxId, + private Mono from(MailboxId mailboxId, MailboxPath mailboxPath, MailboxCounters.Sanitized mailboxCounters, MailboxACL resolvedAcl, Optional> userMailboxesMetadata, QuotaLoader quotaLoader, - MailboxSession mailboxSession) throws MailboxException { + MailboxSession mailboxSession) { boolean isOwner = mailboxPath.belongsTo(mailboxSession); Optional role = Role.from(mailboxPath.getName()) .filter(any -> mailboxPath.belongsTo(mailboxSession)); @@ -175,26 +184,27 @@ private Mailbox from(MailboxId mailboxId, .removeEntriesFor(mailboxPath.getUser()); Username username = mailboxSession.getUser(); - Quotas quotas = quotaLoader.getQuotas(mailboxPath); - - return Mailbox.builder() - .id(mailboxId) - .name(getName(mailboxPath, mailboxSession)) - .parentId(getParentIdFromMailboxPath(mailboxPath, userMailboxesMetadata, mailboxSession).orElse(null)) - .role(role) - .unreadMessages(mailboxCounters.getUnseen()) - .totalMessages(mailboxCounters.getCount()) - .sortOrder(SortOrder.getSortOrder(role)) - .sharedWith(rights) - .mayAddItems(rights.mayAddItems(username).orElse(isOwner)) - .mayCreateChild(rights.mayCreateChild(username).orElse(isOwner)) - .mayDelete(rights.mayDelete(username).orElse(isOwner)) - .mayReadItems(rights.mayReadItems(username).orElse(isOwner)) - .mayRemoveItems(rights.mayRemoveItems(username).orElse(isOwner)) - .mayRename(rights.mayRename(username).orElse(isOwner)) - .namespace(getNamespace(mailboxPath, isOwner)) - .quotas(quotas) - .build(); + return Mono.zip( + quotaLoader.getQuotas(mailboxPath), + getParentIdFromMailboxPath(mailboxPath, userMailboxesMetadata, mailboxSession)) + .map(tuple -> Mailbox.builder() + .id(mailboxId) + .name(getName(mailboxPath, mailboxSession)) + .parentId(tuple.getT2().orElse(null)) + .role(role) + .unreadMessages(mailboxCounters.getUnseen()) + .totalMessages(mailboxCounters.getCount()) + .sortOrder(SortOrder.getSortOrder(role)) + .sharedWith(rights) + .mayAddItems(rights.mayAddItems(username).orElse(isOwner)) + .mayCreateChild(rights.mayCreateChild(username).orElse(isOwner)) + .mayDelete(rights.mayDelete(username).orElse(isOwner)) + .mayReadItems(rights.mayReadItems(username).orElse(isOwner)) + .mayRemoveItems(rights.mayRemoveItems(username).orElse(isOwner)) + .mayRename(rights.mayRename(username).orElse(isOwner)) + .namespace(getNamespace(mailboxPath, isOwner)) + .quotas(tuple.getT1()) + .build()); } private MailboxNamespace getNamespace(MailboxPath mailboxPath, boolean isOwner) { @@ -215,21 +225,21 @@ String getName(MailboxPath mailboxPath, MailboxSession mailboxSession) { } @VisibleForTesting - Optional getParentIdFromMailboxPath(MailboxPath mailboxPath, Optional> userMailboxesMetadata, - MailboxSession mailboxSession) throws MailboxException { + Mono> getParentIdFromMailboxPath(MailboxPath mailboxPath, Optional> userMailboxesMetadata, + MailboxSession mailboxSession) { List levels = mailboxPath.getHierarchyLevels(mailboxSession.getPathDelimiter()); if (levels.size() <= 1) { - return Optional.empty(); + return Mono.just(Optional.empty()); } MailboxPath parent = levels.get(levels.size() - 2); - return userMailboxesMetadata.map(list -> retrieveParentFromMetadata(parent, list)) + return userMailboxesMetadata.map(list -> Mono.just(retrieveParentFromMetadata(parent, list))) .orElseGet(Throwing.supplier(() -> retrieveParentFromBackend(mailboxSession, parent)).sneakyThrow()); } - private Optional retrieveParentFromBackend(MailboxSession mailboxSession, MailboxPath parent) throws MailboxException { - return Optional.of( - mailboxManager.getMailbox(parent, mailboxSession) - .getId()); + private Mono> retrieveParentFromBackend(MailboxSession mailboxSession, MailboxPath parent) { + return Mono.from(mailboxManager.getMailboxReactive(parent, mailboxSession)) + .map(MessageManager::getId) + .map(Optional::of); } private Optional retrieveParentFromMetadata(MailboxPath parent, List list) { diff --git a/server/protocols/jmap-draft/src/test/java/org/apache/james/jmap/draft/methods/SetMailboxesUpdateProcessorTest.java b/server/protocols/jmap-draft/src/test/java/org/apache/james/jmap/draft/methods/SetMailboxesUpdateProcessorTest.java index 79859c44909..7faa903cbbe 100644 --- a/server/protocols/jmap-draft/src/test/java/org/apache/james/jmap/draft/methods/SetMailboxesUpdateProcessorTest.java +++ b/server/protocols/jmap-draft/src/test/java/org/apache/james/jmap/draft/methods/SetMailboxesUpdateProcessorTest.java @@ -45,6 +45,8 @@ import org.junit.Test; import org.mockito.Mockito; +import reactor.core.publisher.Mono; + public class SetMailboxesUpdateProcessorTest { private MailboxManager mockedMailboxManager; @@ -79,7 +81,7 @@ public void processShouldReturnNotUpdatedWhenMailboxExceptionOccured() throws Ex when(mockBuilder.session(mockedMailboxSession)) .thenReturn(mockBuilder); when(mockBuilder.build()) - .thenReturn(Optional.of(mailbox)); + .thenReturn(Mono.just(mailbox)); when(mockedMailboxFactory.builder()) .thenReturn(mockBuilder); diff --git a/server/protocols/jmap-draft/src/test/java/org/apache/james/jmap/draft/model/MailboxFactoryTest.java b/server/protocols/jmap-draft/src/test/java/org/apache/james/jmap/draft/model/MailboxFactoryTest.java index 22e6c8f3844..e2ee24157af 100644 --- a/server/protocols/jmap-draft/src/test/java/org/apache/james/jmap/draft/model/MailboxFactoryTest.java +++ b/server/protocols/jmap-draft/src/test/java/org/apache/james/jmap/draft/model/MailboxFactoryTest.java @@ -82,7 +82,8 @@ public void mailboxFromMailboxIdShouldReturnAbsentWhenDoesntExist() throws Excep Optional mailbox = sut.builder() .id(InMemoryId.of(123)) .session(mailboxSession) - .build(); + .build() + .blockOptional(); assertThat(mailbox).isEmpty(); } @@ -96,7 +97,8 @@ public void mailboxFromMailboxIdShouldReturnPresentWhenExists() throws Exception Optional mailbox = sut.builder() .id(mailboxId) .session(mailboxSession) - .build(); + .build() + .blockOptional(); assertThat(mailbox).isPresent(); assertThat(mailbox.get().getId()).isEqualTo(mailboxId); @@ -134,7 +136,7 @@ public void getParentIdFromMailboxPathShouldReturNullWhenRootMailbox() throws Ex MailboxPath mailboxPath = MailboxPath.forUser(user, "mailbox"); mailboxManager.createMailbox(mailboxPath, mailboxSession); - Optional id = sut.getParentIdFromMailboxPath(mailboxPath, Optional.empty(), mailboxSession); + Optional id = sut.getParentIdFromMailboxPath(mailboxPath, Optional.empty(), mailboxSession).block(); assertThat(id).isEmpty(); } @@ -147,7 +149,7 @@ public void getParentIdFromMailboxPathShouldReturnParentIdWhenChildMailbox() thr MailboxPath mailboxPath = parentMailboxPath.child("mailbox", '.'); mailboxManager.createMailbox(mailboxPath, mailboxSession); - Optional id = sut.getParentIdFromMailboxPath(mailboxPath, Optional.empty(), mailboxSession); + Optional id = sut.getParentIdFromMailboxPath(mailboxPath, Optional.empty(), mailboxSession).block(); assertThat(id).contains(parentId); } @@ -162,7 +164,7 @@ public void getParentIdFromMailboxPathShouldReturnParentIdWhenChildOfChildMailbo mailboxManager.createMailbox(mailboxPath, mailboxSession); - Optional id = sut.getParentIdFromMailboxPath(mailboxPath, Optional.empty(), mailboxSession); + Optional id = sut.getParentIdFromMailboxPath(mailboxPath, Optional.empty(), mailboxSession).block(); assertThat(id).contains(parentId); } @@ -185,7 +187,7 @@ MailboxMetaData.Children.CHILDREN_ALLOWED_BUT_UNKNOWN, MailboxMetaData.Selectabi .count(0) .unseen(0) .build()))), - mailboxSession); + mailboxSession).block(); assertThat(id).contains(parentId); } @@ -198,7 +200,7 @@ public void getNamespaceShouldReturnPersonalNamespaceWhenUserMailboxPathAndUserM .id(mailboxId.get()) .session(mailboxSession) .build() - .get(); + .block(); assertThat(retrievedMailbox.getNamespace()) .isEqualTo(MailboxNamespace.personal()); @@ -227,7 +229,7 @@ public void buildShouldRelyOnPreloadedMailboxes() throws Exception { .unseen(0) .build())))) .build() - .get(); + .block(); assertThat(retrievedMailbox.getParentId()) .contains(preLoadedId); @@ -248,7 +250,7 @@ public void getNamespaceShouldReturnDelegatedNamespaceWhenUserMailboxPathAndUser .id(mailboxId.get()) .session(otherMailboxSession) .build() - .get(); + .block(); assertThat(retrievedMailbox.getNamespace()) .isEqualTo(MailboxNamespace.delegated(user)); @@ -263,7 +265,7 @@ public void ownerShouldHaveFullRightsViaMayProperties() throws Exception { .id(mailboxId.get()) .session(mailboxSession) .build() - .get(); + .block(); softly.assertThat(retrievedMailbox.isMayAddItems()).isTrue(); softly.assertThat(retrievedMailbox.isMayCreateChild()).isTrue(); @@ -288,7 +290,7 @@ public void delegatedUserShouldHaveMayAddItemsWhenAllowedToInsert() throws Excep .id(mailboxId.get()) .session(otherMailboxSession) .build() - .get(); + .block(); softly.assertThat(retrievedMailbox.isMayAddItems()).isTrue(); softly.assertThat(retrievedMailbox.isMayCreateChild()).isFalse(); @@ -313,7 +315,7 @@ public void delegatedUserShouldHaveMayReadItemsWhenAllowedToRead() throws Except .id(mailboxId.get()) .session(otherMailboxSession) .build() - .get(); + .block(); softly.assertThat(retrievedMailbox.isMayAddItems()).isFalse(); softly.assertThat(retrievedMailbox.isMayCreateChild()).isFalse(); @@ -338,7 +340,7 @@ public void delegatedUserShouldHaveMayRemoveItemsWhenAllowedToRemoveItems() thro .id(mailboxId.get()) .session(otherMailboxSession) .build() - .get(); + .block(); softly.assertThat(retrievedMailbox.isMayAddItems()).isFalse(); softly.assertThat(retrievedMailbox.isMayCreateChild()).isFalse(); @@ -367,7 +369,8 @@ public void mailboxFromMetaDataShouldReturnPresentStoredValue() throws Exception Optional mailbox = sut.builder() .mailboxMetadata(metaData) .session(mailboxSession) - .build(); + .build() + .blockOptional(); softly.assertThat(mailbox).isPresent(); softly.assertThat(mailbox).map(Mailbox::getId).contains(metaData.getId()); diff --git a/server/protocols/jmap-rfc-8621/src/main/scala/org/apache/james/jmap/method/EmailQueryMethod.scala b/server/protocols/jmap-rfc-8621/src/main/scala/org/apache/james/jmap/method/EmailQueryMethod.scala index ac2e32d8ed8..6f08d44edb0 100644 --- a/server/protocols/jmap-rfc-8621/src/main/scala/org/apache/james/jmap/method/EmailQueryMethod.scala +++ b/server/protocols/jmap-rfc-8621/src/main/scala/org/apache/james/jmap/method/EmailQueryMethod.scala @@ -107,8 +107,7 @@ class EmailQueryMethod @Inject() (serializer: EmailQuerySerializer, val condition: FilterCondition = request.filter.get.asInstanceOf[FilterCondition] val mailboxId: MailboxId = condition.inMailbox.get val after: ZonedDateTime = condition.after.get.asUTC - SMono.fromCallable(() => mailboxManager.getMailbox(mailboxId, mailboxSession)) - .subscribeOn(Schedulers.elastic()) + SMono(mailboxManager.getMailboxReactive(mailboxId, mailboxSession)) .`then`(SFlux.fromPublisher( emailQueryView.listMailboxContentSinceReceivedAt(mailboxId, after, JavaLimit.from(limitToUse.value))) .drop(position.value) @@ -122,8 +121,7 @@ class EmailQueryMethod @Inject() (serializer: EmailQuerySerializer, private def queryViewForListingSortedBySentAt(mailboxSession: MailboxSession, position: Position, limitToUse: Limit, request: EmailQueryRequest): SMono[Seq[MessageId]] = { val mailboxId: MailboxId = request.filter.get.asInstanceOf[FilterCondition].inMailbox.get - SMono.fromCallable(() => mailboxManager.getMailbox(mailboxId, mailboxSession)) - .subscribeOn(Schedulers.elastic()) + SMono(mailboxManager.getMailboxReactive(mailboxId, mailboxSession)) .`then`(SFlux.fromPublisher( emailQueryView.listMailboxContent(mailboxId, JavaLimit.from(limitToUse.value))) .drop(position.value) diff --git a/server/protocols/webadmin/webadmin-jmap/src/main/java/org/apache/james/webadmin/data/jmap/EmailQueryViewPopulator.java b/server/protocols/webadmin/webadmin-jmap/src/main/java/org/apache/james/webadmin/data/jmap/EmailQueryViewPopulator.java index 2ea9f167a9c..7a77db56617 100644 --- a/server/protocols/webadmin/webadmin-jmap/src/main/java/org/apache/james/webadmin/data/jmap/EmailQueryViewPopulator.java +++ b/server/protocols/webadmin/webadmin-jmap/src/main/java/org/apache/james/webadmin/data/jmap/EmailQueryViewPopulator.java @@ -189,8 +189,7 @@ private Flux listUsersMailboxes(MailboxSession session) { } private Mono retrieveMailbox(MailboxSession session, MailboxMetaData mailboxMetadata) { - return Mono.fromCallable(() -> mailboxManager.getMailbox(mailboxMetadata.getId(), session)) - .subscribeOn(Schedulers.elastic()); + return Mono.from(mailboxManager.getMailboxReactive(mailboxMetadata.getId(), session)); } private Flux listAllMessages(MessageManager messageManager, MailboxSession session) { diff --git a/server/protocols/webadmin/webadmin-jmap/src/main/java/org/apache/james/webadmin/data/jmap/MessageFastViewProjectionCorrector.java b/server/protocols/webadmin/webadmin-jmap/src/main/java/org/apache/james/webadmin/data/jmap/MessageFastViewProjectionCorrector.java index da0a6d70016..054f85d428f 100644 --- a/server/protocols/webadmin/webadmin-jmap/src/main/java/org/apache/james/webadmin/data/jmap/MessageFastViewProjectionCorrector.java +++ b/server/protocols/webadmin/webadmin-jmap/src/main/java/org/apache/james/webadmin/data/jmap/MessageFastViewProjectionCorrector.java @@ -211,8 +211,7 @@ private Flux listUsersMailboxes(MailboxSession session) { } private Mono retrieveMailbox(MailboxSession session, MailboxMetaData mailboxMetadata) { - return Mono.fromCallable(() -> mailboxManager.getMailbox(mailboxMetadata.getId(), session)) - .subscribeOn(Schedulers.elastic()); + return Mono.from(mailboxManager.getMailboxReactive(mailboxMetadata.getId(), session)); } private Flux listAllMailboxMessages(MessageManager messageManager, MailboxSession session) { From 2be07f81813a71d2a1831c8dfc819ada21427a17 Mon Sep 17 00:00:00 2001 From: Benoit Tellier Date: Fri, 7 May 2021 18:39:37 +0700 Subject: [PATCH 05/15] [REFACTORING] Wrap blocking calls in JWTAuthenticationStrategy It was forcing to schedule the whole chain on an elastic scheduler... --- .../org/apache/james/jmap/jwt/JWTAuthenticationStrategy.java | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/server/protocols/jmap/src/main/java/org/apache/james/jmap/jwt/JWTAuthenticationStrategy.java b/server/protocols/jmap/src/main/java/org/apache/james/jmap/jwt/JWTAuthenticationStrategy.java index 4314c31d1ac..55352af550e 100644 --- a/server/protocols/jmap/src/main/java/org/apache/james/jmap/jwt/JWTAuthenticationStrategy.java +++ b/server/protocols/jmap/src/main/java/org/apache/james/jmap/jwt/JWTAuthenticationStrategy.java @@ -36,6 +36,7 @@ import com.google.common.collect.ImmutableMap; import reactor.core.publisher.Mono; +import reactor.core.scheduler.Schedulers; import reactor.netty.http.server.HttpServerRequest; public class JWTAuthenticationStrategy implements AuthenticationStrategy { @@ -61,7 +62,7 @@ public Mono createMailboxSession(HttpServerRequest httpRequest) return Mono.fromCallable(() -> authHeaders(httpRequest)) .filter(header -> header.startsWith(AUTHORIZATION_HEADER_PREFIX)) .map(header -> header.substring(AUTHORIZATION_HEADER_PREFIX.length())) - .map(userJWTToken -> { + .flatMap(userJWTToken -> Mono.fromCallable(() -> { if (!tokenManager.verify(userJWTToken)) { throw new UnauthorizedException("Failed Jwt verification"); } @@ -74,7 +75,7 @@ public Mono createMailboxSession(HttpServerRequest httpRequest) } return username; - }) + }).subscribeOn(Schedulers.elastic())) .map(mailboxManager::createSystemSession); } From 09dcd0780e52b0a704c77d55e647f7028012dee2 Mon Sep 17 00:00:00 2001 From: Benoit Tellier Date: Fri, 7 May 2021 18:40:26 +0700 Subject: [PATCH 06/15] [REFACTORING] Wrap blocking calls in UserProvisioner It was forcing to schedule the whole chain on an elastic scheduler... --- .../java/org/apache/james/jmap/http/UserProvisioner.java | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/http/UserProvisioner.java b/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/http/UserProvisioner.java index 1ac9aacf6c0..599228f5958 100644 --- a/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/http/UserProvisioner.java +++ b/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/http/UserProvisioner.java @@ -35,6 +35,7 @@ import com.google.common.annotations.VisibleForTesting; import reactor.core.publisher.Mono; +import reactor.core.scheduler.Schedulers; public class UserProvisioner { private final UsersRepository usersRepository; @@ -49,7 +50,9 @@ public class UserProvisioner { public Mono provisionUser(MailboxSession session) { if (session != null && !usersRepository.isReadOnly()) { - return Mono.fromRunnable(() -> createAccountIfNeeded(session)); + return Mono.fromRunnable(() -> createAccountIfNeeded(session)) + .subscribeOn(Schedulers.elastic()) + .then(); } return Mono.empty(); } From 27f51fe865628a373f586f1e938f5437186b9f5f Mon Sep 17 00:00:00 2001 From: Benoit Tellier Date: Fri, 7 May 2021 18:51:53 +0700 Subject: [PATCH 07/15] [REFACTORING] Dequeuer should not aggressively enforce a scheduler --- .../apache/james/queue/rabbitmq/Dequeuer.java | 22 +++++++++---------- 1 file changed, 10 insertions(+), 12 deletions(-) diff --git a/server/queue/queue-rabbitmq/src/main/java/org/apache/james/queue/rabbitmq/Dequeuer.java b/server/queue/queue-rabbitmq/src/main/java/org/apache/james/queue/rabbitmq/Dequeuer.java index 0241c460347..d3f0e66a6c9 100644 --- a/server/queue/queue-rabbitmq/src/main/java/org/apache/james/queue/rabbitmq/Dequeuer.java +++ b/server/queue/queue-rabbitmq/src/main/java/org/apache/james/queue/rabbitmq/Dequeuer.java @@ -42,7 +42,6 @@ import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; -import reactor.core.scheduler.Schedulers; import reactor.rabbitmq.AcknowledgableDelivery; import reactor.rabbitmq.ConsumeOptions; import reactor.rabbitmq.Receiver; @@ -105,21 +104,20 @@ public void close() { } Flux deQueue() { - return flux.flatMapSequential(response -> loadItem(response).subscribeOn(Schedulers.elastic())) - .concatMap(item -> filterIfDeleted(item).subscribeOn(Schedulers.elastic())); + return flux.flatMapSequential(this::loadItem) + .concatMap(this::filterIfDeleted); } private Mono filterIfDeleted(RabbitMQMailQueueItem item) { return mailQueueView.isPresent(item.getEnqueueId()) - .flatMap(isPresent -> keepWhenPresent(item, isPresent)); - } - - private Mono keepWhenPresent(RabbitMQMailQueueItem item, Boolean isPresent) { - if (isPresent) { - return Mono.just(item); - } - item.done(true); - return Mono.empty(); + .handle((isPresent, sink) -> { + if (isPresent) { + sink.next(item); + } else { + item.done(true); + sink.complete(); + } + }); } private Mono loadItem(AcknowledgableDelivery response) { From 5c31e369e0ea922c195a7e90065d62511c6dfc94 Mon Sep 17 00:00:00 2001 From: Benoit Tellier Date: Sat, 8 May 2021 14:50:59 +0700 Subject: [PATCH 08/15] [REFACTORING] RabbitMQ receivers can be blocking Detected thanks to https://github.com/reactor/BlockHound They are blocking in case of connection recovery. Doing that on the parallel pool is bad as it "steals" threads dedicated to non blocking tasks and slows the entire application down. --- .../src/main/java/org/apache/james/events/GroupRegistration.java | 1 + .../java/org/apache/james/events/KeyRegistrationHandler.java | 1 + 2 files changed, 2 insertions(+) diff --git a/event-bus/distributed/src/main/java/org/apache/james/events/GroupRegistration.java b/event-bus/distributed/src/main/java/org/apache/james/events/GroupRegistration.java index d0013ff3891..a33638bf18b 100644 --- a/event-bus/distributed/src/main/java/org/apache/james/events/GroupRegistration.java +++ b/event-bus/distributed/src/main/java/org/apache/james/events/GroupRegistration.java @@ -133,6 +133,7 @@ private Disposable consumeWorkQueue() { .publishOn(Schedulers.parallel()) .filter(delivery -> Objects.nonNull(delivery.getBody())) .flatMap(this::deliver, EventBus.EXECUTION_RATE) + .subscribeOn(Schedulers.elastic()) .subscribe(); } diff --git a/event-bus/distributed/src/main/java/org/apache/james/events/KeyRegistrationHandler.java b/event-bus/distributed/src/main/java/org/apache/james/events/KeyRegistrationHandler.java index a6088beb56a..d62f59fbb6e 100644 --- a/event-bus/distributed/src/main/java/org/apache/james/events/KeyRegistrationHandler.java +++ b/event-bus/distributed/src/main/java/org/apache/james/events/KeyRegistrationHandler.java @@ -96,6 +96,7 @@ void start() { receiverSubscriber = Optional.of(receiver.consumeAutoAck(registrationQueue.asString(), new ConsumeOptions().qos(EventBus.EXECUTION_RATE)) .subscribeOn(Schedulers.parallel()) .flatMap(this::handleDelivery, EventBus.EXECUTION_RATE) + .subscribeOn(Schedulers.elastic()) .subscribe()); } From 6af0f800827fdadc971aa249b59910e67b739d00 Mon Sep 17 00:00:00 2001 From: Benoit Tellier Date: Sat, 8 May 2021 14:49:59 +0700 Subject: [PATCH 09/15] [REFACTORING] GroupRegistration was doing acks on the parallel pool Detected thanks to https://github.com/reactor/BlockHound ACKs are blocking (as we send a message over the network to RabbitMQ). Doing that on the parallel pool is bad as it "steals" threads dedicated to non blocking tasks and slows the entire application down. --- .../main/java/org/apache/james/events/GroupRegistration.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/event-bus/distributed/src/main/java/org/apache/james/events/GroupRegistration.java b/event-bus/distributed/src/main/java/org/apache/james/events/GroupRegistration.java index a33638bf18b..4a83f8186a9 100644 --- a/event-bus/distributed/src/main/java/org/apache/james/events/GroupRegistration.java +++ b/event-bus/distributed/src/main/java/org/apache/james/events/GroupRegistration.java @@ -145,7 +145,7 @@ private Mono deliver(AcknowledgableDelivery acknowledgableDelivery) { .flatMap(event -> delayGenerator.delayIfHaveTo(currentRetryCount) .flatMap(any -> runListener(event)) .onErrorResume(throwable -> retryHandler.handleRetry(event, currentRetryCount, throwable)) - .then(Mono.fromRunnable(acknowledgableDelivery::ack))) + .then(Mono.fromRunnable(acknowledgableDelivery::ack).subscribeOn(Schedulers.elastic()))) .onErrorResume(e -> { LOGGER.error("Unable to process delivery for group {}", group, e); return Mono.fromRunnable(() -> acknowledgableDelivery.nack(!REQUEUE)); From 8d808b8ccf30d5db1d889f764780ccfd6d7da0e5 Mon Sep 17 00:00:00 2001 From: Benoit Tellier Date: Fri, 7 May 2021 18:47:23 +0700 Subject: [PATCH 10/15] JAMES-3148 DeleteMessageListener can easily be fully reactive --- .../mailbox/cassandra/DeleteMessageListener.java | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/DeleteMessageListener.java b/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/DeleteMessageListener.java index 2deefadde90..25af4c35594 100644 --- a/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/DeleteMessageListener.java +++ b/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/DeleteMessageListener.java @@ -62,6 +62,7 @@ import org.apache.james.mailbox.model.MessageRange; import org.apache.james.mailbox.store.mail.MessageMapper; import org.apache.james.util.streams.Limit; +import org.reactivestreams.Publisher; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; @@ -77,7 +78,7 @@ * Mailbox listener failures lead to eventBus retrying their execution, it ensures the result of the deletion to be * idempotent. */ -public class DeleteMessageListener implements EventListener.GroupEventListener { +public class DeleteMessageListener implements EventListener.ReactiveGroupEventListener { private static final Optional ALL_MAILBOXES = Optional.empty(); public static class DeleteMessageListenerGroup extends Group { @@ -138,21 +139,20 @@ public boolean isHandling(Event event) { } @Override - public void event(Event event) { + public Publisher reactiveEvent(Event event) { if (event instanceof Expunged) { Expunged expunged = (Expunged) event; - handleMessageDeletion(expunged) - .block(); + return handleMessageDeletion(expunged); } if (event instanceof MailboxDeletion) { MailboxDeletion mailboxDeletion = (MailboxDeletion) event; CassandraId mailboxId = (CassandraId) mailboxDeletion.getMailboxId(); - handleMailboxDeletion(mailboxId) - .block(); + return handleMailboxDeletion(mailboxId); } + return Mono.empty(); } private Mono handleMailboxDeletion(CassandraId mailboxId) { From c27b529f28ef4f4510e376c29add6f05d5389621 Mon Sep 17 00:00:00 2001 From: Benoit Tellier Date: Sat, 8 May 2021 16:48:14 +0700 Subject: [PATCH 11/15] JAMES-2393 Allow writing reactive eventSourcing subscribers --- .../apache/james/eventsourcing/EventBus.scala | 25 ++++++------------- .../james/eventsourcing/Subscriber.scala | 25 +++++++++++++++++++ .../eventsourcing/acl/AclV2DAOSubscriber.java | 12 +++++---- .../acl/UserRightsDAOSubscriber.java | 13 ++++++---- 4 files changed, 48 insertions(+), 27 deletions(-) diff --git a/event-sourcing/event-sourcing-core/src/main/scala/org/apache/james/eventsourcing/EventBus.scala b/event-sourcing/event-sourcing-core/src/main/scala/org/apache/james/eventsourcing/EventBus.scala index ca8d753643c..e932252aa6c 100644 --- a/event-sourcing/event-sourcing-core/src/main/scala/org/apache/james/eventsourcing/EventBus.scala +++ b/event-sourcing/event-sourcing-core/src/main/scala/org/apache/james/eventsourcing/EventBus.scala @@ -19,11 +19,9 @@ package org.apache.james.eventsourcing import javax.inject.Inject - import org.apache.james.eventsourcing.eventstore.{EventStore, EventStoreFailedException} import org.reactivestreams.Publisher import org.slf4j.LoggerFactory - import reactor.core.scala.publisher.{SFlux, SMono} object EventBus { @@ -32,27 +30,20 @@ object EventBus { class EventBus @Inject() (eventStore: EventStore, subscribers: Set[Subscriber]) { @throws[EventStoreFailedException] - def publish(events: Iterable[Event]): SMono[Void] = { + def publish(events: Iterable[Event]): SMono[Void] = SMono(eventStore.appendAll(events)) .`then`(runHandlers(events, subscribers)) - } - - def runHandlers(events: Iterable[Event], subscribers: Set[Subscriber]): SMono[Void] = { + def runHandlers(events: Iterable[Event], subscribers: Set[Subscriber]): SMono[Void] = SFlux.fromIterable(events.flatMap((event: Event) => subscribers.map(subscriber => (event, subscriber)))) - .flatMap(infos => runHandler(infos._1, infos._2)) + .flatMapSequential(infos => runHandler(infos._1, infos._2)) .`then`() .`then`(SMono.empty) - } - - def runHandler(event: Event, subscriber: Subscriber): Publisher[Void] = SMono.fromCallable(() => handle(event, subscriber)).`then`(SMono.empty) - private def handle(event : Event, subscriber: Subscriber) : Unit = { - try { - subscriber.handle(event) - } catch { - case e: Exception => + def runHandler(event: Event, subscriber: Subscriber): Publisher[Void] = + SMono(ReactiveSubscriber.asReactiveSubscriber(subscriber).handleReactive(event)) + .onErrorResume(e => { EventBus.LOGGER.error("Error while calling {} for {}", subscriber, event, e) - } - } + SMono.empty + }) } \ No newline at end of file diff --git a/event-sourcing/event-sourcing-core/src/main/scala/org/apache/james/eventsourcing/Subscriber.scala b/event-sourcing/event-sourcing-core/src/main/scala/org/apache/james/eventsourcing/Subscriber.scala index ad6cf0a604c..3b14fef183d 100644 --- a/event-sourcing/event-sourcing-core/src/main/scala/org/apache/james/eventsourcing/Subscriber.scala +++ b/event-sourcing/event-sourcing-core/src/main/scala/org/apache/james/eventsourcing/Subscriber.scala @@ -18,6 +18,31 @@ * ***************************************************************/ package org.apache.james.eventsourcing +import org.reactivestreams.Publisher +import reactor.core.scala.publisher.SMono +import reactor.core.scheduler.Schedulers + trait Subscriber { def handle(event: Event) : Unit +} + +trait ReactiveSubscriber extends Subscriber { + def handleReactive(event: Event): Publisher[Void] + + override def handle(event: Event) : Unit = SMono(handleReactive(event)).block() +} + +object ReactiveSubscriber { + def asReactiveSubscriber(subscriber: Subscriber): ReactiveSubscriber = subscriber match { + case reactive: ReactiveSubscriber => reactive + case nonReactive => new ReactiveSubscriberWrapper(nonReactive) + } +} + +class ReactiveSubscriberWrapper(delegate: Subscriber) extends ReactiveSubscriber { + override def handle(event: Event) : Unit = delegate.handle(event) + + def handleReactive(event: Event): Publisher[Void] = SMono.fromCallable(() => handle(event)) + .subscribeOn(Schedulers.elastic()) + .`then`() } \ No newline at end of file diff --git a/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/eventsourcing/acl/AclV2DAOSubscriber.java b/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/eventsourcing/acl/AclV2DAOSubscriber.java index 6c579ea5c78..1b3ecada56d 100644 --- a/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/eventsourcing/acl/AclV2DAOSubscriber.java +++ b/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/eventsourcing/acl/AclV2DAOSubscriber.java @@ -20,12 +20,13 @@ package org.apache.james.mailbox.cassandra.mail.eventsourcing.acl; import org.apache.james.eventsourcing.Event; -import org.apache.james.eventsourcing.Subscriber; +import org.apache.james.eventsourcing.ReactiveSubscriber; import org.apache.james.mailbox.cassandra.mail.CassandraACLDAOV2; import reactor.core.publisher.Flux; +import reactor.core.publisher.Mono; -public class AclV2DAOSubscriber implements Subscriber { +public class AclV2DAOSubscriber implements ReactiveSubscriber { private final CassandraACLDAOV2 acldaov2; public AclV2DAOSubscriber(CassandraACLDAOV2 acldaov2) { @@ -33,15 +34,16 @@ public AclV2DAOSubscriber(CassandraACLDAOV2 acldaov2) { } @Override - public void handle(Event event) { + public Mono handleReactive(Event event) { if (event instanceof ACLUpdated) { ACLUpdated aclUpdated = (ACLUpdated) event; - Flux.fromStream( + return Flux.fromStream( aclUpdated.getAclDiff() .commands()) .flatMap(command -> acldaov2.updateACL(aclUpdated.mailboxId(), command)) - .blockLast(); + .then(); } + return Mono.empty(); } } diff --git a/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/eventsourcing/acl/UserRightsDAOSubscriber.java b/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/eventsourcing/acl/UserRightsDAOSubscriber.java index 4bdfb19be4c..32eec224927 100644 --- a/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/eventsourcing/acl/UserRightsDAOSubscriber.java +++ b/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/eventsourcing/acl/UserRightsDAOSubscriber.java @@ -20,10 +20,13 @@ package org.apache.james.mailbox.cassandra.mail.eventsourcing.acl; import org.apache.james.eventsourcing.Event; -import org.apache.james.eventsourcing.Subscriber; +import org.apache.james.eventsourcing.ReactiveSubscriber; import org.apache.james.mailbox.cassandra.mail.CassandraUserMailboxRightsDAO; +import org.reactivestreams.Publisher; -public class UserRightsDAOSubscriber implements Subscriber { +import reactor.core.publisher.Mono; + +public class UserRightsDAOSubscriber implements ReactiveSubscriber { private final CassandraUserMailboxRightsDAO userRightsDAO; public UserRightsDAOSubscriber(CassandraUserMailboxRightsDAO userRightsDAO) { @@ -31,11 +34,11 @@ public UserRightsDAOSubscriber(CassandraUserMailboxRightsDAO userRightsDAO) { } @Override - public void handle(Event event) { + public Publisher handleReactive(Event event) { if (event instanceof ACLUpdated) { ACLUpdated aclUpdated = (ACLUpdated) event; - userRightsDAO.update(aclUpdated.mailboxId(), aclUpdated.getAclDiff()) - .block(); + return userRightsDAO.update(aclUpdated.mailboxId(), aclUpdated.getAclDiff()); } + return Mono.empty(); } } From dd2d5b4e4dc5faa02355cce967d1444b0c2c1ac7 Mon Sep 17 00:00:00 2001 From: Benoit Tellier Date: Fri, 7 May 2021 18:31:12 +0700 Subject: [PATCH 12/15] JAMES-2979 MailDispatcher: Use immediate schedule to avoid allocating threads upon nested block calls This test mimics the previous behaviour of LocalDelivery ``` @Test void testNestedBlocksWithElasticScheduler() { // mono1 corresponds to the append to the mailbox (blocking) Mono mono1 = Mono.fromRunnable(() -> { System.out.println("mono1 running on " + Thread.currentThread().getName()); }); // mono2 corresponds to the retries perforned in MailDispatcher Mono mono2 = Mono.fromRunnable(() -> { System.out.println("mono2 running on " + Thread.currentThread().getName()); mono1.subscribeOn(Schedulers.elastic()).block(); }); // This is a spooler thread running LocalDelivery System.out.println("Current thread " + Thread.currentThread().getName()); mono2.subscribeOn(Schedulers.elastic()).block(); } ``` Output: ``` Current thread main mono2 running on elastic-2 mono1 running on elastic-3 ``` One thread doing the work and two waiting... With the new paradigm (waiting for a fully reactive version) ``` @Test void testNestedBlockWithImmediateScheduler() { // mono1 corresponds to the append to the mailbox (blocking) Mono mono1 = Mono.fromRunnable(() -> { System.out.println("mono1 running on " + Thread.currentThread().getName()); }); // mono2 corresponds to the retries perforned in MailDispatcher Mono mono2 = Mono.fromRunnable(() -> { System.out.println("mono2 running on " + Thread.currentThread().getName()); mono1.subscribeOn(Schedulers.immediate()).block(); }); // This is a spooler thread running LocalDelivery System.out.println("Current thread " + Thread.currentThread().getName()); mono2.subscribeOn(Schedulers.immediate()).block(); } ``` We get... ``` Current thread main mono2 running on main mono1 running on main ``` Way better! We do mobilize only one thread that would anyway be blocked. --- .../james/transport/mailets/delivery/MailDispatcher.java | 6 +----- 1 file changed, 1 insertion(+), 5 deletions(-) diff --git a/server/mailet/mailets/src/main/java/org/apache/james/transport/mailets/delivery/MailDispatcher.java b/server/mailet/mailets/src/main/java/org/apache/james/transport/mailets/delivery/MailDispatcher.java index ec66410ac91..2a7bf6aba50 100644 --- a/server/mailet/mailets/src/main/java/org/apache/james/transport/mailets/delivery/MailDispatcher.java +++ b/server/mailet/mailets/src/main/java/org/apache/james/transport/mailets/delivery/MailDispatcher.java @@ -43,7 +43,6 @@ import com.google.common.collect.ImmutableMap; import reactor.core.publisher.Mono; -import reactor.core.scheduler.Scheduler; import reactor.core.scheduler.Schedulers; import reactor.util.retry.Retry; @@ -90,13 +89,11 @@ public MailDispatcher build() { private final MailStore mailStore; private final boolean consume; private final MailetContext mailetContext; - private final Scheduler scheduler; private MailDispatcher(MailStore mailStore, boolean consume, MailetContext mailetContext) { this.mailStore = mailStore; this.consume = consume; this.mailetContext = mailetContext; - this.scheduler = Schedulers.elastic(); } public void dispatch(Mail mail) throws MessagingException { @@ -142,7 +139,7 @@ private List deliver(Mail mail, MimeMessage message) { Map> savedHeaders = saveHeaders(mail, recipient); addSpecificHeadersForRecipient(mail, message, recipient); - storeMailWithRetry(mail, recipient).block(); + storeMailWithRetry(mail, recipient).subscribeOn(Schedulers.immediate()).block(); restoreHeaders(mail.getMessage(), savedHeaders); } catch (Exception ex) { @@ -156,7 +153,6 @@ private List deliver(Mail mail, MimeMessage message) { private Mono storeMailWithRetry(Mail mail, MailAddress recipient) { return Mono.fromRunnable((ThrowingRunnable)() -> mailStore.storeMail(recipient, mail)) .doOnError(error -> LOGGER.warn("Error While storing mail. This error will be retried.", error)) - .subscribeOn(scheduler) .retryWhen(Retry.backoff(RETRIES, FIRST_BACKOFF).maxBackoff(MAX_BACKOFF).scheduler(Schedulers.elastic())) .then(); } From 1fe6ae24a8f028e704d62c0feefa204da489e892 Mon Sep 17 00:00:00 2001 From: Benoit Tellier Date: Sat, 8 May 2021 14:59:24 +0700 Subject: [PATCH 13/15] JAMES-3138 ListeningCurrentQuotaUpdater should base itself on the mailboxPath to determine quotaRoot --- .../store/quota/ListeningCurrentQuotaUpdater.java | 5 ++--- .../mailbox/store/quota/StoreQuotaManager.java | 1 - .../quota/ListeningCurrentQuotaUpdaterTest.java | 14 +++++++++++++- 3 files changed, 15 insertions(+), 5 deletions(-) diff --git a/mailbox/store/src/main/java/org/apache/james/mailbox/store/quota/ListeningCurrentQuotaUpdater.java b/mailbox/store/src/main/java/org/apache/james/mailbox/store/quota/ListeningCurrentQuotaUpdater.java index cfe2d801582..e6b158e34cc 100644 --- a/mailbox/store/src/main/java/org/apache/james/mailbox/store/quota/ListeningCurrentQuotaUpdater.java +++ b/mailbox/store/src/main/java/org/apache/james/mailbox/store/quota/ListeningCurrentQuotaUpdater.java @@ -45,7 +45,6 @@ import com.google.common.collect.ImmutableSet; import reactor.core.publisher.Mono; -import reactor.core.scheduler.Schedulers; public class ListeningCurrentQuotaUpdater implements EventListener.ReactiveGroupEventListener, QuotaUpdater { public static class ListeningCurrentQuotaUpdaterGroup extends Group { @@ -82,11 +81,11 @@ public boolean isHandling(Event event) { public Publisher reactiveEvent(Event event) { if (event instanceof Added) { Added addedEvent = (Added) event; - return Mono.from(quotaRootResolver.getQuotaRootReactive(addedEvent.getMailboxId())) + return Mono.from(quotaRootResolver.getQuotaRootReactive(addedEvent.getMailboxPath())) .flatMap(quotaRoot -> handleAddedEvent(addedEvent, quotaRoot)); } else if (event instanceof Expunged) { Expunged expungedEvent = (Expunged) event; - return Mono.from(quotaRootResolver.getQuotaRootReactive(expungedEvent.getMailboxId())) + return Mono.from(quotaRootResolver.getQuotaRootReactive(expungedEvent.getMailboxPath())) .flatMap(quotaRoot -> handleExpungedEvent(expungedEvent, quotaRoot)); } else if (event instanceof MailboxDeletion) { MailboxDeletion mailboxDeletionEvent = (MailboxDeletion) event; diff --git a/mailbox/store/src/main/java/org/apache/james/mailbox/store/quota/StoreQuotaManager.java b/mailbox/store/src/main/java/org/apache/james/mailbox/store/quota/StoreQuotaManager.java index aa010bce42a..24ca588ed44 100644 --- a/mailbox/store/src/main/java/org/apache/james/mailbox/store/quota/StoreQuotaManager.java +++ b/mailbox/store/src/main/java/org/apache/james/mailbox/store/quota/StoreQuotaManager.java @@ -37,7 +37,6 @@ import org.reactivestreams.Publisher; import reactor.core.publisher.Mono; -import reactor.core.scheduler.Schedulers; /** * Default implementation for the Quota Manager. diff --git a/mailbox/store/src/test/java/org/apache/james/mailbox/store/quota/ListeningCurrentQuotaUpdaterTest.java b/mailbox/store/src/test/java/org/apache/james/mailbox/store/quota/ListeningCurrentQuotaUpdaterTest.java index 26452e2a726..f71cba29c84 100644 --- a/mailbox/store/src/test/java/org/apache/james/mailbox/store/quota/ListeningCurrentQuotaUpdaterTest.java +++ b/mailbox/store/src/test/java/org/apache/james/mailbox/store/quota/ListeningCurrentQuotaUpdaterTest.java @@ -46,6 +46,7 @@ import org.apache.james.mailbox.events.MailboxEvents.Expunged; import org.apache.james.mailbox.events.MailboxEvents.MailboxDeletion; import org.apache.james.mailbox.model.MailboxId; +import org.apache.james.mailbox.model.MailboxPath; import org.apache.james.mailbox.model.MessageMetaData; import org.apache.james.mailbox.model.QuotaOperation; import org.apache.james.mailbox.model.QuotaRoot; @@ -67,6 +68,7 @@ class ListeningCurrentQuotaUpdaterTest { static final MailboxId MAILBOX_ID = TestId.of(42); static final String BENWA = "benwa"; static final Username USERNAME_BENWA = Username.of(BENWA); + static final MailboxPath MAILBOX_PATH = MailboxPath.forUser(USERNAME_BENWA, "path"); static final QuotaRoot QUOTA_ROOT = QuotaRoot.quotaRoot(BENWA, Optional.empty()); static final QuotaOperation QUOTA = new QuotaOperation(QUOTA_ROOT, QuotaCountUsage.count(2), QuotaSizeUsage.size(2 * SIZE)); @@ -80,8 +82,10 @@ void setUp() { mockedCurrentQuotaManager = mock(CurrentQuotaManager.class); EventBus eventBus = mock(EventBus.class); when(eventBus.dispatch(any(Event.class), anySet())).thenReturn(Mono.empty()); + QuotaManager quotaManager = mock(QuotaManager.class); + when(quotaManager.getQuotasReactive(eq(QUOTA_ROOT))).thenReturn(Mono.empty()); testee = new ListeningCurrentQuotaUpdater(mockedCurrentQuotaManager, mockedQuotaRootResolver, - eventBus, mock(QuotaManager.class)); + eventBus, quotaManager); } @Test @@ -94,11 +98,13 @@ void deserializeListeningCurrentQuotaUpdaterGroup() throws Exception { void addedEventShouldIncreaseCurrentQuotaValues() throws Exception { Added added = mock(Added.class); when(added.getMailboxId()).thenReturn(MAILBOX_ID); + when(added.getMailboxPath()).thenReturn(MAILBOX_PATH); when(added.getMetaData(MessageUid.of(36))).thenReturn(new MessageMetaData(MessageUid.of(36), ModSeq.first(),new Flags(), SIZE, new Date(), new DefaultMessageId())); when(added.getMetaData(MessageUid.of(38))).thenReturn(new MessageMetaData(MessageUid.of(38), ModSeq.first(),new Flags(), SIZE, new Date(), new DefaultMessageId())); when(added.getUids()).thenReturn(Lists.newArrayList(MessageUid.of(36), MessageUid.of(38))); when(added.getUsername()).thenReturn(USERNAME_BENWA); when(mockedQuotaRootResolver.getQuotaRootReactive(eq(MAILBOX_ID))).thenReturn(Mono.just(QUOTA_ROOT)); + when(mockedQuotaRootResolver.getQuotaRootReactive(eq(MAILBOX_PATH))).thenReturn(Mono.just(QUOTA_ROOT)); when(mockedCurrentQuotaManager.increase(QUOTA)).thenAnswer(any -> Mono.empty()); testee.event(added); @@ -114,6 +120,8 @@ void expungedEventShouldDecreaseCurrentQuotaValues() throws Exception { when(expunged.getUids()).thenReturn(Lists.newArrayList(MessageUid.of(36), MessageUid.of(38))); when(expunged.getMailboxId()).thenReturn(MAILBOX_ID); when(expunged.getUsername()).thenReturn(USERNAME_BENWA); + when(expunged.getMailboxPath()).thenReturn(MAILBOX_PATH); + when(mockedQuotaRootResolver.getQuotaRootReactive(eq(MAILBOX_PATH))).thenReturn(Mono.just(QUOTA_ROOT)); when(mockedQuotaRootResolver.getQuotaRootReactive(eq(MAILBOX_ID))).thenReturn(Mono.just(QUOTA_ROOT)); when(mockedCurrentQuotaManager.decrease(QUOTA)).thenAnswer(any -> Mono.empty()); @@ -128,6 +136,8 @@ void emptyExpungedEventShouldNotTriggerDecrease() throws Exception { when(expunged.getUids()).thenReturn(Lists.newArrayList()); when(expunged.getMailboxId()).thenReturn(MAILBOX_ID); when(expunged.getUsername()).thenReturn(USERNAME_BENWA); + when(expunged.getMailboxPath()).thenReturn(MAILBOX_PATH); + when(mockedQuotaRootResolver.getQuotaRootReactive(eq(MAILBOX_PATH))).thenReturn(Mono.just(QUOTA_ROOT)); when(mockedQuotaRootResolver.getQuotaRootReactive(eq(MAILBOX_ID))).thenReturn(Mono.just(QUOTA_ROOT)); testee.event(expunged); @@ -141,6 +151,8 @@ void emptyAddedEventShouldNotTriggerDecrease() throws Exception { when(added.getUids()).thenReturn(Lists.newArrayList()); when(added.getMailboxId()).thenReturn(MAILBOX_ID); when(added.getUsername()).thenReturn(USERNAME_BENWA); + when(added.getMailboxPath()).thenReturn(MAILBOX_PATH); + when(mockedQuotaRootResolver.getQuotaRootReactive(eq(MAILBOX_PATH))).thenReturn(Mono.just(QUOTA_ROOT)); when(mockedQuotaRootResolver.getQuotaRootReactive(eq(MAILBOX_ID))).thenReturn(Mono.just(QUOTA_ROOT)); testee.event(added); From fac7763850f6dfde5cccbcf6a81b309259a7f466 Mon Sep 17 00:00:00 2001 From: Benoit Tellier Date: Sat, 8 May 2021 15:00:30 +0700 Subject: [PATCH 14/15] JAMES-3407 Applicative read-repair: draw a random number only if needed --- .../mailbox/cassandra/mail/CassandraMailboxMapper.java | 3 ++- .../mailbox/cassandra/mail/CassandraMessageMapper.java | 7 +++++-- 2 files changed, 7 insertions(+), 3 deletions(-) diff --git a/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/CassandraMailboxMapper.java b/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/CassandraMailboxMapper.java index d372af74ddd..17fcae8bc94 100644 --- a/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/CassandraMailboxMapper.java +++ b/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/CassandraMailboxMapper.java @@ -147,7 +147,8 @@ private Mono performPathReadRepair(Mailbox mailboxPathEntry) { } private boolean shouldReadRepair() { - return secureRandom.nextFloat() < cassandraConfiguration.getMailboxReadRepair(); + return cassandraConfiguration.getMailboxReadRepair() > 0 + && secureRandom.nextFloat() < cassandraConfiguration.getMailboxReadRepair(); } diff --git a/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/CassandraMessageMapper.java b/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/CassandraMessageMapper.java index 0aa353a97ea..f929b830980 100644 --- a/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/CassandraMessageMapper.java +++ b/mailbox/cassandra/src/main/java/org/apache/james/mailbox/cassandra/mail/CassandraMessageMapper.java @@ -180,9 +180,12 @@ private Mono fixCounters(Mailbox mailbox) { } private boolean shouldReadRepair(MailboxCounters counters) { + boolean activated = cassandraConfiguration.getMailboxCountersReadRepairChanceMax() != 0 || cassandraConfiguration.getMailboxCountersReadRepairChanceOneHundred() != 0; double ponderedReadRepairChance = cassandraConfiguration.getMailboxCountersReadRepairChanceOneHundred() * (100.0 / counters.getUnseen()); - return secureRandom.nextFloat() < Math.min( - cassandraConfiguration.getMailboxCountersReadRepairChanceMax(), ponderedReadRepairChance); + return activated && + secureRandom.nextFloat() < Math.min( + cassandraConfiguration.getMailboxCountersReadRepairChanceMax(), + ponderedReadRepairChance); } @Override From cfdd144a669c92d344f6eaee7411d15d12a5edd4 Mon Sep 17 00:00:00 2001 From: Benoit Tellier Date: Mon, 10 May 2021 13:23:59 +0700 Subject: [PATCH 15/15] JAMES-1965 MessageFullViewFactory: Avoid performing HTML text extraction if not needed --- .../draft/model/message/view/MessageFullViewFactory.java | 8 +++----- 1 file changed, 3 insertions(+), 5 deletions(-) diff --git a/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/model/message/view/MessageFullViewFactory.java b/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/model/message/view/MessageFullViewFactory.java index d391e0314da..73ca18b662f 100644 --- a/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/model/message/view/MessageFullViewFactory.java +++ b/server/protocols/jmap-draft/src/main/java/org/apache/james/jmap/draft/model/message/view/MessageFullViewFactory.java @@ -102,8 +102,7 @@ public Mono fromMetaDataWithContent(MetaDataWithContent message private Mono fromMetaDataWithContent(MetaDataWithContent message, Message mimeMessage) throws IOException { MessageContent messageContent = messageContentExtractor.extract(mimeMessage); Optional htmlBody = messageContent.getHtmlBody(); - Optional mainTextContent = messageContent.extractMainTextContent(htmlTextExtractor); - Optional textBody = computeTextBodyIfNeeded(messageContent, mainTextContent); + Optional textBody = computeTextBodyIfNeeded(messageContent); return retrieveProjection(messageContent, message.getMessageId(), () -> MessageFullView.hasAttachment(getAttachments(message.getAttachments()))) @@ -196,10 +195,9 @@ private MetaDataWithContent toMetaDataWithContent(Collection mess .build(); } - private Optional computeTextBodyIfNeeded(MessageContent messageContent, Optional mainTextContent) { + private Optional computeTextBodyIfNeeded(MessageContent messageContent) { return messageContent.getTextBody() - .map(Optional::of) - .orElse(mainTextContent); + .or(() -> messageContent.extractMainTextContent(htmlTextExtractor)); } private Optional mainTextContent(MessageContent messageContent) {