#282: fix a chained fleet_ask killing its own async ticket #289

Closed
agent wants to merge 0 commits from worker/fix-282-chained-ask-e6d0bb-8 into main
Member

fleetd #282 — a second fleet_ask in the same resumed turn kills its own async ticket

Step 1: reachability — confirmed, driven through the public API

I traced the exact sequence the ticket describes and confirmed it is reachable purely through
MessageService's public methods (sendAsync → ask → answer → ask again), no reaching into
Rendezvous or the task maps:

  1. sendAsync creates task1, send()'s worker thread opens waiter w1 and registers it in
    asyncTasksByWaiter (:802).
  2. Worker ask("Q1") → markAsyncQuestion finds task1 via asyncTasksByWaiter, stamps
    task1.turnId = turnId1, asyncTasksByTurn[turnId1] = task1. resolveQuestion completes w1
    with QUESTION; sendAsync's lambda sees QUESTION and correctly skips finishing the task.
  3. Lead answer(turnId1, "a1") opens a new waiter w2 — and, before my fix, never registered
    it in asyncTasksByWaiter. clearAsyncQuestion(turnId1, false) keeps asyncTasksByTurn[turnId1]
    pointing at task1, then blocks on w2.get().
  4. Worker, still in the resumed turn, ask("Q2") before replying. markAsyncQuestion looks up
    asyncTasksByWaiter.get(w2) — before the fix, null (never registered) — so turnId2 never
    re-associates with task1. resolveQuestion still completes w2 with QUESTION regardless
    (it only needs a live rendezvous waiter, not a Task).
  5. answer() wakes with Resolution(QUESTION, "Q2", turnId2) and, before the fix, called
    finishAsyncTask(turnId1, result) unconditionally — completing task1.future with a
    QUESTION "reply". Reply.completed() is false for QUESTION, so fleet_poll reported
    Phase.FAILED while the worker was alive and mid-conversation with the primary.

Step 2: fix — both halves, as specified

  • Guard answer()'s finishAsyncTask the same way sendAsync's own lambda already guards
    itself (MessageService.java around the old :1000): skip the call when
    result.outcome() == Outcome.QUESTION. On its own this stops the ticket from being killed —
    it leaves the ticket's future open, but (without the second change) with no Task re-armed
    under the new turnId, so a later answer(turnId2, ...) would find nothing in
    asyncTasksByTurn and never actually finish the ticket either.
  • Register answer()'s new waiter in asyncTasksByWaiter, looking up the Task via
    asyncTasksByTurn.get(turnId) right after opening the waiter (mirroring send()'s :802,
    which has the Task passed in directly instead of looked up). On its own, without the first
    change, markAsyncQuestion would successfully re-associate task1 with turnId2, but
    finishAsyncTask(turnId1, result) would still unconditionally complete the ticket with the
    QUESTION "reply" from the old turnId1 right before that association could ever be used.
    Both changes together make the chained ask actually work: the ticket stays open, surfaces the
    second question, and a subsequent answer(turnId2, ...) finishes it for real.

Also fixed, as a direct consequence of the above becoming reachable for the first time: when
markAsyncQuestion re-arms a Task under a new turnId, it now drops the stale
asyncTasksByTurn entry for the old turnId (previously only ever written, never removed on
reassignment) — otherwise every chained ask leaks one entry in that map forever, since nothing
else in this codebase prunes it.

