From b8b25cf74cacafbb86ab83cc5ff8fcb57238839c Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Fri, 4 Sep 2026 13:23:01 +0700 Subject: [PATCH] #307: an ask() timeout no longer strands the worker's real reply MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit MessageService.reply()'s async-recovery path (askAnsweredAsyncTasks) required a live Task.turnId, but ask()'s own TimeoutException handler calls clearAsyncQuestion(turnId, true) — deliberately forgetting turnId so hasAsyncQuestion() stops reporting the target BUSY. That made a worker's eventual real fleet_reply, after an unanswered fleet_ask, fall through to the inbox: fleet_poll{ticket} stayed PENDING forever and was later force-failed with the false reason "session released before it replied". Fix: a new Task.askTimedOut marker is set (markAskTimedOut) right before the turnId is forgotten, and askAnsweredAsyncTasks accepts it in place of a live turnId. The marker never touches asyncTasksByTurn, so the BUSY-release behaviour (invariant 1) is untouched. The existing ambiguity guard (candidates.size() > 1 -> inbox, never guess) still applies unchanged, but is now genuinely reachable rather than pure defence in depth, since an ask timeout frees its target for a fresh, independent delegation — the affected javadocs are updated to say so. Tests: MessageServiceTest.aReplyAfterAnAskTimeoutStillCompletesTheAsyncTicket (positive, mutation-proven) and .twoAskTimedOutTicketsOnOneTargetFallBackToTheInboxRatherThanGuess (negative/ambiguity). FleetMcpTest's unansweredAsyncAskReturnsTheTicketToPending was renamed and its final assertion updated — it had pinned the old (buggy) inbox-stranding behaviour as expected. --- .../dev/ltms/fleet/msg/MessageService.java | 138 ++++++++++++------ .../java/dev/ltms/fleet/mcp/FleetMcpTest.java | 11 +- .../ltms/fleet/msg/MessageServiceTest.java | 70 +++++++++ 3 files changed, 173 insertions(+), 46 deletions(-) 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");