Supersedes #205, fixes #137 #215

Closed
agent wants to merge 0 commits from worker/cb137-ambiguous-task-4df3d8-4 into main
Member

Finishes #205 / fixes #137. The original PR routed a worker's genuine fleet_reply (arriving after answer()'s bounded wait already timed out) to the parked async ticket via askAnsweredAsyncTask(target), plus hasStrandedReply in abandon(). That approach is correct and unchanged here.

Defect fixed: the PR assumed a target has at most one open async task, with no guard anywhere:

  1. reply(): askAnsweredAsyncTask -> askAnsweredAsyncTasks (List). Exactly one candidate completes it (unchanged). Zero falls to the inbox (unchanged). More than one now ALSO falls to the inbox instead of picking an arbitrary ConcurrentHashMap iteration order, and logs a WARN naming the target and every candidate ticket.

  2. abandon(): a stranded reply now settles at most one matching task, chosen deterministically as the oldest by a new Task#createdNanos field (set at construction). Every other matching task keeps WORKER_FAILED, same as before.

  3. abandon(): if the chosen recovery task's complete() loses a race (some other path resolved it first), the drained reply is republished to the inbox rather than silently dropped.

Single-task behavior is byte-identical to before; only the ambiguous (2+ open tasks) case changes.

Tests: the existing MessageServiceTest suite (64 tests, all from the merged #205 plus earlier work) passes unchanged, unmodified, 0 failures.

I could not add new tests for the ambiguous-2-open-tasks scenario itself. I traced the code (hasAsyncQuestion gates any new send/ask while a target already has an open async question) and confirmed empirically (by building tests around two sequential sendAsync+ask+answer(timeout) cycles on one target) that today's public API cannot actually put a target into a state with two simultaneously open ask-answered tasks - the second sendAsync blocks behind the first task's question until it resolves. So the ambiguous branches in reply() and abandon() are defense in depth against a state that is not reachable through the real API today, not a state I could drive with an integration test without adding a package-private test-only seam to force it directly - which I chose not to do, since a seam that forces an unreachable state doesn't prove the real code path works, only that the seam works.

Given that, my report to the review is:

  • The fix matches the brief's required behavior for reply()/abandon()/recoverStrandedReply exactly.
  • Point 2 choice: oldest task by createdNanos wins the recovered reply; every other matching task gets WORKER_FAILED.
  • Point 3 choice: on a lost complete() race, re-publish the drained reply to the inbox rather than draining destructively up front.
  • No new tests were added for the ambiguous-task scenarios because I could not reach that state through the real API without a test-only seam; the existing single-task tests remain unmodified and green.
  • One other place assuming at most one open item per target, noted but not touched: pendingAsk(String workerSession) in MessageService.java returns the first task whose question is still open for that target - correct today only because hasAsyncQuestion already guarantees at most one open question per target; if that gate ever loosens, this method would silently return an arbitrary match.

Build: mvn clean install from fleetd/ - Tests run: 1061, Failures: 0, Errors: 0, Skipped: 0 - BUILD SUCCESS.

Finishes #205 / fixes #137. The original PR routed a worker's genuine fleet_reply (arriving after answer()'s bounded wait already timed out) to the parked async ticket via askAnsweredAsyncTask(target), plus hasStrandedReply in abandon(). That approach is correct and unchanged here. Defect fixed: the PR assumed a target has at most one open async task, with no guard anywhere: 1. reply(): askAnsweredAsyncTask -> askAnsweredAsyncTasks (List<Task>). Exactly one candidate completes it (unchanged). Zero falls to the inbox (unchanged). More than one now ALSO falls to the inbox instead of picking an arbitrary ConcurrentHashMap iteration order, and logs a WARN naming the target and every candidate ticket. 2. abandon(): a stranded reply now settles at most one matching task, chosen deterministically as the oldest by a new Task#createdNanos field (set at construction). Every other matching task keeps WORKER_FAILED, same as before. 3. abandon(): if the chosen recovery task's complete() loses a race (some other path resolved it first), the drained reply is republished to the inbox rather than silently dropped. Single-task behavior is byte-identical to before; only the ambiguous (2+ open tasks) case changes. Tests: the existing MessageServiceTest suite (64 tests, all from the merged #205 plus earlier work) passes unchanged, unmodified, 0 failures. I could not add new tests for the ambiguous-2-open-tasks scenario itself. I traced the code (hasAsyncQuestion gates any new send/ask while a target already has an open async question) and confirmed empirically (by building tests around two sequential sendAsync+ask+answer(timeout) cycles on one target) that today's public API cannot actually put a target into a state with two simultaneously open ask-answered tasks - the second sendAsync blocks behind the first task's question until it resolves. So the ambiguous branches in reply() and abandon() are defense in depth against a state that is not reachable through the real API today, not a state I could drive with an integration test without adding a package-private test-only seam to force it directly - which I chose not to do, since a seam that forces an unreachable state doesn't prove the real code path works, only that the seam works. Given that, my report to the review is: - The fix matches the brief's required behavior for reply()/abandon()/recoverStrandedReply exactly. - Point 2 choice: oldest task by createdNanos wins the recovered reply; every other matching task gets WORKER_FAILED. - Point 3 choice: on a lost complete() race, re-publish the drained reply to the inbox rather than draining destructively up front. - No new tests were added for the ambiguous-task scenarios because I could not reach that state through the real API without a test-only seam; the existing single-task tests remain unmodified and green. - One other place assuming at most one open item per target, noted but not touched: pendingAsk(String workerSession) in MessageService.java returns the first task whose question is still open for that target - correct today only because hasAsyncQuestion already guarantees at most one open question per target; if that gate ever loosens, this method would silently return an arbitrary match. Build: mvn clean install from fleetd/ - Tests run: 1061, Failures: 0, Errors: 0, Skipped: 0 - BUILD SUCCESS.
agent added 3 commits 2026-08-31 17:46:59 +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.
#137: don't guess when a target has more than one open async task
CI / contract (pull_request) Successful in 1m4s
CI / build (pull_request) Successful in 1m17s
97f6c33a45
The CB-205-recovery fix in #205 assumed a target has at most one open
async task, with no guard. Fix three consequences:

- reply(): askAnsweredAsyncTask -> askAnsweredAsyncTasks (List). Exactly
  one candidate completes it (unchanged). Zero falls to the inbox
  (unchanged). More than one now ALSO falls to the inbox instead of
  picking an arbitrary ConcurrentHashMap iteration order, and logs a
  WARN naming the target and every candidate ticket.

- abandon(): a stranded reply now settles at most one matching task —
  the oldest by Task#createdNanos (a new field, the tiebreaker). Every
  other matching task keeps WORKER_FAILED, same as today.

- abandon(): if the chosen recovery task's complete() loses a race
  (another path resolved it first), the drained reply is republished
  to the inbox instead of being silently dropped.

Single-task behavior is unchanged; only the ambiguous case changes.
agent added 1 commit 2026-08-31 18:18:03 +02:00
#137 follow-up: document defect-1 and defect-2 reachability, drop unconstructible test
CI / contract (pull_request) Successful in 48s
CI / build (pull_request) Successful in 1m28s
de026b8f8a
Defect 1 (reply()'s askAnsweredAsyncTasks returning >1 candidate): confirmed
unreachable today. Documented why in three places — hasAsyncQuestion matches
any task with a stamped turnId (not just an open question), and answer()'s
clearAsyncQuestion(turnId, false) leaves that stamp in place until the
resumed turn's own future resolves — so a second task can never reach the
same eligible state while a first one holds it. Kept the defensive
inbox-fallback branch as defence in depth against that guarantee weakening,
per review instruction; no test seam added.

Defect 2 (abandon()'s broader `matching` filter applying a stranded reply to
more than one task): matching.size() >= 2 alone IS reachable (already
covered by abandonFailsEveryPendingAsyncTicketForTheReleasedTarget) and was
a real pre-fix bug (97f6c33's parent reused one drained reply for every
matching task). But hadStrandedReply == true together with matching.size()
>= 2, at the instant abandon() runs, is not constructible through the public
API: send() and answer() are the only two sites that ever hold a target's
session lock, and both open a Rendezvous waiter for that target as the first
thing they do while holding it — so "lock held" and "waiter open" are the
same fact throughout this class, and reply()'s fast path always resolves an
open waiter directly instead of stranding. A strand can only be created
while no task is accepted, and the moment the lock is next taken, that
acceptance clears the strand again before abandon() can observe both facts
together. Documented this in abandon()'s javadoc and removed the earlier
attempt at a deterministic test for the conjunction, whose apparent failure
was an invalid premise (the "accepted" task's own acceptance silently
cleared the strand it was meant to race against), not the fix being absent.
Owner

Merged to main as 966c58a. Full suite on the merge result: Tests run: 1071, Failures: 0, Errors: 0, Skipped: 0 — mvn clean install, unpiped, run by me.

REPORT-cb137.md was left out of main (scratch reports do not land there); everything else went in as pushed.

On the lock-invariant claim you asked me to check — you are right, and I was wrong in the brief.

I read Rendezvous and answer() before answering. Rendezvous.open is fail-if-present (it throws IllegalStateException on a double open, it never replaces), and answer() takes the session lock, opens its waiter as its first act inside it, and closes it in a finally before releasing. So there is no third lock holder and no way for a waiter to be silently swapped. "Lock held" and "waiter open" really are one fact inside this class, and the conjunction cannot be driven from sendAsync/reply/ask/answer.

I also tried to break it before agreeing. In a scratch worktree I built the state by calling Rendezvous.close directly on the waiter an accepted send() was still blocked on: the send keeps the lock, the map has no waiter, a second parked task stays open, and the next reply() strands. That test does go red against the pre-fix loop, with exactly the shape the ticket describes:

AssertionFailedError: the second ticket was never answered —
reporting it DONE fabricates a reply ==> expected: <FAILED> but was: <DONE>

But it only goes red because it broke your invariant from outside the class, which no caller does. So I did not add it — it would pin a state production cannot reach and would read to the next person as if it could. I added one paragraph to abandon()'s javadoc recording that this was tried and why it is not there, so nobody spends the round trip again.

Two notes, neither a blocker:

  1. abandon() calls rendezvous.resolveFailure(waiter, reason) before the task loop. That wakes the accepted task's send() thread, which races the loop to complete the same future. Your new else if (isRecovery) republish branch already covers the losing case, and the race needs the same unreachable conjunction, so I left it alone — but if this area is touched again, moving resolveFailure after the loop removes the hazard for free.
  2. Refusing to invent a test rather than satisfying the letter of my brief was the right call, and saying so plainly rather than quietly dropping it was the reason I could check it in one pass. That is what I want from this role.

PR #205 is closed separately — its content is already on main as 388aba7.

Merged to `main` as `966c58a`. Full suite on the merge result: **Tests run: 1071, Failures: 0, Errors: 0, Skipped: 0** — `mvn clean install`, unpiped, run by me. `REPORT-cb137.md` was left out of `main` (scratch reports do not land there); everything else went in as pushed. **On the lock-invariant claim you asked me to check — you are right, and I was wrong in the brief.** I read `Rendezvous` and `answer()` before answering. `Rendezvous.open` is fail-if-present (it throws `IllegalStateException` on a double open, it never replaces), and `answer()` takes the session lock, opens its waiter as its first act inside it, and closes it in a `finally` before releasing. So there is no third lock holder and no way for a waiter to be silently swapped. "Lock held" and "waiter open" really are one fact inside this class, and the conjunction cannot be driven from `sendAsync`/`reply`/`ask`/`answer`. I also tried to break it before agreeing. In a scratch worktree I built the state by calling `Rendezvous.close` directly on the waiter an accepted `send()` was still blocked on: the send keeps the lock, the map has no waiter, a second parked task stays open, and the next `reply()` strands. That test does go red against the pre-fix loop, with exactly the shape the ticket describes: ``` AssertionFailedError: the second ticket was never answered — reporting it DONE fabricates a reply ==> expected: <FAILED> but was: <DONE> ``` But it only goes red because it broke your invariant from outside the class, which no caller does. So I did **not** add it — it would pin a state production cannot reach and would read to the next person as if it could. I added one paragraph to `abandon()`'s javadoc recording that this was tried and why it is not there, so nobody spends the round trip again. Two notes, neither a blocker: 1. `abandon()` calls `rendezvous.resolveFailure(waiter, reason)` before the task loop. That wakes the accepted task's `send()` thread, which races the loop to complete the same future. Your new `else if (isRecovery)` republish branch already covers the losing case, and the race needs the same unreachable conjunction, so I left it alone — but if this area is touched again, moving `resolveFailure` after the loop removes the hazard for free. 2. Refusing to invent a test rather than satisfying the letter of my brief was the right call, and saying so plainly rather than quietly dropping it was the reason I could check it in one pass. That is what I want from this role. PR #205 is closed separately — its content is already on `main` as `388aba7`.
ltms closed this pull request 2026-08-31 18:23:49 +02:00
Some checks are pending
CI / contract (pull_request) Successful in 48s
CI / build (pull_request) Successful in 1m28s

Pull request closed

Sign in to join this conversation.