From fa39a5f55e40a0a6275350a6308779fed849f177 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Fri, 4 Sep 2026 11:28:12 +0700 Subject: [PATCH] #280: sweep a GONE/NEVER_READY target's lapsed fleet_ask once, not just on transition MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit FleetHealthMonitor's fire-once-per-transition rule (CB-580) means a target that is genuinely mid-fleet_ask when health first classifies it GONE is correctly skipped (sweepAsking=false). But nothing re-fires abandon() once that ask lapses on its own 55-115s later: FleetHealth.decide keeps reporting GONE every tick, and reportTransition's previous==next guard never lets the sweep run again. The ticket then sat PENDING forever, the same destination #275 fixed for an explicit teardown, reached here by a health guess instead. SessionManager.reapIdle only reaps READY/DONE sessions (SessionManager.java:844), and a session mid-turn (including mid-ask) stays BUSY the whole time (onDelivered sets BUSY, nothing clears it until the turn completes) — so SessionReaper never releases such a session and onRelease's sweepAsking=true path is never reached. Fix: schedule one bounded, delayed re-check per terminal transition (not a per-tick retry — that shape was rejected by CB-580). It fires failTerminalTarget again after a delay that exceeds the worst-case ask-lapse window, and only if the target is still classified in the same terminal state at that time, so a recovered or since-released target is never reached into. sweepAsking stays false throughout, so a ticket whose ask has not yet lapsed is still never touched — same invariant abandonDoesNotFailAnAsyncTicketWaitingForAnAnswer pins. Proven with a mutation: neutering recheckTerminalTarget's body made delayedRecheckSweepsATicketWhoseAskLapsedAfterGoneWasFirstObserved fail with "expected: but was: ", all 20 other FleetHealthMonitorTest cases still green; restored and reran clean (1300 tests, 0 failures). --- .../ltms/fleet/health/FleetHealthMonitor.java | 51 ++++++++ .../fleet/health/FleetHealthMonitorTest.java | 118 ++++++++++++++++++ 2 files changed, 169 insertions(+) diff --git a/fleetd/src/main/java/dev/ltms/fleet/health/FleetHealthMonitor.java b/fleetd/src/main/java/dev/ltms/fleet/health/FleetHealthMonitor.java index ee61559..6ec2e49 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/health/FleetHealthMonitor.java +++ b/fleetd/src/main/java/dev/ltms/fleet/health/FleetHealthMonitor.java @@ -28,6 +28,14 @@ public final class FleetHealthMonitor { static final int MAX_FAIL_TARGET_ATTEMPTS = 3; // CB-641: Match the injector's 60s readiness gate so health allows a full first boot. static final long READINESS_GRACE_NANOS = TimeUnit.SECONDS.toNanos(60); + /** + * fleetd #280: how long after a terminal transition to wait before the one bounded re-check + * fires. Must exceed the worst-case reverse-rendezvous {@code fleet_ask} window (55-115s, see + * {@code FleetMcp.ASK_DEFAULT_TIMEOUT_MS} / {@code FleetApp.MAX_ASK_TIMEOUT_MS}) so that, if the + * target was genuinely {@code ASKING} when {@code state} was first observed, its own ask has had + * time to lapse (clearing {@code Task#question} back to {@code null}) before this fires. + */ + static final long ASK_LAPSE_RECHECK_DELAY_SECONDS = 120; private final AgentControl agents; private final Supplier> roster; @@ -171,6 +179,7 @@ public final class FleetHealthMonitor { // member stayed terminal. if (terminal(next)) { failTerminalTarget(target, next); + scheduleTerminalRecheck(target, next); } } @@ -191,6 +200,48 @@ public final class FleetHealthMonitor { target, state, MAX_FAIL_TARGET_ATTEMPTS, last); } + /** + * fleetd #280: schedule the one bounded, delayed follow-up for a terminal transition — never a + * per-tick retry (CB-580 rejected that shape; {@link #reportTransition} still fires + * {@link #failTerminalTarget} exactly once per transition, unconditionally on the tick loop). + * This is a single one-shot task, scheduled once per transition into GONE/NEVER_READY, so a + * member stuck terminal for the rest of its life gets exactly one extra attempt, not one per + * tick. See {@link #recheckTerminalTarget} for why the extra attempt is safe. + */ + private void scheduleTerminalRecheck(String target, HealthState state) { + if (scheduler.isShutdown()) return; + try { + scheduler.schedule(() -> recheckTerminalTarget(target, state), + ASK_LAPSE_RECHECK_DELAY_SECONDS, TimeUnit.SECONDS); + } catch (RuntimeException e) { + log.warn("fleet health: could not schedule terminal re-check for member={} state={}", + target, state, e); + } + } + + /** + * fleetd #280: the delayed re-check {@link #scheduleTerminalRecheck} scheduled for one terminal + * transition. By now, a {@code fleet_ask} that was still open when {@code state} was first + * observed has had time to lapse on its own (see {@link #ASK_LAPSE_RECHECK_DELAY_SECONDS}), + * clearing {@code Task#question} back to {@code null} — which is exactly what + * {@link MessageService#abandon(String, String, boolean)}'s {@code sweepAsking=false} filter + * needs to finally match it. Calling {@link #failTerminalTarget} again is safe only because + * {@code sweepAsking} stays {@code false}: a task genuinely still {@code ASKING} is skipped + * exactly as it was on the very first attempt — this never fails a ticket whose ask has not yet + * lapsed. + * + *

