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..0f5abd3 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/ReplyPushLoop.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/ReplyPushLoop.java @@ -2,6 +2,7 @@ package dev.ltms.fleet.msg; import dev.ltms.fleet.herdr.AgentControl; import dev.ltms.fleet.herdr.AgentStatus; +import dev.ltms.fleet.herdr.HerdrException; import dev.ltms.fleet.mcp.PrimaryRegistry; import dev.ltms.fleet.metrics.FleetMetrics; import dev.ltms.fleet.metrics.Metrics; @@ -13,6 +14,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 +377,80 @@ 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. + * + *

fleetd #368 review — only a positive "gone" reading forgets the binding. + * The first version of this method treated any {@code RuntimeException} from the probe + * as death, which is the #359 mistake repeated: a transient socket blip or a codec error on a + * perfectly live lead would silently and permanently unbind it, with no re-record ever coming. + * That is destructive on one bad reading, exactly what #359 shipped a two-reading guard to avoid + * for the analogous lead-tab-liveness question. {@link #isLive} now matches + * {@code AgentControl.agentCall}'s own narrower rule (see its {@code agent_not_found} check): only + * that specific, affirmative "herdr has no such agent" signal counts as gone. Every other failure + * — timeout, transport error, a decode error — is treated as still live and the binding is left + * alone, because guessing wrong here is unrecoverable while guessing "live" merely costs one more + * retry on the next tick, which {@link #decide} already tolerates. + */ + 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 {@code lead} should still be trusted: {@code false} only when herdr affirmatively + * reports the terminal gone ({@code agent_not_found}), never on a merely inconclusive failure. + * + *

fleetd #368 review: an earlier version returned {@code false} for any {@code RuntimeException}, + * which made a transient herdr hiccup on a live lead indistinguishable from the lead actually + * being dead — and the caller's response to {@code false} ({@code forgetDelegation}) is + * destructive and permanent. Narrowed to the one code {@code AgentControl.agentCall} itself + * already treats as a genuine, resolvable absence (see its {@code agent_not_found} handling) — + * every other {@code RuntimeException} is treated as "still live" and the binding survives to be + * probed again next time, which costs nothing worse than one more retry. + */ + private boolean isLive(String lead) { + try { + agents.status(lead); + return true; + } catch (RuntimeException e) { + boolean gone = e instanceof HerdrException he && "agent_not_found".equals(he.code()); + if (gone) { + log.debug("push: lead {} no longer exists ({})", lead, e.toString()); + } else { + log.debug("push: liveness check for lead {} was inconclusive ({}); treating as live " + + "rather than risk destroying a live binding", lead, e.toString()); + } + return !gone; + } + } + // --- public entrypoints ---------------------------------------------------------------------- /** @@ -385,7 +461,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 +489,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 +524,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 +552,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 +574,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..2543044 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,101 @@ 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"); + } + + /** + * fleetd #368 review, must-fix: the first version of {@code isLive} treated any + * {@code RuntimeException} from the liveness probe as "the lead is gone" — indistinguishable + * from a transient herdr hiccup (a socket blip, a decode error) on a lead that is actually + * still live. The consequence of that misdiagnosis is destructive and permanent + * ({@code forgetDelegation}), which is the exact #359 mistake repeated two days later: a single + * bad reading must never destroy a live binding. This pins the narrower rule — only an + * affirmative {@code agent_not_found} may forget a binding; a merely inconclusive failure must + * leave the binding alone, and the lead must still be nudged once the probe recovers. + */ + @Test + void aTransientLivenessFailureMustNotForgetABindingToAStillLiveLead() throws Exception { + // OTHER_PRIMARY is delegated to and genuinely live — its FIRST agent.get call fails with a + // transient, non-agent_not_found HerdrException (a transport-level failure, code null, + // exactly what a socket blip looks like), then succeeds on every call after. + registry.recordDelegation(WORKER, OTHER_PRIMARY); + + var rec = new FlakyThenLiveHerdrClient(OTHER_PRIMARY); + agents = new AgentControl(rec); + inbox.publish(WORKER, "m1", "hello"); + + loop(2, 50).onReplyQueued(WORKER); + + assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS), + "the nudge must still reach the live lead once the transient failure clears"); + assertEquals(List.of(OTHER_PRIMARY), rec.promptTargets(), + "the nudge must go to the lead that was only transiently unreachable, not the " + + "unrelated pinned primary"); + assertEquals(OTHER_PRIMARY, registry.nudgeTargetFor(WORKER).orElseThrow(), + "a merely transient failure must not forget the binding to a lead that is actually " + + "still live"); + } + // --- nudge format -------------------------------------------------------------------------- @Test @@ -1098,4 +1195,101 @@ 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() { + } + } + + /** + * Fake herdr client for fleetd #368 review: {@code flakyTarget}'s FIRST {@code agent.get} call + * fails with a transient, non-{@code agent_not_found} {@code HerdrException} — a transport-level + * failure (code {@code null}), exactly what a socket blip or a decode error on a perfectly live + * lead looks like — then succeeds ({@code idle}) on every call after. Used to prove a merely + * inconclusive failure must not be treated as the lead being gone. + */ + private static final class FlakyThenLiveHerdrClient implements HerdrClient { + private final String flakyTarget; + private final AtomicInteger getCalls = new AtomicInteger(); + private final List promptTargets = Collections.synchronizedList(new ArrayList<>()); + volatile CountDownLatch sendLatch = new CountDownLatch(1); + + FlakyThenLiveHerdrClient(String flakyTarget) { + this.flakyTarget = flakyTarget; + } + + @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 (flakyTarget.equals(target) && getCalls.getAndIncrement() == 0) { + throw new HerdrException("herdr socket read timed out"); // transport failure, code == 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); + } + + @Override + public void close() { + } + } }