diff --git a/doc/telemetry.md b/doc/telemetry.md index 6c4e29c..9f2d241 100644 --- a/doc/telemetry.md +++ b/doc/telemetry.md @@ -82,9 +82,9 @@ Therefore `.limit` gives the total number of workers available. And `.usage` gives the metrics about how the workers are used. The usage metrics have an attribute `state` that identifies if the metric is for `used` or `free` amount of workers. Additional `aer.rabbitmq.worker.usage` records a third state `waiting`. -And `aer.taskmanager.client.queue.usage` has states `used` and `waiting` that provides information on tasks send to the TaskManager. -This shows the number of tasks on the worker queue that are not yet picked up by any worker. - +The `aer.taskmanager.queue` metrics has state attributes `used` and `waiting` that provides information on tasks picked up by the TaskManager. +Waiting tasks are only interesting for queues that use `eagerFetch`, because in that configuration all tasks are always picked up from the queue. +To know the number of waiting tasks for other queues use the metric `aer.rabbitmq.client.queue`. | metric name | type | description | |-------------------------------------------------------|-----------|----------------------------------------------------------------------------| @@ -97,13 +97,14 @@ This shows the number of tasks on the worker queue that are not yet picked up by | `aer.taskmanager.worker.usage`2 | gauge | Weighted usage of workers based on tasks send to workers. | | `aer.taskmanager.workerpool.worker.limit`1 | gauge | TaskManager internal total number of workers. | | `aer.taskmanager.workerpool.worker.usage`2 | gauge | TaskManager internal usage of workers. | -| `aer.taskmanager.client.queue.usage`2 | gauge | TaskManager internal metrics on client queue usage. | +| `aer.taskmanager.queue`4 | gauge | TaskManager internal metrics on client queue usage. | | `aer.rabbitmq.worker.limit`1 | gauge | Total number of workers available as reported by RabbitMQ | -| `aer.rabbitmq.worker.usage`2 | gauge | Usage o the workers based on the messages on the RabbitMQ worker queue. | +| `aer.rabbitmq.worker.usage`2 | gauge | Usage of the workers based on the messages on the RabbitMQ worker queue. | +| `aer.rabbitmq.client.queue`3 | gauge | Number of messages on per client queue as reported by RabbitMQ. | | `aer.taskmanager.dispatched`1 | histogram | The number of tasks dispatched. | | `aer.taskmanager.dispatched.wait`1 | histogram | The average wait time of tasks dispatched. | -| `aer.taskmanager.dispatched.queue`3 | histogram | The number of tasks dispatched per client queue. | -| `aer.taskmanager.dispatched.queue.wait`3 | histogram | The average wait time of tasks dispatched per client queue. | +| `aer.taskmanager.dispatched.queue`4 | histogram | The number of tasks dispatched per client queue. | +| `aer.taskmanager.dispatched.queue.wait`4 | histogram | The average wait time of tasks dispatched per client queue. | Basically there are 3 metric groups that report similar information. First the `aer.taskmanager.worker.*` metrics are a weighted value based on when tasks are send to the workers. @@ -120,9 +121,10 @@ The workers have different attributes to distinguish specific metrics. * 1 have attribute `worker_type`. * 2 have attribute `worker_type` and `state`. `state` can have the value `used`, `free` or `waiting`. * 3 have attribute `worker_type` and `queue_name`. +* 4 have attribute `worker_type`, `queue_name` and `state`. `state` can have the value `used`, `free` or `waiting`. `worker_type` is the type of worker, e.g. `ops`. -`queue_name` is the originating queue the task initially was put on, e.g. `...calculator_ui_small`. +`queue_name` is the originating queue the task initially was put on, e.g. `calculator_ui_small`. > [!NOTE] > The metrics in the TaskManager operate on a time frame of 1 minute. diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/TaskManager.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/TaskManager.java index 1e8ab79..c7e7f09 100644 --- a/source/taskmanager/src/main/java/nl/aerius/taskmanager/TaskManager.java +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/TaskManager.java @@ -42,6 +42,7 @@ import nl.aerius.taskmanager.domain.TaskSchedule; import nl.aerius.taskmanager.metrics.OpenTelemetryMetrics; import nl.aerius.taskmanager.metrics.PerformanceMetricsReporter; +import nl.aerius.taskmanager.metrics.RabbitMQClientQueueReporter; import nl.aerius.taskmanager.metrics.RabbitMQUsageMetricsProvider; import nl.aerius.taskmanager.metrics.TaskManagerMetricsRegister; import nl.aerius.taskmanager.metrics.TaskManagerUsageMetricsProvider; @@ -167,6 +168,7 @@ public TaskScheduleBucket(final QueueConfig queueConfig) throws InterruptedExcep workerSizeObserverProxy.addObserver(workerQueueName, taskManagerMetricsRegister); workerSizeObserverProxy.addObserver(workerQueueName, workerPool); workerSizeObserverProxy.addObserver(workerQueueName, watchDog); + workerSizeObserverProxy.addClientObserver(workerQueueName, new RabbitMQClientQueueReporter(OpenTelemetryMetrics.METER, workerQueueName)); // startup Guard should be the last observer added as it will unlock the task dispatcher workerSizeObserverProxy.addObserver(workerQueueName, startupGuard); diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/adaptor/ClientQueueObserver.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/adaptor/ClientQueueObserver.java new file mode 100644 index 0000000..7225865 --- /dev/null +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/adaptor/ClientQueueObserver.java @@ -0,0 +1,41 @@ +/* + * 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.adaptor; + +import nl.aerius.taskmanager.domain.RabbitMQQueueStatus; + +/** + * Observer that listens to client queue updates. + */ +public interface ClientQueueObserver { + + /** + * Update for the given client queue. + * + * @param clientQueueName name of the client queue + * @param status queue metrics + */ + void onClientQueueUpdate(String clientQueueName, RabbitMQQueueStatus status); + + /** + * Returns true if this observer should receive updates for the given client queue. + * + * @param clientQueueName client queue to check + * @return true if should receive updates + */ + boolean filter(String clientQueueName); +} 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 3e681f1..7ce1aef 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 @@ -31,6 +31,14 @@ public interface WorkerSizeProviderProxy { */ void addObserver(String workerQueueName, WorkerSizeObserver workerSizeObserver); + /** + * Add a ClientQueueObserver. + * + * @param workerQueueName name of the worker queue the client queues are related too + * @param observer observer to be informed + */ + public void addClientObserver(final String workerQueueName, final ClientQueueObserver observer); + /** * Removes the observer for the worker queue. * diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/OpenTelemetryMetrics.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/OpenTelemetryMetrics.java index e26a662..c394a83 100644 --- a/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/OpenTelemetryMetrics.java +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/OpenTelemetryMetrics.java @@ -66,7 +66,7 @@ public static Attributes queueAttributes(final String workerQueueName, final Str .build(); } - private static String onlyLastPart(final String name) { + public static String onlyLastPart(final String name) { final int lastIndex = name.lastIndexOf('.'); return (lastIndex < 0 ? name : name.substring(lastIndex + 1)).toLowerCase(Locale.ROOT); diff --git a/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/RabbitMQClientQueueReporter.java b/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/RabbitMQClientQueueReporter.java new file mode 100644 index 0000000..acba9b6 --- /dev/null +++ b/source/taskmanager/src/main/java/nl/aerius/taskmanager/metrics/RabbitMQClientQueueReporter.java @@ -0,0 +1,61 @@ +/* + * 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.metrics; + +import java.util.HashMap; +import java.util.Map; + +import io.opentelemetry.api.metrics.Meter; + +import nl.aerius.taskmanager.adaptor.ClientQueueObserver; +import nl.aerius.taskmanager.domain.RabbitMQQueueStatus; + +/** + * Reports RabbitMQ client queue number of messages on the queue. + * + * The report registers the metrics using the last part of the worker queue name. Typical worker queue name is 'aerius.worker.some_name'. + * The last part of the worker queue is used as attribute when reporting the metric. + * All client queues related to this worker queue will be reported on. + * The specific client metric is reported with the last part of the client queue name as attribute value. + * The metrics are reported as gauge metric to the metric 'aer.rabbitmq.client_queue'. + */ +public class RabbitMQClientQueueReporter implements ClientQueueObserver { + + private static final String CLIENT_QUEUE_METRIC = "aer.rabbitmq.client.queue"; + private final String workerQueue; + private final UsageMetricsReporter clientQueueReporter; + private final Map queueCounters = new HashMap<>(); + + public RabbitMQClientQueueReporter(final Meter meter, final String workerQueueName) { + this.workerQueue = OpenTelemetryMetrics.onlyLastPart(workerQueueName); + clientQueueReporter = new UsageMetricsReporter(meter, CLIENT_QUEUE_METRIC, "Number of messages on the RabbitMQ client queues"); + } + + @Override + public void onClientQueueUpdate(final String clientQueueName, final RabbitMQQueueStatus value) { + if (!queueCounters.containsKey(clientQueueName)) { + clientQueueReporter.addMetrics(workerQueue, () -> queueCounters.get(clientQueueName), + OpenTelemetryMetrics.queueAttributes(workerQueue, clientQueueName)); + } + queueCounters.put(clientQueueName, Double.valueOf(value.messages())); + } + + @Override + public boolean filter(final String queueName) { + return queueName.contains(workerQueue) && !queueName.contains("worker"); + } +} 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 6614294..f9025ca 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 @@ -47,8 +47,8 @@ public UsageMetricsReporter(final Meter meter, final String metricName, final St } private void recordMetrics(final ObservableDoubleMeasurement measurement) { - for (final Entry> workerMetrics : metricsMap.entrySet()) { - for (final UsageMetric metric : workerMetrics.getValue()) { + for (final Entry> entry : metricsMap.entrySet()) { + for (final UsageMetric metric : entry.getValue()) { measurement.record(metric.metricSupplier().getAsDouble(), metric.attributes()); } 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 0f29a54..e1a1dea 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 @@ -110,25 +110,25 @@ public void close() { /** * Retrieves the status for all queues from the RabbitMQ admin api. */ - public Map getWorkerQueueStates() { + public Map getQueueStates() { + final Map queueStates = new HashMap<>(); + 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<>(); + if (jsonObject instanceof final ArrayNode array) { for (int i = 0; i < array.size(); i++) { final JsonNode jsonNode = array.get(i); queueStates.put(getJsonString(jsonNode, "name"), getQueueStatus(jsonNode)); } - return queueStates; + } else { + LOG.error("Queue configuration from RabbitMQ admin json returned an unexpected value."); } } catch (final URISyntaxException | IOException e) { LOG.info("Error getting RabbitMQ status from admin api: {}", e.getMessage()); } - return new HashMap<>(); + return queueStates; } private static RabbitMQQueueStatus getQueueStatus(final JsonNode jsonNode) { 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 d49238e..351a164 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 @@ -28,6 +28,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import nl.aerius.taskmanager.adaptor.ClientQueueObserver; import nl.aerius.taskmanager.adaptor.WorkerProducer.WorkerMetrics; import nl.aerius.taskmanager.adaptor.WorkerSizeObserver; import nl.aerius.taskmanager.adaptor.WorkerSizeProviderProxy; @@ -56,6 +57,7 @@ public class RabbitMQWorkerSizeProvider implements WorkerSizeProviderProxy { private final long refreshRateSeconds; private final Map observers = new HashMap<>(); + private final Map clientObservers = new HashMap<>(); private final RabbitMQQueueMonitor monitor; private boolean running; @@ -72,17 +74,22 @@ public RabbitMQWorkerSizeProvider(final ScheduledExecutorService executorService } @Override - public void addObserver(final String queueName, final WorkerSizeObserver observer) { - observers.computeIfAbsent(queueName, k -> new WorkerSizeObserverComposite()).add(observer); + public void addObserver(final String workerQueueName, final WorkerSizeObserver observer) { + observers.computeIfAbsent(workerQueueName, k -> new WorkerSizeObserverComposite()).add(observer); if (observer instanceof WorkerMetrics) { - eventProducer.addMetrics(queueName, (WorkerMetrics) observer); + eventProducer.addMetrics(workerQueueName, (WorkerMetrics) observer); } } + @Override + public void addClientObserver(final String workerQueueName, final ClientQueueObserver observer) { + clientObservers.put(workerQueueName, observer); + } + @Override public boolean removeObserver(final String queueName) { eventProducer.removeMetrics(queueName); - return observers.remove(queueName) != null; + return observers.remove(queueName) != null || clientObservers.remove(queueName) != null; } @Override @@ -106,8 +113,9 @@ public void shutdown() { private void updateWorkerQueueState() { if (running) { try { - final Map queueStates = new HashMap<>(monitor.getWorkerQueueStates()); + final Map queueStates = new HashMap<>(monitor.getQueueStates()); observers.forEach((q, v) -> updateWorkerQueueState(q, queueStates.get(q))); + clientObservers.forEach((q, v) -> updateClientQueueState(v, queueStates)); } catch (final RuntimeException e) { LOG.error("Runtime error during updateWorkerQueueState", e); } @@ -120,6 +128,10 @@ private void updateWorkerQueueState(final String queueName, final RabbitMQQueueS }); } + private void updateClientQueueState(final ClientQueueObserver observer, final Map queueStates) { + queueStates.entrySet().stream().filter(e -> observer.filter(e.getKey())).forEach(e -> observer.onClientQueueUpdate(e.getKey(), e.getValue())); + } + private static class WorkerSizeObserverComposite implements WorkerSizeObserver { private final List observers = new ArrayList<>(); 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 2bcd6d3..2511693 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 @@ -31,7 +31,7 @@ */ class PriorityTaskSchedulerMetrics { - private static final String METRIC_PREFIX = "aer.taskmanager.client.queue"; + private static final String METRIC_NAME = "aer.taskmanager.queue.usage"; private static final String DESCRIPTION = "Number of tasks running on client queues"; private final Map usageMetrics = new HashMap<>(); @@ -64,7 +64,7 @@ private ObservableDoubleGauge createMetric(final IntSupplier countSupplier, fina final Attributes queueAttributes = OpenTelemetryMetrics.queueAttributes(workerQueueName, clientQueueName, "state", state); return OpenTelemetryMetrics.METER - .gaugeBuilder(METRIC_PREFIX) + .gaugeBuilder(METRIC_NAME) .setDescription(DESCRIPTION) .buildWithCallback(result -> result.record(countSupplier.getAsInt(), queueAttributes)); } diff --git a/source/taskmanager/src/test/java/nl/aerius/taskmanager/metrics/RabbitMQClientQueueReporterTest.java b/source/taskmanager/src/test/java/nl/aerius/taskmanager/metrics/RabbitMQClientQueueReporterTest.java new file mode 100644 index 0000000..6be76d3 --- /dev/null +++ b/source/taskmanager/src/test/java/nl/aerius/taskmanager/metrics/RabbitMQClientQueueReporterTest.java @@ -0,0 +1,92 @@ +/* + * 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.metrics; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.any; +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.util.function.Consumer; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.ArgumentCaptor; +import org.mockito.Captor; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +import io.opentelemetry.api.common.AttributeKey; +import io.opentelemetry.api.common.Attributes; +import io.opentelemetry.api.metrics.DoubleGaugeBuilder; +import io.opentelemetry.api.metrics.Meter; +import io.opentelemetry.api.metrics.ObservableDoubleGauge; +import io.opentelemetry.api.metrics.ObservableDoubleMeasurement; + +import nl.aerius.taskmanager.domain.RabbitMQQueueStatus; + +/** + * Test class for {@link RabbitMQClientQueueReporter}. + */ +@ExtendWith(MockitoExtension.class) +class RabbitMQClientQueueReporterTest { + + private @Mock Meter meter; + private Consumer recordMetrics; + + private @Captor ArgumentCaptor metricCaptor; + private @Captor ArgumentCaptor attributesCaptor; + + private RabbitMQClientQueueReporter reporter; + + @BeforeEach + void beforeEach() { + final DoubleGaugeBuilder builder = mock(DoubleGaugeBuilder.class); + final ObservableDoubleGauge gauge = mock(ObservableDoubleGauge.class); + doReturn(builder).when(meter).gaugeBuilder(any()); + doReturn(builder).when(builder).setDescription(any()); + doAnswer(a -> { + recordMetrics = a.getArgument(0); + return gauge; + }).when(builder).buildWithCallback(any()); + reporter = new RabbitMQClientQueueReporter(meter, "worker.test"); + } + + @Test + void testOnClientQueueUpdate() { + reporter.onClientQueueUpdate("aerius.test.some_client_queue", new RabbitMQQueueStatus(12, 20, 3)); + final ObservableDoubleMeasurement measurement = mock(ObservableDoubleMeasurement.class); + + recordMetrics.accept(measurement); + verify(measurement).record(metricCaptor.capture(), attributesCaptor.capture()); + + assertEquals(20, metricCaptor.getValue().intValue(), "Expected to report the number of messages"); + assertEquals("some_client_queue", attributesCaptor.getValue().get(AttributeKey.stringKey("queue_name")), + "Expected the client queue are attribute."); + } + + @Test + void testFitler() { + assertTrue(reporter.filter("aerius.test.some_client_queue"), "Should return true if queue name contains worker 'test' name"); + assertFalse(reporter.filter("aerius.other.some_client_queue"), "Should return false if queue name doesn't contain worker 'test' name"); + } +} 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 81223e4..eace2f7 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 @@ -54,7 +54,7 @@ protected JsonNode getJsonResultFromApi(final String apiPath) throws IOException } }; try { - final RabbitMQQueueStatus status = rpm.getWorkerQueueStates().get(QUEUENAME); + final RabbitMQQueueStatus status = rpm.getQueueStates().get(QUEUENAME); assertEquals(51, status.consumers(), "Number of workers"); assertEquals(10, status.messages(), "Number of messages"); 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 e8ae5e0..775a18d 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 @@ -61,7 +61,7 @@ void setUp() throws Exception { @Test @Timeout(value = 10, unit = TimeUnit.SECONDS) void testTriggerWorkerQueueState() throws InterruptedException, IOException { - doReturn(Map.of(TEST_QUEUE, new RabbitMQQueueStatus(1, 2, 3))).when(mockMonitor).getWorkerQueueStates(); + doReturn(Map.of(TEST_QUEUE, new RabbitMQQueueStatus(1, 2, 3))).when(mockMonitor).getQueueStates(); final CountDownLatch latch = new CountDownLatch(1); final WorkerSizeObserver observer = mock(WorkerSizeObserver.class); @@ -72,7 +72,7 @@ void testTriggerWorkerQueueState() throws InterruptedException, IOException { provider.addObserver(TEST_QUEUE, observer); provider.start(); latch.await(); - verify(mockMonitor).getWorkerQueueStates(); + verify(mockMonitor).getQueueStates(); verify(observer).onNumberOfWorkersUpdate(any()); }