diff --git a/bridged/bridged.example.yaml b/bridged/bridged.example.yaml index a7e0927..01d5080 100644 --- a/bridged/bridged.example.yaml +++ b/bridged/bridged.example.yaml @@ -84,6 +84,26 @@ bind: # tabPrefix: "lead:" # `lead: opus-5.0` ⇒ a lead named opus-5.0 (case-insensitive; default "lead:") # intervalSeconds: 10 # rescan cadence, and the worst case before a new tab is recognised +# CB-551: IDLE-LEAD HEARTBEAT — nudge the single lead back to work when it has been continuously +# idle (no open bridge_send driving it) past the quiet period. The fleet is one lead + architects + +# workers, so a lead that stalls is a single point of failure; the ReplyPushLoop only nudges when a +# reply lands, and this timer catches the gap where nothing lands and the lead just sits idle. +# +# Opt-in on purpose — it SPENDS the operator's subscription on its own initiative (each nudge starts +# a lead turn nobody asked for), so upgrading the daemon must never switch it on for you. Absent +# block = feature off, exactly as before. +# +# Three knobs, each with a default that errs on the side of not burning context: +# idleAfterSeconds: 300 # how long the lead must stay idle before the FIRST nudge (default 300 — +# # absorbs normal post-turn pauses; re-prompting every pause burns context) +# backoffMs: 60000 # re-check cadence / spacing between nudges past the quiet period (default 60000) +# quietNudgeCap: 3 # cap on consecutive nudges that find NOTHING pending, then it stops +# # until real state appears (default 3 — never nag an empty fleet forever) +# leadHeartbeat: +# idleAfterSeconds: 300 +# backoffMs: 60000 +# quietNudgeCap: 3 + # herdr Unix socket. Omit to use the client default # (${HERDR_SOCKET_PATH:-~/.config/herdr/herdr.sock}). herdrSocket: ~/.config/herdr/herdr.sock diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java index 34e828a..9e6b474 100644 --- a/bridged/src/main/java/dev/ltms/bridged/Bridged.java +++ b/bridged/src/main/java/dev/ltms/bridged/Bridged.java @@ -28,6 +28,7 @@ import dev.ltms.bridged.msg.InMemoryReplyInbox; import dev.ltms.bridged.msg.MessageService; import dev.ltms.bridged.msg.Rendezvous; import dev.ltms.bridged.msg.ReplyInbox; +import dev.ltms.bridged.msg.LeadHeartbeatLoop; import dev.ltms.bridged.msg.ReplyPushLoop; import dev.ltms.bridged.rest.BridgedApp; import dev.ltms.bridged.session.GitWorktrees; @@ -296,6 +297,23 @@ public final class Bridged { Metrics metrics = BridgedMetrics.create(sessions, replyInbox); var pushLoop = new ReplyPushLoop(primaryRegistry, agents, replyInbox, pushScheduler, maxReminders, backoffMs, metrics); + // CB-551: idle-lead heartbeat. Opt-in; absent `leadHeartbeat:` this is never constructed, so + // an upgraded daemon cannot silently start spending subscription on nudging an idle lead. + // It has its own single-thread scheduler and holds its own scheduler shutdown via close(). + final LeadHeartbeatLoop heartbeat; + var heartbeatScheduler = Executors.newSingleThreadScheduledExecutor(r -> + Thread.ofVirtual().name("bridge-heartbeat-").unstarted(r)); + if (cfg.leadHeartbeat() != null) { + var hb = cfg.leadHeartbeat(); + heartbeat = new LeadHeartbeatLoop(primaryRegistry, agents, replyInbox, sessions::roster, + pushLoop, heartbeatScheduler, System::nanoTime, + TimeUnit.SECONDS.toNanos(hb.idleAfterSeconds()), hb.backoffMs(), hb.quietNudgeCap(), + metrics); + heartbeat.start(); + } else { + heartbeat = null; + heartbeatScheduler.shutdownNow(); + } MessageService messages = new MessageService(agents, injector, rendezvous, replyInbox, pushLoop, metrics); @@ -346,6 +364,7 @@ public final class Bridged { poller.stop(); messages.close(); pushLoop.close(); + if (heartbeat != null) heartbeat.close(); // CB-551: stop the idle-lead heartbeat scheduler mcp.close(); if (reaper != null) reaper.stop(); // Release the broker connection last among message resources (no-op for the in-memory inbox). diff --git a/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java b/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java index 65e42fd..c538320 100644 --- a/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java +++ b/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java @@ -53,6 +53,8 @@ import java.util.Set; * to a slot. Nothing here spawns a slot. * @param leadScan opt-in discovery of leads by tab label (CB-531); {@code null} ⇒ no scanning, * and only {@code leaders:}/{@code primary:} name a lead + * @param leadHeartbeat opt-in idle-lead heartbeat (CB-551); {@code null} ⇒ off, and an upgraded + * daemon never nudges an idle lead on its own initiative * @param placement how to choose a worker profile for an unqualified spawn: * {@code fixed} (default), {@code round-robin}, or {@code weighted} * @param auth API authentication mode ({@code null} → {@code loopback-trust}, the @@ -75,6 +77,7 @@ public record BridgedConfig( Map leaders, Map architects, LeadScan leadScan, + LeadHeartbeat leadHeartbeat, String placement, Auth auth) { @@ -441,6 +444,40 @@ public record BridgedConfig( } } + /** + * Opt-in idle-lead heartbeat (CB-551): when the single lead has been continuously idle past + * {@code idleAfterSeconds} with no open {@code bridge_send} driving it, nudge it back to work. + * + *

