Compare commits

...

17 Commits

Author SHA1 Message Date
Dai Ha 8b4320ed24 fleetd #553: split sentHandled's two meanings so onDelivered's own throw still completes the future
CI / contract (pull_request) Successful in 54s
CI / build (pull_request) Successful in 1m43s
Lead review of PR #557 (ticket comment 16916) found one path left open: sentHandled
is set to true BEFORE onDelivered() runs (correctly, per the earlier fix), so when
onDelivered() itself throws on the normal path, the finally's 'if (sent != null &&
!sentHandled)' guard skipped the whole recovery -- completion included -- and left
sent.delivered() pending forever for a message that really was delivered.

sentHandled must guard only the onDelivered RE-CALL (the permanent-suppression
hazard), never the future completion, since CompletableFuture.complete/
completeExceptionally are idempotent and a no-op on the already-handled path.
Split the one flag's two jobs: the outer 'if (sent != null)' now always runs the
recovery block, and '!sentHandled' moved onto just the onDelivered call inside it.

Added anOnDeliveredThrowOnTheNormalPathStillCompletesTheDeliveryFuture, proven with
the lead's own mutation (reverting !sentHandled onto the outer if): the new test
goes red while anOnDeliveredThrowAfterItsOwnRegistrationDoesNotRunASecondTime stays
green, showing the two concerns are genuinely separate.
2026-09-12 15:45:36 +07:00
Dai Ha d4a51c6274 fleetd #553: register the rendezvous waiter in onStatus's finally backstop, not just the delivery future
CI / contract (pull_request) Successful in 50s
CI / build (pull_request) Successful in 2m12s
The previous try/finally around onStatus's post-monitor region completed
sent.delivered() but never registered sent.token().waiter() when an earlier
listener threw. That waiter is registered only by turnListener.onDelivered(),
inside the very if (sent != null) block the finally backstops, so a caller
was told its send landed and then waited out its full timeout for an answer
that could never resolve (worse than a plain hang).

The finally now does that block's whole job on the unhandled path: it calls
onDelivered() (when sendError == null) before completing the future, guarded
by its own try/catch(Throwable) so a failure there cannot mask the original
throwable. A sentHandled flag, set true at the START of the normal block
(before any side effect), tells the finally whether that already ran, so a
throw partway through onDelivered cannot trigger a second, late captureBaseline
that would permanently suppress the turn's completion.

Also removed a leftover duplicated forget.accept(target) call (with a stray
'MUTATION-TEST-3' comment) in the notReady block — residue from the previous
worker's own mutation testing that was not fully reverted.
2026-09-12 15:35:53 +07:00
ltms 93a9ed3f83 Merge #549: widen Injector's delivery catch to Throwable (#546)
CI / contract (push) Successful in 56s
CI / build (push) Successful in 1m43s
Closes the re-delivery window that merging #543 opened. I caused that; this closes it the same day.

Verified by me on the branch at 87871ea, base 0b032f5 (not stale).

Production diff is 11 lines: `catch (RuntimeException)` -> `catch (Throwable)` at Injector.java:391,
and the `sendError` local widened to `Throwable` so it compiles. I checked every use of `sendError`
myself — :533 `getMessage()` and :534 `completeExceptionally(Throwable)` — so the wider type
reaches nothing that needed the narrow one.

Build on the branch: exit 0, Tests run: 1719, Failures: 0, Errors: 0, Skipped: 0 (1716 on main plus
the 3 new tests). #459's javadoc reference gate: exit 0, 0 reference errors. Gitea CI run 1803 on
87871ea: success.

Three mutations of my own, none of them the ones the worker used:
1. Reverted the catch to `RuntimeException`. RED: `anErrorFromSendDoesNotRedeliverOnASecondRound`
   "expected: <1> but was: <2>" prompt calls, and `anErrorFromSendRemovesTheMessage...` reporting
   the Error escaping `onStatus`. That second message is the defect itself, stated by the test.
2. Deleted the `t.queue.poll()` in the catch arm. RED on the new test AND on the pre-existing
   `sendFailureDropsMessageAndFailsItsFuture` — so the new test is not carrying that behaviour alone.
3. Wrote `DELIVERED` instead of `NOT_DELIVERED` at :400, the catch arm only. RED on the new test and
   on `aHerdrExceptionFromSendStillProducesNotDeliveredUnchanged`, which is the control that proves
   the widening did not quietly change the ordinary path.

All three restored; sha256 back to 97c560b6e33fc49a1772abec92e5bbab613f8991deba2220d30221a2f546ba14.
Green control after the restores: InjectorTest 37/37, exit 0.

