Merge #503: Injector's readiness-grace warn prints measured elapsed time, never arithmetic on constants (fleetd #501)
Verified by the lead, not taken from the worker's report. Trial-merged onto main (4f28da6) and built the merge, because a clean auto-merge is not a compiling merge: mvn clean install -> Tests run: 1689, Failures: 0, Errors: 0, Skipped: 0, BUILD SUCCESS InjectorTest -> Tests run: 34, Failures: 0 Gitea CI onac351ee: success The worker mutated the two log arguments. I mutated the half it did not: the STAMP SITE. Changing the first-sample guard so the clock is re-stamped on every non-ready poll gave 1 failure, InjectorTest.readinessGraceExpiryLogsTheMeasuredElapsedTimeNotArithmeticOnConstants, on a clock-read counter assertion: "the readiness-not-ready branch should read the clock exactly twice — once to stamp the first non-ready sample, once at grace expiry". Two greps proved the mutant applied; restored, shasum byte-identical, control 34/0 green. Behaviour preserved: the restructure from `else if (p != null && ++t.notReadySincePoll >= N)` to a nested `else if (p != null) { ... }` keeps the short-circuit, so the counter still increments only when p != null. All three reset sites now clear notReadySinceMillis alongside notReadySincePoll. Two notes for the record, neither a blocker. The poll-count half is an equivalent mutant and the worker said so instead of reporting a kill it did not get. That is the right call and the code comment states the limit honestly. The clock is injected through a package-private constructor overload while the public constructors default it to System::currentTimeMillis. My brief prescribed that shape. With exactly one production construction site (Fleetd.java:531) a required parameter — fleetd #415's antidote — would have been just as cheap and would match LeadRollover and the new Fleetd.awaitHerdr. The silent-survivor risk #415 names is not present here, because the field is final and both public constructors delegate, so a future constructor cannot compile without supplying it. Recording the choice so the next person does not read it as an oversight.
This commit was merged in pull request #503.
This commit is contained in:
@@ -15,6 +15,7 @@ import java.util.Set;
|
|||||||
import java.util.concurrent.CompletableFuture;
|
import java.util.concurrent.CompletableFuture;
|
||||||
import java.util.concurrent.ConcurrentHashMap;
|
import java.util.concurrent.ConcurrentHashMap;
|
||||||
import java.util.function.Consumer;
|
import java.util.function.Consumer;
|
||||||
|
import java.util.function.LongSupplier;
|
||||||
import java.util.function.Predicate;
|
import java.util.function.Predicate;
|
||||||
import java.util.stream.Collectors;
|
import java.util.stream.Collectors;
|
||||||
|
|
||||||
@@ -94,6 +95,14 @@ public final class Injector {
|
|||||||
private final TurnListener turnListener;
|
private final TurnListener turnListener;
|
||||||
private final Predicate<String> ready; // CB-113: a target is deliverable only when available
|
private final Predicate<String> ready; // CB-113: a target is deliverable only when available
|
||||||
private final Consumer<String> forget; // CB-114: clear a gone worker's readiness/presence
|
private final Consumer<String> 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<String, Target> targets = new ConcurrentHashMap<>();
|
private final ConcurrentHashMap<String, Target> targets = new ConcurrentHashMap<>();
|
||||||
|
|
||||||
/** Delivery only; completion signalling is a no-op and every target is treated as available. */
|
/** 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<String> ready,
|
public Injector(AgentControl agents, TurnListener turnListener, Predicate<String> ready,
|
||||||
Consumer<String> forget) {
|
Consumer<String> 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<String> ready,
|
||||||
|
Consumer<String> forget, LongSupplier nowMillis) {
|
||||||
this.agents = agents;
|
this.agents = agents;
|
||||||
this.router = null;
|
this.router = null;
|
||||||
this.turnListener = turnListener;
|
this.turnListener = turnListener;
|
||||||
this.ready = ready;
|
this.ready = ready;
|
||||||
this.forget = forget;
|
this.forget = forget;
|
||||||
|
this.nowMillis = nowMillis;
|
||||||
}
|
}
|
||||||
|
|
||||||
public Injector(HerdrRouter router, TurnListener turnListener, Predicate<String> ready,
|
public Injector(HerdrRouter router, TurnListener turnListener, Predicate<String> ready,
|
||||||
Consumer<String> forget) {
|
Consumer<String> 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<String> ready,
|
||||||
|
Consumer<String> forget, LongSupplier nowMillis) {
|
||||||
this.agents = null;
|
this.agents = null;
|
||||||
this.router = router;
|
this.router = router;
|
||||||
this.turnListener = turnListener;
|
this.turnListener = turnListener;
|
||||||
this.ready = ready;
|
this.ready = ready;
|
||||||
this.forget = forget;
|
this.forget = forget;
|
||||||
|
this.nowMillis = nowMillis;
|
||||||
}
|
}
|
||||||
|
|
||||||
private AgentControl agentsFor(String target) {
|
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 unknownSinceTurn; // consecutive `unknown` samples while a delegation is outstanding (CB-109)
|
||||||
int unknownSincePostTurn; // the same, for the post-turn housekeeping phase (fleetd #306)
|
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)
|
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 postTurnPending; // completion observed; adapter housekeeping has not started yet
|
||||||
boolean awaitingPostTurnPickup;
|
boolean awaitingPostTurnPickup;
|
||||||
boolean postTurnObserved;
|
boolean postTurnObserved;
|
||||||
@@ -302,6 +326,7 @@ public final class Injector {
|
|||||||
t.unknownSinceTurn = 0;
|
t.unknownSinceTurn = 0;
|
||||||
t.unknownSincePostTurn = 0;
|
t.unknownSincePostTurn = 0;
|
||||||
t.notReadySincePoll = 0;
|
t.notReadySincePoll = 0;
|
||||||
|
t.notReadySinceMillis = 0;
|
||||||
if (t.awaitingCompletion) t.turnObserved = true;
|
if (t.awaitingCompletion) t.turnObserved = true;
|
||||||
} else if (status.injectable()) { // IDLE or BLOCKED
|
} else if (status.injectable()) { // IDLE or BLOCKED
|
||||||
t.unknownSinceTurn = 0;
|
t.unknownSinceTurn = 0;
|
||||||
@@ -353,6 +378,7 @@ public final class Injector {
|
|||||||
Pending p = t.queue.peek();
|
Pending p = t.queue.peek();
|
||||||
if (p != null && ready.test(target)) {
|
if (p != null && ready.test(target)) {
|
||||||
t.notReadySincePoll = 0;
|
t.notReadySincePoll = 0;
|
||||||
|
t.notReadySinceMillis = 0;
|
||||||
try {
|
try {
|
||||||
agentsFor(target).send(target, p.text());
|
agentsFor(target).send(target, p.text());
|
||||||
t.queue.poll();
|
t.queue.poll();
|
||||||
@@ -370,23 +396,52 @@ public final class Injector {
|
|||||||
sent = p;
|
sent = p;
|
||||||
sendError = e;
|
sendError = e;
|
||||||
}
|
}
|
||||||
} else if (p != null && ++t.notReadySincePoll >= READINESS_GRACE_POLLS) {
|
} else if (p != null) {
|
||||||
// The worker has been idle-but-not-ready for the whole grace: its Claude
|
// fleetd #501: stamp the wall-clock time of the FIRST non-ready sample in
|
||||||
// never connected the bridge MCP (crashed during boot, or wedged on a
|
// this streak, so the expiry log below can print how long the target
|
||||||
// startup prompt). The readiness gate would hold this message forever, so
|
// actually sat non-ready — not just how many polls that took.
|
||||||
// fail every queued message and release the target (CB-114) instead of
|
if (t.notReadySincePoll == 0) {
|
||||||
// polling it indefinitely with the caller's future never completing.
|
t.notReadySinceMillis = nowMillis.getAsLong();
|
||||||
notReady = new ArrayList<>(t.queue);
|
}
|
||||||
for (Pending pending : notReady) {
|
if (++t.notReadySincePoll >= READINESS_GRACE_POLLS) {
|
||||||
pending.state = Pending.State.NOT_DELIVERED;
|
// 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;
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -19,6 +19,8 @@ import java.util.Set;
|
|||||||
import java.util.concurrent.CompletableFuture;
|
import java.util.concurrent.CompletableFuture;
|
||||||
import java.util.concurrent.ExecutionException;
|
import java.util.concurrent.ExecutionException;
|
||||||
import java.util.concurrent.TimeUnit;
|
import java.util.concurrent.TimeUnit;
|
||||||
|
import java.util.concurrent.atomic.AtomicInteger;
|
||||||
|
import java.util.function.LongSupplier;
|
||||||
|
|
||||||
import static org.junit.jupiter.api.Assertions.*;
|
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<ILoggingEvent> 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<ILoggingEvent> 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
|
@Test
|
||||||
void aWorkerThatBecomesReadyWithinTheGraceIsDeliveredNormally() {
|
void aWorkerThatBecomesReadyWithinTheGraceIsDeliveredNormally() {
|
||||||
// The readiness grace must not fail a worker that is merely slow to boot: once it becomes
|
// The readiness grace must not fail a worker that is merely slow to boot: once it becomes
|
||||||
|
|||||||
Reference in New Issue
Block a user