Compare commits

..

8 Commits

Author SHA1 Message Date
Dai Ha ca47e90c01 fleetd#323: close the reload-classifier drift with a reflection coverage test
CI / contract (pull_request) Successful in 1m11s
CI / build (pull_request) Successful in 2m6s
ConfigRef.sameLaunchSettings' javadoc claimed it compares every component
the launcher reads at spawn. It missed ideProjectDir, ideOpenCommand and
autoCompactWindow, and changedDeferredKeys separately missed worktreeGroup
(baked into the same GitWorktrees as worktreeRoot, Fleetd.java:251). A
reload that changed only one of those keys reported "config reloaded" with
nothing deferred, and the running daemon kept the old value.

Fix the four instances, and add ConfigRefProfileCoverageTest: it enumerates
every FleetConfig.Profile record component by reflection, mutates each one
not in the new ConfigRef.LAUNCH_SETTINGS_EXCLUDED set on a base profile,
and asserts sameLaunchSettings actually notices — so a fifth missed field
fails the build by name instead of drifting silently. It also prints its
own denominator (26 components, 23 compared, 3 excluded) per the ticket's
requirement that a checker must be able to state what it checked.

Also add the `profile` field itself to the comparison (it was neither
compared nor excluded before this fix — the coverage test surfaced it).

Rewrote the sameLaunchSettings javadoc to describe what the coverage test
actually guarantees instead of repeating the unchecked claim.
2026-09-04 14:26:09 +07:00
Dai Ha d05205d1eb Add a hunter skill: a sweep and a diff review are different jobs with different output contracts
CI / contract (push) Successful in 47s
CI / build (push) Failing after 1m49s
2026-09-04 14:13:08 +07:00
Dai Ha c801851c66 Correct the isLoopback javadoc: after #305 narrowing this range refuses a caller, it does not promote one
CI / contract (push) Successful in 1m16s
CI / build (push) Successful in 1m47s
2026-09-04 14:10:10 +07:00
Dai Ha de70aa38f1 Merge #317: an unresolved caller is refused, never promoted to primary
CI / contract (push) Successful in 1m1s
CI / build (push) Successful in 1m56s
2026-09-04 14:03:35 +07:00
Dai Ha 53a533afb4 #317: refuse an unresolved caller instead of promoting it to primary
CI / contract (pull_request) Successful in 47s
CI / build (pull_request) Successful in 1m52s
ConnectionIdentity.resolve() called pids.pidForLocalPort(remotePort),
which returns -1 both on a real failure and (silently, no log line)
when lsof just finds no matching process. terminalForPid(-1) then
matches no pane, so CallerResolver's loopback-trust fallback could not
tell that caller apart from a genuine primary and handed it
Principal.primary(...) — granting SPAWN, STOP, SEND and DRAIN to a
worker whose PID lookup failed. This is the escalation PaneLocator's
own javadoc already names; CB-161's ancestry walk only helps once a
candidate pid exists, and a failed lookup has none.

Fix: ConnectionIdentity.Caller gets a resolved() predicate (pid > 0),
centralised next to the -1 sentinel it tests for the same reason
isLoopback() is centralised (fleetd #305: two independent copies of
one rule already drifted once). CallerResolver's loopback-trust
fallback now requires c.resolved() before granting PRIMARY; an
unresolved caller gets Principal.anonymous() — the same already-tested
"authenticated as nothing" outcome used everywhere else in that
method, so the refusal is a clean, named, unsurprising result rather
than something that looks like a bug.

Also logs the previously-silent "lsof ran clean, found no match" case
in LsofPeerPidLookup at DEBUG, since that (not a slow lsof — the
waitFor result was already discarded) is the likelier real trigger.

loopbackTrustTreatsANonWorkerLoopbackCallerAsThePrimary is untouched
and still green: a real pid that owns no pane (the actual primary) is
still resolved() and still PRIMARY. Token mode is unaffected — it
never consults c.pid() at all.

Mutation-tested: reverting only the CallerResolver.java guard
reproduces the escalation exactly (aFailedPeerPidLookupIsRefusedNotPromotedToPrimary
fails with "expected: <ANONYMOUS> but was: <PRIMARY>").
2026-09-04 13:58:14 +07:00
Dai Ha 77ad88631b Merge #315: the fixed placement policy honours the retry loop's unreachable set
CI / build (push) Successful in 1m22s
CI / contract (push) Successful in 1m55s
2026-09-04 13:57:24 +07:00
Dai Ha d88017807b #315: fix self-contradicting javadoc left by the previous commit
CI / contract (pull_request) Successful in 1m21s
CI / build (pull_request) Successful in 2m31s
FixedPlacementPolicy's class javadoc still opened with "This ignores caps
and reachability" after the previous commit added reachability as the
fourth carve-out that is explicitly NOT ignored — caught by a shape-check
survey run against this same file as part of #315's own request ("look in
placement/ ... for the same shape: a caller/comment that documents an
expectation ... where an implementation does not meet it"). Reworded the
opening sentence: fixed still ignores caps (maxLoad) by design, but
reachability is now a narrower, per-call retry exclusion, not an ignored
concern.
2026-09-04 13:55:28 +07:00
Dai Ha 2159a5a94a #315: FixedPlacementPolicy now honors the retry loop's unreachable set
CI / contract (pull_request) Successful in 53s
CI / build (pull_request) Successful in 2m33s
CompositePeerLauncher.spawn retries a failed candidate on the next one and
rebuilds PlacementContext "so the policy excludes this profile" (its own
comment), but FixedPlacementPolicy.select never read ctx.unreachable(). Under
the default `fixed` placement policy (used when `placement` is unset or set
to `fixed`), every retry re-picked the same dead default and a second,
healthy, configured profile was never tried. This also covers the wiring-bug
branch (a candidate profile with no owning adapter), which hit the exact same
symptom for the same reason.

Not live on this fleet: fleetd.yaml sets placement: weighted, which already
consults ctx.unreachable() via PlacementPolicyUtil.available(). This is live
only for a deployment that leaves placement unset or sets it to fixed.

Fix is in FixedPlacementPolicy: consult ctx.unreachable() in the same two
places it already consults quarantined/coolingOff (the default check and the
fallback walk over candidates()), and add a fourth reason to the "no
candidate remains" exception. Considered fixing this in
CompositePeerLauncher's retry loop instead (break when select() returns an
already-unreachable profile), but that only fails faster on the same dead
profile — it cannot make the loop advance to a different candidate, because
only the policy decides which candidate is next. The defect is that one
policy implementation does not honor the loop's stated contract, so the fix
belongs in that policy, matching how weighted/round-robin already behave.

Also fixed: the "no reachable worker profile" exception message said
"trying N candidate(s)" where N was unreachable.size(), a count of DISTINCT
profiles (a HashSet dedupes a profile added twice), under wording that reads
as a count of attempts. Reworded to "N distinct candidate(s)" so the count
matches what is measured and the profile list that follows it.

