From 5f0e10177cca7ab2282c6f20076fc5047e13e9aa Mon Sep 17 00:00:00 2001 From: Hilbrand Bouwkamp Date: Tue, 15 Sep 2026 11:02:38 +0200 Subject: [PATCH 1/2] Refactored querying RabbitMQ query status from RabbitMQ admin api This changes the way the scheduled RabbitMQ admin api is queried to get the latest stats on queries. It used to query per worker queue, but this change queries the data for all queues at once. This will make it easier to also get status on all client queues. This change is in preparation to create telemetry metrics on client queues. To keep things simple the changes for client queue metrics is not included in this pr. Most of the changes in the pr are due to using a record to pass the states instead of the individual values to methods. It also changes the frequency the api is queried to each 20 seconds. It was 60 seconds, but it did that for each worker queue. With this change it only needs 1 call and we'll get slightly more detailed data by querying a bit more often. --- .../ConnectionConfiguration.java | 2 +- .../nl/aerius/taskmanager/QueueWatchDog.java | 5 +- .../nl/aerius/taskmanager/StartupGuard.java | 3 +- .../nl/aerius/taskmanager/WorkerPool.java | 7 +- .../adaptor/WorkerSizeObserver.java | 10 +-- .../domain/RabbitMQQueueStatus.java | 23 ++++++ .../taskmanager/metrics/LoadMetric.java | 10 +++ .../metrics/RabbitMQUsageMetricsProvider.java | 19 ++--- .../metrics/TaskManagerMetricsRegister.java | 13 +-- .../TaskManagerUsageMetricsWrapper.java | 14 ++-- .../metrics/UsageMetricsReporter.java | 9 +-- ...er.java => WorkerUsageMetricsWrapper.java} | 4 +- .../taskmanager/mq/RabbitMQQueueMonitor.java | 61 +++++++++++--- .../mq/RabbitMQWorkerSizeProvider.java | 81 +++++++++---------- .../priorityqueue/PriorityTaskScheduler.java | 2 +- .../PriorityTaskSchedulerMetrics.java | 36 +++------ .../aerius/taskmanager/QueueWatchDogTest.java | 7 +- .../aerius/taskmanager/StartupGuardTest.java | 8 +- .../taskmanager/TaskDispatcherTest.java | 17 ++-- .../aerius/taskmanager/TaskManagerTest.java | 3 +- .../nl/aerius/taskmanager/WorkerPoolTest.java | 17 ++-- .../TaskManagerMetricsRegisterTest.java | 11 ++- .../mq/RabbitMQQueueMonitorTest.java | 26 ++++-- .../mq/RabbitMQWorkerSizeProviderTest.java | 21 +++-- .../nl/aerius/taskmanager/mq/queue_aerius.txt | 1 + .../mq/queue_aerius.worker.ops.txt | 2 +- 26 files changed, 251 insertions(+), 161 deletions(-) create mode 100644 source/taskmanager/src/main/java/nl/aerius/taskmanager/domain/RabbitMQQueueStatus.java rename source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/{UsageMetricsWrapper.java => WorkerUsageMetricsWrapper.java} (94%) create mode 100644 source/taskmanager/src/test/resources/nl/aerius/taskmanager/mq/queue_aerius.txt diff --git a/source/taskmanager-client/src/main/java/nl/aerius/taskmanager/client/configuration/ConnectionConfiguration.java b/source/taskmanager-client/src/main/java/nl/aerius/taskmanager/client/configuration/ConnectionConfiguration.java index ac6fca33..9d87959e 100644 --- a/source/taskmanager-client/src/main/java/nl/aerius/taskmanager/client/configuration/ConnectionConfiguration.java +++ b/source/taskmanager-client/src/main/java/nl/aerius/taskmanager/client/configuration/ConnectionConfiguration.java @@ -43,7 +43,7 @@ public final class ConnectionConfiguration { /** * Default refresh time in seconds. */ - private static final int DEFAULT_MANAGEMENT_REFRESH_RATE = 60; //seconds + private static final int DEFAULT_MANAGEMENT_REFRESH_RATE = 20; //seconds /** * Default wait time before retrying to connect. diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/QueueWatchDog.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/QueueWatchDog.java index defb72e4..bd51951c 100644 --- a/source/taskmanager/src/main/java/nl/aerius/taskmanager/QueueWatchDog.java +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/QueueWatchDog.java @@ -29,6 +29,7 @@ import nl.aerius.taskmanager.adaptor.WorkerProducer.WorkerProducerHandler; import nl.aerius.taskmanager.adaptor.WorkerSizeObserver; import nl.aerius.taskmanager.domain.QueueWatchDogListener; +import nl.aerius.taskmanager.domain.RabbitMQQueueStatus; /** * WatchDog to detect dead messages. Dead messages are messages once put on the queue, but those messages have gone. For example because @@ -73,9 +74,9 @@ public void onWorkerFinished(final String messageId, final Map m } @Override - public void onNumberOfWorkersUpdate(final int numberOfWorkers, final int numberOfMessages, final int numberOfMessagesInProgress) { + public void onNumberOfWorkersUpdate(final RabbitMQQueueStatus queueStatus) { synchronized (runningTasks) { - if (isItDead(!runningTasks.isEmpty(), numberOfMessages)) { + if (isItDead(!runningTasks.isEmpty(), queueStatus.messages())) { LOG.info("It looks like some tasks are zombies on {} worker queue. All tasks in state running are released (running:{}).", workerQueueName, runningTasks.size()); runningTasks.clear(); diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/StartupGuard.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/StartupGuard.java index f8e547c7..36d25ce6 100644 --- a/source/taskmanager/src/main/java/nl/aerius/taskmanager/StartupGuard.java +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/StartupGuard.java @@ -20,6 +20,7 @@ import java.util.concurrent.locks.ReentrantLock; import nl.aerius.taskmanager.adaptor.WorkerSizeObserver; +import nl.aerius.taskmanager.domain.RabbitMQQueueStatus; /** * Class to be used at startup. The Scheduler should not start before the number of messages on the queue is zero. @@ -48,7 +49,7 @@ public void waitForOpen() throws InterruptedException { } @Override - public void onNumberOfWorkersUpdate(final int numberOfWorkers, final int numberOfMessages, final int numberOfMessagesInProgress) { + public void onNumberOfWorkersUpdate(final RabbitMQQueueStatus queueStatus) { lock.lock(); try { if (!open) { diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/WorkerPool.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/WorkerPool.java index 6fec2682..cd4af83d 100644 --- a/source/taskmanager/src/main/java/nl/aerius/taskmanager/WorkerPool.java +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/WorkerPool.java @@ -30,6 +30,7 @@ import nl.aerius.taskmanager.adaptor.WorkerProducer.WorkerProducerHandler; import nl.aerius.taskmanager.adaptor.WorkerSizeObserver; import nl.aerius.taskmanager.domain.QueueWatchDogListener; +import nl.aerius.taskmanager.domain.RabbitMQQueueStatus; import nl.aerius.taskmanager.domain.Task; import nl.aerius.taskmanager.domain.TaskRecord; import nl.aerius.taskmanager.domain.WorkerUpdateHandler; @@ -178,13 +179,13 @@ public void reserveWorker() { } @Override - public void onNumberOfWorkersUpdate(final int numberOfWorkers, final int numberOfMessages, final int numberOfMessagesInProgress) { + public void onNumberOfWorkersUpdate(final RabbitMQQueueStatus queueStatus) { synchronized (this) { if (!firstUpdateReceived) { - initialUnaccountedWorkers = numberOfMessages; + initialUnaccountedWorkers = queueStatus.messages(); firstUpdateReceived = true; } - updateNumberOfWorkers(numberOfWorkers); + updateNumberOfWorkers(queueStatus.consumers()); } } diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/adaptor/WorkerSizeObserver.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/adaptor/WorkerSizeObserver.java index 21963c00..8ba2a183 100644 --- a/source/taskmanager/src/main/java/nl/aerius/taskmanager/adaptor/WorkerSizeObserver.java +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/adaptor/WorkerSizeObserver.java @@ -16,17 +16,17 @@ */ package nl.aerius.taskmanager.adaptor; +import nl.aerius.taskmanager.domain.RabbitMQQueueStatus; + /** * Interface called to observer the number of workers. */ public interface WorkerSizeObserver { /** - * Gives the number of workers processes connected on the queue. + * Gives metrics on a RabbitMQ queue. * - * @param numberOfWorkers number of number of workers processes - * @param numberOfMessages Total number of messages on the queue - * @param numberOfMessagesInProgress Number of messages being processed by the workers + * @param queueStatus RabbitMQ status metrics */ - void onNumberOfWorkersUpdate(final int numberOfWorkers, final int numberOfMessages, int numberOfMessagesInProgress); + void onNumberOfWorkersUpdate(RabbitMQQueueStatus queueStatus); } diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/domain/RabbitMQQueueStatus.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/domain/RabbitMQQueueStatus.java new file mode 100644 index 00000000..0feb8391 --- /dev/null +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/domain/RabbitMQQueueStatus.java @@ -0,0 +1,23 @@ +/* + * 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; + +/** + * Data record containing several collected metrics from the admin API of a single RabbitMQ queue. + */ +public record RabbitMQQueueStatus(int consumers, int messages, int unacknowledged) { +} diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/LoadMetric.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/LoadMetric.java index 1bca565c..50bc016f 100644 --- a/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/LoadMetric.java +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/LoadMetric.java @@ -52,6 +52,16 @@ class LoadMetric { private final ToDoubleBiFunction sumFunction; private final Object lock = new Object(); + /** + * Constructor + * + * @param countFunction Returns the value to use a count in the previous time frame. It gets 2 parameters + * - numberOfWorkers: The number of workers available in the last time period. + * - usedWorkers: The number of workers used in the last time period. + * @param sumFunction Returns the average value to report based on the 2 parameters: + * - total: the sum of all counts since the last time this metric was requested. + * - totalMeasureTime: the total time since the last time this metric was requested. + */ public LoadMetric(final ToDoubleBiFunction countFunction, final ToDoubleBiFunction sumFunction) { this.countFunction = countFunction; this.sumFunction = sumFunction; diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/RabbitMQUsageMetricsProvider.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/RabbitMQUsageMetricsProvider.java index 3993c0d2..5d0ab2d5 100644 --- a/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/RabbitMQUsageMetricsProvider.java +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/RabbitMQUsageMetricsProvider.java @@ -17,24 +17,21 @@ package nl.aerius.taskmanager.metrics; import nl.aerius.taskmanager.adaptor.WorkerSizeObserver; +import nl.aerius.taskmanager.domain.RabbitMQQueueStatus; public class RabbitMQUsageMetricsProvider implements WorkerSizeObserver, UsageMetricsProvider { private final String workerQueueName; - private int numberOfWorkers; - private int numberOfMessages; - private int numberOfMessagesInProgress; + private RabbitMQQueueStatus queueStatus = new RabbitMQQueueStatus(0,0,0); public RabbitMQUsageMetricsProvider(final String workerQueueName) { this.workerQueueName = workerQueueName; } @Override - public void onNumberOfWorkersUpdate(final int numberOfWorkers, final int numberOfMessages, final int numberOfMessagesInProgress) { - this.numberOfWorkers = numberOfWorkers; - this.numberOfMessages = numberOfMessages; - this.numberOfMessagesInProgress = numberOfMessagesInProgress; + public void onNumberOfWorkersUpdate(final RabbitMQQueueStatus queueStatus) { + this.queueStatus = queueStatus; } @Override @@ -44,22 +41,22 @@ public String getWorkerQueueName() { @Override public int getNumberOfWorkers() { - return numberOfWorkers; + return queueStatus.consumers(); } @Override public int getNumberOfUsedWorkers() { - return numberOfMessagesInProgress; + return queueStatus.unacknowledged(); } @Override public int getNumberOfFreeWorkers() { - return Math.max(0, numberOfWorkers - numberOfMessagesInProgress); + return Math.max(0, getNumberOfWorkers() - getNumberOfUsedWorkers()); } @Override public int getNumberOfWaiting() { - return Math.max(0, numberOfMessages - numberOfMessagesInProgress); + return Math.max(0, queueStatus.messages() - getNumberOfUsedWorkers()); } } diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/TaskManagerMetricsRegister.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/TaskManagerMetricsRegister.java index 00e828f8..f5afb418 100644 --- a/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/TaskManagerMetricsRegister.java +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/TaskManagerMetricsRegister.java @@ -25,6 +25,7 @@ import nl.aerius.taskmanager.adaptor.WorkerProducer.WorkerProducerHandler; import nl.aerius.taskmanager.adaptor.WorkerSizeObserver; import nl.aerius.taskmanager.domain.QueueWatchDogListener; +import nl.aerius.taskmanager.domain.RabbitMQQueueStatus; /** * This class provides the input for the {@link TaskManagerUsageMetricsProvider}. It will register updates on the amount of worker/workers from @@ -57,12 +58,14 @@ public void onWorkerFinished(final String messageId, final Map m } @Override - public synchronized void onNumberOfWorkersUpdate(final int numberOfWorkers, final int numberOfMessages, final int numberOfMessagesInProgress) { - this.numberOfWorkers = numberOfWorkers; - if (!startupGuard.isOpen() && numberOfMessages > 0) { + public void onNumberOfWorkersUpdate(final RabbitMQQueueStatus queueStatus) { + this.numberOfWorkers = queueStatus.consumers(); + final int messages = queueStatus.messages(); + + if (!startupGuard.isOpen() && messages > 0) { LOG.info("Queue {} will be started with {} messages already on the queue.", taskManagerUsageMetricsProvider.getWorkerQueueName(), - numberOfMessages); - taskManagerUsageMetricsProvider.register(numberOfMessages, numberOfWorkers); + messages); + taskManagerUsageMetricsProvider.register(messages, numberOfWorkers); } else { taskManagerUsageMetricsProvider.register(0, numberOfWorkers); } diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/TaskManagerUsageMetricsWrapper.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/TaskManagerUsageMetricsWrapper.java index 3863899e..3b8c43ba 100644 --- a/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/TaskManagerUsageMetricsWrapper.java +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/TaskManagerUsageMetricsWrapper.java @@ -20,19 +20,19 @@ /** * Class that wraps all Task Manager usage metrics reporters. TaskSchedulerBuckets should add the metric providers to this class. - * Each of those providers added is for a specific worker queue. The {@link UsageMetricsWrapper} will manager metrics per worker queue. + * Each of those providers added is for a specific worker queue. The {@link WorkerUsageMetricsWrapper} will manager metrics per worker queue. */ public class TaskManagerUsageMetricsWrapper { - private final UsageMetricsWrapper rabbitMQUsageMetrics; - private final UsageMetricsWrapper workerPoolUsageMetrics; - private final UsageMetricsWrapper taskManagerUsageMetrics; + private final WorkerUsageMetricsWrapper rabbitMQUsageMetrics; + private final WorkerUsageMetricsWrapper workerPoolUsageMetrics; + private final WorkerUsageMetricsWrapper taskManagerUsageMetrics; private final UsageMetricsReporter loadUsageMetricsReporter; public TaskManagerUsageMetricsWrapper(final Meter meter) { - rabbitMQUsageMetrics = new UsageMetricsWrapper(meter, "aer.rabbitmq", true); - workerPoolUsageMetrics = new UsageMetricsWrapper(meter, "aer.taskmanager.workerpool", false); - taskManagerUsageMetrics = new UsageMetricsWrapper(meter, "aer.taskmanager", false); + rabbitMQUsageMetrics = new WorkerUsageMetricsWrapper(meter, "aer.rabbitmq", true); + workerPoolUsageMetrics = new WorkerUsageMetricsWrapper(meter, "aer.taskmanager.workerpool", false); + taskManagerUsageMetrics = new WorkerUsageMetricsWrapper(meter, "aer.taskmanager", false); loadUsageMetricsReporter = new UsageMetricsReporter(meter, "aer.taskmanager.work.load", "Report average worker load"); } diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/UsageMetricsReporter.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/UsageMetricsReporter.java index d3e72626..6614294a 100644 --- a/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/UsageMetricsReporter.java +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/UsageMetricsReporter.java @@ -17,15 +17,12 @@ package nl.aerius.taskmanager.metrics; import java.util.ArrayList; -import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.Map.Entry; +import java.util.concurrent.ConcurrentHashMap; import java.util.function.DoubleSupplier; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - import io.opentelemetry.api.common.Attributes; import io.opentelemetry.api.metrics.Meter; import io.opentelemetry.api.metrics.ObservableDoubleGauge; @@ -37,11 +34,9 @@ */ class UsageMetricsReporter { - private static final Logger LOG = LoggerFactory.getLogger(UsageMetricsReporter.class); - private record UsageMetric(DoubleSupplier metricSupplier, Attributes attributes) {} - private final Map> metricsMap = new HashMap<>(); + private final Map> metricsMap = new ConcurrentHashMap<>(); private final ObservableDoubleGauge gauge; public UsageMetricsReporter(final Meter meter, final String metricName, final String description) { diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/UsageMetricsWrapper.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/WorkerUsageMetricsWrapper.java similarity index 94% rename from source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/UsageMetricsWrapper.java rename to source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/WorkerUsageMetricsWrapper.java index 464be144..2e267d99 100644 --- a/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/UsageMetricsWrapper.java +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/WorkerUsageMetricsWrapper.java @@ -23,12 +23,12 @@ * The metrics supported are limit (i.e. number of workers available). * */ -class UsageMetricsWrapper { +class WorkerUsageMetricsWrapper { private final boolean hasWaiting; private final UsageMetricsReporter limitReporter; private final UsageMetricsReporter usageReporter; - public UsageMetricsWrapper(final Meter meter, final String metricPrefix, final boolean hasWaiting) { + public WorkerUsageMetricsWrapper(final Meter meter, final String metricPrefix, final boolean hasWaiting) { this.hasWaiting = hasWaiting; limitReporter = new UsageMetricsReporter(meter, metricPrefix + ".worker.limit", "Report nunber of workers available"); usageReporter = new UsageMetricsReporter(meter, metricPrefix + ".worker.usage", "Report worker usage"); diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/mq/RabbitMQQueueMonitor.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/mq/RabbitMQQueueMonitor.java index 6cc14946..b2e444d1 100644 --- a/source/taskmanager/src/main/java/nl/aerius/taskmanager/mq/RabbitMQQueueMonitor.java +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/mq/RabbitMQQueueMonitor.java @@ -22,6 +22,8 @@ import java.net.URISyntaxException; import java.net.URL; import java.nio.charset.StandardCharsets; +import java.util.HashMap; +import java.util.Map; import java.util.concurrent.TimeUnit; import org.apache.http.HttpHost; @@ -44,16 +46,17 @@ import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; +import com.fasterxml.jackson.databind.node.ArrayNode; -import nl.aerius.taskmanager.adaptor.WorkerSizeObserver; import nl.aerius.taskmanager.client.configuration.ConnectionConfiguration; +import nl.aerius.taskmanager.domain.RabbitMQQueueStatus; /** * RabbitMQ implementation to manage implementation specific part of the worker pool. This covers managing the total size of available workers and * informing the worker pool when a worker is finished. *

