#307: an ask() timeout no longer strands the worker's real reply #314

Closed
agent wants to merge 0 commits from worker/fix-307-275890-6 into main
Member

Fixes fleetd #307.

The bug

MessageService.reply() recovers a stranded reply into its async ticket only when the task still has a live Task.turnId (askAnsweredAsyncTasks). answer()'s own bounded-wait timeout keeps turnId (clearAsyncQuestion(turnId, false)), so that recovery path works. ask()'s own timeout (no answer from the primary at all) calls clearAsyncQuestion(ticket.turnId(), true), which nulls turnId and drops the task from asyncTasksByTurn — deliberately, so hasAsyncQuestion(target) stops reporting the target BUSY. But that also erased the only signal askAnsweredAsyncTasks had that the worker's eventual real fleet_reply still belonged to this task. The reply then fell straight to the inbox: fleet_poll{ticket} stayed PENDING forever and was later force-failed by abandon() with the false reason "session released before it replied".

The trap I avoided

The issue explicitly warns against the obvious fix of passing forgetTurn=false on the ask timeout, since that would keep the target BUSY forever (every later fleet_send to it refused). I did not do that — forgetTurn=true on the ask timeout is untouched.

The fix

Added a second, independent signal on Task: askTimedOut (a volatile boolean), set by a new markAskTimedOut(turnId) call before clearAsyncQuestion(turnId, true) erases turnId in ask()'s TimeoutException catch. askAnsweredAsyncTasks now accepts (turnId != null || askTimedOut) instead of requiring turnId != null alone. The marker never touches asyncTasksByTurn, so the BUSY-release behaviour (issue invariant 1) is untouched — verified by the existing unansweredAsyncQuestionReturnsTheTicketToPendingAndReleasesItsTarget test, still green.

The existing ambiguity guard (candidates.size() > 1 → fall back to the inbox rather than guess — issue invariant 2) is unchanged in code, but I found it is now genuinely reachable, not just defence in depth: an ask timeout frees its target, so a completely independent, fresh delegation can be dispatched to the same target and itself ask-and-lapse before the first worker's real reply arrives, landing two askTimedOut tasks on one target at once. I updated the two javadoc blocks that previously claimed "returns at most one entry today — verified, not assumed" to describe this, and added a test that exercises it (see below).

Invariant 3 (a reply with a live rendezvous waiter keeps the existing fast path) is untouched — reply()'s first branch (rendezvous.resolve) is unchanged.

Call sites searched before deciding

grep -n "clearAsyncQuestion\|finishAsyncTask\|TimeoutException" src/main/java/dev/ltms/fleet/msg/MessageService.java — three TimeoutException catches in the file: send() (unrelated — no async-question bookkeeping there), ask() (the one this issue is about, now fixed), and answer() (see the "shape check" note below). Two forgetTurn=true call sites for clearAsyncQuestion: ask()'s NO_WAITER branch and ask()'s TIMEOUT branch. I checked the NO_WAITER branch and it is not the same bug: markAsyncQuestion only ever attaches a Task there when rendezvous.currentWaiter(workerSession) is already non-null, and in that case the forward rendezvous waiter is genuinely still open (unlike the timeout path, where send() already returned once the question resolved), so a real fleet_reply there resolves through the ordinary fast path (rendezvous.resolve) in reply(), never needing the async-recovery path at all. I left it unchanged.

