#137: complete an async ticket's reply after answer() times out #205

Closed
agent wants to merge 0 commits from worker/cb-137-ask-ticket-e7760c-2 into main
Member

Fixes fleetd issue #137: after a fleet_ask round-trip, a worker's final fleet_reply was orphaning its async ticket and later reporting a false failure.

Root cause

MessageService.answer() (behind fleet_send{turnId}) blocks the primary only for its own bounded MCP-call window (25s default, 120s max) — much shorter than a worker's resumed turn can genuinely take (more edits, a build, a commit, a push, opening a PR). When that window times out, answer() closes its rendezvous waiter. The worker's eventual fleet_reply then has no live waiter to resolve, falls back to the session inbox, and the async ticket's future is never completed — fleet_poll{ticket} stays PENDING forever until fleet_stop's abandon() forces it FAILED with the misleading "the worker session was released before it replied; worktree=... branch=... snapshot=..." text, even though the reply genuinely arrived.

Fix (MessageService.java)

  1. reply(): before falling back to the inbox, look for the async task this exact turn belongs to (already answered — question cleared, turnId still stamped — but not yet resolved) and complete it directly with the real reply. fleet_poll{ticket} now returns the actual reply.
  2. abandon(): defense in depth, independent of (1) — never write a false WORKER_FAILED once a reply reached the inbox for the target; recover and use its real content instead. This also means the misleading worktree/branch/snapshot hint never prints once a reply exists.

Tests

Two new tests in MessageServiceTest drive the full real delegation path (async send → ask → answer with a short timeout that genuinely expires → reply → poll/abandon) rather than calling a reply sink directly. Verified both fail with the fix disabled (PENDING forever / false failure) and pass with it restored — full transcript in REPORT-cb137.md.

Build

mvn clean install
Tests run: 1039, Failures: 0, Errors: 0, Skipped: 0
BUILD SUCCESS

(main was 1037 tests; +2 new tests here, all green)

See REPORT-cb137.md in this PR for the full writeup, including what I could not verify (no IDE tooling, no live daemon dogfood).

Fixes fleetd issue #137: after a `fleet_ask` round-trip, a worker's final `fleet_reply` was orphaning its async ticket and later reporting a false failure. ## Root cause `MessageService.answer()` (behind `fleet_send{turnId}`) blocks the primary only for its own bounded MCP-call window (25s default, 120s max) — much shorter than a worker's resumed turn can genuinely take (more edits, a build, a commit, a push, opening a PR). When that window times out, `answer()` closes its rendezvous waiter. The worker's eventual `fleet_reply` then has no live waiter to resolve, falls back to the session inbox, and the async ticket's future is never completed — `fleet_poll{ticket}` stays `PENDING` forever until `fleet_stop`'s `abandon()` forces it `FAILED` with the misleading "the worker session was released before it replied; worktree=... branch=... snapshot=..." text, even though the reply genuinely arrived. ## Fix (`MessageService.java`) 1. `reply()`: before falling back to the inbox, look for the async task this exact turn belongs to (already answered — question cleared, `turnId` still stamped — but not yet resolved) and complete it directly with the real reply. `fleet_poll{ticket}` now returns the actual reply. 2. `abandon()`: defense in depth, independent of (1) — never write a false `WORKER_FAILED` once a reply reached the inbox for the target; recover and use its real content instead. This also means the misleading worktree/branch/snapshot hint never prints once a reply exists. ## Tests Two new tests in `MessageServiceTest` drive the full real delegation path (async send → ask → answer with a short timeout that genuinely expires → reply → poll/abandon) rather than calling a reply sink directly. Verified both fail with the fix disabled (`PENDING` forever / false failure) and pass with it restored — full transcript in `REPORT-cb137.md`. ## Build ``` mvn clean install Tests run: 1039, Failures: 0, Errors: 0, Skipped: 0 BUILD SUCCESS ``` (main was 1037 tests; +2 new tests here, all green) See `REPORT-cb137.md` in this PR for the full writeup, including what I could not verify (no IDE tooling, no live daemon dogfood).
agent added 1 commit 2026-08-31 09:12:48 +02:00
#137: complete an async ticket's own reply after answer() times out
CI / contract (pull_request) Successful in 49s
CI / build (pull_request) Successful in 1m44s
ee5f8b932b
fleet_send{turnId} (MessageService.answer) blocks the primary only for its
own bounded MCP-call window (25s default, 120s max) — far shorter than a
worker's resumed turn can genuinely take. When that window expires, answer()
closes its rendezvous waiter, so the worker's eventual fleet_reply has no
live waiter to resolve and falls back to the session inbox. The async
ticket's future was never completed by that path, so fleet_poll{ticket}
stayed PENDING until fleet_stop's abandon() forced it FAILED with a
misleading "the worker session was released before it replied" reason,
even though the reply had genuinely arrived.

