Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -148,7 +148,7 @@ public TaskScheduleBucket(final QueueConfig queueConfig) throws InterruptedExcep
taskScheduler = schedulerFactory.createScheduler(queueConfig);
workerProducer = factory.createWorkerProducer(queueConfig);
final WorkerPool workerPool = new WorkerPool(workerQueueName, workerProducer, taskScheduler);
final TaskManagerUsageMetricsProvider taskManagerUsageMetrics = new TaskManagerUsageMetricsProvider(workerQueueName);
final TaskManagerUsageMetricsProvider taskManagerUsageMetrics = new TaskManagerUsageMetricsProvider(workerQueueName, taskScheduler);
final TaskManagerMetricsRegister taskManagerMetricsRegister = new TaskManagerMetricsRegister(taskManagerUsageMetrics, startupGuard);
final PerformanceMetricsReporter reporter = new PerformanceMetricsReporter(scheduledExecutorService, queueConfig.queueName(),
OpenTelemetryMetrics.METER);
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
/*
* Copyright (c) Contributors to the project
*
* This program is free software: you can redistribute it and/or modify
* it under the terms of the GNU Affero General Public License as published by
* the Free Software Foundation, either version 3 of the License, or
* (at your option) any later version.
*
* This program is distributed in the hope that it will be useful,
* but WITHOUT ANY WARRANTY; without even the implied warranty of
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
* GNU Affero General Public License for more details.
*
* You should have received a copy of the GNU Affero General Public License
* along with this program. If not, see http://www.gnu.org/licenses/.
*/
package nl.aerius.taskmanager.domain;

/**
* Interface to check if a queue is empty.
*/
public interface QueueEmptyCheck {

/**
* @return Return true if the queue is empty
*/
boolean isQueueEmpty();

}
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,8 @@

import java.util.function.ToDoubleBiFunction;

import nl.aerius.taskmanager.domain.QueueEmptyCheck;

/**
* The {@link UsageMetricsProvider} for weighted load. limit and usage metrics for worker queues.
* The values are calculated averages over the time between the last measurement point and the moment the metric value is requested.
Expand All @@ -33,9 +35,13 @@ public class TaskManagerUsageMetricsProvider implements UsageMetricsProvider {
private final LoadMetric used;
private final LoadMetric free;

public TaskManagerUsageMetricsProvider(final String workerQueueName) {
public TaskManagerUsageMetricsProvider(final String workerQueueName, final QueueEmptyCheck queueEmptyCheck) {
this.workerQueueName = workerQueueName;
load = new LoadMetric((numberOfWorkers, usedWorkers) -> (numberOfWorkers > 0 ? (usedWorkers / (double) numberOfWorkers) : 0), LOAD_SUM_FUNCTION);
final ToDoubleBiFunction<Integer, Integer> loadCountFuction = (numberOfWorkers, usedWorkers) -> numberOfWorkers == 0
? (queueEmptyCheck.isQueueEmpty() ? 0 : 1.0)
: (usedWorkers / (double) numberOfWorkers);

load = new LoadMetric(loadCountFuction, LOAD_SUM_FUNCTION);
limit = new LoadMetric((numberOfWorkers, usedWorkers) -> numberOfWorkers, COUNT_SUM_FUNCTION);
used = new LoadMetric((numberOfWorkers, usedWorkers) -> usedWorkers, COUNT_SUM_FUNCTION);
free = new LoadMetric((numberOfWorkers, usedWorkers) -> numberOfWorkers - usedWorkers, COUNT_SUM_FUNCTION);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
package nl.aerius.taskmanager.scheduler;

import nl.aerius.taskmanager.domain.QueueConfig;
import nl.aerius.taskmanager.domain.QueueEmptyCheck;
import nl.aerius.taskmanager.domain.QueueWatchDogListener;
import nl.aerius.taskmanager.domain.Task;
import nl.aerius.taskmanager.domain.TaskQueue;
Expand All @@ -27,7 +28,7 @@
* Interface for the scheduling algorithm. The implementation should maintain an internal list of all tasks added and return the task to be processed
* in {@link #getNextTask()} based on whatever priority algorithm the scheduler implements.
*/
public interface TaskScheduler<T extends TaskQueue> extends WorkerUpdateHandler, QueueWatchDogListener {
public interface TaskScheduler<T extends TaskQueue> extends WorkerUpdateHandler, QueueWatchDogListener, QueueEmptyCheck {

/**
* Adds a Task to the scheduler to being processed. The scheduler will return this task in {@link #getNextTask()}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -111,7 +111,7 @@ public boolean containsAll(final Collection<?> arg0) {

@Override
public boolean isEmpty() {
throw new UnsupportedOperationException("Not implemented");
return queue.isEmpty() && groupedQueue.isEmpty();
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -117,6 +117,11 @@ public Task getNextTask() throws InterruptedException {
return task;
}

@Override
public boolean isQueueEmpty() {
return queue.isEmpty();
}

private void obtainTask() {
final Task task = queue.poll();

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,11 @@ public Task getNextTask() throws InterruptedException {
return tasks.take();
}

@Override
public boolean isQueueEmpty() {
return tasks.isEmpty();
}

@Override
public void updateQueue(final PriorityTaskQueue queue) {
// Not used
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,16 +23,21 @@
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;

import nl.aerius.taskmanager.domain.QueueEmptyCheck;

/**
* Test class for {@link TaskManagerUsageMetricsProvider}.
*/
class TaskManagerUsageMetricsProviderTest {

private boolean queuEmpty;
private final QueueEmptyCheck emptyChecker = () -> queuEmpty;
private TaskManagerUsageMetricsProvider provider;

@BeforeEach
void beforeEach() {
provider = new TaskManagerUsageMetricsProvider("TEST");
provider = new TaskManagerUsageMetricsProvider("TEST", emptyChecker);
queuEmpty = true;
}

@Test
Expand All @@ -55,6 +60,12 @@ void testNumberOfFreeWorkers() throws InterruptedException {
assertMetricAndZero(3, 10, provider::getNumberOfFreeWorkers, 7, "Expected 7 worker to be free.", 10);
}

@Test
void testZeroWorkersNoEmptyQueue() throws InterruptedException {
queuEmpty = false;
assertMetric(10, 0, provider::getLoad, 100, "Expected a load of 100%", 100);
}

private void assertMetricAndZero(final int numberOfUsed, final int numberOfWorkers, final DoubleSupplier supplier, final int expected,
final String description, final int expectedAfterReset) throws InterruptedException {
assertMetric(numberOfUsed, numberOfWorkers, supplier, expected, description, expectedAfterReset);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,9 @@
package nl.aerius.taskmanager.scheduler.priorityqueue;

import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.Mockito.doReturn;
import static org.mockito.Mockito.mock;

Expand Down Expand Up @@ -49,11 +51,13 @@ void testPoll() {
queue.add(task5);

assertNotNull(queue.peek(), "Queue should have task at the queue.");
assertFalse(queue.isEmpty(), "Queue should not be empty.");
assertEquals(task1, queue.poll(), "Poll should return first task added.");
assertEquals(task3, queue.poll(), "Poll should return first task of task with correlationId 2.");
assertEquals(task2, queue.poll(), "Poll should return second task of correlationId 1, because it's at the front of the queue.");
assertEquals(task4, queue.poll(), "Poll should return second task of correlationId 2, because it should be at front of the queue.");
assertEquals(task5, queue.poll(), "Poll should return last added task.");
assertTrue(queue.isEmpty(), "Queue should be empty after all tasks taken from the queue.");
}

private static Task mockTask(final String correlationId) {
Expand Down
Loading