Merge CB-577 follow-up: drop the target-keyed async index
CI / contract (push) Successful in 43s
CI / build (push) Successful in 55s

asyncTasksByWaiter correlates an async question by the exact rendezvous
waiter, so the target-keyed set it replaced can no longer decide
anything. Keeping it meant two indexes of the same fact, one of them
ambiguous whenever a target has two accepted tickets.

The race the old test modelled by reflection is gone with it: identity
keys make 'some other task reached this target' unrepresentable, so
there is no longer a wrong task for the question to land on.

Also carries the criterion-1 doc correction, which is identical to
aac29d6 on a different parent.
This commit is contained in:
Dai Ha
2026-08-15 06:40:31 +02:00
2 changed files with 1 additions and 39 deletions
@@ -9,7 +9,6 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.List;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CompletionException;
@@ -171,8 +170,6 @@ public final class MessageService {
private final Metrics metrics; // CB-502: nullable — no registry in unit tests
private final ConcurrentHashMap<String, ReentrantLock> sessionLocks = new ConcurrentHashMap<>();
private final ConcurrentHashMap<String, Task> tasks = new ConcurrentHashMap<>();
/** Async tasks that have accepted delivery for a target. */
private final ConcurrentHashMap<String, Set<Task>> asyncTasksByTarget = new ConcurrentHashMap<>();
/** Async task that owns each exact forward rendezvous waiter. */
private final ConcurrentHashMap<CompletableFuture<Rendezvous.Resolution>, Task> asyncTasksByWaiter =
new ConcurrentHashMap<>();
@@ -548,13 +545,7 @@ public final class MessageService {
tasks.put(ticket, task);
asyncExecutor.submit(() -> {
try {
Runnable trackingAccepted = () -> {
if (onAccepted != null) {
onAccepted.run();
}
asyncTasksByTarget.computeIfAbsent(target, _ -> ConcurrentHashMap.newKeySet()).add(task);
};
Reply result = send(target, content, ASYNC_TIMEOUT_MS, trackingAccepted, task);
Reply result = send(target, content, ASYNC_TIMEOUT_MS, onAccepted, task);
if (result.outcome() == Outcome.QUESTION) {
// Keep the accepted owner until answer() finishes it. markAsyncQuestion may run
// just after resolveQuestion wakes this thread.
@@ -563,7 +554,6 @@ public final class MessageService {
}
} catch (Throwable t) {
task.future.completeExceptionally(t);
untrackAsyncTarget(task);
}
});
pruneTerminalTickets();
@@ -643,7 +633,6 @@ public final class MessageService {
if (forgetTurn) {
asyncTasksByTurn.remove(turnId, task);
task.turnId = null;
untrackAsyncTarget(task);
}
}
}
@@ -651,7 +640,6 @@ public final class MessageService {
/** Complete and detach an async ticket after its worker's actual terminal reply. */
private void finishAsyncTask(Task task, Reply result) {
task.future.complete(result);
untrackAsyncTarget(task);
if (task.turnId != null) {
asyncTasksByTurn.remove(task.turnId, task);
}
@@ -670,17 +658,6 @@ public final class MessageService {
return asyncTasksByTurn.values().stream().anyMatch(task -> target.equals(task.target));
}
/** Stop tracking a task once it no longer owns an accepted target turn. */
private void untrackAsyncTarget(Task task) {
Set<Task> targetTasks = asyncTasksByTarget.get(task.target);
if (targetTasks != null) {
targetTasks.remove(task);
if (targetTasks.isEmpty()) {
asyncTasksByTarget.remove(task.target, targetTasks);
}
}
}
/** Release the async executor. */
public void close() {
asyncExecutor.shutdown();
@@ -10,9 +10,6 @@ import dev.ltms.bridged.inject.Injector;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import java.lang.reflect.Field;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
@@ -697,22 +694,10 @@ class MessageServiceTest {
}
@Test
@SuppressWarnings("unchecked")
void asyncQuestionBelongsToTheTaskThatOwnsItsForwardWaiter() throws Exception {
String first = messages.sendAsync(T, "first task");
awaitWaiting();
// Model the resolveQuestion/markAsyncQuestion race: another accepted task reached the target set.
Field byTargetField = MessageService.class.getDeclaredField("asyncTasksByTarget");
byTargetField.setAccessible(true);
Map<String, Set<Object>> byTarget = (Map<String, Set<Object>>) byTargetField.get(messages);
Set<Object> targetTasks = byTarget.get(T);
Class<?> taskClass = Class.forName(MessageService.class.getName() + "$Task");
var constructor = taskClass.getDeclaredConstructor(String.class);
constructor.setAccessible(true);
targetTasks.clear();
targetTasks.add(constructor.newInstance(T));
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
assertEquals(MessageService.Phase.ASKING, awaitTicketPhase(first, MessageService.Phase.ASKING).phase());