Prune non-waiting tasks from queue during distribution Make the replication distributor also prune the tasks which have been removed from the persistent waiting directory but that are still in the queue. Since these tasks are not in the waiting directory, they will not run once they could run anyway. Cleaning these tasks up gives a more accurate picture of the replication work still needing to be done. Presumably these tasks have been removed externally by another actor (likely because they completed on another master in the cluster). This keeps the scheduled and waiting lists approximately synchronized across all the nodes in a multi-master cluster. Change-Id: If99ed75041fa07632fdc0f23473dad4d8942a3e1
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 5d555f1..1524c2a 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/replication/Destination.java +++ b/src/main/java/com/googlesource/gerrit/plugins/replication/Destination.java
@@ -64,6 +64,7 @@ import java.io.UnsupportedEncodingException; import java.net.URLEncoder; import java.util.HashMap; +import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.Set; @@ -578,6 +579,17 @@ } } + public Set<String> getPrunableTaskNames() { + Set<String> names = new HashSet<>(); + for (PushOne push : pending.values()) { + if (!replicationTasksStorage.get().isWaiting(push)) { + repLog.debug("No longer isWaiting, can prune " + push.getURI()); + names.add(push.toString()); + } + } + return names; + } + boolean wouldPushProject(Project.NameKey project) { if (!shouldReplicate(project)) { return false;
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 148d5ba..0dcfc95 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/replication/ReplicationQueue.java +++ b/src/main/java/com/googlesource/gerrit/plugins/replication/ReplicationQueue.java
@@ -180,6 +180,30 @@ } } + private void pruneCompleted() { + // Queue tasks have wrappers around them so workQueue.getTasks() does not return the PushOnes. + // We also cannot access them by taskId since PushOnes don't have a taskId, they do have + // and Id, but it not the id assigned to the task in the queues. The tasks in the queue + // do use the same name as returned by toString() though, so that be used to correlate + // PushOnes with queue tasks despite their wrappers. + Set<String> prunableTaskNames = new HashSet<>(); + for (Destination destination : destinations.get().getAll(FilterType.ALL)) { + prunableTaskNames.addAll(destination.getPrunableTaskNames()); + } + + for (WorkQueue.Task<?> task : workQueue.getTasks()) { + WorkQueue.Task.State state = task.getState(); + if (state == WorkQueue.Task.State.SLEEPING || state == WorkQueue.Task.State.READY) { + if (task instanceof WorkQueue.ProjectTask) { + if (prunableTaskNames.contains(task.toString())) { + repLog.debug("Pruning externally completed task:" + task); + task.cancel(false); + } + } + } + } + } + @Override public void onProjectDeleted(ProjectDeletedListener.Event event) { Project.NameKey p = Project.nameKey(event.getProjectName()); @@ -239,6 +263,7 @@ } try { firePendingEvents(); + pruneCompleted(); } catch (Exception e) { repLog.error("error distributing tasks", e); }
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 7ba6c8a..992cf5b 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/replication/ReplicationTasksStorage.java +++ b/src/main/java/com/googlesource/gerrit/plugins/replication/ReplicationTasksStorage.java
@@ -178,6 +178,10 @@ } } + public boolean isWaiting(PushOne push) { + return push.getRefs().stream().map(ref -> new Task(push, ref)).anyMatch(Task::isWaiting); + } + public void finish(PushOne push) { UriLock lock = new UriLock(push); for (ReplicateRefUpdate r : list(lock.runningDir)) { @@ -276,6 +280,10 @@ this(new UriLock(r), r.ref); } + public Task(PushOne push, String ref) { + this(new UriLock(push), ref); + } + public Task(UriLock lock, String ref) { update = new ReplicateRefUpdate(lock.update, ref); json = GSON.toJson(update) + "\n"; @@ -311,6 +319,10 @@ rename(running, waiting); } + public boolean isWaiting() { + return Files.exists(waiting); + } + public void finish() { if (disableDeleteForTesting) { logger.atFine().log("DELETE %s %s DISABLED", running, updateLog());