Missing-guard shape elsewhere in msg/: I grepped every .future.complete( /
finishAsyncTask( call site in msg/. The only other completions are reply()'s
orphan.future.complete(new Reply(REPLIED, content)) and abandon()'s
task.future.complete(outcome) — both always construct a terminal (REPLIED/WORKER_FAILED)
Reply themselves rather than forwarding a Rendezvous.Resolution's outcome, so neither can ever
carry QUESTION through. I did not find a second instance of this exact shape.

Step 3 & 4: test on the real path, proven to catch the regression

Added secondFleetAskInTheSameResumedTurnDoesNotKillTheAsyncTicket to
MessageServiceTest, driven entirely through sendAsync/ask/answer/poll (no reaching into
Rendezvous or the task maps) — asks twice in one resumed turn, and asserts the ticket surfaces
the second question (Phase.ASKING, "Q2") instead of FAILED, then finishes normally
(Phase.DONE, "done") once the primary answers the second question and the worker replies for
real.

Reverted the production change (kept the test), reran it, and it went red:

org.opentest4j.AssertionFailedError: expected: <ASKING> but was: <FAILED>
	at dev.ltms.fleet.msg.MessageServiceTest.awaitTicketPhase(MessageServiceTest.java:1468)
	at dev.ltms.fleet.msg.MessageServiceTest.secondFleetAskInTheSameResumedTurnDoesNotKillTheAsyncTicket(MessageServiceTest.java:1019)

Restored the production fix; reran — green again.

Build

cd fleetd && mvn clean install, unpiped, three full clean runs after restoring the fix:
Tests run: 1284, Failures: 0, Errors: 0, Skipped: 0 / BUILD SUCCESS (reproduced twice at the
full-suite level, plus 5 consecutive isolated runs of MessageServiceTest alone — 71 tests,
0 failures each time — to rule out flakiness in the new test itself).

One caveat for the reviewer: an earlier, unsaved mvn clean install run (before I started
capturing output to a file) exited 1 with truncated tool output I could not fully read. Two
subsequent full clean builds captured to a log both came back BUILD SUCCESS with the identical
Tests run: 1284, Failures: 0, Errors: 0 count, and MessageServiceTest alone was green 5/5
times in isolation, so this looks like a pre-existing flaky test elsewhere in the 1284-test suite
unrelated to this change — but I could not name the specific test from that first run's output,
so I'm reporting it rather than asserting it away.

Files changed

  • fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java
  • fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java
## fleetd #282 — a second `fleet_ask` in the same resumed turn kills its own async ticket ### Step 1: reachability — confirmed, driven through the public API I traced the exact sequence the ticket describes and confirmed it is reachable purely through `MessageService`'s public methods (`sendAsync` → `ask` → `answer` → `ask` again), no reaching into `Rendezvous` or the task maps: 1. `sendAsync` creates `task1`, `send()`'s worker thread opens waiter `w1` and registers it in `asyncTasksByWaiter` (`:802`). 2. Worker `ask("Q1")` → `markAsyncQuestion` finds `task1` via `asyncTasksByWaiter`, stamps `task1.turnId = turnId1`, `asyncTasksByTurn[turnId1] = task1`. `resolveQuestion` completes `w1` with `QUESTION`; `sendAsync`'s lambda sees `QUESTION` and correctly skips finishing the task. 3. Lead `answer(turnId1, "a1")` opens a **new** waiter `w2` — and, before my fix, never registered it in `asyncTasksByWaiter`. `clearAsyncQuestion(turnId1, false)` keeps `asyncTasksByTurn[turnId1]` pointing at `task1`, then blocks on `w2.get()`. 4. Worker, still in the resumed turn, `ask("Q2")` before replying. `markAsyncQuestion` looks up `asyncTasksByWaiter.get(w2)` — before the fix, `null` (never registered) — so `turnId2` never re-associates with `task1`. `resolveQuestion` still completes `w2` with `QUESTION` regardless (it only needs a live rendezvous waiter, not a `Task`). 5. `answer()` wakes with `Resolution(QUESTION, "Q2", turnId2)` and, before the fix, called `finishAsyncTask(turnId1, result)` **unconditionally** — completing `task1.future` with a `QUESTION` "reply". `Reply.completed()` is `false` for `QUESTION`, so `fleet_poll` reported `Phase.FAILED` while the worker was alive and mid-conversation with the primary. ### Step 2: fix — both halves, as specified - **Guard `answer()`'s `finishAsyncTask`** the same way `sendAsync`'s own lambda already guards itself (`MessageService.java` around the old `:1000`): skip the call when `result.outcome() == Outcome.QUESTION`. On its own this stops the ticket from being killed — it leaves the ticket's future open, but (without the second change) with no `Task` re-armed under the new `turnId`, so a later `answer(turnId2, ...)` would find nothing in `asyncTasksByTurn` and never actually finish the ticket either. - **Register `answer()`'s new waiter in `asyncTasksByWaiter`**, looking up the `Task` via `asyncTasksByTurn.get(turnId)` right after opening the waiter (mirroring `send()`'s `:802`, which has the `Task` passed in directly instead of looked up). On its own, without the first change, `markAsyncQuestion` would successfully re-associate `task1` with `turnId2`, but `finishAsyncTask(turnId1, result)` would still unconditionally complete the ticket with the `QUESTION` "reply" from the *old* `turnId1` right before that association could ever be used. Both changes together make the chained ask actually work: the ticket stays open, surfaces the second question, and a subsequent `answer(turnId2, ...)` finishes it for real. Also fixed, as a direct consequence of the above becoming reachable for the first time: when `markAsyncQuestion` re-arms a `Task` under a new `turnId`, it now drops the stale `asyncTasksByTurn` entry for the *old* `turnId` (previously only ever written, never removed on reassignment) — otherwise every chained ask leaks one entry in that map forever, since nothing else in this codebase prunes it. **Missing-guard shape elsewhere in `msg/`:** I grepped every `.future.complete(` / `finishAsyncTask(` call site in `msg/`. The only other completions are `reply()`'s `orphan.future.complete(new Reply(REPLIED, content))` and `abandon()`'s `task.future.complete(outcome)` — both always construct a terminal (`REPLIED`/`WORKER_FAILED`) `Reply` themselves rather than forwarding a `Rendezvous.Resolution`'s outcome, so neither can ever carry `QUESTION` through. I did not find a second instance of this exact shape. ### Step 3 & 4: test on the real path, proven to catch the regression Added `secondFleetAskInTheSameResumedTurnDoesNotKillTheAsyncTicket` to `MessageServiceTest`, driven entirely through `sendAsync`/`ask`/`answer`/`poll` (no reaching into `Rendezvous` or the task maps) — asks twice in one resumed turn, and asserts the ticket surfaces the second question (`Phase.ASKING`, `"Q2"`) instead of `FAILED`, then finishes normally (`Phase.DONE`, `"done"`) once the primary answers the second question and the worker replies for real. Reverted the production change (kept the test), reran it, and it went red: ``` org.opentest4j.AssertionFailedError: expected: <ASKING> but was: <FAILED> at dev.ltms.fleet.msg.MessageServiceTest.awaitTicketPhase(MessageServiceTest.java:1468) at dev.ltms.fleet.msg.MessageServiceTest.secondFleetAskInTheSameResumedTurnDoesNotKillTheAsyncTicket(MessageServiceTest.java:1019) ``` Restored the production fix; reran — green again. ### Build `cd fleetd && mvn clean install`, unpiped, three full clean runs after restoring the fix: `Tests run: 1284, Failures: 0, Errors: 0, Skipped: 0` / `BUILD SUCCESS` (reproduced twice at the full-suite level, plus 5 consecutive isolated runs of `MessageServiceTest` alone — 71 tests, 0 failures each time — to rule out flakiness in the new test itself). One caveat for the reviewer: an earlier, unsaved `mvn clean install` run (before I started capturing output to a file) exited 1 with truncated tool output I could not fully read. Two subsequent full clean builds captured to a log both came back `BUILD SUCCESS` with the identical `Tests run: 1284, Failures: 0, Errors: 0` count, and `MessageServiceTest` alone was green 5/5 times in isolation, so this looks like a pre-existing flaky test elsewhere in the 1284-test suite unrelated to this change — but I could not name the specific test from that first run's output, so I'm reporting it rather than asserting it away. ### Files changed - `fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java` - `fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java`
agent added 1 commit 2026-09-04 05:53:34 +02:00
#282: don't let a chained fleet_ask kill its own async ticket
CI / contract (pull_request) Successful in 1m6s
CI / build (pull_request) Successful in 1m38s
6ed70700a0
answer() opened a fresh forward waiter but, unlike send(), never
registered it in asyncTasksByWaiter. So when a worker chained a
second fleet_ask inside the same resumed turn (before calling
fleet_reply), markAsyncQuestion had no Task to re-associate, and
answer() then completed the async ticket's future with the second
QUESTION as if it were a terminal reply — fleet_poll reported FAILED
while the worker was still alive and mid-conversation.

Fix: register answer()'s waiter in asyncTasksByWaiter (mirroring
send()) so a chained ask can re-arm the ticket under its new turnId,
and guard answer()'s finishAsyncTask call the same way sendAsync's
own lambda already does (skip on Outcome.QUESTION). Also drop the
stale asyncTasksByTurn entry left behind when markAsyncQuestion
re-arms a task under a new turnId, a leak the fix makes reachable
for the first time.

Reachability confirmed by driving the exact sequence through the
public API (sendAsync -> ask -> answer -> ask again) in a new test;
reverting the production change makes it fail with
"expected: <ASKING> but was: <FAILED>", confirming it catches the
regression.
ltms closed this pull request 2026-09-04 06:00:07 +02:00
Some checks are pending
CI / contract (pull_request) Successful in 1m6s
CI / build (pull_request) Successful in 1m38s

Pull request closed

Sign in to join this conversation.