From f74fff66db33398d0536667a85af6a14b5407912 Mon Sep 17 00:00:00 2001 From: Bennett Zhu Date: Tue, 2 Jun 2026 16:00:08 -0400 Subject: [PATCH] Cache continued tasks during queue scheduling Queue maintenance calls QueueTaskDispatcher.canTake for each buildable item and idle executor. ContinuedTask.Scheduler previously scanned every buildable item on each call to find continued tasks, making queue maintenance quadratic in the number of buildables. Maintain a queue-id keyed cache of buildable continued tasks via QueueListener events. The first scheduler call initializes the cache from the current buildable queue, and later queue transitions add or remove entries. Cache entries keep only a weak task reference plus display name and assigned label so stale entries can be dropped during the next scheduler pass. Add coverage for a cached task that stops being continued. Fixes #560 --- .../durabletask/executors/ContinuedTask.java | 101 +++++++++++++++--- .../executors/ContinuedTaskTest.java | 64 +++++++++++ 2 files changed, 151 insertions(+), 14 deletions(-) diff --git a/src/main/java/org/jenkinsci/plugins/durabletask/executors/ContinuedTask.java b/src/main/java/org/jenkinsci/plugins/durabletask/executors/ContinuedTask.java index 722753b3..6b1eb173 100644 --- a/src/main/java/org/jenkinsci/plugins/durabletask/executors/ContinuedTask.java +++ b/src/main/java/org/jenkinsci/plugins/durabletask/executors/ContinuedTask.java @@ -29,7 +29,12 @@ import hudson.model.Node; import hudson.model.Queue; import hudson.model.queue.CauseOfBlockage; +import hudson.model.queue.QueueListener; import hudson.model.queue.QueueTaskDispatcher; +import java.lang.ref.WeakReference; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentMap; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.logging.Logger; import org.kohsuke.accmod.Restricted; import org.kohsuke.accmod.restrictions.NoExternalUse; @@ -62,33 +67,101 @@ private static boolean isContinued(Queue.Task task) { LOGGER.finer(() -> item.task + " is a continued task, so we are not blocking it"); return null; } - for (Queue.BuildableItem other : Queue.getInstance().getBuildableItems()) { - if (isContinued(other.task)) { - Label label = other.task.getAssignedLabel(); - if (label == null || label.matches(node)) { // conservative; might actually go to a different node - LOGGER.fine(() -> "blocking " + item.task + " in favor of " + other.task); - return new HoldOnPlease(other.task); - } else { - LOGGER.finer(() -> other.task + "’s label " + label + " does not match " + node); - } + BuildableContinuedTasks.initialize(); + for (ContinuedItem continued : BuildableContinuedTasks.values()) { + Queue.Task task = continued.task.get(); + if (task == null || !isContinued(task)) { + BuildableContinuedTasks.remove(continued.id); + continue; + } + Label label = continued.label; + if (label == null || label.matches(node)) { // conservative; might actually go to a different node + LOGGER.fine(() -> "blocking " + item.task + " in favor of " + continued.fullDisplayName); + return new HoldOnPlease(continued.fullDisplayName); } else { - LOGGER.finer(() -> other.task + " is not continued, so it would not block " + item.task); + LOGGER.finer(() -> continued.fullDisplayName + "’s label " + label + " does not match " + node); } } LOGGER.finer(() -> "no reason to block " + item.task); return null; } + private static final class BuildableContinuedTasks { + + private static final ConcurrentMap ITEMS = new ConcurrentHashMap<>(); + private static final AtomicBoolean INITIALIZED = new AtomicBoolean(); + + private static synchronized void initialize() { + if (INITIALIZED.compareAndSet(false, true)) { + for (Queue.BuildableItem item : Queue.getInstance().getBuildableItems()) { + addIfContinued(item); + } + } + } + + private static Iterable values() { + return ITEMS.values(); + } + + private static synchronized void addIfContinued(Queue.BuildableItem item) { + if (isContinued(item.task)) { + ITEMS.put(item.getId(), new ContinuedItem(item)); + } + } + + private static synchronized void remove(long id) { + ITEMS.remove(id); + } + + } + + private static final class ContinuedItem { + + private final long id; + private final WeakReference task; + private final String fullDisplayName; + private final Label label; + + ContinuedItem(Queue.BuildableItem item) { + id = item.getId(); + task = new WeakReference<>(item.task); + fullDisplayName = item.task.getFullDisplayName(); + label = item.task.getAssignedLabel(); + } + + } + private static final class HoldOnPlease extends CauseOfBlockage { - private final Queue.Task task; + private final String fullDisplayName; - HoldOnPlease(Queue.Task task) { - this.task = task; + HoldOnPlease(String fullDisplayName) { + this.fullDisplayName = fullDisplayName; } @Override public String getShortDescription() { - return Messages.ContinuedTask__should_be_allowed_to_run_first(task.getFullDisplayName()); + return Messages.ContinuedTask__should_be_allowed_to_run_first(fullDisplayName); + } + + } + + @Restricted(NoExternalUse.class) // implementation + @Extension public static final class Listener extends QueueListener { + + @Override public void onEnterBuildable(Queue.BuildableItem bi) { + BuildableContinuedTasks.addIfContinued(bi); + } + + @Override public void onEnterBlocked(Queue.BlockedItem bi) { + BuildableContinuedTasks.remove(bi.getId()); + } + + @Override public void onLeaveBuildable(Queue.BuildableItem bi) { + BuildableContinuedTasks.remove(bi.getId()); + } + + @Override public void onLeft(Queue.LeftItem li) { + BuildableContinuedTasks.remove(li.getId()); } } diff --git a/src/test/java/org/jenkinsci/plugins/durabletask/executors/ContinuedTaskTest.java b/src/test/java/org/jenkinsci/plugins/durabletask/executors/ContinuedTaskTest.java index 8296b4dd..5e202360 100644 --- a/src/test/java/org/jenkinsci/plugins/durabletask/executors/ContinuedTaskTest.java +++ b/src/test/java/org/jenkinsci/plugins/durabletask/executors/ContinuedTaskTest.java @@ -30,6 +30,9 @@ import hudson.model.FreeStyleBuild; import hudson.model.FreeStyleProject; import hudson.model.Label; +import hudson.model.Node; +import hudson.model.Queue; +import hudson.model.queue.CauseOfBlockage; import hudson.model.queue.QueueTaskFuture; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; @@ -39,10 +42,14 @@ import org.jvnet.hudson.test.junit.jupiter.WithJenkins; import java.io.IOException; +import java.util.Calendar; +import java.util.Collections; import java.util.concurrent.atomic.AtomicInteger; import java.util.logging.Level; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNull; @WithJenkins class ContinuedTaskTest { @@ -82,6 +89,28 @@ public boolean perform(AbstractBuild build, Launcher launcher, BuildListen assertEquals(1, cntB.get()); } + @Test + void cacheDropsTasksThatAreNoLongerContinued() throws Exception { + Label label = Label.get("continued-cache"); + Node node = j.createSlave(label); + ContinuedTask.Scheduler.Listener listener = new ContinuedTask.Scheduler.Listener(); + ContinuedTask.Scheduler scheduler = new ContinuedTask.Scheduler(); + MutableTestTask continued = new MutableTestTask(new AtomicInteger(), true, label); + Queue.BuildableItem continuedItem = buildableItem(continued); + Queue.BuildableItem regularItem = buildableItem(new LabelledTask(new AtomicInteger(), label)); + + listener.onEnterBuildable(continuedItem); + CauseOfBlockage blockage = scheduler.canTake(node, regularItem); + assertNotNull(blockage); + + continued.continued = false; + assertNull(scheduler.canTake(node, regularItem)); + } + + private static Queue.BuildableItem buildableItem(Queue.Task task) { + return new Queue.BuildableItem(new Queue.WaitingItem(Calendar.getInstance(), task, Collections.emptyList())); + } + private static final class TestTask extends MockTask implements ContinuedTask { private final boolean continued; @@ -106,4 +135,39 @@ public String toString() { } } + private static final class MutableTestTask extends MockTask implements ContinuedTask { + private boolean continued; + private final Label label; + + MutableTestTask(AtomicInteger cnt, boolean continued, Label label) { + super(cnt); + this.continued = continued; + this.label = label; + } + + @Override + public boolean isContinued() { + return continued; + } + + @Override + public Label getAssignedLabel() { + return label; + } + } + + private static final class LabelledTask extends MockTask { + private final Label label; + + LabelledTask(AtomicInteger cnt, Label label) { + super(cnt); + this.label = label; + } + + @Override + public Label getAssignedLabel() { + return label; + } + } + }