From ee5f8b932b4e6cb7b20beccb59d33ede70e92401 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Mon, 31 Aug 2026 14:12:23 +0700 Subject: [PATCH 1/3] #137: complete an async ticket's own reply after answer() times out MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- REPORT-cb137.md | 166 ++++++++++++++++++ .../dev/ltms/fleet/msg/MessageService.java | 92 +++++++++- .../ltms/fleet/msg/MessageServiceTest.java | 62 +++++++ 3 files changed, 312 insertions(+), 8 deletions(-) create mode 100644 REPORT-cb137.md diff --git a/REPORT-cb137.md b/REPORT-cb137.md new file mode 100644 index 0000000..4a94222 --- /dev/null +++ b/REPORT-cb137.md @@ -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: but was: +dev.ltms.fleet.msg.MessageServiceTest.fleetStopAfterAnOrphanedReplyDoesNotFailTheTicket + org.opentest4j.AssertionFailedError: a reply already arrived, so nothing here is a genuine failure + ==> expected: but was: +``` + +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) diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java index d1e8d77..7aab4b1 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java @@ -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 not 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 not a + * failure — the reply is held for later drain. * *

Do NOT use this for mid-turn questions. {@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 { *

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 + *

#137 defence in depth. {@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 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. diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java index eca9721..bf27236 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java @@ -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 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 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"); From 97f6c33a45555c4313e1aa876c428e2667eefeeb Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Mon, 31 Aug 2026 22:46:23 +0700 Subject: [PATCH 2/3] #137: don't guess when a target has more than one open async task MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- .../dev/ltms/fleet/msg/MessageService.java | 108 ++++++++++++++---- 1 file changed, 85 insertions(+), 23 deletions(-) diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java index 7aab4b1..ba069ef 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java @@ -9,6 +9,7 @@ import dev.ltms.fleet.metrics.Metrics; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.util.ArrayList; import java.util.List; import java.util.UUID; import java.util.concurrent.CompletableFuture; @@ -167,6 +168,12 @@ public final class MessageService { private static final class Task { private final String ticket; private final String target; + /** + * When this task was created (#137 fix): the tiebreaker for which of several open tasks on + * one target gets a recovered reply in {@link #abandon} — the oldest, since it is the one + * that has been waiting longest. + */ + private final long createdNanos; private final CompletableFuture future = new CompletableFuture<>(); /** * When {@link #future} resolved, or {@code null} while it is still pending — the clock @@ -184,6 +191,7 @@ public final class MessageService { private Task(String ticket, String target, LongSupplier nowNanos) { this.ticket = ticket; this.target = target; + this.createdNanos = nowNanos.getAsLong(); future.whenComplete((reply, ex) -> completedNanos = nowNanos.getAsLong()); } } @@ -375,6 +383,16 @@ public final class MessageService { * {@link Rendezvous#resolveQuestion} must keep today's {@code NO_WAITER} behaviour — questions * are interactive and must never be queued. * + *

Ambiguous match also falls to the inbox. A target may have more than one + * open async task waiting on an answered turn — a lead can, and does, issue a second + * {@code fleet_send{wait:false}} at a member that is still busy. Returning whichever candidate a + * {@code ConcurrentHashMap} iteration reaches first would let a genuine reply complete the + * wrong ticket — silently handing the lead something that reads like a correct answer to + * a delegation the worker never touched, which is worse than a failure because the lead acts on + * it. When more than one candidate exists, guessing is not safe: fall back to the inbox exactly + * as the zero-candidate case does, and let {@link #abandon} apply the eventual recovery + * deterministically instead. + * * @return always {@code true} — the reply resolved a live send, completed a parked ticket, or * was queued */ @@ -392,13 +410,21 @@ public final class MessageService { // 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); + List candidates = askAnsweredAsyncTasks(session); + if (candidates.size() == 1) { + Task orphan = candidates.get(0); + if (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 } - count(FleetMetrics.REPLIES, "path", "async-recovered"); - return true; // the ticket itself took it — no inbox stranding at all + } else if (candidates.size() > 1) { + List tickets = candidates.stream().map(t -> t.ticket).toList(); + log.warn("reply from {} matches {} open async tickets {} — cannot tell which one it " + + "answers, queuing to the inbox instead of guessing", session, candidates.size(), + tickets); } inbox.publish(session, UUID.randomUUID().toString(), content); // CB-640: record the stranding itself (not just the reply text) so fleet health can see a @@ -414,21 +440,27 @@ public final class MessageService { } /** - * 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 + * Every 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). Empty 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}). + * + *

Usually holds at most one entry, but not always: a target has no built-in limit of one open + * async task (see the {@link #reply} javadoc), so more than one task can independently reach this + * exact state at once. The caller decides what "more than one" means — {@link #reply} treats it + * as unresolvable and falls back to the inbox; {@link #abandon} picks the oldest deterministically. */ - private Task askAnsweredAsyncTask(String target) { + private List askAnsweredAsyncTasks(String target) { + List candidates = new ArrayList<>(); for (Task task : tasks.values()) { if (target.equals(task.target) && task.question == null && task.turnId != null && !task.future.isDone()) { - return task; + candidates.add(task); } } - return null; + return candidates; } /** Record a counter sample when a registry is wired; a no-op in unit tests. */ @@ -473,16 +505,24 @@ public final class MessageService { * outcome is counted, so a torn-down delegation stops being invisible to {@code /metrics}. * *

#137 defence in depth. {@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 + * {@code fleet_reply} straight to the async ticket it belongs to whenever exactly one is still + * parked waiting for it (see {@link #askAnsweredAsyncTasks}), 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). * + *

At most one task gets the recovered reply. A stranded reply is one worker + * answer, so it can settle at most one open task on this target — never every open task, and + * never a guess. When more than one task is still open here (the same "more than one open async + * task per target" situation {@link #reply} defers on), the recovered reply goes to the + * oldest (lowest {@link Task#createdNanos}) — it has been waiting longest, so it is the + * one most likely to be what the reply actually answers. Every other open task keeps the ordinary + * {@code WORKER_FAILED} path it would take without a stranded reply at all. + * * @return true if a live waiter or an async task was failed (never true for one recovered as a * reply — see the note above) */ @@ -494,21 +534,39 @@ public final class MessageService { CompletableFuture 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 + + List matching = new ArrayList<>(); for (Task task : tasks.values()) { - if (!target.equals(task.target) || task.question != null || task.future.isDone()) { - continue; + if (target.equals(task.target) && task.question == null && !task.future.isDone()) { + matching.add(task); } - if (hadStrandedReply && recovered == null) { - recovered = recoverStrandedReply(target); + } + Task recoveryTask = null; + if (hadStrandedReply && !matching.isEmpty()) { + recoveryTask = matching.get(0); + for (Task candidate : matching) { + if (candidate.createdNanos < recoveryTask.createdNanos) { + recoveryTask = candidate; + } } - Reply outcome = recovered != null ? recovered : new Reply(Outcome.WORKER_FAILED, reason); + } + Reply recovered = recoveryTask != null ? recoverStrandedReply(target) : null; + for (Task task : matching) { + boolean isRecovery = task == recoveryTask && recovered != null; + Reply outcome = isRecovery ? 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); } + } else if (isRecovery) { + // The recovered reply was already drained out of the inbox, but this task resolved + // through another path (e.g. a concurrent reply() or a second abandon() racing this + // one) between us choosing it and completing it here. Put the reply back rather than + // lose it silently — it may still belong to some other still-open task, or the next + // caller that drains this target's inbox. + inbox.publish(target, UUID.randomUUID().toString(), recovered.text()); } } if (failed) { @@ -523,6 +581,10 @@ public final class MessageService { * 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). + * + *

This does drain (removes the messages from the inbox) before the caller knows whether the + * task it is recovering for will actually accept them — {@link #abandon} is the one that puts a + * reply back if its {@code complete} call turns out to lose the race. */ private Reply recoverStrandedReply(String target) { var messages = drainReplies(target); From de026b8f8abf4d8e3894745bfcb2639d82b2dd08 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Mon, 31 Aug 2026 23:17:50 +0700 Subject: [PATCH 3/3] #137 follow-up: document defect-1 and defect-2 reachability, drop unconstructible test MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- .../dev/ltms/fleet/msg/MessageService.java | 73 +++++++++++++++---- .../ltms/fleet/msg/MessageServiceTest.java | 26 +++++++ 2 files changed, 85 insertions(+), 14 deletions(-) diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java index ba069ef..522a8e7 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java @@ -383,9 +383,12 @@ public final class MessageService { * {@link Rendezvous#resolveQuestion} must keep today's {@code NO_WAITER} behaviour — questions * are interactive and must never be queued. * - *

Ambiguous match also falls to the inbox. A target may have more than one - * open async task waiting on an answered turn — a lead can, and does, issue a second - * {@code fleet_send{wait:false}} at a member that is still busy. Returning whichever candidate a + *

Ambiguous match also falls to the inbox. {@link #askAnsweredAsyncTasks} + * cannot actually return more than one entry today (see its own javadoc for why — in short, + * {@link #hasAsyncQuestion} keeps a target BUSY, so no second task can reach this state, for as + * long as an earlier one's {@code turnId} is still stamped). That is an emergent guarantee from + * two other facts, not one this method enforces, so this branch stays in as defence in depth + * rather than being removed as dead code: if it ever weakens, returning whichever candidate a * {@code ConcurrentHashMap} iteration reaches first would let a genuine reply complete the * wrong ticket — silently handing the lead something that reads like a correct answer to * a delegation the worker never touched, which is worse than a failure because the lead acts on @@ -447,10 +450,21 @@ public final class MessageService { * 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}). * - *

Usually holds at most one entry, but not always: a target has no built-in limit of one open - * async task (see the {@link #reply} javadoc), so more than one task can independently reach this - * exact state at once. The caller decides what "more than one" means — {@link #reply} treats it - * as unresolvable and falls back to the inbox; {@link #abandon} picks the oldest deterministically. + *

Returns at most one entry today — verified, not assumed. {@link #send} + * refuses to open a waiter on {@code target} while {@link #hasAsyncQuestion} is true, and that + * check matches ANY task whose {@code turnId} is still stamped in {@code asyncTasksByTurn} — + * not only while its question is still open. {@link #answer} deliberately leaves that stamp in + * place ({@code clearAsyncQuestion(turnId, false)}) until the resumed turn's own future actually + * resolves, at which point {@link #finishAsyncTask} both removes the stamp AND completes that + * task's future in the same call. So a second task can never reach "{@code turnId} stamped, future + * still open" — the exact pair this method matches on — while a first one already holds it: by + * the time the stamp is gone, so is the eligibility. This is an emergent property of those two + * facts holding together, not something this method (or its callers) enforces on its own — flip + * {@code forgetTurn} to {@code true} in that one {@link #answer} call and it silently stops being + * true, with nothing left to fail loudly. The callers below still handle "more than one" as + * defence in depth against exactly that, not because they exercise it today: {@link #reply} + * treats it as unresolvable and falls back to the inbox; {@link #abandon} would pick the oldest + * deterministically (its own {@code matching} list has no such guarantee — see its javadoc). */ private List askAnsweredAsyncTasks(String target) { List candidates = new ArrayList<>(); @@ -515,13 +529,44 @@ public final class MessageService { * 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). * - *

At most one task gets the recovered reply. A stranded reply is one worker - * answer, so it can settle at most one open task on this target — never every open task, and - * never a guess. When more than one task is still open here (the same "more than one open async - * task per target" situation {@link #reply} defers on), the recovered reply goes to the - * oldest (lowest {@link Task#createdNanos}) — it has been waiting longest, so it is the - * one most likely to be what the reply actually answers. Every other open task keeps the ordinary - * {@code WORKER_FAILED} path it would take without a stranded reply at all. + *

At most one task gets the recovered reply — and here, unlike {@link #reply}'s + * {@link #askAnsweredAsyncTasks}, {@code matching.size() >= 2} alone is reachable today. + * This method's {@code matching} filter has no {@code turnId != null} requirement, so it matches + * any plain (never-asked) open task too — and {@link #sendAsync} does not limit a target to one + * of those: a second {@code fleet_send{wait:false}} at a target that is still busy returns its own + * ticket immediately and simply parks its {@link #send} behind the target's session lock for up + * to {@link #ASYNC_TIMEOUT_MS}, exactly as {@code abandonFailsEveryPendingAsyncTicketForTheReleasedTarget} + * already proves. Before this fix, the loop below drained the strand once and then reused that + * same {@code Reply} for every task it walked past — so two open tasks really did both + * complete {@code REPLIED} with the same text (see the pre-fix loop in commit 97f6c33's parent). + * A stranded reply is one worker answer, so it can settle at most one open task on this target — + * never every open task, and never a guess. When more than one task is still open here, the + * recovered reply goes to the oldest (lowest {@link Task#createdNanos}) — it has been + * waiting longest, so it is the one most likely to be what the reply actually answers. Every + * other open task keeps the ordinary {@code WORKER_FAILED} path it would take without a stranded + * reply at all. + * + *

{@code matching.size() >= 2} together with {@code hadStrandedReply} is a different + * question, and today it is defence in depth rather than a path this codebase's public API can + * drive. This class has exactly two sites that ever acquire a target's entry in + * {@code sessionLocks} — {@link #send} and {@link #answer} — and both open a {@link Rendezvous} + * waiter for that same target as the very first thing they do after acquiring the lock, then hold + * lock and waiter together for the rest of their critical section ({@link #send} also clears + * {@link #strandedReplies} right there, the instant it opens its waiter — before it ever enqueues + * delivery). So "the session lock is held" and "a live waiter is open for it" are the same fact + * throughout this class, and {@link #reply}'s fast path always resolves a currently-open waiter + * directly rather than stranding. The two facts this method wants therefore cannot be produced + * side by side: while the lock is held, a real reply resolves the open waiter directly and never + * reaches {@link #strandedReplies}; the instant the lock is free, any parked matching task's own + * {@link #send} that is scheduled next wins it and, by opening its waiter, clears the strand again + * before this method ever runs. There is no way to hold that lock open-but-unaccepted from outside + * {@link #send}/{@link #answer} to freeze a window in between. Constructing both facts at once + * through {@code sendAsync}/{@code reply}/{@code ask}/{@code answer} would need a race against + * virtual-thread scheduling, not a deterministic sequence — so the oldest-wins code below stays as + * defence in depth against a regression to that mechanism (e.g. clearing {@link #strandedReplies} + * on a narrower condition than "any acceptance"), not because today's test suite exercises the + * conjunction. {@code matching.size() >= 2} alone, without a strand, is exactly what + * {@code abandonFailsEveryPendingAsyncTicketForTheReleasedTarget} already covers. * * @return true if a live waiter or an async task was failed (never true for one recovered as a * reply — see the note above) diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java index bf27236..95f1e92 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java @@ -688,6 +688,32 @@ class MessageServiceTest { assertFailedTicket(third, "agent target term_a not found"); } + // --- #137 follow-up: abandon() must not guess when more than one task is open --------------- + // + // A test combining a genuine stranded reply (hasStrandedReply(T)==true) with two simultaneously + // open matching tasks was attempted here and removed after investigation showed the combination + // is not reachable through the public API today, not merely hard to time right: + // + // This class has exactly two call sites that ever hold a target's entry in the session-lock map + // (send() and answer()), and both open a Rendezvous waiter for that same target as the first thing + // they do after acquiring the lock, holding lock and waiter together for their whole critical + // section. So "the lock is held" and "a live waiter is open" are the same fact throughout this + // class. reply()'s fast path always resolves a currently-open waiter directly instead of + // stranding — so a strand can only be created while NO task is accepted (lock free), and the + // instant the lock is next taken (by any parked matching task's own send(), the moment it is + // scheduled), that acceptance clears strandedReplies again (see send()'s CB-640 comment) before + // abandon() can ever observe both facts together. Confirmed empirically too: an earlier version of + // this test stranded a reply, then created an "accepted" task (awaitWaiting()) followed by a + // "parked" one — and the accepted task's own acceptance silently cleared the strand it was + // supposed to be racing against, so the parked task came back WORKER_FAILED instead of DONE, not + // because the fix was missing but because the test's premise could not be constructed. + // + // The reachable half — matching.size() >= 2 alone, no strand — is exactly what + // abandonFailsEveryPendingAsyncTicketForTheReleasedTarget already covers (all fail, none guess). + // The oldest-wins code in abandon() stays as defence in depth (see its own javadoc) against a + // regression that would make the conjunction reachable, e.g. clearing strandedReplies on a + // narrower condition than "any acceptance" — not because this suite exercises it today. + @Test void abandonDoesNotFailAnAsyncTicketWaitingForAnAnswer() throws Exception { String ticket = messages.sendAsync(T, "task that asks");