fleetd #368: a stale lead delegation binding must not shadow the primary fallback
PrimaryRegistry.forgetDelegation only fires when a WORKER is released, never when the delegating LEAD terminal itself disappears (closed, crashed, or relaunched). A stale, non-null leadByTarget entry always beat nudgeTargetFor's single-primary fallback, so a dead lead silently swallowed every reply nudge for its workers. Fix: ReplyPushLoop now verifies (via the same agents.status check decide() already uses every tick) that a recorded delegating lead is actually live before trusting it. A dead lead is treated as if never recorded — self-healing the binding (mirroring AgentControl.paneByTerminal's self-heal on agent_not_found) and falling through to PrimaryRegistry's existing fallback.
This commit is contained in:
@@ -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).
|
||||
*
|
||||
* <p><strong>The bug this closes.</strong> {@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.
|
||||
*
|
||||
* <p><strong>The fix.</strong> 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<String> resolveLiveLead(String target) {
|
||||
Optional<String> 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<String> profiles, int remainingCoolOffSeconds) {
|
||||
Map<String, List<String>> 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;
|
||||
|
||||
@@ -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.
|
||||
*
|
||||
* <p>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<String> 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<String, Object> p = params instanceof Map ? (Map<String, Object>) 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<String> promptTargets() {
|
||||
return List.copyOf(promptTargets);
|
||||
}
|
||||
|
||||
long sendCount() {
|
||||
return promptTargets.size();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user