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. This cherry-pick is modified to also introduce the failover strategy for projectSharded distribution strategy. Change-Id: Id8ff42afdcd40a315fbab1ecb272b5e87bd0880d (cherry picked from commit 28ee207f72e3acb7910ab39d0a033c9da38b699c)
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 c895763..873ec82 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; @@ -604,9 +605,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 @@ -623,9 +622,7 @@ pendingPushOp.canceledByReplication(); queue.pending.remove(uri); - pushOp.addRefBatches(pendingPushOp.getRefs()); - pushOp.addStates(pendingPushOp.getStates()); - pendingPushOp.removeStates(); + consolidateOnto(pushOp, pendingPushOp); } } @@ -646,11 +643,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(); @@ -671,6 +678,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 f9e1081..ad516a6 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 ddff448..50d6504 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/replication/UrlDistributionStrategy.java +++ b/src/main/java/com/googlesource/gerrit/plugins/replication/UrlDistributionStrategy.java
@@ -16,6 +16,7 @@ import com.google.common.collect.ImmutableList; import com.google.gerrit.entities.Project; +import com.google.gerrit.entities.Project.NameKey; import java.util.Arrays; import java.util.Comparator; import java.util.List; @@ -44,17 +45,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 (project, candidates) -> { - if (candidates.isEmpty()) { - return List.of(); + return new Instance() { + private final AtomicInteger index = new AtomicInteger(); + + @Override + public List<URIish> select(NameKey project, 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; + } }; } }, @@ -76,13 +96,32 @@ PROJECT_SHARDED("projectSharded") { @Override public Instance newInstance() { - return (project, candidates) -> { - if (candidates.isEmpty()) { - return List.of(); + return new Instance() { + @Override + public List<URIish> select(NameKey project, List<URIish> candidates) { + if (candidates.isEmpty()) { + return List.of(); + } + ImmutableList<URIish> sorted = sortedByUrl(candidates); + return List.of(sorted.get(Math.floorMod(project.hashCode(), sorted.size()))); } - ImmutableList<URIish> sorted = - ImmutableList.sortedCopyOf(Comparator.comparing(URIish::toString), candidates); - return List.of(sorted.get(Math.floorMod(project.hashCode(), sorted.size()))); + + @Override + public URIish failover(List<URIish> candidates, URIish failed) { + if (candidates.size() < 2) { + return failed; + } + ImmutableList<URIish> sorted = sortedByUrl(candidates); + int failedIndex = sorted.indexOf(failed); + if (failedIndex < 0) { + return failed; + } + return sorted.get((failedIndex + 1) % sorted.size()); + } + + private ImmutableList<URIish> sortedByUrl(List<URIish> candidates) { + return ImmutableList.sortedCopyOf(Comparator.comparing(URIish::toString), candidates); + } }; } }; @@ -118,5 +157,13 @@ * @return the subset of candidates to push to. */ List<URIish> select(Project.NameKey project, List<URIish> candidates); + + /** + * If a push to any URI returned by {@link #select(Project.NameKey, 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 ac71ed6..f49a5f5 100644 --- a/src/main/resources/Documentation/config.md +++ b/src/main/resources/Documentation/config.md
@@ -626,7 +626,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. `projectSharded` : Push to one URL, chosen so that a given project always maps to the same @@ -644,6 +647,12 @@ list of URLs, so it does not depend on the order in which the URLs are configured and is identical on every host reading the same config. Adding or removing a URL remaps projects across the remaining URLs. + + 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. + Has no effect if only one URL is configured. Defaults to `all`.