Add config option to disable replication event emission The plugin posts a stream event at every point of the replication lifecycle (ref-replicated, replication-scheduled, replication-done, project-deletion-*). When no listener is configured these events are pure overhead. The cost is especially visible in multi-primary deployments running an events-sharing plugin that fans every replication event out to all the other primaries, even though nothing on either end consumes it. Introduce 'replication.emitEvents' (default true, so existing behavior is preserved) to suppress emission at the source. Change-Id: I32972d371ec83856e6783164bad9a229d1c8f934
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 54410f2..9620609 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,7 +156,7 @@ 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; @@ -187,7 +187,7 @@ GroupBackend groupBackend, ReplicationStateListeners stateLog, GroupIncludeCache groupIncludeCache, - DynamicItem<EventDispatcher> eventDispatcher, + EventDispatcher eventDispatcher, Provider<ReplicationTasksStorage> rts, CredentialsFactory credentialsFactory, @Assisted DestinationConfiguration cfg) { @@ -939,7 +939,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"); } @@ -952,7 +952,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/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/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 5d7fd19..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) {
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 fd3318a..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
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;