From 388aba76324d5f32ebe539a760a51d63ed9a7af0 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Mon, 31 Aug 2026 14:24:54 +0700 Subject: [PATCH] #137: an answered turn's reply completes its own ticket, instead of a false failure A wait:false delegation whose worker used fleet_ask ended with fleet_poll{ticket} reporting 'the worker session was released before it replied' -- naming a worktree, a branch and a snapshot commit, so it read as lost work. The worker had in fact replied in full. The ticket guessed the ask rendezvous detached the turn. The real cause is narrower: answer() (behind fleet_send{turnId}) waits only for the lead's own bounded MCP call window. A resumed turn doing real work -- edits, a build, a push, a PR -- routinely outlives it. On timeout answer()'s finally closed the waiter, so the worker's later fleet_reply found none and fell to the session inbox, leaving the ticket's future unresolved until fleet_stop forced it FAILED. reply() now looks for the async task parked on this exact answered turn and completes it with the real reply. That is safe against completing the wrong ticket: answer() calls clearAsyncQuestion(turnId, false), so the task keeps its turnId and stays in asyncTasksByTurn, and hasAsyncQuestion therefore still makes send() return BUSY for a second async send to that target. At most one candidate task can exist per target. abandon() keeps an independent check: if a reply was stranded, a released session reports REPLIED with that text rather than a failure -- so the recovery hint that implies lost work never prints once a reply exists. Verified before merging: build green unpiped, and both new tests drive the full delegation path (send async, ask, answer, reply, poll the ticket) rather than handing a Reply to a sink, which is the trap this ticket called out. Co-authored-by: fleetd worker --- .../dev/ltms/fleet/msg/MessageService.java | 92 +++++++++++++++++-- .../ltms/fleet/msg/MessageServiceTest.java | 62 +++++++++++++ 2 files changed, 146 insertions(+), 8 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 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");