Compare commits

...

22 Commits

Author SHA1 Message Date
Dai Ha dbf6fef0e9 config: guard FleetConfig.withDefaults() against silently dropping a component
CI / contract (pull_request) Successful in 1m13s
CI / build (pull_request) Successful in 1m53s
Adding a component to FleetConfig follows an established pattern: the
record grows by one arg, and a back-compat constructor is added at the
OLD arity so existing callers keep compiling. That back-compat
constructor also silently captures withDefaults()'s own literal-arity
'return new FleetConfig(...)' call the next time this happens, since
that call is now a legal overload match too. It compiles, every other
test passes, and the new component is defaulted away on every load().
This is not hypothetical - it happened live while building the (now
parked) idle-sleep-guard PR, caught only because that branch's own new
tests asserted on the new field.

Add a reflective test that builds a FleetConfig through the true
canonical constructor (resolved by record-component types, not arg
count - the same pattern ConfigRefTopLevelReportingCoverageTest already
uses in this file) with a real, non-null value in every component, runs
the real withDefaults(), and asserts every value survives unchanged.
Never hardcodes the arity - it enumerates
FleetConfig.class.getRecordComponents() - so it keeps working as the
record grows. No back-compat constructor is touched or removed.
2026-09-05 05:51:59 +07:00
Dai Ha b6b88c5f1c #334: pin the fresh-owner gate on ask()'s timeout teardown
CI / build (push) Successful in 2m7s
CI / contract (push) Successful in 34m37s
2026-09-04 17:12:49 +07:00
Dai Ha 86dddfe240 Merge #334: ask()'s timeout closes the turn before it forgets the task mapping 2026-09-04 17:08:40 +07:00
Dai Ha 0d5944af63 fleetd #334: close ask()'s turn before forgetting its Task, closing the last stranding window
CI / build (pull_request) Successful in 2m14s
CI / contract (pull_request) Successful in 19m14s
ask()'s TimeoutException catch used to run clearAsyncQuestion(turnId, true) -- forgetting the
Task's asyncTasksByTurn mapping -- before rendezvous.closeAsk(turnId) ran in the shared finally.
Between those two calls the ask was still "answerable" (askSession(turnId) non-null) but the Task
mapping was already gone, so a racing answer() call found task == null, skipped
finishAsyncTask, and stranded the async ticket at PENDING even though answer() itself reported a
result. #329 fixed one step of this same race; this closes the remaining one.

The fix reorders the fresh owner's teardown: closeAsk runs first, then markAskTimedOut and
clearAsyncQuestion. A racing answer() call now either sees the ask still open (and the Task
mapping guaranteed intact) or sees it already closed (STALE_TURN, before it ever reaches
asyncTasksByTurn). It also gates the whole block by ticket.fresh(), matching the invariant the
finally block already states ("only the fresh owner tears down the shared turn") -- a duplicate
coalesced ask() timing out no longer forgets bookkeeping the fresh owner still needs.