Guarded on "target is still classified {@code state}." Without this guard, + * a member that recovered (or was released and dropped from the roster) between the transition + * and this re-check would still take a blind {@code failTarget} call — reaching into whatever + * brand-new, unrelated turn it has since picked up and failing it too. {@link #states} already + * carries the live classification (updated every tick, pruned to the current roster on release), + * so a stale or recovered target simply reads as a mismatch here and this is a no-op. + */ + void recheckTerminalTarget(String target, HealthState state) { + if (states.get(target) != state) return; + failTerminalTarget(target, state); + } + private static boolean terminal(HealthState state) { return state == HealthState.GONE || state == HealthState.NEVER_READY; } diff --git a/fleetd/src/test/java/dev/ltms/fleet/health/FleetHealthMonitorTest.java b/fleetd/src/test/java/dev/ltms/fleet/health/FleetHealthMonitorTest.java index 08360a7..6e83b8d 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/health/FleetHealthMonitorTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/health/FleetHealthMonitorTest.java @@ -24,7 +24,9 @@ import org.slf4j.LoggerFactory; import java.util.List; import java.util.Map; import java.util.Set; +import java.util.concurrent.CompletableFuture; import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledFuture; import java.util.concurrent.ScheduledThreadPoolExecutor; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicLong; @@ -434,4 +436,120 @@ class FleetHealthMonitorTest { logger.detachAppender(appender); } } + + // --- fleetd #280: a GONE/NEVER_READY guess whose fleet_ask lapses AFTER the first sweep must + // still be swept, without ever reaching into a target that has since recovered or left the + // roster. See FleetHealthMonitor.recheckTerminalTarget's javadoc for the full reachability chain. + + @Test void terminalTransitionSchedulesExactlyOneDelayedRecheck() { + ScheduledThreadPoolExecutor scheduler = new ScheduledThreadPoolExecutor(1); + FleetHealthMonitor monitor = monitor(new FakeHerdr(), + List.of(member("term_a", MemberSession.State.BUSY, 0, 0)), scheduler, () -> 1, 600, + (_, _) -> { }); + + monitor.reportTransition("term_a", HealthState.GONE); + + // Driven via reportTransition directly (not tick()), so the queue holds only the recheck. + assertEquals(1, scheduler.getQueue().size()); + ScheduledFuture scheduled = (ScheduledFuture) scheduler.getQueue().peek(); + assertTrue(scheduled.getDelay(TimeUnit.SECONDS) > 100, + "the delay must clear the worst-case fleet_ask lapse window (up to 115s)"); + + // An unchanged tick must not queue a second one (CB-580's fire-once rule extends to this). + monitor.reportTransition("term_a", HealthState.GONE); + assertEquals(1, scheduler.getQueue().size()); + monitor.stop(); + } + + @Test void recheckIsANoOpOnceTheTargetHasRecovered() { + RecordingFailTarget failTarget = new RecordingFailTarget(); + FleetHealthMonitor monitor = monitorWith(failTarget); + monitor.reportTransition("term_a", HealthState.GONE); + assertEquals(1, failTarget.calls.size()); + + monitor.reportTransition("term_a", HealthState.IDLE); // recovered before the recheck fired + monitor.recheckTerminalTarget("term_a", HealthState.GONE); + + assertEquals(1, failTarget.calls.size(), "a recovered target must not be reached into again"); + monitor.stop(); + } + + @Test void recheckIsANoOpForATargetItNeverObserved() { + // Mirrors "left the roster": tick() prunes states.keySet() to the current roster on release + // (see FleetHealthMonitor.tick), so a target this monitor never recorded is the same case. + RecordingFailTarget failTarget = new RecordingFailTarget(); + FleetHealthMonitor monitor = monitorWith(failTarget); + + monitor.recheckTerminalTarget("term_never_seen", HealthState.GONE); + + assertEquals(0, failTarget.calls.size(), "an untracked/released target must not be reached into"); + monitor.stop(); + } + + /** + * The scenario from the ticket, end to end, driven through the real {@link MessageService}: a + * target's ask is still genuinely open when health first observes GONE (sweep must skip it, + * exactly as {@code abandonDoesNotFailAnAsyncTicketWaitingForAnAnswer} pins), the ask then lapses + * on its own, an unchanged tick still must not refire, and only the delayed recheck sweeps the + * now-lapsed ticket to FAILED. + */ + @Test void delayedRecheckSweepsATicketWhoseAskLapsedAfterGoneWasFirstObserved() throws Exception { + FakeHerdr herdr = new FakeHerdr().withAgent("worker", "term_a", "pane-term_a", "tab_a") + .readText("$ prompt"); + AgentControl agents = new AgentControl(herdr); + Rendezvous rendezvous = new Rendezvous(); + Injector injector = new Injector(agents); + InMemoryReplyInbox inbox = new InMemoryReplyInbox(); + inbox.own("term_a"); + MessageService messages = new MessageService(agents, injector, rendezvous, inbox); + FleetHealthMonitor monitor = new FleetHealthMonitor(agents, + () -> List.of(member("term_a", MemberSession.State.READY, 0, 0)), messages, + new ScheduledThreadPoolExecutor(1), () -> 1, 60, 600, messages::abandon); + + String ticket = messages.sendAsync("term_a", "task that asks"); + long deadline = System.currentTimeMillis() + 2000; + while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) { + Thread.sleep(5); + } + assertTrue(rendezvous.isWaiting("term_a"), "the async send should have opened its waiter"); + injector.onStatus("term_a", AgentStatus.IDLE); // deliver the task + injector.onStatus("term_a", AgentStatus.WORKING); // the worker picks it up + + // The worker asks, with a short timeout so its own fleet_ask lapses quickly in test time. + CompletableFuture ask = CompletableFuture.supplyAsync( + () -> messages.ask("term_a", "which config?", 200)); + awaitPhase(messages, ticket, MessageService.Phase.ASKING); + + monitor.reportTransition("term_a", HealthState.GONE); + assertEquals(MessageService.Phase.ASKING, messages.poll(ticket).phase(), + "the first sweep must not fail a ticket that is still genuinely being asked"); + + // The worker's own fleet_ask now lapses on its own — task.question clears to null. + assertEquals(MessageService.AskOutcome.TIMED_OUT, ask.get(5, TimeUnit.SECONDS).outcome()); + + // An unchanged tick still must not refire (CB-580). + monitor.reportTransition("term_a", HealthState.GONE); + assertEquals(MessageService.Phase.PENDING, messages.poll(ticket).phase()); + + // The delayed recheck scheduled for the original transition finally sweeps it. + monitor.recheckTerminalTarget("term_a", HealthState.GONE); + assertEquals(MessageService.Phase.FAILED, messages.poll(ticket).phase()); + + monitor.stop(); + } + + private static MessageService.TaskView awaitPhase(MessageService messages, String ticket, + MessageService.Phase phase) throws InterruptedException { + long deadline = System.currentTimeMillis() + 2000; + MessageService.TaskView view; + do { + view = messages.poll(ticket); + if (view.phase() == phase) { + return view; + } + Thread.sleep(5); + } while (System.currentTimeMillis() < deadline); + assertEquals(phase, view.phase()); + return view; + } }