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..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 = 60; //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/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/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/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/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 6cc14946..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 @@ -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,40 @@ public void close() { } } - public void updateWorkerQueueState(final String queueName, final WorkerSizeObserver observer) { - // 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); - + /** + * Retrieves the status for all queues from the RabbitMQ admin api. + */ + public Map getWorkerQueueStates() { try { - final JsonNode jsonObject = getJsonResultFromApi(apiPath); + final JsonNode jsonObject = getJsonResultFromApi("/api/queues"); if (jsonObject == 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"); + } 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..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,7 +23,6 @@ import java.util.Map; import java.util.Optional; import java.util.concurrent.ScheduledExecutorService; -import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; import org.slf4j.Logger; @@ -33,6 +32,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. @@ -46,83 +46,47 @@ 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 BrokerConnectionFactory factory; - 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 Map monitors = new HashMap<>(); - + private final RabbitMQQueueMonitor monitor; 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; - channelQueueEventsWatcher = new RabbitMQChannelQueueEventsWatcher(factory, this); + this.monitor = monitor; refreshRateSeconds = factory.getConnectionConfiguration().getBrokerManagementRefreshRate(); eventProducer = new RabbitMQWorkerEventProducer(executorService, factory); - refreshDelayBeforeUpdateSeconds = Math.min(refreshRateSeconds / 2, DELAY_BEFORE_UPDATE_TIME_SECONDS); } @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 public void start() throws IOException { - channelQueueEventsWatcher.start(); eventProducer.start(); if (refreshRateSeconds > 0) { running = true; @@ -132,42 +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 { - monitors.forEach((k, v) -> triggerWorkerQueueState(k)); + 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) { - // 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); - - lastRuns.put(queueName, executorService.schedule(updateTask, refreshDelayBeforeUpdateSeconds, TimeUnit.SECONDS)); - } - } - - private void updateWorkerQueueState(final String queueName) { - synchronized (sync) { - Optional.ofNullable(monitors.get(queueName)).ifPresent(m -> m.updateWorkerQueueState(queueName, observers.get(queueName))); - lastRuns.remove(queueName); - } + 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 { @@ -178,10 +128,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/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 c71fe0fe..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,42 +21,44 @@ import java.io.IOException; import java.io.InputStream; import java.io.InputStreamReader; -import java.util.concurrent.atomic.AtomicInteger; 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}. */ 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() { + void testGetWorkerQueueStates() { 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 = rpm.getWorkerQueueStates().get(QUEUENAME); + + 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 780104e1..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 @@ -18,54 +18,62 @@ 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; import java.io.IOException; +import java.util.Map; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; 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 @BeforeEach void setUp() throws Exception { - brokerManagementRefreshRate = 5; + brokerManagementRefreshRate = 1; super.setUp(); - provider = new RabbitMQWorkerSizeProvider(executor, factory); + provider = new RabbitMQWorkerSizeProvider(executor, factory, mockMonitor); } @Test @Timeout(value = 10, unit = TimeUnit.SECONDS) - void testTriggerWorkerQueueState() throws InterruptedException { + 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 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); - // Call twice, which should result in only 1 call to updateWorkerQueueState - provider.triggerWorkerQueueState(TEST_QUEUE); - provider.triggerWorkerQueueState(TEST_QUEUE); + }).when(observer).onNumberOfWorkersUpdate(any()); + provider.addObserver(TEST_QUEUE, observer); + provider.start(); latch.await(); - verify(mockMonitor).updateWorkerQueueState(eq(TEST_QUEUE), any()); + verify(mockMonitor).getWorkerQueueStates(); + 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 deleted file mode 100644 index c7b2815a..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":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