Compare commits

...

23 Commits

Author SHA1 Message Date
Dai Ha ee5f8b932b #137: complete an async ticket's own reply after answer() times out
CI / contract (pull_request) Successful in 49s
CI / build (pull_request) Successful in 1m44s
fleet_send{turnId} (MessageService.answer) blocks the primary only for its
own bounded MCP-call window (25s default, 120s max) — far shorter than a
worker's resumed turn can genuinely take. When that window expires, answer()
closes its rendezvous waiter, so the worker's eventual fleet_reply has no
live waiter to resolve and falls back to the session inbox. The async
ticket's future was never completed by that path, so fleet_poll{ticket}
stayed PENDING until fleet_stop's abandon() forced it FAILED with a
misleading "the worker session was released before it replied" reason,
even though the reply had genuinely arrived.

- MessageService.reply(): before falling to the inbox, look for the async
  task this exact turn belongs to (already answered — question cleared,
  turnId still stamped — but not yet resolved) and complete it directly
  with the real reply, so fleet_poll{ticket} returns it.
- MessageService.abandon(): defense in depth, independent of the above —
  never write a false WORKER_FAILED once a reply reached the inbox for
  this target; recover and use its real content instead.
- Two new tests drive the full delegation path (async send -> ask ->
  answer with a short timeout -> reply -> poll/abandon), not a reply sink
  directly; both fail with the fix disabled and pass with it restored.
2026-08-31 14:12:23 +07:00
Dai Ha 4ac688b6d9 #164: classify a backend-error scrape as WORKER_FAILED, carrying the whole pane
CI / contract (push) Successful in 48s
CI / build (push) Successful in 1m22s
main already shipped the core of #164 in 3bfa828: the MIN_TURN_NANOS floor and
the hard fail on an empty or unreadable scrape. This adds the one case that was
still resolving as a success -- a scrape that reads cleanly but whose content is
the backend's own rejection (e.g. "API Error: 400 invalid request body").

The BACKEND_ERROR pattern is deliberately narrow. A growing list of ad-hoc error
strings rots as backends change their wording, and broader backend-error
surfacing is #164 point 3.

Because the pattern is a heuristic, it also matches a member that forgot
fleet_reply while reporting *about* a backend error. So the failure reason
carries the whole pane tail, not just the matched line: a genuine backend error
reads as before, and a false positive keeps its report instead of losing it.

Checked before merging: 3bfa828 is an ancestor of main; the branch was current
with main; mvn clean install green unpiped (1037 tests, 0 [ERROR] lines); and
each of the 5 new tests fails with the fix commented out.

Co-authored-by: fleetd worker <worker@ltms.dev>
2026-08-31 10:56:40 +07:00
ltms a814d1ef00 #197: measure the async ticket TTL from completion, not from creation
CI / build (push) Successful in 1m15s
CI / contract (push) Successful in 1m17s
Lead-authored and lead-verified: full mvn clean install green at 1032 tests, and the new test proved by restoring the old comparison and watching it fail.
2026-08-31 05:10:46 +02:00
Dai Ha ea98856130 #197: measure the async ticket TTL from completion, not from creation
CI / contract (pull_request) Successful in 54s
CI / build (pull_request) Successful in 1m32s
pruneTerminalTickets compared the cutoff against createdNanos, so the real
window to collect a reply was "TTL minus however long the task ran". A
delegation that ran longer than the 10-minute TTL was already past the cutoff
the moment it finished, so the next prune destroyed its reply.

That is the normal case here, not an edge case. Real delegated work runs well
past ten minutes. Three workers in one session did, and two of their complete
reports were lost. The reply lives only in Task.future, so pruning it discards
the worker's whole report, and fleet_poll{target} returns [] rather than
holding it — there is no fallback.

Task now stamps completedNanos from a whenComplete hook registered in its
constructor, so every completion path stamps it (a reply, the completion
fallback, a timeout, a failure, an abandon on teardown) without each one having
to remember to. The stamp is a boxed Long, not a long with a sentinel:
System.nanoTime may return any value, so no number can mean "not stamped yet".
A task that is done but not yet stamped is left for the next sweep.

createdNanos had no other reader and is removed.

The TTL still bounds tasks — an uncollected finished ticket is still evicted
once the TTL passes since it finished. Both halves are pinned by a test, and
the first one was proved by restoring the old comparison and watching it fail.