Deliberately opt-in ({@code null} ⇒ off, exactly like {@code leadScan:}). The heartbeat + * spends the operator's model subscription on its own initiative — it prompts the lead to start + * a turn nobody asked for — which is precisely the kind of behaviour a daemon upgrade must never + * acquire silently. Absent the block, nothing is constructed and nothing fires. + * + * @param idleAfterSeconds how long the lead must be continuously injectable before the + * first nudge (default 300). Why: a lead that just finished a turn sits + * momentarily idle; re-prompting instantly would turn every natural pause + * into a new turn and burn the lead's context for nothing. Context + * exhaustion is what actually kills a long-running lead, so every nudge + * must earn its cost — 5 minutes absorbs normal pauses without stalling. + * @param backoffMs re-check cadence once the quiet period has passed (default 60000). The + * loop samples on this interval; 60s is gentle enough not to burn context + * while prompt enough to unstick a stalled lead promptly. Also the spacing + * between nudges to a lead that is idle past the quiet period. + * @param quietNudgeCap cap on consecutive nudges that find no pending fleet state + * (default 3), before the loop stops until real state appears again. Why: + * each such nudge costs the lead a turn just to read "nothing pending"; + * three is enough to tell it it may stand down without nagging forever, and + * it is the bound that stops an idle fleet from being a subscription burner. + */ + @JsonIgnoreProperties(ignoreUnknown = true) + public record LeadHeartbeat(Integer idleAfterSeconds, Long backoffMs, Integer quietNudgeCap) { + public LeadHeartbeat { + idleAfterSeconds = (idleAfterSeconds == null || idleAfterSeconds <= 0) ? 300 : idleAfterSeconds; + backoffMs = (backoffMs == null || backoffMs <= 0) ? 60_000L : backoffMs; + quietNudgeCap = (quietNudgeCap == null || quietNudgeCap < 0) ? 3 : quietNudgeCap; + } + } + /** * The terminal → lead-name map that {@link dev.ltms.bridged.auth.CallerResolver} resolves * against, merging the {@code leaders:} registry with the legacy singular {@code primary:} pin. @@ -571,7 +608,7 @@ public record BridgedConfig( private static final Set KNOWN_TOP_LEVEL_KEYS = Set.of( "bind", "herdrSocket", "worker", "workers", "defaultWorker", "guard", "worktreeRoot", "lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs", "broker", "primary", "leaders", - "architects", "leadScan", "placement", "auth"); + "architects", "leadScan", "leadHeartbeat", "placement", "auth"); /** Load and validate config from {@code path}. */ public static BridgedConfig load(Path path) { @@ -743,9 +780,11 @@ public record BridgedConfig( // leadScan is left as-is: null is "off", and LeadScan's own compact constructor defaults the // fields of a block that IS present. Defaulting it here would switch the feature on for // every config that never mentioned it. + // leadHeartbeat is left as-is for the same reason (CB-551): null is "off", and LeadHeartbeat's + // own compact constructor defaults the fields of a block that IS present. // architects is left as-is: null is "none configured", and Architect's fields have no // defaults to fill. Defaulting it here would change nothing, so leave the call natural. - return new BridgedConfig(b, herdrSocket, worker, workers, defaultWorker, g, worktreeRoot, l, timeout, pollMs, broker, primary, leaders, architects, leadScan, placementOrDefault, a); + return new BridgedConfig(b, herdrSocket, worker, workers, defaultWorker, g, worktreeRoot, l, timeout, pollMs, broker, primary, leaders, architects, leadScan, leadHeartbeat, placementOrDefault, a); } /** diff --git a/bridged/src/main/java/dev/ltms/bridged/metrics/BridgedMetrics.java b/bridged/src/main/java/dev/ltms/bridged/metrics/BridgedMetrics.java index 7e53f35..3329879 100644 --- a/bridged/src/main/java/dev/ltms/bridged/metrics/BridgedMetrics.java +++ b/bridged/src/main/java/dev/ltms/bridged/metrics/BridgedMetrics.java @@ -26,6 +26,8 @@ public final class BridgedMetrics { public static final String REPLIES = "bridged_replies_total"; /** Counter: push-loop nudges to the primary, by outcome. */ public static final String PUSH_NUDGES = "bridged_push_nudges_total"; + /** Counter: idle-lead heartbeat nudges to the lead, by outcome (CB-551). */ + public static final String HEARTBEAT_NUDGES = "bridged_lead_heartbeat_nudges_total"; /** Counter: spawn attempts by peer kind and outcome. */ public static final String SPAWNS = "bridged_spawns_total"; /** Counter: herdr socket calls by method and outcome. */ @@ -55,6 +57,9 @@ public final class BridgedMetrics { "Worker replies by delivery path (rendezvous=resolved an open send, inbox=stranded and held)."); m.describe(PUSH_NUDGES, "counter", "CB-307 push-loop nudges to the primary (delivered|exhausted)."); + m.describe(HEARTBEAT_NUDGES, "counter", + "CB-551 idle-lead heartbeat nudges (delivered|failed|exhausted). Quiet-cap exhaustion " + + "means the lead idled with nothing pending and was told to stand down."); m.describe(SPAWNS, "counter", "Worker spawn attempts by peer kind and outcome (ready|timeout|guard_rejected)."); m.describe(HERDR_CALLS, "counter", diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/LeadHeartbeatLoop.java b/bridged/src/main/java/dev/ltms/bridged/msg/LeadHeartbeatLoop.java new file mode 100644 index 0000000..dca7d65 --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/msg/LeadHeartbeatLoop.java @@ -0,0 +1,351 @@ +package dev.ltms.bridged.msg; + +import dev.ltms.bridged.herdr.AgentControl; +import dev.ltms.bridged.herdr.AgentStatus; +import dev.ltms.bridged.mcp.PrimaryRegistry; +import dev.ltms.bridged.metrics.BridgedMetrics; +import dev.ltms.bridged.metrics.Metrics; +import dev.ltms.bridged.session.WorkerSession; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; +import java.util.function.LongSupplier; +import java.util.function.Supplier; + +/** + * CB-551: an opt-in heartbeat that nudges the single idle lead back to work once it has been + * continuously idle past a quiet period with no open {@code bridge_send} driving it. + * + *

Why this exists: the fleet is ONE lead + architects + workers, so an idle, stalled lead is a + * single point of failure for the fleet's progress. {@link ReplyPushLoop} nudges the lead only when + * a worker reply lands; this loop is the timer that catches the gap where nothing lands and the + * lead simply sits idle with nobody to prompt it onward. + * + *

Mechanism mirrors {@link ReplyPushLoop}: a pure, unit-testable decision function + * ({@link #decide(long, Long, int, AgentStatus, boolean, boolean, FleetState)}) plus a thin + * scheduler around it. The loop is entirely opt-in ({@code leadHeartbeat:} in config); with + * no such block it is never constructed, so upgrading the daemon cannot silently acquire a behaviour + * that spends the operator's model subscription on its own initiative (constraint 1). + * + *

Four invariants keep it from becoming a runaway subscription burner: + *

    + *
  1. Status-gated — a {@code WORKING} lead is making progress and is never touched; only an + * injectable (idle/done/blocked) lead is even considered (constraint 2).
  2. + *
  3. Debounced — a nudge only fires once the lead has been continuously idle past + * {@code idleAfterSeconds}, so a lead that just finished a turn is not re-prompted into every + * natural pause (constraint 3).
  4. + *
  5. Bounded when quiet — consecutive nudges that find no pending fleet state are capped at + * {@code quietNudgeCap}, then the loop stops nudging until real state returns (constraint 4).
  6. + *
  7. Never races {@link ReplyPushLoop} — while that loop is actively nudging any target this + * loop stands down, so two competing injections never start two turns in the same pane + * (constraint 6).
  8. + *
+ */ +public final class LeadHeartbeatLoop { + + private static final Logger log = LoggerFactory.getLogger(LeadHeartbeatLoop.class); + + /** Sentinel for {@link #idleSinceNanos}: the lead is not currently in an idle stretch. */ + private static final long NOT_IDLE = Long.MIN_VALUE; + + private final PrimaryRegistry primaryRegistry; + private final AgentControl agents; + private final ReplyInbox inbox; + private final Supplier> roster; + private final ReplyPushLoop pushLoop; + private final ScheduledExecutorService scheduler; + private final LongSupplier clock; + private final long idleAfterNanos; + private final long backoffMs; + private final int quietNudgeCap; + private final Metrics metrics; // CB-512 pattern: nullable — no registry in unit tests + + /** When the current idle stretch began (nanos), or {@link #NOT_IDLE}. Single scheduler thread only. */ + private long idleSinceNanos = NOT_IDLE; + /** Consecutive nudges that found no pending fleet state. Single scheduler thread only. */ + private int quietCount = 0; + + /** Constructor with an injectable clock and no metric registry (unit tests, or wiring that opts out). */ + public LeadHeartbeatLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox, + Supplier> roster, ReplyPushLoop pushLoop, + ScheduledExecutorService scheduler, LongSupplier clock, + long idleAfterNanos, long backoffMs, int quietNudgeCap) { + this(primaryRegistry, agents, inbox, roster, pushLoop, scheduler, clock, + idleAfterNanos, backoffMs, quietNudgeCap, null); + } + + /** As above, with a metric registry (the CB-512 pattern) so nudge outcomes are counted. */ + public LeadHeartbeatLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox, + Supplier> roster, ReplyPushLoop pushLoop, + ScheduledExecutorService scheduler, LongSupplier clock, + long idleAfterNanos, long backoffMs, int quietNudgeCap, Metrics metrics) { + this.primaryRegistry = primaryRegistry; + this.agents = agents; + this.inbox = inbox; + this.roster = roster; + this.pushLoop = pushLoop; + this.scheduler = scheduler; + this.clock = clock; + this.idleAfterNanos = idleAfterNanos; + this.backoffMs = backoffMs; + this.quietNudgeCap = quietNudgeCap; + this.metrics = metrics; + } + + /** Count one nudge outcome when a registry is wired; a no-op in unit tests. */ + private void countNudge(String outcome) { + if (metrics != null) { + metrics.inc(BridgedMetrics.HEARTBEAT_NUDGES, "outcome", outcome); + } + } + + // --- decision logic (package-private for unit-testing) -------------------------------------- + + /** The action the loop should take on one tick. */ + enum Action { + /** Send a heartbeat nudge to the lead now. */ + INJECT, + /** Lead is idle but not yet past the quiet period — keep waiting, do nothing. */ + WAIT_IDLE, + /** Lead is making progress (or its status cannot be read) — reset the idle window and re-arm. */ + LEAD_BUSY, + /** Idle past the quiet period with nothing pending and the cap exhausted — stop until new state appears. */ + QUIET_DONE, + /** {@link ReplyPushLoop} is actively nudging — stand aside rather than start a competing turn. */ + STAND_DOWN + } + + /** The outcome of one decision: the action plus the state to persist for the next tick. */ + record Decision(Action action, Long idleSinceNanos, int quietCount) {} + + /** + * Pure decision function: given the current loop state and fleet/lead facts, return what to do + * next and the exact state to carry forward. Pure means no I/O and no mutation — the caller + * ({@link #tick}) applies {@link Decision} to its fields. This is what makes the loop unit-testable + * with a fake clock and no sleeping. + * + * @param nowNanos the current time (an injected clock in tests, {@link System#nanoTime()} live) + * @param idleSinceNanos when the current idle stretch began, or {@code null} if the lead is not idle + * @param quietCount consecutive nudges so far that found no pending fleet state + * @param status the lead's herdr status + * @param pushLoopActive whether {@link ReplyPushLoop} is currently nudging some target (constraint 6) + * @param leadKnown whether a lead terminal is known to nudge at all + * @param fleet a snapshot of the pending fleet state (constraint 5) + * @return the action to take and the state to persist + */ + Decision decide(long nowNanos, Long idleSinceNanos, int quietCount, AgentStatus status, + boolean pushLoopActive, boolean leadKnown, FleetState fleet) { + // Constraint 6: while ReplyPushLoop is actively nudging the lead, injecting a second, + // competing prompt into the same pane would start a second turn — racing loops multiply + // turns and context burn. Stand aside, and treat the active push as real state (re-arm the + // quiet counter), because the reply that drove it is exactly the kind of new state that + // should reset the cap. + if (pushLoopActive) { + return new Decision(Action.STAND_DOWN, idleSinceNanos, 0); + } + // Constraint 2: a WORKING lead is making progress and must NOT be touched; an unreadable + // status (read failure, or the agent is gone) is safest treated the same way — never inject + // into a state we cannot read. Either way, reset the idle window and the quiet counter: the + // lead was / may be active, so the next idle stretch must count its own quiet period fresh. + if (status == null || !status.injectable()) { + return new Decision(Action.LEAD_BUSY, null, 0); + } + if (idleSinceNanos == null) { + // The lead just became injectable — record the start of an idle stretch and wait out the + // debounce quiet period before ever nudging (constraint 3). + return new Decision(Action.WAIT_IDLE, nowNanos, quietCount); + } + if (nowNanos - idleSinceNanos < idleAfterNanos) { + // Still within the quiet period: the lead that just finished a turn sits momentarily idle + // and must not be re-prompted into every natural pause. + return new Decision(Action.WAIT_IDLE, idleSinceNanos, quietCount); + } + // Past the quiet period with an injectable lead: it is a genuine candidate for a nudge. Two + // gating facts decide whether and how: + if (!leadKnown) { + // No lead terminal is known yet (e.g. the registry has not learned one) — there is nobody + // to nudge. Keep waiting; the window stays open so discovery re-arms it without a fresh + // quiet period. + return new Decision(Action.WAIT_IDLE, idleSinceNanos, quietCount); + } + if (fleet.hasPending()) { + // Real fleet state is waiting — a worker reply or a DONE session. This is new state, so + // it resets the quiet counter (constraint 4) and the lead is nudged to go collect it. + return new Decision(Action.INJECT, idleSinceNanos, 0); + } + if (quietCount < quietNudgeCap) { + // Nothing is pending, but the cap is not exhausted: nudge anyway, telling the lead + // exactly that nothing is waiting so it can choose to stand down rather than hunt + // (constraint 5). Count it toward the consecutive-quiet cap. + return new Decision(Action.INJECT, idleSinceNanos, quietCount + 1); + } + // Nothing pending and the cap is exhausted: stop nudging until real state appears again + // (constraint 4). The loop still ticks on backoff so a genuinely new reply or session change + // re-arms it — QUIET_DONE stops injection, not observation. + return new Decision(Action.QUIET_DONE, idleSinceNanos, quietCount); + } + + // --- loop ---------------------------------------------------------------------------------- + + /** + * Start the heartbeat. Schedules the first tick after one backoff so a freshly-booting daemon + * does not evaluate the lead's idle state before the fleet has settled. + */ + public void start() { + log.info("idle-lead heartbeat: on — nudge lead after {}s idle (recheck {}ms, quiet cap {})", + TimeUnit.NANOSECONDS.toSeconds(idleAfterNanos), backoffMs, quietNudgeCap); + scheduler.schedule(this::tick, backoffMs, TimeUnit.MILLISECONDS); + } + + /** One loop tick, every {@link #backoffMs} — the thin scheduler around {@link #decide}. */ + private void tick() { + boolean leadKnown = primaryRegistry.primaryTerminal().isPresent(); + FleetState fleet = snapshot(inbox, roster); + AgentStatus status = AgentStatus.UNKNOWN; + if (leadKnown) { + try { + status = agents.status(primaryRegistry.primaryTerminal().orElseThrow()); + } catch (RuntimeException e) { + // A failed status read degrades to "unknown" — decide() treats that like a busy lead + // and never injects into a state it cannot read. Retry on the next backoff. + log.debug("idle-heartbeat: status check failed for lead, will retry: {}", e.toString()); + } + } + + Decision d = decide(clock.getAsLong(), + idleSinceNanos == NOT_IDLE ? null : idleSinceNanos, + quietCount, status, pushLoop.isActive(), leadKnown, fleet); + applyDecision(d); + switch (d.action()) { + case INJECT -> injectNudge(fleet); + case QUIET_DONE -> countNudge("exhausted"); + case WAIT_IDLE, LEAD_BUSY, STAND_DOWN -> { /* nothing to inject, nothing to count */ } + } + scheduleNext(); + } + + /** Persist the state a decision returned, so the next tick starts from it. */ + private void applyDecision(Decision d) { + idleSinceNanos = d.idleSinceNanos() == null ? NOT_IDLE : d.idleSinceNanos(); + quietCount = d.quietCount(); + } + + /** Send the nudge to the known lead. */ + private void injectNudge(FleetState fleet) { + var lead = primaryRegistry.primaryTerminal(); + if (lead.isEmpty()) { + return; // the lead disappeared between the decision and the injection + } + String leadTerminal = lead.get(); + try { + agents.send(leadTerminal, fleet.nudgeText()); + log.debug("idle-heartbeat: nudge sent to lead {} (quiet nudges so far in this stretch: {})", + leadTerminal, quietCount); + countNudge("delivered"); + } catch (RuntimeException e) { + log.warn("idle-heartbeat: failed to nudge lead {}: {}", leadTerminal, e.toString()); + countNudge("failed"); + } + } + + /** Schedule the next tick on the scheduler thread pool. */ + private void scheduleNext() { + scheduler.schedule(this::tick, backoffMs, TimeUnit.MILLISECONDS); + } + + // --- lifecycle ----------------------------------------------------------------------------- + + /** Shut down the scheduler; outstanding ticks are cancelled. */ + public void stop() { + scheduler.shutdownNow(); + } + + /** @see #stop() */ + public void close() { + stop(); + } + + /** + * A read-only snapshot of the fleet facts the heartbeat reports to an idle lead. Carries exact + * counts and names the targets that hold replies, so the nudge is actionable rather than a + * content-free "keep working" that would cost the lead a turn just to discover what there is to + * do (constraint 5). + * + * @param pendingReplies total undrained worker replies across owned targets + * @param doneSessions session count in the DONE state — finished turns awaiting teardown + * @param liveWorkers sessions still registered (acquired minus released) + * @param replyTargets the terminals whose inboxes currently hold at least one reply + */ + record FleetState(int pendingReplies, int doneSessions, int liveWorkers, List replyTargets) { + + /** True when the lead has something real to collect — a reply or a session awaiting teardown. */ + boolean hasPending() { + return pendingReplies > 0 || doneSessions > 0; + } + + /** The nudge body, phrased for the two cases the heartbeat distinguishes. */ + String nudgeText() { + StringBuilder sb = new StringBuilder( + "Heartbeat: you are idle and no bridge_send is waiting on you."); + if (hasPending()) { + sb.append(" The fleet has state to collect: ").append(pendingDetail()); + } else { + sb.append(" Fleet is quiet — nothing pending to collect (") + .append(pendingReplies).append(" worker replies, ") + .append(doneSessions).append(" DONE sessions); ") + .append(liveWorkers).append(" workers live. You may stand down until new state re-arms me."); + } + return sb.toString(); + } + + private String pendingDetail() { + StringBuilder sb = new StringBuilder(); + sb.append(pendingReplies).append(" worker ") + .append(pendingReplies == 1 ? "reply" : "replies") + .append(" pending collection"); + if (!replyTargets.isEmpty()) { + // Render each as the exact command so the lead can act without parsing: the nearest + // analogue to ReplyPushLoop's bridge_poll(target=...) nudge. + sb.append(" (") + .append(String.join(", ", + replyTargets.stream().map(t -> "bridge_poll(target=" + t + ")").toList())) + .append(")"); + } + sb.append(", ").append(doneSessions).append(" DONE session").append(doneSessions == 1 ? "" : "s") + .append(" awaiting teardown") + .append(", ").append(liveWorkers).append(" worker").append(liveWorkers == 1 ? "" : "s") + .append(" live"); + return sb.toString(); + } + } + + /** + * Build a {@link FleetState} snapshot over the live roster and the reply inbox — the same + * sources the metrics collector uses, so the heartbeat reports facts the operator can cross-check + * against {@code /metrics}. Package-private and pure (reads only) so it is unit-testable. + */ + static FleetState snapshot(ReplyInbox inbox, Supplier> roster) { + List sessions = roster.get(); + int replies = 0; + int done = 0; + List replyTargets = new ArrayList<>(); + for (WorkerSession s : sessions) { + if (s.terminalId() == null) { + continue; + } + int depth = inbox.peek(s.terminalId()).size(); + if (depth > 0) { + replies += depth; + replyTargets.add(s.terminalId()); + } + if (s.state() == WorkerSession.State.DONE) { + done++; + } + } + return new FleetState(replies, done, sessions.size(), List.copyOf(replyTargets)); + } +} diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java b/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java index 2ff0712..890a0c0 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java @@ -175,6 +175,17 @@ public final class ReplyPushLoop { // --- lifecycle ----------------------------------------------------------------------------- + /** + * Whether any reminder loop is currently active for some target (CB-551). The idle-lead heartbeat + * uses this to stand aside: while the push loop is actively nudging the lead, a concurrent + * heartbeat injection would start a second competing turn in the same pane — racing loops multiply + * turns and context burn. "Active" means a schedule exists in {@link #activeTargets}; the set is + * bounded by what has been triggered, not by any persistent state. + */ + public boolean isActive() { + return !activeTargets.isEmpty(); + } + /** Shut down the scheduler. Outstanding reminders are cancelled. */ public void stop() { scheduler.shutdownNow(); diff --git a/bridged/src/test/java/dev/ltms/bridged/config/BridgedConfigTest.java b/bridged/src/test/java/dev/ltms/bridged/config/BridgedConfigTest.java index fb8e17c..a5bdb27 100644 --- a/bridged/src/test/java/dev/ltms/bridged/config/BridgedConfigTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/config/BridgedConfigTest.java @@ -132,6 +132,7 @@ class BridgedConfigTest { leaders: {} architects: {} leadScan: {} + leadHeartbeat: {} placement: fixed auth: {} """).isEmpty(), "the known-key set must not drift from the record components"); @@ -180,6 +181,51 @@ class BridgedConfigTest { assertEquals(30, scan.intervalSeconds()); } + // ── CB-551: the idle-lead heartbeat ───────────────────────────────────────────────────────── + + /** + * Acceptance (f): absent config ⇒ the heartbeat is OFF completely. It spends the operator's + * subscription on its own initiative, so an upgraded daemon must never switch it on unasked. + */ + @Test + void leadHeartbeatIsOffUnlessTheBlockIsPresent(@TempDir Path dir) throws Exception { + Path f = dir.resolve("no-heartbeat.yaml"); + Files.writeString(f, "bind:\n port: 8080\n"); + + assertNull(BridgedConfig.load(f).leadHeartbeat(), + "no leadHeartbeat: block ⇒ the loop is never constructed and never fires (CB-551)"); + } + + @Test + void leadHeartbeatDefaultsItsFieldsWhenTheBlockIsPresentButBare(@TempDir Path dir) throws Exception { + Path f = dir.resolve("bare-heartbeat.yaml"); + Files.writeString(f, "bind:\n port: 8080\nleadHeartbeat: {}\n"); + + BridgedConfig.LeadHeartbeat hb = BridgedConfig.load(f).leadHeartbeat(); + assertNotNull(hb); + assertEquals(300, hb.idleAfterSeconds(), "default quiet period"); + assertEquals(60_000L, hb.backoffMs(), "default backoff"); + assertEquals(3, hb.quietNudgeCap(), "default quiet-nudge cap"); + } + + @Test + void leadHeartbeatReadsExplicitValues(@TempDir Path dir) throws Exception { + Path f = dir.resolve("heartbeat.yaml"); + Files.writeString(f, """ + bind: + port: 8080 + leadHeartbeat: + idleAfterSeconds: 600 + backoffMs: 45000 + quietNudgeCap: 5 + """); + + BridgedConfig.LeadHeartbeat hb = BridgedConfig.load(f).leadHeartbeat(); + assertEquals(600, hb.idleAfterSeconds()); + assertEquals(45_000L, hb.backoffMs()); + assertEquals(5, hb.quietNudgeCap()); + } + /** * The hazard the guard exists for: bridged writes worker tab labels and reads lead tab labels. * Overlap the two and every worker it spawns is read back as a lead. @@ -835,6 +881,10 @@ class BridgedConfigTest { terminal: term_abc123 pushReminders: 5 pushBackoffMs: 15000 + leadHeartbeat: + idleAfterSeconds: 600 + backoffMs: 45000 + quietNudgeCap: 5 """); BridgedConfig cfg = BridgedConfig.load(f); @@ -860,6 +910,9 @@ class BridgedConfigTest { assertEquals("term_abc123", cfg.primary().terminal()); assertEquals(5, cfg.primary().remindersOrDefault()); assertEquals(15000L, cfg.primary().backoffMsOrDefault()); + assertEquals(600, cfg.leadHeartbeat().idleAfterSeconds(), "leadHeartbeat binds at the top level"); + assertEquals(45_000L, cfg.leadHeartbeat().backoffMs()); + assertEquals(5, cfg.leadHeartbeat().quietNudgeCap()); } @Test diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/LeadHeartbeatLoopTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/LeadHeartbeatLoopTest.java new file mode 100644 index 0000000..68c69ee --- /dev/null +++ b/bridged/src/test/java/dev/ltms/bridged/msg/LeadHeartbeatLoopTest.java @@ -0,0 +1,215 @@ +package dev.ltms.bridged.msg; + +import dev.ltms.bridged.herdr.AgentStatus; +import dev.ltms.bridged.mcp.PrimaryRegistry; +import dev.ltms.bridged.session.WorkerSession; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import java.util.List; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; + +import static org.junit.jupiter.api.Assertions.*; + +/** + * Unit tests for the CB-551 idle-lead heartbeat: the pure {@link LeadHeartbeatLoop#decide} decision + * function, the nudge text, and the {@link LeadHeartbeatLoop#snapshot} fleet snapshot. + * + *

All decision tests call {@code decide} directly with explicit nanoTime values from an injected + * clock — no sleeping, no scheduler races. This mirrors how {@code ReplyPushLoopTest} pins the pure + * decision before exercising the loop. + */ +class LeadHeartbeatLoopTest { + + /** Arbitrary nanoTime origin for the fake clock. */ + private static final long NOW = 1_000_000_000L; + /** 300s (the default quiet period) in nanos — under it the lead is "within the quiet period". */ + private static final long IDLE_AFTER_NANOS = TimeUnit.SECONDS.toNanos(300); + /** 400s of idle — clearly past the quiet period. */ + private static final long IDLE_PAST = NOW - TimeUnit.SECONDS.toNanos(400); + /** 10s of idle — clearly within the quiet period. */ + private static final long IDLE_WITHIN = NOW - TimeUnit.SECONDS.toNanos(10); + + private ScheduledExecutorService scheduler; + + @BeforeEach + void setUp() { + scheduler = Executors.newSingleThreadScheduledExecutor(); + } + + @AfterEach + void tearDown() { + scheduler.shutdownNow(); + } + + /** A fleet with nothing pending. */ + private static LeadHeartbeatLoop.FleetState quietFleet() { + return new LeadHeartbeatLoop.FleetState(0, 0, 2, List.of()); + } + + /** A fleet with a pending reply and a DONE session. */ + private static LeadHeartbeatLoop.FleetState pendingFleet() { + return new LeadHeartbeatLoop.FleetState(2, 1, 3, List.of("term_a", "term_b")); + } + + /** The loop under test; the scheduler is never invoked on the pure decide path. */ + private static LeadHeartbeatLoop loop(int quietCap) { + return new LeadHeartbeatLoop( + new PrimaryRegistry("term_lead"), null /*agents — unused on the decide path*/, + null /*inbox*/, List::of, null /*pushLoop*/, null /*scheduler*/, () -> 0L, + IDLE_AFTER_NANOS, 1_000L, quietCap); + } + + // ── (a) a working lead is never injected ─────────────────────────────────────────────────── + + @Test + void aWorkingLeadIsNeverInjected() { + LeadHeartbeatLoop.Decision d = loop(3).decide( + NOW, IDLE_PAST, 0, AgentStatus.WORKING, false, true, quietFleet()); + assertEquals(LeadHeartbeatLoop.Action.LEAD_BUSY, d.action(), + "a WORKING lead is making progress and must not be touched"); + assertNull(d.idleSinceNanos(), "a busy lead resets the idle window"); + assertEquals(0, d.quietCount(), "a busy lead re-arms the quiet counter"); + } + + @Test + void anUnreadableStatusIsNeverInjectedEither() { + // A failed status read (or a gone agent) must degrade to "do not inject", never hammer the pane. + LeadHeartbeatLoop.Decision d = loop(3).decide( + NOW, IDLE_PAST, 0, AgentStatus.UNKNOWN, false, true, pendingFleet()); + assertEquals(LeadHeartbeatLoop.Action.LEAD_BUSY, d.action(), + "never inject into a state the loop cannot read"); + } + + // ── (b) an idle lead within the quiet period is not yet injected ─────────────────────────── + + @Test + void justBecameIdleStartsTheDebounceWindow() { + LeadHeartbeatLoop.Decision d = loop(3).decide( + NOW, null, 0, AgentStatus.IDLE, false, true, quietFleet()); + assertEquals(LeadHeartbeatLoop.Action.WAIT_IDLE, d.action(), + "the first injectable tick only records the start of the idle stretch"); + assertEquals(NOW, d.idleSinceNanos(), "the idle window opens at the moment the lead became injectable"); + } + + @Test + void idleWithinQuietPeriodIsNotInjected() { + LeadHeartbeatLoop.Decision d = loop(3).decide( + NOW, IDLE_WITHIN, 0, AgentStatus.IDLE, false, true, quietFleet()); + assertEquals(LeadHeartbeatLoop.Action.WAIT_IDLE, d.action(), + "a lead idle for 10s (< 300s) has just finished a turn — do not re-prompt it"); + } + + // ── (c) an idle lead past the quiet period is injected ───────────────────────────────────── + + @Test + void idlePastQuietPeriodIsInjected() { + LeadHeartbeatLoop.Decision d = loop(3).decide( + NOW, IDLE_PAST, 0, AgentStatus.IDLE, false, true, quietFleet()); + assertEquals(LeadHeartbeatLoop.Action.INJECT, d.action(), + "a lead continuously idle past the quiet period is the reason to nudge"); + } + + @Test + void blockedAndDoneAreInjectableViewsOfIdle() { + assertEquals(LeadHeartbeatLoop.Action.INJECT, loop(3).decide( + NOW, IDLE_PAST, 0, AgentStatus.BLOCKED, false, true, quietFleet()).action()); + assertEquals(LeadHeartbeatLoop.Action.INJECT, loop(3).decide( + NOW, IDLE_PAST, 0, AgentStatus.DONE, false, true, quietFleet()).action()); + } + + @Test + void nothingIsInjectedWhenNoLeadIsKnown() { + var d = loop(3).decide(NOW, IDLE_PAST, 0, AgentStatus.IDLE, false, false, pendingFleet()); + assertEquals(LeadHeartbeatLoop.Action.WAIT_IDLE, d.action(), + "with no known lead there is nobody to nudge — keep waiting until one is discovered"); + assertEquals(IDLE_PAST, d.idleSinceNanos(), "the idle window stays open so discovery re-arms it"); + } + + // ── (d) the quiet-nudge cap stops the loop ────────────────────────────────────────────────── + + @Test + void quietNudgeCapStopsTheLoopWhenNothingIsPending() { + LeadHeartbeatLoop.Decision d = loop(3).decide( + NOW, IDLE_PAST, 3, AgentStatus.IDLE, false, true, quietFleet()); + assertEquals(LeadHeartbeatLoop.Action.QUIET_DONE, d.action(), + "3 consecutive nothing-pending nudges have already happened — stop nagging an empty fleet"); + } + + // ── (e) new pending state resets the cap and re-arms the loop ─────────────────────────────── + + @Test + void newPendingStateResetsTheQuietCap() { + LeadHeartbeatLoop.Decision d = loop(3).decide( + NOW, IDLE_PAST, 3, AgentStatus.IDLE, false, true, pendingFleet()); + assertEquals(LeadHeartbeatLoop.Action.INJECT, d.action(), + "real state appearing re-arms the loop past an exhausted cap"); + assertEquals(0, d.quietCount(), "the pending state resets the consecutive-quiet counter"); + } + + // ── constraint 6: never race ReplyPushLoop ───────────────────────────────────────────────── + + @Test + void standsDownWhileReplyPushLoopIsActive() { + LeadHeartbeatLoop.Decision d = loop(3).decide( + NOW, IDLE_PAST, 2, AgentStatus.IDLE, true, true, quietFleet()); + assertEquals(LeadHeartbeatLoop.Action.STAND_DOWN, d.action(), + "a second injection would start a competing turn — stand aside instead"); + assertEquals(0, d.quietCount(), "the active push is real state, so it re-arms the cap"); + } + + // ── nudge text (constraint 5: the nudge must carry state) ────────────────────────────────── + + @Test + void pendingNudgeCarriesTheFleetFacts() { + String t = pendingFleet().nudgeText(); + assertTrue(t.startsWith("Heartbeat:"), "the nudge identifies itself as a heartbeat"); + assertTrue(t.contains("2 worker replies pending collection"), t); + assertTrue(t.contains("bridge_poll(target=term_a)"), t); + assertTrue(t.contains("bridge_poll(target=term_b)"), t); + assertTrue(t.contains("1 DONE session awaiting teardown"), t); + assertTrue(t.contains("3 workers live"), t); + } + + @Test + void nothingPendingNudgeSaysExactlyThatSoTheLeadCanStandDown() { + String t = quietFleet().nudgeText(); + assertTrue(t.contains("nothing pending to collect"), t); + assertTrue(t.contains("0 worker replies"), t); + assertTrue(t.contains("0 DONE sessions"), t); + assertTrue(t.contains("2 workers live"), t); + assertTrue(t.contains("stand down"), "allow the lead to stand down rather than hunt"); + } + + // ── snapshot over roster + inbox ─────────────────────────────────────────────────────────── + + @Test + void snapshotCountsRepliesDoneSessionsAndLiveWorkers() { + InMemoryReplyInbox inbox = new InMemoryReplyInbox(); + inbox.own("term_w1"); + inbox.publish("term_w1", "m1", "hello"); + WorkerSession done = new WorkerSession("p1", "term_w1", "prof", "/cwd", null, + 0, 0, 0, WorkerSession.State.DONE, null, null); + WorkerSession ready = new WorkerSession("p2", "term_w2", "prof", "/cwd", null, + 0, 0, 0, WorkerSession.State.READY, null, null); + + LeadHeartbeatLoop.FleetState fs = LeadHeartbeatLoop.snapshot(inbox, () -> List.of(done, ready)); + + assertEquals(1, fs.pendingReplies(), "one undrained reply on term_w1"); + assertEquals(1, fs.doneSessions(), "term_w1 is DONE, awaiting teardown"); + assertEquals(2, fs.liveWorkers(), "both sessions are still registered"); + assertEquals(List.of("term_w1"), fs.replyTargets(), "the target with a pending reply is named"); + assertTrue(fs.hasPending(), "a pending reply counts as fleet state to collect"); + } + + @Test + void snapshotWithEmptyRosterHasNoPendingState() { + InMemoryReplyInbox inbox = new InMemoryReplyInbox(); + LeadHeartbeatLoop.FleetState fs = LeadHeartbeatLoop.snapshot(inbox, List::of); + assertFalse(fs.hasPending()); + assertEquals(0, fs.liveWorkers()); + } +} diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java index 7cff0d8..5051450 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/ReplyPushLoopTest.java @@ -179,6 +179,20 @@ class ReplyPushLoopTest { // --- nudge format -------------------------------------------------------------------------- + @Test + void isActiveReflectsALiveScheduleForTheHeartbeatStandDown() { + var rec = recordingClient(); + agents = new AgentControl(rec); + inbox.publish(WORKER, "m1", "hello"); + ReplyPushLoop loop = loop(1, 100_000); // long backoff so the tick cannot fire mid-test + + loop.onReplyQueued(WORKER); + assertTrue(loop.isActive(), "CB-551: the heartbeat must stand aside while a reminder is live"); + + loop.stop(); + assertFalse(loop.isActive(), "stopping clears the active schedule"); + } + @Test void nudgeFormatIsCorrect() { String nudge = ReplyPushLoop.NUDGE_FORMAT.formatted(WORKER, WORKER);