- MessageService.reply(): before falling to the inbox, look for the async
  task this exact turn belongs to (already answered — question cleared,
  turnId still stamped — but not yet resolved) and complete it directly
  with the real reply, so fleet_poll{ticket} returns it.
- MessageService.abandon(): defense in depth, independent of the above —
  never write a false WORKER_FAILED once a reply reached the inbox for
  this target; recover and use its real content instead.
- Two new tests drive the full delegation path (async send -> ask ->
  answer with a short timeout -> reply -> poll/abandon), not a reply sink
  directly; both fail with the fix disabled and pass with it restored.
Owner

Reviewed. This solves a real problem and the reasoning in the comments is good — a worker's genuine fleet_reply landing after answer() gave up should not leave the ticket PENDING until abandon() mislabels it "released before it replied". The hasStrandedReply check before writing a failure is the right instinct.

Not merging yet. Two defects, both from the same missing premise: a target can have more than one open async task at a time. sendAsync (line ~809) just mints a ticket and does tasks.put(ticket, task) — there is no guard, no "one open task per target" check, nothing that rejects a second fleet_send{wait:false} at the same member. In practice this fleet does queue a second send at a busy member.

1. askAnsweredAsyncTask returns an arbitrary task

for (Task task : tasks.values()) {
    if (target.equals(task.target) && task.question == null && task.turnId != null
            && !task.future.isDone()) {
        return task;
    }
}

tasks is a ConcurrentHashMap; values() has no defined order. With two open ask-answered tasks on one target, the worker's reply completes whichever the iteration happens to reach first. The lead then reads a real reply under the wrong ticket.

2. abandon applies one recovered reply to every task

recovered is drained once and then reused for every matching task in the loop, so two open tasks both complete REPLIED with the same text. One worker answer is fabricated as the answer to two different delegations.

Both failures are quiet and hand the lead something that reads like a correct answer. That is worse than a failure, because the lead acts on it — the same class of defect #164 exists to remove.

3. Smaller: the recovery drain is destructive

recoverStrandedReply calls drainReplies(target), which empties the inbox. If the following task.future.complete(...) returns false, the reply is drained and never delivered. The isDone() check above makes the window small, but the consequence is a lost reply, so it is worth closing.

What I would like instead