Adds aLateAnswerDuringAskTimeoutTeardownStillCompletesTheAsyncTicket, which pins the exact window
with a new test-only hook (askTimeoutRaceHookForTest) and proves both invariants: a late answer()
racing the timeout sees STALE_TURN, and the async ticket still resolves DONE from the worker's
real reply. Reverting the reorder (verified locally, not committed) makes this test fail with
"expected STALE_TURN but was TIMED_OUT_WORKING".
2026-09-04 17:06:15 +07:00
Dai Ha 4a5030a5c6 #348: drop a chrome skip that cannot fire, and pin the live pattern shape
CI / contract (push) Successful in 2m10s
CI / build (push) Successful in 2m12s
2026-09-04 16:57:36 +07:00
Dai Ha c1c8794c48 Merge #348: a member's prose about a usage limit no longer quarantines a credential 2026-09-04 16:53:19 +07:00
Dai Ha 65a78932c1 Merge #335: a per-task cleanup throw in abandon() no longer strands the tasks behind it
CI / contract (push) Successful in 1m24s
CI / build (push) Successful in 2m8s
2026-09-04 16:44:04 +07:00
Dai Ha f429ca1a50 Avoid exhaustion cooldown for member prose 2026-09-04 16:43:16 +07:00
Dai Ha 73aab3f83e Merge #342: teardown resolves the pane's real tab instead of trusting the delegate's placement config
CI / contract (push) Successful in 53s
CI / build (push) Failing after 1m47s
2026-09-04 16:38:31 +07:00
Dai Ha 887aca0183 fleetd#335: abandon()'s per-task cleanup and sendAsync's terminal hook must not swallow throws
CI / contract (pull_request) Successful in 1m2s
CI / build (pull_request) Successful in 2m21s
Site 1 (abandon()'s matching loop, reachable): the recovery/put-back branch calls
inbox.publish, which AmqpReplyInbox implements as a real broker round trip that
throws IllegalStateException on an unroutable/unconfirmed/interrupted publish.
An uncaught throw there aborted the loop, stranding every task after it in
`matching` PENDING forever. Fixed by recording each task's own future.complete()
result before any cleanup runs, then wrapping the cleanup in try/catch so one
task's failure cannot stop its siblings from getting their outcome. Reaching the
throwing branch by real timing needs a race the file's own #137 follow-up already
found unreachable through the public API, so the reproducing test uses a
test-only hook (same technique as the existing fleetd #324/#329 hooks) to inject
the throw at that exact point.

Site 2 (sendAsync's task.future.whenComplete, reachable): the returned stage is
discarded, so an uncaught throw from pushLoop.onTicketTerminal vanished with no
log line. Reproduced for real: Fleetd's shutdown hook runs messages.close()
(stops the async executor from taking new work, but does not cancel a send
already in flight) before pushLoop.close() (shuts its scheduler down
immediately) — a ticket completing in that window makes onTicketTerminal's own
scheduler.schedule(...) throw a genuine RejectedExecutionException. Fixed with a
try/catch(Throwable) plus log.error inside the whenComplete action.

Site 3 (the two `finally { asyncTasksByWaiter.remove(reply); rendezvous.close(...);
}` blocks in send() and answer()): read Rendezvous.close/closeAsk and the
ConcurrentHashMap operations behind them — both are plain map ops on a non-null
key with no user-overridable code, so neither can throw. Left unchanged; not a
defect.

Mutation-proven: reverting either fix reproduces the failure it exists to catch
— removing site 1's try/catch aborts abandon() with the injected exception
(MessageServiceTest#aPerTaskCleanupFailureDoesNotStrandTheRemainingMatchingTasks
errors); removing site 2's try/catch leaves the RejectedExecutionException
unlogged (MessageServiceTest#aTicketTerminalPushFailureDoesNotVanishSilently
fails its log assertion). Full suite: mvn clean install, Tests run: 1365,
Failures: 0, Errors: 0, BUILD SUCCESS.
2026-09-04 16:38:14 +07:00
Dai Ha 3fd23ecafa fleetd #342: base tab-cleanup teardown on the pane's real placement, not a delegate's static config
CI / contract (pull_request) Successful in 1m22s
CI / build (pull_request) Failing after 1m38s
HerdrPeerLauncher.stop() used to gate spaces.locatePane() on usesTabPlacement(),
which reads the delegate's OWN configured profiles. When CompositePeerLauncher's
single-daemon stop() shortcut hands a pane to a delegate that never spawned it
(spawnedBy empty after a daemon restart, herdrDaemonCount()==1), that delegate's
placement config says nothing true about how the pane was actually placed, and a
dedicated tab could be skipped and leaked.

Resolve the tab unconditionally instead — WorkspaceControl#locatePane already
tolerates a missing pane by returning null — and let the existing single-occupant
check (tabPaneCount()==1) be the only thing that decides whether to close it, same
as it already protects a shared tab regardless of declared placement.

Adds a mixed-placement CompositePeerLauncherTest (every existing stop-fallback test
configured both adapters as tab placement, so the mis-routing never showed) and
updates FleetAppTest#stopWorkerInPanePlacementClosesOnlyThePane, whose old
assertion (no pane.get on pane placement) documented exactly the skip this fix
removes.
2026-09-04 16:36:11 +07:00
Dai Ha f379847942 Merge #345: the timeout path's use of Cancellation.DELIVERED is now pinned
CI / contract (push) Successful in 47s
CI / build (push) Successful in 2m8s
2026-09-04 16:32:45 +07:00
Dai Ha ea12107497 fleetd #345: test timeout cancellation race
CI / contract (pull_request) Successful in 1m28s
CI / build (pull_request) Successful in 2m3s
2026-09-04 16:30:46 +07:00
Dai Ha 591df91de1 Merge #337: the deferred key set proves its own reporting coverage too
CI / contract (push) Successful in 54s
CI / build (push) Failing after 1m57s
2026-09-04 16:16:20 +07:00
Dai Ha 6a814176f0 #339: the start-of-line check must skip terminal chrome
CI / build (push) Failing after 1m30s
CI / contract (push) Successful in 1m54s
#339 stopped a member's own prose about an error from recording a credential
outage, by requiring the pattern at the start of its matched line. A bare
lookingAt also rejected a genuine error line rendered as

    | 503 Service Unavailable: upstream credential rejected

The send still failed, but the outage was never recorded. That is the false
negative #339's own invariant 3 named as worse than the false positive it set
out to fix: an unrecorded outage leaves the fleet spawning into a dead
credential.

Measured with a throwaway probe on the raw-scrape path, whose own comment says
to expect leading chrome there: kind=FAILED, sinkNotified=0.

startsWithBackendError now skips a leading run of non-letter, non-digit
characters before the check. That keeps #339's intent: prose still does not
match, because there the pattern sits after words rather than after chrome.
The worker's own prose test still passes.

Mutation: restoring the bare lookingAt fails the new test.
2026-09-04 16:13:19 +07:00
Dai Ha d11d1d157c Merge #339: a backend-error text match must be at the start of its line before it records a credential outage 2026-09-04 16:08:46 +07:00
Dai Ha d057d56156 #341: pin the noise control, and the reverse policy order
The fix reported every distinct unprotected name, but nothing held it there.

Mutation: replacing .filter(unprotectedGapNamesWarned::add) with a filter that
adds and always returns true - so every name is logged on every spawn - left
all 1358 tests green. The Set behaved; nothing proved this class used it as a
guard rather than as a record.

Two tests added:

- theSameUnprotectedNameIsWarnedAboutOnlyOnceAcrossSpawns pins invariant 1, the
  noise control. It now fails on that mutation, showing both duplicate WARNs.
- anAllowListWarnDoesNotSuppressALaterDenyByDefaultWarnForADifferentName covers
  the reverse policy order. The defect was found going deny-by-default then
  allow-list; a guard fixed in one direction is not fixed in the other.
2026-09-04 16:07:50 +07:00
Dai Ha d703ce1313 #337: extend ConfigRefTopLevelReportingCoverageTest to DEFERRED_KEYS
CI / contract (pull_request) Successful in 1m23s
CI / build (pull_request) Successful in 1m27s
ConfigRefTopLevelReportingCoverageTest (added by #333) proved every COLD_KEYS
and SPLIT_KEYS member has a real comparison behind it, but left
DEFERRED_TOP_LEVEL_KEYS unexercised. Re-measured by mutation (drop each
key's branch from changedDeferredKeys, run the suite, restore): 6 of the 11
deferred keys had no behavioural test naming them — guard, leadHeartbeat,
worktreeRoot, spawnReadyTimeoutMs, spawnReadyPollMs, quarantineCooldownSeconds
— which corrects the issue's own guessed list in two ways: lifecycle is
actually covered (ConfigRefTest.aDeferredChangeIsAppliedAndReported), and
worktreeRoot was missing from the issue's list entirely.

Promoted the test-side DEFERRED_TOP_LEVEL_KEYS copy into ConfigRef.DEFERRED_KEYS
(package-private, alongside COLD_KEYS/SPLIT_KEYS) so the reflective test reads
the same set changedDeferredKeys is compared against, and made
changedDeferredKeys package-private so the test can call it directly. Every
DEFERRED_KEYS component turned out to be a scalar or a simple record, so no
exclusion set was needed.

Mutation proof: dropping guard's branch from changedDeferredKeys leaves the
whole suite green except the new
everyDeferredKeyIsActuallyReportedByChangedDeferredKeys test, which fails
naming guard exactly.
2026-09-04 16:06:09 +07:00
Dai Ha b32a30fd47 Merge #341: warn once per distinct unprotected credential name, not once per launcher 2026-09-04 16:02:51 +07:00
Dai Ha 464dbc0930 fleetd#341: a per-name guard so a later spawn's different unprotected name still warns
CI / contract (pull_request) Successful in 1m1s
CI / build (pull_request) Successful in 1m29s
unprotectedGapLogged was one AtomicBoolean guarding two WARN branches in
logCredentialGap that name different env var names (the allow-list
keptByDerivedList branch, and warnGapUnprotected's deny-by-default /
non-zsh-fallback branch). memberCredentials is a live, re-read-per-spawn
supplier, so between two spawns a policy reload can change which names are
in the gap: spawn 1 warns about name A and trips the shared flag, and
spawn 2's gap containing a different name B never gets its WARN.

Replace the AtomicBoolean with unprotectedGapNamesWarned, a
ConcurrentHashMap-backed Set<String> guard keyed per name (same shape as
OpenCodeLauncher.modelCheckSkippedWarned), so each distinct credential-shaped
name is warned about exactly once, ever, regardless of which branch or
which spawn first reports it. allowListGapLogged (the separate INFO guard,
#192) is untouched. Neither WARN's wording changed.
2026-09-04 15:59:43 +07:00
Dai Ha a8cadd9150 Merge #338: a timed-out queued send cancels its message instead of leaving it to be delivered later
CI / contract (push) Successful in 1m23s
CI / build (push) Successful in 2m12s
2026-09-04 15:57:39 +07:00
Dai Ha 57b8c0b56d fleetd #339: guard backend error sink
CI / contract (pull_request) Successful in 50s
CI / build (pull_request) Successful in 1m50s
2026-09-04 15:57:18 +07:00
13 changed files with 1248 additions and 164 deletions
@@ -154,10 +154,13 @@ import java.util.function.Supplier;
* {@link 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} member has any reporting code
* behind it at all — {@code ConfigRefTopLevelReportingCoverageTest} is what fleetd #333 added for
* that, after measuring that a {@code SPLIT_KEYS} entry with its reporting branch deleted passes
* both this file's own "kept in step" assert and {@code ConfigRefTopLevelCoverageTest} unchanged.
* triaged; it does NOT prove a {@code SPLIT_KEYS}/{@code COLD_KEYS}/{@code DEFERRED_KEYS} member has
* any reporting code behind it at all — {@code ConfigRefTopLevelReportingCoverageTest} is what
* fleetd #333 added for that, after measuring that a {@code SPLIT_KEYS} entry with its reporting
* branch deleted passes both this file's own "kept in step" assert and
* {@code ConfigRefTopLevelCoverageTest} unchanged. fleetd #337 extended it to {@code DEFERRED_KEYS}
* after measuring the same one-way gap there directly: dropping {@code guard}'s branch out of
* {@link #changedDeferredKeys} while {@code "guard"} stayed in the set left the whole suite green.
*
* <p><strong>A cold change refuses the whole reload.</strong> Not the hot half applied and the cold
* half warned about: that would leave the running daemon in a state matching no file on disk, which
@@ -193,6 +196,25 @@ public final class ConfigRef implements Supplier<FleetConfig> {
*/
static final Set<String> SPLIT_KEYS = Set.of("health", "coordinator", "fleet");
/**
* Top-level keys {@link #changedDeferredKeys} compares — see the class doc's Deferred bullet.
* Promoted here from a test-side copy in {@code ConfigRefTopLevelCoverageTest} by fleetd #337,
* the same reason {@link #COLD_KEYS} and {@link #SPLIT_KEYS} live here rather than in a test: a
* second, hand-maintained copy of this set is exactly the kind of thing that silently drifts
* from the method it is supposed to describe. {@code spawnReadyTimeoutMs} and
* {@code spawnReadyPollMs} are compared together in one branch and reported under the combined
* label {@code "spawnReady*"}; {@code profiles} is compared twice over (added/removed names,
* then an existing profile's launch settings) — see {@link #changedDeferredKeys}.
*
* <p>Package-private (not {@code private}) so {@code ConfigRefTopLevelCoverageTest} and
* {@code ConfigRefTopLevelReportingCoverageTest} can both read it, the same way they already
* read {@link #COLD_KEYS} and {@link #SPLIT_KEYS}.
*/
static final Set<String> DEFERRED_KEYS = Set.of(
"guard", "worktreeRoot", "worktreeGroup", "primary", "configReload",
"leadHeartbeat", "lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs",
"quarantineCooldownSeconds", "profiles");
private final Path path;
private final AtomicReference<FleetConfig> current;
@@ -350,8 +372,16 @@ public final class ConfigRef implements Supplier<FleetConfig> {
return changed;
}
/** Changed keys that were accepted but whose effect waits for a restart. */
private static List<String> changedDeferredKeys(FleetConfig old, FleetConfig fresh) {
/**
* Changed keys that were accepted but whose effect waits for a restart.
*
* <p>Package-private (not {@code private}) so {@code ConfigRefTopLevelReportingCoverageTest}
* can call it directly with a reflection-built {@code FleetConfig} pair, the same reason
* {@link #changedColdKeys} and {@link #changedSplitKeys} already are (fleetd #333, extended to
* this method by fleetd #337 — membership in {@link #DEFERRED_KEYS} proved nothing about this
* method on its own until then; see that test's class doc).
*/
static List<String> changedDeferredKeys(FleetConfig old, FleetConfig fresh) {
List<String> changed = new ArrayList<>();
if (!Objects.equals(old.lifecycle(), fresh.lifecycle())) {
changed.add("lifecycle");
@@ -1,10 +1,12 @@
package dev.ltms.fleet.inject;
/**
* Notified when {@link CompletionResolver} actually delivers a typed backend-error classification
* to a waiting send (fleetd#201 / #227) — never on a race that lost. {@link CompletionResolver}
* calls this only after {@code Rendezvous.resolveFailure} returns {@code true} for that exact
* waiter, mirroring the win-only race rule {@link ExhaustionSink} already uses.
* Notified when {@link CompletionResolver} has a backend-error match at the start of a pane line,
* or both a match and its too-fast crash signature, for a waiting send (fleetd#201 / #227). A text
* match inside ordinary pane prose can be a member's report about an error, so it fails the send
* without notifying this sink.
* {@link CompletionResolver} calls this only after {@code Rendezvous.resolveFailure} returns
* {@code true} for that exact waiter, mirroring the win-only race rule {@link ExhaustionSink} uses.
*
* <p>The public send result is unchanged by this classification — it is still a failed send
* ({@code Rendezvous.Kind#FAILED}); this sink is the internal seam a later stage (fleetd#201 Unit
@@ -380,16 +380,19 @@ public final class CompletionResolver implements TurnListener {
+ "matched the profile's exhausted pattern): {}", target, reason);
// CB-578 stage B: only on the resolution that actually won the race — a late
// duplicate must never quarantine a credential twice for one refusal.
exhaustionSink.onExhausted(target, reason);
if (startsWithExhaustion(matchedLine, exhausted)) {
exhaustionSink.onExhausted(target, reason);
}
}
return;
}
// fleetd#164 (part 2) / fleetd#201: a scrape that read cleanly and produced content still
// isn't a real reply when that content is the backend's own rejection (e.g. an HTTP 400
// before the worker did any work). Classify it as a failure naming the member, rather than
// handing the caller a scrape that reads like a completed answer, and — only on the
// resolution that actually wins the race, mirroring the exhaustion sink above — notify the
// typed backend-error sink so a later stage can act on repeated failures.
// handing the caller a scrape that reads like a completed answer. A text match alone is not
// enough to notify the typed backend-error sink: this assistant block can be a member's
// normal prose about an error. A line that starts with the error match is stronger evidence;
// the too-fast path below also has its crash signature before it records a credential failure.
String backendError = firstMatchingLine(assistantBlock, backendErrorPatternOrFallback(target));
if (backendError != null) {
// Carry the whole scrape, not just the matched line. The pattern is a heuristic: a member
@@ -401,9 +404,9 @@ public final class CompletionResolver implements TurnListener {
if (rendezvous.resolveFailure(waiter, reason)) {
inFlight.remove(target, turn);
log.warn("failing send to {} via turn-stall fallback: {}", target, reason);
// fleetd#201 Unit 1: only on the resolution that actually won the race — a late
// duplicate must never double-count one backend failure.
backendErrorSink.onBackendError(target, backendError, reason);
if (startsWithBackendError(backendError, backendErrorPatternOrFallback(target))) {
backendErrorSink.onBackendError(target, backendError, reason);
}
}
return;
}
@@ -477,7 +480,9 @@ public final class CompletionResolver implements TurnListener {
+ "usable assistant block; no fleet_reply): {}", target, reason);
// CB-578 stage B: only on the resolution that actually won the race — a late
// duplicate must never quarantine a credential twice for one refusal.
exhaustionSink.onExhausted(target, reason);
if (startsWithExhaustion(matchedLine, exhausted)) {
exhaustionSink.onExhausted(target, reason);
}
}
return true;
}
@@ -492,8 +497,9 @@ public final class CompletionResolver implements TurnListener {
if (rendezvous.resolveFailure(waiter, reason)) {
inFlight.remove(target, turn);
log.warn("failing send to {} via turn-stall fallback from the raw scrape: {}", target, reason);
// fleetd#201 Unit 1: only on the resolution that actually won the race.
backendErrorSink.onBackendError(target, backendError, reason);
if (startsWithBackendError(backendError, backendErrorPatternOrFallback(target))) {
backendErrorSink.onBackendError(target, backendError, reason);
}
}
return true;
}
@@ -540,10 +546,10 @@ public final class CompletionResolver implements TurnListener {
* {@link #MIN_TURN_NANOS} — a crash signature (e.g. a backend HTTP 400 before the worker did
* anything) that a bare {@code BUSY -> DONE} transition cannot be told apart from a genuinely
* fast completion. Runs the same backend-error classification the normal and raw-scrape paths
* apply, against whatever is on screen right now: a match is a typed failure that notifies
* {@link #backendErrorSink} (only on the resolution that wins the race); a non-match stays the
* original generic too-fast failure, naming the member and both timings, with whatever the pane
* shows appended so the caller sees the cause, not just "it failed".
* apply, against whatever is on screen right now: a match together with the too-fast crash
* signature notifies {@link #backendErrorSink} (only on the resolution that wins the race). A
* non-match stays the original generic too-fast failure, naming the member and both timings,
* with whatever the pane shows appended so the caller sees the cause, not just "it failed".
*/
private void failTooFast(String target, InFlight turn, CompletableFuture<Rendezvous.Resolution> waiter,
long elapsedNanos) {
@@ -611,6 +617,68 @@ public final class CompletionResolver implements TurnListener {
return null;
}
/**
* True when nothing before the match on this pane line ends a sentence — that is, the match is
* still inside the line's first sentence rather than inside prose a member wrote about it.
* Used to decide whether an exhaustion match may quarantine a credential (fleetd #348).
*
* <p><strong>Why this is looser than {@link #startsWithBackendError}.</strong> An
* {@code exhaustedPattern} is written per profile and may name only the decisive words of a
* provider message — {@code "usage limit has been reached"} without its leading {@code "The"}.
* A start-of-line check would then reject the genuine refusal. That is the false negative
* fleetd #348's invariant 1 calls the worse direction: an unrecorded exhaustion leaves the
* fleet spawning into a credential with no capacity, and a quarantine runs 1800s against the
* backend-error cooldown's fixed 60s.
*
* <p>This rule accepts a superset of what a start-of-line check accepts: if the match begins
* right after the chrome, there is nothing in front of it, so there is no sentence ending
* either. So moving to it cannot add a false negative.
*
* <p><strong>No chrome skipping here, deliberately.</strong> The first version of this method
* copied {@code startsWithBackendError}'s leading-chrome loop. Measured on merge: deleting that
* loop left all 1369 tests green, and it must — the scan only looks for {@code . ! ?}, and no
* terminal chrome character is one of those. A step that cannot change the result is worse than
* no step, because the next reader takes it as evidence that chrome was handled.
*
* <p>It stays a heuristic. Prose whose <em>first</em> sentence carries the pattern still
* notifies the sink, and a genuine refusal behind an earlier full stop (a hostname, a version
* number) still does not. Both are known and neither is fixed here.
*/
private static boolean startsWithExhaustion(String line, Pattern pattern) {
var matcher = pattern.matcher(line);
if (!matcher.find()) {
return false;
}
for (int prefix = 0; prefix < matcher.start(); prefix++) {
if (".!?".indexOf(line.charAt(prefix)) >= 0) {
return false;
}
}
return true;
}
/**
* True when the error pattern begins the matched pane line, rather than appearing in prose.
*
* <p>Leading terminal chrome is skipped first — box-drawing characters, bullets, gutter bars and
* spaces. #339 introduced this check with a bare {@code lookingAt}, and that rejected a genuine
* error line rendered as {@code "| 503 Service Unavailable: ..."}: the send still failed, but the
* credential outage was never recorded. That is the false negative #339's own invariant 3 called
* worse than the false positive it set out to fix — measured with a throwaway probe on the
* raw-scrape path, which is exactly the path whose comment says to expect leading chrome.
*
* <p>Skipping only a leading run of non-letter, non-digit characters keeps the fix's intent. A
* member's prose ({@code "I checked the retry path. An API Error: makes it back off."}) still
* does not match, because there the pattern sits after words, not after chrome.
*/
private static boolean startsWithBackendError(String line, Pattern pattern) {
int i = 0;
while (i < line.length() && !Character.isLetterOrDigit(line.charAt(i))) {
i++;
}
return pattern.matcher(line.substring(i)).lookingAt();
}
/**
* Coverage summary for the CB-578 stage A exhausted-pattern classification, logged at startup
* the way {@link dev.ltms.fleet.health.FleetHealthMonitor#coverage} is — so an operator can
@@ -937,14 +937,20 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
*/
@Override
public void stop(String idOrPane) {
// Teardown knows only the pane, not which profile spawned it. Attempt tab cleanup when any
// profile uses tab placement (so the bridge may have created a dedicated peer tab); the
// single-occupant check below is what actually protects the user's shared tabs.
// Teardown knows only the pane, not which profile spawned it — and, when this call is
// routed here through CompositePeerLauncher's single-daemon stop() shortcut (fleetd #342),
// not even which adapter's config actually governed the spawn: the shortcut can hand the
// pane to a delegate that never spawned it, whose own profiles say nothing about how THIS
// pane was placed. So the decision to look for a tab to clean up is made from the pane's
// actual state, not from this delegate's static profile config: resolve the tab
// unconditionally and let {@link WorkspaceControl#locatePane} tolerate "not found" (it
// returns null rather than throwing); the single-occupant check below is what actually
// protects the user's shared tabs, exactly as it always has.
String paneId = paneByAgentId.remove(idOrPane);
if (paneId == null) {
paneId = idOrPane; // raw-pane fallback (reap, gate timeout, pane-addressed callers)
}
WorkspaceControl.PaneLocation loc = usesTabPlacement() ? spaces.locatePane(paneId) : null;
WorkspaceControl.PaneLocation loc = spaces.locatePane(paneId);
try {
agents.close(paneId);
} catch (HerdrException e) {
@@ -985,11 +991,6 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
}
}
/** Whether any configured profile places peers in their own tab (so tabs may need cleanup). */
private boolean usesTabPlacement() {
return profiles.values().stream().anyMatch(FleetConfig.Profile::tabPlacement);
}
/** True when a herdr error means the target is already gone (safe to treat as done). */
private static boolean isAlreadyGone(HerdrException e) {
return e.code() != null && e.code().endsWith("_not_found");
@@ -1665,30 +1666,59 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
/**
* Guards the {@code effectiveAllowed == null} branch of {@link #logCredentialGap} — the
* genuinely-unprotected report (deny-by-default, and the allow-list non-zsh fallback) — to one
* WARN per launcher instance, not one per spawn.
* genuinely-unprotected report (deny-by-default, and the allow-list non-zsh fallback) — AND
* the {@code effectiveAllowed != null} / {@code keptByDerivedList} branch, the allow-list case
* where a name is in the gap but the derived allow-list keeps it anyway. Both branches log the
* same severity (WARN) about the same fact — a name genuinely reaching a member pane
* unprotected — so they share this one guard, keyed per NAME rather than per launcher instance:
* each credential-shaped name that is ever reported unprotected gets exactly one WARN, however
* many spawns see it and whichever of the two branches first reports it.
*
* <p>fleetd #341: {@code memberCredentials} is a live, re-read-per-spawn supplier, so the
* policy — and so the gap's actual member names — can change between two spawns on the same
* launcher. Before this fix the guard was a single {@code AtomicBoolean} tripped by either
* branch: spawn 1 could warn about name A and trip the flag, and a later spawn's gap containing
* a different name B would never be reported, even though B is just as unprotected as A was.
* {@code AtomicBoolean} could not express "once per distinct name" at all — only "once, ever,
* for whichever name got there first" — so this is a {@code Set<String>} guard instead, the
* same shape {@link OpenCodeLauncher#modelCheckSkippedWarned} already uses for its own
* once-per-distinct-thing WARN. {@link #add}'s return value (true only the first time a name is
* added) is what turns "log the whole gap" into "log only the names never warned about before".
*
* <p>Bounded by construction: every name added here first passed {@link
* #CREDENTIAL_SHAPED_NAME}'s filter over {@link #hostEnvNames}, i.e. it is an actual
* environment variable name from the daemon's own process — a small, OS-bounded set (the host
* environment has, in practice, tens to a few hundred entries), not an attacker- or
* request-controlled input. So this set cannot grow past "however many distinct credential-
* shaped names this host's environment has ever held across this launcher's lifetime," which is
* effectively fixed for the life of one daemon process — no separate cap is needed.
*
* <p>CB-633 follow-up (#192): kept SEPARATE from {@link #allowListGapLogged} on purpose.
* {@code memberCredentials} is a live, re-read-per-spawn supplier, so the policy can change
* between two spawns on the same launcher. A single shared flag would let a harmless allow-list
* INFO on spawn 1 permanently suppress the real deny-by-default WARN a later spawn deserves —
* the report that matters most getting hidden by the report that doesn't. Two flags mean each
* report kind fires exactly once, independent of what the other kind already logged.
* report kind (WARN vs. INFO) fires independently of what the other kind already logged; within
* the WARN kind itself, the set above further separates by name, for the same reason.
*/
private final AtomicBoolean unprotectedGapLogged = new AtomicBoolean();
private final Set<String> unprotectedGapNamesWarned = ConcurrentHashMap.newKeySet();
/**
* Guards the {@code effectiveAllowed != null} branch of {@link #logCredentialGap} — the
* allow-list-scrub-covered report — to one INFO per launcher instance. See {@link
* #unprotectedGapLogged}'s javadoc for why this is a separate flag rather than a shared one.
* #unprotectedGapNamesWarned}'s javadoc for why this is a separate flag rather than a shared
* one; unlike that guard it stays a per-instance {@code AtomicBoolean}, not a per-name set —
* fleetd #341 fixed the WARN-vs-WARN suppression, not this INFO's own one-shot shape, which was
* not reported as broken and is out of that ticket's scope.
*/
private final AtomicBoolean allowListGapLogged = new AtomicBoolean();
/**
* fleetd #185 stage 2: guards {@link #warnUnknownMemberEnvironment} to one WARN per launcher
* instance, not one per spawn — the same one-per-instance shape as {@link #unprotectedGapLogged}
* and {@link #allowListGapLogged}, kept as its own flag for the same reason those two are split:
* this mode is orthogonal to which of the other two branches would otherwise have fired.
* instance, not one per spawn — the same one-shot shape {@link #unprotectedGapNamesWarned} and
* {@link #allowListGapLogged} guard their own branches with, kept as its own flag for the same
* reason those two are split: this mode is orthogonal to which of the other two branches would
* otherwise have fired.
*/
private final AtomicBoolean unknownMemberEnvironmentWarned = new AtomicBoolean();
@@ -1816,14 +1846,21 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
List<String> blankedByScrub = gap.stream()
.filter(name -> !MemberEnvAllowList.keeps(effectiveAllowed, name))
.toList();
if (!keptByDerivedList.isEmpty() && unprotectedGapLogged.compareAndSet(false, true)) {
// fleetd #341: filter to names this guard has never warned about before — not just
// "isEmpty" on the whole branch — so a name this spawn's gap shares with an EARLIER
// spawn's (already-warned) gap does not re-print, while a name unique to THIS gap still
// does, whichever of the two WARN branches reported it first.
List<String> newlyUnprotected = keptByDerivedList.stream()
.filter(unprotectedGapNamesWarned::add)
.toList();
if (!newlyUnprotected.isEmpty()) {
log.warn("memberCredentials gap: {} credential-shaped env var name(s) are on neither "
+ "known: nor allow: — the derived allow-list keeps them anyway (a profile's "
+ "gitTokenEnv/gitHostEnv/tokenEnv/env: names one, or this spawn injects it), "
+ "so every member pane inherits them UNBLOCKED — {}. Add each to "
+ "memberCredentials.known (or .allow if a member legitimately needs it), or "
+ "remove it from whatever profile setting derives it in.",
keptByDerivedList.size(), keptByDerivedList);
newlyUnprotected.size(), newlyUnprotected);
}
if (!blankedByScrub.isEmpty() && allowListGapLogged.compareAndSet(false, true)) {
log.info("memberCredentials gap: {} credential-shaped env var name(s) are on neither "
@@ -1836,12 +1873,18 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
/** The deny-by-default (and allow-list non-zsh fallback) WARN — unchanged byte-for-byte by #192. */
private void warnGapUnprotected(List<String> gap) {
if (unprotectedGapLogged.compareAndSet(false, true)) {
// fleetd #341: same "only the names never warned before" filter as the sibling branch in
// logCredentialGap above — see unprotectedGapNamesWarned's javadoc. Both branches share
// this one guard because both report the exact same fact (a name reaching a member pane
// unprotected) at the exact same severity; keying it by name is what lets a later spawn's
// DIFFERENT name still get its own WARN after an earlier spawn's already fired.
List<String> newlyUnprotected = gap.stream().filter(unprotectedGapNamesWarned::add).toList();
if (!newlyUnprotected.isEmpty()) {
log.warn("memberCredentials gap: {} credential-shaped env var name(s) are on neither "
+ "known: nor allow: — every member pane inherits them UNBLOCKED — {}. "
+ "Add each to memberCredentials.known (blocked by default) or .allow "
+ "(if a member legitimately needs it).",
gap.size(), gap);
newlyUnprotected.size(), newlyUnprotected);
}
}
@@ -709,27 +709,47 @@ public final class MessageService {
boolean isRecovery = task == recoveryTask && recovered != null;
Reply outcome = isRecovery ? recovered : new Reply(Outcome.WORKER_FAILED, reason);
String turnId = task.turnId;
if (task.future.complete(outcome)) {
if (outcome.outcome() == Outcome.WORKER_FAILED) {
asyncFailed = true;
// fleetd #335: completing THIS task's future must not depend on any other task's
// cleanup succeeding — every task in `matching` is owed its own outcome regardless of
// what happens below, so decide and record that before doing anything that can throw.
boolean completedHere = task.future.complete(outcome);
if (completedHere && outcome.outcome() == Outcome.WORKER_FAILED) {
asyncFailed = true;
}
try {
if (abandonCleanupHookForTest != null) {
// Test-only (fleetd #335 site 1): see the field's own javadoc.
abandonCleanupHookForTest.run();
}
if (turnId != null) {
// #275: whether this task was swept out of ASKING or was already answered and
// only waiting on its resumed turn's real reply (#137), nothing will ever
// complete this turnId now — drop it from this class's own bookkeeping AND the
// reverse-rendezvous itself, so hasAsyncQuestion(target) stops reporting a turn
// that is actually done, and a late answer() sees it as lapsed rather than
// resolving a question nothing is listening for any more.
asyncTasksByTurn.remove(turnId, task);
rendezvous.closeAsk(turnId);
if (completedHere) {
if (turnId != null) {
// #275: whether this task was swept out of ASKING or was already answered and
// only waiting on its resumed turn's real reply (#137), nothing will ever
// complete this turnId now — drop it from this class's own bookkeeping AND the
// reverse-rendezvous itself, so hasAsyncQuestion(target) stops reporting a turn
// that is actually done, and a late answer() sees it as lapsed rather than
// resolving a question nothing is listening for any more.
asyncTasksByTurn.remove(turnId, task);
rendezvous.closeAsk(turnId);
}
} else if (isRecovery) {
// The recovered reply was already drained out of the inbox, but this task resolved
// through another path (e.g. a concurrent reply() or a second abandon() racing this
// one) between us choosing it and completing it here. Put the reply back rather than
// lose it silently — it may still belong to some other still-open task, or the next
// caller that drains this target's inbox.
inbox.publish(target, UUID.randomUUID().toString(), recovered.text());
}
} else if (isRecovery) {
// The recovered reply was already drained out of the inbox, but this task resolved
// through another path (e.g. a concurrent reply() or a second abandon() racing this
// one) between us choosing it and completing it here. Put the reply back rather than
// lose it silently — it may still belong to some other still-open task, or the next
// caller that drains this target's inbox.
inbox.publish(target, UUID.randomUUID().toString(), recovered.text());
} catch (RuntimeException e) {
// fleetd #335: inbox.publish reaches a broker (AmqpReplyInbox throws
// IllegalStateException on an unroutable/unconfirmed/interrupted publish) and this
// loop has no other teardown path — a caller on the release path, or the health
// monitor's GONE/NEVER_READY sweep. Losing this exception uncaught would abort the
// loop and leave every task still to come in `matching` PENDING forever (fleetd
// #335 site 1). completedHere is already recorded above, so only this task's
// best-effort bookkeeping is lost — log it and let the loop reach the rest.
log.error("abandon: per-task cleanup failed for ticket {} (target {}, turnId {})",
task.ticket, target, turnId, e);
}
}
if (failed) {
@@ -871,6 +891,10 @@ public final class MessageService {
boolean wasDelivered = delivery.completion().isDone()
&& !delivery.completion().isCompletedExceptionally();
if (!wasDelivered) {
if (timeoutCancellationRaceHookForTest != null) {
// Test-only (fleetd #345): see the field's own javadoc.
timeoutCancellationRaceHookForTest.run();
}
// The target monitor makes cancellation atomic with onStatus picking this
// Pending up. If pickup won, report TIMED_OUT_WORKING because the text landed.
wasDelivered = injector.cancel(delivery) == Injector.Cancellation.DELIVERED;
@@ -941,13 +965,41 @@ public final class MessageService {
return new AskResult(AskOutcome.ANSWERED, answer);
} catch (TimeoutException e) {
log.debug("fleet_ask from {} went unanswered in {}ms", workerSession, timeoutMillis);
// fleetd #307: mark the task BEFORE clearAsyncQuestion(forgetTurn=true) below drops it out of
// asyncTasksByTurn and nulls its turnId — that forgetting is deliberate and stays (it is
// what keeps the target from staying BUSY forever), but it would otherwise also erase
// askAnsweredAsyncTasks' only signal that the worker's eventual real fleet_reply still
// belongs to this task, stranding it in the inbox with a false "never replied" verdict.
markAskTimedOut(ticket.turnId());
clearAsyncQuestion(ticket.turnId(), true);
// Only the fresh owner tears down the shared turn (mirrors the finally block below).
// A duplicate's own timeoutMillis says nothing about whether the SHARED ask is actually
// done — it must leave the close/forget bookkeeping to the fresh owner, exactly as it
// already leaves closeAsk to it.
if (ticket.fresh()) {
// fleetd #334: close the ask turn BEFORE forgetting this task's turnId mapping below.
// Before this fix the order was reversed — the mapping was forgotten here first, and
// rendezvous.closeAsk only ran afterward, in the shared finally. A primary's answer()
// call racing this exact timeout could then find rendezvous.askSession(turnId) still
// non-null (the ask still "answerable") after the Task mapping was already gone:
// answer()'s own asyncTasksByTurn lookup returned null, its task != null guard skipped
// the completion, and the async ticket sat at PENDING forever even though answer()
// itself reported the worker's real reply. Closing here first removes that window:
// any answer() call that still observes askSession(turnId) != null is necessarily
// racing a point BEFORE the forgetting below runs (both happen on this one thread, in
// this order, with nothing that yields in between), so the Task mapping is still there
// for it to find; any call that observes askSession(turnId) == null now correctly
// bails out STALE_TURN (see answer()'s own top check) before ever reaching
// asyncTasksByTurn. rendezvous.closeAsk is idempotent — a no-op once the turn is
// already removed, see its own javadoc — so the shared finally below re-running it
// for this same fresh call is harmless.
rendezvous.closeAsk(ticket.turnId());
if (askTimeoutRaceHookForTest != null) {
// Test-only (fleetd #334): see the field's own javadoc.
askTimeoutRaceHookForTest.run();
}
// fleetd #307: mark the task BEFORE clearAsyncQuestion(forgetTurn=true) below drops it
// out of asyncTasksByTurn and nulls its turnId — that forgetting is deliberate and
// stays (it is what keeps the target from staying BUSY forever), but it would
// otherwise also erase askAnsweredAsyncTasks' only signal that the worker's eventual
// real fleet_reply still belongs to this task, stranding it in the inbox with a false
// "never replied" verdict.
markAskTimedOut(ticket.turnId());
clearAsyncQuestion(ticket.turnId(), true);
}
return new AskResult(AskOutcome.TIMED_OUT, null);
} catch (ExecutionException e) {
Throwable cause = e.getCause();
@@ -1046,19 +1098,28 @@ public final class MessageService {
// the chained ask deliberately left open.
//
// A null task is NOT only "this was never an async ticket". That reading was in this
// comment when #329 merged and it is wrong. A genuine async ticket also lands here
// with task == null, because ask()'s timeout path runs clearAsyncQuestion(turnId,
// true) — which drops the asyncTasksByTurn entry — in its catch block, while
// rendezvous.closeAsk(turnId) runs later, in its finally. Between those two the ask
// is still answerable but the map entry is already gone, so the lookup at :991
// returns null and this ticket is never completed. Measured on 2026-09-04: a probe
// firing only that first half before answer() runs printed
// "answer=REPLIED phase=PENDING reply=null" — the same stranded ticket #329 set out
// to fix, one step earlier in the same race. The probe used forgetTurnForTest, which
// omits ask()'s markAskTimedOut; that cannot change the outcome, because askTimedOut
// is read only by askAnsweredAsyncTasks, and reply() never reaches it while this
// method's own waiter is live. So #329 narrows this window rather than closing it.
// Open as fleetd #334 — do not read this guard as complete.
// comment when #329 merged and it is wrong; it is still not the whole story after
// #334. A genuine async ticket can still land here with task == null — a blocking
// (wait:true) send's fleet_ask never has a Task at all, so that case is expected and
// fine. What #334 fixed was a SECOND, unintended way to get here with task == null:
// ask()'s timeout path used to run clearAsyncQuestion(turnId, true) — which drops the
// asyncTasksByTurn entry — in its catch block, while rendezvous.closeAsk(turnId) ran
// later, in its finally. Between those two the ask was still answerable but the map
// entry was already gone, so the lookup at :1053 returned null and this ticket was
// never completed. Measured on 2026-09-04: a probe firing only that first half before
// answer() ran printed "answer=REPLIED phase=PENDING reply=null" — the same stranded
// ticket #329 set out to fix, one step earlier in the same race; the probe used
// forgetTurnForTest, which omits ask()'s markAskTimedOut, and that omission does not
// change the outcome, because askTimedOut is read only by askAnsweredAsyncTasks, and
// reply() never reaches it while this method's own waiter is live. #334's fix
// reorders ask()'s timeout catch to run closeAsk before the forgetting (see the
// fresh-owner block there), which removes this path entirely rather than narrowing
// it further: once closeAsk has run, rendezvous.askSession(turnId) is null and
// answer() returns STALE_TURN from its own top check, before it ever reaches this
// lookup — see aLateAnswerDuringAskTimeoutTeardownStillCompletesTheAsyncTicket in
// MessageServiceTest, which pins the exact window with askTimeoutRaceHookForTest.
// So by the time this line runs, task == null means only the ordinary blocking-send
// case (or #329's own already-fixed race elsewhere) — not this one.
if (result.outcome() != Outcome.QUESTION && task != null) {
finishAsyncTask(task, result);
}
@@ -1118,7 +1179,18 @@ public final class MessageService {
// class javadoc on sendAsync/CB-107.
task.future.whenComplete((reply, ex) -> {
boolean failed = ex != null || reply == null || !reply.completed();
pushLoop.onTicketTerminal(ticket, target, failed);
try {
pushLoop.onTicketTerminal(ticket, target, failed);
} catch (Throwable t) {
// fleetd #335 (site 2): this stage's own CompletableFuture is discarded, so an
// uncaught throw here (e.g. a RejectedExecutionException from
// ReplyPushLoop's scheduler, already shut down while this in-flight send's
// whenComplete fires during the daemon's own shutdown sequence — messages.close()
// only stops accepting NEW async work, it does not cancel a delivery already
// running) vanishes with no log line and no metric, and the push loop never learns
// the ticket went terminal — the exact thing this hook exists to tell it.
log.error("push loop failed to learn ticket {} (target {}) went terminal", ticket, target, t);
}
});
}
asyncExecutor.submit(() -> {
@@ -1370,6 +1442,24 @@ public final class MessageService {
this.afterFinishAsyncTaskCompleteHookForTest = hook;
}
/**
* Null in production; test seam for fleetd #345 — invoked in {@link #send}'s timeout path after
* {@link Injector.Delivery#completion()} reports incomplete and before {@link Injector#cancel}
* takes the target monitor. A test installs this to make {@code onStatus} pick the exact queued
* delivery up in that window, so {@code cancel} returns {@link Injector.Cancellation#DELIVERED}.
* This deterministically covers the caller's need to use that result rather than relying on a
* timing-sensitive real race.
*/
private volatile Runnable timeoutCancellationRaceHookForTest;
/**
* Test-only (fleetd #345): install {@link #timeoutCancellationRaceHookForTest}. Package-private
* so the test, in the same package, can reach it without widening any production API.
*/
void setTimeoutCancellationRaceHookForTest(Runnable hook) {
this.timeoutCancellationRaceHookForTest = hook;
}
/**
* Null in production; test seam for fleetd #329 (F3) — invoked from {@link #reply} right after
* the single local read of {@code orphan.turnId} passes its null-check and before that (now-local)
@@ -1399,6 +1489,56 @@ public final class MessageService {
clearAsyncQuestion(turnId, true);
}
/**
* Null in production; test seam for fleetd #335 (site 1) — invoked from {@link #abandon(String,
* String, boolean)}'s per-task loop, once per task, right before that task's own cleanup
* (turnId bookkeeping, or the stranded-reply put-back) runs. A test installs this to inject a
* throw at that exact point deterministically.
*
* <p>The one production call there that can really throw is {@code inbox.publish} in the
* put-back branch — {@link AmqpReplyInbox#publish} reaches a broker and throws {@link
* IllegalStateException} on an unroutable, unconfirmed, or interrupted publish — but reaching
* that branch requires a second completion of the very task {@code abandon} is about to
* complete to win the race first (see the branch's own comment), and the #137 follow-up
* investigation above already found the combination this needs (a stranded reply coinciding
* with an open matching task) unreachable through the public API, not merely hard to time.
* This hook reproduces the resulting shape — a per-task cleanup throw — directly, the same
* technique {@link #finishAsyncTaskRaceHook} and {@link
* #afterFinishAsyncTaskCompleteHookForTest} already use for their own hard-to-time races.
*/
private volatile Runnable abandonCleanupHookForTest;
/**
* Test-only (fleetd #335, site 1): install {@link #abandonCleanupHookForTest}. Package-private
* so the test, in the same package, can reach it without widening any production API.
*/
void setAbandonCleanupHookForTest(Runnable hook) {
this.abandonCleanupHookForTest = hook;
}
/**
* Null in production; test seam for fleetd #334 — invoked from {@link #ask}'s {@code
* TimeoutException} catch, only for the fresh owner, right after {@code rendezvous.closeAsk}
* has run and before {@link #markAskTimedOut} / {@link #clearAsyncQuestion} forget this task's
* turnId mapping. A test installs this to call {@link #answer} for the very same {@code turnId}
* synchronously from inside that exact window, deterministically reproducing the race a real
* concurrent {@code answer()} call could otherwise only win by timing luck: with the ask already
* closed, that call must see {@code rendezvous.askSession(turnId) == null} and return {@link
* Outcome#STALE_TURN} immediately, never reaching {@code asyncTasksByTurn} at all — proving the
* window fleetd #334 describes (mapping forgotten while the ask was still "answerable") is
* closed, rather than merely narrowed the way fleetd #329 narrowed the sibling race in {@link
* #finishAsyncTask}.
*/
private volatile Runnable askTimeoutRaceHookForTest;
/**
* Test-only (fleetd #334): install {@link #askTimeoutRaceHookForTest}. Package-private so the
* test, in the same package, can reach it without widening any production API.
*/
void setAskTimeoutRaceHookForTest(Runnable hook) {
this.askTimeoutRaceHookForTest = hook;
}
/** A new send must not open a waiter while an async ticket owns this worker's paused turn. */
private boolean hasAsyncQuestion(String target) {
return asyncTasksByTurn.values().stream().anyMatch(task -> target.equals(task.target));
@@ -47,21 +47,21 @@ class ConfigRefTopLevelCoverageTest {
private static final RecordComponent[] COMPONENTS = FleetConfig.class.getRecordComponents();
/**
* Top-level components whose change {@code ConfigRef.changedDeferredKeys} reads and reports on
* — verified by reading that method as of fleetd #330, not derived from this test.
* {@code spawnReadyTimeoutMs}/{@code spawnReadyPollMs} are compared together and reported under
* one combined label ({@code "spawnReady*"}); {@code profiles} is compared twice over — once for
* added/removed profile names, once for an existing profile's launch settings — and that second
* comparison excludes {@code weight}/{@code maxLoad}/{@code credentialId} as hot sub-fields,
* which is what {@link ConfigRefProfileCoverageTest} exists to keep honest at the sub-field
* level. {@code profiles} itself still belongs here, not in the hot-exclusion set below: most of
* a profile's fields are NOT read live, so citing "read live off the config supplier" for the
* whole top-level key would be false.
* Top-level components whose change {@code ConfigRef.changedDeferredKeys} reads and reports on.
* Fleetd #337 promoted this out of a hand-maintained copy here into {@link ConfigRef#DEFERRED_KEYS}
* itself, the same reason {@link ConfigRef#COLD_KEYS} and {@link ConfigRef#SPLIT_KEYS} are
* production constants rather than test-side copies: two lists that are supposed to describe the
* same method are exactly the shape that silently drifts apart. {@code spawnReadyTimeoutMs}/
* {@code spawnReadyPollMs} are compared together and reported under one combined label
* ({@code "spawnReady*"}); {@code profiles} is compared twice over — once for added/removed
* profile names, once for an existing profile's launch settings — and that second comparison
* excludes {@code weight}/{@code maxLoad}/{@code credentialId} as hot sub-fields, which is what
* {@link ConfigRefProfileCoverageTest} exists to keep honest at the sub-field level. {@code
* profiles} itself still belongs here, not in the hot-exclusion set below: most of a profile's
* fields are NOT read live, so citing "read live off the config supplier" for the whole
* top-level key would be false.
*/
private static final Set<String> DEFERRED_TOP_LEVEL_KEYS = Set.of(
"guard", "worktreeRoot", "worktreeGroup", "primary", "configReload",
"leadHeartbeat", "lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs",
"quarantineCooldownSeconds", "profiles");
private static final Set<String> DEFERRED_TOP_LEVEL_KEYS = ConfigRef.DEFERRED_KEYS;
/**
* The escape hatch: top-level components with no reload bookkeeping at all, because every read
@@ -16,61 +16,68 @@ import static org.junit.jupiter.api.Assertions.assertEquals;
/**
* fleetd #333, finding F2: {@link ConfigRefTopLevelCoverageTest} proves every {@link FleetConfig}
* top-level component sits in exactly one of {@link ConfigRef#COLD_KEYS}, {@code
* DEFERRED_TOP_LEVEL_KEYS}, {@link ConfigRef#SPLIT_KEYS} or the hot-excluded set. It does
* <strong>not</strong> prove that a key's membership in {@code COLD_KEYS} or {@code SPLIT_KEYS}
* corresponds to any actual comparison in {@link ConfigRef}: a key can sit in either set with no
* branch in {@code changedColdKeys}/{@code changedSplitKeys} checking it, and both
* {@link ConfigRefTopLevelCoverageTest} and the "kept in step" {@code assert} inside each of those
* methods stay green, because neither one reads the method body — the coverage test only reads
* set membership, and the assert only checks that reported entries are a SUBSET of the set, never
* that every set member produced a reported entry.
* top-level component sits in exactly one of {@link ConfigRef#COLD_KEYS}, {@link
* ConfigRef#DEFERRED_KEYS}, {@link ConfigRef#SPLIT_KEYS} or the hot-excluded set. It does
* <strong>not</strong> prove that a key's membership in one of the first three sets corresponds to
* any actual comparison in {@link ConfigRef}: a key can sit in a set with no branch in {@code
* changedColdKeys}/{@code changedSplitKeys}/{@code changedDeferredKeys} checking it, and both
* {@link ConfigRefTopLevelCoverageTest} and the "kept in step" {@code assert} inside the first two
* of those methods stay green, because neither one reads the method body — the coverage test only
* reads set membership, and the assert only checks that reported entries are a SUBSET of the set,
* never that every set member produced a reported entry. {@code changedDeferredKeys} does not even
* have a "kept in step" assert of its own.
*
* <p>Measured directly, live, while fixing fleetd #333: dropping the {@code coordinator} branch out
* of {@code ConfigRef.changedSplitKeys} while leaving {@code "coordinator"} in
* {@link ConfigRef#SPLIT_KEYS} left {@link ConfigRefTopLevelCoverageTest} and the in-method assert
* both green — only a hand-written behavioural case in {@link ConfigRefTest} caught it, because it
* happened to name that exact key. This is the {@link ConfigRefProfileCoverageTest} mechanism one
* level up, generalised over every {@code COLD_KEYS}/{@code SPLIT_KEYS} member rather than one
* hand-picked field: enumerate {@link FleetConfig}'s own record components by reflection, build "a
* config where only {@code <key>} differs" for each cold/split key, and call the real
* {@link ConfigRef#changedColdKeys}/{@link ConfigRef#changedSplitKeys} methods (made
* package-private for exactly this, the same reason {@link ConfigRef#sameLaunchSettings} already
* is) to prove each one is actually reported — not assumed from a set literal.
* level up, generalised over every {@code COLD_KEYS}/{@code SPLIT_KEYS}/{@code DEFERRED_KEYS}
* member rather than one hand-picked field: enumerate {@link FleetConfig}'s own record components by
* reflection, build "a config where only {@code <key>} differs" for each key, and call the real
* {@link ConfigRef#changedColdKeys}/{@link ConfigRef#changedSplitKeys}/
* {@link ConfigRef#changedDeferredKeys} methods (all package-private for exactly this, the same
* reason {@link ConfigRef#sameLaunchSettings} already is) to prove each one is actually reported —
* not assumed from a set literal.
*
* <h2>What this deliberately does NOT cover</h2>
* {@code DEFERRED_TOP_LEVEL_KEYS} is not exercised here. That bucket carries the identical
* one-way risk in principle — a key added to it with no matching branch in
* {@code changedDeferredKeys} would pass {@link ConfigRefTopLevelCoverageTest} exactly the way
* {@code coordinator} passed it above — but this test stays narrow to {@code COLD_KEYS} and
* {@code SPLIT_KEYS} for two reasons. First, that is where fleetd #333 actually found and measured
* the gap (F1 was a live instance of it). Second, most of {@code DEFERRED_TOP_LEVEL_KEYS} already
* carries an individual behavioural test in {@link ConfigRefTest} naming it by key — {@code
* worktreeGroup}, {@code primary}, {@code configReload}, {@code profiles}' launch settings /
* weight-maxLoad / exhaustedPattern / errorPattern / ideProjectDir — which is the same protection
* this class gives {@code COLD_KEYS}/{@code SPLIT_KEYS}, just written by hand per key instead of
* generated by reflection over the whole set. {@code guard}, {@code lifecycle},
* {@code leadHeartbeat}, {@code spawnReadyTimeoutMs}/{@code spawnReadyPollMs} and
* {@code quarantineCooldownSeconds} do NOT have a dedicated behavioural test naming them, so the
* one-way gap fleetd's own memory notes ("pre-existing on COLD_KEYS and on the test's own
* DEFERRED_TOP_LEVEL_KEYS") is real and not fully closed by this class — extending this mechanism's
* {@code BASE}/{@code ALT} map to cover every top-level component and adding a third
* {@code everyDeferredKeyIsActuallyReportedByChangedDeferredKeys} test is the natural next step, left
* for whoever next finds a deferred key with the same shape as this ticket's {@code fleet.leaders}.
* <h2>fleetd #337 — DEFERRED_KEYS was the gap left open here</h2>
* This class originally covered {@code COLD_KEYS} and {@code SPLIT_KEYS} only — {@code
* DEFERRED_TOP_LEVEL_KEYS} (now {@link ConfigRef#DEFERRED_KEYS}) carried the identical one-way risk
* in principle, unexercised. fleetd #337 measured the real consequence rather than assuming it from
* the shape of the gap: dropping {@code guard}'s comparison out of {@code changedDeferredKeys} while
* {@code "guard"} stayed in the set left all 1355 tests green — the same failure mode {@code
* coordinator} demonstrated for {@code SPLIT_KEYS} in fleetd #333, now confirmed for {@code
* DEFERRED_KEYS} too. Re-deriving the full list by mutation (drop each key's branch in turn, run the
* suite, restore) found six of the eleven {@code DEFERRED_KEYS} members with no behavioural test in
* {@link ConfigRefTest} naming them: {@code guard}, {@code leadHeartbeat}, {@code worktreeRoot},
* {@code spawnReadyTimeoutMs}, {@code spawnReadyPollMs} and {@code quarantineCooldownSeconds}. That
* list corrects fleetd #333's own guess at it in two ways the mutation proved and a reading did not:
* {@code lifecycle} is NOT on it — {@code ConfigRefTest.aDeferredChangeIsAppliedAndReported} already
* names it, and dropping its branch fails that test — and {@code worktreeRoot} IS on it, which #333
* never named at all. The other five {@code DEFERRED_KEYS} members ({@code lifecycle}, {@code
* worktreeGroup}, {@code primary}, {@code configReload}, {@code profiles}) already had a hand-written
* case each. {@link #everyDeferredKeyIsActuallyReportedByChangedDeferredKeys} below now covers all
* eleven the reflective way, so the six with no hand-written test are no longer silently unpinned —
* every {@code DEFERRED_KEYS} component turned out to be a scalar or a simple record, so, unlike
* {@code fleet.leaders} in fleetd #333, none needed an exclusion set: {@link #BASE}/{@link #ALT} give
* every top-level component (not only {@code COLD_KEYS}/{@code SPLIT_KEYS}) a real, distinct value.
*/
class ConfigRefTopLevelReportingCoverageTest {
private static final RecordComponent[] COMPONENTS = FleetConfig.class.getRecordComponents();
/**
* One valid value per top-level {@link FleetConfig} component — "the a value". Components not
* exercised by either test below ({@code profiles}, {@code guard}, {@code worktreeRoot}, …) are
* left {@code null}/empty; {@link FleetConfig}'s compact constructor only normalizes
* {@code profiles}, so every other field accepts {@code null} unmutated.
* One valid value per top-level {@link FleetConfig} component — "the a value". fleetd #337 gave
* every {@code DEFERRED_KEYS} component a real value here too (previously left {@code null} on
* both sides, which meant {@code mutate(key)} produced no actual difference for any of them);
* only {@code placement}, {@code memberCredentials} and {@code memberLoginShell} — the
* hot-excluded set, never compared by any {@code changed*Keys} method — stay {@code null}.
* {@link FleetConfig}'s compact constructor only normalizes {@code profiles}, so every other
* field accepts whatever is put here unmutated.
*/
private static final Map<String, Object> BASE = baseValues();
/** The same shape, each value distinct from {@link #BASE} — "the b value" — for COLD_KEYS/SPLIT_KEYS only. */
/** The same shape, each value distinct from {@link #BASE} — "the b value". */
private static final Map<String, Object> ALT = altValues();
private static Map<String, Object> baseValues() {
@@ -79,26 +86,26 @@ class ConfigRefTopLevelReportingCoverageTest {
v.put("herdrSocket", "~/.config/herdr/a.sock");
v.put("memberHerdrSocket", "~/.config/herdr/member-a.sock");
v.put("profiles", Map.of());
v.put("guard", null);
v.put("worktreeRoot", null);
v.put("lifecycle", null);
v.put("spawnReadyTimeoutMs", null);
v.put("spawnReadyPollMs", null);
v.put("guard", new FleetConfig.Guard(List.of("host-a")));
v.put("worktreeRoot", "/wt/a");
v.put("lifecycle", new FleetConfig.Lifecycle(300, 5, 30, false));
v.put("spawnReadyTimeoutMs", 5000);
v.put("spawnReadyPollMs", 100);
v.put("broker", new FleetConfig.Broker("amqp://a", null, 1));
v.put("primary", null);
v.put("primary", new FleetConfig.Primary("term-a", 1, 1000));
v.put("fleet", new FleetConfig.Fleet(
Map.of("opus", new FleetConfig.Leader("sonnet", "lead: opus-a", 1, null, 10,
"claude", null, null, null)),
Map.of(), Map.of(), Map.of(), Map.of(), "{role}: {profile} #{n}"));
v.put("leadHeartbeat", null);
v.put("leadHeartbeat", new FleetConfig.LeadHeartbeat(300, 60_000L, 3));
v.put("health", new FleetConfig.Health(true, 30, 600, null, null));
v.put("placement", null);
v.put("auth", new FleetConfig.Auth("loopback-trust", null));
v.put("configReload", null);
v.put("quarantineCooldownSeconds", null);
v.put("configReload", new FleetConfig.ConfigReload(true, 10));
v.put("quarantineCooldownSeconds", 1800);
v.put("memberCredentials", null);
v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-a", null, "self-a", 1));
v.put("worktreeGroup", null);
v.put("worktreeGroup", "group-a");
v.put("memberLoginShell", null);
assertNamesMatchComponents(v);
return v;
@@ -109,14 +116,18 @@ class ConfigRefTopLevelReportingCoverageTest {
v.put("bind", new FleetConfig.Bind("127.0.0.2", 8766));
v.put("herdrSocket", "~/.config/herdr/b.sock");
v.put("memberHerdrSocket", "~/.config/herdr/member-b.sock");
v.put("profiles", Map.of());
v.put("guard", null);
v.put("worktreeRoot", null);
v.put("lifecycle", null);
v.put("spawnReadyTimeoutMs", null);
v.put("spawnReadyPollMs", null);
// A single added profile — enough to trip the "added/removed" comparison in
// ConfigRef.changedDeferredKeys, which is all this mechanism needs to prove "profiles" has
// a branch behind it; the launch-settings comparison already has its own hand-written cases
// in ConfigRefTest (changingAProfilesLaunchSettingsIsReportedAsDeferred and siblings).
v.put("profiles", Map.of("sonnet", minimalProfile("sonnet")));
v.put("guard", new FleetConfig.Guard(List.of("host-b")));
v.put("worktreeRoot", "/wt/b");
v.put("lifecycle", new FleetConfig.Lifecycle(600, 10, 60, true));
v.put("spawnReadyTimeoutMs", 10_000);
v.put("spawnReadyPollMs", 200);
v.put("broker", new FleetConfig.Broker("amqp://b", null, 2));
v.put("primary", null);
v.put("primary", new FleetConfig.Primary("term-b", 2, 2000));
// Differs from BASE.fleet only in fleet.leaders (a different tab for "opus") — the frozen
// sub-field ConfigRef.changedSplitKeys actually compares. A Fleet that instead differed only
// in tabLabel would correctly NOT be reported (see
@@ -126,20 +137,27 @@ class ConfigRefTopLevelReportingCoverageTest {
Map.of("opus", new FleetConfig.Leader("sonnet", "lead: opus-b", 1, null, 10,
"claude", null, null, null)),
Map.of(), Map.of(), Map.of(), Map.of(), "{role}: {profile} #{n}"));
v.put("leadHeartbeat", null);
v.put("leadHeartbeat", new FleetConfig.LeadHeartbeat(600, 120_000L, 5));
v.put("health", new FleetConfig.Health(false, 90, 900, null, null));
v.put("placement", null);
v.put("auth", new FleetConfig.Auth("token", "TOKEN_ENV"));
v.put("configReload", null);
v.put("quarantineCooldownSeconds", null);
v.put("configReload", new FleetConfig.ConfigReload(false, 20));
v.put("quarantineCooldownSeconds", 3600);
v.put("memberCredentials", null);
v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-b", null, "self-b", 2));
v.put("worktreeGroup", null);
v.put("worktreeGroup", "group-b");
v.put("memberLoginShell", null);
assertNamesMatchComponents(v);
return v;
}
/** A minimal, otherwise-null {@link FleetConfig.Profile} — just enough to name one in a map. */
private static FleetConfig.Profile minimalProfile(String name) {
return new FleetConfig.Profile(name, null, null, null, null, null, null, null, null, null,
null, null, null, null, null, null, null, null, null, null, null, null, null, null,
null, null);
}
private static void assertNamesMatchComponents(Map<String, Object> values) {
Set<String> names = new TreeSet<>();
for (RecordComponent rc : COMPONENTS) {
@@ -216,4 +234,51 @@ class ConfigRefTopLevelReportingCoverageTest {
+ "entry from changedSplitKeys — a set entry with no comparison behind it, "
+ "exactly the fleetd #333 F2 shape: " + uncovered);
}
/**
* For most {@link ConfigRef#DEFERRED_KEYS} members, {@code changedDeferredKeys} reports the key
* name verbatim — the default this map assumes. Two entries don't: {@code spawnReadyTimeoutMs}
* and {@code spawnReadyPollMs} are compared together in one branch and reported under the
* combined label {@code "spawnReady*"} (see {@link ConfigRef#changedDeferredKeys}). {@code
* profiles} keeps the default: mutating it here only exercises the added/removed comparison
* (see {@link #altValues}), which reports {@code "profiles (added/removed: …)"} — starts with
* {@code "profiles"}, same as the default would expect.
*/
private static final Map<String, String> DEFERRED_REPORT_PREFIX = Map.of(
"spawnReadyTimeoutMs", "spawnReady*",
"spawnReadyPollMs", "spawnReady*");
/**
* fleetd #337: the same mechanism applied to {@link ConfigRef#DEFERRED_KEYS}, closing the gap
* this class's own javadoc left open since fleetd #333. Mutate each deferred key in isolation
* and prove {@code changedDeferredKeys} actually names it (message starting with the key's
* expected report prefix — see {@link #DEFERRED_REPORT_PREFIX}), not just that {@code
* DEFERRED_KEYS} claims it does. This is the exact check that fails for {@code guard} the way
* {@code coordinator} failed {@link #everySplitKeyIsActuallyReportedByChangedSplitKeys} in
* fleetd #333 — verified live: dropping {@code guard}'s branch from {@code changedDeferredKeys}
* while {@code "guard"} stayed in {@code DEFERRED_KEYS} left the whole 1355-test suite green,
* and this test is what now catches it (it fails naming {@code guard} with that mutation in
* place).
*/
@Test
void everyDeferredKeyIsActuallyReportedByChangedDeferredKeys() throws ReflectiveOperationException {
FleetConfig base = configOf(BASE);
List<String> uncovered = new ArrayList<>();
for (String key : new TreeSet<>(ConfigRef.DEFERRED_KEYS)) {
FleetConfig mutated = mutate(key);
List<String> deferred = ConfigRef.changedDeferredKeys(base, mutated);
String prefix = DEFERRED_REPORT_PREFIX.getOrDefault(key, key);
if (deferred.stream().noneMatch(s -> s.startsWith(prefix))) {
uncovered.add(key);
}
}
System.out.printf(
"ConfigRef.changedDeferredKeys reporting coverage — %d DEFERRED_KEYS, %d verified%n",
ConfigRef.DEFERRED_KEYS.size(), ConfigRef.DEFERRED_KEYS.size() - uncovered.size());
assertEquals(List.of(), uncovered,
"these keys are in ConfigRef.DEFERRED_KEYS but mutating them alone produces no "
+ "matching entry from changedDeferredKeys — a set entry with no comparison "
+ "behind it, exactly the fleetd #333 F2 shape, confirmed here for "
+ "DEFERRED_KEYS by fleetd #337: " + uncovered);
}
}
@@ -0,0 +1,192 @@
package dev.ltms.fleet.config;
import org.junit.jupiter.api.Test;
import java.lang.reflect.Constructor;
import java.lang.reflect.RecordComponent;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Locale;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
import java.util.TreeSet;
import static org.junit.jupiter.api.Assertions.assertEquals;
/**
* Guards against a "defect factory" built into this file's own established pattern, found live
* while building the (parked) idle-sleep-guard PR: every time a component is added to
* {@link FleetConfig}, the record grows by one arg AND a new back-compat constructor is added at
* the OLD arity, so existing callers keep compiling. That is correct and required — see the
* constructor ladder just below the record header. But {@link #withDefaults()}'s own {@code return
* new FleetConfig(...)} call sits in this same file, written at a literal argument count. The very
* next time a component is added, the freshly-added back-compat constructor at the OLD arity
* silently captures that stale call, because it is now a legal overload at that arg count too. It
* compiles. Every other test passes, because nothing else exercises the new field. The new
* component is defaulted away — {@code null}, or whatever that back-compat overload defaults it to
* — on every {@link FleetConfig#load}. Measured, not theoretical: this exact sequence happened
* live when the {@code idleSleepGuard} component was added on a sibling branch; it was caught only
* because that branch's own new tests happened to assert on the new field's value.
*
* <p>This test proves the opposite property, and does it in a way that survives the next field
* being added without being rewritten: reflectively enumerate {@link FleetConfig}'s own record
* components (never a hardcoded count — the arity is exactly what changes over time), build one
* config through the true canonical constructor with a real, distinctive, non-null value in EVERY
* component (reusing the exact reflective-construction pattern
* {@link ConfigRefTopLevelReportingCoverageTest} already established for this file:
* {@code getDeclaredConstructor(exact record-component types)}, which resolves the canonical
* constructor by its true shape, not by binding to whichever overload happens to match arg count —
* the same way Jackson resolves it), call the real {@link FleetConfig#withDefaults()}, and assert
* every one of those values survives unchanged.
*
* <p>Why this is a valid check for every component, not just some: {@link #withDefaults()}'s own
* comments document that it only ever REPLACES a component when the incoming value is {@code null}
* (or blank, for {@code placement}) — {@code broker}/{@code primary}/{@code leadHeartbeat}/
* {@code configReload}/{@code coordinator}/{@code worktreeGroup}/{@code memberLoginShell} are left
* as-is unconditionally, and {@code bind}/{@code guard}/{@code lifecycle}/{@code auth}/
* {@code fleet}/{@code quarantineCooldownSeconds}/{@code memberCredentials}/{@code placement} are
* replaced only on null/blank input. A value that is never null or blank going in must therefore
* never change coming out, for every current component. No exclusion is needed today.
*
* <p>{@link #EXCLUDED_FROM_SURVIVAL_CHECK} exists anyway, kept deliberately empty and size-pinned
* by {@link #exclusionListSizeIsPinned()}: a future component that {@code withDefaults()} is
* <em>documented</em> to transform unconditionally (unlike every field today) would legitimately
* need one. Pinning the size at 0 means growing that set to make a failure go away is itself a
* visible diff to this test, not a silent one — a checker that can be silenced by adding to its
* own escape hatch is not a checker.
*/
class FleetConfigWithDefaultsPreservesEveryComponentTest {
private static final RecordComponent[] COMPONENTS = FleetConfig.class.getRecordComponents();
/** See the class javadoc — deliberately empty today; grow it only with a matching justification. */
private static final Set<String> EXCLUDED_FROM_SURVIVAL_CHECK = Set.of();
/** One real, distinctive, non-null (non-blank where blankness would mean "unset") value per component. */
private static Map<String, Object> baseValues() {
Map<String, Object> v = new LinkedHashMap<>();
v.put("bind", new FleetConfig.Bind("127.0.0.1", 8765));
v.put("herdrSocket", "~/.config/herdr/guard.sock");
v.put("memberHerdrSocket", "~/.config/herdr/member-guard.sock");
v.put("profiles", Map.of("sonnet", minimalProfile("sonnet")));
v.put("guard", new FleetConfig.Guard(List.of("host-guard")));
v.put("worktreeRoot", "/wt/guard");
v.put("lifecycle", new FleetConfig.Lifecycle(300, 5, 30, true));
v.put("spawnReadyTimeoutMs", 12_345);
v.put("spawnReadyPollMs", 234);
v.put("broker", new FleetConfig.Broker("amqp://guard", null, 7));
v.put("primary", new FleetConfig.Primary("term-guard", 4, 4000));
v.put("fleet", new FleetConfig.Fleet(
Map.of("opus", new FleetConfig.Leader("sonnet", "lead: opus-guard", 1, null, 10,
"claude", null, null, null)),
Map.of(), Map.of(), Map.of(), Map.of(), "{role}: {profile} #{n}"));
v.put("leadHeartbeat", new FleetConfig.LeadHeartbeat(301, 61_000L, 4));
v.put("health", new FleetConfig.Health(true, 31, 601, 61, null));
v.put("placement", "round-robin");
v.put("auth", new FleetConfig.Auth("loopback-trust", null));
v.put("configReload", new FleetConfig.ConfigReload(true, 11));
v.put("quarantineCooldownSeconds", 1801);
v.put("memberCredentials", new FleetConfig.MemberCredentials(
FleetConfig.MemberCredentials.POLICY_DENY_BY_DEFAULT,
List.of("git"), List.of("git", "ssh"), null));
v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-guard", null, "self-guard", 3));
v.put("worktreeGroup", "group-guard");
v.put("memberLoginShell", "/bin/zsh");
assertNamesMatchComponents(v);
return v;
}
/** A minimal, otherwise-null {@link FleetConfig.Profile} — just enough to name one in a map. */
private static FleetConfig.Profile minimalProfile(String name) {
return new FleetConfig.Profile(name, null, null, null, null, null, null, null, null, null,
null, null, null, null, null, null, null, null, null, null, null, null, null, null,
null, null);
}
/**
* Guards {@link #baseValues()} itself against drifting from the record's real shape — the same
* assurance {@link ConfigRefTopLevelReportingCoverageTest} already relies on. This is what makes
* "no hardcoded arity" true in practice: forgetting to add a new component here fails this
* assertion by name, rather than silently checking one component fewer than the record has.
*/
private static void assertNamesMatchComponents(Map<String, Object> values) {
Set<String> names = new TreeSet<>();
for (RecordComponent rc : COMPONENTS) {
names.add(rc.getName());
}
assertEquals(names, new TreeSet<>(values.keySet()),
"this test's value map has drifted from FleetConfig's actual top-level components — "
+ "update baseValues() alongside the record");
}
/**
* Builds a {@link FleetConfig} through the TRUE canonical constructor — resolved by the record's
* own component types, not by argument count — so this never accidentally exercises a
* back-compat overload the way a literal {@code new FleetConfig(...)} call risks doing.
*/
private static FleetConfig configOf(Map<String, Object> values) throws ReflectiveOperationException {
Class<?>[] types = Arrays.stream(COMPONENTS).map(RecordComponent::getType).toArray(Class<?>[]::new);
Object[] args = Arrays.stream(COMPONENTS).map(rc -> values.get(rc.getName())).toArray();
Constructor<FleetConfig> ctor = FleetConfig.class.getDeclaredConstructor(types);
return ctor.newInstance(args);
}
@Test
void exclusionListSizeIsPinned() {
assertEquals(0, EXCLUDED_FROM_SURVIVAL_CHECK.size(),
"EXCLUDED_FROM_SURVIVAL_CHECK grew from 0 — every entry needs a justification in "
+ "this test class's javadoc AND this assertion re-pinned to the new size; a "
+ "growing exclusion list that silences failures on its own is not a guard");
}
/**
* The mutation this is built to catch: make {@code withDefaults()}'s final constructor call
* literal at some arg count, add one more component to the record with a new back-compat
* constructor at the old arity, and the stale call silently rebinds. Every component here is
* real and non-null (non-blank for the one String — {@code placement} — where blank has
* meaning), so none of it should be replaced by {@code withDefaults()}; any component that
* comes back different was silently dropped.
*/
@Test
void everyComponentGivenARealValueSurvivesWithDefaults() throws ReflectiveOperationException {
Map<String, Object> base = baseValues();
FleetConfig config = configOf(base);
FleetConfig defaulted = config.withDefaults();
List<String> dropped = new ArrayList<>();
int checked = 0;
for (RecordComponent rc : COMPONENTS) {
String name = rc.getName();
if (EXCLUDED_FROM_SURVIVAL_CHECK.contains(name)) {
continue;
}
checked++;
Object expected = base.get(name);
Object actual;
try {
actual = rc.getAccessor().invoke(defaulted);
} catch (ReflectiveOperationException e) {
throw new RuntimeException("failed to read FleetConfig." + name + "()", e);
}
if (!Objects.equals(expected, actual)) {
dropped.add(String.format(Locale.ROOT,
"%s: withDefaults() was given a real, non-null value (%s) for '%s' but "
+ "returned %s — a component silently dropped by withDefaults(), the "
+ "shape of the defect this test exists to catch (its final "
+ "\"return new FleetConfig(...)\" call binding to a back-compat "
+ "constructor instead of the true canonical one)",
name, expected, name, actual));
}
}
System.out.printf(Locale.ROOT,
"FleetConfig.withDefaults() component-survival coverage — %d components, %d checked, "
+ "%d excluded, %d survived%n",
COMPONENTS.length, checked, EXCLUDED_FROM_SURVIVAL_CHECK.size(), checked - dropped.size());
assertEquals(List.of(), dropped,
"withDefaults() silently dropped these real, given components: " + dropped);
}
}
@@ -484,6 +484,27 @@ class CompletionResolverTest {
// --- CB-578 stage A: backend-exhausted classification ---------------------------------
@Test
void aNormalMemberReportMentioningTheExhaustionPatternDoesNotNotifyTheSink() {
String block = "⏺ I reviewed capacity handling. The usage limit has been reached means no more work can start.\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached");
java.util.List<String> notified = new java.util.ArrayList<>();
ExhaustionSink sink = (target, reason, profile) -> notified.add(target + ": " + reason);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, sink);
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
assertEquals(Rendezvous.Kind.BACKEND_EXHAUSTED, waiter.getNow(null).kind(),
"a matching report still fails the send as exhausted");
assertTrue(waiter.getNow(null).text().contains("I reviewed capacity handling."),
"the exhausted result keeps the whole matched pane line");
assertTrue(notified.isEmpty(),
"a normal report mentioning an exhaustion pattern must not quarantine a credential");
}
@Test
void classifiesAMatchingScrapeAsBackendExhaustedInsteadOfACompletedReply() {
String block = "⏺ Working on it...\nThe usage limit has been reached. Try again later.\n❯ ";
@@ -534,6 +555,53 @@ class CompletionResolverTest {
"the sink is told the matched reason: " + notified.get(0));
}
@Test
void aRealExhaustionBehindTerminalChromeStillNotifiesTheSink() {
String block = "⏺ │ The usage limit has been reached. Try again later.\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached");
java.util.List<String> notified = new java.util.ArrayList<>();
ExhaustionSink sink = (target, reason, profile) -> notified.add(target + ": " + reason);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, sink);
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
assertEquals(Rendezvous.Kind.BACKEND_EXHAUSTED, waiter.getNow(null).kind(),
"a real exhaustion must still fail the send as exhausted");
assertEquals(1, notified.size(),
"a real exhaustion behind terminal chrome must reach the sink");
}
/**
* The live fleet configures {@code exhaustedPattern: "The usage limit has been reached"} — with
* the leading {@code "The"}. Every other test here uses a pattern without it, which is the shape
* that made fleetd #348 need a looser rule than a start-of-line check. This pins the deployed
* shape as well, so a later tightening of {@link CompletionResolver} cannot silently stop
* recording the exhaustion this fleet actually reports.
*
* <p>What it does not prove: that this is the only pattern shape an operator will write.
*/
@Test
void anExhaustionPatternCarryingItsLeadingWordsStillNotifiesTheSink() {
String block = "⏺ │ The usage limit has been reached. Try again later.\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
ExhaustedPatternLookup patterns = target -> Pattern.compile("The usage limit has been reached");
java.util.List<String> notified = new java.util.ArrayList<>();
ExhaustionSink sink = (target, reason, profile) -> notified.add(target + ": " + reason);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, sink);
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
assertEquals(Rendezvous.Kind.BACKEND_EXHAUSTED, waiter.getNow(null).kind(),
"the live pattern shape must still fail the send as exhausted");
assertEquals(1, notified.size(),
"the live pattern shape must still reach the sink");
}
@Test
void aLosingBackendExhaustedClassificationNeverNotifiesTheExhaustionSink() {
// The waiter was already resolved (e.g. by the worker's own reply) before this scrape landed —
@@ -804,7 +872,28 @@ class CompletionResolverTest {
// --- fleetd#201 Unit 1: target-keyed backend-error pattern + typed sink ----------------------
@Test
void aConfiguredBackendErrorPatternClassifiesAMatchAsAFailureAndNotifiesTheSinkOnce() {
void aNormalMemberReportMentioningTheFallbackErrorPatternFailsButDoesNotNotifyTheSink() {
String block = "⏺ I checked the retry path. An API Error: makes it back off.\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
java.util.List<String> notified = new java.util.ArrayList<>();
BackendErrorSink sink = (target, matchedLine, reason) -> notified.add(target + ": " + matchedLine);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none(), BackendErrorPatternLookup.legacy(), sink);
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(),
"a scrape mentioning the pattern must still fail the send");
assertTrue(waiter.getNow(null).text().contains("I checked the retry path. An API Error: makes it back off."),
"the failure must keep the whole pane tail");
assertTrue(notified.isEmpty(),
"a normal report mentioning the fallback pattern must not record a credential failure");
}
@Test
void aConfiguredBackendErrorPatternAtTheStartOfALineClassifiesAMatchAndNotifiesTheSinkOnce() {
String block = "⏺ 503 Service Unavailable: upstream credential rejected\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
@@ -924,7 +1013,7 @@ class CompletionResolverTest {
}
@Test
void aConfiguredPatternAlsoClassifiesTheRawScrapeFallbackAndNotifiesTheSink() {
void aConfiguredPatternAtTheStartOfALineAlsoClassifiesTheRawScrapeFallbackAndNotifiesTheSink() {
// No ⏺ marker and leading TUI chrome ⇒ lastAssistantBlock() yields "", so classification must
// fall back to the raw scrape (fleetd#211) — and it must use the configured pattern too.
String block = """
@@ -949,6 +1038,40 @@ class CompletionResolverTest {
assertTrue(notified.get(0).contains("503 Service Unavailable"), notified.get(0));
}
/**
* fleetd #339 follow-up: a genuine backend error rendered behind terminal chrome must still
* record the credential outage. #339 added a start-of-line check to stop a member's own prose
* being counted as an outage, and a bare {@code lookingAt} also rejected this — the send failed
* but the sink never fired. #339's invariant 3 named that direction as the worse one: a real
* outage going unrecorded leaves the fleet spawning into a dead credential.
*
* <p>The raw-scrape path is where this matters, because its own comment says to expect leading
* TUI chrome there.
*/
@Test
void aRealErrorBehindTerminalChromeStillNotifiesTheSink() {
String block = """
╭──────────────────────────────────────╮
│ 503 Service Unavailable: upstream credential rejected
""";
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
BackendErrorPatternLookup patterns = target -> Pattern.compile("(?i)503 Service Unavailable");
java.util.List<String> notified = new java.util.ArrayList<>();
BackendErrorSink sink = (target, matchedLine, reason) -> notified.add(target + ": " + matchedLine);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none(), patterns, sink);
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(),
"a real backend error must still fail the send");
assertEquals(1, notified.size(),
"a real error line behind box chrome is still a real outage — it must reach the sink, "
+ "or the fleet keeps spawning into a dead credential");
}
// --- fleetd#201 Unit 1: classification inside the fleetd#164 MIN_TURN_NANOS floor -------------
@Test
@@ -299,6 +299,41 @@ class CompositePeerLauncherTest {
"two adapter kinds sharing one daemon keep the fallback route");
}
@Test
void stopThroughTheSingleDaemonShortcutStillClosesTheTabWhenTheFallbackDelegateUsesPanePlacement() {
// fleetd #342: the single-daemon shortcut (spawnedBy empty, herdrDaemonCount()==1) always
// routes stop() through delegates.getFirst() — here the claude adapter, configured for
// PANE placement (its own profiles never create a dedicated tab). The pane being torn
// down here actually belongs to the opencode adapter's TAB placement, sharing the same
// herdr daemon — mixing placements is the point: every existing stop-fallback test in this
// class configures BOTH adapters as "tab", so usesTabPlacement() was true either way and
// the mis-routing never showed.
//
// Before the fix, HerdrPeerLauncher#stop gated tab resolution on usesTabPlacement() of the
// delegate it happened to be called through, so the wrongly-routed (pane-placement) claude
// adapter never even looked for a tab to close, and the now-empty tab leaked with nothing
// to reap it. The fix (fleetd #342) resolves the pane's real tab unconditionally, so the
// decision follows the pane's actual placement rather than the fallback delegate's static
// config.
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile claudePane = new FleetConfig.Profile("claude", "http://gx00.gw:8000",
"coder", null, "FLEETD_WORKER_TOKEN", List.of("claude"), "pane", "fleetd-workers",
"w #{n}", null, null, null);
ClaudeCodeLauncher claude = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of("claude", claudePane), "claude", _ -> null);
PeerLauncher composite = new CompositePeerLauncher(List.of(claude, opencodeAdapter(herdr)), "claude");
// "w9:pW" was never spawned through this composite instance, so spawnedBy has no entry for
// it (the same in-memory-cache-miss shape a daemon restart leaves behind) — stop() falls
// through to the single-daemon shortcut and hands it to delegates.getFirst() (claude).
composite.stop("w9:pW");
assertTrue(herdr.called("pane.close"), "the pane itself is still closed on every routed path");
assertTrue(herdr.called("tab.close"),
"the pane's real (sole-occupant) tab must be closed even though the fallback routed "
+ "through a delegate configured for pane placement");
}
@Test
void listKeepsBothPanesWhenTwoDaemonsShareAPaneId() {
// herdr pane ids are per-daemon counters, so two daemons really can both hold w1:p1 on
@@ -23,6 +23,7 @@ import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.Function;
import java.util.function.Supplier;
@@ -370,6 +371,161 @@ class HerdrPeerLauncherAllowListWiringTest {
"expected the pre-existing 'scrub blanks them' INFO unchanged, got: " + messages);
}
/**
* fleetd #341: {@code unprotectedGapLogged} guarded TWO WARN branches that name DIFFERENT env
* var names — the allow-list branch ({@code keptByDerivedList}, below) and the deny-by-default
* / non-zsh-fallback branch ({@link HerdrPeerLauncher#warnGapUnprotected}). {@code
* memberCredentials} is a live, re-read-per-spawn supplier, so the policy can change between
* two spawns on the same launcher instance — a config reload needs no restart. Spawn 1 runs
* under {@code deny-by-default} with a gap of {@code SPAWN_ONE_UNCOVERED_TOKEN}, which trips
* the (before this fix) SHARED one-shot flag. The policy is then reloaded to {@code
* allow-list}; spawn 2's gap is {@code FLEETD_WORKER_TOKEN} instead — the test profile's own
* {@code tokenEnv}, which the derived allow-list keeps even though it is on neither {@code
* known:} nor {@code allow:}, so it is genuinely unprotected and deserves its own WARN. Before
* this fix that WARN never fires, because the shared flag was already {@code true} — the
* operator is never told {@code FLEETD_WORKER_TOKEN} reaches every member pane unblocked. Real
* path: two real {@link HerdrPeerLauncher#spawn} calls on ONE launcher instance, with mutable
* {@code memberCredentials}/host-env suppliers standing in for a live config reload between
* spawns.
*/
@Test
void aDifferentUnprotectedGapOnALaterSpawnIsNotSuppressedByAnEarlierSpawnsWarn() {
FakeHerdr herdr = new FakeHerdr();
AtomicReference<FleetConfig.MemberCredentials> credsState = new AtomicReference<>(
new FleetConfig.MemberCredentials(null, List.of(), List.of(), null)); // deny-by-default
AtomicReference<Set<String>> hostEnvState =
new AtomicReference<>(Set.of("SPAWN_ONE_UNCOVERED_TOKEN"));
WiringLauncher launcher = new WiringLauncher(herdr, credsState::get, "/bin/zsh", hostEnvState::get);
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
Level original = logger.getLevel();
logger.setLevel(Level.WARN);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
try {
// Spawn 1: deny-by-default, gap = {SPAWN_ONE_UNCOVERED_TOKEN} — the effectiveAllowed ==
// null branch, via warnGapUnprotected.
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
// Live policy reload to allow-list, with a DIFFERENT gap name.
credsState.set(new FleetConfig.MemberCredentials(
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST, List.of(), List.of(), null));
hostEnvState.set(Set.of("FLEETD_WORKER_TOKEN"));
// Spawn 2: allow-list, gap = {FLEETD_WORKER_TOKEN} — kept by the derived allow-list
// (the profile's own tokenEnv), so it is the effectiveAllowed != null / keptByDerivedList
// branch, at the SAME log line HerdrPeerLauncher:1819 guards with the shared flag.
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
} finally {
logger.detachAppender(appender);
logger.setLevel(original);
}
List<String> messages = appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList();
assertTrue(messages.stream().anyMatch(
m -> m.contains("UNBLOCKED") && m.contains("SPAWN_ONE_UNCOVERED_TOKEN")),
"spawn 1's deny-by-default gap must still warn — got: " + messages);
assertTrue(messages.stream().anyMatch(
m -> m.contains("UNBLOCKED") && m.contains("FLEETD_WORKER_TOKEN")),
"spawn 2's gap names a DIFFERENT env var than spawn 1 (FLEETD_WORKER_TOKEN, not "
+ "SPAWN_ONE_UNCOVERED_TOKEN) — it must still be warned about even though a "
+ "flag already fired once for spawn 1's unrelated name. Before fleetd #341's "
+ "fix this WARN never fires because unprotectedGapLogged was already true. "
+ "Got: " + messages);
}
/**
* fleetd #341 follow-up: the OTHER half of the guard's contract. The set exists to report every
* distinct name, but it must still report each one only ONCE — the noise control is the reason
* a guard is here at all, and the ticket named it as invariant 1. Two spawns, same policy, same
* gap name: exactly one WARN mentioning it.
*
* <p>Measured before this test existed: replacing {@code .filter(unprotectedGapNamesWarned::add)}
* with a filter that adds and always returns {@code true} — so every name is logged on every
* spawn — left all 1358 tests green. The fix was correct and nothing held it there. That is the
* "a test on the seam does not prove the caller" shape: the {@code Set} behaves, and nothing
* proved this class used it as a guard rather than as a record.
*/
@Test
void theSameUnprotectedNameIsWarnedAboutOnlyOnceAcrossSpawns() {
FakeHerdr herdr = new FakeHerdr();
AtomicReference<FleetConfig.MemberCredentials> credsState = new AtomicReference<>(
new FleetConfig.MemberCredentials(null, List.of(), List.of(), null)); // deny-by-default
AtomicReference<Set<String>> hostEnvState =
new AtomicReference<>(Set.of("REPEATED_UNCOVERED_TOKEN"));
WiringLauncher launcher = new WiringLauncher(herdr, credsState::get, "/bin/zsh", hostEnvState::get);
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
Level original = logger.getLevel();
logger.setLevel(Level.WARN);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
try {
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
// Same policy, same gap, second spawn. Nothing new to tell the operator.
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
} finally {
logger.detachAppender(appender);
logger.setLevel(original);
}
List<String> messages = appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList();
long mentioning = messages.stream()
.filter(m -> m.contains("REPEATED_UNCOVERED_TOKEN"))
.count();
assertEquals(1, mentioning,
"one unchanged unprotected name across two spawns must produce exactly one WARN — "
+ "the set is a guard, not just a record. Got: " + messages);
}
/**
* fleetd #341 follow-up: the reverse policy order. The original defect was found going
* deny-by-default then allow-list, and a guard that is fixed in one direction is not
* necessarily fixed in the other — "ask which states still OPEN the gate". Here spawn 1 runs
* under {@code allow-list} (the {@code keptByDerivedList} branch) and spawn 2 under
* {@code deny-by-default} ({@link HerdrPeerLauncher#warnGapUnprotected}), with a different name
* each time. Both must be reported.
*/
@Test
void anAllowListWarnDoesNotSuppressALaterDenyByDefaultWarnForADifferentName() {
FakeHerdr herdr = new FakeHerdr();
AtomicReference<FleetConfig.MemberCredentials> credsState = new AtomicReference<>(
new FleetConfig.MemberCredentials(
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST, List.of(), List.of(), null));
AtomicReference<Set<String>> hostEnvState =
new AtomicReference<>(Set.of("FLEETD_WORKER_TOKEN"));
WiringLauncher launcher = new WiringLauncher(herdr, credsState::get, "/bin/zsh", hostEnvState::get);
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
Level original = logger.getLevel();
logger.setLevel(Level.WARN);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
try {
// Spawn 1: allow-list, gap kept by the derived list — the keptByDerivedList WARN.
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
// Live reload the OTHER way: back to deny-by-default, with a different name.
credsState.set(new FleetConfig.MemberCredentials(null, List.of(), List.of(), null));
hostEnvState.set(Set.of("LATER_UNCOVERED_TOKEN"));
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
} finally {
logger.detachAppender(appender);
logger.setLevel(original);
}
List<String> messages = appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList();
assertTrue(messages.stream().anyMatch(
m -> m.contains("UNBLOCKED") && m.contains("FLEETD_WORKER_TOKEN")),
"spawn 1's allow-list gap must warn — got: " + messages);
assertTrue(messages.stream().anyMatch(
m -> m.contains("UNBLOCKED") && m.contains("LATER_UNCOVERED_TOKEN")),
"spawn 2's deny-by-default gap names a different variable and must still be warned "
+ "about, even though an allow-list WARN already fired. Got: " + messages);
}
/**
* fleetd #185 stage 2: with {@code memberHerdrSocket:} configured, member panes run under a
* different OS user — {@link HerdrPeerLauncher#hostEnvNames} describes fleetd's own process, not
@@ -377,6 +377,126 @@ class MessageServiceTest {
assertEquals(MessageService.Outcome.QUESTION, q.outcome());
}
/**
* fleetd #334. {@code ask()}'s {@code TimeoutException} catch used to forget this task's
* {@code turnId} mapping ({@code clearAsyncQuestion(turnId, true)}) BEFORE closing the ask
* ({@code rendezvous.closeAsk}, in the shared {@code finally}). A primary's {@code answer()}
* call racing that exact window found the ask still "answerable" ({@code
* rendezvous.askSession(turnId)} still non-null) while the {@code Task} was already forgotten,
* so its {@code task != null} guard skipped the completion and the async ticket sat at
* {@code PENDING} forever even though {@code answer()} itself reported a result. The fix
* (closing the ask first) makes this window impossible: a racing {@code answer()} call either
* still finds the ask open (and the {@code Task} mapping guaranteed intact) or finds it already
* closed (and bails {@code STALE_TURN} before ever touching the {@code Task}). This test pins
* the exact window with {@code askTimeoutRaceHookForTest} and proves both invariants the ticket
* named: (1) a late/racing {@code answer()} sees the ask as already lapsed ({@code STALE_TURN}),
* never made answerable again, and (2) the async ticket still resolves {@code DONE} once the
* worker's real {@code fleet_reply} lands — it is never stranded {@code PENDING}.
*/
/**
* fleetd #334 gated the ask-timeout teardown on {@code ticket.fresh()}, matching the {@code
* finally} block that already did. This pins that gate. A coalesced duplicate passes its own
* {@code timeoutMillis}, which says nothing about whether the shared ask is done — so a
* duplicate timing out first must leave the fresh owner's still-open ask answerable.
*
* <p>Measured on merge: without this test, removing the {@code ticket.fresh()} gate left all
* 1371 tests green. The gate shipped with the reorder and nothing held it there.
*
* <p>What this does not prove: anything about the ordering inside the gate — that is
* {@code aLateAnswerDuringAskTimeoutTeardownStillCompletesTheAsyncTicket}'s job.
*/
@Test
void aCoalescedDuplicateAskTimingOutLeavesTheFreshOwnersAskOpen() throws Exception {
String ticket = messages.sendAsync(T, "long task");
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE); // deliver
injector.onStatus(T, AgentStatus.WORKING); // worker picks it up, then pauses to ask
CompletableFuture<MessageService.AskResult> fresh =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
MessageService.TaskView asking = null;
long deadline = System.currentTimeMillis() + 2000;
while ((asking == null || asking.phase() != MessageService.Phase.ASKING)
&& System.currentTimeMillis() < deadline) {
asking = messages.poll(ticket);
//noinspection BusyWait
Thread.sleep(5);
}
assertNotNull(asking, "the fresh owner's question must surface before the duplicate asks");
String turnId = asking.turnId();
assertNotNull(turnId, "an ASKING view carries the turnId to answer on");
// A coalesced duplicate on the same session, with its own much shorter timeout.
MessageService.AskResult duplicate = messages.ask(T, "which config file?", 100);
assertEquals(MessageService.AskOutcome.TIMED_OUT, duplicate.outcome(),
"the duplicate's own timeout elapses first");
CompletableFuture<MessageService.Reply> answered =
CompletableFuture.supplyAsync(() -> messages.answer(turnId, "fleetd.yaml", 500));
MessageService.AskResult a = fresh.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.AskOutcome.ANSWERED, a.outcome(),
"a duplicate's timeout must not lapse the ask the fresh owner still holds");
assertEquals("fleetd.yaml", a.answer());
assertFalse(answered.get(5, TimeUnit.SECONDS).outcome() == MessageService.Outcome.STALE_TURN,
"the answer must not be rejected as stale");
}
@Test
void aLateAnswerDuringAskTimeoutTeardownStillCompletesTheAsyncTicket() throws Exception {
String ticket = messages.sendAsync(T, "long task");
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE); // deliver
injector.onStatus(T, AgentStatus.WORKING); // worker picks it up, then pauses to ask
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 200));
// Wait for the question to actually surface (poll sees ASKING) before racing the timeout.
MessageService.TaskView asking = null;
long deadline = System.currentTimeMillis() + 2000;
while ((asking == null || asking.phase() != MessageService.Phase.ASKING)
&& System.currentTimeMillis() < deadline) {
asking = messages.poll(ticket);
//noinspection BusyWait
Thread.sleep(5);
}
assertNotNull(asking, "the question must surface before the ask times out");
String turnId = asking.turnId();
assertNotNull(turnId, "an ASKING view carries the turnId to answer on");
CompletableFuture<MessageService.Reply> lateAnswer = new CompletableFuture<>();
messages.setAskTimeoutRaceHookForTest(() ->
lateAnswer.complete(messages.answer(turnId, "too late", 500)));
try {
MessageService.AskResult a = ask.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.AskOutcome.TIMED_OUT, a.outcome());
MessageService.Reply late = lateAnswer.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.STALE_TURN, late.outcome(),
"a late answer racing the timeout teardown must see the ask as already lapsed");
// The worker resumes on its own (per the ask() contract) and eventually sends its real
// fleet_reply; the async ticket must still resolve with it, not strand at PENDING.
assertTrue(messages.reply(T, "real result"), "the worker's real reply must still be accepted");
} finally {
messages.setAskTimeoutRaceHookForTest(null);
}
MessageService.TaskView done = null;
deadline = System.currentTimeMillis() + 2000;
while ((done == null || done.phase() == MessageService.Phase.PENDING)
&& System.currentTimeMillis() < deadline) {
done = messages.poll(ticket);
//noinspection BusyWait
Thread.sleep(5);
}
assertNotNull(done);
assertEquals(MessageService.Phase.DONE, done.phase(), "the async ticket must not be stranded PENDING");
assertEquals("real result", done.reply());
}
@Test
void answeringAnUnknownTurnIsStale() {
MessageService.Reply r = messages.answer(T + "#999", "too late", 500);
@@ -410,6 +530,29 @@ class MessageServiceTest {
"a delivered send whose worker never replies times out as still working");
}
/**
* fleetd #345. This forces the injector to pick up the exact pending delivery after {@code send}
* first observes its completion as incomplete, but before {@code cancel} takes the target monitor.
* The timeout must use {@link Injector.Cancellation#DELIVERED} from {@code cancel} and report
* {@link MessageService.Outcome#TIMED_OUT_WORKING}, because the text landed.
*
* <p>What this does not prove: that this precise interleaving happens by itself under production
* timing. The test forces it through a test-only hook; it proves the timeout caller handles the
* injector result when the interleaving occurs.
*/
@Test
void sendTimeoutUsesCancellationDeliveredWhenPickupWinsTheRace() {
messages.setTimeoutCancellationRaceHookForTest(() -> injector.onStatus(T, AgentStatus.IDLE));
try {
MessageService.Reply reply = messages.send(T, "race delivery", 50);
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, reply.outcome(),
"cancel reporting DELIVERED means the worker received the timed-out message");
} finally {
messages.setTimeoutCancellationRaceHookForTest(null);
}
}
@Test
void answerTimesOutWhenTheResumedWorkerNeverReplies() throws Exception {
CompletableFuture<MessageService.Reply> send = sendAsync();
@@ -785,6 +928,49 @@ class MessageServiceTest {
assertFailedTicket(third, "agent target term_a not found");
}
// --- fleetd #335 (site 1): a per-task cleanup failure inside the abandon() loop must not -----
// strand the tasks that come after it. abandon()'s own comment on the loop documents the one
// real production call that can throw there (inbox.publish, in the stranded-reply put-back
// branch, reached when a concurrent reply() or a second abandon() races this one) — but
// reaching that branch requires the exact combination the #137 follow-up above already found
// unreachable through the public API. abandonCleanupHookForTest reproduces the resulting SHAPE
// (one task's cleanup throws) directly instead, the same technique this file already uses for
// fleetd #324/#329's own hard-to-time races.
@Test
void aPerTaskCleanupFailureDoesNotStrandTheRemainingMatchingTasks() throws Exception {
ListAppender<ILoggingEvent> appender = attachMessageServiceLog();
try {
String first = messages.sendAsync(T, "first task");
awaitWaiting(); // first task owns the target lock and rendezvous waiter
String second = messages.sendAsync(T, "second task"); // parked on the same lock
String third = messages.sendAsync(T, "third task"); // parked too — the whole sweep must survive
java.util.concurrent.atomic.AtomicInteger calls = new java.util.concurrent.atomic.AtomicInteger();
messages.setAbandonCleanupHookForTest(() -> {
if (calls.getAndIncrement() == 0) {
throw new RuntimeException("PROBE-335-SITE1");
}
});
assertTrue(messages.abandon(T, "agent target term_a not found"));
// Every task in the loop still gets its own outcome — the one whose cleanup threw
// included — even though the loop had no way to know in advance which one that would be.
assertFailedTicket(first, "agent target term_a not found");
assertFailedTicket(second, "agent target term_a not found");
assertFailedTicket(third, "agent target term_a not found");
assertTrue(appender.list.stream().anyMatch(e ->
e.getLevel() == Level.ERROR
&& e.getThrowableProxy() != null
&& "PROBE-335-SITE1".equals(e.getThrowableProxy().getMessage())),
"a per-task cleanup failure must still reach the log, not vanish silently");
} finally {
messages.setAbandonCleanupHookForTest(null);
detachMessageServiceLog(appender);
}
}
// --- #137 follow-up: abandon() must not guess when more than one task is open ---------------
//
// A test combining a genuine stranded reply (hasStrandedReply(T)==true) with two simultaneously
@@ -1472,6 +1658,42 @@ class MessageServiceTest {
}
}
// --- fleetd #335 (site 2): task.future.whenComplete's own returned stage is discarded, so an --
// uncaught throw from ReplyPushLoop.onTicketTerminal used to vanish with no log line and no
// metric. The real production trigger is the daemon's own shutdown sequence (Fleetd's shutdown
// hook): messages.close() only stops the async executor from taking NEW work — it does not
// cancel a send already in flight — while pushLoop.close() shuts its scheduler down immediately
// right after, so a ticket that completes in that narrow window has onTicketTerminal's own
// scheduler.schedule(...) throw a real RejectedExecutionException. Reproduced here by shutting
// the very same scheduler down before the ticket resolves — no test-only hook needed, this
// reachable path throws for real.
@Test
void aTicketTerminalPushFailureDoesNotVanishSilently() throws Exception {
ListAppender<ILoggingEvent> appender = attachMessageServiceLog();
try (var wiring = wireWithPushLoop(1, 50)) {
String ticket = wiring.service().sendAsync(T, "long task");
awaitWaiting();
wiring.scheduler().shutdownNow(); // simulate pushLoop.close() racing an in-flight send
injectDelivery();
assertTrue(rendezvous.resolve(T, "async result"));
// The ticket's own outcome must be unaffected by the swallowed exception — finishAsyncTask
// completes task.future before whenComplete's action (and thus onTicketTerminal) ever runs.
MessageService.TaskView done = awaitTicketPhaseOn(wiring.service(), ticket, MessageService.Phase.DONE);
assertEquals("async result", done.reply());
assertTrue(appender.list.stream().anyMatch(e ->
e.getLevel() == Level.ERROR
&& e.getFormattedMessage().contains(ticket)
&& e.getThrowableProxy() != null
&& "java.util.concurrent.RejectedExecutionException"
.equals(e.getThrowableProxy().getClassName())),
"onTicketTerminal throwing must still reach the log, not vanish silently");
} finally {
detachMessageServiceLog(appender);
}
}
// --- CB-582: fleet_ask question-open nudges --------------------------------------------------
@Test
@@ -711,13 +711,21 @@ class FleetAppTest {
@Test
void stopWorkerInPanePlacementClosesOnlyThePane() throws Exception {
FakeHerdr herdr = new FakeHerdr();
// fleetd #342: tab-cleanup resolution is no longer skipped based on a profile's declared
// placement — a stop() routed through the wrong delegate (e.g. CompositePeerLauncher's
// single-daemon shortcut after a daemon restart) could carry a placement config that says
// nothing true about how THIS pane was actually placed. So the pane's real tab is now
// always resolved, and the single-occupant check below is what protects a pane-placement
// peer's shared tab, exactly as it always protected a tab-placement one. Model that
// realistically: the peer's pane was split into an existing tab that already held another
// occupant, so the tab must never be closed.
FakeHerdr herdr = new FakeHerdr().withWorkerTabPaneCount(2);
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"), "pane");
assertEquals(204, req(port, "DELETE", "/members/w9:pW").statusCode());
assertTrue(herdr.called("pane.close"));
assertFalse(herdr.called("tab.close"), "pane placement owns no tab to close");
assertFalse(herdr.called("pane.get"), "no tab resolution in pane placement");
assertTrue(herdr.called("pane.get"), "tab resolution now always runs, regardless of placement");
assertFalse(herdr.called("tab.close"), "pane placement's shared tab must never be closed");
}
@Test