Retry failovers by applying url distribution strategy Previously, a push with urlDistributionStrategy=roundRobin retried the same URL after a transport error - a single unreachable replica could block replication for any push the rotation happened to land on it. Fix it to apply the distribution strategy on retries. Change-Id: Id8ff42afdcd40a315fbab1ecb272b5e87bd0880d
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`.