When the connection shuts down this pool manager will shutdown. */ -public class RabbitMQQueueMonitor { +class RabbitMQQueueMonitor { private static final Logger LOG = LoggerFactory.getLogger(RabbitMQQueueMonitor.class); @@ -104,28 +107,66 @@ public void close() { } } - public void updateWorkerQueueState(final String queueName, final WorkerSizeObserver observer) { + /** + * Retrieves the queue status for the given queue from the RabbitMQ admin api. + * + * @param queueName name of the queue to get the statuss + * @return Status of the queue + */ + public RabbitMQQueueStatus getWorkerQueueState(final String queueName) { // Use RabbitMQ HTTP-API. // URL: [host]:[port]/api/queues/[virtualHost]/[QueueName] final String virtualHost = configuration.getBrokerVirtualHost().replace("/", "%2f"); final String apiPath = String.format("/api/queues/%s/%s", virtualHost, queueName); try { - final JsonNode jsonObject = getJsonResultFromApi(apiPath); + final JsonNode jsonNode = getJsonResultFromApi(apiPath); - if (jsonObject == null) { + if (jsonNode == null) { LOG.error("Queue configuration from RabbitMQ admin json get call returned null."); } else { - final int numberOfWorkers = getJsonIntPrimitive(jsonObject, "consumers"); - final int numberOfMessages = getJsonIntPrimitive(jsonObject, "messages"); - final int numberOfMessagesInProgress = getJsonIntPrimitive(jsonObject, "messages_unacknowledged"); + return getQueueStatus(jsonNode); + } + } catch (final URISyntaxException | IOException e) { + LOG.info("Error getting RabbitMQ status from admin api: {}", e.getMessage()); + } + return null; + } + + /** + * Retrieves the status for all queues from the RabbitMQ admin api. + */ + public Map getWorkerQueueStates() { + try { + final JsonNode jsonObject = getJsonResultFromApi("/api/queues"); + + if (jsonObject == null) { + LOG.error("Queue configuration from RabbitMQ admin json get call returned null."); + } if (jsonObject instanceof final ArrayNode array) { + final Map queueStates = new HashMap<>(); - observer.onNumberOfWorkersUpdate(numberOfWorkers, numberOfMessages, numberOfMessagesInProgress); - LOG.trace("[{}] active workers:{}", queueName, numberOfWorkers); + for (int i = 0; i < array.size(); i++) { + final JsonNode jsonNode = array.get(i); + queueStates.put(getJsonString(jsonNode, "name"), getQueueStatus(jsonNode)); + } + return queueStates; } } catch (final URISyntaxException | IOException e) { LOG.info("Error getting RabbitMQ status from admin api: {}", e.getMessage()); } + return new HashMap<>(); + } + + private static RabbitMQQueueStatus getQueueStatus(final JsonNode jsonNode) { + final int consumers = getJsonIntPrimitive(jsonNode, "consumers"); + final int mesages = getJsonIntPrimitive(jsonNode, "messages"); + final int unacknowledged = getJsonIntPrimitive(jsonNode, "messages_unacknowledged"); + + return new RabbitMQQueueStatus(consumers, mesages, unacknowledged); + } + + private static String getJsonString(final JsonNode jsonObject, final String key) { + return jsonObject == null || !jsonObject.has(key) ? "" : jsonObject.get(key).asText(); } private static int getJsonIntPrimitive(final JsonNode jsonObject, final String key) { diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/mq/RabbitMQWorkerSizeProvider.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/mq/RabbitMQWorkerSizeProvider.java index 4dd4a28d..ddccbe04 100644 --- a/source/taskmanager/src/main/java/nl/aerius/taskmanager/mq/RabbitMQWorkerSizeProvider.java +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/mq/RabbitMQWorkerSizeProvider.java @@ -25,6 +25,7 @@ import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; +import java.util.function.Supplier; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -33,6 +34,7 @@ import nl.aerius.taskmanager.adaptor.WorkerSizeObserver; import nl.aerius.taskmanager.adaptor.WorkerSizeProviderProxy; import nl.aerius.taskmanager.client.BrokerConnectionFactory; +import nl.aerius.taskmanager.domain.RabbitMQQueueStatus; /** * Provider to use different means to get information about the size and utilisation of the workers. @@ -53,7 +55,6 @@ public class RabbitMQWorkerSizeProvider implements WorkerSizeProviderProxy { private static final long DELAY_BEFORE_UPDATE_TIME_SECONDS = 15; private final ScheduledExecutorService executorService; - private final BrokerConnectionFactory factory; private final RabbitMQChannelQueueEventsWatcher channelQueueEventsWatcher; private final RabbitMQWorkerEventProducer eventProducer; /** @@ -69,13 +70,19 @@ public class RabbitMQWorkerSizeProvider implements WorkerSizeProviderProxy { private final Map> lastRuns = new HashMap<>(); private final Map observers = new HashMap<>(); - private final Map monitors = new HashMap<>(); - + private final RabbitMQQueueMonitor monitor; + // Map to keep track of the last known states of the queues as retrieved form the RabbitMQ admin API. + private Map lastKnownQueueStates = new HashMap<>(); private boolean running; public RabbitMQWorkerSizeProvider(final ScheduledExecutorService executorService, final BrokerConnectionFactory factory) { + this(executorService, factory, new RabbitMQQueueMonitor(factory.getConnectionConfiguration())); + } + + RabbitMQWorkerSizeProvider(final ScheduledExecutorService executorService, final BrokerConnectionFactory factory, + final RabbitMQQueueMonitor monitor) { this.executorService = executorService; - this.factory = factory; + this.monitor = monitor; channelQueueEventsWatcher = new RabbitMQChannelQueueEventsWatcher(factory, this); refreshRateSeconds = factory.getConnectionConfiguration().getBrokerManagementRefreshRate(); eventProducer = new RabbitMQWorkerEventProducer(executorService, factory); @@ -83,41 +90,17 @@ public RabbitMQWorkerSizeProvider(final ScheduledExecutorService executorService } @Override - public void addObserver(final String workerQueueName, final WorkerSizeObserver observer) { - if (!observers.containsKey(workerQueueName)) { - if (refreshRateSeconds > 0) { - final RabbitMQQueueMonitor monitor = new RabbitMQQueueMonitor(factory.getConnectionConfiguration()); - - putMonitor(workerQueueName, monitor); - } else { - LOG.info("Not monitoring RabbitMQ admin api because refresh delay was {} seconds", refreshRateSeconds); - } - } - observers.computeIfAbsent(workerQueueName, k -> new WorkerSizeObserverComposite()).add(observer); + public void addObserver(final String queueName, final WorkerSizeObserver observer) { + observers.computeIfAbsent(queueName, k -> new WorkerSizeObserverComposite()).add(observer); if (observer instanceof WorkerMetrics) { - eventProducer.addMetrics(workerQueueName, (WorkerMetrics) observer); + eventProducer.addMetrics(queueName, (WorkerMetrics) observer); } } - /** - * Store the monitor. Should only be called outside of this class from unit tests to add a mock monitor. - * - * @param workerQueueName - * @param monitor - */ - void putMonitor(final String workerQueueName, final RabbitMQQueueMonitor monitor) { - monitors.put(workerQueueName, monitor); - } - @Override - public boolean removeObserver(final String workerQueueName) { - final RabbitMQQueueMonitor monitor = monitors.remove(workerQueueName); - - if (monitor != null) { - monitor.shutdown(); - } - eventProducer.removeMetrics(workerQueueName); - return observers.remove(workerQueueName) != null; + public boolean removeObserver(final String queueName) { + eventProducer.removeMetrics(queueName); + return observers.remove(queueName) != null; } @Override @@ -142,7 +125,8 @@ public void shutdown() { private void updateWorkerQueueState() { if (running) { try { - monitors.forEach((k, v) -> triggerWorkerQueueState(k)); + lastKnownQueueStates = new HashMap<>(monitor.getWorkerQueueStates()); + observers.forEach((q, v) -> scheduledUpdateWorkerQueueState(q, () -> lastKnownQueueStates.get(q))); } catch (final RuntimeException e) { LOG.error("Runtime error during updateWorkerQueueState", e); } @@ -151,22 +135,37 @@ private void updateWorkerQueueState() { @Override public void triggerWorkerQueueState(final String queueName) { + if (!observers.containsKey(queueName)) { + return; + } + scheduledUpdateWorkerQueueState(queueName, () -> monitor.getWorkerQueueState(queueName)); + } + + private void scheduledUpdateWorkerQueueState(final String queueName, final Supplier statusSupplier) { // This uses a delayed update. It schedules a task to run in x-seconds. // If a new update is received before the schedule has run it will cancel the current schedule and reschedule. // This is mainly for when multiple events are triggered to not trigger a call for every event, // and also to manage the events trigger in combination with the scheduled process. synchronized (sync) { Optional.ofNullable(lastRuns.get(queueName)).ifPresent(f -> f.cancel(false)); - final Runnable updateTask = () -> updateWorkerQueueState(queueName); + final Runnable updateTask = () -> updateWorkerQueueState(queueName, statusSupplier); lastRuns.put(queueName, executorService.schedule(updateTask, refreshDelayBeforeUpdateSeconds, TimeUnit.SECONDS)); } } - private void updateWorkerQueueState(final String queueName) { + private void updateWorkerQueueState(final String queueName, final Supplier statusSupplier) { synchronized (sync) { - Optional.ofNullable(monitors.get(queueName)).ifPresent(m -> m.updateWorkerQueueState(queueName, observers.get(queueName))); - lastRuns.remove(queueName); + try { + Optional.ofNullable(statusSupplier.get()) + .ifPresent(s -> { + lastKnownQueueStates.put(queueName, s); + lastRuns.remove(queueName); + Optional.ofNullable(observers.get(queueName)).ifPresent(observer -> observer.onNumberOfWorkersUpdate(s)); + }); + } catch (final RuntimeException e) { + LOG.error("RuntimeException during updateWorkerQueueState", e); + } } } @@ -178,10 +177,10 @@ public void add(final WorkerSizeObserver observer) { } @Override - public void onNumberOfWorkersUpdate(final int numberOfWorkers, final int numberOfMessages, final int numberOfMessagesInProgress) { + public void onNumberOfWorkersUpdate(final RabbitMQQueueStatus queueStatus) { for (final WorkerSizeObserver observer : observers) { try { - observer.onNumberOfWorkersUpdate(numberOfWorkers, numberOfMessages, numberOfMessagesInProgress); + observer.onNumberOfWorkersUpdate(queueStatus); } catch (final RuntimeException e) { LOG.error("RuntimeException during onNumberOfWorkersUpdate in {}", observer.getClass(), e); } 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 cf78d920..681fcae5 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 @@ -203,7 +203,7 @@ public void updateQueue(final PriorityTaskQueue priorityTaskQueue) { try { final String queueName = priorityTaskQueue.getQueueName(); if (!priorityQueueMap.containsKey(queueName)) { - metrics.addMetric(() -> priorityQueueMap.onWorkerByQueue(queueName), workerQueueName, queueName); + metrics.addMetricUsed(() -> priorityQueueMap.onWorkerByQueue(queueName), workerQueueName, queueName); if (this.queue instanceof final GroupedPriorityQueue gpq) { metrics.addMetricWaiting(() -> gpq.getGroupSize(), workerQueueName, queueName); } diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/scheduler/priorityqueue/PriorityTaskSchedulerMetrics.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/scheduler/priorityqueue/PriorityTaskSchedulerMetrics.java index 14ccebbf..2bcd6d36 100644 --- a/source/taskmanager/src/main/java/nl/aerius/taskmanager/scheduler/priorityqueue/PriorityTaskSchedulerMetrics.java +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/scheduler/priorityqueue/PriorityTaskSchedulerMetrics.java @@ -21,6 +21,7 @@ import java.util.Optional; import java.util.function.IntSupplier; +import io.opentelemetry.api.common.Attributes; import io.opentelemetry.api.metrics.ObservableDoubleGauge; import nl.aerius.taskmanager.metrics.OpenTelemetryMetrics; @@ -30,15 +31,9 @@ */ class PriorityTaskSchedulerMetrics { - /** - * @deprecated replaced by "aer.taskmanager.client.queue" metric. - */ - @Deprecated - private static final String METRIC_PREFIX_LEGACY = "aer.taskmanager.running_client_size"; private static final String METRIC_PREFIX = "aer.taskmanager.client.queue"; private static final String DESCRIPTION = "Number of tasks running on client queues"; - private final Map metrics = new HashMap<>(); private final Map usageMetrics = new HashMap<>(); private final Map waitingMetrics = new HashMap<>(); @@ -49,19 +44,8 @@ class PriorityTaskSchedulerMetrics { * @param workerQueueName worker queue name * @param clientQueueName client queue name */ - public void addMetric(final IntSupplier countSupplier, final String workerQueueName, final String clientQueueName) { - metrics.put(clientQueueName, OpenTelemetryMetrics.METER - .gaugeBuilder(METRIC_PREFIX_LEGACY) - .setDescription(DESCRIPTION) - .buildWithCallback( - result -> result.record(countSupplier.getAsInt(), - OpenTelemetryMetrics.queueAttributes(workerQueueName, clientQueueName, "state", "used")))); - metrics.put(clientQueueName, OpenTelemetryMetrics.METER - .gaugeBuilder(METRIC_PREFIX) - .setDescription(DESCRIPTION) - .buildWithCallback( - result -> result.record(countSupplier.getAsInt(), - OpenTelemetryMetrics.queueAttributes(workerQueueName, clientQueueName, "state", "used")))); + public void addMetricUsed(final IntSupplier countSupplier, final String workerQueueName, final String clientQueueName) { + usageMetrics.put(clientQueueName, createMetric(countSupplier, workerQueueName, clientQueueName, "used")); } /** @@ -72,12 +56,17 @@ public void addMetric(final IntSupplier countSupplier, final String workerQueueN * @param clientQueueName client queue name */ public void addMetricWaiting(final IntSupplier countSupplier, final String workerQueueName, final String clientQueueName) { - waitingMetrics.put(clientQueueName, OpenTelemetryMetrics.METER + waitingMetrics.put(clientQueueName, createMetric(countSupplier, workerQueueName, clientQueueName, "waiting")); + } + + private ObservableDoubleGauge createMetric(final IntSupplier countSupplier, final String workerQueueName, final String clientQueueName, + final String state) { + final Attributes queueAttributes = OpenTelemetryMetrics.queueAttributes(workerQueueName, clientQueueName, "state", state); + + return OpenTelemetryMetrics.METER .gaugeBuilder(METRIC_PREFIX) .setDescription(DESCRIPTION) - .buildWithCallback( - result -> result.record(countSupplier.getAsInt(), - OpenTelemetryMetrics.queueAttributes(workerQueueName, clientQueueName, "state", "waiting")))); + .buildWithCallback(result -> result.record(countSupplier.getAsInt(), queueAttributes)); } /** @@ -86,7 +75,6 @@ public void addMetricWaiting(final IntSupplier countSupplier, final String worke * @param clienQueueName */ public void removeMetric(final String clienQueueName) { - Optional.ofNullable(metrics.remove(clienQueueName)).ifPresent(ObservableDoubleGauge::close); Optional.ofNullable(usageMetrics.remove(clienQueueName)).ifPresent(ObservableDoubleGauge::close); Optional.ofNullable(waitingMetrics.remove(clienQueueName)).ifPresent(ObservableDoubleGauge::close); } diff --git a/source/taskmanager/src/test/java/nl/aerius/taskmanager/QueueWatchDogTest.java b/source/taskmanager/src/test/java/nl/aerius/taskmanager/QueueWatchDogTest.java index 4bd67607..715d843e 100644 --- a/source/taskmanager/src/test/java/nl/aerius/taskmanager/QueueWatchDogTest.java +++ b/source/taskmanager/src/test/java/nl/aerius/taskmanager/QueueWatchDogTest.java @@ -31,6 +31,7 @@ import org.junit.jupiter.params.provider.MethodSource; import nl.aerius.taskmanager.domain.QueueWatchDogListener; +import nl.aerius.taskmanager.domain.RabbitMQQueueStatus; /** * Test class for {@link QueueWatchDog}. @@ -55,16 +56,16 @@ protected LocalDateTime now() { IntStream.range(0, runningWorkers).forEach(i -> qwd.onWorkDispatched(String.valueOf(i), null)); IntStream.range(0, finishedWorkers).forEach(i -> qwd.onWorkerFinished(String.valueOf(i), null)); - qwd.onNumberOfWorkersUpdate(0, numberOfMessages, 0); + qwd.onNumberOfWorkersUpdate(new RabbitMQQueueStatus(0, numberOfMessages, 0)); // reset should never trigger the first time the problem was reported. verify(listener, never()).reset(); // Fast forward 20 minutes to trigger reset if there is a problem. now.set(now.get().plusMinutes(20)); - qwd.onNumberOfWorkersUpdate(0, numberOfMessages, 0); + qwd.onNumberOfWorkersUpdate(new RabbitMQQueueStatus(0, numberOfMessages, 0)); verify(listener, times(expected)).reset(); // Call update again. This should not trigger reset again because we just called reset. - qwd.onNumberOfWorkersUpdate(0, numberOfMessages, 0); + qwd.onNumberOfWorkersUpdate(new RabbitMQQueueStatus(0, numberOfMessages, 0)); verify(listener, times(expected)).reset(); } diff --git a/source/taskmanager/src/test/java/nl/aerius/taskmanager/StartupGuardTest.java b/source/taskmanager/src/test/java/nl/aerius/taskmanager/StartupGuardTest.java index 4c526bc9..e5f68abc 100644 --- a/source/taskmanager/src/test/java/nl/aerius/taskmanager/StartupGuardTest.java +++ b/source/taskmanager/src/test/java/nl/aerius/taskmanager/StartupGuardTest.java @@ -25,6 +25,8 @@ import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Timeout; +import nl.aerius.taskmanager.domain.RabbitMQQueueStatus; + /** * Test class for {@link StartupGuard} */ @@ -35,9 +37,9 @@ void testOpen() { final StartupGuard guard = new StartupGuard(); assertFalse(guard.isOpen(), "Guard should not be open."); - guard.onNumberOfWorkersUpdate(0, 1, 0); + guard.onNumberOfWorkersUpdate(new RabbitMQQueueStatus(0, 1, 0)); assertTrue(guard.isOpen(), "Guard should be open when onNumberOfWorkersUpdate is called."); - guard.onNumberOfWorkersUpdate(0, 1, 0); + guard.onNumberOfWorkersUpdate(new RabbitMQQueueStatus(0, 1, 0)); assertTrue(guard.isOpen(), "Guard should still remain open onNumberOfWorkersUpdate has been called."); } @@ -60,7 +62,7 @@ void testWaitForOpen() throws InterruptedException { // First wait for first semaphore to be unlocked. waitForStart.acquire(); assertFalse(guard.isOpen(), "Guard should not be open."); - guard.onNumberOfWorkersUpdate(1, 1, 0); + guard.onNumberOfWorkersUpdate(new RabbitMQQueueStatus(1, 1, 0)); // Wait for semaphore that is called after waitForOpen is unlocked. waitForOpen.acquire(); assertTrue(guard.isOpen(), "Guard should now be open."); diff --git a/source/taskmanager/src/test/java/nl/aerius/taskmanager/TaskDispatcherTest.java b/source/taskmanager/src/test/java/nl/aerius/taskmanager/TaskDispatcherTest.java index 7cb3b5b7..25f7f567 100644 --- a/source/taskmanager/src/test/java/nl/aerius/taskmanager/TaskDispatcherTest.java +++ b/source/taskmanager/src/test/java/nl/aerius/taskmanager/TaskDispatcherTest.java @@ -39,6 +39,7 @@ import nl.aerius.taskmanager.TaskDispatcher.State; import nl.aerius.taskmanager.domain.QueueConfig; +import nl.aerius.taskmanager.domain.RabbitMQQueueStatus; import nl.aerius.taskmanager.domain.Task; import nl.aerius.taskmanager.domain.TaskConsumer; @@ -80,17 +81,17 @@ void after() throws InterruptedException { @Timeout(value = 3, unit = TimeUnit.SECONDS) void testNoFreeWorkers() { // Add Worker which will unlock - workerPool.onNumberOfWorkersUpdate(1, 0, 0); + workerPool.onNumberOfWorkersUpdate(new RabbitMQQueueStatus(1, 0, 0)); executor.execute(dispatcher); await().until(() -> dispatcher.getState() == State.WAIT_FOR_TASK); // Remove worker, 1 worker locked but at this point no actual workers available. - workerPool.onNumberOfWorkersUpdate(0, 0, 0); + workerPool.onNumberOfWorkersUpdate(new RabbitMQQueueStatus(0, 0, 0)); // Send task, should get NoFreeWorkersException in dispatcher. forwardTaskAsync(createTask(), null); // Dispatcher should go back to wait for worker to become available. await().until(() -> dispatcher.getState() == State.WAIT_FOR_WORKER); assertEquals(0, workerPool.getReportedWorkerSize(), "WorkerPool should be empty"); - workerPool.onNumberOfWorkersUpdate(1, 0, 0); + workerPool.onNumberOfWorkersUpdate(new RabbitMQQueueStatus(1, 0, 0)); assertEquals(1, workerPool.getReportedWorkerSize(), "WorkerPool should have 1 running"); } @@ -100,7 +101,7 @@ void testForwardTest() { final Task task = createTask(); final Future future = forwardTaskAsync(task, null); executor.execute(dispatcher); - workerPool.onNumberOfWorkersUpdate(1, 0, 0); //add worker which will unlock + workerPool.onNumberOfWorkersUpdate(new RabbitMQQueueStatus(1, 0, 0)); //add worker which will unlock await().until(() -> dispatcher.getState() == State.WAIT_FOR_WORKER); await().until(future::isDone); assertFalse(future.isCancelled(), "Taskconsumer must be unlocked at this point without error"); @@ -114,7 +115,7 @@ void testForwardDuplicateTask() { executor.execute(dispatcher); final Future future = forwardTaskAsync(task, null); await().until(() -> dispatcher.getState() == State.WAIT_FOR_WORKER); - workerPool.onNumberOfWorkersUpdate(2, 0, 0); //add worker which will unlock + workerPool.onNumberOfWorkersUpdate(new RabbitMQQueueStatus(2, 0, 0)); //add worker which will unlock // Now force the issue. assertSame(TaskDispatcher.State.WAIT_FOR_TASK, dispatcher.getState(), "Taskdispatcher must be waiting for task"); // Forwarding same Task object, so same id. @@ -134,19 +135,19 @@ void testExceptionDuringForward() { final Future future = forwardTaskAsync(task, null); await().until(() -> dispatcher.getState() == State.WAIT_FOR_WORKER); // Now open up a worker - workerPool.onNumberOfWorkersUpdate(1, 0, 0); + workerPool.onNumberOfWorkersUpdate(new RabbitMQQueueStatus(1, 0, 0)); // At this point the exception should be thrown. This could be the case when rabbitmq connection is lost for a second. // Wait for it to be unlocked again await().until(() -> dispatcher.getState() == State.WAIT_FOR_TASK); //simulate workerpool being reset - workerPool.onNumberOfWorkersUpdate(0, 0, 0); + workerPool.onNumberOfWorkersUpdate(new RabbitMQQueueStatus(0, 0, 0)); //now stop throwing exception to indicate connection is restored again workerProducer.setShutdownExceptionOnForward(false); //simulate connection being restored by first forwarding task again forwardTaskAsync(task, future); await().until(() -> dispatcher.getState() == State.WAIT_FOR_WORKER); //now simulate the worker being back - workerPool.onNumberOfWorkersUpdate(1, 0, 0); + workerPool.onNumberOfWorkersUpdate(new RabbitMQQueueStatus(1, 0, 0)); //should now be unlocked, but waiting for worker to be done await().until(() -> dispatcher.getState() == State.WAIT_FOR_WORKER); workerPool.onWorkerFinished(task.getId(), Map.of()); diff --git a/source/taskmanager/src/test/java/nl/aerius/taskmanager/TaskManagerTest.java b/source/taskmanager/src/test/java/nl/aerius/taskmanager/TaskManagerTest.java index 27a49edf..179b8eca 100644 --- a/source/taskmanager/src/test/java/nl/aerius/taskmanager/TaskManagerTest.java +++ b/source/taskmanager/src/test/java/nl/aerius/taskmanager/TaskManagerTest.java @@ -40,6 +40,7 @@ import nl.aerius.taskmanager.adaptor.WorkerSizeProviderProxy; import nl.aerius.taskmanager.domain.PriorityTaskQueue; import nl.aerius.taskmanager.domain.PriorityTaskSchedule; +import nl.aerius.taskmanager.domain.RabbitMQQueueStatus; import nl.aerius.taskmanager.domain.RabbitMQQueueType; import nl.aerius.taskmanager.scheduler.priorityqueue.PriorityTaskSchedulerFileHandler; @@ -66,7 +67,7 @@ void setUp() throws IOException { doAnswer(a -> { // This will unblock the startup guard - ((WorkerSizeObserver) a.getArgument(1)).onNumberOfWorkersUpdate(0, 0, 0); + ((WorkerSizeObserver) a.getArgument(1)).onNumberOfWorkersUpdate(new RabbitMQQueueStatus(0, 0, 0)); return null; }).when(workerSizeProvider).addObserver(any(), any()); taskManager = new TaskManager<>(executor, scheduledExecutorService, factory, schedulerFactory, workerSizeProvider); diff --git a/source/taskmanager/src/test/java/nl/aerius/taskmanager/WorkerPoolTest.java b/source/taskmanager/src/test/java/nl/aerius/taskmanager/WorkerPoolTest.java index aa3d567e..f3155ddd 100644 --- a/source/taskmanager/src/test/java/nl/aerius/taskmanager/WorkerPoolTest.java +++ b/source/taskmanager/src/test/java/nl/aerius/taskmanager/WorkerPoolTest.java @@ -36,6 +36,7 @@ import nl.aerius.taskmanager.domain.ForwardTaskHandler; import nl.aerius.taskmanager.domain.Message; import nl.aerius.taskmanager.domain.QueueConfig; +import nl.aerius.taskmanager.domain.RabbitMQQueueStatus; import nl.aerius.taskmanager.domain.Task; import nl.aerius.taskmanager.domain.TaskConsumer; import nl.aerius.taskmanager.domain.WorkerUpdateHandler; @@ -77,7 +78,7 @@ public void messageDelivered(final Message message) { @Test void testWorkerPoolSizing() throws IOException { assertEquals(0, workerPool.getReportedWorkerSize(), "Check if workerPool size is empty at start"); - workerPool.onNumberOfWorkersUpdate(10, 0, 0); + workerPool.onNumberOfWorkersUpdate(new RabbitMQQueueStatus(10, 0, 0)); assertEquals(10, workerPool.getReportedWorkerSize(), "Check if workerPool size is changed after sizing"); assertEquals(10, numberOfWorkers, "Check if workerPool change handler called."); workerPool.reserveWorker(); @@ -90,10 +91,10 @@ void testWorkerPoolSizing() throws IOException { @Test void testWorkerPoolSizingWithInitialSize() throws IOException { - workerPool.onNumberOfWorkersUpdate(10, 5, 0); + workerPool.onNumberOfWorkersUpdate(new RabbitMQQueueStatus(10, 5, 0)); assertEquals(5, workerPool.getNumberOfUsedWorkers(), "Check if workerPool size is 5"); assertEquals(10, workerPool.getNumberOfWorkers(), "Internal worker size should match reported number of workers"); - workerPool.onNumberOfWorkersUpdate(10, 5, 0); + workerPool.onNumberOfWorkersUpdate(new RabbitMQQueueStatus(10, 5, 0)); assertEquals(5, workerPool.getNumberOfUsedWorkers(), "Check if workerPool size is still 5"); assertEquals(10, workerPool.getNumberOfWorkers(), "Internal worker size should still match reported number of workers"); IntStream.range(1, 6).forEach(a -> workerPool.onWorkerFinished("", null)); @@ -109,12 +110,12 @@ void testNoFreeWorkers() { @Test void testWorkerPoolScaleDown() throws IOException { - workerPool.onNumberOfWorkersUpdate(5, 0, 0); + workerPool.onNumberOfWorkersUpdate(new RabbitMQQueueStatus(5, 0, 0)); final Task task1 = createAndSendTaskToWorker(); final Task task2 = createAndSendTaskToWorker(); final Task task3 = createAndSendTaskToWorker(); assertEquals(5, workerPool.getReportedWorkerSize(), "Check if workerPool size is same after 2 workers running"); - workerPool.onNumberOfWorkersUpdate(1, 0, 0); + workerPool.onNumberOfWorkersUpdate(new RabbitMQQueueStatus(1, 0, 0)); assertEquals(3, workerPool.getNumberOfWorkers(), "Workpool size should match number of running tasks, since new total is lower than currently running"); assertEquals(1, workerPool.getReportedWorkerSize(), "Check if current workerPool size is same after decreasing # workers"); @@ -128,7 +129,7 @@ void testWorkerPoolScaleDown() throws IOException { @Test void testReleaseTaskTwice() throws IOException { - workerPool.onNumberOfWorkersUpdate(2, 0, 0); + workerPool.onNumberOfWorkersUpdate(new RabbitMQQueueStatus(2, 0, 0)); final Task task1 = createAndSendTaskToWorker(); final String id = task1.getId(); workerPool.releaseWorker(id); @@ -141,14 +142,14 @@ void testReleaseTaskTwice() throws IOException { @Test void testMessageDeliverd() throws IOException { - workerPool.onNumberOfWorkersUpdate(1, 0, 0); + workerPool.onNumberOfWorkersUpdate(new RabbitMQQueueStatus(1, 0, 0)); createAndSendTaskToWorker(); assertNotSame(0, message.getDeliveryTag(), "Check if message is delivered"); } @Test void testReset() throws IOException { - workerPool.onNumberOfWorkersUpdate(5, 0, 0); + workerPool.onNumberOfWorkersUpdate(new RabbitMQQueueStatus(5, 0, 0)); createAndSendTaskToWorker(); createAndSendTaskToWorker(); assertEquals(2, workerPool.getNumberOfUsedWorkers(), "Should report 2 workers running."); diff --git a/source/taskmanager/src/test/java/nl/aerius/taskmanager/metrics/TaskManagerMetricsRegisterTest.java b/source/taskmanager/src/test/java/nl/aerius/taskmanager/metrics/TaskManagerMetricsRegisterTest.java index 2608ca9d..bcb53e89 100644 --- a/source/taskmanager/src/test/java/nl/aerius/taskmanager/metrics/TaskManagerMetricsRegisterTest.java +++ b/source/taskmanager/src/test/java/nl/aerius/taskmanager/metrics/TaskManagerMetricsRegisterTest.java @@ -33,6 +33,7 @@ import nl.aerius.taskmanager.StartupGuard; import nl.aerius.taskmanager.client.TaskMetrics; +import nl.aerius.taskmanager.domain.RabbitMQQueueStatus; /** * Test class for {@link TaskManagerMetricsRegister}. @@ -64,7 +65,7 @@ void testOnWorkDispatched() { verifyTaskManagerUsageMetricsProvider(1, 0, startUpNrOfMessagesCaptor); register.onWorkDispatched("1", createMap(QUEUE_1, 100L)); register.onWorkDispatched("2", createMap(QUEUE_2, 200L)); - register.onNumberOfWorkersUpdate(10, 2, 2); + register.onNumberOfWorkersUpdate(new RabbitMQQueueStatus(10, 2, 2)); // Should have called register 4 times (startup, 2 for dispatch and 1 for update. // But total delta should be +2 for the 2 dispatched messages. verifyTaskManagerUsageMetricsProvider(4, 2, lastNrOfMessagesCaptor); @@ -78,7 +79,7 @@ void testOnWorkerFinished() { register.onWorkDispatched("1", createMap(QUEUE_1, 100L)); register.onWorkerFinished("1", createMap(QUEUE_1, 100L)); register.onWorkerFinished("2", createMap(QUEUE_2, 200L)); - register.onNumberOfWorkersUpdate(10, 2, 2); + register.onNumberOfWorkersUpdate(new RabbitMQQueueStatus(10, 2, 2)); // Should have called register 5 times (startup, 1 for dispatch, 2 for finished and 1 for update. // But total delta should be 0 for 2 startup + 1 dispatch - 2 for finish.. @@ -104,8 +105,10 @@ private void verifyTaskManagerUsageMetricsProvider(final int times, final int su } private void startUp(final int numberOfWorkers, final int numberOfMessages) { - register.onNumberOfWorkersUpdate(numberOfWorkers, numberOfMessages, 0); - startupGuard.onNumberOfWorkersUpdate(numberOfWorkers, numberOfMessages, 0); + final RabbitMQQueueStatus status = new RabbitMQQueueStatus(numberOfWorkers, numberOfMessages, 0); + + register.onNumberOfWorkersUpdate(status); + startupGuard.onNumberOfWorkersUpdate(status); } private Map createMap(final String queueName, final long duration) { diff --git a/source/taskmanager/src/test/java/nl/aerius/taskmanager/mq/RabbitMQQueueMonitorTest.java b/source/taskmanager/src/test/java/nl/aerius/taskmanager/mq/RabbitMQQueueMonitorTest.java index c71fe0fe..04ea5ca4 100644 --- a/source/taskmanager/src/test/java/nl/aerius/taskmanager/mq/RabbitMQQueueMonitorTest.java +++ b/source/taskmanager/src/test/java/nl/aerius/taskmanager/mq/RabbitMQQueueMonitorTest.java @@ -21,15 +21,15 @@ import java.io.IOException; import java.io.InputStream; import java.io.InputStreamReader; -import java.util.concurrent.atomic.AtomicInteger; +import java.util.function.Function; import org.junit.jupiter.api.Test; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; -import nl.aerius.taskmanager.adaptor.WorkerSizeObserver; import nl.aerius.taskmanager.client.configuration.ConnectionConfiguration; +import nl.aerius.taskmanager.domain.RabbitMQQueueStatus; /** * Test class for {@link RabbitMQQueueMonitor}. @@ -37,26 +37,38 @@ class RabbitMQQueueMonitorTest { private static final String DUMMY = "dummy"; + private static final String QUEUENAME = "aerius.worker.ops"; private final ObjectMapper objectMapper = new ObjectMapper(); @Test void testGetWorkerQueueState() { + assertRabbitMQQueueMonitor("queue_aerius.worker.ops.txt", 4, 3, 5, rpm -> rpm.getWorkerQueueState(DUMMY)); + } + + @Test + void testGetWorkerQueueStates() { + assertRabbitMQQueueMonitor("queue_aerius.txt", 51, 10, 30, rpm -> rpm.getWorkerQueueStates().get(QUEUENAME)); + } + + private void assertRabbitMQQueueMonitor(final String filename, final int expectedConsumers, final int expectedMessages, + final int expectedUnacknowledged, final Function collector) { final ConnectionConfiguration configuration = ConnectionConfiguration.builder() .brokerHost(DUMMY).brokerPort(0).brokerUsername(DUMMY).brokerPassword(DUMMY).build(); - final AtomicInteger workerSize = new AtomicInteger(); - final WorkerSizeObserver mwps = (numberOfWorkers, numberOfMessages, numberOfMessagesInProgress) -> workerSize.set(numberOfWorkers); final RabbitMQQueueMonitor rpm = new RabbitMQQueueMonitor(configuration) { @Override protected JsonNode getJsonResultFromApi(final String apiPath) throws IOException { - try (final InputStream fr = getClass().getResourceAsStream("queue_aerius.worker.ops.txt"); + try (final InputStream fr = getClass().getResourceAsStream(filename); final InputStreamReader is = new InputStreamReader(fr)) { return objectMapper.readTree(is); } } }; try { - rpm.updateWorkerQueueState(DUMMY, mwps); - assertEquals(4, workerSize.get(), "Number of workers"); + final RabbitMQQueueStatus status = collector.apply(rpm); + + assertEquals(expectedConsumers, status.consumers(), "Number of workers"); + assertEquals(expectedMessages, status.messages(), "Number of messages"); + assertEquals(expectedUnacknowledged, status.unacknowledged(), "Number of unacknowledged"); } finally { rpm.shutdown(); } diff --git a/source/taskmanager/src/test/java/nl/aerius/taskmanager/mq/RabbitMQWorkerSizeProviderTest.java b/source/taskmanager/src/test/java/nl/aerius/taskmanager/mq/RabbitMQWorkerSizeProviderTest.java index 780104e1..4579c67d 100644 --- a/source/taskmanager/src/test/java/nl/aerius/taskmanager/mq/RabbitMQWorkerSizeProviderTest.java +++ b/source/taskmanager/src/test/java/nl/aerius/taskmanager/mq/RabbitMQWorkerSizeProviderTest.java @@ -18,8 +18,8 @@ import static org.junit.jupiter.api.Assertions.assertFalse; import static org.mockito.ArgumentMatchers.any; -import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.verify; @@ -30,16 +30,23 @@ import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Timeout; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; import nl.aerius.taskmanager.adaptor.WorkerSizeObserver; +import nl.aerius.taskmanager.domain.RabbitMQQueueStatus; /** * Test class for {@link RabbitMQWorkerSizeProvider} */ +@ExtendWith(MockitoExtension.class) class RabbitMQWorkerSizeProviderTest extends AbstractRabbitMQTest { private static final String TEST_QUEUE = "test"; + private @Mock RabbitMQQueueMonitor mockMonitor; + private RabbitMQWorkerSizeProvider provider; @Override @@ -47,25 +54,27 @@ class RabbitMQWorkerSizeProviderTest extends AbstractRabbitMQTest { void setUp() throws Exception { brokerManagementRefreshRate = 5; super.setUp(); - provider = new RabbitMQWorkerSizeProvider(executor, factory); + provider = new RabbitMQWorkerSizeProvider(executor, factory, mockMonitor); } @Test @Timeout(value = 10, unit = TimeUnit.SECONDS) void testTriggerWorkerQueueState() throws InterruptedException { + doReturn(new RabbitMQQueueStatus(1, 2, 3)).when(mockMonitor).getWorkerQueueState(TEST_QUEUE); final CountDownLatch latch = new CountDownLatch(1); - final RabbitMQQueueMonitor mockMonitor = mock(RabbitMQQueueMonitor.class); + final WorkerSizeObserver observer = mock(WorkerSizeObserver.class); doAnswer(inv -> { latch.countDown(); return null; - }).when(mockMonitor).updateWorkerQueueState(eq(TEST_QUEUE), any()); - provider.putMonitor(TEST_QUEUE, mockMonitor); + }).when(observer).onNumberOfWorkersUpdate(any()); + provider.addObserver(TEST_QUEUE, observer); // Call twice, which should result in only 1 call to updateWorkerQueueState provider.triggerWorkerQueueState(TEST_QUEUE); provider.triggerWorkerQueueState(TEST_QUEUE); latch.await(); - verify(mockMonitor).updateWorkerQueueState(eq(TEST_QUEUE), any()); + verify(mockMonitor).getWorkerQueueState(TEST_QUEUE); + verify(observer).onNumberOfWorkersUpdate(any()); } @Test diff --git a/source/taskmanager/src/test/resources/nl/aerius/taskmanager/mq/queue_aerius.txt b/source/taskmanager/src/test/resources/nl/aerius/taskmanager/mq/queue_aerius.txt new file mode 100644 index 00000000..5f0a8f42 --- /dev/null +++ b/source/taskmanager/src/test/resources/nl/aerius/taskmanager/mq/queue_aerius.txt @@ -0,0 +1 @@ +[{"arguments":{"x-queue-type":"classic"},"auto_delete":false,"consumer_capacity":1.0,"consumer_utilisation":1.0,"consumers":40,"durable":false,"effective_policy_definition":{},"exclusive":false,"internal":false,"internal_owner":false,"memory":13984,"message_bytes":0,"message_bytes_paged_out":0,"message_bytes_persistent":0,"message_bytes_ram":0,"message_bytes_ready":0,"message_bytes_unacknowledged":0,"messages":0,"messages_details":{"rate":0.0},"messages_paged_out":0,"messages_persistent":0,"messages_ram":0,"messages_ready":0,"messages_ready_details":{"rate":0.0},"messages_ready_ram":0,"messages_unacknowledged":15,"messages_unacknowledged_details":{"rate":0.0},"messages_unacknowledged_ram":0,"name":"aerius.workers.chunker","node":"rabbit@","reductions":38260,"reductions_details":{"rate":0.0},"state":"running","storage_version":2,"type":"classic","vhost":"/"},{"arguments":{"x-queue-type":"classic"},"auto_delete":false,"consumer_capacity":1.0,"consumer_utilisation":1.0,"consumers":51,"durable":false,"effective_policy_definition":{},"exclusive":false,"internal":false,"internal_owner":false,"memory":13984,"message_bytes":0,"message_bytes_paged_out":0,"message_bytes_persistent":0,"message_bytes_ram":0,"message_bytes_ready":0,"message_bytes_unacknowledged":0,"messages":10,"messages_details":{"rate":0.0},"messages_paged_out":0,"messages_persistent":0,"messages_ram":0,"messages_ready":20,"messages_ready_details":{"rate":0.0},"messages_ready_ram":0,"messages_unacknowledged":30,"messages_unacknowledged_details":{"rate":0.0},"messages_unacknowledged_ram":0,"name":"aerius.worker.ops","node":"rabbit","reductions":38262,"reductions_details":{"rate":0.0},"state":"running","storage_version":2,"type":"classic","vhost":"/"}] diff --git a/source/taskmanager/src/test/resources/nl/aerius/taskmanager/mq/queue_aerius.worker.ops.txt b/source/taskmanager/src/test/resources/nl/aerius/taskmanager/mq/queue_aerius.worker.ops.txt index c7b2815a..802f7144 100644 --- a/source/taskmanager/src/test/resources/nl/aerius/taskmanager/mq/queue_aerius.worker.ops.txt +++ b/source/taskmanager/src/test/resources/nl/aerius/taskmanager/mq/queue_aerius.worker.ops.txt @@ -1 +1 @@ -{"message_stats":{"ack":59,"ack_details":{"rate":1.8},"deliver":59,"deliver_details":{"rate":2.4},"deliver_get":59,"deliver_get_details":{"rate":2.4},"publish":59,"publish_details":{"rate":1.8}},"messages":3,"messages_details":{"rate":0.0},"messages_ready":0,"messages_ready_details":{"rate":0.0},"messages_unacknowledged":3,"messages_unacknowledged_details":{"rate":0.0},"policy":"","exclusive_consumer_tag":"","consumers":4,"memory":123400,"backing_queue_status":{"q1":0,"q2":0,"delta":["delta","undefined",0,"undefined"],"q3":0,"q4":0,"len":0,"pending_acks":3,"target_ram_count":"infinity","ram_msg_count":0,"ram_ack_count":3,"next_seq_id":59,"persistent_count":3,"avg_ingress_rate":0.4028468488540645,"avg_egress_rate":0.4028468488540645,"avg_ack_ingress_rate":0.4028468488540645,"avg_ack_egress_rate":0.3222774790832516},"status":"running","incoming":[{"stats":{"publish":59,"publish_details":{"rate":1.8}},"exchange":{"name":"","vhost":"/"}}],"deliveries":[{"stats":{"deliver_get":15,"deliver_get_details":{"rate":0.6},"deliver":15,"deliver_details":{"rate":0.6},"ack":15,"ack_details":{"rate":0.4}},"channel_details":{"name":"127.0.0.1:42327 -> 127.0.0.1:5672 (6)","number":6,"connection_name":"127.0.0.1:42327 -> 127.0.0.1:5672","peer_port":42327,"peer_host":"127.0.0.1"}},{"stats":{"deliver_get":15,"deliver_get_details":{"rate":0.8},"deliver":15,"deliver_details":{"rate":0.8},"ack":15,"ack_details":{"rate":0.6}},"channel_details":{"name":"127.0.0.1:42327 -> 127.0.0.1:5672 (5)","number":5,"connection_name":"127.0.0.1:42327 -> 127.0.0.1:5672","peer_port":42327,"peer_host":"127.0.0.1"}},{"stats":{"deliver_get":15,"deliver_get_details":{"rate":0.4},"deliver":15,"deliver_details":{"rate":0.4},"ack":15,"ack_details":{"rate":0.2}},"channel_details":{"name":"127.0.0.1:42327 -> 127.0.0.1:5672 (3)","number":3,"connection_name":"127.0.0.1:42327 -> 127.0.0.1:5672","peer_port":42327,"peer_host":"127.0.0.1"}},{"stats":{"deliver_get":14,"deliver_get_details":{"rate":0.0},"deliver":14,"deliver_details":{"rate":0.0},"ack":14,"ack_details":{"rate":0.0}},"channel_details":{"name":"127.0.0.1:42327 -> 127.0.0.1:5672 (1)","number":1,"connection_name":"127.0.0.1:42327 -> 127.0.0.1:5672","peer_port":42327,"peer_host":"127.0.0.1"}}],"consumer_details":[{"channel_details":{"name":"127.0.0.1:42327 -> 127.0.0.1:5672 (1)","number":1,"connection_name":"127.0.0.1:42327 -> 127.0.0.1:5672","peer_port":42327,"peer_host":"127.0.0.1"},"queue":{"name":"aerius.worker.ops","vhost":"/"},"consumer_tag":"aerius.ops.worker","exclusive":false,"ack_required":true,"arguments":{}},{"channel_details":{"name":"127.0.0.1:42327 -> 127.0.0.1:5672 (3)","number":3,"connection_name":"127.0.0.1:42327 -> 127.0.0.1:5672","peer_port":42327,"peer_host":"127.0.0.1"},"queue":{"name":"aerius.worker.ops","vhost":"/"},"consumer_tag":"aerius.ops.worker","exclusive":false,"ack_required":true,"arguments":{}},{"channel_details":{"name":"127.0.0.1:42327 -> 127.0.0.1:5672 (5)","number":5,"connection_name":"127.0.0.1:42327 -> 127.0.0.1:5672","peer_port":42327,"peer_host":"127.0.0.1"},"queue":{"name":"aerius.worker.ops","vhost":"/"},"consumer_tag":"aerius.ops.worker","exclusive":false,"ack_required":true,"arguments":{}},{"channel_details":{"name":"127.0.0.1:42327 -> 127.0.0.1:5672 (6)","number":6,"connection_name":"127.0.0.1:42327 -> 127.0.0.1:5672","peer_port":42327,"peer_host":"127.0.0.1"},"queue":{"name":"aerius.worker.ops","vhost":"/"},"consumer_tag":"aerius.ops.worker","exclusive":false,"ack_required":true,"arguments":{}}],"name":"aerius.worker.ops","vhost":"/","durable":true,"auto_delete":false,"arguments":{},"node":"rabbit@LEGEND"} \ No newline at end of file +{"message_stats":{"ack":59,"ack_details":{"rate":1.8},"deliver":59,"deliver_details":{"rate":2.4},"deliver_get":59,"deliver_get_details":{"rate":2.4},"publish":59,"publish_details":{"rate":1.8}},"messages":3,"messages_details":{"rate":0.0},"messages_ready":0,"messages_ready_details":{"rate":0.0},"messages_unacknowledged":5,"messages_unacknowledged_details":{"rate":0.0},"policy":"","exclusive_consumer_tag":"","consumers":4,"memory":123400,"backing_queue_status":{"q1":0,"q2":0,"delta":["delta","undefined",0,"undefined"],"q3":0,"q4":0,"len":0,"pending_acks":3,"target_ram_count":"infinity","ram_msg_count":0,"ram_ack_count":3,"next_seq_id":59,"persistent_count":3,"avg_ingress_rate":0.4028468488540645,"avg_egress_rate":0.4028468488540645,"avg_ack_ingress_rate":0.4028468488540645,"avg_ack_egress_rate":0.3222774790832516},"status":"running","incoming":[{"stats":{"publish":59,"publish_details":{"rate":1.8}},"exchange":{"name":"","vhost":"/"}}],"deliveries":[{"stats":{"deliver_get":15,"deliver_get_details":{"rate":0.6},"deliver":15,"deliver_details":{"rate":0.6},"ack":15,"ack_details":{"rate":0.4}},"channel_details":{"name":"127.0.0.1:42327 -> 127.0.0.1:5672 (6)","number":6,"connection_name":"127.0.0.1:42327 -> 127.0.0.1:5672","peer_port":42327,"peer_host":"127.0.0.1"}},{"stats":{"deliver_get":15,"deliver_get_details":{"rate":0.8},"deliver":15,"deliver_details":{"rate":0.8},"ack":15,"ack_details":{"rate":0.6}},"channel_details":{"name":"127.0.0.1:42327 -> 127.0.0.1:5672 (5)","number":5,"connection_name":"127.0.0.1:42327 -> 127.0.0.1:5672","peer_port":42327,"peer_host":"127.0.0.1"}},{"stats":{"deliver_get":15,"deliver_get_details":{"rate":0.4},"deliver":15,"deliver_details":{"rate":0.4},"ack":15,"ack_details":{"rate":0.2}},"channel_details":{"name":"127.0.0.1:42327 -> 127.0.0.1:5672 (3)","number":3,"connection_name":"127.0.0.1:42327 -> 127.0.0.1:5672","peer_port":42327,"peer_host":"127.0.0.1"}},{"stats":{"deliver_get":14,"deliver_get_details":{"rate":0.0},"deliver":14,"deliver_details":{"rate":0.0},"ack":14,"ack_details":{"rate":0.0}},"channel_details":{"name":"127.0.0.1:42327 -> 127.0.0.1:5672 (1)","number":1,"connection_name":"127.0.0.1:42327 -> 127.0.0.1:5672","peer_port":42327,"peer_host":"127.0.0.1"}}],"consumer_details":[{"channel_details":{"name":"127.0.0.1:42327 -> 127.0.0.1:5672 (1)","number":1,"connection_name":"127.0.0.1:42327 -> 127.0.0.1:5672","peer_port":42327,"peer_host":"127.0.0.1"},"queue":{"name":"aerius.worker.ops","vhost":"/"},"consumer_tag":"aerius.ops.worker","exclusive":false,"ack_required":true,"arguments":{}},{"channel_details":{"name":"127.0.0.1:42327 -> 127.0.0.1:5672 (3)","number":3,"connection_name":"127.0.0.1:42327 -> 127.0.0.1:5672","peer_port":42327,"peer_host":"127.0.0.1"},"queue":{"name":"aerius.worker.ops","vhost":"/"},"consumer_tag":"aerius.ops.worker","exclusive":false,"ack_required":true,"arguments":{}},{"channel_details":{"name":"127.0.0.1:42327 -> 127.0.0.1:5672 (5)","number":5,"connection_name":"127.0.0.1:42327 -> 127.0.0.1:5672","peer_port":42327,"peer_host":"127.0.0.1"},"queue":{"name":"aerius.worker.ops","vhost":"/"},"consumer_tag":"aerius.ops.worker","exclusive":false,"ack_required":true,"arguments":{}},{"channel_details":{"name":"127.0.0.1:42327 -> 127.0.0.1:5672 (6)","number":6,"connection_name":"127.0.0.1:42327 -> 127.0.0.1:5672","peer_port":42327,"peer_host":"127.0.0.1"},"queue":{"name":"aerius.worker.ops","vhost":"/"},"consumer_tag":"aerius.ops.worker","exclusive":false,"ack_required":true,"arguments":{}}],"name":"aerius.worker.ops","vhost":"/","durable":true,"auto_delete":false,"arguments":{},"node":"rabbit@LEGEND"} \ No newline at end of file From 326e021eea59e506e80aa21a032a357680cb3dc4 Mon Sep 17 00:00:00 2001 From: Hilbrand Bouwkamp Date: Mon, 21 Sep 2026 13:04:33 +0200 Subject: [PATCH 2/2] Review comment: Removed update channels state triggered by RabbitMQ events, and reduced generic update frequency to 10 seconds 10 seconds update frequency is enough to handle worker changes. Removing the additional event based trigger reduces complexity of the code. --- .../ConnectionConfiguration.java | 2 +- .../adaptor/WorkerSizeProviderProxy.java | 7 -- .../mq/RabbitMQChannelQueueEventsWatcher.java | 119 ------------------ .../taskmanager/mq/RabbitMQQueueMonitor.java | 26 ---- .../mq/RabbitMQWorkerSizeProvider.java | 65 ++-------- ...RabbitMQChannelQueueEventsWatcherTest.java | 108 ---------------- .../mq/RabbitMQQueueMonitorTest.java | 22 +--- .../mq/RabbitMQWorkerSizeProviderTest.java | 13 +- .../mq/queue_aerius.worker.ops.txt | 1 - 9 files changed, 21 insertions(+), 342 deletions(-) delete mode 100644 source/taskmanager/src/main/java/nl/aerius/taskmanager/mq/RabbitMQChannelQueueEventsWatcher.java delete mode 100644 source/taskmanager/src/test/java/nl/aerius/taskmanager/mq/RabbitMQChannelQueueEventsWatcherTest.java delete mode 100644 source/taskmanager/src/test/resources/nl/aerius/taskmanager/mq/queue_aerius.worker.ops.txt diff --git a/source/taskmanager-client/src/main/java/nl/aerius/taskmanager/client/configuration/ConnectionConfiguration.java b/source/taskmanager-client/src/main/java/nl/aerius/taskmanager/client/configuration/ConnectionConfiguration.java index 9d87959e..dd11a057 100644 --- a/source/taskmanager-client/src/main/java/nl/aerius/taskmanager/client/configuration/ConnectionConfiguration.java +++ b/source/taskmanager-client/src/main/java/nl/aerius/taskmanager/client/configuration/ConnectionConfiguration.java @@ -43,7 +43,7 @@ public final class ConnectionConfiguration { /** * Default refresh time in seconds. */ - private static final int DEFAULT_MANAGEMENT_REFRESH_RATE = 20; //seconds + private static final int DEFAULT_MANAGEMENT_REFRESH_RATE = 10; //seconds /** * Default wait time before retrying to connect. diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/adaptor/WorkerSizeProviderProxy.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/adaptor/WorkerSizeProviderProxy.java index 7da65323..3e681f12 100644 --- a/source/taskmanager/src/main/java/nl/aerius/taskmanager/adaptor/WorkerSizeProviderProxy.java +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/adaptor/WorkerSizeProviderProxy.java @@ -39,13 +39,6 @@ public interface WorkerSizeProviderProxy { */ boolean removeObserver(String workerQueueName); - /** - * Triggers to get the worker queue state. - * - * @param queueName name of the worker queue - */ - void triggerWorkerQueueState(final String queueName); - /** * Starts the worker size provider. * diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/mq/RabbitMQChannelQueueEventsWatcher.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/mq/RabbitMQChannelQueueEventsWatcher.java deleted file mode 100644 index 6ac6e999..00000000 --- a/source/taskmanager/src/main/java/nl/aerius/taskmanager/mq/RabbitMQChannelQueueEventsWatcher.java +++ /dev/null @@ -1,119 +0,0 @@ -/* - * 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.mq; - -import java.io.IOException; -import java.util.Map; -import java.util.concurrent.TimeoutException; - -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -import com.rabbitmq.client.AMQP; -import com.rabbitmq.client.Channel; -import com.rabbitmq.client.Connection; -import com.rabbitmq.client.DefaultConsumer; -import com.rabbitmq.client.Envelope; -import com.rabbitmq.client.ShutdownSignalException; - -import nl.aerius.taskmanager.adaptor.WorkerSizeProviderProxy; -import nl.aerius.taskmanager.client.BrokerConnectionFactory; - -/** - * Watches the RabbitMQ event channel for changes in consumers to dynamically monitor the number of workers added or removed. - */ -class RabbitMQChannelQueueEventsWatcher { - - private static final String AMQ_RABBITMQ_EVENT = "amq.rabbitmq.event"; - private static final String CHANNEL_PATTERN = "consumer.*"; - private static final String HEADER_PARAM_QUEUE = "queue"; - - private static final Logger LOG = LoggerFactory.getLogger(RabbitMQChannelQueueEventsWatcher.class); - - private final BrokerConnectionFactory factory; - private final WorkerSizeProviderProxy proxy; - private Channel channel; - - /** - * Constructor. - * - * @param factory connection factory - * @param proxy proxy to get observers for specific worker queues - */ - public RabbitMQChannelQueueEventsWatcher(final BrokerConnectionFactory factory, final WorkerSizeProviderProxy proxy) { - this.factory = factory; - this.proxy = proxy; - } - - /** - * Start the watcher. - * - * @throws IOException Throws IOException in case of communication problems with RabbitMQ - */ - public void start() { - try { - final Connection c = factory.getConnection(); - c.addShutdownListener(this::handleShutdownSignal); - channel = c.createChannel(); - final String q = channel.queueDeclare().getQueue(); - channel.queueBind(q, AMQ_RABBITMQ_EVENT, CHANNEL_PATTERN); - channel.basicConsume(q, true, createConsumer()); - } catch (final IOException e) { - LOG.error("Failed to bind to RabbitMQ event queue. No Queue Event watch not available. Message: {}", e.getMessage()); - } - } - - private void handleShutdownSignal(final ShutdownSignalException sse) { - if (sse != null && sse.isInitiatedByApplication()) { - return; - } - LOG.debug("Channel RabbitMQChannelQueueEventsWatcher was shut down."); - // restart - try { - if (channel != null) { - channel.abort(); - } - } catch (final IOException e) { - // Eat error when closing channel. - } - start(); - LOG.info("Restarted RabbitMQChannelQueueEventsWatcher"); - } - - private DefaultConsumer createConsumer() { - return new DefaultConsumer(channel) { - @Override - public void handleDelivery(final String consumerTag, final Envelope envelope, final AMQP.BasicProperties properties, final byte[] body) - throws IOException { - final Map headers = properties.getHeaders(); - final Object queue = headers.get(HEADER_PARAM_QUEUE); - final String queueName = queue == null ? null : queue.toString(); - - LOG.trace("Event: {} - queue: {}", envelope.getRoutingKey(), queueName); - proxy.triggerWorkerQueueState(queueName); - } - }; - } - - public void shutdown() { - try { - channel.close(); - } catch (final IOException | TimeoutException e) { - LOG.trace("Channel watcher shutdown failed", e); - } - } -} diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/mq/RabbitMQQueueMonitor.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/mq/RabbitMQQueueMonitor.java index b2e444d1..0f29a543 100644 --- a/source/taskmanager/src/main/java/nl/aerius/taskmanager/mq/RabbitMQQueueMonitor.java +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/mq/RabbitMQQueueMonitor.java @@ -107,32 +107,6 @@ public void close() { } } - /** - * Retrieves the queue status for the given queue from the RabbitMQ admin api. - * - * @param queueName name of the queue to get the statuss - * @return Status of the queue - */ - public RabbitMQQueueStatus getWorkerQueueState(final String queueName) { - // Use RabbitMQ HTTP-API. - // URL: [host]:[port]/api/queues/[virtualHost]/[QueueName] - final String virtualHost = configuration.getBrokerVirtualHost().replace("/", "%2f"); - final String apiPath = String.format("/api/queues/%s/%s", virtualHost, queueName); - - try { - final JsonNode jsonNode = getJsonResultFromApi(apiPath); - - if (jsonNode == null) { - LOG.error("Queue configuration from RabbitMQ admin json get call returned null."); - } else { - return getQueueStatus(jsonNode); - } - } catch (final URISyntaxException | IOException e) { - LOG.info("Error getting RabbitMQ status from admin api: {}", e.getMessage()); - } - return null; - } - /** * Retrieves the status for all queues from the RabbitMQ admin api. */ diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/mq/RabbitMQWorkerSizeProvider.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/mq/RabbitMQWorkerSizeProvider.java index ddccbe04..d49238eb 100644 --- a/source/taskmanager/src/main/java/nl/aerius/taskmanager/mq/RabbitMQWorkerSizeProvider.java +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/mq/RabbitMQWorkerSizeProvider.java @@ -23,9 +23,7 @@ import java.util.Map; import java.util.Optional; import java.util.concurrent.ScheduledExecutorService; -import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; -import java.util.function.Supplier; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -48,31 +46,17 @@ public class RabbitMQWorkerSizeProvider implements WorkerSizeProviderProxy { /** * Delay first read from the RabittMQ admin to give the taskmanager some time to start up and register all observers. */ - private static final int INITIAL_DELAY_SECONDS = 10; - /** - * The minimum time before the RabbitMQ management api is fetched again to get an update on the queue state. - */ - private static final long DELAY_BEFORE_UPDATE_TIME_SECONDS = 15; + private static final int INITIAL_DELAY_SECONDS = 5; private final ScheduledExecutorService executorService; - private final RabbitMQChannelQueueEventsWatcher channelQueueEventsWatcher; private final RabbitMQWorkerEventProducer eventProducer; /** * The time in seconds between each scheduled update. */ private final long refreshRateSeconds; - /** - * The time delay in seconds before the update call is made. - */ - private final long refreshDelayBeforeUpdateSeconds; - - private final Object sync = new Object(); - private final Map> lastRuns = new HashMap<>(); private final Map observers = new HashMap<>(); private final RabbitMQQueueMonitor monitor; - // Map to keep track of the last known states of the queues as retrieved form the RabbitMQ admin API. - private Map lastKnownQueueStates = new HashMap<>(); private boolean running; public RabbitMQWorkerSizeProvider(final ScheduledExecutorService executorService, final BrokerConnectionFactory factory) { @@ -83,10 +67,8 @@ public RabbitMQWorkerSizeProvider(final ScheduledExecutorService executorService final RabbitMQQueueMonitor monitor) { this.executorService = executorService; this.monitor = monitor; - channelQueueEventsWatcher = new RabbitMQChannelQueueEventsWatcher(factory, this); refreshRateSeconds = factory.getConnectionConfiguration().getBrokerManagementRefreshRate(); eventProducer = new RabbitMQWorkerEventProducer(executorService, factory); - refreshDelayBeforeUpdateSeconds = Math.min(refreshRateSeconds / 2, DELAY_BEFORE_UPDATE_TIME_SECONDS); } @Override @@ -105,7 +87,6 @@ public boolean removeObserver(final String queueName) { @Override public void start() throws IOException { - channelQueueEventsWatcher.start(); eventProducer.start(); if (refreshRateSeconds > 0) { running = true; @@ -115,58 +96,28 @@ public void start() throws IOException { @Override public void shutdown() { + running = false; for (final String key : new ArrayList<>(observers.keySet())) { removeObserver(key); } eventProducer.shutdown(); - channelQueueEventsWatcher.shutdown(); } private void updateWorkerQueueState() { if (running) { try { - lastKnownQueueStates = new HashMap<>(monitor.getWorkerQueueStates()); - observers.forEach((q, v) -> scheduledUpdateWorkerQueueState(q, () -> lastKnownQueueStates.get(q))); + final Map queueStates = new HashMap<>(monitor.getWorkerQueueStates()); + observers.forEach((q, v) -> updateWorkerQueueState(q, queueStates.get(q))); } catch (final RuntimeException e) { LOG.error("Runtime error during updateWorkerQueueState", e); } } } - @Override - public void triggerWorkerQueueState(final String queueName) { - if (!observers.containsKey(queueName)) { - return; - } - scheduledUpdateWorkerQueueState(queueName, () -> monitor.getWorkerQueueState(queueName)); - } - - private void scheduledUpdateWorkerQueueState(final String queueName, final Supplier statusSupplier) { - // This uses a delayed update. It schedules a task to run in x-seconds. - // If a new update is received before the schedule has run it will cancel the current schedule and reschedule. - // This is mainly for when multiple events are triggered to not trigger a call for every event, - // and also to manage the events trigger in combination with the scheduled process. - synchronized (sync) { - Optional.ofNullable(lastRuns.get(queueName)).ifPresent(f -> f.cancel(false)); - final Runnable updateTask = () -> updateWorkerQueueState(queueName, statusSupplier); - - lastRuns.put(queueName, executorService.schedule(updateTask, refreshDelayBeforeUpdateSeconds, TimeUnit.SECONDS)); - } - } - - private void updateWorkerQueueState(final String queueName, final Supplier statusSupplier) { - synchronized (sync) { - try { - Optional.ofNullable(statusSupplier.get()) - .ifPresent(s -> { - lastKnownQueueStates.put(queueName, s); - lastRuns.remove(queueName); - Optional.ofNullable(observers.get(queueName)).ifPresent(observer -> observer.onNumberOfWorkersUpdate(s)); - }); - } catch (final RuntimeException e) { - LOG.error("RuntimeException during updateWorkerQueueState", e); - } - } + private void updateWorkerQueueState(final String queueName, final RabbitMQQueueStatus queueStatus) { + Optional.ofNullable(queueStatus).ifPresent(s -> { + Optional.ofNullable(observers.get(queueName)).ifPresent(observer -> observer.onNumberOfWorkersUpdate(queueStatus)); + }); } private static class WorkerSizeObserverComposite implements WorkerSizeObserver { diff --git a/source/taskmanager/src/test/java/nl/aerius/taskmanager/mq/RabbitMQChannelQueueEventsWatcherTest.java b/source/taskmanager/src/test/java/nl/aerius/taskmanager/mq/RabbitMQChannelQueueEventsWatcherTest.java deleted file mode 100644 index 1d6119ba..00000000 --- a/source/taskmanager/src/test/java/nl/aerius/taskmanager/mq/RabbitMQChannelQueueEventsWatcherTest.java +++ /dev/null @@ -1,108 +0,0 @@ -/* - * 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.mq; - -import static org.mockito.ArgumentMatchers.any; -import static org.mockito.ArgumentMatchers.eq; -import static org.mockito.Mockito.doAnswer; -import static org.mockito.Mockito.doReturn; -import static org.mockito.Mockito.times; -import static org.mockito.Mockito.verify; - -import java.io.IOException; -import java.util.HashMap; -import java.util.Map; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; - -import org.junit.jupiter.api.AfterAll; -import org.junit.jupiter.api.BeforeAll; -import org.junit.jupiter.api.BeforeEach; -import org.junit.jupiter.api.Test; -import org.mockito.Mockito; - -import com.rabbitmq.client.AMQP.BasicProperties; -import com.rabbitmq.client.AMQP.Queue; -import com.rabbitmq.client.Channel; -import com.rabbitmq.client.Connection; -import com.rabbitmq.client.Consumer; -import com.rabbitmq.client.Envelope; - -import nl.aerius.taskmanager.adaptor.WorkerSizeProviderProxy; -import nl.aerius.taskmanager.client.BrokerConnectionFactory; - -/** - * Test class for {@link RabbitMQChannelQueueEventsWatcher}. - */ -class RabbitMQChannelQueueEventsWatcherTest { - private static final String TEST_QUEUENAME = "test"; - private static final String HEADER_PARAM_QUEUE = "queue"; - - private static ExecutorService executor; - - private RabbitMQChannelQueueEventsWatcher watcher; - private Channel mockChannel; - - private WorkerSizeProviderProxy proxy; - - @BeforeAll - static void setupClass() { - executor = Executors.newSingleThreadExecutor(); - } - - @AfterAll - static void afterClass() { - executor.shutdown(); - } - - @BeforeEach - void setUp() throws Exception { - final Connection mockConnection = Mockito.mock(Connection.class); - mockChannel = Mockito.mock(Channel.class); - doReturn(mockChannel).when(mockConnection).createChannel(); - final Queue.DeclareOk mockDeclareOk = Mockito.mock(Queue.DeclareOk.class); - doReturn(mockDeclareOk).when(mockChannel).queueDeclare(); - proxy = Mockito.mock(WorkerSizeProviderProxy.class); - watcher = new RabbitMQChannelQueueEventsWatcher(new BrokerConnectionFactory(executor) { - @Override - protected Connection createNewConnection() throws IOException { - return mockConnection; - } - }, proxy); - } - - @Test - void testReceiving() throws IOException { - assertDeltaCheck("consumer.created"); - verify(proxy).triggerWorkerQueueState(TEST_QUEUENAME); - assertDeltaCheck("consumer.removed"); - verify(proxy, times(2)).triggerWorkerQueueState(TEST_QUEUENAME); - } - - private void assertDeltaCheck(final String event) throws IOException { - doAnswer(i -> { - final Envelope envelope = Mockito.mock(Envelope.class); - doReturn(event).when(envelope).getRoutingKey(); - final Map headers = new HashMap<>(); - - headers.put(HEADER_PARAM_QUEUE, TEST_QUEUENAME); - ((Consumer) i.getArgument(2)).handleDelivery(null, envelope, new BasicProperties().builder().headers(headers).build(), null); - return null; - }).when(mockChannel).basicConsume(any(), eq(true), any()); - watcher.start(); - } -} diff --git a/source/taskmanager/src/test/java/nl/aerius/taskmanager/mq/RabbitMQQueueMonitorTest.java b/source/taskmanager/src/test/java/nl/aerius/taskmanager/mq/RabbitMQQueueMonitorTest.java index 04ea5ca4..81223e4a 100644 --- a/source/taskmanager/src/test/java/nl/aerius/taskmanager/mq/RabbitMQQueueMonitorTest.java +++ b/source/taskmanager/src/test/java/nl/aerius/taskmanager/mq/RabbitMQQueueMonitorTest.java @@ -21,7 +21,6 @@ import java.io.IOException; import java.io.InputStream; import java.io.InputStreamReader; -import java.util.function.Function; import org.junit.jupiter.api.Test; @@ -36,39 +35,30 @@ */ class RabbitMQQueueMonitorTest { + private static final String FILENAME = "queue_aerius.txt"; private static final String DUMMY = "dummy"; private static final String QUEUENAME = "aerius.worker.ops"; private final ObjectMapper objectMapper = new ObjectMapper(); - @Test - void testGetWorkerQueueState() { - assertRabbitMQQueueMonitor("queue_aerius.worker.ops.txt", 4, 3, 5, rpm -> rpm.getWorkerQueueState(DUMMY)); - } - @Test void testGetWorkerQueueStates() { - assertRabbitMQQueueMonitor("queue_aerius.txt", 51, 10, 30, rpm -> rpm.getWorkerQueueStates().get(QUEUENAME)); - } - - private void assertRabbitMQQueueMonitor(final String filename, final int expectedConsumers, final int expectedMessages, - final int expectedUnacknowledged, final Function collector) { final ConnectionConfiguration configuration = ConnectionConfiguration.builder() .brokerHost(DUMMY).brokerPort(0).brokerUsername(DUMMY).brokerPassword(DUMMY).build(); final RabbitMQQueueMonitor rpm = new RabbitMQQueueMonitor(configuration) { @Override protected JsonNode getJsonResultFromApi(final String apiPath) throws IOException { - try (final InputStream fr = getClass().getResourceAsStream(filename); + try (final InputStream fr = getClass().getResourceAsStream(FILENAME); final InputStreamReader is = new InputStreamReader(fr)) { return objectMapper.readTree(is); } } }; try { - final RabbitMQQueueStatus status = collector.apply(rpm); + final RabbitMQQueueStatus status = rpm.getWorkerQueueStates().get(QUEUENAME); - assertEquals(expectedConsumers, status.consumers(), "Number of workers"); - assertEquals(expectedMessages, status.messages(), "Number of messages"); - assertEquals(expectedUnacknowledged, status.unacknowledged(), "Number of unacknowledged"); + assertEquals(51, status.consumers(), "Number of workers"); + assertEquals(10, status.messages(), "Number of messages"); + assertEquals(30, status.unacknowledged(), "Number of unacknowledged"); } finally { rpm.shutdown(); } diff --git a/source/taskmanager/src/test/java/nl/aerius/taskmanager/mq/RabbitMQWorkerSizeProviderTest.java b/source/taskmanager/src/test/java/nl/aerius/taskmanager/mq/RabbitMQWorkerSizeProviderTest.java index 4579c67d..e8ae5e06 100644 --- a/source/taskmanager/src/test/java/nl/aerius/taskmanager/mq/RabbitMQWorkerSizeProviderTest.java +++ b/source/taskmanager/src/test/java/nl/aerius/taskmanager/mq/RabbitMQWorkerSizeProviderTest.java @@ -24,6 +24,7 @@ import static org.mockito.Mockito.verify; import java.io.IOException; +import java.util.Map; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; @@ -52,15 +53,15 @@ class RabbitMQWorkerSizeProviderTest extends AbstractRabbitMQTest { @Override @BeforeEach void setUp() throws Exception { - brokerManagementRefreshRate = 5; + brokerManagementRefreshRate = 1; super.setUp(); provider = new RabbitMQWorkerSizeProvider(executor, factory, mockMonitor); } @Test @Timeout(value = 10, unit = TimeUnit.SECONDS) - void testTriggerWorkerQueueState() throws InterruptedException { - doReturn(new RabbitMQQueueStatus(1, 2, 3)).when(mockMonitor).getWorkerQueueState(TEST_QUEUE); + void testTriggerWorkerQueueState() throws InterruptedException, IOException { + doReturn(Map.of(TEST_QUEUE, new RabbitMQQueueStatus(1, 2, 3))).when(mockMonitor).getWorkerQueueStates(); final CountDownLatch latch = new CountDownLatch(1); final WorkerSizeObserver observer = mock(WorkerSizeObserver.class); @@ -69,11 +70,9 @@ void testTriggerWorkerQueueState() throws InterruptedException { return null; }).when(observer).onNumberOfWorkersUpdate(any()); provider.addObserver(TEST_QUEUE, observer); - // Call twice, which should result in only 1 call to updateWorkerQueueState - provider.triggerWorkerQueueState(TEST_QUEUE); - provider.triggerWorkerQueueState(TEST_QUEUE); + provider.start(); latch.await(); - verify(mockMonitor).getWorkerQueueState(TEST_QUEUE); + verify(mockMonitor).getWorkerQueueStates(); verify(observer).onNumberOfWorkersUpdate(any()); } diff --git a/source/taskmanager/src/test/resources/nl/aerius/taskmanager/mq/queue_aerius.worker.ops.txt b/source/taskmanager/src/test/resources/nl/aerius/taskmanager/mq/queue_aerius.worker.ops.txt deleted file mode 100644 index 802f7144..00000000 --- a/source/taskmanager/src/test/resources/nl/aerius/taskmanager/mq/queue_aerius.worker.ops.txt +++ /dev/null @@ -1 +0,0 @@ -{"message_stats":{"ack":59,"ack_details":{"rate":1.8},"deliver":59,"deliver_details":{"rate":2.4},"deliver_get":59,"deliver_get_details":{"rate":2.4},"publish":59,"publish_details":{"rate":1.8}},"messages":3,"messages_details":{"rate":0.0},"messages_ready":0,"messages_ready_details":{"rate":0.0},"messages_unacknowledged":5,"messages_unacknowledged_details":{"rate":0.0},"policy":"","exclusive_consumer_tag":"","consumers":4,"memory":123400,"backing_queue_status":{"q1":0,"q2":0,"delta":["delta","undefined",0,"undefined"],"q3":0,"q4":0,"len":0,"pending_acks":3,"target_ram_count":"infinity","ram_msg_count":0,"ram_ack_count":3,"next_seq_id":59,"persistent_count":3,"avg_ingress_rate":0.4028468488540645,"avg_egress_rate":0.4028468488540645,"avg_ack_ingress_rate":0.4028468488540645,"avg_ack_egress_rate":0.3222774790832516},"status":"running","incoming":[{"stats":{"publish":59,"publish_details":{"rate":1.8}},"exchange":{"name":"","vhost":"/"}}],"deliveries":[{"stats":{"deliver_get":15,"deliver_get_details":{"rate":0.6},"deliver":15,"deliver_details":{"rate":0.6},"ack":15,"ack_details":{"rate":0.4}},"channel_details":{"name":"127.0.0.1:42327 -> 127.0.0.1:5672 (6)","number":6,"connection_name":"127.0.0.1:42327 -> 127.0.0.1:5672","peer_port":42327,"peer_host":"127.0.0.1"}},{"stats":{"deliver_get":15,"deliver_get_details":{"rate":0.8},"deliver":15,"deliver_details":{"rate":0.8},"ack":15,"ack_details":{"rate":0.6}},"channel_details":{"name":"127.0.0.1:42327 -> 127.0.0.1:5672 (5)","number":5,"connection_name":"127.0.0.1:42327 -> 127.0.0.1:5672","peer_port":42327,"peer_host":"127.0.0.1"}},{"stats":{"deliver_get":15,"deliver_get_details":{"rate":0.4},"deliver":15,"deliver_details":{"rate":0.4},"ack":15,"ack_details":{"rate":0.2}},"channel_details":{"name":"127.0.0.1:42327 -> 127.0.0.1:5672 (3)","number":3,"connection_name":"127.0.0.1:42327 -> 127.0.0.1:5672","peer_port":42327,"peer_host":"127.0.0.1"}},{"stats":{"deliver_get":14,"deliver_get_details":{"rate":0.0},"deliver":14,"deliver_details":{"rate":0.0},"ack":14,"ack_details":{"rate":0.0}},"channel_details":{"name":"127.0.0.1:42327 -> 127.0.0.1:5672 (1)","number":1,"connection_name":"127.0.0.1:42327 -> 127.0.0.1:5672","peer_port":42327,"peer_host":"127.0.0.1"}}],"consumer_details":[{"channel_details":{"name":"127.0.0.1:42327 -> 127.0.0.1:5672 (1)","number":1,"connection_name":"127.0.0.1:42327 -> 127.0.0.1:5672","peer_port":42327,"peer_host":"127.0.0.1"},"queue":{"name":"aerius.worker.ops","vhost":"/"},"consumer_tag":"aerius.ops.worker","exclusive":false,"ack_required":true,"arguments":{}},{"channel_details":{"name":"127.0.0.1:42327 -> 127.0.0.1:5672 (3)","number":3,"connection_name":"127.0.0.1:42327 -> 127.0.0.1:5672","peer_port":42327,"peer_host":"127.0.0.1"},"queue":{"name":"aerius.worker.ops","vhost":"/"},"consumer_tag":"aerius.ops.worker","exclusive":false,"ack_required":true,"arguments":{}},{"channel_details":{"name":"127.0.0.1:42327 -> 127.0.0.1:5672 (5)","number":5,"connection_name":"127.0.0.1:42327 -> 127.0.0.1:5672","peer_port":42327,"peer_host":"127.0.0.1"},"queue":{"name":"aerius.worker.ops","vhost":"/"},"consumer_tag":"aerius.ops.worker","exclusive":false,"ack_required":true,"arguments":{}},{"channel_details":{"name":"127.0.0.1:42327 -> 127.0.0.1:5672 (6)","number":6,"connection_name":"127.0.0.1:42327 -> 127.0.0.1:5672","peer_port":42327,"peer_host":"127.0.0.1"},"queue":{"name":"aerius.worker.ops","vhost":"/"},"consumer_tag":"aerius.ops.worker","exclusive":false,"ack_required":true,"arguments":{}}],"name":"aerius.worker.ops","vhost":"/","durable":true,"auto_delete":false,"arguments":{},"node":"rabbit@LEGEND"} \ No newline at end of file