Compare commits
29 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 8b4320ed24 | |||
| d4a51c6274 | |||
| 93a9ed3f83 | |||
| 0b032f5a1a | |||
| fad99c4c5e | |||
| 87871eaefb | |||
| a476a14f1c | |||
| 5eb4267a4a | |||
| 7611b69667 | |||
| cc302fe4af | |||
| 4a8a780274 | |||
| 343ce0f4c0 | |||
| cec3e191d4 | |||
| 1850a5f324 | |||
| e4eb3dbed4 | |||
| 57cd96f5e6 | |||
| 202e37e3b3 | |||
| 90253f832d | |||
| f1640f5dcc | |||
| c7903c1efe | |||
| 7d711942fe | |||
| bec87f987c | |||
| 7f8a8829f9 | |||
| af9589783e | |||
| 190436c9cf | |||
| 8ea5c2bb1f | |||
| 8335b12562 | |||
| d25c863118 | |||
| 7c34e8f4f9 |
@@ -39,6 +39,10 @@ jobs:
|
||||
# that runs does so against the fake UDS herdr and fake ccs/claude stubs.
|
||||
run: mvn -B clean install
|
||||
|
||||
- name: javadoc reference lint
|
||||
working-directory: fleetd
|
||||
run: mvn -B -DskipTests javadoc:javadoc -Ddoclint=reference
|
||||
|
||||
# Deliberately NOT actions/upload-artifact: this Gitea instance presents as GHES, and
|
||||
# @actions/artifact v2+ (i.e. upload-artifact@v4) refuses to run there —
|
||||
# "GHESNotSupportedError ... not currently supported on GHES", which red-Xes an otherwise
|
||||
|
||||
@@ -696,11 +696,7 @@ public final class Fleetd {
|
||||
FleetMcp mcp = new FleetMcp(messages, workers, sessions, identity, presence,
|
||||
primaryRegistry, callers, FleetMcp.AuthorizationMode.ENFORCED, metrics,
|
||||
capacitySource(config, cfg, profile -> liveCountRef.get().apply(profile)),
|
||||
new FleetMcp.HealthCoverageSource(() -> {
|
||||
var health = config.get().health();
|
||||
return FleetHealthMonitor.coverage(health != null && health.isEnabled(),
|
||||
health != null && health.notifications() != null && health.notifications().configured());
|
||||
}),
|
||||
healthCoverageSource(config),
|
||||
quarantineSource,
|
||||
leadMailbox,
|
||||
outageSource,
|
||||
@@ -1038,6 +1034,38 @@ public final class Fleetd {
|
||||
cfg.profiles()::keySet, System::nanoTime);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #426: package-private factory for {@code fleet_list}'s {@code healthCoverage} source,
|
||||
* extracted out of {@code main} for the same reason {@link #capacitySource} and {@link
|
||||
* #quarantineSource} were — and the same reason {@link #exhaustedPatternCoverageLine}/{@link
|
||||
* #errorPatternCoverageLine} exist: {@link FleetHealthMonitor#coverage}'s three-branch method
|
||||
* is easy to pin directly (a plain {@code (boolean, boolean) -> String} call), but that proves
|
||||
* nothing about whether <em>this call site</em> pairs the right boolean with the right meaning.
|
||||
* fleetd #415's measured lesson is the reason this matters here — swapping the two arguments at
|
||||
* a call site like this one compiled clean and left the full suite green, because every existing
|
||||
* test exercised the method in both directions without ever exercising the pairing.
|
||||
*
|
||||
* <p>{@code #407}'s "keep the config invalid, assert on the log line before the throw" option
|
||||
* does not apply to this call site: the five reporters #407 covers all run in {@code main}
|
||||
* <em>before</em> {@code cfg.validateAll()} (line ~171), so an invalid config still exercises
|
||||
* them. This call site is built during {@code FleetMcp} construction, which runs only after
|
||||
* {@code UnixSocketHerdrClient.connect} has already opened a real herdr socket (line ~188) —
|
||||
* reaching it at all means main already performed real I/O, which the no-socket constraint on
|
||||
* this ticket rules out. So the pairing is pinned by extracting it to this directly-callable
|
||||
* factory instead, the same shape {@link #capacitySource}/{@link #quarantineSource} already use.
|
||||
*
|
||||
* <p>Reads {@code config.get().health()} live (health.notifications is a {@code SPLIT_KEYS}
|
||||
* entry — see {@link ConfigRef#SPLIT_KEYS}), so a hot-reloaded notifications block changes what
|
||||
* {@code fleet_list} reports without a restart, exactly like {@link #capacitySource}'s maxLoad.
|
||||
*/
|
||||
static FleetMcp.HealthCoverageSource healthCoverageSource(ConfigRef config) {
|
||||
return new FleetMcp.HealthCoverageSource(() -> {
|
||||
var health = config.get().health();
|
||||
return FleetHealthMonitor.coverage(health != null && health.isEnabled(),
|
||||
health != null && health.notifications() != null && health.notifications().configured());
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #248: package-private factory for the member worktree/branch lookup {@link
|
||||
* CompletionResolver} uses to name a fallback report's worktree and branch (fleetd#241).
|
||||
|
||||
@@ -208,7 +208,7 @@ import java.util.function.Supplier;
|
||||
* that is correctly hot and for a key nobody triaged. Three times now — {@code worktreeGroup} (#323),
|
||||
* {@code primary}/{@code configReload} (#326), and {@code fleet.leaders} sitting in the escape hatch
|
||||
* (#333) — the second kind hid among the first. A top-level coverage checker in the
|
||||
* {@link ConfigRefProfileCoverageTest} shape (one level up, over {@code FleetConfig} itself rather
|
||||
* {@code ConfigRefProfileCoverageTest} shape (one level up, over {@code FleetConfig} itself rather
|
||||
* than {@code FleetConfig.Profile}) proves this file's four classes exhaust the record's components
|
||||
* — see {@code ConfigRefTopLevelCoverageTest}. That test proves the record's <em>shape</em> is fully
|
||||
* triaged; it does NOT prove a {@code SPLIT_KEYS}/{@code COLD_KEYS}/{@code DEFERRED_KEYS} member has
|
||||
|
||||
@@ -1670,7 +1670,7 @@ public record FleetConfig(
|
||||
* <p>A herdr pane runs a login shell that re-sources the operator's own secret store, so a
|
||||
* member inherits every credential the operator's shell holds — measured at 31 names on this
|
||||
* host, of which only one ({@code GITEA_ACCESS_TOKEN}) used to be blocked, and that block was a
|
||||
* single name hardcoded in {@link HerdrPeerLauncher} rather than driven by config (gitea issue
|
||||
* single name hardcoded in {@link dev.ltms.fleet.member.HerdrPeerLauncher} rather than driven by config (gitea issue
|
||||
* #82). This record replaces that hardcoded shadow with a config-driven one.
|
||||
*
|
||||
* <p><b>deny-by-default, not a deny-list.</b> A deny-list (block these specific names, let
|
||||
|
||||
@@ -306,7 +306,7 @@ public final class Injector {
|
||||
if (t == null) return;
|
||||
|
||||
Pending sent = null;
|
||||
RuntimeException sendError = null;
|
||||
Throwable sendError = null;
|
||||
boolean turnCompleted = false;
|
||||
boolean turnFailed = false;
|
||||
boolean resubmit = false;
|
||||
@@ -388,9 +388,14 @@ public final class Injector {
|
||||
t.turnObserved = false;
|
||||
t.injectableSincePickup = 0;
|
||||
sent = p;
|
||||
} catch (RuntimeException e) {
|
||||
} catch (Throwable e) {
|
||||
// Delivery failed at herdr; drop the poisoned message and surface it
|
||||
// rather than blocking the queue behind it.
|
||||
// rather than blocking the queue behind it. Catches Throwable, not just
|
||||
// RuntimeException: fleetd #546 — an Error escaping this send (e.g. a
|
||||
// NoClassDefFoundError, see #413) would otherwise leave the entry QUEUED
|
||||
// at the head of t.queue. Line :378 peeks rather than polls, so the next
|
||||
// onStatus round would re-enter this try and send the same text again,
|
||||
// typing the same brief into the member's pane a second time.
|
||||
t.queue.poll();
|
||||
p.state = Pending.State.NOT_DELIVERED;
|
||||
sent = p;
|
||||
@@ -489,54 +494,163 @@ public final class Injector {
|
||||
|
||||
// Fire listeners / herdr calls after releasing the monitor so nothing runs on the poller
|
||||
// thread while it holds the target lock.
|
||||
if (resubmit) {
|
||||
try {
|
||||
agentsFor(target).submit(target); // nudge a raced Enter so the pending paste submits
|
||||
} catch (RuntimeException e) {
|
||||
log.debug("resubmit to {} failed (will retry next poll): {}", target, e.getMessage());
|
||||
//
|
||||
// fleetd #553: wrapped in try/finally. By this point, if `sent != null`, the delivery has
|
||||
// already happened inside the monitor above — off the queue, p.state == DELIVERED, text
|
||||
// typed into the target's pane — so `sent`'s future MUST be completed one way or another,
|
||||
// in every path out of this region, or the caller (a blocking fleet_send, or an async
|
||||
// ticket) waits forever on a message it actually received. But the finally must not
|
||||
// swallow whatever escaped: a listener that throws is a defect in THAT listener, and it
|
||||
// must still reach StatusPoller's catch (Throwable) so log.error fires — converting a
|
||||
// loud listener bug into a silently orphaned future would be worse than the bug itself.
|
||||
//
|
||||
// There are TWO futures at stake here, not one: `sent.delivered()` (the delivery future)
|
||||
// and `sent.token().waiter()` (the rendezvous waiter a blocking fleet_send actually waits
|
||||
// on for the worker's ANSWER). The waiter is registered only by turnListener.onDelivered(),
|
||||
// called from inside the `if (sent != null)` block below. Completing `sent.delivered()`
|
||||
// without also calling onDelivered leaves the waiter unregistered forever — a hang with a
|
||||
// success receipt, which is worse than the plain hang this ticket is about. So the
|
||||
// `finally` below does not just complete the future: on the path where the `if (sent !=
|
||||
// null)` block never ran, it does that block's whole job — onDelivered, then complete.
|
||||
//
|
||||
// `sentHandled` is NOT allowed to move earlier than this: an earlier comment on #553
|
||||
// proposed running `if (sent != null)` first, before turnCompleted/turnFailed, and that
|
||||
// was withdrawn — onDelivered WRITES CompletionResolver's inFlight record for this turn,
|
||||
// onTurnComplete READS it for the PREVIOUS turn, and running onDelivered first makes
|
||||
// onTurnComplete resolve the NEW turn's waiter with the PREVIOUS turn's stale output
|
||||
// (CB-116). The `if (sent != null)` block stays last; the `finally` is a backstop for it,
|
||||
// not a replacement.
|
||||
boolean sentHandled = false;
|
||||
try {
|
||||
if (resubmit) {
|
||||
try {
|
||||
agentsFor(target).submit(target); // nudge a raced Enter so the pending paste submits
|
||||
} catch (Throwable e) {
|
||||
// fleetd #553: widened from RuntimeException (same reasoning as #549 at :382) —
|
||||
// this is the FIRST block after the monitor, so an Error escaping it used to skip
|
||||
// every block below it, including the `sent` completion.
|
||||
log.debug("resubmit to {} failed (will retry next poll): {}", target, e.getMessage());
|
||||
}
|
||||
}
|
||||
}
|
||||
if (notReady != null) {
|
||||
// Worker never became available: forget its (never-set) readiness, unblock every queued
|
||||
// caller, and route the awaiting send through the same failure path as a stalled turn so
|
||||
// a blocking or async waiter resolves WORKER_FAILED rather than riding out the timeout.
|
||||
forget.accept(target);
|
||||
RuntimeException cause = new IllegalStateException(
|
||||
target + " never became available (no bridge MCP connection within the boot window)");
|
||||
for (Pending p : notReady) {
|
||||
p.delivered().completeExceptionally(cause);
|
||||
if (notReady != null) {
|
||||
// Worker never became available: unblock every queued caller FIRST — these messages'
|
||||
// fate (NOT_DELIVERED, off the queue) was already decided inside the monitor above, so
|
||||
// fleetd #553 completes every one of them before calling forget.accept or
|
||||
// turnListener.onTurnFailed below. Either of those is a listener/consumer callback and
|
||||
// can throw (an ordinary RuntimeException is enough — the same reasoning as the rest of
|
||||
// this ticket): completing the futures first means such a throw can no longer leave any
|
||||
// of them permanently pending, regardless of which one throws or in which order.
|
||||
RuntimeException cause = new IllegalStateException(
|
||||
target + " never became available (no bridge MCP connection within the boot window)");
|
||||
for (Pending p : notReady) {
|
||||
p.delivered().completeExceptionally(cause);
|
||||
}
|
||||
// route the awaiting send through the same failure path as a stalled turn so a blocking
|
||||
// or async waiter resolves WORKER_FAILED rather than riding out the timeout.
|
||||
forget.accept(target);
|
||||
turnListener.onTurnFailed(target);
|
||||
}
|
||||
turnListener.onTurnFailed(target);
|
||||
}
|
||||
if (turnCompleted) {
|
||||
if (startPostTurn) {
|
||||
boolean started = turnListener.onTurnCompleteWithPostAction(target);
|
||||
synchronized (t) {
|
||||
t.postTurnPending = false;
|
||||
if (started) {
|
||||
t.awaitingPostTurnPickup = true;
|
||||
t.injectableSincePostTurnPickup = 0;
|
||||
if (turnCompleted) {
|
||||
if (startPostTurn) {
|
||||
// fleetd #553: the listener call is wrapped so `t.postTurnPending` (set true inside
|
||||
// the monitor above, before this call) is always reset. Before this wrapping, a
|
||||
// RuntimeException from onTurnCompleteWithPostAction skipped the reset below,
|
||||
// permanently wedging the target: postTurnPending stayed true forever, so this
|
||||
// target's delivery guard (:376) would never again pass and no further message to it
|
||||
// would ever be delivered — a target-level lockup, not just one skipped future.
|
||||
// `started` defaults to false so an exception is treated as "the post-turn action did
|
||||
// not start" rather than falsely arming the post-turn pickup latch for an action that
|
||||
// never ran.
|
||||
boolean started = false;
|
||||
try {
|
||||
started = turnListener.onTurnCompleteWithPostAction(target);
|
||||
} finally {
|
||||
synchronized (t) {
|
||||
t.postTurnPending = false;
|
||||
if (started) {
|
||||
t.awaitingPostTurnPickup = true;
|
||||
t.injectableSincePostTurnPickup = 0;
|
||||
}
|
||||
if (t.queue.isEmpty() && !t.awaitingPostTurnPickup) {
|
||||
targets.remove(target, t);
|
||||
}
|
||||
}
|
||||
}
|
||||
if (t.queue.isEmpty() && !t.awaitingPostTurnPickup) {
|
||||
targets.remove(target, t);
|
||||
} else {
|
||||
turnListener.onTurnComplete(target);
|
||||
}
|
||||
}
|
||||
if (turnFailed) {
|
||||
turnListener.onTurnFailed(target);
|
||||
}
|
||||
if (sent != null) {
|
||||
// fleetd #553: record the INTENT to handle before any side effect, not the fact of
|
||||
// having completed it. If onDelivered() throws part way through — e.g. after
|
||||
// CompletionResolver's captureBaseline has already done its inFlight.put — a flag
|
||||
// set only after this block would still read false, and the `finally` below would
|
||||
// call onDelivered() a SECOND time. A second captureBaseline runs later, once
|
||||
// onTurnComplete has already thrown and possibly scraped the pane, and can snapshot
|
||||
// a pane that already absorbed this turn's output — which means the pane's tail
|
||||
// never differs from that baseline again and CompletionResolver.resolve's own
|
||||
// suppression at :364-369 (`return` keeping the in-flight record) drops every future
|
||||
// completion for this turn, permanently. So this flag must go true FIRST.
|
||||
sentHandled = true;
|
||||
if (sendError != null) {
|
||||
log.warn("inject to {} failed, dropped message: {}", target, sendError.getMessage());
|
||||
sent.delivered().completeExceptionally(sendError);
|
||||
} else {
|
||||
// Baseline the pane's pre-turn content so a misattributed completion (no new output)
|
||||
// can't resolve this send with the previous turn's stale answer (CB-115).
|
||||
turnListener.onDelivered(target, sent.token());
|
||||
sent.delivered().complete(null);
|
||||
}
|
||||
}
|
||||
} finally {
|
||||
// Backstop (fleetd #553): if an earlier block in this try threw before the `if (sent !=
|
||||
// null)` block above ran, `sentHandled` is still false here, and this does that block's
|
||||
// WHOLE job — not just the future completion. Skipping onDelivered() here would leave
|
||||
// sent.token().waiter() never registered with CompletionResolver, so a blocking
|
||||
// fleet_send would be told its message was delivered and then wait out its full timeout
|
||||
// for an answer that can never resolve — worse than the plain hang, because it now looks
|
||||
// like success. onDelivered() runs only when sendError == null: nothing was delivered on
|
||||
// the error path, so there is nothing to register.
|
||||
//
|
||||
// fleetd #553 (lead review, comment 16916): `sentHandled` guards ONLY the onDelivered
|
||||
// re-call below, never the future completion outside this inner try. One boolean cannot
|
||||
// carry both meanings — "onDelivered was called" and "the future has been dealt with" —
|
||||
// because they come apart exactly when onDelivered throws PART WAY THROUGH: sentHandled
|
||||
// is already true (set before the call, correctly — see the `if (sent != null)` block
|
||||
// above), so a guard on the outer `if` here would skip this whole recovery, including the
|
||||
// completion, and leave sent.delivered() pending forever even though the message really
|
||||
// was typed into the pane and taken off the queue. So `sent != null` alone gates whether
|
||||
// this target has anything to finish; `!sentHandled` gates only the onDelivered re-call
|
||||
// inside. CompletableFuture.complete/completeExceptionally are idempotent — on the
|
||||
// ordinary path (sentHandled == true, no throw) the `if (sent != null)` block above has
|
||||
// already completed this future, so the calls below are a no-op returning false.
|
||||
//
|
||||
// This recovery is wrapped in its own try/catch(Throwable) that swallows only ITS OWN
|
||||
// throwable and logs at WARN — the future is still completed either way — while the
|
||||
// ORIGINAL throwable from the try block above is left alone to keep unwinding out of
|
||||
// this method to StatusPoller's catch (Throwable), so a listener bug stays loud.
|
||||
if (sent != null) {
|
||||
try {
|
||||
if (!sentHandled && sendError == null) {
|
||||
turnListener.onDelivered(target, sent.token());
|
||||
}
|
||||
} catch (Throwable recoveryError) {
|
||||
log.warn("fleetd #553 backstop: onDelivered failed for {} while recovering from "
|
||||
+ "an earlier listener failure; completing its delivery future "
|
||||
+ "anyway: {}",
|
||||
target, recoveryError.getMessage());
|
||||
} finally {
|
||||
// No-op (returns false) on the ordinary path, where the `if (sent != null)` block
|
||||
// above already completed this future — see the comment above this block.
|
||||
if (sendError != null) {
|
||||
sent.delivered().completeExceptionally(sendError);
|
||||
} else {
|
||||
sent.delivered().complete(null);
|
||||
}
|
||||
}
|
||||
} else {
|
||||
turnListener.onTurnComplete(target);
|
||||
}
|
||||
}
|
||||
if (turnFailed) {
|
||||
turnListener.onTurnFailed(target);
|
||||
}
|
||||
if (sent != null) {
|
||||
if (sendError != null) {
|
||||
log.warn("inject to {} failed, dropped message: {}", target, sendError.getMessage());
|
||||
sent.delivered().completeExceptionally(sendError);
|
||||
} else {
|
||||
// Baseline the pane's pre-turn content so a misattributed completion (no new output)
|
||||
// can't resolve this send with the previous turn's stale answer (CB-115).
|
||||
turnListener.onDelivered(target, sent.token());
|
||||
sent.delivered().complete(null);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -64,34 +64,39 @@ public final class StatusPoller {
|
||||
}
|
||||
|
||||
private void loop() {
|
||||
while (running) {
|
||||
Set<String> active = injector.activeTargets();
|
||||
for (String target : active) {
|
||||
if (!running) return;
|
||||
try {
|
||||
// herdr's agent_status can misreport a settled worker as `unknown`; refine it
|
||||
// against the pane content before it drives delivery/completion (CB-115).
|
||||
// CB-185: refine THROUGH the same control the raw status came from — a router
|
||||
// splits lead/member targets across two herdr daemons, and reading a lead's pane
|
||||
// through the (fixed) member refiner never finds it, wedging that lead at UNKNOWN.
|
||||
AgentControl control = router != null ? router.agentsFor(target) : agents;
|
||||
AgentStatus status = refiner.refine(target, control.status(target), control);
|
||||
injector.onStatus(target, status);
|
||||
} catch (HerdrException e) {
|
||||
// The worker's agent is gone — stop trying and unblock its waiters.
|
||||
if (e.code() != null && e.code().endsWith("_not_found")) {
|
||||
log.debug("target {} gone; dropping its queue", target);
|
||||
injector.drop(target, e);
|
||||
} else {
|
||||
log.debug("status poll for {} failed (will retry): {}", target, e.getMessage());
|
||||
try {
|
||||
while (running) {
|
||||
Set<String> active = injector.activeTargets();
|
||||
for (String target : active) {
|
||||
if (!running) return;
|
||||
try {
|
||||
// herdr's agent_status can misreport a settled worker as `unknown`; refine it
|
||||
// against the pane content before it drives delivery/completion (CB-115).
|
||||
// CB-185: refine THROUGH the same control the raw status came from — a router
|
||||
// splits lead/member targets across two herdr daemons, and reading a lead's pane
|
||||
// through the (fixed) member refiner never finds it, wedging that lead at UNKNOWN.
|
||||
AgentControl control = router != null ? router.agentsFor(target) : agents;
|
||||
AgentStatus status = refiner.refine(target, control.status(target), control);
|
||||
injector.onStatus(target, status);
|
||||
} catch (HerdrException e) {
|
||||
// The worker's agent is gone — stop trying and unblock its waiters.
|
||||
if (e.code() != null && e.code().endsWith("_not_found")) {
|
||||
log.debug("target {} gone; dropping its queue", target);
|
||||
injector.drop(target, e);
|
||||
} else {
|
||||
log.debug("status poll for {} failed (will retry): {}", target, e.getMessage());
|
||||
}
|
||||
} catch (Throwable e) {
|
||||
log.error("unexpected failure polling {}; skipping this round", target, e);
|
||||
}
|
||||
} catch (RuntimeException e) {
|
||||
// Never let one target's unexpected error (e.g. an odd agent.get shape) kill
|
||||
// the single poller thread and stall injection for every worker.
|
||||
log.warn("unexpected error polling {}; skipping this round", target, e);
|
||||
}
|
||||
sleep();
|
||||
}
|
||||
sleep();
|
||||
} finally {
|
||||
if (running) {
|
||||
log.error("status poller loop exited unexpectedly; it can be restarted");
|
||||
}
|
||||
running = false;
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -246,7 +246,7 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
* neither. So before this method existed with a requeue step, it dropped {@link #held}'s entries
|
||||
* for {@code target} while the broker still considered them outstanding: never acked, never
|
||||
* nacked, never requeued, and no longer reachable by {@link #peek} — permanently invisible. This
|
||||
* is unlike {@link #handleRecovery} and {@link #close()}, whose bare {@code held.clear()} is
|
||||
* is unlike {@link RecoveryListener#handleRecovery(Recoverable)} and {@link #close()}, whose bare {@code held.clear()} is
|
||||
* correct because each has already made the broker requeue (a real connection drop, or
|
||||
* {@code channel.close()} respectively) before clearing local state.
|
||||
*
|
||||
|
||||
@@ -713,7 +713,7 @@ public final class MessageService {
|
||||
* failure.
|
||||
*
|
||||
* <p><strong>Without {@code sweepAsking} on the release path, a target torn down while
|
||||
* genuinely {@code ASKING} was unrecoverable.</strong> {@link #resolveQuestion} had already
|
||||
* genuinely {@code ASKING} was unrecoverable.</strong> {@link Rendezvous#resolveQuestion(String, String, String)} had already
|
||||
* closed the forward waiter the instant the question surfaced (so the {@code waiter} branch
|
||||
* below finds nothing to fail), the {@code question == null} guard excluded the task from
|
||||
* {@code matching} (so the loop below skipped it too), and the worker's own {@code fleet_ask}
|
||||
|
||||
@@ -44,7 +44,7 @@ import java.util.stream.Collectors;
|
||||
* still holds an unacked message ({@link #pendingReplies}), tickets not yet collected
|
||||
* ({@link #pendingTickets}), and open questions not yet answered or lapsed
|
||||
* ({@link #pendingQuestions}) — and sends at most one combined nudge per tick
|
||||
* ({@link #injectNudge(String, int, int, int)}). Work that arrives while the lead is busy is
|
||||
* ({@link #injectNudge(String, int, int, int, int, int)}). Work that arrives while the lead is busy is
|
||||
* never lost: it is re-read fresh on every tick until the lead is injectable or its own reminder
|
||||
* cap ({@link #maxReminders}) is reached — each source spends from its own budget, so one source
|
||||
* exhausting its cap does not stop nudges about the others (post-CB-590 regression fix; see
|
||||
|
||||
@@ -1074,6 +1074,29 @@ public final class SessionManager implements TurnListener {
|
||||
* members were live) and must still produce the line. Both {@link #drainSnapshot} passes (the
|
||||
* main snapshot and the straggler sweep) are folded into the one line: a caller reading two
|
||||
* lines could not tell a two-pass drain from two separate drains.
|
||||
*
|
||||
* <p><strong>Non-goal: this line must never move into a {@code finally} block, and this method
|
||||
* must never grow one around it.</strong> "Every time" above means every time the drain
|
||||
* <em>finishes</em>, not every time this method exits. The absence of the line is the signal
|
||||
* that the drain died, so a {@code finally} would destroy the signal and print confident
|
||||
* partial counts in the same edit — the line would appear after a drain that threw, carrying
|
||||
* whatever {@code tally} it had reached. Both halves of the value are lost at once. The line
|
||||
* has to be the last statement of the successful path and reachable only from it.
|
||||
*
|
||||
* <p>This is written down because it is the obvious review comment ("shouldn't we always log
|
||||
* the drain result?"), it sounds like thoroughness, and the paragraph above reads as an
|
||||
* invitation to it. Raised by the fleet01 lead on 2026-09-12, from their 2026-09-10 incident:
|
||||
* we only know that drain died on that host because it <em>threw</em>, and a
|
||||
* {@code NoClassDefFoundError} reached the JVM's uncaught handler. A drain that hung on one
|
||||
* session, or returned early on a condition rather than an exception, would leave no stack
|
||||
* trace, no {@code ERROR} token and no priority — only a missing line. That makes the loud
|
||||
* variant the one we have seen and the quiet variants the ones this line exists to catch.
|
||||
*
|
||||
* <p>Related: {@code released} and {@code abandoned} are counted incrementally inside {@link
|
||||
* #drainSnapshot}'s loop and folded with {@link DrainTally#plus}, rather than derived from a
|
||||
* collection read at the end, for the same reason. If a partial report is ever wanted it must
|
||||
* be a different line with a different verb. One line must not serve both, or a reader cannot
|
||||
* tell a finished drain from an interrupted one by its wording.
|
||||
*/
|
||||
void drainAll(long timeoutNanos) {
|
||||
long deadline = System.nanoTime() + timeoutNanos;
|
||||
|
||||
@@ -60,14 +60,21 @@ public final class SessionReaper {
|
||||
}
|
||||
|
||||
private void loop() {
|
||||
while (running) {
|
||||
try {
|
||||
sessions.reapIdle(idleTtlNanos);
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("session reaper iteration failed; continuing", e);
|
||||
try {
|
||||
while (running) {
|
||||
try {
|
||||
sessions.reapIdle(idleTtlNanos);
|
||||
} catch (Throwable e) {
|
||||
log.error("session reaper iteration failed; continuing", e);
|
||||
}
|
||||
maybeSweepWipRefs();
|
||||
sleep();
|
||||
}
|
||||
maybeSweepWipRefs();
|
||||
sleep();
|
||||
} finally {
|
||||
if (running) {
|
||||
log.error("session reaper loop exited unexpectedly; it can be restarted");
|
||||
}
|
||||
running = false;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -90,8 +97,8 @@ public final class SessionReaper {
|
||||
log.info("refs/wip retention sweep deleted {} snapshot ref(s) older than 24h whose "
|
||||
+ "content was already reachable from main", deleted);
|
||||
}
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("refs/wip retention sweep failed; continuing", e);
|
||||
} catch (Throwable e) {
|
||||
log.error("refs/wip retention sweep failed; continuing", e);
|
||||
}
|
||||
// Set even when the sweep threw, so a broken repo is retried on the slow cadence rather
|
||||
// than hammering git on every 5-second iteration.
|
||||
|
||||
@@ -1,14 +1,12 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
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.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.herdr.HerdrException;
|
||||
import dev.ltms.fleet.testing.CapturedLog;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.function.LongSupplier;
|
||||
@@ -92,49 +90,42 @@ class FleetdAwaitHerdrTest {
|
||||
|
||||
@Test
|
||||
void answeredLogsNothingAndSaysReap() {
|
||||
ListAppender<ILoggingEvent> events = attach();
|
||||
try {
|
||||
try (CapturedLog log = attach()) {
|
||||
boolean shouldReap = Fleetd.logHerdrWaitOutcomeAndShouldReap(
|
||||
new Fleetd.HerdrAwaitOutcome(Fleetd.HerdrWaitResult.ANSWERED, 0L));
|
||||
|
||||
assertTrue(shouldReap, "only ANSWERED should tell main to reap orphan workers");
|
||||
assertEquals(0, events.list.size(), "the answered path logs nothing itself");
|
||||
} finally {
|
||||
detach(events);
|
||||
assertEquals(0, log.events().size(), "the answered path logs nothing itself");
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void deadlinePassedLogsConfiguredAndMeasuredElapsedTogether() {
|
||||
ListAppender<ILoggingEvent> events = attach();
|
||||
try {
|
||||
try (CapturedLog log = attach()) {
|
||||
boolean shouldReap = Fleetd.logHerdrWaitOutcomeAndShouldReap(
|
||||
new Fleetd.HerdrAwaitOutcome(Fleetd.HerdrWaitResult.DEADLINE_PASSED, 30_500_000_000L));
|
||||
|
||||
assertFalse(shouldReap, "a deadline-passed wait must not tell main to reap");
|
||||
assertEquals(1, events.list.size());
|
||||
ILoggingEvent event = events.list.getFirst();
|
||||
assertEquals(1, log.events().size());
|
||||
ILoggingEvent event = log.events().getFirst();
|
||||
assertEquals(Level.WARN, event.getLevel());
|
||||
assertEquals("herdr did not answer within the configured wait (configured=30s "
|
||||
+ "elapsed=30500ms) — starting anyway; /healthz will report degraded until it "
|
||||
+ "comes up. Orphaned worker panes (if any) were NOT reaped.",
|
||||
event.getFormattedMessage());
|
||||
} finally {
|
||||
detach(events);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void interruptedLogsItsOwnMessageAndNeverClaimsTheBudgetElapsed() {
|
||||
ListAppender<ILoggingEvent> events = attach();
|
||||
try {
|
||||
try (CapturedLog log = attach()) {
|
||||
// 3ms: the ticket's own example of "a few milliseconds in", not the 30s budget.
|
||||
boolean shouldReap = Fleetd.logHerdrWaitOutcomeAndShouldReap(
|
||||
new Fleetd.HerdrAwaitOutcome(Fleetd.HerdrWaitResult.INTERRUPTED, 3_000_000L));
|
||||
|
||||
assertFalse(shouldReap, "an interrupted wait must not tell main to reap");
|
||||
assertEquals(1, events.list.size());
|
||||
ILoggingEvent event = events.list.getFirst();
|
||||
assertEquals(1, log.events().size());
|
||||
ILoggingEvent event = log.events().getFirst();
|
||||
assertEquals(Level.WARN, event.getLevel());
|
||||
String message = event.getFormattedMessage();
|
||||
assertEquals("herdr wait was interrupted before the configured wait ran out "
|
||||
@@ -143,8 +134,6 @@ class FleetdAwaitHerdrTest {
|
||||
message);
|
||||
assertFalse(message.contains("did not answer"),
|
||||
"an interrupted wait must not be reported as if herdr failed to answer within the budget");
|
||||
} finally {
|
||||
detach(events);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -197,16 +186,7 @@ class FleetdAwaitHerdrTest {
|
||||
}
|
||||
}
|
||||
|
||||
private static ListAppender<ILoggingEvent> attach() {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class);
|
||||
logger.setLevel(Level.DEBUG);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
return appender;
|
||||
}
|
||||
|
||||
private static void detach(ListAppender<ILoggingEvent> appender) {
|
||||
((Logger) LoggerFactory.getLogger(Fleetd.class)).detachAppender(appender);
|
||||
private static CapturedLog attach() {
|
||||
return CapturedLog.at(Fleetd.class, Level.DEBUG);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,160 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.mcp.FleetMcp;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #426: {@code FleetHealthMonitor.coverage} had zero references anywhere in the test
|
||||
* tree — not the method, not either output string, not the field it populates. {@code
|
||||
* FleetHealthMonitorCoverageTest} (package {@code dev.ltms.fleet.health}) pins the three-branch
|
||||
* method itself; that is the easy half.
|
||||
*
|
||||
* <p>The half that actually matters is this one: {@link Fleetd#healthCoverageSource} is the exact
|
||||
* call site {@code Fleetd.main} wires into {@code FleetMcp}'s constructor, and it is what feeds
|
||||
* {@code fleet_list}'s {@code healthCoverage} field (see {@code FleetMcp#listFleet}'s {@code
|
||||
* result.put("healthCoverage", healthCoverage.value().get())}). Measured precedent on fleetd #423
|
||||
* (for #415): swapping two arguments at a call site like this one — recreating #415's defect with
|
||||
* the keys exchanged — compiled with 0 errors and ran the ENTIRE suite (1506 tests) green. A
|
||||
* thoroughly-tested method proves nothing about whether the call site pairs its arguments correctly;
|
||||
* only a test that drives the call site itself can catch that.
|
||||
*
|
||||
* <p><strong>Why this is not driven through a real {@code Fleetd.main} the way #407 drives its five
|
||||
* reporters</strong> (keep the config invalid, assert on the log line emitted before {@code
|
||||
* cfg.validateAll()} throws): all five of #407's reporters run in {@code main} before {@code
|
||||
* validateAll()} (line ~171). {@link Fleetd#healthCoverageSource} is built during {@code FleetMcp}
|
||||
* construction, which happens only after {@code UnixSocketHerdrClient.connect} has already opened a
|
||||
* real herdr socket (line ~188) and after {@code sessions}/{@code workers} are constructed. Reaching
|
||||
* this call site by actually running {@code main} would require a real socket connect — banned by
|
||||
* this ticket's hard constraints — so #407's option 1 does not apply here. Instead {@link
|
||||
* Fleetd#healthCoverageSource} is extracted to a directly-callable package-private factory, the same
|
||||
* shape {@link Fleetd#capacitySource} and {@link Fleetd#quarantineSource} already use for the same
|
||||
* reason (see {@code FleetdCapacitySourceWiringTest}, the direct precedent this test follows).
|
||||
*
|
||||
* <p>Uses a real {@link FleetConfig#load} + {@link ConfigRef} (no socket, no port bind, no spawn,
|
||||
* nothing written outside {@code @TempDir}) so the fixture goes through the actual YAML parser and
|
||||
* {@code FleetConfig.Health}/{@code Notifications} records, not a hand-built stand-in that could
|
||||
* silently drift from what the parser actually produces.
|
||||
*/
|
||||
class FleetdHealthCoverageSourceWiringTest {
|
||||
|
||||
private static final String BASE = """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
herdrSocket: ~/.config/herdr/herdr.sock
|
||||
""";
|
||||
|
||||
private static final String HEALTH_DISABLED_WITH_WEBHOOK = BASE + """
|
||||
health:
|
||||
enabled: false
|
||||
notifications:
|
||||
mode: webhook
|
||||
""";
|
||||
|
||||
private static final String HEALTH_DETECTION_ONLY = BASE + """
|
||||
health:
|
||||
enabled: true
|
||||
""";
|
||||
|
||||
private static final String HEALTH_FULL = BASE + """
|
||||
health:
|
||||
enabled: true
|
||||
notifications:
|
||||
mode: webhook
|
||||
""";
|
||||
|
||||
private static final String HEALTH_ABSENT = BASE;
|
||||
|
||||
@Test
|
||||
@DisplayName("enabled: false reports off, even with a webhook configured")
|
||||
void disabledHealthReportsOff(@TempDir Path dir) throws Exception {
|
||||
FleetMcp.HealthCoverageSource source = sourceFor(dir, HEALTH_DISABLED_WITH_WEBHOOK);
|
||||
|
||||
assertEquals("off", source.value().get(),
|
||||
"health.enabled: false must report 'off' regardless of notifications — flipping "
|
||||
+ "the 'enabled' argument at the HealthCoverageSource call site would report "
|
||||
+ "'full' here instead");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("enabled with no notifications reports detection-only")
|
||||
void enabledWithoutNotificationsReportsDetectionOnly(@TempDir Path dir) throws Exception {
|
||||
FleetMcp.HealthCoverageSource source = sourceFor(dir, HEALTH_DETECTION_ONLY);
|
||||
|
||||
assertEquals("detection-only", source.value().get());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("enabled with a webhook configured reports full")
|
||||
void enabledWithNotificationsReportsFull(@TempDir Path dir) throws Exception {
|
||||
FleetMcp.HealthCoverageSource source = sourceFor(dir, HEALTH_FULL);
|
||||
|
||||
assertEquals("full", source.value().get(),
|
||||
"health.enabled: true with notifications.mode: webhook must report 'full' — "
|
||||
+ "swapping 'full' and 'detection-only' at the call site, or breaking the "
|
||||
+ "enabled/notificationConfigured argument pairing, would report "
|
||||
+ "'detection-only' here instead");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("an absent health: block reports off")
|
||||
void absentHealthBlockReportsOff(@TempDir Path dir) throws Exception {
|
||||
FleetMcp.HealthCoverageSource source = sourceFor(dir, HEALTH_ABSENT);
|
||||
|
||||
assertEquals("off", source.value().get());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #426, the live-wiring half: {@code health.notifications} is a {@code
|
||||
* ConfigRef.SPLIT_KEYS} entry, and {@link Fleetd#healthCoverageSource} reads {@code
|
||||
* config.get().health()} live (not the frozen startup {@code cfg}) — exactly like {@link
|
||||
* Fleetd#capacitySource}'s {@code maxLoad} ({@code FleetdCapacitySourceWiringTest}'s {@code
|
||||
* reloadedMaxLoadStillChangesWhatFleetListReports}). A hot-reloaded notifications block must
|
||||
* change what {@code fleet_list} reports without a restart; a fix that froze the whole source
|
||||
* against the startup snapshot would silently break that and every other test above would stay
|
||||
* green, since none of them reload.
|
||||
*/
|
||||
@Test
|
||||
@DisplayName("a hot notifications reload still changes what fleet_list reports")
|
||||
void reloadedNotificationsStillChangeWhatFleetListReports(@TempDir Path dir) throws Exception {
|
||||
Path file = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(file, HEALTH_DETECTION_ONLY);
|
||||
FleetConfig cfg = FleetConfig.load(file);
|
||||
ConfigRef config = new ConfigRef(file, cfg);
|
||||
|
||||
FleetMcp.HealthCoverageSource source = Fleetd.healthCoverageSource(config);
|
||||
assertEquals("detection-only", source.value().get(),
|
||||
"sanity: detection-only before any reload");
|
||||
|
||||
Files.writeString(file, HEALTH_FULL);
|
||||
assertTrue(config.reload().applied());
|
||||
// The live snapshot now carries the webhook — proves the reload really happened and this
|
||||
// test is not accidentally passing because nothing changed.
|
||||
assertTrue(config.get().health().notifications() != null
|
||||
&& config.get().health().notifications().configured(),
|
||||
"sanity: the reloaded config really carries a configured webhook");
|
||||
|
||||
assertEquals("full", source.value().get(),
|
||||
"the SAME HealthCoverageSource instance must reflect a reloaded notifications "
|
||||
+ "block without a restart — health.notifications is read live off "
|
||||
+ "config.get(), exactly like capacitySource's maxLoad");
|
||||
}
|
||||
|
||||
private static FleetMcp.HealthCoverageSource sourceFor(Path dir, String yaml) throws Exception {
|
||||
Path file = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(file, yaml);
|
||||
FleetConfig cfg = FleetConfig.load(file);
|
||||
ConfigRef config = new ConfigRef(file, cfg);
|
||||
return Fleetd.healthCoverageSource(config);
|
||||
}
|
||||
}
|
||||
@@ -1,15 +1,12 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
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 dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.msg.LeadMailbox;
|
||||
import dev.ltms.fleet.testing.CapturedLog;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
@@ -52,29 +49,22 @@ class FleetdLeadMailboxSelectionTest {
|
||||
}
|
||||
}
|
||||
|
||||
private static ListAppender<ILoggingEvent> captureFleetdLogs() {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
return appender;
|
||||
}
|
||||
|
||||
private static String joined(ListAppender<ILoggingEvent> appender, Level level) {
|
||||
return appender.list.stream().filter(e -> e.getLevel() == level)
|
||||
private static String joined(CapturedLog captured, Level level) {
|
||||
return captured.events().stream().filter(e -> e.getLevel() == level)
|
||||
.map(ILoggingEvent::getFormattedMessage).reduce("", (a, b) -> a + "\n" + b);
|
||||
}
|
||||
|
||||
@Test
|
||||
void noCoordinatorBlockLeavesTheFeatureOffSilently() {
|
||||
var appender = captureFleetdLogs();
|
||||
var opener = new RecordingOpener();
|
||||
try (var captured = CapturedLog.of(Fleetd.class)) {
|
||||
var opener = new RecordingOpener();
|
||||
|
||||
assertNull(Fleetd.openLeadMailbox(null, Map.of(), opener));
|
||||
assertNull(Fleetd.openLeadMailbox(null, Map.of(), opener));
|
||||
|
||||
assertNull(opener.offeredUri, "nothing configured means nothing is opened");
|
||||
assertEquals("", joined(appender, Level.WARN),
|
||||
"an opt-in feature nobody asked for must not warn on every boot");
|
||||
assertNull(opener.offeredUri, "nothing configured means nothing is opened");
|
||||
assertEquals("", joined(captured, Level.WARN),
|
||||
"an opt-in feature nobody asked for must not warn on every boot");
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -114,31 +104,33 @@ class FleetdLeadMailboxSelectionTest {
|
||||
|
||||
@Test
|
||||
void warnsAndStaysOffWhenSelfIdIsMissing() {
|
||||
var appender = captureFleetdLogs();
|
||||
var opener = new RecordingOpener();
|
||||
var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, null, null, null);
|
||||
try (var captured = CapturedLog.of(Fleetd.class)) {
|
||||
var opener = new RecordingOpener();
|
||||
var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, null, null, null);
|
||||
|
||||
assertNull(Fleetd.openLeadMailbox(coordinator, Map.of(), opener));
|
||||
assertNull(Fleetd.openLeadMailbox(coordinator, Map.of(), opener));
|
||||
|
||||
assertNull(opener.offeredUri, "a mailbox with no owning coord-id has no queue to declare");
|
||||
String warns = joined(appender, Level.WARN);
|
||||
assertTrue(warns.contains("coordinator.selfId"), () -> "say which key is missing: " + warns);
|
||||
assertFalse(warns.contains(SECRET), () -> "the URI's password must never be logged: " + warns);
|
||||
assertNull(opener.offeredUri, "a mailbox with no owning coord-id has no queue to declare");
|
||||
String warns = joined(captured, Level.WARN);
|
||||
assertTrue(warns.contains("coordinator.selfId"), () -> "say which key is missing: " + warns);
|
||||
assertFalse(warns.contains(SECRET), () -> "the URI's password must never be logged: " + warns);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void warnsAndStaysOffWhenTheBrokerIsUnreachableAtBoot() {
|
||||
var appender = captureFleetdLogs();
|
||||
var opener = new RecordingOpener();
|
||||
opener.unreachable = true;
|
||||
var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, "mac-opus", null, null);
|
||||
try (var captured = CapturedLog.of(Fleetd.class)) {
|
||||
var opener = new RecordingOpener();
|
||||
opener.unreachable = true;
|
||||
var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, "mac-opus", null, null);
|
||||
|
||||
assertNull(Fleetd.openLeadMailbox(coordinator, Map.of(), opener),
|
||||
"a down coordination broker turns the feature off; it must never take the daemon down");
|
||||
assertNull(Fleetd.openLeadMailbox(coordinator, Map.of(), opener),
|
||||
"a down coordination broker turns the feature off; it must never take the daemon down");
|
||||
|
||||
String warns = joined(appender, Level.WARN);
|
||||
assertTrue(warns.contains("coord.example"), () -> "name the host that failed: " + warns);
|
||||
assertFalse(warns.contains(SECRET), () -> "with credentials stripped: " + warns);
|
||||
assertTrue(warns.contains("Connection refused"), () -> "and the real reason: " + warns);
|
||||
String warns = joined(captured, Level.WARN);
|
||||
assertTrue(warns.contains("coord.example"), () -> "name the host that failed: " + warns);
|
||||
assertFalse(warns.contains(SECRET), () -> "with credentials stripped: " + warns);
|
||||
assertTrue(warns.contains("Connection refused"), () -> "and the real reason: " + warns);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,15 +1,13 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
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 dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.msg.AmqpReplyInbox;
|
||||
import dev.ltms.fleet.msg.InMemoryReplyInbox;
|
||||
import dev.ltms.fleet.msg.ReplyInbox;
|
||||
import dev.ltms.fleet.testing.CapturedLog;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.net.ServerSocket;
|
||||
import java.util.List;
|
||||
@@ -50,43 +48,34 @@ class FleetdReplyInboxSelectionTest {
|
||||
}
|
||||
}
|
||||
|
||||
private static ListAppender<ILoggingEvent> attach() {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class);
|
||||
private static CapturedLog attach() {
|
||||
// logback-test.xml pins dev.ltms.fleet to WARN; raise it so INFO selection lines are captured.
|
||||
logger.setLevel(Level.INFO);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
return appender;
|
||||
return CapturedLog.at(Fleetd.class, Level.INFO);
|
||||
}
|
||||
|
||||
private static void detach(ListAppender<ILoggingEvent> appender) {
|
||||
((Logger) LoggerFactory.getLogger(Fleetd.class)).detachAppender(appender);
|
||||
}
|
||||
|
||||
private static void assertNoLogContains(ListAppender<ILoggingEvent> appender, String secret) {
|
||||
assertTrue(appender.list.stream().noneMatch(e -> e.getFormattedMessage().contains(secret)),
|
||||
private static void assertNoLogContains(List<ILoggingEvent> events, String secret) {
|
||||
assertTrue(events.stream().noneMatch(e -> e.getFormattedMessage().contains(secret)),
|
||||
"no log line may contain the resolved URI's password");
|
||||
}
|
||||
|
||||
@Test
|
||||
void uriEnvSetAndPresentSelectsAmqpWithTheResolvedUri() {
|
||||
FleetConfig.Broker broker = new FleetConfig.Broker(null, "LAVINMQ_URI", null);
|
||||
recording(broker, Map.of("LAVINMQ_URI", RESOLVED_URI), false, (opener, appender, inbox) -> {
|
||||
recording(broker, Map.of("LAVINMQ_URI", RESOLVED_URI), false, (opener, events, inbox) -> {
|
||||
assertEquals(RESOLVED_URI, opener.offeredUri,
|
||||
"the daemon must connect with the value resolved from uriEnv — selection, not just parse");
|
||||
assertEquals(opener.inbox, inbox, "the AMQP opener's inbox is what is selected");
|
||||
assertNoLogContains(appender, SECRET);
|
||||
assertNoLogContains(events, SECRET);
|
||||
});
|
||||
}
|
||||
|
||||
@Test
|
||||
void uriEnvSetButVariableMissingFallsBackToInMemoryAndWarns() {
|
||||
FleetConfig.Broker broker = new FleetConfig.Broker(null, "LAVINMQ_URI", null);
|
||||
recording(broker, Map.of(), false, (opener, appender, inbox) -> {
|
||||
recording(broker, Map.of(), false, (opener, events, inbox) -> {
|
||||
assertInstanceOf(InMemoryReplyInbox.class, inbox);
|
||||
assertNull(opener.offeredUri, "AMQP must never be attempted when the variable is missing");
|
||||
assertTrue(hasWarnContaining(appender, "LAVINMQ_URI") && hasWarnContaining(appender, "DISABLED"),
|
||||
assertTrue(hasWarnContaining(events, "LAVINMQ_URI") && hasWarnContaining(events, "DISABLED"),
|
||||
"a missing uriEnv variable must warn loudly, not fail silently");
|
||||
});
|
||||
}
|
||||
@@ -94,10 +83,10 @@ class FleetdReplyInboxSelectionTest {
|
||||
@Test
|
||||
void uriEnvSetButVariableBlankFallsBackToInMemoryAndWarns() {
|
||||
FleetConfig.Broker broker = new FleetConfig.Broker("amqp://user:lame@old:5672/", "LAVINMQ_URI", null);
|
||||
recording(broker, Map.of("LAVINMQ_URI", " "), false, (opener, appender, inbox) -> {
|
||||
recording(broker, Map.of("LAVINMQ_URI", " "), false, (opener, events, inbox) -> {
|
||||
assertInstanceOf(InMemoryReplyInbox.class, inbox);
|
||||
assertNull(opener.offeredUri, "a blank env value must not select AMQP, not even via the literal uri");
|
||||
assertTrue(hasWarnContaining(appender, "LAVINMQ_URI"),
|
||||
assertTrue(hasWarnContaining(events, "LAVINMQ_URI"),
|
||||
"a blank uriEnv value must warn, and must not fall back to the literal uri");
|
||||
});
|
||||
}
|
||||
@@ -106,22 +95,22 @@ class FleetdReplyInboxSelectionTest {
|
||||
void bothUriAndUriEnvSetUriEnvWinsDeterministically() {
|
||||
FleetConfig.Broker broker
|
||||
= new FleetConfig.Broker("amqp://user:oldpw@old.example:5672/", "LAVINMQ_URI", null);
|
||||
recording(broker, Map.of("LAVINMQ_URI", RESOLVED_URI), false, (opener, appender, inbox) -> {
|
||||
recording(broker, Map.of("LAVINMQ_URI", RESOLVED_URI), false, (opener, events, inbox) -> {
|
||||
assertEquals(RESOLVED_URI, opener.offeredUri,
|
||||
"uriEnv must win over uri, deterministically, every run");
|
||||
assertTrue(appender.list.stream().anyMatch(e -> e.getFormattedMessage().contains("broker.uri is ignored")),
|
||||
assertTrue(events.stream().anyMatch(e -> e.getFormattedMessage().contains("broker.uri is ignored")),
|
||||
"must log that the literal uri is ignored when uriEnv is set");
|
||||
assertNoLogContains(appender, SECRET);
|
||||
assertNoLogContains(appender, "oldpw");
|
||||
assertNoLogContains(events, SECRET);
|
||||
assertNoLogContains(events, "oldpw");
|
||||
});
|
||||
}
|
||||
|
||||
@Test
|
||||
void unreachableBrokerStartsDaemonWithInMemoryInboxAndLoudWarning() {
|
||||
FleetConfig.Broker broker = new FleetConfig.Broker(null, "LAVINMQ_URI", null);
|
||||
recording(broker, Map.of("LAVINMQ_URI", RESOLVED_URI), true, (opener, appender, inbox) -> {
|
||||
recording(broker, Map.of("LAVINMQ_URI", RESOLVED_URI), true, (opener, events, inbox) -> {
|
||||
assertInstanceOf(InMemoryReplyInbox.class, inbox, "an unreachable broker must NOT stop the daemon");
|
||||
String warn = appender.list.stream()
|
||||
String warn = events.stream()
|
||||
.filter(e -> e.getLevel() == Level.WARN)
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.reduce("", (a, b) -> a + "\n" + b)
|
||||
@@ -129,16 +118,16 @@ class FleetdReplyInboxSelectionTest {
|
||||
assertTrue(warn.contains("durable") && warn.contains("soft-state"),
|
||||
"the warning must say exactly what was lost: durable delivery off, replies soft-state");
|
||||
assertTrue(!warn.contains(SECRET), "the failing URI must be logged with credentials stripped");
|
||||
assertNoLogContains(appender, SECRET);
|
||||
assertNoLogContains(events, SECRET);
|
||||
});
|
||||
}
|
||||
|
||||
@Test
|
||||
void noBrokerConfiguredStaysQuietInMemory() {
|
||||
FleetConfig.Broker broker = null;
|
||||
recording(broker, Map.of(), false, (opener, appender, inbox) -> {
|
||||
recording(broker, Map.of(), false, (opener, events, inbox) -> {
|
||||
assertInstanceOf(InMemoryReplyInbox.class, inbox);
|
||||
assertTrue(appender.list.stream().noneMatch(e -> e.getLevel() == Level.WARN),
|
||||
assertTrue(events.stream().noneMatch(e -> e.getLevel() == Level.WARN),
|
||||
"no broker configured must keep the existing QUIET in-memory path — no warning");
|
||||
});
|
||||
}
|
||||
@@ -152,21 +141,18 @@ class FleetdReplyInboxSelectionTest {
|
||||
closedPort = s.getLocalPort();
|
||||
}
|
||||
FleetConfig.Broker broker = new FleetConfig.Broker(null, "LAVINMQ_URI", null);
|
||||
ListAppender<ILoggingEvent> appender = attach();
|
||||
try {
|
||||
try (CapturedLog log = attach()) {
|
||||
ReplyInbox inbox = Fleetd.selectReplyInbox(
|
||||
broker, Map.of("LAVINMQ_URI", "amqp://user:" + SECRET + "@127.0.0.1:" + closedPort + "/vh"),
|
||||
AmqpReplyInbox::open);
|
||||
assertInstanceOf(InMemoryReplyInbox.class, inbox,
|
||||
"a genuinely unreachable broker (real AmqpReplyInbox::open) must fall back to in-memory");
|
||||
} finally {
|
||||
detach(appender);
|
||||
assertNoLogContains(log.events(), SECRET);
|
||||
}
|
||||
assertNoLogContains(appender, SECRET);
|
||||
}
|
||||
|
||||
private boolean hasWarnContaining(ListAppender<ILoggingEvent> appender, String fragment) {
|
||||
return appender.list.stream().anyMatch(e ->
|
||||
private boolean hasWarnContaining(List<ILoggingEvent> events, String fragment) {
|
||||
return events.stream().anyMatch(e ->
|
||||
e.getLevel() == Level.WARN && e.getFormattedMessage().contains(fragment));
|
||||
}
|
||||
|
||||
@@ -174,18 +160,14 @@ class FleetdReplyInboxSelectionTest {
|
||||
private void recording(FleetConfig.Broker broker, Map<String, String> env, boolean unreachable, Check check) {
|
||||
RecordingAmqp opener = new RecordingAmqp();
|
||||
opener.unreachable = unreachable;
|
||||
ListAppender<ILoggingEvent> appender = attach();
|
||||
ReplyInbox inbox;
|
||||
try {
|
||||
inbox = Fleetd.selectReplyInbox(broker, env, opener);
|
||||
} finally {
|
||||
detach(appender);
|
||||
try (CapturedLog log = attach()) {
|
||||
ReplyInbox inbox = Fleetd.selectReplyInbox(broker, env, opener);
|
||||
check.run(opener, log.events(), inbox);
|
||||
}
|
||||
check.run(opener, appender, inbox);
|
||||
}
|
||||
|
||||
@FunctionalInterface
|
||||
private interface Check {
|
||||
void run(RecordingAmqp opener, ListAppender<ILoggingEvent> appender, ReplyInbox inbox);
|
||||
void run(RecordingAmqp opener, List<ILoggingEvent> events, ReplyInbox inbox);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,15 +1,12 @@
|
||||
package dev.ltms.fleet.auth;
|
||||
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.LoggerContext;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import dev.ltms.fleet.testing.CapturedLog;
|
||||
import org.junit.jupiter.api.AfterEach;
|
||||
import org.junit.jupiter.api.BeforeEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
@@ -24,28 +21,21 @@ import static org.junit.jupiter.api.Assertions.*;
|
||||
class AuditLogTest {
|
||||
|
||||
private final ObjectMapper mapper = new ObjectMapper();
|
||||
private ListAppender<ILoggingEvent> appender;
|
||||
private ch.qos.logback.classic.Logger auditLogger;
|
||||
private CapturedLog auditLog;
|
||||
|
||||
@BeforeEach
|
||||
void attach() {
|
||||
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
|
||||
auditLogger = ctx.getLogger("audit");
|
||||
appender = new ListAppender<>();
|
||||
appender.setContext(ctx);
|
||||
appender.start();
|
||||
auditLogger.addAppender(appender);
|
||||
auditLogger.setLevel(Level.INFO);
|
||||
auditLog = CapturedLog.at("audit", Level.INFO);
|
||||
}
|
||||
|
||||
@AfterEach
|
||||
void detach() {
|
||||
auditLogger.detachAppender(appender);
|
||||
auditLog.close();
|
||||
}
|
||||
|
||||
private JsonNode onlyRecord() throws Exception {
|
||||
assertEquals(1, appender.list.size(), "exactly one audit line expected");
|
||||
String line = appender.list.getFirst().getFormattedMessage();
|
||||
assertEquals(1, auditLog.events().size(), "exactly one audit line expected");
|
||||
String line = auditLog.events().getFirst().getFormattedMessage();
|
||||
return mapper.readTree(line); // throws if the line is not valid JSON
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,40 @@
|
||||
package dev.ltms.fleet.health;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
|
||||
/**
|
||||
* fleetd #426: {@link FleetHealthMonitor#coverage} had zero references anywhere in the test tree —
|
||||
* not the method, not either output string, not the field it populates. This pins the method's own
|
||||
* three branches directly.
|
||||
*
|
||||
* <p>This is the easy half. It proves the method words each combination correctly, but it proves
|
||||
* nothing about whether {@code Fleetd.java}'s two call sites pass the right argument in the right
|
||||
* position — see {@code FleetdHealthCoverageSourceWiringTest} (package {@code dev.ltms.fleet}) for
|
||||
* the half that actually guards the call site, following the measured fleetd #415 lesson that a
|
||||
* thoroughly-tested method and an untested argument pairing at its call site are different risks.
|
||||
*
|
||||
* <p><strong>The three output strings are load-bearing and must not change here.</strong> {@code
|
||||
* "detection-only"} is read live off a running daemon's {@code fleet_list} today (measured
|
||||
* 2026-09-12) — this test intentionally asserts the exact literal strings so a future edit to the
|
||||
* wording trips it here first.
|
||||
*/
|
||||
class FleetHealthMonitorCoverageTest {
|
||||
|
||||
@Test
|
||||
void disabledIsOffRegardlessOfNotificationConfig() {
|
||||
assertEquals("off", FleetHealthMonitor.coverage(false, false));
|
||||
assertEquals("off", FleetHealthMonitor.coverage(false, true));
|
||||
}
|
||||
|
||||
@Test
|
||||
void enabledWithoutNotificationsIsDetectionOnly() {
|
||||
assertEquals("detection-only", FleetHealthMonitor.coverage(true, false));
|
||||
}
|
||||
|
||||
@Test
|
||||
void enabledWithNotificationsIsFull() {
|
||||
assertEquals("full", FleetHealthMonitor.coverage(true, true));
|
||||
}
|
||||
}
|
||||
@@ -1,9 +1,7 @@
|
||||
package dev.ltms.fleet.inject;
|
||||
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.LoggerContext;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
@@ -12,8 +10,8 @@ import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.msg.Rendezvous;
|
||||
import dev.ltms.fleet.msg.TestTurnTokens;
|
||||
import dev.ltms.fleet.msg.TurnToken;
|
||||
import dev.ltms.fleet.testing.CapturedLog;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.Set;
|
||||
import java.util.regex.Pattern;
|
||||
@@ -470,15 +468,7 @@ class CompletionResolverTest {
|
||||
// CB-564: this used to be a bare DEBUG "failed send to X via turn-stall fallback" — a symptom
|
||||
// with no cause, and below the level anyone watching for member health would see. A fail that
|
||||
// resolves a caller's blocked send is at least WARN and must carry the reason.
|
||||
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
|
||||
ch.qos.logback.classic.Logger resolverLog =
|
||||
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(CompletionResolver.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.setContext(ctx);
|
||||
appender.start();
|
||||
resolverLog.addAppender(appender);
|
||||
resolverLog.setLevel(Level.WARN);
|
||||
try {
|
||||
try (CapturedLog log = CapturedLog.at(CompletionResolver.class, Level.WARN)) {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("stuck on an error screen");
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
@@ -486,7 +476,7 @@ class CompletionResolverTest {
|
||||
|
||||
resolver.fail("term_a", null);
|
||||
|
||||
String warn = appender.list.stream()
|
||||
String warn = log.events().stream()
|
||||
.filter(e -> e.getLevel().equals(Level.WARN))
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.findFirst()
|
||||
@@ -494,8 +484,6 @@ class CompletionResolverTest {
|
||||
assertTrue(warn.contains("term_a"), "the log names the target: " + warn);
|
||||
assertTrue(warn.contains("stuck on an error screen"), "the log carries the reason: " + warn);
|
||||
assertTrue(waiter.isDone());
|
||||
} finally {
|
||||
resolverLog.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -1,16 +1,18 @@
|
||||
package dev.ltms.fleet.inject;
|
||||
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.LoggerContext;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.AgentStatus;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.herdr.HerdrException;
|
||||
import dev.ltms.fleet.msg.Rendezvous;
|
||||
import dev.ltms.fleet.msg.TestTurnTokens;
|
||||
import dev.ltms.fleet.msg.TurnToken;
|
||||
import dev.ltms.fleet.testing.CapturedLog;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
@@ -493,23 +495,15 @@ class InjectorTest {
|
||||
void readinessGraceExpiryIsLogged() {
|
||||
// CB-562: the grace-expiry path used to clear the queue silently, so a message that never
|
||||
// reached the worker's pane surfaced elsewhere as an unrelated turn-stall failure. Assert the
|
||||
// expiry now names the real cause. (ListAppender capture pattern mirrors AuditLogTest.)
|
||||
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 {
|
||||
// expiry now names the real cause. (CapturedLog pattern mirrors AuditLogTest.)
|
||||
try (CapturedLog log = CapturedLog.at(Injector.class, Level.WARN)) {
|
||||
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()
|
||||
String warn = log.events().stream()
|
||||
.filter(e -> e.getLevel().equals(Level.WARN))
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.findFirst()
|
||||
@@ -518,8 +512,6 @@ class InjectorTest {
|
||||
assertTrue(warn.contains("never reached"), "the log names the real cause: " + warn);
|
||||
assertTrue(warn.contains("1 queued message"),
|
||||
"the log carries the failed message count: " + warn);
|
||||
} finally {
|
||||
injectorLog.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -533,22 +525,14 @@ class InjectorTest {
|
||||
// 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 {
|
||||
try (CapturedLog log = CapturedLog.at(Injector.class, Level.WARN)) {
|
||||
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()
|
||||
String warn = log.events().stream()
|
||||
.filter(e -> e.getLevel().equals(Level.WARN))
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.findFirst()
|
||||
@@ -557,8 +541,6 @@ class InjectorTest {
|
||||
"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);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -584,22 +566,14 @@ class InjectorTest {
|
||||
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 {
|
||||
try (CapturedLog log = CapturedLog.at(Injector.class, Level.WARN)) {
|
||||
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()
|
||||
String warn = log.events().stream()
|
||||
.filter(e -> e.getLevel().equals(Level.WARN))
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.findFirst()
|
||||
@@ -609,8 +583,6 @@ class InjectorTest {
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -647,21 +619,13 @@ class InjectorTest {
|
||||
// CB-564: a vanished worker used to drop its queue with no log at all — the only trace was
|
||||
// whatever failed downstream (e.g. a caller's send timing out with no clue why). Assert the
|
||||
// drop itself now names the cause and the number of messages it failed.
|
||||
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 {
|
||||
try (CapturedLog log = CapturedLog.at(Injector.class, Level.WARN)) {
|
||||
Injector inj = new Injector(new AgentControl(herdr), TurnListener.NOOP, _ -> true, _ -> {
|
||||
});
|
||||
inj.enqueue(T, "orphan", TestTurnTokens.inert(T));
|
||||
inj.drop(T, new HerdrException("worker gone", "pane_not_found", null));
|
||||
|
||||
String warn = appender.list.stream()
|
||||
String warn = log.events().stream()
|
||||
.filter(e -> e.getLevel().equals(Level.WARN))
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.findFirst()
|
||||
@@ -669,8 +633,6 @@ class InjectorTest {
|
||||
assertTrue(warn.contains(T), "the log names the target terminal: " + warn);
|
||||
assertTrue(warn.contains("1 message"), "the log carries the failed message count: " + warn);
|
||||
assertTrue(warn.contains("worker gone"), "the log carries the real cause: " + warn);
|
||||
} finally {
|
||||
injectorLog.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -706,4 +668,533 @@ class InjectorTest {
|
||||
ExecutionException ex = assertThrows(ExecutionException.class, f::get);
|
||||
assertInstanceOf(HerdrException.class, ex.getCause());
|
||||
}
|
||||
|
||||
/**
|
||||
* A {@link HerdrClient} that throws a non-{@link RuntimeException} {@link Error} from {@code
|
||||
* agent.prompt} instead of delegating — the fleetd #546 case: a stray non-RuntimeException
|
||||
* throwable (e.g. a {@code NoClassDefFoundError}, fleetd #413) escaping the send seam at
|
||||
* {@code Injector.java:383}. Records every call it sees itself, including the ones it throws
|
||||
* for, since the delegate's own recording is never reached for {@code agent.prompt} — so a test
|
||||
* can assert on exactly what this fake actually received.
|
||||
*/
|
||||
private static final class ErrorOnPrompt implements HerdrClient {
|
||||
private final FakeHerdr delegate;
|
||||
private final List<FakeHerdr.Call> calls = new java.util.concurrent.CopyOnWriteArrayList<>();
|
||||
|
||||
private ErrorOnPrompt(FakeHerdr delegate) {
|
||||
this.delegate = delegate;
|
||||
}
|
||||
|
||||
List<FakeHerdr.Call> calls() {
|
||||
return calls;
|
||||
}
|
||||
|
||||
@Override
|
||||
public JsonNode call(String method, Object params) {
|
||||
calls.add(new FakeHerdr.Call(method, params));
|
||||
if (method.equals("agent.prompt")) {
|
||||
throw new AssertionError("simulated non-RuntimeException send failure (fleetd #546)");
|
||||
}
|
||||
return delegate.call(method, params);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
delegate.close();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Mirrors StatusPoller's own per-target {@code catch (Throwable)} (fleetd #538 / PR #543): the
|
||||
* production polling loop already swallows whatever escapes one target's round and comes back
|
||||
* for the next one. A unit test that calls {@code onStatus} directly (bypassing StatusPoller)
|
||||
* needs the same survival so it can observe what a SECOND round does, regardless of whether
|
||||
* fleetd #546's fix is present.
|
||||
*/
|
||||
private static void pollOnceSurviving(Injector inj, String target, AgentStatus status) {
|
||||
try {
|
||||
inj.onStatus(target, status);
|
||||
} catch (Throwable ignored) {
|
||||
// matches StatusPoller.loop's own catch (Throwable) added by PR #543
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void anErrorFromSendRemovesTheMessageAndMarksItNotDelivered() {
|
||||
// fleetd #546, acceptance test 1: an Error (not a RuntimeException) escaping the send seam
|
||||
// at Injector.java:383 must still be caught, the message dropped from the queue, and its
|
||||
// state set to NOT_DELIVERED — never left QUEUED at the head.
|
||||
ErrorOnPrompt throwing = new ErrorOnPrompt(new FakeHerdr());
|
||||
Injector inj = new Injector(new AgentControl(throwing));
|
||||
Injector.Delivery delivery = inj.enqueue(T, "brief", TestTurnTokens.inert(T));
|
||||
|
||||
assertDoesNotThrow(() -> inj.onStatus(T, AgentStatus.IDLE),
|
||||
"fleetd #546: an Error from the send seam must be caught inside onStatus, not "
|
||||
+ "escape it");
|
||||
|
||||
assertTrue(delivery.completion().isCompletedExceptionally(),
|
||||
"the delivery's future must surface the send failure");
|
||||
assertEquals(Injector.Cancellation.NOT_DELIVERED, inj.cancel(delivery),
|
||||
"the message must be dropped and marked NOT_DELIVERED, not left QUEUED at the "
|
||||
+ "head of the queue");
|
||||
}
|
||||
|
||||
@Test
|
||||
void anErrorFromSendDoesNotRedeliverOnASecondRound() {
|
||||
// fleetd #546, acceptance test 2 — the test this ticket exists for. Before the fix, an
|
||||
// Error at the send seam left the message QUEUED (Injector.java:378 peeks, not polls), so a
|
||||
// second onStatus round re-entered the same try and sent the same text again: the member's
|
||||
// pane got the same brief typed into it twice.
|
||||
ErrorOnPrompt throwing = new ErrorOnPrompt(new FakeHerdr());
|
||||
Injector inj = new Injector(new AgentControl(throwing));
|
||||
inj.enqueue(T, "brief", TestTurnTokens.inert(T));
|
||||
|
||||
pollOnceSurviving(inj, T, AgentStatus.IDLE);
|
||||
pollOnceSurviving(inj, T, AgentStatus.IDLE);
|
||||
|
||||
long promptCalls = throwing.calls().stream().filter(c -> c.method().equals("agent.prompt")).count();
|
||||
assertEquals(1, promptCalls, "fleetd #546: the poisoned text must be sent exactly once — a "
|
||||
+ "second onStatus round must not re-enter send for the same message");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aHerdrExceptionFromSendStillProducesNotDeliveredUnchanged() {
|
||||
// fleetd #546, acceptance test 3 (control): HerdrException extends RuntimeException, so it
|
||||
// was already caught before this ticket's widening. This pins that the ordinary path is
|
||||
// unchanged — still dropped, still NOT_DELIVERED, still surfaced to the caller.
|
||||
FakeHerdr failing = new FakeHerdr().agentSendFailsWith("send_failed");
|
||||
Injector inj = new Injector(new AgentControl(failing));
|
||||
Injector.Delivery delivery = inj.enqueue(T, "boom", TestTurnTokens.inert(T));
|
||||
|
||||
inj.onStatus(T, AgentStatus.IDLE);
|
||||
|
||||
assertTrue(delivery.completion().isCompletedExceptionally(),
|
||||
"a HerdrException at the send seam must still surface to the caller, unchanged by "
|
||||
+ "fleetd #546's widening");
|
||||
assertEquals(Injector.Cancellation.NOT_DELIVERED, inj.cancel(delivery),
|
||||
"a HerdrException must still be dropped and marked NOT_DELIVERED, unchanged by the "
|
||||
+ "wider Throwable catch");
|
||||
}
|
||||
|
||||
// --- fleetd #553: a throwable from any listener callback in onStatus must not skip the
|
||||
// delivered-future completion. The whole post-monitor region is now wrapped in try/finally. ---
|
||||
|
||||
@Test
|
||||
void theOriginalThrowableFromAListenerStillEscapesOnStatus() {
|
||||
// fleetd #553's control: the entire point of the finally backstop is that it completes a
|
||||
// future WITHOUT swallowing whatever escaped. A fix built the wrong way (e.g. catching
|
||||
// Throwable in the finally, or wrapping the region in try/catch instead of try/finally)
|
||||
// could make every other test in this group pass while still converting a loud listener
|
||||
// bug into a silent one — which the ticket calls a worse outcome than the bug itself. This
|
||||
// is deliberately its own test, not folded into another one's assertion.
|
||||
RuntimeException boom = new RuntimeException("fleetd #553 control");
|
||||
TurnListener throwing = new TurnListener() {
|
||||
@Override
|
||||
public void onTurnComplete(String target) {
|
||||
throw boom;
|
||||
}
|
||||
};
|
||||
Injector inj = new Injector(new AgentControl(herdr), throwing);
|
||||
inj.enqueue(T, "only", TestTurnTokens.inert(T));
|
||||
inj.onStatus(T, AgentStatus.IDLE); // deliver
|
||||
inj.onStatus(T, AgentStatus.WORKING); // pickup
|
||||
|
||||
RuntimeException thrown = assertThrows(RuntimeException.class,
|
||||
() -> inj.onStatus(T, AgentStatus.IDLE),
|
||||
"onStatus must still propagate the listener's own throwable, unmodified");
|
||||
assertSame(boom, thrown, "must be the EXACT throwable, not a wrapper or a different instance");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aRuntimeExceptionFromOnTurnCompleteStillCompletesTheNextDelivery() {
|
||||
// fleetd #553, acceptance test 1: onTurnComplete fires for "first"'s completed turn INSIDE
|
||||
// the same onStatus call that then peeks and delivers "second" — the turnCompleted check
|
||||
// (Injector.java, inside the released-pickup branch) clears awaitingCompletion just before
|
||||
// the delivery guard right below it is evaluated, so both happen in one round when there is
|
||||
// no post-turn action. Before fleetd #553, a RuntimeException thrown here unwound straight
|
||||
// out of onStatus and skipped the `if (sent != null)` block, leaving "second"'s future
|
||||
// pending forever even though it was already off the queue, marked DELIVERED, and typed
|
||||
// into the pane.
|
||||
RuntimeException boom = new RuntimeException("boom from onTurnComplete");
|
||||
TurnListener throwing = new TurnListener() {
|
||||
@Override
|
||||
public void onTurnComplete(String target) {
|
||||
throw boom;
|
||||
}
|
||||
};
|
||||
Injector inj = new Injector(new AgentControl(herdr), throwing);
|
||||
inj.enqueue(T, "first", TestTurnTokens.inert(T));
|
||||
Injector.Delivery second = inj.enqueue(T, "second", TestTurnTokens.inert(T));
|
||||
|
||||
inj.onStatus(T, AgentStatus.IDLE); // delivers "first"
|
||||
inj.onStatus(T, AgentStatus.WORKING); // picked up
|
||||
|
||||
RuntimeException thrown = assertThrows(RuntimeException.class,
|
||||
() -> inj.onStatus(T, AgentStatus.IDLE),
|
||||
"the throwable from onTurnComplete must still escape onStatus");
|
||||
assertSame(boom, thrown);
|
||||
|
||||
assertEquals(List.of("first", "second"), sent(),
|
||||
"fleetd #553: \"second\" must still be sent even though onTurnComplete threw for "
|
||||
+ "\"first\"'s completion");
|
||||
assertTrue(second.completion().isDone() && !second.completion().isCompletedExceptionally(),
|
||||
"fleetd #553: \"second\"'s delivered future must still complete normally despite "
|
||||
+ "the throw");
|
||||
}
|
||||
|
||||
private static final int READINESS_GRACE_SAMPLES = 240; // Injector.READINESS_GRACE_POLLS is private
|
||||
|
||||
@Test
|
||||
void aRuntimeExceptionFromForgetStillCompletesTheReadinessFailureFutures() {
|
||||
// fleetd #553: forget.accept fires as part of the CB-114 readiness-grace path (production
|
||||
// wires it to presence::forget). Before fleetd #553, that block called forget.accept
|
||||
// BEFORE completing the queued messages' futures, so a RuntimeException from forget left
|
||||
// every one of them pending forever even though they were already marked NOT_DELIVERED and
|
||||
// dropped from the queue. fleetd #553 completes those futures first, so forget.accept
|
||||
// throwing afterward can no longer un-complete them.
|
||||
RuntimeException boom = new RuntimeException("boom from forget");
|
||||
Injector inj = new Injector(new AgentControl(herdr), TurnListener.NOOP, _ -> false, _ -> {
|
||||
throw boom;
|
||||
});
|
||||
Injector.Delivery delivery = inj.enqueue(T, "task", TestTurnTokens.inert(T));
|
||||
|
||||
// Build the not-ready streak up to (but not past) the grace threshold.
|
||||
for (int i = 0; i < READINESS_GRACE_SAMPLES - 1; i++) inj.onStatus(T, AgentStatus.IDLE);
|
||||
assertFalse(delivery.completion().isDone(), "must not be resolved before the grace expires");
|
||||
|
||||
RuntimeException thrown = assertThrows(RuntimeException.class,
|
||||
() -> inj.onStatus(T, AgentStatus.IDLE),
|
||||
"the throwable from forget.accept must still escape onStatus");
|
||||
assertSame(boom, thrown);
|
||||
|
||||
assertTrue(delivery.completion().isCompletedExceptionally(),
|
||||
"fleetd #553: the queued message's future must still be completed even though "
|
||||
+ "forget.accept threw");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aRuntimeExceptionFromOnTurnFailedStillLeavesTheReadinessFailureFuturesCompleted() {
|
||||
// fleetd #553, acceptance test: same readiness-grace path as the forget test above, but the
|
||||
// throw comes from onTurnFailed instead. NOTE: onTurnFailed is the LAST statement in this
|
||||
// block (both before and after this ticket), so the queued messages' futures are already
|
||||
// completed by the time it runs, regardless of this fix — this test cannot be made to fail
|
||||
// against the pre-fleetd-553 code the way the onTurnComplete/forget/postAction tests can; I
|
||||
// could not find a construction where onTurnFailed's own throw is what stands between a
|
||||
// future and its completion (flagged to the lead via fleet_ask; no reply arrived before the
|
||||
// ~55s window closed, so recorded here instead). It still pins a real invariant this ticket
|
||||
// cares about — an ordinary RuntimeException from this callback must not un-complete a
|
||||
// future that was already decided — so it stays as a regression lock, not a bug-fix proof.
|
||||
RuntimeException boom = new RuntimeException("boom from onTurnFailed");
|
||||
TurnListener throwing = new TurnListener() {
|
||||
@Override
|
||||
public void onTurnComplete(String target) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onTurnFailed(String target) {
|
||||
throw boom;
|
||||
}
|
||||
};
|
||||
Injector inj = new Injector(new AgentControl(herdr), throwing, _ -> false, _ -> {
|
||||
});
|
||||
Injector.Delivery delivery = inj.enqueue(T, "task", TestTurnTokens.inert(T));
|
||||
|
||||
for (int i = 0; i < READINESS_GRACE_SAMPLES - 1; i++) inj.onStatus(T, AgentStatus.IDLE);
|
||||
|
||||
RuntimeException thrown = assertThrows(RuntimeException.class,
|
||||
() -> inj.onStatus(T, AgentStatus.IDLE),
|
||||
"the throwable from onTurnFailed must still escape onStatus");
|
||||
assertSame(boom, thrown);
|
||||
|
||||
assertTrue(delivery.completion().isCompletedExceptionally(),
|
||||
"the queued message's future must remain completed despite onTurnFailed throwing");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aRuntimeExceptionFromOnTurnCompleteWithPostActionStillUnwedgesTheTarget() {
|
||||
// fleetd #553, acceptance test: before this ticket, a RuntimeException from
|
||||
// onTurnCompleteWithPostAction skipped the `t.postTurnPending = false` reset that used to
|
||||
// sit right after it (with nothing guarding it), permanently wedging the target —
|
||||
// postTurnPending stayed true forever, so the delivery guard never passed again and
|
||||
// "second" (already queued) was never delivered, no matter how many further onStatus
|
||||
// rounds ran. fleetd #553 wraps that call so the reset always runs.
|
||||
RuntimeException boom = new RuntimeException("boom from onTurnCompleteWithPostAction");
|
||||
class ThrowingPostTurn implements TurnListener {
|
||||
@Override
|
||||
public void onTurnComplete(String target) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean hasPostTurnAction(String target) {
|
||||
return true;
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean onTurnCompleteWithPostAction(String target) {
|
||||
throw boom;
|
||||
}
|
||||
}
|
||||
Injector inj = new Injector(new AgentControl(herdr), new ThrowingPostTurn());
|
||||
inj.enqueue(T, "first", TestTurnTokens.inert(T));
|
||||
Injector.Delivery second = inj.enqueue(T, "second", TestTurnTokens.inert(T));
|
||||
|
||||
inj.onStatus(T, AgentStatus.IDLE); // delivers "first"
|
||||
inj.onStatus(T, AgentStatus.WORKING); // picked up
|
||||
|
||||
RuntimeException thrown = assertThrows(RuntimeException.class,
|
||||
() -> inj.onStatus(T, AgentStatus.IDLE), // "first"'s turn completes; post-action throws
|
||||
"the throwable from onTurnCompleteWithPostAction must still escape onStatus");
|
||||
assertSame(boom, thrown);
|
||||
assertEquals(List.of("first"), sent(), "\"second\" must not be sent in the SAME round as the throw");
|
||||
|
||||
inj.onStatus(T, AgentStatus.IDLE); // a later, unrelated round
|
||||
assertEquals(List.of("first", "second"), sent(),
|
||||
"fleetd #553: the target must not stay wedged — \"second\" must still be delivered "
|
||||
+ "once postTurnPending is reset despite the earlier throw");
|
||||
assertTrue(second.completion().isDone() && !second.completion().isCompletedExceptionally());
|
||||
}
|
||||
|
||||
/**
|
||||
* A {@link HerdrClient} that throws a non-{@link RuntimeException} {@link Error} from {@code
|
||||
* agent.send_keys} instead of delegating — the fleetd #553 case for the resubmit nudge: its own
|
||||
* catch (Injector.java, the first block after the monitor) was {@code RuntimeException}-only,
|
||||
* the same one-class-too-narrow shape #546 fixed at the send seam. Records every call it sees
|
||||
* itself, mirroring {@code ErrorOnPrompt} above, since the delegate's own recording is never
|
||||
* reached for {@code agent.send_keys}.
|
||||
*/
|
||||
private static final class ErrorOnSendKeys implements HerdrClient {
|
||||
private final FakeHerdr delegate;
|
||||
private final List<FakeHerdr.Call> calls = new java.util.concurrent.CopyOnWriteArrayList<>();
|
||||
|
||||
private ErrorOnSendKeys(FakeHerdr delegate) {
|
||||
this.delegate = delegate;
|
||||
}
|
||||
|
||||
List<FakeHerdr.Call> calls() {
|
||||
return calls;
|
||||
}
|
||||
|
||||
@Override
|
||||
public JsonNode call(String method, Object params) {
|
||||
calls.add(new FakeHerdr.Call(method, params));
|
||||
if (method.equals("agent.send_keys")) {
|
||||
throw new AssertionError("simulated non-RuntimeException resubmit failure (fleetd #553)");
|
||||
}
|
||||
return delegate.call(method, params);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
delegate.close();
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void anErrorFromTheResubmitNudgeDoesNotPreventTheSentCompletion() {
|
||||
// fleetd #553, acceptance test: the resubmit nudge's own catch was RuntimeException-only —
|
||||
// the same one-class-too-narrow shape #546 fixed at the send seam (now :382). It sits FIRST
|
||||
// after the monitor, so before this ticket an Error escaping it skipped every block below in
|
||||
// THAT round, including any `sent` completion that round might otherwise produce. Widened to
|
||||
// Throwable (matching #549's own widening) so it can no longer escape onStatus at all — the
|
||||
// delivery that already completed at send time, and the worker's later pickup, are both
|
||||
// unaffected by the Error in between.
|
||||
ErrorOnSendKeys throwing = new ErrorOnSendKeys(new FakeHerdr());
|
||||
Injector inj = new Injector(new AgentControl(throwing));
|
||||
Injector.Delivery delivery = inj.enqueue(T, "task", TestTurnTokens.inert(T));
|
||||
|
||||
inj.onStatus(T, AgentStatus.IDLE); // delivers "task"; awaiting pickup
|
||||
assertTrue(delivery.completion().isDone() && !delivery.completion().isCompletedExceptionally(),
|
||||
"the delivery completes at send time, before any resubmit nudge is attempted");
|
||||
|
||||
assertDoesNotThrow(() -> inj.onStatus(T, AgentStatus.IDLE), // still idle -> resubmit; Error thrown+caught
|
||||
"fleetd #553: an Error from the resubmit nudge must be caught inside onStatus, not "
|
||||
+ "escape it");
|
||||
|
||||
inj.onStatus(T, AgentStatus.WORKING); // the worker's real pickup must still resolve cleanly
|
||||
long promptCalls = throwing.calls().stream().filter(c -> c.method().equals("agent.prompt")).count();
|
||||
assertEquals(1, promptCalls,
|
||||
"sent exactly once, unaffected by the resubmit Error in between");
|
||||
}
|
||||
|
||||
@Test
|
||||
void ordinarySuccessStillCompletesExactlyOnceAfterOnDelivered() {
|
||||
// Acceptance: the ordinary success path is unchanged by fleetd #553 — the future completes
|
||||
// normally exactly once, and onDelivered still runs before it (CB-115's pane baseline).
|
||||
List<String> events = new ArrayList<>();
|
||||
TurnListener listener = new TurnListener() {
|
||||
@Override
|
||||
public void onTurnComplete(String target) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onDelivered(String target, dev.ltms.fleet.msg.TurnToken token) {
|
||||
events.add("onDelivered");
|
||||
}
|
||||
};
|
||||
Injector inj = new Injector(new AgentControl(herdr), listener);
|
||||
CompletableFuture<Void> f = inj.enqueue(T, "hello", TestTurnTokens.inert(T)).completion();
|
||||
f.whenComplete((v, ex) -> events.add("completed"));
|
||||
|
||||
inj.onStatus(T, AgentStatus.IDLE);
|
||||
|
||||
assertEquals(List.of("onDelivered", "completed"), events,
|
||||
"onDelivered must still run before the future completes, unchanged by fleetd #553");
|
||||
assertTrue(f.isDone() && !f.isCompletedExceptionally());
|
||||
}
|
||||
|
||||
// Acceptance: the ordinary failure path (a HerdrException at the send seam still completes the
|
||||
// future exceptionally with that exception) is already covered, unchanged, by the existing
|
||||
// aHerdrExceptionFromSendStillProducesNotDeliveredUnchanged test above (fleetd #546) — its
|
||||
// codepath is untouched by fleetd #553's try/finally, since sendError != null is set before the
|
||||
// try block begins and that branch never threw to begin with.
|
||||
|
||||
// --- fleetd #553, sharpened acceptance (ticket comments 16870/16884/16890/16903): completing
|
||||
// sent.delivered() is only HALF the job. See the two tests below. ---
|
||||
|
||||
@Test
|
||||
void aRuntimeExceptionFromOnTurnCompleteStillRegistersTheRendezvousWaiterForTheNextDelivery() {
|
||||
// fleetd #553's real invariant, not just the completion of sent.delivered(). There are TWO
|
||||
// futures at stake: sent.delivered() (the delivery future) and sent.token().waiter() (the
|
||||
// rendezvous waiter a blocking fleet_send actually waits on for the worker's ANSWER). The
|
||||
// waiter is registered only by turnListener.onDelivered() — production wires this to
|
||||
// CompletionResolver.captureBaseline, which does the inFlight.put() — and that call lives
|
||||
// inside the very `if (sent != null)` block a plain "complete the future" backstop does not
|
||||
// reach. A finally that only completes sent.delivered() converts a hang into a HANG WITH A
|
||||
// SUCCESS RECEIPT: the caller is told the send landed, then waits out its full timeout for
|
||||
// an answer that can never resolve, because CompletionResolver.resolve (:301-308) finds no
|
||||
// inFlight entry and returns silently. This test is red on that half-fix, and green only
|
||||
// once the finally also runs onDelivered() when the normal block never got the chance.
|
||||
//
|
||||
// The scenario: onTurnComplete throws while completing "first"'s turn, in the SAME onStatus
|
||||
// round that (per Injector.java :361-391) then peeks and delivers "second" — the normal
|
||||
// case, not a corner, since the :365 assignment is what lets the :376 delivery guard pass.
|
||||
// turnListener.onTurnComplete runs BEFORE the `if (sent != null)` block, so the throw here
|
||||
// means onDelivered for "second" is never reached on the normal path — only the finally
|
||||
// backstop can register it.
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(new FakeHerdr()), rendezvous,
|
||||
ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
RuntimeException boom = new RuntimeException("boom from onTurnComplete");
|
||||
TurnListener throwingOnTurnCompleteOnly = new TurnListener() {
|
||||
@Override
|
||||
public void onDelivered(String target, TurnToken token) {
|
||||
resolver.onDelivered(target, token); // the real registration, exactly as production wires it
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onTurnComplete(String target) {
|
||||
throw boom; // "first"'s completion callback, thrown before "second" is delivered
|
||||
}
|
||||
};
|
||||
Injector inj = new Injector(new AgentControl(herdr), throwingOnTurnCompleteOnly);
|
||||
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.open("session-a");
|
||||
TurnToken secondToken = new TurnToken(T, waiter);
|
||||
inj.enqueue(T, "first", TestTurnTokens.inert(T));
|
||||
Injector.Delivery second = inj.enqueue(T, "second", secondToken);
|
||||
|
||||
inj.onStatus(T, AgentStatus.IDLE); // delivers "first"
|
||||
inj.onStatus(T, AgentStatus.WORKING); // picked up
|
||||
|
||||
RuntimeException thrown = assertThrows(RuntimeException.class,
|
||||
() -> inj.onStatus(T, AgentStatus.IDLE), // "first" completes (throws); "second" delivered
|
||||
"the original throwable from onTurnComplete must still escape onStatus");
|
||||
assertSame(boom, thrown, "must be the EXACT throwable, not a wrapper or a different instance");
|
||||
|
||||
assertEquals(List.of("first", "second"), sent(),
|
||||
"\"second\" must still be sent even though onTurnComplete threw for \"first\"");
|
||||
assertTrue(second.completion().isDone() && !second.completion().isCompletedExceptionally(),
|
||||
"\"second\"'s delivery future must still complete despite the throw");
|
||||
|
||||
CompletionResolver.InFlight inFlight = resolver.inFlight(T);
|
||||
assertNotNull(inFlight,
|
||||
"fleetd #553: the finally must register the new turn's waiter (via onDelivered), not "
|
||||
+ "just complete its delivery future — otherwise a blocking fleet_send is told "
|
||||
+ "its message landed and then waits out the full timeout for an answer that "
|
||||
+ "can never resolve");
|
||||
assertSame(waiter, inFlight.waiter(),
|
||||
"the registered waiter must be exactly \"second\"'s waiter, not some other one");
|
||||
}
|
||||
|
||||
@Test
|
||||
void anOnDeliveredThrowAfterItsOwnRegistrationDoesNotRunASecondTime() {
|
||||
// fleetd #553 comment 16903: `sentHandled` must be set to true BEFORE onDelivered() runs,
|
||||
// not after. Setting it after would mean a throw from onDelivered() PART WAY THROUGH — e.g.
|
||||
// after CompletionResolver's own captureBaseline has already done its inFlight.put() — still
|
||||
// leaves the flag false, so the finally backstop (seeing "not handled") calls onDelivered() a
|
||||
// SECOND time. That second captureBaseline runs later, over a pane that may already have
|
||||
// absorbed this turn's output, and CompletionResolver.resolve's own baseline-equals-tail
|
||||
// suppression (:364-369, which deliberately keeps the in-flight record on a match) can then
|
||||
// drop every future completion for the turn, permanently. This assertion is red on a
|
||||
// `sentHandled = true` placed AFTER the onDelivered() call (onDelivered runs twice) and green
|
||||
// on the correct placement (runs exactly once) — regardless of what onDelivered itself did.
|
||||
AtomicInteger onDeliveredCalls = new AtomicInteger();
|
||||
RuntimeException boom = new RuntimeException("boom from onDelivered, after its own registration ran");
|
||||
TurnListener listener = new TurnListener() {
|
||||
@Override
|
||||
public void onTurnComplete(String target) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onDelivered(String target, TurnToken token) {
|
||||
onDeliveredCalls.incrementAndGet(); // stand-in for CompletionResolver's inFlight.put
|
||||
throw boom; // then fail, as if a LATER step inside onDelivered blew up
|
||||
}
|
||||
};
|
||||
Injector inj = new Injector(new AgentControl(herdr), listener);
|
||||
inj.enqueue(T, "task", TestTurnTokens.inert(T));
|
||||
|
||||
RuntimeException thrown = assertThrows(RuntimeException.class,
|
||||
() -> inj.onStatus(T, AgentStatus.IDLE), // delivers "task"; onDelivered throws
|
||||
"the original throwable from onDelivered must still escape onStatus");
|
||||
assertSame(boom, thrown);
|
||||
|
||||
assertEquals(1, onDeliveredCalls.get(),
|
||||
"fleetd #553: onDelivered must be called exactly once — a `sentHandled` flag set "
|
||||
+ "AFTER the call (rather than before) would leave it false here and the "
|
||||
+ "finally backstop would call onDelivered a second time");
|
||||
}
|
||||
|
||||
@Test
|
||||
void anOnDeliveredThrowOnTheNormalPathStillCompletesTheDeliveryFuture() {
|
||||
// fleetd #553, lead review of PR #557 (ticket comment 16916): `sentHandled` must guard ONLY
|
||||
// the onDelivered RE-CALL in the finally, never the future completion alongside it. Gating
|
||||
// BOTH behind `!sentHandled` (the shape this test is red against) misses the one path the
|
||||
// flag's own correct placement creates: `sentHandled` is set to true FIRST, before the
|
||||
// onDelivered() call, inside the `if (sent != null)` block above (see the previous test —
|
||||
// that placement is right and must not change). So when onDelivered() itself throws on that
|
||||
// NORMAL path, `sentHandled` already reads true by the time control reaches the finally, and
|
||||
// a `sent != null && !sentHandled` guard around the WHOLE recovery — completion included —
|
||||
// skips it entirely. The message was typed into the target's pane and taken off the queue
|
||||
// inside the monitor, same as any other delivery, so its future is left pending forever: the
|
||||
// exact defect this ticket exists to close, just reached from a different throwing call.
|
||||
//
|
||||
// The fix splits the one flag's two jobs: `!sentHandled` keeps gating only the onDelivered
|
||||
// call (so the exactly-once guarantee in the test above still holds — CompletableFuture.
|
||||
// complete/completeExceptionally are idempotent, so completing unconditionally here is a
|
||||
// no-op on the ordinary path, where the `if (sent != null)` block already completed it.
|
||||
RuntimeException boom = new RuntimeException("boom from onDelivered on the normal path");
|
||||
TurnListener listener = new TurnListener() {
|
||||
@Override
|
||||
public void onTurnComplete(String target) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onDelivered(String target, TurnToken token) {
|
||||
throw boom;
|
||||
}
|
||||
};
|
||||
Injector inj = new Injector(new AgentControl(herdr), listener);
|
||||
CompletableFuture<Void> delivered = inj.enqueue(T, "task", TestTurnTokens.inert(T)).completion();
|
||||
|
||||
RuntimeException thrown = assertThrows(RuntimeException.class,
|
||||
() -> inj.onStatus(T, AgentStatus.IDLE), // delivers "task"; onDelivered throws
|
||||
"the original throwable from onDelivered must still escape onStatus");
|
||||
assertSame(boom, thrown);
|
||||
|
||||
assertTrue(delivered.isDone(),
|
||||
"fleetd #553: \"task\" was actually delivered — typed into the pane and taken off "
|
||||
+ "the queue inside the monitor — so its delivery future must be completed on "
|
||||
+ "every path out of onStatus, including the one where onDelivered itself is "
|
||||
+ "what threw. Leaving it pending here is a hang, not a fix");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,99 @@
|
||||
package dev.ltms.fleet.inject;
|
||||
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.msg.TestTurnTokens;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.lang.reflect.Field;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotSame;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
class StatusPollerResilienceTest {
|
||||
|
||||
@Test
|
||||
void anErrorForOneTargetDoesNotStopPollingTheNextTarget() throws Exception {
|
||||
FakeHerdr fake = new FakeHerdr().withAgent("worker", "term_b", "w2:p8", "w2:t8");
|
||||
CountDownLatch errorThrown = new CountDownLatch(1);
|
||||
AgentControl agents = new AgentControl(new ErrorOnceForFirstTarget(fake, errorThrown));
|
||||
Injector injector = new Injector(agents);
|
||||
StatusPoller poller = new StatusPoller(agents, injector, 1);
|
||||
poller.start();
|
||||
try {
|
||||
injector.enqueue("term_a", "first", TestTurnTokens.inert("term_a"));
|
||||
assertTrue(errorThrown.await(2, TimeUnit.SECONDS),
|
||||
"the first target must throw its test Error");
|
||||
|
||||
CompletableFuture<Void> delivered =
|
||||
injector.enqueue("term_b", "second", TestTurnTokens.inert("term_b")).completion();
|
||||
delivered.get(2, TimeUnit.SECONDS);
|
||||
} finally {
|
||||
poller.stop();
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void anAbnormalExitClearsRunningSoStartCreatesANewLoop() throws Exception {
|
||||
StatusPoller poller = new StatusPoller(new AgentControl(new FakeHerdr()), new Injector(new AgentControl(new FakeHerdr())), -1);
|
||||
poller.start();
|
||||
Thread first = threadOf(poller);
|
||||
first.join(2000);
|
||||
assertFalse(runningOf(poller), "an abnormal loop exit must clear running");
|
||||
|
||||
poller.start();
|
||||
Thread restarted = threadOf(poller);
|
||||
try {
|
||||
assertNotSame(first, restarted, "start() must create a new loop after an abnormal exit");
|
||||
restarted.join(2000);
|
||||
} finally {
|
||||
poller.stop();
|
||||
}
|
||||
}
|
||||
|
||||
private static Thread threadOf(StatusPoller poller) throws ReflectiveOperationException {
|
||||
Field field = StatusPoller.class.getDeclaredField("thread");
|
||||
field.setAccessible(true);
|
||||
return (Thread) field.get(poller);
|
||||
}
|
||||
|
||||
private static boolean runningOf(StatusPoller poller) throws ReflectiveOperationException {
|
||||
Field field = StatusPoller.class.getDeclaredField("running");
|
||||
field.setAccessible(true);
|
||||
return field.getBoolean(poller);
|
||||
}
|
||||
|
||||
private static final class ErrorOnceForFirstTarget implements HerdrClient {
|
||||
private final FakeHerdr delegate;
|
||||
private final CountDownLatch errorThrown;
|
||||
private final AtomicBoolean first = new AtomicBoolean(true);
|
||||
|
||||
private ErrorOnceForFirstTarget(FakeHerdr delegate, CountDownLatch errorThrown) {
|
||||
this.delegate = delegate;
|
||||
this.errorThrown = errorThrown;
|
||||
}
|
||||
|
||||
@Override
|
||||
public JsonNode call(String method, Object params) {
|
||||
if (method.equals("agent.get") && params instanceof Map<?, ?> map
|
||||
&& "w2:p7".equals(map.get("target")) && first.compareAndSet(true, false)) {
|
||||
errorThrown.countDown();
|
||||
throw new AssertionError("test Error from the first poll target");
|
||||
}
|
||||
return delegate.call(method, params);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
delegate.close();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,23 +1,22 @@
|
||||
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;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.herdr.HerdrException;
|
||||
import dev.ltms.fleet.testing.CapturedLog;
|
||||
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;
|
||||
import java.nio.file.Path;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.LongSupplier;
|
||||
@@ -862,25 +861,16 @@ class LeadRolloverTest {
|
||||
|
||||
// ---- 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 CapturedLog attachLog() {
|
||||
return CapturedLog.at(LeadRollover.class, Level.DEBUG);
|
||||
}
|
||||
|
||||
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()
|
||||
private static ILoggingEvent lastEventContaining(List<ILoggingEvent> events, String substring) {
|
||||
return events.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()));
|
||||
+ "\"; got: " + events.stream().map(ILoggingEvent::getFormattedMessage).toList()));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -914,15 +904,14 @@ class LeadRolloverTest {
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
LeadRollover rollover = newRollover(flipsAfterClear, config, () -> clock.addAndGet(500));
|
||||
|
||||
ListAppender<ILoggingEvent> events = attachLog();
|
||||
try {
|
||||
try (CapturedLog log = attachLog()) {
|
||||
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");
|
||||
ILoggingEvent event = lastEventContaining(log.events(), "NOT sending bootstrapText");
|
||||
assertEquals(Level.WARN, event.getLevel());
|
||||
String message = event.getFormattedMessage();
|
||||
assertTrue(message.contains("configured=1s"), "must label the configured budget: " + message);
|
||||
@@ -934,8 +923,6 @@ class LeadRolloverTest {
|
||||
+ message);
|
||||
assertFalse(message.contains("within 1s"), "must not present the configured budget as if "
|
||||
+ "it were the measured wait duration: " + message);
|
||||
} finally {
|
||||
detachLog(events);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -954,15 +941,14 @@ class LeadRolloverTest {
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
LeadRollover rollover = newRollover(herdr, config, () -> clock.addAndGet(500));
|
||||
|
||||
ListAppender<ILoggingEvent> events = attachLog();
|
||||
try {
|
||||
try (CapturedLog log = attachLog()) {
|
||||
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");
|
||||
ILoggingEvent event = lastEventContaining(log.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)");
|
||||
@@ -982,8 +968,6 @@ class LeadRolloverTest {
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1000,15 +984,14 @@ class LeadRolloverTest {
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
LeadRollover rollover = newRollover(herdr, config, () -> clock.addAndGet(500));
|
||||
|
||||
ListAppender<ILoggingEvent> events = attachLog();
|
||||
try {
|
||||
try (CapturedLog log = attachLog()) {
|
||||
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");
|
||||
ILoggingEvent event = lastEventContaining(log.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
|
||||
@@ -1018,8 +1001,6 @@ class LeadRolloverTest {
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1040,8 +1021,7 @@ class LeadRolloverTest {
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
LeadRollover rollover = newRollover(herdr, config, () -> clock.addAndGet(500));
|
||||
|
||||
ListAppender<ILoggingEvent> events = attachLog();
|
||||
try {
|
||||
try (CapturedLog log = attachLog()) {
|
||||
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
|
||||
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
|
||||
|
||||
@@ -1049,7 +1029,7 @@ class LeadRolloverTest {
|
||||
+ "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");
|
||||
ILoggingEvent event = lastEventContaining(log.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);
|
||||
@@ -1058,8 +1038,6 @@ class LeadRolloverTest {
|
||||
+ "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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,17 +1,17 @@
|
||||
package dev.ltms.fleet.msg;
|
||||
|
||||
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.rabbitmq.client.Channel;
|
||||
import com.rabbitmq.client.ConnectionFactory;
|
||||
import com.rabbitmq.client.impl.DefaultExceptionHandler;
|
||||
import dev.ltms.fleet.testing.CapturedLog;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.lang.reflect.Proxy;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
@@ -32,21 +32,17 @@ class AmqpConnectionFailureLoggerTest {
|
||||
assertEquals(AmqpConnectionFailureLogger.REPLY_INBOX, inboxHandler.connectionName());
|
||||
assertEquals(AmqpConnectionFailureLogger.LEAD_MAILBOX, mailboxHandler.connectionName());
|
||||
|
||||
ListAppender<ILoggingEvent> inboxEvents = attach(AmqpReplyInbox.class);
|
||||
ListAppender<ILoggingEvent> mailboxEvents = attach(LeadMailbox.class);
|
||||
IllegalStateException inboxFailure = new IllegalStateException("inbox failure");
|
||||
IllegalStateException mailboxFailure = new IllegalStateException("mailbox failure");
|
||||
try {
|
||||
try (CapturedLog inboxLog = attach(AmqpReplyInbox.class);
|
||||
CapturedLog mailboxLog = attach(LeadMailbox.class)) {
|
||||
inboxHandler.handleUnexpectedConnectionDriverException(null, inboxFailure);
|
||||
mailboxHandler.handleConnectionRecoveryException(null, mailboxFailure);
|
||||
|
||||
assertError(inboxEvents, "AMQP connection fleetd-reply-inbox: An unexpected connection driver error occurred",
|
||||
assertError(inboxLog.events(), "AMQP connection fleetd-reply-inbox: An unexpected connection driver error occurred",
|
||||
inboxFailure, "inbox failure line");
|
||||
assertError(mailboxEvents, "AMQP connection fleetd-lead-mailbox: Caught an exception during connection recovery!",
|
||||
assertError(mailboxLog.events(), "AMQP connection fleetd-lead-mailbox: Caught an exception during connection recovery!",
|
||||
mailboxFailure, "mailbox recovery line");
|
||||
} finally {
|
||||
detach(AmqpReplyInbox.class, inboxEvents);
|
||||
detach(LeadMailbox.class, mailboxEvents);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -54,17 +50,14 @@ class AmqpConnectionFailureLoggerTest {
|
||||
void connectionResetKeepsForgivingHandlerWarningSemantics() {
|
||||
AmqpConnectionFailureLogger handler = new AmqpConnectionFailureLogger(
|
||||
AmqpConnectionFailureLogger.REPLY_INBOX, LoggerFactory.getLogger(AmqpReplyInbox.class));
|
||||
ListAppender<ILoggingEvent> events = attach(AmqpReplyInbox.class);
|
||||
try {
|
||||
try (CapturedLog log = attach(AmqpReplyInbox.class)) {
|
||||
handler.handleUnexpectedConnectionDriverException(null, new IOException("Connection reset"));
|
||||
assertEquals(1, events.list.size(), "the handler must still log a reset");
|
||||
ILoggingEvent event = events.list.getFirst();
|
||||
assertEquals(1, log.events().size(), "the handler must still log a reset");
|
||||
ILoggingEvent event = log.events().getFirst();
|
||||
assertEquals(Level.WARN, event.getLevel(), "ForgivingExceptionHandler logs connection resets at WARN");
|
||||
assertEquals("AMQP connection fleetd-reply-inbox: An unexpected connection driver error occurred "
|
||||
+ "(Exception message: Connection reset)", event.getFormattedMessage());
|
||||
assertTrue(event.getThrowableProxy() == null, "ForgivingExceptionHandler does not attach a reset stack trace");
|
||||
} finally {
|
||||
detach(AmqpReplyInbox.class, events);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -103,17 +96,8 @@ class AmqpConnectionFailureLoggerTest {
|
||||
"all exception-handling methods must remain inherited from DefaultExceptionHandler");
|
||||
}
|
||||
|
||||
private static ListAppender<ILoggingEvent> attach(Class<?> owner) {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(owner);
|
||||
logger.setLevel(Level.DEBUG);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
return appender;
|
||||
}
|
||||
|
||||
private static void detach(Class<?> owner, ListAppender<ILoggingEvent> appender) {
|
||||
((Logger) LoggerFactory.getLogger(owner)).detachAppender(appender);
|
||||
private static CapturedLog attach(Class<?> owner) {
|
||||
return CapturedLog.at(owner, Level.DEBUG);
|
||||
}
|
||||
|
||||
private static AmqpConnectionFailureLogger installedStrictHandler(ConnectionFactory factory, String connection) {
|
||||
@@ -123,9 +107,9 @@ class AmqpConnectionFailureLoggerTest {
|
||||
return assertInstanceOf(AmqpConnectionFailureLogger.class, factory.getExceptionHandler());
|
||||
}
|
||||
|
||||
private static void assertError(ListAppender<ILoggingEvent> events, String message, Throwable cause, String name) {
|
||||
assertEquals(1, events.list.size(), name);
|
||||
ILoggingEvent event = events.list.getFirst();
|
||||
private static void assertError(List<ILoggingEvent> events, String message, Throwable cause, String name) {
|
||||
assertEquals(1, events.size(), name);
|
||||
ILoggingEvent event = events.getFirst();
|
||||
assertEquals(Level.ERROR, event.getLevel(), name);
|
||||
assertEquals(message, event.getFormattedMessage(), name);
|
||||
assertEquals(cause.toString(), event.getThrowableProxy().getClassName() + ": "
|
||||
|
||||
@@ -1,16 +1,13 @@
|
||||
package dev.ltms.fleet.session;
|
||||
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.Logger;
|
||||
import ch.qos.logback.classic.LoggerContext;
|
||||
import ch.qos.logback.classic.spi.IThrowableProxy;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import dev.ltms.fleet.testing.CapturedLog;
|
||||
import org.junit.jupiter.api.AfterEach;
|
||||
import org.junit.jupiter.api.BeforeEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.nio.charset.StandardCharsets;
|
||||
@@ -489,36 +486,33 @@ class GitWorktreesTest {
|
||||
// ---- CB-189: broader remote-URL coverage — every remote, both fetch and push URLs, any
|
||||
// non-SSH scheme. Reporting only, additive to the origin/https strip-and-refuse tests above. ----
|
||||
|
||||
private Logger reportingLogger;
|
||||
private ListAppender<ILoggingEvent> reportingAppender;
|
||||
private CapturedLog reportingLog;
|
||||
|
||||
/** {@link GitWorktrees}'s own logger, captured fresh for each test so assertions never see a
|
||||
* message left over from a previous test. */
|
||||
* message left over from a previous test. fleetd #529: {@link CapturedLog#close} restores the
|
||||
* level it captured here (whatever it truly was before this test, not just WARN), so a test
|
||||
* below that further lowers the level to INFO for its own assertion (via {@link
|
||||
* #reportingLog}'s {@link CapturedLog#setLevel}) can never leak that INFO pin past its own
|
||||
* {@code @AfterEach} — every test's window is self-contained. */
|
||||
@BeforeEach
|
||||
void attachReportingLogCapture() {
|
||||
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
|
||||
reportingLogger = ctx.getLogger(GitWorktrees.class);
|
||||
reportingLogger.setLevel(Level.WARN);
|
||||
reportingAppender = new ListAppender<>();
|
||||
reportingAppender.setContext(ctx);
|
||||
reportingAppender.start();
|
||||
reportingLogger.addAppender(reportingAppender);
|
||||
reportingLog = CapturedLog.at(GitWorktrees.class, Level.WARN);
|
||||
}
|
||||
|
||||
@AfterEach
|
||||
void detachReportingLogCapture() {
|
||||
reportingLogger.detachAppender(reportingAppender);
|
||||
reportingLog.close();
|
||||
}
|
||||
|
||||
private List<String> capturedMessages() {
|
||||
return reportingAppender.list.stream().map(ILoggingEvent::getFormattedMessage).toList();
|
||||
return reportingLog.events().stream().map(ILoggingEvent::getFormattedMessage).toList();
|
||||
}
|
||||
|
||||
/** Asserts {@code secret} appears in no captured message, and in no attached exception's
|
||||
* message either — the constraint is that a credential must never reach a log, however it
|
||||
* would have gotten there. */
|
||||
private void assertNoLeak(String secret) {
|
||||
for (ILoggingEvent event : reportingAppender.list) {
|
||||
for (ILoggingEvent event : reportingLog.events()) {
|
||||
assertFalse(event.getFormattedMessage().contains(secret),
|
||||
"log message leaked a credential (" + secret + "): " + event.getFormattedMessage());
|
||||
IThrowableProxy thrown = event.getThrowableProxy();
|
||||
@@ -596,7 +590,7 @@ class GitWorktreesTest {
|
||||
|
||||
new GitWorktrees(tmp.resolve("wts").toString()).add(repo.toString(), "cb-189-d", "HEAD");
|
||||
|
||||
assertTrue(reportingAppender.list.isEmpty(),
|
||||
assertTrue(reportingLog.events().isEmpty(),
|
||||
"expected no report for an ssh remote and a credential-free https remote, got:\n"
|
||||
+ capturedMessages());
|
||||
}
|
||||
@@ -812,7 +806,7 @@ class GitWorktreesTest {
|
||||
/** Criterion 1, all three present: the summary names the denominator and every neutralized file. */
|
||||
@Test
|
||||
void isolateToolSurfaceLogsAllThreeConfigsNeutralized(@TempDir Path tmp) throws Exception {
|
||||
reportingLogger.setLevel(Level.INFO);
|
||||
reportingLog.setLevel(Level.INFO);
|
||||
Path repo = initRepoWithAllThreeConfigs(tmp.resolve("repo"));
|
||||
|
||||
new GitWorktrees(tmp.resolve("wts").toString()).add(repo.toString(), "cb-134-log-all", "HEAD");
|
||||
@@ -827,7 +821,7 @@ class GitWorktreesTest {
|
||||
/** Criterion 1, two absent: the summary must still name the denominator and say why. */
|
||||
@Test
|
||||
void isolateToolSurfaceLogsAbsentConfigsWithReason(@TempDir Path tmp) throws Exception {
|
||||
reportingLogger.setLevel(Level.INFO);
|
||||
reportingLog.setLevel(Level.INFO);
|
||||
Path repo = initRepo(tmp.resolve("repo")); // only .mcp.json + README committed
|
||||
|
||||
new GitWorktrees(tmp.resolve("wts").toString()).add(repo.toString(), "cb-134-log-partial", "HEAD");
|
||||
@@ -1435,7 +1429,7 @@ class GitWorktreesTest {
|
||||
/** Criterion 2: both candidates present — the summary line names both and the denominator. */
|
||||
@Test
|
||||
void overlayParityLogsBothCopiedWhenBothCandidatesArePresent(@TempDir Path tmp) throws Exception {
|
||||
reportingLogger.setLevel(Level.INFO);
|
||||
reportingLog.setLevel(Level.INFO);
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
Files.writeString(repo.resolve(".env"), "A=1\n");
|
||||
Files.writeString(repo.resolve(".envrc"), "export A=1\n");
|
||||
@@ -1452,7 +1446,7 @@ class GitWorktreesTest {
|
||||
* denominator, and why the other candidate was not copied. */
|
||||
@Test
|
||||
void overlayParityLogsOneCopiedOneAbsent(@TempDir Path tmp) throws Exception {
|
||||
reportingLogger.setLevel(Level.INFO);
|
||||
reportingLog.setLevel(Level.INFO);
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
Files.writeString(repo.resolve(".env"), "A=1\n");
|
||||
// .envrc deliberately not created — the absent candidate.
|
||||
@@ -1472,7 +1466,7 @@ class GitWorktreesTest {
|
||||
* change, silently, unless this line told it so beforehand. */
|
||||
@Test
|
||||
void overlayParityLogsSkipWorktreeConsequenceForATrackedFile(@TempDir Path tmp) throws Exception {
|
||||
reportingLogger.setLevel(Level.INFO);
|
||||
reportingLog.setLevel(Level.INFO);
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
Files.writeString(repo.resolve(".env"), "A=1\n");
|
||||
git(repo, "add", ".env");
|
||||
@@ -1495,7 +1489,7 @@ class GitWorktreesTest {
|
||||
/** Criterion 5: null and empty overlay lists return quietly — no exception, no log noise. */
|
||||
@Test
|
||||
void overlayParityWithNoCandidatesLogsNothing(@TempDir Path tmp) throws Exception {
|
||||
reportingLogger.setLevel(Level.INFO);
|
||||
reportingLog.setLevel(Level.INFO);
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
Path wt = bareWorktree(repo, tmp.resolve("wt"), "cb134-empty");
|
||||
GitWorktrees worktrees = new GitWorktrees(tmp.resolve("wts").toString());
|
||||
@@ -1503,7 +1497,7 @@ class GitWorktreesTest {
|
||||
worktrees.overlayParity(repo.toString(), wt.toString(), null);
|
||||
worktrees.overlayParity(repo.toString(), wt.toString(), List.of());
|
||||
|
||||
assertTrue(reportingAppender.list.isEmpty(),
|
||||
assertTrue(reportingLog.events().isEmpty(),
|
||||
"a null/empty overlay must log nothing, got:\n" + capturedMessages());
|
||||
}
|
||||
|
||||
@@ -1683,7 +1677,7 @@ class GitWorktreesTest {
|
||||
* denominator, what was seeded, and what was kept because the repo already had it. */
|
||||
@Test
|
||||
void seedSkillsLogsSeededAndKept(@TempDir Path tmp) throws Exception {
|
||||
reportingLogger.setLevel(Level.INFO);
|
||||
reportingLog.setLevel(Level.INFO);
|
||||
Path repo = tmp.resolve("repo");
|
||||
Files.createDirectories(repo);
|
||||
git(repo, "init", "-q", "-b", "main");
|
||||
|
||||
@@ -1,9 +1,7 @@
|
||||
package dev.ltms.fleet.session;
|
||||
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.LoggerContext;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import dev.ltms.fleet.auth.MemberRegistry;
|
||||
import dev.ltms.fleet.auth.MemberLifecycle;
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
@@ -25,6 +23,7 @@ import dev.ltms.fleet.peer.SpawnRequest;
|
||||
import dev.ltms.fleet.placement.BackendQuarantine;
|
||||
import dev.ltms.fleet.placement.PlacementDecision;
|
||||
import dev.ltms.fleet.placement.PlacementPolicies;
|
||||
import dev.ltms.fleet.testing.CapturedLog;
|
||||
import org.junit.jupiter.api.AfterAll;
|
||||
import org.junit.jupiter.api.BeforeAll;
|
||||
import org.junit.jupiter.api.MethodOrderer;
|
||||
@@ -123,54 +122,8 @@ class SessionManagerTest {
|
||||
return new SessionManager(workers, worktrees, clock);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #525: captures a logger's output and, on {@link #close}, restores <em>both</em> the
|
||||
* appender and the level to what they were before. A bare {@code addAppender}/{@code
|
||||
* setLevel} pair whose {@code finally} only detaches the appender leaves the level pinned —
|
||||
* {@code ch.qos.logback.classic.Logger} instances are cached per class and shared across the
|
||||
* whole JVM, so a level set by one test in this class is still in effect for every test that
|
||||
* runs after it, in this class or any other. try-with-resources makes "restored the appender
|
||||
* but not the level" impossible to write, because there is only one thing to close.
|
||||
*/
|
||||
private static final class CapturedLog implements AutoCloseable {
|
||||
private final ch.qos.logback.classic.Logger logger;
|
||||
private final Level originalLevel;
|
||||
private final ListAppender<ILoggingEvent> appender;
|
||||
|
||||
private CapturedLog(Class<?> loggerClass, Level pinnedLevel) {
|
||||
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
|
||||
this.logger = (ch.qos.logback.classic.Logger) LoggerFactory.getLogger(loggerClass);
|
||||
this.originalLevel = logger.getLevel();
|
||||
this.appender = new ListAppender<>();
|
||||
appender.setContext(ctx);
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
if (pinnedLevel != null) {
|
||||
logger.setLevel(pinnedLevel);
|
||||
}
|
||||
}
|
||||
|
||||
/** Capture {@code loggerClass}'s output, pinning its level to {@code pinnedLevel} for the
|
||||
* duration of the try-with-resources block. */
|
||||
static CapturedLog at(Class<?> loggerClass, Level pinnedLevel) {
|
||||
return new CapturedLog(loggerClass, pinnedLevel);
|
||||
}
|
||||
|
||||
/** Capture {@code loggerClass}'s output without changing its level. */
|
||||
static CapturedLog of(Class<?> loggerClass) {
|
||||
return new CapturedLog(loggerClass, null);
|
||||
}
|
||||
|
||||
List<ILoggingEvent> events() {
|
||||
return appender.list;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
logger.detachAppender(appender);
|
||||
logger.setLevel(originalLevel);
|
||||
}
|
||||
}
|
||||
// fleetd #529: CapturedLog moved to dev.ltms.fleet.testing.CapturedLog (imported above) so
|
||||
// every test file shares one implementation instead of each hand-rolling its own capture.
|
||||
|
||||
/**
|
||||
* CB-581: a {@link Worktrees} test double whose {@code hasUncommitted} and {@code remove} can
|
||||
|
||||
@@ -0,0 +1,91 @@
|
||||
package dev.ltms.fleet.session;
|
||||
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
import dev.ltms.fleet.member.ClaudeCodeLauncher;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.lang.reflect.Field;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotSame;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
class SessionReaperResilienceTest {
|
||||
|
||||
@Test
|
||||
void anErrorInOneIterationDoesNotStopTheNextIteration() throws Exception {
|
||||
CountDownLatch errorThrown = new CountDownLatch(1);
|
||||
CountDownLatch nextIteration = new CountDownLatch(1);
|
||||
AtomicBoolean first = new AtomicBoolean(true);
|
||||
LongSupplier clock = () -> {
|
||||
if (first.compareAndSet(true, false)) {
|
||||
errorThrown.countDown();
|
||||
throw new AssertionError("test Error from the first reap iteration");
|
||||
}
|
||||
nextIteration.countDown();
|
||||
return System.nanoTime();
|
||||
};
|
||||
SessionReaper reaper = new SessionReaper(sessionManager(clock), 60, 1);
|
||||
reaper.start();
|
||||
try {
|
||||
assertTrue(errorThrown.await(2, TimeUnit.SECONDS),
|
||||
"the first reap iteration must throw its test Error");
|
||||
assertTrue(nextIteration.await(2, TimeUnit.SECONDS),
|
||||
"the reaper must continue to the next iteration after an Error");
|
||||
} finally {
|
||||
reaper.stop();
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void anAbnormalExitClearsRunningSoStartCreatesANewLoop() throws Exception {
|
||||
SessionReaper reaper = new SessionReaper(sessionManager(System::nanoTime), 60, -1);
|
||||
reaper.start();
|
||||
Thread first = threadOf(reaper);
|
||||
first.join(2000);
|
||||
assertFalse(runningOf(reaper), "an abnormal loop exit must clear running");
|
||||
|
||||
reaper.start();
|
||||
Thread restarted = threadOf(reaper);
|
||||
try {
|
||||
assertNotSame(first, restarted, "start() must create a new loop after an abnormal exit");
|
||||
restarted.join(2000);
|
||||
} finally {
|
||||
reaper.stop();
|
||||
}
|
||||
}
|
||||
|
||||
private static SessionManager sessionManager(LongSupplier clock) {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FleetConfig.Profile cfg = new FleetConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
|
||||
List.of("ccs", "ltms-local"), "tab", "fleetd-workers",
|
||||
"worker: {profile} #{n}", null, null, null);
|
||||
ClaudeCodeLauncher launcher = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
|
||||
return new SessionManager(launcher, new FakeWorktrees(), clock);
|
||||
}
|
||||
|
||||
private static Thread threadOf(SessionReaper reaper) throws ReflectiveOperationException {
|
||||
Field field = SessionReaper.class.getDeclaredField("thread");
|
||||
field.setAccessible(true);
|
||||
return (Thread) field.get(reaper);
|
||||
}
|
||||
|
||||
private static boolean runningOf(SessionReaper reaper) throws ReflectiveOperationException {
|
||||
Field field = SessionReaper.class.getDeclaredField("running");
|
||||
field.setAccessible(true);
|
||||
return field.getBoolean(reaper);
|
||||
}
|
||||
}
|
||||
@@ -7,13 +7,17 @@ import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.LoggerContext;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import dev.ltms.fleet.member.ClaudeCodeLauncher;
|
||||
import dev.ltms.fleet.msg.TestTurnTokens;
|
||||
import dev.ltms.fleet.peer.MemberRole;
|
||||
import dev.ltms.fleet.testing.CapturedLog;
|
||||
import org.junit.jupiter.api.AfterAll;
|
||||
import org.junit.jupiter.api.BeforeAll;
|
||||
import org.junit.jupiter.api.MethodOrderer;
|
||||
import org.junit.jupiter.api.Order;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.TestMethodOrder;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.List;
|
||||
@@ -27,9 +31,41 @@ import static org.junit.jupiter.api.Assertions.*;
|
||||
* CB-301-ext acceptance tests for worktree provisioning and config-parity overlay.
|
||||
* No live git — every Worktrees call is handled by {@link FakeWorktrees} and every herdr
|
||||
* call by {@link FakeHerdr}, matching the project's fake-based test style.
|
||||
*
|
||||
* <p>fleetd #529: only {@link #releasePreservesDirtyWorktreeAndLogsWarn} (explicitly
|
||||
* {@link Order#value() @Order(1)}) and the proving test right after it
|
||||
* ({@link #sharedSessionManagerLoggerLevelIsRestoredAfterDirtyWorktreeReleasePinsWarn},
|
||||
* {@code @Order(2)}) care about method order — every other test here has no {@code @Order} and so
|
||||
* runs after both (JUnit 5's {@link MethodOrderer.OrderAnnotation} gives an unannotated method the
|
||||
* lowest priority), in whatever relative order it already ran in.
|
||||
*/
|
||||
@TestMethodOrder(MethodOrderer.OrderAnnotation.class)
|
||||
class WorktreeSessionManagerTest {
|
||||
|
||||
/**
|
||||
* fleetd #529: the level {@link SessionManager}'s logger had when this class started, captured
|
||||
* before any test here touches it, then forced to a distinctive, known value (TRACE) so {@link
|
||||
* #sharedSessionManagerLoggerLevelIsRestoredAfterDirtyWorktreeReleasePinsWarn} can tell "the
|
||||
* level came back to what it was" apart from "the level happens to already be WARN because
|
||||
* some other test class in this JVM fork (surefire reuses forks by default) left it there".
|
||||
*/
|
||||
private static ch.qos.logback.classic.Level sessionManagerLevelBeforeThisClass;
|
||||
|
||||
@BeforeAll
|
||||
static void pinSessionManagerLoggerToAKnownBaseline() {
|
||||
ch.qos.logback.classic.Logger sessionLog =
|
||||
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(SessionManager.class);
|
||||
sessionManagerLevelBeforeThisClass = sessionLog.getLevel();
|
||||
sessionLog.setLevel(ch.qos.logback.classic.Level.TRACE);
|
||||
}
|
||||
|
||||
@AfterAll
|
||||
static void restoreSessionManagerLoggerLevel() {
|
||||
ch.qos.logback.classic.Logger sessionLog =
|
||||
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(SessionManager.class);
|
||||
sessionLog.setLevel(sessionManagerLevelBeforeThisClass);
|
||||
}
|
||||
|
||||
private static MemberRegistry members() {
|
||||
return new MemberRegistry(new FleetConfig.Fleet(Map.of(),
|
||||
Map.of("architect", new FleetConfig.Slot("ltms-local")),
|
||||
@@ -254,6 +290,7 @@ class WorktreeSessionManagerTest {
|
||||
* path, the session, and the cause an operator needs to find the work.
|
||||
*/
|
||||
@Test
|
||||
@Order(1)
|
||||
void releasePreservesDirtyWorktreeAndLogsWarn() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt")
|
||||
@@ -262,21 +299,13 @@ class WorktreeSessionManagerTest {
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("cb-576", null));
|
||||
|
||||
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
|
||||
ch.qos.logback.classic.Logger sessionLog =
|
||||
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(SessionManager.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.setContext(ctx);
|
||||
appender.start();
|
||||
sessionLog.addAppender(appender);
|
||||
sessionLog.setLevel(Level.WARN);
|
||||
try {
|
||||
try (CapturedLog log = CapturedLog.at(SessionManager.class, Level.WARN)) {
|
||||
sessions.release(s.paneId());
|
||||
|
||||
assertTrue(herdr.called("pane.close"), "release still tears the worker pane down");
|
||||
assertTrue(worktrees.removeCalls().isEmpty(),
|
||||
"a dirty worktree is never removed — it holds the only copy of the work");
|
||||
String warn = appender.list.stream()
|
||||
String warn = log.events().stream()
|
||||
.filter(e -> e.getLevel().equals(Level.WARN))
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.filter(m -> m.contains("dirty worktree"))
|
||||
@@ -285,11 +314,38 @@ class WorktreeSessionManagerTest {
|
||||
assertTrue(warn.contains(s.worktree()), "the WARN names the worktree path: " + warn);
|
||||
assertTrue(warn.contains(s.terminalId()), "the WARN names the session: " + warn);
|
||||
assertTrue(warn.contains("COMPLETED"), "the WARN names the release cause: " + warn);
|
||||
} finally {
|
||||
sessionLog.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #529 proving test: pins that the leak this ticket fixes stays fixed. {@link
|
||||
* #releasePreservesDirtyWorktreeAndLogsWarn} above (which runs immediately before this, via
|
||||
* {@code @Order}) pins the shared {@link SessionManager} logger to WARN through a {@link
|
||||
* CapturedLog}; if {@link CapturedLog#close} only detached the appender — the original bug —
|
||||
* the level would still read WARN here instead of the {@code TRACE} baseline this class's
|
||||
* {@code @BeforeAll} set. Runs at {@code @Order(2)}, guaranteed after {@code @Order(1)} and
|
||||
* before every other (unannotated) test in this class.
|
||||
*
|
||||
* <p>This proves only the WITHIN-CLASS case: JUnit 5's {@code @TestMethodOrder} orders methods
|
||||
* inside one class, not test classes relative to each other, and surefire's default class order
|
||||
* is not something a single test can force. The cross-class leak fleetd #525 measured — this
|
||||
* class's {@code releasePreservesDirtyWorktreeAndLogsWarn} pinning WARN and bleeding into a
|
||||
* later-running {@code SessionManagerTest} in the same fork — is fixed by the same {@link
|
||||
* CapturedLog} mechanism proven here, but that cross-class ordering itself is NOT asserted by
|
||||
* any test and remains unproven by construction.
|
||||
*/
|
||||
@Test
|
||||
@Order(2)
|
||||
void sharedSessionManagerLoggerLevelIsRestoredAfterDirtyWorktreeReleasePinsWarn() {
|
||||
ch.qos.logback.classic.Logger sessionLog =
|
||||
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(SessionManager.class);
|
||||
assertEquals(ch.qos.logback.classic.Level.TRACE, sessionLog.getLevel(),
|
||||
"releasePreservesDirtyWorktreeAndLogsWarn pins the shared SessionManager logger to "
|
||||
+ "WARN; its cleanup must restore the level it captured (TRACE, set by this "
|
||||
+ "class's @BeforeAll) rather than leaving WARN pinned for every test that "
|
||||
+ "runs after it");
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-576 review (fleetd #116). A worktree that is already gone (operator cleanup,
|
||||
* {@code git worktree prune}, an earlier half-completed release) must not break teardown.
|
||||
|
||||
@@ -0,0 +1,112 @@
|
||||
package dev.ltms.fleet.testing;
|
||||
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.Logger;
|
||||
import ch.qos.logback.classic.LoggerContext;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.List;
|
||||
|
||||
/**
|
||||
* fleetd #529 (promoted from {@code SessionManagerTest}, merged in #527 for fleetd #525): captures
|
||||
* a logger's output and, on {@link #close}, restores <em>both</em> the appender and the level to
|
||||
* what they were before.
|
||||
*
|
||||
* <p>{@code ch.qos.logback.classic.Logger} instances are cached per class and shared across the
|
||||
* whole JVM — and surefire reuses forks by default — so a bare {@code addAppender}/{@code
|
||||
* setLevel} pair whose {@code finally} only detaches the appender leaves the level pinned for
|
||||
* every test that runs after it, in the same class or in a completely unrelated one sharing the
|
||||
* fork. try-with-resources makes "restored the appender but not the level" impossible to write,
|
||||
* because there is only one thing to close.
|
||||
*
|
||||
* <p>New code must use this rather than hand-rolling the {@code ListAppender} + {@code setLevel} +
|
||||
* {@code finally detachAppender} pattern: use {@link #at} or {@link #of}. It is <em>not</em> yet
|
||||
* the only instance of the pattern in this test tree, and the earlier wording here said it was —
|
||||
* which would leave a reader who greps unable to tell a leftover from a violation.
|
||||
*
|
||||
* <p>Measured on main at af95897 (2026-09-12): nine test files still hand-roll it, with 42
|
||||
* {@code setLevel} calls on a raw logback {@code Logger} between them — {@code
|
||||
* FleetdStartupReportTest}, {@code GitHostShapeReportTest}, {@code MemberCredentialsGapReportTest},
|
||||
* {@code MemberTrustModelReportTest}, {@code FleetHealthMonitorTest}, {@code
|
||||
* ClaudeCodeLauncherTest}, {@code HerdrPeerLauncherAllowListWiringTest}, {@code
|
||||
* HerdrPeerLauncherCharterTest} and {@code OpenCodeLauncherTest}. Every one of them pairs its pin
|
||||
* with a restore, so none is the fleetd #525 leak and none was in fleetd #529's scope, which was
|
||||
* the 19 <em>unrestored</em> pins only. They are unmigrated, not broken.
|
||||
*
|
||||
* <p>Re-measure with the two commands below, from the repo root. A file that appears in the first
|
||||
* list and not the second still hand-rolls the pattern. When the first list comes back empty, this
|
||||
* paragraph is spent and the sentence above can go back to saying "the one way" — delete the
|
||||
* paragraph then rather than updating the count.
|
||||
*
|
||||
* <pre>{@code
|
||||
* grep -rlE '\.setLevel\(' fleetd/src/test/java --include='*.java' | grep -v CapturedLog.java
|
||||
* grep -rl 'CapturedLog' fleetd/src/test/java --include='*.java'
|
||||
* }</pre>
|
||||
*/
|
||||
public final class CapturedLog implements AutoCloseable {
|
||||
private final Logger logger;
|
||||
private final Level originalLevel;
|
||||
private final ListAppender<ILoggingEvent> appender;
|
||||
|
||||
private CapturedLog(Logger logger, Level pinnedLevel) {
|
||||
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
|
||||
this.logger = logger;
|
||||
this.originalLevel = logger.getLevel();
|
||||
this.appender = new ListAppender<>();
|
||||
appender.setContext(ctx);
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
if (pinnedLevel != null) {
|
||||
logger.setLevel(pinnedLevel);
|
||||
}
|
||||
}
|
||||
|
||||
/** Capture {@code loggerClass}'s output, pinning its level to {@code pinnedLevel} for the
|
||||
* duration of the try-with-resources block. */
|
||||
public static CapturedLog at(Class<?> loggerClass, Level pinnedLevel) {
|
||||
return new CapturedLog((Logger) LoggerFactory.getLogger(loggerClass), pinnedLevel);
|
||||
}
|
||||
|
||||
/** Capture {@code loggerClass}'s output without changing its level. */
|
||||
public static CapturedLog of(Class<?> loggerClass) {
|
||||
return new CapturedLog((Logger) LoggerFactory.getLogger(loggerClass), null);
|
||||
}
|
||||
|
||||
/**
|
||||
* Capture the named logger's output, pinning its level to {@code pinnedLevel}. For a logger
|
||||
* obtained in production code via {@code LoggerFactory.getLogger("some-name")} rather than a
|
||||
* class — e.g. {@code AuditLog}'s {@code "audit"} logger — where {@link #at(Class, Level)}
|
||||
* would capture the wrong {@code Logger} instance (class name and logger name are different
|
||||
* strings and resolve to different cached loggers).
|
||||
*/
|
||||
public static CapturedLog at(String loggerName, Level pinnedLevel) {
|
||||
return new CapturedLog((Logger) LoggerFactory.getLogger(loggerName), pinnedLevel);
|
||||
}
|
||||
|
||||
/** Capture the named logger's output without changing its level. See {@link #at(String, Level)}. */
|
||||
public static CapturedLog of(String loggerName) {
|
||||
return new CapturedLog((Logger) LoggerFactory.getLogger(loggerName), null);
|
||||
}
|
||||
|
||||
public List<ILoggingEvent> events() {
|
||||
return appender.list;
|
||||
}
|
||||
|
||||
/**
|
||||
* Re-pin the level while this capture is still open — for example to lower it further for one
|
||||
* assertion inside a test whose fixture already pinned a coarser baseline in {@code @BeforeEach}.
|
||||
* This does not change what {@link #close} restores: that is always the level captured when
|
||||
* this instance was created, never a value set through this method.
|
||||
*/
|
||||
public void setLevel(Level level) {
|
||||
logger.setLevel(level);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
logger.detachAppender(appender);
|
||||
logger.setLevel(originalLevel);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,96 @@
|
||||
package dev.ltms.fleet.testing;
|
||||
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.Logger;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
|
||||
/**
|
||||
* fleetd #537: pins {@link CapturedLog#close}'s own contract — the appender detach half, the level
|
||||
* restore half, and that {@link CapturedLog#setLevel} does not change what {@code close} restores.
|
||||
* Before this, only the level-restore half was pinned (by {@code
|
||||
* WorktreeSessionManagerTest.sharedSessionManagerLoggerLevelIsRestoredAfterDirtyWorktreeReleasePinsWarn}).
|
||||
* Measured: deleting {@code logger.detachAppender(appender);} from {@code close()} still left
|
||||
* {@code mvn clean install} green — 1701 tests, 0 failures — before this file existed.
|
||||
*
|
||||
* <p>Every test below uses a logger name no production class uses, and unique per test, so this
|
||||
* file cannot become the next entry in fleetd #525's leak family: {@link CapturedLog}'s own
|
||||
* javadoc re-measure commands (top of that file) would otherwise need to start naming this class.
|
||||
*/
|
||||
class CapturedLogTest {
|
||||
|
||||
/**
|
||||
* The appender must be detached on close: an event logged through the raw logger after close
|
||||
* must not land in {@link CapturedLog#events()}. Asserting on the observable list (rather than
|
||||
* {@code logger.iteratorForAppenders()}) is what the ticket asked for, and it is also what a
|
||||
* real leak would actually break — a later test's own {@code ListAppender} silently gaining
|
||||
* events emitted by code under test that has nothing to do with it.
|
||||
*/
|
||||
@Test
|
||||
void closeDetachesTheAppenderSoALaterLogIsNotCaptured() {
|
||||
String loggerName = "capturedlog-test-only.appender-detach";
|
||||
Logger rawLogger = (Logger) LoggerFactory.getLogger(loggerName);
|
||||
|
||||
CapturedLog log = CapturedLog.of(loggerName);
|
||||
rawLogger.info("while open");
|
||||
int eventsWhileOpen = log.events().size();
|
||||
assertEquals(1, eventsWhileOpen, "the event logged while open must be captured");
|
||||
|
||||
log.close();
|
||||
rawLogger.info("after close");
|
||||
|
||||
assertEquals(eventsWhileOpen, log.events().size(),
|
||||
"close() must detach the appender: an event logged after close must not be "
|
||||
+ "captured, but the captured list grew from " + eventsWhileOpen + " to "
|
||||
+ log.events().size());
|
||||
}
|
||||
|
||||
/**
|
||||
* The helper's own headline contract, pinned in one place independent of any production
|
||||
* class's behaviour: {@code close()} restores the level the logger had before {@link
|
||||
* CapturedLog#at} pinned it.
|
||||
*/
|
||||
@Test
|
||||
void closeRestoresTheLevelCapturedAtOpen() {
|
||||
String loggerName = "capturedlog-test-only.level-restore";
|
||||
Logger rawLogger = (Logger) LoggerFactory.getLogger(loggerName);
|
||||
rawLogger.setLevel(Level.DEBUG);
|
||||
|
||||
CapturedLog log = CapturedLog.at(loggerName, Level.ERROR);
|
||||
assertEquals(Level.ERROR, rawLogger.getLevel(), "the pinned level took effect while open");
|
||||
|
||||
log.close();
|
||||
|
||||
assertEquals(Level.DEBUG, rawLogger.getLevel(),
|
||||
"close() must restore the level captured when at() was called (DEBUG), not leave "
|
||||
+ "the pinned level (ERROR) in place");
|
||||
}
|
||||
|
||||
/**
|
||||
* {@link CapturedLog#setLevel}'s javadoc claims that re-pinning the level mid-capture does not
|
||||
* change what {@code close()} restores — that restore always uses the level captured when the
|
||||
* instance was created, never a value set through {@code setLevel}. Nothing checked this
|
||||
* before: open with a pinned WARN, call {@code setLevel(TRACE)}, close, and the result must be
|
||||
* the level from BEFORE {@code at} — neither WARN nor TRACE.
|
||||
*/
|
||||
@Test
|
||||
void setLevelDuringCaptureDoesNotChangeWhatCloseRestores() {
|
||||
String loggerName = "capturedlog-test-only.setlevel-no-effect";
|
||||
Logger rawLogger = (Logger) LoggerFactory.getLogger(loggerName);
|
||||
rawLogger.setLevel(Level.DEBUG);
|
||||
|
||||
CapturedLog log = CapturedLog.at(loggerName, Level.WARN);
|
||||
log.setLevel(Level.TRACE);
|
||||
assertEquals(Level.TRACE, rawLogger.getLevel(), "setLevel took effect immediately");
|
||||
|
||||
log.close();
|
||||
|
||||
assertEquals(Level.DEBUG, rawLogger.getLevel(),
|
||||
"close() must restore the level captured at open (DEBUG) regardless of any "
|
||||
+ "later setLevel() call: it must be neither WARN (the level pinned by "
|
||||
+ "at()) nor TRACE (the level set via setLevel() mid-capture), but got "
|
||||
+ rawLogger.getLevel());
|
||||
}
|
||||
}
|
||||
+259
-26
@@ -42,6 +42,14 @@
|
||||
# 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.
|
||||
# 9. fleetd #512 — the ERROR-line count above is blind by construction to the exact failure #493
|
||||
# is about: an uncaught exception in a shutdown thread never passes through the logger, so it
|
||||
# never carries an ERROR (or SEVERE) token that any count could see. This script now also
|
||||
# greps the previous daemon's shutdown window for that exception's real shape, and separately
|
||||
# asserts that SessionManager's drain-complete line (fleetd #522) is present there — its
|
||||
# absence is the real signal, because a drain that dies on its first session prints nothing
|
||||
# else either. Warns loudly; never fails the redeploy, because by the time this is detectable
|
||||
# the new daemon is already up and healthy.
|
||||
#
|
||||
# Usage:
|
||||
# scripts/redeploy-fleetd.sh # build, confirm, restart, verify
|
||||
@@ -88,13 +96,23 @@ SYSTEMD_UNIT='fleetd'
|
||||
# any of the prose detail strings and needs no escaping in a `case`/glob pattern.
|
||||
SUPERVISOR_DETAIL_SEP=$'\x1f'
|
||||
|
||||
# fleetd #492 follow-up: set by systemd_loaded/systemd_installed when the underlying `systemctl`
|
||||
# call could not answer cleanly — it exited non-zero AND wrote something to stderr, which is a real
|
||||
# tool failure (e.g. it cannot reach the user bus over a non-lingering ssh session), never the same
|
||||
# fact as a clean negative answer ("not active", no stderr). Initialized here, not just inside the
|
||||
# probes, so detect_supervisor can read them under `set -u` even before either probe has ever run,
|
||||
# and so a test that stubs a probe with a plain `return 0`/`return 1` body (leaving these untouched)
|
||||
# reads a deterministic 0 rather than whatever a previous probe call left behind.
|
||||
# fleetd #492 follow-up, refined by fleetd #545: three states, not two, set by
|
||||
# systemd_loaded/systemd_installed —
|
||||
# 0 = no error, the probe ran and gave a clean answer.
|
||||
# 1 = the probe RAN and answered badly: `systemctl` exited non-zero AND wrote something to
|
||||
# stderr, a real tool failure (e.g. it cannot reach the user bus over a non-lingering ssh
|
||||
# session), never the same fact as a clean negative answer ("not active", no stderr).
|
||||
# 2 = the probe could not even be SET UP: the `mktemp` call that makes a place to capture
|
||||
# `systemctl`'s stderr failed before `systemctl` ever ran. This is a fleetd #545 fix: on GNU
|
||||
# coreutils (every Linux distribution) a template with no `X`s made `mktemp` fail every
|
||||
# single time, and the two states were folded into one flag and one message that named
|
||||
# cause 1 ("systemctl exited non-zero and reported an error on stderr") for a failure that
|
||||
# was actually cause 2 — systemctl was never executed at all. One flag with two meanings
|
||||
# needing different messages was the defect; a third value is the fix, not a second flag.
|
||||
# Initialized here, not just inside the probes, so detect_supervisor can read them under `set -u`
|
||||
# even before either probe has ever run, and so a test that stubs a probe with a plain
|
||||
# `return 0`/`return 1` body (leaving these untouched) reads a deterministic 0 rather than whatever
|
||||
# a previous probe call left behind.
|
||||
SYSTEMD_LOADED_ERRORED=0
|
||||
SYSTEMD_INSTALLED_ERRORED=0
|
||||
# fleetd #492 follow-up: SUPERVISOR_UNCLEAR_DETAIL is the specific supervisor/reason that
|
||||
@@ -230,12 +248,17 @@ launchd_loaded() { launchctl list "$LAUNCHD_LABEL" >/dev/null 2>&1; }
|
||||
# fleetd #492 follow-up: both functions used to throw `systemctl`'s stderr straight into
|
||||
# /dev/null, which meant "systemctl answered no" and "systemctl could not answer at all" (e.g. it
|
||||
# cannot reach the user bus over a non-lingering ssh session) looked identical — both a plain
|
||||
# nonzero exit. They now capture stderr separately and set their own *_ERRORED flag ONLY when the
|
||||
# nonzero exit. They now capture stderr separately and set their own *_ERRORED flag to 1 when the
|
||||
# call exited non-zero AND wrote something to stderr — a real tool failure, never a clean "not
|
||||
# installed"/"not active" answer (which exits non-zero with empty stderr). detect_supervisor reads
|
||||
# the flag right after calling the probe, so a probe that could not answer routes to "unclear",
|
||||
# never silently becomes "none".
|
||||
#
|
||||
# fleetd #545: the flag has a third value, 2, set when the `mktemp` call that sets up the probe's
|
||||
# own stderr capture fails, before `systemctl` ever runs — see the SYSTEMD_LOADED_ERRORED /
|
||||
# SYSTEMD_INSTALLED_ERRORED comment above their initialization for why this is a third value on the
|
||||
# same flag, not a second flag.
|
||||
#
|
||||
# "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.
|
||||
@@ -243,8 +266,8 @@ systemd_installed() {
|
||||
SYSTEMD_INSTALLED_ERRORED=0
|
||||
command -v systemctl >/dev/null 2>&1 || return 1
|
||||
local err_file out rc=0
|
||||
if ! err_file="$(mktemp -t systemd-installed-err)"; then
|
||||
SYSTEMD_INSTALLED_ERRORED=1
|
||||
if ! err_file="$(mktemp -t systemd-installed-err.XXXXXX)"; then
|
||||
SYSTEMD_INSTALLED_ERRORED=2
|
||||
return 1
|
||||
fi
|
||||
out="$(systemctl --user list-unit-files "$SYSTEMD_UNIT.service" --no-legend 2>"$err_file")" || rc=$?
|
||||
@@ -267,8 +290,8 @@ systemd_loaded() {
|
||||
SYSTEMD_LOADED_ERRORED=0
|
||||
command -v systemctl >/dev/null 2>&1 || return 1
|
||||
local err_file rc=0
|
||||
if ! err_file="$(mktemp -t systemd-loaded-err)"; then
|
||||
SYSTEMD_LOADED_ERRORED=1
|
||||
if ! err_file="$(mktemp -t systemd-loaded-err.XXXXXX)"; then
|
||||
SYSTEMD_LOADED_ERRORED=2
|
||||
return 1
|
||||
fi
|
||||
systemctl --user is-active "$SYSTEMD_UNIT" >/dev/null 2>"$err_file" || rc=$?
|
||||
@@ -279,6 +302,49 @@ systemd_loaded() {
|
||||
return "$rc"
|
||||
}
|
||||
|
||||
# fleetd #504: the "loaded but not currently running" branches in the main stop step (case
|
||||
# launchd/systemd, reached when $OLD_PID is empty) used to run `launchctl unload`/`systemctl --user
|
||||
# stop` with `2>/dev/null || true` and then print `ok` unconditionally — the exact conflation
|
||||
# systemd_loaded/systemd_installed above already fixed on the READ side (fleetd #492 follow-up): a
|
||||
# genuine "already stopped" answer (nonzero exit, nothing on stderr) is harmless, but a real tool
|
||||
# failure (nonzero exit WITH a stderr message — e.g. launchd or the systemd user bus is
|
||||
# unreachable) is not, and reporting `ok` on THAT means the start step below can register a fresh
|
||||
# load on top of a supervisor that never actually let go: the exact two-daemons failure fleetd #492
|
||||
# exists to prevent, reached from the one state (already odd) where a false `ok` is least
|
||||
# affordable. These two functions apply the same "capture stderr separately, flag only a nonzero
|
||||
# exit WITH stderr as a real failure" pattern to the WRITE side. No ${VAR:-default} anywhere here —
|
||||
# see the #492 follow-up constraints comment above detect_supervisor for why a default would hide a
|
||||
# lost value instead of surfacing it (fleetd #497's defect class).
|
||||
unload_launchd_if_loaded() {
|
||||
local err_file rc=0
|
||||
if ! err_file="$(mktemp -t launchd-unload-err.XXXXXX)"; then
|
||||
die "could not create a temp file to capture 'launchctl unload' stderr — cannot tell a real
|
||||
failure from a clean already-unloaded answer, so refusing to guess. The daemon's
|
||||
supervision state was NOT touched."
|
||||
fi
|
||||
launchctl unload -w "$LAUNCHD_PLIST" >/dev/null 2>"$err_file" || rc=$?
|
||||
if [ "$rc" -ne 0 ] && [ -s "$err_file" ]; then
|
||||
die "'launchctl unload -w $LAUNCHD_PLIST' failed: $(cat "$err_file")
|
||||
The daemon may still be under supervision; investigate before retrying."
|
||||
fi
|
||||
rm -f "$err_file"
|
||||
}
|
||||
|
||||
stop_systemd_if_loaded() {
|
||||
local err_file rc=0
|
||||
if ! err_file="$(mktemp -t systemd-stop-err.XXXXXX)"; then
|
||||
die "could not create a temp file to capture 'systemctl --user stop' stderr — cannot tell a
|
||||
real failure from a clean already-stopped answer, so refusing to guess. The daemon's
|
||||
supervision state was NOT touched."
|
||||
fi
|
||||
systemctl --user stop "$SYSTEMD_UNIT" >/dev/null 2>"$err_file" || rc=$?
|
||||
if [ "$rc" -ne 0 ] && [ -s "$err_file" ]; then
|
||||
die "'systemctl --user stop $SYSTEMD_UNIT' failed: $(cat "$err_file")
|
||||
The daemon may still be under supervision; investigate before retrying."
|
||||
fi
|
||||
rm -f "$err_file"
|
||||
}
|
||||
|
||||
# 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
|
||||
@@ -298,6 +364,14 @@ systemd_loaded() {
|
||||
# "none" now means only: neither supervisor is installed, neither is loaded, and neither probe
|
||||
# errored.
|
||||
#
|
||||
# fleetd #545: *_ERRORED carries a THIRD state (2 = the probe's own mktemp setup failed, before
|
||||
# `systemctl` ever ran — see the flag's own comment above its initialization), and it must never be
|
||||
# reported with the same detail text as state 1 (`systemctl` ran and answered badly on stderr). The
|
||||
# two are different facts about different failures, and conflating them makes the "unclear" message
|
||||
# assert a cause ("systemctl exited non-zero and reported an error on stderr") that was never
|
||||
# measured when the real cause was state 2. detect_supervisor below picks the detail text off the
|
||||
# flag's value, not off a single "errored at all" boolean.
|
||||
#
|
||||
# fleetd #492 follow-up — constraints every caller of this function depends on (learned the hard
|
||||
# way: an earlier version of this fix set a SUPERVISOR_UNCLEAR_DETAIL global from inside here and
|
||||
# it was silently lost, because every real call site invokes this as `$(detect_supervisor)`):
|
||||
@@ -306,10 +380,12 @@ systemd_loaded() {
|
||||
# never assigned to a global: a global set inside a `$( )` subshell dies with that subshell.
|
||||
# This function packs BOTH values (kind and detail) onto that one stdout line, joined by
|
||||
# $SUPERVISOR_DETAIL_SEP, and the caller unpacks them on its own side of the subshell boundary.
|
||||
# 2. This script runs under `set -euo pipefail` (line 50), so an unset variable is a loud
|
||||
# failure. Do not add a `${VAR:-default}` anywhere downstream to paper over a value that
|
||||
# should always be there — that hides a lost value instead of surfacing it (fleetd #497's
|
||||
# defect class).
|
||||
# 2. This script runs under `set -u` (part of the `set -euo pipefail` at the top of the file), so
|
||||
# an unset variable is a loud failure. Do not add a `${VAR:-default}` anywhere downstream to
|
||||
# paper over a value that should always be there — that hides a lost value instead of
|
||||
# surfacing it (fleetd #497's defect class). The mechanism is named rather than cited by line
|
||||
# number on purpose: a line number in a comment goes stale on the next insert above it, and
|
||||
# this one already had — it said line 50 while the `set` line was at 54.
|
||||
# 3. Every `case` on this function's return value needs an explicit final `*)` arm, chosen by
|
||||
# whether that caller ACTS on the value (`die` — an unrecognised value must never be silently
|
||||
# driven) or only DISPLAYS it (`echo`/`warn` and continue — a diagnostic must not go silent on
|
||||
@@ -329,7 +405,10 @@ detect_supervisor() {
|
||||
launchd_installed && li=1
|
||||
systemd_installed && si=1
|
||||
|
||||
if [ "$SYSTEMD_LOADED_ERRORED" = 1 ] || [ "$SYSTEMD_INSTALLED_ERRORED" = 1 ]; then
|
||||
if [ "$SYSTEMD_LOADED_ERRORED" = 2 ] || [ "$SYSTEMD_INSTALLED_ERRORED" = 2 ]; then
|
||||
detail="the systemd --user probe for '$SYSTEMD_UNIT' could not even be set up (a temp file to capture systemctl's stderr could not be created) — systemctl was never run, so this says nothing about systemd, the user bus, or the unit itself"
|
||||
kind="unclear"
|
||||
elif [ "$SYSTEMD_LOADED_ERRORED" = 1 ] || [ "$SYSTEMD_INSTALLED_ERRORED" = 1 ]; then
|
||||
detail="the systemd --user probe for '$SYSTEMD_UNIT' could not answer cleanly (systemctl exited non-zero and reported an error on stderr, not a clean negative — e.g. it cannot reach the user bus)"
|
||||
kind="unclear"
|
||||
elif [ "$ld" = 1 ] && [ "$sd" = 1 ]; then
|
||||
@@ -490,6 +569,114 @@ classify_amqp_connection_errors() {
|
||||
REDEPLOY_UNEXPLAINED_ERRORS=$((REDEPLOY_UNEXPLAINED_ERRORS + pending_inbox + pending_lead_mailbox))
|
||||
}
|
||||
|
||||
# fleetd #512 part 2 — the negative check. #493's failure (an uncaught exception in a shutdown
|
||||
# thread) never passes through the logger: the JVM's default uncaught-exception handler prints
|
||||
# straight to stderr, so the line never carries a level, so classify_amqp_connection_errors's
|
||||
# ERROR/SEVERE token scan is structurally blind to it — measured on two hosts, including one where
|
||||
# even a syslog PRIORITY filter is blind to it too (fd 1 and fd 2 collapse to one socket there, so
|
||||
# every uncaught-exception line lands at priority 6/info). The fix is to grep the shape instead of
|
||||
# the level: `Exception in thread` at the start of a line (the handler's own banner) or
|
||||
# `NoClassDefFoundError` anywhere in it (the one real instance seen so far, but not the only shape
|
||||
# this could take). Kept as its own function, never folded into classify_amqp_connection_errors —
|
||||
# this is not an AMQP concern, and the two must stay independently readable and independently
|
||||
# testable.
|
||||
#
|
||||
# Sets REDEPLOY_UNCAUGHT_EXCEPTION_COUNT (lines matched) and REDEPLOY_UNCAUGHT_EXCEPTION_SAMPLE
|
||||
# (the first matching line, "" if none) so a caller can report both a count and a concrete quote
|
||||
# without re-reading the file. Pure: reads $1, sets globals, no side effects.
|
||||
scan_uncaught_exceptions() {
|
||||
local log_file="$1" line
|
||||
REDEPLOY_UNCAUGHT_EXCEPTION_COUNT=0
|
||||
REDEPLOY_UNCAUGHT_EXCEPTION_SAMPLE=""
|
||||
while IFS= read -r line || [ -n "$line" ]; do
|
||||
case "$line" in
|
||||
'Exception in thread'*|*NoClassDefFoundError*)
|
||||
REDEPLOY_UNCAUGHT_EXCEPTION_COUNT=$((REDEPLOY_UNCAUGHT_EXCEPTION_COUNT + 1))
|
||||
[ -n "$REDEPLOY_UNCAUGHT_EXCEPTION_SAMPLE" ] || REDEPLOY_UNCAUGHT_EXCEPTION_SAMPLE="$line"
|
||||
;;
|
||||
esac
|
||||
done < "$log_file"
|
||||
}
|
||||
|
||||
# fleetd #512 part 2 — the positive check. fleetd #522 added a `log.info` at the very end of
|
||||
# SessionManager.drainAll's normal path (never in a `finally` — see the ticket discussion for why
|
||||
# that distinction matters): "drain complete: released=N abandoned=M (still BUSY at the shutdown
|
||||
# deadline)", printed once, on every successful drain, including the all-zero case. A drain that
|
||||
# dies partway through never reaches that statement, so the line's ABSENCE is a real signal — unlike
|
||||
# the ERROR-count check above, this one does not depend on the failure happening to throw.
|
||||
#
|
||||
# Sets REDEPLOY_DRAIN_COMPLETE_LINE to the matching line (last one, though drainAll runs at most
|
||||
# once per shutdown so there should never be more than one) or "" if absent. Pure, same shape as
|
||||
# scan_uncaught_exceptions above.
|
||||
find_drain_complete_line() {
|
||||
local log_file="$1"
|
||||
REDEPLOY_DRAIN_COMPLETE_LINE="$(grep -F 'drain complete: released=' "$log_file" | tail -1 || true)"
|
||||
}
|
||||
|
||||
# fleetd #512 part 2 — THE TRAP, and the reason this is one function instead of two independent
|
||||
# checks the caller ORs together. Absence of the drain-complete line has TWO causes that need
|
||||
# OPPOSITE handling, and a naive "line absent -> the drain died" reading collapses them exactly the
|
||||
# way this whole ticket exists to stop: the line is emitted by the daemon being STOPPED, which is
|
||||
# running the OLD jar. Until a redeploy has landed fleetd #522 once, every previous daemon predates
|
||||
# the line and cannot emit it no matter how cleanly it drained — so on the very first redeploy after
|
||||
# #522 merged, "absent" means "too old to know how", not "died". Only once BOTH signals — this
|
||||
# line's absence AND scan_uncaught_exceptions' result — have been read together can the three real
|
||||
# outcomes be told apart:
|
||||
#
|
||||
# complete -> the line is present: the drain finished. Name the counts it reported.
|
||||
# died -> the line is absent AND an uncaught-exception shape was found: the drain died. Name
|
||||
# what was found.
|
||||
# unknown -> the line is absent AND no exception shape either: cannot tell. Say so, and say why
|
||||
# (predates the line, or failed without throwing) — never worded as a pass or a
|
||||
# failure, and never reassuring: "ok, no ERROR lines" one level up is the exact mistake
|
||||
# this ticket exists to fix, and this outcome must not reproduce it.
|
||||
#
|
||||
# A fourth case, n/a, covers a cold start or a "loaded but wasn't running" restart: no previous
|
||||
# daemon was actually stopped THIS run, so there is no shutdown window in $log_file to have an
|
||||
# opinion about at all — scanning it anyway would read the NEW daemon's own startup lines and could
|
||||
# misreport "cannot tell" on every clean cold start. had_previous_daemon carries that fact in from
|
||||
# the caller (it already knows $OLD_PID) rather than this function re-deriving it from log content.
|
||||
#
|
||||
# Same shape as swap_if_built/refuse_drain_gate (fleetd #521/#528): the decision (which of the four
|
||||
# outcomes applies) and the action (which ok/warn line to print, and setting REDEPLOY_DRAIN_STATE
|
||||
# for the "result" section below to consult) live together in ONE function that the main flow calls
|
||||
# unconditionally — there is no guard left in the main flow to remove, invert, or bypass
|
||||
# independently of this function. Never calls die(): #512's own decision is to warn loudly and let
|
||||
# the redeploy stand, because by the time this is detectable the new daemon is already up and
|
||||
# healthy and failing here would give the operator nothing to do differently.
|
||||
report_shutdown_drain() {
|
||||
local log_file="$1" had_previous_daemon="$2"
|
||||
REDEPLOY_UNCAUGHT_EXCEPTION_COUNT=0
|
||||
REDEPLOY_UNCAUGHT_EXCEPTION_SAMPLE=""
|
||||
REDEPLOY_DRAIN_COMPLETE_LINE=""
|
||||
|
||||
if [ "$had_previous_daemon" != 1 ]; then
|
||||
REDEPLOY_DRAIN_STATE="n/a"
|
||||
ok "no previous daemon was running before this restart — nothing to check for a died shutdown drain"
|
||||
return 0
|
||||
fi
|
||||
|
||||
find_drain_complete_line "$log_file"
|
||||
scan_uncaught_exceptions "$log_file"
|
||||
|
||||
if [ -n "$REDEPLOY_DRAIN_COMPLETE_LINE" ]; then
|
||||
REDEPLOY_DRAIN_STATE="complete"
|
||||
ok "previous daemon's shutdown drain finished: $REDEPLOY_DRAIN_COMPLETE_LINE"
|
||||
elif [ "$REDEPLOY_UNCAUGHT_EXCEPTION_COUNT" -gt 0 ]; then
|
||||
REDEPLOY_DRAIN_STATE="died"
|
||||
warn "previous daemon's shutdown drain DIED — no drain-complete line, and an uncaught exception"
|
||||
warn "was found in its shutdown window ($REDEPLOY_UNCAUGHT_EXCEPTION_COUNT line(s)):"
|
||||
warn " $REDEPLOY_UNCAUGHT_EXCEPTION_SAMPLE"
|
||||
warn "Some sessions from the PREVIOUS daemon may not have been released."
|
||||
else
|
||||
REDEPLOY_DRAIN_STATE="unknown"
|
||||
warn "cannot tell whether the previous daemon's shutdown drain finished — no drain-complete line"
|
||||
warn "and no uncaught-exception shape either. This is NOT a pass and NOT a failure: it means"
|
||||
warn "either that daemon predates fleetd #522's drain-complete log line, or its drain failed"
|
||||
warn "without throwing (hung, or returned early)."
|
||||
fi
|
||||
}
|
||||
|
||||
# fleetd #517: extracted so the suite can call this decision directly, the same way #510 extracted
|
||||
# wait_for_daemon_exit so its ordering became checkable. Before this, the only test of the drain-gate
|
||||
# abort message was a grep of this script's own source for the wording — so mutating the `if` below
|
||||
@@ -519,6 +706,32 @@ drain_gate_refusal() {
|
||||
fi
|
||||
}
|
||||
|
||||
# fleetd #528 — drain_gate_refusal above is well tested (four cases, all direct), but nothing made
|
||||
# the MAIN FLOW's abort actually consult it. Before this, the main flow read
|
||||
# `die "$(drain_gate_refusal "$DO_BUILD" "$JAR_STAGED")"` directly, and mutating that one line to a
|
||||
# flat `die "aborted — nothing changed"` left the whole suite at exit 0 with zero FAIL lines and
|
||||
# byte-identical output to a clean run — every one of drain_gate_refusal's own tests still passed,
|
||||
# because they call the predicate directly and never touch this call site. That silently reinstated
|
||||
# the exact defect #517 was filed to fix. Same shape as #521/#526's should_swap/swap_if_built: a
|
||||
# predicate alone is not enough, because a test proving the predicate is right cannot also prove the
|
||||
# main flow consults it. So the decision (drain_gate_refusal) and the action (die) now live together
|
||||
# in ONE function, and the main flow calls it unconditionally instead of building the die() call
|
||||
# itself — there is no guard left in the main flow to remove, invert, or bypass independently of this
|
||||
# function. drain_gate_refusal stays separate and separately tested because the message-selection
|
||||
# logic is worth naming and testing on its own; refuse_drain_gate is the only thing that ever dies.
|
||||
#
|
||||
# What the behavioural tests above still cannot pin on their own: deleting the call to this function
|
||||
# from the main flow altogether — they call refuse_drain_gate directly, never through the main flow,
|
||||
# because sourcing stops before the main flow ever runs (see the SOURCED guard below). That gap is
|
||||
# closed the same way swap_if_built's is: test_refuse_drain_gate_call_site_present greps this script
|
||||
# for the real invocation, the same shape test_swap_ordered_after_wait_and_before_start already uses
|
||||
# for the swap call. Deliberately NOT written out here as a literal quoted string, so this comment
|
||||
# itself can never become a second match for that test's needle.
|
||||
refuse_drain_gate() {
|
||||
local do_build="$1" staged_path="$2"
|
||||
die "$(drain_gate_refusal "$do_build" "$staged_path")"
|
||||
}
|
||||
|
||||
# CB-600: sourceable for testing. When this file is SOURCED (not executed) it stops here — nothing
|
||||
# below runs — so a test harness can `source` it to call check_log_path_matches_plist (or the
|
||||
# other pure helpers above) against a throwaway plist fixture without ever reaching the mutating
|
||||
@@ -643,7 +856,7 @@ if [ "$DO_BUILD" = 1 ]; then
|
||||
# fleetd #493: wipe a leftover staged jar from a previous failed/interrupted run BEFORE doing
|
||||
# anything else, so that run's leftovers can never be mistaken for this run's output.
|
||||
rm -f "$JAR_STAGED"
|
||||
BUILD_LOG="$(mktemp -t fleetd-build)"
|
||||
BUILD_LOG="$(mktemp -t fleetd-build.XXXXXX)"
|
||||
echo " log: $BUILD_LOG"
|
||||
if ! mvn -f "$MODULE/pom.xml" clean install > "$BUILD_LOG" 2>&1; then
|
||||
grep -E 'ERROR|BUILD FAILURE|Tests run:.*Failures: [1-9]|Tests run:.*Errors: [1-9]' "$BUILD_LOG" \
|
||||
@@ -677,10 +890,11 @@ if [ -n "$OLD_PID" ] && [ "$ASSUME_YES" = 0 ]; then
|
||||
echo
|
||||
read -r -p " Fleet drained? type yes to restart: " reply
|
||||
if [ "$reply" != "yes" ]; then
|
||||
# fleetd #493 / #517: "nothing changed" would be a lie once a build has run and staged a jar —
|
||||
# see drain_gate_refusal above for the full decision and why each of its four cases reads the
|
||||
# way it does.
|
||||
die "$(drain_gate_refusal "$DO_BUILD" "$JAR_STAGED")"
|
||||
# fleetd #493 / #517 / #528: "nothing changed" would be a lie once a build has run and staged a
|
||||
# jar — see drain_gate_refusal above for the full decision and why each of its four cases reads
|
||||
# the way it does. refuse_drain_gate composes that message AND calls die itself, so this guard
|
||||
# has nothing left of its own to get wrong beyond whether it calls refuse_drain_gate at all.
|
||||
refuse_drain_gate "$DO_BUILD" "$JAR_STAGED"
|
||||
fi
|
||||
fi
|
||||
|
||||
@@ -736,16 +950,20 @@ if [ -n "$OLD_PID" ]; then
|
||||
elif [ "$SUPERVISOR_KIND" = "launchd" ]; 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.
|
||||
# fleetd #504: unload_launchd_if_loaded (above) tolerates a genuine already-unloaded answer but
|
||||
# dies on a real `launchctl` failure — never a bare `|| true` that would print `ok` either way.
|
||||
say "stop"
|
||||
RESTART_MARK="$(wc -l < "$OUT" 2>/dev/null || echo 0)"
|
||||
launchctl unload -w "$LAUNCHD_PLIST" 2>/dev/null || true
|
||||
unload_launchd_if_loaded
|
||||
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.
|
||||
# fleetd #504: stop_systemd_if_loaded (above) tolerates that genuine no-op but dies on a real
|
||||
# `systemctl` failure — never a bare `|| true` that would print `ok` either way.
|
||||
say "stop"
|
||||
RESTART_MARK="$(wc -l < "$OUT" 2>/dev/null || echo 0)"
|
||||
systemctl --user stop "$SYSTEMD_UNIT" 2>/dev/null || true
|
||||
stop_systemd_if_loaded
|
||||
ok "systemd --user unit stopped (was already not running)"
|
||||
else
|
||||
RESTART_MARK="$(wc -l < "$OUT" 2>/dev/null || echo 0)"
|
||||
@@ -868,11 +1086,19 @@ tail -n "+$((RESTART_MARK + 1))" "$OUT" 2>/dev/null \
|
||||
|
||||
# Errors since the restart, anchored to the marker so old noise cannot leak in. Keep the fresh
|
||||
# region in a file because the classifier must preserve the order of errors and recoveries.
|
||||
FRESH_LOG="$(mktemp -t fleetd-fresh-log)"
|
||||
FRESH_LOG="$(mktemp -t fleetd-fresh-log.XXXXXX)"
|
||||
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 #512 part 2: the previous daemon's shutdown drain, checked in the same fresh-log region —
|
||||
# see report_shutdown_drain above for the full decision (four outcomes, one of them a deliberate
|
||||
# "cannot tell"). HAD_OLD_PID crosses in whether a previous daemon was actually stopped this run;
|
||||
# see the function's own comment for why that matters.
|
||||
say "previous daemon's shutdown drain"
|
||||
HAD_OLD_PID=0; [ -n "$OLD_PID" ] && HAD_OLD_PID=1
|
||||
report_shutdown_drain "$FRESH_LOG" "$HAD_OLD_PID"
|
||||
|
||||
# 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
|
||||
@@ -882,7 +1108,14 @@ assert_single_daemon "$(running_pid)"
|
||||
say "result"
|
||||
ok "pid $NEW_PID, jar $(jar_id)"
|
||||
if [ "$REDEPLOY_ERROR_COUNT" -eq 0 ]; then
|
||||
ok "no ERROR lines since restart"
|
||||
# fleetd #512 item 4: this line must not print when the shutdown-drain check above found the
|
||||
# previous daemon's drain died, or could not tell — either would make "no ERROR lines" read as a
|
||||
# clean bill of health it is not (an uncaught exception never carries an ERROR token to begin
|
||||
# with, so this count alone cannot see that failure). "complete" and "n/a" are the only two
|
||||
# outcomes report_shutdown_drain sets that mean nothing is wrong there.
|
||||
if [ "$REDEPLOY_DRAIN_STATE" = "complete" ] || [ "$REDEPLOY_DRAIN_STATE" = "n/a" ]; then
|
||||
ok "no ERROR lines since restart"
|
||||
fi
|
||||
elif [ "$REDEPLOY_UNEXPLAINED_ERRORS" -eq 0 ]; then
|
||||
ok "$REDEPLOY_RECOVERED_AMQP_ERRORS AMQP connection reset ERROR lines recovered since restart"
|
||||
else
|
||||
|
||||
@@ -128,8 +128,85 @@ STUB
|
||||
launchd_loaded() { return 1; }
|
||||
result="$(PATH="$bin_dir:$PATH" detect_supervisor)"
|
||||
assert_equals "unclear" "$(supervisor_kind_of "$result")" "a systemd probe error must read as unclear, not none"
|
||||
printf '%s' "$(supervisor_detail_of "$result")" | grep -qF "$SYSTEMD_UNIT" \
|
||||
local detail
|
||||
detail="$(supervisor_detail_of "$result")"
|
||||
printf '%s' "$detail" | grep -qF "$SYSTEMD_UNIT" \
|
||||
|| fail "detail does not name the systemd unit whose probe errored"
|
||||
# fleetd #545: this is the PROBE-RAN-AND-ANSWERED-BADLY case (systemctl actually executed and
|
||||
# wrote to stderr) — it must carry that story and never the SET-UP-FAILED story (mktemp never
|
||||
# even ran here), or the two "unclear" causes have collapsed back into one message that asserts a
|
||||
# cause it did not measure, which is the exact defect this ticket exists to fix.
|
||||
printf '%s' "$detail" | grep -qF "systemctl exited non-zero and reported an error on stderr" \
|
||||
|| fail "detail does not say systemctl ran and answered with stderr: $detail"
|
||||
printf '%s' "$detail" | grep -qF "could not even be set up" \
|
||||
&& fail "detail wrongly claims the probe could not be set up, but systemctl actually ran and answered on stderr: $detail"
|
||||
return 0
|
||||
}
|
||||
|
||||
# fleetd #545 — the companion case to the probe-error test above: here `mktemp` itself fails
|
||||
# (whatever the reason — the historical bug was a GNU-mktemp-rejects-a-template-with-no-Xs case,
|
||||
# but this stub simulates ANY reason the probe's own stderr-capture temp file cannot be created,
|
||||
# e.g. a full or unwritable temp dir) and `systemctl` is never invoked at all. Before this ticket,
|
||||
# this collapsed into the SAME "systemctl exited non-zero and reported an error on stderr" detail
|
||||
# as the sibling test above, which asserts a cause (systemctl ran and answered badly) that was
|
||||
# never measured, because systemctl never ran. This proves the SET-UP-FAILED detail is distinct and
|
||||
# does not claim systemctl said anything.
|
||||
test_detect_supervisor_systemd_probe_setup_failure_is_unclear() {
|
||||
# Re-source first for the same reason test_detect_supervisor_systemd_probe_error_is_unclear does:
|
||||
# restore the REAL probe bodies before driving them through a stub PATH.
|
||||
source "$ROOT/scripts/redeploy-fleetd.sh"
|
||||
local bin_dir result rc=0
|
||||
bin_dir="$TMP/stub-bin-mktemp-fails"
|
||||
mkdir -p "$bin_dir"
|
||||
# A systemctl stub that would fail loudly if it were ever actually invoked — proves the mktemp
|
||||
# failure short-circuits the probe before systemctl runs, not merely that this test forgot to
|
||||
# supply a working systemctl.
|
||||
cat > "$bin_dir/systemctl" <<'STUB'
|
||||
#!/usr/bin/env bash
|
||||
echo "systemctl must never run when mktemp already failed" >&2
|
||||
exit 1
|
||||
STUB
|
||||
chmod +x "$bin_dir/systemctl"
|
||||
cat > "$bin_dir/mktemp" <<'STUB'
|
||||
#!/usr/bin/env bash
|
||||
echo "mktemp: cannot create temp file" >&2
|
||||
exit 1
|
||||
STUB
|
||||
chmod +x "$bin_dir/mktemp"
|
||||
|
||||
PATH="$bin_dir:$PATH" systemd_loaded && rc=0 || rc=$?
|
||||
[ "$rc" -ne 0 ] \
|
||||
|| fail "systemd_loaded must not report loaded=true when its own mktemp setup failed"
|
||||
assert_equals "2" "$SYSTEMD_LOADED_ERRORED" \
|
||||
"systemd_loaded must flag a SETUP failure (2), distinct from a probe-answered-with-stderr failure (1)"
|
||||
|
||||
launchd_installed() { return 1; }
|
||||
launchd_loaded() { return 1; }
|
||||
result="$(PATH="$bin_dir:$PATH" detect_supervisor)"
|
||||
assert_equals "unclear" "$(supervisor_kind_of "$result")" "a systemd probe setup failure must read as unclear, not none"
|
||||
local detail
|
||||
detail="$(supervisor_detail_of "$result")"
|
||||
printf '%s' "$detail" | grep -qF "could not even be set up" \
|
||||
|| fail "detail does not say the probe could not be SET UP: $detail"
|
||||
printf '%s' "$detail" | grep -qF "systemctl exited non-zero and reported an error on stderr" \
|
||||
&& fail "detail wrongly asserts systemctl exited non-zero and reported an error on stderr, but systemctl was never run: $detail"
|
||||
return 0
|
||||
}
|
||||
|
||||
# fleetd #545 — source-text check: every `mktemp -t` template in redeploy-fleetd.sh must contain an
|
||||
# `X` placeholder. BSD mktemp (macOS) tolerates a bare template with no `X`s and just appends its
|
||||
# own random suffix, which is exactly why six such sites survived undetected here — GNU mktemp
|
||||
# (every Linux distribution) refuses a template with fewer than three `X`s and exits non-zero. There
|
||||
# is no BSD-vs-GNU seam to stub on this Mac, so this is a source-text check rather than a
|
||||
# behavioural one, the same shape as test_refuse_drain_gate_call_site_present above. Anchored on
|
||||
# `mktemp -t ` (with the trailing space) so it inspects only the `-t`-style templates this ticket is
|
||||
# about, never the `mktemp -d` calls this file and test-probe-member-credentials.sh already use
|
||||
# (both already carry their own `XXXXXX` and are a different mktemp mode entirely).
|
||||
test_mktemp_dash_t_templates_have_x_placeholders() {
|
||||
local src="$ROOT/scripts/redeploy-fleetd.sh" bad
|
||||
bad="$(grep -n 'mktemp -t ' "$src" | grep -v 'XXX' || true)"
|
||||
[ -z "$bad" ] \
|
||||
|| fail "mktemp -t template(s) with no X placeholder (fails under GNU coreutils): $bad"
|
||||
}
|
||||
|
||||
# fleetd #492 follow-up (Item 1): this must go through the REAL call-site shape at :437-440, not a
|
||||
@@ -521,6 +598,210 @@ test_drain_gate_refusal_no_build_staged_absent() {
|
||||
assert_equals "aborted — nothing changed" "$result" "no-build+staged-absent refusal wording"
|
||||
}
|
||||
|
||||
# fleetd #528 — the four tests above pin drain_gate_refusal(), and that is ALL they pin: they call
|
||||
# the predicate directly and never touch the main flow's call site. That was measured to be not
|
||||
# enough, the same way test_should_swap_true_when_build_ran/test_should_swap_false_when_build_skipped
|
||||
# were not enough for #521: with the main flow reading `die "$(drain_gate_refusal "$DO_BUILD"
|
||||
# "$JAR_STAGED")"`, replacing that whole line with a flat `die "aborted — nothing changed"` left this
|
||||
# suite at exit 0 with zero FAIL lines and byte-identical output to a clean run. Nothing above could
|
||||
# tell the difference, because none of it calls anything at or above the call site itself.
|
||||
#
|
||||
# So these four call refuse_drain_gate() — the function the main flow actually calls, holding the
|
||||
# composed message and the die() together — with die() stubbed to RECORD whether it was called and
|
||||
# with what message, instead of exiting the process. That fails if refuse_drain_gate stops consulting
|
||||
# drain_gate_refusal, mangles what it passes it, or simply never calls die.
|
||||
#
|
||||
# What none of these four can catch: deleting the `refuse_drain_gate "$DO_BUILD" "$JAR_STAGED"` line
|
||||
# from the main flow altogether — see the comment above refuse_drain_gate in redeploy-fleetd.sh for
|
||||
# why no test in this file can do better than that (sourcing stops before the main flow runs).
|
||||
DIED_CALLED=0
|
||||
DIED_MESSAGE=""
|
||||
stub_die_recorder() {
|
||||
DIED_CALLED=0
|
||||
DIED_MESSAGE=""
|
||||
die() { DIED_CALLED=1; DIED_MESSAGE="$*"; }
|
||||
}
|
||||
|
||||
test_refuse_drain_gate_build_ran_staged_present() {
|
||||
source "$ROOT/scripts/redeploy-fleetd.sh"
|
||||
local dir staged
|
||||
dir="$TMP/refuse-drain-build-staged"; mkdir -p "$dir"
|
||||
staged="$dir/fleetd-new.jar"
|
||||
printf 'staged jar bytes' > "$staged"
|
||||
stub_die_recorder
|
||||
refuse_drain_gate 1 "$staged"
|
||||
[ "$DIED_CALLED" = 1 ] \
|
||||
|| fail "refuse_drain_gate build-ran+staged-present must call die, and did not"
|
||||
printf '%s' "$DIED_MESSAGE" | grep -qF "$staged" \
|
||||
|| fail "refuse_drain_gate build-ran+staged-present die message does not name the staged jar"
|
||||
printf '%s' "$DIED_MESSAGE" | grep -qF 'Rerun WITHOUT --no-build' \
|
||||
|| fail "refuse_drain_gate build-ran+staged-present die message is missing the rerun instruction"
|
||||
if printf '%s' "$DIED_MESSAGE" | grep -qF 'nothing changed'; then
|
||||
fail "refuse_drain_gate build-ran+staged-present must not claim nothing changed — the jar already moved"
|
||||
fi
|
||||
source "$ROOT/scripts/redeploy-fleetd.sh"
|
||||
}
|
||||
|
||||
test_refuse_drain_gate_build_ran_staged_absent() {
|
||||
source "$ROOT/scripts/redeploy-fleetd.sh"
|
||||
local dir
|
||||
dir="$TMP/refuse-drain-build-no-staged"; mkdir -p "$dir"
|
||||
stub_die_recorder
|
||||
refuse_drain_gate 1 "$dir/fleetd-new.jar"
|
||||
[ "$DIED_CALLED" = 1 ] \
|
||||
|| fail "refuse_drain_gate build-ran+staged-absent must call die, and did not"
|
||||
assert_equals "aborted — nothing changed" "$DIED_MESSAGE" "refuse_drain_gate build-ran+staged-absent die message"
|
||||
source "$ROOT/scripts/redeploy-fleetd.sh"
|
||||
}
|
||||
|
||||
test_refuse_drain_gate_no_build_staged_present() {
|
||||
source "$ROOT/scripts/redeploy-fleetd.sh"
|
||||
local dir staged
|
||||
dir="$TMP/refuse-drain-no-build-staged"; mkdir -p "$dir"
|
||||
staged="$dir/fleetd-new.jar"
|
||||
printf 'leftover staged jar bytes' > "$staged"
|
||||
stub_die_recorder
|
||||
refuse_drain_gate 0 "$staged"
|
||||
[ "$DIED_CALLED" = 1 ] \
|
||||
|| fail "refuse_drain_gate no-build+staged-present must call die, and did not"
|
||||
assert_equals "aborted — nothing changed" "$DIED_MESSAGE" "refuse_drain_gate no-build+staged-present die message"
|
||||
source "$ROOT/scripts/redeploy-fleetd.sh"
|
||||
}
|
||||
|
||||
test_refuse_drain_gate_no_build_staged_absent() {
|
||||
source "$ROOT/scripts/redeploy-fleetd.sh"
|
||||
local dir
|
||||
dir="$TMP/refuse-drain-no-build-no-staged"; mkdir -p "$dir"
|
||||
stub_die_recorder
|
||||
refuse_drain_gate 0 "$dir/fleetd-new.jar"
|
||||
[ "$DIED_CALLED" = 1 ] \
|
||||
|| fail "refuse_drain_gate no-build+staged-absent must call die, and did not"
|
||||
assert_equals "aborted — nothing changed" "$DIED_MESSAGE" "refuse_drain_gate no-build+staged-absent die message"
|
||||
source "$ROOT/scripts/redeploy-fleetd.sh"
|
||||
}
|
||||
|
||||
# fleetd #528 — closes the one gap the four behavioural tests above cannot: they call
|
||||
# refuse_drain_gate directly, and sourcing stops before the main flow ever runs (the SOURCED guard),
|
||||
# so none of them can prove the main flow still CALLS refuse_drain_gate at all. Same shape as
|
||||
# test_swap_ordered_after_wait_and_before_start: a source-text grep for the real call site. This is
|
||||
# what actually kills the item-1 mutation from the ticket — replacing the main flow's call with a
|
||||
# flat `die "aborted — nothing changed"` removes this exact needle, where none of the behavioural
|
||||
# tests above would even notice.
|
||||
#
|
||||
# The grep ends `|| true`: this file runs under `set -euo pipefail`, so an ABSENT needle would fail
|
||||
# the assignment and `set -e` would kill the whole suite before the `[ -n ... ] || fail` guard below
|
||||
# ever ran — the exact dead-check shape fleetd #528 also flags as a sweep finding (see the PR body).
|
||||
test_refuse_drain_gate_call_site_present() {
|
||||
local src="$ROOT/scripts/redeploy-fleetd.sh" call_line
|
||||
call_line="$(grep -Fn 'refuse_drain_gate "$DO_BUILD" "$JAR_STAGED"' "$src" | head -1 | cut -d: -f1 || true)"
|
||||
[ -n "$call_line" ] \
|
||||
|| fail "could not find the main flow's refuse_drain_gate call site in redeploy-fleetd.sh"
|
||||
}
|
||||
|
||||
# fleetd #504 — the "loaded but not currently running" branches for launchd/systemd used to run
|
||||
# `launchctl unload`/`systemctl --user stop` with `2>/dev/null || true` and print `ok`
|
||||
# unconditionally, so a real supervisor failure (e.g. it cannot reach launchd/the systemd user bus)
|
||||
# read exactly like a harmless already-stopped answer. unload_launchd_if_loaded/
|
||||
# stop_systemd_if_loaded (redeploy-fleetd.sh, right after systemd_loaded) apply systemd_loaded's own
|
||||
# "capture stderr separately — only a non-zero exit WITH stderr is a real failure" pattern to the
|
||||
# WRITE side. Both call the real `launchctl`/`systemctl` binaries directly (they are not overridable
|
||||
# wrapper functions the way launchd_loaded/systemd_loaded are), so these tests put a stub binary
|
||||
# first on PATH — the same technique test_detect_supervisor_systemd_probe_error_is_unclear above
|
||||
# already uses for `systemctl`.
|
||||
test_unload_launchd_if_loaded_dies_on_real_failure() {
|
||||
local bin_dir output rc=0 saved_plist="$LAUNCHD_PLIST"
|
||||
bin_dir="$TMP/stub-bin-launchctl-error"; mkdir -p "$bin_dir"
|
||||
cat > "$bin_dir/launchctl" <<'STUB'
|
||||
#!/usr/bin/env bash
|
||||
echo "Could not find specified service" >&2
|
||||
exit 1
|
||||
STUB
|
||||
chmod +x "$bin_dir/launchctl"
|
||||
LAUNCHD_PLIST="$TMP/fake-fail.plist"
|
||||
output="$(PATH="$bin_dir:$PATH" unload_launchd_if_loaded 2>&1)" || rc=$?
|
||||
LAUNCHD_PLIST="$saved_plist"
|
||||
[ "$rc" -ne 0 ] \
|
||||
|| fail "unload_launchd_if_loaded must die when launchctl exits non-zero AND writes to stderr"
|
||||
printf '%s' "$output" | grep -qF 'launchctl unload' \
|
||||
|| fail "die message does not name the failing launchctl unload command"
|
||||
}
|
||||
|
||||
# Captured via $(...) rather than called bare: unload_launchd_if_loaded's own die() does a hard
|
||||
# `exit`, and calling it directly at this level would let a regression that makes it die on this
|
||||
# clean-negative case kill the WHOLE suite before the `|| fail` below ever ran — printing die's own
|
||||
# message instead of this test's. Inside a command substitution, that `exit` only ends the subshell
|
||||
# (a-guard-is-defeated-by-its-calling-context: the same reason the *_dies_on_real_failure tests
|
||||
# above capture this way), so this test's own message is what actually reaches the report.
|
||||
test_unload_launchd_if_loaded_tolerates_clean_negative() {
|
||||
local bin_dir saved_plist="$LAUNCHD_PLIST" output rc=0
|
||||
bin_dir="$TMP/stub-bin-launchctl-noop"; mkdir -p "$bin_dir"
|
||||
cat > "$bin_dir/launchctl" <<'STUB'
|
||||
#!/usr/bin/env bash
|
||||
exit 1
|
||||
STUB
|
||||
chmod +x "$bin_dir/launchctl"
|
||||
LAUNCHD_PLIST="$TMP/fake-noop.plist"
|
||||
output="$(PATH="$bin_dir:$PATH" unload_launchd_if_loaded 2>&1)" || rc=$?
|
||||
LAUNCHD_PLIST="$saved_plist"
|
||||
[ "$rc" -eq 0 ] \
|
||||
|| fail "unload_launchd_if_loaded must tolerate a clean already-unloaded answer (non-zero exit, empty stderr): $output"
|
||||
}
|
||||
|
||||
test_stop_systemd_if_loaded_dies_on_real_failure() {
|
||||
local bin_dir output rc=0
|
||||
bin_dir="$TMP/stub-bin-systemctl-stop-error"; mkdir -p "$bin_dir"
|
||||
cat > "$bin_dir/systemctl" <<'STUB'
|
||||
#!/usr/bin/env bash
|
||||
echo "Failed to connect to bus: No such file or directory" >&2
|
||||
exit 1
|
||||
STUB
|
||||
chmod +x "$bin_dir/systemctl"
|
||||
output="$(PATH="$bin_dir:$PATH" stop_systemd_if_loaded 2>&1)" || rc=$?
|
||||
[ "$rc" -ne 0 ] \
|
||||
|| fail "stop_systemd_if_loaded must die when systemctl exits non-zero AND writes to stderr"
|
||||
printf '%s' "$output" | grep -qF 'systemctl --user stop' \
|
||||
|| fail "die message does not name the failing systemctl --user stop command"
|
||||
}
|
||||
|
||||
# Same subshell-capture reasoning as test_unload_launchd_if_loaded_tolerates_clean_negative above:
|
||||
# stop_systemd_if_loaded's own die() does a hard `exit`, so this must run inside $(...) or a
|
||||
# regression here would kill the whole suite with die's message instead of this test's.
|
||||
test_stop_systemd_if_loaded_tolerates_clean_negative() {
|
||||
local bin_dir output rc=0
|
||||
bin_dir="$TMP/stub-bin-systemctl-stop-noop"; mkdir -p "$bin_dir"
|
||||
cat > "$bin_dir/systemctl" <<'STUB'
|
||||
#!/usr/bin/env bash
|
||||
exit 1
|
||||
STUB
|
||||
chmod +x "$bin_dir/systemctl"
|
||||
output="$(PATH="$bin_dir:$PATH" stop_systemd_if_loaded 2>&1)" || rc=$?
|
||||
[ "$rc" -eq 0 ] \
|
||||
|| fail "stop_systemd_if_loaded must tolerate a clean already-stopped answer (non-zero exit, empty stderr): $output"
|
||||
}
|
||||
|
||||
# Closes the same gap test_refuse_drain_gate_call_site_present closes for the drain gate: the four
|
||||
# tests above call unload_launchd_if_loaded/stop_systemd_if_loaded directly, and sourcing stops
|
||||
# before the main flow ever runs (the SOURCED guard), so none of them can prove the main flow still
|
||||
# CALLS these two functions instead of the original bare `2>/dev/null || true`. A source-text check,
|
||||
# like test_swap_ordered_after_wait_and_before_start. The call-site needle is anchored (`^ name$`)
|
||||
# so it cannot be satisfied by the comment lines above each call site that merely mention the
|
||||
# function by name.
|
||||
test_stop_branches_call_tolerant_helpers_not_bare_or_true() {
|
||||
local src="$ROOT/scripts/redeploy-fleetd.sh" unload_call_line stop_call_line
|
||||
unload_call_line="$(grep -n '^ unload_launchd_if_loaded$' "$src" | head -1 | cut -d: -f1 || true)"
|
||||
stop_call_line="$(grep -n '^ stop_systemd_if_loaded$' "$src" | head -1 | cut -d: -f1 || true)"
|
||||
[ -n "$unload_call_line" ] \
|
||||
|| fail "could not find the main flow's call to unload_launchd_if_loaded in redeploy-fleetd.sh"
|
||||
[ -n "$stop_call_line" ] \
|
||||
|| fail "could not find the main flow's call to stop_systemd_if_loaded in redeploy-fleetd.sh"
|
||||
if grep -qF 'launchctl unload -w "$LAUNCHD_PLIST" 2>/dev/null || true' "$src"; then
|
||||
fail "the bare 'launchctl unload ... 2>/dev/null || true' defect (fleetd #504) is back in redeploy-fleetd.sh"
|
||||
fi
|
||||
if grep -qF 'systemctl --user stop "$SYSTEMD_UNIT" 2>/dev/null || true' "$src"; then
|
||||
fail "the bare 'systemctl --user stop ... 2>/dev/null || true' defect (fleetd #504) is back in redeploy-fleetd.sh"
|
||||
fi
|
||||
}
|
||||
|
||||
test_no_errors() {
|
||||
cat > "$TMP/no-errors.log" <<'LOG'
|
||||
2026-09-05 12:00:00 INFO fleetd listening
|
||||
@@ -720,12 +1001,172 @@ test_unattributable_quiet_mutation_is_caught() {
|
||||
printf 'Unattributable mutation: FAIL: cross-unattributable recovered: expected 0, got 2\n'
|
||||
}
|
||||
|
||||
# fleetd #512 part 2 — the negative check (scan_uncaught_exceptions). The heart of this half of the
|
||||
# ticket: a fixture with the uncaught-exception shape and NO line carrying an ERROR token at all,
|
||||
# proving the scan finds it without one. A fixture that also carried an ERROR line would pass for
|
||||
# the wrong reason.
|
||||
test_scan_uncaught_exceptions_finds_shape_without_error_token() {
|
||||
cat > "$TMP/scan-died.log" <<'LOG'
|
||||
2026-09-12 10:15:00 INFO fleetd listening on 127.0.0.1:8765
|
||||
Exception in thread "Thread-0" java.lang.NoClassDefFoundError: reactor/core/Exceptions
|
||||
at dev.ltms.fleet.session.SessionManager.drainAll(SessionManager.java:1081)
|
||||
LOG
|
||||
local error_count
|
||||
error_count="$(grep -c ' ERROR ' "$TMP/scan-died.log" || true)"
|
||||
[ "$error_count" = "0" ] \
|
||||
|| fail "test fixture error: scan-died.log unexpectedly carries an ERROR token"
|
||||
scan_uncaught_exceptions "$TMP/scan-died.log"
|
||||
assert_equals 1 "$REDEPLOY_UNCAUGHT_EXCEPTION_COUNT" "scan must find the exception without an ERROR token"
|
||||
printf '%s' "$REDEPLOY_UNCAUGHT_EXCEPTION_SAMPLE" | grep -qF 'NoClassDefFoundError' \
|
||||
|| fail "scan did not capture the matching line as the sample"
|
||||
}
|
||||
|
||||
test_scan_uncaught_exceptions_clean_control() {
|
||||
cat > "$TMP/scan-clean.log" <<'LOG'
|
||||
2026-09-12 10:15:00 INFO fleetd listening on 127.0.0.1:8765
|
||||
2026-09-12 10:15:05 INFO dev.ltms.fleet.session.SessionManager - drain complete: released=0 abandoned=0 (still BUSY at the shutdown deadline)
|
||||
LOG
|
||||
scan_uncaught_exceptions "$TMP/scan-clean.log"
|
||||
assert_equals 0 "$REDEPLOY_UNCAUGHT_EXCEPTION_COUNT" "clean control must find no uncaught exception"
|
||||
assert_equals "" "$REDEPLOY_UNCAUGHT_EXCEPTION_SAMPLE" "clean control sample must be empty"
|
||||
}
|
||||
|
||||
# fleetd #512 part 2 — the positive check (find_drain_complete_line). Both halves of #522's line:
|
||||
# present, and absent.
|
||||
test_find_drain_complete_line_present() {
|
||||
cat > "$TMP/drain-line-present.log" <<'LOG'
|
||||
2026-09-12 10:15:05 INFO dev.ltms.fleet.session.SessionManager - drain complete: released=2 abandoned=1 (still BUSY at the shutdown deadline)
|
||||
LOG
|
||||
find_drain_complete_line "$TMP/drain-line-present.log"
|
||||
printf '%s' "$REDEPLOY_DRAIN_COMPLETE_LINE" | grep -qF 'released=2 abandoned=1' \
|
||||
|| fail "find_drain_complete_line did not capture the present line"
|
||||
}
|
||||
|
||||
test_find_drain_complete_line_absent() {
|
||||
cat > "$TMP/drain-line-absent.log" <<'LOG'
|
||||
2026-09-12 10:15:00 INFO fleetd listening on 127.0.0.1:8765
|
||||
LOG
|
||||
find_drain_complete_line "$TMP/drain-line-absent.log"
|
||||
assert_equals "" "$REDEPLOY_DRAIN_COMPLETE_LINE" "find_drain_complete_line must report empty when absent"
|
||||
}
|
||||
|
||||
# fleetd #512 part 2 — report_shutdown_drain, the composite decision+action function the main flow
|
||||
# calls unconditionally (same shape as swap_if_built/refuse_drain_gate, #521/#528). These four cover
|
||||
# the four outcomes named in the ticket's "trap": complete, died, unknown ("cannot tell" — neither a
|
||||
# pass nor a failure), and n/a (no previous daemon was actually stopped this run).
|
||||
#
|
||||
# Deliberately NOT run inside `$(...)`: report_shutdown_drain sets REDEPLOY_DRAIN_STATE as a global
|
||||
# side effect that these tests need to read back afterward, and a command substitution forks a
|
||||
# subshell that global assignment would not survive (the exact trap documented above
|
||||
# detect_supervisor in redeploy-fleetd.sh, for the same reason). Plain output redirection to a file
|
||||
# does not fork a subshell, so it is used to capture what was printed instead.
|
||||
test_report_shutdown_drain_died_without_error_token() {
|
||||
cat > "$TMP/drain-died.log" <<'LOG'
|
||||
2026-09-12 10:15:00 INFO fleetd listening on 127.0.0.1:8765
|
||||
2026-09-12 10:15:05 INFO dev.ltms.fleet.Fleetd - shutting down
|
||||
Exception in thread "Thread-0" java.lang.NoClassDefFoundError: reactor/core/Exceptions
|
||||
at dev.ltms.fleet.session.SessionManager.drainAll(SessionManager.java:1081)
|
||||
LOG
|
||||
local error_count
|
||||
error_count="$(grep -c ' ERROR ' "$TMP/drain-died.log" || true)"
|
||||
[ "$error_count" = "0" ] \
|
||||
|| fail "test fixture error: drain-died.log unexpectedly carries an ERROR token"
|
||||
|
||||
report_shutdown_drain "$TMP/drain-died.log" 1 > "$TMP/drain-died-output" 2>&1
|
||||
assert_equals "died" "$REDEPLOY_DRAIN_STATE" "died fixture must set REDEPLOY_DRAIN_STATE=died"
|
||||
assert_equals 1 "$REDEPLOY_UNCAUGHT_EXCEPTION_COUNT" "died fixture uncaught-exception count"
|
||||
grep -qF 'NoClassDefFoundError' "$TMP/drain-died-output" \
|
||||
|| fail "report_shutdown_drain did not report the uncaught-exception shape it found"
|
||||
grep -qF 'DIED' "$TMP/drain-died-output" \
|
||||
|| fail "report_shutdown_drain did not report the drain as DIED"
|
||||
}
|
||||
|
||||
test_report_shutdown_drain_complete_control() {
|
||||
cat > "$TMP/drain-complete.log" <<'LOG'
|
||||
2026-09-12 10:15:00 INFO fleetd listening on 127.0.0.1:8765
|
||||
2026-09-12 10:15:05 INFO dev.ltms.fleet.Fleetd - shutting down
|
||||
2026-09-12 10:15:05 INFO dev.ltms.fleet.session.SessionManager - drain complete: released=3 abandoned=0 (still BUSY at the shutdown deadline)
|
||||
LOG
|
||||
report_shutdown_drain "$TMP/drain-complete.log" 1 > "$TMP/drain-complete-output" 2>&1
|
||||
assert_equals "complete" "$REDEPLOY_DRAIN_STATE" "complete-control fixture must set REDEPLOY_DRAIN_STATE=complete"
|
||||
assert_equals 0 "$REDEPLOY_UNCAUGHT_EXCEPTION_COUNT" "complete-control fixture must find no uncaught exception"
|
||||
grep -qF 'released=3 abandoned=0' "$TMP/drain-complete-output" \
|
||||
|| fail "report_shutdown_drain did not report the drain-complete counts"
|
||||
}
|
||||
|
||||
test_report_shutdown_drain_unknown_cannot_tell() {
|
||||
cat > "$TMP/drain-unknown.log" <<'LOG'
|
||||
2026-09-12 10:15:00 INFO fleetd listening on 127.0.0.1:8765
|
||||
2026-09-12 10:15:05 INFO dev.ltms.fleet.Fleetd - shutting down
|
||||
LOG
|
||||
report_shutdown_drain "$TMP/drain-unknown.log" 1 > "$TMP/drain-unknown-output" 2>&1
|
||||
assert_equals "unknown" "$REDEPLOY_DRAIN_STATE" "cannot-tell fixture must set REDEPLOY_DRAIN_STATE=unknown"
|
||||
grep -qF 'cannot tell' "$TMP/drain-unknown-output" \
|
||||
|| fail "report_shutdown_drain did not say it could not tell"
|
||||
if grep -qF ' ok' "$TMP/drain-unknown-output"; then
|
||||
fail "cannot-tell outcome must not be printed via ok() — it is neither a pass nor a failure"
|
||||
fi
|
||||
}
|
||||
|
||||
# A cold start (or a restart where nothing was actually stopped) has no previous-daemon shutdown
|
||||
# window to have an opinion about at all. This fixture's log content looks exactly like a died drain
|
||||
# — proving the had_previous_daemon=0 gate is actually consulted, not merely documented: without it,
|
||||
# this would misreport "died" or "unknown" on every clean cold start.
|
||||
test_report_shutdown_drain_no_previous_daemon_is_na() {
|
||||
cat > "$TMP/drain-na.log" <<'LOG'
|
||||
Exception in thread "Thread-0" java.lang.NoClassDefFoundError: reactor/core/Exceptions
|
||||
LOG
|
||||
report_shutdown_drain "$TMP/drain-na.log" 0 > "$TMP/drain-na-output" 2>&1
|
||||
assert_equals "n/a" "$REDEPLOY_DRAIN_STATE" "no-previous-daemon fixture must set REDEPLOY_DRAIN_STATE=n/a even though the log content looks like a died drain"
|
||||
grep -qF 'nothing to check' "$TMP/drain-na-output" \
|
||||
|| fail "report_shutdown_drain did not report that there was nothing to check"
|
||||
}
|
||||
|
||||
# fleetd #512 — closes the gap none of the seven tests above can: they call report_shutdown_drain
|
||||
# directly, and sourcing stops before the main flow ever runs (the SOURCED guard), so none of them
|
||||
# can prove the main flow still calls it at all. Same shape as test_refuse_drain_gate_call_site_present
|
||||
# and test_swap_ordered_after_wait_and_before_start: a source-text grep for the real call site, plus
|
||||
# an ordering check against its neighbours in the verify/result flow.
|
||||
test_report_shutdown_drain_call_site_present() {
|
||||
local src="$ROOT/scripts/redeploy-fleetd.sh" call_line
|
||||
call_line="$(grep -Fn 'report_shutdown_drain "$FRESH_LOG" "$HAD_OLD_PID"' "$src" | head -1 | cut -d: -f1 || true)"
|
||||
[ -n "$call_line" ] \
|
||||
|| fail "could not find the main flow's report_shutdown_drain call site in redeploy-fleetd.sh"
|
||||
}
|
||||
|
||||
test_report_shutdown_drain_ordered_after_classify_and_before_result() {
|
||||
local src="$ROOT/scripts/redeploy-fleetd.sh" classify_line drain_line result_line
|
||||
classify_line="$(grep -Fn 'classify_amqp_connection_errors "$FRESH_LOG"' "$src" | tail -1 | cut -d: -f1 || true)"
|
||||
drain_line="$(grep -Fn 'report_shutdown_drain "$FRESH_LOG" "$HAD_OLD_PID"' "$src" | head -1 | cut -d: -f1 || true)"
|
||||
result_line="$(grep -Fn 'say "result"' "$src" | head -1 | cut -d: -f1 || true)"
|
||||
[ -n "$classify_line" ] || fail "could not find the classify_amqp_connection_errors call site"
|
||||
[ -n "$drain_line" ] || fail "could not find the report_shutdown_drain call site"
|
||||
[ -n "$result_line" ] || fail "could not find the result section"
|
||||
[ "$drain_line" -gt "$classify_line" ] \
|
||||
|| fail "report_shutdown_drain (line $drain_line) is not after classify_amqp_connection_errors (line $classify_line)"
|
||||
[ "$drain_line" -lt "$result_line" ] \
|
||||
|| fail "report_shutdown_drain (line $drain_line) is not before the result section (line $result_line)"
|
||||
}
|
||||
|
||||
# fleetd #512 item 4 — the summary line must not read as reassurance when the shutdown-drain check
|
||||
# found something wrong (or could not tell). Sourcing stops before the main flow runs, so this is a
|
||||
# source-text check like test_drain_gate_abort_message_says_no_no_build above.
|
||||
test_no_error_lines_message_gated_by_drain_state() {
|
||||
local src="$ROOT/scripts/redeploy-fleetd.sh" block
|
||||
block="$(grep -B2 -F 'ok "no ERROR lines since restart"' "$src")"
|
||||
[ -n "$block" ] || fail "could not find the 'no ERROR lines since restart' line in redeploy-fleetd.sh"
|
||||
printf '%s' "$block" | grep -qF 'REDEPLOY_DRAIN_STATE' \
|
||||
|| fail "'no ERROR lines since restart' is not guarded by the shutdown-drain outcome (fleetd #512 item 4)"
|
||||
}
|
||||
|
||||
test_detect_supervisor_launchd_only
|
||||
test_detect_supervisor_systemd_only
|
||||
test_detect_supervisor_none
|
||||
test_detect_supervisor_systemd_installed_not_loaded_is_unclear
|
||||
test_detect_supervisor_launchd_installed_not_loaded_is_unclear
|
||||
test_detect_supervisor_systemd_probe_error_is_unclear
|
||||
test_detect_supervisor_systemd_probe_setup_failure_is_unclear
|
||||
test_mktemp_dash_t_templates_have_x_placeholders
|
||||
test_require_drivable_supervisor_refuses_ambiguous
|
||||
test_require_drivable_supervisor_refuses_unclear
|
||||
test_require_drivable_supervisor_accepts_known_kinds
|
||||
@@ -753,6 +1194,16 @@ test_drain_gate_refusal_build_ran_staged_present
|
||||
test_drain_gate_refusal_build_ran_staged_absent
|
||||
test_drain_gate_refusal_no_build_staged_present
|
||||
test_drain_gate_refusal_no_build_staged_absent
|
||||
test_refuse_drain_gate_build_ran_staged_present
|
||||
test_refuse_drain_gate_build_ran_staged_absent
|
||||
test_refuse_drain_gate_no_build_staged_present
|
||||
test_refuse_drain_gate_no_build_staged_absent
|
||||
test_refuse_drain_gate_call_site_present
|
||||
test_unload_launchd_if_loaded_dies_on_real_failure
|
||||
test_unload_launchd_if_loaded_tolerates_clean_negative
|
||||
test_stop_systemd_if_loaded_dies_on_real_failure
|
||||
test_stop_systemd_if_loaded_tolerates_clean_negative
|
||||
test_stop_branches_call_tolerant_helpers_not_bare_or_true
|
||||
test_no_errors
|
||||
test_recovery_patterns_match_source
|
||||
test_attributed_recovered_connection_error
|
||||
@@ -764,4 +1215,15 @@ test_other_error_is_unexplained
|
||||
test_recovery_requirement_mutation_is_caught
|
||||
test_shared_counter_mutation_is_caught
|
||||
test_unattributable_quiet_mutation_is_caught
|
||||
test_scan_uncaught_exceptions_finds_shape_without_error_token
|
||||
test_scan_uncaught_exceptions_clean_control
|
||||
test_find_drain_complete_line_present
|
||||
test_find_drain_complete_line_absent
|
||||
test_report_shutdown_drain_died_without_error_token
|
||||
test_report_shutdown_drain_complete_control
|
||||
test_report_shutdown_drain_unknown_cannot_tell
|
||||
test_report_shutdown_drain_no_previous_daemon_is_na
|
||||
test_report_shutdown_drain_call_site_present
|
||||
test_report_shutdown_drain_ordered_after_classify_and_before_result
|
||||
test_no_error_lines_message_gated_by_drain_state
|
||||
printf 'PASS: redeploy log classifier\n'
|
||||
|
||||
Reference in New Issue
Block a user