#275: abandon() sweeps an ASKING ticket only on a definite teardown
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.
This commit is contained in:
@@ -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
|
||||
});
|
||||
|
||||
@@ -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.
|
||||
*
|
||||
* <p><strong>Deliberately still {@code question == null} only (fleetd #275).</strong> 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).
|
||||
*
|
||||
* <p>{@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.
|
||||
*
|
||||
* <p><strong>Without {@code sweepAsking} on the release path, a target torn down while
|
||||
* genuinely {@code ASKING} was unrecoverable.</strong> {@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<Task> 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
|
||||
|
||||
@@ -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<MessageService.AskResult> 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
|
||||
|
||||
Reference in New Issue
Block a user