My third mutation did not apply on its first attempt — it asserted a unique match on
`p.state = Pending.State.NOT_DELIVERED;`, which occurs twice (:400 and :419), so the script wrote
nothing and the test run came back exit 0. That is not a surviving mutant, it is a non-result
wearing the same clothes. The pristine-anchor count catching it is the only reason I noticed.

Deliberately NOT fixed here, filed as #551: the catch arm assumes that reaching it means nothing was
sent, and nothing establishes that. `agent.prompt` pastes and submits in one call, and every failure
in the response half of `UnixSocketHerdrClient.call()` — dropped connection, malformed line, error
result — is a `HerdrException`, which is a `RuntimeException`, which this catch arm already caught
before today. So "records NOT_DELIVERED for a delivery that happened" is older than this PR and is
not created by it. The fleet01 lead argued it was a trap inside this change and asked to be argued
out of it before the merge; the measurement above is the argument, and their underlying diagnosis is
right and is now #551 with their wording on it.
2026-09-12 09:39:41 +02:00
ltms 0b032f5a1a Merge #548: fix mktemp -t templates for GNU coreutils, split the unclear-supervisor detail (#545)
CI / contract (push) Successful in 55s
CI / build (push) Successful in 2m16s
Verified by me on the branch, not on the worker's report.

Source audit, on a scratch worktree at a476a14:
- Every `mktemp -t` site in scripts/ now carries an X placeholder. The only remaining
  `mktemp -t` text with no X is a prose comment in the test file, not a call.

Suite, macOS (/bin/bash 3.2.57 and env bash 5.3.9):
- exit 0, anchored `^FAIL:` count 0. Unanchored `FAIL:` count 3 — the suite's own internal
  mutation-cell fixture lines, same as main.
- Test functions defined vs invoked: 67/67, `comm -3` empty.

Three mutations of my own, none of them one the worker used:
1. Removed `.XXXXXX` from the `fleetd-fresh-log` site (line 1089). Pristine anchor count went
   1 -> 0, so the mutation really applied. Suite exit 1, FAIL named that exact line.
2. Made the state-2 branch in `detect_supervisor` unreachable (`= 2` -> `= 9`). Suite exit 1:
   "a systemd probe setup failure must read as unclear, not none: expected unclear, got none".
3. Made `systemd_loaded` set the old value (`=2` -> `=1`) on setup failure. Suite exit 1:
   "must flag a SETUP failure (2), distinct from a probe-answered-with-stderr failure (1)".
All three restored; `shasum -a 256` back to 77fe15e5945c7d4ef9b1a2cd8f46e6f1d0be5004f595ea03964d1d7c16e859f7,
the same hash the worker reported independently. Green control after the restores: exit 0, 0
anchored FAILs.

A fourth attempt did not count. A perl `\Q...\E` pattern silently interpolated the shell
variables in it, so the file was never changed and the suite's exit 0 meant nothing. The proof
cell caught it: the pristine anchor count was still 1 after the "mutation". A mutation that did
not apply is not a surviving mutant.

