diff --git a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java index 5cb7cf8..cb6028a 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java +++ b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java @@ -904,6 +904,32 @@ public final class Fleetd { throw (T) t; } + /** + * The {@link PrimaryRegistry} lookup for "which terminal currently hosts the lead named + * {@code name}" — the inverse of {@code liveLeadTerminals} (terminal id → lead name), read live + * on every call so a lead discovered, rolled, or lost since the last call is reflected without + * a restart. Returns {@code null} when no currently recognised lead carries that name — a name + * that is not a lead at all (an architect slot, a collaborator), or a lead whose tab the scan + * cannot currently place (just rolled, off-host, non-herdr). + * + * @param liveLeadTerminals terminal id → lead name for every CURRENTLY recognised lead, normally + * the same {@code leads} supplier {@code main} already builds for + * {@code HerdrRouter}/{@link #leadSeatLookup} + */ + static Function currentTerminalForName(Supplier> liveLeadTerminals) { + return name -> { + if (name == null) { + return null; + } + for (var entry : liveLeadTerminals.get().entrySet()) { + if (name.equals(entry.getValue())) { + return entry.getKey(); + } + } + return null; + }; + } + /** * fleetd #480: construct the {@link LeadRollover} executor only when {@code leadRollover:} is * present at startup — the same presence gate {@code leadHeartbeat:} uses just above this diff --git a/fleetd/src/main/java/dev/ltms/fleet/FleetdAssembly.java b/fleetd/src/main/java/dev/ltms/fleet/FleetdAssembly.java index 72d6f1c..541e0bb 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/FleetdAssembly.java +++ b/fleetd/src/main/java/dev/ltms/fleet/FleetdAssembly.java @@ -393,7 +393,8 @@ final class FleetdAssembly { ports.leadMailboxOpener()); // CB-307: learn the primary's terminal from orchestration tool calls (or pin from config). String pinnedPrimaryTerminal = cfg.primary() != null ? cfg.primary().terminal() : null; - PrimaryRegistry primaryRegistry = new PrimaryRegistry(pinnedPrimaryTerminal); + PrimaryRegistry primaryRegistry = new PrimaryRegistry(pinnedPrimaryTerminal, + Fleetd.currentTerminalForName(leads)); // CB-532: `primary.terminal` is superseded and no longer needed for either of its jobs. if (pinnedPrimaryTerminal != null && !pinnedPrimaryTerminal.isBlank()) { log.warn("primary.terminal is DEPRECATED (CB-532) and can be deleted: identity now comes " diff --git a/fleetd/src/main/java/dev/ltms/fleet/lead/LeadRollover.java b/fleetd/src/main/java/dev/ltms/fleet/lead/LeadRollover.java index daeabda..1a0dcf6 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/lead/LeadRollover.java +++ b/fleetd/src/main/java/dev/ltms/fleet/lead/LeadRollover.java @@ -93,7 +93,7 @@ import java.util.function.Supplier; * *

Identity is resolved by the caller, never looked up here — a second fleetd #480 * correction. The first version resolved the pane to clear via {@code - * PrimaryRegistry#primaryTerminal()}. That is correct for a background loop with no caller (see + * PrimaryRegistry#currentPrimaryTerminal()}. That is correct for a background loop with no caller (see * {@code dev.ltms.fleet.msg.LeadHeartbeatLoop}), but wrong here and a violation of this project's * own charter invariant 3 — "identity comes from the connection, never an argument." This daemon * can hold more than one labelled lead tab (see {@code LeadLauncher}'s fleetd #359 two-reading diff --git a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java index e835ed5..9cb1099 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java @@ -489,7 +489,13 @@ public final class FleetMcp { // session lock and queued delivery — via the accepted-delivery callback, never at // request time. A concurrent sender that times out BUSY therefore cannot steal a // live turn's reply routing without ever owning the turn. - Runnable onAccepted = () -> primaryRegistry.recordDelegation(target, callerTerminal); + // Only a lead's name is ever resolvable back to a current terminal (PrimaryRegistry + // only looks it up among currently recognised leads) — an architect or collaborator + // name would never match there anyway, but passing null for them keeps the intent + // explicit rather than relying on that lookup to filter it out. + String delegatorName = caller.isPrimary() ? caller.name() : null; + Runnable onAccepted = () -> + primaryRegistry.recordDelegation(target, callerTerminal, delegatorName); // wait defaults to true (block for the reply); wait:false is fire-and-poll. return Boolean.FALSE.equals(a.get("wait")) ? sendAsync(messages, target, content, onAccepted, workers.profiles(), caller) @@ -731,7 +737,7 @@ public final class FleetMcp { */ static void recordPrimarySingleton(PrimaryRegistry registry, String callerTerminal, Principal caller) { if (caller != null && caller.isPrimary()) { - registry.record(callerTerminal); + registry.record(callerTerminal, caller.name()); } } diff --git a/fleetd/src/main/java/dev/ltms/fleet/mcp/PrimaryRegistry.java b/fleetd/src/main/java/dev/ltms/fleet/mcp/PrimaryRegistry.java index 3682bee..3f2a61c 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/mcp/PrimaryRegistry.java +++ b/fleetd/src/main/java/dev/ltms/fleet/mcp/PrimaryRegistry.java @@ -6,6 +6,7 @@ import org.slf4j.LoggerFactory; import java.util.Optional; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicReference; +import java.util.function.Function; /** * Single-slot, thread-safe registry for the primary's herdr {@code terminal_id}. @@ -18,13 +19,25 @@ import java.util.concurrent.atomic.AtomicReference; *

The push loop ({@code ReplyPushLoop}) uses {@link #isKnown()} to decide * whether active nudging is possible; an empty registry means the primary is * off-host or non-herdr and delivery falls back to pull. + * + *

A learned terminal can go stale; a configured lead's name cannot. A lead that + * is rolled (a fresh pane replacing the old one) keeps its name but gets a new {@code terminal_id}. + * So every terminal this class learns — the singleton and each per-target delegation — is recorded + * together with the delegating lead's name, when the caller carries one. {@link + * #currentPrimaryTerminal()} and {@link #nudgeTargetFor(String)} resolve that name back to a + * terminal through the live {@code currentTerminalForName} lookup before falling back to the + * terminal that was actually recorded. A caller with no name (an unnamed primary, an architect, a + * collaborator — none of those are leads a lookup keyed on lead names can resolve) is tracked by + * terminal alone, exactly as before this indirection existed. */ public final class PrimaryRegistry { private static final Logger log = LoggerFactory.getLogger(PrimaryRegistry.class); private final AtomicReference terminal = new AtomicReference<>(); + private final AtomicReference primaryName = new AtomicReference<>(); private final boolean pinned; + private final Function currentTerminalForName; /** * CB-532: worker terminal → the lead that delegated to it. The single slot above answers "who is @@ -33,12 +46,32 @@ public final class PrimaryRegistry { * other lead's delegations. This map answers the question that actually matters — "who is * waiting on THIS worker" — and is what lets {@code primary.terminal} be retired. */ - private final ConcurrentHashMap leadByTarget = new ConcurrentHashMap<>(); + private final ConcurrentHashMap leadByTarget = new ConcurrentHashMap<>(); + + /** A recorded delegator: the terminal learned from call traffic, and its name, if it has one. */ + private record Delegation(String terminal, String name) { + } /** * @param pinnedTerminal an optional pinned terminal from config ({@code null}/blank = unpinned) */ public PrimaryRegistry(String pinnedTerminal) { + this(pinnedTerminal, name -> null); + } + + /** + * As above, with a live {@code lead name → current terminal} lookup — normally the inverse of + * the same {@code terminal_id → lead name} supplier {@code CallerResolver} and the lead-tab + * scan already read. A lookup that cannot place a name (it is not a currently recognised lead, + * or no lookup is wired) returns {@code null}, and every resolution here falls back to the + * terminal that was actually recorded. + * + * @param pinnedTerminal an optional pinned terminal from config ({@code null}/blank = + * unpinned) + * @param currentTerminalForName lead name → its current terminal, or {@code null} if that name + * is not a currently recognised lead + */ + public PrimaryRegistry(String pinnedTerminal, Function currentTerminalForName) { if (pinnedTerminal != null && !pinnedTerminal.isBlank()) { this.terminal.set(pinnedTerminal); this.pinned = true; @@ -46,19 +79,31 @@ public final class PrimaryRegistry { } else { this.pinned = false; } + this.currentTerminalForName = currentTerminalForName != null ? currentTerminalForName : name -> null; } /** - * Record a terminal_id. No-op when: + * Record a terminal_id, with no lead name. No-op when: *

*/ public void record(String terminalId) { + record(terminalId, null); + } + + /** + * As {@link #record(String)}, additionally recording the caller's name — present for a + * configured lead, {@code null} for an unnamed primary. The name is what lets {@link + * #currentPrimaryTerminal()} keep nudging the same lead across a roll even though its terminal + * changed. + */ + public void record(String terminalId, String name) { if (pinned) return; if (terminalId == null || terminalId.isBlank()) return; String prev = terminal.getAndSet(terminalId); + primaryName.set(blankToNull(name)); if (prev == null) { log.debug("primary terminal learned: {}", terminalId); } else if (!prev.equals(terminalId)) { @@ -67,7 +112,8 @@ public final class PrimaryRegistry { } /** - * Record that {@code leadTerminal} owns the accepted delegation of worker {@code target} (CB-532). + * Record that {@code leadTerminal} owns the accepted delegation of worker {@code target} + * (CB-532), with no lead name. * *

Called from the {@code MessageService} accepted-delivery hook — only after a send has won * the session's send lock and queued delivery — where both halves are known (CB-548). It is @@ -77,10 +123,19 @@ public final class PrimaryRegistry { * lead that most recently delegated to it, which is the one waiting. */ public void recordDelegation(String target, String leadTerminal) { + recordDelegation(target, leadTerminal, null); + } + + /** + * As {@link #recordDelegation(String, String)}, additionally recording the delegating lead's + * name when the caller carries one. See {@link #record(String, String)} for why the name + * matters. + */ + public void recordDelegation(String target, String leadTerminal, String leadName) { if (target == null || target.isBlank() || leadTerminal == null || leadTerminal.isBlank()) { return; } - leadByTarget.put(target, leadTerminal); + leadByTarget.put(target, new Delegation(leadTerminal, blankToNull(leadName))); } /** Forget a worker's delegating lead — call on release, so a torn-down session leaks nothing. */ @@ -99,19 +154,59 @@ public final class PrimaryRegistry { * recorded delegation there is no right answer, so this returns empty rather than guessing — * delivery degrades to pull, which is exactly what the durable inbox is for, instead of * interrupting the wrong lead with someone else's result. + * + *

A delegation recorded with a name is resolved to that lead's current terminal + * first — see {@link #currentTerminalForName} — so a lead that has since been rolled is still + * reachable here, not just the pane that delegated the work originally. */ public Optional nudgeTargetFor(String target) { - String lead = target == null ? null : leadByTarget.get(target); - return lead != null ? Optional.of(lead) : Optional.ofNullable(terminal.get()); + Delegation delegation = target == null ? null : leadByTarget.get(target); + if (delegation != null) { + return Optional.of(resolveCurrent(delegation.terminal(), delegation.name())); + } + return currentPrimaryTerminal(); } - /** The known primary terminal, or empty if not yet learned (and not pinned). */ + /** + * The known primary terminal, or empty if not yet learned (and not pinned) — the raw value as + * it was recorded, with no attempt to resolve a named lead's current pane. Callers that need a + * nudge destination which survives a lead roll want {@link #currentPrimaryTerminal()} instead. + */ public Optional primaryTerminal() { return Optional.ofNullable(terminal.get()); } + /** + * The terminal to nudge for the singleton primary right now: the recorded name resolved to its + * current terminal when one was recorded and is still a recognised lead, otherwise the terminal + * that was actually recorded — empty only when nothing has been learned or pinned at all. + */ + public Optional currentPrimaryTerminal() { + String learned = terminal.get(); + if (learned == null) { + return Optional.empty(); + } + return Optional.of(resolveCurrent(learned, primaryName.get())); + } + /** {@code true} once a terminal has been recorded (or was pinned at construction). */ public boolean isKnown() { return terminal.get() != null; } + + /** + * {@code learnedTerminal}, unless {@code name} is non-null and {@code currentTerminalForName} + * currently places that name at a different, live terminal — in which case the live one wins. + */ + private String resolveCurrent(String learnedTerminal, String name) { + if (name == null) { + return learnedTerminal; + } + String current = currentTerminalForName.apply(name); + return current != null && !current.isBlank() ? current : learnedTerminal; + } + + private static String blankToNull(String s) { + return s == null || s.isBlank() ? null : s; + } } diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/LeadHeartbeatLoop.java b/fleetd/src/main/java/dev/ltms/fleet/msg/LeadHeartbeatLoop.java index 2a5e211..af156f5 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/LeadHeartbeatLoop.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/LeadHeartbeatLoop.java @@ -307,12 +307,12 @@ public final class LeadHeartbeatLoop { * (mirroring {@link ReplyPushLoop#tick(String)}) so tests can drive it directly with a fake clock and a * fake {@link AgentControl} instead of racing the scheduler thread. */ void tick() { - boolean leadKnown = primaryRegistry.primaryTerminal().isPresent(); + boolean leadKnown = primaryRegistry.currentPrimaryTerminal().isPresent(); FleetState fleet = snapshot(inbox, roster); AgentStatus status = AgentStatus.UNKNOWN; LeadContextGauge.Reading reading = LeadContextGauge.Reading.unknown(); if (leadKnown) { - String leadTerminal = primaryRegistry.primaryTerminal().orElseThrow(); + String leadTerminal = primaryRegistry.currentPrimaryTerminal().orElseThrow(); try { status = agents.status(leadTerminal); } catch (RuntimeException e) { @@ -367,7 +367,7 @@ public final class LeadHeartbeatLoop { // itself should be built from. Otherwise a HIGH stretch that is still latched would never see the // notice at all, defeating the very check this fixes. String notice = contextNotice(contextHighNudge, reading, contextNotified, requireOperatorConfirm); - var lead = primaryRegistry.primaryTerminal(); + var lead = primaryRegistry.currentPrimaryTerminal(); boolean sent = lead.isPresent() && trySend(lead.get(), fleet.nudgeText() + notice, notice); // The latch becomes true only when all three hold: decide() chose to notify, a notice was // actually included in the text, and the send reached the pane without throwing. Whenever no 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 b79c51e..74f048e 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/ReplyPushLoop.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/ReplyPushLoop.java @@ -439,6 +439,15 @@ public final class ReplyPushLoop { * — 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. + * + *

The fallback is probed too. {@code PrimaryRegistry.nudgeTargetFor} already + * resolves a named delegator to its current terminal before this method ever sees it, which + * keeps a rolled lead's per-target binding live. What that resolution cannot fix is a caller + * that was never recorded with a name at all — an unnamed primary, or a lead whose tab the + * scanner cannot currently see — where the fallback it returns is still the raw terminal last + * learned from call traffic. This method returns that fallback only after the same liveness + * check, and gives up for this tick (an empty result, exactly like "no lead known at all") rather + * than hand a caller a second stale address un-probed. */ private Optional resolveLiveLead(String target) { Optional lead = primaryRegistry.nudgeTargetFor(target); @@ -448,7 +457,12 @@ public final class ReplyPushLoop { 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); + Optional fallback = primaryRegistry.nudgeTargetFor(target); + if (fallback.isEmpty() || isLive(fallback.get())) { + return fallback; + } + log.debug("push: fallback lead {} for {} is also not live, skipping this tick", fallback.get(), target); + return Optional.empty(); } /** diff --git a/fleetd/src/test/java/dev/ltms/fleet/mcp/PrimaryRegistryTest.java b/fleetd/src/test/java/dev/ltms/fleet/mcp/PrimaryRegistryTest.java index 5132c23..6d53ff1 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/mcp/PrimaryRegistryTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/mcp/PrimaryRegistryTest.java @@ -169,4 +169,101 @@ class PrimaryRegistryTest { assertTrue(reg.nudgeTargetFor("term_worker").isEmpty()); assertTrue(reg.nudgeTargetFor(null).isEmpty()); } + + // ── fleetd #737 unit 3: a named lead's terminal is resolved live, not just recorded ───────── + + /** + * The whole point of carrying a name: a lead that has been rolled keeps its name but gets a + * fresh terminal. {@code currentTerminalForName} stands in for the live lead-tab scan here — + * it reports the lead now sits on a different terminal than the one that was recorded — and + * {@code nudgeTargetFor} must follow the name to that current terminal, not the stale one. + */ + @Test + void nudgeTargetForFollowsARolledLeadsNameToItsCurrentTerminal() { + var reg = new PrimaryRegistry(null, name -> "opus".equals(name) ? "term_opus_after_roll" : null); + reg.recordDelegation("term_worker", "term_opus_before_roll", "opus"); + + assertEquals("term_opus_after_roll", reg.nudgeTargetFor("term_worker").orElseThrow(), + "the name must be resolved to the lead's CURRENT terminal, not the one recorded " + + "at delegation time"); + } + + /** + * The lookup cannot place every name — an architect/collaborator name (never a lead), or a lead + * whose tab the scan cannot currently see (just rolled, off-host, non-herdr). Either way the + * terminal actually recorded is still the right thing to try, exactly as before this unit. + */ + @Test + void nudgeTargetForFallsBackToTheRecordedTerminalWhenTheNameCannotBePlaced() { + var reg = new PrimaryRegistry(null, name -> null); // nothing is ever currently recognised + reg.recordDelegation("term_worker", "term_lead_recorded", "opus"); + + assertEquals("term_lead_recorded", reg.nudgeTargetFor("term_worker").orElseThrow()); + } + + /** The 2-arg {@code recordDelegation} overload records no name, so resolution never applies. */ + @Test + void recordDelegationWithNoNameIsNeverResolvedByLookup() { + var reg = new PrimaryRegistry(null, name -> { + throw new AssertionError("a delegation recorded with no name must never consult the lookup"); + }); + reg.recordDelegation("term_worker", "term_lead"); + + assertEquals("term_lead", reg.nudgeTargetFor("term_worker").orElseThrow()); + } + + /** As {@link #nudgeTargetForFollowsARolledLeadsNameToItsCurrentTerminal}, for the singleton. */ + @Test + void currentPrimaryTerminalFollowsARolledLeadsNameToItsCurrentTerminal() { + var reg = new PrimaryRegistry(null, name -> "sol".equals(name) ? "term_sol_after_roll" : null); + reg.record("term_sol_before_roll", "sol"); + + assertEquals("term_sol_after_roll", reg.currentPrimaryTerminal().orElseThrow()); + assertEquals("term_sol_before_roll", reg.primaryTerminal().orElseThrow(), + "primaryTerminal() stays the raw recorded value — currentPrimaryTerminal() is the " + + "one that resolves live"); + } + + @Test + void currentPrimaryTerminalFallsBackWhenTheNameCannotBePlaced() { + var reg = new PrimaryRegistry(null, name -> null); + reg.record("term_sol", "sol"); + + assertEquals("term_sol", reg.currentPrimaryTerminal().orElseThrow()); + } + + @Test + void currentPrimaryTerminalWithNoNameRecordedIsTheRawTerminal() { + var reg = new PrimaryRegistry(null, name -> { + throw new AssertionError("no name was ever recorded, the lookup must not be consulted"); + }); + reg.record("term_x"); + + assertEquals("term_x", reg.currentPrimaryTerminal().orElseThrow()); + } + + @Test + void currentPrimaryTerminalIsEmptyWhenNothingWasEverLearned() { + var reg = new PrimaryRegistry(null, name -> "anything"); + assertTrue(reg.currentPrimaryTerminal().isEmpty()); + } + + /** A pin never carries a name, so a pinned registry's singleton resolution is always a no-op. */ + @Test + void currentPrimaryTerminalForAPinIsNeverResolvedByLookup() { + var reg = new PrimaryRegistry("term_pinned", name -> { + throw new AssertionError("a pin carries no name, the lookup must not be consulted"); + }); + + assertEquals("term_pinned", reg.currentPrimaryTerminal().orElseThrow()); + } + + /** {@code nudgeTargetFor}'s fallback to the singleton is the resolved one, not the raw one. */ + @Test + void nudgeTargetForWithNoDelegationFallsBackToTheResolvedSingleton() { + var reg = new PrimaryRegistry(null, name -> "opus".equals(name) ? "term_opus_after_roll" : null); + reg.record("term_opus_before_roll", "opus"); + + assertEquals("term_opus_after_roll", reg.nudgeTargetFor("term_never_seen").orElseThrow()); + } } diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/LeadHeartbeatLoopTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/LeadHeartbeatLoopTest.java index 4cee349..e245d1c 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/LeadHeartbeatLoopTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/LeadHeartbeatLoopTest.java @@ -648,6 +648,42 @@ class LeadHeartbeatLoopTest { "the notice text must appear exactly once across all three sends: " + herdr.sentTexts()); } + // ── fleetd #737 unit 3: tick() nudges the lead's CURRENT terminal, not the learned one ─────── + + /** + * {@code tick()} reads {@code primaryRegistry.currentPrimaryTerminal()} both to check the lead's + * status and to send the nudge. Here the registry learned the lead's terminal under its name + * before a roll; {@code currentTerminalForName} stands in for the live lead-tab scan and reports + * the lead now sits on a different terminal. A correct tick must follow the name and nudge the + * new terminal — nudging the old one would mean the heartbeat lost the lead across its own roll. + */ + @Test + void tickNudgesTheLeadsCurrentTerminalAfterARoll() { + var herdr = new FailableHerdrClient("term_lead_after_roll"); + var now = new AtomicLong(NOW); + InMemoryReplyInbox inbox = new InMemoryReplyInbox(); + inbox.own(WORKER); + inbox.publish(WORKER, "m1", "hello"); + @SuppressWarnings("unchecked") + List[] rosterBox = new List[]{List.of(new MemberSession("p1", WORKER, "prof", + MemberRole.DEV, "/cwd", null, 0, 0, 0, MemberSession.State.READY, null, null))}; + AgentControl agents = new AgentControl(herdr); + PrimaryRegistry registry = new PrimaryRegistry(null, + name -> "opus".equals(name) ? "term_lead_after_roll" : null); + registry.record("term_lead_before_roll", "opus"); + ReplyPushLoop pushLoop = new ReplyPushLoop(registry, agents, inbox, scheduler, 5, 100_000); + LeadHeartbeatLoop loop = new LeadHeartbeatLoop(registry, agents, inbox, () -> rosterBox[0], pushLoop, + scheduler, now::get, IDLE_AFTER_NANOS, 100_000L, 0); + + loop.tick(); // opens the idle window + now.addAndGet(TimeUnit.SECONDS.toNanos(400)); + loop.tick(); // past the quiet period, pending reply -> INJECT + + assertEquals(List.of("term_lead_after_roll"), herdr.promptTargets(), + "the heartbeat must read and nudge the lead's CURRENT terminal, not the one learned " + + "before the roll"); + } + /** * Fake herdr client for the four tests above: always reports {@code lead} as IDLE, records the * {@code text} of every {@code agent.prompt} call, and can be told to throw on the very next @@ -657,6 +693,7 @@ class LeadHeartbeatLoopTest { private static final ObjectMapper MAPPER = new ObjectMapper(); private final String lead; private final List sentTexts = new ArrayList<>(); + private final List promptTargets = new ArrayList<>(); private boolean throwOnNextSend = false; FailableHerdrClient(String lead) { @@ -671,6 +708,11 @@ class LeadHeartbeatLoopTest { return List.copyOf(sentTexts); } + /** Every terminal an {@code agent.prompt} call named, in call order. */ + List promptTargets() { + return List.copyOf(promptTargets); + } + @Override @SuppressWarnings("unchecked") public JsonNode call(String method, Object params) { @@ -681,12 +723,13 @@ class LeadHeartbeatLoopTest { .put("agent_status", "idle")); } if ("agent.prompt".equals(method)) { + Map p = params instanceof Map ? (Map) params : Map.of(); if (throwOnNextSend) { throwOnNextSend = false; throw new RuntimeException("simulated transient herdr send failure"); } - Map p = params instanceof Map ? (Map) params : Map.of(); sentTexts.add(String.valueOf(p.get("text"))); + promptTargets.add(String.valueOf(p.get("target"))); } return MAPPER.createObjectNode(); } 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 e07b414..5aeaa3b 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/ReplyPushLoopTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/ReplyPushLoopTest.java @@ -304,6 +304,40 @@ class ReplyPushLoopTest { + "still live"); } + // --- fleetd #737 unit 3: the fallback nudgeTargetFor returns must be probed too -------------- + + /** + * {@code resolveLiveLead} forgets a dead per-target binding and asks {@code PrimaryRegistry} + * again for a fallback. That fallback can be dead too — here PRIMARY, the pinned singleton, is + * affirmatively gone alongside DEAD_LEAD. A correct {@code resolveLiveLead} probes it exactly + * like the first lead and gives up for this tick rather than trust it unchecked. + * + *

This is checked through {@code onReplyQueued}/{@code isActive} rather than a send count: + * {@link ReplyPushLoop#onReplyQueued} only registers pending work and starts a schedule once + * {@code resolveLiveLead} returns a present value — a dead fallback that was trusted unprobed + * would already make this true, synchronously, with no tick or send needed to observe it. Pairs + * with {@link #aStaleLeadBindingFallsBackToTheLiveLeadInsteadOfNudgingADeadTerminal} as the + * positive control: same stale-DEAD_LEAD setup, but there the fallback (PRIMARY) is live and the + * nudge does fire — proving this test's "nothing happens" result comes from the fallback being + * dead, not from the assertion being unable to observe a nudge at all. + */ + @Test + void aDoublyDeadFallbackIsNeverTrustedAndStartsNoSchedule() { + registry.recordDelegation(WORKER, DEAD_LEAD); + + var rec = new AllDeadHerdrClient(Set.of(DEAD_LEAD, PRIMARY)); + agents = new AgentControl(rec); + inbox.publish(WORKER, "m1", "hello"); + + var loop = loop(1, 50); + loop.onReplyQueued(WORKER); + assertFalse(loop.isActive(), + "both the per-target binding and the fallback are dead, so resolveLiveLead must " + + "return empty and onReplyQueued must never register pending work or start " + + "a schedule — an unprobed fallback would start one here"); + assertEquals(0, rec.promptTargets().size(), "nobody live was found, so nothing was ever sent"); + } + // --- nudge format -------------------------------------------------------------------------- @Test @@ -1248,6 +1282,52 @@ class ReplyPushLoopTest { } } + /** + * Fake herdr client for fleetd #737 unit 3: every terminal named in {@code deadTargets} reports + * {@code agent_not_found} from {@code agent.get} — unlike {@link DeadLeadHerdrClient}, which can + * only make one terminal dead, this can make a per-target binding AND its fallback dead in the + * same test. {@code agent.prompt} is recorded unconditionally (no liveness check of its own), + * so a test can tell "resolveLiveLead probed and correctly found nobody live" (no prompt call) + * apart from "resolveLiveLead trusted a dead fallback and sent into it anyway" (a prompt call to + * a terminal this fake has already declared gone). + */ + private static final class AllDeadHerdrClient implements HerdrClient { + private final Set deadTargets; + private final List promptTargets = Collections.synchronizedList(new ArrayList<>()); + + AllDeadHerdrClient(Set deadTargets) { + this.deadTargets = deadTargets; + } + + @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 (deadTargets.contains(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"))); + } + return MAPPER.createObjectNode(); + } + + List promptTargets() { + return List.copyOf(promptTargets); + } + + @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