From a55ccc6f17691cb865c9bc207bbb0d08db9e5d59 Mon Sep 17 00:00:00 2001 From: Hilbrand Bouwkamp Date: Tue, 15 Sep 2026 17:41:28 +0200 Subject: [PATCH] If no workers are present but there are tasks on the queue report load as 100% to trigger scaling from zero. --- .../nl/aerius/taskmanager/TaskManager.java | 2 +- .../taskmanager/domain/QueueEmptyCheck.java | 29 +++++++++++++++++++ .../TaskManagerUsageMetricsProvider.java | 10 +++++-- .../taskmanager/scheduler/TaskScheduler.java | 3 +- .../priorityqueue/GroupedPriorityQueue.java | 2 +- .../priorityqueue/PriorityTaskScheduler.java | 5 ++++ .../aerius/taskmanager/MockTaskScheduler.java | 5 ++++ .../TaskManagerUsageMetricsProviderTest.java | 13 ++++++++- .../GroupedPriorityQueueTest.java | 4 +++ 9 files changed, 67 insertions(+), 6 deletions(-) create mode 100644 source/taskmanager/src/main/java/nl/aerius/taskmanager/domain/QueueEmptyCheck.java diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/TaskManager.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/TaskManager.java index 55c6ed1..1e8ab79 100644 --- a/source/taskmanager/src/main/java/nl/aerius/taskmanager/TaskManager.java +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/TaskManager.java @@ -148,7 +148,7 @@ public TaskScheduleBucket(final QueueConfig queueConfig) throws InterruptedExcep taskScheduler = schedulerFactory.createScheduler(queueConfig); workerProducer = factory.createWorkerProducer(queueConfig); final WorkerPool workerPool = new WorkerPool(workerQueueName, workerProducer, taskScheduler); - final TaskManagerUsageMetricsProvider taskManagerUsageMetrics = new TaskManagerUsageMetricsProvider(workerQueueName); + final TaskManagerUsageMetricsProvider taskManagerUsageMetrics = new TaskManagerUsageMetricsProvider(workerQueueName, taskScheduler); final TaskManagerMetricsRegister taskManagerMetricsRegister = new TaskManagerMetricsRegister(taskManagerUsageMetrics, startupGuard); final PerformanceMetricsReporter reporter = new PerformanceMetricsReporter(scheduledExecutorService, queueConfig.queueName(), OpenTelemetryMetrics.METER); diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/domain/QueueEmptyCheck.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/domain/QueueEmptyCheck.java new file mode 100644 index 0000000..ceb5cb8 --- /dev/null +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/domain/QueueEmptyCheck.java @@ -0,0 +1,29 @@ +/* + * Copyright (c) Contributors to the project + * + * This program is free software: you can redistribute it and/or modify + * it under the terms of the GNU Affero General Public License as published by + * the Free Software Foundation, either version 3 of the License, or + * (at your option) any later version. + * + * This program is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU Affero General Public License for more details. + * + * You should have received a copy of the GNU Affero General Public License + * along with this program. If not, see http://www.gnu.org/licenses/. + */ +package nl.aerius.taskmanager.domain; + +/** + * Interface to check if a queue is empty. + */ +public interface QueueEmptyCheck { + + /** + * @return Return true if the queue is empty + */ + boolean isQueueEmpty(); + +} diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/TaskManagerUsageMetricsProvider.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/TaskManagerUsageMetricsProvider.java index fc856f8..00bc8f8 100644 --- a/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/TaskManagerUsageMetricsProvider.java +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/TaskManagerUsageMetricsProvider.java @@ -18,6 +18,8 @@ import java.util.function.ToDoubleBiFunction; +import nl.aerius.taskmanager.domain.QueueEmptyCheck; + /** * The {@link UsageMetricsProvider} for weighted load. limit and usage metrics for worker queues. * The values are calculated averages over the time between the last measurement point and the moment the metric value is requested. @@ -33,9 +35,13 @@ public class TaskManagerUsageMetricsProvider implements UsageMetricsProvider { private final LoadMetric used; private final LoadMetric free; - public TaskManagerUsageMetricsProvider(final String workerQueueName) { + public TaskManagerUsageMetricsProvider(final String workerQueueName, final QueueEmptyCheck queueEmptyCheck) { this.workerQueueName = workerQueueName; - load = new LoadMetric((numberOfWorkers, usedWorkers) -> (numberOfWorkers > 0 ? (usedWorkers / (double) numberOfWorkers) : 0), LOAD_SUM_FUNCTION); + final ToDoubleBiFunction loadCountFuction = (numberOfWorkers, usedWorkers) -> numberOfWorkers == 0 + ? (queueEmptyCheck.isQueueEmpty() ? 0 : 1.0) + : (usedWorkers / (double) numberOfWorkers); + + load = new LoadMetric(loadCountFuction, LOAD_SUM_FUNCTION); limit = new LoadMetric((numberOfWorkers, usedWorkers) -> numberOfWorkers, COUNT_SUM_FUNCTION); used = new LoadMetric((numberOfWorkers, usedWorkers) -> usedWorkers, COUNT_SUM_FUNCTION); free = new LoadMetric((numberOfWorkers, usedWorkers) -> numberOfWorkers - usedWorkers, COUNT_SUM_FUNCTION); diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/scheduler/TaskScheduler.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/scheduler/TaskScheduler.java index 8495f1f..e419d53 100644 --- a/source/taskmanager/src/main/java/nl/aerius/taskmanager/scheduler/TaskScheduler.java +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/scheduler/TaskScheduler.java @@ -17,6 +17,7 @@ package nl.aerius.taskmanager.scheduler; import nl.aerius.taskmanager.domain.QueueConfig; +import nl.aerius.taskmanager.domain.QueueEmptyCheck; import nl.aerius.taskmanager.domain.QueueWatchDogListener; import nl.aerius.taskmanager.domain.Task; import nl.aerius.taskmanager.domain.TaskQueue; @@ -27,7 +28,7 @@ * Interface for the scheduling algorithm. The implementation should maintain an internal list of all tasks added and return the task to be processed * in {@link #getNextTask()} based on whatever priority algorithm the scheduler implements. */ -public interface TaskScheduler extends WorkerUpdateHandler, QueueWatchDogListener { +public interface TaskScheduler extends WorkerUpdateHandler, QueueWatchDogListener, QueueEmptyCheck { /** * Adds a Task to the scheduler to being processed. The scheduler will return this task in {@link #getNextTask()} diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/scheduler/priorityqueue/GroupedPriorityQueue.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/scheduler/priorityqueue/GroupedPriorityQueue.java index 2986873..b390a0e 100644 --- a/source/taskmanager/src/main/java/nl/aerius/taskmanager/scheduler/priorityqueue/GroupedPriorityQueue.java +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/scheduler/priorityqueue/GroupedPriorityQueue.java @@ -111,7 +111,7 @@ public boolean containsAll(final Collection arg0) { @Override public boolean isEmpty() { - throw new UnsupportedOperationException("Not implemented"); + return queue.isEmpty() && groupedQueue.isEmpty(); } @Override diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/scheduler/priorityqueue/PriorityTaskScheduler.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/scheduler/priorityqueue/PriorityTaskScheduler.java index cf78d92..9b91620 100644 --- a/source/taskmanager/src/main/java/nl/aerius/taskmanager/scheduler/priorityqueue/PriorityTaskScheduler.java +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/scheduler/priorityqueue/PriorityTaskScheduler.java @@ -117,6 +117,11 @@ public Task getNextTask() throws InterruptedException { return task; } + @Override + public boolean isQueueEmpty() { + return queue.isEmpty(); + } + private void obtainTask() { final Task task = queue.poll(); diff --git a/source/taskmanager/src/test/java/nl/aerius/taskmanager/MockTaskScheduler.java b/source/taskmanager/src/test/java/nl/aerius/taskmanager/MockTaskScheduler.java index ce1b9b8..76fb2ff 100644 --- a/source/taskmanager/src/test/java/nl/aerius/taskmanager/MockTaskScheduler.java +++ b/source/taskmanager/src/test/java/nl/aerius/taskmanager/MockTaskScheduler.java @@ -49,6 +49,11 @@ public Task getNextTask() throws InterruptedException { return tasks.take(); } + @Override + public boolean isQueueEmpty() { + return tasks.isEmpty(); + } + @Override public void updateQueue(final PriorityTaskQueue queue) { // Not used diff --git a/source/taskmanager/src/test/java/nl/aerius/taskmanager/metrics/TaskManagerUsageMetricsProviderTest.java b/source/taskmanager/src/test/java/nl/aerius/taskmanager/metrics/TaskManagerUsageMetricsProviderTest.java index e8f56ee..cff8fd0 100644 --- a/source/taskmanager/src/test/java/nl/aerius/taskmanager/metrics/TaskManagerUsageMetricsProviderTest.java +++ b/source/taskmanager/src/test/java/nl/aerius/taskmanager/metrics/TaskManagerUsageMetricsProviderTest.java @@ -23,16 +23,21 @@ import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import nl.aerius.taskmanager.domain.QueueEmptyCheck; + /** * Test class for {@link TaskManagerUsageMetricsProvider}. */ class TaskManagerUsageMetricsProviderTest { + private boolean queuEmpty; + private final QueueEmptyCheck emptyChecker = () -> queuEmpty; private TaskManagerUsageMetricsProvider provider; @BeforeEach void beforeEach() { - provider = new TaskManagerUsageMetricsProvider("TEST"); + provider = new TaskManagerUsageMetricsProvider("TEST", emptyChecker); + queuEmpty = true; } @Test @@ -55,6 +60,12 @@ void testNumberOfFreeWorkers() throws InterruptedException { assertMetricAndZero(3, 10, provider::getNumberOfFreeWorkers, 7, "Expected 7 worker to be free.", 10); } + @Test + void testZeroWorkersNoEmptyQueue() throws InterruptedException { + queuEmpty = false; + assertMetric(10, 0, provider::getLoad, 100, "Expected a load of 100%", 100); + } + private void assertMetricAndZero(final int numberOfUsed, final int numberOfWorkers, final DoubleSupplier supplier, final int expected, final String description, final int expectedAfterReset) throws InterruptedException { assertMetric(numberOfUsed, numberOfWorkers, supplier, expected, description, expectedAfterReset); diff --git a/source/taskmanager/src/test/java/nl/aerius/taskmanager/scheduler/priorityqueue/GroupedPriorityQueueTest.java b/source/taskmanager/src/test/java/nl/aerius/taskmanager/scheduler/priorityqueue/GroupedPriorityQueueTest.java index c229b10..aaf508b 100644 --- a/source/taskmanager/src/test/java/nl/aerius/taskmanager/scheduler/priorityqueue/GroupedPriorityQueueTest.java +++ b/source/taskmanager/src/test/java/nl/aerius/taskmanager/scheduler/priorityqueue/GroupedPriorityQueueTest.java @@ -17,7 +17,9 @@ package nl.aerius.taskmanager.scheduler.priorityqueue; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.mock; @@ -49,11 +51,13 @@ void testPoll() { queue.add(task5); assertNotNull(queue.peek(), "Queue should have task at the queue."); + assertFalse(queue.isEmpty(), "Queue should not be empty."); assertEquals(task1, queue.poll(), "Poll should return first task added."); assertEquals(task3, queue.poll(), "Poll should return first task of task with correlationId 2."); assertEquals(task2, queue.poll(), "Poll should return second task of correlationId 1, because it's at the front of the queue."); assertEquals(task4, queue.poll(), "Poll should return second task of correlationId 2, because it should be at front of the queue."); assertEquals(task5, queue.poll(), "Poll should return last added task."); + assertTrue(queue.isEmpty(), "Queue should be empty after all tasks taken from the queue."); } private static Task mockTask(final String correlationId) {