fleetd #808: a late fleet_reply claims its turn's waiter so the completion scrape cannot double-publish
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 55s
CI / build (pull_request) Failing after 2m19s
CI / shell-tests (push) Failing after 9s
CI / contract (push) Successful in 46s
CI / build (push) Failing after 2m3s

A send that times out leaves its captured Rendezvous waiter open for the CB-106
completion fallback (strandLateResolution). If the member then calls fleet_reply,
reply() found no live waiter and queued the answer to the inbox, but never touched
that captured waiter — so the eventual completion fallback still resolved it and
published a second, redundant scrape of the same turn.

reply() now looks up the exact waiter strandLateResolution is tracking for the
target (lateWaiters, a new per-target index onto a per-turn CompletableFuture) and
claims it with Rendezvous.resolveLateReply before falling back to the inbox. The
CompletableFuture's own single-winner complete() is the discriminator: whichever
side resolves it first decides what the other sees. A later completion scrape
against an already-claimed waiter is a no-op in CompletionResolver.resolve's
existing waiter.isDone() check, so nothing is scraped or published a second time.
A turn that never gets an explicit reply is unaffected — its own waiter is still
open when the completion fallback fires, exactly as #801 already fixed.
This commit is contained in:
Dai Ha
2026-10-07 07:02:42 +02:00
parent a2cbf50961
commit 0ff0d6cb02
3 changed files with 133 additions and 0 deletions
@@ -361,6 +361,16 @@ public final class MessageService {
* target ({@link #send} opening a fresh waiter) or a teardown ({@link #abandon}).
*/
private final ConcurrentHashMap<String, Boolean> queuedDeliveries = new ConcurrentHashMap<>();
/**
* The captured waiter a timed-out send left open for a late turn-completion scrape
* ({@link #strandLateResolution}), keyed by target so a late {@code fleet_reply} can find and
* claim the exact one its own turn belongs to ({@link #reply}). The waiter's own {@code
* complete()} call is the discriminator: whichever caller resolves it first — this reply, or
* the eventual completion scrape — decides what the other one sees, and only the winner's kind
* ever reaches the inbox. An entry removes itself once its waiter resolves.
*/
private final ConcurrentHashMap<String, CompletableFuture<Rendezvous.Resolution>> lateWaiters =
new ConcurrentHashMap<>();
private final AtomicLong ticketSeq = new AtomicLong();
/**
* Minted once per {@code MessageService} instance and folded into every ticket id (see
@@ -620,6 +630,14 @@ public final class MessageService {
+ "answers, queuing to the inbox instead of guessing", session, candidates.size(),
tickets);
}
// Claim the matching timed-out send's captured waiter, if one is still open, before
// publishing this reply. The waiter is this exact turn's own CompletableFuture, so
// completing it here means a completion-fallback scrape that resolves the same waiter
// afterward finds it already answered and does not publish a second entry for this turn.
CompletableFuture<Rendezvous.Resolution> lateWaiter = lateWaiters.get(session);
if (lateWaiter != null) {
rendezvous.resolveLateReply(lateWaiter, content);
}
inbox.publish(session, UUID.randomUUID().toString(), content);
// CB-640: record the stranding itself (not just the reply text) so fleet health can see a
// worker whose replies keep missing their waiter, not only the queue depth this leaves behind.
@@ -816,6 +834,7 @@ public final class MessageService {
// CB-640: the session is gone — nothing will ever accept or deliver into it now.
strandedReplies.remove(target);
queuedDeliveries.remove(target);
lateWaiters.remove(target);
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(target);
boolean failed = waiter != null && !waiter.isDone() && rendezvous.resolveFailure(waiter, reason);
boolean asyncFailed = false;
@@ -1091,7 +1110,9 @@ public final class MessageService {
* racing publish of the same reply.
*/
private void strandLateResolution(String target, CompletableFuture<Rendezvous.Resolution> waiter) {
lateWaiters.put(target, waiter);
waiter.whenComplete((resolution, error) -> {
lateWaiters.remove(target, waiter);
if (resolution != null && resolution.kind() == Rendezvous.Kind.COMPLETION) {
try {
// inbox.publish reaches a broker and can throw. An exception thrown inside a
@@ -304,6 +304,20 @@ public final class Rendezvous {
return waiter != null && waiter.complete(new Resolution(Kind.BACKEND_EXHAUSTED, reason));
}
/**
* Resolve a specific captured {@code waiter} as a reply — for a {@code fleet_reply} that
* arrives after the send which opened {@code waiter} already gave up on it, so the ordinary
* session-keyed {@link #resolve} finds no live waiter to match. Like
* {@link #resolveCompletion(CompletableFuture, String)} it targets the exact captured waiter
* (CB-116): a no-op if that waiter was already resolved — first resolution wins, and whichever
* one wins is what any later resolution attempt on the same waiter will see.
*
* @return true if this call resolved the waiter, false if it was null or already resolved
*/
public boolean resolveLateReply(CompletableFuture<Resolution> waiter, String content) {
return waiter != null && waiter.complete(new Resolution(Kind.REPLY, content));
}
private boolean complete(String session, Resolution resolution) {
ForwardWaiter waiter = waiters.get(session);
return waiter != null && waiter.future().complete(resolution);
@@ -123,6 +123,23 @@ class MessageServiceTest {
return drained;
}
/**
* Polls {@code target}'s inbox for {@code millis} and returns {@code true} only if nothing new
* ever arrived. Used to assert an absence against an asynchronous completion race: the window
* has to be long enough for a virtual thread already scheduled to actually run.
*/
private boolean noNewReplyWithin(String target, long millis) throws InterruptedException {
long deadline = System.currentTimeMillis() + millis;
while (System.currentTimeMillis() < deadline) {
if (!messages.drainReplies(target).isEmpty()) {
return false;
}
//noinspection BusyWait
Thread.sleep(5);
}
return true;
}
/**
* fleetd #801: {@code send}'s TimeoutException branch never completes {@code reply} itself, and
* {@link Rendezvous#close} only deregisters it from future session lookups — it does not cancel
@@ -153,6 +170,87 @@ class MessageServiceTest {
assertEquals("late answer, nobody was listening", drained.get(0).content());
}
/**
* fleetd #808: once a send has timed out, an explicit {@code fleet_reply} for that same turn
* must land in the inbox exactly once — not twice, with the second copy a completion scrape of
* the same answer the reply already delivered.
*/
@Test
void anExplicitReplyAfterATimeoutLandsExactlyOnceNotTwice() throws Exception {
CompletableFuture<MessageService.Reply> send = sendAsync("do the task", 150);
awaitWaiting();
herdr.readText("$ prompt");
injector.onStatus(T, AgentStatus.IDLE); // deliver
injector.onStatus(T, AgentStatus.WORKING); // worker starts — no reply yet
MessageService.Reply timedOut = send.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, timedOut.outcome(),
"the caller gives up while the worker is still delivered and working");
assertEquals(MessageService.ReplyOutcome.QUEUED, messages.reply(T, "structured answer"),
"no live waiter is open by now, so the explicit reply is held for a later drain");
var afterReply = awaitDrained(T);
assertEquals(1, afterReply.size(), "the explicit reply must land in the inbox exactly once");
assertEquals("structured answer", afterReply.get(0).content());
// The worker's turn then finishes for real; the completion fallback races the same waiter
// the reply already claimed, and must not publish a second entry for it.
herdr.readText("late scrape that must not also publish");
injector.onStatus(T, AgentStatus.IDLE);
assertTrue(noNewReplyWithin(T, 500),
"a completion scrape for a turn that already got an explicit fleet_reply must not "
+ "publish a second entry");
}
/**
* fleetd #808: the suppression above is per-turn, not per-target. A second, independent turn on
* the same member that never calls {@code fleet_reply} must still have its completion scrape
* land, even though an earlier turn on that same target already got an explicit reply.
*/
@Test
void theSuppressionIsPerTurnNotPerTarget() throws Exception {
// First turn: times out, then gets an explicit fleet_reply.
CompletableFuture<MessageService.Reply> first = sendAsync("first task", 150);
awaitWaiting();
herdr.readText("$ prompt 1");
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, first.get(5, TimeUnit.SECONDS).outcome());
assertEquals(MessageService.ReplyOutcome.QUEUED, messages.reply(T, "first answer"));
var firstDrain = awaitDrained(T);
assertEquals(1, firstDrain.size());
assertEquals("first answer", firstDrain.get(0).content());
// Finishing turn 1 for real also frees the worker for turn 2, and exercises that turn 1's
// own late scrape stays suppressed right up to that point.
herdr.readText("first turn's late scrape — must stay suppressed");
injector.onStatus(T, AgentStatus.IDLE);
assertTrue(noNewReplyWithin(T, 300), "turn 1's own scrape must still be suppressed");
// Second, independent turn on the same target: times out, and no fleet_reply is ever called.
CompletableFuture<MessageService.Reply> second = sendAsync("second task", 150);
awaitWaiting();
herdr.readText("$ prompt 2");
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, second.get(5, TimeUnit.SECONDS).outcome());
herdr.readText("second turn's late scrape — must still land");
injector.onStatus(T, AgentStatus.IDLE);
var secondDrain = awaitDrained(T);
assertEquals(1, secondDrain.size(),
"a later, independent turn's completion must still land even though an earlier "
+ "turn on the same target already got an explicit reply");
assertEquals("second turn's late scrape — must still land", secondDrain.get(0).content());
}
/**
* fleetd #801 defect 2: the publish inside {@code strandLateResolution}'s {@code whenComplete}
* callback runs with no caller left to surface a throw to — {@code AmqpReplyInbox.publish} can