Merge branch 'stable-3.14' * stable-3.14: Retry failovers by applying url distribution strategy Add projectSharded URL selection per remote Change-Id: Id760dc515503ea27f6720ddf9edec00b0e159e19
diff --git a/src/main/java/com/googlesource/gerrit/plugins/replication/AutoReloadConfigDecorator.java b/src/main/java/com/googlesource/gerrit/plugins/replication/AutoReloadConfigDecorator.java index 0c77b16..e202729 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/replication/AutoReloadConfigDecorator.java +++ b/src/main/java/com/googlesource/gerrit/plugins/replication/AutoReloadConfigDecorator.java
@@ -23,6 +23,7 @@ import com.google.inject.Singleton; import com.googlesource.gerrit.plugins.replication.api.ReplicationConfig; import java.nio.file.Path; +import java.time.Duration; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; @@ -130,6 +131,21 @@ } @Override + public Duration getAutoRepairInterval() { + return currentConfig.getAutoRepairInterval(); + } + + @Override + public int getAutoRepairMaxAttempts() { + return currentConfig.getAutoRepairMaxAttempts(); + } + + @Override + public int getAutoRepairConcurrencyLimit() { + return currentConfig.getAutoRepairConcurrencyLimit(); + } + + @Override public Config getConfig() { return currentConfig.getConfig(); }
diff --git a/src/main/java/com/googlesource/gerrit/plugins/replication/AutoRepairHandler.java b/src/main/java/com/googlesource/gerrit/plugins/replication/AutoRepairHandler.java new file mode 100644 index 0000000..fa97cf2 --- /dev/null +++ b/src/main/java/com/googlesource/gerrit/plugins/replication/AutoRepairHandler.java
@@ -0,0 +1,126 @@ +// 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.googlesource.gerrit.plugins.replication.ReplicationQueue.repLog; + +import com.google.gerrit.entities.Project; +import com.google.gerrit.extensions.annotations.PluginName; +import com.google.gerrit.server.git.WorkQueue; +import com.google.inject.Inject; +import com.google.inject.Singleton; +import com.googlesource.gerrit.plugins.replication.ProjectRepairer.Action; +import com.googlesource.gerrit.plugins.replication.api.ReplicationConfig; +import com.googlesource.gerrit.plugins.replication.events.dispatcher.EventDispatcher; +import java.io.ByteArrayOutputStream; +import java.io.InterruptedIOException; +import java.nio.charset.StandardCharsets; +import java.util.Collections; +import java.util.List; +import java.util.Set; +import java.util.concurrent.Future; +import java.util.concurrent.ScheduledExecutorService; +import org.eclipse.jgit.transport.URIish; + +@Singleton +public class AutoRepairHandler { + static final String MISSING_NECESSARY_OBJECTS = "missing necessary objects"; + + private final AutoRepairTracker tracker; + private final ProjectRepairer projectRepairer; + private final ReplicationStarter replicationStarter; + private final EventDispatcher eventDispatcher; + private final ScheduledExecutorService executor; + + @Inject + AutoRepairHandler( + AutoRepairTracker tracker, + ProjectRepairer projectRepairer, + ReplicationStarter replicationStarter, + EventDispatcher eventDispatcher, + WorkQueue workQueue, + ReplicationConfig replicationConfig, + @PluginName String pluginName) { + this.tracker = tracker; + this.projectRepairer = projectRepairer; + this.replicationStarter = replicationStarter; + this.eventDispatcher = eventDispatcher; + this.executor = + workQueue.createQueue( + replicationConfig.getAutoRepairConcurrencyLimit(), pluginName + "_auto-repair"); + } + + public static boolean isMissingNecessaryObjectsError(String message) { + return message != null && message.contains(MISSING_NECESSARY_OBJECTS); + } + + public void handle( + Project.NameKey project, + URIish uri, + String remoteName, + UrlDistributionStrategy urlDistributionStrategy) { + if (!tracker.tryBeginRepair(project, uri, remoteName, urlDistributionStrategy)) { + return; + } + repLog.atInfo().log("Scheduling auto-repair for project %s to %s", project.get(), uri); + @SuppressWarnings("unused") + Future<?> possiblyIgnoredError = executor.submit(new AutoRepairTask(project, uri)); + } + + private class AutoRepairTask implements Runnable { + private final Project.NameKey project; + private final URIish uri; + + AutoRepairTask(Project.NameKey project, URIish uri) { + this.project = project; + this.uri = uri; + } + + @Override + public void run() { + ByteArrayOutputStream buf = new ByteArrayOutputStream(); + boolean isRepaired; + try { + isRepaired = projectRepairer.repair(project, uri, buf, Action.all()); + } catch (InterruptedIOException e) { + repLog.atWarning().withCause(e).log( + "Auto-repair interrupted for project %s to %s", project.get(), uri); + return; + } + (isRepaired ? repLog.atInfo() : repLog.atWarning()) + .log( + "Auto-repair %s for project %s to %s:%s", + isRepaired ? "succeeded" : "failed", + project.get(), + uri, + buf.toString(StandardCharsets.UTF_8)); + if (isRepaired) { + replicationStarter.start( + uri.toString(), + Set.of(), + new ReplicationFilter(List.of(project.get()), Collections.emptyList()), + /* now= */ true, + /* wait= */ false, + new PushResultProcessing.GitUpdateProcessing(eventDispatcher)); + repLog.atInfo().log("Scheduled full replication of %s to %s", project.get(), uri); + } + } + + @Override + public String toString() { + return "auto-repair " + project.get() + " to " + uri; + } + } +}
diff --git a/src/main/java/com/googlesource/gerrit/plugins/replication/AutoRepairTracker.java b/src/main/java/com/googlesource/gerrit/plugins/replication/AutoRepairTracker.java new file mode 100644 index 0000000..53d9c9e --- /dev/null +++ b/src/main/java/com/googlesource/gerrit/plugins/replication/AutoRepairTracker.java
@@ -0,0 +1,118 @@ +// 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.googlesource.gerrit.plugins.replication.ReplicationQueue.repLog; + +import com.google.gerrit.entities.Project; +import com.google.inject.Inject; +import com.google.inject.Singleton; +import com.googlesource.gerrit.plugins.replication.api.ReplicationConfig; +import java.time.Duration; +import java.time.Instant; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentMap; +import org.eclipse.jgit.transport.URIish; + +/** + * Tracks per-project, per-destination auto-repair attempts and enforces interval and attempt + * limits. + * + * <p>State is kept in memory and is reset whenever the plugin is reloaded or Gerrit restarts. + */ +@Singleton +public class AutoRepairTracker { + private final ReplicationConfig replicationConfig; + private final ConcurrentMap<DestinationKey, State> statesByDestination = + new ConcurrentHashMap<>(); + + @Inject + AutoRepairTracker(ReplicationConfig replicationConfig) { + this.replicationConfig = replicationConfig; + } + + public boolean isEnabled() { + return replicationConfig.getAutoRepairMaxAttempts() > 0; + } + + /** + * Returns whether a new auto-repair attempt is allowed for the project on the given destination. + * + * <p>If allowed, records the attempt before returning {@code true}. + */ + public boolean tryBeginRepair( + Project.NameKey project, + URIish uri, + String remoteName, + UrlDistributionStrategy urlDistributionStrategy) { + if (!isEnabled()) { + return false; + } + + if (!ProjectRepairer.canCopy(uri)) { + repLog.atWarning().log( + "Skipping auto-repair for %s to %s: only plain SSH destinations are supported", + project.get(), uri); + return false; + } + + DestinationKey key = DestinationKey.create(project, uri, remoteName, urlDistributionStrategy); + return statesByDestination.computeIfAbsent(key, k -> new State()).isRepairAllowed(project, uri); + } + + private record DestinationKey(Project.NameKey project, String destination) { + static DestinationKey create( + Project.NameKey project, + URIish uri, + String remoteName, + UrlDistributionStrategy urlDistributionStrategy) { + return new DestinationKey( + project, + urlDistributionStrategy == UrlDistributionStrategy.ROUND_ROBIN + ? remoteName + : uri.toString()); + } + } + + private class State { + private int repairAttemptCount; + private Instant lastAttempt = Instant.EPOCH; + + synchronized boolean isRepairAllowed(Project.NameKey project, URIish uri) { + int maxAttempts = replicationConfig.getAutoRepairMaxAttempts(); + if (repairAttemptCount >= maxAttempts) { + repLog.atWarning().log( + "Skipping auto-repair for %s to %s: reached max attempts (%d)", + project.get(), uri, maxAttempts); + return false; + } + Duration interval = replicationConfig.getAutoRepairInterval(); + Instant now = Instant.now(); + if (!interval.isZero()) { + Instant nextAllowed = lastAttempt.plus(interval); + if (now.isBefore(nextAllowed)) { + repLog.atInfo().log( + "Skipping auto-repair for %s to %s: interval not elapsed." + + " Next attempt allowed at %s", + project.get(), uri, nextAllowed); + return false; + } + } + lastAttempt = now; + repairAttemptCount++; + return true; + } + } +}
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 e05cc2b..f17b50b 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; @@ -85,7 +85,6 @@ import java.util.concurrent.locks.ReadWriteLock; import java.util.function.Function; import java.util.function.Supplier; -import java.util.regex.Pattern; import java.util.stream.Collectors; import org.eclipse.jgit.lib.Constants; import org.eclipse.jgit.lib.Ref; @@ -157,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; @@ -188,7 +187,7 @@ GroupBackend groupBackend, ReplicationStateListeners stateLog, GroupIncludeCache groupIncludeCache, - DynamicItem<EventDispatcher> eventDispatcher, + EventDispatcher eventDispatcher, Provider<ReplicationTasksStorage> rts, CredentialsFactory credentialsFactory, @Assisted DestinationConfiguration cfg) { @@ -287,9 +286,21 @@ return queue; } + /** + * Whether this destination performs pushes. + * + * <p>When {@code remote.NAME.threads} is set to 0, no pushes will be done, but the tasks are + * still persisted to the {@link ReplicationTasksStorage}. + */ + public boolean isPushEnabled() { + return config.getPoolThreads() > 0; + } + public void start(WorkQueue workQueue) { - String poolName = "ReplicateTo-" + config.getRemoteConfig().getName(); - pool = workQueue.createQueue(config.getPoolThreads(), poolName); + if (isPushEnabled()) { + String poolName = "ReplicateTo-" + config.getRemoteConfig().getName(); + pool = workQueue.createQueue(config.getPoolThreads(), poolName); + } } public int shutdown() { @@ -438,10 +449,6 @@ return false; } - void schedule(Project.NameKey project, Set<String> refs, URIish uri, ReplicationState state) { - schedule(project, refs, uri, state, false); - } - void scheduleFromStorage( Project.NameKey project, Set<String> refs, URIish uri, ReplicationState state) { schedule(project, refs, uri, state, false, true); @@ -459,6 +466,9 @@ ReplicationState state, boolean now, boolean fromStorage) { + if (!isPushEnabled()) { + return; + } ImmutableSet.Builder<String> toSchedule = ImmutableSet.builder(); for (String ref : refs) { if (!shouldReplicate(project, ref, state)) { @@ -548,17 +558,21 @@ } void scheduleDeleteProject(URIish uri, Project.NameKey project, ProjectDeletionState state) { - repLog.atFine().log("scheduling deletion of project %s at %s", project, uri); - @SuppressWarnings("unused") - ScheduledFuture<?> ignored = - pool.schedule(deleteProjectFactory.create(uri, project, state), 0, TimeUnit.SECONDS); - state.setScheduled(uri); + if (isPushEnabled()) { + repLog.atFine().log("scheduling deletion of project %s at %s", project, uri); + @SuppressWarnings("unused") + ScheduledFuture<?> ignored = + pool.schedule(deleteProjectFactory.create(uri, project, state), 0, TimeUnit.SECONDS); + state.setScheduled(uri); + } } void scheduleUpdateHead(URIish uri, Project.NameKey project, String newHead) { - @SuppressWarnings("unused") - ScheduledFuture<?> ignored = - pool.schedule(updateHeadFactory.create(uri, project, newHead), 0, TimeUnit.SECONDS); + if (isPushEnabled()) { + @SuppressWarnings("unused") + ScheduledFuture<?> ignored = + pool.schedule(updateHeadFactory.create(uri, project, newHead), 0, TimeUnit.SECONDS); + } } /** @@ -567,7 +581,7 @@ * <p>If the reason for rescheduling is to avoid a collision with an in-flight push to the same * URI, we don't mark the operation as "retrying," and we schedule using the replication delay, * rather than the retry delay. Otherwise, the operation is marked as "retrying" and scheduled to - * run following the minutes count determined by class attribute retryDelay. + * run following the retry delay (in seconds) determined by class attribute retryDelay. * * <p>In case the PushOp instance to be scheduled has same URI than one marked as "retrying," it * adds to the one pending the refs list of the parameter instance. @@ -583,6 +597,9 @@ * @param pushOp The PushOp instance to be scheduled. */ void reschedule(PushOne pushOp, RetryReason reason) { + if (!isPushEnabled()) { + return; + } RescheduleStatus status = new RescheduleStatus(); stateLock.withWriteLock( pushOp.getURI(), @@ -656,7 +673,7 @@ replicationTasksStorage.get().reset(pushOp); @SuppressWarnings("unused") ScheduledFuture<?> ignored2 = - pool.schedule(pushOp, config.getRetryDelay(), TimeUnit.MINUTES); + pool.schedule(pushOp, config.getRetryDelay(), TimeUnit.SECONDS); } } else { pushOp.canceledByReplication(); @@ -708,7 +725,7 @@ queue.pending.put(newUri, replacement); @SuppressWarnings("unused") ScheduledFuture<?> ignored = - pool.schedule(replacement, config.getRetryDelay(), TimeUnit.MINUTES); + pool.schedule(replacement, config.getRetryDelay(), TimeUnit.SECONDS); return true; }); repLog.atInfo().log( @@ -803,6 +820,10 @@ if (PushOne.ALL_REFS.equals(ref)) { return true; } + if (isRefExcluded(ref)) { + repLog.atFine().log("Skipping push of ref %s; it matches excludedRefsPattern", ref); + return false; + } for (RefSpec s : config.getRemoteConfig().getPushRefSpecs()) { if (s.matchSource(ref)) { return true; @@ -929,6 +950,10 @@ return config.getRemoteConfig().getName(); } + UrlDistributionStrategy getUrlDistributionStrategy() { + return config.getUrlDistributionStrategy(); + } + public int getMaxRetries() { return config.getMaxRetries(); } @@ -953,8 +978,8 @@ return config.replicateNoteDbMetaRefs(); } - ImmutableList<Pattern> excludedRefsPattern() { - return config.excludedRefsPattern(); + boolean isRefExcluded(String ref) { + return config.isRefExcluded(ref); } boolean storeRefLog() { @@ -991,7 +1016,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"); } @@ -1004,7 +1029,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 03ba914..02389d6 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/replication/DestinationConfiguration.java +++ b/src/main/java/com/googlesource/gerrit/plugins/replication/DestinationConfiguration.java
@@ -32,6 +32,7 @@ public class DestinationConfiguration implements RemoteConfiguration { static final int DEFAULT_REPLICATION_DELAY = 15; static final int DEFAULT_RESCHEDULE_DELAY = 3; + static final int DEFAULT_REPLICATION_RETRY_MINUTES = 1; static final int DEFAULT_DRAIN_QUEUE_ATTEMPTS = 0; private static final int DEFAULT_SLOW_LATENCY_THRESHOLD_SECS = 900; @@ -73,7 +74,7 @@ projects = ImmutableList.copyOf(cfg.getStringList("remote", name, "projects")); excludeProjects = ImmutableList.copyOf(cfg.getStringList("remote", name, "excludeProjects")); adminUrls = ImmutableList.copyOf(cfg.getStringList("remote", name, "adminUrl")); - retryDelay = Math.max(0, getInt(remoteConfig, cfg, "replicationretry", 1)); + retryDelay = getRetryDelaySeconds(remoteConfig, cfg); drainQueueAttempts = Math.max(0, getInt(remoteConfig, cfg, "drainQueueAttempts", DEFAULT_DRAIN_QUEUE_ATTEMPTS)); poolThreads = Math.max(0, getInt(remoteConfig, cfg, "threads", 1)); @@ -231,6 +232,33 @@ return cfg.getInt("remote", rc.getName(), name, defValue); } + /** + * Parses {@code remote.NAME.replicationRetry} into seconds. For backwards compatibility a value + * without a time-unit suffix keeps its historical meaning of minutes -- including a negative bare + * number, which the historical {@code cfg.getInt(...)} parsing also accepted and which is clamped + * to zero below -- so existing configurations are unchanged. A value with a time-unit suffix -- + * e.g. {@code 30 s}, {@code 90 s} or {@code 2 m} -- is honoured as written, letting an admin + * configure a sub-minute offline-retry backoff. The result is clamped to {@code [0, + * Integer.MAX_VALUE]} so it never overflows the {@code int} the scheduler expects. + */ + private static int getRetryDelaySeconds(RemoteConfig rc, Config cfg) { + String value = cfg.getString("remote", rc.getName(), "replicationRetry"); + long defaultSeconds = TimeUnit.MINUTES.toSeconds(DEFAULT_REPLICATION_RETRY_MINUTES); + long seconds; + if (value == null || value.trim().isEmpty()) { + seconds = defaultSeconds; + } else if (value.trim().matches("-?[0-9]+")) { + seconds = TimeUnit.MINUTES.toSeconds(Long.parseLong(value.trim())); + } else { + // Use the config overload so an invalid value is reported against + // remote.NAME.replicationRetry rather than just the raw string. + seconds = + ConfigUtil.getTimeUnit( + cfg, "remote", rc.getName(), "replicationRetry", defaultSeconds, TimeUnit.SECONDS); + } + return (int) Math.max(0, Math.min(seconds, Integer.MAX_VALUE)); + } + @Override public int getSlowLatencyThreshold() { return slowLatencyThreshold;
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..2bf2ae4 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,10 +50,10 @@ 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) + .create(null, PushOne.ALL_REFS, Set.of(), ReplicationFilter.all(), state, false) .schedule(30, TimeUnit.SECONDS)); } }
diff --git a/src/main/java/com/googlesource/gerrit/plugins/replication/ProjectRepairer.java b/src/main/java/com/googlesource/gerrit/plugins/replication/ProjectRepairer.java index 5866569..d2e0edf 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/replication/ProjectRepairer.java +++ b/src/main/java/com/googlesource/gerrit/plugins/replication/ProjectRepairer.java
@@ -17,24 +17,49 @@ import static com.googlesource.gerrit.plugins.replication.ReplicationQueue.repLog; import com.google.common.base.Strings; +import com.google.gerrit.common.Nullable; import com.google.gerrit.entities.Project; import com.google.gerrit.server.git.GitRepositoryManager; import com.google.inject.Inject; import com.google.inject.Singleton; import com.googlesource.gerrit.plugins.replication.api.ReplicationConfig; import java.io.IOException; +import java.io.InterruptedIOException; import java.io.OutputStream; +import java.nio.file.DirectoryIteratorException; +import java.nio.file.DirectoryStream; import java.nio.file.Files; +import java.nio.file.NoSuchFileException; import java.nio.file.Path; +import java.time.Duration; +import java.time.Instant; import java.util.ArrayList; +import java.util.Collection; import java.util.List; +import java.util.Optional; +import java.util.UUID; +import org.eclipse.jgit.internal.storage.file.PackFile; +import org.eclipse.jgit.internal.storage.pack.PackExt; import org.eclipse.jgit.lib.Repository; import org.eclipse.jgit.transport.URIish; +import org.eclipse.jgit.util.FileUtils; import org.eclipse.jgit.util.QuotedString; import org.eclipse.jgit.util.io.StreamCopyThread; @Singleton public class ProjectRepairer { + public enum Action { + COPY_LOOSE_OBJECTS, + COPY_PACKS; + + public static List<Action> all() { + return List.of(values()); + } + } + + private static final String OBJECTS_DIR = "objects/"; + private static final String PACK_DIR = OBJECTS_DIR + "pack/"; + private static final String PACK_GLOB = "pack-*." + PackExt.PACK.getExtension(); private final GitRepositoryManager gitManager; private final ReplicationConfig replicationConfig; @@ -44,67 +69,164 @@ this.replicationConfig = replicationConfig; } - public boolean repair(Project.NameKey project, URIish uri, OutputStream out, boolean copyPacks) { - if (copyPacks && !copyPackTo(project, uri, out)) { - repLog.atSevere().log("Repair failed for %s on %s", project.get(), uri); + public boolean repair( + Project.NameKey project, URIish uri, OutputStream out, Collection<Action> actions) + throws InterruptedIOException { + if (actions.isEmpty()) { + return true; + } + + Path objectsDir = objectsDir(project, uri); + if (objectsDir == null) { return false; } - return true; + Snapshot.sweepStale(objectsDir.getParent()); + + boolean isRepaired = true; + for (Action action : actions) { + isRepaired &= repair(project, uri, out, objectsDir, action); + } + return isRepaired; + } + + private boolean repair( + Project.NameKey project, URIish uri, OutputStream out, Path objectsDir, Action action) + throws InterruptedIOException { + boolean isRepaired = + switch (action) { + case COPY_LOOSE_OBJECTS -> copyLooseObjectsTo(objectsDir, uri, out); + case COPY_PACKS -> copyPacksTo(objectsDir, uri, out); + }; + if (!isRepaired) { + repLog.atSevere().log("Repair (%s) failed for %s on %s", action, project.get(), uri); + } + return isRepaired; } public static boolean canCopy(URIish uri) { return AdminApiFactory.isSSH(uri) && !AdminApiFactory.isGerrit(uri); } - private boolean copyPackTo(Project.NameKey project, URIish uri, OutputStream out) { + @Nullable + private Path objectsDir(Project.NameKey project, URIish uri) { if (Strings.isNullOrEmpty(uri.getHost())) { repLog.atSevere().log("Cannot repair %s: URI has no host: %s", project.get(), uri); - return false; + return null; } if (Strings.isNullOrEmpty(uri.getPath())) { repLog.atSevere().log("Cannot repair %s: URI has no path: %s", project.get(), uri); - return false; + return null; } - Path packDir; try (Repository repo = gitManager.openRepository(project)) { - packDir = repo.getDirectory().toPath().resolve("objects").resolve("pack"); + return repo.getDirectory().toPath().resolve("objects"); } catch (IOException e) { repLog.atSevere().withCause(e).log("Cannot open repository %s for repair", project.get()); + return null; + } + } + + private boolean copyLooseObjectsTo(Path objectsDir, URIish uri, OutputStream out) + throws InterruptedIOException { + if (!Files.isDirectory(objectsDir)) { + repLog.atSevere().log("No objects directory %s", objectsDir); return false; } + try (Snapshot snapshot = new Snapshot(objectsDir.getParent())) { + linkLooseObjects(objectsDir, snapshot); + return copy(snapshot.dir, uri, out, OBJECTS_DIR) == 0; + } catch (InterruptedIOException e) { + // An interrupt must abort the repair rather than report a snapshot failure + throw e; + } catch (IOException e) { + repLog.atSevere().withCause(e).log("Cannot snapshot %s", objectsDir); + return false; + } + } + + private static void linkLooseObjects(Path objectsDir, Snapshot snapshot) throws IOException { + try (DirectoryStream<Path> fanoutDirs = Files.newDirectoryStream(objectsDir)) { + for (Path fanoutDir : fanoutDirs) { + if (fanoutDir.getFileName().toString().length() == 2 && Files.isDirectory(fanoutDir)) { + linkFanoutDir(fanoutDir, snapshot); + } + } + } + } + + private static void linkFanoutDir(Path fanoutDir, Snapshot snapshot) throws IOException { + String name = fanoutDir.getFileName().toString(); + try (DirectoryStream<Path> objects = Files.newDirectoryStream(fanoutDir)) { + snapshot.createSubdir(name); + for (Path object : objects) { + if (Files.isRegularFile(object)) { + snapshot.linkIfExists(object, name); + } + } + } catch (NoSuchFileException e) { + // ignore + } + } + + private boolean copyPacksTo(Path objectsDir, URIish uri, OutputStream out) + throws InterruptedIOException { + Path packDir = objectsDir.resolve("pack"); if (!Files.isDirectory(packDir)) { - repLog.atSevere().log("No objects/pack directory for project %s", project.get()); + repLog.atSevere().log("No objects/pack directory %s", packDir); return false; } - return copyInOrder(packDir, uri, out); - } - - private boolean copyInOrder(Path packDir, URIish uri, OutputStream out) { - try { - return copy(packDir, uri, out, "*.pack") == 0 - && copy(packDir, uri, out, "*.idx", "*.bitmap", "*.rev") == 0; - } catch (InterruptedException e) { - repLog.atWarning().withCause(e).log("Interrupted during copy to %s", uri); + try (Snapshot snapshot = new Snapshot(objectsDir.getParent())) { + linkPacks(packDir, snapshot); + return copy(snapshot.dir, uri, out, PACK_DIR, PACK_GLOB) == 0 + && copy(snapshot.dir, uri, out, PACK_DIR) == 0; + } catch (InterruptedIOException e) { + // An interrupt must abort the repair rather than report a snapshot failure + throw e; + } catch (IOException e) { + repLog.atSevere().withCause(e).log("Cannot snapshot %s", packDir); return false; } } - private int copy(Path src, URIish uri, OutputStream out, String... includes) - throws InterruptedException { + private static void linkPacks(Path packDir, Snapshot snapshot) throws IOException { + try (DirectoryStream<Path> packs = Files.newDirectoryStream(packDir, PACK_GLOB)) { + for (Path pack : packs) { + linkPackSet(new PackFile(pack.toFile()), snapshot); + } + } + } + + private static void linkPackSet(PackFile pack, Snapshot snapshot) throws IOException { + Optional<Path> packLink = snapshot.linkIfExists(pack.toPath()); + if (packLink.isEmpty()) { + return; + } + if (snapshot.linkIfExists(pack.create(PackExt.INDEX).toPath()).isEmpty()) { + Files.delete(packLink.get()); + return; + } + snapshot.linkIfExists(pack.create(PackExt.BITMAP_INDEX).toPath()); + snapshot.linkIfExists(pack.create(PackExt.REVERSE_INDEX).toPath()); + } + + private int copy(Path src, URIish uri, OutputStream out, String destDir, String... includes) + throws InterruptedIOException { List<String> cmd = new ArrayList<>(); cmd.add(replicationConfig.getRsyncPath()); - cmd.add("-avP"); + cmd.add("-av"); + cmd.add("--progress"); cmd.add("-e"); cmd.add(buildSshTransport(uri)); - for (String inc : includes) { - cmd.add("--include=" + inc); + if (includes.length > 0) { + for (String inc : includes) { + cmd.add("--include=" + inc); + } + cmd.add("--exclude=*"); } - cmd.add("--exclude=*"); cmd.add(src.toAbsolutePath().normalize() + "/"); - cmd.add(buildCopyDestination(uri)); + cmd.add(buildCopyDestination(uri, destDir)); repLog.atInfo().log("Running repair cmd: %s", String.join(" ", cmd)); @@ -119,7 +241,7 @@ } StreamCopyThread outStream = new StreamCopyThread(p.getInputStream(), out); - outStream.setName("copy-packs-output"); + outStream.setName("repair-copy-output"); outStream.start(); try { int code = p.waitFor(); @@ -130,20 +252,25 @@ return code; } catch (InterruptedException e) { p.destroyForcibly(); - outStream.halt(); - return -1; + try { + outStream.halt(); + } catch (InterruptedException ignored) { + // ignore + } + throw (InterruptedIOException) + new InterruptedIOException("Interrupted during copy to " + uri).initCause(e); } } - private static String buildCopyDestination(URIish uri) { + private static String buildCopyDestination(URIish uri, String destDir) { String host = uri.getHost(); String path = uri.getPath(); - String remotePackPath = QuotedString.BOURNE.quote(path + "/objects/pack/"); + String remotePath = QuotedString.BOURNE.quote(path + "/" + destDir); String user = uri.getUser(); if (user != null && !user.isEmpty()) { - return user + "@" + host + ":" + remotePackPath; + return user + "@" + host + ":" + remotePath; } - return host + ":" + remotePackPath; + return host + ":" + remotePath; } private static String buildSshTransport(URIish uri) { @@ -154,4 +281,70 @@ } return sb.toString(); } + + private static final class Snapshot implements AutoCloseable { + private static final String SNAPSHOT_PREFIX = "replication-repair-snapshot-"; + private static final Duration MAX_AGE = Duration.ofDays(1); + + private final Path dir; + + private Snapshot(Path parentDir) throws IOException { + this.dir = Files.createDirectory(parentDir.resolve(SNAPSHOT_PREFIX + UUID.randomUUID())); + } + + private static void sweepStale(Path parentDir) { + Instant cutoff = Instant.now().minus(MAX_AGE); + try (DirectoryStream<Path> snapshotDirs = + Files.newDirectoryStream(parentDir, SNAPSHOT_PREFIX + "*")) { + for (Path snapshotDir : snapshotDirs) { + try { + if (Files.getLastModifiedTime(snapshotDir).toInstant().isBefore(cutoff)) { + repLog.atWarning().log("Deleting stale repair snapshot %s", snapshotDir); + delete(snapshotDir); + } + } catch (IOException e) { + repLog.atWarning().withCause(e).log("Cannot check repair snapshot %s", snapshotDir); + } + } + } catch (IOException | DirectoryIteratorException e) { + repLog.atWarning().withCause(e).log("Cannot sweep repair snapshots in %s", parentDir); + } + } + + private void createSubdir(String name) throws IOException { + Files.createDirectory(dir.resolve(name)); + } + + private Optional<Path> linkIfExists(Path src) throws IOException { + return linkInto(src, dir); + } + + private void linkIfExists(Path src, String subdir) throws IOException { + linkInto(src, dir.resolve(subdir)); + } + + private static Optional<Path> linkInto(Path src, Path destDir) throws IOException { + Path link = destDir.resolve(src.getFileName().toString()); + try { + Files.createLink(link, src); + return Optional.of(link); + } catch (NoSuchFileException e) { + return Optional.empty(); + } + } + + @Override + public void close() { + delete(dir); + } + + private static void delete(Path snapshotDir) { + try { + FileUtils.delete( + snapshotDir.toFile(), FileUtils.RECURSIVE | FileUtils.SKIP_MISSING | FileUtils.RETRY); + } catch (IOException e) { + repLog.atSevere().withCause(e).log("Cannot delete repair snapshot %s", snapshotDir); + } + } + } }
diff --git a/src/main/java/com/googlesource/gerrit/plugins/replication/PushAll.java b/src/main/java/com/googlesource/gerrit/plugins/replication/PushAll.java index 8958a2a..80fc210 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/replication/PushAll.java +++ b/src/main/java/com/googlesource/gerrit/plugins/replication/PushAll.java
@@ -29,7 +29,8 @@ public interface Factory { PushAll create( - String urlMatch, + @Assisted("urlMatch") String urlMatch, + @Assisted("refName") String refName, Set<String> remotesToConsider, ReplicationFilter filter, ReplicationState state, @@ -41,6 +42,7 @@ private final ReplicationQueue replication; private final String urlMatch; private final Set<String> remotesToConsider; + private final String refName; private final ReplicationFilter filter; private final ReplicationState state; private final boolean now; @@ -51,7 +53,8 @@ ProjectCache projectCache, ReplicationQueue rq, ReplicationStateListeners stateLog, - @Assisted @Nullable String urlMatch, + @Assisted("urlMatch") @Nullable String urlMatch, + @Assisted("refName") String refName, @Assisted Set<String> remotesToConsider, @Assisted ReplicationFilter filter, @Assisted ReplicationState state, @@ -61,6 +64,7 @@ this.replication = rq; this.stateLog = stateLog; this.urlMatch = urlMatch; + this.refName = refName; this.remotesToConsider = remotesToConsider; this.filter = filter; this.state = state; @@ -76,10 +80,10 @@ try { for (Project.NameKey nameKey : projectCache.all()) { if (filter.matches(nameKey)) { - replication.scheduleFullSync(nameKey, urlMatch, remotesToConsider, state, now); + replication.scheduleFullSync(nameKey, urlMatch, refName, remotesToConsider, state, now); } } - } catch (Exception e) { + } catch (RuntimeException e) { stateLog.error("Cannot enumerate known projects", e, state); } state.markAllPushTasksScheduled(); @@ -87,7 +91,9 @@ @Override public String toString() { - String s = "Replicate All Projects"; + String refs = PushOne.ALL_REFS.equals(refName) ? "All Refs" : "[" + refName + "]"; + String s = "Replicate " + refs + " for All Projects"; + if (urlMatch != null) { s = s + " to " + urlMatch; }
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 600ce8c..07cb830 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/replication/PushOne.java +++ b/src/main/java/com/googlesource/gerrit/plugins/replication/PushOne.java
@@ -142,7 +142,9 @@ private final CreateProjectTask.Factory createProjectFactory; private final AtomicBoolean canceledWhileRunning; private final TransportFactory transportFactory; + private final AutoRepairHandler autoRepairHandler; private DynamicItem<ReplicationPushFilter> replicationPushFilter; + private volatile String batchProgress = ""; @Inject PushOne( @@ -159,6 +161,7 @@ ProjectCache pc, CreateProjectTask.Factory cpf, TransportFactory tf, + AutoRepairHandler autoRepairHandler, @Assisted Project.NameKey d, @Assisted URIish u) { gitManager = grm; @@ -181,6 +184,7 @@ canceledWhileRunning = new AtomicBoolean(false); maxRetries = p.getMaxRetries(); transportFactory = tf; + this.autoRepairHandler = autoRepairHandler; } @Inject(optional = true) @@ -220,7 +224,8 @@ @Override public String toString() { - String print = "[" + HexFormat.fromInt(id) + "] push " + uri + " " + getLimitedRefs(); + String print = + "[" + HexFormat.fromInt(id) + "] push " + uri + " " + getLimitedRefs() + batchProgress; if (retryCount > 0) { print = "(retry " + retryCount + ") " + print; @@ -603,21 +608,26 @@ List<List<RemoteRefUpdate>> batches = Lists.partition(todo, batchSize); repLog.atInfo().log("Push to %s in %d batches", uri, batches.size()); AggregatedPushResult result = new AggregatedPushResult(); - int completedBatch = 1; - for (List<RemoteRefUpdate> batch : batches) { - repLog.atInfo().log( - "Pushing %d/%d batches for replication to %s", completedBatch, batches.size(), uri); - result.addResult(tn.push(NullProgressMonitor.INSTANCE, batch)); - - // check if push should be no longer continued - if (wasCanceled()) { + try { + int completedBatch = 1; + for (List<RemoteRefUpdate> batch : batches) { + batchProgress = " (batch " + completedBatch + "/" + batches.size() + ")"; repLog.atInfo().log( - "Push for replication to %s was canceled after %d completed batch and thus won't be" - + " rescheduled", - uri, completedBatch); - break; + "Pushing %d/%d batches for replication to %s", completedBatch, batches.size(), uri); + result.addResult(tn.push(NullProgressMonitor.INSTANCE, batch)); + + // check if push should be no longer continued + if (wasCanceled()) { + repLog.atInfo().log( + "Push for replication to %s was canceled after %d completed batch and thus won't be" + + " rescheduled", + uri, completedBatch); + break; + } + completedBatch++; } - completedBatch++; + } finally { + batchProgress = ""; } return result; } @@ -785,7 +795,7 @@ return !(noPerms && RefNames.REFS_CONFIG.equals(ref)) && !ref.startsWith(RefNames.REFS_CACHE_AUTOMERGE) && !(!pool.replicateNoteDbMetaRefs() && RefNames.isNoteDbMetaRef(ref)) - && pool.excludedRefsPattern().stream().noneMatch(p -> p.matcher(ref).matches()); + && !pool.isRefExcluded(ref); } private Map<String, Ref> listRemote(Transport tn) @@ -832,6 +842,7 @@ throws UpdateRefFailureException { Set<String> doneRefs = new HashSet<>(); boolean anyRefFailed = false; + boolean autoRepairNeeded = false; RemoteRefUpdate.Status lastRefStatusError = RemoteRefUpdate.Status.OK; for (RemoteRefUpdate u : refUpdates) { @@ -877,6 +888,10 @@ || UPDATE_REF_FAILURE.equals(u.getMessage())) { throw new UpdateRefFailureException(uri, u.getMessage()); } else { + if (!autoRepairNeeded + && AutoRepairHandler.isMissingNecessaryObjectsError(u.getMessage())) { + autoRepairNeeded = true; + } stateLog.error( String.format( "Failed replicate of %s to %s, reason: %s", @@ -911,6 +926,10 @@ projectName.get(), entry.getKey(), uri, RefPushResult.NOT_ATTEMPTED, null); } } + if (autoRepairNeeded) { + autoRepairHandler.handle( + projectName, uri, pool.getRemoteConfigName(), pool.getUrlDistributionStrategy()); + } stateMap.clear(); }
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 79bdbf3..108250b 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/replication/RemoteConfiguration.java +++ b/src/main/java/com/googlesource/gerrit/plugins/replication/RemoteConfiguration.java
@@ -161,6 +161,16 @@ } /** + * Whether a ref is excluded from replication by any of the {@link #excludedRefsPattern()} + * + * @param ref name of the ref to check + * @return true if the ref should not be replicated, false otherwise + */ + default boolean isRefExcluded(String ref) { + return excludedRefsPattern().stream().anyMatch(p -> p.matcher(ref).matches()); + } + + /** * reflog storage flag for newly created repositories * * @return true if new repositories should store ref-updates in their reflog
diff --git a/src/main/java/com/googlesource/gerrit/plugins/replication/RepairCommand.java b/src/main/java/com/googlesource/gerrit/plugins/replication/RepairCommand.java index 0371c9c..7aa7e94 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/replication/RepairCommand.java +++ b/src/main/java/com/googlesource/gerrit/plugins/replication/RepairCommand.java
@@ -21,9 +21,12 @@ import com.google.gerrit.sshd.CommandMetaData; import com.google.gerrit.sshd.SshCommand; import com.google.inject.Inject; +import com.googlesource.gerrit.plugins.replication.ProjectRepairer.Action; import com.googlesource.gerrit.plugins.replication.api.ReplicationConfig; import java.io.IOException; +import java.io.InterruptedIOException; import java.io.OutputStream; +import java.util.ArrayList; import java.util.Collection; import java.util.Collections; import java.util.HashSet; @@ -46,13 +49,26 @@ usage = "substring URL must match (or * to match everything)") private String urlMatch; + private final List<Action> actions = new ArrayList<>(); + + @Option( + name = "--copy-loose-objects", + usage = "rsync loose object files to SSH destinations before triggering replication") + void setCopyLooseObjects(@SuppressWarnings("unused") boolean arg) { + actions.add(Action.COPY_LOOSE_OBJECTS); + } + @Option( name = "--copy-packs", usage = "rsync objects/pack files to SSH destinations before triggering replication") - private boolean copyPacks; + void setCopyPacks(@SuppressWarnings("unused") boolean arg) { + actions.add(Action.COPY_PACKS); + } @Option(name = "--full", usage = "run all supported repair actions (default)") - private boolean full; + void setFull(@SuppressWarnings("unused") boolean arg) { + actions.addAll(Action.all()); + } @Inject private ProjectCache projectCache; @Inject private ReplicationDestinations destinations; @@ -62,7 +78,7 @@ private final Object outputLock = new Object(); @Override - protected void run() throws Failure { + protected void run() throws Failure, InterruptedIOException { Project.NameKey project = Project.nameKey(projectName); try { if (projectCache.get(project).isEmpty()) { @@ -72,17 +88,18 @@ throw die(e); } - if (!copyPacks) { - full = true; - } - - Set<URIish> failedUris = repair(project); + Set<URIish> failedUris = repair(project, repairActions()); if (!failedUris.isEmpty()) { throw new UnloggedFailure(1, "Repair failed for " + failedUris.size() + " destination(s)"); } } - private Set<URIish> repair(Project.NameKey project) throws Failure { + private Collection<Action> repairActions() { + return actions.isEmpty() ? Action.all() : actions; + } + + private Set<URIish> repair(Project.NameKey project, Collection<Action> actions) + throws Failure, InterruptedIOException { Set<URIish> copyTargets = new HashSet<>(); Collection<URIish> destUris = destinations @@ -91,7 +108,7 @@ for (URIish uri : destUris) { if (!ProjectRepairer.canCopy(uri)) { writeStdErrSync( - "Warning: skipping " + uri + " as copy-packs only supports plain SSH destinations"); + "Warning: skipping " + uri + " as repair only supports plain SSH destinations"); continue; } copyTargets.add(uri); @@ -105,11 +122,12 @@ OutputStream out = getFlushingOutputStream(); for (URIish uri : copyTargets) { writeStdOutSync("\nRepairing " + uri + " ..."); - if (projectRepairer.repair(project, uri, out, full || copyPacks)) { + if (projectRepairer.repair(project, uri, out, actions)) { writeStdOutSync( "\nRunning replication start for " + project.get() + " to " + uri.toString() + " ..."); replicationStarter.start( uri.toString(), + PushOne.ALL_REFS, Set.of(), new ReplicationFilter(List.of(project.get()), Collections.emptyList()), /* now= */ true,
diff --git a/src/main/java/com/googlesource/gerrit/plugins/replication/ReplicationConfigImpl.java b/src/main/java/com/googlesource/gerrit/plugins/replication/ReplicationConfigImpl.java index f625232..72cd417 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/replication/ReplicationConfigImpl.java +++ b/src/main/java/com/googlesource/gerrit/plugins/replication/ReplicationConfigImpl.java
@@ -24,11 +24,15 @@ import com.google.inject.Inject; import com.googlesource.gerrit.plugins.replication.api.ReplicationConfig; import java.nio.file.Path; +import java.time.Duration; import org.eclipse.jgit.lib.Config; public class ReplicationConfigImpl implements ReplicationConfig { private static final int DEFAULT_SSH_CONNECTION_TIMEOUT_MS = 2 * 60 * 1000; // 2 minutes private static final String DEFAULT_RSYNC_PATH = "rsync"; + private static final Duration DEFAULT_AUTO_REPAIR_INTERVAL = Duration.ofDays(3); + private static final int DEFAULT_AUTO_REPAIR_MAX_ATTEMPTS = 0; + private static final int DEFAULT_AUTO_REPAIR_CONCURRENCY_LIMIT = 2; private final SitePaths site; private final MergedConfigResource configResource; @@ -39,6 +43,9 @@ private final int maxRefsToShow; private int sshCommandTimeout; private int sshConnectionTimeout; + private final Duration autoRepairInterval; + private final int autoRepairMaxAttempts; + private final int autoRepairConcurrencyLimit; private final Path pluginDataDir; private final Config config; @@ -63,6 +70,24 @@ "sshConnectionTimeout", DEFAULT_SSH_CONNECTION_TIMEOUT_MS, MILLISECONDS); + this.autoRepairInterval = + Duration.ofSeconds( + ConfigUtil.getTimeUnit( + config, + "replication", + null, + "autoRepairInterval", + DEFAULT_AUTO_REPAIR_INTERVAL.getSeconds(), + SECONDS)); + this.autoRepairMaxAttempts = + config.getInt("replication", "autoRepairMaxAttempts", DEFAULT_AUTO_REPAIR_MAX_ATTEMPTS); + this.autoRepairConcurrencyLimit = + Math.max( + 1, + config.getInt( + "replication", + "autoRepairConcurrencyLimit", + DEFAULT_AUTO_REPAIR_CONCURRENCY_LIMIT)); this.pluginDataDir = pluginDataDir; this.useLegacyCredentials = config.getBoolean("gerrit", "useLegacyCredentials", false); } @@ -152,4 +177,19 @@ String rsyncPath = getConfig().getString("replication", null, "rsyncPath"); return Strings.isNullOrEmpty(rsyncPath) ? DEFAULT_RSYNC_PATH : rsyncPath; } + + @Override + public Duration getAutoRepairInterval() { + return autoRepairInterval; + } + + @Override + public int getAutoRepairMaxAttempts() { + return autoRepairMaxAttempts; + } + + @Override + public int getAutoRepairConcurrencyLimit() { + return autoRepairConcurrencyLimit; + } }
diff --git a/src/main/java/com/googlesource/gerrit/plugins/replication/ReplicationLogFile.java b/src/main/java/com/googlesource/gerrit/plugins/replication/ReplicationLogFile.java index 622034a..7e51f9f 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/replication/ReplicationLogFile.java +++ b/src/main/java/com/googlesource/gerrit/plugins/replication/ReplicationLogFile.java
@@ -22,8 +22,11 @@ import com.google.gerrit.util.logging.JsonLogEntry; import com.google.gson.annotations.SerializedName; import com.google.inject.Inject; +import java.util.HashMap; +import java.util.Map; import org.apache.log4j.PatternLayout; import org.apache.log4j.spi.LoggingEvent; +import org.apache.log4j.spi.ThrowableInformation; public class ReplicationLogFile extends PluginLogFile { @@ -43,13 +46,18 @@ private class ReplicationJsonLogEntry extends JsonLogEntry { public String timestamp; public String message; + public Map<String, String> exception; @SerializedName("@version") - public final int version = 1; + public final int version = 2; public ReplicationJsonLogEntry(LoggingEvent event) { timestamp = timestampFormatter.format(event.getTimeStamp()); message = (String) event.getMessage(); + + if (event.getThrowableInformation() != null) { + this.exception = getException(event.getThrowableInformation()); + } } } @@ -57,5 +65,25 @@ public JsonLogEntry toJsonLogEntry(LoggingEvent event) { return new ReplicationJsonLogEntry(event); } + + private Map<String, String> getException(ThrowableInformation throwable) { + HashMap<String, String> exceptionInformation = new HashMap<>(); + + String throwableName = throwable.getThrowable().getClass().getCanonicalName(); + if (throwableName != null) { + exceptionInformation.put("exception_class", throwableName); + } + + String throwableMessage = throwable.getThrowable().getMessage(); + if (throwableMessage != null) { + exceptionInformation.put("exception_message", throwableMessage); + } + + String[] stackTrace = throwable.getThrowableStrRep(); + if (stackTrace != null) { + exceptionInformation.put("stacktrace", String.join("\n", stackTrace)); + } + return exceptionInformation; + } } }
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..711228b 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()) @@ -106,9 +110,21 @@ bind(ReplicationQueue.class).in(Scopes.SINGLETON); bind(ReplicationDestinations.class).to(DestinationsCollection.class); + bind(AutoRepairHandler.class).in(Scopes.SINGLETON); + bind(AutoRepairTracker.class).in(Scopes.SINGLETON); bind(ProjectRepairer.class).in(Scopes.SINGLETON); 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..2ff785e 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) { @@ -103,8 +102,11 @@ if (!running) { destinations.get().startup(workQueue); running = true; - replicationTasksStorage.recoverAll(); - synchronizePendingEvents(Prune.FALSE); + Set<String> pushEnabledRemoteNames = getPushEnabledRemoteNames(); + if (!pushEnabledRemoteNames.isEmpty()) { + replicationTasksStorage.recoverAll(r -> pushEnabledRemoteNames.contains(r.remote())); + synchronizePendingEvents(Prune.FALSE); + } fireBeforeStartupEvents(); distributor = new Distributor(workQueue); } @@ -131,13 +133,9 @@ } public void scheduleFullSync( - Project.NameKey project, String urlMatch, ReplicationState state, boolean now) { - scheduleFullSync(project, urlMatch, Set.of(), state, now); - } - - public void scheduleFullSync( Project.NameKey project, String urlMatch, + String refName, Set<String> remotesToConsider, ReplicationState state, boolean now) { @@ -145,7 +143,7 @@ project, urlMatch, remotesToConsider, - Set.of(new GitReferenceUpdated.UpdatedRef(PushOne.ALL_REFS, null, null, null)), + Set.of(new GitReferenceUpdated.UpdatedRef(refName, null, null, null)), state, now); } @@ -155,8 +153,15 @@ fire(event.getProjectName(), event.getUpdatedRefs()); } + private Set<String> getPushEnabledRemoteNames() { + return destinations.get().getAll(FilterType.ALL).stream() + .filter(Destination::isPushEnabled) + .map(Destination::getRemoteConfigName) + .collect(Collectors.toSet()); + } + 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 +204,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 +225,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) { @@ -264,10 +269,13 @@ @Override public void run(ReplicationTasksStorage.ReplicateRefUpdate u) { try { - fireFromStorage(new URIish(u.uri()), Project.nameKey(u.project()), u.refs()); if (Prune.TRUE.equals(prune)) { - taskNamesByReplicateRefUpdate.remove(u); + if (taskNamesByReplicateRefUpdate.remove(u) != null) { + repLog.atFine().log("Task %s is already scheduled, not re-firing", u); + return; + } } + fireFromStorage(new URIish(u.uri()), Project.nameKey(u.project()), u.refs()); } catch (URISyntaxException e) { repLog.atSevere().withCause(e).log( "Encountered malformed URI for persisted event %s", u); @@ -384,7 +392,7 @@ } try { synchronizePendingEvents(Prune.TRUE); - } catch (Exception e) { + } catch (RuntimeException e) { repLog.atSevere().withCause(e).log("error distributing tasks"); } }
diff --git a/src/main/java/com/googlesource/gerrit/plugins/replication/ReplicationStarter.java b/src/main/java/com/googlesource/gerrit/plugins/replication/ReplicationStarter.java index b7baab9..e715747 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/replication/ReplicationStarter.java +++ b/src/main/java/com/googlesource/gerrit/plugins/replication/ReplicationStarter.java
@@ -37,16 +37,17 @@ void start( @Nullable String urlMatch, + String refName, Set<String> remotesToConsider, ReplicationFilter filter, boolean now, boolean wait, - SshOutputCommand sink) { - ReplicationState state = new ReplicationState(new CommandProcessing(sink)); + PushResultProcessing processing) { + ReplicationState state = new ReplicationState(processing); Future<?> future = pushFactory - .create(urlMatch, remotesToConsider, filter, state, now) + .create(urlMatch, refName, remotesToConsider, filter, state, now) .schedule(0, TimeUnit.SECONDS); if (wait) { @@ -67,11 +68,32 @@ try { state.waitForReplication(); } catch (InterruptedException e) { - sink.writeStdErrSync("We are interrupted while waiting replication to complete"); + processing.writeStdErr("We are interrupted while waiting replication to complete"); } } else { - sink.writeStdOutSync("Nothing to replicate"); + processing.writeStdOut("Nothing to replicate"); } } } + + void start( + @Nullable String urlMatch, + String refName, + Set<String> remotesToConsider, + ReplicationFilter filter, + boolean now, + boolean wait, + SshOutputCommand sink) { + start(urlMatch, refName, remotesToConsider, filter, now, wait, new CommandProcessing(sink)); + } + + void start( + @Nullable String urlMatch, + Set<String> remotesToConsider, + ReplicationFilter filter, + boolean now, + boolean wait, + PushResultProcessing processing) { + start(urlMatch, PushOne.ALL_REFS, remotesToConsider, filter, now, wait, processing); + } }
diff --git a/src/main/java/com/googlesource/gerrit/plugins/replication/ReplicationTasksStorage.java b/src/main/java/com/googlesource/gerrit/plugins/replication/ReplicationTasksStorage.java index 7fc4acb..bc0d900 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/replication/ReplicationTasksStorage.java +++ b/src/main/java/com/googlesource/gerrit/plugins/replication/ReplicationTasksStorage.java
@@ -44,6 +44,7 @@ import java.util.HashSet; import java.util.Optional; import java.util.Set; +import java.util.function.Predicate; import java.util.stream.Stream; import org.eclipse.jgit.lib.ObjectId; import org.eclipse.jgit.transport.URIish; @@ -189,8 +190,13 @@ } } + @VisibleForTesting public void recoverAll() { - streamRunning().forEach(r -> new Task(r).recover()); + recoverAll(r -> true); + } + + public void recoverAll(Predicate<ReplicateRefUpdate> shouldRecover) { + streamRunning().filter(shouldRecover).forEach(r -> new Task(r).recover()); } public boolean isWaiting(UriUpdates uriUpdates) {
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 2c084b1..d6d99a8 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/replication/StartCommand.java +++ b/src/main/java/com/googlesource/gerrit/plugins/replication/StartCommand.java
@@ -37,6 +37,9 @@ @Option(name = "--url", metaVar = "SUBSTRING", usage = "substring URL must match (or * to match everything)") private String urlMatch; + @Option(name = "--ref", metaVar = "RefName", usage = "ref to replicate") + private String refName = PushOne.ALL_REFS; + private final Set<String> remotesToConsider = new HashSet<>(); @Option(name = "--remote", metaVar = "REMOTE", usage = "name of remote to replicate to") @@ -68,7 +71,7 @@ ? ReplicationFilter.all() : new ReplicationFilter(projectPatterns, Collections.emptyList()); - replicationStarter.start(urlMatch, remotesToConsider, projectFilter, now, wait, this); + replicationStarter.start(urlMatch, refName, remotesToConsider, projectFilter, now, wait, this); } @Override
diff --git a/src/main/java/com/googlesource/gerrit/plugins/replication/api/ReplicationConfig.java b/src/main/java/com/googlesource/gerrit/plugins/replication/api/ReplicationConfig.java index 1714c1a..6860098 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/replication/api/ReplicationConfig.java +++ b/src/main/java/com/googlesource/gerrit/plugins/replication/api/ReplicationConfig.java
@@ -15,6 +15,7 @@ package com.googlesource.gerrit.plugins.replication.api; import java.nio.file.Path; +import java.time.Duration; import org.eclipse.jgit.lib.Config; /** Configuration of all the replication end points. */ @@ -97,6 +98,30 @@ String getRsyncPath(); /** + * Minimum interval between automatic repair attempts for the same project on the same + * destination. + * + * @return interval as a {@link Duration}, {@link Duration#ZERO} for no minimum interval between + * attempts. + */ + Duration getAutoRepairInterval(); + + /** + * Maximum number of automatic repair attempts per project on each destination. + * + * @return maximum attempts, zero to disable auto-repair. + */ + int getAutoRepairMaxAttempts(); + + /** + * Maximum number of automatic repair tasks allowed to run concurrently. Changing this value + * requires a plugin reload to take effect. + * + * @return concurrency limit, minimum 1. + */ + int getAutoRepairConcurrencyLimit(); + + /** * Current logical version string of the current configuration loaded in memory, depending on the * actual implementation of the configuration on the persistent storage. *
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/cmd-repair.md b/src/main/resources/Documentation/cmd-repair.md index 3f95907..ca20350 100644 --- a/src/main/resources/Documentation/cmd-repair.md +++ b/src/main/resources/Documentation/cmd-repair.md
@@ -11,7 +11,7 @@ ```console ssh -p @SSH_PORT@ @SSH_HOST@ @PLUGIN@ repair [--url <PATTERN>] - [--full | --copy-packs] + [--full | [--copy-loose-objects] [--copy-packs]] <PROJECT> ``` @@ -24,6 +24,13 @@ If no repair action flag is supplied, `--full` is assumed. +For each remote, [remote.NAME.adminUrl](config.md#remote.NAME.adminUrl) is +preferred when set (same as repository creation); otherwise +[remote.NAME.url](config.md#remote.NAME.url) is used. Only plain SSH +destinations are eligible (for example `user@host:/path/to/repo.git`). +Destinations whose URL uses `gerrit+ssh`, HTTP(S), or a local path are +skipped. + REQUIREMENTS ------------ The Gerrit runtime user must have `ssh` on `PATH`, plus `rsync` either on @@ -51,15 +58,13 @@ `--full` : Run every supported repair action. +`--copy-loose-objects` +: rsync the loose objects. Everything else under `objects/`, such as +`pack/` and `info/`, is left untouched. + `--copy-packs` : rsync regular files in `objects/pack/` whose names end with `.pack`, -`.idx`, `.bitmap`, or `.rev` to each matching destination. For each -remote, [remote.NAME.adminUrl](config.md#remote.NAME.adminUrl) is preferred -when set (same as repository creation); otherwise -[remote.NAME.url](config.md#remote.NAME.url) is used. Only plain SSH -destinations are eligible (for example `user@host:/path/to/repo.git`). -Destinations whose URL uses `gerrit+ssh`, HTTP(S), or a local path are -skipped. +`.idx`, `.bitmap`, or `.rev` to each matching destination. `PROJECT` : Exact Gerrit project name. @@ -86,6 +91,12 @@ $ ssh -p @SSH_PORT@ @SSH_HOST@ @PLUGIN@ repair --copy-packs tools/gerrit ``` +Only copy loose objects: + +```console + $ ssh -p @SSH_PORT@ @SSH_HOST@ @PLUGIN@ repair --copy-loose-objects tools/gerrit +``` + Repair only against destinations whose URL mentions `replica1`: ```console
diff --git a/src/main/resources/Documentation/cmd-start.md b/src/main/resources/Documentation/cmd-start.md index e6d9fbd..47989a1 100644 --- a/src/main/resources/Documentation/cmd-start.md +++ b/src/main/resources/Documentation/cmd-start.md
@@ -13,6 +13,7 @@ [--now] [--wait] [--remote <REMOTE> ...] + [--ref <RefName>] {--url <PATTERN> | [--url <PATTERN>] --all | [--url <PATTERN>] <PROJECT PATTERN> ...} ``` @@ -97,6 +98,10 @@ `--all` : Schedule replication for all projects. +`--ref <RefName>` +: Replicate only the single reference specified by `<RefName>`. +If omitted, the command defaults to all references (`refs/*`). + `--url <PATTERN>` : Replicate only to replication destinations whose configuration URL contains the substring `PATTERN`, or whose expanded project @@ -131,6 +136,12 @@ $ ssh -p @SSH_PORT@ @SSH_HOST@ @PLUGIN@ start tools/gerrit ``` +Replicate only the `master` branch of the `tools/gerrit` project: + +```console + $ ssh -p @SSH_PORT@ @SSH_HOST@ @PLUGIN@ start --ref refs/heads/master tools/gerrit +``` + Replicate only projects located in the `documentation` subdirectory: ```console
diff --git a/src/main/resources/Documentation/config.md b/src/main/resources/Documentation/config.md index 083f4ee..014ca15 100644 --- a/src/main/resources/Documentation/config.md +++ b/src/main/resources/Documentation/config.md
@@ -261,14 +261,69 @@ 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 - to SSH destinations. Set this when the Gerrit runtime user's `PATH` - does not contain `rsync`, or to pin a specific build. + [`@PLUGIN@ repair`](cmd-repair.md) command when transferring loose + objects and pack files to SSH destinations. Set this when the Gerrit + runtime user's `PATH` does not contain `rsync`, or to pin a specific + build. Default: `rsync` (resolved via the Gerrit runtime user's `PATH`) +replication.autoRepairInterval +: Minimum interval between automatic repair attempts for the same + project on the same destination. Values are expressed with a time + unit suffix, e.g. `2h`, `30m`, `3d`. If no unit is given the value is + interpreted as seconds. + + Auto-repair is event-driven and each attempt is initiated by a failing + replication push and is not retried on its own. When a push fails with + a `missing necessary objects` error from `git-receive-pack`, the plugin + schedules a repair attempt for the affected SSH destination. Unlike + normal replication failures (which the plugin retries automatically), + repair attempts are not retried by themselves. If a repair completes + but the destination is still broken, the next repair runs only when + another replication push to the same destination fails with the + same missing objects error, and only after this interval has elapsed. + + Default: `3d` + +replication.autoRepairMaxAttempts +: Maximum number of automatic repair attempts per project on each + destination. After this limit is reached for a given project and + destination, further `missing necessary objects` errors for that pair + are logged but no additional automatic repairs are attempted. + + Auto-repair state is kept in memory and is reset whenever the plugin is + reloaded or Gerrit restarts. Enabling auto-repair in a clustered deployment + can lead to redundant repairs as the state is not shared between them. + + Default: `0` (auto-repair disabled) + +replication.autoRepairConcurrencyLimit +: Maximum number of automatic repair tasks allowed to run concurrently. + Each running repair occupies one slot for the duration of its repair + and the follow-up full replication, so this caps the load auto-repair + can put on the host and on destination SSH endpoints at any one time. + Increase it when many destinations are expected to need auto-repair + in parallel. The value is read once at plugin start, so changes only + take effect after a plugin reload. + + Minimum: `1`. Default: `2`. + remote.NAME.url : Address of the remote server to push to. Multiple URLs may be specified within a single remote block, listing different @@ -446,6 +501,11 @@ This is a Gerrit specific extension to the Git remote block. + For backwards compatibility, a bare number is interpreted as + minutes. A time-unit suffix may be appended to configure a + finer granularity, for example `30 s`, `90 s` or `2 m`; the value is + applied at second precision. + By default, 1 minute. remote.NAME.replicationMaxRetries @@ -482,6 +542,17 @@ remote block describes 4 URLs, allocating 4 threads in the pool will permit some level of parallel pushing. + A value of `0` puts this remote into a persist-only mode. No thread + pool is created and this node performs no pushes. Replication tasks + are still written to the replication task storage, so that other + nodes in the cluster sharing the same storage (with `threads` > 0 and + `replication.distributionInterval` enabled) pick them up and perform + the actual pushes. + + > **NOTE**: In persist-only mode, project deletions and HEAD updates + > originating on this node cannot be replicated by other nodes, as + > they are not persisted to the task storage. + By default, 1 thread. remote.NAME.authGroup @@ -666,12 +737,17 @@ By default, true. remote.NAME.excludedRefsPattern -: Refs that match the pattern provided using this config will not be replicated. - This option is useful when admins want to skip replicating certain refs, for - example refs created by plugins. Multiple excludedRefsPattern keys can be - supplied, to specify multiple patterns to match against. +: Refs that match the regular expression provided using this config will not + be replicated. Only Java regular expressions are supported. The expression + must fully match the ref name. This option is useful when admins want to + skip replicating certain refs, for example refs created by plugins. Multiple + `excludedRefsPattern` keys can be supplied, to specify multiple regular + expressions to match against. - Do not exclude any refs pattern by default. + Excluded refs are filtered out before a replication task is scheduled, so they + do not appear in the replication queue or in the persisted task storage. + + Do not exclude any refs by default. remote.NAME.urlDistributionStrategy : URL distribution strategy to use when a remote has multiple configured URLs.
diff --git a/src/test/java/com/googlesource/gerrit/plugins/replication/AutoRepairTrackerTest.java b/src/test/java/com/googlesource/gerrit/plugins/replication/AutoRepairTrackerTest.java new file mode 100644 index 0000000..865ed23 --- /dev/null +++ b/src/test/java/com/googlesource/gerrit/plugins/replication/AutoRepairTrackerTest.java
@@ -0,0 +1,114 @@ +// 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 static org.mockito.Mockito.when; + +import com.google.gerrit.entities.Project; +import com.googlesource.gerrit.plugins.replication.api.ReplicationConfig; +import java.time.Duration; +import org.eclipse.jgit.transport.URIish; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; +import org.mockito.Mock; +import org.mockito.MockitoAnnotations; + +public class AutoRepairTrackerTest { + private static final Project.NameKey PROJECT = Project.nameKey("foo"); + private static final String REMOTE = "mirror"; + + @Mock private ReplicationConfig replicationConfig; + + private AutoCloseable mocks; + private AutoRepairTracker tracker; + private URIish uri1; + private URIish uri2; + + @Before + public void setUp() throws Exception { + mocks = MockitoAnnotations.openMocks(this); + when(replicationConfig.getAutoRepairInterval()).thenReturn(Duration.ofDays(3)); + when(replicationConfig.getAutoRepairMaxAttempts()).thenReturn(5); + + tracker = new AutoRepairTracker(replicationConfig); + uri1 = new URIish("ssh://mirror1.example.com/foo.git"); + uri2 = new URIish("ssh://mirror2.example.com/foo.git"); + } + + @After + public void tearDown() throws Exception { + mocks.close(); + } + + @Test + public void allowsRepairUpToMaxAttemptsPerDestination() { + when(replicationConfig.getAutoRepairInterval()).thenReturn(Duration.ZERO); + for (int i = 0; i < 5; i++) { + assertThat(tryRepair(uri1, UrlDistributionStrategy.ALL)).isTrue(); + } + assertThat(tryRepair(uri1, UrlDistributionStrategy.ALL)).isFalse(); + } + + @Test + public void tracksAttemptsSeparatelyPerDestination() { + when(replicationConfig.getAutoRepairInterval()).thenReturn(Duration.ZERO); + when(replicationConfig.getAutoRepairMaxAttempts()).thenReturn(1); + + assertThat(tryRepair(uri1, UrlDistributionStrategy.ALL)).isTrue(); + assertThat(tryRepair(uri1, UrlDistributionStrategy.ALL)).isFalse(); + assertThat(tryRepair(uri2, UrlDistributionStrategy.ALL)).isTrue(); + } + + @Test + public void roundRobinSharesAttemptsAcrossUrls() { + when(replicationConfig.getAutoRepairInterval()).thenReturn(Duration.ZERO); + when(replicationConfig.getAutoRepairMaxAttempts()).thenReturn(1); + + assertThat(tryRepair(uri1, UrlDistributionStrategy.ROUND_ROBIN)).isTrue(); + assertThat(tryRepair(uri2, UrlDistributionStrategy.ROUND_ROBIN)).isFalse(); + } + + @Test + public void allowsRepairWithoutIntervalWhenIntervalIsZero() { + when(replicationConfig.getAutoRepairInterval()).thenReturn(Duration.ZERO); + assertThat(tryRepair(uri1, UrlDistributionStrategy.ALL)).isTrue(); + assertThat(tryRepair(uri1, UrlDistributionStrategy.ALL)).isTrue(); + } + + @Test + public void disabledWhenMaxAttemptsIsZero() { + when(replicationConfig.getAutoRepairMaxAttempts()).thenReturn(0); + assertThat(tracker.isEnabled()).isFalse(); + assertThat(tryRepair(uri1, UrlDistributionStrategy.ALL)).isFalse(); + } + + @Test + public void enforcesIntervalBetweenAttempts() { + assertThat(tryRepair(uri1, UrlDistributionStrategy.ALL)).isTrue(); + assertThat(tryRepair(uri1, UrlDistributionStrategy.ALL)).isFalse(); + } + + @Test + public void rejectsNonCopyableDestination() throws Exception { + URIish httpUri = new URIish("http://mirror.example.com/foo.git"); + assertThat(tryRepair(httpUri, UrlDistributionStrategy.ALL)).isFalse(); + } + + private boolean tryRepair(URIish uri, UrlDistributionStrategy strategy) { + return tracker.tryBeginRepair(PROJECT, uri, REMOTE, strategy); + } +}
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 f5118dd..e178ab1 100644 --- a/src/test/java/com/googlesource/gerrit/plugins/replication/DestinationConfigurationTest.java +++ b/src/test/java/com/googlesource/gerrit/plugins/replication/DestinationConfigurationTest.java
@@ -15,6 +15,7 @@ package com.googlesource.gerrit.plugins.replication; import static com.google.common.truth.Truth.assertThat; +import static org.junit.Assert.assertThrows; import static org.mockito.Mockito.when; import org.eclipse.jgit.lib.Config; @@ -100,6 +101,96 @@ } @Test + public void shouldDefaultReplicationRetryToOneMinute() { + assertThat(objectUnderTest.getRetryDelay()).isEqualTo(60); + } + + @Test + public void shouldTreatBareReplicationRetryAsMinutes() { + // given + when(cfgMock.getString("remote", REMOTE, "replicationRetry")).thenReturn("2"); + objectUnderTest = new DestinationConfiguration(remoteConfigMock, cfgMock); + + // when / then + assertThat(objectUnderTest.getRetryDelay()).isEqualTo(120); + } + + @Test + public void shouldParseReplicationRetryWithSecondsSuffix() { + // given + when(cfgMock.getString("remote", REMOTE, "replicationRetry")).thenReturn("30 s"); + objectUnderTest = new DestinationConfiguration(remoteConfigMock, cfgMock); + + // when / then + assertThat(objectUnderTest.getRetryDelay()).isEqualTo(30); + } + + @Test + public void shouldParseReplicationRetryWithMinutesSuffix() { + // given + when(cfgMock.getString("remote", REMOTE, "replicationRetry")).thenReturn("2 m"); + objectUnderTest = new DestinationConfiguration(remoteConfigMock, cfgMock); + + // when / then + assertThat(objectUnderTest.getRetryDelay()).isEqualTo(120); + } + + @Test + public void shouldTreatZeroReplicationRetryAsNoDelay() { + // given + when(cfgMock.getString("remote", REMOTE, "replicationRetry")).thenReturn("0"); + objectUnderTest = new DestinationConfiguration(remoteConfigMock, cfgMock); + + // when / then + assertThat(objectUnderTest.getRetryDelay()).isEqualTo(0); + } + + @Test + public void shouldClampNegativeBareReplicationRetryToZero() { + // given: a bare negative was accepted historically (cfg.getInt) and clamped to zero + when(cfgMock.getString("remote", REMOTE, "replicationRetry")).thenReturn("-1"); + objectUnderTest = new DestinationConfiguration(remoteConfigMock, cfgMock); + + // when / then + assertThat(objectUnderTest.getRetryDelay()).isEqualTo(0); + } + + @Test + public void shouldDefaultReplicationRetryWhenEmpty() { + // given + when(cfgMock.getString("remote", REMOTE, "replicationRetry")).thenReturn(" "); + objectUnderTest = new DestinationConfiguration(remoteConfigMock, cfgMock); + + // when / then + assertThat(objectUnderTest.getRetryDelay()).isEqualTo(60); + } + + @Test + public void shouldClampHugeReplicationRetryToIntMax() { + // given: a value whose seconds exceed Integer.MAX_VALUE must not overflow to a negative int + when(cfgMock.getString("remote", REMOTE, "replicationRetry")).thenReturn("999999999999 s"); + objectUnderTest = new DestinationConfiguration(remoteConfigMock, cfgMock); + + // when / then + assertThat(objectUnderTest.getRetryDelay()).isEqualTo(Integer.MAX_VALUE); + } + + @Test + public void shouldRejectInvalidReplicationRetry() { + // given + when(cfgMock.getString("remote", REMOTE, "replicationRetry")).thenReturn("banana"); + + // when + IllegalArgumentException thrown = + assertThrows( + IllegalArgumentException.class, + () -> new DestinationConfiguration(remoteConfigMock, cfgMock)); + + // then: the error names the offending config key + assertThat(thrown).hasMessageThat().contains("remote." + REMOTE + ".replicationRetry"); + } + + @Test public void shouldDefaultUrlDistributionToAll() { assertThat(objectUnderTest.getUrlDistributionStrategy()).isEqualTo(UrlDistributionStrategy.ALL); }
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/PushOneTest.java b/src/test/java/com/googlesource/gerrit/plugins/replication/PushOneTest.java index 339123a..076a903 100644 --- a/src/test/java/com/googlesource/gerrit/plugins/replication/PushOneTest.java +++ b/src/test/java/com/googlesource/gerrit/plugins/replication/PushOneTest.java
@@ -53,7 +53,6 @@ import java.util.concurrent.Callable; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; -import java.util.regex.Pattern; import org.eclipse.jgit.errors.NotSupportedException; import org.eclipse.jgit.errors.RepositoryNotFoundException; import org.eclipse.jgit.errors.TransportException; @@ -297,9 +296,8 @@ @Test public void skipPushingExcludedRefs() throws InterruptedException, IOException { - when(destinationMock.excludedRefsPattern()) - .thenReturn( - ImmutableList.of(Pattern.compile("refs/foo/.*"), Pattern.compile("refs/bar/.*"))); + when(destinationMock.isRefExcluded("refs/foo/test")).thenReturn(true); + when(destinationMock.isRefExcluded("refs/bar/test")).thenReturn(true); PushOne pushOne = Mockito.spy(createPushOne(null)); Ref ref1 = @@ -423,6 +421,7 @@ projectCacheMock, createProjectTaskFactoryMock, transportFactoryMock, + mock(AutoRepairHandler.class), projectNameKey, urish); @@ -495,7 +494,6 @@ private void setupDestinationMock() { destinationMock = mock(Destination.class); when(destinationMock.requestRunway(any())).thenReturn(RunwayStatus.allowed()); - when(destinationMock.excludedRefsPattern()).thenReturn(ImmutableList.of()); } private void setupPermissionBackedMock() {
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 436fb1c..5764521 100644 --- a/src/test/java/com/googlesource/gerrit/plugins/replication/ReplicationDaemon.java +++ b/src/test/java/com/googlesource/gerrit/plugins/replication/ReplicationDaemon.java
@@ -56,15 +56,23 @@ protected static final int TEST_REPLICATION_DELAY_SECONDS = 1; protected static final int TEST_LONG_REPLICATION_DELAY_SECONDS = 30; protected static final int TEST_REPLICATION_RETRY_MINUTES = 1; + // Sub-minute replicationRetry used by the new-project cases so their missing-repository retry + // fires in a second instead of waiting the full minute (see remote.NAME.replicationRetry). + protected static final int TEST_REPLICATION_RETRY_SECONDS = 1; protected static final int TEST_PUSH_TIME_SECONDS = 1; protected static final int TEST_PROJECT_CREATION_SECONDS = 10; protected static final Duration TEST_PUSH_TIMEOUT = Duration.ofSeconds(TEST_REPLICATION_DELAY_SECONDS + TEST_PUSH_TIME_SECONDS); protected static final Duration TEST_PUSH_TIMEOUT_LONG = Duration.ofSeconds(TEST_LONG_REPLICATION_DELAY_SECONDS + TEST_PUSH_TIME_SECONDS); + // A new project's first ref-push fails as REPOSITORY_MISSING, creates the repository, and is + // retried after replicationRetry; the new-project cases set that to + // TEST_REPLICATION_RETRY_SECONDS + // (no longer a full minute), so this timeout budgets the replication delay + retry + push plus a + // cushion for project creation. protected static final Duration TEST_NEW_PROJECT_TIMEOUT = Duration.ofSeconds( - (TEST_REPLICATION_DELAY_SECONDS + TEST_REPLICATION_RETRY_MINUTES * 60) + (TEST_REPLICATION_DELAY_SECONDS + TEST_REPLICATION_RETRY_SECONDS + TEST_PUSH_TIME_SECONDS) + TEST_PROJECT_CREATION_SECONDS); @Inject private ProjectOperations projectOperations;
diff --git a/src/test/java/com/googlesource/gerrit/plugins/replication/ReplicationDistributorIT.java b/src/test/java/com/googlesource/gerrit/plugins/replication/ReplicationDistributorIT.java index dc73036..5256bbe 100644 --- a/src/test/java/com/googlesource/gerrit/plugins/replication/ReplicationDistributorIT.java +++ b/src/test/java/com/googlesource/gerrit/plugins/replication/ReplicationDistributorIT.java
@@ -15,6 +15,7 @@ package com.googlesource.gerrit.plugins.replication; import static com.google.common.truth.Truth.assertThat; +import static com.google.gerrit.testing.GerritJUnit.assertThrows; import com.google.gerrit.acceptance.TestPlugin; import com.google.gerrit.acceptance.UseLocalDisk; @@ -22,6 +23,7 @@ import com.google.gerrit.entities.BranchNameKey; import com.google.gerrit.entities.Project; import com.google.gerrit.server.git.WorkQueue; +import com.googlesource.gerrit.plugins.replication.api.ReplicationConfig.FilterType; import java.time.Duration; import java.util.List; import java.util.Set; @@ -85,6 +87,40 @@ } @Test + public void distributorDoesNotReFirePendingTask() throws Exception { + String remote = "foo"; + String replica = "replica"; + String master = "refs/heads/master"; + String pendingBranch = "refs/heads/pending_branch"; + String otherPrimaryBranch = "refs/heads/other_primary_branch"; + Project.NameKey targetProject = createTestProject(project + replica); + URIish targetUri = new URIish(getProjectUri(targetProject)); + setReplicationDestination(remote, replica, ALL_PROJECTS, TEST_LONG_REPLICATION_DELAY_SECONDS); + reloadConfig(); + + createBranch(BranchNameKey.create(project, pendingBranch)); + assertThat(listWaitingReplicationTasks(pendingBranch)).hasSize(1); + PushOne pendingPush = getPendingPush(remote, targetUri); + assertThat(pendingPush.getStatesByRef(pendingBranch)).hasLength(1); + + createBranch(project, master, otherPrimaryBranch); + tasksStorage.create( + ReplicationTasksStorage.ReplicateRefUpdate.create( + project.get(), Set.of(otherPrimaryBranch), targetUri, remote)); + + WaitUtil.waitUntil( + () -> pendingPush.getStatesByRef(otherPrimaryBranch).length == 1, + Duration.ofSeconds(TEST_DISTRIBUTION_CYCLE_SECONDS)); + + assertThrows( + InterruptedException.class, + () -> + WaitUtil.waitUntil( + () -> pendingPush.getStatesByRef(pendingBranch).length > 1, + Duration.ofSeconds(TEST_DISTRIBUTION_CYCLE_SECONDS))); + } + + @Test public void distributorPrunesTaskFromWorkQueue() throws Exception { createTestProject(project + "replica"); setReplicationDestination("foo", "replica", ALL_PROJECTS, Integer.MAX_VALUE); @@ -100,6 +136,16 @@ .isTrue(); } + private PushOne getPendingPush(String remote, URIish uri) { + return destinationCollection.getAll(FilterType.ALL).stream() + .filter(dest -> remote.equals(dest.getRemoteConfigName())) + .findFirst() + .get() + .getQueue() + .pending + .get(uri); + } + private List<WorkQueue.Task<?>> getProjectTasks() { return getInstance(WorkQueue.class).getTasks().stream() .filter(t -> t instanceof WorkQueue.ProjectTask)
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 a267337..3e9e825 100644 --- a/src/test/java/com/googlesource/gerrit/plugins/replication/ReplicationIT.java +++ b/src/test/java/com/googlesource/gerrit/plugins/replication/ReplicationIT.java
@@ -59,10 +59,18 @@ name = "replication", sysModule = "com.googlesource.gerrit.plugins.replication.TestReplicationModule") public class ReplicationIT extends ReplicationDaemon { - private static final int TEST_REPLICATION_DELAY = 1; - private static final int TEST_REPLICATION_RETRY = 1; private static final Duration TEST_TIMEOUT = - Duration.ofSeconds((TEST_REPLICATION_DELAY + TEST_REPLICATION_RETRY * 60) + 1); + Duration.ofSeconds(TEST_REPLICATION_DELAY_SECONDS + TEST_REPLICATION_RETRY_MINUTES * 60 + 1); + + // Timeout for asserting that a ref is *not* replicated. Unlike TEST_TIMEOUT it + // deliberately excludes the retry cycle: a push that is going to happen at all + // completes within the replication delay plus the push time, so if the ref has + // not appeared within that window (plus a small cushion) it never will + // (destination shut down / project excluded / non-matching remote). Sized in + // seconds rather than the ~60s retry quantum so these negative tests do not each + // burn a full retry window while proving a negative. + private static final Duration TEST_NOT_REPLICATED_TIMEOUT = + Duration.ofSeconds(TEST_REPLICATION_DELAY_SECONDS + TEST_PUSH_TIME_SECONDS + 5); @Inject private DynamicSet<ProjectDeletedListener> deletedListeners; @@ -79,6 +87,7 @@ @Test public void shouldReplicateNewProjectWithoutRefLog() throws Exception { setReplicationDestination("foo", "replica", ALL_PROJECTS); + setSubMinuteReplicationRetry("foo"); reloadConfig(); Project.NameKey sourceProject = createTestProject("no_reflog_project"); @@ -95,6 +104,7 @@ public void shouldCreateNewProjectWithRefLog() throws Exception { config.setBoolean("remote", "foo", "storeRefLog", true); setReplicationDestination("foo", "replica", ALL_PROJECTS); + setSubMinuteReplicationRetry("foo"); reloadConfig(); Project.NameKey sourceProject = createTestProject("reflog_project"); @@ -109,6 +119,14 @@ WaitUtil.waitUntil(() -> nonEmptyProjectExists(replicaProject), TEST_NEW_PROJECT_TIMEOUT); } + // Overrides the minute-scale replicationRetry set by setReplicationDestination with a sub-minute + // value, so the first-ref retry to a just-created project fires in seconds rather than a minute. + private void setSubMinuteReplicationRetry(String remoteName) throws IOException { + config.setString( + "remote", remoteName, "replicationRetry", TEST_REPLICATION_RETRY_SECONDS + "s"); + config.save(); + } + private static Consumer<StoredConfig> assertStoreRefLog(boolean expectedValue) { return conf -> assertThat( @@ -189,6 +207,30 @@ } @Test + public void shouldReplicateOnlySpecificRef() throws Exception { + Project.NameKey targetProject = createTestProject(project + "replica"); + + setReplicationDestination("foo", "replica", ALL_PROJECTS); + reloadConfig(); + + String branch1 = "refs/heads/branch1"; + String branch2 = "refs/heads/branch2"; + createNewBranchWithoutPush("refs/heads/master", branch1); + createNewBranchWithoutPush("refs/heads/master", branch2); + + plugin + .getSysInjector() + .getInstance(ReplicationQueue.class) + .scheduleFullSync(project, null, branch1, Set.of(), new ReplicationState(NO_OP), true); + + try (Repository repo = repoManager.openRepository(targetProject)) { + waitUntil(() -> checkedGetRef(repo, branch1) != null); + assertThat(getRef(repo, branch1)).isNotNull(); + assertThat(getRef(repo, branch2)).isNull(); + } + } + + @Test public void shouldReplicateNewBranchToTwoRemotes() throws Exception { Project.NameKey targetProject1 = createTestProject(project + "replica1"); Project.NameKey targetProject2 = createTestProject(project + "replica2"); @@ -232,7 +274,8 @@ plugin .getSysInjector() .getInstance(ReplicationQueue.class) - .scheduleFullSync(project, urlMatch, new ReplicationState(NO_OP), true); + .scheduleFullSync( + project, urlMatch, PushOne.ALL_REFS, Set.of(), new ReplicationState(NO_OP), true); try (Repository repo = repoManager.openRepository(targetProject)) { waitUntil(() -> checkedGetRef(repo, newRef) != null); @@ -258,7 +301,8 @@ plugin .getSysInjector() .getInstance(ReplicationQueue.class) - .scheduleFullSync(project, urlMatch, new ReplicationState(NO_OP), true); + .scheduleFullSync( + project, urlMatch, PushOne.ALL_REFS, Set.of(), new ReplicationState(NO_OP), true); try (Repository repo = repoManager.openRepository(targetProject)) { waitUntil(() -> checkedGetRef(repo, newRef) != null); @@ -284,6 +328,7 @@ .getInstance(PushAll.Factory.class) .create( null, + PushOne.ALL_REFS, Set.of(), new ReplicationFilter(Arrays.asList(project.get()), null), state, @@ -309,6 +354,7 @@ .getInstance(PushAll.Factory.class) .create( null, + PushOne.ALL_REFS, Set.of(), new ReplicationFilter(Arrays.asList(project.get()), null), state, @@ -429,7 +475,7 @@ InterruptedException.class, () -> { try (Repository repo = repoManager.openRepository(targetProject)) { - waitUntil(() -> checkedGetRef(repo, sourceRef) != null); + waitUntil(() -> checkedGetRef(repo, sourceRef) != null, TEST_NOT_REPLICATED_TIMEOUT); } }); } @@ -518,7 +564,8 @@ plugin .getSysInjector() .getInstance(ReplicationQueue.class) - .scheduleFullSync(project, null, new ReplicationState(NO_OP), true); + .scheduleFullSync( + project, null, PushOne.ALL_REFS, Set.of(), new ReplicationState(NO_OP), true); // Wait for the push to land on both the refs try (Repository r1 = repoManager.openRepository(replica1Project); @@ -543,7 +590,8 @@ plugin .getSysInjector() .getInstance(ReplicationQueue.class) - .scheduleFullSync(project, null, new ReplicationState(NO_OP), true); + .scheduleFullSync( + project, null, PushOne.ALL_REFS, Set.of(), new ReplicationState(NO_OP), true); // Wait for the push to land in at least one replica try (Repository r1 = repoManager.openRepository(replica1Project); @@ -574,7 +622,8 @@ ReplicationQueue queue = plugin.getSysInjector().getInstance(ReplicationQueue.class); // First sync - goes to replica1 (index 0) - queue.scheduleFullSync(project, null, new ReplicationState(NO_OP), true); + queue.scheduleFullSync( + project, null, PushOne.ALL_REFS, Set.of(), new ReplicationState(NO_OP), true); try (Repository r1 = repoManager.openRepository(replica1Project)) { waitUntil(() -> checkedGetRef(r1, branch1) != null); @@ -587,7 +636,8 @@ // 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); + queue.scheduleFullSync( + project, null, PushOne.ALL_REFS, Set.of(), new ReplicationState(NO_OP), true); try (Repository r2 = repoManager.openRepository(replica2Project)) { waitUntil(() -> checkedGetRef(r2, branch1) != null && checkedGetRef(r2, branch2) != null); @@ -615,7 +665,8 @@ plugin .getSysInjector() .getInstance(ReplicationQueue.class) - .scheduleFullSync(project, null, new ReplicationState(NO_OP), true); + .scheduleFullSync( + project, null, PushOne.ALL_REFS, Set.of(), new ReplicationState(NO_OP), true); // Wait for the push to land in at least one replica try (Repository r1 = repoManager.openRepository(replica1Project); @@ -645,7 +696,8 @@ ReplicationQueue queue = plugin.getSysInjector().getInstance(ReplicationQueue.class); - queue.scheduleFullSync(project, null, new ReplicationState(NO_OP), true); + queue.scheduleFullSync( + project, null, PushOne.ALL_REFS, Set.of(), new ReplicationState(NO_OP), true); // The project is pinned to one of the two replicas; discover which one. Project.NameKey pinnedProject; @@ -662,7 +714,8 @@ // 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); + queue.scheduleFullSync( + project, null, PushOne.ALL_REFS, Set.of(), new ReplicationState(NO_OP), true); try (Repository pinnedRepo = repoManager.openRepository(pinnedProject)) { waitUntil(() -> checkedGetRef(pinnedRepo, branch2) != null); @@ -687,7 +740,8 @@ plugin .getSysInjector() .getInstance(ReplicationQueue.class) - .scheduleFullSync(project, null, Set.of("foo"), new ReplicationState(NO_OP), true); + .scheduleFullSync( + project, null, PushOne.ALL_REFS, Set.of("foo"), new ReplicationState(NO_OP), true); try (Repository repo = repoManager.openRepository(targetProject)) { waitUntil(() -> checkedGetRef(repo, newRef) != null); @@ -745,10 +799,10 @@ ReplicationQueue replicationQueue = plugin.getSysInjector().getInstance(ReplicationQueue.class); ReplicationState state = new ReplicationState(NO_OP); - replicationQueue.scheduleFullSync(prj1, null, Set.of("foo"), state, true); - replicationQueue.scheduleFullSync(prj2, null, Set.of("foo"), state, true); - replicationQueue.scheduleFullSync(prj3, null, Set.of("foo"), state, true); - replicationQueue.scheduleFullSync(prj4, null, Set.of("foo"), state, true); + replicationQueue.scheduleFullSync(prj1, null, PushOne.ALL_REFS, Set.of("foo"), state, true); + replicationQueue.scheduleFullSync(prj2, null, PushOne.ALL_REFS, Set.of("foo"), state, true); + replicationQueue.scheduleFullSync(prj3, null, PushOne.ALL_REFS, Set.of("foo"), state, true); + replicationQueue.scheduleFullSync(prj4, null, PushOne.ALL_REFS, Set.of("foo"), state, true); try (Repository excludeRepo1 = repoManager.openRepository(targetPrj1); Repository excludeRepo2 = repoManager.openRepository(targetPrj2); @@ -783,7 +837,8 @@ try (Repository repo = repoManager.openRepository(targetProject)) { assertThrows( InterruptedException.class, - () -> waitUntil(() -> checkedGetRef(repo, sourceRef) != null)); + () -> + waitUntil(() -> checkedGetRef(repo, sourceRef) != null, TEST_NOT_REPLICATED_TIMEOUT)); } } @@ -800,16 +855,23 @@ plugin .getSysInjector() .getInstance(ReplicationQueue.class) - .scheduleFullSync(project, null, Set.of("bar"), new ReplicationState(NO_OP), true); + .scheduleFullSync( + project, null, PushOne.ALL_REFS, Set.of("bar"), new ReplicationState(NO_OP), true); try (Repository repo = repoManager.openRepository(targetProject)) { assertThrows( - InterruptedException.class, () -> waitUntil(() -> checkedGetRef(repo, newRef) != null)); + InterruptedException.class, + () -> waitUntil(() -> checkedGetRef(repo, newRef) != null, TEST_NOT_REPLICATED_TIMEOUT)); } } private void waitUntil(Supplier<Boolean> waitCondition) throws InterruptedException { - WaitUtil.waitUntil(waitCondition, TEST_TIMEOUT); + waitUntil(waitCondition, TEST_TIMEOUT); + } + + private void waitUntil(Supplier<Boolean> waitCondition, Duration timeout) + throws InterruptedException { + WaitUtil.waitUntil(waitCondition, timeout); } private void shutdownDestinations() {
diff --git a/src/test/java/com/googlesource/gerrit/plugins/replication/ReplicationQueueTest.java b/src/test/java/com/googlesource/gerrit/plugins/replication/ReplicationQueueTest.java new file mode 100644 index 0000000..3a10570 --- /dev/null +++ b/src/test/java/com/googlesource/gerrit/plugins/replication/ReplicationQueueTest.java
@@ -0,0 +1,147 @@ +// 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 static java.util.concurrent.TimeUnit.MILLISECONDS; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import com.google.gerrit.entities.Project; +import com.google.gerrit.server.git.WorkQueue; +import com.googlesource.gerrit.plugins.replication.ReplicationTasksStorage.ReplicateRefUpdate; +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.time.Duration; +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.concurrent.ScheduledThreadPoolExecutor; +import org.eclipse.jgit.transport.URIish; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; + +public class ReplicationQueueTest { + private static final String REMOTE = "remote"; + private static final String REF = "refs/heads/master"; + private static final Duration REPLAY_TIMEOUT = Duration.ofSeconds(10); + + private Project.NameKey project; + private URIish uri; + private ReplicateRefUpdate update; + private List<ReplicateRefUpdate> waitingTasks; + private Map<ReplicateRefUpdate, String> taskNamesByReplicateRefUpdate; + private ScheduledThreadPoolExecutor defaultQueue; + private WorkQueue workQueue; + private Destination destination; + private ReplicationQueue replicationQueue; + + @Before + public void setUp() throws Exception { + project = Project.nameKey("project"); + uri = new URIish("git://host/project.git"); + update = ReplicateRefUpdate.create(project.get(), Set.of(REF), uri, REMOTE); + + waitingTasks = new ArrayList<>(); + taskNamesByReplicateRefUpdate = new HashMap<>(); + + destination = mock(Destination.class); + when(destination.isPushEnabled()).thenReturn(true); + when(destination.getRemoteConfigName()).thenReturn(REMOTE); + when(destination.getTaskNamesByReplicateRefUpdate()).thenReturn(taskNamesByReplicateRefUpdate); + + ReplicationDestinations destinations = mock(ReplicationDestinations.class); + when(destinations.getAll(FilterType.ALL)).thenReturn(List.of(destination)); + when(destinations.getDestinations(any(), any(), any())).thenReturn(List.of(destination)); + + ReplicationTasksStorage tasksStorage = mock(ReplicationTasksStorage.class); + when(tasksStorage.streamWaiting()).thenAnswer(invocation -> List.copyOf(waitingTasks).stream()); + + defaultQueue = new ScheduledThreadPoolExecutor(1); + workQueue = mock(WorkQueue.class); + when(workQueue.getDefaultQueue()).thenReturn(defaultQueue); + + replicationQueue = + new ReplicationQueue( + mock(ReplicationConfig.class), + workQueue, + () -> destinations, + mock(EventDispatcher.class), + mock(ReplicationStateListeners.class), + tasksStorage, + mock(ProjectDeletionState.Factory.class)); + } + + @After + public void tearDown() { + defaultQueue.shutdownNow(); + } + + @Test + public void distributorDoesNotFireTaskPendingOnThisNode() throws Exception { + start(); + waitingTasks.add(update); + taskNamesByReplicateRefUpdate.put(update, "pending push task"); + runDistributor(); + + verify(destination, never()).scheduleFromStorage(any(), any(), any(), any()); + } + + @Test + public void distributorFiresTaskNotPendingOnThisNode() throws Exception { + start(); + waitingTasks.add(update); + runDistributor(); + + verify(destination).scheduleFromStorage(eq(project), eq(update.refs()), eq(uri), any()); + } + + @Test + public void startupFiresTaskPendingOnThisNode() throws Exception { + waitingTasks.add(update); + taskNamesByReplicateRefUpdate.put(update, "pending push task"); + start(); + + verify(destination, never()).getTaskNamesByReplicateRefUpdate(); + verify(destination).scheduleFromStorage(eq(project), eq(update.refs()), eq(uri), any()); + } + + private void start() throws Exception { + replicationQueue.start(); + awaitReplayed(); + } + + private void runDistributor() throws Exception { + replicationQueue.new Distributor(workQueue).run(); + awaitReplayed(); + } + + private void awaitReplayed() throws InterruptedException { + long deadline = System.nanoTime() + REPLAY_TIMEOUT.toNanos(); + while (replicationQueue.isReplaying() && System.nanoTime() < deadline) { + MILLISECONDS.sleep(5); + } + assertThat(replicationQueue.isReplaying()).isFalse(); + } +}
diff --git a/src/test/java/com/googlesource/gerrit/plugins/replication/ReplicationStorageIT.java b/src/test/java/com/googlesource/gerrit/plugins/replication/ReplicationStorageIT.java index 643a781..08a7ecb 100644 --- a/src/test/java/com/googlesource/gerrit/plugins/replication/ReplicationStorageIT.java +++ b/src/test/java/com/googlesource/gerrit/plugins/replication/ReplicationStorageIT.java
@@ -67,6 +67,30 @@ } @Test + public void shouldPersistTaskButNotPushWhenRemoteHasZeroThreads() throws Exception { + String remote = "persistOnly"; + Project.NameKey target = createTestProject(project + "replica"); + setReplicationDestination(remote, "replica", ALL_PROJECTS); + config.setInt("remote", remote, "threads", 0); + config.save(); + reloadConfig(); + + String changeRef = createChange().getPatchSet().refName(); + + assertThat(waitingChangeReplicationTasksForRemote(changeRef, remote).count()).isEqualTo(1); + + Destination destination = + destinationCollection.getAll(FilterType.ALL).stream() + .filter(dest -> remote.equals(dest.getRemoteConfigName())) + .findFirst() + .get(); + + assertThat(destination.isPushEnabled()).isFalse(); + assertThat(destination.getQueue().pending).isEmpty(); + assertThat(isPushCompleted(target, changeRef, TEST_PUSH_TIMEOUT)).isFalse(); + } + + @Test public void shouldCreateOneReplicationTaskWhenSchedulingRepoFullSync() throws Exception { createTestProject(project + "replica"); @@ -76,12 +100,52 @@ plugin .getSysInjector() .getInstance(ReplicationQueue.class) - .scheduleFullSync(project, null, new ReplicationState(NO_OP), false); + .scheduleFullSync( + project, null, PushOne.ALL_REFS, Set.of(), new ReplicationState(NO_OP), false); assertThat(listWaitingReplicationTasks(Pattern.quote(PushOne.ALL_REFS))).hasSize(1); } @Test + public void shouldCreateOneReplicationTaskWhenSchedulingSpecificRefSync() throws Exception { + createTestProject(project + "replica"); + + setReplicationDestination("foo", "replica", ALL_PROJECTS, Integer.MAX_VALUE); + reloadConfig(); + + String specificRef = "refs/heads/master"; + + plugin + .getSysInjector() + .getInstance(ReplicationQueue.class) + .scheduleFullSync(project, null, specificRef, Set.of(), new ReplicationState(NO_OP), false); + + assertThat(listWaitingReplicationTasks(Pattern.quote(specificRef))).hasSize(1); + + tasksStorage + .streamWaiting() + .forEach( + (task) -> { + assertThat(task.refs()).containsExactly(specificRef); + }); + } + + @Test + public void shouldNotCreateReplicationTaskForRefMatchingExcludedRefsPattern() throws Exception { + createTestProject(project + "replica"); + setExcludedRefsPattern("foo", "refs/heads/excluded.*"); + + scheduleFullSync("refs/heads/excluded-branch"); + + assertThat(listWaiting()).isEmpty(); + + String replicatedRef = "refs/heads/replicated-branch"; + scheduleFullSync(replicatedRef); + + assertThat(listWaitingReplicationTasks(Pattern.quote(replicatedRef))).hasSize(1); + } + + @Test public void shouldFirePendingOnlyToIncompleteUri() throws Exception { String suffix1 = "replica1"; String suffix2 = "replica2"; @@ -197,7 +261,8 @@ plugin .getSysInjector() .getInstance(ReplicationQueue.class) - .scheduleFullSync(project, urlMatch, new ReplicationState(NO_OP), false); + .scheduleFullSync( + project, urlMatch, PushOne.ALL_REFS, Set.of(), new ReplicationState(NO_OP), false); assertThat(listWaiting()).hasSize(1); tasksStorage @@ -222,7 +287,8 @@ plugin .getSysInjector() .getInstance(ReplicationQueue.class) - .scheduleFullSync(project, urlMatch, new ReplicationState(NO_OP), false); + .scheduleFullSync( + project, urlMatch, PushOne.ALL_REFS, Set.of(), new ReplicationState(NO_OP), false); assertThat(listWaiting()).hasSize(1); tasksStorage @@ -326,6 +392,20 @@ assertThat(listWaitingReplicationTasks(branchToDelete)).hasSize(1); } + private void setExcludedRefsPattern(String remote, String pattern) throws Exception { + setReplicationDestination(remote, "replica", ALL_PROJECTS, Integer.MAX_VALUE); + config.setString("remote", remote, "excludedRefsPattern", pattern); + config.save(); + reloadConfig(); + } + + private void scheduleFullSync(String ref) { + plugin + .getSysInjector() + .getInstance(ReplicationQueue.class) + .scheduleFullSync(project, null, ref, Set.of(), new ReplicationState(NO_OP), false); + } + private boolean isTaskRescheduled(Queue queue, URIish uri) { PushOne pushOne = queue.pending.get(uri); return pushOne == null ? false : pushOne.isRetrying();
diff --git a/src/test/java/com/googlesource/gerrit/plugins/replication/ReplicationTasksStorageTest.java b/src/test/java/com/googlesource/gerrit/plugins/replication/ReplicationTasksStorageTest.java index 6d334fa..cb1d541 100644 --- a/src/test/java/com/googlesource/gerrit/plugins/replication/ReplicationTasksStorageTest.java +++ b/src/test/java/com/googlesource/gerrit/plugins/replication/ReplicationTasksStorageTest.java
@@ -358,6 +358,23 @@ } @Test + public void recoverAllRecoversOnlyUpdatesMatchingFilter() throws Exception { + ReplicateRefUpdate otherRefUpdate = + ReplicateRefUpdate.create(PROJECT, Set.of(REF), URISH, "otherRemote"); + UriUpdates otherUriUpdates = new TestUriUpdates(otherRefUpdate); + storage.create(REF_UPDATE); + storage.create(otherRefUpdate); + storage.start(uriUpdates); + storage.start(otherUriUpdates); + + storage.recoverAll(r -> REMOTE.equals(r.remote())); + + assertThatStream(storage.streamWaiting()).containsExactly(STORED_REF_UPDATE); + assertThatStream(storage.streamRunning()) + .containsExactly(ReplicateRefUpdate.create(otherRefUpdate, otherRefUpdate.sha1())); + } + + @Test public void canCompleteMultipleRecoveredUpdates() throws Exception { ReplicateRefUpdate updateB = ReplicateRefUpdate.create(