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