diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/ReplyPushLoop.java b/fleetd/src/main/java/dev/ltms/fleet/msg/ReplyPushLoop.java index 653d9c9..45d9dd2 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/ReplyPushLoop.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/ReplyPushLoop.java @@ -13,6 +13,7 @@ import java.util.Collection; import java.util.HashSet; import java.util.List; import java.util.Map; +import java.util.Optional; import java.util.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ScheduledExecutorService; @@ -375,6 +376,51 @@ public final class ReplyPushLoop { return Action.WAIT_BUSY; } + /** + * Resolve who to nudge about {@code target}, the way every public entry point below wants it: + * {@link PrimaryRegistry#nudgeTargetFor}, but only after checking the delegating lead it names + * is still actually there (fleetd #368). + * + *

The bug this closes. {@code PrimaryRegistry.forgetDelegation} is wired to + * exactly one event — a worker's release — because that is the only teardown the daemon already + * observes for a session in this map. Nothing removes a binding when the LEAD half goes away: a + * lead that is closed, crashes, or is relaunched leaves {@code leadByTarget} entries pointing at + * a terminal that no longer exists. {@code nudgeTargetFor} falls back to the single known + * primary only when the map holds nothing for {@code target} — a stale non-null entry beats the + * fallback every time, which is exactly backwards: the fallback's own javadoc argues it is safe + * precisely in the case a stale entry now hides. + * + *

The fix. Before trusting a recorded delegation, probe the lead the same + * way {@link #decide} already does every tick ({@code agents.status}) — cheap, since it is a + * local herdr round-trip, and it is the same signal {@code AgentControl.paneByTerminal} already + * trusts to tell a genuinely dead target from a live one. A lead that fails the probe is treated + * as if it had never been recorded: the stale entry is forgotten (self-healing, exactly like + * {@code AgentControl.paneByTerminal} already does on {@code agent_not_found}) and resolution is + * retried, which now reaches the fallback {@code nudgeTargetFor} was built to reach — the same + * empty-map state its javadoc already argues is correct. + */ + private Optional resolveLiveLead(String target) { + Optional lead = primaryRegistry.nudgeTargetFor(target); + if (lead.isEmpty() || isLive(lead.get())) { + return lead; + } + log.debug("push: lead {} delegated to for {} is no longer live, forgetting the stale binding " + + "and falling back", lead.get(), target); + primaryRegistry.forgetDelegation(target); + return primaryRegistry.nudgeTargetFor(target); + } + + /** Whether herdr still reports a status for {@code lead} — false for a closed/dead terminal. */ + private boolean isLive(String lead) { + try { + agents.status(lead); + return true; + } catch (RuntimeException e) { + log.debug("push: liveness check failed for lead {}: {}", lead, e.toString()); + return false; + } + } + // --- public entrypoints ---------------------------------------------------------------------- /** @@ -385,7 +431,7 @@ public final class ReplyPushLoop { * backstop until a lead is recorded. */ public void onReplyQueued(String target) { - var lead = primaryRegistry.nudgeTargetFor(target); + var lead = resolveLiveLead(target); if (lead.isEmpty()) { log.debug("push: no lead is known to be waiting on {}, skipping reminder", target); return; @@ -413,7 +459,7 @@ public final class ReplyPushLoop { * @param failed whether the ticket ended in a failure phase rather than {@code DONE} */ public void onTicketTerminal(String ticket, String target, boolean failed) { - var lead = primaryRegistry.nudgeTargetFor(target); + var lead = resolveLiveLead(target); if (lead.isEmpty()) { log.debug("push: no lead is known to be waiting on ticket {} (target {}), skipping nudge", ticket, target); @@ -448,7 +494,7 @@ public final class ReplyPushLoop { * @param question the question text */ public void onQuestionOpened(String ticket, String target, String turnId, String question) { - var lead = primaryRegistry.nudgeTargetFor(target); + var lead = resolveLiveLead(target); if (lead.isEmpty()) { log.debug("push: no lead is known to be waiting on {}'s question (turnId {}), skipping nudge", target, turnId); @@ -476,7 +522,7 @@ public final class ReplyPushLoop { Collection profiles, int remainingCoolOffSeconds) { Map> targetsByLead = new ConcurrentHashMap<>(); for (String target : targets) { - var lead = primaryRegistry.nudgeTargetFor(target); + var lead = resolveLiveLead(target); if (lead.isEmpty()) { log.warn("push: backend incident {} has no known lead for target {}", incidentId, target); continue; @@ -498,7 +544,7 @@ public final class ReplyPushLoop { * Without an owning lead, emit a warning because no control can act on the target. */ public void onBackendTargetUnmapped(String target, String reason) { - var lead = primaryRegistry.nudgeTargetFor(target); + var lead = resolveLiveLead(target); if (lead.isEmpty()) { log.warn("push: backend target {} could not map to a credential: {}", target, reason); return; diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/ReplyPushLoopTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/ReplyPushLoopTest.java index 4d97c49..bdb5e2a 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/ReplyPushLoopTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/ReplyPushLoopTest.java @@ -4,6 +4,7 @@ import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; import dev.ltms.fleet.herdr.AgentControl; import dev.ltms.fleet.herdr.HerdrClient; +import dev.ltms.fleet.herdr.HerdrException; import dev.ltms.fleet.mcp.PrimaryRegistry; import dev.ltms.fleet.metrics.FleetMetrics; import dev.ltms.fleet.metrics.Metrics; @@ -44,6 +45,7 @@ class ReplyPushLoopTest { private static final String WORKER = "term_worker"; private static final String WORKER2 = "term_worker2"; private static final String OTHER_PRIMARY = "term_other_primary"; + private static final String DEAD_LEAD = "term_dead_lead"; private static final ObjectMapper MAPPER = new ObjectMapper(); private PrimaryRegistry registry; @@ -207,6 +209,68 @@ class ReplyPushLoopTest { "exactly " + cap + " agent.prompt calls (cap=" + cap + ")"); } + // --- fleetd #368: a lead's delegation binding must not outlive the lead --------------------- + + /** + * The bug: {@code PrimaryRegistry.forgetDelegation} is wired to a worker's release, never to + * the delegating lead's own disappearance, so a lead that closed, crashed, or was relaunched + * leaves {@code leadByTarget} pointing at a terminal herdr no longer knows. Before the fix, + * {@code onReplyQueued} took that stale, non-null entry at face value — {@code nudgeTargetFor} + * only ever falls back to the pinned primary when the map holds nothing for the target — so + * the nudge's only schedule ran against the dead terminal forever and the live primary never + * heard about the reply through this path. + * + *

