Compare commits

...

29 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
Dai Ha b9c2cf69f4 Merge #307: a worker's real reply after an ask timeout completes its ticket instead of stranding
CI / build (push) Successful in 2m5s
CI / contract (push) Successful in 2m19s
2026-09-04 13:31:43 +07:00
Dai Ha 8beae50fe7 Merge #308: refuse spawns once the shutdown drain has started, and sweep stragglers
CI / contract (push) Successful in 1m49s
CI / build (push) Successful in 3m0s
2026-09-04 13:28:09 +07:00
Dai Ha b2a58cb966 Merge #309: clean up partial worktree state when git worktree add fails 2026-09-04 13:28:04 +07:00
Dai Ha b8b25cf74c #307: an ask() timeout no longer strands the worker's real reply
CI / contract (pull_request) Successful in 42s
CI / build (pull_request) Successful in 1m30s
MessageService.reply()'s async-recovery path (askAnsweredAsyncTasks)
required a live Task.turnId, but ask()'s own TimeoutException handler
calls clearAsyncQuestion(turnId, true) — deliberately forgetting turnId
so hasAsyncQuestion() stops reporting the target BUSY. That made a
worker's eventual real fleet_reply, after an unanswered fleet_ask, fall
through to the inbox: fleet_poll{ticket} stayed PENDING forever and was
later force-failed with the false reason "session released before it
replied".

Fix: a new Task.askTimedOut marker is set (markAskTimedOut) right
before the turnId is forgotten, and askAnsweredAsyncTasks accepts it in
place of a live turnId. The marker never touches asyncTasksByTurn, so
the BUSY-release behaviour (invariant 1) is untouched. The existing
ambiguity guard (candidates.size() > 1 -> inbox, never guess) still
applies unchanged, but is now genuinely reachable rather than pure
defence in depth, since an ask timeout frees its target for a fresh,
independent delegation — the affected javadocs are updated to say so.

Tests: MessageServiceTest.aReplyAfterAnAskTimeoutStillCompletesTheAsyncTicket
(positive, mutation-proven) and
.twoAskTimedOutTicketsOnOneTargetFallBackToTheInboxRatherThanGuess
(negative/ambiguity). FleetMcpTest's
unansweredAsyncAskReturnsTheTicketToPending was renamed and its final
assertion updated — it had pinned the old (buggy) inbox-stranding
behaviour as expected.
2026-09-04 13:23:01 +07:00
Dai Ha f159ca7d27 #310: log when a reap is skipped because the record changed
CI / contract (push) Successful in 1m7s
CI / build (push) Successful in 1m46s
The compare-and-release declines silently. This race is unobservable by
construction, so a reaper that quietly stops reaping is the hardest kind of
behaviour to diagnose later. One debug line names the pane and the likely
cause.
2026-09-04 13:22:55 +07:00
Dai Ha 83f2aea60f #308: refuse a spawn once the shutdown drain has started, and sweep stragglers
CI / contract (pull_request) Successful in 1m28s
CI / build (pull_request) Successful in 1m36s
drainAll iterated a one-shot registry snapshot with nothing to refuse a new
fleet_spawn while the drain was still running (mcp.close() only runs 8 calls
after sessions.close() in the shutdown hook). A session registered in that
window was never visited by the drain loop: its pane kept running and its
worktree was never preserved, with the in-memory registry gone at exit.

Fix, both mechanisms as the issue asked for (neither alone is complete):

- SessionManager.acquire now checks a `draining` flag, flipped true at the
  very start of drainAll before the registry snapshot is even taken, and
  throws the new ShuttingDownException (invariant 3: fail loudly, say why).
  FleetMcp.spawn and FleetApp.spawnMember surface it as a clean error/503
  rather than an uncaught RuntimeException.
- The flag alone cannot close the whole race: a caller already past the
  check can still be mid-launcher.spawn() (a real herdr round trip) when
  drainAll snapshots the registry. drainAll now re-reads the registry once
  its main pass finishes and drains whatever straggler landed there too,
  bounded by the SAME whole-drain deadline (invariant 1: timeoutNanos stays
  a budget for the whole drain, never extended for a straggler).
- ReleaseCause.SHUTDOWN still preserves worktrees for both the initial pass
  and the sweep (invariant 2, unchanged release() path).

Tests (SessionManagerTest): a guard test proving acquire() throws once
drainAll has started, and a race test using a launcher double that blocks
the second spawn() and the first stop() call to force, deterministically,
the exact interleaving where a spawn passes the guard before drainAll flips
it and only registers after the initial snapshot — proving the post-loop
sweep catches it.

Shape check (SessionManager.java only, not fixed): reapIdle has the same
shape — a decision made from a roster() snapshot, then acted on via
release(s.paneId()) with no re-check of the session's current state.
2026-09-04 13:21:53 +07:00
Dai Ha 002329adb5 #309: clean partial worktrees after add failure
CI / contract (pull_request) Successful in 1m23s
CI / build (pull_request) Successful in 2m2s
2026-09-04 13:18:57 +07:00
Dai Ha a49671ceb9 #310: prevent idle reap from stopping delivered workers
CI / contract (pull_request) Successful in 1m19s
CI / build (pull_request) Successful in 2m1s
2026-09-04 13:17:51 +07:00
Dai Ha 76672ff016 #306: the post-turn phase gets the same unknown-stall escape as a turn
CI / build (push) Successful in 1m43s
CI / contract (push) Successful in 1m50s
Four latches gate delivery in Injector, and only awaitingCompletion had a way
out of a sustained unknown streak. CB-109 added that escape because a worker
stuck in a state herdr cannot classify never produces a working->idle
boundary. The same is true during post-turn housekeeping, but the escape was
never extended there.

awaitingPostTurnPickup and postTurnObserved are both released only on an
injectable sample, so a worker that goes unknown and stays there wedges: the
target is polled forever, every later message to it is blocked by the delivery
gate, and no onTurnFailed fires, so the session sits at DONE and looks healthy.
The counter did not even increment, since ++unknownSinceTurn sits inside the
awaitingCompletion short-circuit.

postTurnPending needs no escape; it is cleared on the line after the listener
call that sets it.

The escape does not set turnFailed. The delegated turn already completed and
its waiter already resolved — what is outstanding is the /clear. Failing the
turn would drive SessionManager.onFailed on a session that genuinely finished.

Not reachable in the live configuration: the path needs lifecycle.clearAfterTurn,
which fleetd.yaml does not set. It becomes reachable as soon as anyone turns
that supported knob on.

Fixes #306
2026-09-04 13:06:29 +07:00
Dai Ha 9379f92c23 #305: one definition of loopback, so a worker cannot become the primary
CI / contract (push) Successful in 52s
CI / build (push) Successful in 1m40s
ConnectionIdentity and CallerResolver each kept their own isLoopback. They
drifted: the identity resolver accepted only 127.0.0.1, the authorization
check accepted all of 127.0.0.0/8.

A caller from 127.0.0.2 therefore had its identity resolution skipped, so it
carried no terminal, and CallerResolver reads a missing terminal as "not a
worker" — which under loopback-trust, the default mode, is the primary. A
worker got spawn, stop, send and drain. The skip also happens before the PID
ancestry walk, so that defence is bypassed too.

Being strict in ConnectionIdentity was not the safe direction. That predicate
decides whether identity is resolved at all, and resolution is what demotes a
worker, so every address it excluded was one where a worker became the lead.

Measured, not assumed: on Linux the whole 127.0.0.0/8 is bound to lo, and
binding a source of 127.0.0.2 on the fleet host succeeds (curl rc=7, the
connect refused rather than the bind). On macOS the source bind fails (rc=45),
so this workstation was never exposed.

The shared predicate also accepts the IPv4-mapped IPv6 form, which neither
copy handled. That one failed in the safe direction: a primary on
::ffff:127.0.0.1 was refused as anonymous.

No transport-level test binds a real 127.0.0.2 source — it cannot run on
macOS. The reasoning is recorded on the issue.

Fixes #305
2026-09-04 12:58:21 +07:00
Dai Ha 21ff63b11d #304: the member routes report a herdr failure instead of a bare 500
CI / build (push) Successful in 1m34s
CI / contract (push) Successful in 1m53s
POST /members and DELETE /members/{paneId} were the two routes in FleetApp
with no catch (HerdrException). FleetApp has no Javalin exception mapper, so
the exception escaped as the default 500 with the body "Server Error" — no
herdr code, no herdr message. fleet_spawn and fleet_stop catch the same
exception and report a named error, so this was the same one-door-guarded
shape as #297.

Both now go through the existing herdrError mapper: 404 when herdr says the
target is gone, 502 otherwise. That is an answer the caller can act on.

stopMember matters more than spawnMember. SessionManager.release deregisters
the session, notifies the release listener and preserves a dirty worktree
before it calls launcher.stop, so a throw from that stop arrives after the
teardown the caller asked for has already happened. A bare 500 told the caller
to retry and carried nothing to explain what went wrong.

The two existing tests that asserted 500 now assert 502 and check the error
body. Neither was about the status code: one guards that a failed teardown is
not reported as a successful 204, the other that a failed spawn still closes
its tab. Both properties are unchanged.

Fixes #304
2026-09-04 12:48:21 +07:00
Dai Ha 21844b54d7 #302: fleet_reply refuses blank content too, matching fleet_send
CI / contract (push) Successful in 1m15s
CI / build (push) Successful in 3m34s
MessageService.reply now throws on blank content. FleetMcp.reply guarded only
against null, and its handler is a bare BiFunction with no try/catch, so a
whitespace-only fleet_reply left the handler as an uncaught
IllegalArgumentException instead of the clean tool error null already got.
fleet_send has always used isBlank here; reply now matches it.

The worker found that null/isBlank difference and reported it as a correction
to my ticket, which had quoted the guard wrongly. It was right: I grepped the
error string and assumed the condition matched its sibling.
2026-09-04 12:36:23 +07:00
Dai Ha 4769481515 Merge #302: a reply with no content is refused, not silently delivered
REST read content with .path("content").asText(""), so a body missing the key
became an empty string that resolved the lead's waiter. The turn completed and
the lead saw a member that finished and reported nothing, indistinguishable
from one that genuinely said nothing. The guard went into MessageService.reply,
which both doors call, rather than being written a second time in FleetApp.
2026-09-04 12:33:47 +07:00
Dai Ha bbbb4c1eb3 fleetd #302: require content in MessageService.reply so a REST reply with no content cannot silently resolve a waiter
CI / build (pull_request) Successful in 1m44s
CI / contract (pull_request) Successful in 3m17s
2026-09-04 12:31:09 +07:00
Dai Ha 38c248e617 Merge #298: a released reply is requeued, not dropped
CI / contract (push) Successful in 1m27s
CI / build (push) Successful in 1m46s
release() cancelled the consumer and dropped its local record of deliveries the
broker still held as outstanding. Cancelling a consumer does not requeue them:
they stay unacked on the still-open channel until it or the connection closes.
So a held-but-undrained reply became permanently unreachable — a worker's
report lost with no error and no log line. It now nacks with requeue, after
cancelling, so a later own() can still receive it.
2026-09-04 12:21:30 +07:00
Dai Ha 9debc0de27 #297: render both profiles doors from one body builder, not two copies
CI / contract (push) Successful in 1m23s
CI / build (push) Successful in 1m31s
The merged change shared the QuarantineSource and OutageSource instances
between fleet_profiles and GET /profiles, so the two doors read identical
facts. It then rendered those facts through a character-for-character copy of
the loop, in a different file. Shared inputs do not make duplicated computation
safe: a later edit to the row shape lands on one door and not the other, and
the two disagree about a live outage. That is what #284 was.

The ticket caused this. It said 'read from the same shared instances' and 'do
not change FleetMcp', and together those made copying the loop the only legal
move. Extracting FleetMcp.profilesView and calling it from both is what the
ticket should have asked for.
2026-09-04 12:19:02 +07:00
Dai Ha f34361b263 Merge #297: REST reports what MCP reports
GET /agents and GET /members now map a herdr transport failure into the same
{error, detail} envelope every other handler in FleetApp uses, instead of
letting it escape to Javalin's default handling. GET /profiles now reports the
quarantined and coolingOff states, read from the same shared sources FleetMcp
reads. REST is the door a lead falls back to when its MCP mount drops, so it
was weakest exactly when it was load-bearing.
2026-09-04 12:16:00 +07:00
Dai Ha 85c90d440a #297: map HerdrException on GET /agents and /members; GET /profiles reports quarantine + cool-off
CI / contract (pull_request) Successful in 45s
CI / build (pull_request) Successful in 1m33s
Two REST-only visibility gaps, both against the same shared instances FleetMcp
reads (BackendQuarantine/BackendOutagePolicy), never recomputed:

- GET /agents and GET /members let a HerdrException escape uncaught, outside
  the {error, detail} envelope every other failure path in FleetApp uses.
  Both now route through the existing herdrError() helper, matching healthz/
  sessionStatus. GET /members is the endpoint's own comment names as the
  out-of-band path a lead falls back to when its MCP mount drops.
