Merge #307: a worker's real reply after an ask timeout completes its ticket instead of stranding
CI / build (push) Successful in 2m5s
CI / contract (push) Successful in 2m19s

This commit is contained in:
Dai Ha
2026-09-04 13:31:43 +07:00
3 changed files with 173 additions and 46 deletions
@@ -187,6 +187,24 @@ public final class MessageService {
private volatile Long completedNanos;
private volatile Reply question;
private volatile String turnId;
/**
* Set when this task's {@code fleet_ask} lapsed with no answer (fleetd #307):
* {@link #clearAsyncQuestion} then forgets {@link #turnId} (nulls it and drops the task from
* {@code asyncTasksByTurn}) so {@link #hasAsyncQuestion} stops reporting the target BUSY — a
* later {@code fleet_send} to it must be accepted, not refused. But the worker's turn is
* still genuinely live: it resumed on its own and will eventually call its real
* {@code fleet_reply}. Losing {@link #turnId} loses {@link #askAnsweredAsyncTasks}' only
* signal that such a reply belongs to this task, so that reply used to fall straight to the
* inbox and strand — {@code fleet_poll} stayed {@code PENDING} forever, later force-failed by
* {@link #abandon} with the misleading "session released before it replied". This flag is a
* second, independent signal that survives the forgetting: {@link #askAnsweredAsyncTasks}
* accepts it in place of a live {@link #turnId}, without ever re-adding the task to
* {@code asyncTasksByTurn} (so the BUSY release is untouched). Cleared implicitly once
* {@link #future} resolves — every match in {@link #askAnsweredAsyncTasks} already requires
* {@code !future.isDone()}, so a task that recovered (or was later failed by
* {@link #abandon}) can never match again regardless of this flag's value.
*/
private volatile boolean askTimedOut;
private Task(String ticket, String target, LongSupplier nowNanos) {
this.ticket = ticket;
@@ -395,18 +413,16 @@ public final class MessageService {
* {@link Rendezvous#resolveQuestion} must keep today's {@code NO_WAITER} behaviour — questions
* are interactive and must never be queued.
*
* <p><strong>Ambiguous match also falls to the inbox.</strong> {@link #askAnsweredAsyncTasks}
* cannot actually return more than one entry today (see its own javadoc for why — in short,
* {@link #hasAsyncQuestion} keeps a target BUSY, so no second task can reach this state, for as
* long as an earlier one's {@code turnId} is still stamped). That is an emergent guarantee from
* two other facts, not one this method enforces, so this branch stays in as defence in depth
* rather than being removed as dead code: if it ever weakens, returning whichever candidate a
* {@code ConcurrentHashMap} iteration reaches first would let a genuine reply complete the
* <em>wrong</em> ticket — silently handing the lead something that reads like a correct answer to
* a delegation the worker never touched, which is worse than a failure because the lead acts on
* it. When more than one candidate exists, guessing is not safe: fall back to the inbox exactly
* as the zero-candidate case does, and let {@link #abandon} apply the eventual recovery
* deterministically instead.
* <p><strong>Ambiguous match also falls to the inbox.</strong> {@link #askAnsweredAsyncTasks} can
* return more than one entry — a reachable state, not a hypothetical one (see its own javadoc:
* an {@code fleet_ask} that lapsed with no answer, fleetd #307, frees the target for a completely fresh
* delegation, which can itself go on to ask-and-lapse before the first worker's real reply
* arrives). Returning whichever candidate a {@code ConcurrentHashMap} iteration reaches first
* would let a genuine reply complete the <em>wrong</em> ticket — silently handing the lead
* something that reads like a correct answer to a delegation the worker never touched, which is
* worse than a failure because the lead acts on it. When more than one candidate exists, guessing
* is not safe: fall back to the inbox exactly as the zero-candidate case does, and let
* {@link #abandon} apply the eventual recovery deterministically instead.
*
* <p><strong>{@code content} is required (fleetd #302).</strong> Both doors that reach this
* method must reject a missing/blank reply the same way, so the check lives here rather than in
@@ -434,15 +450,18 @@ public final class MessageService {
count(FleetMetrics.REPLIES, "path", "rendezvous");
return true; // a live send took it — unchanged fast path
}
// #137: no live rendezvous waiter, but this may be the worker's real fleet_reply resuming a
// turn that {@link #answer} already gave up waiting on. answer()'s own bounded wait (the
// primary's fleet_send{turnId} call, capped well under a minute) can time out and close its
// waiter long before the worker — now actually resuming real work — finishes and replies. That
// reply used to have nowhere to land but the session inbox, leaving the async ticket's future
// unresolved forever: fleet_poll{ticket} stayed PENDING until fleet_stop's abandon() forced it
// FAILED with a misleading "session released before it replied" reason, even though the reply
// had, in fact, arrived. Completing the matching ticket directly here means fleet_poll{ticket}
// sees the real reply instead.
// #137/fleetd #307: no live rendezvous waiter, but this may be the worker's real fleet_reply resuming
// a turn that either answer() (#137) or ask() (fleetd #307) already gave up waiting on:
// - answer()'s own bounded wait (the primary's fleet_send{turnId} call, capped well under a
// minute) can time out and close its waiter long before the worker — now actually resuming
// real work — finishes and replies.
// - ask()'s own wait for the primary can time out first, with the worker resuming on its own
// and finishing unanswered.
// Either way that reply used to have nowhere to land but the session inbox, leaving the async
// ticket's future unresolved forever: fleet_poll{ticket} stayed PENDING until fleet_stop's
// abandon() forced it FAILED with a misleading "session released before it replied" reason,
// even though the reply had, in fact, arrived. Completing the matching ticket directly here
// means fleet_poll{ticket} sees the real reply instead.
List<Task> candidates = askAnsweredAsyncTasks(session);
if (candidates.size() == 1) {
Task orphan = candidates.get(0);
@@ -473,34 +492,42 @@ public final class MessageService {
}
/**
* Every still-open async task on {@code target} whose {@code fleet_ask} was already answered —
* its {@link Task#turnId} is stamped but its {@link Task#question} was cleared by {@link #answer}
* — yet whose future is not resolved yet (#137). Empty if no such task exists, including the
* common case where {@code target}'s worker never used {@code fleet_ask} at all (a task that was
* never asked has {@code turnId == null}, so it can never match here and only ever completes
* through the ordinary rendezvous fast path in {@link #reply}).
* Every still-open async task on {@code target} whose worker is genuinely expected to send a
* real {@code fleet_reply} next with nothing left registered to catch it: either its
* {@code fleet_ask} was already answered — {@link Task#turnId} is stamped but {@link
* Task#question} was cleared by {@link #answer} — or its {@code fleet_ask} lapsed unanswered and
* {@link Task#askTimedOut} marks that (fleetd #307; {@link Task#turnId} is {@code null} by then, forgotten
* so the target is not left BUSY — see {@link Task#askTimedOut}'s own javadoc). Either way the
* task's future is not resolved yet. Empty if no such task exists, including the common case
* where {@code target}'s worker never used {@code fleet_ask} at all (a task that was never asked
* has both {@code turnId == null} and {@code askTimedOut == false}, so it can never match here and
* only ever completes through the ordinary rendezvous fast path in {@link #reply}).
*
* <p><strong>Returns at most one entry today — verified, not assumed.</strong> {@link #send}
* refuses to open a waiter on {@code target} while {@link #hasAsyncQuestion} is true, and that
* check matches ANY task whose {@code turnId} is still stamped in {@code asyncTasksByTurn} —
* not only while its question is still open. {@link #answer} deliberately leaves that stamp in
* place ({@code clearAsyncQuestion(turnId, false)}) until the resumed turn's own future actually
* resolves, at which point {@link #finishAsyncTask} both removes the stamp AND completes that
* task's future in the same call. So a second task can never reach "{@code turnId} stamped, future
* still open" — the exact pair this method matches on — while a first one already holds it: by
* the time the stamp is gone, so is the eligibility. This is an emergent property of those two
* facts holding together, not something this method (or its callers) enforces on its own — flip
* {@code forgetTurn} to {@code true} in that one {@link #answer} call and it silently stops being
* true, with nothing left to fail loudly. The callers below still handle "more than one" as
* defence in depth against exactly that, not because they exercise it today: {@link #reply}
* treats it as unresolvable and falls back to the inbox; {@link #abandon} would pick the oldest
* deterministically (its own {@code matching} list has no such guarantee — see its javadoc).
* <p><strong>Can return more than one entry — reachable, not just defence in depth.</strong>
* {@link #send} refuses to open a waiter on {@code target} while {@link #hasAsyncQuestion} is
* true, and that check matches ANY task whose {@code turnId} is still stamped in
* {@code asyncTasksByTurn}. While a task's {@code turnId} stays stamped — {@link #answer} leaves
* it in place ({@code clearAsyncQuestion(turnId, false)}) until {@link #finishAsyncTask} removes
* the stamp and completes the future in the same call — no second task on the same target can
* reach an eligible state, because {@link #send} would refuse it as BUSY first. That single-task
* guarantee holds only for the {@code turnId}-stamped half of this method's match: an
* {@link Task#askTimedOut} task is, by construction, no longer stamped in {@code asyncTasksByTurn}
* (that is the whole point of forgetting {@code turnId} in {@link #clearAsyncQuestion}), so the
* target is free the moment one ask lapses. A fresh, independent {@code sendAsync} to the same
* target can then be dispatched, itself pause on {@code fleet_ask}, and itself time out — landing
* a second {@code askTimedOut} task on the very target the first one is still waiting to answer
* for. Two (or more) genuinely open tasks on one target is therefore a real, reachable state
* today, not a hypothetical: {@link #reply} treats it as unresolvable and falls back to the
* inbox rather than guess which task a reply belongs to (guessing wrong would hand the lead a
* plausible-looking answer to a delegation the worker never touched — worse than a failure,
* because the lead acts on it); {@link #abandon} instead picks the oldest deterministically (its
* own {@code matching} list has a different, wider match — see its javadoc).
*/
private List<Task> askAnsweredAsyncTasks(String target) {
List<Task> candidates = new ArrayList<>();
for (Task task : tasks.values()) {
if (target.equals(task.target) && task.question == null && task.turnId != null
&& !task.future.isDone()) {
if (target.equals(task.target) && task.question == null && !task.future.isDone()
&& (task.turnId != null || task.askTimedOut)) {
candidates.add(task);
}
}
@@ -898,6 +925,12 @@ public final class MessageService {
return new AskResult(AskOutcome.ANSWERED, answer);
} catch (TimeoutException e) {
log.debug("fleet_ask from {} went unanswered in {}ms", workerSession, timeoutMillis);
// fleetd #307: mark the task BEFORE clearAsyncQuestion(forgetTurn=true) below drops it out of
// asyncTasksByTurn and nulls its turnId — that forgetting is deliberate and stays (it is
// what keeps the target from staying BUSY forever), but it would otherwise also erase
// askAnsweredAsyncTasks' only signal that the worker's eventual real fleet_reply still
// belongs to this task, stranding it in the inbox with a false "never replied" verdict.
markAskTimedOut(ticket.turnId());
clearAsyncQuestion(ticket.turnId(), true);
return new AskResult(AskOutcome.TIMED_OUT, null);
} catch (ExecutionException e) {
@@ -1162,6 +1195,23 @@ public final class MessageService {
return task;
}
/**
* Mark {@code turnId}'s task as having a {@code fleet_ask} that lapsed with no answer (fleetd #307), so
* {@link #askAnsweredAsyncTasks} still recognizes the worker's eventual real {@code fleet_reply}
* as belonging to it after {@link #clearAsyncQuestion}'s {@code forgetTurn=true} erases
* {@link Task#turnId} — see {@link Task#askTimedOut}. Must be called before that forgetting, while
* {@code turnId} can still resolve the task in {@code asyncTasksByTurn}; a lookup afterward would
* find nothing. Only when it matches the task's current turn — same guard as
* {@link #clearAsyncQuestion} — so a chained second {@code fleet_ask} (#282) that already moved
* the task to a fresh {@code turnId} cannot mark it for a turn that is no longer its own.
*/
private void markAskTimedOut(String turnId) {
Task task = asyncTasksByTurn.get(turnId);
if (task != null && turnId.equals(task.turnId)) {
task.askTimedOut = true;
}
}
/** Clear an answered or lapsed question, but only when it matches the ticket's current turn. */
private void clearAsyncQuestion(String turnId, boolean forgetTurn) {
// CB-582: tell the push loop first — like ticketCollected, a removal for a turnId it never
@@ -189,7 +189,7 @@ class FleetMcpTest {
}
@Test
void unansweredAsyncAskReturnsTheTicketToPending() throws Exception {
void unansweredAsyncAskReturnsTheTicketToPendingThenAWorkersLateReplyStillCompletesIt() throws Exception {
McpSchema.CallToolResult accepted = FleetMcp.sendAsync(messages, "term_a", "do it", null, Set.of());
String ticket = textOf(accepted).substring(textOf(accepted).indexOf("ticket=") + "ticket=".length()).trim();
@@ -203,8 +203,15 @@ class FleetMcpTest {
assertTrue(textOf(ask).contains("no answer"), textOf(ask));
assertTrue(textOf(FleetMcp.poll(messages, ticket, null)).startsWith("[pending"));
// fleetd #307: the worker resumed on its own after the primary never answered, and its real
// fleet_reply must complete its OWN async ticket — not strand in the inbox with
// fleet_poll{ticket} stuck PENDING forever and later force-failed with a false "session
// released before it replied" reason. This used to land in the inbox instead (see the old
// assertion this replaced: messages.drainReplies("term_a").getFirst()...) — that was the bug.
FleetMcp.reply(messages, "term_a", "finished after timeout");
assertEquals("finished after timeout", messages.drainReplies("term_a").getFirst().content());
assertEquals("finished after timeout", textOf(FleetMcp.poll(messages, ticket, null)));
assertTrue(messages.drainReplies("term_a").isEmpty(),
"the reply completed its own ticket directly and never touched the inbox");
}
@Test
@@ -957,6 +957,76 @@ class MessageServiceTest {
assertEquals(MessageService.Phase.DONE, awaitTicketPhase(next, MessageService.Phase.DONE).phase());
}
/**
* fleetd #307: a worker's {@code fleet_ask} can time out because the primary never answers —
* distinct from {@link #aReplyAfterAnswerTimesOutStillCompletesTheAsyncTicket}, where the
* primary DID answer and only its own bounded wait for the resumed turn expired.
* {@code ask()}'s timeout path deliberately forgets the task's {@code turnId} (so
* {@code hasAsyncQuestion} stops reporting the target BUSY — see
* {@code unansweredAsyncQuestionReturnsTheTicketToPendingAndReleasesItsTarget} above), which used
* to also erase the one signal {@code askAnsweredAsyncTasks} needed to recognize the worker's
* eventual real {@code fleet_reply}. That reply then had nowhere to land but the inbox, and
* {@code fleet_poll{ticket}} stayed PENDING forever — later force-failed with the false reason
* "session released before it replied", even though the worker had, in fact, replied.
*/
@Test
void aReplyAfterAnAskTimeoutStillCompletesTheAsyncTicket() throws Exception {
String ticket = messages.sendAsync(T, "task that asks then finishes alone");
awaitWaiting();
injectDelivery();
assertEquals(MessageService.AskOutcome.TIMED_OUT,
messages.ask(T, "which config?", 200).outcome());
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket).phase(),
"only the question wait ended; the delegated turn may still finish");
// The worker keeps working past the timeout and only now calls fleet_reply — with no live
// rendezvous waiter open (ask()'s timeout already closed it) and no new send() having
// reopened one for this target.
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE);
assertEquals("PR opened: https://example/pulls/42", done.reply(),
"fleet_poll{ticket} must return the worker's real reply, not stay pending forever");
assertEquals("reply", done.replySource());
assertFalse(messages.hasStrandedReply(T),
"the reply completed its own ticket directly and never touched the inbox");
}
/**
* fleetd #307's ambiguity guard: an ask timeout frees its target ({@code hasAsyncQuestion}
* becomes false the instant it lapses — proven above), so a second, independent delegation can
* be dispatched to the same target and itself go on to ask-and-lapse before the first worker's
* real reply ever arrives. Two open tasks are then both eligible candidates on one target with
* no live waiter to disambiguate them. A reply arriving now must not guess which one it answers
* — guessing wrong would hand the lead a plausible-looking answer to a delegation the worker
* never touched, worse than a failure because the lead acts on it — so it must fall back to the
* inbox exactly as the zero-candidate case does.
*/
@Test
void twoAskTimedOutTicketsOnOneTargetFallBackToTheInboxRatherThanGuess() throws Exception {
String ticket1 = messages.sendAsync(T, "first task that asks");
awaitWaiting();
injectDelivery();
assertEquals(MessageService.AskOutcome.TIMED_OUT, messages.ask(T, "Q1?", 200).outcome());
String ticket2 = messages.sendAsync(T, "second task that asks");
awaitWaiting();
injectDelivery();
assertEquals(MessageService.AskOutcome.TIMED_OUT, messages.ask(T, "Q2?", 200).outcome());
assertTrue(messages.reply(T, "which task does this answer?"));
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket1).phase(),
"an ambiguous reply must not guess ticket1");
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket2).phase(),
"an ambiguous reply must not guess ticket2");
assertTrue(messages.hasStrandedReply(T));
var drained = messages.drainReplies(T);
assertEquals(1, drained.size());
assertEquals("which task does this answer?", drained.get(0).content());
}
@Test
void asyncQuestionBelongsToTheTaskThatOwnsItsForwardWaiter() throws Exception {
String first = messages.sendAsync(T, "first task");