Merge "Destination: remove unused method"
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 6a90f50..c57f256 100644
--- a/src/main/java/com/googlesource/gerrit/plugins/replication/Destination.java
+++ b/src/main/java/com/googlesource/gerrit/plugins/replication/Destination.java
@@ -43,7 +43,6 @@
import com.google.gerrit.server.account.GroupBackends;
import com.google.gerrit.server.account.GroupIncludeCache;
import com.google.gerrit.server.account.ListGroupMembership;
-import com.google.gerrit.server.events.EventDispatcher;
import com.google.gerrit.server.git.GitRepositoryManager;
import com.google.gerrit.server.git.PerThreadRequestScope;
import com.google.gerrit.server.git.WorkQueue;
@@ -68,6 +67,7 @@
import com.googlesource.gerrit.plugins.replication.events.ProjectDeletionState;
import com.googlesource.gerrit.plugins.replication.events.RefReplicatedEvent;
import com.googlesource.gerrit.plugins.replication.events.ReplicationScheduledEvent;
+import com.googlesource.gerrit.plugins.replication.events.dispatcher.EventDispatcher;
import java.io.IOException;
import java.net.URISyntaxException;
import java.util.Collection;
@@ -156,8 +156,9 @@
private volatile ScheduledExecutorService pool;
private final PerThreadRequestScope.Scoper threadScoper;
private final DestinationConfiguration config;
- private final DynamicItem<EventDispatcher> eventDispatcher;
+ private final EventDispatcher eventDispatcher;
private final Provider<ReplicationTasksStorage> replicationTasksStorage;
+ private final UrlDistributionStrategy.Instance urlDistributor;
protected enum RetryReason {
TRANSPORT_ERROR,
@@ -186,7 +187,7 @@
GroupBackend groupBackend,
ReplicationStateListeners stateLog,
GroupIncludeCache groupIncludeCache,
- DynamicItem<EventDispatcher> eventDispatcher,
+ EventDispatcher eventDispatcher,
Provider<ReplicationTasksStorage> rts,
CredentialsFactory credentialsFactory,
@Assisted DestinationConfiguration cfg) {
@@ -200,6 +201,7 @@
this.replicationTasksStorage = rts;
this.credentialsFactory = credentialsFactory;
config = cfg;
+ urlDistributor = cfg.getUrlDistributionStrategy().newInstance();
ImmutableList<String> projects = cfg.getProjects();
int numStripes = projects.isEmpty() ? MAX_STRIPES : Math.min(projects.size(), MAX_STRIPES);
@@ -797,6 +799,14 @@
return r;
}
+ List<URIish> getDistributedUris(Project.NameKey project, String urlMatch) {
+ return getDistributedUris(getURIs(project, urlMatch));
+ }
+
+ List<URIish> getDistributedUris(List<URIish> candidates) {
+ return urlDistributor.select(candidates);
+ }
+
URIish getURI(URIish template, Project.NameKey project) throws URISyntaxException {
return getURI(template, project, config.getRemoteNameStyle(), config.isSingleProjectMatch());
}
@@ -925,7 +935,7 @@
ReplicationScheduledEvent event =
new ReplicationScheduledEvent(project.get(), ref, pushOp.getURI());
try {
- eventDispatcher.get().postEvent(BranchNameKey.create(project, ref), event);
+ eventDispatcher.postEvent(BranchNameKey.create(project, ref), event);
} catch (PermissionBackendException e) {
repLog.atSevere().withCause(e).log("error posting event");
}
@@ -938,7 +948,7 @@
RefReplicatedEvent event =
new RefReplicatedEvent(project.get(), ref, pushOp.getURI(), RefPushResult.FAILED, status);
try {
- eventDispatcher.get().postEvent(BranchNameKey.create(project, ref), event);
+ eventDispatcher.postEvent(BranchNameKey.create(project, ref), event);
} catch (PermissionBackendException e) {
repLog.atSevere().withCause(e).log("error posting event");
}
diff --git a/src/main/java/com/googlesource/gerrit/plugins/replication/DestinationConfiguration.java b/src/main/java/com/googlesource/gerrit/plugins/replication/DestinationConfiguration.java
index fc73d15..03ba914 100644
--- a/src/main/java/com/googlesource/gerrit/plugins/replication/DestinationConfiguration.java
+++ b/src/main/java/com/googlesource/gerrit/plugins/replication/DestinationConfiguration.java
@@ -61,6 +61,7 @@
private final String uploadPack;
private final String receivePack;
private final String gitPath;
+ private final UrlDistributionStrategy urlDistributionStrategy;
protected DestinationConfiguration(RemoteConfig remoteConfig, Config cfg) {
this.remoteConfig = remoteConfig;
@@ -132,6 +133,9 @@
uploadPack = cfg.getString("remote", name, "uploadpack");
receivePack = cfg.getString("remote", name, "receivepack");
gitPath = cfg.getString("remote", name, "gitPath");
+ urlDistributionStrategy =
+ UrlDistributionStrategy.fromConfig(
+ cfg.getString("remote", name, "urlDistributionStrategy"));
}
@Override
@@ -259,6 +263,11 @@
return gitPath;
}
+ @Override
+ public UrlDistributionStrategy getUrlDistributionStrategy() {
+ return urlDistributionStrategy;
+ }
+
private ImmutableList<Pattern> getExcludedRefsPattern(Config cfg, String name) {
List<Pattern> patterns = new ArrayList<>();
for (String regex : cfg.getStringList("remote", name, "excludedRefsPattern")) {
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 6eb4d3f..4896dcb 100644
--- a/src/main/java/com/googlesource/gerrit/plugins/replication/DestinationsCollection.java
+++ b/src/main/java/com/googlesource/gerrit/plugins/replication/DestinationsCollection.java
@@ -103,6 +103,7 @@
}
boolean adminURLUsed = false;
+ List<URIish> validUris = new ArrayList<>();
for (String url : config.getAdminUrls()) {
if (Strings.isNullOrEmpty(url)) {
@@ -134,16 +135,17 @@
}
}
if (matchesConfigUrl || Destination.matches(uri, urlMatch)) {
- uris.put(config, uri);
+ validUris.add(uri);
adminURLUsed = true;
}
}
if (!adminURLUsed) {
for (URIish uri : config.getURIs(projectName, urlMatch)) {
- uris.put(config, uri);
+ validUris.add(uri);
}
}
+ config.getDistributedUris(validUris).forEach(uri -> uris.put(config, uri));
}
return uris;
}
diff --git a/src/main/java/com/googlesource/gerrit/plugins/replication/OnStartStop.java b/src/main/java/com/googlesource/gerrit/plugins/replication/OnStartStop.java
index c779857..ecc55a2 100644
--- a/src/main/java/com/googlesource/gerrit/plugins/replication/OnStartStop.java
+++ b/src/main/java/com/googlesource/gerrit/plugins/replication/OnStartStop.java
@@ -16,12 +16,11 @@
import com.google.common.util.concurrent.Atomics;
import com.google.gerrit.extensions.events.LifecycleListener;
-import com.google.gerrit.extensions.registration.DynamicItem;
import com.google.gerrit.extensions.systemstatus.ServerInformation;
-import com.google.gerrit.server.events.EventDispatcher;
import com.google.inject.Inject;
import com.googlesource.gerrit.plugins.replication.PushResultProcessing.GitUpdateProcessing;
import com.googlesource.gerrit.plugins.replication.api.ReplicationConfig;
+import com.googlesource.gerrit.plugins.replication.events.dispatcher.EventDispatcher;
import java.util.Set;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
@@ -32,14 +31,14 @@
private final ServerInformation srvInfo;
private final PushAll.Factory pushAll;
private final ReplicationConfig config;
- private final DynamicItem<EventDispatcher> eventDispatcher;
+ private final EventDispatcher eventDispatcher;
@Inject
protected OnStartStop(
ServerInformation srvInfo,
PushAll.Factory pushAll,
ReplicationConfig config,
- DynamicItem<EventDispatcher> eventDispatcher) {
+ EventDispatcher eventDispatcher) {
this.srvInfo = srvInfo;
this.pushAll = pushAll;
this.config = config;
@@ -51,7 +50,7 @@
public void start() {
if (srvInfo.getState() == ServerInformation.State.STARTUP
&& config.isReplicateAllOnPluginStart()) {
- ReplicationState state = new ReplicationState(new GitUpdateProcessing(eventDispatcher.get()));
+ ReplicationState state = new ReplicationState(new GitUpdateProcessing(eventDispatcher));
pushAllFuture.set(
pushAll
.create(null, Set.of(), ReplicationFilter.all(), state, false)
diff --git a/src/main/java/com/googlesource/gerrit/plugins/replication/PushResultProcessing.java b/src/main/java/com/googlesource/gerrit/plugins/replication/PushResultProcessing.java
index 8d279d5..509fbcc 100644
--- a/src/main/java/com/googlesource/gerrit/plugins/replication/PushResultProcessing.java
+++ b/src/main/java/com/googlesource/gerrit/plugins/replication/PushResultProcessing.java
@@ -16,12 +16,12 @@
import com.google.common.flogger.FluentLogger;
import com.google.gerrit.exceptions.StorageException;
-import com.google.gerrit.server.events.EventDispatcher;
import com.google.gerrit.server.events.RefEvent;
import com.google.gerrit.server.permissions.PermissionBackendException;
import com.googlesource.gerrit.plugins.replication.ReplicationState.RefPushResult;
import com.googlesource.gerrit.plugins.replication.events.RefReplicatedEvent;
import com.googlesource.gerrit.plugins.replication.events.RefReplicationDoneEvent;
+import com.googlesource.gerrit.plugins.replication.events.dispatcher.EventDispatcher;
import java.lang.ref.WeakReference;
import java.util.concurrent.atomic.AtomicBoolean;
import org.eclipse.jgit.transport.RemoteRefUpdate;
diff --git a/src/main/java/com/googlesource/gerrit/plugins/replication/RemoteConfiguration.java b/src/main/java/com/googlesource/gerrit/plugins/replication/RemoteConfiguration.java
index bcd9ace..79bdbf3 100644
--- a/src/main/java/com/googlesource/gerrit/plugins/replication/RemoteConfiguration.java
+++ b/src/main/java/com/googlesource/gerrit/plugins/replication/RemoteConfiguration.java
@@ -166,4 +166,9 @@
* @return true if new repositories should store ref-updates in their reflog
*/
boolean storeRefLog();
+
+ /** Returns the URL distribution strategy configured for this remote. */
+ default UrlDistributionStrategy getUrlDistributionStrategy() {
+ return UrlDistributionStrategy.ALL;
+ }
}
diff --git a/src/main/java/com/googlesource/gerrit/plugins/replication/ReplicationModule.java b/src/main/java/com/googlesource/gerrit/plugins/replication/ReplicationModule.java
index bae633d..60c0254 100644
--- a/src/main/java/com/googlesource/gerrit/plugins/replication/ReplicationModule.java
+++ b/src/main/java/com/googlesource/gerrit/plugins/replication/ReplicationModule.java
@@ -40,6 +40,9 @@
import com.googlesource.gerrit.plugins.replication.events.RefReplicatedEvent;
import com.googlesource.gerrit.plugins.replication.events.RefReplicationDoneEvent;
import com.googlesource.gerrit.plugins.replication.events.ReplicationScheduledEvent;
+import com.googlesource.gerrit.plugins.replication.events.dispatcher.EventDispatcher;
+import com.googlesource.gerrit.plugins.replication.events.dispatcher.ForwardingEventDispatcher;
+import com.googlesource.gerrit.plugins.replication.events.dispatcher.NoopEventDispatcher;
import org.apache.http.impl.client.CloseableHttpClient;
import org.eclipse.jgit.transport.SshSessionFactory;
@@ -55,6 +58,7 @@
@Override
protected void configure() {
install(configModule);
+ bindEventDispatcher();
bind(ObservableQueue.class).to(ReplicationQueue.class);
bind(LifecycleListener.class)
.annotatedWith(UniqueAnnotations.create())
@@ -111,4 +115,14 @@
install(new FactoryModuleBuilder().build(Destination.Factory.class));
install(new FactoryModuleBuilder().build(ProjectDeletionState.Factory.class));
}
+
+ private void bindEventDispatcher() {
+ boolean emitEvents =
+ configModule.getReplicationConfig().getBoolean("replication", "emitEvents", true);
+ if (emitEvents) {
+ bind(EventDispatcher.class).to(ForwardingEventDispatcher.class).in(Scopes.SINGLETON);
+ } else {
+ bind(EventDispatcher.class).to(NoopEventDispatcher.class).in(Scopes.SINGLETON);
+ }
+ }
}
diff --git a/src/main/java/com/googlesource/gerrit/plugins/replication/ReplicationQueue.java b/src/main/java/com/googlesource/gerrit/plugins/replication/ReplicationQueue.java
index 7b6079a..2652cb9 100644
--- a/src/main/java/com/googlesource/gerrit/plugins/replication/ReplicationQueue.java
+++ b/src/main/java/com/googlesource/gerrit/plugins/replication/ReplicationQueue.java
@@ -25,8 +25,6 @@
import com.google.gerrit.extensions.events.HeadUpdatedListener;
import com.google.gerrit.extensions.events.LifecycleListener;
import com.google.gerrit.extensions.events.ProjectDeletedListener;
-import com.google.gerrit.extensions.registration.DynamicItem;
-import com.google.gerrit.server.events.EventDispatcher;
import com.google.gerrit.server.extensions.events.GitReferenceUpdated;
import com.google.gerrit.server.git.WorkQueue;
import com.google.gerrit.util.logging.NamedFluentLogger;
@@ -37,6 +35,7 @@
import com.googlesource.gerrit.plugins.replication.api.ReplicationConfig;
import com.googlesource.gerrit.plugins.replication.api.ReplicationConfig.FilterType;
import com.googlesource.gerrit.plugins.replication.events.ProjectDeletionState;
+import com.googlesource.gerrit.plugins.replication.events.dispatcher.EventDispatcher;
import java.net.URISyntaxException;
import java.util.Collection;
import java.util.HashSet;
@@ -65,7 +64,7 @@
private final ReplicationConfig replConfig;
private final WorkQueue workQueue;
- private final DynamicItem<EventDispatcher> dispatcher;
+ private final EventDispatcher dispatcher;
private final Provider<ReplicationDestinations> destinations; // For Guice circular dependency
private final ReplicationTasksStorage replicationTasksStorage;
private final ProjectDeletionState.Factory projectDeletionStateFactory;
@@ -84,7 +83,7 @@
ReplicationConfig rc,
WorkQueue wq,
Provider<ReplicationDestinations> rd,
- DynamicItem<EventDispatcher> dis,
+ EventDispatcher dis,
ReplicationStateListeners sl,
ReplicationTasksStorage rts,
ProjectDeletionState.Factory pd) {
@@ -156,7 +155,7 @@
}
private void fire(String projectName, Set<UpdatedRef> updatedRefs) {
- ReplicationState state = new ReplicationState(new GitUpdateProcessing(dispatcher.get()));
+ ReplicationState state = new ReplicationState(new GitUpdateProcessing(dispatcher));
fire(Project.nameKey(projectName), null, updatedRefs, state, false);
state.markAllPushTasksScheduled();
}
@@ -199,7 +198,7 @@
}
private void fireFromStorage(URIish uri, Project.NameKey project, ImmutableSet<String> refNames) {
- ReplicationState state = new ReplicationState(new GitUpdateProcessing(dispatcher.get()));
+ ReplicationState state = new ReplicationState(new GitUpdateProcessing(dispatcher));
for (Destination dest : destinations.get().getDestinations(uri, project, refNames)) {
dest.scheduleFromStorage(project, refNames, uri, state);
}
@@ -220,7 +219,7 @@
boolean now) {
boolean withoutState = state == null;
if (withoutState) {
- state = new ReplicationState(new GitUpdateProcessing(dispatcher.get()));
+ state = new ReplicationState(new GitUpdateProcessing(dispatcher));
}
Set<String> refNamesToPush = new HashSet<>();
for (String refName : refNames) {
@@ -231,7 +230,7 @@
}
}
if (!refNamesToPush.isEmpty()) {
- for (URIish uri : cfg.getURIs(project, urlMatch)) {
+ for (URIish uri : cfg.getDistributedUris(project, urlMatch)) {
replicationTasksStorage.create(
ReplicateRefUpdate.create(
project.get(), refNamesToPush, uri, cfg.getRemoteConfigName()));
diff --git a/src/main/java/com/googlesource/gerrit/plugins/replication/StartCommand.java b/src/main/java/com/googlesource/gerrit/plugins/replication/StartCommand.java
index 57718bd..2c084b1 100644
--- a/src/main/java/com/googlesource/gerrit/plugins/replication/StartCommand.java
+++ b/src/main/java/com/googlesource/gerrit/plugins/replication/StartCommand.java
@@ -34,7 +34,7 @@
@Option(name = "--all", usage = "push all known projects")
private boolean all;
- @Option(name = "--url", metaVar = "PATTERN", usage = "pattern to match URL on")
+ @Option(name = "--url", metaVar = "SUBSTRING", usage = "substring URL must match (or * to match everything)")
private String urlMatch;
private final Set<String> remotesToConsider = new HashSet<>();
diff --git a/src/main/java/com/googlesource/gerrit/plugins/replication/UrlDistributionStrategy.java b/src/main/java/com/googlesource/gerrit/plugins/replication/UrlDistributionStrategy.java
new file mode 100644
index 0000000..46183eb
--- /dev/null
+++ b/src/main/java/com/googlesource/gerrit/plugins/replication/UrlDistributionStrategy.java
@@ -0,0 +1,84 @@
+// 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 java.util.Arrays;
+import java.util.List;
+import java.util.concurrent.atomic.AtomicInteger;
+import org.eclipse.jgit.transport.URIish;
+
+/**
+ * URL distribution strategy used when a remote has multiple configured URLs.
+ *
+ * <p>Each enum constant acts as a factory: call {@link #newInstance()} to obtain a stateful
+ * executor. Callers (e.g. {@link Destination}) hold the {@link Instance}, while the enum constant
+ * itself remains stateless and safe to use in equality checks.
+ *
+ * <p>Configured via {@code remote.NAME.urlDistribution} in {@code replication.config}.
+ */
+public enum UrlDistributionStrategy {
+ /** Push to all configured URLs. */
+ ALL("all") {
+ @Override
+ public Instance newInstance() {
+ return candidates -> candidates;
+ }
+ },
+
+ /**
+ * 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.
+ */
+ ROUND_ROBIN("roundRobin") {
+ @Override
+ public Instance newInstance() {
+ final AtomicInteger index = new AtomicInteger();
+ return candidates -> {
+ if (candidates.isEmpty()) {
+ return List.of();
+ }
+ return List.of(candidates.get(Math.floorMod(index.getAndIncrement(), candidates.size())));
+ };
+ }
+ };
+
+ public final String configKey;
+
+ UrlDistributionStrategy(String key) {
+ configKey = key;
+ }
+
+ /** Creates a new stateful executor for this distribution strategy. */
+ public abstract Instance newInstance();
+
+ /**
+ * Returns the distribution strategy for the given config value, or {@link #ALL} if the value is
+ * unrecognized or absent.
+ */
+ public static UrlDistributionStrategy fromConfig(String value) {
+ return Arrays.stream(values())
+ .filter(candidate -> candidate.configKey.equals(value))
+ .findFirst()
+ .orElse(ALL);
+ }
+
+ /** A stateful executor for a {@link UrlDistributionStrategy} strategy. */
+ @FunctionalInterface
+ public interface Instance {
+ List<URIish> select(List<URIish> candidates);
+ }
+}
diff --git a/src/main/java/com/googlesource/gerrit/plugins/replication/events/ProjectDeletionState.java b/src/main/java/com/googlesource/gerrit/plugins/replication/events/ProjectDeletionState.java
index e091915..4d6be7e 100644
--- a/src/main/java/com/googlesource/gerrit/plugins/replication/events/ProjectDeletionState.java
+++ b/src/main/java/com/googlesource/gerrit/plugins/replication/events/ProjectDeletionState.java
@@ -20,11 +20,10 @@
import static com.googlesource.gerrit.plugins.replication.events.ProjectDeletionState.ProjectDeletionStatus.TO_PROCESS;
import com.google.gerrit.entities.Project;
-import com.google.gerrit.extensions.registration.DynamicItem;
-import com.google.gerrit.server.events.EventDispatcher;
import com.google.gerrit.server.events.ProjectEvent;
import com.google.inject.Inject;
import com.google.inject.assistedinject.Assisted;
+import com.googlesource.gerrit.plugins.replication.events.dispatcher.EventDispatcher;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import org.eclipse.jgit.transport.URIish;
@@ -34,14 +33,13 @@
ProjectDeletionState create(Project.NameKey project);
}
- private final DynamicItem<EventDispatcher> eventDispatcher;
+ private final EventDispatcher eventDispatcher;
private final Project.NameKey project;
private final ConcurrentMap<URIish, ProjectDeletionStatus> statusByURI =
new ConcurrentHashMap<>();
@Inject
- public ProjectDeletionState(
- DynamicItem<EventDispatcher> eventDispatcher, @Assisted Project.NameKey project) {
+ public ProjectDeletionState(EventDispatcher eventDispatcher, @Assisted Project.NameKey project) {
this.eventDispatcher = eventDispatcher;
this.project = project;
}
@@ -70,7 +68,7 @@
private void setStatusAndBroadcastEvent(
URIish uri, ProjectDeletionStatus status, ProjectEvent event) {
statusByURI.put(uri, status);
- eventDispatcher.get().postEvent(project, event);
+ eventDispatcher.postEvent(project, event);
}
public void notifyIfDeletionDoneOnAllNodes() {
@@ -80,9 +78,7 @@
.noneMatch(s -> s.equals(TO_PROCESS) || s.equals(SCHEDULED))) {
statusByURI.clear();
- eventDispatcher
- .get()
- .postEvent(project, new ProjectDeletionReplicationDoneEvent(project.get()));
+ eventDispatcher.postEvent(project, new ProjectDeletionReplicationDoneEvent(project.get()));
}
}
}
diff --git a/src/main/java/com/googlesource/gerrit/plugins/replication/events/dispatcher/EventDispatcher.java b/src/main/java/com/googlesource/gerrit/plugins/replication/events/dispatcher/EventDispatcher.java
new file mode 100644
index 0000000..03dd588
--- /dev/null
+++ b/src/main/java/com/googlesource/gerrit/plugins/replication/events/dispatcher/EventDispatcher.java
@@ -0,0 +1,34 @@
+// 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.events.dispatcher;
+
+import com.google.gerrit.entities.BranchNameKey;
+import com.google.gerrit.entities.Project;
+import com.google.gerrit.server.events.Event;
+import com.google.gerrit.server.events.ProjectEvent;
+import com.google.gerrit.server.events.RefEvent;
+import com.google.gerrit.server.permissions.PermissionBackendException;
+
+/**
+ * Plugin-local indirection over {@link com.google.gerrit.server.events.EventDispatcher}, so that
+ * event emission can be turned off entirely via the {@code replication.emitEvents} config option.
+ */
+public interface EventDispatcher {
+ void postEvent(BranchNameKey branchName, RefEvent event) throws PermissionBackendException;
+
+ void postEvent(Project.NameKey projectName, ProjectEvent event);
+
+ void postEvent(Event event) throws PermissionBackendException;
+}
diff --git a/src/main/java/com/googlesource/gerrit/plugins/replication/events/dispatcher/ForwardingEventDispatcher.java b/src/main/java/com/googlesource/gerrit/plugins/replication/events/dispatcher/ForwardingEventDispatcher.java
new file mode 100644
index 0000000..ce19814
--- /dev/null
+++ b/src/main/java/com/googlesource/gerrit/plugins/replication/events/dispatcher/ForwardingEventDispatcher.java
@@ -0,0 +1,50 @@
+// 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.events.dispatcher;
+
+import com.google.gerrit.entities.BranchNameKey;
+import com.google.gerrit.entities.Project;
+import com.google.gerrit.extensions.registration.DynamicItem;
+import com.google.gerrit.server.events.Event;
+import com.google.gerrit.server.events.ProjectEvent;
+import com.google.gerrit.server.events.RefEvent;
+import com.google.gerrit.server.permissions.PermissionBackendException;
+import com.google.inject.Inject;
+
+public class ForwardingEventDispatcher implements EventDispatcher {
+ protected final DynamicItem<com.google.gerrit.server.events.EventDispatcher> delegate;
+
+ @Inject
+ public ForwardingEventDispatcher(
+ DynamicItem<com.google.gerrit.server.events.EventDispatcher> delegate) {
+ this.delegate = delegate;
+ }
+
+ @Override
+ public void postEvent(BranchNameKey branchName, RefEvent event)
+ throws PermissionBackendException {
+ delegate.get().postEvent(branchName, event);
+ }
+
+ @Override
+ public void postEvent(Project.NameKey projectName, ProjectEvent event) {
+ delegate.get().postEvent(projectName, event);
+ }
+
+ @Override
+ public void postEvent(Event event) throws PermissionBackendException {
+ delegate.get().postEvent(event);
+ }
+}
diff --git a/src/main/java/com/googlesource/gerrit/plugins/replication/events/dispatcher/NoopEventDispatcher.java b/src/main/java/com/googlesource/gerrit/plugins/replication/events/dispatcher/NoopEventDispatcher.java
new file mode 100644
index 0000000..c75b28c
--- /dev/null
+++ b/src/main/java/com/googlesource/gerrit/plugins/replication/events/dispatcher/NoopEventDispatcher.java
@@ -0,0 +1,32 @@
+// 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.events.dispatcher;
+
+import com.google.gerrit.entities.BranchNameKey;
+import com.google.gerrit.entities.Project;
+import com.google.gerrit.server.events.Event;
+import com.google.gerrit.server.events.ProjectEvent;
+import com.google.gerrit.server.events.RefEvent;
+
+public class NoopEventDispatcher implements EventDispatcher {
+ @Override
+ public void postEvent(BranchNameKey branchName, RefEvent event) {}
+
+ @Override
+ public void postEvent(Project.NameKey projectName, ProjectEvent event) {}
+
+ @Override
+ public void postEvent(Event event) {}
+}
diff --git a/src/main/resources/Documentation/config.md b/src/main/resources/Documentation/config.md
index 821f56a..cef009c 100644
--- a/src/main/resources/Documentation/config.md
+++ b/src/main/resources/Documentation/config.md
@@ -261,6 +261,19 @@
When not set, defaults to the plugin's data directory.
+replication.emitEvents
+: Whether to emit replication events to Gerrit's event dispatcher.
+
+ When disabled, all replication events are dropped before reaching the
+ event bus, which also prevents any downstream listeners (stream-events,
+ plugins, etc.) from receiving them. The replication work itself is not
+ affected.
+
+ The value is read at plugin load time; toggling it requires a plugin
+ reload.
+
+ Default: true
+
replication.rsyncPath
: Path to the `rsync` binary on the host running Gerrit, used by the
`@PLUGIN@ repair --copy-packs` command when transferring pack files
@@ -673,6 +686,24 @@
Do not exclude any refs pattern by default.
+remote.NAME.urlDistributionStrategy
+: URL distribution strategy to use when a remote has multiple configured URLs.
+ Applies to both push URLs (`url`) and admin URLs (`adminUrl`).
+
+ `all` (default)
+ : Push to all configured URLs simultaneously. This is the standard
+ replication behaviour.
+
+ `roundRobin`
+ : Push to one URL at a time, rotating through the list on each push event.
+ Particularly useful in high-availability setups where multiple replica
+ 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.
+
+ Defaults to `all`.
+
Directory `replication`
--------------------
The optional directory `$site_path/etc/replication` contains Git-style
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 646f915..c7b4174 100644
--- a/src/test/java/com/googlesource/gerrit/plugins/replication/DestinationConfigurationTest.java
+++ b/src/test/java/com/googlesource/gerrit/plugins/replication/DestinationConfigurationTest.java
@@ -98,4 +98,20 @@
// then
assertThat(actual).isEqualTo(globalPushBatchSize);
}
+
+ @Test
+ public void shouldDefaultUrlDistributionToAll() {
+ assertThat(objectUnderTest.getUrlDistributionStrategy()).isEqualTo(UrlDistributionStrategy.ALL);
+ }
+
+ @Test
+ public void shouldSetUrlDistributionToRoundRobinWhenConfigured() {
+ // given
+ when(cfgMock.getString("remote", REMOTE, "urlDistributionStrategy")).thenReturn("roundRobin");
+ objectUnderTest = new DestinationConfiguration(remoteConfigMock, cfgMock);
+
+ // when / then
+ assertThat(objectUnderTest.getUrlDistributionStrategy())
+ .isEqualTo(UrlDistributionStrategy.ROUND_ROBIN);
+ }
}
diff --git a/src/test/java/com/googlesource/gerrit/plugins/replication/GitUpdateProcessingTest.java b/src/test/java/com/googlesource/gerrit/plugins/replication/GitUpdateProcessingTest.java
index 43d97c1..af2f142 100644
--- a/src/test/java/com/googlesource/gerrit/plugins/replication/GitUpdateProcessingTest.java
+++ b/src/test/java/com/googlesource/gerrit/plugins/replication/GitUpdateProcessingTest.java
@@ -19,12 +19,12 @@
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
-import com.google.gerrit.server.events.EventDispatcher;
import com.google.gerrit.server.permissions.PermissionBackendException;
import com.googlesource.gerrit.plugins.replication.PushResultProcessing.GitUpdateProcessing;
import com.googlesource.gerrit.plugins.replication.ReplicationState.RefPushResult;
import com.googlesource.gerrit.plugins.replication.events.RefReplicatedEvent;
import com.googlesource.gerrit.plugins.replication.events.RefReplicationDoneEvent;
+import com.googlesource.gerrit.plugins.replication.events.dispatcher.EventDispatcher;
import java.net.URISyntaxException;
import org.eclipse.jgit.transport.RemoteRefUpdate;
import org.eclipse.jgit.transport.URIish;
diff --git a/src/test/java/com/googlesource/gerrit/plugins/replication/ReplicationDaemon.java b/src/test/java/com/googlesource/gerrit/plugins/replication/ReplicationDaemon.java
index 11b7718..436fb1c 100644
--- a/src/test/java/com/googlesource/gerrit/plugins/replication/ReplicationDaemon.java
+++ b/src/test/java/com/googlesource/gerrit/plugins/replication/ReplicationDaemon.java
@@ -308,4 +308,10 @@
config.setBoolean("remote", remoteName, "replicateProjectDeletions", replicateProjectDeletion);
config.save();
}
+
+ protected void setUrlDistribution(String remoteName, UrlDistributionStrategy distribution)
+ throws IOException {
+ config.setString("remote", remoteName, "urlDistributionStrategy", distribution.configKey);
+ config.save();
+ }
}
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 229808c..d436f42 100644
--- a/src/test/java/com/googlesource/gerrit/plugins/replication/ReplicationIT.java
+++ b/src/test/java/com/googlesource/gerrit/plugins/replication/ReplicationIT.java
@@ -504,6 +504,102 @@
}
@Test
+ public void shouldReplicateToAllUrlsByDefault() 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);
+ 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 on both the refs
+ try (Repository r1 = repoManager.openRepository(replica1Project);
+ Repository r2 = repoManager.openRepository(replica2Project)) {
+ waitUntil(() -> checkedGetRef(r1, newRef) != null && checkedGetRef(r2, newRef) != null);
+ }
+ }
+
+ @Test
+ public void shouldReplicateToOnlyOneUrlWhenRoundRobinEnabled() 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.ROUND_ROBIN);
+ 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 shouldRotateAcrossUrlsOnConsecutivePushesWhenRoundRobinEnabled() 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.ROUND_ROBIN);
+ reloadConfig();
+
+ String branch1 = "refs/heads/branch1";
+ String branch2 = "refs/heads/branch2";
+ createNewBranchWithoutPush("refs/heads/master", branch1);
+
+ ReplicationQueue queue = plugin.getSysInjector().getInstance(ReplicationQueue.class);
+
+ // First sync - goes to replica1 (index 0)
+ queue.scheduleFullSync(project, null, new ReplicationState(NO_OP), true);
+
+ try (Repository r1 = repoManager.openRepository(replica1Project)) {
+ waitUntil(() -> checkedGetRef(r1, branch1) != null);
+ }
+
+ // replica2 should not have received branch1 yet
+ try (Repository r2 = repoManager.openRepository(replica2Project)) {
+ assertThat(checkedGetRef(r2, branch1)).isNull();
+ }
+
+ // Second sync - goes to replica2 (index 1), includes branch1 and branch2
+ createNewBranchWithoutPush("refs/heads/master", branch2);
+ queue.scheduleFullSync(project, null, new ReplicationState(NO_OP), true);
+
+ try (Repository r2 = repoManager.openRepository(replica2Project)) {
+ waitUntil(() -> checkedGetRef(r2, branch1) != null && checkedGetRef(r2, branch2) != null);
+ }
+
+ // replica1 should not have received branch2 (it was only in the second sync)
+ try (Repository r1 = repoManager.openRepository(replica1Project)) {
+ assertThat(checkedGetRef(r1, branch2)).isNull();
+ }
+ }
+
+ @Test
public void shouldReplicateToMatchingRemote() throws Exception {
Project.NameKey targetProject = createTestProject(project + "replica");