CB-568c: fail every queued async ticket
This commit is contained in:
@@ -302,9 +302,13 @@ public final class MessageService {
|
||||
public boolean abandon(String target, String reason) {
|
||||
CompletableFuture<Rendezvous.Resolution> 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);
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user