Tests: two new failover tests next to the three existing ones in
CompositePeerLauncherTest (which all use PlacementPolicies.weighted(), which
is why this had no coverage) — one pinned to PlacementPolicies.fixed() for
the unreachable-default case, one for the wiring-bug (no adapter) case.
Mutation-proofed: reverted FixedPlacementPolicy.java, both new tests failed
with the exact bug ("no reachable worker profile available after trying 1
distinct candidate(s): a" / "...c"), then restored the fix.
2026-09-04 13:52:29 +07:00
16 changed files with 642 additions and 366 deletions
+102
View File
@@ -0,0 +1,102 @@
---
name: hunter
description: Defect-hunt procedure for a fleetd worker — sweep an assigned package for real bugs and report several ranked findings without fixing anything. Load this when the lead asks you to hunt or audit a scope rather than review one diff. Do NOT load `reviewer` for this; the two want different output.
---
# Hunter worker — procedure
The turn contract (one `fleet_reply`, `fleet_ask` for the lead's decisions, honest reporting,
never merge) is in **`CLAUDE.md` → Bridge communication → Worker** and already applies.
**This skill is not `reviewer`.** `reviewer` judges one diff and reports the *single* most
important issue in about 90 words. A hunt sweeps a whole package and reports *several* findings
in a long structured form. Loading both gives you two contradictory output contracts, and the
usual result is a worker that writes a good report into its terminal and ends the turn without
sending it. Load exactly one.
## 0. Read this before you read code: how the report gets home
Your terminal reaches nobody. The lead sees **only** the text inside your `fleet_reply` call.
A long report is exactly the case where this goes wrong, so plan for it:
- **Write the report into the `fleet_reply` argument itself.** Do not compose it in your terminal
and then summarise it into the call.
- If the report is long, **send it anyway** — one `fleet_reply` with everything.
- If you end the turn without replying, the bridge scrapes your pane instead. That scrape carries
at most the last 4000 characters, and on a hunt it usually captures the tail of the lead's own
brief rather than your findings. The lead then has nothing and has to ask you again.
## 1. Change nothing
A hunt is read-only. Do not edit a production file, do not "quickly fix" what you find, and do
not run a formatter. You may run the build and tests to *check* a claim, and you should say so
when you did.
## 2. Read the whole scope first
Read every file in the assigned package before you judge any of it. A defect that a caller
elsewhere in the same package makes unreachable is not a defect, and you cannot know that from
one file.
Stay inside the scope. If a defect there depends on a class outside it, read that class to
confirm — but the defect itself must live in the scope you were given.
## 3. The bar — this matters more than the count
**Name the path into the bad state.** Say which caller, in which state, reaches it. A defect on
paper is not a reachable defect. If you cannot name that path, keep the finding but mark it
`unproven` and say exactly what you could not check. Do not drop it, and do not dress it up.
**Say which direction the harm goes.** Data loss, privilege escalation and silent wrong answers
are worth reporting even when the window is narrow. A finding whose worst outcome is a worse log
line is not worth a block.
Two workers once ran the same scope: the one that applied the direction-of-harm filter found ten
real defects, the one that did not found none. Fewer findings the lead can act on beat many the
lead has to triage.
## 4. Shapes that have produced real merged fixes here
Read for these first:
1. **A one-way gate.** A guard added after an incident closes only the direction that incident
came from. Do not only ask what closes the gate — ask **which states still open it**.
2. **A value read once, then used later to authorise something destructive**, after something
else has had a chance to change it.
3. **A failure downgraded to a value that looks like a legitimate result** — `-1`, `null`, an
empty list, `false` — which a caller then trusts.
4. **A lock held for one half of a read-modify-write and not the other**, or two collections
updated under different locks.
5. **A comment or javadoc stating an invariant the code no longer keeps.** Comments are
load-bearing in this repo; a stale one has already caused a bug.
## 5. What you cannot check, and must not claim you did
- `fleetd/fleetd.yaml` is gitignored and **absent from your worktree**. You cannot read it. If a
finding depends on live configuration, name the key and say you could not check it.
- `.mcp.json`, `opencode.json` and `.autoenv` in your worktree are neutralised stubs, not the
repo's real files.
- The `wiki/` submodule pointer is months old. Do not cite it.
Reporting a fact you took from the lead's brief as something you measured yourself is a false
report, even when the fact is correct. Say where each fact came from.
## 6. The report — what goes in `fleet_reply`
One block per finding, most severe first:
```
FINDING N — <one line>
file:line
Path in: <which caller, in which state, reaches this>
Direction: <data loss | escalation | silent wrong answer | outage | ...>
Window/trigger: <when it actually happens>
Confidence: <confirmed by reading | unproven — say what you could not check>
Why nothing else catches it: <the guard or test you checked, and why it misses>
```
End with one line naming every file you read, so the lead knows the denominator.
**Nothing clears the bar?** Reply `NO FINDINGS`, name the files you read, and say what you ruled
out. A clean sweep is a valid result; an invented defect is worse than none.
+5
View File
@@ -10,6 +10,11 @@ never merge) is in **`CLAUDE.md` → Bridge communication → Worker** and alrea
skill is only the *review procedure*: how to work the scope, and the exact shape of what you
send back.
**Wrong skill for a sweep.** This one reviews *one* diff or scope and reports the *single* most
important issue. If the lead asked you to hunt or audit a whole package for several defects, load
`hunter` instead and ignore this file — the two want different output, and following both is how a
worker ends its turn with a good report that never gets sent.
## 1. Read the whole scope before you judge
The delegation names your scope — a file, a diff, a PR, a function. **Read all of it first.**
+7 -2
View File
@@ -200,8 +200,13 @@ must obey belongs in the charter, not here.
adapter, with a message naming the credential and the remaining seconds ("cooling off after
repeated backend errors") — distinct wording from a quarantine refusal, so don't conflate the
two when reading a spawn failure.
- **Skills available to delegate:** `implementer` (worktree → commit → push → own PR) and
`reviewer` (scoped review → one structured finding). Name one in every delegation.
- **Skills available to delegate:** `implementer` (worktree → commit → push → own PR),
`reviewer` (one diff → one structured finding) and `hunter` (sweep a package → several ranked
findings, change nothing). Name exactly one in every delegation. **`reviewer` and `hunter` are
not interchangeable** — `reviewer` caps the answer at one finding in about 90 words, so naming
it for a multi-finding sweep hands the worker two contradictory output contracts. That has
already cost three workers' turns: each wrote a good report to its terminal and ended the turn
with no `fleet_reply`, and the scrape returned the tail of the brief instead.
- **Primary-side skills** (not delegation playbooks — a worker cannot use them):
`port-to-opencode` (make an OpenCode session a participant in this workspace) and
`fleets-status` (report every fleet that shares one LavinMQ instance).
@@ -236,7 +236,15 @@ public final class CallerResolver {
// loopback-trust: same-host callers that are not workers are the primary. A non-loopback
// caller is anonymous even here — and startup refuses that combination anyway
// (FleetConfig.validateAuthExposure), so this is defence in depth, not the control.
return isLoopback(remoteAddr) ? Principal.primary(c.pid()) : Principal.anonymous();
//
// fleetd #317: "not a worker" must not be conflated with "identity unresolved". The real
// primary is a real process — its pid resolves (c.resolved()), it just owns no herdr pane.
// A caller whose peer-PID lookup failed (LsofPeerPidLookup's -1 sentinel — on any failure,
// silently including "lsof found no match") has no such pid, and PaneLocator's own javadoc
// already names what happens if that case is handed the primary role: a worker→primary
// escalation. So an unresolved caller is refused (ANONYMOUS — the same clean, already-tested
// "authenticated as nothing" outcome used everywhere else in this method), never promoted.
return isLoopback(remoteAddr) && c.resolved() ? Principal.primary(c.pid()) : Principal.anonymous();
}
private boolean presentedTokenMatches(String authorizationHeader) {
@@ -39,12 +39,18 @@ import java.util.function.Supplier;
* keeps the old value until a restart: {@code lifecycle:}, {@code leadHeartbeat:},
* {@code spawnReadyTimeoutMs} / {@code spawnReadyPollMs}, {@code quarantineCooldownSeconds}
* (CB-578 stage B — baked once into the {@code BackendQuarantine} built at startup),
* {@code guard:}, {@code worktreeRoot:}, adding or removing a profile (a new backend needs its own launcher,
* {@code guard:}, {@code worktreeRoot:} and {@code worktreeGroup:} (both baked once into the
* {@code GitWorktrees} built at {@code Fleetd.java:251} and never rebuilt — fleetd #323
* instance 2 found {@code worktreeGroup} missing from this list and from
* {@link #changedDeferredKeys}), adding or removing a profile (a new backend needs its own launcher,
* which is constructed once), <em>and an existing profile's launch settings</em> —
* {@code model}, {@code baseUrl}, {@code argv}, {@code env}, {@code mcpUrl},
* {@code exhaustedPattern} (CB-578 stage A — compiled once into {@code Fleetd.main}'s
* pattern map at startup), {@code errorPattern} (fleetd #201 Unit 5 — compiled once into
* {@code Fleetd.main}'s backend-error pattern map at startup, the same way), and the rest.
* {@code Fleetd.main}'s backend-error pattern map at startup, the same way),
* {@code ideProjectDir} / {@code ideOpenCommand} / {@code autoCompactWindow} (fleetd #323
* instance 1 — all three are read at spawn off the same frozen profile map and were missing
* from {@link #sameLaunchSettings}), and the rest of {@link #sameLaunchSettings}.
* {@code credentialId} (CB-578 stage B) is NOT on
* this list — it is read live off the config supplier at every quarantine check and
* exhaustion event, exactly like {@code weight} / {@code maxLoad}, so it is hot instead.
@@ -221,6 +227,13 @@ public final class ConfigRef implements Supplier<FleetConfig> {
if (!Objects.equals(old.worktreeRoot(), fresh.worktreeRoot())) {
changed.add("worktreeRoot");
}
// Baked into the same GitWorktrees as worktreeRoot (Fleetd.java:251) and never rebuilt
// either — see the class doc. Missing this check was fleetd #323 instance 2: a reload
// that changed only worktreeGroup reported "config reloaded" with nothing deferred, and
// newly provisioned worktrees kept the old sharing behaviour.
if (!Objects.equals(old.worktreeGroup(), fresh.worktreeGroup())) {
changed.add("worktreeGroup");
}
if (!Objects.equals(old.spawnReadyTimeoutMs(), fresh.spawnReadyTimeoutMs())
|| !Objects.equals(old.spawnReadyPollMs(), fresh.spawnReadyPollMs())) {
changed.add("spawnReady*");
@@ -265,14 +278,34 @@ public final class ConfigRef implements Supplier<FleetConfig> {
}
/**
* Whether two versions of a profile would launch a peer identically. Compares every component
* the launcher reads at spawn; {@code weight}, {@code maxLoad} and {@code credentialId} are
* excluded because those are read live (by the placement policy and, for credentialId, by
* {@code CompositePeerLauncher}/the CB-578 stage B exhaustion sink) and really do take effect on
* the next spawn.
* {@link FleetConfig.Profile} record components deliberately left out of
* {@link #sameLaunchSettings} because they are read <em>live</em>, not baked in at spawn — see
* the class doc's <em>Hot</em> bullet. {@code weight} and {@code maxLoad} are read live by the
* placement policy on every spawn; {@code credentialId} is read live by
* {@code CompositePeerLauncher} and the CB-578 stage B exhaustion sink. Nothing else is
* excluded — see {@code sameLaunchSettingsComparesEveryProfileComponentOrExcludesIt} in
* {@code ConfigRefProfileCoverageTest}, which enumerates every {@code Profile} record component
* by reflection and fails the build if one is neither compared below nor named here.
*/
private static boolean sameLaunchSettings(FleetConfig.Profile a, FleetConfig.Profile b) {
return Objects.equals(a.baseUrl(), b.baseUrl())
static final Set<String> LAUNCH_SETTINGS_EXCLUDED = Set.of("weight", "maxLoad", "credentialId");
/**
* Whether two versions of a profile would launch a peer identically.
*
* <p>This must compare every {@link FleetConfig.Profile} record component except the three in
* {@link #LAUNCH_SETTINGS_EXCLUDED}. That is not a claim this javadoc can make good on by
* itself — a javadoc saying "compares every component" is exactly what fleetd #323 found to be
* false for three fields (and a sibling method's field list, for a fourth). The actual
* guarantee comes from {@code ConfigRefProfileCoverageTest}: it enumerates every record
* component of {@code FleetConfig.Profile} by reflection, mutates each one not in
* {@code LAUNCH_SETTINGS_EXCLUDED} on a base profile, and asserts this method reports a
* difference — so a new component that is neither compared here nor added to
* {@code LAUNCH_SETTINGS_EXCLUDED} (with a reason) fails that test by name, rather than
* silently reporting "config reloaded" for a value the daemon never picked up.
*/
static boolean sameLaunchSettings(FleetConfig.Profile a, FleetConfig.Profile b) {
return Objects.equals(a.profile(), b.profile())
&& Objects.equals(a.baseUrl(), b.baseUrl())
&& Objects.equals(a.model(), b.model())
&& Objects.equals(a.configDir(), b.configDir())
&& Objects.equals(a.tokenEnv(), b.tokenEnv())
@@ -298,6 +331,16 @@ public final class ConfigRef implements Supplier<FleetConfig> {
// fleetd #201 Unit 5: errorPattern is compiled once into Fleetd.main's backend-error
// pattern map at startup (see BackendErrorPatternLookup wiring), the same way
// exhaustedPattern is — a reload never re-reads it either.
&& Objects.equals(a.errorPattern(), b.errorPattern());
&& Objects.equals(a.errorPattern(), b.errorPattern())
// fleetd #323 instance 1: ideProjectDir and ideOpenCommand are read at spawn off the
// same frozen profile map as ideMcpUrl above (ClaudeCodeLauncher.java:267/269,
// OpenCodeLauncher.java:474/480/486) and were missing from this comparison.
&& Objects.equals(a.ideProjectDir(), b.ideProjectDir())
&& Objects.equals(a.ideOpenCommand(), b.ideOpenCommand())
// fleetd #323 instance 1: autoCompactWindow is read at spawn the same way
// (ClaudeCodeLauncher.java:926, OpenCodeLauncher.java:650). Comparing it here only
// makes the reload REPORT that a restart is needed — it deliberately does not make
// autoCompactWindow take effect live, which is a separate, larger change.
&& Objects.equals(a.autoCompactWindow(), b.autoCompactWindow());
}
}
@@ -36,6 +36,25 @@ public final class ConnectionIdentity {
* primary / an off-host client) and its {@code pid} (or {@code -1} if not resolvable).
*/
public record Caller(String terminal, long pid) {
/**
* Whether the OS peer-PID lookup actually succeeded — {@code false} means {@code pid} is
* the {@code -1} sentinel, not a real process id, so this caller's identity could not be
* established at all. That is a different fact from a real pid that simply owns no worker
* pane (the primary's own connection): the primary is {@code resolved()} and has a
* {@code null terminal}; an unresolvable caller is {@code !resolved()} and also has a
* {@code null terminal}. The two look identical through {@link #terminal} alone, which is
* exactly how fleetd #317 happened — a failed {@code lsof} lookup and a genuine primary both
* fell through to {@code Principal.primary(...)}.
*
* <p>Centralised here, next to the sentinel it tests, for the same reason
* {@link ConnectionIdentity#isLoopback} is centralised rather than left for each caller to
* reimplement: a raw {@code pid > 0} check duplicated at every call site is precisely the
* "one rule, two copies" shape that let #305 drift.
*/
public boolean resolved() {
return pid > 0;
}
}
/** Resolve the caller's terminal and PID from one peer-PID lookup. */
@@ -71,13 +90,27 @@ public final class ConnectionIdentity {
* {@code 127.0.0.1:8765} with a source address of {@code 127.0.0.2} — measured on the Linux
* fleet host, where binding that source succeeds.
*
* <p><strong>Being strict here does not make the daemon safer; it makes it unsafe.</strong>
* That reads backwards, so it is worth stating plainly. This predicate does not decide whether
* a caller is trusted — it decides whether the caller's identity is <em>resolved at all</em>.
* Returning false means {@link #resolve} answers "no terminal", and downstream a caller with no
* terminal is treated as the primary under loopback-trust. So every address excluded here is an
* address on which a worker silently becomes the lead. Widening a check normally weakens it;
* widening this one is what closes the hole.
* <p><strong>What excluding an address costs, stated as it is today.</strong> This paragraph
* used to say that narrowing this range turned a worker into the lead, and that widening the
* check was what closed the hole. That was true only while there were <em>two</em> definitions
* that disagreed: {@code ConnectionIdentity} skipped the identity lookup for {@code 127.0.0.2}
* while {@code CallerResolver} read the same address as loopback and granted the primary role.
* #305 removed the second copy, and with one shared definition the old sentence no longer holds.
*
* <p>Measured on 2026-09-04 by narrowing this method back to exactly {@code 127.0.0.1} and
* running {@code CallerResolverTest} and {@code ConnectionIdentityTest}: a caller from
* {@code 127.0.0.2} then resolves to {@code ANONYMOUS}, not {@code PRIMARY} — for a worker
* ({@code aWorkerOnAnyLoopbackSourceAddressIsStillAWorkerNotThePrimary}) and for a non-worker
* ({@code aNonWorkerOnAnyLoopbackSourceAddressIsStillThePrimary}) alike. Excluding an address
* now <em>refuses</em> its caller; it does not promote one.
*
* <p>So keep the whole range, but for the plain reason: a genuine worker or primary that
* connects from {@code 127.0.0.2} must be identifiable at all, and narrowing this predicate
* locks it out. That is an outage, and an outage is the direction to fail in — which is exactly
* why the range must not be narrowed casually and also why doing so is no longer a security
* hole. This predicate still does not decide whether a caller is trusted; it decides whether the
* caller's identity is <em>resolved at all</em>. What makes an unresolved caller safe is
* {@link Caller#resolved()} (#317), not this method.
*/
public static boolean isLoopback(String addr) {
if (addr == null) {
@@ -41,6 +41,14 @@ public final class LsofPeerPidLookup implements PeerPidLookup {
if (!p.waitFor(2, TimeUnit.SECONDS)) {
p.destroyForcibly();
}
if (found < 0) {
// fleetd #317: this is the silent path — lsof ran clean and simply reported no
// matching process (e.g. queried before the OS socket table settles). Previously
// this logged nothing at all, which is exactly why the escalation went unnoticed;
// the exception path below already logs. A caller now refused because of this is
// still refused (never promoted) — this line only makes the refusal diagnosable.
log.debug("lsof peer-pid lookup for port {} found no matching process", port);
}
return found;
} catch (Exception e) {
log.debug("lsof peer-pid lookup for port {} failed: {}", port, e.getMessage());
@@ -385,9 +385,12 @@ public final class CompositePeerLauncher implements PeerLauncher {
}
}
// unreachable.size() counts DISTINCT profiles, not attempts (a HashSet dedupes a profile
// added twice) — say "distinct" so the count matches the sentence and the profile list that
// follows, rather than reading as a count of attempts made (fleetd #315).
throw new PeerUnreachableException(
"no reachable worker profile available after trying " + unreachable.size()
+ " candidate(s): " + String.join(", ", unreachable));
+ " distinct candidate(s): " + String.join(", ", unreachable));
}
/**
@@ -22,7 +22,6 @@ import java.util.concurrent.ConcurrentSkipListMap;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.util.concurrent.atomic.AtomicReference;
/**
* AMQP-backed {@link ReplyInbox} (CB-307 Stage 2): genuine cross-restart durability behind the same
@@ -94,23 +93,8 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
private final Channel channel;
/** All channel operations (publish/declare/ack/cancel) serialize on this — a Channel is not thread-safe. */
private final Object channelLock = new Object();
/**
* target → (msgId → held delivery). Per-target map is guarded by synchronizing on itself.
*
* <p><strong>CB-318 tombstone.</strong> The value {@link #RELEASED} is a reserved sentinel: it
* marks a target whose {@link #release} has already run, so {@link #deliverCallback} can tell a
* delivery landing after release() apart from a fresh target it has never seen. See both methods'
* javadoc for why a plain {@code held.remove(target)} is not enough.
*/
/** target → (msgId → held delivery). Per-target map is guarded by synchronizing on itself. */
private final ConcurrentHashMap<String, LinkedHashMap<String, Held>> held = new ConcurrentHashMap<>();
/**
* CB-318 sentinel stored in {@link #held} for a target whose {@link #release} has already run.
* Never mutated — every read site compares it by reference ({@code ==}) before touching it as a
* map, because it is a single object shared across every released target and calling a mutator on
* it would corrupt state for all of them.
*/
private static final LinkedHashMap<String, Held> RELEASED = new LinkedHashMap<>();
/** Targets whose queue is declared and consumer is running, mapped to their broker consumer tag. */
private final ConcurrentHashMap<String, String> consumerTags = new ConcurrentHashMap<>();
@@ -217,12 +201,6 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
channel.queueDeclare(queue, true, false, false, null); // durable, non-exclusive, keep on idle
String tag = channel.basicConsume(queue, false, deliverCallback(target), _ -> { });
consumerTags.put(target, tag);
// CB-318: drop a stale RELEASED tombstone from a prior ownership of this same target
// string, so a delivery under this fresh consumer is held normally instead of being
// nacked forever by deliverCallback's RELEASED check. Safe to do here, still under
// channelLock: no delivery for the consumer tag just registered above can reach
// deliverCallback before this basicConsume call returns.
held.remove(target, RELEASED);
log.debug("AMQP inbox owns queue {} for target {}", queue, target);
} catch (IOException e) {
throw new IllegalStateException("cannot own queue " + queue, e);
@@ -265,32 +243,6 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
* later connection drop even though this release did not manage to requeue it immediately. A
* failed {@code basicCancel} still throws, unchanged from before this fix — that failure means
* the consumer may still be attached, so best-effort requeue is not attempted underneath it.
*
* <p><strong>CB-318: {@code held.remove(target)} alone leaves a second window open.</strong> The
* bullet above already explains why cancelling first does not save a tag from going stale — but
* that only accounts for a delivery landing before this method starts touching {@link #held}.
* {@code basicCancel} stops <em>new</em> dispatches; it does not flush one already handed to the
* consumer work pool. So a delivery can still land on that pool's thread and reach
* {@link #deliverCallback} at any point during, or after, this method's body — and a plain
* {@code held.remove(target)} does nothing to stop it: {@code deliverCallback}'s
* {@code computeIfAbsent} finds the key gone and happily creates a brand-new map under it, which
* this method — already past its {@code remove} — never looks at again. That entry then sits
* delivered-but-unacked on {@link #channel} until the whole inbox closes: never requeued, never
* redelivered, and {@link #peek} is never called again for a target nothing owns any more.
*
* <p>The fix is {@link #held}{@code .compute(target, ...)} instead of {@code remove}: it takes
* whatever was held (to nack, same as before) and, in the same atomic step, leaves the
* {@link #RELEASED} tombstone behind instead of an absent key. {@code computeIfAbsent} and
* {@code compute} calls for the same key are mutually exclusive in {@link ConcurrentHashMap} —
* whichever of this call and a concurrent {@code deliverCallback} runs first is fully visible to
* the other, with no gap between them. So a delivery that loses the race sees a real map here and
* gets nacked by the loop below, same as always; a delivery that wins the race (runs first) is
* itself nacked by that same loop, once it settles into {@code held}. A delivery that arrives once
* this method has stored {@link #RELEASED} finds it via {@code computeIfAbsent} and refuses itself
* — see {@link #deliverCallback}. Either way nothing is silently retained forever, satisfying the
* ticket's invariant against dropping a message. This closes the window rather than merely
* narrowing it — correctness does not depend on how much time elapses between the swap and this
* method returning.
*/
@Override
public void release(String target) {
@@ -303,13 +255,8 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
throw new IllegalStateException("cannot cancel consumer for " + target, e);
}
}
AtomicReference<LinkedHashMap<String, Held>> previouslyHeld = new AtomicReference<>();
held.compute(target, (_, v) -> {
previouslyHeld.set(v);
return RELEASED;
});
var perTarget = previouslyHeld.get();
if (perTarget != null && perTarget != RELEASED) {
var perTarget = held.remove(target);
if (perTarget != null) {
synchronized (perTarget) {
for (Held h : perTarget.values()) {
try {
@@ -379,7 +326,7 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
@Override
public List<InboxMessage> peek(String target) {
var perTarget = held.get(target);
if (perTarget == null || perTarget == RELEASED) {
if (perTarget == null) {
return List.of();
}
synchronized (perTarget) {
@@ -390,7 +337,7 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
@Override
public void ack(String target, String msgId) {
var perTarget = held.get(target);
if (perTarget == null || perTarget == RELEASED) {
if (perTarget == null) {
return;
}
Held h;
@@ -423,20 +370,6 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
}
String content = new String(delivery.getBody(), StandardCharsets.UTF_8);
var perTarget = held.computeIfAbsent(target, _ -> new LinkedHashMap<>());
if (perTarget == RELEASED) {
// CB-318: release() already ran for this target and left the RELEASED tombstone in
// held (see release()'s javadoc) — computeIfAbsent() is guaranteed to see it rather
// than recreate a fresh map, because ConcurrentHashMap serializes compute/
// computeIfAbsent calls for the same key against each other. Refuse the delivery
// instead of holding it somewhere release() will never look at again: requeue it, the
// same way release() nacks its own held entries, so a later owner (or a connection
// drop) can still recover it. This does not need channelLock across a broker round
// trip — basicNack, like the duplicate-ack case just below, does not wait for one.
synchronized (channelLock) {
channel.basicNack(tag, false, true);
}
return;
}
boolean duplicate;
synchronized (perTarget) {
if (perTarget.containsKey(msgId)) {
@@ -5,10 +5,14 @@ import java.util.List;
/**
* Backward-compatible placement: an unqualified spawn always resolves to the configured default
* profile, exactly as {@code CompositePeerLauncher} did before CB-518. This ignores caps and
* reachability so that a pre-existing config behaves identically after upgrade.
* profile, exactly as {@code CompositePeerLauncher} did before CB-518. This ignores caps
* ({@code maxLoad}) so that a pre-existing config behaves identically after upgrade — capacity
* gating for automatic placement is deliberately out of scope for {@code fixed}, exactly as it
* always has been. Reachability is a narrower exception (fleetd #315, below): a profile is never
* checked for reachability up front, only skipped once it has already failed in <em>this same</em>
* spawn call's retry loop — see the unreachable case below.
*
* <p>Three exceptions walk past the default instead of returning it unconditionally:
* <p>Four exceptions walk past the default instead of returning it unconditionally:
* <ul>
* <li>Quarantine (CB-578 stage B): a quarantined default is a credential that just refused on
* a usage limit, not a transient capacity or reachability concern.
@@ -16,13 +20,21 @@ import java.util.List;
* ({@code BackendOutagePolicy}) — a separate, shorter-lived source from quarantine. When a
* profile is both quarantined and cooling off, only the quarantine reason is reported
* (exhaustion takes priority), matching {@code CompositePeerLauncher}'s explicit-spawn order.
* <li>Unreachable (fleetd #315): {@code CompositePeerLauncher.spawn} retries a failed candidate
* on the next one and rebuilds the {@link PlacementContext} so {@code ctx.unreachable()}
* names every profile that already failed with {@code PeerUnreachableException} in this same
* call. Without this check {@code select} kept handing back the same dead default forever —
* the retry loop's own comment says "so the policy excludes this profile", and this is what
* makes that true for {@code fixed} too, matching {@code weighted}/{@code round-robin}
* (both filter on {@code ctx.unreachable()} via {@link PlacementPolicyUtil#available}).
* <li>Weight 0 (CB-554): {@code fixed} is still automatic selection, so a profile the operator
* marked "never auto-select me" ({@code weight <= 0}) must be skipped here exactly as
* {@code weighted}/{@code round-robin} skip it — an explicit {@code fleet_spawn} naming
* the profile is unaffected, only this automatic fallback walk.
* </ul>
* A fleet where nothing is ever quarantined, cooling off, or weight-0 never exercises any of these
* paths, so today's behaviour is unchanged.
* A fleet where nothing is ever quarantined, cooling off, unreachable, or weight-0 never exercises
* any of these paths, so today's behaviour is unchanged — in particular, the very first selection
* of a spawn call always sees an empty {@code unreachable} set, so the first choice is untouched.
*/
final class FixedPlacementPolicy implements PlacementPolicy {
@@ -30,12 +42,12 @@ final class FixedPlacementPolicy implements PlacementPolicy {
public PlacementCandidate select(PlacementContext ctx) {
String d = ctx.defaultProfile();
if (d != null && !d.isBlank() && !ctx.quarantined().contains(d) && !ctx.coolingOff().contains(d)
&& !weightExcluded(ctx, d)) {
&& !ctx.unreachable().contains(d) && !weightExcluded(ctx, d)) {
return new PlacementCandidate(d, null, 1.0f, null);
}
for (PlacementCandidate c : ctx.candidates()) {
if (!ctx.quarantined().contains(c.profile()) && !ctx.coolingOff().contains(c.profile())
&& !c.excluded()) {
&& !ctx.unreachable().contains(c.profile()) && !c.excluded()) {
return new PlacementCandidate(c.profile(), null, c.weight(), c.maxLoad());
}
}
@@ -44,8 +56,9 @@ final class FixedPlacementPolicy implements PlacementPolicy {
// Exhaustion quarantine takes priority: reported only when quarantine is absent, so the
// message never claims "cooling off" for a profile that is really backend-exhausted.
boolean dCoolingOff = !dQuarantined && ctx.coolingOff().contains(d);
boolean dUnreachable = ctx.unreachable().contains(d);
boolean dWeightExcluded = weightExcluded(ctx, d);
if (dQuarantined || dCoolingOff || dWeightExcluded) {
if (dQuarantined || dCoolingOff || dUnreachable || dWeightExcluded) {
List<String> reasons = new ArrayList<>();
if (dQuarantined) {
reasons.add("is quarantined (backend exhausted)");
@@ -53,6 +66,9 @@ final class FixedPlacementPolicy implements PlacementPolicy {
if (dCoolingOff) {
reasons.add("is cooling off after repeated backend errors");
}
if (dUnreachable) {
reasons.add("is unreachable");
}
if (dWeightExcluded) {
reasons.add("has weight 0 (excluded from automatic selection)");
}
@@ -62,7 +78,7 @@ final class FixedPlacementPolicy implements PlacementPolicy {
}
if (!ctx.candidates().isEmpty()) {
throw new PlacementException("all worker profiles are excluded from automatic "
+ "selection (quarantined, cooling off, or weight-0)");
+ "selection (quarantined, cooling off, unreachable, or weight-0)");
}
throw new PlacementException("no worker profiles configured");
}
@@ -103,6 +103,54 @@ class CallerResolverTest {
assertEquals(Role.PRIMARY, p.role(), "the historical behaviour, now an explicit choice");
}
// ── fleetd #317: an unresolvable caller must never be promoted to the primary ──────────────────
// #305 closed the trigger where a resolved pid matched no pane *and* had no ancestry walk to
// save it. This is the other trigger PaneLocator's javadoc names: the pid never resolves at
// all — LsofPeerPidLookup returns -1 on any failure, including (silently) "lsof found no
// match" — so there is no candidate pid for an ancestry walk to even attempt.
/**
* The failing-without-the-fix case. Before #317's fix, {@code c.terminal() == null} was the
* only test in the loopback-trust fallback, and an unresolved pid produces exactly that same
* {@code null} terminal as a genuine primary — so this caller was handed
* {@code Principal.primary(...)}, a real worker's failed lookup becoming indistinguishable from
* the lead.
*/
@Test
void aFailedPeerPidLookupIsRefusedNotPromotedToPrimary() {
ConnectionIdentity unresolved = new ConnectionIdentity(new PaneLocator(herdr), _ -> -1);
Principal p = new CallerResolver(unresolved).resolve("127.0.0.1", 55555, null);
assertEquals(Role.ANONYMOUS, p.role(),
"an unresolvable caller must never be silently promoted to the primary");
}
/**
* The companion invariant #317 must not break: a caller whose lookup genuinely succeeded, and
* who simply owns no herdr pane — the real primary's own connection — is still the primary.
* This is {@link #loopbackTrustTreatsANonWorkerLoopbackCallerAsThePrimary} pinned again here,
* named for #317 and placed next to the test it must be distinguished from: same {@code null}
* terminal, opposite verdict, because {@code Caller.resolved()} tells them apart.
*/
@Test
void aRealPidThatOwnsNoPaneIsStillThePrimaryNotRefused() {
Principal p = new CallerResolver(nonWorkerIdentity()).resolve("127.0.0.1", 55555, null);
assertEquals(Role.PRIMARY, p.role());
}
/** #317 point 4: token mode never consults {@code c.pid()}, so a failed lookup must not change it. */
@Test
void tokenModeIsUndisturbedByAnUnresolvedLookup() {
ConnectionIdentity unresolved = new ConnectionIdentity(new PaneLocator(herdr), _ -> -1);
CallerResolver r = new CallerResolver(unresolved, true, "s3cret");
assertEquals(Role.ANONYMOUS, r.resolve("127.0.0.1", 55555, null).role(),
"no credential is still just ANONYMOUS, as before #317 — unchanged by the lookup failing");
assertEquals(Role.PRIMARY, r.resolve("127.0.0.1", 55555, "Bearer s3cret").role(),
"a valid token still authenticates the primary even though the peer-pid lookup failed");
}
@Test
void tokenModeRefusesANonWorkerCallerThatPresentsNoToken() {
Principal p = new CallerResolver(nonWorkerIdentity(), true, "s3cret")
@@ -0,0 +1,187 @@
package dev.ltms.fleet.config;
import org.junit.jupiter.api.Test;
import java.lang.reflect.Constructor;
import java.lang.reflect.RecordComponent;
import java.util.Arrays;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.TreeSet;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #323: {@code ConfigRef.sameLaunchSettings} javadoc used to claim it "compares every
* component the launcher reads at spawn". It did not — {@code ideProjectDir}, {@code
* ideOpenCommand} and {@code autoCompactWindow} were all baked in at daemon startup (see
* {@code ClaudeCodeLauncher}/{@code OpenCodeLauncher}) and missing from the comparison, so a reload
* that changed only one of them reported "config reloaded" and the running daemon kept the old
* value.
*
* <p>This class is the mechanism the issue asked for: it enumerates every record component of
* {@link FleetConfig.Profile} by reflection and proves — by actually mutating a base profile one
* field at a time and calling the real method — that each component is either compared by
* {@link ConfigRef#sameLaunchSettings} or named in {@link ConfigRef#LAUNCH_SETTINGS_EXCLUDED} with
* a reason. A new profile field that is neither fails this test by name, not a hand-maintained list
* going stale.
*/
class ConfigRefProfileCoverageTest {
private static final RecordComponent[] COMPONENTS = FleetConfig.Profile.class.getRecordComponents();
/**
* One valid, non-blank value per record component — "the a value". None of these trip any
* defaulting/normalization in {@code Profile}'s compact constructor (see {@code
* FleetConfig.java}), so what goes in is what {@code sameLaunchSettings} sees back out.
*/
private static final Map<String, Object> BASE = baseValues();
/** 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() {
Map<String, Object> v = new LinkedHashMap<>();
v.put("profile", "sonnet");
v.put("baseUrl", "http://gx00.gw:8000");
v.put("model", "sonnet");
v.put("configDir", "/config/a");
v.put("tokenEnv", "TOKEN_A");
v.put("argv", List.of("claude", "--flag-a"));
v.put("placement", "tab");
v.put("workspace", "workspace-a");
v.put("tabLabel", "label-a");
v.put("mcpUrl", "http://mcp-a");
v.put("cwd", "/cwd/a");
v.put("parityOverlay", List.of(".env", ".env.a"));
v.put("gitTokenEnv", "GIT_TOKEN_A");
v.put("gitHostEnv", "GITEA_HOST_A");
v.put("kind", "claude-code");
v.put("env", Map.of("K", "A"));
v.put("weight", 1.0f);
v.put("maxLoad", 5);
v.put("subscription", Boolean.TRUE);
v.put("exhaustedPattern", "usage limit a");
v.put("credentialId", "cred-a");
v.put("ideMcpUrl", "http://ide-mcp-a");
v.put("ideProjectDir", "modules/a");
v.put("ideOpenCommand", "open-cmd-a {dir}");
v.put("autoCompactWindow", 150000);
v.put("errorPattern", "error a");
assertNamesMatchComponents(v);
return v;
}
private static Map<String, Object> altValues() {
Map<String, Object> v = new LinkedHashMap<>();
v.put("profile", "sonnet-b");
v.put("baseUrl", "http://gx01.gw:8000");
v.put("model", "haiku");
v.put("configDir", "/config/b");
v.put("tokenEnv", "TOKEN_B");
v.put("argv", List.of("claude", "--flag-b"));
v.put("placement", "weighted");
v.put("workspace", "workspace-b");
v.put("tabLabel", "label-b");
v.put("mcpUrl", "http://mcp-b");
v.put("cwd", "/cwd/b");
v.put("parityOverlay", List.of(".env", ".env.b"));
v.put("gitTokenEnv", "GIT_TOKEN_B");
v.put("gitHostEnv", "GITEA_HOST_B");
v.put("kind", "opencode");
v.put("env", Map.of("K", "B"));
v.put("weight", 2.0f);
v.put("maxLoad", 9);
v.put("subscription", Boolean.FALSE);
v.put("exhaustedPattern", "usage limit b");
v.put("credentialId", "cred-b");
v.put("ideMcpUrl", "http://ide-mcp-b");
v.put("ideProjectDir", "modules/b");
v.put("ideOpenCommand", "open-cmd-b {dir}");
v.put("autoCompactWindow", 250000);
v.put("errorPattern", "error b");
assertNamesMatchComponents(v);
return v;
}
private static void assertNamesMatchComponents(Map<String, Object> values) {
Set<String> componentNames = new TreeSet<>();
for (RecordComponent rc : COMPONENTS) {
componentNames.add(rc.getName());
}
assertEquals(componentNames, new TreeSet<>(values.keySet()),
"this test's value map has drifted from FleetConfig.Profile's actual components — "
+ "update BASE/ALT alongside the record");
}
private static FleetConfig.Profile profileOf(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.Profile> ctor = FleetConfig.Profile.class.getDeclaredConstructor(types);
return ctor.newInstance(args);
}
/** {@code BASE} with exactly one named component swapped for its {@code ALT} value. */
private static FleetConfig.Profile mutate(String componentName) throws ReflectiveOperationException {
Map<String, Object> values = new LinkedHashMap<>(BASE);
values.put(componentName, ALT.get(componentName));
return profileOf(values);
}
/**
* The mechanism fleetd #323 asked for: enumerate {@link FleetConfig.Profile}'s record
* components, mutate each non-excluded one, and prove {@code sameLaunchSettings} actually
* notices — not just that some hand-maintained list claims it does. Prints the denominator
* (total / compared / excluded) the issue required: a checker that cannot state its own
* denominator is the failure this repo keeps hitting.
*/
@Test
void sameLaunchSettingsComparesEveryProfileComponentOrExcludesIt() throws ReflectiveOperationException {
int total = COMPONENTS.length;
Set<String> excluded = ConfigRef.LAUNCH_SETTINGS_EXCLUDED;
Set<String> allNames = new TreeSet<>();
for (RecordComponent rc : COMPONENTS) {
allNames.add(rc.getName());
}
assertTrue(allNames.containsAll(excluded),
"ConfigRef.LAUNCH_SETTINGS_EXCLUDED names a component that does not exist on "
+ "FleetConfig.Profile — check for a typo: " + excluded);
FleetConfig.Profile base = profileOf(BASE);
List<String> uncovered = new java.util.ArrayList<>();
int compared = 0;
for (RecordComponent rc : COMPONENTS) {
String name = rc.getName();
if (excluded.contains(name)) {
continue;
}
FleetConfig.Profile mutated = mutate(name);
if (ConfigRef.sameLaunchSettings(base, mutated)) {
uncovered.add(name);
} else {
compared++;
}
}
System.out.printf(
"ConfigRef.sameLaunchSettings coverage — %d Profile components total, %d compared, "
+ "%d excluded (%s)%n",
total, compared, excluded.size(), excluded);
assertEquals(List.of(), uncovered,
"these FleetConfig.Profile components changed but ConfigRef.sameLaunchSettings "
+ "reported no difference — add each one to the comparison (it is read at "
+ "spawn and baked in until a restart) or to ConfigRef.LAUNCH_SETTINGS_EXCLUDED "
+ "with a reason it is genuinely read live: " + uncovered);
assertEquals(total, compared + excluded.size(),
"every FleetConfig.Profile record component must be either compared or excluded — "
+ total + " components, " + compared + " compared, " + excluded.size()
+ " excluded");
}
}
@@ -434,6 +434,78 @@ class ConfigRefTest {
assertEquals("provider 5xx", ref.get().profiles().get("sonnet").errorPattern());
}
/**
* fleetd #323 instance 1: {@code ideProjectDir} is read at spawn off the frozen profile map
* (see {@code ClaudeCodeLauncher}/{@code OpenCodeLauncher}) exactly like {@code model}, but was
* missing from {@code sameLaunchSettings} — a reload changing only this field used to report a
* bare "config reloaded" and the running daemon kept launching with the old value.
*/
@Test
void changingAProfilesIdeProjectDirIsReportedAsDeferred(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
bind:
host: 127.0.0.1
port: 8765
herdrSocket: ~/.config/herdr/herdr.sock
profiles:
sonnet:
baseUrl: http://gx00.gw:8000
model: sonnet
ideProjectDir: fleetd
guard:
offSubscriptionHosts:
- gx00.gw
""");
ConfigRef ref = refFor(f);
Files.writeString(f, """
bind:
host: 127.0.0.1
port: 8765
herdrSocket: ~/.config/herdr/herdr.sock
profiles:
sonnet:
baseUrl: http://gx00.gw:8000
model: sonnet
ideProjectDir: fleetd-renamed
guard:
offSubscriptionHosts:
- gx00.gw
""");
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied());
assertEquals(1, out.deferred().size(), out.deferred().toString());
assertTrue(out.deferred().getFirst().contains("sonnet"), out.deferred().toString());
assertTrue(out.deferred().getFirst().contains("launch settings"), out.deferred().toString());
// The snapshot still carries the new value — a restart is what makes it take effect.
assertEquals("fleetd-renamed", ref.get().profiles().get("sonnet").ideProjectDir());
}
/**
* fleetd #323 instance 2: {@code worktreeGroup} is baked into the same {@code GitWorktrees}
* as {@code worktreeRoot} (Fleetd.java:251) and never rebuilt, but only {@code worktreeRoot}
* was on {@code changedDeferredKeys} — a reload changing only the group reported a bare
* "config reloaded" and newly provisioned worktrees kept the old sharing behaviour.
*/
@Test
void changingWorktreeGroupIsReportedAsDeferred(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml("worktreeGroup: devgroup\n"));
ConfigRef ref = refFor(f);
Files.writeString(f, yaml("worktreeGroup: devgroup2\n"));
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied());
assertEquals(java.util.List.of("worktreeGroup"), out.deferred());
assertTrue(out.summary().contains("needs a restart") || out.summary().contains("need a restart"),
out.summary());
// The snapshot still carries the new value — a restart is what makes it take effect.
assertEquals("devgroup2", ref.get().worktreeGroup());
}
@Test
void aFixedRefHasNoFileAndRefusesToReload() {
FleetConfig cfg = new FleetConfig(null, null, null, null, null, null,
@@ -43,6 +43,22 @@ class ConnectionIdentityTest {
assertNull(with(_ -> 999_999).callerTerminal("127.0.0.1", 55555));
}
@Test
void callerIsUnresolvedWhenThePeerPidLookupFails() {
// fleetd #317: LsofPeerPidLookup returns -1 on any failure — a fork error, or (silently)
// simply no matching lsof line. Caller.resolved() is the one place that sentinel is tested.
ConnectionIdentity.Caller c = with(_ -> -1).resolve("127.0.0.1", 55555);
assertFalse(c.resolved(), "a -1 pid means the lookup failed, not that this pid owns no pane");
}
@Test
void callerIsResolvedWhenThePidIsRealEvenThoughItOwnsNoPane() {
// The primary's own connection: a real, lsof-found pid that just isn't a worker pane. This
// must read as "resolved" — the distinction #317 turns on.
ConnectionIdentity.Caller c = with(_ -> 999_999).resolve("127.0.0.1", 55555);
assertTrue(c.resolved());
}
@Test
void resolvesTheCallersPidAndCwd() {
// CB-112: the primary maps to no pane, but its PID and cwd are still readable.
@@ -615,6 +615,67 @@ class CompositePeerLauncherTest {
assertEquals(1, adapter.spawnCount("b"));
}
/**
* fleetd #315: {@code CompositePeerLauncher.spawn} rebuilds the {@link PlacementContext} after
* every failed attempt "so the policy excludes this profile" (see the comment at the retry call
* site) — but {@code FixedPlacementPolicy} never read {@code ctx.unreachable()}, so under the
* default {@code fixed} placement every retry re-picked the same dead default and a second,
* healthy, configured profile was never tried. This is the same scenario as
* {@link #failoverRetriesNextCandidateWhenProfileIsUnreachable}, but pinned to {@code fixed()}
* instead of {@code weighted()} — the three existing failover tests all use {@code weighted()},
* which is exactly why nobody caught this: the retry loop's contract has no coverage under its
* own default policy.
*/
@Test
void failoverRetriesNextCandidateUnderFixedPlacementWhenProfileIsUnreachable() {
FakeHerdr herdr = new FakeHerdr();
Map<String, FleetConfig.Profile> profiles = ordered(
"a", stubWorker("a"),
"b", stubWorker("b"));
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "a", Set.of("a"));
CompositePeerLauncher composite = new CompositePeerLauncher(
List.of(adapter), "a", profiles, PlacementPolicies.fixed(), _ -> 0);
PeerHandle h = composite.spawn(new SpawnRequest(null, null, null));
assertEquals("b", h.profile(),
"fixed placement must fail over from the unreachable default a to the healthy b");
assertEquals(1, adapter.spawnCount("a"), "a was tried once and failed");
assertEquals(1, adapter.spawnCount("b"), "b was tried once and succeeded");
}
/**
* fleetd #315: the same fix — {@code FixedPlacementPolicy} consulting {@code ctx.unreachable()}
* — also covers the wiring-bug branch in {@code CompositePeerLauncher.spawn}: a profile that
* placement is allowed to choose (it is in the configured candidate list) but that no delegate
* declares ({@code byProfile.get(chosen.profile()) == null}). That branch adds the profile to
* {@code unreachable} and {@code continue}s without ever calling a launcher, so before this fix
* {@code fixed} handed back the same adapterless profile on every remaining attempt too.
*/
@Test
void failoverSkipsAConfiguredProfileNoAdapterDeclaresUnderFixedPlacement() {
FakeHerdr herdr = new FakeHerdr();
// Placement's candidate list has three profiles, in this order (LinkedHashMap preserves it,
// and the fixed default resolves to the first — see the `ordered` helper's own javadoc).
Map<String, FleetConfig.Profile> profiles = new LinkedHashMap<>();
profiles.put("c", stubWorker("c"));
profiles.put("a", stubWorker("a"));
profiles.put("b", stubWorker("b"));
// The adapter only declares a and b — c is a configured profile with no owning adapter,
// the "wiring bug" the comment in CompositePeerLauncher.spawn calls out.
Map<String, FleetConfig.Profile> adapterProfiles = new LinkedHashMap<>();
adapterProfiles.put("a", stubWorker("a"));
adapterProfiles.put("b", stubWorker("b"));
StubLauncher adapter = new StubLauncher("claude", herdr, adapterProfiles, "a", Set.of());
CompositePeerLauncher composite = new CompositePeerLauncher(
List.of(adapter), "a", profiles, PlacementPolicies.fixed(), _ -> 0);
PeerHandle h = composite.spawn(new SpawnRequest(null, null, null));
assertEquals("a", h.profile(),
"c has no adapter, so fixed placement must skip it and land on the next candidate, a");
assertEquals(0, adapter.spawnCount("c"), "c is never spawned — no adapter owns it");
assertEquals(1, adapter.spawnCount("a"));
}
@Test
void explicitSpawnAtMaxLoadThrowsPlacementExceptionNamingProfileLiveAndCap() {
FakeHerdr herdr = new FakeHerdr();
@@ -1,264 +0,0 @@
package dev.ltms.fleet.msg;
import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.DeliverCallback;
import com.rabbitmq.client.Delivery;
import com.rabbitmq.client.Envelope;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.Timeout;
import java.lang.reflect.InvocationHandler;
import java.lang.reflect.Proxy;
import java.nio.charset.StandardCharsets;
import java.time.Duration;
import java.util.List;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicReference;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* CB-318: a delivery landing on the consumer work-pool thread <em>after</em>
* {@link AmqpReplyInbox#release} has already swapped the target's {@code held} entry for its
* tombstone, but <em>before</em> {@code release()} itself returns, must be nacked-with-requeue —
* never silently retained in a fresh map {@code release()} has already stopped looking at.
*
* <p><strong>This forces the actual interleaving, not a sequence of calls.</strong> {@code release()}
* runs on its own thread and is made to block <em>inside</em> its nack loop, via a fake
* {@link Channel} whose {@code basicNack} blocks on its first invocation. That block is only
* reachable after {@code release()}'s {@code held.compute(...)} has already swapped in the
* {@code RELEASED} tombstone (the compute call happens strictly before the loop that calls
* {@code basicNack}), so observing it is direct, ordering-guaranteed proof that the tombstone is in
* place and {@code release()} has not yet returned — still holding {@code channelLock} — when a
* second thread fires {@code own()}'s captured {@link DeliverCallback} for a brand-new message on the
* same target. No mocking library is on the classpath, so the fake broker is a {@link Proxy}, the
* same pattern {@code AmqpReplyInboxRecoveryRaceTest} already uses.
*
* <p><strong>What this does and does not prove.</strong> It proves that a delivery whose
* {@code computeIfAbsent} call is ordered strictly after {@code release()}'s tombstone swap — while
* {@code release()} is still running — is nacked-with-requeue rather than silently parked forever.
* It does not drive a real broker: {@code basicNack} here is a recorded call on a fake channel, not a
* verified requeue-and-redeliver. That half of the contract (a nacked-with-requeue delivery really
* does come back to a later owner) is already covered against a real broker by
* {@code AmqpReplyInboxContractTest.releaseCancelsConsumerAndRequeuesHeldDeliveryForRecovery}, which
* this test does not duplicate.
*/
class AmqpReplyInboxReleaseRaceTest {
@Test
@Timeout(15)
void deliveryArrivingWhileReleaseIsStillRunningIsNackedNotStranded() throws Exception {
String target = "worker-release-race";
List<long[]> nacks = new CopyOnWriteArrayList<>(); // {deliveryTag, requeue(1/0)}
List<Long> acks = new CopyOnWriteArrayList<>();
AtomicReference<DeliverCallback> deliverCallback = new AtomicReference<>();
CountDownLatch nackStarted = new CountDownLatch(1);
CountDownLatch releaseMayFinishNack = new CountDownLatch(1);
AtomicInteger nackCallCount = new AtomicInteger();
Channel consumeChannel = fakeConsumeChannel(deliverCallback, nacks, acks, nackCallCount,
nackStarted, releaseMayFinishNack);
Channel publishChannel = fakeInertChannel();
Connection connection = fakeConnection(consumeChannel, publishChannel);
AmqpReplyInbox inbox = new AmqpReplyInbox(connection, AmqpReplyInbox.DEFAULT_PREFETCH);
inbox.own(target);
assertTrue(deliverCallback.get() != null, "own() must have registered a DeliverCallback");
// Seed one already-held delivery (m0) so release()'s nack loop has something to iterate, and
// therefore somewhere to block, before it can return.
deliverCallback.get().handle("ctag", delivery(1L, "m0", "first"));
AtomicReference<Throwable> releaseError = new AtomicReference<>();
Thread releaseThread = new Thread(() -> {
try {
inbox.release(target);
} catch (Throwable t) {
releaseError.set(t);
}
}, "release-under-test");
releaseThread.start();
// This latch only fires from inside the fake channel's basicNack — i.e. from inside
// release()'s nack loop, which release()'s code only reaches AFTER held.compute(...) has
// already swapped in RELEASED. Waiting for it is direct proof the swap has happened and
// release() has not yet returned (it is stuck mid-loop, still holding channelLock).
assertTrue(nackStarted.await(10, TimeUnit.SECONDS),
"release() never reached its nack loop — it may not have started");
// The exact interleaving CB-318 describes: a delivery for a NEW message on the same target
// lands on the "consumer work-pool thread" (this second thread) while release() is still
// running. With the pre-fix code (a bare held.remove(target)) this created a brand-new map
// under computeIfAbsent that release() — already past its remove — never looks at again.
AtomicReference<Throwable> deliveryError = new AtomicReference<>();
Thread deliveryThread = new Thread(() -> {
try {
deliverCallback.get().handle("ctag", delivery(2L, "m1", "second"));
} catch (Throwable t) {
deliveryError.set(t);
}
}, "concurrent-delivery");
deliveryThread.start();
// Head start for the delivery thread to reach (and, on the fixed code, block on)
// channelLock — release() still holds it at this point, so a correct fix cannot have
// resolved m1's nack yet. Purely in-memory work (computeIfAbsent, a reference compare)
// separates deliveryThread.start() from that block point, so 300ms is a large margin, not a
// tight timing assumption.
Thread.sleep(300);
assertEquals(1, nackCallCount.get(),
"the concurrent delivery must not resolve its nack before release() gives up "
+ "channelLock — if this is 2 already, the interleaving below is not being "
+ "tested, only a sequential call");
releaseMayFinishNack.countDown(); // let release() finish nacking m0 and return
assertTrue(releaseThread.join(Duration.ofSeconds(10)), "release() did not finish");
assertTrue(deliveryThread.join(Duration.ofSeconds(10)), "the concurrent delivery did not finish");
assertNull(releaseError.get(), "release() threw: " + releaseError.get());
assertNull(deliveryError.get(), "the concurrent delivery threw: " + deliveryError.get());
assertEquals(2, nacks.size(),
"both the pre-held m0 and the concurrently-arriving m1 must be nacked, got: "
+ nacks.stream().map(n -> "[tag=" + n[0] + " requeue=" + n[1] + "]").toList());
assertTrue(nacks.stream().allMatch(n -> n[1] == 1L),
"invariant 1 (never drop): every nack must set requeue=true");
assertTrue(nacks.stream().anyMatch(n -> n[0] == 1L), "m0's delivery tag must be nacked");
assertTrue(nacks.stream().anyMatch(n -> n[0] == 2L),
"m1 — delivered while release() was still running, after the tombstone swap — must be "
+ "nacked, not silently retained in a map release() will never look at again");
assertTrue(acks.isEmpty(), "invariant 1 (never drop): a held reply must never be basicAck'd");
}
private static Delivery delivery(long tag, String msgId, String body) {
Envelope envelope = new Envelope(tag, false, "", "irrelevant");
AMQP.BasicProperties props = new AMQP.BasicProperties.Builder().messageId(msgId).build();
return new Delivery(envelope, props, body.getBytes(StandardCharsets.UTF_8));
}
/** A {@link Proxy}-backed consume {@link Channel}: blocks the FIRST {@code basicNack} call on
* {@code releaseMayFinishNack}, after signalling {@code nackStarted} — everything else records
* the call and returns a harmless default, matching the style already used by
* {@code AmqpReplyInboxRecoveryRaceTest}. */
private static Channel fakeConsumeChannel(AtomicReference<DeliverCallback> deliverCallback,
List<long[]> nacks, List<Long> acks,
AtomicInteger nackCallCount,
CountDownLatch nackStarted,
CountDownLatch releaseMayFinishNack) {
InvocationHandler handler = (proxy, method, args) -> {
String name = method.getName();
if (name.equals("basicConsume")) {
deliverCallback.set((DeliverCallback) args[2]);
return "ctag";
}
if (name.equals("basicNack")) {
long tag = (long) args[0];
boolean requeue = (boolean) args[2];
if (nackCallCount.incrementAndGet() == 1) {
nackStarted.countDown();
if (!releaseMayFinishNack.await(10, TimeUnit.SECONDS)) {
throw new IllegalStateException("test never released the nack latch");
}
}
nacks.add(new long[] {tag, requeue ? 1L : 0L});
return null;
}
if (name.equals("basicAck")) {
acks.add((long) args[0]);
return null;
}
if (name.equals("equals")) {
return proxy == args[0];
}
if (name.equals("hashCode")) {
return System.identityHashCode(proxy);
}
if (name.equals("toString")) {
return "FakeConsumeChannel";
}
return defaultValue(method.getReturnType());
};
return (Channel) Proxy.newProxyInstance(AmqpReplyInboxReleaseRaceTest.class.getClassLoader(),
new Class<?>[] {Channel.class}, handler);
}
/** A {@link Proxy}-backed {@link Channel} that answers every call with a harmless default — used
* as the publish channel, which this test never actually publishes on. */
private static Channel fakeInertChannel() {
InvocationHandler handler = (proxy, method, args) -> {
String name = method.getName();
if (name.equals("equals")) {
return proxy == args[0];
}
if (name.equals("hashCode")) {
return System.identityHashCode(proxy);
}
if (name.equals("toString")) {
return "FakeInertChannel";
}
return defaultValue(method.getReturnType());
};
return (Channel) Proxy.newProxyInstance(AmqpReplyInboxReleaseRaceTest.class.getClassLoader(),
new Class<?>[] {Channel.class}, handler);
}
/** A {@link Proxy}-backed {@link Connection} handing out {@code first} then {@code second} from
* successive {@code createChannel()} calls, matching {@link AmqpReplyInbox}'s constructor. */
private static Connection fakeConnection(Channel first, Channel second) {
AtomicInteger calls = new AtomicInteger();
InvocationHandler handler = (proxy, method, args) -> {
String name = method.getName();
if (name.equals("createChannel") && (args == null || args.length == 0)) {
return calls.getAndIncrement() == 0 ? first : second;
}
if (name.equals("equals")) {
return proxy == args[0];
}
if (name.equals("hashCode")) {
return System.identityHashCode(proxy);
}
if (name.equals("toString")) {
return "FakeConnection";
}
return defaultValue(method.getReturnType());
};
return (Connection) Proxy.newProxyInstance(AmqpReplyInboxReleaseRaceTest.class.getClassLoader(),
new Class<?>[] {Connection.class}, handler);
}
private static Object defaultValue(Class<?> type) {
if (!type.isPrimitive() || type == void.class) {
return null;
}
if (type == boolean.class) {
return Boolean.FALSE;
}
if (type == long.class) {
return 0L;
}
if (type == short.class) {
return (short) 0;
}
if (type == byte.class) {
return (byte) 0;
}
if (type == char.class) {
return (char) 0;
}
if (type == double.class) {
return 0.0d;
}
if (type == float.class) {
return 0.0f;
}
return 0;
}
}