Merge "Retry failovers by applying url distribution strategy"
diff --git a/src/main/java/com/googlesource/gerrit/plugins/replication/Destination.java b/src/main/java/com/googlesource/gerrit/plugins/replication/Destination.java
index c57f256..568bd0d 100644
--- a/src/main/java/com/googlesource/gerrit/plugins/replication/Destination.java
+++ b/src/main/java/com/googlesource/gerrit/plugins/replication/Destination.java
@@ -139,6 +139,7 @@
     boolean isRescheduled;
     boolean isFailed;
     RemoteRefUpdate.Status failedStatus;
+    URIish failoverTo;
   }
 
   private final ReplicationStateListener stateLog;
@@ -600,9 +601,7 @@
               // second one fails, it will also be rescheduled and then,
               // here, find out replication to its URI is already pending
               // for retry (blocking).
-              pendingPushOp.addRefBatches(pushOp.getRefs());
-              pendingPushOp.addStates(pushOp.getStates());
-              pushOp.removeStates();
+              consolidateOnto(pendingPushOp, pushOp);
 
             } else {
               // The one pending is one that is NOT retrying, it was just
@@ -619,9 +618,7 @@
               pendingPushOp.canceledByReplication();
               queue.pending.remove(uri);
 
-              pushOp.addRefBatches(pendingPushOp.getRefs());
-              pushOp.addStates(pendingPushOp.getStates());
-              pendingPushOp.removeStates();
+              consolidateOnto(pushOp, pendingPushOp);
             }
           }
 
@@ -642,11 +639,21 @@
                         : REJECTED_OTHER_REASON;
                 status.isFailed = true;
                 if (pushOp.setToRetry()) {
-                  status.isRescheduled = true;
-                  replicationTasksStorage.get().reset(pushOp);
-                  @SuppressWarnings("unused")
-                  ScheduledFuture<?> ignored2 =
-                      pool.schedule(pushOp, config.getRetryDelay(), TimeUnit.MINUTES);
+                  URIish nextUri =
+                      urlDistributor.failover(
+                          getURIs(pushOp.getProjectNameKey(), null), pushOp.getURI());
+                  if (!nextUri.equals(pushOp.getURI())) {
+                    // Defer actual failoverTo until after this write lock is released
+                    pushOp.canceledByReplication();
+                    queue.pending.remove(uri);
+                    status.failoverTo = nextUri;
+                  } else {
+                    status.isRescheduled = true;
+                    replicationTasksStorage.get().reset(pushOp);
+                    @SuppressWarnings("unused")
+                    ScheduledFuture<?> ignored2 =
+                        pool.schedule(pushOp, config.getRetryDelay(), TimeUnit.MINUTES);
+                  }
                 } else {
                   pushOp.canceledByReplication();
                   pushOp.retryDone();
@@ -667,6 +674,51 @@
     if (status.isRescheduled) {
       postReplicationScheduledEvent(pushOp);
     }
+    if (status.failoverTo != null) {
+      failoverTo(pushOp, status.failoverTo);
+    }
+  }
+
+  private void failoverTo(PushOne pushOp, URIish newUri) {
+    PushOne replacement = opFactory.create(pushOp.getProjectNameKey(), newUri);
+    replacement.addRefBatches(pushOp.getRefs());
+    replacement.addStates(pushOp.getStates());
+    replacement.setToRetryWithCount(pushOp.getRetryCount());
+    pushOp.removeStates();
+
+    boolean installed =
+        stateLock.withWriteLock(
+            newUri,
+            () -> {
+              for (ReplicationTasksStorage.ReplicateRefUpdate u :
+                  replacement.getReplicateRefUpdates()) {
+                replicationTasksStorage.get().create(u);
+              }
+              PushOne existing = getPendingPush(newUri);
+              if (existing != null) {
+                consolidateOnto(existing, replacement);
+                replicationTasksStorage.get().finish(pushOp);
+                return false;
+              }
+              replicationTasksStorage.get().finish(pushOp);
+              queue.pending.put(newUri, replacement);
+              @SuppressWarnings("unused")
+              ScheduledFuture<?> ignored =
+                  pool.schedule(replacement, config.getRetryDelay(), TimeUnit.MINUTES);
+              return true;
+            });
+    repLog.atInfo().log(
+        "Failover: replication for %s rerouted from %s to %s (retry %d)",
+        pushOp.getProjectNameKey(), pushOp.getURI(), newUri, replacement.getRetryCount());
+    if (installed) {
+      postReplicationScheduledEvent(replacement);
+    }
+  }
+
+  private static void consolidateOnto(PushOne into, PushOne from) {
+    into.addRefBatches(from.getRefs());
+    into.addStates(from.getStates());
+    from.removeStates();
   }
 
   RunwayStatus requestRunway(PushOne op) {
diff --git a/src/main/java/com/googlesource/gerrit/plugins/replication/PushOne.java b/src/main/java/com/googlesource/gerrit/plugins/replication/PushOne.java
index 0e04e1e..600ce8c 100644
--- a/src/main/java/com/googlesource/gerrit/plugins/replication/PushOne.java
+++ b/src/main/java/com/googlesource/gerrit/plugins/replication/PushOne.java
@@ -263,11 +263,19 @@
   }
 
   boolean setToRetry() {
+    return setToRetryWithCount(retryCount + 1);
+  }
+
+  boolean setToRetryWithCount(int count) {
     retrying = true;
-    retryCount++;
+    retryCount = count;
     return maxRetries == 0 || retryCount <= maxRetries;
   }
 
