Compare commits
5 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| ac351ee1de | |||
| 8f59019305 | |||
| e966cbadf9 | |||
| c87cc25aa6 | |||
| 3fb331145a |
@@ -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<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
|
||||
/**
|
||||
* 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<>();
|
||||
|
||||
/** 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,
|
||||
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.router = null;
|
||||
this.turnListener = turnListener;
|
||||
this.ready = ready;
|
||||
this.forget = forget;
|
||||
this.nowMillis = nowMillis;
|
||||
}
|
||||
|
||||
public Injector(HerdrRouter router, TurnListener turnListener, Predicate<String> ready,
|
||||
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.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;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -381,13 +381,18 @@ public final class LeadRollover {
|
||||
*/
|
||||
private void runRollover(PendingRollover p, FleetConfig.LeadRollover cfg) {
|
||||
String lead = p.leadTerminal();
|
||||
boolean turnSettled = waitUntilAtTurnBoundary(lead, cfg.turnSettleSeconds());
|
||||
if (!turnSettled) {
|
||||
log.warn("lead-rollover: pane {} did not reach a turn boundary (IDLE or DONE) within {}s "
|
||||
+ "after confirm() — refusing to send /clear at all; the calling lead's "
|
||||
+ "own turn is still live and clearing it now would destroy live context "
|
||||
+ "(token={})",
|
||||
lead, cfg.turnSettleSeconds(), p.token());
|
||||
long rollStartMillis = nowMillis.getAsLong();
|
||||
TurnSettleResult turnResult = waitUntilAtTurnBoundary(lead, cfg.turnSettleSeconds());
|
||||
if (!turnResult.settled()) {
|
||||
// fleetd #494 follow-up: this line had the SAME defect as the /clear-timeout line below
|
||||
// — cfg.turnSettleSeconds() is the CONFIGURED budget, not how long this wait actually
|
||||
// ran. Print the measured elapsed time alongside it, labelled, exactly like the /clear
|
||||
// path already does.
|
||||
log.warn("lead-rollover: pane {} did not reach a turn boundary (IDLE or DONE) after "
|
||||
+ "confirm() — refusing to send /clear at all; the calling lead's own "
|
||||
+ "turn is still live and clearing it now would destroy live context "
|
||||
+ "(token={}, configured={}s elapsed={}ms)",
|
||||
lead, p.token(), cfg.turnSettleSeconds(), turnResult.elapsedMillis());
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -395,15 +400,22 @@ public final class LeadRollover {
|
||||
// /clear is housekeeping, not a delegated turn, and routing it through Injector wedges the
|
||||
// pane forever (see this class's javadoc).
|
||||
agents.send(lead, "/clear");
|
||||
boolean clearSettled = waitForClearPickupAndSettle(lead, cfg.clearSettleSeconds());
|
||||
if (!clearSettled) {
|
||||
log.warn("lead-rollover: pane {} did not reach a turn boundary (IDLE or DONE) within {}s "
|
||||
+ "after /clear — NOT sending bootstrapText (token={})",
|
||||
lead, cfg.clearSettleSeconds(), p.token());
|
||||
ClearSettleResult clearResult = waitForClearPickupAndSettle(lead, cfg.clearSettleSeconds());
|
||||
if (!clearResult.settled()) {
|
||||
// fleetd #494: cfg.clearSettleSeconds() is the CONFIGURED budget, not how long the wait
|
||||
// actually ran — an operator reading only that number wrongly believes it is a measured
|
||||
// duration. Print the measured elapsed time and nudge count alongside it, each labelled,
|
||||
// so the two can be compared at a glance.
|
||||
log.warn("lead-rollover: pane {} did not reach a turn boundary (IDLE or DONE) after "
|
||||
+ "/clear — NOT sending bootstrapText (token={}, configured={}s "
|
||||
+ "elapsed={}ms nudges={})",
|
||||
lead, p.token(), cfg.clearSettleSeconds(), clearResult.elapsedMillis(),
|
||||
clearResult.nudges());
|
||||
return;
|
||||
}
|
||||
agents.send(lead, cfg.bootstrapTextFor(p.handoverPath()));
|
||||
log.info("lead-rollover: rolled token={} lead={}", p.token(), lead);
|
||||
long rollElapsedMillis = nowMillis.getAsLong() - rollStartMillis;
|
||||
log.info("lead-rollover: rolled token={} lead={} elapsedMs={}", p.token(), lead, rollElapsedMillis);
|
||||
}
|
||||
|
||||
/** Drop a pending request without rolling. @return whether a pending request existed for {@code token} */
|
||||
@@ -475,9 +487,15 @@ public final class LeadRollover {
|
||||
* gate exists to prevent. Do not "simplify" this back to {@code injectable()}. ({@link
|
||||
* #waitForClearPickupAndSettle} keeps the same exclusion of {@code BLOCKED}, for the same
|
||||
* reason, on the second wait.)
|
||||
*
|
||||
* @return a {@link TurnSettleResult} whose {@code settled()} is {@code true} once a real
|
||||
* boundary was observed, {@code false} if {@code settleSeconds} elapses first.
|
||||
* {@code elapsedMillis()} is a MEASURED value from the injected {@link #nowMillis}
|
||||
* clock, never the configured {@code settleSeconds} budget (fleetd #494 follow-up).
|
||||
*/
|
||||
private boolean waitUntilAtTurnBoundary(String target, int settleSeconds) {
|
||||
long deadline = nowMillis.getAsLong() + TimeUnit.SECONDS.toMillis(settleSeconds);
|
||||
private TurnSettleResult waitUntilAtTurnBoundary(String target, int settleSeconds) {
|
||||
long startMillis = nowMillis.getAsLong();
|
||||
long deadline = startMillis + TimeUnit.SECONDS.toMillis(settleSeconds);
|
||||
while (nowMillis.getAsLong() < deadline) {
|
||||
AgentStatus status;
|
||||
try {
|
||||
@@ -488,13 +506,20 @@ public final class LeadRollover {
|
||||
status = null;
|
||||
}
|
||||
if (status == AgentStatus.IDLE || status == AgentStatus.DONE) {
|
||||
return true;
|
||||
return new TurnSettleResult(true, nowMillis.getAsLong() - startMillis);
|
||||
}
|
||||
settleSleeper.run();
|
||||
}
|
||||
return false;
|
||||
return new TurnSettleResult(false, nowMillis.getAsLong() - startMillis);
|
||||
}
|
||||
|
||||
/**
|
||||
* The measured outcome of {@link #waitUntilAtTurnBoundary} — fleetd #494 follow-up. The sibling
|
||||
* of {@link ClearSettleResult} for the FIRST wait, which never nudges, so it carries no nudge
|
||||
* count.
|
||||
*/
|
||||
private record TurnSettleResult(boolean settled, long elapsedMillis) {}
|
||||
|
||||
/**
|
||||
* The SECOND wait in {@link #runRollover} — after {@code /clear} has been sent, waits for it to
|
||||
* settle, bounded by {@code settleSeconds}. <strong>fleetd #489 — the paste-race fix.</strong>
|
||||
@@ -524,8 +549,9 @@ public final class LeadRollover {
|
||||
* Claude Code prompt is a no-op, so repeating it is safe;</li>
|
||||
* <li>the {@code PICKUP_GRACE_POLLS}th consecutive such poll, with {@code WORKING} still never
|
||||
* observed, releases rather than wedges the roll instead of nudging again — the same
|
||||
* choice {@code Injector} makes — and returns {@code true} anyway, logged at {@code info}
|
||||
* so an operator can see which path ran;</li>
|
||||
* choice {@code Injector} makes — and returns {@code settled() == true} anyway, logged at
|
||||
* {@code warn} with the measured elapsed time (fleetd #494) so an operator can see which
|
||||
* path ran and how long it actually took;</li>
|
||||
* <li>once {@code WORKING} has been observed, nudging stops and this instead waits for a real
|
||||
* {@code working → IDLE/DONE} completion boundary before returning {@code true}.</li>
|
||||
* </ul>
|
||||
@@ -541,15 +567,20 @@ public final class LeadRollover {
|
||||
* swallowed and logged at {@code debug}, exactly like {@code Injector.java:437-442} — a failed
|
||||
* nudge must not abort the roll.
|
||||
*
|
||||
* @return {@code true} once {@code /clear} has settled, or once the nudge budget was exhausted
|
||||
* with no pickup ever observed (released rather than wedged); {@code false} if {@code
|
||||
* settleSeconds} elapses first — the caller must NOT send {@code bootstrapText} in that
|
||||
* case, exactly as before this fix
|
||||
* @return a {@link ClearSettleResult} whose {@code settled()} is {@code true} once {@code
|
||||
* /clear} has settled, or once the nudge budget was exhausted with no pickup ever
|
||||
* observed (released rather than wedged); {@code false} if {@code settleSeconds} elapses
|
||||
* first — the caller must NOT send {@code bootstrapText} in that case, exactly as before
|
||||
* this fix. {@code elapsedMillis()} and {@code nudges()} are MEASURED values (from the
|
||||
* injected {@link #nowMillis} clock and an actual nudge count), never the configured
|
||||
* {@code settleSeconds} budget (fleetd #494).
|
||||
*/
|
||||
private boolean waitForClearPickupAndSettle(String target, int settleSeconds) {
|
||||
long deadline = nowMillis.getAsLong() + TimeUnit.SECONDS.toMillis(settleSeconds);
|
||||
private ClearSettleResult waitForClearPickupAndSettle(String target, int settleSeconds) {
|
||||
long startMillis = nowMillis.getAsLong();
|
||||
long deadline = startMillis + TimeUnit.SECONDS.toMillis(settleSeconds);
|
||||
boolean pickedUp = false; // a WORKING sample has been observed since /clear was sent
|
||||
int idlePollsAwaitingPickup = 0;
|
||||
int nudges = 0;
|
||||
while (nowMillis.getAsLong() < deadline) {
|
||||
AgentStatus status;
|
||||
try {
|
||||
@@ -563,26 +594,56 @@ public final class LeadRollover {
|
||||
pickedUp = true;
|
||||
} else if (status == AgentStatus.IDLE || status == AgentStatus.DONE) {
|
||||
if (pickedUp) {
|
||||
return true; // a real WORKING -> IDLE/DONE completion boundary
|
||||
// a real WORKING -> IDLE/DONE completion boundary
|
||||
return new ClearSettleResult(true, nowMillis.getAsLong() - startMillis, nudges);
|
||||
}
|
||||
if (++idlePollsAwaitingPickup >= PICKUP_GRACE_POLLS) {
|
||||
log.info("lead-rollover: /clear on {} was never observed as WORKING after {} "
|
||||
long elapsedMillis = nowMillis.getAsLong() - startMillis;
|
||||
// fleetd #494: this release trades a possibly-unsubmitted /clear for progress
|
||||
// instead of wedging the roll — that trade is deliberate and stays. But it is
|
||||
// also exactly the case that reported false success in the real incident (the
|
||||
// whole roll "succeeded" after 438ms of a 20s budget), so raise it to WARN and
|
||||
// print the MEASURED elapsed time next to the target pane, not just the count.
|
||||
//
|
||||
// fleetd #494 follow-up (2nd pass): BOTH numbers in this line must come from
|
||||
// the loop's own counters, never from the PICKUP_GRACE_POLLS constant.
|
||||
// `idlePollsAwaitingPickup` and `nudges` each have exactly one write site in
|
||||
// this loop, on the same branch, so on this branch they cannot differ from
|
||||
// PICKUP_GRACE_POLLS / PICKUP_GRACE_POLLS - 1 today — no test can prove the
|
||||
// difference on this line, and printing the counters does not change that.
|
||||
// What it does buy: one source of truth instead of two, so a later change to
|
||||
// the loop (an early return, a second increment site, a different exit
|
||||
// condition) cannot leave this message reporting a number the loop no longer
|
||||
// produces. The place where `nudges` genuinely varies with the run — and is
|
||||
// covered by a test that can tell it apart from a constant — is the
|
||||
// /clear-timeout warn in runRollover, which prints clearResult.nudges().
|
||||
log.warn("lead-rollover: /clear on {} was never observed as WORKING after {} "
|
||||
+ "consecutive IDLE/DONE polls ({} of those were nudged) — "
|
||||
+ "releasing rather than wedging the roll",
|
||||
target, PICKUP_GRACE_POLLS, PICKUP_GRACE_POLLS - 1);
|
||||
return true;
|
||||
+ "releasing rather than wedging the roll (elapsed={}ms)",
|
||||
target, idlePollsAwaitingPickup, nudges, elapsedMillis);
|
||||
return new ClearSettleResult(true, elapsedMillis, nudges);
|
||||
}
|
||||
try {
|
||||
agents.submit(target); // nudge a raced Enter (CB-113) so /clear actually submits
|
||||
} catch (RuntimeException e) {
|
||||
log.debug("lead-rollover: resubmit to {} failed (will retry next poll): {}",
|
||||
target, e.getMessage());
|
||||
} finally {
|
||||
nudges++; // an attempted nudge, whether or not the submit call itself threw
|
||||
}
|
||||
}
|
||||
// AgentStatus.BLOCKED or UNKNOWN (or an unreadable status, above): neither a pickup
|
||||
// signal nor a boundary — keep polling without nudging or releasing.
|
||||
settleSleeper.run();
|
||||
}
|
||||
return false;
|
||||
return new ClearSettleResult(false, nowMillis.getAsLong() - startMillis, nudges);
|
||||
}
|
||||
|
||||
/**
|
||||
* The measured outcome of {@link #waitForClearPickupAndSettle} — fleetd #494. Carries the
|
||||
* MEASURED elapsed time (from the injected {@link #nowMillis} clock) and nudge count alongside
|
||||
* the settle/timeout decision, so callers can log them instead of the configured budget, which
|
||||
* is not how long the wait actually ran.
|
||||
*/
|
||||
private record ClearSettleResult(boolean settled, long elapsedMillis, int nudges) {}
|
||||
}
|
||||
|
||||
@@ -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<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
|
||||
void aWorkerThatBecomesReadyWithinTheGraceIsDeliveredNormally() {
|
||||
// The readiness grace must not fail a worker that is merely slow to boot: once it becomes
|
||||
|
||||
@@ -1,5 +1,9 @@
|
||||
package dev.ltms.fleet.lead;
|
||||
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.Logger;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
@@ -9,6 +13,7 @@ import dev.ltms.fleet.herdr.HerdrException;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.nio.file.Files;
|
||||
@@ -854,4 +859,207 @@ class LeadRolloverTest {
|
||||
assertFalse(bootstrapSent.contains("\"handover.md\""),
|
||||
"must not name the raw relative configured value in the text actually sent");
|
||||
}
|
||||
|
||||
// ---- fleetd #494: the log lines must print MEASURED values, never the configured budget --
|
||||
|
||||
private static ListAppender<ILoggingEvent> attachLog() {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(LeadRollover.class);
|
||||
logger.setLevel(Level.DEBUG);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
return appender;
|
||||
}
|
||||
|
||||
private static void detachLog(ListAppender<ILoggingEvent> appender) {
|
||||
((Logger) LoggerFactory.getLogger(LeadRollover.class)).detachAppender(appender);
|
||||
}
|
||||
|
||||
private static ILoggingEvent lastEventContaining(ListAppender<ILoggingEvent> events, String substring) {
|
||||
return events.list.stream()
|
||||
.filter(e -> e.getFormattedMessage().contains(substring))
|
||||
.reduce((_, b) -> b)
|
||||
.orElseThrow(() -> new AssertionError("no log event contained \"" + substring
|
||||
+ "\"; got: " + events.list.stream().map(ILoggingEvent::getFormattedMessage).toList()));
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[fleetd #494] the /clear-timeout warn line prints the MEASURED elapsed time and "
|
||||
+ "nudge count next to the configured budget, never the configured value alone")
|
||||
void clearTimeoutLogPrintsMeasuredElapsedAndNudgesNotJustConfigured() throws IOException {
|
||||
// Idle until /clear is sent, then permanently WORKING (a genuinely stuck /clear that never
|
||||
// reaches a completion boundary) — isolates the SECOND wait exactly like
|
||||
// clearThatNeverSettlesAfterwardsNeverSendsBootstrapText, but with a clock that ADVANCES on
|
||||
// every read so the measured elapsed time is a deterministic, non-zero value distinct from
|
||||
// the configured budget — pinning fleetd #494's fix, not just its absence of a hang.
|
||||
FakeHerdr fake = new FakeHerdr();
|
||||
HerdrClient flipsAfterClear = new HerdrClient() {
|
||||
@Override
|
||||
public JsonNode call(String method, Object params) throws HerdrException {
|
||||
JsonNode result = fake.call(method, params);
|
||||
if ("agent.prompt".equals(method) && String.valueOf(params).contains("/clear")) {
|
||||
fake.agentStatus("working");
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
fake.close();
|
||||
}
|
||||
};
|
||||
Path handover = writeHandover("handover contents");
|
||||
FleetConfig.LeadRollover config =
|
||||
new FleetConfig.LeadRollover(handover.toString(), true, 3600, 20, 1 /*clearSettleSeconds*/, "boot text");
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
LeadRollover rollover = newRollover(flipsAfterClear, config, () -> clock.addAndGet(500));
|
||||
|
||||
ListAppender<ILoggingEvent> events = attachLog();
|
||||
try {
|
||||
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
|
||||
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
|
||||
|
||||
assertTrue(decision.accepted(), "every synchronous gate passes; the refusal is logged "
|
||||
+ "only, deep inside the deferred continuation");
|
||||
|
||||
ILoggingEvent event = lastEventContaining(events, "NOT sending bootstrapText");
|
||||
assertEquals(Level.WARN, event.getLevel());
|
||||
String message = event.getFormattedMessage();
|
||||
assertTrue(message.contains("configured=1s"), "must label the configured budget: " + message);
|
||||
assertTrue(message.contains("elapsed=1500ms"), "must print the MEASURED elapsed time — "
|
||||
+ "with this fixture's advancing clock, the wait actually ran 1500ms against a "
|
||||
+ "1s(=1000ms) configured budget: " + message);
|
||||
assertTrue(message.contains("nudges=0"), "must print the measured nudge count (0 here — "
|
||||
+ "the pane was WORKING throughout, never IDLE/DONE, so no nudge was ever sent): "
|
||||
+ message);
|
||||
assertFalse(message.contains("within 1s"), "must not present the configured budget as if "
|
||||
+ "it were the measured wait duration: " + message);
|
||||
} finally {
|
||||
detachLog(events);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[fleetd #494] the /clear pickup-grace release line is WARN (was INFO) and prints "
|
||||
+ "the measured elapsed time next to the target pane")
|
||||
void clearGraceReleaseLogIsWarnWithMeasuredElapsed() throws IOException {
|
||||
// Default idle throughout — no WORKING sample is ever observed, so the pickup-grace wait
|
||||
// exhausts PICKUP_GRACE_POLLS and releases rather than wedging (see
|
||||
// clearPickupIsNudgedBeforeBootstrapTextWhenPaneStaysIdle for the un-logged half of this
|
||||
// scenario). A self-advancing clock makes the measured elapsed time deterministic and
|
||||
// provably distinct from a bare poll/nudge count.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Path handover = writeHandover("handover contents");
|
||||
FleetConfig.LeadRollover config = cfg(handover.toString()); // turnSettleSeconds=clearSettleSeconds=20
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
LeadRollover rollover = newRollover(herdr, config, () -> clock.addAndGet(500));
|
||||
|
||||
ListAppender<ILoggingEvent> events = attachLog();
|
||||
try {
|
||||
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
|
||||
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
|
||||
|
||||
assertTrue(decision.accepted(), "expected approval; got: " + decision.reason()
|
||||
+ " / " + decision.detail());
|
||||
|
||||
ILoggingEvent event = lastEventContaining(events, "releasing rather than wedging the roll");
|
||||
assertEquals(Level.WARN, event.getLevel(), "the grace-limit release must be WARN, not "
|
||||
+ "INFO — it is exactly the case that reported false success in the real incident "
|
||||
+ "this fix comes from (a roll that 'succeeded' after 438ms of a 20s budget)");
|
||||
String message = event.getFormattedMessage();
|
||||
assertTrue(message.contains(LEAD), "must name the target pane: " + message);
|
||||
assertTrue(message.contains("elapsed=4500ms"), "must print the MEASURED elapsed time — "
|
||||
+ "with this fixture's advancing clock, the wait ran 4500ms before releasing: "
|
||||
+ message);
|
||||
// fleetd #494 follow-up (2nd pass): both numbers here are DELIBERATE plain literals,
|
||||
// not derived from LeadRollover.PICKUP_GRACE_POLLS. A version of this assertion that
|
||||
// reads "(" + (LeadRollover.PICKUP_GRACE_POLLS - 1) + " of those were nudged)" builds
|
||||
// its expectation the same way the production code used to build the log line, so it
|
||||
// cannot tell a fixed constant apart from the measured counter — proved by reverting
|
||||
// the production fix and re-running: that mutant stayed green under the old assertion.
|
||||
// If PICKUP_GRACE_POLLS ever changes, THIS TEST MUST FAIL and a human must look at the
|
||||
// new message and update the literals below, not just re-derive them.
|
||||
assertTrue(message.contains("after 8 consecutive IDLE/DONE polls (7 of those were nudged)"),
|
||||
"must print the measured poll count and nudge count as plain numbers, not the "
|
||||
+ "PICKUP_GRACE_POLLS constant standing in for either: " + message);
|
||||
} finally {
|
||||
detachLog(events);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[fleetd #494] the success line prints the measured elapsed time for the whole roll")
|
||||
void successLogPrintsMeasuredElapsedForTheWholeRoll() throws IOException {
|
||||
// Same fixture as clearGraceReleaseLogIsWarnWithMeasuredElapsed: default idle throughout, so
|
||||
// the grace release fires and the roll still goes on to send bootstrapText and log success.
|
||||
// This is deliberately the SAME shape as the real incident (a roll that "succeeds" quickly)
|
||||
// — the missing signal was the elapsed time on this exact line.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Path handover = writeHandover("handover contents");
|
||||
FleetConfig.LeadRollover config = cfg(handover.toString());
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
LeadRollover rollover = newRollover(herdr, config, () -> clock.addAndGet(500));
|
||||
|
||||
ListAppender<ILoggingEvent> events = attachLog();
|
||||
try {
|
||||
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
|
||||
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
|
||||
|
||||
assertTrue(decision.accepted(), "expected approval; got: " + decision.reason()
|
||||
+ " / " + decision.detail());
|
||||
|
||||
ILoggingEvent event = lastEventContaining(events, "lead-rollover: rolled");
|
||||
assertEquals(Level.INFO, event.getLevel());
|
||||
String message = event.getFormattedMessage();
|
||||
// fleetd #494 follow-up: waitUntilAtTurnBoundary now also reads the injected clock one
|
||||
// extra time (to compute ITS OWN measured elapsed on the success path), so the shared
|
||||
// fixture clock advances by one more 500ms tick before the roll finishes than it did
|
||||
// before that follow-up — 7000ms, not 6500ms.
|
||||
assertTrue(message.contains("elapsedMs=7000"), "must print the MEASURED elapsed time for "
|
||||
+ "the whole roll — with this fixture's advancing clock, the full roll (turn-settle "
|
||||
+ "wait + /clear wait + bootstrapText) took 7000ms: " + message);
|
||||
} finally {
|
||||
detachLog(events);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[fleetd #494 follow-up] the turn-settle timeout warn line prints the MEASURED "
|
||||
+ "elapsed time next to the configured budget, never the configured value alone")
|
||||
void turnTimeoutLogPrintsMeasuredElapsedNotJustConfigured() throws IOException {
|
||||
// The brief that named the four items this ticket fixed left this exact sibling line out —
|
||||
// "did not reach a turn boundary (IDLE or DONE) within {}s after confirm()" — even though it
|
||||
// has the identical defect shape one method below. Same fixture shape as
|
||||
// clearTimeoutLogPrintsMeasuredElapsedAndNudgesNotJustConfigured, but for the FIRST wait: the
|
||||
// calling lead's own pane never goes idle, so waitUntilAtTurnBoundary times out.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
herdr.agentStatus("working"); // the calling lead's own pane — never goes idle in this test
|
||||
Path handover = writeHandover("handover contents");
|
||||
FleetConfig.LeadRollover config =
|
||||
new FleetConfig.LeadRollover(handover.toString(), true, 3600, 1 /*turnSettleSeconds*/, 20, "text");
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
LeadRollover rollover = newRollover(herdr, config, () -> clock.addAndGet(500));
|
||||
|
||||
ListAppender<ILoggingEvent> events = attachLog();
|
||||
try {
|
||||
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
|
||||
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
|
||||
|
||||
assertTrue(decision.accepted(), "every synchronous gate should pass; the refusal happens "
|
||||
+ "only inside the deferred continuation, which this test's synchronous runner "
|
||||
+ "has already run to completion by the time confirm() returns");
|
||||
|
||||
ILoggingEvent event = lastEventContaining(events, "refusing to send /clear at all");
|
||||
assertEquals(Level.WARN, event.getLevel());
|
||||
String message = event.getFormattedMessage();
|
||||
assertTrue(message.contains("configured=1s"), "must label the configured budget: " + message);
|
||||
assertTrue(message.contains("elapsed=1500ms"), "must print the MEASURED elapsed time — "
|
||||
+ "with this fixture's advancing clock, the wait actually ran 1500ms against a "
|
||||
+ "1s(=1000ms) configured budget: " + message);
|
||||
assertFalse(message.contains("within 1s"), "must not present the configured budget as if "
|
||||
+ "it were the measured wait duration: " + message);
|
||||
} finally {
|
||||
detachLog(events);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+55
-218
@@ -5,7 +5,7 @@
|
||||
# A merge is not a deployment: the running daemon holds the jar it was started with, so code merged
|
||||
# to main does nothing until this runs. See CLAUDE.md -> "Redeploying the daemon".
|
||||
#
|
||||
# This script exists to turn eight remembered traps into one auditable command:
|
||||
# This script exists to turn six remembered traps into one auditable command:
|
||||
#
|
||||
# 1. A piped `mvn` hides BUILD FAILURE behind a zero exit, so the build here is never piped.
|
||||
# 2. The daemon must start from a LOGIN shell, or the tokens it hands to members are empty:
|
||||
@@ -27,17 +27,6 @@
|
||||
# restart of the OLD jar. So this script detects whether the agent is loaded and, only then,
|
||||
# swaps `kill` + manual `nohup` for `launchctl unload`/`load` — the one supervisor in control
|
||||
# at any moment is whichever one you asked to act, never both.
|
||||
# 7. fleetd #492 — a systemd --user unit is a THIRD possible supervisor (seen on a second host):
|
||||
# Restart=on-failure treats this JVM's SIGTERM exit code (143, per CB-594 above) as a failure
|
||||
# too, so a bare `kill` there would race systemd's own restart of the OLD jar exactly like
|
||||
# launchd would. This script now tells launchd, systemd, and "genuinely unsupervised" apart as
|
||||
# three different answers, drives whichever one it finds through its own control plane
|
||||
# (`launchctl` / `systemctl --user`), and REFUSES outright — never falls back to `kill` — when
|
||||
# it finds a supervision signal it cannot map to exactly one of the two it knows how to drive.
|
||||
# A wrong guess here is how two daemons end up running against one herdr session.
|
||||
# 8. fleetd #492 — a post-restart check counts running fleetd processes and fails the whole run if
|
||||
# more than one is alive. That is the one thing none of the checks above (healthz 200, jar id,
|
||||
# the fresh "listening" line) can see: every one of them is satisfied by EITHER daemon.
|
||||
#
|
||||
# Usage:
|
||||
# scripts/redeploy-fleetd.sh # build, confirm, restart, verify
|
||||
@@ -66,18 +55,13 @@ HEALTH_WAIT=60 # seconds to wait for /healthz to answer after start
|
||||
LAUNCHD_LABEL='dev.ltms.fleetd'
|
||||
LAUNCHD_PLIST="$HOME/Library/LaunchAgents/$LAUNCHD_LABEL.plist"
|
||||
|
||||
# fleetd #492: the systemd --user unit this script must not fight with either (see trap 7 above).
|
||||
# Measured on the second host: `systemctl --user cat fleetd` names the unit "fleetd" (not
|
||||
# "dev.ltms.fleetd" — systemd user units here are not namespaced the way the launchd label is).
|
||||
SYSTEMD_UNIT='fleetd'
|
||||
|
||||
DO_BUILD=1; ASSUME_YES=0; CHECK_ONLY=0
|
||||
for arg in "$@"; do
|
||||
case "$arg" in
|
||||
--yes|-y) ASSUME_YES=1 ;;
|
||||
--no-build) DO_BUILD=0 ;;
|
||||
--check) CHECK_ONLY=1 ;;
|
||||
-h|--help) sed -n '3,48p' "${BASH_SOURCE[0]}"; exit 0 ;;
|
||||
-h|--help) sed -n '3,37p' "${BASH_SOURCE[0]}"; exit 0 ;;
|
||||
*) echo "unknown option: $arg (try --help)" >&2; exit 2 ;;
|
||||
esac
|
||||
done
|
||||
@@ -95,91 +79,6 @@ running_pid() { pgrep -f "$PATTERN" || true; }
|
||||
launchd_installed() { [ -f "$LAUNCHD_PLIST" ]; }
|
||||
launchd_loaded() { launchctl list "$LAUNCHD_LABEL" >/dev/null 2>&1; }
|
||||
|
||||
# fleetd #492: same two questions for systemd --user. Kept as separate, overridable functions
|
||||
# (never an inline `systemctl` call at each use site) so a test on a box with no systemd at all
|
||||
# (this repo is developed on macOS) can substitute each one independently — the same seam
|
||||
# launchd_installed/launchd_loaded above already use.
|
||||
#
|
||||
# "installed": a unit FILE by this name exists, regardless of its current state — the systemd
|
||||
# analogue of the plist file existing on disk. `list-unit-files` reads unit definitions without
|
||||
# depending on runtime state, so this stays read-only and safe under --check.
|
||||
systemd_installed() {
|
||||
command -v systemctl >/dev/null 2>&1 \
|
||||
&& systemctl --user list-unit-files "$SYSTEMD_UNIT.service" --no-legend 2>/dev/null | grep -q .
|
||||
}
|
||||
# "loaded": systemd currently supervises this unit as an active job — the systemd analogue of
|
||||
# `launchctl list <label>` succeeding. Measured on the second host: `systemctl --user is-active
|
||||
# fleetd` -> "active".
|
||||
systemd_loaded() {
|
||||
command -v systemctl >/dev/null 2>&1 && systemctl --user is-active "$SYSTEMD_UNIT" >/dev/null 2>&1
|
||||
}
|
||||
|
||||
# fleetd #492: three real answers, not two — launchd, systemd, or genuinely unsupervised — plus a
|
||||
# fourth, "ambiguous", for the one case this script cannot tell apart: both signals firing at once.
|
||||
# That is exactly "I cannot tell who supervises this process", and guessing wrong here is how two
|
||||
# daemons end up running against one herdr session (see trap 7 in the header). Pure and
|
||||
# side-effect-free: reads the two probes above and decides — never mutates anything, so it is safe
|
||||
# under --check and testable by overriding launchd_loaded/systemd_loaded after sourcing.
|
||||
detect_supervisor() {
|
||||
local ld=0 sd=0
|
||||
launchd_loaded && ld=1
|
||||
systemd_loaded && sd=1
|
||||
if [ "$ld" = 1 ] && [ "$sd" = 1 ]; then
|
||||
echo "ambiguous"
|
||||
elif [ "$ld" = 1 ]; then
|
||||
echo "launchd"
|
||||
elif [ "$sd" = 1 ]; then
|
||||
echo "systemd"
|
||||
else
|
||||
echo "none"
|
||||
fi
|
||||
}
|
||||
|
||||
# fleetd #492: turns anything detect_supervisor returns that is NOT exactly one of the two
|
||||
# supervisors this script knows how to drive into a die() — never a fall-through to the `kill`
|
||||
# path. Kept as its own function so a test can call it directly (in a subshell, since it die()s)
|
||||
# without running the whole report-state flow or needing a real launchd/systemd.
|
||||
require_drivable_supervisor() {
|
||||
local kind="$1"
|
||||
case "$kind" in
|
||||
launchd|systemd|none) ;;
|
||||
ambiguous)
|
||||
die "both launchd ($LAUNCHD_LABEL) and systemd --user ($SYSTEMD_UNIT) report themselves as
|
||||
loaded for this daemon at the same time. This script cannot tell which one actually
|
||||
supervises the running process, and driving either alone risks the OTHER reviving the
|
||||
OLD jar out from under it — the exact failure this ticket (fleetd #492) exists to
|
||||
prevent. Stop one of the two supervisors by hand, confirm only one remains loaded, then
|
||||
rerun." ;;
|
||||
*)
|
||||
die "detect_supervisor returned an unrecognized value '$kind' — refusing to guess which
|
||||
supervisor, if any, controls this daemon." ;;
|
||||
esac
|
||||
}
|
||||
|
||||
# fleetd #492: the exact symptom a racing supervisor produces — count how many fleetd processes are
|
||||
# alive right now. Takes the pid list as a parameter (rather than calling running_pid() itself) so a
|
||||
# test can pass a canned two-line string without a real second process running. Pure except for the
|
||||
# die() in assert_single_daemon below.
|
||||
count_daemon_pids() {
|
||||
local pids="$1"
|
||||
if [ -z "$pids" ]; then
|
||||
echo 0
|
||||
else
|
||||
printf '%s\n' "$pids" | grep -c .
|
||||
fi
|
||||
}
|
||||
assert_single_daemon() {
|
||||
local pids="$1" count
|
||||
count="$(count_daemon_pids "$pids")"
|
||||
if [ "$count" -gt 1 ]; then
|
||||
die "more than one fleetd process is running after this restart (pids: $(printf '%s' "$pids" | tr '\n' ' ')).
|
||||
This is the exact failure a racing supervisor produces: the OLD jar was revived by its
|
||||
supervisor while this script started a NEW copy. Two daemons on one herdr session kill
|
||||
each other's members. Investigate with 'pgrep -f \"$PATTERN\"' and stop the wrong one by
|
||||
hand — do not assume either pid is the one you want."
|
||||
fi
|
||||
}
|
||||
|
||||
# CB-600: the script computes its own log path from where it sits on disk (REPO, above); the
|
||||
# plist hard-codes an absolute StandardOutPath. Nothing forced the two to agree — if this script
|
||||
# were ever run from a checkout other than the one the loaded plist names, launchd would start and
|
||||
@@ -283,44 +182,24 @@ fi
|
||||
ok "jar on disk: $(jar_id) ($([ -f "$JAR" ] && date -r "$JAR" '+%Y-%m-%d %H:%M:%S' || echo 'none'))"
|
||||
ok "HEAD: $(git -C "$REPO" log --oneline -1)"
|
||||
|
||||
# CB-594 / fleetd #492: supervision state. Installed and loaded are different facts — a
|
||||
# copied-but-never-loaded plist (or an unloaded systemd unit) supervises nothing, and a loaded
|
||||
# label/unit with no file backing it is still what its supervisor will act on.
|
||||
# CB-594: supervision state. Installed and loaded are different facts — a copied-but-never-loaded
|
||||
# plist supervises nothing, and a loaded label with no file backing it (rare, but possible after an
|
||||
# edited/moved plist) is still what launchd will act on.
|
||||
if launchd_installed; then
|
||||
ok "launchd agent installed: $LAUNCHD_PLIST"
|
||||
else
|
||||
warn "launchd agent NOT installed."
|
||||
warn "launchd agent NOT installed (no supervision — a crash will not restart the daemon)."
|
||||
fi
|
||||
if systemd_installed; then
|
||||
ok "systemd --user unit installed: $SYSTEMD_UNIT"
|
||||
else
|
||||
warn "systemd --user unit NOT installed ($SYSTEMD_UNIT)."
|
||||
fi
|
||||
|
||||
# fleetd #492: decide which of the two (if either) actually supervises this daemon, and refuse
|
||||
# outright — before touching anything — if that cannot be told apart (see require_drivable_
|
||||
# supervisor above). --check reaches this same line, so a host with an undrivable supervisor is
|
||||
# reported as a failure even in --check, without ever reaching the build/stop/start steps.
|
||||
SUPERVISOR_KIND="$(detect_supervisor)"
|
||||
require_drivable_supervisor "$SUPERVISOR_KIND"
|
||||
ok "supervisor detected: $SUPERVISOR_KIND"
|
||||
SUPERVISED=0
|
||||
case "$SUPERVISOR_KIND" in
|
||||
launchd)
|
||||
SUPERVISED=1
|
||||
ok "launchd agent loaded ($LAUNCHD_LABEL) — launchd supervises this daemon"
|
||||
# CB-600: fail loudly here, before ANY other check runs, if this script and the loaded plist
|
||||
# would read different log files — every check after this point is worthless otherwise.
|
||||
check_log_path_matches_plist "$OUT" "$LAUNCHD_PLIST"
|
||||
;;
|
||||
systemd)
|
||||
SUPERVISED=1
|
||||
ok "systemd --user unit active ($SYSTEMD_UNIT) — systemd supervises this daemon"
|
||||
;;
|
||||
none)
|
||||
warn "no supervisor loaded — this script is the only thing that will restart the daemon."
|
||||
;;
|
||||
esac
|
||||
if launchd_loaded; then
|
||||
SUPERVISED=1
|
||||
ok "launchd agent loaded ($LAUNCHD_LABEL) — launchd supervises this daemon"
|
||||
# CB-600: fail loudly here, before ANY other check runs, if this script and the loaded plist
|
||||
# would read different log files — every check after this point is worthless otherwise.
|
||||
check_log_path_matches_plist "$OUT" "$LAUNCHD_PLIST"
|
||||
else
|
||||
warn "launchd agent not loaded — this script is the only thing that will restart the daemon."
|
||||
fi
|
||||
|
||||
# The trap with no log line. Checked in a LOGIN shell, because that is how the daemon is started
|
||||
# below. Never prints the value — only whether it resolved.
|
||||
@@ -402,39 +281,25 @@ fi
|
||||
|
||||
# ------------------------------------------------------------------ stop
|
||||
#
|
||||
# CB-594 / fleetd #492: when SUPERVISED, the supervisor owns the stop — never a raw `kill` here. A
|
||||
# bare SIGTERM makes this JVM exit 143 even with its shutdown hook running to completion (verified
|
||||
# separately: a throwaway Java process with an equivalent shutdown hook, sent SIGTERM from a login
|
||||
# shell that could `wait` on it directly, reported exit code 143 every time — never 0). launchd's
|
||||
# KeepAlive.SuccessfulExit=false and systemd's Restart=on-failure both treat any nonzero exit as a
|
||||
# crash and restart the OLD jar, which would race this script's own restart of the NEW one.
|
||||
# `launchctl unload` avoids that race by deregistering the job first, so no KeepAlive is left
|
||||
# armed when the process actually stops. `systemctl --user stop` needs no such dance: unlike
|
||||
# KeepAlive, systemd's Restart= does not fire on a deliberate stop, only on an unexpected exit of
|
||||
# an active unit.
|
||||
# CB-594: when SUPERVISED, launchd owns the stop — never a raw `kill` here. A bare SIGTERM makes
|
||||
# this JVM exit 143 even with its shutdown hook running to completion (verified separately: a
|
||||
# throwaway Java process with an equivalent shutdown hook, sent SIGTERM from a login shell that
|
||||
# could `wait` on it directly, reported exit code 143 every time — never 0). launchd's
|
||||
# KeepAlive.SuccessfulExit=false treats any nonzero exit as a crash and restarts the OLD jar,
|
||||
# which would race this script's own restart of the NEW one. `launchctl unload` avoids that race
|
||||
# by deregistering the job first, so no KeepAlive is left armed when the process actually stops.
|
||||
|
||||
if [ -n "$OLD_PID" ]; then
|
||||
say "stop"
|
||||
RESTART_MARK="$(wc -l < "$OUT" 2>/dev/null || echo 0)" # verify a FRESH line appears later
|
||||
case "$SUPERVISOR_KIND" in
|
||||
launchd)
|
||||
echo " supervision is ON (launchd): using 'launchctl unload' (not kill) so launchd's own"
|
||||
echo " KeepAlive cannot restart the OLD jar out from under this script — see the CB-594"
|
||||
echo " comment above."
|
||||
launchctl unload -w "$LAUNCHD_PLIST" \
|
||||
|| die "launchctl unload failed — the daemon may still be under supervision; investigate before retrying"
|
||||
;;
|
||||
systemd)
|
||||
echo " supervision is ON (systemd --user): using 'systemctl --user stop' (not kill) so"
|
||||
echo " systemd's own Restart=on-failure cannot restart the OLD jar out from under this"
|
||||
echo " script — see the fleetd #492 comment above."
|
||||
systemctl --user stop "$SYSTEMD_UNIT" \
|
||||
|| die "'systemctl --user stop $SYSTEMD_UNIT' failed — the daemon may still be under supervision; investigate before retrying"
|
||||
;;
|
||||
none)
|
||||
kill "$OLD_PID"
|
||||
;;
|
||||
esac
|
||||
if [ "$SUPERVISED" = 1 ]; then
|
||||
echo " supervision is ON: using 'launchctl unload' (not kill) so launchd's own KeepAlive"
|
||||
echo " cannot restart the OLD jar out from under this script — see the CB-594 comment above."
|
||||
launchctl unload -w "$LAUNCHD_PLIST" \
|
||||
|| die "launchctl unload failed — the daemon may still be under supervision; investigate before retrying"
|
||||
else
|
||||
kill "$OLD_PID"
|
||||
fi
|
||||
for _ in $(seq "$STOP_WAIT"); do
|
||||
[ -z "$(running_pid)" ] && break
|
||||
sleep 1
|
||||
@@ -445,20 +310,13 @@ if [ -n "$OLD_PID" ]; then
|
||||
leave worktrees and panes behind. Investigate, then kill -9 by hand if you accept that."
|
||||
fi
|
||||
ok "pid $OLD_PID exited"
|
||||
elif [ "$SUPERVISOR_KIND" = "launchd" ]; then
|
||||
elif [ "$SUPERVISED" = 1 ]; then
|
||||
# Loaded but not currently running (e.g. throttled after a crash loop). Unload it anyway so the
|
||||
# start step below does a clean load, never a load stacked on an already-loaded label.
|
||||
say "stop"
|
||||
RESTART_MARK="$(wc -l < "$OUT" 2>/dev/null || echo 0)"
|
||||
launchctl unload -w "$LAUNCHD_PLIST" 2>/dev/null || true
|
||||
ok "launchd agent unloaded (was already not running)"
|
||||
elif [ "$SUPERVISOR_KIND" = "systemd" ]; then
|
||||
# Same case for systemd: the unit is known/active-capable but not currently running. `stop` on an
|
||||
# already-stopped unit is a harmless no-op — kept for symmetry with the launchd branch above.
|
||||
say "stop"
|
||||
RESTART_MARK="$(wc -l < "$OUT" 2>/dev/null || echo 0)"
|
||||
systemctl --user stop "$SYSTEMD_UNIT" 2>/dev/null || true
|
||||
ok "systemd --user unit stopped (was already not running)"
|
||||
else
|
||||
RESTART_MARK="$(wc -l < "$OUT" 2>/dev/null || echo 0)"
|
||||
fi
|
||||
@@ -466,48 +324,34 @@ fi
|
||||
# ------------------------------------------------------------------ start
|
||||
# Unsupervised: login shell (zsh -l) is what puts the secrets on the daemon's environment, and cwd
|
||||
# must be fleetd/ because the daemon resolves fleetd.yaml, logs/ and target/ relative to it.
|
||||
# Supervised (launchd): launchd does both — deploy/dev.ltms.fleetd.plist points ProgramArguments at
|
||||
# Supervised: launchd does both — deploy/dev.ltms.fleetd.plist points ProgramArguments at
|
||||
# scripts/fleetd-launchd-wrapper.sh (CB-594), which is what execs the login shell in launchd's
|
||||
# place, and WorkingDirectory in the plist already pins fleetd/.
|
||||
# Supervised (systemd --user): the unit does both too — measured on the second host, ExecStart is
|
||||
# `/bin/zsh -lc "exec java -jar target/fleetd.jar fleetd.yaml"` (a login shell, same reason as
|
||||
# above) and WorkingDirectory is already pinned to fleetd/.
|
||||
|
||||
say "start"
|
||||
case "$SUPERVISOR_KIND" in
|
||||
launchd)
|
||||
echo " supervision is ON (launchd): using 'launchctl load' so launchd starts and keeps"
|
||||
echo " supervising this process, instead of a manual nohup that launchd would know nothing"
|
||||
echo " about."
|
||||
# CB-600: 'launchctl unload -w' above already persisted Disabled=true for this label. A load -w
|
||||
# that succeeds clears it; a load -w that FAILS leaves the agent both stopped and disabled — worse
|
||||
# than before this script ran, because a later reboot or login will not bring it back either. One
|
||||
# retry covers a transient race (e.g. launchd not yet fully done deregistering); if it still fails,
|
||||
# die with the exact recovery command rather than a bare "failed".
|
||||
if ! launchctl load -w "$LAUNCHD_PLIST" 2>/dev/null; then
|
||||
warn "launchctl load failed on the first attempt — retrying once after a short pause"
|
||||
sleep 2
|
||||
launchctl load -w "$LAUNCHD_PLIST" || die "launchctl load failed twice.
|
||||
The agent is now STOPPED and DISABLED — it will NOT come back on its own, not even after a
|
||||
reboot or login, because 'launchctl unload -w' above persisted Disabled=true and load -w
|
||||
never got the chance to clear it. Recover with:
|
||||
launchctl load -w \"$LAUNCHD_PLIST\"
|
||||
If that still fails, check 'launchctl list $LAUNCHD_LABEL', validate the plist with
|
||||
'plutil -lint \"$LAUNCHD_PLIST\"', and check $OUT before assuming a retry will succeed."
|
||||
fi
|
||||
;;
|
||||
systemd)
|
||||
echo " supervision is ON (systemd --user): using 'systemctl --user start' so systemd starts"
|
||||
echo " and keeps supervising this process, instead of a manual nohup it would know nothing"
|
||||
echo " about."
|
||||
systemctl --user start "$SYSTEMD_UNIT" || die "'systemctl --user start $SYSTEMD_UNIT' failed.
|
||||
Check 'systemctl --user status $SYSTEMD_UNIT' and $OUT before assuming a retry will succeed."
|
||||
;;
|
||||
none)
|
||||
# Absolute jar path so `ps` names which checkout is running.
|
||||
( cd "$MODULE" && zsh -lc "nohup java -jar '$JAR' >> fleetd.out 2>&1 &" )
|
||||
;;
|
||||
esac
|
||||
if [ "$SUPERVISED" = 1 ]; then
|
||||
echo " supervision is ON: using 'launchctl load' so launchd starts and keeps supervising this"
|
||||
echo " process, instead of a manual nohup that launchd would know nothing about."
|
||||
# CB-600: 'launchctl unload -w' above already persisted Disabled=true for this label. A load -w
|
||||
# that succeeds clears it; a load -w that FAILS leaves the agent both stopped and disabled — worse
|
||||
# than before this script ran, because a later reboot or login will not bring it back either. One
|
||||
# retry covers a transient race (e.g. launchd not yet fully done deregistering); if it still fails,
|
||||
# die with the exact recovery command rather than a bare "failed".
|
||||
if ! launchctl load -w "$LAUNCHD_PLIST" 2>/dev/null; then
|
||||
warn "launchctl load failed on the first attempt — retrying once after a short pause"
|
||||
sleep 2
|
||||
launchctl load -w "$LAUNCHD_PLIST" || die "launchctl load failed twice.
|
||||
The agent is now STOPPED and DISABLED — it will NOT come back on its own, not even after a
|
||||
reboot or login, because 'launchctl unload -w' above persisted Disabled=true and load -w
|
||||
never got the chance to clear it. Recover with:
|
||||
launchctl load -w \"$LAUNCHD_PLIST\"
|
||||
If that still fails, check 'launchctl list $LAUNCHD_LABEL', validate the plist with
|
||||
'plutil -lint \"$LAUNCHD_PLIST\"', and check $OUT before assuming a retry will succeed."
|
||||
fi
|
||||
else
|
||||
# Absolute jar path so `ps` names which checkout is running.
|
||||
( cd "$MODULE" && zsh -lc "nohup java -jar '$JAR' >> fleetd.out 2>&1 &" )
|
||||
fi
|
||||
|
||||
for _ in $(seq 10); do
|
||||
NEW_PID="$(running_pid)"
|
||||
@@ -567,13 +411,6 @@ FRESH_LOG="$(mktemp -t fleetd-fresh-log)"
|
||||
trap 'rm -f "$FRESH_LOG"' EXIT
|
||||
tail -n "+$((RESTART_MARK + 1))" "$OUT" > "$FRESH_LOG" 2>/dev/null || true
|
||||
classify_amqp_connection_errors "$FRESH_LOG"
|
||||
|
||||
# fleetd #492: checked here, after healthz and the fresh-log check have both had time to run, so a
|
||||
# supervisor that revives the OLD jar a few seconds late is caught too. Every check above (healthz
|
||||
# 200, jar id, the fresh 'listening' line) is satisfied by EITHER daemon if two are alive — this is
|
||||
# the only one that can tell.
|
||||
assert_single_daemon "$(running_pid)"
|
||||
|
||||
say "result"
|
||||
ok "pid $NEW_PID, jar $(jar_id)"
|
||||
if [ "$REDEPLOY_ERROR_COUNT" -eq 0 ]; then
|
||||
|
||||
@@ -25,78 +25,6 @@ classify_fixture() {
|
||||
classify_amqp_connection_errors "$TMP/$name"
|
||||
}
|
||||
|
||||
|
||||
# fleetd #492 — supervisor detection. Detect_supervisor() reads launchd_loaded/systemd_loaded, so
|
||||
# each test overrides BOTH pairs (installed + loaded) explicitly, rather than relying on either
|
||||
# being naturally absent: this machine may itself be running a real fleetd under launchd right now
|
||||
# (see CLAUDE.md/MEMORY.md — launchd supervision has been live here since 2026-08-26), so leaving
|
||||
# launchd_loaded unmocked in a "systemd only" test would silently read this host's own live state
|
||||
# instead of the fixture.
|
||||
test_detect_supervisor_launchd_only() {
|
||||
launchd_installed() { return 0; }
|
||||
launchd_loaded() { return 0; }
|
||||
systemd_installed() { return 1; }
|
||||
systemd_loaded() { return 1; }
|
||||
assert_equals "launchd" "$(detect_supervisor)" "launchd-only detection"
|
||||
}
|
||||
|
||||
test_detect_supervisor_systemd_only() {
|
||||
launchd_installed() { return 1; }
|
||||
launchd_loaded() { return 1; }
|
||||
systemd_installed() { return 0; }
|
||||
systemd_loaded() { return 0; }
|
||||
assert_equals "systemd" "$(detect_supervisor)" "systemd-only detection"
|
||||
}
|
||||
|
||||
test_detect_supervisor_none() {
|
||||
launchd_installed() { return 1; }
|
||||
launchd_loaded() { return 1; }
|
||||
systemd_installed() { return 1; }
|
||||
systemd_loaded() { return 1; }
|
||||
assert_equals "none" "$(detect_supervisor)" "unsupervised detection"
|
||||
}
|
||||
|
||||
# The heart of the ticket: a supervisor this script cannot drive must refuse, never fall through to
|
||||
# `kill`. require_drivable_supervisor die()s, so it is invoked inside a command substitution — that
|
||||
# forks a subshell, so its exit() only ends the subshell and this test script keeps running under
|
||||
# `set -e`.
|
||||
test_require_drivable_supervisor_refuses_ambiguous() {
|
||||
local output rc=0
|
||||
output="$(require_drivable_supervisor "ambiguous" 2>&1)" || rc=$?
|
||||
[ "$rc" -ne 0 ] || fail "require_drivable_supervisor accepted an ambiguous (undrivable) supervisor"
|
||||
printf '%s' "$output" | grep -qF "$LAUNCHD_LABEL" \
|
||||
|| fail "refusal message does not name the launchd label it found"
|
||||
printf '%s' "$output" | grep -qF "$SYSTEMD_UNIT" \
|
||||
|| fail "refusal message does not name the systemd unit it found"
|
||||
}
|
||||
|
||||
test_require_drivable_supervisor_accepts_known_kinds() {
|
||||
require_drivable_supervisor "launchd" || fail "refused a drivable launchd supervisor"
|
||||
require_drivable_supervisor "systemd" || fail "refused a drivable systemd supervisor"
|
||||
require_drivable_supervisor "none" || fail "refused the unsupervised case"
|
||||
}
|
||||
|
||||
# fleetd #492 — the one-daemon check. Two live pids is the exact symptom a racing supervisor
|
||||
# produces, and none of the other post-restart checks (healthz, jar id, the fresh log line) can see
|
||||
# it because either daemon alone satisfies them.
|
||||
test_count_daemon_pids() {
|
||||
assert_equals 0 "$(count_daemon_pids "")" "count of an empty pid list"
|
||||
assert_equals 1 "$(count_daemon_pids "4242")" "count of a single pid"
|
||||
assert_equals 2 "$(count_daemon_pids "$(printf '4242\n4343\n')")" "count of two pids"
|
||||
}
|
||||
|
||||
test_assert_single_daemon_accepts_one_pid() {
|
||||
assert_single_daemon "4242" || fail "assert_single_daemon rejected a single running pid"
|
||||
}
|
||||
|
||||
test_assert_single_daemon_rejects_two_pids() {
|
||||
local output rc=0
|
||||
output="$(assert_single_daemon "$(printf '4242\n4343\n')" 2>&1)" || rc=$?
|
||||
[ "$rc" -ne 0 ] || fail "assert_single_daemon accepted two simultaneously running pids"
|
||||
printf '%s' "$output" | grep -qF '4242' || fail "refusal message does not list the pids it found"
|
||||
printf '%s' "$output" | grep -qF '4343' || fail "refusal message does not list the pids it found"
|
||||
}
|
||||
|
||||
test_no_errors() {
|
||||
cat > "$TMP/no-errors.log" <<'LOG'
|
||||
2026-09-05 12:00:00 INFO fleetd listening
|
||||
@@ -296,14 +224,6 @@ test_unattributable_quiet_mutation_is_caught() {
|
||||
printf 'Unattributable mutation: FAIL: cross-unattributable recovered: expected 0, got 2\n'
|
||||
}
|
||||
|
||||
test_detect_supervisor_launchd_only
|
||||
test_detect_supervisor_systemd_only
|
||||
test_detect_supervisor_none
|
||||
test_require_drivable_supervisor_refuses_ambiguous
|
||||
test_require_drivable_supervisor_accepts_known_kinds
|
||||
test_count_daemon_pids
|
||||
test_assert_single_daemon_accepts_one_pid
|
||||
test_assert_single_daemon_rejects_two_pids
|
||||
test_no_errors
|
||||
test_recovery_patterns_match_source
|
||||
test_attributed_recovered_connection_error
|
||||
|
||||
Reference in New Issue
Block a user