Skip to content
Open
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
180 changes: 107 additions & 73 deletions core/src/main/java/org/jboss/pnc/rex/core/QueueManagerImpl.java
Original file line number Diff line number Diff line change
Expand Up @@ -13,32 +13,36 @@
import org.jboss.pnc.rex.core.api.TaskRegistry;
import org.jboss.pnc.rex.core.counter.Counter;
import org.jboss.pnc.rex.core.counter.MaxConcurrent;
import org.jboss.pnc.rex.core.counter.Running;
import org.jboss.pnc.rex.core.counter.RunningTracker;
import org.jboss.pnc.rex.core.delegates.FaultToleranceDecorator;
import org.jboss.pnc.rex.model.Task;

import jakarta.enterprise.context.ApplicationScoped;
import jakarta.transaction.Transactional;

import java.util.*;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.stream.Collectors;

import static jakarta.transaction.Transactional.TxType.MANDATORY;
import static java.util.stream.Collectors.groupingBy;

@Slf4j
@ApplicationScoped
public class QueueManagerImpl implements QueueManager {

private static final String DEFAULT_QUEUE_NAMING = "DEFAULT";
private final Counter max;
private final Counter running;
private final RunningTracker running;
private final TaskRegistry container;
private final TaskController controller;
private final FaultToleranceDecorator ft;

public QueueManagerImpl(@MaxConcurrent Counter max,
@Running Counter running,
RunningTracker running,
TaskRegistry container,
TaskController controller,
FaultToleranceDecorator ft) {
Expand All @@ -54,71 +58,109 @@ public QueueManagerImpl(@MaxConcurrent Counter max,
public void poke() {
log.info("QUEUE: Poking Task queues");
Map<String, Long> maxEntries = max.entries();
Map<String, Long> runningEntries = running.entries();

Set<String> consideredQueues = new HashSet<>();
for (var queue : maxEntries.keySet()) {
Long maxValue = maxEntries.get(queue);
Long runningValue = runningEntries.get(queue);

if (runningValue >= maxValue) {
log.debug("QUEUE '{}': Maximum number of parallel tasks reached.({} out of {})",
queue == null ? DEFAULT_QUEUE_NAMING : queue,
runningValue,
maxValue);
if (running.requiresReservation()) {
pokeWithReservation(queue, maxValue);
} else {
consideredQueues.add(queue);
pokeLegacy(queue, maxValue);
}
}
}

for (var queue : consideredQueues) {
Long maxValue = maxEntries.get(queue);
VersionedValue<Long> runningMetadata = running.getMetadataValue(queue);
Long runningValue = runningMetadata.getValue();
/**
* Historical order for optimistic compare-and-swap counters: dequeue first, then increment.
* The shared-key version check serializes admission across nodes (a stale increment fails and
* the whole transaction is retried), so no over-admission can happen.
*/
private void pokeLegacy(String queue, Long maxValue) {
long runningValue = running.get(queue);
if (runningValue >= maxValue) {
log.debug("QUEUE '{}': Maximum number of parallel tasks reached.({} out of {})",
queue == null ? DEFAULT_QUEUE_NAMING : queue,
runningValue,
maxValue);
return;
}

long freeSpace = maxValue - runningValue;
List<Task> randomEnqueuedTasks = container.getEnqueuedTasksByQueueName(queue, freeSpace);
if (randomEnqueuedTasks.isEmpty()) {
continue;
}
long freeSpace = maxValue - runningValue;
List<Task> randomEnqueuedTasks = container.getEnqueuedTasksByQueueName(queue, freeSpace);
if (randomEnqueuedTasks.isEmpty()) {
return;
}

log.info("QUEUE '{}': Free space of {} found. Scheduling {} task(s) of {}",
queue == null ? DEFAULT_QUEUE_NAMING : queue,
freeSpace,
randomEnqueuedTasks.size(),
randomEnqueuedTasks.stream().map(Task::getName).collect(Collectors.toList())
);
log.info("QUEUE '{}': Free space of {} found. Scheduling {} task(s) of {}",
queue == null ? DEFAULT_QUEUE_NAMING : queue,
freeSpace,
randomEnqueuedTasks.size(),
randomEnqueuedTasks.stream().map(Task::getName).collect(Collectors.toList())
);

randomEnqueuedTasks.forEach(task -> controller.dequeue(task.getName()));

randomEnqueuedTasks.forEach(task -> controller.dequeue(task.getName()));
long newRunning = running.reserve(queue, randomEnqueuedTasks.size());
log.info("QUEUE '{}': Increased running counter to {}.",
queue == null ? DEFAULT_QUEUE_NAMING : queue,
newRunning);
}

log.info("QUEUE '{}': Increasing running counter. ({} to {}) [ISPN-VERSION:{}]",
/**
* Order for atomic, non-transactional counters (e.g. StrongCounter): reserve capacity first,
* release any overshoot, then dequeue only the amount we are allowed to start. This prevents
* two nodes from both observing the same free space and over-admitting.
*/
private void pokeWithReservation(String queue, Long maxValue) {
long runningValue = running.get(queue);
if (runningValue >= maxValue) {
log.debug("QUEUE '{}': Maximum number of parallel tasks reached.({} out of {})",
queue == null ? DEFAULT_QUEUE_NAMING : queue,
runningValue,
(runningValue + randomEnqueuedTasks.size()),
runningMetadata.getVersion());
if (!running.replaceValue(queue, runningMetadata, runningValue + randomEnqueuedTasks.size())) {
RuntimeException e = new ConcurrentModificationException("Running counter was modified concurrently.");
log.error("QUEUE '{}': Concurrent modification detected.", queue == null ? DEFAULT_QUEUE_NAMING : queue, e);
throw e;
}
maxValue);
return;
}

long freeSpace = maxValue - runningValue;
List<Task> candidates = container.getEnqueuedTasksByQueueName(queue, freeSpace);
if (candidates.isEmpty()) {
return;
}

int want = candidates.size();
long afterReserve = running.reserve(queue, want);
int allowed = want;
if (afterReserve > maxValue) {
long giveBack = Math.min(afterReserve - maxValue, want);
running.release(queue, giveBack);
allowed = (int) (want - giveBack);
}
if (allowed <= 0) {
return;
}

List<Task> toSchedule = candidates.subList(0, allowed);
log.info("QUEUE '{}': Reserved {} slot(s). Scheduling {} task(s) of {}",
queue == null ? DEFAULT_QUEUE_NAMING : queue,
allowed,
toSchedule.size(),
toSchedule.stream().map(Task::getName).collect(Collectors.toList())
);

try {
toSchedule.forEach(task -> controller.dequeue(task.getName()));
} catch (RuntimeException e) {
// The transaction rolls back every dequeue done in this invocation; release the whole
// reservation so the (non-transactional) counter does not drift.
running.release(queue, allowed);
throw e;
}
}

@Override
@Transactional(MANDATORY)
public void decreaseRunningCounter(@Nullable String name) {
VersionedValue<Long> runningMetadata = running.getMetadataValue(name);
long runningValue = runningMetadata.getValue() - 1;
log.info("QUEUE '{}': Decreasing running counter by one. ({} to {}) [ISPN-VERSION:{}]",
name == null ? DEFAULT_QUEUE_NAMING : name,
runningMetadata.getValue(),
runningValue,
runningMetadata.getVersion());
if (!running.replaceValue(name, runningMetadata, runningValue)) {
RuntimeException e = new ConcurrentModificationException("Running counter was modified concurrently.");
log.error("QUEUE: Concurrent modification detected.", e);
throw e;
}
running.decrement(name);
log.info("QUEUE '{}': Decreased running counter by one.", name == null ? DEFAULT_QUEUE_NAMING : name);
}

@Override
Expand Down Expand Up @@ -146,35 +188,27 @@ public Long getMaximumConcurrency(@Nullable String name) {

@Override
public Long getRunningCounter(@Nullable String name) {
VersionedValue<Long> meta = running.getMetadataValue(name);
return meta == null ? null : meta.getValue();
return running.peek(name);
}

@Override
@Transactional(MANDATORY)
public void synchronizeRunningCounter() {
Map<String, List<Task>> tasksByQueue = container.getTasks(false, false, true, false, false, null)
.stream()
.collect(groupingBy(Task::getQueue));


Set<String> existingQueues = running.entries().keySet();
for (String queue : existingQueues) {
var runningValue = running.getMetadataValue(queue);
var tasksInQueue = tasksByQueue.get(queue);

long actualValue;
if (tasksInQueue == null) {
actualValue = 0L;
} else {
actualValue = tasksInQueue.size();
}
// null-safe grouping: DEFAULT queue tasks have a null queue, which Collectors.groupingBy rejects
Map<String, List<Task>> tasksByQueue = new HashMap<>();
for (Task task : container.getTasks(false, false, true, false, false, null)) {
tasksByQueue.computeIfAbsent(task.getQueue(), q -> new ArrayList<>()).add(task);
}

if (!runningValue.getValue().equals(actualValue)) {
log.info("Synchronizing running counter. Mismatch between active tasks and counter found. Previous value '{}' -> new value '{}'", runningValue.getValue(), actualValue);
running.replaceValue(queue, runningValue, actualValue);
}
// reconcile every known queue. Max counters survive a restart (persistent cache), while the
// running counters may not (e.g. volatile StrongCounters), so union both sources.
Set<String> queues = new HashSet<>(max.entries().keySet());
queues.addAll(running.queues());

for (String queue : queues) {
List<Task> tasksInQueue = tasksByQueue.get(queue);
long actualValue = tasksInQueue == null ? 0L : tasksInQueue.size();
running.reconcile(queue, actualValue);
}
}

Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,112 @@
/*
* SPDX-FileCopyrightText: Copyright © 2021 Red Hat, Inc., and individual contributors as indicated by the @author tags.
* SPDX-License-Identifier: Apache-2.0
*/
package org.jboss.pnc.rex.core.counter;

import jakarta.annotation.Nullable;
import lombok.extern.slf4j.Slf4j;
import org.infinispan.client.hotrod.VersionedValue;

import java.util.ConcurrentModificationException;
import java.util.Set;

/**
* Optimistic compare-and-swap implementation of {@link RunningTracker}, backed by the existing
* {@code rex-counter} cache through the {@link Running} {@link Counter}.
* <p>
* This preserves the historical behavior: every mutation reads the versioned value and writes it
* back with {@code replaceWithVersion}; on a version conflict it throws
* {@link ConcurrentModificationException} so the surrounding {@code internal-retry} guarded
* transaction is replayed.
* <p>
* Instantiated by {@link RunningTrackerProducer} (not a CDI bean itself) so the active
* implementation can be selected by configuration.
*/
@Slf4j
public class CasRunningTracker implements RunningTracker {

private static final String DEFAULT_QUEUE_NAMING = "DEFAULT";

private final Counter counter;

public CasRunningTracker(@Running Counter counter) {
this.counter = counter;
}

@Override
public long get(@Nullable String queue) {
VersionedValue<Long> meta = counter.getMetadataValue(queue);
return meta == null ? 0L : meta.getValue();
}

@Override
@Nullable
public Long peek(@Nullable String queue) {
VersionedValue<Long> meta = counter.getMetadataValue(queue);
return meta == null ? null : meta.getValue();
}

@Override
public long reserve(@Nullable String queue, long n) {
VersionedValue<Long> meta = counter.getMetadataValue(queue);
long newValue = meta.getValue() + n;
if (!counter.replaceValue(queue, meta, newValue)) {
throw concurrentModification(queue);
}
return newValue;
}

@Override
public void release(@Nullable String queue, long n) {
VersionedValue<Long> meta = counter.getMetadataValue(queue);
long newValue = meta.getValue() - n;
if (!counter.replaceValue(queue, meta, newValue)) {
throw concurrentModification(queue);
}
}

@Override
public void decrement(@Nullable String queue) {
release(queue, 1);
}

@Override
public void reconcile(@Nullable String queue, long actual) {
VersionedValue<Long> meta = counter.getMetadataValue(queue);
if (meta == null) {
counter.initialize(queue, actual);
return;
}
if (meta.getValue() != actual) {
log.info("QUEUE '{}': Synchronizing running counter. Mismatch between active tasks and counter found. "
+ "Previous value '{}' -> new value '{}'",
queue == null ? DEFAULT_QUEUE_NAMING : queue, meta.getValue(), actual);
counter.replaceValue(queue, meta, actual);
}
}

@Override
public void initialize(@Nullable String queue, long value) {
counter.initialize(queue, value);
}

@Override
public Set<String> queues() {
return counter.entries().keySet();
}

@Override
public boolean requiresReservation() {
// Optimistic CAS on the shared counter key already serializes admission across nodes;
// keep the historical dequeue-then-increment order.
return false;
}

private ConcurrentModificationException concurrentModification(String queue) {
ConcurrentModificationException e =
new ConcurrentModificationException("Running counter was modified concurrently.");
log.error("QUEUE '{}': Concurrent modification detected.", queue == null ? DEFAULT_QUEUE_NAMING : queue, e);
return e;
}
}
Loading
Loading