When the match is ambiguous, do today's thing. Guessing is the only unsafe option here:

  • In reply(...), collect all candidates rather than returning the first. Exactly one → complete it, as now. Zero → fall through to the inbox, as now. More than one → fall through to the inbox as well, and log a WARN naming the target and the candidate tickets. Stranding a reply is recoverable; resolving the wrong ticket is not.
  • In abandon(...), apply the recovered reply to at most one task, chosen deterministically (the oldest by the task's nowNanos stamp — say which you picked). Every other matching task keeps the WORKER_FAILED path it has today.
  • In recoverStrandedReply, only drain once you know a task will accept it, or re-publish if complete returns false. Either is fine; say which you chose.

Tests to add

  • Two open ask-answered tasks on one target, then a fleet_reply: neither ticket is completed, the reply is in the inbox, and the WARN names both tickets.
  • Two open tasks on one target, then abandon: exactly one is REPLIED, the other is WORKER_FAILED, and the reply text appears once.
  • The existing single-task behaviour is unchanged (pin your current tests).
  • Watch each new test fail before keeping it.

The one-task path you built is right and should not change — this is only about refusing to guess when there is more than one.

Reviewed. This solves a real problem and the reasoning in the comments is good — a worker's genuine `fleet_reply` landing after `answer()` gave up should not leave the ticket PENDING until `abandon()` mislabels it "released before it replied". The `hasStrandedReply` check before writing a failure is the right instinct. **Not merging yet.** Two defects, both from the same missing premise: **a target can have more than one open async task at a time.** `sendAsync` (line ~809) just mints a ticket and does `tasks.put(ticket, task)` — there is no guard, no "one open task per target" check, nothing that rejects a second `fleet_send{wait:false}` at the same member. In practice this fleet does queue a second send at a busy member. ### 1. `askAnsweredAsyncTask` returns an arbitrary task ```java for (Task task : tasks.values()) { if (target.equals(task.target) && task.question == null && task.turnId != null && !task.future.isDone()) { return task; } } ``` `tasks` is a `ConcurrentHashMap`; `values()` has no defined order. With two open ask-answered tasks on one target, the worker's reply completes whichever the iteration happens to reach first. The lead then reads a real reply **under the wrong ticket**. ### 2. `abandon` applies one recovered reply to every task `recovered` is drained once and then reused for every matching task in the loop, so two open tasks both complete `REPLIED` with the same text. One worker answer is fabricated as the answer to two different delegations. Both failures are quiet and hand the lead something that *reads like* a correct answer. That is worse than a failure, because the lead acts on it — the same class of defect #164 exists to remove. ### 3. Smaller: the recovery drain is destructive `recoverStrandedReply` calls `drainReplies(target)`, which empties the inbox. If the following `task.future.complete(...)` returns false, the reply is drained and never delivered. The `isDone()` check above makes the window small, but the consequence is a lost reply, so it is worth closing. ## What I would like instead **When the match is ambiguous, do today's thing.** Guessing is the only unsafe option here: - In `reply(...)`, collect *all* candidates rather than returning the first. Exactly one → complete it, as now. Zero → fall through to the inbox, as now. **More than one → fall through to the inbox as well, and log a WARN naming the target and the candidate tickets.** Stranding a reply is recoverable; resolving the wrong ticket is not. - In `abandon(...)`, apply the recovered reply to **at most one** task, chosen deterministically (the oldest by the task's `nowNanos` stamp — say which you picked). Every other matching task keeps the `WORKER_FAILED` path it has today. - In `recoverStrandedReply`, only drain once you know a task will accept it, or re-publish if `complete` returns false. Either is fine; say which you chose. ## Tests to add - Two open ask-answered tasks on one target, then a `fleet_reply`: neither ticket is completed, the reply is in the inbox, and the WARN names both tickets. - Two open tasks on one target, then `abandon`: exactly one is `REPLIED`, the other is `WORKER_FAILED`, and the reply text appears once. - The existing single-task behaviour is unchanged (pin your current tests). - Watch each new test fail before keeping it. The one-task path you built is right and should not change — this is only about refusing to guess when there is more than one.
Owner

Closing: this work is already on main as 388aba7 ("#137: an answered turn's reply completes its own ticket, instead of a false failure"). I checked the content, not just the title — main carries askAnsweredAsyncTask, the async-recovered metric label, and both tests (aReplyAfterAnswerTimesOutStillCompletesTheAsyncTicket, fleetStopAfterAnOrphanedReplyDoesNotFailTheTicket).

The follow-up work — do not guess when a target has more than one open async task — continues in #215.

Closing: this work is already on `main` as `388aba7` ("#137: an answered turn's reply completes its own ticket, instead of a false failure"). I checked the content, not just the title — `main` carries `askAnsweredAsyncTask`, the `async-recovered` metric label, and both tests (`aReplyAfterAnswerTimesOutStillCompletesTheAsyncTicket`, `fleetStopAfterAnOrphanedReplyDoesNotFailTheTicket`). The follow-up work — do not guess when a target has more than one open async task — continues in #215.
ltms closed this pull request 2026-08-31 18:02:04 +02:00
Some checks are pending
CI / contract (pull_request) Successful in 49s
CI / build (pull_request) Successful in 1m44s

Pull request closed

Sign in to join this conversation.