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 42df198..2636d69 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/LeadHeartbeatLoop.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/LeadHeartbeatLoop.java @@ -276,8 +276,10 @@ public final class LeadHeartbeatLoop { scheduler.schedule(this::tick, backoffMs, TimeUnit.MILLISECONDS); } - /** One loop tick, every {@link #backoffMs} — the thin scheduler around {@link #decide}. */ - private void tick() { + /** One loop tick, every {@link #backoffMs} — the thin scheduler around {@link #decide}. Package-private + * (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(); FleetState fleet = snapshot(inbox, roster); AgentStatus status = AgentStatus.UNKNOWN; @@ -303,29 +305,59 @@ public final class LeadHeartbeatLoop { reading.state(), contextNotified); applyDecision(d); switch (d.action()) { - case INJECT -> injectNudge(fleet, reading); - case QUIET_DONE -> countNudge("exhausted"); - case WAIT_IDLE, LEAD_BUSY, STAND_DOWN -> { /* nothing to inject, nothing to count */ } + case INJECT -> injectNudge(d, fleet, reading); + case QUIET_DONE -> { + countNudge("exhausted"); + contextNotified = d.contextNotified(); + } + case WAIT_IDLE, LEAD_BUSY, STAND_DOWN -> contextNotified = d.contextNotified(); } scheduleNext(); } - /** Persist the state a decision returned, so the next tick starts from it. */ + /** + * Persist the idle/quiet state a decision returned, so the next tick starts from it. + * + *

fleetd #609 review: the context latch ({@link #contextNotified}) is deliberately not + * set here any more. Setting it from the decision unconditionally — before {@link #injectNudge} even + * tries to send — is exactly the review's blocker: a decision to notify is not the same fact as "the + * notice reached the pane". Every branch of {@link #tick} now assigns {@link #contextNotified} itself, + * once it knows whether a send happened and whether it carried the notice (see {@link #injectNudge}). + */ private void applyDecision(Decision d) { idleSinceNanos = d.idleSinceNanos() == null ? NOT_IDLE : d.idleSinceNanos(); quietCount = d.quietCount(); - contextNotified = d.contextNotified(); } - /** Send the nudge to the known lead, with the fleetd #609 context notice appended when it applies. */ - private void injectNudge(FleetState fleet, LeadContextGauge.Reading reading) { + /** + * Send the nudge to the known lead, with the fleetd #609 context notice appended when it applies, and + * persist the context latch based on what actually happened this tick — not merely what {@code d} + * chose to attempt. + */ + private void injectNudge(Decision d, FleetState fleet, LeadContextGauge.Reading reading) { + // fleetd #609 review: build the notice from the latch as it stood BEFORE this tick's decision — + // d.contextNotified() is the value to persist once delivery is confirmed, not the value the text + // 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); var lead = primaryRegistry.primaryTerminal(); - if (lead.isEmpty()) { - return; // the lead disappeared between the decision and the injection - } - String leadTerminal = lead.get(); - String notice = contextNotice(contextHighNudge, reading); - String text = fleet.nudgeText() + notice; + 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 + // notice was attempted (disabled, not HIGH, or already latched), nothing was promised to the lead + // this tick, so apply the decision's own carried-forward value unconditionally — that is how the + // OK-only re-arm rule and STAND_DOWN's "don't burn the notice" rule keep working through this path + // too. A lead that disappeared between the decision and the send (lead.isEmpty()) is treated the + // same as a failed send: nothing reached the pane, so the latch must not be set. + contextNotified = notice.isEmpty() ? d.contextNotified() : sent; + } + + /** + * Attempt one herdr send and count its outcome. Returns whether {@code agents.send} returned without + * throwing — the caller ({@link #injectNudge}) needs this to decide whether the fleetd #609 context + * latch may be persisted as set. + */ + private boolean trySend(String leadTerminal, String text, String notice) { try { agents.send(leadTerminal, text); log.debug("idle-heartbeat: nudge sent to lead {} (quiet nudges so far in this stretch: {})", @@ -333,9 +365,11 @@ public final class LeadHeartbeatLoop { // fleetd #609: a nudge that carries the context notice is counted under its own outcome so // it is visible in /metrics — one count per nudge either way, never two. countNudge(notice.isEmpty() ? "sent" : "sent_context"); + return true; } catch (RuntimeException e) { log.warn("idle-heartbeat: failed to nudge lead {}: {}", leadTerminal, e.toString()); countNudge("failed"); + return false; } } @@ -351,7 +385,21 @@ public final class LeadHeartbeatLoop { * FleetState#nudgeText()}), or {@code ""} when disabled or the state is not {@code HIGH} */ static String contextNotice(boolean enabled, LeadContextGauge.Reading reading) { - if (!enabled || reading.state() != LeadContextGauge.State.HIGH) { + return contextNotice(enabled, reading, false); + } + + /** + * fleetd #609 review: as {@link #contextNotice(boolean, LeadContextGauge.Reading)}, but also gated on + * {@code alreadyNotified} — the context latch as it stood before the current tick's decision. + * Without this gate, every pending-driven {@code INJECT} that lands while the context stays {@code + * HIGH} would re-append the full notice on top of an already-latched stretch, making the notice's own + * closing sentence ("You will not be told again until your context reads ok.") false. {@link + * #injectNudge} is the only caller that passes a non-default {@code alreadyNotified}. + * + * @param alreadyNotified whether the lead has already been told about the current HIGH stretch + */ + static String contextNotice(boolean enabled, LeadContextGauge.Reading reading, boolean alreadyNotified) { + if (!enabled || alreadyNotified || reading.state() != LeadContextGauge.State.HIGH) { return ""; } StringBuilder sb = new StringBuilder(" Your own context is nearly full"); 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 b03603d..d327861 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/LeadHeartbeatLoopTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/LeadHeartbeatLoopTest.java @@ -1,6 +1,10 @@ package dev.ltms.fleet.msg; +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import dev.ltms.fleet.herdr.AgentControl; import dev.ltms.fleet.herdr.AgentStatus; +import dev.ltms.fleet.herdr.HerdrClient; import dev.ltms.fleet.lead.LeadContextGauge; import dev.ltms.fleet.mcp.PrimaryRegistry; import dev.ltms.fleet.peer.MemberRole; @@ -9,10 +13,13 @@ import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import java.util.ArrayList; import java.util.List; +import java.util.Map; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; import static org.junit.jupiter.api.Assertions.*; @@ -411,4 +418,204 @@ class LeadHeartbeatLoopTest { assertFalse(notice.toLowerCase().contains("null"), notice); assertTrue(notice.contains("2 compactions"), notice); } + + // ── fleetd #609 review: the latch must mean "the notice reached the pane" ──────────────────── + // + // These four drive LeadHeartbeatLoop.tick() directly (package-private, same reasoning as + // ReplyPushLoop#tick(String) being directly testable) against a real AgentControl wrapping a + // FailableHerdrClient, so the send path (agents.send -> herdr -> possible throw) is exercised + // for real rather than assumed from decide()'s Decision alone. + + private static final String LEAD = "term_lead"; + private static final String WORKER = "term_w1"; + + /** A HIGH reading with a fixed token/compaction count, for the four tests below. */ + private static LeadContextGauge.Reading highReading() { + return new LeadContextGauge.Reading(LeadContextGauge.State.HIGH, 260_771L, 1); + } + + /** + * Builds a real {@link LeadHeartbeatLoop} wired to {@code herdr} via a real {@link AgentControl}, + * a mutable fake clock, and a mutable roster so a test can change fleet state between ticks. The + * lead is always reported IDLE by {@code herdr}, so every tick's outcome is governed only by the + * idle-window/quiet-cap/context gates under test. + */ + private static LeadHeartbeatLoop tickableLoop(FailableHerdrClient herdr, AtomicLong now, + List[] rosterBox, InMemoryReplyInbox inbox, + int quietNudgeCap, ScheduledExecutorService scheduler) { + AgentControl agents = new AgentControl(herdr); + PrimaryRegistry registry = new PrimaryRegistry(LEAD); + ReplyPushLoop pushLoop = new ReplyPushLoop(registry, agents, inbox, scheduler, 5, 100_000); + return new LeadHeartbeatLoop(registry, agents, inbox, () -> rosterBox[0], pushLoop, scheduler, + now::get, IDLE_AFTER_NANOS, 100_000L, quietNudgeCap, null, + new LeadHeartbeatLoop.LeadContextSource(t -> highReading()), true); + } + + @Test + void iAFailedSendDoesNotConsumeTheNotice() { + var herdr = new FailableHerdrClient(LEAD); + var now = new AtomicLong(NOW); + @SuppressWarnings("unchecked") + List[] rosterBox = new List[]{List.of()}; // quiet: nothing pending + InMemoryReplyInbox inbox = new InMemoryReplyInbox(); + LeadHeartbeatLoop loop = tickableLoop(herdr, now, rosterBox, inbox, 0, scheduler); + + loop.tick(); // first injectable tick: only opens the idle window (WAIT_IDLE) + now.addAndGet(TimeUnit.SECONDS.toNanos(400)); // now clearly past the quiet period + + herdr.throwOnNextSend(); + loop.tick(); // quiet fleet, quiet cap exhausted (0), context HIGH, latch clear -> INJECT, send throws + + assertEquals(0, herdr.sentTexts().size(), "the failed send must not have recorded any text"); + + loop.tick(); // same inputs — the latch must still be clear, so this must INJECT and send again + assertEquals(1, herdr.sentTexts().size(), + "a retried tick with the latch still clear must attempt the send again"); + assertTrue(herdr.sentTexts().get(0).contains("Your own context is nearly full"), + "the retried, successful send must carry the notice: " + herdr.sentTexts().get(0)); + } + + @Test + void jASuccessfulSendDoesConsumeIt() { + var herdr = new FailableHerdrClient(LEAD); + var now = new AtomicLong(NOW); + @SuppressWarnings("unchecked") + List[] rosterBox = new List[]{List.of()}; + InMemoryReplyInbox inbox = new InMemoryReplyInbox(); + LeadHeartbeatLoop loop = tickableLoop(herdr, now, rosterBox, inbox, 0, scheduler); + + loop.tick(); // opens the idle window + now.addAndGet(TimeUnit.SECONDS.toNanos(400)); + + loop.tick(); // INJECT, send succeeds -> latch set + assertEquals(1, herdr.sentTexts().size()); + assertTrue(herdr.sentTexts().get(0).contains("Your own context is nearly full"), herdr.sentTexts().get(0)); + + loop.tick(); // same inputs — the lead was already told this stretch + assertEquals(1, herdr.sentTexts().size(), + "the next tick with the same inputs must not send a second notice"); + } + + @Test + void kAPendingDrivenInjectWithTheLatchAlreadySetSendsNoNotice() { + var herdr = new FailableHerdrClient(LEAD); + var now = new AtomicLong(NOW); + @SuppressWarnings("unchecked") + List[] rosterBox = new List[]{List.of()}; + InMemoryReplyInbox inbox = new InMemoryReplyInbox(); + LeadHeartbeatLoop loop = tickableLoop(herdr, now, rosterBox, inbox, 0, scheduler); + + loop.tick(); + now.addAndGet(TimeUnit.SECONDS.toNanos(400)); + loop.tick(); // latches the notice (quiet, HIGH, cap exhausted -> the forced context INJECT) + assertEquals(1, herdr.sentTexts().size()); + assertTrue(herdr.sentTexts().get(0).contains("Your own context is nearly full")); + + // Now make the fleet have real pending state, so the NEXT INJECT is pending-driven, not the + // forced-context route — with the latch already set from the tick above. + inbox.own(WORKER); + inbox.publish(WORKER, "m1", "hello"); + rosterBox[0] = List.of(new MemberSession("p1", WORKER, "prof", MemberRole.DEV, "/cwd", null, + 0, 0, 0, MemberSession.State.READY, null, null)); + + loop.tick(); + + assertEquals(2, herdr.sentTexts().size(), "the pending-driven tick must still send a nudge"); + assertFalse(herdr.sentTexts().get(1).contains("Your own context is nearly full"), + "a pending-driven INJECT while the latch is already set must carry no context notice: " + + herdr.sentTexts().get(1)); + } + + @Test + void lTheNoticeAppearsExactlyOnceAcrossThreeDifferentlyDrivenInjects() { + var herdr = new FailableHerdrClient(LEAD); + var now = new AtomicLong(NOW); + @SuppressWarnings("unchecked") + List[] rosterBox = new List[]{List.of()}; + InMemoryReplyInbox inbox = new InMemoryReplyInbox(); + // quietNudgeCap=1 so a still-not-exhausted quiet nudge is available as the "forced" route below, + // distinct from both the pending-driven route and the exhausted-cap forced-context route. + LeadHeartbeatLoop loop = tickableLoop(herdr, now, rosterBox, inbox, 1, scheduler); + + loop.tick(); // opens the idle window + now.addAndGet(TimeUnit.SECONDS.toNanos(400)); + + // 1) pending-driven INJECT: real fleet state present. Sets the latch and carries the notice. + inbox.own(WORKER); + inbox.publish(WORKER, "m1", "hello"); + rosterBox[0] = List.of(new MemberSession("p1", WORKER, "prof", MemberRole.DEV, "/cwd", null, + 0, 0, 0, MemberSession.State.READY, null, null)); + loop.tick(); + assertEquals(1, herdr.sentTexts().size()); + assertTrue(herdr.sentTexts().get(0).contains("Your own context is nearly full"), herdr.sentTexts().get(0)); + + // 2) "forced" INJECT: nothing pending, but the quiet cap (1) is not yet exhausted, so decide() + // nudges anyway. The latch is already set, so no notice. + inbox.ack(WORKER, "m1"); + rosterBox[0] = List.of(); + loop.tick(); + assertEquals(2, herdr.sentTexts().size(), "the quiet-cap-not-yet-exhausted nudge must still fire"); + assertFalse(herdr.sentTexts().get(1).contains("Your own context is nearly full"), herdr.sentTexts().get(1)); + + // 3) pending-driven INJECT again. Still latched, still no notice. + inbox.publish(WORKER, "m2", "hello again"); + rosterBox[0] = List.of(new MemberSession("p1", WORKER, "prof", MemberRole.DEV, "/cwd", null, + 0, 0, 0, MemberSession.State.READY, null, null)); + loop.tick(); + assertEquals(3, herdr.sentTexts().size()); + assertFalse(herdr.sentTexts().get(2).contains("Your own context is nearly full"), herdr.sentTexts().get(2)); + + long noticeCount = herdr.sentTexts().stream() + .filter(t -> t.contains("Your own context is nearly full")).count(); + assertEquals(1, noticeCount, + "the notice text must appear exactly once across all three sends: " + herdr.sentTexts()); + } + + /** + * 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 + * {@code agent.prompt} call — standing in for one transient herdr send failure. + */ + private static final class FailableHerdrClient implements HerdrClient { + private static final ObjectMapper MAPPER = new ObjectMapper(); + private final String lead; + private final List sentTexts = new ArrayList<>(); + private boolean throwOnNextSend = false; + + FailableHerdrClient(String lead) { + this.lead = lead; + } + + void throwOnNextSend() { + throwOnNextSend = true; + } + + List sentTexts() { + return List.copyOf(sentTexts); + } + + @Override + @SuppressWarnings("unchecked") + public JsonNode call(String method, Object params) { + if ("agent.get".equals(method)) { + return MAPPER.createObjectNode() + .set("agent", MAPPER.createObjectNode() + .put("terminal_id", lead) + .put("agent_status", "idle")); + } + if ("agent.prompt".equals(method)) { + 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"))); + } + return MAPPER.createObjectNode(); + } + + @Override + public void close() { + } + } }