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 3f32664..1bca565 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 @@ -50,6 +50,7 @@ class LoadMetric { private int numberOfWorkers; private final ToDoubleBiFunction countFunction; private final ToDoubleBiFunction sumFunction; + private final Object lock = new Object(); public LoadMetric(final ToDoubleBiFunction countFunction, final ToDoubleBiFunction sumFunction) { this.countFunction = countFunction; @@ -62,23 +63,23 @@ public LoadMetric(final ToDoubleBiFunction countFunction, fina * @param deltaUsedWorkers number of jobs on the workers being added or subtracted. * @param numberOfWorkers Number of available workers */ - public synchronized void register(final int deltaUsedWorkers, final int numberOfWorkers) { - this.numberOfWorkers = numberOfWorkers; - final long newLast = System.currentTimeMillis(); - final long delta = newLast - last; - - total += delta * countFunction.applyAsDouble(numberOfWorkers, usedWorkers); - totalMeasureTime += delta; - last = newLast; - usedWorkers = Math.max(0, usedWorkers + deltaUsedWorkers); + public void register(final int deltaUsedWorkers, final int numberOfWorkers) { + synchronized (lock) { + if (deltaUsedWorkers != 0 || this.numberOfWorkers != numberOfWorkers) { + updateState(deltaUsedWorkers, numberOfWorkers); + } + } } + /** * Resets the metric state. Sets running workers to 0, and resets the average load time by calling process. */ - public synchronized void reset() { - usedWorkers = 0; - process(); + public void reset() { + synchronized(lock) { + usedWorkers = 0; + process(); + } } /** @@ -86,13 +87,27 @@ public synchronized void reset() { * * @return Average load of the workers since the last time this method was called */ - public synchronized double process() { - // Call register here to set the end time this moment. This will calculate workers running up till now as being active. - register(0, numberOfWorkers); - final double averageTotal = totalMeasureTime > 0 ? sumFunction.applyAsDouble(total, totalMeasureTime) : 0; + public double process() { + synchronized(lock) { + // Call register here to set the end time this moment. This will calculate workers running up till now as being active. + updateState(0, numberOfWorkers); + final double averageTotal = totalMeasureTime > 0 ? sumFunction.applyAsDouble(total, totalMeasureTime) : 0; - totalMeasureTime = 0; - total = 0; - return averageTotal; + totalMeasureTime = 0; + total = 0; + return averageTotal; + } + } + + private void updateState(final int deltaUsedWorkers, final int updatedNumberOfWorkers) { + final long newLast = System.currentTimeMillis(); + final long delta = newLast - last; + + // First calculate the total for passed time period with the worker values up to this time. + total += delta * countFunction.applyAsDouble(numberOfWorkers, usedWorkers); + totalMeasureTime += delta; + last = newLast; + numberOfWorkers = updatedNumberOfWorkers; + usedWorkers = Math.max(0, usedWorkers + deltaUsedWorkers); } } 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 56e3132..00e828f 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 @@ -17,8 +17,6 @@ package nl.aerius.taskmanager.metrics; import java.util.Map; -import java.util.Set; -import java.util.concurrent.ConcurrentHashMap; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -40,9 +38,6 @@ public class TaskManagerMetricsRegister implements WorkerProducerHandler, Worker private final TaskManagerUsageMetricsProvider taskManagerUsageMetricsProvider; private final StartupGuard startupGuard; - // Keep track of dispatched tasks, because when taskmanager restarts it should not register tasks already on the queue - // as it doesn't have any metrics on it anymore. - private final Set dispatchedTasks = ConcurrentHashMap.newKeySet(); private int numberOfWorkers; @@ -53,19 +48,12 @@ public TaskManagerMetricsRegister(final TaskManagerUsageMetricsProvider taskMana @Override public void onWorkDispatched(final String messageId, final Map messageMetaData) { - synchronized (dispatchedTasks) { - dispatchedTasks.add(messageId); - taskManagerUsageMetricsProvider.register(1, numberOfWorkers); - } + taskManagerUsageMetricsProvider.register(1, numberOfWorkers); } @Override public void onWorkerFinished(final String messageId, final Map messageMetaData) { - synchronized (dispatchedTasks) { - if (dispatchedTasks.remove(messageId)) { - taskManagerUsageMetricsProvider.register(-1, numberOfWorkers); - } - } + taskManagerUsageMetricsProvider.register(-1, numberOfWorkers); } @Override @@ -75,14 +63,13 @@ public synchronized void onNumberOfWorkersUpdate(final int numberOfWorkers, fina LOG.info("Queue {} will be started with {} messages already on the queue.", taskManagerUsageMetricsProvider.getWorkerQueueName(), numberOfMessages); taskManagerUsageMetricsProvider.register(numberOfMessages, numberOfWorkers); + } else { + taskManagerUsageMetricsProvider.register(0, numberOfWorkers); } } @Override public void reset() { - synchronized (dispatchedTasks) { - dispatchedTasks.clear(); - taskManagerUsageMetricsProvider.reset(); - } + taskManagerUsageMetricsProvider.reset(); } } diff --git a/source/taskmanager/src/test/java/nl/aerius/taskmanager/metrics/LoadMetricTest.java b/source/taskmanager/src/test/java/nl/aerius/taskmanager/metrics/LoadMetricTest.java new file mode 100644 index 0000000..c24c625 --- /dev/null +++ b/source/taskmanager/src/test/java/nl/aerius/taskmanager/metrics/LoadMetricTest.java @@ -0,0 +1,96 @@ +/* + * 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.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; + +import java.util.function.ToDoubleBiFunction; + +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; + +/** + * Test class for {@link LoadMetric}. + */ +@ExtendWith(MockitoExtension.class) +class LoadMetricTest { + + @Mock ToDoubleBiFunction countFunction; + @Captor ArgumentCaptor countUsedWorkersCaptor; + @Captor ArgumentCaptor countWorkersCaptor; + + private LoadMetric loadMetric; + + @BeforeEach + void beforeEach() { + doAnswer(a -> ((Integer) a.getArgument(0)).doubleValue()).when(countFunction).applyAsDouble(any(), any()); + + loadMetric = new LoadMetric(countFunction, (total, time) -> total); + } + + /** + * Register + */ + @Test + void testProcess() throws InterruptedException { + loadMetric.register(5, 10); + loadMetric.process(); + verify(countFunction, times(2)).applyAsDouble(countWorkersCaptor.capture(),countUsedWorkersCaptor.capture()); + assertEquals(10, countWorkersCaptor.getValue(), "Should get original number of workers"); + assertEquals(5, countUsedWorkersCaptor.getValue(), "Should get original number of used workers"); + } + + @Test + void testWorkersChanged() throws InterruptedException { + loadMetric.register(0, 10); + loadMetric.register(0, 10); + // 2nd call to register should not trigger updating internal state. + verify(countFunction, times(1)).applyAsDouble(countWorkersCaptor.capture(), countUsedWorkersCaptor.capture()); + assertEquals(0, countWorkersCaptor.getValue(), "Should get inital number of workers, which was 0"); + assertEquals(0, countUsedWorkersCaptor.getValue(), "Should get inital number of used workers, which was 0"); + loadMetric.register(0, 5); + // call with changed number of workers should trigger update state. + verify(countFunction, times(2)).applyAsDouble(any(), any()); + // process call should also trigger update. + loadMetric.process(); + verify(countFunction, times(3)).applyAsDouble(any(), any()); + } + + @Test + void testReset() throws InterruptedException { + // call 2 times. because first time used workers is initialized. + loadMetric.register(5, 10); + loadMetric.register(6, 10); + verify(countFunction, times(2)).applyAsDouble(any(), countUsedWorkersCaptor.capture()); + assertEquals(5, countUsedWorkersCaptor.getValue(), "Should get the number of used workers of the first call"); + loadMetric.reset(); + verify(countFunction, times(3)).applyAsDouble(any(), countUsedWorkersCaptor.capture()); + assertEquals(0, countUsedWorkersCaptor.getValue(), "Reset has set used number to 0, so that is at is expected here."); + loadMetric.register(8, 10); + verify(countFunction, times(4)).applyAsDouble(any(), countUsedWorkersCaptor.capture()); + assertEquals(0, countUsedWorkersCaptor.getValue(), "This call should get previous value of 0 used workers"); + } +} 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 af8805a..2608ca9 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 @@ -45,7 +45,8 @@ public class TaskManagerMetricsRegisterTest { private @Mock TaskManagerUsageMetricsProvider taskManagerUsageMetricsProvider; - private @Captor ArgumentCaptor taskManagerUsagerMetricsProviderCaptor; + private @Captor ArgumentCaptor startUpNrOfMessagesCaptor; + private @Captor ArgumentCaptor lastNrOfMessagesCaptor; private StartupGuard startupGuard; private TaskManagerMetricsRegister register; @@ -59,18 +60,29 @@ void beforeEach() { @Test void testOnWorkDispatched() { startUp(10, 0); + // Should have called register only once, but with 0 delta + verifyTaskManagerUsageMetricsProvider(1, 0, startUpNrOfMessagesCaptor); register.onWorkDispatched("1", createMap(QUEUE_1, 100L)); register.onWorkDispatched("2", createMap(QUEUE_2, 200L)); - verifyTaskManagerUsageMetricsProvider(2, 2); + register.onNumberOfWorkersUpdate(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); } @Test void testOnWorkerFinished() { - startUp(2, 0); + startUp(2, 2); + // Should have called register only once, but with 2 as delta as it should match total number of messages. + verifyTaskManagerUsageMetricsProvider(1, 2, startUpNrOfMessagesCaptor); register.onWorkDispatched("1", createMap(QUEUE_1, 100L)); register.onWorkerFinished("1", createMap(QUEUE_1, 100L)); register.onWorkerFinished("2", createMap(QUEUE_2, 200L)); - verifyTaskManagerUsageMetricsProvider(2, 0); + register.onNumberOfWorkersUpdate(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.. + verifyTaskManagerUsageMetricsProvider(5, 1, lastNrOfMessagesCaptor); } @Test @@ -85,10 +97,9 @@ void testReset() throws InterruptedException { verify(taskManagerUsageMetricsProvider, times(1)).reset(); } - private void verifyTaskManagerUsageMetricsProvider(final int times, final int sum) { - verify(taskManagerUsageMetricsProvider, times(times)).register(taskManagerUsagerMetricsProviderCaptor.capture(), - anyInt()); - assertEquals(sum, taskManagerUsagerMetricsProviderCaptor.getAllValues().stream().mapToInt(Integer::intValue).sum(), + private void verifyTaskManagerUsageMetricsProvider(final int times, final int sum, final ArgumentCaptor captor) { + verify(taskManagerUsageMetricsProvider, times(times)).register(captor.capture(), anyInt()); + assertEquals(sum, captor.getAllValues().stream().mapToInt(Integer::intValue).sum(), "Should have registered a total sum of " + sum); }