From 75f57cdba7b5421a429fb37e7a9fb16a6507f1b7 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 15 Aug 2026 06:15:21 +0200 Subject: [PATCH] CB-568c: fail every queued async ticket --- .../main/java/dev/ltms/bridged/msg/MessageService.java | 10 +++++++--- .../java/dev/ltms/bridged/msg/MessageServiceTest.java | 2 ++ 2 files changed, 9 insertions(+), 3 deletions(-) diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java index f949539..c50675e 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java @@ -302,9 +302,13 @@ public final class MessageService { public boolean abandon(String target, String reason) { CompletableFuture waiter = rendezvous.currentWaiter(target); boolean failed = waiter != null && !waiter.isDone() && rendezvous.resolveFailure(waiter, reason); - boolean asyncFailed = tasks.values().stream() - .filter(task -> target.equals(task.target) && task.question == null) - .anyMatch(task -> task.future.complete(new Reply(Outcome.WORKER_FAILED, reason))); + boolean asyncFailed = false; + for (Task task : tasks.values()) { + if (target.equals(task.target) && task.question == null + && task.future.complete(new Reply(Outcome.WORKER_FAILED, reason))) { + asyncFailed = true; + } + } if (failed) { log.warn("abandoning the blocked send to {}: {}", target, reason); } diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java index 7072a7d..ecb1b5d 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java @@ -645,11 +645,13 @@ class MessageServiceTest { String first = messages.sendAsync(T, "first task"); awaitWaiting(); // first task owns the target lock and rendezvous waiter String second = messages.sendAsync(T, "second task"); // parked on the same lock, not yet queued + String third = messages.sendAsync(T, "third task"); // a second queued ticket proves the full sweep assertTrue(messages.abandon(T, "agent target term_a not found")); assertFailedTicket(first, "agent target term_a not found"); assertFailedTicket(second, "agent target term_a not found"); + assertFailedTicket(third, "agent target term_a not found"); } @Test