fleetd #609 review: the context latch must mean the notice reached the pane
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 49s
CI / build (pull_request) Failing after 2m6s

Fixes the PR #610 review blocker: LeadHeartbeatLoop committed contextNotified
before injectNudge attempted the send, so a transient herdr failure marked the
lead as told when nothing reached its pane, and contextNotice() carried no
latch at all, so a pending-driven INJECT re-appended the notice on every tick
while the context stayed HIGH.

- injectNudge now reports whether agents.send succeeded and persists
  contextNotified only when a notice was actually included in the text and
  the send did not throw. The latch is split out of applyDecision (kept for
  idleSinceNanos/quietCount, applied unconditionally as before) so it is
  written on the success path only, once per branch in tick().
- contextNotice gained an overloaded 3-arg form gated on the latch as it
  stood before the tick's decision; the existing 2-arg form delegates to it
  with alreadyNotified=false, so all pre-existing callers/tests are unchanged.
- tick() is now package-private (mirrors ReplyPushLoop#tick(String)) so tests
  can drive the real send path with a fake AgentControl instead of only the
  pure decide() function.
- Added tests I-L covering: a failed send does not consume the notice and
  retries; a successful send does; the text is gated when the latch is
  already set; and the notice appears exactly once across three differently
  driven INJECTs.

Both required mutations verified red and reverted:
1. Setting the latch from the Decision regardless of send outcome -> test I
   (iAFailedSendDoesNotConsumeTheNotice) fails.
2. Dropping the latch argument at the contextNotice call site -> tests K
   (kAPendingDrivenInjectWithTheLatchAlreadySetSendsNoNotice) and L
   (lTheNoticeAppearsExactlyOnceAcrossThreeDifferentlyDrivenInjects) fail.
This commit is contained in:
Dai Ha
2026-09-20 17:34:05 +07:00
parent d7390ccd37
commit 89cb8ff79b
2 changed files with 271 additions and 16 deletions
@@ -276,8 +276,10 @@ public final class LeadHeartbeatLoop {
scheduler.schedule(this::tick, backoffMs, TimeUnit.MILLISECONDS); scheduler.schedule(this::tick, backoffMs, TimeUnit.MILLISECONDS);
} }
/** One loop tick, every {@link #backoffMs} — the thin scheduler around {@link #decide}. */ /** One loop tick, every {@link #backoffMs} — the thin scheduler around {@link #decide}. Package-private
private void tick() { * (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.primaryTerminal().isPresent();
FleetState fleet = snapshot(inbox, roster); FleetState fleet = snapshot(inbox, roster);
AgentStatus status = AgentStatus.UNKNOWN; AgentStatus status = AgentStatus.UNKNOWN;
@@ -303,29 +305,59 @@ public final class LeadHeartbeatLoop {
reading.state(), contextNotified); reading.state(), contextNotified);
applyDecision(d); applyDecision(d);
switch (d.action()) { switch (d.action()) {
case INJECT -> injectNudge(fleet, reading); case INJECT -> injectNudge(d, fleet, reading);
case QUIET_DONE -> countNudge("exhausted"); case QUIET_DONE -> {
case WAIT_IDLE, LEAD_BUSY, STAND_DOWN -> { /* nothing to inject, nothing to count */ } countNudge("exhausted");
contextNotified = d.contextNotified();
}
case WAIT_IDLE, LEAD_BUSY, STAND_DOWN -> contextNotified = d.contextNotified();
} }
scheduleNext(); 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.
*
* <p>fleetd #609 review: the context latch ({@link #contextNotified}) is deliberately <em>not</em>
* 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) { private void applyDecision(Decision d) {
idleSinceNanos = d.idleSinceNanos() == null ? NOT_IDLE : d.idleSinceNanos(); idleSinceNanos = d.idleSinceNanos() == null ? NOT_IDLE : d.idleSinceNanos();
quietCount = d.quietCount(); 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(); var lead = primaryRegistry.primaryTerminal();
if (lead.isEmpty()) { boolean sent = lead.isPresent() && trySend(lead.get(), fleet.nudgeText() + notice, notice);
return; // the lead disappeared between the decision and the injection // 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
String leadTerminal = lead.get(); // notice was attempted (disabled, not HIGH, or already latched), nothing was promised to the lead
String notice = contextNotice(contextHighNudge, reading); // this tick, so apply the decision's own carried-forward value unconditionally — that is how the
String text = fleet.nudgeText() + notice; // 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 { try {
agents.send(leadTerminal, text); agents.send(leadTerminal, text);
log.debug("idle-heartbeat: nudge sent to lead {} (quiet nudges so far in this stretch: {})", 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 // 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. // it is visible in /metrics — one count per nudge either way, never two.
countNudge(notice.isEmpty() ? "sent" : "sent_context"); countNudge(notice.isEmpty() ? "sent" : "sent_context");
return true;
} catch (RuntimeException e) { } catch (RuntimeException e) {
log.warn("idle-heartbeat: failed to nudge lead {}: {}", leadTerminal, e.toString()); log.warn("idle-heartbeat: failed to nudge lead {}: {}", leadTerminal, e.toString());
countNudge("failed"); 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} * FleetState#nudgeText()}), or {@code ""} when disabled or the state is not {@code HIGH}
*/ */
static String contextNotice(boolean enabled, LeadContextGauge.Reading reading) { 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 <em>before</em> 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 ""; return "";
} }
StringBuilder sb = new StringBuilder(" Your own context is nearly full"); StringBuilder sb = new StringBuilder(" Your own context is nearly full");
@@ -1,6 +1,10 @@
package dev.ltms.fleet.msg; 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.AgentStatus;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.lead.LeadContextGauge; import dev.ltms.fleet.lead.LeadContextGauge;
import dev.ltms.fleet.mcp.PrimaryRegistry; import dev.ltms.fleet.mcp.PrimaryRegistry;
import dev.ltms.fleet.peer.MemberRole; 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.BeforeEach;
import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Test;
import java.util.ArrayList;
import java.util.List; import java.util.List;
import java.util.Map;
import java.util.concurrent.Executors; import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import static org.junit.jupiter.api.Assertions.*; import static org.junit.jupiter.api.Assertions.*;
@@ -411,4 +418,204 @@ class LeadHeartbeatLoopTest {
assertFalse(notice.toLowerCase().contains("null"), notice); assertFalse(notice.toLowerCase().contains("null"), notice);
assertTrue(notice.contains("2 compactions"), 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<MemberSession>[] 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<MemberSession>[] 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<MemberSession>[] 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<MemberSession>[] 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<MemberSession>[] 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<String> sentTexts = new ArrayList<>();
private boolean throwOnNextSend = false;
FailableHerdrClient(String lead) {
this.lead = lead;
}
void throwOnNextSend() {
throwOnNextSend = true;
}
List<String> 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<String, Object> p = params instanceof Map ? (Map<String, Object>) params : Map.of();
sentTexts.add(String.valueOf(p.get("text")));
}
return MAPPER.createObjectNode();
}
@Override
public void close() {
}
}
} }