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]

Reply via email to