+  int getRetryCount() {
+    return retryCount;
+  }
+
   void retryDone() {
     this.retrying = false;
   }
diff --git a/src/main/java/com/googlesource/gerrit/plugins/replication/UrlDistributionStrategy.java b/src/main/java/com/googlesource/gerrit/plugins/replication/UrlDistributionStrategy.java
index 46183eb..1230e36 100644
--- a/src/main/java/com/googlesource/gerrit/plugins/replication/UrlDistributionStrategy.java
+++ b/src/main/java/com/googlesource/gerrit/plugins/replication/UrlDistributionStrategy.java
@@ -41,17 +41,36 @@
    * Push to one URL at a time, rotating through the list on each push event. Particularly useful
    * when multiple replica hosts share a single backend (likely via NFS): pushing to all URLs would
    * cause redundant writes to the same underlying storage, while round-robin distributes load
-   * evenly and ensures each push is executed exactly once.
+   * evenly and ensures each push is executed exactly once. On transport failure {@link
+   * Instance#failover} hands the push over to the next URL in the rotation.
    */
   ROUND_ROBIN("roundRobin") {
     @Override
     public Instance newInstance() {
-      final AtomicInteger index = new AtomicInteger();
-      return candidates -> {
-        if (candidates.isEmpty()) {
-          return List.of();
+      return new Instance() {
+        private final AtomicInteger index = new AtomicInteger();
+
+        @Override
+        public List<URIish> select(List<URIish> candidates) {
+          if (candidates.isEmpty()) {
+            return List.of();
+          }
+          return List.of(candidates.get(Math.floorMod(index.getAndIncrement(), candidates.size())));
         }
-        return List.of(candidates.get(Math.floorMod(index.getAndIncrement(), candidates.size())));
+
+        @Override
+        public URIish failover(List<URIish> candidates, URIish failed) {
+          if (candidates.size() < 2) {
+            return failed;
+          }
+          for (int attempt = 0; attempt < candidates.size(); attempt++) {
+            URIish next = candidates.get(Math.floorMod(index.getAndIncrement(), candidates.size()));
+            if (!next.equals(failed)) {
+              return next;
+            }
+          }
+          return failed;
+        }
       };
     }
   };
@@ -79,6 +98,15 @@
   /** A stateful executor for a {@link UrlDistributionStrategy} strategy. */
   @FunctionalInterface
   public interface Instance {
+    /** Select the URLs to push to for this scheduling event. */
     List<URIish> select(List<URIish> candidates);
+
+    /**
+     * If a push to any URI returned by {@link #select(List)} fails, {@link #failover(List,
+     * URIish)}} is invoked to select the next URI for retry.
+     */
+    default URIish failover(List<URIish> candidates, URIish failed) {
+      return failed;
+    }
   }
 }
diff --git a/src/main/resources/Documentation/config.md b/src/main/resources/Documentation/config.md
index cef009c..cffaa8c 100644
--- a/src/main/resources/Documentation/config.md
+++ b/src/main/resources/Documentation/config.md
@@ -700,7 +700,10 @@
     hosts share a single NFS backend. Pushing to all URLs simultaneously
     would cause redundant writes to the same underlying storage; round-robin
     distributes load evenly and ensures each push is written exactly once.
-    Has no effect if only one URL is configured.
+    On a transport error during retry, the next push attempt fails over to a
+    different URL in the rotation (bounded by `replicationRetry`), so a single
+    unreachable host does not block replication. Has no effect if only one URL
+    is configured.
 
   Defaults to `all`.