From f4176ae455f97ca19dd2234a72e91353708a76bd Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Mon, 5 Oct 2026 05:29:06 +0200 Subject: [PATCH 1/2] fleetd #737 unit 3: probe the fallback terminal and resolve it by lead name ReplyPushLoop.resolveLiveLead dropped a dead per-target delegation and retried PrimaryRegistry.nudgeTargetFor, but returned that fallback without checking isLive. The fallback is now probed the same way the first lead is, and the method returns empty rather than trust a dead terminal. PrimaryRegistry records a delegating lead's name alongside its learned terminal and resolves the name back to its current terminal at nudge time, through a name-to-terminal lookup backed by the live lead-tab scan. A name with no current match falls back to the terminal that was actually learned, so an unnamed primary, an off-host lead, or a non-herdr lead keeps working exactly as before. LeadHeartbeatLoop now reads the resolved current terminal instead of the raw learned one. --- .../src/main/java/dev/ltms/fleet/Fleetd.java | 26 +++++ .../java/dev/ltms/fleet/FleetdAssembly.java | 3 +- .../java/dev/ltms/fleet/mcp/FleetMcp.java | 10 +- .../dev/ltms/fleet/mcp/PrimaryRegistry.java | 109 ++++++++++++++++-- .../dev/ltms/fleet/msg/LeadHeartbeatLoop.java | 6 +- .../dev/ltms/fleet/msg/ReplyPushLoop.java | 16 ++- .../ltms/fleet/mcp/PrimaryRegistryTest.java | 97 ++++++++++++++++ .../ltms/fleet/msg/LeadHeartbeatLoopTest.java | 45 +++++++- .../dev/ltms/fleet/msg/ReplyPushLoopTest.java | 80 +++++++++++++ 9 files changed, 377 insertions(+), 15 deletions(-) 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/mcp/FleetMcp.java b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java index eae1b71..3e400e5 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) @@ -730,7 +736,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: *

    *
  • the registry is pinned (config override), *
  • {@code terminalId} is {@code null} or blank (non-herdr caller). *
*/ 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 From 803c91ea6c47b646770c75a3da555f247678c1c4 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Mon, 5 Oct 2026 05:37:32 +0200 Subject: [PATCH 2/2] fleetd #737 unit 3: name the resolving accessor in LeadRollover's javadoc The heartbeat loop now resolves its nudge target through PrimaryRegistry#currentPrimaryTerminal(), so the sentence describing what a background loop with no caller uses named the raw accessor instead. --- fleetd/src/main/java/dev/ltms/fleet/lead/LeadRollover.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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