From 4887d03d88d717ee8280bed9c6a7515779b47c13 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Fri, 4 Sep 2026 10:13:56 +0700 Subject: [PATCH] #275: abandon() sweeps an ASKING ticket only on a definite teardown MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Confirmed reachable: a target torn down for good (fleet_stop / the idle reaper) while its async ticket sits in fleet_ask (Phase.ASKING) got permanently stuck. resolveQuestion already closes the forward waiter, the question == null guard excluded the task from abandon()'s sweep, and by the time the worker's own fleet_ask lapses (~55-115s) the released session no longer appears in FleetHealthMonitor's roster, so nothing ever calls abandon() again. fleet_poll{ticket} then reports PENDING forever. Add abandon(target, reason, sweepAsking) — sessions.onRelease (a definite teardown: the pane is being stopped right now) passes true and now fails the ASKING ticket and closes its reverse-rendezvous ask. FleetHealthMonitor's health-classification call keeps the 2-arg overload (sweepAsking=false): a GONE/NEVER_READY reading is a guess from the live agent list, not a teardown it performed, and abandonDoesNotFailAnAsyncTicketWaitingForAnAnswer already covers why an active ask must survive that guess (the primary may be mid-answer for the same turn). hasOrphanedDelegation is left unchanged for the same reason — it must not flag a live, active ask as orphaned. Proven with a test driving the real public sequence (sendAsync -> ask -> abandon(..., true)), not a hand-built task map; reverted the widening to confirm it goes red, then restored it. --- .../src/main/java/dev/ltms/fleet/Fleetd.java | 6 +- .../dev/ltms/fleet/msg/MessageService.java | 61 ++++++++++++++++++- .../ltms/fleet/msg/MessageServiceTest.java | 50 +++++++++++++++ 3 files changed, 113 insertions(+), 4 deletions(-) diff --git a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java index 6106af5..9730fb4 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java +++ b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java @@ -593,7 +593,11 @@ public final class Fleetd { if (detail.agentSessionId() != null) { reason += " agentSessionId=" + detail.agentSessionId(); } - messages.abandon(detail.terminalId(), reason); + // fleetd #275: this is an explicit teardown (fleet_stop, or the idle reaper) — the + // worker's pane is being stopped right now, so an open fleet_ask has no turn left to + // resume into. Sweep it too, unlike FleetHealthMonitor's health-classification call + // (see MessageService.abandon's javadoc for why those two must differ). + messages.abandon(detail.terminalId(), reason, true); replyInbox.release(detail.terminalId()); primaryRegistry.forgetDelegation(detail.terminalId()); // CB-532: don't leak the lead binding }); 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 7138991..5d92c2d 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java @@ -360,6 +360,18 @@ public final class MessageService { * {@link #ask} clears the ticket's question and returns it to {@code PENDING}, but {@link #send} * already closed the forward waiter the instant the question surfaced, so the target has * neither an accepted nor a queued delivery left to show for it. + * + *