Shape check (reported only, not fixed, per the issue's instructions)

  • answer()'s own TimeoutException catch (around line 1014 after this change) does not call finishAsyncTask/clearAsyncQuestion, unlike its own success path a few lines above (finishAsyncTask(turnId, result)). I checked this and it is intentional, not a new defect: clearAsyncQuestion(turnId, false) already ran earlier in answer(), before the reply.get(...) wait that can time out, so turnId is deliberately still live when this catch runs — that is exactly the "answer() times out" row of the issue's own table, the one that already worked before this PR. Reporting it since it matches the literal "cleanup on one path but not its sibling" shape the issue asked me to look for, even though it turned out correct.

Anything wrong in the issue itself

Nothing — I read reply(), askAnsweredAsyncTasks, ask(), answer(), hasAsyncQuestion, and abandon() in full and the issue's description of the code (the table, the trap, the invariants) matched what is actually there.

Tests

  • MessageServiceTest.aReplyAfterAnAskTimeoutStillCompletesTheAsyncTicket — positive case (mutation-proven, see below).
  • MessageServiceTest.twoAskTimedOutTicketsOnOneTargetFallBackToTheInboxRatherThanGuess — negative/ambiguity case: two independent ask-timed-out tasks on one target, reply must fall back to the inbox.
  • FleetMcpTest.unansweredAsyncAskReturnsTheTicketToPendingThenAWorkersLateReplyStillCompletesIt (renamed from unansweredAsyncAskReturnsTheTicketToPending) — this existing test had pinned the old, buggy behaviour as its expected outcome (messages.drainReplies("term_a").getFirst()...); updated its final assertions to the fixed behaviour (the ticket completes with the reply; the inbox stays empty).

Mutation proof

Reverted MessageService.java only (tests kept), ran mvn -Dtest=MessageServiceTest,FleetMcpTest test:

[ERROR] Tests run: 73, Failures: 1, Errors: 0, Skipped: 0 -- in dev.ltms.fleet.msg.MessageServiceTest
[ERROR] dev.ltms.fleet.msg.MessageServiceTest.aReplyAfterAnAskTimeoutStillCompletesTheAsyncTicket -- Time elapsed: 3.220 s <<< FAILURE!
org.opentest4j.AssertionFailedError: expected: <DONE> but was: <PENDING>

[ERROR] Tests run: 71, Failures: 1, Errors: 0, Skipped: 0 -- in dev.ltms.fleet.mcp.FleetMcpTest
[ERROR] dev.ltms.fleet.mcp.FleetMcpTest.unansweredAsyncAskReturnsTheTicketToPendingThenAWorkersLateReplyStillCompletesIt -- Time elapsed: 0.057 s <<< FAILURE!
org.opentest4j.AssertionFailedError: expected: <finished after timeout> but was: <[pending — worker idle]>

Restored the fix (git apply of the saved patch), re-ran the full build: mvn clean install → Tests run: 1318, Failures: 0, Errors: 0, Skipped: 0 / BUILD SUCCESS.

Build

cd fleetd && mvn clean install, run unpiped, full output read:

[INFO] Tests run: 1318, Failures: 0, Errors: 0, Skipped: 0
[INFO] BUILD SUCCESS

Worktree: /Users/dai.ha/LTMS/.bridged-worktrees/1ea531-6
Branch: worker/fix-307-275890-6
Files changed: fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java, fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java, fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java

Fixes fleetd #307. ## The bug `MessageService.reply()` recovers a stranded reply into its async ticket only when the task still has a live `Task.turnId` (`askAnsweredAsyncTasks`). `answer()`'s own bounded-wait timeout keeps `turnId` (`clearAsyncQuestion(turnId, false)`), so that recovery path works. `ask()`'s own timeout (no answer from the primary at all) calls `clearAsyncQuestion(ticket.turnId(), true)`, which nulls `turnId` and drops the task from `asyncTasksByTurn` — deliberately, so `hasAsyncQuestion(target)` stops reporting the target BUSY. But that also erased the only signal `askAnsweredAsyncTasks` had that the worker's eventual real `fleet_reply` still belonged to this task. The reply then fell straight to the inbox: `fleet_poll{ticket}` stayed `PENDING` forever and was later force-failed by `abandon()` with the false reason "session released before it replied". ## The trap I avoided The issue explicitly warns against the obvious fix of passing `forgetTurn=false` on the ask timeout, since that would keep the target BUSY forever (every later `fleet_send` to it refused). I did not do that — `forgetTurn=true` on the ask timeout is untouched. ## The fix Added a second, independent signal on `Task`: `askTimedOut` (a `volatile boolean`), set by a new `markAskTimedOut(turnId)` call **before** `clearAsyncQuestion(turnId, true)` erases `turnId` in `ask()`'s `TimeoutException` catch. `askAnsweredAsyncTasks` now accepts `(turnId != null || askTimedOut)` instead of requiring `turnId != null` alone. The marker never touches `asyncTasksByTurn`, so the BUSY-release behaviour (issue invariant 1) is untouched — verified by the existing `unansweredAsyncQuestionReturnsTheTicketToPendingAndReleasesItsTarget` test, still green. The existing ambiguity guard (`candidates.size() > 1` → fall back to the inbox rather than guess — issue invariant 2) is unchanged in code, but I found it is now genuinely *reachable*, not just defence in depth: an ask timeout frees its target, so a completely independent, fresh delegation can be dispatched to the same target and itself ask-and-lapse before the first worker's real reply arrives, landing two `askTimedOut` tasks on one target at once. I updated the two javadoc blocks that previously claimed "returns at most one entry today — verified, not assumed" to describe this, and added a test that exercises it (see below). Invariant 3 (a reply with a live rendezvous waiter keeps the existing fast path) is untouched — `reply()`'s first branch (`rendezvous.resolve`) is unchanged. ## Call sites searched before deciding `grep -n "clearAsyncQuestion\|finishAsyncTask\|TimeoutException" src/main/java/dev/ltms/fleet/msg/MessageService.java` — three `TimeoutException` catches in the file: `send()` (unrelated — no async-question bookkeeping there), `ask()` (the one this issue is about, now fixed), and `answer()` (see the "shape check" note below). Two `forgetTurn=true` call sites for `clearAsyncQuestion`: `ask()`'s `NO_WAITER` branch and `ask()`'s `TIMEOUT` branch. I checked the `NO_WAITER` branch and it is not the same bug: `markAsyncQuestion` only ever attaches a `Task` there when `rendezvous.currentWaiter(workerSession)` is already non-null, and in that case the forward rendezvous waiter is genuinely still open (unlike the timeout path, where `send()` already returned once the question resolved), so a real `fleet_reply` there resolves through the ordinary fast path (`rendezvous.resolve`) in `reply()`, never needing the async-recovery path at all. I left it unchanged. ## Shape check (reported only, not fixed, per the issue's instructions) - `answer()`'s own `TimeoutException` catch (around line 1014 after this change) does not call `finishAsyncTask`/`clearAsyncQuestion`, unlike its own success path a few lines above (`finishAsyncTask(turnId, result)`). I checked this and it is intentional, not a new defect: `clearAsyncQuestion(turnId, false)` already ran earlier in `answer()`, before the `reply.get(...)` wait that can time out, so `turnId` is deliberately still live when this catch runs — that is exactly the "answer() times out" row of the issue's own table, the one that already worked before this PR. Reporting it since it matches the literal "cleanup on one path but not its sibling" shape the issue asked me to look for, even though it turned out correct. ## Anything wrong in the issue itself Nothing — I read `reply()`, `askAnsweredAsyncTasks`, `ask()`, `answer()`, `hasAsyncQuestion`, and `abandon()` in full and the issue's description of the code (the table, the trap, the invariants) matched what is actually there. ## Tests - `MessageServiceTest.aReplyAfterAnAskTimeoutStillCompletesTheAsyncTicket` — positive case (mutation-proven, see below). - `MessageServiceTest.twoAskTimedOutTicketsOnOneTargetFallBackToTheInboxRatherThanGuess` — negative/ambiguity case: two independent ask-timed-out tasks on one target, reply must fall back to the inbox. - `FleetMcpTest.unansweredAsyncAskReturnsTheTicketToPendingThenAWorkersLateReplyStillCompletesIt` (renamed from `unansweredAsyncAskReturnsTheTicketToPending`) — this existing test had pinned the *old, buggy* behaviour as its expected outcome (`messages.drainReplies("term_a").getFirst()...`); updated its final assertions to the fixed behaviour (the ticket completes with the reply; the inbox stays empty). ## Mutation proof Reverted `MessageService.java` only (tests kept), ran `mvn -Dtest=MessageServiceTest,FleetMcpTest test`: ``` [ERROR] Tests run: 73, Failures: 1, Errors: 0, Skipped: 0 -- in dev.ltms.fleet.msg.MessageServiceTest [ERROR] dev.ltms.fleet.msg.MessageServiceTest.aReplyAfterAnAskTimeoutStillCompletesTheAsyncTicket -- Time elapsed: 3.220 s <<< FAILURE! org.opentest4j.AssertionFailedError: expected: <DONE> but was: <PENDING> [ERROR] Tests run: 71, Failures: 1, Errors: 0, Skipped: 0 -- in dev.ltms.fleet.mcp.FleetMcpTest [ERROR] dev.ltms.fleet.mcp.FleetMcpTest.unansweredAsyncAskReturnsTheTicketToPendingThenAWorkersLateReplyStillCompletesIt -- Time elapsed: 0.057 s <<< FAILURE! org.opentest4j.AssertionFailedError: expected: <finished after timeout> but was: <[pending — worker idle]> ``` Restored the fix (`git apply` of the saved patch), re-ran the full build: `mvn clean install` → `Tests run: 1318, Failures: 0, Errors: 0, Skipped: 0` / `BUILD SUCCESS`. ## Build `cd fleetd && mvn clean install`, run unpiped, full output read: ``` [INFO] Tests run: 1318, Failures: 0, Errors: 0, Skipped: 0 [INFO] BUILD SUCCESS ``` Worktree: `/Users/dai.ha/LTMS/.bridged-worktrees/1ea531-6` Branch: `worker/fix-307-275890-6` Files changed: `fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java`, `fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java`, `fleetd/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java`
agent added 1 commit 2026-09-04 08:24:19 +02:00
#307: an ask() timeout no longer strands the worker's real reply
CI / contract (pull_request) Successful in 42s
CI / build (pull_request) Successful in 1m30s
b8b25cf74c
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.
ltms closed this pull request 2026-09-04 08:34:31 +02:00
Some checks are pending
CI / contract (pull_request) Successful in 42s
CI / build (pull_request) Successful in 1m30s

Pull request closed

Sign in to join this conversation.