Compare commits
2 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| ee5f8b932b | |||
| 4ac688b6d9 |
+166
@@ -0,0 +1,166 @@
|
||||
# CB-137 / fleetd issue #137 — report
|
||||
|
||||
## Real root cause (not the hypothesis in the ticket)
|
||||
|
||||
I read `MessageService.java` and `Rendezvous.java` before changing anything. The mechanism is real,
|
||||
but the exact place it happens is `MessageService.answer()`, not "the reply goes to the inbox on
|
||||
purpose" in general.
|
||||
|
||||
1. A lead delegates with `fleet_send{wait:false}` → `sendAsync()` creates a `Task` and runs `send()`
|
||||
on a background virtual thread with a 30-minute internal budget (`ASYNC_TIMEOUT_MS`).
|
||||
2. The worker calls `fleet_ask`. That resolves the open rendezvous waiter with `Kind.QUESTION`, so
|
||||
`send()` returns immediately and the `Task` is left open (its `future` stays unresolved — see the
|
||||
comment in `sendAsync`'s lambda: "Keep the accepted owner until answer() finishes it").
|
||||
3. The lead answers with `fleet_send{turnId, content}`. This calls `FleetMcp.answer()` →
|
||||
`MessageService.answer(turnId, content, timeout)`. The `timeout` here is **not** the generous
|
||||
30-minute async budget — it is the MCP tool's own bounded wait: `DEFAULT_TIMEOUT_MS = 25_000`,
|
||||
clamped to at most `MAX_TIMEOUT_MS = 120_000` (`FleetMcp.java:71-72,495,512`). This is the same
|
||||
~60–120s window every blocking `fleet_send` call is capped at (documented elsewhere as "the
|
||||
caller's own MCP client call timeout").
|
||||
4. `answer()` opens a **fresh** rendezvous waiter for the worker session and blocks on it for at most
|
||||
that window. If the worker's resumed turn takes longer than that to actually finish (very
|
||||
plausible — the resumed turn can mean more edits, a build, a commit, a push, opening a PR), the
|
||||
wait times out. On timeout, `answer()`'s `finally` block unconditionally calls
|
||||
`rendezvous.close(workerSession, reply)`, **removing the waiter from the map**, and returns
|
||||
`Outcome.TIMED_OUT_WORKING` to the lead.
|
||||
5. The worker keeps working, unaware anything happened, and eventually calls `fleet_reply`. That
|
||||
reaches `MessageService.reply(session, content)`, which tries `rendezvous.resolve(session,
|
||||
content)` — but the waiter was already closed in step 4, so `resolve` returns `false`. `reply()`
|
||||
then falls back to `inbox.publish(...)` and marks `strandedReplies.put(session, true)`
|
||||
(CB-640 bookkeeping) — the reply is safely held, but **the async `Task`'s `future` is never
|
||||
completed**.
|
||||
6. `fleet_poll{ticket}` keeps returning `PENDING` forever (the `Task` never resolves) — until the
|
||||
lead eventually calls `fleet_stop`. That fires `sessions.onRelease` → `messages.abandon(target,
|
||||
reason)` (`Fleetd.java:481-497`), where `reason` is built with the exact text from the bug report
|
||||
("the worker session was released before it replied; worktree=... branch=... snapshot=...",
|
||||
`Fleetd.java:484-487`). `abandon()`'s loop finds the still-open `Task` (`question == null`, future
|
||||
not done) and completes it as `WORKER_FAILED` with that misleading reason — even though the
|
||||
worker's real reply is sitting, intact, in the inbox the whole time.
|
||||
|
||||
So: the reported behaviour is correct, and the specific trigger is `answer()`'s own bounded wait
|
||||
being shorter than the worker's real resumed-turn time — not anything to do with the ~55s
|
||||
`fleet_ask` window itself (that part, issue #61, is untouched).
|
||||
|
||||
## Fix
|
||||
|
||||
Two changes in `fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java`, both scoped to the
|
||||
ticket/reply routing and the terminal-state text — `fleet_ask`'s own window and mechanics are
|
||||
untouched.
|
||||
|
||||
**1. `reply()` — priority 1 (the ticket resolves with the real reply).**
|
||||
Before falling back to the inbox, `reply()` now looks for an async `Task` that is specifically in the
|
||||
"already answered but not yet resolved" state (`question == null`, `turnId != null` — set once
|
||||
`answer()` has cleared the question but before anything completed the future, `!future.isDone()`).
|
||||
If one exists for this `target`, the worker's reply completes that `Task`'s future directly as
|
||||
`Outcome.REPLIED` with the real content, and the reply never touches the inbox at all. A task that
|
||||
was never asked has `turnId == null` and can never match, so ordinary (no-`fleet_ask`) delegations
|
||||
are unaffected — they already resolve through the pre-existing rendezvous fast path.
|
||||
|
||||
I chose this over leaving `answer()`'s own timeout behaviour untouched and instead keeping its
|
||||
rendezvous waiter open in the background: that alternative works but reopens the "at most one
|
||||
waiter per session" invariant (`Rendezvous.open` throws on a double-open) to a new class of races
|
||||
with a fresh send arriving mid-window. The `send()` path already guards against sending into an
|
||||
answered-but-still-resolving worker via `hasAsyncQuestion(target)` (checks `asyncTasksByTurn`,
|
||||
which still holds the task until it resolves), so routing through `reply()` gets the same protection
|
||||
without touching `answer()`'s waiter lifecycle at all — the smaller, safer diff.
|
||||
|
||||
**2. `abandon()` — priority 3 (required independently, "even if you fix (1)").**
|
||||
Before marking any of a released target's still-open tasks `WORKER_FAILED`, `abandon()` now checks
|
||||
`hasStrandedReply(target)` (the existing CB-640 fact — true whenever the *last* `reply()` for this
|
||||
target fell through to the inbox). If true, it drains the inbox (`recoverStrandedReply`) and — if it
|
||||
actually finds a message — completes the task as `REPLIED` with that real content instead of writing
|
||||
the failure. This is deliberately a **separate** check from fix 1: fix 1 already prevents the
|
||||
inbox-stranding from happening in the exact scenario this ticket describes, so by the time
|
||||
`abandon()` runs the task is normally already resolved and `abandon()`'s `complete()` call is a
|
||||
harmless no-op. This second check exists so that if some *other* future path ever strands a reply
|
||||
in the inbox without resolving its ticket, `abandon()` still refuses to report a false failure —
|
||||
"if a reply reached any sink for that turn, the terminal state is done," per the ticket. I verified
|
||||
both are required by disabling each independently and confirming the two new tests fail (see below).
|
||||
|
||||
**Priority 4 (the snapshot/worktree hint).** Handled as a consequence of both fixes rather than a
|
||||
separate branch: once a task resolves as `REPLIED` (via either fix), `abandon()` never calls
|
||||
`new Reply(Outcome.WORKER_FAILED, reason)` for that task at all, so the "the worker session was
|
||||
released before it replied; worktree=... branch=... snapshot=..." text is never constructed or
|
||||
attached to that ticket's outcome. It still appears, correctly, for a task that never got a reply
|
||||
(the existing `abandonFailsEveryPendingAsyncTicketForTheReleasedTarget` /
|
||||
`anAbandonedAsyncTaskPollsAsFailedNotPending` tests still pass unchanged).
|
||||
|
||||
**Priority 2** was not needed — fix 1 makes `fleet_poll{ticket}` return the actual reply (the
|
||||
higher-priority option), so I did not fall back to "the ticket merely resolves as done with no
|
||||
content."
|
||||
|
||||
## Tests — driven through the real delegation path, not the reply sink directly
|
||||
|
||||
Both new tests in `fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java` go through
|
||||
`sendAsync` → `injectDelivery` → `ask` → `answer` (with a short timeout, so it genuinely times out,
|
||||
mirroring the ~25–120s real MCP-call bound vs. a longer resumed turn) → `reply` → `poll`/`abandon`.
|
||||
No test constructs a `Reply` and hands it to a sink directly.
|
||||
|
||||
- `aReplyAfterAnswerTimesOutStillCompletesTheAsyncTicket` — asserts `fleet_poll{ticket}` (via
|
||||
`messages.poll`) reaches `Phase.DONE` with the worker's actual reply text and
|
||||
`replySource() == "reply"`, and that `hasStrandedReply(T)` stays `false` (proves the reply never
|
||||
touched the inbox at all — fix 1 caught it).
|
||||
- `fleetStopAfterAnOrphanedReplyDoesNotFailTheTicket` — same setup, then calls `abandon(T, "the
|
||||
worker session was released before it replied")` (what `fleet_stop` triggers) and asserts it
|
||||
returns `false` (no failure recorded) and the ticket still polls `DONE` with the real reply.
|
||||
|
||||
**Proof both fail without the change.** I temporarily short-circuited both new private methods
|
||||
(`askAnsweredAsyncTask` → always `null`, `recoverStrandedReply` → always `null`) — i.e. disabled
|
||||
both fixes — and ran just these two tests:
|
||||
|
||||
```
|
||||
[ERROR] Tests run: 2, Failures: 2, Errors: 0, Skipped: 0
|
||||
dev.ltms.fleet.msg.MessageServiceTest.aReplyAfterAnswerTimesOutStillCompletesTheAsyncTicket
|
||||
org.opentest4j.AssertionFailedError: expected: <DONE> but was: <PENDING>
|
||||
dev.ltms.fleet.msg.MessageServiceTest.fleetStopAfterAnOrphanedReplyDoesNotFailTheTicket
|
||||
org.opentest4j.AssertionFailedError: a reply already arrived, so nothing here is a genuine failure
|
||||
==> expected: <false> but was: <true>
|
||||
```
|
||||
|
||||
This is the exact bug: the ticket stays `PENDING` forever, and `abandon()` reports `true` (a
|
||||
failure) even though a reply had already arrived. I then restored both fixes (verified with
|
||||
`grep -n "TEMP #137-proof"` finding nothing) and re-ran — both pass.
|
||||
|
||||
## Build
|
||||
|
||||
Ran from `fleetd/`, unpiped, full output read (not `| tail`):
|
||||
|
||||
```
|
||||
mvn clean install
|
||||
...
|
||||
[INFO] Tests run: 1039, Failures: 0, Errors: 0, Skipped: 0
|
||||
[INFO] BUILD SUCCESS
|
||||
[INFO] Total time: 36.315 s
|
||||
```
|
||||
|
||||
Main was at 1037 tests; this branch adds the 2 new tests above → 1039, all green, `exit=0`.
|
||||
|
||||
## What I could NOT check
|
||||
|
||||
- No IDE tooling is mounted for me (worker), so no `ide_diagnostics`/IntelliJ inspection pass — only
|
||||
`mvn clean install` (compiler + full test suite), as the worker procedure allows.
|
||||
- I cannot restart the daemon or dogfood this live — I have no forge/daemon control. This is
|
||||
unverified against a real herdr pane, a real MCP client's ~60s call cap, or a real worker session;
|
||||
everything above is verified only through the JUnit fixture's simulated timing
|
||||
(`FakeHerdr`/`injector.onStatus`/direct `messages.answer(...,150)` calls), not a live fleet.
|
||||
A primary should still consider a short live dogfood (an async delegation that asks, gets answered,
|
||||
and takes longer than ~2 minutes to reply) before calling this closed.
|
||||
- I did not touch, and did not re-verify, the `fleet_ask` ~55s window itself (issue #61) — out of
|
||||
scope per the brief.
|
||||
|
||||
## Scope note (not investigated further)
|
||||
|
||||
`answer()`'s nested/double-`fleet_ask` case (the worker asks a second question before ever
|
||||
replying to the first answer) has some pre-existing behaviour around which `turnId` a `QUESTION`
|
||||
resolution gets attributed to that I did not fully untangle — it predates this change, my fix does
|
||||
not touch it, and it is unrelated to the reported defect. Flagging only; not investigated further.
|
||||
|
||||
## Handoff
|
||||
|
||||
- Branch: `worker/cb-137-ask-ticket-e7760c-2`
|
||||
- Worktree root: `/Users/dai.ha/LTMS/.bridged-worktrees/734324-2`
|
||||
- Files changed:
|
||||
- `fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java`
|
||||
- `fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java`
|
||||
- `REPORT-cb137.md` (this file)
|
||||
- Build: `Tests run: 1039, Failures: 0, Errors: 0, Skipped: 0` / `BUILD SUCCESS` (verbatim above)
|
||||
-148
@@ -1,148 +0,0 @@
|
||||
# CB-164 report — worker/cb-164-rebase-885863-8
|
||||
|
||||
## Headline finding — please read this before the diff
|
||||
|
||||
**`main` already ships a full fix for issue #164 points 1 and 2, done separately from the
|
||||
rescued branch.** Commit `3bfa828` ("fleetd#164: an empty or suspiciously fast scrape must
|
||||
fail, never resolve as a success") is already an ancestor of `origin/main` (verified with
|
||||
`git merge-base --is-ancestor 3bfa828 origin/main`). It was written the same day as the
|
||||
rescued commit (2026-08-28), by the same author, but on `main` directly. No later commit
|
||||
touches `CompletionResolver.java` after it.
|
||||
|
||||
What `main` already has, before any change of mine:
|
||||
|
||||
- `CompletionResolver.MIN_TURN_NANOS` = 2 seconds — the exact "too fast to be a real
|
||||
completion" floor the brief described, just under a different name and a different
|
||||
mechanism (an injectable `LongSupplier nowNanos` clock read at delivery and again at
|
||||
resolution, stored on the `InFlight` record — no interface or caller changes needed).
|
||||
- A hard fail on any empty or unreadable scrape, via the existing `fail()` →
|
||||
`Rendezvous.resolveFailure()` → `MessageService.Outcome.WORKER_FAILED` path — the same
|
||||
pattern already used for the CB-109 wedge case. `FleetMcp.formatReply` already renders
|
||||
`WORKER_FAILED` as `"[worker failed — turn ended in an unrecoverable state]\n" + text`,
|
||||
never as blank content. I read this chain end to end and confirmed it (not just from the
|
||||
commit message).
|
||||
|
||||
So **"Part 2" as described in the brief — the one called out as the most important — was
|
||||
already done.** What I did **not** find on `main`: the narrow `BACKEND_ERROR` pattern
|
||||
(`(?i)\bAPI Error\s*:`) from the rescued branch. That is genuinely new.
|
||||
|
||||
### Why I did not cherry-pick 851ebca, and what I did instead
|
||||
|
||||
`851ebca` implements the same floor-check idea a second time, but via a structurally
|
||||
different, more invasive mechanism: it threads a real measured `elapsedNanos` through
|
||||
`TurnListener.onTurnComplete`, `Injector` (a new `turnStartedAtNanos` field, a new
|
||||
constructor, a new `LongSupplier` clock), and `Fleetd`'s anonymous `TurnListener`. Cherry-
|
||||
picking that wholesale onto current `main` would have created **two competing timing
|
||||
mechanisms for the same floor check** in the same class — main's already-shipped
|
||||
delivery-clock read, and a second, parameter-threaded one from the rescued branch — with
|
||||
no clear answer for which one should win if they ever disagreed. Reintroducing the
|
||||
Injector/TurnListener/Fleetd signature changes for that also seemed hard to justify: main's
|
||||
already-tested approach measures essentially the same interval (delivery to resolution)
|
||||
with no interface changes at all.
|
||||
|
||||
Given that, I read the instruction *"if you cannot tell which behaviour a hunk is meant to
|
||||
have, keep both behaviours and say so"* as pointing at genuinely uncertain hunks — not at
|
||||
knowingly wiring in two redundant implementations of the identical check. So instead of a
|
||||
mechanical cherry-pick, I hand-ported only the parts of `851ebca` that are **not** already
|
||||
on `main`:
|
||||
|
||||
1. Added `CompletionResolver.BACKEND_ERROR` — the exact narrow pattern from the brief,
|
||||
`(?i)\bAPI Error\s*:`, with the same "keep it narrow" comment.
|
||||
2. Added a classification check for it, placed right after the existing `BACKEND_EXHAUSTED`
|
||||
pattern-match block (same position, same shape: `firstMatchingLine` against
|
||||
`assistantBlock`, then `fail(target, turn, reason)` naming the member and the matched
|
||||
line) and before the success path (`fleetd/src/main/java/dev/ltms/fleet/inject/
|
||||
CompletionResolver.java`).
|
||||
|
||||
I did **not** port:
|
||||
- `MIN_COMPLETED_TURN_NANOS` / the `elapsedNanos` parameter threading through
|
||||
`TurnListener`/`Injector`/`Fleetd` — redundant with `main`'s already-shipped
|
||||
`MIN_TURN_NANOS`/`LongSupplier` mechanism, and touching those three extra files for no
|
||||
behavioural gain seemed like avoidable risk.
|
||||
- The rescued branch's `visibleTurn = assistantBlock.isBlank() ? raw : assistantBlock`
|
||||
fallback (falls back to the raw pane when the parsed assistant block is blank, so a
|
||||
TUI-hidden error/exhausted line still classifies). This is a real, separate improvement,
|
||||
but on `main`'s current check order the empty-tail fail already fires *before* the
|
||||
exhausted/backend-error pattern checks run, so the fallback would be dead code unless I
|
||||
also reordered those checks ahead of the empty-tail fail. That reorder is a bigger,
|
||||
separate behavioural change than either "Part 1" or "Part 2" asked for, so I left it out
|
||||
and am flagging it here rather than making that call silently. **Out-of-scope note for
|
||||
the lead**, not fixed: a raw-screen-only exhausted/backend-error line (no `⏺` marker, so
|
||||
`lastAssistantBlock` returns blank) is still swallowed into the generic "empty scrape"
|
||||
failure rather than being classified specifically.
|
||||
|
||||
I did not use `fleet_ask` for this: the brief's own instruction for exactly this kind of
|
||||
conflict ("keep both, say so") already pre-authorized me to decide and document rather than
|
||||
block, and a live `fleet_ask` has only a ~55s window per prior fleet notes, so I judged
|
||||
documenting clearly here was more reliable than gambling on that window. If this call is
|
||||
wrong, it's easy to undo — the diff is 3 files, +85/-0 lines.
|
||||
|
||||
## Files changed (worktree-relative)
|
||||
|
||||
- `fleetd/src/main/java/dev/ltms/fleet/inject/CompletionResolver.java` — added the
|
||||
`BACKEND_ERROR` pattern constant and the classification check (+17 lines).
|
||||
- `fleetd/src/test/java/dev/ltms/fleet/inject/CompletionResolverTest.java` — 3 new unit
|
||||
tests for the `BACKEND_ERROR` classification (+46 lines).
|
||||
- `fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java` — 1 new end-to-end test
|
||||
driving the same scenario through `MessageService` (+22 lines).
|
||||
|
||||
No changes to `Fleetd.java`, `Injector.java`, or `TurnListener.java` — see above for why.
|
||||
|
||||
## Build — verbatim result
|
||||
|
||||
Full unpiped `mvn clean install` from `fleetd/`, exit code `0`:
|
||||
|
||||
```
|
||||
[INFO] Tests run: 1036, Failures: 0, Errors: 0, Skipped: 0
|
||||
[INFO] BUILD SUCCESS
|
||||
```
|
||||
|
||||
No `[ERROR]` lines anywhere in the full log (checked with `grep -n "\[ERROR\]"` over the
|
||||
whole log, not a piped/truncated view).
|
||||
|
||||
## Sabotage proof for the new tests
|
||||
|
||||
Commented out the new check in `CompletionResolver.java`:
|
||||
|
||||
```java
|
||||
// String backendError = firstMatchingLine(assistantBlock, BACKEND_ERROR);
|
||||
// if (backendError != null) {
|
||||
// fail(target, turn, "member " + target + " ended on a backend error: " + backendError);
|
||||
// return;
|
||||
// }
|
||||
```
|
||||
|
||||
Ran just the 4 new tests:
|
||||
|
||||
```
|
||||
mvn -q test -Dtest='CompletionResolverTest#classifiesABackendErrorLineAsAFailureInsteadOfACompletedReply+theBackendErrorReasonNamesTheMemberAndCarriesTheMatchedLine+aCaseInsensitiveApiErrorLineIsStillClassifiedAsABackendError,MessageServiceTest#backendErrorScrapeThroughMessageServiceFailsInsteadOfBecomingReplyText'
|
||||
```
|
||||
|
||||
All 4 failed, exactly as expected without the fix:
|
||||
|
||||
```
|
||||
[ERROR] Tests run: 4, Failures: 4, Errors: 0, Skipped: 0
|
||||
[ERROR] CompletionResolverTest.aCaseInsensitiveApiErrorLineIsStillClassifiedAsABackendError:612
|
||||
the pattern is case-insensitive ==> expected: <FAILED> but was: <COMPLETION>
|
||||
[ERROR] CompletionResolverTest.classifiesABackendErrorLineAsAFailureInsteadOfACompletedReply:582
|
||||
a backend rejection is a failure, not a completed reply ==> expected: <FAILED> but was: <COMPLETION>
|
||||
[ERROR] CompletionResolverTest.theBackendErrorReasonNamesTheMemberAndCarriesTheMatchedLine:597
|
||||
the failure names the member: API Error: 400 invalid request body ==> expected: <true> but was: <false>
|
||||
[ERROR] MessageServiceTest.backendErrorScrapeThroughMessageServiceFailsInsteadOfBecomingReplyText:104
|
||||
a backend rejection must use the caller's failure outcome, not a completed reply
|
||||
==> expected: <WORKER_FAILED> but was: <COMPLETED_UNREPLIED>
|
||||
```
|
||||
|
||||
Then restored the fix (put the block back exactly as committed) and re-ran the full
|
||||
`mvn clean install` — green again, see above (`Tests run: 1036, Failures: 0`).
|
||||
|
||||
## What I could not confirm
|
||||
|
||||
- I have no IDE tooling and no way to run the daemon live — only `mvn` in this worktree.
|
||||
Everything reported above is from that command, nothing else.
|
||||
- I did not attempt the "visibleTurn/raw-fallback" reorder described above — flagged as an
|
||||
open, separate item, not fixed, not tested.
|
||||
- Points 3 (broader backend-error surfacing) and 4 (spawn-time profile quarantine) of the
|
||||
issue are untouched, as instructed.
|
||||
- I have not verified this against a live daemon or a real crashed backend — only against
|
||||
the unit/integration test fixtures in this repo.
|
||||
@@ -366,21 +366,40 @@ public final class MessageService {
|
||||
}
|
||||
|
||||
/**
|
||||
* Route a worker's explicit {@code fleet_reply}: resolve an open send, or queue it in the
|
||||
* inbox if no send is currently open. Unlike the bare {@link Rendezvous#resolve}, a no-waiter
|
||||
* result is <em>not</em> a failure — the reply is held for later drain.
|
||||
* Route a worker's explicit {@code fleet_reply}: resolve an open send, complete an async ticket
|
||||
* still parked waiting on this exact turn's answer, or — only once neither applies — queue it in
|
||||
* the inbox. Unlike the bare {@link Rendezvous#resolve}, a no-waiter result is <em>not</em> a
|
||||
* failure — the reply is held for later drain.
|
||||
*
|
||||
* <p><strong>Do NOT use this for mid-turn questions.</strong> {@code fleet_ask} /
|
||||
* {@link Rendezvous#resolveQuestion} must keep today's {@code NO_WAITER} behaviour — questions
|
||||
* are interactive and must never be queued.
|
||||
*
|
||||
* @return always {@code true} — the reply either resolved a live send or was queued
|
||||
* @return always {@code true} — the reply resolved a live send, completed a parked ticket, or
|
||||
* was queued
|
||||
*/
|
||||
public boolean reply(String session, String content) {
|
||||
if (rendezvous.resolve(session, content)) {
|
||||
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.
|
||||
Task orphan = askAnsweredAsyncTask(session);
|
||||
if (orphan != null && orphan.future.complete(new Reply(Outcome.REPLIED, content))) {
|
||||
if (orphan.turnId != null) {
|
||||
asyncTasksByTurn.remove(orphan.turnId, orphan);
|
||||
}
|
||||
count(FleetMetrics.REPLIES, "path", "async-recovered");
|
||||
return true; // the ticket itself took it — no inbox stranding at all
|
||||
}
|
||||
inbox.publish(session, UUID.randomUUID().toString(), content);
|
||||
// CB-640: record the stranding itself (not just the reply text) so fleet health can see a
|
||||
// worker whose replies keep missing their waiter, not only the queue depth this leaves behind.
|
||||
@@ -394,6 +413,24 @@ public final class MessageService {
|
||||
return true; // held, not lost
|
||||
}
|
||||
|
||||
/**
|
||||
* The 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). {@code null} 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}).
|
||||
*/
|
||||
private Task askAnsweredAsyncTask(String target) {
|
||||
for (Task task : tasks.values()) {
|
||||
if (target.equals(task.target) && task.question == null && task.turnId != null
|
||||
&& !task.future.isDone()) {
|
||||
return task;
|
||||
}
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
/** Record a counter sample when a registry is wired; a no-op in unit tests. */
|
||||
private void count(String name, String... labels) {
|
||||
if (metrics != null) {
|
||||
@@ -435,19 +472,43 @@ public final class MessageService {
|
||||
* <p>Resolving the waiter as a failure — rather than letting it time out — also means the
|
||||
* outcome is counted, so a torn-down delegation stops being invisible to {@code /metrics}.
|
||||
*
|
||||
* @return true if a live waiter was failed
|
||||
* <p><strong>#137 defence in depth.</strong> {@link #reply} already hands a worker's real
|
||||
* {@code fleet_reply} straight to the async ticket it belongs to whenever one is still parked
|
||||
* waiting for it (see {@link #askAnsweredAsyncTask}), so by the time a session is released its
|
||||
* tasks are normally already resolved — this loop's {@code complete} calls are then harmless
|
||||
* no-ops (a {@link CompletableFuture} can only resolve once). But should some other path someday
|
||||
* strand a reply in the inbox without completing its ticket, checking
|
||||
* {@link #hasStrandedReply(String)} here — before ever writing a failure — means a torn-down
|
||||
* session whose worker in fact replied is still reported {@code REPLIED} with that reply's own
|
||||
* text, never the misleading "the worker session was released before it replied" (which also
|
||||
* means the snapshot/worktree recovery hint that follows it never prints once a reply exists).
|
||||
*
|
||||
* @return true if a live waiter or an async task was failed (never true for one recovered as a
|
||||
* reply — see the note above)
|
||||
*/
|
||||
public boolean abandon(String target, String reason) {
|
||||
boolean hadStrandedReply = hasStrandedReply(target);
|
||||
// CB-640: the session is gone — nothing will ever accept or deliver into it now.
|
||||
strandedReplies.remove(target);
|
||||
queuedDeliveries.remove(target);
|
||||
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(target);
|
||||
boolean failed = waiter != null && !waiter.isDone() && rendezvous.resolveFailure(waiter, reason);
|
||||
boolean asyncFailed = false;
|
||||
Reply recovered = null; // lazily drained at most once, only if a task actually needs it
|
||||
for (Task task : tasks.values()) {
|
||||
if (target.equals(task.target) && task.question == null
|
||||
&& task.future.complete(new Reply(Outcome.WORKER_FAILED, reason))) {
|
||||
asyncFailed = true;
|
||||
if (!target.equals(task.target) || task.question != null || task.future.isDone()) {
|
||||
continue;
|
||||
}
|
||||
if (hadStrandedReply && recovered == null) {
|
||||
recovered = recoverStrandedReply(target);
|
||||
}
|
||||
Reply outcome = recovered != null ? recovered : new Reply(Outcome.WORKER_FAILED, reason);
|
||||
if (task.future.complete(outcome)) {
|
||||
if (outcome.outcome() == Outcome.WORKER_FAILED) {
|
||||
asyncFailed = true;
|
||||
} else if (task.turnId != null) {
|
||||
asyncTasksByTurn.remove(task.turnId, task);
|
||||
}
|
||||
}
|
||||
}
|
||||
if (failed) {
|
||||
@@ -456,6 +517,21 @@ public final class MessageService {
|
||||
return failed || asyncFailed;
|
||||
}
|
||||
|
||||
/**
|
||||
* Drain {@code target}'s inbox and hand its content back as a {@link Outcome#REPLIED} result
|
||||
* (#137 defence in depth for {@link #abandon}) — {@code null} if it turned out empty (the
|
||||
* stranding fact raced away, e.g. a lead's own {@code fleet_poll} on the raw session already
|
||||
* drained it first). When more than one message is queued, only the newest is the worker's actual
|
||||
* final answer ({@link #drainReplies} returns them oldest-first).
|
||||
*/
|
||||
private Reply recoverStrandedReply(String target) {
|
||||
var messages = drainReplies(target);
|
||||
if (messages.isEmpty()) {
|
||||
return null;
|
||||
}
|
||||
return new Reply(Outcome.REPLIED, messages.get(messages.size() - 1).content());
|
||||
}
|
||||
|
||||
/**
|
||||
* Acknowledge a specific reply by {@code msgId} for {@code target}. Removes it from the inbox
|
||||
* so that a subsequent drain or peek no longer returns it.
|
||||
|
||||
@@ -710,6 +710,68 @@ class MessageServiceTest {
|
||||
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
|
||||
}
|
||||
|
||||
// --- #137: a fleet_ask round-trip must not orphan the ticket's own reply -------------------
|
||||
//
|
||||
// The primary's fleet_send{turnId} answer call is itself bounded (a real MCP call, capped well
|
||||
// under a minute) — far shorter than a resumed turn can genuinely take to finish real work. These
|
||||
// drive the exact real delegation path (async send -> worker asks -> primary answers -> primary's
|
||||
// own wait gives up -> worker's real fleet_reply arrives afterwards) rather than calling a reply
|
||||
// sink directly, since the bug is specifically about which sink the resumed turn's reply reaches.
|
||||
|
||||
@Test
|
||||
void aReplyAfterAnswerTimesOutStillCompletesTheAsyncTicket() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "task that asks");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
CompletableFuture<MessageService.AskResult> ask =
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
|
||||
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
|
||||
|
||||
// The primary answers, but its own bounded wait for the worker's resumed turn is short and
|
||||
// expires before the worker (still genuinely working) gets back to it.
|
||||
MessageService.Reply answerReply = messages.answer(asking.turnId(), "config.yaml", 150);
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome(),
|
||||
"the primary's own bounded wait gives up before the worker finishes resuming");
|
||||
|
||||
// The worker keeps working past that window and only now calls fleet_reply.
|
||||
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");
|
||||
}
|
||||
|
||||
@Test
|
||||
void fleetStopAfterAnOrphanedReplyDoesNotFailTheTicket() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "task that asks");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
CompletableFuture<MessageService.AskResult> ask =
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
|
||||
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
|
||||
|
||||
MessageService.Reply answerReply = messages.answer(asking.turnId(), "config.yaml", 150);
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome());
|
||||
|
||||
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
|
||||
|
||||
// fleet_stop tears the worker's session down right after the reply landed — this must never
|
||||
// report the misleading "the worker session was released before it replied": a reply is
|
||||
// exactly what happened.
|
||||
assertFalse(messages.abandon(T, "the worker session was released before it replied"),
|
||||
"a reply already arrived, so nothing here is a genuine failure");
|
||||
|
||||
MessageService.TaskView view = awaitTicketPhase(ticket, MessageService.Phase.DONE);
|
||||
assertEquals("PR opened: https://example/pulls/42", view.reply());
|
||||
}
|
||||
|
||||
@Test
|
||||
void unansweredAsyncQuestionReturnsTheTicketToPendingAndReleasesItsTarget() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "task that asks");
|
||||
|
||||
Reference in New Issue
Block a user