Deliberately still {@code question == null} only (fleetd #275). This + * method must not also report a still-{@link Phase#ASKING} task as orphaned: the worker may + * genuinely be waiting on a live primary that is about to (or already mid-{@link #answer}) + * answer it, and {@link dev.ltms.fleet.health.FleetHealthMonitor} would classify that as + * {@code DELEGATION_ORPHANED} on nothing more than an active, healthy conversation. {@link + * #abandon(String, String, boolean)}'s {@code sweepAsking} path fixes the actual reachable gap + * (a target torn down for good while genuinely {@code ASKING}) at the point of teardown itself, + * by completing the task's future right there — so by the time this method would ever see it, + * {@code task.future.isDone()} is already {@code true} and it is excluded regardless of this + * guard. Widening this check instead of that one would trade a real fix for false positives on + * every ordinary in-flight question. */ public boolean hasOrphanedDelegation(String target) { if (target == null || hasAcceptedDelivery(target) || hasQueuedDelivery(target)) { @@ -580,6 +592,39 @@ public final class MessageService { * reply — see the note above) */ public boolean abandon(String target, String reason) { + return abandon(target, reason, false); + } + + /** + * As {@link #abandon(String, String)}, with control over whether a task still paused in + * {@code fleet_ask} ({@link Phase#ASKING}) is swept too (fleetd #275). + * + *

{@code sweepAsking} must be {@code true} only when the caller has independent, certain + * knowledge that {@code target} can never resume its turn — today that is only + * {@code sessions.onRelease}'s teardown (an explicit {@code fleet_stop}, or the idle reaper): + * the worker's pane is being stopped right now, so whatever it was mid-{@code fleet_ask} about + * has no turn left to resume into. {@link dev.ltms.fleet.health.FleetHealthMonitor}'s + * health-classification call keeps passing {@code false} (via {@link #abandon(String, String)}): + * a GONE/NEVER_READY reading is the daemon's best guess from the live agent list, not a teardown + * it performed itself, and {@code abandonDoesNotFailAnAsyncTicketWaitingForAnAnswer} documents + * why an active ask must survive that guess — the primary may already be mid-{@link #answer} for + * the very same turn, and completing it here first would preempt a real answer with a misleading + * failure. + * + *

Without {@code sweepAsking} on the release path, a target torn down while + * genuinely {@code ASKING} was unrecoverable. {@link #resolveQuestion} had already + * closed the forward waiter the instant the question surfaced (so the {@code waiter} branch + * below finds nothing to fail), the {@code question == null} guard excluded the task from + * {@code matching} (so the loop below skipped it too), and the worker's own {@code fleet_ask} + * clears {@link Task#question} back to {@code null} only once it lapses (the reverse-rendezvous + * window — up to {@code FleetMcp.ASK_DEFAULT_TIMEOUT_MS} / {@code FleetApp.MAX_ASK_TIMEOUT_MS}, + * 55–115s) — by which point the released session no longer appears in {@code sessions.roster()} + * for {@link dev.ltms.fleet.health.FleetHealthMonitor} to ever re-observe, so nothing was ever + * left to call {@link #abandon} on this target again. The ticket then sat in {@link #tasks} + * forever: not terminal, so {@link #pruneTerminalTickets} never dropped it, and + * {@code fleet_poll} reported it stuck at {@link Phase#PENDING} for good. + */ + public boolean abandon(String target, String reason, boolean sweepAsking) { boolean hadStrandedReply = hasStrandedReply(target); // CB-640: the session is gone — nothing will ever accept or deliver into it now. strandedReplies.remove(target); @@ -590,7 +635,8 @@ public final class MessageService { List matching = new ArrayList<>(); for (Task task : tasks.values()) { - if (target.equals(task.target) && task.question == null && !task.future.isDone()) { + if (target.equals(task.target) && (sweepAsking || task.question == null) + && !task.future.isDone()) { matching.add(task); } } @@ -607,11 +653,20 @@ public final class MessageService { for (Task task : matching) { boolean isRecovery = task == recoveryTask && recovered != null; Reply outcome = isRecovery ? recovered : new Reply(Outcome.WORKER_FAILED, reason); + String turnId = task.turnId; if (task.future.complete(outcome)) { if (outcome.outcome() == Outcome.WORKER_FAILED) { asyncFailed = true; - } else if (task.turnId != null) { - asyncTasksByTurn.remove(task.turnId, task); + } + if (turnId != null) { + // #275: whether this task was swept out of ASKING or was already answered and + // only waiting on its resumed turn's real reply (#137), nothing will ever + // complete this turnId now — drop it from this class's own bookkeeping AND the + // reverse-rendezvous itself, so hasAsyncQuestion(target) stops reporting a turn + // that is actually done, and a late answer() sees it as lapsed rather than + // resolving a question nothing is listening for any more. + asyncTasksByTurn.remove(turnId, task); + rendezvous.closeAsk(turnId); } } else if (isRecovery) { // The recovered reply was already drained out of the inbox, but this task resolved 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 7e7af3c..7ce8c22 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java @@ -828,6 +828,56 @@ class MessageServiceTest { assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome()); } + // --- fleetd #275: a target torn down FOR GOOD while genuinely ASKING must not orphan -------- + // + // sessions.onRelease (fleet_stop, or the idle reaper) is the one abandon() caller that knows + // for certain the target can never resume: its pane is being stopped right now. Unlike the + // health-classification caller above (a GONE/NEVER_READY guess, not a teardown it performed), + // it must sweep an ASKING ticket right here — see MessageService.abandon(String, String, + // boolean)'s javadoc for the full reachability chain this closes: without this, the forward + // waiter is already closed by the time the question surfaces, the ASKING guard skips the task, + // and by the time the worker's own fleet_ask lapses (~55-115s later) the released session no + // longer appears in FleetHealthMonitor's roster for anything to ever sweep it again — leaving + // fleet_poll{ticket} stuck PENDING forever. + + @Test + void abandonWithSweepAskingFailsATornDownTargetsAskingTicket() throws Exception { + String ticket = messages.sendAsync(T, "task that asks"); + awaitWaiting(); + injectDelivery(); + + CompletableFuture ask = + CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 300)); + MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING); + + assertTrue(messages.abandon(T, "the worker session was released before it replied", true), + "a released target's open ask can never resume, so it must fail right here"); + + MessageService.TaskView failed = awaitTicketPhase(ticket, MessageService.Phase.FAILED); + assertEquals("the worker session was released before it replied", failed.detail()); + + // The reverse-rendezvous ask is torn down too: the worker's still-blocked fleet_ask rides + // out its own timeout (nothing completed its answer future), and a late answer() for the + // same turnId must see it as lapsed rather than resolving a question nobody is waiting on. + assertEquals(MessageService.AskOutcome.TIMED_OUT, ask.get(5, TimeUnit.SECONDS).outcome()); + assertEquals(MessageService.Outcome.STALE_TURN, + messages.answer(asking.turnId(), "config.yaml", 200).outcome()); + } + + @Test + void abandonWithoutSweepAskingBehavesLikeTheTwoArgOverload() throws Exception { + String ticket = messages.sendAsync(T, "task that asks"); + awaitWaiting(); + injectDelivery(); + + CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000)); + awaitTicketPhase(ticket, MessageService.Phase.ASKING); + + assertFalse(messages.abandon(T, "agent target term_a not found", false), + "sweepAsking=false must match the plain abandon(target, reason) overload"); + assertEquals(MessageService.Phase.ASKING, messages.poll(ticket).phase()); + } + // --- #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