Add projectSharded URL selection per remote Extend `remote.NAME.urlDistributionStrategy` with a `projectSharded` value. Like `roundRobin` it picks a single URL per push, but derives it from the project name, so a project always replicates to the same URL instead of rotating on every event. Replication tasks are coalesced per (project, URI), so rotating URLs leaves successive updates to one project as separate tasks racing against the shared NFS backend. Pinning collapses them into one task. The URL is chosen by hashing the project name over the sorted candidate list, so every host agrees regardless of config order. Change-Id: I91118a47e0e7da64995a9a444d4baf56ec8b8386
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 1c2bc94..c895763 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/replication/Destination.java +++ b/src/main/java/com/googlesource/gerrit/plugins/replication/Destination.java
@@ -803,11 +803,11 @@ } List<URIish> getDistributedUris(Project.NameKey project, String urlMatch) { - return getDistributedUris(getURIs(project, urlMatch)); + return getDistributedUris(project, getURIs(project, urlMatch)); } - List<URIish> getDistributedUris(List<URIish> candidates) { - return urlDistributor.select(candidates); + List<URIish> getDistributedUris(Project.NameKey project, List<URIish> candidates) { + return urlDistributor.select(project, candidates); } URIish getURI(URIish template, Project.NameKey project) throws URISyntaxException {
diff --git a/src/main/java/com/googlesource/gerrit/plugins/replication/DestinationsCollection.java b/src/main/java/com/googlesource/gerrit/plugins/replication/DestinationsCollection.java index 82f33d7..58ebed8 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/replication/DestinationsCollection.java +++ b/src/main/java/com/googlesource/gerrit/plugins/replication/DestinationsCollection.java
@@ -138,7 +138,7 @@ validUris.add(uri); } } - config.getDistributedUris(validUris).forEach(uri -> uris.put(config, uri)); + config.getDistributedUris(projectName, validUris).forEach(uri -> uris.put(config, uri)); } return uris; }
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..ddff448 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/replication/UrlDistributionStrategy.java +++ b/src/main/java/com/googlesource/gerrit/plugins/replication/UrlDistributionStrategy.java
@@ -14,7 +14,10 @@ package com.googlesource.gerrit.plugins.replication; +import com.google.common.collect.ImmutableList; +import com.google.gerrit.entities.Project; import java.util.Arrays; +import java.util.Comparator; import java.util.List; import java.util.concurrent.atomic.AtomicInteger; import org.eclipse.jgit.transport.URIish; @@ -33,7 +36,7 @@ ALL("all") { @Override public Instance newInstance() { - return candidates -> candidates; + return (project, candidates) -> candidates; } }, @@ -47,13 +50,41 @@ @Override public Instance newInstance() { final AtomicInteger index = new AtomicInteger(); - return candidates -> { + return (project, candidates) -> { if (candidates.isEmpty()) { return List.of(); } return List.of(candidates.get(Math.floorMod(index.getAndIncrement(), candidates.size()))); }; } + }, + + /** + * Push to exactly one URL, chosen by hashing the project name so that a given project always maps + * to the same URL. Like {@link #ROUND_ROBIN} this writes each push only once, which matters when + * the replica hosts share a single backend, but it additionally keeps consecutive updates for one + * project on the same URL. Because replication tasks are coalesced per (project, URI), rotating + * URLs would let successive updates for one project run as separate tasks racing against the same + * backend; pinning the project collapses them into a single task and keeps the receiving host's + * caches warm. + * + * <p>The candidates are sorted before indexing so that the mapping does not depend on the order + * in which the URLs happen to be configured. Together with {@link Project.NameKey#hashCode()}, + * which is the specified {@link String#hashCode()} of the project name, this makes every host + * reading the same config agree on the mapping. + */ + PROJECT_SHARDED("projectSharded") { + @Override + public Instance newInstance() { + return (project, candidates) -> { + if (candidates.isEmpty()) { + return List.of(); + } + ImmutableList<URIish> sorted = + ImmutableList.sortedCopyOf(Comparator.comparing(URIish::toString), candidates); + return List.of(sorted.get(Math.floorMod(project.hashCode(), sorted.size()))); + }; + } }; public final String configKey; @@ -79,6 +110,13 @@ /** A stateful executor for a {@link UrlDistributionStrategy} strategy. */ @FunctionalInterface public interface Instance { - List<URIish> select(List<URIish> candidates); + /** + * Selects the URLs to push to out of the candidates for the given project. + * + * @param project project being replicated, used by project-affine strategies. + * @param candidates URLs the project could be pushed to. + * @return the subset of candidates to push to. + */ + List<URIish> select(Project.NameKey project, List<URIish> candidates); } }
diff --git a/src/main/resources/Documentation/config.md b/src/main/resources/Documentation/config.md index 73b3805..ac71ed6 100644 --- a/src/main/resources/Documentation/config.md +++ b/src/main/resources/Documentation/config.md
@@ -628,6 +628,24 @@ distributes load evenly and ensures each push is written exactly once. 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 + URL. Like `roundRobin` each push is written exactly once, and load is + spread across the URLs, but consecutive updates for one project always + go to the same URL. + + Prefer this over `roundRobin` when replica hosts share a single backend. + Replication tasks are coalesced per (project, URL), so rotating URLs + makes successive updates for one project run as separate tasks that race + against the same backend; pinning a project to one URL collapses them + into a single task and keeps the receiving host's caches warm. + + The mapping is computed from a hash of the project name over the sorted + 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. + Has no effect if only one URL is configured. + Defaults to `all`. Directory `replication`
diff --git a/src/test/java/com/googlesource/gerrit/plugins/replication/DestinationConfigurationTest.java b/src/test/java/com/googlesource/gerrit/plugins/replication/DestinationConfigurationTest.java index c7b4174..f5118dd 100644 --- a/src/test/java/com/googlesource/gerrit/plugins/replication/DestinationConfigurationTest.java +++ b/src/test/java/com/googlesource/gerrit/plugins/replication/DestinationConfigurationTest.java
@@ -114,4 +114,16 @@ assertThat(objectUnderTest.getUrlDistributionStrategy()) .isEqualTo(UrlDistributionStrategy.ROUND_ROBIN); } + + @Test + public void shouldSetUrlDistributionToProjectShardedWhenConfigured() { + // given + when(cfgMock.getString("remote", REMOTE, "urlDistributionStrategy")) + .thenReturn("projectSharded"); + objectUnderTest = new DestinationConfiguration(remoteConfigMock, cfgMock); + + // when / then + assertThat(objectUnderTest.getUrlDistributionStrategy()) + .isEqualTo(UrlDistributionStrategy.PROJECT_SHARDED); + } }
diff --git a/src/test/java/com/googlesource/gerrit/plugins/replication/ReplicationIT.java b/src/test/java/com/googlesource/gerrit/plugins/replication/ReplicationIT.java index 3b3686f..40f485f 100644 --- a/src/test/java/com/googlesource/gerrit/plugins/replication/ReplicationIT.java +++ b/src/test/java/com/googlesource/gerrit/plugins/replication/ReplicationIT.java
@@ -551,6 +551,81 @@ } @Test + public void shouldReplicateToOnlyOneUrlWhenProjectShardedEnabled() throws Exception { + Project.NameKey replica1Project = createTestProject(project + "replica1"); + Project.NameKey replica2Project = createTestProject(project + "replica2"); + + setReplicationDestination( + "foo", List.of("replica1", "replica2"), ALL_PROJECTS, TEST_REPLICATION_DELAY_SECONDS); + setUrlDistribution("foo", UrlDistributionStrategy.PROJECT_SHARDED); + reloadConfig(); + + String newRef = "refs/heads/newForTest"; + createNewBranchWithoutPush("refs/heads/master", newRef); + + plugin + .getSysInjector() + .getInstance(ReplicationQueue.class) + .scheduleFullSync(project, null, new ReplicationState(NO_OP), true); + + // Wait for the push to land in at least one replica + try (Repository r1 = repoManager.openRepository(replica1Project); + Repository r2 = repoManager.openRepository(replica2Project)) { + waitUntil(() -> checkedGetRef(r1, newRef) != null || checkedGetRef(r2, newRef) != null); + + // Exactly one replica should have received the push + boolean r1HasRef = checkedGetRef(r1, newRef) != null; + boolean r2HasRef = checkedGetRef(r2, newRef) != null; + assertThat(r1HasRef ^ r2HasRef).isTrue(); + } + } + + @Test + public void shouldReuseSameUrlOnConsecutivePushesWhenProjectShardedEnabled() throws Exception { + Project.NameKey replica1Project = createTestProject(project + "replica1"); + Project.NameKey replica2Project = createTestProject(project + "replica2"); + + setReplicationDestination( + "foo", List.of("replica1", "replica2"), ALL_PROJECTS, TEST_REPLICATION_DELAY_SECONDS); + setUrlDistribution("foo", UrlDistributionStrategy.PROJECT_SHARDED); + reloadConfig(); + + String branch1 = "refs/heads/branch1"; + String branch2 = "refs/heads/branch2"; + createNewBranchWithoutPush("refs/heads/master", branch1); + + ReplicationQueue queue = plugin.getSysInjector().getInstance(ReplicationQueue.class); + + queue.scheduleFullSync(project, null, new ReplicationState(NO_OP), true); + + // The project is pinned to one of the two replicas; discover which one. + Project.NameKey pinnedProject; + Project.NameKey otherProject; + try (Repository r1 = repoManager.openRepository(replica1Project); + Repository r2 = repoManager.openRepository(replica2Project)) { + waitUntil(() -> checkedGetRef(r1, branch1) != null || checkedGetRef(r2, branch1) != null); + + boolean replica1IsPinned = checkedGetRef(r1, branch1) != null; + assertThat(replica1IsPinned ^ (checkedGetRef(r2, branch1) != null)).isTrue(); + pinnedProject = replica1IsPinned ? replica1Project : replica2Project; + otherProject = replica1IsPinned ? replica2Project : replica1Project; + } + + // Second sync must land on the same replica rather than rotating to the other one + createNewBranchWithoutPush("refs/heads/master", branch2); + queue.scheduleFullSync(project, null, new ReplicationState(NO_OP), true); + + try (Repository pinnedRepo = repoManager.openRepository(pinnedProject)) { + waitUntil(() -> checkedGetRef(pinnedRepo, branch2) != null); + } + + try (Repository otherRepo = repoManager.openRepository(otherProject)) { + assertThat(checkedGetRef(otherRepo, branch1)).isNull(); + assertThat(checkedGetRef(otherRepo, branch2)).isNull(); + } + } + + @Test public void shouldReplicateToMatchingRemote() throws Exception { Project.NameKey targetProject = createTestProject(project + "replica");
diff --git a/src/test/java/com/googlesource/gerrit/plugins/replication/UrlDistributionStrategyTest.java b/src/test/java/com/googlesource/gerrit/plugins/replication/UrlDistributionStrategyTest.java new file mode 100644 index 0000000..2ded31e --- /dev/null +++ b/src/test/java/com/googlesource/gerrit/plugins/replication/UrlDistributionStrategyTest.java
@@ -0,0 +1,170 @@ +// Copyright (C) 2026 The Android Open Source Project +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package com.googlesource.gerrit.plugins.replication; + +import static com.google.common.truth.Truth.assertThat; + +import com.google.common.collect.ImmutableList; +import com.google.gerrit.entities.Project; +import java.net.URISyntaxException; +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashSet; +import java.util.List; +import java.util.Random; +import java.util.Set; +import org.eclipse.jgit.transport.URIish; +import org.junit.Test; + +public class UrlDistributionStrategyTest { + private static final Project.NameKey PROJECT = Project.nameKey("some/project"); + private static final Project.NameKey OTHER_PROJECT = Project.nameKey("some/other/project"); + + private static ImmutableList<URIish> uris(String... hosts) throws URISyntaxException { + ImmutableList.Builder<URIish> uris = ImmutableList.builder(); + for (String host : hosts) { + uris.add(new URIish("git://" + host + "/${name}.git")); + } + return uris.build(); + } + + @Test + public void shouldResolveConfigKeys() { + assertThat(UrlDistributionStrategy.fromConfig("all")).isEqualTo(UrlDistributionStrategy.ALL); + assertThat(UrlDistributionStrategy.fromConfig("roundRobin")) + .isEqualTo(UrlDistributionStrategy.ROUND_ROBIN); + assertThat(UrlDistributionStrategy.fromConfig("projectSharded")) + .isEqualTo(UrlDistributionStrategy.PROJECT_SHARDED); + } + + @Test + public void shouldFallBackToAllForUnknownConfigKey() { + assertThat(UrlDistributionStrategy.fromConfig("bogus")).isEqualTo(UrlDistributionStrategy.ALL); + assertThat(UrlDistributionStrategy.fromConfig(null)).isEqualTo(UrlDistributionStrategy.ALL); + } + + @Test + public void allShouldSelectEveryCandidate() throws Exception { + ImmutableList<URIish> candidates = uris("replica1", "replica2", "replica3"); + UrlDistributionStrategy.Instance distributor = UrlDistributionStrategy.ALL.newInstance(); + + assertThat(distributor.select(PROJECT, candidates)).isEqualTo(candidates); + } + + @Test + public void roundRobinShouldRotateOnConsecutiveSelections() throws Exception { + ImmutableList<URIish> candidates = uris("replica1", "replica2"); + UrlDistributionStrategy.Instance distributor = + UrlDistributionStrategy.ROUND_ROBIN.newInstance(); + + assertThat(distributor.select(PROJECT, candidates)).containsExactly(candidates.get(0)); + assertThat(distributor.select(PROJECT, candidates)).containsExactly(candidates.get(1)); + assertThat(distributor.select(PROJECT, candidates)).containsExactly(candidates.get(0)); + } + + @Test + public void roundRobinShouldRotateRegardlessOfProject() throws Exception { + ImmutableList<URIish> candidates = uris("replica1", "replica2"); + UrlDistributionStrategy.Instance distributor = + UrlDistributionStrategy.ROUND_ROBIN.newInstance(); + + assertThat(distributor.select(PROJECT, candidates)).containsExactly(candidates.get(0)); + assertThat(distributor.select(OTHER_PROJECT, candidates)).containsExactly(candidates.get(1)); + } + + @Test + public void projectShardedShouldSelectExactlyOneCandidate() throws Exception { + ImmutableList<URIish> candidates = uris("replica1", "replica2", "replica3"); + UrlDistributionStrategy.Instance distributor = + UrlDistributionStrategy.PROJECT_SHARDED.newInstance(); + + List<URIish> selected = distributor.select(PROJECT, candidates); + + assertThat(selected).hasSize(1); + assertThat(candidates).containsAtLeastElementsIn(selected); + } + + @Test + public void projectShardedShouldPinProjectToSameCandidateAcrossSelections() throws Exception { + ImmutableList<URIish> candidates = uris("replica1", "replica2", "replica3"); + UrlDistributionStrategy.Instance distributor = + UrlDistributionStrategy.PROJECT_SHARDED.newInstance(); + + List<URIish> first = distributor.select(PROJECT, candidates); + for (int i = 0; i < 10; i++) { + assertThat(distributor.select(PROJECT, candidates)).isEqualTo(first); + } + } + + @Test + public void projectShardedShouldAgreeAcrossInstances() throws Exception { + ImmutableList<URIish> candidates = uris("replica1", "replica2", "replica3"); + + List<URIish> first = + UrlDistributionStrategy.PROJECT_SHARDED.newInstance().select(PROJECT, candidates); + List<URIish> second = + UrlDistributionStrategy.PROJECT_SHARDED.newInstance().select(PROJECT, candidates); + + assertThat(second).isEqualTo(first); + } + + @Test + public void projectShardedShouldIgnoreCandidateOrdering() throws Exception { + ImmutableList<URIish> candidates = + uris("replica1", "replica2", "replica3", "replica4", "replica5"); + UrlDistributionStrategy.Instance distributor = + UrlDistributionStrategy.PROJECT_SHARDED.newInstance(); + + List<URIish> expected = distributor.select(PROJECT, candidates); + + List<URIish> shuffled = new ArrayList<>(candidates); + Random random = new Random(42); + for (int i = 0; i < 20; i++) { + Collections.shuffle(shuffled, random); + assertThat(distributor.select(PROJECT, shuffled)).isEqualTo(expected); + } + } + + @Test + public void projectShardedShouldSpreadProjectsOverAllCandidates() throws Exception { + ImmutableList<URIish> candidates = uris("replica1", "replica2", "replica3"); + UrlDistributionStrategy.Instance distributor = + UrlDistributionStrategy.PROJECT_SHARDED.newInstance(); + + Set<URIish> selected = new HashSet<>(); + for (int i = 0; i < 100; i++) { + selected.addAll(distributor.select(Project.nameKey("project" + i), candidates)); + } + + assertThat(selected).containsExactlyElementsIn(candidates); + } + + @Test + public void projectShardedShouldSelectOnlyCandidateWhenSingleUrlConfigured() throws Exception { + ImmutableList<URIish> candidates = uris("replica1"); + UrlDistributionStrategy.Instance distributor = + UrlDistributionStrategy.PROJECT_SHARDED.newInstance(); + + assertThat(distributor.select(PROJECT, candidates)).isEqualTo(candidates); + assertThat(distributor.select(OTHER_PROJECT, candidates)).isEqualTo(candidates); + } + + @Test + public void shouldSelectNothingWhenThereAreNoCandidates() { + for (UrlDistributionStrategy strategy : UrlDistributionStrategy.values()) { + assertThat(strategy.newInstance().select(PROJECT, List.of())).isEmpty(); + } + } +}