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