diff --git a/fleetd/src/main/java/dev/ltms/fleet/inject/Injector.java b/fleetd/src/main/java/dev/ltms/fleet/inject/Injector.java index 08c154a..be97219 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/inject/Injector.java +++ b/fleetd/src/main/java/dev/ltms/fleet/inject/Injector.java @@ -15,6 +15,7 @@ import java.util.Set; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; import java.util.function.Consumer; +import java.util.function.LongSupplier; import java.util.function.Predicate; import java.util.stream.Collectors; @@ -94,6 +95,14 @@ public final class Injector { private final TurnListener turnListener; private final Predicate ready; // CB-113: a target is deliverable only when available private final Consumer forget; // CB-114: clear a gone worker's readiness/presence + /** + * Wall-clock source for the readiness-grace elapsed time logged in {@link #onStatus} (fleetd + * #501). Production constructors default this to {@code System::currentTimeMillis}; the + * package-private constructors below take it explicitly so a test can supply a stub whose + * advance does not track {@link #POLL_INTERVAL_MILLIS} — copying the shape {@code LeadRollover} + * already uses for the same purpose. + */ + private final LongSupplier nowMillis; private final ConcurrentHashMap targets = new ConcurrentHashMap<>(); /** Delivery only; completion signalling is a no-op and every target is treated as available. */ @@ -125,20 +134,34 @@ public final class Injector { */ public Injector(AgentControl agents, TurnListener turnListener, Predicate ready, Consumer forget) { + this(agents, turnListener, ready, forget, System::currentTimeMillis); + } + + /** Full constructor — for tests: an injectable wall-clock supplier (fleetd #501). */ + Injector(AgentControl agents, TurnListener turnListener, Predicate ready, + Consumer forget, LongSupplier nowMillis) { this.agents = agents; this.router = null; this.turnListener = turnListener; this.ready = ready; this.forget = forget; + this.nowMillis = nowMillis; } public Injector(HerdrRouter router, TurnListener turnListener, Predicate ready, Consumer forget) { + this(router, turnListener, ready, forget, System::currentTimeMillis); + } + + /** Full constructor — for tests: an injectable wall-clock supplier (fleetd #501). */ + Injector(HerdrRouter router, TurnListener turnListener, Predicate ready, + Consumer forget, LongSupplier nowMillis) { this.agents = null; this.router = router; this.turnListener = turnListener; this.ready = ready; this.forget = forget; + this.nowMillis = nowMillis; } private AgentControl agentsFor(String target) { @@ -209,6 +232,7 @@ public final class Injector { int unknownSinceTurn; // consecutive `unknown` samples while a delegation is outstanding (CB-109) int unknownSincePostTurn; // the same, for the post-turn housekeeping phase (fleetd #306) int notReadySincePoll; // consecutive injectable samples a queued message waited on the readiness gate (CB-114) + long notReadySinceMillis; // wall-clock time of the FIRST non-ready sample in the current notReadySincePoll streak (fleetd #501); reset alongside it boolean postTurnPending; // completion observed; adapter housekeeping has not started yet boolean awaitingPostTurnPickup; boolean postTurnObserved; @@ -302,6 +326,7 @@ public final class Injector { t.unknownSinceTurn = 0; t.unknownSincePostTurn = 0; t.notReadySincePoll = 0; + t.notReadySinceMillis = 0; if (t.awaitingCompletion) t.turnObserved = true; } else if (status.injectable()) { // IDLE or BLOCKED t.unknownSinceTurn = 0; @@ -353,6 +378,7 @@ public final class Injector { Pending p = t.queue.peek(); if (p != null && ready.test(target)) { t.notReadySincePoll = 0; + t.notReadySinceMillis = 0; try { agentsFor(target).send(target, p.text()); t.queue.poll(); @@ -370,23 +396,52 @@ public final class Injector { sent = p; sendError = e; } - } else if (p != null && ++t.notReadySincePoll >= READINESS_GRACE_POLLS) { - // The worker has been idle-but-not-ready for the whole grace: its Claude - // never connected the bridge MCP (crashed during boot, or wedged on a - // startup prompt). The readiness gate would hold this message forever, so - // fail every queued message and release the target (CB-114) instead of - // polling it indefinitely with the caller's future never completing. - notReady = new ArrayList<>(t.queue); - for (Pending pending : notReady) { - pending.state = Pending.State.NOT_DELIVERED; + } else if (p != null) { + // fleetd #501: stamp the wall-clock time of the FIRST non-ready sample in + // this streak, so the expiry log below can print how long the target + // actually sat non-ready — not just how many polls that took. + if (t.notReadySincePoll == 0) { + t.notReadySinceMillis = nowMillis.getAsLong(); + } + if (++t.notReadySincePoll >= READINESS_GRACE_POLLS) { + // The worker has been idle-but-not-ready for the whole grace: its Claude + // never connected the bridge MCP (crashed during boot, or wedged on a + // startup prompt). The readiness gate would hold this message forever, so + // fail every queued message and release the target (CB-114) instead of + // polling it indefinitely with the caller's future never completing. + notReady = new ArrayList<>(t.queue); + for (Pending pending : notReady) { + pending.state = Pending.State.NOT_DELIVERED; + } + // fleetd #501: t.notReadySincePoll — the loop's own counter, already in + // scope — is printed here instead of the READINESS_GRACE_POLLS constant. + // On this branch the counter has JUST reached the threshold, so the two + // agree by construction and no test can tell them apart. Printed anyway: + // it gives this line one source of truth instead of two, so a later + // change to the loop above cannot leave this message reporting a number + // the loop no longer produces. + // + // elapsedMillis is a different case: it is NOT equal-by-construction to + // the truth. notReadySincePoll only increments on a sample that reaches + // this branch (p != null, not ready) — a poll that misses that condition + // advances real time without advancing the counter — and this loop's real + // period is not guaranteed to equal POLL_INTERVAL_MILLIS (load, or a host + // sleep, can widen the real gap far past it). READINESS_GRACE_POLLS * + // POLL_INTERVAL_MILLIS / 1000 is arithmetic on two constants, not a + // measurement, so it stays here only as the labelled CONFIGURED budget, + // never presented as elapsed time. + long elapsedMillis = nowMillis.getAsLong() - t.notReadySinceMillis; + log.warn("readiness grace for {} expired after {} polls (configured={} " + + "polls/{}s elapsed={}ms): target never became " + + "deliverable, so failing {} queued message(s) that " + + "never reached its pane", + target, t.notReadySincePoll, READINESS_GRACE_POLLS, + READINESS_GRACE_POLLS * POLL_INTERVAL_MILLIS / 1000, elapsedMillis, + notReady.size()); + t.queue.clear(); + t.notReadySincePoll = 0; + t.notReadySinceMillis = 0; } - log.warn("readiness grace for {} expired after {} polls ({}s): target never " - + "became deliverable, so failing {} queued message(s) that never " - + "reached its pane", - target, READINESS_GRACE_POLLS, - READINESS_GRACE_POLLS * POLL_INTERVAL_MILLIS / 1000, notReady.size()); - t.queue.clear(); - t.notReadySincePoll = 0; } } } diff --git a/fleetd/src/test/java/dev/ltms/fleet/inject/InjectorTest.java b/fleetd/src/test/java/dev/ltms/fleet/inject/InjectorTest.java index 77d4056..ed18607 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/inject/InjectorTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/inject/InjectorTest.java @@ -19,6 +19,8 @@ import java.util.Set; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.function.LongSupplier; import static org.junit.jupiter.api.Assertions.*; @@ -521,6 +523,97 @@ class InjectorTest { } } + @Test + void readinessGraceExpiryLogsTheMeasuredPollCountNextToTheConfiguredBudget() { + // fleetd #501, defect 1: READINESS_GRACE_POLLS (240) used to be printed twice — once as + // "after {} polls" and once inside the parenthesised budget — even though the loop's own + // counter (Target.notReadySincePoll) was in scope at the same call site. On THIS branch the + // counter has just reached the threshold, so it equals the constant by construction and this + // test cannot tell the two apart — it only pins that the message still carries both a poll + // count and a labelled configured budget, using literal numbers (240, 60), never + // READINESS_GRACE_POLLS or POLL_INTERVAL_MILLIS, so the assertion can't silently track a + // constant change instead of catching a real regression. + LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory(); + ch.qos.logback.classic.Logger injectorLog = + (ch.qos.logback.classic.Logger) LoggerFactory.getLogger(Injector.class); + ListAppender appender = new ListAppender<>(); + appender.setContext(ctx); + appender.start(); + injectorLog.addAppender(appender); + injectorLog.setLevel(Level.WARN); + try { + Injector inj = new Injector(new AgentControl(herdr), TurnListener.NOOP, _ -> false, _ -> { + }); + inj.enqueue(T, "task", TestTurnTokens.inert(T)); + + for (int i = 0; i < READINESS_SAMPLES; i++) inj.onStatus(T, AgentStatus.IDLE); + + String warn = appender.list.stream() + .filter(e -> e.getLevel().equals(Level.WARN)) + .map(ILoggingEvent::getFormattedMessage) + .findFirst() + .orElse("no grace-expiry WARN logged"); + assertTrue(warn.contains("after 240 polls"), + "must print the measured poll count as a plain number: " + warn); + assertTrue(warn.contains("configured=240 polls/60s"), + "must print the configured budget, clearly labelled: " + warn); + } finally { + injectorLog.detachAppender(appender); + } + } + + @Test + void readinessGraceExpiryLogsTheMeasuredElapsedTimeNotArithmeticOnConstants() { + // fleetd #501, defect 2: the old line computed "({}s)" as READINESS_GRACE_POLLS * + // POLL_INTERVAL_MILLIS / 1000 — arithmetic on two constants, never a measurement, and wrong + // in the direction that says everything ran on schedule. This stub clock returns two FIXED + // values (1_000ms at the first non-ready sample, 318_412ms at the poll that trips the grace) + // whose difference — 317_412ms — does NOT equal 240 * POLL_INTERVAL_MILLIS (=60_000ms). + // Asserting on that literal, non-derived number is what makes this test able to fail if the + // production code goes back to printing the constant-arithmetic value instead of the + // injected clock's measurement. + long[] readings = {1_000L, 318_412L}; + AtomicInteger call = new AtomicInteger(0); + LongSupplier stubClock = () -> { + int i = call.getAndIncrement(); + if (i >= readings.length) { + throw new AssertionError("nowMillis read more times than this fixture expects (" + i + + "); the readiness-not-ready branch should read the clock exactly twice — " + + "once to stamp the first non-ready sample, once at grace expiry"); + } + return readings[i]; + }; + + LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory(); + ch.qos.logback.classic.Logger injectorLog = + (ch.qos.logback.classic.Logger) LoggerFactory.getLogger(Injector.class); + ListAppender appender = new ListAppender<>(); + appender.setContext(ctx); + appender.start(); + injectorLog.addAppender(appender); + injectorLog.setLevel(Level.WARN); + try { + Injector inj = new Injector(new AgentControl(herdr), TurnListener.NOOP, _ -> false, _ -> { + }, stubClock); + inj.enqueue(T, "task", TestTurnTokens.inert(T)); + + for (int i = 0; i < READINESS_SAMPLES; i++) inj.onStatus(T, AgentStatus.IDLE); + + String warn = appender.list.stream() + .filter(e -> e.getLevel().equals(Level.WARN)) + .map(ILoggingEvent::getFormattedMessage) + .findFirst() + .orElse("no grace-expiry WARN logged"); + assertTrue(warn.contains("elapsed=317412ms"), "must print the MEASURED elapsed time from " + + "the injected clock (318412 - 1000 = 317412), not an arithmetic value: " + warn); + assertFalse(warn.contains("elapsed=60000ms"), "must not print " + + "READINESS_GRACE_POLLS * POLL_INTERVAL_MILLIS (240 * 250 = 60000ms) as if it " + + "were the measured elapsed time: " + warn); + } finally { + injectorLog.detachAppender(appender); + } + } + @Test void aWorkerThatBecomesReadyWithinTheGraceIsDeliveredNormally() { // The readiness grace must not fail a worker that is merely slow to boot: once it becomes