This drives {@link ReplyPushLoop#onReplyQueued(String)} itself (not {@code PrimaryRegistry} + * directly), because the registry lookup was never the defect — the caller trusting it without + * checking liveness was. A test that only asserted on {@code PrimaryRegistry.nudgeTargetFor} + * would pass whether or not {@code ReplyPushLoop} ever adopted the fix. + */ + @Test + void aStaleLeadBindingFallsBackToTheLiveLeadInsteadOfNudgingADeadTerminal() throws Exception { + // PRIMARY is the single known (pinned) lead — set up in @BeforeEach via `registry`. + // DEAD_LEAD is a second lead that once delegated to WORKER and is now gone: herdr reports + // agent_not_found for it, exactly as it would for a closed/crashed/relaunched terminal. + registry.recordDelegation(WORKER, DEAD_LEAD); + + var rec = new DeadLeadHerdrClient(DEAD_LEAD); + agents = new AgentControl(rec); + inbox.publish(WORKER, "m1", "hello"); + + loop(1, 50).onReplyQueued(WORKER); + + assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS), + "the nudge should still reach the live primary, not silently vanish with the dead lead"); + assertEquals(List.of(PRIMARY), rec.promptTargets(), + "the nudge must be sent to the live primary, never to the dead lead's terminal"); + assertEquals(PRIMARY, registry.nudgeTargetFor(WORKER).orElseThrow(), + "the stale binding must be forgotten (self-healed) once found dead, exactly like " + + "AgentControl.paneByTerminal already does on agent_not_found"); + } + + /** + * Same dead binding, but with no pinned primary to fall back to (the multi-lead, no-fallback + * case {@code PrimaryRegistry.nudgeTargetFor}'s own javadoc already covers): the loop must + * never nudge the dead terminal, and must not spin — no schedule starts at all once the stale + * binding resolves to empty, same as if the map had never held an entry for this target. + */ + @Test + void aStaleLeadBindingWithNoFallbackNeverNudgesTheDeadTerminal() throws Exception { + var unpinned = new PrimaryRegistry(null); + unpinned.recordDelegation(WORKER, DEAD_LEAD); + + var rec = new DeadLeadHerdrClient(DEAD_LEAD); + agents = new AgentControl(rec); + inbox.publish(WORKER, "m1", "hello"); + + var loop = new ReplyPushLoop(unpinned, agents, inbox, scheduler, 1, 50); + loop.onReplyQueued(WORKER); + + Thread.sleep(200); + assertEquals(0, rec.sendCount(), "no lead is live to nudge, so nothing should ever be sent"); + assertTrue(unpinned.nudgeTargetFor(WORKER).isEmpty(), + "the stale binding must be forgotten even when there is no fallback to hand back"); + } + // --- nudge format -------------------------------------------------------------------------- @Test @@ -1098,4 +1162,54 @@ class ReplyPushLoopTest { public void close() { } } + + /** + * Fake herdr client for fleetd #368: {@code deadTarget} is a terminal herdr genuinely no + * longer knows about — {@code agent.get} fails with {@code agent_not_found} exactly as + * {@code AgentControl.agentCall} expects for a real dead/closed pane (see its javadoc). Every + * other target reports {@code idle} (injectable). Records the {@code target} named by every + * {@code agent.prompt} call, so a test can prove which terminal actually got nudged. + */ + private static final class DeadLeadHerdrClient implements HerdrClient { + private final String deadTarget; + private final List promptTargets = Collections.synchronizedList(new ArrayList<>()); + volatile CountDownLatch sendLatch = new CountDownLatch(1); + + DeadLeadHerdrClient(String deadTarget) { + this.deadTarget = deadTarget; + } + + @Override + @SuppressWarnings("unchecked") + public JsonNode call(String method, Object params) { + Map p = params instanceof Map ? (Map) params : Map.of(); + if ("agent.get".equals(method)) { + String target = String.valueOf(p.get("target")); + if (deadTarget.equals(target)) { + throw new HerdrException("no such agent: " + target, "agent_not_found", null); + } + return MAPPER.createObjectNode() + .set("agent", MAPPER.createObjectNode() + .put("terminal_id", target) + .put("agent_status", "idle")); + } + if ("agent.prompt".equals(method)) { + promptTargets.add(String.valueOf(p.get("target"))); + sendLatch.countDown(); + } + return MAPPER.createObjectNode(); + } + + List promptTargets() { + return List.copyOf(promptTargets); + } + + long sendCount() { + return promptTargets.size(); + } + + @Override + public void close() { + } + } }