Introduce minStartForQueue (capacity reservation) This change introduces the ability to reserve a minimum number of threads on a work queue for a specific project, ensuring they can start even under heavy load. Change-Id: Idbc863a0f6e453fabb4c11b00fbf6bb453fc865f
diff --git a/src/main/java/com/googlesource/gerrit/plugins/quota/MinStartForQueueQuota.java b/src/main/java/com/googlesource/gerrit/plugins/quota/MinStartForQueueQuota.java new file mode 100644 index 0000000..33b9a3f --- /dev/null +++ b/src/main/java/com/googlesource/gerrit/plugins/quota/MinStartForQueueQuota.java
@@ -0,0 +1,54 @@ +// Copyright (C) 2025 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.quota; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.util.Optional; +import java.util.regex.Matcher; +import java.util.regex.Pattern; + +public class MinStartForQueueQuota { + public static final Logger log = LoggerFactory.getLogger(MinStartForQueueQuota.class); + // 10 SSH-Interactive-Worker + public static final Pattern CONFIG_PATTERN = Pattern.compile("(\\d+)\\s+(.+)"); + + public static Optional<TaskQuota> build(QuotaSection qs, String cfg) { + Matcher matcher = CONFIG_PATTERN.matcher(cfg); + + if (qs instanceof GlobalQuotaSection || qs.isFallbackQuota()) { + log.warn("minStartForQueueQuota is not applicable in global and fallback quota sections"); + return Optional.empty(); + } + + if (matcher.matches()) { + int limit = Integer.parseInt(matcher.group(1)); + String queue = matcher.group(2); + QueueManager.registerReservation( + queue, + new QueueManager.Reservation( + limit, + task -> { + return task.getQueueName().equalsIgnoreCase(queue) + && TaskQuotas.estimateProject(task).map(qs::matches).orElse(false); + })); + } else { + log.error("Invalid configuration entry [{}]", cfg); + } + + return Optional.empty(); + } +}
diff --git a/src/main/java/com/googlesource/gerrit/plugins/quota/NamespacedQuotaSection.java b/src/main/java/com/googlesource/gerrit/plugins/quota/NamespacedQuotaSection.java index fd7f729..ad8a90c 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/quota/NamespacedQuotaSection.java +++ b/src/main/java/com/googlesource/gerrit/plugins/quota/NamespacedQuotaSection.java
@@ -17,12 +17,17 @@ import com.google.gerrit.entities.Project; import org.eclipse.jgit.lib.Config; -public record NamespacedQuotaSection(Config cfg, String namespace, String resolvedNamespace) +public record NamespacedQuotaSection( + Config cfg, String namespace, String resolvedNamespace, boolean isFallBack) implements QuotaSection { public static final String QUOTA = "quota"; public NamespacedQuotaSection(Config cfg, String namespace) { - this(cfg, namespace, namespace); + this(cfg, namespace, false); + } + + public NamespacedQuotaSection(Config cfg, String namespace, boolean isFallBack) { + this(cfg, namespace, namespace, isFallBack); } public String getNamespace() { @@ -42,4 +47,9 @@ public String subSection() { return namespace(); } + + @Override + public boolean isFallbackQuota() { + return isFallBack; + } }
diff --git a/src/main/java/com/googlesource/gerrit/plugins/quota/QueueManager.java b/src/main/java/com/googlesource/gerrit/plugins/quota/QueueManager.java index c762b8f..4b0010d 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/quota/QueueManager.java +++ b/src/main/java/com/googlesource/gerrit/plugins/quota/QueueManager.java
@@ -15,32 +15,95 @@ package com.googlesource.gerrit.plugins.quota; import com.google.gerrit.server.git.WorkQueue; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.util.ArrayList; +import java.util.Arrays; +import java.util.List; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; +import java.util.function.Predicate; import java.util.*; -import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.atomic.AtomicBoolean; public class QueueManager { - public record QueueInfo(int maxThreads, Set<Integer> runningTasks) { + public static final Logger log = LoggerFactory.getLogger(QueueManager.class); + + public static class QueueInfo { + public final int maxThreads; + public int spareThreads; + public final Map<Integer, WorkQueue.Task<?>> runningTaskById; + public final List<Reservation> reservations; + public QueueInfo(int maxThreads) { - this(maxThreads, new HashSet<>()); + this.maxThreads = maxThreads; + this.spareThreads = maxThreads; + this.runningTaskById = new HashMap<>(); + this.reservations = new ArrayList<>(); } - public boolean run(int taskId) { - if (runningTasks().size() < maxThreads()) { - runningTasks().add(taskId); + public boolean run(WorkQueue.Task<?> task) { + if (runningTaskById.size() >= maxThreads) { + return false; + } + + if (runningTaskById.put(task.getTaskId(), task) != null) { return true; } - return false; + + if (!reservations.isEmpty() && !canAllocate()) { + runningTaskById.remove((task.getTaskId())); + return false; + } + + return true; } - public void complete(int taskId) { - runningTasks().remove(taskId); + public void complete(WorkQueue.Task<?> task) { + runningTaskById.remove(task.getTaskId()); } public boolean ensureIdle(int threads) { - return maxThreads() - runningTasks().size() >= threads; + return maxThreads - runningTaskById.size() >= threads; + } + + public void addReservation(Reservation incomingReservation) { + reservations.add(incomingReservation); + spareThreads -= incomingReservation.reservedCapacity(); + } + + public boolean canAllocate() { + int spareAllocations = 0; + Map<Reservation, Integer> allocationsByReservation = new HashMap<>(); + + for (WorkQueue.Task<?> runningTask : runningTaskById.values()) { + boolean allocatedToReservation = false; + for (Reservation reservation : reservations) { + if (reservation.matches(runningTask)) { + int currentAllocation = allocationsByReservation.getOrDefault(reservation, 0); + if (currentAllocation < reservation.reservedCapacity()) { + allocationsByReservation.put(reservation, currentAllocation + 1); + allocatedToReservation = true; + break; + } + } + } + + if (!allocatedToReservation) { + spareAllocations++; + } + } + + return spareAllocations <= spareThreads; + } + } + + public record Reservation(int reservedCapacity, Predicate<WorkQueue.Task<?>> taskMatcher) { + public boolean matches(WorkQueue.Task<?> task) { + return taskMatcher.test(task); } } @@ -79,6 +142,35 @@ infoByQueue.put(q, new QueueInfo(c)); } + public static void registerReservation(String qName, Reservation reservation) { + Queue q = Queue.fromKey(qName); + if (q == Queue.UNKNOWN) { + return; + } + + QueueInfo queueInfo = infoByQueue.get(q); + int capacityToReserve = queueInfo.spareThreads - 1; + if (capacityToReserve < 1) { + log.error( + "Cannot enforce reservation for queue '{}' Requested: {} threads. No threads reserved.", + qName, + reservation.reservedCapacity()); + return; + } + + if (reservation.reservedCapacity() > capacityToReserve) { + log.warn( + "Partial reservation enforced for queue '{}'. Requested: {}, Actual reserved: {}", + qName, + reservation.reservedCapacity(), + capacityToReserve); + queueInfo.addReservation(new Reservation(capacityToReserve, reservation.taskMatcher)); + return; + } + + queueInfo.addReservation(reservation); + } + public static boolean acquire(WorkQueue.Task<?> task) { Queue q = Queue.fromKey(task.getQueueName()); if (q == Queue.UNKNOWN) { @@ -89,7 +181,7 @@ infoByQueue.computeIfPresent( q, (queue, info) -> { - acquired.setPlain(info.run(task.getTaskId())); + acquired.setPlain(info.run(task)); return info; }); @@ -101,7 +193,7 @@ infoByQueue.computeIfPresent( q, (queue, info) -> { - info.complete(task.getTaskId()); + info.complete(task); return info; }); }
diff --git a/src/main/java/com/googlesource/gerrit/plugins/quota/QuotaFinder.java b/src/main/java/com/googlesource/gerrit/plugins/quota/QuotaFinder.java index 9584461..4e3e251 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/quota/QuotaFinder.java +++ b/src/main/java/com/googlesource/gerrit/plugins/quota/QuotaFinder.java
@@ -44,7 +44,7 @@ String prefix = n.substring(0, n.length() - 3); Matcher m = Pattern.compile("^" + prefix + "([^/]+)/.*$").matcher(p); if (m.matches()) { - return new NamespacedQuotaSection(cfg, n, prefix + m.group(1) + "/*"); + return new NamespacedQuotaSection(cfg, n, prefix + m.group(1) + "/*", false); } } else if (n.endsWith("/*")) { if (p.startsWith(n.substring(0, n.length() - 1))) { @@ -70,7 +70,7 @@ } public QuotaSection getFallbackNamespacedQuota(Config cfg) { - return new NamespacedQuotaSection(cfg, "*"); + return new NamespacedQuotaSection(cfg, "*", true); } public List<NamespacedQuotaSection> getQuotaNamespaces(Config cfg) {
diff --git a/src/main/java/com/googlesource/gerrit/plugins/quota/QuotaSection.java b/src/main/java/com/googlesource/gerrit/plugins/quota/QuotaSection.java index 0a9d0fd..8de2ff8 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/quota/QuotaSection.java +++ b/src/main/java/com/googlesource/gerrit/plugins/quota/QuotaSection.java
@@ -55,11 +55,15 @@ .flatMap( type -> Arrays.stream(cfg().getStringList(section(), subSection(), type.key)) - .map(type.processor) + .map(cfg -> type.processor.apply(this, cfg)) .flatMap(Optional::stream)) .toList(); } + default boolean isFallbackQuota() { + return false; + } + Config cfg(); String section();
diff --git a/src/main/java/com/googlesource/gerrit/plugins/quota/SoftMaxPerUserForQueue.java b/src/main/java/com/googlesource/gerrit/plugins/quota/SoftMaxPerUserForQueue.java index 3442ec6..de7c630 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/quota/SoftMaxPerUserForQueue.java +++ b/src/main/java/com/googlesource/gerrit/plugins/quota/SoftMaxPerUserForQueue.java
@@ -75,7 +75,7 @@ })); } - public static Optional<TaskQuota> build(String cfg) { + public static Optional<TaskQuota> build(QuotaSection qs, String cfg) { Matcher matcher = CONFIG_PATTERN.matcher(cfg); return matcher.find() ? Optional.of(
diff --git a/src/main/java/com/googlesource/gerrit/plugins/quota/TaskQuotaForTaskForQueue.java b/src/main/java/com/googlesource/gerrit/plugins/quota/TaskQuotaForTaskForQueue.java index f8ae3cc..ac57a00 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/quota/TaskQuotaForTaskForQueue.java +++ b/src/main/java/com/googlesource/gerrit/plugins/quota/TaskQuotaForTaskForQueue.java
@@ -38,7 +38,7 @@ return super.isApplicable(task) && task.getQueueName().equals(queueName); } - public static Optional<TaskQuota> build(String cfg) { + public static Optional<TaskQuota> build(QuotaSection qs, String cfg) { Matcher matcher = CONFIG_PATTERN.matcher(cfg); if (matcher.matches()) { return Optional.of(
diff --git a/src/main/java/com/googlesource/gerrit/plugins/quota/TaskQuotaForTaskForQueueForUser.java b/src/main/java/com/googlesource/gerrit/plugins/quota/TaskQuotaForTaskForQueueForUser.java index 7582619..7afa870 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/quota/TaskQuotaForTaskForQueueForUser.java +++ b/src/main/java/com/googlesource/gerrit/plugins/quota/TaskQuotaForTaskForQueueForUser.java
@@ -44,7 +44,7 @@ return taskUser.find() && user.equals(taskUser.group(1)) && super.isApplicable(task); } - public static Optional<TaskQuota> build(String config) { + public static Optional<TaskQuota> build(QuotaSection qs, String config) { Matcher matcher = CONFIG_PATTERN.matcher(config); if (matcher.matches()) { return Optional.of(
diff --git a/src/main/java/com/googlesource/gerrit/plugins/quota/TaskQuotaKeys.java b/src/main/java/com/googlesource/gerrit/plugins/quota/TaskQuotaKeys.java index 1d840e3..ff0a4cb 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/quota/TaskQuotaKeys.java +++ b/src/main/java/com/googlesource/gerrit/plugins/quota/TaskQuotaKeys.java
@@ -15,10 +15,11 @@ package com.googlesource.gerrit.plugins.quota; import java.util.Optional; -import java.util.function.Function; +import java.util.function.BiFunction; public enum TaskQuotaKeys { MAX_START_FOR_TASK_FOR_QUEUE("maxStartForTaskForQueue", TaskQuotaForTaskForQueue::build), + MIN_START_FOR_TASK_FOR_QUEUE("minStartForQueue", MinStartForQueueQuota::build), MAX_START_FOR_TASK_FOR_USER_FOR_QUEUE( "maxStartForTaskForUserForQueue", TaskQuotaForTaskForQueueForUser::build), MAX_START_PER_USER_FOR_TASK_FOR_QUEUE( @@ -26,9 +27,9 @@ SOFT_MAX_START_FOR_QUEUE_PER_USER("softMaxStartPerUserForQueue", SoftMaxPerUserForQueue::build); public final String key; - public final Function<String, Optional<TaskQuota>> processor; + public final BiFunction<QuotaSection, String, Optional<TaskQuota>> processor; - TaskQuotaKeys(String key, Function<String, Optional<TaskQuota>> processor) { + TaskQuotaKeys(String key, BiFunction<QuotaSection, String, Optional<TaskQuota>> processor) { this.key = key; this.processor = processor; }
diff --git a/src/main/java/com/googlesource/gerrit/plugins/quota/TaskQuotaPerUserForTaskForQueue.java b/src/main/java/com/googlesource/gerrit/plugins/quota/TaskQuotaPerUserForTaskForQueue.java index 6d7d076..a7299a2 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/quota/TaskQuotaPerUserForTaskForQueue.java +++ b/src/main/java/com/googlesource/gerrit/plugins/quota/TaskQuotaPerUserForTaskForQueue.java
@@ -39,7 +39,7 @@ perUserTaskQuota.release(task); } - public static Optional<TaskQuota> build(String cfg) { + public static Optional<TaskQuota> build(QuotaSection qs, String cfg) { Matcher matcher = CONFIG_PATTERN.matcher(cfg); if (matcher.matches()) { return Optional.of(
diff --git a/src/main/java/com/googlesource/gerrit/plugins/quota/TaskQuotas.java b/src/main/java/com/googlesource/gerrit/plugins/quota/TaskQuotas.java index 3817aba..0c1b8b4 100644 --- a/src/main/java/com/googlesource/gerrit/plugins/quota/TaskQuotas.java +++ b/src/main/java/com/googlesource/gerrit/plugins/quota/TaskQuotas.java
@@ -37,7 +37,7 @@ private final Map<Integer, List<TaskQuota>> quotasByTask = new ConcurrentHashMap<>(); private final Map<QuotaSection, List<TaskQuota>> quotasByNamespace = new HashMap<>(); private final List<TaskQuota> globalQuotas = new ArrayList<>(); - private final Pattern PROJECT_PATTERN = Pattern.compile("\\s+/?(.*)\\s+(\\(\\S+\\))$"); + private static final Pattern PROJECT_PATTERN = Pattern.compile("\\s+/?(.*)\\s+(\\(\\S+\\))$"); private final Config quotaConfig; @Inject @@ -125,7 +125,7 @@ .ifPresent(quotas -> quotas.forEach(q -> q.onStop(task))); } - private Optional<Project.NameKey> estimateProject(WorkQueue.Task<?> task) { + public static Optional<Project.NameKey> estimateProject(WorkQueue.Task<?> task) { Matcher matcher = PROJECT_PATTERN.matcher(task.toString()); return matcher.find() ? Optional.of(Project.NameKey.parse(matcher.group(1))) : Optional.empty();
diff --git a/src/main/resources/Documentation/config.md b/src/main/resources/Documentation/config.md index 6ae88af..f2ceb15 100644 --- a/src/main/resources/Documentation/config.md +++ b/src/main/resources/Documentation/config.md
@@ -301,6 +301,20 @@ maxStartPerUserForTaskForQueue = 20 uploadpack SSH-Interactive-Worker ``` +We can also reserve a certain amount of the queue's capacity for specific project +namespaces using `minStartForQueue`. + +``` + [quota "android"] + minStartForQueue = 5 SSH-Interactive-Worker +``` + +The configuration reserves 5 threads from the interactive queue exclusively for tasks +related to the android project. One must ensure the total number of reserved threads does +not exceed the queue's capacity. If this limit is surpassed, some of the configured +minStarts will not be enforced and will be logged. Additionally, note that +`minStartForQueue` cannot be defined in the global or fallback quota sections. + Currently supported tasks: * `uploadpack`: Maps directly to git-upload-pack operations (used during Git
diff --git a/src/test/java/com/googlesource/gerrit/plugins/quota/QueueManagerTest.java b/src/test/java/com/googlesource/gerrit/plugins/quota/QueueManagerTest.java index 92930ec..1740561 100644 --- a/src/test/java/com/googlesource/gerrit/plugins/quota/QueueManagerTest.java +++ b/src/test/java/com/googlesource/gerrit/plugins/quota/QueueManagerTest.java
@@ -52,10 +52,11 @@ QueueInfo info = QueueManager.infoByQueue.get(TEST_QUEUE); assertEquals( - "QueueInfo should store the correct maxThreads capacity.", MAX_CAPACITY, info.maxThreads()); - assertNotNull("QueueInfo's runningTasks set should not be null.", info.runningTasks()); + "QueueInfo should store the correct maxThreads capacity.", MAX_CAPACITY, info.maxThreads); + assertNotNull("QueueInfo's runningTasksById set should not be null.", info.runningTaskById); assertTrue( - "QueueInfo's runningTasks set should be empty initially.", info.runningTasks().isEmpty()); + "QueueInfo's runningTasksById set should be empty initially.", + info.runningTaskById.isEmpty()); } @Test @@ -77,8 +78,8 @@ QueueInfo info = QueueManager.infoByQueue.get(TEST_QUEUE); - assertEquals("Running tasks count should be 2.", 2, info.runningTasks().size()); - assertTrue("Task 1 ID should be registered.", info.runningTasks().contains(101)); + assertEquals("Running tasks count should be 2.", 2, info.runningTaskById.size()); + assertTrue("Task 1 ID should be registered.", info.runningTaskById.containsKey(101)); } @Test @@ -92,8 +93,8 @@ QueueInfo info = QueueManager.infoByQueue.get(TEST_QUEUE); - assertEquals("Running tasks count should remain 1.", 1, info.runningTasks().size()); - assertFalse("Task 2 ID should not be registered.", info.runningTasks().contains(102)); + assertEquals("Running tasks count should remain 1.", 1, info.runningTaskById.size()); + assertFalse("Task 2 ID should not be registered.", info.runningTaskById.containsKey(102)); } @Test @@ -103,11 +104,11 @@ QueueManager.acquire(task); QueueInfo info = QueueManager.infoByQueue.get(TEST_QUEUE); - assertTrue("Task should be running before release.", info.runningTasks().contains(201)); + assertTrue("Task should be running before release.", info.runningTaskById.containsKey(201)); QueueManager.release(task); - assertFalse("Task should be removed after release.", info.runningTasks().contains(201)); - assertTrue("Running tasks set should be empty.", info.runningTasks().isEmpty()); + assertFalse("Task should be removed after release.", info.runningTaskById.containsKey(201)); + assertTrue("Running tasks set should be empty.", info.runningTaskById.isEmpty()); } @Test @@ -126,10 +127,11 @@ QueueManager.acquire(runningTask); QueueInfo info = QueueManager.infoByQueue.get(TEST_QUEUE); - int initialSize = info.runningTasks().size(); + int initialSize = info.runningTaskById.size(); QueueManager.release(releasedTask); - assertEquals("Running tasks count should not change.", initialSize, info.runningTasks().size()); + assertEquals( + "Running tasks count should not change.", initialSize, info.runningTaskById.size()); } @Test @@ -201,7 +203,7 @@ assertEquals( "The number of running tasks must not exceed the max capacity due to race conditions.", MAX_CAPACITY, - info.runningTasks().size()); + info.runningTaskById.size()); assertEquals( "The correct number of unique task IDs should be acquired.", MAX_CAPACITY, acquired.size());