diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java
index 902f9a4..a1bd20e 100644
--- a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java
+++ b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java
@@ -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.
*
- *
Ambiguous match also falls to the inbox. {@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
- * wrong 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.
+ *
Ambiguous match also falls to the inbox. {@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 wrong 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.
*
*
{@code content} is required (fleetd #302). 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 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}).
*
- * Returns at most one entry today — verified, not assumed. {@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).
+ *
Can return more than one entry — reachable, not just defence in depth.
+ * {@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 askAnsweredAsyncTasks(String target) {
List 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
diff --git a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java
index e6acc30..0e08820 100644
--- a/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java
+++ b/fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java
@@ -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
diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java
index 6e3580c..7dd20c8 100644
--- a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java
+++ b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java
@@ -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");