Avoid duplicated sync tasks during index sync with more than 2 peers

Change-Id: Ie771f420e41f7f006a9cc7c2105f9e566b3a90b2
diff --git a/src/main/java/com/ericsson/gerrit/plugins/highavailability/indexsync/IndexSyncRunner.java b/src/main/java/com/ericsson/gerrit/plugins/highavailability/indexsync/IndexSyncRunner.java
index 919b139..69bc2e4 100644
--- a/src/main/java/com/ericsson/gerrit/plugins/highavailability/indexsync/IndexSyncRunner.java
+++ b/src/main/java/com/ericsson/gerrit/plugins/highavailability/indexsync/IndexSyncRunner.java
@@ -47,6 +47,7 @@
 import dev.failsafe.function.CheckedSupplier;
 import java.io.IOException;
 import java.util.ArrayList;
+import java.util.HashSet;
 import java.util.List;
 import java.util.Set;
 import java.util.concurrent.ExecutionException;
@@ -106,48 +107,51 @@
     if (peers.size() == 0) {
       return false;
     }
-
-    ChangeIndexer indexer = changeIndexerFactory.create(executor, changeIndexes, false);
     boolean failed = false;
+    Set<String> updatedChanges = new HashSet<>();
     for (PeerInfo peer : peers) {
-      if (syncFrom(peer, indexer)) {
-        log.atFine().log("Finished indexSync for %s", peer.getDirectUrl());
-      } else {
-        log.atSevere().log("Failed to sync index with %s", peer.getDirectUrl());
+      try {
+        updatedChanges.addAll(collectUpdatedChanges(peer));
+      } catch (IOException e) {
+        log.atSevere().withCause(e).log("Error while querying changes from %s", peer);
         failed = true;
       }
     }
 
+    if (syncIndex(updatedChanges)) {
+      log.atFine().log("Finished indexSync");
+    } else {
+      log.atSevere().log("Failed to index out of sync changes");
+      failed = true;
+    }
+
     return !failed;
   }
 
-  private boolean syncFrom(PeerInfo peer, ChangeIndexer indexer) {
-    log.atFine().log("Syncing index with %s", peer.getDirectUrl());
+  private List<String> collectUpdatedChanges(PeerInfo peer) throws IOException {
+    log.atFine().log("Collecting out of sync changes from %s", peer.getDirectUrl());
     String peerUrl = peer.getDirectUrl();
     String uri =
         Joiner.on("/").join(peerUrl, pluginRelativePath, "query/changes.updated.since", age);
     HttpGet queryRequest = new HttpGet(uri);
-    List<String> ids;
-    try {
-      log.atFine().log("Executing %s", queryRequest);
-      ids = httpClient.execute(queryRequest, queryChangesResponseHandler);
-    } catch (IOException e) {
-      log.atSevere().withCause(e).log("Error while querying changes from %s", uri);
-      return false;
-    }
+    log.atFine().log("Executing %s", queryRequest);
+    return httpClient.execute(queryRequest, queryChangesResponseHandler);
+  }
 
+  private boolean syncIndex(Set<String> updatedChanges) {
+    ChangeIndexer indexer = changeIndexerFactory.create(executor, changeIndexes, false);
     try {
-      List<ListenableFuture<Boolean>> indexingTasks = new ArrayList<>(ids.size());
-      for (String id : ids) {
+      List<ListenableFuture<Boolean>> indexingTasks = new ArrayList<>(updatedChanges.size());
+      for (String id : updatedChanges) {
         indexingTasks.add(indexAsync(id, indexer));
       }
       Futures.allAsList(indexingTasks).get();
     } catch (InterruptedException | ExecutionException e) {
-      log.atSevere().withCause(e).log("Error while reindexing %s", ids);
+      log.atSevere().withCause(e).log("Error while reindexing %s", updatedChanges);
       return false;
     }
 
-    syncChangeDeletions(ids, indexer);
+    syncChangeDeletions(updatedChanges, indexer);
 
     return true;
   }
@@ -167,7 +171,7 @@
     return indexer.asyncReindexIfStale(projectName, Change.id(changeNumber));
   }
 
-  private void syncChangeDeletions(List<String> theirChanges, ChangeIndexer indexer) {
+  private void syncChangeDeletions(Set<String> theirChanges, ChangeIndexer indexer) {
     Set<String> ourChanges = queryLocalIndex();
     for (String d : Sets.difference(ourChanges, ImmutableSet.copyOf(theirChanges))) {
       deleteIfMissingInNoteDb(d, indexer);