- GET /profiles omitted the two outage states fleet_profiles already reports:
  quarantined (CB-578 stage B) and coolingOff (fleetd #201 Unit 5). FleetApp
  now takes the SAME FleetMcp.QuarantineSource/OutageSource instances Fleetd
  wires into FleetMcp (extracted to local vars in Fleetd.java so both doors
  share one object, not two independently-built copies of the same rule).

FleetMcp itself is unchanged. Item 3 of the ticket (a capacity block on
GET /members) is explicitly out of scope and was not added.
2026-09-04 12:12:35 +07:00
Dai Ha e97502d550 drainAll: record that the timeout is a whole-drain budget, not a per-session grace
CI / build (push) Successful in 1m27s
CI / contract (push) Successful in 1m48s
The javadoc said 'for each session that is BUSY, poll up to timeoutNanos',
which reads as a per-session grace period. The deadline is taken once, before
the loop, so the first BUSY session can spend all of it. That is deliberate and
is the safer of the two designs: the drain is one phase of a shutdown sequence
that must finish inside launchd's exit window, and a per-session grace would
overrun it and get the daemon SIGKILLed part-way through, leaving the sessions
not yet reached with no clean release, no preserved-worktree log and no
snapshot. Found by a read-only hunt that read the code correctly and drew the
opposite conclusion from the wording.
2026-09-04 12:09:44 +07:00
Dai Ha aef14ff46e Merge #296: a failed spawn no longer leaks the pane it opened
Two exits created a pane and left it running. The readiness gate propagated an
unrelated herdr error without teardown, and spawnAsPane never closed the pane it
split when the peer failed to start. Neither could be cleaned up by the caller:
SessionManager.acquire never learns the pane id, because spawn throws before it
returns one. It removed the worktree anyway, so the leak was a live backend with
a deleted cwd, invisible to fleet_list and holding a seat nothing decremented.
2026-09-04 12:08:23 +07:00
Dai Ha cba516bda4 fleetd #296: close panes on failed spawn
CI / contract (pull_request) Successful in 1m38s
CI / build (pull_request) Successful in 2m7s
2026-09-04 12:05:26 +07:00
31 changed files with 1780 additions and 150 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).
@@ -625,6 +625,19 @@ public final class Fleetd {
log.info("auth: loopback-trust (any loopback non-worker caller is the primary)");
}
// fleetd #297: named once and reused verbatim below for FleetApp's GET /profiles, rather than
// built a second time — two independently-constructed sources reading the SAME BackendQuarantine
// / BackendOutagePolicy would still be able to drift (e.g. a future edit to the credentialIdFor
// closure in only one of the two places), exactly the shape #284 was.
FleetMcp.QuarantineSource quarantineSource = new FleetMcp.QuarantineSource(profile -> {
var configured = config.get().profiles().get(profile);
return configured == null ? null : configured.effectiveCredentialId();
}, quarantine);
FleetMcp.OutageSource outageSource = new FleetMcp.OutageSource(profile -> {
var configured = config.get().profiles().get(profile);
return configured == null ? null : configured.effectiveCredentialId();
}, outagePolicy);
FleetMcp mcp = new FleetMcp(messages, workers, sessions, identity, presence,
primaryRegistry, callers, metrics, new FleetMcp.CapacitySource(profile -> liveCountRef.get().apply(profile),
profile -> {
@@ -636,15 +649,9 @@ public final class Fleetd {
return FleetHealthMonitor.coverage(health != null && health.isEnabled(),
health != null && health.notifications() != null && health.notifications().configured());
}),
new FleetMcp.QuarantineSource(profile -> {
var configured = config.get().profiles().get(profile);
return configured == null ? null : configured.effectiveCredentialId();
}, quarantine),
quarantineSource,
leadMailbox,
new FleetMcp.OutageSource(profile -> {
var configured = config.get().profiles().get(profile);
return configured == null ? null : configured.effectiveCredentialId();
}, outagePolicy),
outageSource,
new FleetMcp.LeadSeatSource(leadSeatLookup(() -> config.get().profiles(), leaders, leads)));
// CB-637: the receive half. Only constructed when a lead mailbox actually opened — with no
@@ -715,9 +722,12 @@ public final class Fleetd {
// GET /sessions must merge across both, or a down/unpolled member daemon is invisible.
// fleetd #111: live (re-read-per-request) memberCredentials view for GET /member-credentials —
// same hot-reload shape as the memberCredentials supplier passed to ClaudeCodeLauncher above.
// fleetd #297: quarantineSource/outageSource are the SAME instances passed to FleetMcp above —
// GET /profiles must report the identical quarantine/cool-off facts as fleet_profiles.
Javalin app = new FleetApp(herdr, memberHerdr, workers, sessions, messages, presence, mcp.servlet(),
callers, metrics, deliverable,
() -> MemberCredentialPolicyView.of(config.get().memberCredentials())).build();
() -> MemberCredentialPolicyView.of(config.get().memberCredentials()),
quarantineSource, outageSource).build();
app.start(cfg.bind().host(), cfg.bind().port());
log.info("fleetd listening on {}:{}, herdr socket {}",
cfg.bind().host(), cfg.bind().port(), socket);
@@ -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) {
@@ -262,11 +270,15 @@ public final class CallerResolver {
return token.isEmpty() ? null : token;
}
/**
* fleetd #305: delegates to {@link ConnectionIdentity#isLoopback}. This used to be a second,
* independent copy of the same rule, and the two drifted: this one accepted all of
* {@code 127.0.0.0/8}, {@code ConnectionIdentity}'s accepted only {@code 127.0.0.1}. A caller
* from {@code 127.0.0.2} therefore had its identity skipped (so it had no terminal) and was
* then read as loopback here — which under loopback-trust is the primary. Sharing the inputs
* would not have prevented that; only sharing the computation does.
*/
private static boolean isLoopback(String remoteAddr) {
if (remoteAddr == null) {
return false;
}
return remoteAddr.equals("127.0.0.1") || remoteAddr.equals("::1")
|| remoteAddr.equals("0:0:0:0:0:0:0:1") || remoteAddr.startsWith("127.");
return ConnectionIdentity.isLoopback(remoteAddr);
}
}
@@ -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());
}
}
@@ -157,6 +157,7 @@ public final class Injector {
boolean awaitingCompletion; // a delivered message's turn is not yet known-complete
boolean turnObserved; // saw a real `working` sample since that delivery (turn ran)
int unknownSinceTurn; // consecutive `unknown` samples while a delegation is outstanding (CB-109)
int unknownSincePostTurn; // the same, for the post-turn housekeeping phase (fleetd #306)
int notReadySincePoll; // consecutive injectable samples a queued message waited on the readiness gate (CB-114)
boolean postTurnPending; // completion observed; adapter housekeeping has not started yet
boolean awaitingPostTurnPickup;
@@ -217,10 +218,12 @@ public final class Injector {
t.awaitingPickup = false;
t.injectableSincePickup = 0;
t.unknownSinceTurn = 0;
t.unknownSincePostTurn = 0;
t.notReadySincePoll = 0;
if (t.awaitingCompletion) t.turnObserved = true;
} else if (status.injectable()) { // IDLE or BLOCKED
t.unknownSinceTurn = 0;
t.unknownSincePostTurn = 0;
if (t.awaitingPostTurnPickup) {
if (++t.injectableSincePostTurnPickup >= PICKUP_GRACE_POLLS) {
t.awaitingPostTurnPickup = false;
@@ -315,6 +318,24 @@ public final class Injector {
t.unknownSinceTurn = 0;
turnFailed = true;
}
// fleetd #306: the same escape for the post-turn housekeeping phase. Four latches
// gate delivery (awaitingCompletion, postTurnPending, awaitingPostTurnPickup,
// postTurnObserved) and only the first had a way out of a sustained unknown streak —
// a gate that closed one direction only. The other two below are released here as
// well; postTurnPending needs no escape because it is cleared unconditionally on the
// line after the listener call that sets it.
//
// This does NOT set turnFailed. The delegated turn already completed and its waiter
// already resolved — what is outstanding is adapter housekeeping (the `/clear`).
// Reporting a turn failure here would drive SessionManager.onFailed on a session
// that genuinely finished its work, which is a worse lie than the wedge.
if ((t.awaitingPostTurnPickup || t.postTurnObserved)
&& ++t.unknownSincePostTurn >= TURN_STALL_GRACE_POLLS) {
t.awaitingPostTurnPickup = false;
t.postTurnObserved = false;
t.injectableSincePostTurnPickup = 0;
t.unknownSincePostTurn = 0;
}
}
// Reclaim the entry once the worker is fully quiescent (nothing queued, no pickup or
@@ -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. */
@@ -60,7 +79,44 @@ public final class ConnectionIdentity {
return pid > 0 ? cwds.cwdForPid(pid) : null;
}
private static boolean isLoopback(String addr) {
return "127.0.0.1".equals(addr) || "::1".equals(addr) || "0:0:0:0:0:0:0:1".equals(addr);
/**
* Whether {@code addr} is a same-host address, and therefore one whose peer PID is worth
* looking up. <strong>This is the one definition of loopback in the daemon</strong> —
* {@code CallerResolver} calls it rather than keeping its own, because the two used to differ
* and that difference was a privilege escalation (fleetd #305).
*
* <p>The whole of {@code 127.0.0.0/8} counts, not just {@code 127.0.0.1}. On Linux every
* address in that range is bound to {@code lo} by default, so a process can connect to
* {@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>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) {
return false;
}
String a = addr.startsWith("::ffff:") ? addr.substring(7) : addr; // IPv4-mapped IPv6
return a.startsWith("127.") || "::1".equals(a) || "0:0:0:0:0:0:0:1".equals(a);
}
}
@@ -22,6 +22,7 @@ import dev.ltms.fleet.placement.BackendQuarantine;
import dev.ltms.fleet.placement.PlacementException;
import dev.ltms.fleet.session.SessionManager;
import dev.ltms.fleet.session.MemberSession;
import dev.ltms.fleet.session.ShuttingDownException;
import dev.ltms.fleet.session.WorktreeRequest;
import dev.ltms.fleet.peer.MemberRole;
import dev.ltms.fleet.peer.PeerLauncher;
@@ -813,7 +814,12 @@ public final class FleetMcp {
return error("fleet_reply is for workers only — could not identify the calling worker "
+ "from the connection");
}
if (content == null) {
// fleetd #302: isBlank, not == null, to match fleet_send's own guard above. MessageService
// .reply now REJECTS blank content, and this handler is a bare BiFunction with no try/catch
// around it — so a whitespace-only fleet_reply would leave here as an uncaught
// IllegalArgumentException instead of this clean tool error. Null and whitespace are the
// same mistake by the caller and must get the same answer.
if (isBlank(content)) {
return error("content is required");
}
messages.reply(callerTerminal, content);
@@ -957,6 +963,10 @@ public final class FleetMcp {
return text(json(memberView(member)));
} catch (GuardException e) {
return error("subscription boundary: " + e.getMessage());
} catch (ShuttingDownException e) {
// fleetd #308: the daemon's shutdown drain has already started — refuse loudly rather
// than register a session drainAll will never see again.
return error("shutting down: " + e.getMessage());
} catch (PlacementException e) {
// CB-599: no candidate had capacity (maxLoad, quarantine, or all-exhausted) — distinct
// from "profile does not exist" below.
@@ -1015,6 +1025,25 @@ public final class FleetMcp {
* both maps at once when it is both exhaustion-quarantined AND cooling off.
*/
static McpSchema.CallToolResult profiles(PeerLauncher workers, QuarantineSource quarantine, OutageSource outage) {
return text(json(profilesView(workers, quarantine, outage)));
}
/**
* The body both front doors answer {@code profiles} with: the configured profile names, the
* default, and the two independent outage states — {@code quarantined} (the backend reported it
* out of capacity) and {@code coolingOff} (the credential threw repeated non-exhaustion backend
* errors). Each map is present only when at least one profile is in that state, and a profile
* can appear in both at once, because the two checks are separate.
*
* <p>fleetd #297: extracted so {@code fleet_profiles} and {@code GET /profiles} render from ONE
* body builder rather than two copies. Passing both doors the same {@link QuarantineSource} and
* {@link OutageSource} instances is necessary but not sufficient: with the loop written out
* twice, a later edit to the row shape — a renamed key, an added field — lands on one door and
* not the other, and the two then disagree about a live outage. That is exactly what fleetd
* #284 was, where one rule computed in two places was widened in only one and a single response
* contradicted itself. Shared inputs do not make duplicated computation safe.
*/
public static Map<String, Object> profilesView(PeerLauncher workers, QuarantineSource quarantine, OutageSource outage) {
Map<String, Object> result = new LinkedHashMap<>();
result.put("profiles", workers.profiles());
result.put("default", workers.defaultProfile() == null ? "" : workers.defaultProfile());
@@ -1046,7 +1075,7 @@ public final class FleetMcp {
if (!coolingOff.isEmpty()) {
result.put("coolingOff", coolingOff);
}
return text(json(result));
return result;
}
/**
@@ -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));
}
/**
@@ -697,7 +697,20 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
if (paneId == null) {
throw new IllegalStateException("pane.split returned no pane — cannot start a peer");
}
Agent peer = startUniquelyNamed(cfg, argv, paneId).agent();
Agent peer;
try {
peer = startUniquelyNamed(cfg, argv, paneId).agent();
} catch (RuntimeException e) {
// The peer never started — don't leave the pane we just created orphaned.
// Best-effort cleanup; never let it mask the real spawn failure.
try {
stop(paneId);
} catch (RuntimeException cleanup) {
log.warn("failed to close orphaned pane {} after spawn error: {}",
paneId, cleanup.getMessage());
}
throw e;
}
log.info("{} started pane={} terminal={}", namePrefix, peer.paneId(), peer.terminalId());
return peer;
}
@@ -996,7 +1009,8 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
* that gap: it stops waiting immediately (never burns the rest of the timeout), runs the same
* teardown the timeout path below runs, and throws with a message that says the backend exited
* rather than that the pane was slow. Any other {@link HerdrException} still propagates
* unchanged — this gate does not know how to recover from it.
* unchanged — this gate does not interpret or recover from it, but it still closes the pane
* it opened before handing the exception to its caller.
*/
private void waitUntilInjectableOrThrow(String paneId) {
long start = nowMillis.getAsLong();
@@ -1010,7 +1024,15 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
if (isAlreadyGone(e)) {
failFastOnGoneBackend(paneId, e, nowMillis.getAsLong() - start);
}
throw e; // any other herdr failure is not ours to interpret — let it propagate
// This gate must not interpret an unrelated herdr error, but the caller does not
// receive paneId when spawn throws. Close the pane here before propagating e unchanged.
try {
stop(paneId);
} catch (RuntimeException cleanup) {
log.warn("failed to close orphaned pane {} after readiness-gate error: {}",
paneId, cleanup.getMessage());
}
throw e;
}
lastStatus = sample.status();
if (lastStatus.injectable() || refinedInjectable(paneId, sample)) {
@@ -187,6 +187,24 @@ public final class MessageService {
private volatile Long completedNanos;
private volatile Reply question;
private volatile String turnId;
/**
* Set when this task's {@code fleet_ask} lapsed with no answer (fleetd #307):
* {@link #clearAsyncQuestion} then forgets {@link #turnId} (nulls it and drops the task from
* {@code asyncTasksByTurn}) so {@link #hasAsyncQuestion} stops reporting the target BUSY — a
* later {@code fleet_send} to it must be accepted, not refused. But the worker's turn is
* still genuinely live: it resumed on its own and will eventually call its real
* {@code fleet_reply}. Losing {@link #turnId} loses {@link #askAnsweredAsyncTasks}' only
* signal that such a reply belongs to this task, so that reply used to fall straight to the
* inbox and strand — {@code fleet_poll} stayed {@code PENDING} forever, later force-failed by
* {@link #abandon} with the misleading "session released before it replied". This flag is a
* second, independent signal that survives the forgetting: {@link #askAnsweredAsyncTasks}
* accepts it in place of a live {@link #turnId}, without ever re-adding the task to
* {@code asyncTasksByTurn} (so the BUSY release is untouched). Cleared implicitly once
* {@link #future} resolves — every match in {@link #askAnsweredAsyncTasks} already requires
* {@code !future.isDone()}, so a task that recovered (or was later failed by
* {@link #abandon}) can never match again regardless of this flag's value.
*/
private volatile boolean askTimedOut;
private Task(String ticket, String target, LongSupplier nowNanos) {
this.ticket = ticket;
@@ -395,36 +413,55 @@ public final class MessageService {
* {@link Rendezvous#resolveQuestion} must keep today's {@code NO_WAITER} behaviour — questions
* are interactive and must never be queued.
*
* <p><strong>Ambiguous match also falls to the inbox.</strong> {@link #askAnsweredAsyncTasks}
* cannot actually return more than one entry today (see its own javadoc for why — in short,
* {@link #hasAsyncQuestion} keeps a target BUSY, so no second task can reach this state, for as
* long as an earlier one's {@code turnId} is still stamped). That is an emergent guarantee from
* two other facts, not one this method enforces, so this branch stays in as defence in depth
* rather than being removed as dead code: if it ever weakens, returning whichever candidate a
* {@code ConcurrentHashMap} iteration reaches first would let a genuine reply complete the
* <em>wrong</em> ticket — silently handing the lead something that reads like a correct answer to
* a delegation the worker never touched, which is worse than a failure because the lead acts on
* it. When more than one candidate exists, guessing is not safe: fall back to the inbox exactly
* as the zero-candidate case does, and let {@link #abandon} apply the eventual recovery
* deterministically instead.
* <p><strong>Ambiguous match also falls to the inbox.</strong> {@link #askAnsweredAsyncTasks} can
* return more than one entry — a reachable state, not a hypothetical one (see its own javadoc:
* an {@code fleet_ask} that lapsed with no answer, fleetd #307, frees the target for a completely fresh
* delegation, which can itself go on to ask-and-lapse before the first worker's real reply
* arrives). Returning whichever candidate a {@code ConcurrentHashMap} iteration reaches first
* would let a genuine reply complete the <em>wrong</em> ticket — silently handing the lead
* something that reads like a correct answer to a delegation the worker never touched, which is
* worse than a failure because the lead acts on it. When more than one candidate exists, guessing
* is not safe: fall back to the inbox exactly as the zero-candidate case does, and let
* {@link #abandon} apply the eventual recovery deterministically instead.
*
* <p><strong>{@code content} is required (fleetd #302).</strong> Both doors that reach this
* method must reject a missing/blank reply the same way, so the check lives here rather than in
* either caller: {@code FleetMcp.reply} already refuses a {@code null} content before it ever
* calls this method (its own required-arg guard), and no test or production call site anywhere
* in the codebase relies on replying with empty content — confirmed by searching every call site
* of this method before adding the check, not assumed. Without this guard, a REST {@code
* POST /sessions/{id}/reply} whose body omits {@code content} (or a client library that maps a
* missing field to {@code ""}) used to reach {@link Rendezvous#resolve} with an empty string,
* silently completing the lead's blocking wait with nothing — indistinguishable from a worker
* that genuinely replied with nothing, which is worse than a loud failure because it destroys the
* information that the reply never arrived.
*
* @throws IllegalArgumentException if {@code content} is {@code null} or blank — the caller must
* report this as a client error (REST: 400 {@code bad_request}) rather than resolve
* anything
* @return always {@code true} — the reply resolved a live send, completed a parked ticket, or
* was queued
*/
public boolean reply(String session, String content) {
if (content == null || content.isBlank()) {
throw new IllegalArgumentException("content is required");
}
if (rendezvous.resolve(session, content)) {
count(FleetMetrics.REPLIES, "path", "rendezvous");
return true; // a live send took it — unchanged fast path
}
// #137: no live rendezvous waiter, but this may be the worker's real fleet_reply resuming a
// turn that {@link #answer} already gave up waiting on. answer()'s own bounded wait (the
// primary's fleet_send{turnId} call, capped well under a minute) can time out and close its
// waiter long before the worker — now actually resuming real work — finishes and replies. That
// reply used to have nowhere to land but the session inbox, leaving the async ticket's future
// unresolved forever: fleet_poll{ticket} stayed PENDING until fleet_stop's abandon() forced it
// FAILED with a misleading "session released before it replied" reason, even though the reply
// had, in fact, arrived. Completing the matching ticket directly here means fleet_poll{ticket}
// sees the real reply instead.
// #137/fleetd #307: no live rendezvous waiter, but this may be the worker's real fleet_reply resuming
// a turn that either answer() (#137) or ask() (fleetd #307) already gave up waiting on:
// - answer()'s own bounded wait (the primary's fleet_send{turnId} call, capped well under a
// minute) can time out and close its waiter long before the worker — now actually resuming
// real work — finishes and replies.
// - ask()'s own wait for the primary can time out first, with the worker resuming on its own
// and finishing unanswered.
// Either way that reply used to have nowhere to land but the session inbox, leaving the async
// ticket's future unresolved forever: fleet_poll{ticket} stayed PENDING until fleet_stop's
// abandon() forced it FAILED with a misleading "session released before it replied" reason,
// even though the reply had, in fact, arrived. Completing the matching ticket directly here
// means fleet_poll{ticket} sees the real reply instead.
List<Task> candidates = askAnsweredAsyncTasks(session);
if (candidates.size() == 1) {
Task orphan = candidates.get(0);
@@ -455,34 +492,42 @@ public final class MessageService {
}
/**
* Every still-open async task on {@code target} whose {@code fleet_ask} was already answered —
* its {@link Task#turnId} is stamped but its {@link Task#question} was cleared by {@link #answer}
* — yet whose future is not resolved yet (#137). Empty if no such task exists, including the
* common case where {@code target}'s worker never used {@code fleet_ask} at all (a task that was
* never asked has {@code turnId == null}, so it can never match here and only ever completes
* through the ordinary rendezvous fast path in {@link #reply}).
* Every still-open async task on {@code target} whose worker is genuinely expected to send a
* real {@code fleet_reply} next with nothing left registered to catch it: either its
* {@code fleet_ask} was already answered — {@link Task#turnId} is stamped but {@link
* Task#question} was cleared by {@link #answer} — or its {@code fleet_ask} lapsed unanswered and
* {@link Task#askTimedOut} marks that (fleetd #307; {@link Task#turnId} is {@code null} by then, forgotten
* so the target is not left BUSY — see {@link Task#askTimedOut}'s own javadoc). Either way the
* task's future is not resolved yet. Empty if no such task exists, including the common case
* where {@code target}'s worker never used {@code fleet_ask} at all (a task that was never asked
* has both {@code turnId == null} and {@code askTimedOut == false}, so it can never match here and
* only ever completes through the ordinary rendezvous fast path in {@link #reply}).
*
* <p><strong>Returns at most one entry today — verified, not assumed.</strong> {@link #send}
* refuses to open a waiter on {@code target} while {@link #hasAsyncQuestion} is true, and that
* check matches ANY task whose {@code turnId} is still stamped in {@code asyncTasksByTurn} —
* not only while its question is still open. {@link #answer} deliberately leaves that stamp in
* place ({@code clearAsyncQuestion(turnId, false)}) until the resumed turn's own future actually
* resolves, at which point {@link #finishAsyncTask} both removes the stamp AND completes that
* task's future in the same call. So a second task can never reach "{@code turnId} stamped, future
* still open" — the exact pair this method matches on — while a first one already holds it: by
* the time the stamp is gone, so is the eligibility. This is an emergent property of those two
* facts holding together, not something this method (or its callers) enforces on its own — flip
* {@code forgetTurn} to {@code true} in that one {@link #answer} call and it silently stops being
* true, with nothing left to fail loudly. The callers below still handle "more than one" as
* defence in depth against exactly that, not because they exercise it today: {@link #reply}
* treats it as unresolvable and falls back to the inbox; {@link #abandon} would pick the oldest
* deterministically (its own {@code matching} list has no such guarantee — see its javadoc).
* <p><strong>Can return more than one entry — reachable, not just defence in depth.</strong>
* {@link #send} refuses to open a waiter on {@code target} while {@link #hasAsyncQuestion} is
* true, and that check matches ANY task whose {@code turnId} is still stamped in
* {@code asyncTasksByTurn}. While a task's {@code turnId} stays stamped — {@link #answer} leaves
* it in place ({@code clearAsyncQuestion(turnId, false)}) until {@link #finishAsyncTask} removes
* the stamp and completes the future in the same call — no second task on the same target can
* reach an eligible state, because {@link #send} would refuse it as BUSY first. That single-task
* guarantee holds only for the {@code turnId}-stamped half of this method's match: an
* {@link Task#askTimedOut} task is, by construction, no longer stamped in {@code asyncTasksByTurn}
* (that is the whole point of forgetting {@code turnId} in {@link #clearAsyncQuestion}), so the
* target is free the moment one ask lapses. A fresh, independent {@code sendAsync} to the same
* target can then be dispatched, itself pause on {@code fleet_ask}, and itself time out — landing
* a second {@code askTimedOut} task on the very target the first one is still waiting to answer
* for. Two (or more) genuinely open tasks on one target is therefore a real, reachable state
* today, not a hypothetical: {@link #reply} treats it as unresolvable and falls back to the
* inbox rather than guess which task a reply belongs to (guessing wrong would hand the lead a
* plausible-looking answer to a delegation the worker never touched — worse than a failure,
* because the lead acts on it); {@link #abandon} instead picks the oldest deterministically (its
* own {@code matching} list has a different, wider match — see its javadoc).
*/
private List<Task> askAnsweredAsyncTasks(String target) {
List<Task> candidates = new ArrayList<>();
for (Task task : tasks.values()) {
if (target.equals(task.target) && task.question == null && task.turnId != null
&& !task.future.isDone()) {
if (target.equals(task.target) && task.question == null && !task.future.isDone()
&& (task.turnId != null || task.askTimedOut)) {
candidates.add(task);
}
}
@@ -880,6 +925,12 @@ 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);
return new AskResult(AskOutcome.TIMED_OUT, null);
} catch (ExecutionException e) {
@@ -1144,6 +1195,23 @@ public final class MessageService {
return task;
}
/**
* Mark {@code turnId}'s task as having a {@code fleet_ask} that lapsed with no answer (fleetd #307), so
* {@link #askAnsweredAsyncTasks} still recognizes the worker's eventual real {@code fleet_reply}
* as belonging to it after {@link #clearAsyncQuestion}'s {@code forgetTurn=true} erases
* {@link Task#turnId} — see {@link Task#askTimedOut}. Must be called before that forgetting, while
* {@code turnId} can still resolve the task in {@code asyncTasksByTurn}; a lookup afterward would
* find nothing. Only when it matches the task's current turn — same guard as
* {@link #clearAsyncQuestion} — so a chained second {@code fleet_ask} (#282) that already moved
* the task to a fresh {@code turnId} cannot mark it for a turn that is no longer its own.
*/
private void markAskTimedOut(String turnId) {
Task task = asyncTasksByTurn.get(turnId);
if (task != null && turnId.equals(task.turnId)) {
task.askTimedOut = true;
}
}
/** Clear an answered or lapsed question, but only when it matches the ticket's current turn. */
private void clearAsyncQuestion(String turnId, boolean forgetTurn) {
// CB-582: tell the push loop first — like ticketCollected, a removal for a turnId it never
@@ -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");
}
@@ -9,6 +9,7 @@ import dev.ltms.fleet.auth.Principal;
import dev.ltms.fleet.guard.GuardException;
import dev.ltms.fleet.metrics.FleetMetrics;
import dev.ltms.fleet.metrics.Metrics;
import dev.ltms.fleet.mcp.FleetMcp;
import dev.ltms.fleet.herdr.Agent;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.herdr.HerdrException;
@@ -18,6 +19,7 @@ import dev.ltms.fleet.peer.PeerUnreachableException;
import dev.ltms.fleet.placement.PlacementException;
import dev.ltms.fleet.msg.MessageService;
import dev.ltms.fleet.session.SessionManager;
import dev.ltms.fleet.session.ShuttingDownException;
import dev.ltms.fleet.peer.MemberRole;
import dev.ltms.fleet.session.MemberSession;
import dev.ltms.fleet.session.WorktreeRequest;
@@ -86,6 +88,12 @@ public final class FleetApp {
// absent() (the honest "no policy configured" view) for every constructor that does not wire
// a real one, so existing legacy call sites keep building without knowing this field exists.
private final Supplier<MemberCredentialPolicyView> memberCredentials;
// fleetd #297: the SAME shared sources FleetMcp.profiles/fleet_profiles reads (BackendQuarantine
// and BackendOutagePolicy are each one instance for the whole daemon — see Fleetd wiring) so
// GET /profiles cannot drift from fleet_profiles about which profile is quarantined/cooling off.
// .none() (the honest "feature not wired" view) for every constructor that does not pass one.
private final FleetMcp.QuarantineSource quarantine;
private final FleetMcp.OutageSource outage;
private final ObjectMapper mapper = new ObjectMapper();
/**
@@ -146,6 +154,23 @@ public final class FleetApp {
MessageService messages, MemberPresence presence,
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics,
Predicate<String> deliverable, Supplier<MemberCredentialPolicyView> memberCredentials) {
this(herdr, memberHerdr, workers, sessions, messages, presence, mcpServlet, auth, metrics,
deliverable, memberCredentials, FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none());
}
/**
* @param quarantine the SAME {@link FleetMcp.QuarantineSource} instance passed to {@code
* FleetMcp} (fleetd #297), so {@code GET /profiles} reports the identical
* exhaustion-quarantine facts as {@code fleet_profiles} rather than a second,
* independently-computed copy
* @param outage the SAME {@link FleetMcp.OutageSource} instance passed to {@code FleetMcp} —
* see {@code quarantine}; a SEPARATE check from it, never merged in
*/
public FleetApp(HerdrClient herdr, HerdrClient memberHerdr, PeerLauncher workers, SessionManager sessions,
MessageService messages, MemberPresence presence,
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics,
Predicate<String> deliverable, Supplier<MemberCredentialPolicyView> memberCredentials,
FleetMcp.QuarantineSource quarantine, FleetMcp.OutageSource outage) {
this.herdr = herdr;
this.memberHerdr = memberHerdr != null ? memberHerdr : herdr;
this.workers = workers;
@@ -156,6 +181,8 @@ public final class FleetApp {
this.auth = auth;
this.metrics = metrics;
this.memberCredentials = memberCredentials != null ? memberCredentials : MemberCredentialPolicyView::absent;
this.quarantine = quarantine != null ? quarantine : FleetMcp.QuarantineSource.none();
this.outage = outage != null ? outage : FleetMcp.OutageSource.none();
}
/** Wire routes onto a fresh, unstarted Javalin instance. Caller starts it. */
@@ -338,8 +365,15 @@ public final class FleetApp {
if (!allow(ctx, routeAction("GET /agents"), null)) {
return;
}
ctx.status(200).json(Map.of("agents",
workers.list().stream().map(Agent.class::cast).map(FleetApp::view).toList()));
try {
ctx.status(200).json(Map.of("agents",
workers.list().stream().map(Agent.class::cast).map(FleetApp::view).toList()));
} catch (HerdrException e) {
// fleetd #297: workers.list() reaches herdr — a transport failure must land in the same
// {error, detail} envelope every other failure path here uses, not escape as a bare
// exception and leave Javalin's default handling to respond outside the JSON contract.
herdrError(ctx, e);
}
}
/** CB-304: bridge-owned roster merged with live herdr status by paneId. */
@@ -347,40 +381,57 @@ public final class FleetApp {
if (!allow(ctx, routeAction("GET /members"), null)) {
return;
}
// CB-519: the registry key is a host-unique id, not the pane coordinate — join on terminal.
Map<String, Agent> live = workers.list().stream()
.map(Agent.class::cast)
.filter(a -> a.terminalId() != null)
.collect(Collectors.toMap(Agent::terminalId, Function.identity(), (_, b) -> b));
// fleetd #209: this REST roster reports agentSessionId via SessionManager.rosterView, so it
// uses the resolving roster read (caller-driven, not a timer) rather than the plain one.
List<Map<String, Object>> out = sessions.rosterResolved().stream()
.map(s -> SessionManager.rosterView(s, live.get(s.terminalId())))
.toList();
Map<String, Object> body = new LinkedHashMap<>();
// fleetd #199: the endpoint became /members in the CB-634 rename but the body key stayed
// "workers", so a caller that read "members" saw an empty fleet and reported no members at
// all. "members" is the canonical key; "workers" stays as a deprecated alias so an existing
// REST consumer keeps working — the out-of-band path a lead falls back to when its MCP mount
// drops reads this endpoint. Drop the alias once nothing reads it.
body.put("members", out);
body.put("workers", out);
// CB-586: operator visibility for the refs/wip snapshot store without shelling into the
// repo — how many snapshot refs exist and roughly what they cost. Present only once a
// worktree session has established the repo, so a never-snapshotted fleet reports nothing.
sessions.wipRefs().ifPresent(st -> body.put("wipRefs",
Map.of("count", st.count(), "costBytes", st.costBytes())));
ctx.status(200).json(body);
try {
// CB-519: the registry key is a host-unique id, not the pane coordinate — join on terminal.
Map<String, Agent> live = workers.list().stream()
.map(Agent.class::cast)
.filter(a -> a.terminalId() != null)
.collect(Collectors.toMap(Agent::terminalId, Function.identity(), (_, b) -> b));
// fleetd #209: this REST roster reports agentSessionId via SessionManager.rosterView, so it
// uses the resolving roster read (caller-driven, not a timer) rather than the plain one.
List<Map<String, Object>> out = sessions.rosterResolved().stream()
.map(s -> SessionManager.rosterView(s, live.get(s.terminalId())))
.toList();
Map<String, Object> body = new LinkedHashMap<>();
// fleetd #199: the endpoint became /members in the CB-634 rename but the body key stayed
// "workers", so a caller that read "members" saw an empty fleet and reported no members at
// all. "members" is the canonical key; "workers" stays as a deprecated alias so an existing
// REST consumer keeps working — the out-of-band path a lead falls back to when its MCP mount
// drops reads this endpoint. Drop the alias once nothing reads it.
body.put("members", out);
body.put("workers", out);
// CB-586: operator visibility for the refs/wip snapshot store without shelling into the
// repo — how many snapshot refs exist and roughly what they cost. Present only once a
// worktree session has established the repo, so a never-snapshotted fleet reports nothing.
sessions.wipRefs().ifPresent(st -> body.put("wipRefs",
Map.of("count", st.count(), "costBytes", st.costBytes())));
ctx.status(200).json(body);
} catch (HerdrException e) {
// fleetd #297: same reasoning as agents() above — this is the out-of-band roster a lead
// falls back to when its MCP mount drops, so it must stay inside the JSON error contract
// exactly when herdr is briefly unreachable, not escape as a bare exception.
herdrError(ctx, e);
}
}
/** The configured worker profiles and which one a no-argument spawn uses. */
/**
* The configured worker profiles, which one a no-argument spawn uses, and (fleetd #297) the two
* outage states {@code fleet_profiles} already reports: {@code quarantined} (CB-578 stage B —
* the backend reported it out of capacity) and {@code coolingOff} (fleetd #201 Unit 5 — the
* credential threw repeated non-exhaustion backend errors). Both are read from the SAME shared
* {@link FleetMcp.QuarantineSource}/{@link FleetMcp.OutageSource} instances {@code FleetMcp}
* reads, never recomputed, so the two doors cannot disagree about which profile is down and why.
* Independent checks, so a profile can appear in both maps at once; each map is present only
* when at least one profile is in that state.
*/
private void profiles(Context ctx) {
if (!allow(ctx, routeAction("GET /profiles"), null)) {
return;
}
ctx.status(200).json(Map.of(
"profiles", workers.profiles(),
"default", workers.defaultProfile() == null ? "" : workers.defaultProfile()));
// fleetd #297: ONE body builder, shared with fleet_profiles. Handing both doors the same
// QuarantineSource/OutageSource instances stops them reading different facts; rendering
// through the same method stops them reporting those facts differently. Both are needed.
ctx.status(200).json(FleetMcp.profilesView(workers, quarantine, outage));
}
/**
@@ -451,6 +502,11 @@ public final class FleetApp {
ctx.status(201).json(view(member));
} catch (GuardException e) {
ctx.status(403).json(Map.of("error", "subscription_boundary", "detail", e.getMessage()));
} catch (ShuttingDownException e) {
// fleetd #308: the daemon's shutdown drain has already started — 503, not a bare 500,
// so this reads the same as PlacementException below: valid request, refused because
// of a transient daemon state rather than a bad argument.
ctx.status(503).json(Map.of("error", "shutting_down", "detail", e.getMessage()));
} catch (PlacementException e) {
// CB-599: no candidate had capacity (maxLoad, quarantine, or all-exhausted) — a benign,
// likely-transient refusal, distinct from "profile does not exist" below. 503: the
@@ -460,6 +516,12 @@ public final class FleetApp {
ctx.status(400).json(Map.of("error", "unknown_profile", "detail", e.getMessage()));
} catch (PeerUnreachableException e) {
ctx.status(502).json(Map.of("error", "spawn_timeout", "detail", e.getMessage()));
} catch (HerdrException e) {
// fleetd #304: not every herdr failure on the spawn path is a readiness timeout, so
// PeerUnreachableException above does not cover this. Without this catch the exception
// escapes to Javalin's default 500, while fleet_spawn reports the same failure as a
// clean named error (FleetMcp.spawn) — the #297 one-door-guarded shape.
herdrError(ctx, e);
}
}
@@ -480,13 +542,29 @@ public final class FleetApp {
return (s == null || s.isBlank()) ? null : s;
}
/** Tear a worker down by pane id. */
/**
* Tear a worker down by pane id.
*
* <p>fleetd #304: the {@code HerdrException} catch is not cosmetic. {@code release} deregisters
* the session, notifies the release listener and preserves a dirty worktree <em>before</em> it
* calls {@code launcher.stop}, so a throw from that stop arrives after the teardown the caller
* asked for has already happened. Letting it escape gave Javalin's default 500, which tells the
* caller to retry — and the retry finds nothing in the registry, reaches the same stop, and
* throws again, so it can never succeed. {@code herdrError} instead answers 404 ("the pane is
* gone, stop retrying") or 502 ("herdr is upstream and broken, a retry may help"), matching what
* {@code fleet_stop} reports for the same failure.
*/
private void stopMember(Context ctx) {
String paneId = ctx.pathParam("paneId");
if (!allow(ctx, routeAction("DELETE /members/{paneId}"), paneId)) {
return;
}
sessions.release(paneId);
try {
sessions.release(paneId);
} catch (HerdrException e) {
herdrError(ctx, e);
return;
}
ctx.status(204);
}
@@ -634,7 +712,19 @@ public final class FleetApp {
ctx.status(400).json(Map.of("error", "bad_request", "detail", "body must be JSON"));
return;
}
messages.reply(id, content);
// fleetd #302: content is required. `.path("content").asText("")` above turns a missing key
// into "" rather than throwing, so without this check an empty/blank reply used to reach
// messages.reply(...) and silently resolve the lead's waiter — the same class of bug as the
// sibling "content is required" guards on sendMessage/askMessage below, except this one wrote
// a WRONG value instead of failing loudly. The check lives in MessageService.reply so both
// this door and FleetMcp.reply inherit the same rule; this catch only translates it into the
// {error, detail} envelope this file uses everywhere else.
try {
messages.reply(id, content);
} catch (IllegalArgumentException e) {
ctx.status(400).json(Map.of("error", "bad_request", "detail", e.getMessage()));
return;
}
ctx.status(200).json(Map.of("sessionId", id, "delivered", true));
}
@@ -93,6 +93,9 @@ public final class GitWorktrees implements Worktrees {
/** OS group name for {@link #shareWithGroup} (fleetd #185 stage 3); {@code null} ⇒ feature off. */
private final String group;
private final Consumer<String> afterWorktreeAdded;
/** How the initial {@code git worktree add} command runs. Package-private test seam for an
* interrupted command after Git has made worktree state. */
private final Function<String[], String> worktreeAddRunner;
/** How {@link #shareWithGroup}'s processes (git config / chgrp / chmod / find) actually run.
* Defaults to the real {@link #exec(String...)}. Package-private test seam so a unit test can
* prove "no group configured ⇒ zero processes spawned" and inspect exactly what a configured
@@ -132,7 +135,13 @@ public final class GitWorktrees implements Worktrees {
/** Test seam combining a configurable {@code group} with {@link #afterWorktreeAdded}. */
GitWorktrees(String configuredRoot, String group, Consumer<String> afterWorktreeAdded) {
this(configuredRoot, group, afterWorktreeAdded, null);
this(configuredRoot, group, afterWorktreeAdded, null, null);
}
/** Test seam for changing how {@link #shareWithGroup}'s processes run. */
GitWorktrees(String configuredRoot, String group, Consumer<String> afterWorktreeAdded,
Function<String[], String> shareGroupRunner) {
this(configuredRoot, group, afterWorktreeAdded, shareGroupRunner, null);
}
/**
@@ -141,13 +150,15 @@ public final class GitWorktrees implements Worktrees {
* exactly what commands a configured group runs, without a real second OS user/group.
*
* @param shareGroupRunner {@code null} ⇒ the real {@link #exec(String...)}.
* @param worktreeAddRunner {@code null} ⇒ the real {@link #exec(String...)}.
*/
GitWorktrees(String configuredRoot, String group, Consumer<String> afterWorktreeAdded,
Function<String[], String> shareGroupRunner) {
Function<String[], String> shareGroupRunner, Function<String[], String> worktreeAddRunner) {
this.configuredRoot = configuredRoot;
this.group = (group == null || group.isBlank()) ? null : group;
this.afterWorktreeAdded = afterWorktreeAdded == null ? _ -> {} : afterWorktreeAdded;
this.shareGroupRunner = shareGroupRunner != null ? shareGroupRunner : this::exec;
this.worktreeAddRunner = worktreeAddRunner != null ? worktreeAddRunner : this::exec;
}
@Override
@@ -166,8 +177,8 @@ public final class GitWorktrees implements Worktrees {
String wt = path.toAbsolutePath().toString();
log.info("adding worktree branch={} path={} base={}", branch, wt, base);
removeUserInfoFromHttpsOrigin(repoRoot);
exec("git", "-C", repoRoot, "worktree", "add", wt, "-b", branch, base);
try {
worktreeAddRunner.apply(new String[] {"git", "-C", repoRoot, "worktree", "add", wt, "-b", branch, base});
afterWorktreeAdded.accept(wt);
requireCredentialFreeHttpsOrigin(wt);
configureEnvironmentCredentialHelper(repoRoot, wt);
@@ -181,10 +192,11 @@ public final class GitWorktrees implements Worktrees {
}
/**
* {@code add()} has already created the worktree and its branch by the time any step from
* {@link #afterWorktreeAdded} through {@link #isolateToolSurface} can throw — including
* {@code git worktree add} may have created the worktree and its branch by the time it, or any
* later step through {@link #isolateToolSurface}, throws. This includes
* {@link #requireCredentialFreeHttpsOrigin}, an intended security refusal, not only an IO
* accident. Without this, {@code add()} never returns, so its caller
* accident. A Git-reported {@code worktree add} failure usually creates nothing, but an
* interrupted command can leave partial state. Without cleanup, {@code add()} never returns, so its caller
* ({@code SessionManager#acquireWithWorktree}) never receives a path to register or clean up:
* its local {@code path} stays null, the {@code if (path != null)} guard in its own catch block
* never runs, and the worktree directory and branch leak on disk forever with nothing tracking
@@ -204,6 +216,10 @@ public final class GitWorktrees implements Worktrees {
* used in {@code SessionManager#acquireWithWorktree}'s own catch block.
*/
private void cleanupAfterAddFailure(String repoRoot, String worktreePath, String branch, RuntimeException original) {
if (!Files.exists(Path.of(worktreePath))) {
log.debug("provisioning failed before worktree {} existed; nothing to clean up", worktreePath);
return;
}
log.warn("provisioning failed for branch={} path={}: {} — cleaning up before rethrowing",
branch, worktreePath, original.getMessage());
try {
@@ -21,6 +21,7 @@ import java.util.Optional;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.Consumer;
import java.util.function.LongSupplier;
@@ -61,6 +62,8 @@ public final class SessionManager implements TurnListener {
private final LongSupplier nowNanos;
private final int contextCap;
private final boolean clearAfterTurn;
/** Null in production; test seam for the interval before an idle session's conditional release. */
private final Consumer<MemberSession> beforeIdleRelease;
private volatile MemberLifecycle memberLifecycle = MemberLifecycle.NONE;
/**
* CB-586: the repo root the fleet actually works in, remembered the first time a worktree
@@ -71,6 +74,16 @@ public final class SessionManager implements TurnListener {
*/
private volatile String fleetRepoRoot;
/**
* fleetd #308: flips true the instant {@link #drainAll} starts, before its registry snapshot
* is even taken — so a spawn already in flight sees the refusal as early as a plain flag can
* make it. This alone cannot close the race completely: a caller that read {@code false} just
* before the flip can still land in the registry after the snapshot. {@link #drainAll}'s
* post-loop sweep is what catches that straggler; the two mechanisms are deliberately paired,
* see {@link #drainAll}'s javadoc.
*/
private final AtomicBoolean draining = new AtomicBoolean(false);
/** CB-520: notified with a terminalId on every acquire; no-op until wired. */
private final List<Consumer<String>> acquireListeners = new java.util.concurrent.CopyOnWriteArrayList<>();
/** CB-516: notified with a {@link ReleaseDetail} on every release; no-op until wired. */
@@ -102,13 +115,23 @@ public final class SessionManager implements TurnListener {
}
public SessionManager(PeerLauncher launcher, Worktrees worktrees, LongSupplier nowNanos,
int contextCap, boolean clearAfterTurn) {
int contextCap, boolean clearAfterTurn) {
this(launcher, worktrees, nowNanos, contextCap, clearAfterTurn, null);
}
/**
* Package-private constructor for a deterministic reap/delivery race test. Production callers
* use the constructor above, whose null hook adds no callback or lock to an ordinary reap.
*/
SessionManager(PeerLauncher launcher, Worktrees worktrees, LongSupplier nowNanos,
int contextCap, boolean clearAfterTurn, Consumer<MemberSession> beforeIdleRelease) {
this.launcher = launcher;
this.worktrees = worktrees;
this.presence = new PresenceFleet(this);
this.nowNanos = nowNanos;
this.contextCap = contextCap;
this.clearAfterTurn = clearAfterTurn;
this.beforeIdleRelease = beforeIdleRelease;
}
/**
@@ -181,6 +204,14 @@ public final class SessionManager implements TurnListener {
public MemberSession acquire(String profile, MemberRole role, String requestedCwd, String callerCwd,
String ownerTerminal, WorktreeRequest wt,
String sessionName, String resumeSessionId) {
// fleetd #308: refuse before anything else runs — no slot reservation, no launcher spawn —
// so a caller learns the daemon is going down instead of getting a session drainAll will
// never see again. Checked here because every other acquire(...) overload delegates to
// this one, so this is the single point every spawn path passes through.
if (draining.get()) {
throw new ShuttingDownException("fleetd is shutting down; refusing to spawn a session "
+ "the shutdown drain would never see");
}
MemberRole memberRole = (role == null) ? MemberRole.DEV : role;
requireResumeCapability(profile, resumeSessionId);
// CB-619 / fleetd #123: an explicit profile bypasses placement (CompositePeerLauncher only
@@ -272,10 +303,32 @@ public final class SessionManager implements TurnListener {
*/
private void release(String paneId, ReleaseCause cause) {
MemberSession removed = registry.remove(paneId);
releaseRemoved(paneId, removed, handles.remove(paneId), cause);
}
/**
* Tear a session down only while {@code expected} is still its registry value. A lifecycle
* transition replaces the immutable record, so this prevents a reap based on an old READY or
* DONE record from stopping a worker that delivery has made BUSY.
*/
private boolean releaseIfCurrent(MemberSession expected, ReleaseCause cause) {
if (!registry.remove(expected.paneId(), expected)) {
// A lifecycle transition replaced the record between the caller's check and this remove.
// Log it: this race is by definition unobservable otherwise, and a reaper that silently
// declines to reap is the hardest kind of behaviour to diagnose after the fact.
log.debug("skipping reap of pane={}: its registry record changed after the idle check "
+ "(most likely a delivery made it BUSY)", expected.paneId());
return false;
}
releaseRemoved(expected.paneId(), expected, handles.remove(expected.paneId()), cause);
return true;
}
private void releaseRemoved(String paneId, MemberSession removed, PeerHandle removedHandle,
ReleaseCause cause) {
// fleetd #209: remove right alongside the registry entry so a released session's handle is
// never leaked — but keep the local reference below, so the id can still be resolved for
// the ReleaseDetail this teardown notifies with.
PeerHandle removedHandle = handles.remove(paneId);
boolean preserveWorktree = cause == ReleaseCause.SHUTDOWN;
String snapshotRef = null;
if (removed != null) {
@@ -846,14 +899,18 @@ public final class SessionManager implements TurnListener {
}
long idleNanos = now - s.lastActivityAtNanos();
if (idleNanos > idleTtlNanos) {
log.debug("reaping idle session terminal={} pane={}: idle {}s exceeds the {}s ttl",
s.terminalId(), s.paneId(), TimeUnit.NANOSECONDS.toSeconds(idleNanos),
TimeUnit.NANOSECONDS.toSeconds(idleTtlNanos));
// CB-581: one session that fails to release must not abort the whole reaping pass —
// match drainAll's per-session try/catch so the rest of the roster still gets reaped.
try {
release(s.paneId());
reaped++;
if (beforeIdleRelease != null) {
beforeIdleRelease.accept(s);
}
if (releaseIfCurrent(s, ReleaseCause.COMPLETED)) {
log.debug("reaping idle session terminal={} pane={}: idle {}s exceeds the {}s ttl",
s.terminalId(), s.paneId(), TimeUnit.NANOSECONDS.toSeconds(idleNanos),
TimeUnit.NANOSECONDS.toSeconds(idleTtlNanos));
reaped++;
}
} catch (RuntimeException e) {
log.warn("reap failed for pane={} terminal={} worktree={}; continuing with "
+ "remaining sessions", s.paneId(), s.terminalId(), s.worktree(), e);
@@ -864,20 +921,60 @@ public final class SessionManager implements TurnListener {
}
/**
* Gracefully drain all registered sessions on daemon shutdown. For each session that is
* {@code BUSY}, poll up to {@code timeoutNanos} for it to leave {@code BUSY}, then release it
* regardless. Non-busy sessions are released immediately. A failure releasing one session is
* logged and does not abort the rest.
* Gracefully drain all registered sessions on daemon shutdown. Non-busy sessions are released
* immediately; a {@code BUSY} one is polled until it leaves {@code BUSY}, then released
* regardless. A failure releasing one session is logged and does not abort the rest.
*
* <p>{@code timeoutNanos} is a budget for the WHOLE drain, not a grace period per session: the
* deadline is taken once, before the loop. So the first BUSY session can spend all of it, and a
* later BUSY one is then released with no wait at all. That is deliberate. This drain is only
* one phase of shutdown — {@code Fleetd} closes the message service, the push loop, the
* heartbeat, MCP and the router after it — and the whole sequence has to finish inside
* launchd's exit window. A per-session grace would let N busy members drain for N * the
* timeout, overrun that window, and get the daemon SIGKILLed part-way through; the members not
* yet reached would then get no clean release, no preserved-worktree log, and no snapshot.
* Cutting one turn short is the cheaper failure, and it is not silent: an abandoned BUSY
* session is preserved, snapshotted, and logged at WARN by {@code logPreservedForShutdown}.
*
* <p>CB-544: this is a {@link ReleaseCause#SHUTDOWN} release — the worker's pane is stopped
* (the process must end) but its worktree is preserved and its path logged. Shutdown is never
* a reason to delete a worker's only copy of its uncommitted work. A session still {@code BUSY}
* when the timeout expired is abandoned mid-turn and logged loudly so an operator can find its
* kept worktree.
*
* <p>fleetd #308: {@code roster()} is a one-shot snapshot (see its javadoc), and nothing used
* to stop a new session from registering after it was taken — {@link #acquire} stayed open for
* as long as this drain waited on a {@code BUSY} session, up to the whole {@code timeoutNanos}
* budget. Two things close that window, deliberately paired because neither alone is complete:
* {@link #draining} is flipped true before the snapshot is even taken, so {@link #acquire}
* refuses (invariant 3: loudly, via {@link ShuttingDownException}) as much of the window as a
* plain flag can close; and the sweep below re-reads the registry once the initial snapshot has
* fully drained and drains whatever a straggler — a caller that read the flag as {@code false}
* a moment before it flipped — still managed to register. The sweep shares the same
* {@code deadline} rather than getting its own: {@code timeoutNanos} is a budget for the WHOLE
* drain (see above), and a straggler must not buy the drain more time than the flag it lost the
* race against would have. In the ordinary case the sweep finds nothing and costs one empty
* {@link #roster()} call.
*/
void drainAll(long timeoutNanos) {
long deadline = System.nanoTime() + timeoutNanos;
for (MemberSession s : roster()) {
draining.set(true);
drainSnapshot(roster(), deadline);
List<MemberSession> stragglers = roster();
if (!stragglers.isEmpty()) {
log.warn("drain sweep found {} session(s) registered after the drain snapshot was "
+ "taken (raced past the shutdown guard); draining them too", stragglers.size());
drainSnapshot(stragglers, deadline);
}
}
/**
* Drain exactly the sessions in {@code snapshot}, waiting out a {@code BUSY} one against the
* shared whole-drain {@code deadline} before releasing it. Shared by {@link #drainAll}'s main
* pass and its post-loop straggler sweep (fleetd #308) so both honor the same one budget.
*/
private void drainSnapshot(List<MemberSession> snapshot, long deadline) {
for (MemberSession s : snapshot) {
try {
if (s.state() == MemberSession.State.BUSY) {
while (System.nanoTime() < deadline) {
@@ -0,0 +1,18 @@
package dev.ltms.fleet.session;
/**
* Thrown by {@link SessionManager#acquire} when a spawn is requested after the daemon's shutdown
* drain has already begun (fleetd #308).
*
* <p>{@link SessionManager#drainAll} snapshots the registry once and tears down exactly what is
* in that snapshot. A session registered after the snapshot is invisible to the drain loop: its
* pane is left running and its worktree is never preserved, and nothing else ever reclaims
* either — the daemon's in-memory registry dies with the process. Refusing the spawn here,
* loudly, is what stops that session from ever being created in the first place, rather than
* silently handing the caller a session the daemon can no longer manage.
*/
public final class ShuttingDownException extends RuntimeException {
public ShuttingDownException(String message) {
super(message);
}
}
@@ -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")
@@ -413,4 +461,29 @@ class CallerResolverTest {
assertThrows(IllegalArgumentException.class, () -> new CallerResolver(id, true, null));
assertThrows(IllegalArgumentException.class, () -> new CallerResolver(id, true, " "));
}
@Test
void aWorkerOnAnyLoopbackSourceAddressIsStillAWorkerNotThePrimary() {
// fleetd #305: the escalation. ConnectionIdentity used to accept only 127.0.0.1, so a
// worker connecting from 127.0.0.2 resolved to no terminal, and this resolver's own
// (wider) loopback check then made it the PRIMARY — granting spawn, stop, send and drain.
// Measured on the Linux fleet host: binding a source of 127.0.0.2 succeeds there, so the
// path is real and not theoretical.
CallerResolver r = new CallerResolver(workerIdentity(), false, null);
for (String src : new String[]{"127.0.0.1", "127.0.0.2", "127.1.2.3", "::ffff:127.0.0.2"}) {
Principal p = r.resolve(src, 55555, null);
assertEquals(Role.WORKER, p.role(), "a worker must stay a worker from source " + src);
assertEquals("term_a", p.terminal(), "worker terminal from source " + src);
}
}
@Test
void aNonWorkerOnAnyLoopbackSourceAddressIsStillThePrimary() {
// The other direction of the same fix: widening the identity check must not demote a
// legitimate same-host primary that happens to connect from another 127.* address.
CallerResolver r = new CallerResolver(nonWorkerIdentity(), false, null);
for (String src : new String[]{"127.0.0.1", "127.0.0.2", "::ffff:127.0.0.1"}) {
assertEquals(Role.PRIMARY, r.resolve(src, 55555, null).role(), "source " + src);
}
}
}
@@ -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,
@@ -40,6 +40,7 @@ public final class FakeHerdr implements HerdrClient {
private final Map<String, List<String>> extraTabs = new LinkedHashMap<>();
private int agentNameTakenFor = 0;
private int agentPaneBusyFor = 0;
private String agentStartErrorCode = null;
private int workerTabPaneCount = 1;
private String paneCloseErrorCode = null;
private final Map<String, String> paneCloseErrorCodeFor = new ConcurrentHashMap<>();
@@ -84,6 +85,12 @@ public final class FakeHerdr implements HerdrClient {
return this;
}
/** Make every {@code agent.start} call fail with this herdr error code. */
public FakeHerdr agentStartFailsWith(String code) {
this.agentStartErrorCode = code;
return this;
}
/** Make the worker tab (w9:t2) report this many panes in {@code tab.list} (default 1). */
public FakeHerdr withWorkerTabPaneCount(int n) {
this.workerTabPaneCount = n;
@@ -311,6 +318,10 @@ public final class FakeHerdr implements HerdrClient {
+ required + "`", "invalid_request", null);
}
}
if (agentStartErrorCode != null) {
throw new HerdrException("herdr error [" + agentStartErrorCode + "]: agent.start failed",
agentStartErrorCode, null);
}
long starts = calls.stream().filter(c -> c.method().equals("agent.start")).count();
if (starts <= agentPaneBusyFor) {
throw new HerdrException(
@@ -297,6 +297,68 @@ class InjectorTest {
assertTrue(inj.activeTargets().isEmpty(), "the wedged target is reclaimed, not polled forever");
}
/** A listener whose post-turn housekeeping always starts, as SessionManager's does with clearAfterTurn on. */
private static final class PostTurnListener implements TurnListener {
@Override public void onTurnComplete(String target) { }
@Override public boolean hasPostTurnAction(String target) { return true; }
@Override public boolean onTurnCompleteWithPostAction(String target) { return true; }
}
@Test
void aWorkerThatWedgesInUnknownAwaitingPostTurnPickupIsReleased() {
// fleetd #306: the post-turn phase had no way out of a sustained unknown streak, so the
// pickup latch stayed set, the target was polled forever, and every later message to it was
// blocked by the delivery gate — while the session still looked healthy.
Captor cap = new Captor();
Injector inj = new Injector(new AgentControl(herdr), new PostTurnListener());
inj.enqueue(T, "task", TestTurnTokens.inert(T));
inj.onStatus(T, AgentStatus.IDLE); // deliver
inj.onStatus(T, AgentStatus.WORKING); // turn starts
inj.onStatus(T, AgentStatus.IDLE); // turn completes; housekeeping dispatched
for (int i = 0; i < STALL_SAMPLES; i++) inj.onStatus(T, AgentStatus.UNKNOWN); // then wedges
assertTrue(inj.activeTargets().isEmpty(),
"a target wedged awaiting post-turn pickup must be reclaimed, not polled forever");
assertEquals(List.of(), cap.failed,
"the delegated turn already completed — a stuck /clear must not be reported as a failed turn");
}
@Test
void aWorkerThatWedgesInUnknownAfterPickingUpTheResetIsReleased() {
// The sibling latch. postTurnObserved is set when the reset is seen picked up (WORKING) and
// is cleared only on a later injectable sample, so a wedge right after pickup sticks too.
Captor cap = new Captor();
Injector inj = new Injector(new AgentControl(herdr), new PostTurnListener());
inj.enqueue(T, "task", TestTurnTokens.inert(T));
inj.onStatus(T, AgentStatus.IDLE);
inj.onStatus(T, AgentStatus.WORKING);
inj.onStatus(T, AgentStatus.IDLE); // turn complete; reset dispatched
inj.onStatus(T, AgentStatus.WORKING); // reset picked up -> postTurnObserved
for (int i = 0; i < STALL_SAMPLES; i++) inj.onStatus(T, AgentStatus.UNKNOWN);
assertTrue(inj.activeTargets().isEmpty(), "a wedge after reset pickup must also be reclaimed");
assertEquals(List.of(), cap.failed, "still not a turn failure");
}
@Test
void aBriefUnknownDuringPostTurnHousekeepingDoesNotDropTheLatch() {
// The other direction: the escape must not fire on a glitch, or the queued next delegation
// would overtake housekeeping that is still running.
Injector inj = new Injector(new AgentControl(herdr), new PostTurnListener());
inj.enqueue(T, "first", TestTurnTokens.inert(T));
inj.enqueue(T, "second", TestTurnTokens.inert(T));
inj.onStatus(T, AgentStatus.IDLE);
inj.onStatus(T, AgentStatus.WORKING);
inj.onStatus(T, AgentStatus.IDLE); // first completes; reset dispatched
for (int i = 0; i < 10; i++) inj.onStatus(T, AgentStatus.UNKNOWN); // well under the grace
assertFalse(inj.activeTargets().isEmpty(), "a brief glitch must not release the post-turn latch");
assertEquals(List.of("first"), sent(), "the queued delegation must not overtake housekeeping");
}
@Test
void aTransientUnknownGlitchNeitherFailsNorBlocksCompletion() {
Captor cap = new Captor();
@@ -20,6 +20,17 @@ class ConnectionIdentityTest {
assertEquals("term_a", with(_ -> FakeHerdr.WORKER_PID).callerTerminal("127.0.0.1", 55555));
}
@Test
void resolvesWorkerFromAnyLoopbackSourceAddressNotJust127001() {
// fleetd #305. On Linux the whole 127.0.0.0/8 is bound to lo, so a worker can connect with
// a source address of 127.0.0.2. If identity resolution skips that address the caller has
// no terminal, and a caller with no terminal is the primary under loopback-trust — so this
// must resolve the worker, not null.
assertEquals("term_a", with(_ -> FakeHerdr.WORKER_PID).callerTerminal("127.0.0.2", 55555));
assertEquals("term_a", with(_ -> FakeHerdr.WORKER_PID).callerTerminal("127.1.2.3", 55555));
assertEquals("term_a", with(_ -> FakeHerdr.WORKER_PID).callerTerminal("::ffff:127.0.0.2", 55555));
}
@Test
void nullForOffHostCaller() {
// A non-loopback peer can't be an on-host worker → treat as primary/unknown.
@@ -32,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.
@@ -189,7 +189,7 @@ class FleetMcpTest {
}
@Test
void unansweredAsyncAskReturnsTheTicketToPending() throws Exception {
void unansweredAsyncAskReturnsTheTicketToPendingThenAWorkersLateReplyStillCompletesIt() throws Exception {
McpSchema.CallToolResult accepted = FleetMcp.sendAsync(messages, "term_a", "do it", null, Set.of());
String ticket = textOf(accepted).substring(textOf(accepted).indexOf("ticket=") + "ticket=".length()).trim();
@@ -203,8 +203,15 @@ class FleetMcpTest {
assertTrue(textOf(ask).contains("no answer"), textOf(ask));
assertTrue(textOf(FleetMcp.poll(messages, ticket, null)).startsWith("[pending"));
// fleetd #307: the worker resumed on its own after the primary never answered, and its real
// fleet_reply must complete its OWN async ticket — not strand in the inbox with
// fleet_poll{ticket} stuck PENDING forever and later force-failed with a false "session
// released before it replied" reason. This used to land in the inbox instead (see the old
// assertion this replaced: messages.drainReplies("term_a").getFirst()...) — that was the bug.
FleetMcp.reply(messages, "term_a", "finished after timeout");
assertEquals("finished after timeout", messages.drainReplies("term_a").getFirst().content());
assertEquals("finished after timeout", textOf(FleetMcp.poll(messages, ticket, null)));
assertTrue(messages.drainReplies("term_a").isEmpty(),
"the reply completed its own ticket directly and never touched the inbox");
}
@Test
@@ -326,6 +333,26 @@ class FleetMcpTest {
assertEquals("orphan", drained.getFirst().content());
}
@Test
void replyWithBlankContentIsACleanToolErrorNotAnUncaughtException() {
// fleetd #302: MessageService.reply now REJECTS blank content by throwing. fleet_reply's
// handler is a bare BiFunction with no try/catch around it, so if this guard only checked
// `== null` (as it did), a whitespace-only reply would leave the handler as an uncaught
// IllegalArgumentException instead of a tool error the caller can read. Null and whitespace
// are the same caller mistake and must get the same answer — the sibling fleet_send guard
// has always used isBlank for exactly this reason.
for (String blank : new String[] {null, "", " ", "\n\t"}) {
McpSchema.CallToolResult res = assertDoesNotThrow(
() -> FleetMcp.reply(messages, "term_a", blank),
"blank content must be refused as a tool error, never thrown out of the handler");
assertEquals(Boolean.TRUE, res.isError(), "blank content is an error result");
assertTrue(textOf(res).contains("content is required"),
"the error names the missing argument: " + textOf(res));
}
assertEquals(0, messages.drainReplies("term_a").size(),
"a refused reply must not reach the inbox");
}
@Test
void bridgePollWithTargetDrainsReplies() {
// A reply with no open send queues it in the inbox.
@@ -1165,7 +1165,8 @@ class ClaudeCodeLauncherTest {
@Test
void spawnLetsAnUnrelatedHerdrErrorPropagateUnchanged() {
// Fix 1 must only special-case a "*_not_found" answer. Any other herdr failure keeps
// propagating as-is — this gate does not know how to recover from it.
// propagating as-is — this gate does not know how to recover from it. The pane still needs
// closing because spawn throws before it can return the pane id to a caller that could stop it.
FakeHerdr herdr = new FakeHerdr();
herdr.agentStatus("unknown");
herdr.agentGetFailsWithAfter(0, "internal_error");
@@ -1182,8 +1183,27 @@ class ClaudeCodeLauncherTest {
() -> svc.spawn(new SpawnRequest(null, null, null)));
assertEquals("internal_error", ex.code());
assertEquals(0, paneCloseCount(herdr, "w9:pRoot_1"),
"an error this gate does not recognize is not this gate's teardown to run");
assertEquals(1, paneCloseCount(herdr, "w9:pRoot_1"),
"the unchanged error leaves spawn without a pane id, so this gate closes its orphaned pane");
}
@Test
void panePlacementClosesTheSplitPaneWhenAgentStartFails() {
FakeHerdr herdr = new FakeHerdr().agentStartFailsWith("internal_error");
FleetConfig.Profile cfg = new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
List.of("claude"), "pane", "fleetd-workers", "w #{n}", null, null, null);
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(
new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
dev.ltms.fleet.herdr.HerdrException ex = assertThrows(
dev.ltms.fleet.herdr.HerdrException.class,
() -> svc.spawn(new SpawnRequest(null, null, null)));
assertEquals("internal_error", ex.code(), "agent.start failure propagates unchanged");
assertEquals(1, paneCloseCount(herdr, "w1:pSplit"),
"the pane split for a peer that never starts is closed instead of left orphaned");
}
// --- fleetd #176 fix 2: corroborated UNKNOWN refinement --------------------------------------
@@ -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();
@@ -957,6 +957,76 @@ class MessageServiceTest {
assertEquals(MessageService.Phase.DONE, awaitTicketPhase(next, MessageService.Phase.DONE).phase());
}
/**
* fleetd #307: a worker's {@code fleet_ask} can time out because the primary never answers —
* distinct from {@link #aReplyAfterAnswerTimesOutStillCompletesTheAsyncTicket}, where the
* primary DID answer and only its own bounded wait for the resumed turn expired.
* {@code ask()}'s timeout path deliberately forgets the task's {@code turnId} (so
* {@code hasAsyncQuestion} stops reporting the target BUSY — see
* {@code unansweredAsyncQuestionReturnsTheTicketToPendingAndReleasesItsTarget} above), which used
* to also erase the one signal {@code askAnsweredAsyncTasks} needed to recognize the worker's
* eventual real {@code fleet_reply}. That reply then had nowhere to land but the inbox, and
* {@code fleet_poll{ticket}} stayed PENDING forever — later force-failed with the false reason
* "session released before it replied", even though the worker had, in fact, replied.
*/
@Test
void aReplyAfterAnAskTimeoutStillCompletesTheAsyncTicket() throws Exception {
String ticket = messages.sendAsync(T, "task that asks then finishes alone");
awaitWaiting();
injectDelivery();
assertEquals(MessageService.AskOutcome.TIMED_OUT,
messages.ask(T, "which config?", 200).outcome());
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket).phase(),
"only the question wait ended; the delegated turn may still finish");
// The worker keeps working past the timeout and only now calls fleet_reply — with no live
// rendezvous waiter open (ask()'s timeout already closed it) and no new send() having
// reopened one for this target.
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE);
assertEquals("PR opened: https://example/pulls/42", done.reply(),
"fleet_poll{ticket} must return the worker's real reply, not stay pending forever");
assertEquals("reply", done.replySource());
assertFalse(messages.hasStrandedReply(T),
"the reply completed its own ticket directly and never touched the inbox");
}
/**
* fleetd #307's ambiguity guard: an ask timeout frees its target ({@code hasAsyncQuestion}
* becomes false the instant it lapses — proven above), so a second, independent delegation can
* be dispatched to the same target and itself go on to ask-and-lapse before the first worker's
* real reply ever arrives. Two open tasks are then both eligible candidates on one target with
* no live waiter to disambiguate them. A reply arriving now must not guess which one it answers
* — guessing wrong would hand the lead a plausible-looking answer to a delegation the worker
* never touched, worse than a failure because the lead acts on it — so it must fall back to the
* inbox exactly as the zero-candidate case does.
*/
@Test
void twoAskTimedOutTicketsOnOneTargetFallBackToTheInboxRatherThanGuess() throws Exception {
String ticket1 = messages.sendAsync(T, "first task that asks");
awaitWaiting();
injectDelivery();
assertEquals(MessageService.AskOutcome.TIMED_OUT, messages.ask(T, "Q1?", 200).outcome());
String ticket2 = messages.sendAsync(T, "second task that asks");
awaitWaiting();
injectDelivery();
assertEquals(MessageService.AskOutcome.TIMED_OUT, messages.ask(T, "Q2?", 200).outcome());
assertTrue(messages.reply(T, "which task does this answer?"));
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket1).phase(),
"an ambiguous reply must not guess ticket1");
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket2).phase(),
"an ambiguous reply must not guess ticket2");
assertTrue(messages.hasStrandedReply(T));
var drained = messages.drainReplies(T);
assertEquals(1, drained.size());
assertEquals("which task does this answer?", drained.get(0).content());
}
@Test
void asyncQuestionBelongsToTheTaskThatOwnsItsForwardWaiter() throws Exception {
String first = messages.sendAsync(T, "first task");
@@ -19,6 +19,10 @@ import dev.ltms.fleet.session.SessionManager;
import dev.ltms.fleet.session.Worktrees;
import dev.ltms.fleet.member.ClaudeCodeLauncher;
import dev.ltms.fleet.member.CompositePeerLauncher;
import dev.ltms.fleet.member.MemberCredentialPolicyView;
import dev.ltms.fleet.mcp.FleetMcp;
import dev.ltms.fleet.placement.BackendOutagePolicy;
import dev.ltms.fleet.placement.BackendQuarantine;
import dev.ltms.fleet.placement.PlacementPolicies;
import io.javalin.Javalin;
import org.junit.jupiter.api.AfterEach;
@@ -32,6 +36,7 @@ import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.TimeUnit;
import java.util.function.Predicate;
import static org.junit.jupiter.api.Assertions.*;
@@ -70,6 +75,19 @@ class FleetAppTest {
private int start(FakeHerdr herdr, String workerBaseUrl, Set<String> allow, String placement,
Worktrees worktrees, Predicate<String> deliverable) {
return start(herdr, workerBaseUrl, allow, placement, worktrees, deliverable,
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none());
}
/**
* fleetd #297: same wiring as above, plus the two SAME shared sources {@code GET /profiles}
* must read — lets a test prove the quarantined/coolingOff facts it reports come from a real
* {@link dev.ltms.fleet.placement.BackendQuarantine}/{@link
* dev.ltms.fleet.placement.BackendOutagePolicy}, exactly like {@code fleet_profiles}'s own tests.
*/
private int start(FakeHerdr herdr, String workerBaseUrl, Set<String> allow, String placement,
Worktrees worktrees, Predicate<String> deliverable,
FleetMcp.QuarantineSource quarantine, FleetMcp.OutageSource outage) {
FleetConfig.Profile wcfg = new FleetConfig.Profile(
"ltms-local", workerBaseUrl, "coder", null, "FLEETD_WORKER_TOKEN", null,
placement, "fleet", "worker: {profile} #{n}", null, null, null);
@@ -90,8 +108,9 @@ class FleetAppTest {
// it directly so the inbox contract holds for those endpoints.
inbox.own("term_a");
MessageService messages = new MessageService(agents, injector, rendezvous, inbox);
app = new FleetApp(herdr, workers, sessions, messages, this.presence, null,
null, null, id -> this.presence.isPresent(id) || deliverable.test(id))
app = new FleetApp(herdr, herdr, workers, sessions, messages, this.presence, null,
null, null, id -> this.presence.isPresent(id) || deliverable.test(id),
MemberCredentialPolicyView::absent, quarantine, outage)
.build().start("127.0.0.1", 0);
return app.port();
}
@@ -164,6 +183,22 @@ class FleetAppTest {
assertEquals("idle", agents.get(0).get("status").asText());
}
/**
* fleetd #297 gap 1: {@code workers.list()} reaches herdr, and a transport failure there must
* land in the same {@code {error, detail}} envelope every other failure path in this file uses
* (see {@code herdrError}), not escape as a bare exception outside the JSON contract.
*/
@Test
void agentsMapsAHerdrFailureToTheJsonErrorEnvelope() throws Exception {
FakeHerdr down = new FakeHerdr().healthy(false);
int port = start(down, "http://gx00.gw:8000", Set.of("gx00.gw"));
HttpResponse<String> res = req(port, "GET", "/agents");
assertEquals(502, res.statusCode(), res.body());
JsonNode body = mapper.readTree(res.body());
assertEquals("herdr_error", body.get("error").asText());
assertTrue(body.has("detail"), res.body());
}
@Test
void spawnWorkerLandsInOwnTabInWorkerSpaceAndInjectsBaseUrl() throws Exception {
FakeHerdr herdr = new FakeHerdr();
@@ -203,6 +238,42 @@ class FleetAppTest {
JsonNode body = mapper.readTree(req(port, "GET", "/profiles").body());
assertEquals("ltms-local", body.get("default").asText());
assertEquals("ltms-local", body.get("profiles").get(0).asText());
assertFalse(body.has("quarantined"), "nothing is quarantined, so the key is omitted: " + body);
assertFalse(body.has("coolingOff"), "nothing is cooling off, so the key is omitted: " + body);
}
/**
* fleetd #297 gap 2: {@code GET /profiles} must report the same two outage states {@code
* fleet_profiles} does — CB-578 stage B exhaustion quarantine and fleetd #201 Unit 5 cool-off —
* reading the SAME shared {@link BackendQuarantine}/{@link BackendOutagePolicy} instances rather
* than recomputing them. The two checks are independent, and this profile is deliberately put in
* both states at once, matching {@code FleetMcpTest}'s own coverage of that overlap.
*/
@Test
void profilesReportsQuarantineAndCoolingOffFromTheSameSharedSources() throws Exception {
FakeHerdr herdr = new FakeHerdr();
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
quarantine.quarantine("shared-openai");
FleetMcp.QuarantineSource quarantineSource = new FleetMcp.QuarantineSource(
profile -> "ltms-local".equals(profile) ? "shared-openai" : null, quarantine);
BackendOutagePolicy outagePolicy = new BackendOutagePolicy(() -> 0L);
outagePolicy.record("shared-openai", "t1", "API Error: rate limited");
outagePolicy.record("shared-openai", "t2", "API Error: rate limited"); // 2nd distinct target starts the incident
FleetMcp.OutageSource outageSource = new FleetMcp.OutageSource(
profile -> "ltms-local".equals(profile) ? "shared-openai" : null, outagePolicy);
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"), "tab", new GitWorktrees(),
ignored -> false, quarantineSource, outageSource);
JsonNode body = mapper.readTree(req(port, "GET", "/profiles").body());
assertTrue(body.has("quarantined"), body.toString());
assertEquals("shared-openai",
body.get("quarantined").get("ltms-local").get("credentialId").asText());
assertEquals(1800,
body.get("quarantined").get("ltms-local").get("quarantinedForSeconds").asLong());
assertTrue(body.has("coolingOff"), body.toString());
assertEquals("shared-openai",
body.get("coolingOff").get("ltms-local").get("credentialId").asText());
assertEquals(60, body.get("coolingOff").get("ltms-local").get("coolingOffForSeconds").asLong());
}
@Test
@@ -239,6 +310,23 @@ class FleetAppTest {
"liveStatus is unknown when herdr has no matching pane");
}
/**
* fleetd #297 gap 1: same reasoning as {@code agentsMapsAHerdrFailureToTheJsonErrorEnvelope} —
* {@code GET /members} is the endpoint's own comment names as "the out-of-band path a lead falls
* back to when its MCP mount drops", so it must stay inside the {@code {error, detail}} envelope
* exactly when herdr is briefly unreachable.
*/
@Test
void membersMapsAHerdrFailureToTheJsonErrorEnvelope() throws Exception {
FakeHerdr down = new FakeHerdr().healthy(false);
int port = start(down, "http://gx00.gw:8000", Set.of("gx00.gw"));
HttpResponse<String> res = req(port, "GET", "/members");
assertEquals(502, res.statusCode(), res.body());
JsonNode body = mapper.readTree(res.body());
assertEquals("herdr_error", body.get("error").asText());
assertTrue(body.has("detail"), res.body());
}
@Test
void spawnWithACwdParamRootsTheWorkerThere() throws Exception {
FakeHerdr herdr = new FakeHerdr();
@@ -327,7 +415,11 @@ class FleetAppTest {
FakeHerdr herdr = new FakeHerdr().agentNameTakenTimes(99);
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
assertEquals(500, req(port, "POST", "/members").statusCode());
// fleetd #304: 502, not Javalin's default 500 — the herdr failure is named, and the body
// carries herdr's own message, matching what fleet_spawn reports for the same failure.
HttpResponse<String> res = req(port, "POST", "/members");
assertEquals(502, res.statusCode());
assertEquals("herdr_error", mapper.readTree(res.body()).get("error").asText());
assertTrue(herdr.called("tab.create"), "a tab was created before the failed start");
assertEquals("w9:t2", params(herdr, "tab.close").get("tab_id"), "orphaned tab must be closed");
}
@@ -444,6 +536,53 @@ class FleetAppTest {
assertEquals("orphan", body.get("replies").get(0).get("content").asText());
}
@Test
void replyWithMissingContentIsRejectedAndDoesNotResolveTheWaiter() throws Exception {
// fleetd #302: `.path("content").asText("")` used to turn a missing "content" key into an
// empty string that reached rendezvous.resolve, silently completing the lead's blocking wait
// with nothing. Prove the fix two ways: the bad call is rejected with 400, AND the real send
// it would have wrongly resolved is still open afterwards — a real reply completes it.
FakeHerdr herdr = new FakeHerdr().agentStatus("idle");
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
var send = java.util.concurrent.CompletableFuture.supplyAsync(() -> {
try { return postMessage(port, "{\"content\":\"review this\",\"timeoutMs\":4000}"); }
catch (Exception e) { throw new RuntimeException(e); }
});
Thread.sleep(200); // let the background send open its rendezvous waiter
HttpResponse<String> badReply = postJson(port, "/sessions/term_a/reply", "{}");
assertEquals(400, badReply.statusCode());
JsonNode err = mapper.readTree(badReply.body());
assertEquals("bad_request", err.get("error").asText());
assertTrue(err.has("detail"));
// The waiter must still be open — a real reply now completes the ORIGINAL send.
HttpResponse<String> goodReply = postJson(port, "/sessions/term_a/reply", "{\"content\":\"LGTM ship it\"}");
assertEquals(200, goodReply.statusCode());
HttpResponse<String> res = send.get(6, java.util.concurrent.TimeUnit.SECONDS);
assertEquals(200, res.statusCode());
assertEquals("LGTM ship it", mapper.readTree(res.body()).get("reply").asText());
}
@Test
void replyWithEmptyOrWhitespaceContentIsRejectedSameAsMissing() throws Exception {
// fleetd #302 sibling case: present-but-blank content is treated the same as a missing key —
// FleetMcp's own required-content guard (fleet_reply's "content is required") makes no
// distinction between the two either, so diverging here would be a new asymmetry.
int port = startHealthy();
HttpResponse<String> empty = postJson(port, "/sessions/term_a/reply", "{\"content\":\"\"}");
assertEquals(400, empty.statusCode());
assertEquals("bad_request", mapper.readTree(empty.body()).get("error").asText());
HttpResponse<String> whitespace = postJson(port, "/sessions/term_a/reply", "{\"content\":\" \"}");
assertEquals(400, whitespace.statusCode());
assertEquals("bad_request", mapper.readTree(whitespace.body()).get("error").asText());
}
@Test
void drainRepliesReturnsEmptyForNoReplies() throws Exception {
int port = startHealthy();
@@ -598,7 +737,12 @@ class FleetAppTest {
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
// A genuine teardown failure must surface, not be reported as a successful 204.
assertEquals(500, req(port, "DELETE", "/members/w9:pW").statusCode());
// fleetd #304: it surfaces as a named 502 rather than Javalin's default 500. The property
// this test guards is "not 204" and the herdr detail reaching the caller — a bare 500 gave
// the body "Server Error" and said nothing about herdr.
HttpResponse<String> res = req(port, "DELETE", "/members/w9:pW");
assertEquals(502, res.statusCode());
assertEquals("herdr_error", mapper.readTree(res.body()).get("error").asText());
assertFalse(herdr.called("tab.close"), "tab is not removed when the pane close failed");
}
@@ -611,4 +755,5 @@ class FleetAppTest {
assertEquals(204, req(port, "DELETE", "/members/w9:pW").statusCode());
assertTrue(herdr.called("tab.close"));
}
}
@@ -403,6 +403,52 @@ class GitWorktreesTest {
assertTrue(heads.isBlank(), "the branch leaked after a post-creation step threw:\n" + heads);
}
/**
* fleetd #309. The add runner is a narrow seam for the case where Git has created state but
* the caller then kills the process. The runner first performs the real add in this throwaway
* repo, then throws the same kind of exception that {@link GitWorktrees#exec} uses for a timeout.
* This proves the failure path cleans both real Git objects without waiting for a slow checkout.
*/
@Test
void addCleansUpWhenTheWorktreeAddRunnerFailsAfterCreatingState(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
String branch = "cb-309-timeout";
WorktreeException timeout = new WorktreeException("command timed out: synthetic git worktree add");
AtomicReference<String> createdPath = new AtomicReference<>();
GitWorktrees worktrees = new GitWorktrees(tmp.resolve("wts").toString(), null, null, null, command -> {
createdPath.set(command[5]);
try {
git(repo, "worktree", "add", command[5], "-b", command[7], command[8]);
} catch (Exception e) {
throw new AssertionError("test setup could not create the worktree", e);
}
throw timeout;
});
WorktreeException thrown = assertThrows(WorktreeException.class,
() -> worktrees.add(repo.toString(), branch, "HEAD"));
assertSame(timeout, thrown, "cleanup must not replace the add failure");
assertNotNull(createdPath.get(), "the add runner must receive the worktree path");
assertFalse(Files.exists(Path.of(createdPath.get())),
"the worktree directory leaked after the add runner failed");
assertFalse(refExists(repo, "refs/heads/" + branch), "the branch leaked after the add runner failed");
}
/** An ordinary Git refusal must not delete the existing branch or log a cleanup warning. */
@Test
void addFailureBeforeCreatingAWorktreeIsQuiet(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
String branch = "already-exists";
git(repo, "branch", branch);
assertThrows(WorktreeException.class, () -> new GitWorktrees(tmp.resolve("wts").toString())
.add(repo.toString(), branch, "HEAD"));
assertTrue(refExists(repo, "refs/heads/" + branch), "the existing branch must remain");
assertTrue(capturedMessages().isEmpty(), "an ordinary Git refusal logged a warning: " + capturedMessages());
}
// ---- CB-189: broader remote-URL coverage — every remote, both fetch and push URLs, any
// non-SSH scheme. Reporting only, additive to the origin/https strip-and-refuse tests above. ----
@@ -26,7 +26,12 @@ import org.slf4j.LoggerFactory;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.function.LongSupplier;
import static org.junit.jupiter.api.Assertions.*;
@@ -178,14 +183,19 @@ class SessionManagerTest {
}
private SessionManager sessionManager(FakeHerdr herdr, LongSupplier clock, int contextCap,
boolean clearAfterTurn) {
boolean clearAfterTurn) {
return sessionManager(herdr, clock, contextCap, clearAfterTurn, null);
}
private SessionManager sessionManager(FakeHerdr herdr, LongSupplier clock, int contextCap,
boolean clearAfterTurn, java.util.function.Consumer<MemberSession> hook) {
FleetConfig.Profile cfg = new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
List.of("ccs", "ltms-local"), "tab", "fleetd-workers",
"worker: {profile} #{n}", null, null, null);
ClaudeCodeLauncher workers = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
return new SessionManager(workers, new GitWorktrees(), clock, contextCap, clearAfterTurn);
return new SessionManager(workers, new GitWorktrees(), clock, contextCap, clearAfterTurn, hook);
}
@Test
@@ -670,6 +680,24 @@ class SessionManagerTest {
"BUSY session remains");
}
@Test
void reapIdleDoesNotReleaseSessionDeliveredAfterItsEligibilityCheck() {
long[] clock = {0};
FakeHerdr herdr = new FakeHerdr();
SessionManager[] manager = new SessionManager[1];
SessionManager sessions = sessionManager(herdr, () -> clock[0], 0, false,
session -> manager[0].onDelivered(session.terminalId(), TestTurnTokens.inert(session.terminalId())));
manager[0] = sessions;
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
sessions.asPresence().markPresent(session.terminalId());
clock[0] = 11;
assertEquals(0, sessions.reapIdle(10), "delivery replaces the idle snapshot before release");
assertEquals(MemberSession.State.BUSY, sessions.get(session.paneId()).orElseThrow().state(),
"a just-delivered session stays registered and busy");
assertFalse(herdr.called("pane.close"), "the busy session pane is not stopped");
}
@Test
void doneSessionPastIdleTtlIsReaped() {
long[] clock = {0};
@@ -824,6 +852,186 @@ class SessionManagerTest {
.count();
}
// --- fleetd #308: a spawn accepted while the shutdown drain is running must not orphan ---
@Test
void acquireRefusesANewSpawnOnceDrainAllHasStarted() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
sessions.drainAll(TimeUnit.MILLISECONDS.toNanos(50)); // empty roster — returns immediately,
// but the shutdown guard it flips must stay tripped for the life of the process.
ShuttingDownException e = assertThrows(ShuttingDownException.class,
() -> sessions.acquire("ltms-local", null, "/caller", "term_primary"),
"a spawn requested after the drain has begun must be refused loudly (invariant 3), "
+ "not silently registered into a registry the drain will never revisit");
assertNotNull(e.getMessage());
assertFalse(e.getMessage().isBlank(), "the refusal must say why, not just that it failed");
assertTrue(sessions.roster().isEmpty(), "the refused spawn must never reach the registry");
}
/**
* fleetd #308: the guard above closes most of the shutdown-race window, but it cannot close
* all of it — a caller that already passed the {@code draining} check before {@code drainAll}
* flips it can still be mid-{@code launcher.spawn()} (a real herdr round trip, not
* instantaneous) when {@code drainAll} takes its registry snapshot. This test forces exactly
* that interleaving with a launcher double that blocks the second {@code spawn()} call and the
* first {@code stop()} call until released, then proves the post-loop sweep in {@code
* drainAll} still finds and tears down the straggler that lands in the registry afterward.
*/
@Test
void drainAllSweepsAStragglerThatRegisteredAfterTheInitialSnapshot() throws Exception {
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile cfg = new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
List.of("ccs", "ltms-local"), "tab", "fleetd-workers",
"worker: {profile} #{n}", null, null, null);
ClaudeCodeLauncher delegate = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
RaceLauncher race = new RaceLauncher(delegate);
SessionManager sessions = new SessionManager(race);
// Registered normally, before the drain starts — the first spawn call, never blocked.
MemberSession ready = sessions.acquire("ltms-local", "/ready", "/caller", "ownerR");
ExecutorService exec = Executors.newFixedThreadPool(2);
try {
// The straggler's acquire() reads `draining == false` (checked before this call ever
// touches the launcher) and then blocks inside its own spawn() — the second spawn call.
Future<MemberSession> straggler = exec.submit(() ->
sessions.acquire("ltms-local", "/late", "/caller", "ownerLate"));
assertTrue(race.enteredSecondSpawn.await(5, TimeUnit.SECONDS),
"the straggler must have passed the shutdown guard and reached spawn() before "
+ "drainAll ever runs");
assertEquals(1, sessions.roster().size(),
"the straggler is still inside spawn() — not registered yet");
// drainAll flips `draining`, snapshots the registry (only `ready` is in it), and starts
// releasing that snapshot — its first release() call stops `ready`'s pane, which this
// launcher double blocks on so the interleaving below is deterministic, not a timing bet.
Future<?> drain = exec.submit(() -> sessions.drainAll(TimeUnit.SECONDS.toNanos(5)));
assertTrue(race.enteredFirstStop.await(5, TimeUnit.SECONDS),
"drainAll must be stopping the ready session's pane — proof its initial "
+ "registry snapshot has already been taken");
// Only now does the straggler's spawn complete and register — strictly after the
// snapshot drainAll's main pass is working from.
race.releaseSecondSpawn.countDown();
MemberSession registered = straggler.get(5, TimeUnit.SECONDS);
// Let drainAll finish releasing `ready`; it then re-checks the registry and must find
// (and drain) the straggler that just landed in it.
race.releaseFirstStop.countDown();
drain.get(5, TimeUnit.SECONDS);
assertTrue(sessions.roster().isEmpty(),
"the post-loop sweep must drain the straggler too, not just the initial snapshot");
assertNotNull(registered.paneId());
long paneCloseCalls = herdr.calls.stream().filter(c -> "pane.close".equals(c.method())).count();
assertEquals(2, paneCloseCalls,
"both ready's pane AND the straggler's pane must actually be stopped — a pane "
+ "left running is exactly the orphan this ticket is about");
} finally {
exec.shutdownNow();
}
}
/**
* Delegates every call while blocking the SECOND {@code spawn()} call and the FIRST
* {@code stop()} call until the test releases them — used to force the fleetd #308 race
* deterministically instead of betting on real thread-scheduling timing.
*/
private static final class RaceLauncher implements PeerLauncher {
private final PeerLauncher delegate;
private final AtomicInteger spawnCalls = new AtomicInteger();
private final AtomicInteger stopCalls = new AtomicInteger();
final CountDownLatch enteredSecondSpawn = new CountDownLatch(1);
final CountDownLatch releaseSecondSpawn = new CountDownLatch(1);
final CountDownLatch enteredFirstStop = new CountDownLatch(1);
final CountDownLatch releaseFirstStop = new CountDownLatch(1);
RaceLauncher(PeerLauncher delegate) {
this.delegate = delegate;
}
private static void awaitOrFail(CountDownLatch latch) {
try {
if (!latch.await(5, TimeUnit.SECONDS)) {
throw new AssertionError("RaceLauncher latch timed out");
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new AssertionError("RaceLauncher latch interrupted", e);
}
}
@Override
public Set<Capability> capabilities() {
return delegate.capabilities();
}
@Override
public Set<Capability> capabilitiesFor(String profileName) {
return delegate.capabilitiesFor(profileName);
}
@Override
public PeerHandle spawn(SpawnRequest req) {
if (spawnCalls.incrementAndGet() == 2) {
enteredSecondSpawn.countDown();
awaitOrFail(releaseSecondSpawn);
}
return delegate.spawn(req);
}
@Override
public Set<String> profiles() {
return delegate.profiles();
}
@Override
public String defaultProfile() {
return delegate.defaultProfile();
}
@Override
public String effectiveCwd(SpawnRequest req) {
return delegate.effectiveCwd(req);
}
@Override
public List<String> parityOverlay(String profileName) {
return delegate.parityOverlay(profileName);
}
@Override
public List<?> list() {
return delegate.list();
}
@Override
public int reapOrphanWorkers() {
return delegate.reapOrphanWorkers();
}
@Override
public void stop(String id) {
if (stopCalls.incrementAndGet() == 1) {
enteredFirstStop.countDown();
awaitOrFail(releaseFirstStop);
}
delegate.stop(id);
}
@Override
public boolean clearContext(String id) {
return delegate.clearContext(id);
}
}
private static List<String> promptTexts(FakeHerdr herdr) {
return herdr.calls.stream()
.filter(c -> "agent.prompt".equals(c.method()))