The measurement macOS cannot make: I ran both arms under GNU coreutils 9.1 in a
debian:bookworm-slim container, with a `systemctl` stub that exits non-zero and writes NOTHING to
stderr — a clean negative answer.

  main (a476a14's base):
    mktemp: too few X's in template 'systemd-loaded-err'
    systemd_loaded rc=1  SYSTEMD_LOADED_ERRORED=1
    detect_supervisor => unclear | "systemctl exited non-zero and reported an error on stderr,
                                   not a clean negative — e.g. it cannot reach the user bus"

  this branch:
    systemd_loaded rc=1  SYSTEMD_LOADED_ERRORED=0
    detect_supervisor => none

So on Linux, main tells the operator that systemctl answered badly, when systemctl ran fine and
gave a clean negative. The message named a cause that was never measured. This branch removes it.

Found while doing this, NOT part of this PR, ticket to follow: the shell suite cannot run on
Linux at all. It dies at the first `shasum` call with "command not found" and exit 127, and the
anchored `^FAIL:` count reads 0 — identical to a green run. Gitea CI never runs this suite, so
nothing caught it.
2026-09-12 09:30:30 +02:00
ltms fad99c4c5e Merge #547: record the finally non-goal on drainAll's completion line
CI / contract (push) Successful in 49s
CI / build (push) Successful in 2m10s
Comment only. No behaviour change.

Verified by me before merging:
- `mvn -B clean install` in a scratch worktree: exit 0, Tests run: 1716, Failures: 0, Errors: 0,
  Skipped: 0. Same count as main, as expected for a javadoc-only change.
- #459's javadoc reference gate: exit 0, 0 reference errors.
- Gitea CI run 1801 on 5eb4267: success.

The constraint came from the fleet01 lead. Their point: the value of the `drain complete` line is
that it is MISSING when a drain does not finish. A `finally` block would print it after a drain
that threw, with partial counts, and destroy both halves at once. The code already avoids this;
what was missing was the sentence that stops a reviewer putting it back.
2026-09-12 09:29:46 +02:00
Dai Ha 87871eaefb fleetd #546: widen Injector's delivery catch to Throwable, stop re-delivery on Error
CI / contract (pull_request) Successful in 1m29s
CI / build (pull_request) Successful in 1m49s
Injector.java:391 caught only RuntimeException around the herdr send seam. PR #543
(fleetd #538) widened StatusPoller's per-target catch to Throwable so the polling
loop now survives an Error there, which means it comes back round — and Injector's
narrower catch let the poisoned message stay QUEUED (the loop peeks, not polls),
so the next round re-sent the same text into the member's pane.

Widen the catch to Throwable, matching #543 one layer down. sendError's declared
type widens from RuntimeException to Throwable to keep compiling; its only consumer
(CompletableFuture.completeExceptionally(Throwable)) already accepts that type, so
no other caller-visible behavior changes. The ordinary HerdrException/RuntimeException
path is unchanged.

Adds three tests: an Error at the send seam is dropped and marked NOT_DELIVERED, a
second onStatus round does not re-send it, and a HerdrException control proves the
ordinary path is untouched.
2026-09-12 14:26:27 +07:00
Dai Ha a476a14f1c fleetd #545: fix mktemp -t templates for GNU coreutils, split unclear-supervisor detail
CI / contract (pull_request) Successful in 1m16s
CI / build (pull_request) Successful in 2m28s
Every mktemp -t template in redeploy-fleetd.sh lacked an X placeholder. BSD mktemp
(macOS) tolerates that and appends its own suffix; GNU mktemp (every Linux
distribution) refuses it and exits non-zero. All six sites now use .XXXXXX.

detect_supervisor's 'unclear' detail used to cover two different facts with one
message that always named 'systemctl exited non-zero and reported an error on
stderr' — even when systemctl was never run, because mktemp failed first. The
SYSTEMD_LOADED_ERRORED/SYSTEMD_INSTALLED_ERRORED flags now carry a third value
(2 = the probe's own mktemp setup failed) alongside the existing 1 (systemctl ran
and answered badly on stderr), and detect_supervisor gives each its own detail
text. kind stays 'unclear' in both cases; require_drivable_supervisor is unchanged.

Tests added to scripts/test-redeploy-fleetd.sh:
- test_mktemp_dash_t_templates_have_x_placeholders: source-text check, fails if
  any mktemp -t template lacks an X.
- test_detect_supervisor_systemd_probe_setup_failure_is_unclear: proves the
  SET-UP-FAILED detail when mktemp itself fails (systemctl never runs).
- test_detect_supervisor_systemd_probe_error_is_unclear: extended with assertions
  that the PROBE-ANSWERED-WITH-STDERR detail is present and the SET-UP-FAILED
  wording is absent, so swapping the two messages fails a test in both
  directions.
2026-09-12 14:22:57 +07:00
Dai Ha 5eb4267a4a fleetd #512 follow-up: record why the drain-complete line must not move into a finally
CI / contract (pull_request) Successful in 1m3s
CI / build (pull_request) Successful in 2m8s
The javadoc said the log.info fires "every time", which reads as an invitation
to the exact edit that destroys it. The absence of the line is the signal that
the drain died, so a finally would remove the signal and print partial counts in
the same change.

Raised by the fleet01 lead from their 2026-09-10 incident: that drain is known to
have died only because it threw and left a stack trace. A drain that hung, or
returned early on a condition, leaves no trace, no ERROR token and no priority --
only a missing line.

Comment only. No behaviour change.
2026-09-12 14:21:15 +07:00
ltms 7611b69667 Merge pull request 'fleetd #504 item 1: stop the false ok on the loaded-but-not-running stop path' (#541) from worker/504-failed-reported-clean-3cfd66-3 into main
CI / build (push) Successful in 1m35s
CI / contract (push) Successful in 1m42s
2026-09-12 09:09:00 +02:00
ltms cc302fe4af Merge pull request 'fleetd #538: recover polling loops after errors' (#543) from worker/538-loop-dies-on-error-4a5eeb-6 into main
CI / contract (push) Successful in 51s
CI / build (push) Successful in 2m10s
2026-09-12 09:02:48 +02:00
ltms 4a8a780274 Merge pull request 'fleetd #426: pin FleetHealthMonitor.coverage and its HealthCoverageSource call site' (#542) from worker/426-health-coverage-ef1fd4-4 into main
CI / contract (push) Successful in 1m16s
CI / build (push) Successful in 2m28s
2026-09-12 08:55:46 +02:00
Dai Ha 343ce0f4c0 fleetd #538: recover polling loops after errors
CI / contract (pull_request) Successful in 44s
CI / build (pull_request) Successful in 2m0s
2026-09-12 13:47:44 +07:00
ltms cec3e191d4 Merge pull request 'fleetd #459: lint Javadoc references in CI' (#539) from worker/459-broken-link-targets-cadc17-5 into main
CI / contract (push) Successful in 1m4s
CI / build (push) Successful in 2m9s
2026-09-12 08:46:23 +02:00
Dai Ha 1850a5f324 fleetd #426: pin FleetHealthMonitor.coverage and its HealthCoverageSource call site
CI / contract (pull_request) Successful in 48s
CI / build (pull_request) Successful in 2m16s
FleetHealthMonitor.coverage had zero references in the test tree — not the
method, not either output string, not the field it populates. Inverting
`enabled`, swapping "full"/"detection-only", or breaking the argument pairing
at the HealthCoverageSource call site in Fleetd.java all shipped a green
build.

Extract the HealthCoverageSource lambda out of Fleetd.main into a
package-private static factory (healthCoverageSource(ConfigRef)), the same
shape capacitySource/quarantineSource already use for the identical
argument-pairing risk (fleetd #415). #407's "keep the config invalid, assert
on the log line before validateAll() throws" option does not apply here: this
call site is built well after validateAll() and after a real herdr socket
connect, so driving it through a real Fleetd.main would require the socket
I/O this ticket's tests must not do.

Add FleetHealthMonitorCoverageTest (the three-branch method itself) and
FleetdHealthCoverageSourceWiringTest (the call site, via a real
FleetConfig.load + ConfigRef against @TempDir fixtures, including a hot
notifications-reload case). Output strings are unchanged — "detection-only"
is still what a live fleet_list reports today.

Measured: all three mutations killed by the new tests.
2026-09-12 13:43:50 +07:00
ltms 57cd96f5e6 Merge pull request 'fleetd #537: pin CapturedLog.close()'s appender-detach and setLevel-immunity contracts' (#540) from worker/537-capturedlog-close-e4c437-2 into main
CI / contract (push) Successful in 47s
CI / build (push) Successful in 2m9s
2026-09-12 08:42:19 +02:00
Dai Ha 202e37e3b3 fleetd #537: pin CapturedLog.close()'s appender-detach and setLevel-immunity contracts
CI / contract (pull_request) Successful in 1m2s
CI / build (pull_request) Successful in 1m36s
Only the level-restore half of close() was pinned before this
(WorktreeSessionManagerTest). Deleting logger.detachAppender(appender)
from close() left mvn clean install green (1701 tests, 0 failures) --
the appender-detach half of the contract was unmeasured.

Adds CapturedLogTest with three tests, each using a logger name no
production class uses:
- closeDetachesTheAppenderSoALaterLogIsNotCaptured: an event logged
  after close() must not land in events().
- closeRestoresTheLevelCapturedAtOpen: the helper's headline contract
  in one place, independent of any production class.
- setLevelDuringCaptureDoesNotChangeWhatCloseRestores: setLevel()'s
  own javadoc claim that close() always restores the level captured
  at construction, never a value set through setLevel() mid-capture.

Test-only change; CapturedLog.java itself is untouched.
2026-09-12 13:36:55 +07:00
Dai Ha 90253f832d fleetd #459: lint Javadoc references in CI
CI / contract (pull_request) Successful in 1m30s
CI / build (pull_request) Successful in 2m9s
2026-09-12 13:36:37 +07:00
19 changed files with 1413 additions and 108 deletions
+4
View File
@@ -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.
@@ -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);
}
}
@@ -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));
}
}
@@ -2,11 +2,15 @@ package dev.ltms.fleet.inject;
import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.spi.ILoggingEvent;
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;
@@ -664,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();
}
}
}
@@ -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);
}
}
@@ -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());
}
}
+43 -17
View File
@@ -96,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
@@ -238,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.
@@ -251,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=$?
@@ -275,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=$?
@@ -302,7 +317,7 @@ systemd_loaded() {
# 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)"; then
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."
@@ -317,7 +332,7 @@ unload_launchd_if_loaded() {
stop_systemd_if_loaded() {
local err_file rc=0
if ! err_file="$(mktemp -t systemd-stop-err)"; then
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."
@@ -349,6 +364,14 @@ stop_systemd_if_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)`):
@@ -382,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
@@ -830,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" \
@@ -1060,7 +1086,7 @@ 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"
+80 -1
View File
@@ -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
@@ -1088,6 +1165,8 @@ 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