This is an automated email from the ASF dual-hosted git repository. btellier pushed a commit to branch master in repository https://gitbox.apache.org/repos/asf/james-project.git
commit 992cec416423cd5a71bfc6d07b33e6bb304fe58c Author: Rene Cordier <[email protected]> AuthorDate: Thu Jun 11 17:47:08 2020 +0700 JAMES-3202 Reindex only outdated documents with the Mode option set to CORRECT in reindexing tasks --- .../mailbox/tools/indexer/ReIndexerPerformer.java | 108 ++++++++++++++------- 1 file changed, 71 insertions(+), 37 deletions(-) diff --git a/mailbox/tools/indexer/src/main/java/org/apache/mailbox/tools/indexer/ReIndexerPerformer.java b/mailbox/tools/indexer/src/main/java/org/apache/mailbox/tools/indexer/ReIndexerPerformer.java index 34d1a95..b854d67 100644 --- a/mailbox/tools/indexer/src/main/java/org/apache/mailbox/tools/indexer/ReIndexerPerformer.java +++ b/mailbox/tools/indexer/src/main/java/org/apache/mailbox/tools/indexer/ReIndexerPerformer.java @@ -19,11 +19,10 @@ package org.apache.mailbox.tools.indexer; -import static org.apache.james.mailbox.store.mail.AbstractMessageMapper.UNLIMITED; - import java.time.Duration; import javax.inject.Inject; +import javax.mail.Flags; import org.apache.james.core.Username; import org.apache.james.mailbox.MailboxManager; @@ -56,24 +55,25 @@ import reactor.core.publisher.Mono; public class ReIndexerPerformer { public static final int MAILBOX_CONCURRENCY = 1; + public static final int ONE = 1; private static class ReIndexingEntry { private final Mailbox mailbox; private final MailboxSession mailboxSession; - private final MailboxMessage message; + private final MessageUid uid; - ReIndexingEntry(Mailbox mailbox, MailboxSession mailboxSession, MailboxMessage message) { + ReIndexingEntry(Mailbox mailbox, MailboxSession mailboxSession, MessageUid uid) { this.mailbox = mailbox; this.mailboxSession = mailboxSession; - this.message = message; + this.uid = uid; } public Mailbox getMailbox() { return mailbox; } - public MailboxMessage getMessage() { - return message; + public MessageUid getUid() { + return uid; } public MailboxSession getMailboxSession() { @@ -149,7 +149,7 @@ public class ReIndexerPerformer { LOGGER.info("Starting a full reindex"); Flux<Either<Failure, ReIndexingEntry>> entriesToIndex = mailboxSessionMapperFactory.getMailboxMapper(mailboxSession).list() - .flatMap(mailbox -> reIndexingEntriesForMailbox(mailbox, mailboxSession), MAILBOX_CONCURRENCY); + .flatMap(mailbox -> reIndexingEntriesForMailbox(mailbox, mailboxSession, runningOptions), MAILBOX_CONCURRENCY); return reIndexMessages(entriesToIndex, runningOptions, reprocessingContext) .doFinally(any -> LOGGER.info("Full reindex finished")); @@ -160,7 +160,7 @@ public class ReIndexerPerformer { Flux<Either<Failure, ReIndexingEntry>> entriesToIndex = mailboxSessionMapperFactory.getMailboxMapper(mailboxSession) .findMailboxById(mailboxId) - .flatMapMany(mailbox -> reIndexingEntriesForMailbox(mailbox, mailboxSession)); + .flatMapMany(mailbox -> reIndexingEntriesForMailbox(mailbox, mailboxSession, runningOptions)); return reIndexMessages(entriesToIndex, runningOptions, reprocessingContext); } @@ -174,7 +174,7 @@ public class ReIndexerPerformer { try { Flux<Either<Failure, ReIndexingEntry>> entriesToIndex = mailboxMapper.findMailboxWithPathLike(mailboxQuery.asUserBound()) - .flatMap(mailbox -> reIndexingEntriesForMailbox(mailbox, mailboxSession), MAILBOX_CONCURRENCY); + .flatMap(mailbox -> reIndexingEntriesForMailbox(mailbox, mailboxSession, runningOptions), MAILBOX_CONCURRENCY); return reIndexMessages(entriesToIndex, runningOptions, reprocessingContext) .doFinally(any -> LOGGER.info("User {} reindex finished", username.asString())); @@ -189,9 +189,9 @@ public class ReIndexerPerformer { return mailboxSessionMapperFactory.getMailboxMapper(mailboxSession) .findMailboxById(mailboxId) - .flatMap(mailbox -> fullyReadMessage(mailboxSession, mailbox, uid) - .map(message -> Either.<Failure, ReIndexingEntry>right(new ReIndexingEntry(mailbox, mailboxSession, message))) - .flatMap(entryOrFailure -> reIndex(entryOrFailure, reprocessingContext))) + .map(mailbox -> new ReIndexingEntry(mailbox, mailboxSession, uid)) + .flatMap(this::fullyReadMessage) + .flatMap(message -> reIndex(message, mailboxSession)) .switchIfEmpty(Mono.just(Result.COMPLETED)); } @@ -218,11 +218,11 @@ public class ReIndexerPerformer { .flatMap(this::createReindexingEntryFromFailure), Flux.fromIterable(previousReIndexingFailures.mailboxFailures()) .flatMap(mailboxId -> mapper.findMailboxById(mailboxId) - .flatMapMany(mailbox -> reIndexingEntriesForMailbox(mailbox, mailboxSession)) + .flatMapMany(mailbox -> reIndexingEntriesForMailbox(mailbox, mailboxSession, runningOptions)) .onErrorResume(e -> { LOGGER.warn("Failed to re-index {}", mailboxId, e); return Mono.just(Either.left(new MailboxFailure(mailboxId))); - }))); + }), MAILBOX_CONCURRENCY)); return reIndexMessages(entriesToIndex, runningOptions, reprocessingContext); } @@ -238,9 +238,9 @@ public class ReIndexerPerformer { }); } - private Mono<MailboxMessage> fullyReadMessage(MailboxSession mailboxSession, Mailbox mailbox, MessageUid mUid) { - return mailboxSessionMapperFactory.getMessageMapper(mailboxSession) - .findInMailboxReactive(mailbox, MessageRange.one(mUid), MessageMapper.FetchType.Full, SINGLE_MESSAGE) + private Mono<MailboxMessage> fullyReadMessage(ReIndexingEntry entry) { + return mailboxSessionMapperFactory.getMessageMapper(entry.getMailboxSession()) + .findInMailboxReactive(entry.getMailbox(), MessageRange.one(entry.getUid()), MessageMapper.FetchType.Full, SINGLE_MESSAGE) .next(); } @@ -249,46 +249,44 @@ public class ReIndexerPerformer { return mailboxSessionMapperFactory.getMailboxMapper(mailboxSession) .findMailboxById(previousFailure.getMailboxId()) - .flatMap(mailbox -> fullyReadMessage(mailboxSession, mailbox, previousFailure.getUid()) - .map(message -> Either.<Failure, ReIndexingEntry>right(new ReIndexingEntry(mailbox, mailboxSession, message)))) + .map(mailbox -> Either.<Failure, ReIndexingEntry>right(new ReIndexingEntry(mailbox, mailboxSession, previousFailure.getUid()))) .onErrorResume(e -> { LOGGER.warn("ReIndexing failed for {}", previousFailure, e); return Mono.just(Either.left(new MessageFailure(previousFailure.getMailboxId(), previousFailure.getUid()))); }); } - private Flux<Either<Failure, ReIndexingEntry>> reIndexingEntriesForMailbox(Mailbox mailbox, MailboxSession mailboxSession) { + private Flux<Either<Failure, ReIndexingEntry>> reIndexingEntriesForMailbox(Mailbox mailbox, MailboxSession mailboxSession, RunningOptions runningOptions) { MessageMapper messageMapper = mailboxSessionMapperFactory.getMessageMapper(mailboxSession); - return messageSearchIndex.deleteAll(mailboxSession, mailbox.getMailboxId()) + return updateSearchIndex(mailbox, mailboxSession, runningOptions) .thenMany(messageMapper.listAllMessageUids(mailbox)) - .flatMap(uid -> reIndexingEntryForUid(mailbox, mailboxSession, messageMapper, uid)) + .map(uid -> Either.<Failure, ReIndexingEntry>right(new ReIndexingEntry(mailbox, mailboxSession, uid))) .onErrorResume(e -> { LOGGER.warn("ReIndexing failed for {}", mailbox.generateAssociatedPath(), e); return Mono.just(Either.left(new MailboxFailure(mailbox.getMailboxId()))); }); } - private Flux<Either<Failure, ReIndexingEntry>> reIndexingEntryForUid(Mailbox mailbox, MailboxSession mailboxSession, MessageMapper messageMapper, MessageUid uid) { - return messageMapper.findInMailboxReactive(mailbox, MessageRange.one(uid), MessageMapper.FetchType.Full, UNLIMITED) - .map(message -> Either.<Failure, ReIndexingEntry>right(new ReIndexingEntry(mailbox, mailboxSession, message))) - .onErrorResume(e -> { - LOGGER.warn("ReIndexing failed for {} {}", mailbox.getMailboxId(), uid, e); - return Mono.just(Either.left(new MessageFailure(mailbox.getMailboxId(), uid))); - }); + private Mono<Void> updateSearchIndex(Mailbox mailbox, MailboxSession mailboxSession, RunningOptions runningOptions) { + if (runningOptions.getMode() == RunningOptions.Mode.REBUILD_ALL) { + return messageSearchIndex.deleteAll(mailboxSession, mailbox.getMailboxId()); + } + return Mono.empty(); } private Mono<Task.Result> reIndexMessages(Flux<Either<Failure, ReIndexingEntry>> entriesToIndex, RunningOptions runningOptions, ReprocessingContext reprocessingContext) { - return entriesToIndex.transform(ReactorUtils.<Either<Failure, ReIndexingEntry>, Task.Result>throttle() + return entriesToIndex.transform( + ReactorUtils.<Either<Failure, ReIndexingEntry>, Task.Result>throttle() .elements(runningOptions.getMessagesPerSecond()) .per(Duration.ofSeconds(1)) - .forOperation(entry -> reIndex(entry, reprocessingContext))) + .forOperation(entry -> reIndex(entry, reprocessingContext, runningOptions))) .reduce(Task::combine) .switchIfEmpty(Mono.just(Result.COMPLETED)); } - private Mono<Task.Result> reIndex(Either<Failure, ReIndexingEntry> failureOrEntry, ReprocessingContext reprocessingContext) { - return toMono(failureOrEntry.map(this::index)) + private Mono<Task.Result> reIndex(Either<Failure, ReIndexingEntry> failureOrEntry, ReprocessingContext reprocessingContext, RunningOptions runningOptions) { + return toMono(failureOrEntry.map(entry -> reIndex(entry, runningOptions))) .map(this::flatten) .map(failureOrTaskResult -> recordIndexingResult(failureOrTaskResult, reprocessingContext)); } @@ -302,15 +300,51 @@ public class ReIndexerPerformer { result -> result.onComplete(reprocessingContext::recordSuccess)); } + private Mono<Either<Failure, Result>> reIndex(ReIndexingEntry entry, RunningOptions runningOptions) { + if (runningOptions.getMode() == RunningOptions.Mode.FIX_OUTDATED) { + return correctIfNeeded(entry); + } + return index(entry); + } + private Mono<Either<Failure, Result>> index(ReIndexingEntry entry) { - return messageSearchIndex.add(entry.getMailboxSession(), entry.getMailbox(), entry.getMessage()) + return fullyReadMessage(entry) + .flatMap(message -> messageSearchIndex.add(entry.getMailboxSession(), entry.getMailbox(), message)) .thenReturn(Either.<Failure, Result>right(Result.COMPLETED)) .onErrorResume(e -> { - LOGGER.warn("ReIndexing failed for {} {}", entry.getMailbox().generateAssociatedPath(), entry.getMessage().getUid(), e); - return Mono.just(Either.left(new MessageFailure(entry.getMailbox().getMailboxId(), entry.getMessage().getUid()))); + LOGGER.warn("ReIndexing failed for {} {}", entry.getMailbox().generateAssociatedPath(), entry.getUid(), e); + return Mono.just(Either.left(new MessageFailure(entry.getMailbox().getMailboxId(), entry.getUid()))); }); } + private Mono<Either<Failure, Result>> correctIfNeeded(ReIndexingEntry entry) { + MessageMapper messageMapper = mailboxSessionMapperFactory.getMessageMapper(entry.getMailboxSession()); + + return messageMapper.findInMailboxReactive(entry.getMailbox(), MessageRange.one(entry.getUid()), MessageMapper.FetchType.Metadata, ONE) + .next() + .flatMap(message -> isIndexUpToDate(entry.getMailbox(), message) + .flatMap(upToDate -> { + if (upToDate) { + return Mono.just(Either.right(Result.COMPLETED)); + } + return correct(entry, message); + })); + } + + private Mono<Either<Failure, Result>> correct(ReIndexingEntry entry, MailboxMessage message) { + return messageSearchIndex.delete(entry.getMailboxSession(), entry.getMailbox(), ImmutableList.of(message.getUid())) + .then(index(entry)); + } + + private Mono<Boolean> isIndexUpToDate(Mailbox mailbox, MailboxMessage message) { + return messageSearchIndex.retrieveIndexedFlags(mailbox, message.getUid()) + .map(flags -> isIndexUpToDate(message, flags)); + } + + private boolean isIndexUpToDate(MailboxMessage message, Flags flags) { + return message.createFlags().equals(flags); + } + private <X, Y> Either<X, Y> flatten(Either<X, Either<X, Y>> nestedEither) { return nestedEither.getOrElseGet(Either::left); } --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