Tests run: 1032, Failures: 0, Errors: 0, Skipped: 0
2026-08-31 10:10:16 +07:00
ltms a49e96835a CB-189: cover every remote, both URLs, and any non-SSH scheme in the credential check
CI / contract (push) Successful in 44s
CI / build (push) Successful in 1m11s
Lead-verified: merged onto current main (which already carries #194 and #196), full `mvn clean install` green at 1030 tests. Confirmed no plain `exec` call that reads a remote URL remains.

Review round 2 closed the half-fix: execRedacted was added but applied only to the new calls, leaving the three pre-existing URL readers (lines 143, 171, 334) still copying stdout into exception messages — the exact hole CB-189 named.

Kept the check before the origin strip, on the worker's reasoning: the credential really is in the config at that moment, so WARN-found followed by INFO-fixed is the full audit trail, whereas moving it after would silence the origin case entirely.
2026-08-31 04:32:14 +02:00
ltms 63c19dcba7 CB-185: fix two blockers to switching on memberHerdrSocket
CI / build (push) Successful in 1m31s
CI / contract (push) Successful in 1m19s
Lead-verified: merged with #194 onto an integration branch off main, full `mvn clean install` green at 1025 tests.

Review round 2 fixed the stale ambiguous-pane message and, more importantly, a real abort: probeOwner called list() unguarded, so one unreachable daemon made panes on a different healthy daemon un-stoppable too — the same bug blocker 1 exists to fix, through a new door. Worker proved it by removing the guard and quoting the failure.
2026-08-31 04:30:57 +02:00
ltms 2823349c8e CB-192: fix false credential-gap WARN under allow-list+zsh, split its log guard
CI / contract (push) Successful in 44s
CI / build (push) Successful in 1m11s
Lead-verified: merged with #196 onto an integration branch off main, full `mvn clean install` green at 1025 tests. Deny-by-default WARN confirmed byte-identical to main.

Review round 2 fixed a defect I found in round 1: the new INFO asserted that gap names were not on the derived allow-list without ever checking, so a name derived from a profile's tokenEnv/gitTokenEnv/env: would be reported as safe while the member actually inherited it. Now split with MemberEnvAllowList.keeps — the same predicate the generated scrub evaluates.
2026-08-31 04:30:48 +02:00
Dai Ha 6fc301d62c CB-189 review fix: redact the three pre-existing URL-reading exec calls
CI / contract (pull_request) Successful in 1m11s
CI / build (pull_request) Successful in 1m42s
Review found that execRedacted was applied only to the new remote-enumeration code
and left three pre-existing calls reading remote.origin.url through the plain,
unredacted exec: removeUserInfoFromHttpsOrigin, requireCredentialFreeHttpsOrigin, and
configureHttpsUrlRewriteForSshOrigin. A non-zero exit or timeout on any of those could
still have copied the credentialed URL into a WorktreeException message. Switches all
three to execRedacted; the set-url write in removeUserInfoFromHttpsOrigin is left on
plain exec with a comment explaining why (it writes the already-stripped URL, not a
read).

Widens the shared exec(Map, boolean, String...) overload to package-private, the same
test-seam pattern already used by the afterWorktreeAdded constructor parameter, and
adds a test that drives it directly with a synthetic failing command whose stdout
carries a marker (passed via env, not argv, so the always-printed command line can't
carry it) and asserts the marker never reaches the exception message.
2026-08-31 09:29:45 +07:00
Dai Ha bc99d64786 CB-185: fix ambiguous-pane message and unreachable-daemon abort in probeOwner (review)
CI / build (pull_request) Successful in 1m6s
CI / contract (pull_request) Successful in 1m5s
Lead review of PR #196 found two issues in CompositePeerLauncher.probeOwner:

1. The "more than one daemon claims this pane" throw kept the old
   pre-fix message ("no owning herdr daemon was recorded"), which was
   only true of the code it replaced. Reworded to say what actually
   happened: N configured herdr daemons report this pane, so it is
   genuinely ambiguous. Updated the one test pinning the old string.

2. probeOwner let list() propagate straight out of the probe loop, so
   one unreachable daemon aborted the whole probe and made a pane on a
   DIFFERENT, healthy daemon un-stoppable too — resurrecting the exact
   bug blocker 1 fixes. Now catches HerdrException per daemon, logs the
   exception class only, and treats that daemon as not knowing the pane
   so probing continues. New test proves this: verified it fails with
   the try/catch removed (HerdrException propagates and the stop that
   should succeed via the healthy daemon throws instead), then restored.

Full mvn clean install: 1019 tests, 0 failures, 0 errors.
2026-08-31 09:28:29 +07:00
Dai Ha 615af4ed0a CB-192 review fix: split the allow-list gap by what the scrub actually keeps
CI / build (pull_request) Successful in 1m5s
CI / contract (pull_request) Successful in 1m20s
Lead review on PR #194 found that the allow-list INFO wording claimed the
whole gap ("credential-shaped names on neither known: nor allow:") is blanked
by the scrub, without checking that against effectiveAllowed. effectiveAllowed
is a SUPERSET of known+allow — MemberEnvAllowList.derive also unions in every
profile's gitTokenEnv/gitHostEnv/tokenEnv/env: keys, and derivedAllowedNames
further unions in the spawn's own env keys — so a gap name can still be kept
by the derived list (e.g. a profile's tokenEnv names it) and reach the member
unblocked while the INFO said "no member pane keeps them". That inversion is
exactly what #192 exists to remove.

logCredentialGap now splits the gap with MemberEnvAllowList.keeps (the same
predicate the generated scrub itself evaluates, so this cannot drift from
what the scrub does): names it keeps get a WARN, guarded by the same
unprotectedGapLogged flag as the deny-by-default case (same severity — a name
reaching a member unprotected is equally serious either way); names it
blanks keep the existing INFO, guarded by allowListGapLogged. The
deny-by-default WARN text and the non-zsh fallback are untouched.

Added allowListWarnsWhenTheDerivedAllowListKeepsAnUncoveredName and
allowListSplitsAMixedGapBetweenTheWarnAndTheInfo to ClaudeCodeLauncherTest.
2026-08-31 09:28:10 +07:00
Dai Ha 045d229728 CB-185: fix two blockers to switching on memberHerdrSocket (#185)
CI / contract (pull_request) Successful in 46s
CI / build (pull_request) Successful in 1m39s
1. CompositePeerLauncher.stop() was permanently un-stoppable for any
   member that survived a daemon restart, because spawnedBy is in-memory
   only. On a cache miss with more than one configured herdr daemon, probe
   each distinct daemon's agent.list() for the pane instead of refusing
   outright: exactly one owner routes and caches; zero owners is treated
   as already-stopped (a no-op, matching the tolerance HerdrPeerLauncher
   already gives an already-gone pane); more than one owner is the
   genuine per-daemon-pane-id ambiguity and still throws.

2. FleetApp#healthz always reported the LEAD daemon's herdr version/
   protocol even when a second (member) daemon was configured, so a
   member-daemon protocol mismatch was invisible behind a green
   /healthz while every spawn silently failed. Added a separate "member"
   key alongside the unchanged "herdr" key, and a "protocolMismatch"
   flag when the two differ. Verified scripts/redeploy-fleetd.sh and
   scripts/rename-checkout.sh only check the HTTP status code and print
   the body verbatim — neither parses a specific field — so adding a key
   is safe.

Both fixes are covered by tests written to fail without the fix
(verified by reverting each fix and watching the new tests fail, then
restoring). Full `mvn clean install`: 1018 tests, 0 failures, 0 errors.
2026-08-31 09:20:14 +07:00
Dai Ha d1fd5700f5 CB-189: cover every remote, both URLs, and any non-SSH scheme in the credential check
CI / build (pull_request) Successful in 1m7s
CI / contract (pull_request) Successful in 1m7s
GitWorktrees only ever inspected origin's HTTPS fetch URL for embedded credentials. A
credential on any other remote, on a pushurl, or on a plain http:// URL passed through
unreported. Adds an additive, reporting-only check that enumerates every remote and both
its fetch and push URLs, flagging non-empty user-info on any non-SSH-family scheme.

The existing origin/https strip-and-refuse behaviour is untouched. The new check is
wrapped so it can never abort a provision, and on failure logs only the exception's
class, never its message, since the enumerating `git remote` call is not redacted.
Also adds execRedacted, an exec variant that never copies captured stdout into a
WorktreeException message, for commands whose stdout may itself be a credentialed URL.
2026-08-31 09:17:54 +07:00
Dai Ha 2d55b0b9a5 #168: correct the audit's MISSING claim — the feature is documented, its index row was not
CI / contract (push) Successful in 41s
CI / build (push) Successful in 1m42s
The memberHerdrSocket section exists at 11-Features.md:2174; what was absent was
its row in the index table. My omission when I added the section. Wiki fixed at
b24965c.
2026-08-31 09:17:41 +07:00
Dai Ha d89ae94a2e CB-192: fix false credential-gap WARN under allow-list+zsh, split its log guard
CI / contract (pull_request) Successful in 1m5s
CI / build (pull_request) Successful in 1m7s
logCredentialGap(creds) always emitted the WARN wording ("every member pane
inherits them UNBLOCKED"), even under memberCredentials.policy: allow-list on
a zsh login shell, where the generated ZDOTDIR scrub genuinely blanks the
name. The line reported the control working as though it were a hole.

Pass an effectiveAllowed set instead: null keeps the WARN (deny-by-default,
and the allow-list non-zsh fallback, where nothing is ever scrubbed); the
derived allow-list set (only reachable after applyEnvironmentAllowListPolicy's
own zsh gate) selects a new INFO wording that says the scrub will blank the
name instead of claiming it is inherited unblocked.

Also split the single credentialGapLogged AtomicBoolean into two guards
(unprotectedGapLogged / allowListGapLogged) — one per report kind. Since
memberCredentials is a live, re-read-per-spawn supplier, a shared flag let a
harmless allow-list INFO on one spawn permanently suppress a later spawn's
real deny-by-default WARN after a policy reload.

Fixes gitea #192.
2026-08-31 09:17:20 +07:00
ltms e6193c4098 Merge pull request '#168: audit current wiki snapshot' (#193) from worker/cb-168-wiki-audit-3ef3d1-3 into main
CI / build (push) Successful in 1m6s
CI / contract (push) Successful in 1m5s
2026-08-31 04:17:07 +02:00
Dai Ha b66f0677ed #168: audit current wiki snapshot
CI / contract (pull_request) Successful in 44s
CI / build (pull_request) Successful in 1m37s
2026-08-31 09:14:22 +07:00
ltms 23ada1981e CB-185: route members to a separate herdr daemon (#186)
CI / build (push) Successful in 1m6s
CI / contract (push) Successful in 10m30s
2026-08-29 01:24:31 +02:00
ltms a22480c117 CB-185: route PaneLocator, StatusRefiner and FleetApp to the right herdr daemon (#188)
CI / build (pull_request) Successful in 1m3s
CI / contract (pull_request) Successful in 1m7s
2026-08-29 01:24:25 +02:00
Dai Ha 24f404f989 CB-185: fix three connection-identity/status/health gaps a second herdr daemon exposes
memberHerdrSocket splits lead operations from member operations onto two herdr
daemons. Three seams still assumed one shared daemon and broke silently when the
two clients differ (all three collapse to today's behaviour when they are the
same object):

1. ConnectionIdentity's PaneLocator was pinned to the member daemon only, so a
   lead's own MCP connection (which lives on the LEAD daemon) resolved to
   terminal == null, breaking fleet_reply/fleet_ask/fleet_whoami for a lead.
   PaneLocator now searches the lead client first, then the member client.

2. StatusPoller's StatusRefiner was pinned to the member daemon, so refining an
   UNKNOWN status for a lead target read the wrong daemon's pane content and
   never left UNKNOWN, wedging status-gated delivery to that lead forever.
   StatusRefiner gained a refine(target, raw, control) overload and the poller
   now refines through the same AgentControl the raw status was sampled from.

3. FleetApp was constructed with the raw lead-only herdr client, so /healthz
   stayed green while the member daemon was down (every spawn then fails
   invisibly) and GET /sessions silently dropped every member workspace.
   FleetApp now takes both clients: healthz requires both to answer, sessions
   merges workspaces from both.

Each fix has a test proven to fail without it (verified by reverting the
production change and re-running): FleetdConnectionIdentityConstructionTest /
FleetdFleetAppConstructionTest assert the actual Fleetd.java wiring (the same
technique as FleetdHerdrControlConstructionTest); StatusPollerRoutingTest and
the new PaneLocatorTest/FleetAppTwoDaemonTest cases exercise the real
production classes end to end rather than a hand-built object graph.
2026-08-29 06:20:32 +07:00
ltms a237fbff9d #185: refuse an unowned paneId when more than one herdr daemon could own it (#187)
CI / build (push) Successful in 1m4s
CI / contract (push) Successful in 1m17s
2026-08-29 01:10:21 +02:00
Ha Trong Dai 17af61e8dd CB-185: route message status by target
CI / contract (pull_request) Successful in 1m4s
CI / build (pull_request) Successful in 1m40s
2026-08-28 09:45:39 +07:00
Ha Trong Dai 6af87b6ad6 CB-185: share routed herdr controls
CI / contract (pull_request) Successful in 46s
CI / build (pull_request) Successful in 1m40s
2026-08-28 09:42:45 +07:00
Ha Trong Dai fc655e78c2 CB-185: route members to separate herdr
CI / contract (pull_request) Successful in 37s
CI / build (pull_request) Successful in 1m28s
2026-08-28 09:37:44 +07:00
31 changed files with 2253 additions and 92 deletions
+166
View File
@@ -0,0 +1,166 @@
# CB-137 / fleetd issue #137 — report
## Real root cause (not the hypothesis in the ticket)
I read `MessageService.java` and `Rendezvous.java` before changing anything. The mechanism is real,
but the exact place it happens is `MessageService.answer()`, not "the reply goes to the inbox on
purpose" in general.
1. A lead delegates with `fleet_send{wait:false}` → `sendAsync()` creates a `Task` and runs `send()`
on a background virtual thread with a 30-minute internal budget (`ASYNC_TIMEOUT_MS`).
2. The worker calls `fleet_ask`. That resolves the open rendezvous waiter with `Kind.QUESTION`, so
`send()` returns immediately and the `Task` is left open (its `future` stays unresolved — see the
comment in `sendAsync`'s lambda: "Keep the accepted owner until answer() finishes it").
3. The lead answers with `fleet_send{turnId, content}`. This calls `FleetMcp.answer()` →
`MessageService.answer(turnId, content, timeout)`. The `timeout` here is **not** the generous
30-minute async budget — it is the MCP tool's own bounded wait: `DEFAULT_TIMEOUT_MS = 25_000`,
clamped to at most `MAX_TIMEOUT_MS = 120_000` (`FleetMcp.java:71-72,495,512`). This is the same
~60–120s window every blocking `fleet_send` call is capped at (documented elsewhere as "the
caller's own MCP client call timeout").
4. `answer()` opens a **fresh** rendezvous waiter for the worker session and blocks on it for at most
that window. If the worker's resumed turn takes longer than that to actually finish (very
plausible — the resumed turn can mean more edits, a build, a commit, a push, opening a PR), the
wait times out. On timeout, `answer()`'s `finally` block unconditionally calls
`rendezvous.close(workerSession, reply)`, **removing the waiter from the map**, and returns
`Outcome.TIMED_OUT_WORKING` to the lead.
5. The worker keeps working, unaware anything happened, and eventually calls `fleet_reply`. That
reaches `MessageService.reply(session, content)`, which tries `rendezvous.resolve(session,
content)` — but the waiter was already closed in step 4, so `resolve` returns `false`. `reply()`
then falls back to `inbox.publish(...)` and marks `strandedReplies.put(session, true)`
(CB-640 bookkeeping) — the reply is safely held, but **the async `Task`'s `future` is never
completed**.
6. `fleet_poll{ticket}` keeps returning `PENDING` forever (the `Task` never resolves) — until the
lead eventually calls `fleet_stop`. That fires `sessions.onRelease` → `messages.abandon(target,
reason)` (`Fleetd.java:481-497`), where `reason` is built with the exact text from the bug report
("the worker session was released before it replied; worktree=... branch=... snapshot=...",
`Fleetd.java:484-487`). `abandon()`'s loop finds the still-open `Task` (`question == null`, future
not done) and completes it as `WORKER_FAILED` with that misleading reason — even though the
worker's real reply is sitting, intact, in the inbox the whole time.
So: the reported behaviour is correct, and the specific trigger is `answer()`'s own bounded wait
being shorter than the worker's real resumed-turn time — not anything to do with the ~55s
`fleet_ask` window itself (that part, issue #61, is untouched).
## Fix
Two changes in `fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java`, both scoped to the
ticket/reply routing and the terminal-state text — `fleet_ask`'s own window and mechanics are
untouched.
**1. `reply()` — priority 1 (the ticket resolves with the real reply).**
Before falling back to the inbox, `reply()` now looks for an async `Task` that is specifically in the
"already answered but not yet resolved" state (`question == null`, `turnId != null` — set once
`answer()` has cleared the question but before anything completed the future, `!future.isDone()`).
If one exists for this `target`, the worker's reply completes that `Task`'s future directly as
`Outcome.REPLIED` with the real content, and the reply never touches the inbox at all. A task that
was never asked has `turnId == null` and can never match, so ordinary (no-`fleet_ask`) delegations
are unaffected — they already resolve through the pre-existing rendezvous fast path.
I chose this over leaving `answer()`'s own timeout behaviour untouched and instead keeping its
rendezvous waiter open in the background: that alternative works but reopens the "at most one
waiter per session" invariant (`Rendezvous.open` throws on a double-open) to a new class of races
with a fresh send arriving mid-window. The `send()` path already guards against sending into an
answered-but-still-resolving worker via `hasAsyncQuestion(target)` (checks `asyncTasksByTurn`,
which still holds the task until it resolves), so routing through `reply()` gets the same protection
without touching `answer()`'s waiter lifecycle at all — the smaller, safer diff.
**2. `abandon()` — priority 3 (required independently, "even if you fix (1)").**
Before marking any of a released target's still-open tasks `WORKER_FAILED`, `abandon()` now checks
`hasStrandedReply(target)` (the existing CB-640 fact — true whenever the *last* `reply()` for this
target fell through to the inbox). If true, it drains the inbox (`recoverStrandedReply`) and — if it
actually finds a message — completes the task as `REPLIED` with that real content instead of writing
the failure. This is deliberately a **separate** check from fix 1: fix 1 already prevents the
inbox-stranding from happening in the exact scenario this ticket describes, so by the time
`abandon()` runs the task is normally already resolved and `abandon()`'s `complete()` call is a
harmless no-op. This second check exists so that if some *other* future path ever strands a reply
in the inbox without resolving its ticket, `abandon()` still refuses to report a false failure —
"if a reply reached any sink for that turn, the terminal state is done," per the ticket. I verified
both are required by disabling each independently and confirming the two new tests fail (see below).
**Priority 4 (the snapshot/worktree hint).** Handled as a consequence of both fixes rather than a
separate branch: once a task resolves as `REPLIED` (via either fix), `abandon()` never calls
`new Reply(Outcome.WORKER_FAILED, reason)` for that task at all, so the "the worker session was
released before it replied; worktree=... branch=... snapshot=..." text is never constructed or
attached to that ticket's outcome. It still appears, correctly, for a task that never got a reply
(the existing `abandonFailsEveryPendingAsyncTicketForTheReleasedTarget` /
`anAbandonedAsyncTaskPollsAsFailedNotPending` tests still pass unchanged).
**Priority 2** was not needed — fix 1 makes `fleet_poll{ticket}` return the actual reply (the
higher-priority option), so I did not fall back to "the ticket merely resolves as done with no
content."
## Tests — driven through the real delegation path, not the reply sink directly
Both new tests in `fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java` go through
`sendAsync` → `injectDelivery` → `ask` → `answer` (with a short timeout, so it genuinely times out,
mirroring the ~25–120s real MCP-call bound vs. a longer resumed turn) → `reply` → `poll`/`abandon`.
No test constructs a `Reply` and hands it to a sink directly.
- `aReplyAfterAnswerTimesOutStillCompletesTheAsyncTicket` — asserts `fleet_poll{ticket}` (via
`messages.poll`) reaches `Phase.DONE` with the worker's actual reply text and
`replySource() == "reply"`, and that `hasStrandedReply(T)` stays `false` (proves the reply never
touched the inbox at all — fix 1 caught it).
- `fleetStopAfterAnOrphanedReplyDoesNotFailTheTicket` — same setup, then calls `abandon(T, "the
worker session was released before it replied")` (what `fleet_stop` triggers) and asserts it
returns `false` (no failure recorded) and the ticket still polls `DONE` with the real reply.
**Proof both fail without the change.** I temporarily short-circuited both new private methods
(`askAnsweredAsyncTask` → always `null`, `recoverStrandedReply` → always `null`) — i.e. disabled
both fixes — and ran just these two tests:
```
[ERROR] Tests run: 2, Failures: 2, Errors: 0, Skipped: 0
dev.ltms.fleet.msg.MessageServiceTest.aReplyAfterAnswerTimesOutStillCompletesTheAsyncTicket
org.opentest4j.AssertionFailedError: expected: <DONE> but was: <PENDING>
dev.ltms.fleet.msg.MessageServiceTest.fleetStopAfterAnOrphanedReplyDoesNotFailTheTicket
org.opentest4j.AssertionFailedError: a reply already arrived, so nothing here is a genuine failure
==> expected: <false> but was: <true>
```
This is the exact bug: the ticket stays `PENDING` forever, and `abandon()` reports `true` (a
failure) even though a reply had already arrived. I then restored both fixes (verified with
`grep -n "TEMP #137-proof"` finding nothing) and re-ran — both pass.
## Build
Ran from `fleetd/`, unpiped, full output read (not `| tail`):
```
mvn clean install
...
[INFO] Tests run: 1039, Failures: 0, Errors: 0, Skipped: 0
[INFO] BUILD SUCCESS
[INFO] Total time: 36.315 s
```
Main was at 1037 tests; this branch adds the 2 new tests above → 1039, all green, `exit=0`.
## What I could NOT check
- No IDE tooling is mounted for me (worker), so no `ide_diagnostics`/IntelliJ inspection pass — only
`mvn clean install` (compiler + full test suite), as the worker procedure allows.
- I cannot restart the daemon or dogfood this live — I have no forge/daemon control. This is
unverified against a real herdr pane, a real MCP client's ~60s call cap, or a real worker session;
everything above is verified only through the JUnit fixture's simulated timing
(`FakeHerdr`/`injector.onStatus`/direct `messages.answer(...,150)` calls), not a live fleet.
A primary should still consider a short live dogfood (an async delegation that asks, gets answered,
and takes longer than ~2 minutes to reply) before calling this closed.
- I did not touch, and did not re-verify, the `fleet_ask` ~55s window itself (issue #61) — out of
scope per the brief.
## Scope note (not investigated further)
`answer()`'s nested/double-`fleet_ask` case (the worker asks a second question before ever
replying to the first answer) has some pre-existing behaviour around which `turnId` a `QUESTION`
resolution gets attributed to that I did not fully untangle — it predates this change, my fix does
not touch it, and it is unrelated to the reported defect. Flagging only; not investigated further.
## Handoff
- Branch: `worker/cb-137-ask-ticket-e7760c-2`
- Worktree root: `/Users/dai.ha/LTMS/.bridged-worktrees/734324-2`
- Files changed:
- `fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java`
- `fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java`
- `REPORT-cb137.md` (this file)
- Build: `Tests run: 1039, Failures: 0, Errors: 0, Skipped: 0` / `BUILD SUCCESS` (verbatim above)
+208
View File
@@ -0,0 +1,208 @@
# Wiki audit for #168
**Source checked:** `.wiki-snapshot/` at `68e32c6` (2026-08-31). I did not use
`wiki/`. Code references below are from the current `fleetd` source tree. A quoted
line is a concrete claim that needs correction, unless the table says `KEEP`.
| Page | Verdict | One-line reason |
|---|---|---|
| `Home.md` | REVISE | Good overview, but it still names the retired product. |
| `_Sidebar.md` | REVISE | The heading still says `claude-bridge`. |
| `1-Architecture.md` | REBUILD | Its component contract mixes current names with removed tools, routes, and planned backends. |
| `2-Message-Server.md` | REBUILD | The claimed MCP schema, mount command, REST/SSE surface, and fallback paths are pre-build design. |
| `3-Approaches.md` | REVISE | Useful research history, but it presents unbuilt AgentAPI as a selectable fallback. |
| `4-Setup.md` | RETIRE | It is an intentional stub that only redirects to chapter 13. |
| `5-Operations.md` | RETIRE | It is an intentional stub that only redirects to chapter 13. |
| `6-Team.md` | REBUILD | It teaches role-addressed sends and a Claude-only team model that the shipped API does not have. |
| `7-Use-Cases.md` | REBUILD | Its flagship flow depends on removed `ccs` profiles and removed send parameters. |
| `8-Roadmap.md` | REBUILD | It is a historical plan, but it presents old implementation choices and planned work as the current stack. |
| `9-Implementation.md` | REBUILD | Its package, class, endpoint, and outcome map has drifted from the source. |
| `10-Cross-Host-Messaging.md` | REVISE | It labels most federation work proposed, but misses the shipped `coordinator:` lead channel. |
| `11-Features.md` | REVISE | It is the right catalogue, but code-path names are old and it misses the second-herdr-daemon capability. |
| `12-Claude-to-OpenCode.md` | REVISE | The porting guide is mostly current, but calls the product and spawned-member path a bridge. |
| `13-User-Guide.md` | REVISE | It is the best operator page, but needs the product rename and the second-herdr-daemon setup. |
## Pages needing work
### `Home.md` — REVISE
- Quote: `# claude-bridge` (line 1) and `` `claude-bridge` keeps`` (line 11).
The product is `fleet` / `fleetd`. The MCP server identifies itself as `fleet` in
`fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java:313-315`.
- Quote: `AgentAPI ... swappable fallback injector` (lines 73-76).
There is no AgentAPI implementation under `fleetd/src/main/java`; the actual
launchers are selected by `Profile.kind` in
`fleetd/src/main/java/dev/ltms/fleet/config/FleetConfig.java:265-270`.
### `_Sidebar.md` — REVISE
- Quote: `### 📖 claude-bridge` (line 1).
Rename it to `fleet`. `FleetMcp` registers the current product-facing tool set at
`fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java:301-326`.
### `1-Architecture.md` — REBUILD
- Quote: `` `claude-bridge` lets`` (line 3). The product was renamed; the MCP
server name is `fleet` (`FleetMcp.java:313-315`).
- Quote: ``fleet_read`` in the tool list (line 102). No such tool is registered.
The complete registered list is `fleet_send` through `fleet_whoami` at
`FleetMcp.java:301-326`; `fleet_read` is absent.
- Quote: `SSE (GET /events)` (line 143). `FleetApp.build()` registers no `/events`
route; its routes are listed at `FleetApp.java:143-159`.
- Quote: `Redis Streams / NATS JetStream, or an embedded queue` (line 106).
The shipped durable inbox is AMQP, configured by `broker`, at
`FleetConfig.java:49-50` and `FleetConfig.java:655-714`.
- Quote: `AgentAPI (fallback)` (line 107). No AgentAPI adapter exists; shipped
launcher kinds are `claude-code` and `opencode` (`FleetConfig.java:265-270`).
### `2-Message-Server.md` — REBUILD
- Quote: `claude mcp add --transport http bridge http://127.0.0.1:8080/mcp`
(line 67). The daemon defaults to port `8765` in `FleetConfig.java:183-187`,
and identifies its server as `fleet` at `FleetMcp.java:313-315`.
- Quote: ``fleet_send(message, target?, {block, timeout_seconds, auto_spawn,
turn_id})`` (line 80). The real parameters are `sessionId`, `content`,
`timeoutMs`, `wait`, `turnId`, and `coordId` (`FleetMcp.java:1096-1108`).
- Quote: ``fleet_read(target, source)`` (line 85). It is not registered; see the
complete registration at `FleetMcp.java:301-326`.
- Quote: `docs/MCP-Contract.md ... normative` (lines 87-88). That is not a valid
reference: only §6 is current, as the current operator guide itself says at
`.wiki-snapshot/13-User-Guide.md:466`.
- Quote: `SSE (GET /events)` (line 45). No route exists in the built REST surface,
`FleetApp.java:143-159`.
### `3-Approaches.md` — REVISE
- Quote: `AgentAPI ... remains a swappable fallback injector` (lines 78-84).
It was never built. The shipped adapter selection is only `claude-code` or
`opencode` (`FleetConfig.java:265-270`). Keep it as discarded research, not an
operational fallback.
- Quote: `claude-bridge` (line 109). Rename the product to `fleet`; the runtime
package is `dev.ltms.fleet`, for example `FleetMcp.java:1`.
### `4-Setup.md` — RETIRE
It is a 25-line redirect and says its procedure was never written (lines 3-9).
Chapter 13 is the maintained install procedure. Keeping a second navigation page
adds no working documentation.
### `5-Operations.md` — RETIRE
It is a 35-line redirect and says its runbook was never written (lines 3-14).
Chapter 13 now owns run and recovery instructions.
### `6-Team.md` — REBUILD
- Quote: `fleet_send {role: w-claude, prompt: A}` (line 98). `fleet_send` accepts
`sessionId` and `content`, not `role` or `prompt` (`FleetMcp.java:1096-1108`).
- Quote: `some on Claude, some on the remote local LLM` (lines 3-5) and `Every
worker is ... Claude Code` (line 25). `opencode` is a first-class launcher kind,
not a Claude worker (`FleetConfig.java:265-270`).
- Quote: `fleetd's concurrency policy` (line 121). The configured capacity control
is per-profile `maxLoad` (`FleetConfig.java:251-264`), not the role routing model
described here.
### `7-Use-Cases.md` — REBUILD
- Quote: `ccs profile` (line 10), `ccs + herdr` (line 22), and `ccs-spawn`
(line 45). The configuration has `profiles` and `fleet`, not `ccs`:
`FleetConfig.java:34-58` and `FleetConfig.java:81-101`.
- Quote: `fleet_send({"to", "kind", "body", "block"})` (lines 55-62).
None of those are the shipped send parameters. The schema is
`FleetMcp.java:1096-1108`.
- Quote: `fleet_list() → { "profiles": ... }` (lines 74-80). `fleet_list` is a
roster view; `fleet_profiles` is the configured-backend view, as registered at
`FleetMcp.java:307-311` and described at `FleetMcp.java:1176-1182`.
### `8-Roadmap.md` — REBUILD
- Quote: `Java 21+` (line 43). The current project guidance and source use Java 25;
the `FleetConfig` source itself uses Java 25 unnamed lambda parameters, for
example `FleetConfig.java:102`.
- Quote: `herdr 0.7.0 / protocol 14` (line 46). The current REST health endpoint
reports the live protocol returned by herdr (`FleetApp.java:240-244`), while the
current operator guide records protocol 19 at
`.wiki-snapshot/13-User-Guide.md:76-85`.
- Quote: `ccs <profile> claude` and `ccs env <profile>` (lines 47-48). Shipped
configuration uses `Profile` records and launcher `kind`,
`FleetConfig.java:313-330` and `FleetConfig.java:265-270`.
- Quote: `Redis Streams via Lettuce` (line 50). The actual durable inbox is AMQP
`broker`, `FleetConfig.java:655-714`.
### `9-Implementation.md` — REBUILD
- Quote: `rest.FleetdApp` and `mcp.BridgeMcp` (lines 29-30). The classes are
`rest.FleetApp` and `mcp.FleetMcp` (`FleetApp.java:46`; `FleetMcp.java:67`).
- Quote: `dev.ltms.fleetd` (line 67). The source package is `dev.ltms.fleet`
(`FleetMcp.java:1`).
- Quote: `WorkerPresence` (line 110). The current class is `MemberPresence`, as
imported and used by `FleetMcp` at `FleetMcp.java:12` and `465-469`.
- Quote: the outcome list ending in `STALE_TURN` (lines 128-131). The code also
has `BACKEND_EXHAUSTED` (`FleetMcp.java:550-554`) and async `ASKING` handling
(`FleetMcp.java:664-668`).
- Quote: `FleetdApp` (line 207) and `FleetdConfig` (line 211). These names do not
resolve; current classes are `FleetApp` and `FleetConfig`.
### `10-Cross-Host-Messaging.md` — REVISE
- Quote: the chapter says the cross-host fabric is proposed except for the
single-host inbox (lines 3-8). Cross-host **lead-to-lead** delivery shipped:
`fleet_send` accepts `coordId` (`FleetMcp.java:1094-1107`) and publishes it at
`FleetMcp.java:616-641`; configuration has `coordinator` at
`FleetConfig.java:74-78` and `99-101`.
- Quote: `bridge.dlx` (line 90). This product name is stale. The shipped lead path
uses `LeadChannel`, not the proposed exchange flow (`FleetMcp.java:95-96` and
`616-641`). Keep the proposed federation design, but add a clear shipped/proposed
boundary for CB-637.
### `11-Features.md` — REVISE
- Quote: `mcp/BridgeMcp` (line 22), `config/FleetdConfig` (lines 25-27), and other
index references. These paths no longer resolve; the source classes are
`mcp/FleetMcp` (`FleetMcp.java:67`) and `config/FleetConfig`
(`FleetConfig.java:81`).
- Quote: `fleet_whoami` returns only `primary` or `worker` (lines 99-100).
It also returns `architect` (`FleetMcp.java:1235-1244`).
- The page needs the missing separate member-herdr-daemon feature listed below.
### `12-Claude-to-OpenCode.md` — REVISE
- Quote: `same bridge mount` (line 5) and `a bridge-spawned worker` (line 94).
Rename the product path to `fleet`. The daemon exposes the MCP server as `fleet`
(`FleetMcp.java:313-315`), and profiles select OpenCode with `kind: opencode`
(`FleetConfig.java:332-335`).
- Quote: the sample mount name is `fleetd` (line 67). The server name is `fleet`;
update the sample to avoid teaching a second product name.
### `13-User-Guide.md` — REVISE
- Quote: `The bridge is the only channel` (line 63). The invariant is correct, but
the product term needs the `fleet` rename. The daemon's MCP server name is
`fleet` (`FleetMcp.java:313-315`).
- Quote: it describes one herdr socket (lines 72-85). It needs the optional
`memberHerdrSocket` setup and two-daemon health meaning. The config key is in
`FleetConfig.java:34-37`, and `/healthz` checks both daemons when configured at
`FleetApp.java:210-245`.
## MISSING
`11-Features.md` has a body section for **routing members through a separate herdr daemon**
(`## memberHerdrSocket`, line 2174), but **no row in the index table** at the top of the page
(lines 20-95). That table is how the page is meant to be read, so a capability absent from it is
effectively undiscoverable. Lead note: this is my own omission — I added the section on 2026-08-31
and did not add the matching row. Fixed in the wiki at `68e32c6`'s successor.
The original audit stated the feature had no entry at all. That was wrong: the section exists. The
gap is the index row. Recorded here rather than silently corrected, because the difference matters —
"undocumented" and "documented but unindexed" are different jobs.
Evidence for the feature itself: `FleetConfig.java:34-37` and `FleetApp.java:103-115`, `210-245`,
and `247-263`.
## Audit method and coverage
I checked all 15 pages. I checked concrete tool, route, config, class, file, and
product-name claims claim-by-claim on 11 pages: Home, Sidebar, 1, 2, 4, 5, 6, 7, 9,
11, and 13. I skimmed the remaining four long historical or research pages (3, 8, 10,
12), then checked their concrete claims that affect the verdict. This is an audit of
the supplied snapshot, not a wiki rewrite.
+3
View File
@@ -114,6 +114,9 @@ bind:
# (${HERDR_SOCKET_PATH:-~/.config/herdr/herdr.sock}).
herdrSocket: ~/.config/herdr/herdr.sock
# Optional socket for member panes. Omit this to use herdrSocket for both leads and members.
# memberHerdrSocket: /Users/member/.config/herdr/herdr.sock
# How member sessions are spawned. Define one or more named profiles (backends) under
# `profiles`; each key is the profile name (also the ccs profile). A profile says only WHICH
# BACKEND — model, CLI adapter, credentials, cost. It says nothing about what a member spawned on
+27 -15
View File
@@ -7,6 +7,7 @@ import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.herdr.HerdrException;
import dev.ltms.fleet.herdr.HerdrRouter;
import dev.ltms.fleet.herdr.LeadTabScanner;
import dev.ltms.fleet.lead.LeadLauncher;
import dev.ltms.fleet.herdr.PaneLocator;
@@ -149,9 +150,12 @@ public final class Fleetd {
: UnixSocketHerdrClient.defaultSocketPath();
UnixSocketHerdrClient herdr = UnixSocketHerdrClient.connect(socket, new com.fasterxml.jackson.databind.ObjectMapper());
AgentControl agents = new AgentControl(herdr);
WorkspaceControl spaces = new WorkspaceControl(herdr);
UnixSocketHerdrClient memberHerdr = cfg.memberHerdrSocket() != null && !cfg.memberHerdrSocket().isBlank()
? UnixSocketHerdrClient.connect(Path.of(cfg.memberHerdrSocket()), new com.fasterxml.jackson.databind.ObjectMapper())
: herdr;
AtomicReference<Supplier<Map<String, String>>> leadsRef = new AtomicReference<>(Map::of);
HerdrRouter router = new HerdrRouter(herdr, memberHerdr,
target -> leadsRef.get().get().containsKey(target));
// CB-402: one adapter per configured peer kind, fronted by a composite router. A profile's
// `kind:` selects its adapter — claude-code (the default) and opencode partition the profile
// set — and the composite dispatches each SPI call to the adapter that owns the profile/pane.
@@ -169,14 +173,14 @@ public final class Fleetd {
// bridge configured with no workers, or opencode-only, still has a well-defined base adapter)
// unless opencode is the only kind configured.
if (!claudeProfiles.isEmpty() || opencodeProfiles.isEmpty()) {
adapters.add(new ClaudeCodeLauncher(agents, spaces, guard,
adapters.add(new ClaudeCodeLauncher(router.memberAgents(), router.memberSpaces(), guard,
claudeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
() -> config.get().fleet(),
() -> config.get().memberCredentials()));
}
if (!opencodeProfiles.isEmpty()) {
adapters.add(new OpenCodeLauncher(agents, spaces,
adapters.add(new OpenCodeLauncher(router.memberAgents(), router.memberSpaces(),
opencodeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
() -> config.get().fleet(),
@@ -271,6 +275,7 @@ public final class Fleetd {
// operational cadence, not identity, so there is no correctness reason to give every
// lead its own scanner.
int scanIntervalSeconds = leaders.values().iterator().next().scanIntervalSeconds();
// This must use the lead daemon: scanning member tabs would demote the lead to a worker.
leads = new LeadTabScanner(herdr, tabToName, Set.of(),
TimeUnit.SECONDS.toNanos(scanIntervalSeconds), System::nanoTime);
log.info("lead scan: tabs {} host a lead (rescan every {}s, shared fleet space)",
@@ -278,13 +283,14 @@ public final class Fleetd {
} else {
leads = () -> leadTerminals;
}
leadsRef.set(leads);
// CB-558: start any declared lead that is not already running. After the scanner is built,
// because both read the same tab labels and the ordering makes that dependency visible; and
// only when herdr answered, because the launcher's whole safety property is that it can
// count live leads first — it must never guess and risk a second orchestrator.
if (herdrUp && !leaders.isEmpty()) {
int launched = new LeadLauncher(agents, spaces, cfg).ensureLeads();
int launched = new LeadLauncher(router.leadAgents(), router.leadSpaces(), cfg).ensureLeads();
if (launched > 0) {
log.info("lead auto-launch: {} lead(s) started", launched);
}
@@ -340,6 +346,7 @@ public final class Fleetd {
+ "BACKEND_EXHAUSTED): {}", credentialId,
cfg.quarantineCooldownSeconds(), profile.profile(), reason);
});
AgentControl agents = router.memberAgents();
CompletionResolver completion = new CompletionResolver(agents, rendezvous, exhaustedPatterns, exhaustionSink);
// CB-113: deliver only to an available worker (its MCP is connected), never its boot window.
// CB-301: the manager's presence bridge records availability and drives SPAWNING → READY.
@@ -381,9 +388,9 @@ public final class Fleetd {
}
};
Predicate<String> deliverable = deliverableTo(presence, leads);
Injector injector = new Injector(agents, turnListener, deliverable,
Injector injector = new Injector(router, turnListener, deliverable,
presence::forget);
StatusPoller poller = new StatusPoller(agents, injector, Injector.POLL_INTERVAL_MILLIS);
StatusPoller poller = new StatusPoller(router, injector, Injector.POLL_INTERVAL_MILLIS);
poller.start();
// CB-307: reply inbox. A broker: block selects the AMQP-backed durable adapter; absent (or
@@ -420,7 +427,7 @@ public final class Fleetd {
// are counted at their single funnel rather than at each of the two caller-facing surfaces.
// CB-512: the push loop takes it too, so nudge outcomes (delivered|exhausted) are counted.
Metrics metrics = FleetMetrics.create(sessions, replyInbox);
var pushLoop = new ReplyPushLoop(primaryRegistry, agents, replyInbox,
var pushLoop = new ReplyPushLoop(primaryRegistry, router.leadAgents(), replyInbox,
pushScheduler, maxReminders, backoffMs, metrics);
// CB-551: idle-lead heartbeat. Opt-in; absent `leadHeartbeat:` this is never constructed, so
// an upgraded daemon cannot silently start spending subscription on nudging an idle lead.
@@ -430,7 +437,7 @@ public final class Fleetd {
Thread.ofVirtual().name("bridge-heartbeat-").unstarted(r));
if (cfg.leadHeartbeat() != null) {
var hb = cfg.leadHeartbeat();
heartbeat = new LeadHeartbeatLoop(primaryRegistry, agents, replyInbox, sessions::roster,
heartbeat = new LeadHeartbeatLoop(primaryRegistry, router.leadAgents(), replyInbox, sessions::roster,
pushLoop, heartbeatScheduler, System::nanoTime,
TimeUnit.SECONDS.toNanos(hb.idleAfterSeconds()), hb.backoffMs(), hb.quietNudgeCap(),
metrics);
@@ -439,7 +446,7 @@ public final class Fleetd {
heartbeat = null;
heartbeatScheduler.shutdownNow();
}
MessageService messages = new MessageService(agents, injector, rendezvous, replyInbox,
MessageService messages = new MessageService(router, injector, rendezvous, replyInbox,
pushLoop, metrics);
// Health is a slow whole-fleet observer. Keep it separate from the 250ms delivery poller.
@@ -491,8 +498,11 @@ public final class Fleetd {
// MCP server face (CB-105): fleet_send/fleet_reply/fleet_status, mounted at /mcp.
// Caller identity is resolved from the connection (peer PID → herdr pane), not arguments.
// CB-185: a caller's pane can live on either daemon (a lead's on the lead daemon, a
// member's on the member daemon) — search both, lead first. Collapses to one scan when
// memberHerdrSocket is unset (herdr == memberHerdr).
ConnectionIdentity identity = new ConnectionIdentity(
new PaneLocator(herdr), new LsofPeerPidLookup(), new LsofProcessCwdLookup());
new PaneLocator(herdr, memberHerdr), new LsofPeerPidLookup(), new LsofProcessCwdLookup());
// CB-501: one resolver behind both entry paths. Worker identity still comes from the
// connection and is never token-gated, so enabling token mode cannot lock the fleet out.
@@ -538,7 +548,7 @@ public final class Fleetd {
if (leadMailbox != null) {
var leadCoordScheduler = Executors.newSingleThreadScheduledExecutor(r ->
Thread.ofVirtual().name("bridge-leadcoord-").unstarted(r));
leadCoordLoop = new LeadCoordLoop(leadMailbox, agents, leads, leadCoordScheduler,
leadCoordLoop = new LeadCoordLoop(leadMailbox, router.leadAgents(), leads, leadCoordScheduler,
LEAD_COORD_INTERVAL_MS);
leadCoordLoop.start();
leadCoordSchedulerRef = leadCoordScheduler;
@@ -589,10 +599,12 @@ public final class Fleetd {
log.debug("lead mailbox close: {}", e.toString());
}
}
herdr.close();
router.close();
}));
Javalin app = new FleetApp(herdr, workers, sessions, messages, presence, mcp.servlet(),
// CB-185: give FleetApp both daemons — /healthz must require both to answer and
// GET /sessions must merge across both, or a down/unpolled member daemon is invisible.
Javalin app = new FleetApp(herdr, memberHerdr, workers, sessions, messages, presence, mcp.servlet(),
callers, metrics, deliverable).build();
app.start(cfg.bind().host(), cfg.bind().port());
log.info("fleetd listening on {}:{}, herdr socket {}",
@@ -71,7 +71,7 @@ public final class ConfigRef implements Supplier<FleetConfig> {
/** Keys that cannot change under a running daemon — see the class doc. */
private static final Set<String> COLD_KEYS =
Set.of("bind", "herdrSocket", "broker", "auth");
Set.of("bind", "herdrSocket", "memberHerdrSocket", "broker", "auth");
private final Path path;
private final AtomicReference<FleetConfig> current;
@@ -190,6 +190,9 @@ public final class ConfigRef implements Supplier<FleetConfig> {
if (!Objects.equals(old.herdrSocket(), fresh.herdrSocket())) {
changed.add("herdrSocket");
}
if (!Objects.equals(old.memberHerdrSocket(), fresh.memberHerdrSocket())) {
changed.add("memberHerdrSocket");
}
if (!Objects.equals(old.broker(), fresh.broker())) {
changed.add("broker");
}
@@ -32,7 +32,8 @@ import java.util.Set;
* silently dropping a whole block is indistinguishable from honouring it.
*
* @param bind REST/MCP listen host:port
* @param herdrSocket path to herdr's Unix socket ({@code null} → client default)
* @param herdrSocket path to the lead herdr Unix socket ({@code null} → client default)
* @param memberHerdrSocket optional member herdr Unix socket ({@code null}/blank → lead socket)
* @param profiles named backend profiles, keyed by profile name (multi-backend fleet). A
* profile answers <em>which backend</em> — model, CLI adapter, credentials,
* cost. It says nothing about what the member spawned on it is for; that is
@@ -80,6 +81,7 @@ import java.util.Set;
public record FleetConfig(
Bind bind,
String herdrSocket,
String memberHerdrSocket,
Map<String, Profile> profiles,
Guard guard,
String worktreeRoot,
@@ -105,7 +107,7 @@ public record FleetConfig(
LeadHeartbeat leadHeartbeat, Health health, String placement, Auth auth,
ConfigReload configReload, Integer quarantineCooldownSeconds,
MemberCredentials memberCredentials) {
this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
this(bind, herdrSocket, null, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth,
configReload, quarantineCooldownSeconds, memberCredentials, null);
}
@@ -116,9 +118,9 @@ public record FleetConfig(
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
LeadHeartbeat leadHeartbeat, Health health, String placement, Auth auth,
ConfigReload configReload, Integer quarantineCooldownSeconds) {
this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
this(bind, herdrSocket, null, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth,
configReload, quarantineCooldownSeconds, null);
configReload, quarantineCooldownSeconds, null, null);
}
/** Default cooldown (CB-578 stage B) when {@code quarantineCooldownSeconds} is absent/non-positive. */
@@ -129,8 +131,8 @@ public record FleetConfig(
String worktreeRoot, Lifecycle lifecycle, Integer spawnReadyTimeoutMs,
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
LeadHeartbeat leadHeartbeat, String placement, Auth auth) {
this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, null, placement, auth, null, null);
this(bind, herdrSocket, null, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, null, placement, auth, null, null, null, null);
}
/** Back-compat form before the optional {@code health:} block was added. */
@@ -138,8 +140,8 @@ public record FleetConfig(
String worktreeRoot, Lifecycle lifecycle, Integer spawnReadyTimeoutMs,
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
LeadHeartbeat leadHeartbeat, String placement, Auth auth, ConfigReload configReload) {
this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, null, placement, auth, configReload, null);
this(bind, herdrSocket, null, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, null, placement, auth, configReload, null, null, null);
}
/** Back-compat form before the CB-578 stage B {@code quarantineCooldownSeconds} field was added. */
@@ -148,9 +150,9 @@ public record FleetConfig(
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
LeadHeartbeat leadHeartbeat, Health health, String placement, Auth auth,
ConfigReload configReload) {
this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
this(bind, herdrSocket, null, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth,
configReload, null);
configReload, null, null, null);
}
/**
@@ -1320,7 +1322,7 @@ public record FleetConfig(
* {@code fleetd.yaml} itself is gitignored.
*/
static final Set<String> KNOWN_TOP_LEVEL_KEYS = Set.of(
"bind", "herdrSocket", "profiles", "guard", "worktreeRoot",
"bind", "herdrSocket", "memberHerdrSocket", "profiles", "guard", "worktreeRoot",
"lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs", "broker", "primary", "fleet",
"leadHeartbeat", "health", "placement", "auth", "configReload", "quarantineCooldownSeconds",
"memberCredentials", "coordinator");
@@ -1941,7 +1943,7 @@ public record FleetConfig(
: new MemberCredentials(null, List.of(), List.of());
// coordinator is left as-is, like broker/primary above: null keeps no LeadMailbox opened,
// and this ticket's Coordinator is config-only anyway (nothing yet reads it at startup).
return new FleetConfig(b, herdrSocket, profiles, g, worktreeRoot, l, timeout, pollMs,
return new FleetConfig(b, herdrSocket, memberHerdrSocket, profiles, g, worktreeRoot, l, timeout, pollMs,
broker, primary, f, leadHeartbeat, health, placementOrDefault, a, configReload,
quarantineCooldown, mc, coordinator);
}
@@ -0,0 +1,40 @@
package dev.ltms.fleet.herdr;
import java.util.Objects;
import java.util.function.Predicate;
/** Routes lead operations and member operations to their owning herdr daemon. */
public final class HerdrRouter implements AutoCloseable {
private final HerdrClient lead;
private final HerdrClient member;
private final AgentControl leadAgents;
private final AgentControl memberAgents;
private final WorkspaceControl leadSpaces;
private final WorkspaceControl memberSpaces;
private final Predicate<String> isLead;
public HerdrRouter(HerdrClient lead, HerdrClient member, Predicate<String> isLead) {
this.lead = Objects.requireNonNull(lead, "lead");
this.member = member != null ? member : lead;
this.isLead = Objects.requireNonNull(isLead, "isLead");
leadAgents = new AgentControl(this.lead);
memberAgents = this.member == this.lead ? leadAgents : new AgentControl(this.member);
leadSpaces = new WorkspaceControl(this.lead);
memberSpaces = this.member == this.lead ? leadSpaces : new WorkspaceControl(this.member);
}
public AgentControl leadAgents() { return leadAgents; }
public WorkspaceControl leadSpaces() { return leadSpaces; }
public AgentControl memberAgents() { return memberAgents; }
public WorkspaceControl memberSpaces() { return memberSpaces; }
public AgentControl agentsFor(String targetId) { return isLead.test(targetId) ? leadAgents : memberAgents; }
HerdrClient leadClient() { return lead; }
HerdrClient memberClient() { return member; }
@Override
public void close() {
lead.close();
if (member != lead) member.close();
}
}
@@ -2,6 +2,7 @@ package dev.ltms.fleet.herdr;
import com.fasterxml.jackson.databind.JsonNode;
import java.util.List;
import java.util.Map;
/**
@@ -13,33 +14,61 @@ import java.util.Map;
* <p>herdr owns the PID→pane truth: {@code pane.process_info} reports each pane's {@code shell_pid}
* and foreground process PIDs. This scans agent panes; a spawn-time {@code pid→terminal} cache is
* the obvious optimization once wired into {@code ClaudeCodeLauncher}.
*
* <p>CB-185 split the fleet across two herdr daemons — lead operations on one, members on the
* other ({@code memberHerdrSocket}). A caller's pane can live on <em>either</em> daemon (a lead's
* MCP connection resolves against the lead daemon; a member's against the member daemon), so this
* must be able to search more than one client. {@link #PaneLocator(HerdrClient, HerdrClient)}
* searches the lead client first, then the member client, and collapses to a single scan when the
* two are the same object (the historical single-daemon deployment).
*/
public final class PaneLocator {
private final HerdrClient herdr;
private final List<HerdrClient> herdrs;
/** Search only this client — the single-daemon deployment. */
public PaneLocator(HerdrClient herdr) {
this.herdr = herdr;
this.herdrs = List.of(herdr);
}
/**
* Search {@code lead} first, then {@code member} — the two-daemon deployment (CB-185). When
* the caller passes the same client for both (no {@code memberHerdrSocket} configured), this
* collapses to one client and one scan, exactly {@link #PaneLocator(HerdrClient)}'s behaviour.
*/
public PaneLocator(HerdrClient lead, HerdrClient member) {
this.herdrs = lead == member ? List.of(lead) : List.of(lead, member);
}
/**
* The {@code terminal_id} of the agent pane whose process tree contains {@code pid}, or
* {@code null} if no agent pane owns it (e.g. the caller is the primary, or off-host).
* {@code null} if no agent pane on any searched daemon owns it (e.g. the caller is the
* primary, or off-host).
*/
public String terminalForPid(long pid) {
if (pid <= 0) {
return null;
}
for (HerdrClient herdr : herdrs) {
String terminal = terminalForPid(herdr, pid);
if (terminal != null) {
return terminal;
}
}
return null;
}
private static String terminalForPid(HerdrClient herdr, long pid) {
for (JsonNode pane : herdr.call("pane.list", Map.of()).path("panes")) {
String paneId = pane.path("pane_id").asText(null);
if (paneId != null && paneOwnsPid(paneId, pid)) {
if (paneId != null && paneOwnsPid(herdr, paneId, pid)) {
return pane.path("terminal_id").asText(null);
}
}
return null;
}
private boolean paneOwnsPid(String paneId, long pid) {
private static boolean paneOwnsPid(HerdrClient herdr, String paneId, long pid) {
JsonNode info;
try {
info = herdr.call("pane.process_info", Map.of("pane_id", paneId)).path("process_info");
@@ -74,6 +74,14 @@ public final class CompletionResolver implements TurnListener {
*/
public static final long MIN_TURN_NANOS = Duration.ofSeconds(2).toNanos();
/**
* fleetd#164 (part 2): one stable, explicit backend-failure marker seen on a Claude Code pane
* when the backend itself rejected the turn (e.g. {@code "API Error: 400 invalid request body"}).
* Kept deliberately narrow — a growing list of ad-hoc error strings rots as backends change their
* wording; broader backend-error surfacing is out of scope here (fleetd#164 point 3).
*/
private static final Pattern BACKEND_ERROR = Pattern.compile("(?i)\\bAPI Error\\s*:");
private static final String CLIPPED_PANE_TAIL_MARKER =
"[Pane tail clipped: member did not call fleet_reply.]";
@@ -278,6 +286,20 @@ public final class CompletionResolver implements TurnListener {
}
return;
}
// fleetd#164 (part 2): a scrape that read cleanly and produced content still isn't a real
// reply when that content is the backend's own rejection (e.g. an HTTP 400 before the worker
// did any work). Classify it as a failure naming the member, rather than handing the caller a
// scrape that reads like a completed answer.
String backendError = firstMatchingLine(assistantBlock, BACKEND_ERROR);
if (backendError != null) {
// Carry the whole scrape, not just the matched line. The pattern is a heuristic: a member
// that forgot fleet_reply while reporting *about* a backend error matches it too. Failing
// is still right — the caller must not read a scrape as an answer — but dropping the rest
// of the pane would destroy the report, which is the same defect fleetd#164 is about.
fail(target, turn, "member " + target + " ended on a backend error: " + backendError
+ "\n--- pane tail ---\n" + tail);
return;
}
String completion = clipped ? tail + "\n" + CLIPPED_PANE_TAIL_MARKER : tail;
if (rendezvous.resolveCompletion(waiter, completion)) {
inFlight.remove(target, turn);
@@ -2,6 +2,7 @@ package dev.ltms.fleet.inject;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.HerdrRouter;
import dev.ltms.fleet.msg.TurnToken;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -89,6 +90,7 @@ public final class Injector {
public static final long POLL_INTERVAL_MILLIS = 250;
private final AgentControl agents;
private final HerdrRouter router;
private final TurnListener turnListener;
private final Predicate<String> ready; // CB-113: a target is deliverable only when available
private final Consumer<String> forget; // CB-114: clear a gone worker's readiness/presence
@@ -124,11 +126,25 @@ public final class Injector {
public Injector(AgentControl agents, TurnListener turnListener, Predicate<String> ready,
Consumer<String> forget) {
this.agents = agents;
this.router = null;
this.turnListener = turnListener;
this.ready = ready;
this.forget = forget;
}
public Injector(HerdrRouter router, TurnListener turnListener, Predicate<String> ready,
Consumer<String> forget) {
this.agents = null;
this.router = router;
this.turnListener = turnListener;
this.ready = ready;
this.forget = forget;
}
private AgentControl agentsFor(String target) {
return router != null ? router.agentsFor(target) : agents;
}
/** A pending message and the future that completes when it has been delivered. */
private record Pending(String text, TurnToken token, CompletableFuture<Void> delivered) {
}
@@ -253,7 +269,7 @@ public final class Injector {
if (p != null && ready.test(target)) {
t.notReadySincePoll = 0;
try {
agents.send(target, p.text());
agentsFor(target).send(target, p.text());
t.queue.poll();
t.awaitingPickup = true;
t.awaitingCompletion = true;
@@ -313,7 +329,7 @@ public final class Injector {
// thread while it holds the target lock.
if (resubmit) {
try {
agents.submit(target); // nudge a raced Enter so the pending paste submits
agentsFor(target).submit(target); // nudge a raced Enter so the pending paste submits
} catch (RuntimeException e) {
log.debug("resubmit to {} failed (will retry next poll): {}", target, e.getMessage());
}
@@ -3,6 +3,7 @@ package dev.ltms.fleet.inject;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.HerdrException;
import dev.ltms.fleet.herdr.HerdrRouter;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -22,6 +23,7 @@ public final class StatusPoller {
private static final Logger log = LoggerFactory.getLogger(StatusPoller.class);
private final AgentControl agents;
private final HerdrRouter router;
private final Injector injector;
private final StatusRefiner refiner;
private final long intervalMillis;
@@ -35,11 +37,24 @@ public final class StatusPoller {
public StatusPoller(AgentControl agents, Injector injector, StatusRefiner refiner,
long intervalMillis) {
this.agents = agents;
this.router = null;
this.injector = injector;
this.refiner = refiner;
this.intervalMillis = intervalMillis;
}
public StatusPoller(HerdrRouter router, Injector injector, long intervalMillis) {
this.agents = null;
this.router = router;
this.injector = injector;
// CB-185: this refiner's own AgentControl (member) is only a default for the legacy 2-arg
// refine() overload — the loop below always calls the 3-arg refine(target, raw, control)
// with the per-target control from router.agentsFor(target), so a lead target is refined
// against the LEAD daemon even though this field points at the member one.
this.refiner = new StatusRefiner(router.memberAgents());
this.intervalMillis = intervalMillis;
}
/** Start the polling loop on a virtual thread. Idempotent. */
public synchronized void start() {
if (running) return;
@@ -56,7 +71,11 @@ public final class StatusPoller {
try {
// herdr's agent_status can misreport a settled worker as `unknown`; refine it
// against the pane content before it drives delivery/completion (CB-115).
AgentStatus status = refiner.refine(target, agents.status(target));
// CB-185: refine THROUGH the same control the raw status came from — a router
// splits lead/member targets across two herdr daemons, and reading a lead's pane
// through the (fixed) member refiner never finds it, wedging that lead at UNKNOWN.
AgentControl control = router != null ? router.agentsFor(target) : agents;
AgentStatus status = refiner.refine(target, control.status(target), control);
injector.onStatus(target, status);
} catch (HerdrException e) {
// The worker's agent is gone — stop trying and unblock its waiters.
@@ -41,15 +41,32 @@ public final class StatusRefiner {
}
/**
* Return a trustworthy status for {@code target}. Any non-{@code UNKNOWN} {@code raw} is returned
* unchanged; an {@code UNKNOWN} triggers a pane read and content classification. A read failure
* leaves it {@code UNKNOWN} (the safe default: no delivery, and the stall path still applies).
* Return a trustworthy status for {@code target}, reading its pane through this refiner's own
* {@link AgentControl}. Equivalent to {@link #refine(String, AgentStatus, AgentControl)} with
* that control — kept for callers that only ever talk to one herdr daemon.
*/
public AgentStatus refine(String target, AgentStatus raw) {
return refine(target, raw, agents);
}
/**
* Return a trustworthy status for {@code target}. Any non-{@code UNKNOWN} {@code raw} is returned
* unchanged; an {@code UNKNOWN} triggers a pane read (through {@code control}) and content
* classification. A read failure leaves it {@code UNKNOWN} (the safe default: no delivery, and
* the stall path still applies).
*
* <p>CB-185: {@code control} must be the {@link AgentControl} for the <em>same</em> daemon the
* raw status was sampled from — a router splits lead and member targets across two herdr
* daemons, and reading a lead's pane through the member client (or vice versa) fails to find
* the pane and leaves the target wedged at {@code UNKNOWN} forever. Callers that route per
* target (e.g. {@code StatusPoller}) must pass that target's control explicitly rather than
* relying on the control fixed at construction.
*/
public AgentStatus refine(String target, AgentStatus raw, AgentControl control) {
if (raw != AgentStatus.UNKNOWN) return raw;
String pane;
try {
pane = agents.read(target, PROBE_SOURCE);
pane = control.read(target, PROBE_SOURCE);
} catch (RuntimeException e) {
log.debug("status refine read for {} failed; leaving UNKNOWN: {}", target, e.getMessage());
return AgentStatus.UNKNOWN;
@@ -3,6 +3,7 @@ package dev.ltms.fleet.member;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.herdr.Agent;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.herdr.HerdrException;
import dev.ltms.fleet.peer.Capability;
import dev.ltms.fleet.peer.MemberRole;
import dev.ltms.fleet.peer.PeerHandle;
@@ -46,9 +47,13 @@ import java.util.stream.Collectors;
* the single adapter that declares it. Profiles partition cleanly across adapters: the
* constructor rejects a name claimed by two.</li>
* <li><strong>By pane id</strong> — {@link #stop} routes to the adapter that spawned that pane
* (recorded at spawn time). A pane the composite never spawned can use the fallback route
* in a one-daemon fleet. With more than one herdr daemon, its owner is unknown, so stop refuses
* the ambiguous id rather than closing a pane on an arbitrary herdr daemon.</li>
* (recorded at spawn time). A pane the composite never spawned, or one whose record was lost
* to a daemon restart (CB-185 blocker 1 — {@link #spawnedBy} is in-memory only), can use the
* fallback route in a one-daemon fleet. With more than one herdr daemon, {@link #probeOwner}
* asks each configured daemon which one actually knows the pane: exactly one match routes
* (and caches); no match is treated as already-gone; more than one match is a genuine
* ambiguity (pane ids are per-daemon counters, so two daemons really can both hold, say,
* {@code w1:p1}) and stop refuses rather than closing a pane on an arbitrary herdr daemon.</li>
* <li><strong>Fleet-wide</strong> — {@link #reapOrphanWorkers} and {@link #capabilities} fan out
* and combine. {@link #list} is deduplicated by (owning daemon, pane id): delegates that share
* one herdr connection report the same global agent set, but two daemons can each hold a pane
@@ -440,12 +445,23 @@ public final class CompositePeerLauncher implements PeerLauncher {
public void stop(String id) {
HerdrPeerLauncher d = spawnedBy.get(id);
if (d == null) {
if (herdrDaemonCount() != 1) {
throw new IllegalArgumentException("ambiguous paneId '" + id
+ "': no owning herdr daemon was recorded");
if (herdrDaemonCount() == 1) {
log.debug("stop({}) — no recorded owner in a single-daemon fleet", id);
d = delegates.getFirst();
} else {
d = probeOwner(id);
if (d == null) {
// No configured herdr daemon has ever heard of this pane. CB-185 blocker 1: this
// is the normal case right after a daemon restart empties spawnedBy for a member
// that has ALREADY been torn down since — the caller retried a stop that already
// succeeded. Nothing to close and no owner to cache; matching the tolerance
// HerdrPeerLauncher#stop already gives an already-gone pane (agent.close swallows
// that as success), stop() here is a no-op rather than a refusal.
log.debug("stop({}) — no configured herdr daemon knows this pane; "
+ "treating as already stopped", id);
return;
}
}
log.debug("stop({}) — no recorded owner in a single-daemon fleet", id);
d = delegates.getFirst();
}
// Drop the owner record only after the delegate accepted the stop. Removing it first meant a
// delegate that threw left the pane alive with its owner forgotten, so the retry fell into
@@ -454,6 +470,61 @@ public final class CompositePeerLauncher implements PeerLauncher {
spawnedBy.remove(id);
}
/**
* CB-185 blocker 1: recover a spawnedBy cache miss by asking every distinct herdr daemon which
* one actually knows {@code id} — the fix for "after a restart, every surviving member becomes
* un-stoppable" (spawnedBy is in-memory only, so a restart empties it, and members intentionally
* outlive the daemon).
*
* <p>Grouped by daemon identity, not by delegate, for the same reason {@link #list()} groups
* that way: two adapters (claude-code, opencode) sharing one herdr connection would otherwise be
* probed twice, and a pane on their shared daemon would look owned by two adapters instead of
* one daemon.
*
* <p>A daemon that fails to answer {@code list()} (e.g. it is down) is treated as "does not know
* this pane" rather than aborting the whole probe — one unreachable daemon must never make a
* pane that a <em>different</em>, healthy daemon actually owns un-stoppable too, which would
* resurrect the exact bug this method exists to fix.
*
* @return the owning delegate — cached into {@link #spawnedBy} so the next call is free — or
* {@code null} when no daemon knows the pane
* @throws IllegalArgumentException when more than one daemon claims the pane: pane ids are
* per-daemon counters, so two daemons really can both hold, say, {@code w1:p1}, and there
* is no way to tell which one the caller means
*/
private HerdrPeerLauncher probeOwner(String id) {
Map<HerdrClient, HerdrPeerLauncher> byDaemon = new IdentityHashMap<>();
for (HerdrPeerLauncher delegate : delegates) {
byDaemon.putIfAbsent(delegate.herdr(), delegate);
}
List<HerdrPeerLauncher> owners = new ArrayList<>();
for (HerdrPeerLauncher representative : byDaemon.values()) {
List<Agent> agents;
try {
agents = representative.list();
} catch (HerdrException e) {
log.warn("stop({}) probe: a configured herdr daemon was unreachable ({}); "
+ "treating it as not knowing this pane", id, e.getClass().getSimpleName());
continue;
}
boolean knows = agents.stream().anyMatch(a -> id.equals(a.paneId()));
if (knows) {
owners.add(representative);
}
}
if (owners.size() > 1) {
throw new IllegalArgumentException("ambiguous paneId '" + id + "': "
+ owners.size() + " configured herdr daemons report this pane — "
+ "no way to tell which one the caller means");
}
if (owners.isEmpty()) {
return null;
}
HerdrPeerLauncher owner = owners.get(0);
spawnedBy.put(id, owner);
return owner;
}
/**
* Count actual herdr daemons, not peer adapter kinds. Identity is intentional: separate client
* objects may represent different daemons even if a client later implements value equality.
@@ -1008,6 +1008,13 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
* applied before the login shell runs and a sourced file can (and did) undo it. The control is
* the ZDOTDIR scrub ({@link #applyEnvironmentAllowListPolicy}); {@code known}/{@code allow}
* remain as reporting only via {@link #logCredentialGap}.
*
* <p>CB-633 follow-up (#192): under {@code allow-list} this method does NOT call {@link
* #logCredentialGap} itself — at this point (called from {@link #baseEnv}, before {@link
* #applyEnvironmentAllowListPolicy} runs) we do not yet know whether the pane's shell is zsh, so
* we cannot yet pick correct wording. That decision, and the call, are deferred entirely to
* {@link #applyEnvironmentAllowListPolicy}, which knows by then whether the scrub will actually
* run.
*/
private void applyMemberCredentialPolicy(Map<String, String> workerEnv) {
FleetConfig.MemberCredentials creds = memberCredentials == null ? null : memberCredentials.get();
@@ -1016,8 +1023,8 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
}
if (!creds.isAllowList()) {
overlayBlockedCredentials(workerEnv, creds);
logCredentialGap(creds, null);
}
logCredentialGap(creds);
}
/** Put {@link #BLOCKED_CREDENTIAL_SENTINEL} over every blocked name in the pane-creation env map. */
@@ -1067,12 +1074,15 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
// logCredentialGap's WARN (below) is the only signal for this path.
warnNonZsh(loginShell);
overlayBlockedCredentials(launch.env(), creds);
logCredentialGap(creds);
logCredentialGap(creds, null);
return null;
}
// Only reached when the scrub is actually about to run — the count below describes that
// scrub, so it must not be logged before this gate (see the non-zsh branch above).
// scrub, so it must not be logged before this gate (see the non-zsh branch above). Same
// reasoning gates logCredentialGap's wording: passing the derived `allowed` set (non-null)
// here, and ONLY here, is what tells it the scrub will really blank an unkept name — #192.
logAllowListCoverage(allowed);
logCredentialGap(creds, allowed);
Path dir = EnvAllowListScrub.generate(Path.of(System.getProperty("java.io.tmpdir")), allowed);
launch.env().put("ZDOTDIR", dir.toAbsolutePath().toString());
log.info("memberCredentials policy=allow-list: profile={} generated ZDOTDIR {} — derived "
@@ -1177,18 +1187,57 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
private static final Pattern CREDENTIAL_SHAPED_NAME =
Pattern.compile("(?i).*(TOKEN|SECRET|_KEY|APIKEY|PASSWORD|CREDENTIAL|AUTH).*");
/** Guards {@link #logCredentialGap} to one WARN per launcher instance, not one per spawn. */
private final AtomicBoolean credentialGapLogged = new AtomicBoolean();
/**
* Guards the {@code effectiveAllowed == null} branch of {@link #logCredentialGap} — the
* genuinely-unprotected report (deny-by-default, and the allow-list non-zsh fallback) — to one
* WARN per launcher instance, not one per spawn.
*
* <p>CB-633 follow-up (#192): kept SEPARATE from {@link #allowListGapLogged} on purpose.
* {@code memberCredentials} is a live, re-read-per-spawn supplier, so the policy can change
* between two spawns on the same launcher. A single shared flag would let a harmless allow-list
* INFO on spawn 1 permanently suppress the real deny-by-default WARN a later spawn deserves —
* the report that matters most getting hidden by the report that doesn't. Two flags mean each
* report kind fires exactly once, independent of what the other kind already logged.
*/
private final AtomicBoolean unprotectedGapLogged = new AtomicBoolean();
/**
* Guards the {@code effectiveAllowed != null} branch of {@link #logCredentialGap} — the
* allow-list-scrub-covered report — to one INFO per launcher instance. See {@link
* #unprotectedGapLogged}'s javadoc for why this is a separate flag rather than a shared one.
*/
private final AtomicBoolean allowListGapLogged = new AtomicBoolean();
/**
* CB-596 criterion 4: a credential-shaped host env var name on neither {@code known} nor
* {@code allow} is not silently allowed — it is reported. {@link #hostEnvNames} enumerates the
* daemon's own environment (see that field's javadoc for why the daemon's env is read rather
* than the spawned pane's, which the daemon has no channel to inspect at spawn time); this logs
* every such NAME, at WARN, at most once per launcher instance — never a value, a prefix of a
* value, or a hash of a value, so the log itself cannot leak anything.
* every such NAME — never a value, a prefix of a value, or a hash of a value, so the log itself
* cannot leak anything.
*
* <p>CB-633 follow-up (#192): {@code effectiveAllowed} picks the wording, and it must NOT be
* picked from {@code creds.isAllowList()} — see {@link #applyEnvironmentAllowListPolicy}'s
* javadoc for the reasoning this mirrors. {@code null} means no scrub-derived allow-list was
* computed for this call — true on the deny-by-default path AND on the allow-list non-zsh
* fallback, where nothing is ever scrubbed — so the whole gap is real and gets the WARN,
* unchanged from before this fix. Non-null means this call came from the zsh branch of {@link
* #applyEnvironmentAllowListPolicy}, reachable ONLY after that method's own zsh gate — but
* {@code effectiveAllowed} is a SUPERSET of {@code known ∪ allow}: {@link MemberEnvAllowList#derive}
* also unions in every profile's {@code gitTokenEnv}/{@code gitHostEnv}/{@code tokenEnv}/
* {@code env:} keys, and {@link #derivedAllowedNames} further unions in this very spawn's own
* env keys — so a name can be in the gap (uncovered by {@code known}/{@code allow}) AND still be
* kept by the derived allow-list, in which case the scrub does NOT blank it and the member DOES
* inherit it. Lead review on #192 caught this: the first cut of this fix reported the WHOLE gap
* as scrub-blanked without checking that, which reported a real leak as safe — the exact
* inversion #192 exists to remove. So on this path the gap is split with {@link
* MemberEnvAllowList#keeps}, the SAME predicate the generated scrub itself evaluates, so this
* split cannot drift from what the scrub actually does: the names it says are kept get the WARN
* (same severity, and same guard, as the deny-by-default case — a name genuinely reaching a
* member unprotected is equally serious whichever path put it there), and the names it says are
* blanked keep the INFO.
*/
private void logCredentialGap(FleetConfig.MemberCredentials creds) {
private void logCredentialGap(FleetConfig.MemberCredentials creds, Set<String> effectiveAllowed) {
Set<String> covered = new HashSet<>(creds.known());
covered.addAll(creds.allow());
List<String> gap = hostEnvNames.get().stream()
@@ -1199,7 +1248,37 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
if (gap.isEmpty()) {
return;
}
if (credentialGapLogged.compareAndSet(false, true)) {
if (effectiveAllowed == null) {
warnGapUnprotected(gap);
return;
}
List<String> keptByDerivedList = gap.stream()
.filter(name -> MemberEnvAllowList.keeps(effectiveAllowed, name))
.toList();
List<String> blankedByScrub = gap.stream()
.filter(name -> !MemberEnvAllowList.keeps(effectiveAllowed, name))
.toList();
if (!keptByDerivedList.isEmpty() && unprotectedGapLogged.compareAndSet(false, true)) {
log.warn("memberCredentials gap: {} credential-shaped env var name(s) are on neither "
+ "known: nor allow: — the derived allow-list keeps them anyway (a profile's "
+ "gitTokenEnv/gitHostEnv/tokenEnv/env: names one, or this spawn injects it), "
+ "so every member pane inherits them UNBLOCKED — {}. Add each to "
+ "memberCredentials.known (or .allow if a member legitimately needs it), or "
+ "remove it from whatever profile setting derives it in.",
keptByDerivedList.size(), keptByDerivedList);
}
if (!blankedByScrub.isEmpty() && allowListGapLogged.compareAndSet(false, true)) {
log.info("memberCredentials gap: {} credential-shaped env var name(s) are on neither "
+ "known: nor allow: — {}. The allow-list scrub blanks them anyway (they "
+ "are not on the derived allow-list), so no member pane keeps them; add "
+ "each to memberCredentials.known or .allow to make that explicit.",
blankedByScrub.size(), blankedByScrub);
}
}
/** The deny-by-default (and allow-list non-zsh fallback) WARN — unchanged byte-for-byte by #192. */
private void warnGapUnprotected(List<String> gap) {
if (unprotectedGapLogged.compareAndSet(false, true)) {
log.warn("memberCredentials gap: {} credential-shaped env var name(s) are on neither "
+ "known: nor allow: — every member pane inherits them UNBLOCKED — {}. "
+ "Add each to memberCredentials.known (blocked by default) or .allow "
@@ -2,6 +2,7 @@ package dev.ltms.fleet.msg;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.HerdrRouter;
import dev.ltms.fleet.inject.Injector;
import dev.ltms.fleet.metrics.FleetMetrics;
import dev.ltms.fleet.metrics.Metrics;
@@ -167,14 +168,23 @@ public final class MessageService {
private final String ticket;
private final String target;
private final CompletableFuture<Reply> future = new CompletableFuture<>();
private final long createdNanos;
/**
* When {@link #future} resolved, or {@code null} while it is still pending — the clock
* {@link #pruneTerminalTickets} measures the TTL from (#197). Deliberately a boxed
* {@code Long} rather than a {@code long} with a sentinel: {@link System#nanoTime} may
* legitimately return any value, zero and negatives included, so no numeric sentinel can mean
* "not stamped yet". Stamped by a {@code whenComplete} hook registered in the constructor, so
* every completion path stamps it — a reply, the completion fallback, a timeout, a failure,
* or an abandon on teardown — without each of those having to remember to.
*/
private volatile Long completedNanos;
private volatile Reply question;
private volatile String turnId;
private Task(String ticket, String target, long createdNanos) {
private Task(String ticket, String target, LongSupplier nowNanos) {
this.ticket = ticket;
this.target = target;
this.createdNanos = createdNanos;
future.whenComplete((reply, ex) -> completedNanos = nowNanos.getAsLong());
}
}
@@ -189,6 +199,7 @@ public final class MessageService {
}
private final AgentControl agents;
private final HerdrRouter router;
private final Injector injector;
private final Rendezvous rendezvous;
private final ReplyInbox inbox;
@@ -257,6 +268,7 @@ public final class MessageService {
MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous, ReplyInbox inbox,
ReplyPushLoop pushLoop, Metrics metrics, LongSupplier nowNanos) {
this.agents = agents;
this.router = null;
this.injector = injector;
this.rendezvous = rendezvous;
this.inbox = inbox;
@@ -265,6 +277,20 @@ public final class MessageService {
this.nowNanos = nowNanos;
}
public MessageService(HerdrRouter router, Injector injector, Rendezvous rendezvous, ReplyInbox inbox,
ReplyPushLoop pushLoop, Metrics metrics) {
this.agents = null;
this.router = router;
this.injector = injector;
this.rendezvous = rendezvous;
this.inbox = inbox;
this.pushLoop = pushLoop;
this.metrics = metrics;
this.nowNanos = System::nanoTime;
}
private AgentControl agentsFor(String target) { return router != null ? router.agentsFor(target) : agents; }
/** Create with an explicit {@link ReplyInbox} and no push loop. */
public MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous, ReplyInbox inbox) {
this(agents, injector, rendezvous, inbox, null);
@@ -277,7 +303,7 @@ public final class MessageService {
/** Current lifecycle status of a worker (the {@code GET /sessions/{id}/status} surface). */
public AgentStatus status(String target) {
return agents.status(target);
return agentsFor(target).status(target);
}
/** Read-only delegation fact for fleet views. */
@@ -340,21 +366,40 @@ public final class MessageService {
}
/**
* Route a worker's explicit {@code fleet_reply}: resolve an open send, or queue it in the
* inbox if no send is currently open. Unlike the bare {@link Rendezvous#resolve}, a no-waiter
* result is <em>not</em> a failure — the reply is held for later drain.
* Route a worker's explicit {@code fleet_reply}: resolve an open send, complete an async ticket
* still parked waiting on this exact turn's answer, or — only once neither applies — queue it in
* the inbox. Unlike the bare {@link Rendezvous#resolve}, a no-waiter result is <em>not</em> a
* failure — the reply is held for later drain.
*
* <p><strong>Do NOT use this for mid-turn questions.</strong> {@code fleet_ask} /
* {@link Rendezvous#resolveQuestion} must keep today's {@code NO_WAITER} behaviour — questions
* are interactive and must never be queued.
*
* @return always {@code true} — the reply either resolved a live send or was queued
* @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 (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.
Task orphan = askAnsweredAsyncTask(session);
if (orphan != null && orphan.future.complete(new Reply(Outcome.REPLIED, content))) {
if (orphan.turnId != null) {
asyncTasksByTurn.remove(orphan.turnId, orphan);
}
count(FleetMetrics.REPLIES, "path", "async-recovered");
return true; // the ticket itself took it — no inbox stranding at all
}
inbox.publish(session, UUID.randomUUID().toString(), content);
// CB-640: record the stranding itself (not just the reply text) so fleet health can see a
// worker whose replies keep missing their waiter, not only the queue depth this leaves behind.
@@ -368,6 +413,24 @@ public final class MessageService {
return true; // held, not lost
}
/**
* The 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). {@code null} 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}).
*/
private Task askAnsweredAsyncTask(String target) {
for (Task task : tasks.values()) {
if (target.equals(task.target) && task.question == null && task.turnId != null
&& !task.future.isDone()) {
return task;
}
}
return null;
}
/** Record a counter sample when a registry is wired; a no-op in unit tests. */
private void count(String name, String... labels) {
if (metrics != null) {
@@ -409,19 +472,43 @@ public final class MessageService {
* <p>Resolving the waiter as a failure — rather than letting it time out — also means the
* outcome is counted, so a torn-down delegation stops being invisible to {@code /metrics}.
*
* @return true if a live waiter was failed
* <p><strong>#137 defence in depth.</strong> {@link #reply} already hands a worker's real
* {@code fleet_reply} straight to the async ticket it belongs to whenever one is still parked
* waiting for it (see {@link #askAnsweredAsyncTask}), so by the time a session is released its
* tasks are normally already resolved — this loop's {@code complete} calls are then harmless
* no-ops (a {@link CompletableFuture} can only resolve once). But should some other path someday
* strand a reply in the inbox without completing its ticket, checking
* {@link #hasStrandedReply(String)} here — before ever writing a failure — means a torn-down
* session whose worker in fact replied is still reported {@code REPLIED} with that reply's own
* text, never the misleading "the worker session was released before it replied" (which also
* means the snapshot/worktree recovery hint that follows it never prints once a reply exists).
*
* @return true if a live waiter or an async task was failed (never true for one recovered as a
* reply — see the note above)
*/
public boolean abandon(String target, String reason) {
boolean hadStrandedReply = hasStrandedReply(target);
// CB-640: the session is gone — nothing will ever accept or deliver into it now.
strandedReplies.remove(target);
queuedDeliveries.remove(target);
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(target);
boolean failed = waiter != null && !waiter.isDone() && rendezvous.resolveFailure(waiter, reason);
boolean asyncFailed = false;
Reply recovered = null; // lazily drained at most once, only if a task actually needs it
for (Task task : tasks.values()) {
if (target.equals(task.target) && task.question == null
&& task.future.complete(new Reply(Outcome.WORKER_FAILED, reason))) {
asyncFailed = true;
if (!target.equals(task.target) || task.question != null || task.future.isDone()) {
continue;
}
if (hadStrandedReply && recovered == null) {
recovered = recoverStrandedReply(target);
}
Reply outcome = recovered != null ? recovered : new Reply(Outcome.WORKER_FAILED, reason);
if (task.future.complete(outcome)) {
if (outcome.outcome() == Outcome.WORKER_FAILED) {
asyncFailed = true;
} else if (task.turnId != null) {
asyncTasksByTurn.remove(task.turnId, task);
}
}
}
if (failed) {
@@ -430,6 +517,21 @@ public final class MessageService {
return failed || asyncFailed;
}
/**
* Drain {@code target}'s inbox and hand its content back as a {@link Outcome#REPLIED} result
* (#137 defence in depth for {@link #abandon}) — {@code null} if it turned out empty (the
* stranding fact raced away, e.g. a lead's own {@code fleet_poll} on the raw session already
* drained it first). When more than one message is queued, only the newest is the worker's actual
* final answer ({@link #drainReplies} returns them oldest-first).
*/
private Reply recoverStrandedReply(String target) {
var messages = drainReplies(target);
if (messages.isEmpty()) {
return null;
}
return new Reply(Outcome.REPLIED, messages.get(messages.size() - 1).content());
}
/**
* Acknowledge a specific reply by {@code msgId} for {@code target}. Removes it from the inbox
* so that a subsequent drain or peek no longer returns it.
@@ -704,7 +806,7 @@ public final class MessageService {
*/
public String sendAsync(String target, String content, Runnable onAccepted) {
String ticket = "task-" + ticketSeq.incrementAndGet();
Task task = new Task(ticket, target, nowNanos.getAsLong());
Task task = new Task(ticket, target, nowNanos);
tasks.put(ticket, task);
if (pushLoop != null) {
// CB-588: task.future only ever completes on a terminal phase (DONE or a failure) — a
@@ -788,7 +890,7 @@ public final class MessageService {
/** Best-effort live worker status for a pending poll; never throws (a lookup error is just noise). */
private String liveStatus(String target) {
try {
return agents.status(target).name().toLowerCase();
return agentsFor(target).status(target).name().toLowerCase();
} catch (RuntimeException e) {
return "unknown";
}
@@ -804,12 +906,25 @@ public final class MessageService {
* (or one the reminder cap already gave up on) is pruned here but never collected there, so it
* lingers in {@code pendingTickets} forever and rides along on every later nudge to the same lead
* — naming a ticket {@code fleet_poll} can no longer find (CB-588 follow-up).
*
* <p>The TTL runs from **completion**, not from creation (#197). It used to compare against
* {@code createdNanos}, which made the real collection window {@code TTL minus however long the
* task ran}: a delegation that took longer than the TTL had its reply destroyed on the first
* sweep after it landed, every time. That is the normal case here — real work runs well past ten
* minutes — and the reply lives only in {@code future}, so pruning it discards the worker's whole
* report with nothing to fall back on. Measuring from completion gives every ticket the same full
* window whatever its runtime, and still bounds {@code tasks}.
*
* <p>A task whose future is done but whose {@code completedNanos} is not stamped yet is left
* alone. That window is the few instructions between {@code complete()} and the constructor's
* {@code whenComplete} hook running; the next sweep collects it.
*/
private void pruneTerminalTickets() {
long cutoff = nowNanos.getAsLong() - TICKET_TTL_NANOS;
tasks.entrySet().removeIf(e -> {
Task t = e.getValue();
boolean expired = t.future.isDone() && t.createdNanos < cutoff;
Long completed = t.completedNanos;
boolean expired = t.future.isDone() && completed != null && completed < cutoff;
if (expired && pushLoop != null) {
pushLoop.ticketCollected(e.getKey());
}
@@ -55,7 +55,8 @@ public final class FleetApp {
/** Context attribute under which the resolved caller is stashed by the auth filter. */
private static final String CALLER = "fleetd.caller";
private final HerdrClient herdr;
private final HerdrClient herdr; // lead daemon
private final HerdrClient memberHerdr; // CB-185: member daemon (same object when unconfigured)
private final PeerLauncher workers;
private final SessionManager sessions; // CB-301: authoritative session registry
private final MessageService messages;
@@ -95,7 +96,23 @@ public final class FleetApp {
MessageService messages, MemberPresence presence,
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics,
Predicate<String> deliverable) {
this(herdr, herdr, workers, sessions, messages, presence, mcpServlet, auth, metrics, deliverable);
}
/**
* @param herdr the lead daemon's client
* @param memberHerdr the member daemon's client (CB-185); pass the same instance as
* {@code herdr} for a single-daemon deployment — {@code healthz}/{@code
* sessions} then make exactly one herdr call each, unchanged from before
* the two-daemon router existed
* @param deliverable the injector's readiness gate, shared so status reports its real result
*/
public FleetApp(HerdrClient herdr, HerdrClient memberHerdr, PeerLauncher workers, SessionManager sessions,
MessageService messages, MemberPresence presence,
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics,
Predicate<String> deliverable) {
this.herdr = herdr;
this.memberHerdr = memberHerdr != null ? memberHerdr : herdr;
this.workers = workers;
this.sessions = sessions;
this.messages = messages;
@@ -190,30 +207,85 @@ public final class FleetApp {
ctx.status(200).contentType("text/plain; version=0.0.4; charset=utf-8").result(metrics.render());
}
/** Liveness + herdr reachability. 200 when herdr answers ping, 503 otherwise. */
/**
* Liveness + herdr reachability. 200 only when BOTH daemons answer ping — 503 otherwise
* (CB-185). With no {@code memberHerdrSocket} configured {@code memberHerdr == herdr}, so this
* makes exactly the one {@code ping} call it always did and reports the same body; with a
* second daemon configured, a member daemon that is down must not be masked by a healthy lead
* daemon — every spawn goes through the member daemon and would otherwise fail silently behind
* a green {@code /healthz}.
*
* <p>CB-185 blocker 2: the {@code herdr} key always carries the <em>lead</em> daemon's
* version/protocol, unchanged, because two consumers — {@code scripts/redeploy-fleetd.sh} and
* {@code scripts/rename-checkout.sh} — read this endpoint already (both only check the HTTP
* status code and print the body verbatim; neither parses a specific field, so adding a key
* alongside {@code herdr} is safe). But it is the <em>member</em> daemon's protocol that decides
* whether a spawn works, so when a second daemon is configured its version/protocol is reported
* too, under a separate {@code member} key — never folded into {@code herdr}, which would make a
* mismatch invisible to whichever consumer only reads that key. If the two protocol numbers
* differ, {@code protocolMismatch: true} calls it out explicitly rather than leaving it to be
* spotted by comparing two numbers by eye.
*/
private void healthz(Context ctx) {
JsonNode pong;
try {
JsonNode pong = herdr.call("ping");
ctx.status(200).json(Map.of(
"status", "ok",
"herdr", Map.of(
"version", pong.path("version").asText(""),
"protocol", pong.path("protocol").asInt())));
pong = herdr.call("ping");
} catch (HerdrException e) {
ctx.status(503).json(Map.of(
"status", "degraded",
"herdr", "unreachable",
"detail", e.getMessage()));
return;
}
Map<String, Object> body = new LinkedHashMap<>();
body.put("status", "ok");
body.put("herdr", Map.of(
"version", pong.path("version").asText(""),
"protocol", pong.path("protocol").asInt()));
if (memberHerdr != herdr) {
JsonNode memberPong;
try {
memberPong = memberHerdr.call("ping");
} catch (HerdrException e) {
ctx.status(503).json(Map.of(
"status", "degraded",
"herdr", "member unreachable",
"detail", e.getMessage()));
return;
}
int leadProtocol = pong.path("protocol").asInt();
int memberProtocol = memberPong.path("protocol").asInt();
body.put("member", Map.of(
"version", memberPong.path("version").asText(""),
"protocol", memberProtocol));
if (leadProtocol != memberProtocol) {
body.put("protocolMismatch", true);
}
}
ctx.status(200).json(body);
}
/** Sessions view derived from herdr {@code workspace.list} (one workspace → one row). */
/**
* Sessions view derived from herdr {@code workspace.list} (one workspace → one row), merged
* across both daemons (CB-185). With no {@code memberHerdrSocket} configured {@code
* memberHerdr == herdr}, so this calls {@code workspace.list} exactly once, same as before the
* router existed; with a second daemon configured, calling it twice would silently drop every
* member workspace (they live on the member daemon only).
*/
private void sessions(Context ctx) {
if (!allow(ctx, Authz.Action.READ, null)) {
return;
}
JsonNode result = herdr.call("workspace.list");
List<Map<String, Object>> out = new ArrayList<>();
collectSessions(herdr, out);
if (memberHerdr != herdr) {
collectSessions(memberHerdr, out);
}
ctx.status(200).json(Map.of("sessions", out));
}
private static void collectSessions(HerdrClient client, List<Map<String, Object>> out) {
JsonNode result = client.call("workspace.list");
for (JsonNode w : result.path("workspaces")) {
out.add(Map.of(
"id", w.path("workspace_id").asText(""),
@@ -222,7 +294,6 @@ public final class FleetApp {
"paneCount", w.path("pane_count").asInt(),
"agentStatus", w.path("agent_status").asText("unknown")));
}
ctx.status(200).json(Map.of("sessions", out));
}
/** Discovery: every agent herdr tracks, keyed by its Claude session UUID. */
@@ -110,6 +110,7 @@ public final class GitWorktrees implements Worktrees {
@Override
public String add(String repoRoot, String branch, String baseRef) {
reportRemoteUrlsWithUserInfo(repoRoot);
String base = (baseRef == null || baseRef.isBlank()) ? "HEAD" : baseRef;
String nonce = nonce();
Path root = resolveRoot(repoRoot);
@@ -139,7 +140,7 @@ public final class GitWorktrees implements Worktrees {
if (exitCode("git", "-C", repoRoot, "config", "--get", "remote.origin.url") != 0) {
return;
}
String origin = exec("git", "-C", repoRoot, "config", "--get", "remote.origin.url").trim();
String origin = execRedacted("git", "-C", repoRoot, "config", "--get", "remote.origin.url").trim();
URI uri;
try {
uri = new URI(origin);
@@ -155,6 +156,9 @@ public final class GitWorktrees implements Worktrees {
throw new WorktreeException("origin URL has invalid HTTPS user info; cannot provision safely");
}
String cleanOrigin = origin.substring(0, schemeEnd) + origin.substring(userInfoEnd + 1);
// Plain exec, not execRedacted, is correct here: this call WRITES cleanOrigin (already
// stripped of user-info above) rather than reading a URL back from stdout, so there is
// nothing secret left in either its argv or its stdout to redact.
exec("git", "-C", repoRoot, "remote", "set-url", "origin", cleanOrigin);
log.info("removed HTTPS user info from forge origin before provisioning worktree");
}
@@ -164,7 +168,7 @@ public final class GitWorktrees implements Worktrees {
if (exitCode("git", "-C", worktreePath, "remote", "get-url", "--all", "origin") != 0) {
return;
}
String origins = exec("git", "-C", worktreePath, "remote", "get-url", "--all", "origin");
String origins = execRedacted("git", "-C", worktreePath, "remote", "get-url", "--all", "origin");
for (String origin : origins.split("\\R")) {
try {
URI uri = new URI(origin);
@@ -177,6 +181,99 @@ public final class GitWorktrees implements Worktrees {
}
}
/**
* Report — never refuse — every remote whose fetch or push URL carries user-info outside SSH.
* A linked worktree shares its parent repository's git config, so a credential on ANY remote
* (not only {@code origin}) or in a {@code pushurl} is just as readable by a member as one on
* {@code origin}'s HTTPS fetch URL — the one case {@link #removeUserInfoFromHttpsOrigin} and
* {@link #requireCredentialFreeHttpsOrigin} already strip and refuse. This check is additive: it
* only logs a warning, it never mutates config and never refuses the provision.
*
* <p>A reporting-only check must never be able to abort a provision — PR #173 shipped one that
* ran unguarded at the top of {@link #add}, and every git call inside it can throw ({@link #exec}
* turns a non-zero exit or its 30-second timeout into a {@link WorktreeException}). Every git call
* here is therefore wrapped, and on failure only the exception's <em>class</em> is logged, never
* its message: the enumerating {@code git remote} call is not redacted, and its stderr is read
* from the very config that may hold the URL this check exists to find.
*/
private void reportRemoteUrlsWithUserInfo(String repoRoot) {
List<String> remotes;
try {
remotes = exec("git", "-C", repoRoot, "remote").lines()
.map(String::trim)
.filter(r -> !r.isBlank())
.toList();
} catch (RuntimeException e) {
log.warn("could not enumerate remotes to check for credentialed URLs in {}: {}",
Path.of(repoRoot).toAbsolutePath().normalize(), e.getClass().getName());
return;
}
for (String remote : remotes) {
reportOneRemoteUrlsWithUserInfo(repoRoot, remote);
}
}
private void reportOneRemoteUrlsWithUserInfo(String repoRoot, String remote) {
boolean leaks = remoteUrlsLeakUserInfo(repoRoot, remote, false)
|| remoteUrlsLeakUserInfo(repoRoot, remote, true);
if (leaks) {
log.warn("member worktree shares a remote URL containing user-info: remote={} repository={}; "
+ "remove credentials from the repository's git config",
remote, Path.of(repoRoot).toAbsolutePath().normalize());
}
}
/**
* True when any resolved fetch (or, if {@code push}, push) URL for {@code remote} carries
* non-empty user-info outside SSH. Never throws — a git failure here is caught, logged (its
* class only, per the javadoc above), and treated as "nothing found", so it cannot abort or
* otherwise affect provisioning. Uses {@link #execRedacted} because the command's stdout is
* itself the URL this check exists to find.
*/
private boolean remoteUrlsLeakUserInfo(String repoRoot, String remote, boolean push) {
try {
String out = push
? execRedacted("git", "-C", repoRoot, "remote", "get-url", "--push", "--all", remote)
: execRedacted("git", "-C", repoRoot, "remote", "get-url", "--all", remote);
return out.lines().anyMatch(url -> !url.isBlank() && urlLeaksUserInfo(url.trim()));
} catch (RuntimeException e) {
log.warn("could not read the {} URL for remote {} to check for credentials: {}",
push ? "push" : "fetch", remote, e.getClass().getName());
return false;
}
}
/**
* True when {@code rawUrl} parses as an absolute URI with a non-SSH-family scheme and non-empty
* user-info. An unparsable or scheme-less URL — including the ssh scp-like shorthand
* ({@code user@host:path}) — is not this check's concern and is treated as "no finding", the
* same way {@link #configureHttpsUrlRewriteForSshOrigin} leaves that shorthand untouched.
*/
private static boolean urlLeaksUserInfo(String rawUrl) {
URI uri;
try {
uri = new URI(rawUrl);
} catch (URISyntaxException e) {
return false;
}
String scheme = uri.getScheme();
if (scheme == null || isSshLikeScheme(scheme)) {
return false;
}
String userInfo = uri.getUserInfo();
return userInfo != null && !userInfo.isEmpty();
}
/**
* SSH-family schemes deliberately excluded from {@link #urlLeaksUserInfo}: there, the user part
* selects an account and authentication itself happens over SSH, so it is not a credential the
* way HTTPS/HTTP user-info is.
*/
private static boolean isSshLikeScheme(String scheme) {
return "ssh".equalsIgnoreCase(scheme) || "git+ssh".equalsIgnoreCase(scheme)
|| "ssh+git".equalsIgnoreCase(scheme);
}
/** Configure a per-worktree helper that supplies a token from the member environment at call time. */
private void configureEnvironmentCredentialHelper(String repoRoot, String worktreePath) {
exec("git", "-C", repoRoot, "config", "extensions.worktreeConfig", "true");
@@ -234,7 +331,7 @@ public final class GitWorktrees implements Worktrees {
if (exitCode("git", "-C", repoRoot, "config", "--get", "remote.origin.url") != 0) {
return;
}
String origin = exec("git", "-C", repoRoot, "config", "--get", "remote.origin.url").trim();
String origin = execRedacted("git", "-C", repoRoot, "config", "--get", "remote.origin.url").trim();
URI uri;
try {
uri = new URI(origin);
@@ -603,6 +700,31 @@ public final class GitWorktrees implements Worktrees {
/** Same as {@link #exec(String...)}, with extra environment variables set on the child process. */
private String exec(Map<String, String> extraEnv, String... command) {
return exec(extraEnv, false, command);
}
/**
* Same as {@link #exec(String...)}, for a command whose stdout may itself carry a credential
* (e.g. {@code git remote get-url}, whose output is a URL). Stdout is still returned normally on
* success — callers still get the URL to inspect — but it is suppressed from BOTH the timeout
* message and the non-zero-exit message, so a failing call here can never copy it into a
* {@link WorktreeException}, and from there into a caller's log.
*/
private String execRedacted(String... command) {
return exec(Map.of(), true, command);
}
/**
* Shared implementation for {@link #exec(Map, String...)} and {@link #execRedacted(String...)}.
* {@code redactOutput} suppresses captured stdout from both failure messages below.
*
* <p>Package-private, not {@code private}: also a test seam, the same way the
* {@code afterWorktreeAdded} constructor parameter is. It lets a test drive the redaction
* guarantee directly — a synthetic failing command whose stdout carries a test marker passed
* through {@code extraEnv} rather than argv — without depending on finding a real git failure
* mode that happens to echo a URL onto stdout before exiting non-zero.
*/
String exec(Map<String, String> extraEnv, boolean redactOutput, String... command) {
String out;
int code;
Process p;
@@ -624,7 +746,8 @@ public final class GitWorktrees implements Worktrees {
try {
if (!p.waitFor(30, TimeUnit.SECONDS)) {
p.destroyForcibly();
throw new WorktreeException("command timed out: " + String.join(" ", command) + "\n" + out);
throw new WorktreeException("command timed out: " + String.join(" ", command)
+ (redactOutput || out.isBlank() ? "" : "\n" + out));
}
code = p.exitValue();
} catch (InterruptedException e) {
@@ -634,7 +757,7 @@ public final class GitWorktrees implements Worktrees {
}
if (code != 0) {
throw new WorktreeException("exit " + code + " for: " + String.join(" ", command)
+ (out.isBlank() ? "" : "\n" + out));
+ (redactOutput || out.isBlank() ? "" : "\n" + out));
}
return out;
}
@@ -0,0 +1,31 @@
package dev.ltms.fleet;
import java.nio.file.Files;
import java.nio.file.Path;
import org.junit.jupiter.api.Test;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* CB-185: {@code ConnectionIdentity} must resolve a caller's pane on EITHER herdr daemon (a
* lead's MCP connection resolves against the lead daemon; a member's against the member daemon).
* Pinning {@code PaneLocator} to {@code memberHerdr} alone — the bug this guards against — leaves
* every lead's own connection unresolvable ({@code callerTerminal == null}) the moment
* {@code memberHerdrSocket} names a second daemon, which breaks {@code fleet_reply}/{@code
* fleet_ask} and {@code fleet_whoami} for a lead. A unit test on {@link
* dev.ltms.fleet.herdr.PaneLocator} alone (see {@code PaneLocatorTest}) proves the class CAN
* search two clients, but not that {@code Fleetd.main} actually wires it that way — hence this
* source-level assertion, the same technique {@code FleetdHerdrControlConstructionTest} uses.
*/
class FleetdConnectionIdentityConstructionTest {
@Test
void connectionIdentitySearchesBothDaemonsNotJustTheMemberOne() throws Exception {
String source = Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java"));
assertFalse(source.contains("new PaneLocator(memberHerdr)"),
"PaneLocator must not be pinned to the member daemon alone — a lead's own "
+ "connection resolves against the LEAD daemon and would never be found");
assertTrue(source.contains("new PaneLocator(herdr, memberHerdr)"),
"PaneLocator must search the lead daemon first, then the member daemon");
}
}
@@ -0,0 +1,29 @@
package dev.ltms.fleet;
import java.nio.file.Files;
import java.nio.file.Path;
import org.junit.jupiter.api.Test;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* CB-185: {@code FleetApp} must be constructed with BOTH herdr clients (the lead's and the
* member's), never the raw lead-only {@code herdr}. Passing only {@code herdr} — the bug this
* guards against — makes {@code GET /healthz} green while the member daemon is down (so every
* spawn fails invisibly) and silently drops every member workspace from {@code GET /sessions}.
* A behavioural test on {@code FleetApp} alone (see {@code FleetAppTwoDaemonTest}) proves the
* class merges/gates correctly when given two clients, but not that {@code Fleetd.main} actually
* passes it two — hence this source-level assertion, mirroring
* {@code FleetdHerdrControlConstructionTest}.
*/
class FleetdFleetAppConstructionTest {
@Test
void fleetAppIsConstructedWithBothHerdrDaemons() throws Exception {
String source = Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java"));
assertFalse(source.contains("new FleetApp(herdr, workers,"),
"FleetApp must not be constructed with the lead-only herdr client");
assertTrue(source.contains("new FleetApp(herdr, memberHerdr, workers,"),
"FleetApp must be constructed with both the lead and the member herdr client");
}
}
@@ -0,0 +1,17 @@
package dev.ltms.fleet;
import java.nio.file.Files;
import java.nio.file.Path;
import org.junit.jupiter.api.Test;
import static org.junit.jupiter.api.Assertions.assertFalse;
class FleetdHerdrControlConstructionTest {
@Test
void fleetdDelegatesStatefulControlsToTheRouter() throws Exception {
// AgentControl caches paneByTerminal, so the router must be its only production factory.
String source = Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java"));
assertFalse(source.contains("new AgentControl("));
assertFalse(source.contains("new WorkspaceControl("));
}
}
@@ -31,6 +31,8 @@ public final class FakeHerdr implements HerdrClient {
*/
public final List<Call> calls = new CopyOnWriteArrayList<>();
private boolean healthy = true;
private String pingVersion = "0.8.0";
private int pingProtocol = 19;
private final List<String> extraWorkspaces = new ArrayList<>();
private final List<String> extraAgents = new ArrayList<>();
/** workspaceId → extra tabs that {@code tab.list} reports for it (CB-558 lead scans). */
@@ -40,6 +42,7 @@ public final class FakeHerdr implements HerdrClient {
private int workerTabPaneCount = 1;
private String paneCloseErrorCode = null;
private String agentSendErrorCode = null;
private boolean noPanes = false;
private volatile String agentStatus = "idle"; // steady-state agent.get status
private volatile String readText = "worker transcript tail"; // canned agent.read output
private int pinnedStarts = 0; // how many upcoming agent.start calls report a fixed pane
@@ -51,6 +54,16 @@ public final class FakeHerdr implements HerdrClient {
return this;
}
/**
* Make {@code ping} report this version/protocol instead of the default 0.8.0/19 — CB-185
* blocker 2's fixture for a lead and a member daemon running mismatched herdr versions.
*/
public FakeHerdr pingReports(String version, int protocol) {
this.pingVersion = version;
this.pingProtocol = protocol;
return this;
}
/** Reject the first {@code n} {@code agent.start} calls with {@code agent_name_taken}. */
public FakeHerdr agentNameTakenTimes(int n) {
this.agentNameTakenFor = n;
@@ -75,6 +88,15 @@ public final class FakeHerdr implements HerdrClient {
return this;
}
/**
* Make {@code pane.list} report no panes at all — models a second herdr daemon (CB-185) that
* simply does not host the pane a {@link PaneLocator} is searching for.
*/
public FakeHerdr withNoPanes() {
this.noPanes = true;
return this;
}
/** Set the {@code agent_status} that {@code agent.get} reports (drives the injector). */
public FakeHerdr agentStatus(String status) {
this.agentStatus = status;
@@ -156,7 +178,8 @@ public final class FakeHerdr implements HerdrClient {
try {
return switch (method) {
case "ping" -> mapper.readTree(
"{\"type\":\"pong\",\"version\":\"0.8.0\",\"protocol\":19}");
("{\"type\":\"pong\",\"version\":\"%s\",\"protocol\":%d}")
.formatted(pingVersion, pingProtocol));
case "workspace.list" -> mapper.readTree(("""
{"type":"workspace_list","workspaces":[
{"workspace_id":"w1","label":"dev-mgnl","focused":true,"pane_count":7,"agent_status":"unknown"},
@@ -268,7 +291,9 @@ public final class FakeHerdr implements HerdrClient {
case "pane.get" -> mapper.readTree("""
{"type":"pane_info","pane":{"pane_id":"w9:pW","workspace_id":"w9",
"tab_id":"w9:t2","agent_status":"idle"}}""");
case "pane.list" -> mapper.readTree("""
case "pane.list" -> noPanes
? mapper.readTree("{\"type\":\"pane_list\",\"panes\":[]}")
: mapper.readTree("""
{"type":"pane_list","panes":[
{"pane_id":"w2:p7","terminal_id":"term_a","workspace_id":"w2","tab_id":"w2:t7","agent":"claude"},
{"pane_id":"w2:p9","terminal_id":"term_shell","workspace_id":"w2","tab_id":"w2:t8"}]}""");
@@ -0,0 +1,31 @@
package dev.ltms.fleet.herdr;
import org.junit.jupiter.api.Test;
import static org.junit.jupiter.api.Assertions.assertSame;
import static org.junit.jupiter.api.Assertions.assertNotSame;
class HerdrRouterTest {
@Test
void absentMemberClientSharesControlsForEveryTarget() {
FakeHerdr client = new FakeHerdr();
HerdrRouter router = new HerdrRouter(client, null, id -> id.equals("lead"));
assertSame(router.leadAgents(), router.memberAgents());
assertSame(router.leadSpaces(), router.memberSpaces());
assertSame(router.leadAgents(), router.agentsFor("lead"));
assertSame(router.leadAgents(), router.agentsFor("member"));
}
@Test
void separateClientsRouteLeadAndMemberTargets() {
FakeHerdr lead = new FakeHerdr();
FakeHerdr member = new FakeHerdr();
HerdrRouter router = new HerdrRouter(lead, member, id -> id.equals("lead"));
assertNotSame(router.leadAgents(), router.memberAgents());
assertNotSame(router.leadSpaces(), router.memberSpaces());
assertSame(router.leadAgents(), router.agentsFor("lead"));
assertSame(router.memberAgents(), router.agentsFor("member"));
}
}
@@ -24,4 +24,43 @@ class PaneLocatorTest {
assertNull(loc.terminalForPid(0));
assertNull(loc.terminalForPid(-1));
}
// --- two-daemon fallback (CB-185) -----------------------------------------
@Test
void fallsBackToTheMemberClientWhenTheLeadHasNoMatch() {
// The caller's pane lives on the member daemon only (e.g. the caller is a spawned
// member) — the lead client reports no panes at all, so the locator must fall back.
HerdrClient lead = new FakeHerdr().withNoPanes();
HerdrClient member = new FakeHerdr();
PaneLocator two = new PaneLocator(lead, member);
assertEquals("term_a", two.terminalForPid(FakeHerdr.WORKER_PID));
}
@Test
void searchesTheLeadClientBeforeTheMemberClient() {
// The caller's pane lives on the LEAD daemon (e.g. the caller is a peer lead) — with two
// daemons, resolving it must not depend on the member client having a matching pane.
HerdrClient lead = new FakeHerdr();
HerdrClient member = new FakeHerdr().withNoPanes();
PaneLocator two = new PaneLocator(lead, member);
assertEquals("term_a", two.terminalForPid(FakeHerdr.WORKER_PID));
}
@Test
void nullWhenNeitherClientHasTheMatch() {
PaneLocator two = new PaneLocator(new FakeHerdr().withNoPanes(), new FakeHerdr().withNoPanes());
assertNull(two.terminalForPid(FakeHerdr.WORKER_PID));
}
@Test
void collapsesToOneScanWhenLeadAndMemberAreTheSameClient() {
// The single-daemon deployment (no memberHerdrSocket configured): the two-arg constructor
// must behave exactly like the one-arg constructor, including making only one herdr call.
FakeHerdr shared = new FakeHerdr();
PaneLocator two = new PaneLocator(shared, shared);
assertEquals("term_a", two.terminalForPid(FakeHerdr.WORKER_PID));
long paneListCalls = shared.calls.stream().filter(c -> c.method().equals("pane.list")).count();
assertEquals(1, paneListCalls, "same-object lead/member must scan exactly once, not twice");
}
}
@@ -566,6 +566,73 @@ class CompletionResolverTest {
assertEquals("The usage limit has been reached.", waiter.getNow(null).text());
}
// --- fleetd#164 (part 2 addendum): narrow BACKEND_ERROR pattern classification ---------
@Test
void classifiesABackendErrorLineAsAFailureInsteadOfACompletedReply() {
String block = "⏺ API Error: 400 invalid request body\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
assertTrue(waiter.isDone(), "a backend-error scrape still resolves the blocked send");
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(),
"a backend rejection is a failure, not a completed reply");
}
@Test
void theBackendErrorReasonNamesTheMemberAndCarriesTheMatchedLine() {
String block = "⏺ API Error: 400 invalid request body\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
String reason = waiter.getNow(null).text();
assertTrue(reason.contains("term_a"), "the failure names the member: " + reason);
assertTrue(reason.contains("API Error: 400 invalid request body"),
"the failure carries the matched backend-error line: " + reason);
}
@Test
void aCaseInsensitiveApiErrorLineIsStillClassifiedAsABackendError() {
String block = "⏺ api error: rate limited\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(), "the pattern is case-insensitive");
}
@Test
void aBackendErrorFailureStillCarriesTheRestOfTheScrape() {
// The pattern is a heuristic: a member that forgot fleet_reply while *reporting on* a backend
// error matches it too. Failing is still correct, but the report itself must survive — losing
// it would be the same information-destroying defect fleetd#164 exists to fix.
String block = "\u23fa I looked into the gateway problem.\n"
+ "The log line was: API Error: 400 invalid request body\n"
+ "The cause is a missing content-type header.\n\u276f ";
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
String reason = waiter.getNow(null).text();
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind());
assertTrue(reason.contains("The cause is a missing content-type header."),
"the failure carries the rest of the pane, not only the matched line: " + reason);
}
@Test
void coverageIsOffWhenNoProfileHasAPatternConfigured() {
assertEquals("off (no profile has an exhaustedPattern configured; profiles: [terra])",
@@ -0,0 +1,75 @@
package dev.ltms.fleet.inject;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.HerdrRouter;
import dev.ltms.fleet.msg.TestTurnTokens;
import org.junit.jupiter.api.Test;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* CB-185: with a router split across two herdr daemons, {@link StatusPoller} must refine a raw
* {@code UNKNOWN} status by reading the pane content from the SAME daemon the status was sampled
* from — the lead daemon for a lead target, the member daemon for a member target. Reading the
* wrong daemon never finds the pane, classification stays {@code UNKNOWN} forever, and the
* status-gated {@link Injector} wedges: a queued message is never delivered.
*
* <p>This exercises the real production classes ({@code StatusPoller(HerdrRouter, ...)},
* {@code Injector(HerdrRouter, ...)}) wired together, not a hand-built object graph — the earlier
* three CB-185 bugs all passed exactly that kind of test while the real wiring stayed broken.
*/
class StatusPollerRoutingTest {
private static final String LEAD_TARGET = "term_a";
@Test
void refinesALeadTargetFromTheLeadDaemonAndDelivers() throws Exception {
// The lead daemon's pane is at a settled idle prompt; the member daemon's pane content is
// unclassifiable garbage. A correct refiner reads the LEAD daemon and delivers.
FakeHerdr leadHerdr = new FakeHerdr().agentStatus("unknown").readText("⏺ answer\n❯ ");
FakeHerdr memberHerdr = new FakeHerdr().agentStatus("unknown")
.readText("garbled ansi noise with no prompt");
HerdrRouter router = new HerdrRouter(leadHerdr, memberHerdr, LEAD_TARGET::equals);
Injector injector = new Injector(router, TurnListener.NOOP, _ -> true, _ -> {
});
StatusPoller poller = new StatusPoller(router, injector, 10);
poller.start();
try {
CompletableFuture<Void> delivered =
injector.enqueue(LEAD_TARGET, "via-poller", TestTurnTokens.inert(LEAD_TARGET));
// Must resolve quickly: refining against the WRONG daemon (member) never classifies
// out of UNKNOWN, so this would time out under the bug.
delivered.get(2, TimeUnit.SECONDS);
} finally {
poller.stop();
}
assertTrue(leadHerdr.called("agent.read"), "refine must probe the LEAD daemon's pane content");
}
@Test
void aLeadTargetNeverDeliversWhenOnlyTheMemberDaemonIsClassifiable() throws Exception {
// Inverted control: the member daemon's content WOULD classify to idle, but this is a lead
// target — a correct implementation must not use it, so delivery must NOT happen.
FakeHerdr leadHerdr = new FakeHerdr().agentStatus("unknown")
.readText("garbled ansi noise with no prompt");
FakeHerdr memberHerdr = new FakeHerdr().agentStatus("unknown").readText("⏺ answer\n❯ ");
HerdrRouter router = new HerdrRouter(leadHerdr, memberHerdr, LEAD_TARGET::equals);
Injector injector = new Injector(router, TurnListener.NOOP, _ -> true, _ -> {
});
StatusPoller poller = new StatusPoller(router, injector, 10);
poller.start();
try {
CompletableFuture<Void> delivered =
injector.enqueue(LEAD_TARGET, "via-poller", TestTurnTokens.inert(LEAD_TARGET));
assertThrows(TimeoutException.class, () -> delivered.get(500, TimeUnit.MILLISECONDS),
"a lead target must never be refined from the member daemon's pane content");
} finally {
poller.stop();
}
}
}
@@ -85,4 +85,34 @@ class StatusRefinerTest {
assertEquals(AgentStatus.UNKNOWN, refiner.refine("term_a", AgentStatus.UNKNOWN));
}
// --- refine(target, raw, control) — CB-185 per-call routing ---------------
@Test
void threeArgRefineReadsThroughTheGivenControlNotTheConstructedOne() {
// The refiner is CONSTRUCTED with one control (standing in for "the member daemon"), but
// a call names a DIFFERENT control (standing in for "the lead daemon") — the read must go
// to the one passed to the call, since that is the daemon the raw status came from.
FakeHerdr constructedWith = new FakeHerdr().readText("nothing recognizable here");
FakeHerdr passedToCall = new FakeHerdr().readText("⏺ answer\n❯ ");
StatusRefiner refiner = new StatusRefiner(new AgentControl(constructedWith));
AgentStatus result = refiner.refine("term_a", AgentStatus.UNKNOWN, new AgentControl(passedToCall));
assertEquals(AgentStatus.IDLE, result, "must classify from the PASSED control's pane content");
assertTrue(passedToCall.called("agent.read"));
assertFalse(constructedWith.called("agent.read"),
"the control fixed at construction must not be read when a call-site control is given");
}
@Test
void twoArgRefineStillReadsTheConstructedControl() {
// The legacy 2-arg overload (single-daemon callers) must keep using the constructed
// control — this is refine(target, raw, control) called with the field as `control`.
FakeHerdr herdr = new FakeHerdr().readText("⏺ answer\n❯ ");
StatusRefiner refiner = new StatusRefiner(new AgentControl(herdr));
assertEquals(AgentStatus.IDLE, refiner.refine("term_a", AgentStatus.UNKNOWN));
assertTrue(herdr.called("agent.read"));
}
}
@@ -1,5 +1,6 @@
package dev.ltms.fleet.member;
import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
@@ -1055,6 +1056,288 @@ class ClaudeCodeLauncherTest {
"a name already on allow: is covered, not a gap");
}
/**
* CB-633 follow-up (#192): the deny-by-default WARN wording is a promise an operator relies on —
* pinned byte-for-byte so a future edit cannot drift it (e.g. while picking wording for the
* allow-list path) without a test noticing.
*/
@Test
void denyByDefaultKeepsTheExactCredentialGapWarn() {
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile cfg = new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
List.of("claude"), "tab", "fleetd-workers", "worker: {profile} #{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, 0, System::currentTimeMillis, () -> {}, null, () -> TEST_MEMBER_CREDENTIALS,
() -> Set.of("A_BRAND_NEW_SECRET_TOKEN"));
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
try {
svc.spawn();
} finally {
logger.detachAppender(appender);
}
assertTrue(appender.list.stream().anyMatch(e ->
("memberCredentials gap: 1 credential-shaped env var name(s) are on neither "
+ "known: nor allow: — every member pane inherits them UNBLOCKED — "
+ "[A_BRAND_NEW_SECRET_TOKEN]. Add each to memberCredentials.known "
+ "(blocked by default) or .allow (if a member legitimately needs it).")
.equals(e.getFormattedMessage())),
"the deny-by-default WARN text must not drift — got: "
+ appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
}
/**
* CB-633 follow-up (#192), defect 1: under {@code policy: allow-list} on a zsh login shell the
* generated ZDOTDIR scrub genuinely blanks an unkept credential-shaped name, so the report must
* not say the pane inherits it UNBLOCKED — that claim is exactly what PR #174 got wrong. This
* goes through the real spawn path (not {@code HerdrPeerLauncherAllowListWiringTest}'s fixture,
* which overrides {@code buildLaunch} and bypasses none of the logic under test here — the
* shell-dependent branch lives in {@code applyEnvironmentAllowListPolicy}, which every spawn
* still passes through).
*/
@Test
void allowListPolicyOnZshReportsTheGapWithoutClaimingItIsUnblocked() {
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile cfg = new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
List.of("claude"), "tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
FleetConfig.MemberCredentials creds = new FleetConfig.MemberCredentials(
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST,
List.of("AI_GATEWAY_TOKEN"), List.of());
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
name -> "SHELL".equals(name) ? "/bin/zsh" : null,
0, System::currentTimeMillis, () -> {}, null, () -> creds,
() -> Set.of("AI_GATEWAY_TOKEN", "A_BRAND_NEW_SECRET_TOKEN"));
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
Level original = logger.getLevel();
logger.setLevel(Level.INFO);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
try {
svc.spawn();
} finally {
logger.detachAppender(appender);
logger.setLevel(original);
}
assertTrue(appender.list.stream().anyMatch(e ->
e.getFormattedMessage().contains("memberCredentials gap")
&& e.getFormattedMessage().contains("A_BRAND_NEW_SECRET_TOKEN")),
"the unkept name must still be reported — got: "
+ appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
assertFalse(appender.list.stream().anyMatch(e ->
e.getFormattedMessage().contains("memberCredentials gap")
&& e.getFormattedMessage().contains("UNBLOCKED")),
"on zsh the scrub genuinely blanks the name, so the report must not claim it is "
+ "inherited unblocked — got: "
+ appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
}
/**
* CB-633 follow-up (#192): the mirror of the zsh test above. On a non-zsh login shell {@code
* ZDOTDIR} is ignored, so no scrub ever runs — the report must keep the WARN wording (a name here
* really is inherited unblocked) rather than claiming a scrub protects it. This is the trap PR
* #174 fell into the other direction: keying the wording on the shell, not on {@code
* creds.isAllowList()}, is what keeps this branch correct.
*/
@Test
void allowListPolicyOnNonZshKeepsTheWarnWording() {
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile cfg = new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
List.of("claude"), "tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
FleetConfig.MemberCredentials creds = new FleetConfig.MemberCredentials(
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST,
List.of("AI_GATEWAY_TOKEN"), List.of());
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
name -> "SHELL".equals(name) ? "/bin/bash" : null,
0, System::currentTimeMillis, () -> {}, null, () -> creds,
() -> Set.of("AI_GATEWAY_TOKEN", "A_BRAND_NEW_SECRET_TOKEN"));
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
try {
svc.spawn();
} finally {
logger.detachAppender(appender);
}
assertTrue(appender.list.stream().anyMatch(e ->
e.getFormattedMessage().contains("memberCredentials gap")
&& e.getFormattedMessage().contains("UNBLOCKED")
&& e.getFormattedMessage().contains("A_BRAND_NEW_SECRET_TOKEN")),
"no scrub runs on a non-zsh shell, so the WARN wording must be kept — got: "
+ appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
assertFalse(appender.list.stream().anyMatch(e ->
e.getFormattedMessage().contains("memberCredentials gap")
&& e.getFormattedMessage().toLowerCase(java.util.Locale.ROOT).contains("scrub")),
"nothing is scrubbed on this path, so the report must not claim otherwise — got: "
+ appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
}
/**
* CB-633 follow-up (#192), defect 2: {@code memberCredentials} is a live, re-read-per-spawn
* supplier, so the policy can change between two spawns on the same launcher. Before this fix a
* single {@code AtomicBoolean} guarded both report kinds, so the harmless allow-list INFO on the
* first spawn would permanently suppress the real deny-by-default WARN a later spawn deserves.
* This goes through {@link ClaudeCodeLauncher#buildLaunch}'s real {@code baseEnv()} path — the
* WARN this test pins fires from {@code applyMemberCredentialPolicy}, which {@code
* HerdrPeerLauncherAllowListWiringTest}'s fixture never reaches at all (see its class javadoc).
*/
@Test
void secondSpawnStillWarnsAfterPolicyChangesFromAllowListToDenyByDefault() {
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile cfg = new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
List.of("claude"), "tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
AtomicReference<FleetConfig.MemberCredentials> creds = new AtomicReference<>(
new FleetConfig.MemberCredentials(FleetConfig.MemberCredentials.POLICY_ALLOW_LIST,
List.of("AI_GATEWAY_TOKEN"), List.of()));
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
name -> "SHELL".equals(name) ? "/bin/zsh" : null,
0, System::currentTimeMillis, () -> {}, null, creds::get,
() -> Set.of("AI_GATEWAY_TOKEN", "A_BRAND_NEW_SECRET_TOKEN"));
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
Level original = logger.getLevel();
logger.setLevel(Level.INFO);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
try {
logger.addAppender(appender);
svc.spawn(); // allow-list + zsh: harmless INFO, sets the allow-list guard only
creds.set(new FleetConfig.MemberCredentials(FleetConfig.MemberCredentials.POLICY_DENY_BY_DEFAULT,
List.of("AI_GATEWAY_TOKEN"), List.of()));
svc.spawn(); // policy reloaded to deny-by-default: this WARN must NOT be suppressed
} finally {
logger.detachAppender(appender);
logger.setLevel(original);
}
assertTrue(appender.list.stream().anyMatch(e ->
e.getFormattedMessage().contains("memberCredentials gap")
&& e.getFormattedMessage().contains("UNBLOCKED")),
"the second spawn's deny-by-default WARN must still fire even though the first "
+ "spawn's allow-list INFO already logged the same underlying gap — got: "
+ appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
}
/**
* CB-633 follow-up (#192), lead-review fix: {@code effectiveAllowed} is a SUPERSET of
* {@code known ∪ allow} — {@code MemberEnvAllowList.derive} also unions in every profile's
* {@code tokenEnv} (among other fields), so a credential-shaped name can be uncovered by
* {@code known:}/{@code allow:} and STILL survive the scrub because a profile's own
* {@code tokenEnv} names it. Here {@code tokenEnv} is deliberately set to a credential-shaped
* name the operator forgot to list — the misconfiguration this report exists to catch. The scrub
* genuinely keeps it, so the report must WARN, not claim (as the pre-lead-review cut of this fix
* did) that "no member pane keeps them".
*/
@Test
void allowListWarnsWhenTheDerivedAllowListKeepsAnUncoveredName() {
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile cfg = new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "SOME_LEAKY_TOKEN",
List.of("claude"), "tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
FleetConfig.MemberCredentials creds = new FleetConfig.MemberCredentials(
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST, List.of(), List.of());
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
name -> "SHELL".equals(name) ? "/bin/zsh" : null,
0, System::currentTimeMillis, () -> {}, null, () -> creds,
() -> Set.of("SOME_LEAKY_TOKEN"));
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
Level original = logger.getLevel();
logger.setLevel(Level.INFO);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
try {
svc.spawn();
} finally {
logger.detachAppender(appender);
logger.setLevel(original);
}
assertTrue(appender.list.stream().anyMatch(e ->
e.getLevel() == Level.WARN
&& e.getFormattedMessage().contains("memberCredentials gap")
&& e.getFormattedMessage().contains("SOME_LEAKY_TOKEN")
&& e.getFormattedMessage().contains("UNBLOCKED")),
"a name kept by the derived allow-list (via this profile's tokenEnv) must still WARN "
+ "— got: " + appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
assertFalse(appender.list.stream().anyMatch(e ->
e.getFormattedMessage().contains("memberCredentials gap")
&& e.getFormattedMessage().contains("SOME_LEAKY_TOKEN")
&& e.getFormattedMessage().contains("blanks them anyway")),
"the scrub does NOT blank this name, so the INFO wording must not claim it does — got: "
+ appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
}
/**
* CB-633 follow-up (#192), lead-review fix: a mixed gap — one name the derived allow-list keeps
* (this profile's {@code tokenEnv}), one it does not — must split cleanly: the WARN names only
* the kept one, the INFO names only the blanked one. Proves the split uses {@code
* MemberEnvAllowList.keeps} per-name rather than an all-or-nothing decision for the whole gap.
*/
@Test
void allowListSplitsAMixedGapBetweenTheWarnAndTheInfo() {
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile cfg = new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "SOME_LEAKY_TOKEN",
List.of("claude"), "tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
FleetConfig.MemberCredentials creds = new FleetConfig.MemberCredentials(
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST, List.of(), List.of());
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
name -> "SHELL".equals(name) ? "/bin/zsh" : null,
0, System::currentTimeMillis, () -> {}, null, () -> creds,
() -> Set.of("SOME_LEAKY_TOKEN", "A_BRAND_NEW_SECRET_TOKEN"));
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
Level original = logger.getLevel();
logger.setLevel(Level.INFO);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
try {
svc.spawn();
} finally {
logger.detachAppender(appender);
logger.setLevel(original);
}
boolean warnNamesOnlyKept = appender.list.stream().anyMatch(e ->
e.getLevel() == Level.WARN
&& e.getFormattedMessage().contains("memberCredentials gap")
&& e.getFormattedMessage().contains("SOME_LEAKY_TOKEN")
&& !e.getFormattedMessage().contains("A_BRAND_NEW_SECRET_TOKEN"));
boolean infoNamesOnlyBlanked = appender.list.stream().anyMatch(e ->
e.getLevel() == Level.INFO
&& e.getFormattedMessage().contains("memberCredentials gap")
&& e.getFormattedMessage().contains("A_BRAND_NEW_SECRET_TOKEN")
&& !e.getFormattedMessage().contains("SOME_LEAKY_TOKEN"));
assertTrue(warnNamesOnlyKept, "the WARN must name the derived-list-kept variable and only it "
+ "— got: " + appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
assertTrue(infoNamesOnlyBlanked, "the INFO must name the scrub-blanked variable and only it "
+ "— got: " + appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
}
/**
* The half of CB-592 that can actually survive the pane's login shell. BRIDGED_MEMBER is a name
* secrets.sh never exports, so nothing overwrites it — measured: GITEA_TOKEN is injected the
@@ -346,13 +346,94 @@ class CompositePeerLauncherTest {
@Test
void stopRejectsAnUnownedPaneIdWhenMultipleDaemonsCouldOwnIt() {
// CB-185 blocker 1: genuine ambiguity — pane ids are per-daemon counters, so two daemons
// can each really hold an agent at "w1:p1". Neither claims ownership through spawnedBy
// (empty, as after a restart), so the probe must find BOTH and refuse rather than guess.
FakeHerdr first = new FakeHerdr().withAgent("x", "term_x", "w1:p1", "w1:t1");
FakeHerdr second = new FakeHerdr().withAgent("y", "term_y", "w1:p1", "w1:t1");
PeerLauncher composite = new CompositePeerLauncher(
List.of(claudeAdapter(new FakeHerdr()), opencodeAdapter(new FakeHerdr())), "claude");
List.of(claudeAdapter(first), opencodeAdapter(second)), "claude");
IllegalArgumentException error = assertThrows(IllegalArgumentException.class,
() -> composite.stop("w1:p1"));
assertEquals("ambiguous paneId 'w1:p1': no owning herdr daemon was recorded", error.getMessage());
assertEquals("ambiguous paneId 'w1:p1': 2 configured herdr daemons report this pane — "
+ "no way to tell which one the caller means", error.getMessage());
}
@Test
void stopOnAPaneNoConfiguredDaemonKnowsIsTreatedAsAlreadyStopped() {
// CB-185 blocker 1, the zero-owner branch: spawnedBy is empty (as after a restart) and
// neither daemon's agent.list mentions this pane at all — it is already gone. A retried
// stop() on an already-gone pane must succeed quietly, not refuse forever.
FakeHerdr first = new FakeHerdr();
FakeHerdr second = new FakeHerdr();
PeerLauncher composite = new CompositePeerLauncher(
List.of(claudeAdapter(first), opencodeAdapter(second)), "claude");
assertDoesNotThrow(() -> composite.stop("w1:p1"));
assertFalse(first.called("pane.close"), "no owner was found, so no delegate is told to close anything");
assertFalse(second.called("pane.close"), "no owner was found, so no delegate is told to close anything");
}
@Test
void stopWithEmptySpawnedByResolvesTheOwnerThroughAProbeAndSkipsTheOtherDaemon() {
// CB-185 blocker 1, the main fix: after a restart spawnedBy is empty for every surviving
// member. stop() must still find the one daemon that actually knows the pane and route
// only to it — never touching the daemon that never held it.
FakeHerdr first = new FakeHerdr().withAgent("x", "term_x", "w1:p1", "w1:t1");
FakeHerdr second = new FakeHerdr();
PeerLauncher composite = new CompositePeerLauncher(
List.of(claudeAdapter(first), opencodeAdapter(second)), "claude");
composite.stop("w1:p1");
assertTrue(first.calls.stream().anyMatch(c -> c.method().equals("pane.close")
&& "w1:p1".equals(((Map<?, ?>) c.params()).get("pane_id"))),
"the daemon that actually knows the pane closes it");
assertFalse(second.called("pane.close"), "the daemon that never held the pane is never touched");
}
@Test
void aProbeSurvivesOneUnreachableDaemonAndStillFindsTheOwnerOnTheOtherOne() {
// CB-185 blocker 1 (lead review): a daemon that is DOWN while we probe must not abort the
// whole probe — the pane the OPERATOR actually wants stopped can live on a different,
// healthy daemon, and that pane must not become un-stoppable because a third one is down.
FakeHerdr down = new FakeHerdr().healthy(false);
FakeHerdr owner = new FakeHerdr().withAgent("x", "term_x", "w1:p1", "w1:t1");
PeerLauncher composite = new CompositePeerLauncher(
List.of(claudeAdapter(down), opencodeAdapter(owner)), "claude");
assertDoesNotThrow(() -> composite.stop("w1:p1"),
"the unreachable daemon must be skipped, not fail the whole stop");
assertTrue(owner.calls.stream().anyMatch(c -> c.method().equals("pane.close")
&& "w1:p1".equals(((Map<?, ?>) c.params()).get("pane_id"))),
"the healthy daemon that actually owns the pane still closes it");
}
@Test
void aProbedOwnerIsCachedSoARetryAfterAFailedStopNeedsNoSecondProbe() {
// CB-185 blocker 1: the probe's whole point is to be cheap on repeat — a failed stop (e.g.
// "pane_busy") must not force another agent.list() round trip on every retry.
FakeHerdr first = new FakeHerdr().withAgent("x", "term_x", "w1:p1", "w1:t1")
.paneCloseFailsWith("pane_busy");
FakeHerdr second = new FakeHerdr();
CompositePeerLauncher composite = new CompositePeerLauncher(
List.of(claudeAdapter(first), opencodeAdapter(second)), "claude");
assertThrows(HerdrException.class, () -> composite.stop("w1:p1"));
long listCallsAfterFirst = first.calls.stream().filter(c -> c.method().equals("agent.list")).count()
+ second.calls.stream().filter(c -> c.method().equals("agent.list")).count();
assertTrue(listCallsAfterFirst > 0, "the first stop needed a probe");
assertThrows(HerdrException.class, () -> composite.stop("w1:p1"),
"still failing on the retry, but through the cached owner");
long listCallsAfterSecond = first.calls.stream().filter(c -> c.method().equals("agent.list")).count()
+ second.calls.stream().filter(c -> c.method().equals("agent.list")).count();
assertEquals(listCallsAfterFirst, listCallsAfterSecond,
"the retry is served from the cache — no additional agent.list probe");
}
@Test
@@ -86,6 +86,28 @@ class MessageServiceTest {
assertTrue(reply.completed(), "a scraped completion still counts as completed");
}
@Test
void backendErrorScrapeThroughMessageServiceFailsInsteadOfBecomingReplyText() throws Exception {
// fleetd#164 (part 2 addendum): a scrape that reads cleanly but is only the backend's own
// rejection (e.g. an HTTP 400) must reach the caller as WORKER_FAILED, not as a completed
// reply whose text happens to be the error line.
CompletableFuture<MessageService.Reply> send = sendAsync();
awaitWaiting();
herdr.readText("$ prompt");
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
herdr.readText("⏺ API Error: 400 invalid request body");
injector.onStatus(T, AgentStatus.IDLE);
MessageService.Reply reply = send.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.WORKER_FAILED, reply.outcome(),
"a backend rejection must use the caller's failure outcome, not a completed reply");
assertFalse(reply.completed(), "plain backend errors are never fallback reply content");
assertTrue(reply.text().contains("API Error: 400 invalid request body"),
"the visible backend error is carried as the failure reason: " + reply.text());
}
@Test
void explicitFleetReplyResolvesAsReplied() throws Exception {
CompletableFuture<MessageService.Reply> send = sendAsync();
@@ -688,6 +710,68 @@ class MessageServiceTest {
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
}
// --- #137: a fleet_ask round-trip must not orphan the ticket's own reply -------------------
//
// The primary's fleet_send{turnId} answer call is itself bounded (a real MCP call, capped well
// under a minute) — far shorter than a resumed turn can genuinely take to finish real work. These
// drive the exact real delegation path (async send -> worker asks -> primary answers -> primary's
// own wait gives up -> worker's real fleet_reply arrives afterwards) rather than calling a reply
// sink directly, since the bug is specifically about which sink the resumed turn's reply reaches.
@Test
void aReplyAfterAnswerTimesOutStillCompletesTheAsyncTicket() throws Exception {
String ticket = messages.sendAsync(T, "task that asks");
awaitWaiting();
injectDelivery();
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
// The primary answers, but its own bounded wait for the worker's resumed turn is short and
// expires before the worker (still genuinely working) gets back to it.
MessageService.Reply answerReply = messages.answer(asking.turnId(), "config.yaml", 150);
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome(),
"the primary's own bounded wait gives up before the worker finishes resuming");
// The worker keeps working past that window and only now calls fleet_reply.
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");
}
@Test
void fleetStopAfterAnOrphanedReplyDoesNotFailTheTicket() throws Exception {
String ticket = messages.sendAsync(T, "task that asks");
awaitWaiting();
injectDelivery();
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
MessageService.Reply answerReply = messages.answer(asking.turnId(), "config.yaml", 150);
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome());
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
// fleet_stop tears the worker's session down right after the reply landed — this must never
// report the misleading "the worker session was released before it replied": a reply is
// exactly what happened.
assertFalse(messages.abandon(T, "the worker session was released before it replied"),
"a reply already arrived, so nothing here is a genuine failure");
MessageService.TaskView view = awaitTicketPhase(ticket, MessageService.Phase.DONE);
assertEquals("PR opened: https://example/pulls/42", view.reply());
}
@Test
void unansweredAsyncQuestionReturnsTheTicketToPendingAndReleasesItsTarget() throws Exception {
String ticket = messages.sendAsync(T, "task that asks");
@@ -1057,6 +1141,74 @@ class MessageServiceTest {
}
}
/**
* #197: the ticket TTL must run from COMPLETION, not from creation.
*
* <p>It used to compare the cutoff against {@code createdNanos}, so the real window to collect a
* reply was {@code TTL minus however long the task ran}. A delegation that ran longer than the
* TTL was already past the cutoff the moment it finished, so the very next prune destroyed its
* reply — and the reply lives only in the task's future, so nothing could get it back. That is
* the normal case for real work here, not an edge case: three workers in one session ran well
* past ten minutes and two of their complete reports were lost this way.
*
* <p>The task below runs for longer than the whole TTL before it replies, which is exactly the
* shape that used to lose everything. Remove the fix and this fails: {@code poll} returns
* {@code null} because the ticket was pruned on arrival.
*/
@Test
void aTaskRunningLongerThanTheTtlStillKeepsItsReport() throws Exception {
java.util.concurrent.atomic.AtomicLong clock = new java.util.concurrent.atomic.AtomicLong(1_000_000_000L);
try (var wiring = wireWithPushLoop(1, 50, clock::get)) {
String slow = wiring.service().sendAsync(T, "a task that takes longer than the TTL");
awaitWaiting();
injectDelivery();
// The worker is still working, and has been for longer than the entire TTL. Nothing may
// be pruned yet — the ticket has not finished, so there is no report to keep or lose.
clock.addAndGet(MessageService.TICKET_TTL_NANOS + TimeUnit.SECONDS.toNanos(30));
// Only now does it reply. Under the old clock this reply was born already expired.
assertTrue(rendezvous.resolve(T, "the long report"));
awaitTicketPhaseOn(wiring.service(), slow, MessageService.Phase.DONE);
// A second delegation runs pruneTerminalTickets before it returns.
wiring.service().sendAsync(T, "an unrelated second task");
MessageService.TaskView view = wiring.service().poll(slow);
assertNotNull(view, "a ticket that completed just now must survive the prune, however "
+ "long its task ran — the TTL is the window to COLLECT the report, not the "
+ "budget for producing it");
assertEquals(MessageService.Phase.DONE, view.phase());
assertEquals("the long report", view.reply(),
"the worker's actual report must still be there, not just the ticket");
}
}
/**
* The other half of #197: the TTL must still bound {@code tasks}. Measuring from completion
* would be a leak if a finished ticket were then kept forever, so this pins the eviction that
* still has to happen — the same ticket, left uncollected for longer than the TTL AFTER it
* finished, is gone.
*/
@Test
void aFinishedTicketIsStillPrunedOnceTheTtlPassesSinceItFinished() throws Exception {
java.util.concurrent.atomic.AtomicLong clock = new java.util.concurrent.atomic.AtomicLong(1_000_000_000L);
try (var wiring = wireWithPushLoop(1, 50, clock::get)) {
String done = wiring.service().sendAsync(T, "a quick task");
awaitWaiting();
injectDelivery();
assertTrue(rendezvous.resolve(T, "quick result"));
awaitTicketPhaseOn(wiring.service(), done, MessageService.Phase.DONE);
// Nobody collected it, and the TTL has now passed since it FINISHED.
clock.addAndGet(MessageService.TICKET_TTL_NANOS + TimeUnit.SECONDS.toNanos(1));
wiring.service().sendAsync(T, "an unrelated second task");
assertNull(wiring.service().poll(done),
"the TTL must still evict an uncollected finished ticket, or tasks grows forever");
}
}
private MessageService.TaskView awaitTicketPhaseOn(MessageService svc, String ticket,
MessageService.Phase phase) throws Exception {
long deadline = System.currentTimeMillis() + 3000;
@@ -0,0 +1,150 @@
package dev.ltms.fleet.rest;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.HerdrClient;
import io.javalin.Javalin;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Test;
import java.net.URI;
import java.net.http.HttpClient;
import java.net.http.HttpRequest;
import java.net.http.HttpResponse;
import static org.junit.jupiter.api.Assertions.*;
/**
* CB-185: with a router split across two herdr daemons (lead + {@code memberHerdrSocket}),
* {@link FleetApp#healthz} must require BOTH daemons to answer and {@link FleetApp#sessions}
* (which the {@code GET /sessions} route calls) must merge workspaces from both — the bug this
* guards against had {@code FleetApp} constructed with the raw lead-only client, so a down member
* daemon was invisible behind a green {@code /healthz} (every spawn then fails) and every member
* workspace was silently dropped from {@code GET /sessions}.
*
* <p>Builds the real {@link FleetApp} directly (not a hand-rolled stand-in) against only the two
* herdr clients — the other collaborators are unused by the two routes under test here.
*/
class FleetAppTwoDaemonTest {
private final HttpClient http = HttpClient.newHttpClient();
private Javalin app;
@AfterEach
void stop() {
if (app != null) app.stop();
}
private int start(HerdrClient lead, HerdrClient member) {
app = new FleetApp(lead, member, null, null, null, null, null, null, null, ignored -> false)
.build().start("127.0.0.1", 0);
return app.port();
}
private HttpResponse<String> get(int port, String path) throws Exception {
HttpRequest req = HttpRequest.newBuilder(URI.create("http://127.0.0.1:" + port + path)).GET().build();
return http.send(req, HttpResponse.BodyHandlers.ofString());
}
@Test
void healthzIsGreenWhenBothDaemonsAnswer() throws Exception {
int port = start(new FakeHerdr(), new FakeHerdr());
assertEquals(200, get(port, "/healthz").statusCode());
}
@Test
void healthzIsDegradedWhenOnlyTheMemberDaemonIsDown() throws Exception {
int port = start(new FakeHerdr(), new FakeHerdr().healthy(false));
HttpResponse<String> res = get(port, "/healthz");
assertEquals(503, res.statusCode(),
"a down MEMBER daemon must not be masked by a healthy lead — every spawn goes "
+ "through the member daemon");
}
@Test
void healthzIsDegradedWhenOnlyTheLeadDaemonIsDown() throws Exception {
int port = start(new FakeHerdr().healthy(false), new FakeHerdr());
assertEquals(503, get(port, "/healthz").statusCode());
}
@Test
void healthzMakesExactlyOneCallWhenLeadAndMemberAreTheSameClient() throws Exception {
// Single-daemon deployment (no memberHerdrSocket) — must be byte-for-byte the old
// behaviour: one ping call, 200 on success.
FakeHerdr shared = new FakeHerdr();
int port = start(shared, shared);
assertEquals(200, get(port, "/healthz").statusCode());
long pings = shared.calls.stream().filter(c -> c.method().equals("ping")).count();
assertEquals(1, pings, "single-daemon deployment must make exactly one ping call");
}
@Test
void sessionsMergesWorkspacesFromBothDaemons() throws Exception {
FakeHerdr lead = new FakeHerdr();
FakeHerdr member = new FakeHerdr().withWorkspace("w9", "member-only-workspace");
int port = start(lead, member);
HttpResponse<String> res = get(port, "/sessions");
assertEquals(200, res.statusCode(), res.body());
assertTrue(res.body().contains("member-only-workspace"),
"GET /sessions must not silently drop the member daemon's workspaces");
}
@Test
void sessionsMakesExactlyOneWorkspaceListCallWhenLeadAndMemberAreTheSameClient() throws Exception {
FakeHerdr shared = new FakeHerdr();
int port = start(shared, shared);
assertEquals(200, get(port, "/sessions").statusCode());
long calls = shared.calls.stream().filter(c -> c.method().equals("workspace.list")).count();
assertEquals(1, calls, "single-daemon deployment must call workspace.list exactly once");
}
// ── CB-185 blocker 2: /healthz must report the MEMBER daemon's protocol too ────────────────
@Test
void healthzReportsBothDaemonsWhenTheirProtocolsDiffer() throws Exception {
FakeHerdr lead = new FakeHerdr().pingReports("0.8.0", 19);
FakeHerdr member = new FakeHerdr().pingReports("0.7.0", 18);
int port = start(lead, member);
HttpResponse<String> res = get(port, "/healthz");
assertEquals(200, res.statusCode(), res.body());
assertTrue(res.body().contains("\"protocol\":19"),
"the herdr key keeps reporting the LEAD's protocol, unchanged: " + res.body());
assertTrue(res.body().contains("\"member\""), "a separate member key is present: " + res.body());
assertTrue(res.body().contains("\"protocol\":18"),
"the member key reports the member daemon's own protocol: " + res.body());
assertTrue(res.body().contains("\"protocolMismatch\":true"),
"a differing protocol is called out explicitly, not left to be spotted by eye: " + res.body());
}
@Test
void healthzReportsBothDaemonsWithNoMismatchWhenProtocolsMatch() throws Exception {
int port = start(new FakeHerdr(), new FakeHerdr());
HttpResponse<String> res = get(port, "/healthz");
assertEquals(200, res.statusCode(), res.body());
assertTrue(res.body().contains("\"member\""), "the member key is present whenever a second daemon "
+ "is configured, even when the protocols happen to agree: " + res.body());
assertFalse(res.body().contains("protocolMismatch"),
"matching protocols must not raise a mismatch flag: " + res.body());
}
@Test
void healthzWithOneDaemonCarriesNoMemberOrMismatchKey() throws Exception {
// The single-daemon deployment (no memberHerdrSocket) must see no change at all beyond the
// historical body: no "member" key, no "protocolMismatch" key. (Map.of()'s own key order is
// JVM-salted regardless of this fix, so this checks content, not exact key order.)
FakeHerdr shared = new FakeHerdr();
int port = start(shared, shared);
HttpResponse<String> res = get(port, "/healthz");
assertEquals(200, res.statusCode());
assertTrue(res.body().contains("\"status\":\"ok\""), res.body());
assertTrue(res.body().contains("\"protocol\":19"), res.body());
assertTrue(res.body().contains("\"version\":\"0.8.0\""), res.body());
assertFalse(res.body().contains("\"member\""), "no second daemon configured, so no member key: " + res.body());
assertFalse(res.body().contains("protocolMismatch"), res.body());
}
}
@@ -1,13 +1,23 @@
package dev.ltms.fleet.session;
import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.LoggerContext;
import ch.qos.logback.classic.spi.IThrowableProxy;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import org.slf4j.LoggerFactory;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.TimeUnit;
@@ -352,6 +362,151 @@ class GitWorktreesTest {
assertEquals("worktree origin contains HTTPS user info; refusing provision", error.getMessage());
}
// ---- CB-189: broader remote-URL coverage — every remote, both fetch and push URLs, any
// non-SSH scheme. Reporting only, additive to the origin/https strip-and-refuse tests above. ----
private Logger reportingLogger;
private ListAppender<ILoggingEvent> reportingAppender;
/** {@link GitWorktrees}'s own logger, captured fresh for each test so assertions never see a
* message left over from a previous test. */
@BeforeEach
void attachReportingLogCapture() {
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
reportingLogger = ctx.getLogger(GitWorktrees.class);
reportingLogger.setLevel(Level.WARN);
reportingAppender = new ListAppender<>();
reportingAppender.setContext(ctx);
reportingAppender.start();
reportingLogger.addAppender(reportingAppender);
}
@AfterEach
void detachReportingLogCapture() {
reportingLogger.detachAppender(reportingAppender);
}
private List<String> capturedMessages() {
return reportingAppender.list.stream().map(ILoggingEvent::getFormattedMessage).toList();
}
/** Asserts {@code secret} appears in no captured message, and in no attached exception's
* message either — the constraint is that a credential must never reach a log, however it
* would have gotten there. */
private void assertNoLeak(String secret) {
for (ILoggingEvent event : reportingAppender.list) {
assertFalse(event.getFormattedMessage().contains(secret),
"log message leaked a credential (" + secret + "): " + event.getFormattedMessage());
IThrowableProxy thrown = event.getThrowableProxy();
if (thrown != null && thrown.getMessage() != null) {
assertFalse(thrown.getMessage().contains(secret),
"logged exception leaked a credential (" + secret + "): " + thrown.getMessage());
}
}
}
/** Gap 1: only {@code origin} was ever inspected. A credential on any other remote's fetch URL
* must now be reported. */
@Test
void aCredentialedUrlOnANonOriginRemoteIsReported(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
git(repo, "remote", "add", "origin", "https://git.ltms.dev/akb/kb.git");
git(repo, "remote", "add", "upstream", "https://leaky-upstream-token@git.ltms.dev/akb/kb.git");
new GitWorktrees(tmp.resolve("wts").toString()).add(repo.toString(), "cb-189-a", "HEAD");
List<String> messages = capturedMessages();
assertTrue(messages.stream().anyMatch(m -> m.contains("remote=upstream")),
"expected a report naming the leaking non-origin remote:\n" + messages);
assertNoLeak("leaky-upstream-token");
assertNoLeak("https://leaky-upstream-token@git.ltms.dev/akb/kb.git");
assertNoLeak("git.ltms.dev");
}
/** Gap 2: push URLs were never inspected. A credential visible only on {@code pushurl} — the
* fetch URL for the same remote stays clean — must now be reported. */
@Test
void aCredentialedPushUrlIsReported(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
git(repo, "remote", "add", "origin", "https://git.ltms.dev/akb/kb.git");
git(repo, "remote", "add", "mirror", "https://git.ltms.dev/akb/mirror.git");
git(repo, "remote", "set-url", "--push", "mirror",
"https://leaky-push-token@git.ltms.dev/akb/mirror.git");
new GitWorktrees(tmp.resolve("wts").toString()).add(repo.toString(), "cb-189-b", "HEAD");
List<String> messages = capturedMessages();
assertTrue(messages.stream().anyMatch(m -> m.contains("remote=mirror")),
"expected a report naming the remote with the leaking pushurl:\n" + messages);
assertNoLeak("leaky-push-token");
assertNoLeak("https://leaky-push-token@git.ltms.dev/akb/mirror.git");
assertNoLeak("git.ltms.dev");
}
/** Gap 3: only {@code https} was handled. A plain {@code http://user:pass@…} remote — worse
* than https, not better — must now be reported. */
@Test
void anHttpUrlWithCredentialsIsReported(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
git(repo, "remote", "add", "origin", "https://git.ltms.dev/akb/kb.git");
git(repo, "remote", "add", "insecure", "http://plainuser:plainpass@git.ltms.dev/akb/kb.git");
new GitWorktrees(tmp.resolve("wts").toString()).add(repo.toString(), "cb-189-c", "HEAD");
List<String> messages = capturedMessages();
assertTrue(messages.stream().anyMatch(m -> m.contains("remote=insecure")),
"expected a report for the credentialed plain-http remote:\n" + messages);
assertNoLeak("plainuser");
assertNoLeak("plainpass");
assertNoLeak("plainuser:plainpass");
assertNoLeak("git.ltms.dev");
}
/** A normal {@code ssh://} remote and a credential-free {@code https://} remote must produce no
* report at all — the check must not cry wolf on ordinary, safe configuration. */
@Test
void anSshRemoteAndACleanHttpsRemoteProduceNoReport(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
git(repo, "remote", "add", "origin", "ssh://git@git.ltms.dev:2224/akb/kb.git");
git(repo, "remote", "add", "clean", "https://git.ltms.dev/akb/kb.git");
new GitWorktrees(tmp.resolve("wts").toString()).add(repo.toString(), "cb-189-d", "HEAD");
assertTrue(reportingAppender.list.isEmpty(),
"expected no report for an ssh remote and a credential-free https remote, got:\n"
+ capturedMessages());
}
/**
* CB-189 review fix. Even a FAILING command that read a credential onto its stdout must never
* let that value reach the thrown {@link WorktreeException}'s message — this is gap 4 from the
* CB-189 issue, and the reason {@link GitWorktrees#execRedacted} exists at all. Drives the
* shared {@code exec}/{@code execRedacted} seam directly (it is package-private for exactly this,
* the same way the {@code afterWorktreeAdded} constructor parameter is a test seam) with a
* synthetic, non-git command whose stdout carries a marker — passed through the environment,
* never through argv, so the marker cannot leak via the command line that IS always printed
* unconditionally in the exception message — and which exits non-zero. This isolates the
* redaction guarantee itself rather than depending on a specific git failure mode that happens to
* echo a URL onto stdout before failing: none of the git subcommands this class actually runs was
* found to have one (a corrupted config makes {@code git config --get} fail before it ever reads
* the target key, so its output never carries the URL either). The marker is generated per-test
* run and injected only by the test, never a real-looking credential, so even a failing assertion
* could not itself print a secret.
*/
@Test
void execRedactedNeverCopiesFailingCommandOutputIntoTheExceptionMessage(@TempDir Path tmp) {
String marker = "cb189-marker-" + System.nanoTime();
GitWorktrees worktrees = new GitWorktrees(tmp.toString());
WorktreeException thrown = assertThrows(WorktreeException.class, () -> worktrees.exec(
Map.of("MARKER", marker), true, "sh", "-c", "echo \"$MARKER\"; exit 7"));
assertNotNull(thrown.getMessage());
assertFalse(thrown.getMessage().contains(marker),
"a failing command's captured stdout leaked into the exception message: "
+ thrown.getMessage());
}
/** Neutralizing must not look like work in progress, or a worker would commit it into its PR. */
@Test
void theNeutralizedConfigIsNotAPendingLocalModification(@TempDir Path tmp) throws Exception {