Compare commits

...

30 Commits

Author SHA1 Message Date
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
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
Dai Ha 11c3ff67b6 fleets-status: fleet01 IS ssh-reachable; correct the 'denied' claim
CI / contract (push) Successful in 46s
CI / build (push) Successful in 1m12s
The skill said SSH to fleet01 is denied, so every report wrote 'not
reachable' for that fleet's daemon facts. That is true only for the user
dai.ha. The host alias fleet01 maps to user ltms and key auth works.

Checked 2026-08-28 while measuring #185: ssh fleet01 connects, and ltms
has passwordless sudo there. So fleet01's PID, uptime, jar and /healthz
can be reported over SSH even though its REST port is unreachable.
2026-08-28 09:06:54 +07:00
Dai Ha d867c87100 #184: correct the false ssh-agent premise in the URL-rewrite javadoc
CI / contract (push) Successful in 46s
CI / build (push) Successful in 1m37s
The javadoc said a member cannot authenticate at all once memberCredentials
blocks SSH_AUTH_SOCK, "there is no private key file on this host, only an
ssh-agent socket". That is wrong, and it was written after looking only in
~/.ssh, which holds nothing but Include lines.

Measured: ssh -G git.ltms.dev resolves an IdentityFile under the shared-env
directory. That file exists, is readable by this user, and has no passphrase.
A live member with SSH_AUTH_SOCK blanked pushed to the forge over SSH.

The rewrite itself is unchanged and still worth having. Only its stated reason
was wrong: it routes a member through its own scoped token instead of the
operator's ssh identity, which is what makes a member's pushes attributable
and revocable. It is not what stands between a member and the forge.
2026-08-28 06:35:46 +07:00
Dai Ha b0c4cedfab #157: redact remote-URL user-info before it reaches a log
CI / contract (push) Successful in 41s
CI / build (push) Successful in 1m25s
The four log lines added with the worktree HTTPS rewrite echoed the origin
URL verbatim, and one of them echoed the ssh:// authority, which carries
user-info. An ssh authority is normally just git@, so in practice this
changes nothing -- but a remote URL is not obviously a credential channel,
and that is precisely why one has leaked here three times (#157, #182).

Redact at the log call, not after it surprises someone.
2026-08-28 06:24:08 +07:00
ltms 85417d5215 #157: rewrite an SSH origin to HTTPS inside the provisioned worktree only
CI / contract (push) Successful in 46s
CI / build (push) Successful in 1m37s
Git never consults a credential.helper for an SSH transport, so #177's helper was inert on this repo — whose origin is ssh://. Once allow-list policy blocks SSH_AUTH_SOCK, a member on an SSH origin cannot authenticate at all: there is no private key file on this host, only an agent socket.

A worktree-scoped `url.<https>.insteadOf <ssh>` gives the member HTTPS for fetch and push while the primary checkout keeps SSH untouched. Host and port are parsed from the origin, never hardcoded — a test with a synthetic host proves it. The scp-like shorthand is left alone deliberately, since its host:path split is defined by ssh_config aliases rather than URI syntax.

Verified by the lead in an independent worktree: Tests run: 986, Failures: 0, Errors: 0, Skipped: 0, BUILD SUCCESS.
2026-08-28 01:22:57 +02:00
Dai Ha 4accc746bd #157: rewrite SSH origin to HTTPS in the worktree so the credential helper is reachable
CI / build (pull_request) Successful in 1m10s
CI / contract (pull_request) Successful in 1m16s
2026-08-28 06:20:10 +07:00
ltms 21c4c8cbef #157: keep the forge token out of git config, via an environment credential helper
CI / build (push) Successful in 1m4s
CI / contract (push) Successful in 47s
The member credential scrub removes environment variables. It cannot remove a token written into git config inside the repo the member works in, so `git remote -v` handed a member a credential it was deliberately not given.

Provisioning now strips HTTPS user info from the origin before `git worktree add`, refuses the worktree if user info survives, and configures a per-worktree credential helper that reads WORKER_GITEA_TOKEN at call time. Nothing is persisted.

The helper emits BOTH username and password, and resets the inherited helper list first. An earlier revision emitted only `username=`, which made git fall through to the next helper — on a Mac that is osxkeychain, so a member would have authenticated with the operator's stored credential while every test passed and `git remote -v` looked clean. See #182.

`worktreeCredentialHelperCompletesWithoutUsingAnInheritedHelper` plants a synthetic operator helper in an isolated global config and proves the worktree helper wins. The worker confirmed it fails when the reset is removed. All credential tests pin GIT_CONFIG_GLOBAL and GIT_CONFIG_SYSTEM so they can neither read nor write real credentials.

Verified by the lead in an independent worktree: Tests run: 973, Failures: 0, Errors: 0, Skipped: 0, BUILD SUCCESS.
2026-08-28 01:10:38 +02:00
ltms 430f5b0dae #164: never resolve a send with an empty scrape or a sub-floor turn
CI / contract (push) Successful in 53s
CI / build (push) Successful in 1m27s
A turn that dies on a backend error produces the same working -> idle transition as a real one, just faster and with nothing on screen. The resolver accepted that as a completed turn and handed the caller HTTP 200 with an empty reply, so a lost turn and a successful empty answer were indistinguishable.

Now: an empty or unreadable scrape fails, naming the member; and a BUSY -> DONE inside MIN_TURN_NANOS (2s) fails as a crash signature.

One existing test encoded the bug — it asserted a failed scrape resolved as a success carrying "" — and has been inverted rather than worked around.

Verified by the lead in an independent worktree: Tests run: 973, Failures: 0, Errors: 0, Skipped: 0, BUILD SUCCESS.
2026-08-28 01:10:27 +02:00
ltms 42731833d0 CB-633: union memberCredentials.allow into the member env allow-list
CI / contract (push) Successful in 42s
CI / build (push) Successful in 1m38s
`policy: allow-list` silently ignored every name an operator wrote under `allow:` unless a profile happened to carry it too, so turning the policy on would have blanked credentials working members depend on. Derivation now unions the operator's list.

`SSH_AUTH_SOCK` stays governed only by `sshAuthSock`, even when listed under `allow:` — it is a live handle to the operator's ssh-agent, not a value.

Adds one INFO line per allow-list spawn, `member credentials: allowed N of M`, emitted only after the shell gate so it can never report coverage on a path where the scrub does not run.

Verified by the lead in an independent worktree: Tests run: 976, Failures: 0, Errors: 0, Skipped: 0, BUILD SUCCESS.
2026-08-28 01:10:18 +02:00
ltms 65acf066ad #154: pin the AMQP reply inbox prefetch bound
CI / contract (push) Successful in 1m26s
CI / build (push) Successful in 1m39s
No behaviour change. #154 supposed that ownership drains a whole queue into the in-memory `held` map, making `x-max-length` and per-message TTL decorative. Measurement says otherwise: `basicConsume` is manual-ack, `deliverCallback` acks only duplicates, and `basicQos` is set on the one shared channel before any consumer starts — so total `held` is bounded by the prefetch window across all targets.

Adds a fake-broker test that drives the real `own()` path and fails if receipt ever starts acking, plus a javadoc line naming the prefetch window at the point of first mention.

Verified by the lead in an independent worktree: Tests run: 971, Failures: 0, Errors: 0, Skipped: 0, BUILD SUCCESS.
2026-08-28 01:06:43 +02:00
Dai Ha ee8f570fd7 #157: isolate worktree credential helpers
CI / contract (pull_request) Successful in 1m5s
CI / build (pull_request) Successful in 2m46s
2026-08-28 06:06:12 +07:00
Dai Ha fa3f910d44 #154: pin AMQP reply inbox prefetch
CI / build (pull_request) Successful in 2m16s
CI / contract (pull_request) Successful in 2m18s
2026-08-28 06:03:54 +07:00
Dai Ha 46ac6e4e38 #157: use environment git credential helper
CI / contract (pull_request) Successful in 1m3s
CI / build (pull_request) Successful in 1m37s
2026-08-28 06:02:09 +07:00
Dai Ha 3bfa82839b fleetd#164: an empty or suspiciously fast scrape must fail, never resolve as a success
CI / contract (pull_request) Successful in 1m3s
CI / build (pull_request) Successful in 1m20s
CompletionResolver.resolve() used to hand the caller a successful "" reply whenever a
turn's scrape came back empty (whether the read failed, or genuinely produced nothing),
making a lost turn indistinguishable from a real empty answer. It also had no way to
tell a crashed backend's near-instant BUSY -> DONE transition apart from a genuine
completion.

Add MIN_TURN_NANOS (2s), a named floor below which a completed turn is treated as a
crash signature and failed rather than resolved as a reply. Fail on any empty scrape
(read failure or a clean-but-empty read) instead of resolving with "". Both failures
name the member and carry whatever is on the pane for context.

Thread an injectable LongSupplier clock through CompletionResolver (matching the
SessionManager/MessageService nowNanos pattern) so the floor is testable without a
real sleep.
2026-08-28 06:01:00 +07:00
Dai Ha 65ccf2e4ad CB-633 follow-up: only log allowed N of M when the scrub actually runs
CI / contract (pull_request) Successful in 53s
CI / build (pull_request) Successful in 1m44s
The coverage line was logged before the zsh gate, so a non-zsh
spawn (where nothing is scrubbed — overlayBlockedCredentials is the
fallback instead) printed 'allowed N of M' as if the derived
allow-list scrub had run. Move the log after the gate so it only
fires on the path that actually generates the ZDOTDIR scrub; the
non-zsh fallback keeps logCredentialGap's WARN as its only signal.

Added a test proving no 'allowed N of M' line is emitted on the
non-zsh fallback, through the real HerdrPeerLauncher#spawn path.
2026-08-28 06:00:38 +07:00
ltms 7a3b27f76f #150: report lead readiness from the delivery gate
CI / contract (push) Successful in 1m6s
CI / build (push) Successful in 1m7s
Share one deliverability predicate between Injector and FleetApp, so the status endpoint reports the same answer the injector acts on instead of re-deriving it from one of that predicate's two inputs.

Verified by the lead in an independent worktree: Tests run: 971, Failures: 0, Errors: 0, Skipped: 0, BUILD SUCCESS.
2026-08-28 00:59:44 +02:00
Dai Ha 82e7be564c CB-633 follow-up: union memberCredentials.allow into the derived env allow-list
CI / build (pull_request) Successful in 1m7s
CI / contract (pull_request) Successful in 1m16s
MemberEnvAllowList.derive only ever looked at profile fields, so
memberCredentials.allow: was silently ignored under
policy: allow-list — turning the policy on would have blanked
credentials working members already depended on.

- derive(profiles, configuredAllow) unions memberCredentials.allow
  into the derived set, with SSH_AUTH_SOCK explicitly excluded from
  that union (it stays governed only by sshAuthSock: allow).
- HerdrPeerLauncher threads MemberCredentials.allowSet() into the
  derivation instead of calling the profiles-only overload.
- Added a per-spawn INFO log 'member credentials: allowed N of M'
  (N/M from the daemon's own env, the existing hostEnvNames proxy),
  never logging a blocked name or a value.
2026-08-28 05:56:27 +07:00
Dai Ha c5e24197bf #150: report lead readiness from delivery gate
CI / contract (pull_request) Successful in 46s
CI / build (pull_request) Successful in 1m37s
2026-08-28 05:54:40 +07:00
Dai Ha bcb402b688 #157: convert forge worktree origins to SSH
CI / contract (pull_request) Successful in 1m10s
CI / build (pull_request) Successful in 1m11s
2026-08-28 05:54:15 +07:00
Dai Ha 7c4170ff6d CB-643: join the message-layer evidence to the health monitor
CI / build (push) Successful in 1m8s
CI / contract (push) Successful in 1m9s
CB-640 published the three message-layer facts and CB-641 wired the herdr
and time ones. This joins them, so every HealthSnapshot field now carries
real evidence and the NOT_YET_OBSERVED placeholder is gone. That constant
is what made 8 of the 9 fault states unreachable, GONE and NEVER_READY
included, which is why CB-580's failTarget never fired.

hasOrphanedDelegation is a true snapshot, but it can read true for one
tick during an ordinary race: an async ticket exists before its virtual
thread reaches rendezvous.open, so for that instant nothing is accepted or
queued behind it. decide maps the field straight to DELEGATION_ORPHANED
with no smoothing, so one racy read would log a fault that clears on the
next tick. The monitor now requires two consecutive observations. That
costs one interval on a real orphan and removes the false positive.

Two tests drive real ticks against a genuinely orphaned ticket (an
unanswered fleet_ask that lapsed back to PENDING), not the seam: one tick
reports nothing, two report once, and a single clean tick in between
resets the streak.

Also correct two config comments. paneProbeIntervalSeconds is parsed and
read by nothing, so its "minimum 60" note promised a floor that does not
exist.

970 tests green.
2026-08-27 22:15:39 +07:00
Dai Ha a6095743f0 Merge CB-641: wire herdr and time evidence into the fleet health monitor 2026-08-27 22:06:50 +07:00
Dai Ha 66e776d178 CB-641: wire herdr health evidence
CI / contract (pull_request) Successful in 1m9s
CI / build (pull_request) Successful in 1m11s
2026-08-27 22:03:20 +07:00
Dai Ha 26f64cba45 Merge CB-640: MessageService evidence accessors for the health monitor
CI / build (push) Successful in 1m4s
CI / contract (push) Successful in 1m28s
2026-08-27 22:02:36 +07:00
Dai Ha 51047848f1 CB-642: make the fleets-status redaction global — a non-global sed leaks a second URI on the same line
CI / contract (push) Successful in 41s
CI / build (push) Successful in 1m29s
2026-08-27 22:01:15 +07:00
Dai Ha d8c0b657e8 Merge CB-642: /fleets-status skill — multi-fleet status over the shared LavinMQ 2026-08-27 22:00:49 +07:00
Dai Ha 312c0584ce CB-640: add MessageService message-layer health evidence accessors
CI / contract (pull_request) Successful in 44s
CI / build (pull_request) Successful in 1m11s
hasQueuedDelivery/hasStrandedReply/hasOrphanedDelegation surface three of the
message-layer facts FleetHealthMonitor needs but currently hardcodes to
NOT_YET_OBSERVED. Additive only — no existing public method's signature or
behavior changes.
2026-08-27 21:58:21 +07:00
35 changed files with 2178 additions and 144 deletions
+20 -4
View File
@@ -27,9 +27,14 @@ can make every broker probe look empty.
reaches the report:
```bash
sed -E 's#://[^@]*@#://<redacted>@#'
sed -E 's#://[^@]*@#://<redacted>@#g'
```
**The `g` flag is not optional.** Without it `sed` replaces only the first match on each line, so a
line carrying two URIs leaks the second one. `scripts/redeploy-fleetd.sh --check` prints lines like
that. Checked on 2026-08-27: without `g`, `amqp://u1:p1@h1/mac and http://u2:p2@h2:15672/api`
redacts the first pair and prints `u2:p2` in the clear.
Keep `pipefail` on when applying that filter. Otherwise the filter can hide a failed probe. Apply
the same no-print rule to the management password below, even though it is not in an AMQP URI.
@@ -44,7 +49,7 @@ Do not copy those checks into new shell code. The script reads `LAVINMQ_URI`, so
```bash
set -o pipefail
scripts/redeploy-fleetd.sh --check 2>&1 \
| sed -E 's#://[^@]*@#://<redacted>@#'
| sed -E 's#://[^@]*@#://<redacted>@#g'
git rev-parse HEAD
```
@@ -108,8 +113,19 @@ PY
**What this tier cannot see:** it proves facts only about the Mac daemon at `127.0.0.1:8765`.
It cannot show the fleet01 daemon, broker queue depth, or broker consumers. The fleet01 REST service
at `10.10.20.13:8765` is not reachable from the Mac, and SSH as `dai.ha@10.10.20.13` is denied.
Say this in the report rather than omitting fleet01.
at `10.10.20.13:8765` is not reachable from the Mac. Say this in the report rather than omitting
fleet01.
**But fleet01 IS reachable over SSH — checked 2026-08-28.** An older version of this line said SSH
was denied. That is true only for the user `dai.ha`. The host alias `fleet01` maps to user `ltms`,
and `ssh fleet01` works with key auth:
```bash
ssh -o BatchMode=yes -o ConnectTimeout=6 fleet01 'echo $(id -un)@$(hostname)'
```
So fleet01's daemon PID, uptime, jar and `/healthz` **can** be reported — over SSH, not over REST.
Do that rather than writing `not reachable`. `ltms` also has passwordless sudo there.
## 3. Tier 2 — the shared broker (run when management access exists)
+11 -9
View File
@@ -93,21 +93,20 @@ bind:
# intervalSeconds → how often a tick runs (default 30). ENFORCED floor of 15: the code computes
# Math.max(15, intervalSeconds), so a lower value is silently raised, not
# rejected.
# workingSuspectAfterSeconds, paneProbeIntervalSeconds → accepted and parsed, but NOT YET READ by
# anything — the dormant monitor only consumes intervalSeconds today (CB-573
# shipped ahead of the evidence publishers these two knobs are for). Setting
# them changes nothing right now, and no minimum is enforced on either, because
# nothing reads them to enforce one. They exist so a later build can start
# honouring them without another config-shape change.
# workingSuspectAfterSeconds → age before a BUSY member is suspected of a stall (default 600).
# ENFORCED floor of 300: a lower value is silently raised.
# paneProbeIntervalSeconds → accepted and parsed, but NOT YET READ by anything. Setting it changes
# nothing right now. It exists so a later build can start honouring it without
# another config-shape change.
# notifications.mode → "webhook" flips what fleet_list REPORTS (healthCoverage: "full" instead
# of "detection-only") — it does NOT make fleetd send any webhook call; no
# delivery mechanism is implemented yet. Any other value, or omitting the
# block, reports "detection-only".
# health:
# enabled: true
# intervalSeconds: 30
# workingSuspectAfterSeconds: 600
# paneProbeIntervalSeconds: 60
# intervalSeconds: 30 # floor 15
# workingSuspectAfterSeconds: 600 # floor 300 — how long BUSY with no activity means STALL_SUSPECTED
# paneProbeIntervalSeconds: 60 # parsed, but nothing reads it yet — changing it changes nothing
# notifications:
# mode: disabled
@@ -115,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
+31 -17
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.
@@ -380,9 +387,10 @@ public final class Fleetd {
sessions.onTurnFailed(target);
}
};
Injector injector = new Injector(agents, turnListener, deliverableTo(presence, leads),
Predicate<String> deliverable = deliverableTo(presence, leads);
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
@@ -419,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.
@@ -429,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);
@@ -438,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.
@@ -449,7 +457,8 @@ public final class Fleetd {
// CB-580: a member found GONE/NEVER_READY must fail whatever ticket is waiting on it,
// through the same idempotent target-wide operation CB-516 already uses on release.
healthMonitor = new FleetHealthMonitor(agents, sessions::roster, messages, healthScheduler,
System::nanoTime, cfg.health().intervalOrDefault(), messages::abandon);
System::nanoTime, cfg.health().intervalOrDefault(),
cfg.health().workingSuspectAfterOrDefault(), messages::abandon);
String coverage = FleetHealthMonitor.coverage(true,
cfg.health().notifications() != null && cfg.health().notifications().configured());
if ("detection-only".equals(coverage)) {
@@ -489,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.
@@ -536,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;
@@ -587,11 +599,13 @@ 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(),
callers, metrics).build();
// 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 {}",
cfg.bind().host(), cfg.bind().port(), 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);
}
/**
@@ -642,6 +644,9 @@ public record FleetConfig(
Integer paneProbeIntervalSeconds, Notifications notifications) {
public boolean isEnabled() { return Boolean.TRUE.equals(enabled); }
public int intervalOrDefault() { return Math.max(15, intervalSeconds == null ? 30 : intervalSeconds); }
public int workingSuspectAfterOrDefault() {
return Math.max(300, workingSuspectAfterSeconds == null ? 600 : workingSuspectAfterSeconds);
}
public record Notifications(String mode) {
public boolean configured() { return "webhook".equalsIgnoreCase(mode); }
}
@@ -1317,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");
@@ -1938,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);
}
@@ -3,6 +3,7 @@ package dev.ltms.fleet.health;
import dev.ltms.fleet.herdr.Agent;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.HerdrException;
import dev.ltms.fleet.msg.MessageService;
import dev.ltms.fleet.session.MemberSession;
import org.slf4j.Logger;
@@ -25,6 +26,8 @@ public final class FleetHealthMonitor {
/** Bounded attempts to run {@link #failTarget} for one transition. Never retried tick-to-tick (CB-580). */
static final int MAX_FAIL_TARGET_ATTEMPTS = 3;
// CB-641: Match the injector's 60s readiness gate so health allows a full first boot.
static final long READINESS_GRACE_NANOS = TimeUnit.SECONDS.toNanos(60);
private final AgentControl agents;
private final Supplier<List<MemberSession>> roster;
@@ -32,12 +35,30 @@ public final class FleetHealthMonitor {
private final ScheduledExecutorService scheduler;
private final LongSupplier clock;
private final long intervalSeconds;
private final long workingSuspectAfterNanos;
private final BiConsumer<String, String> failTarget;
private final Map<String, HealthPrior> priors = new HashMap<>();
private final Map<String, HealthState> states = new HashMap<>();
/**
* CB-643: consecutive ticks on which a target looked like an orphaned delegation. The fact
* {@link MessageService#hasOrphanedDelegation} reports is a true snapshot, but it can read true
* for one tick during an ordinary race — an async ticket exists before its virtual thread has
* reached {@code rendezvous.open()}, so for that instant nothing is accepted or queued behind
* it. {@code decide} maps the field straight to {@code DELEGATION_ORPHANED} with no cross-tick
* smoothing of its own, so a single racy read would log a fault that clears on the next tick.
* Requiring two consecutive observations costs one interval of latency on a real orphan and
* removes that false positive entirely.
*/
private final Map<String, Integer> orphanStreaks = new HashMap<>();
// These facts need the evidence publishers introduced by later M4 units. They are not negatives.
private static final boolean NOT_YET_OBSERVED = false;
/** How many consecutive ticks a target must look orphaned before health reports it (CB-643). */
static final int ORPHAN_CONFIRM_TICKS = 2;
// CB-643: every HealthSnapshot field now carries real evidence. The NOT_YET_OBSERVED placeholder
// that stood in for 7 of the 12 is gone, and with it the reason 8 of the 9 fault states were
// unreachable — GONE and NEVER_READY included, which is what kept CB-580's failTarget from ever
// firing. Do not reintroduce a constant here: a field with no publisher is a dead state, and the
// tests pass either way, so nothing else will tell you.
/**
* @param failTarget CB-568's idempotent target-wide failure operation (e.g. {@code messages::abandon}),
@@ -48,13 +69,14 @@ public final class FleetHealthMonitor {
*/
public FleetHealthMonitor(AgentControl agents, Supplier<List<MemberSession>> roster, MessageService messages,
ScheduledExecutorService scheduler, LongSupplier clock, long intervalSeconds,
BiConsumer<String, String> failTarget) {
long workingSuspectAfterSeconds, BiConsumer<String, String> failTarget) {
this.agents = agents;
this.roster = roster;
this.messages = messages;
this.scheduler = scheduler;
this.clock = clock;
this.intervalSeconds = intervalSeconds;
this.workingSuspectAfterNanos = TimeUnit.SECONDS.toNanos(workingSuspectAfterSeconds);
this.failTarget = Objects.requireNonNull(failTarget, "failTarget");
}
@@ -69,28 +91,50 @@ public final class FleetHealthMonitor {
// Package-private so tests can run one tick without waiting.
void tick() {
try {
List<Agent> agentsNow = agents.list(); // Exactly one list call for this complete observation.
List<MemberSession> rosterNow = roster.get(); // One in-memory roster snapshot for this tick.
List<Agent> agentsNow;
boolean controlLinkDown = false;
try {
agentsNow = agents.list(); // Exactly one list call for this complete observation.
} catch (HerdrException error) {
agentsNow = List.of();
controlLinkDown = true;
log.warn("fleet health control link unavailable; classifying roster", error);
}
Map<String, Agent> live = new HashMap<>();
for (Agent agent : agentsNow) live.put(agent.terminalId(), agent);
HashSet<String> current = new HashSet<>();
long nowNanos = clock.getAsLong();
for (MemberSession session : rosterNow) {
current.add(session.terminalId());
Agent agent = live.get(session.terminalId());
AgentStatus status = agent == null ? AgentStatus.UNKNOWN : agent.status();
boolean accepted = messages.hasAcceptedDelivery(session.terminalId());
HealthSnapshot snapshot = new HealthSnapshot(session.state(), status, accepted, NOT_YET_OBSERVED,
messages.hasInboxMessage(session.terminalId()), agent != null, NOT_YET_OBSERVED,
NOT_YET_OBSERVED, NOT_YET_OBSERVED, NOT_YET_OBSERVED, NOT_YET_OBSERVED, NOT_YET_OBSERVED);
boolean present = agent != null;
boolean targetNotFound = !controlLinkDown && !present
&& session.state() != MemberSession.State.SPAWNING;
boolean readinessGraceElapsed = nowNanos - session.spawnedAtNanos() >= READINESS_GRACE_NANOS;
boolean stalled = session.state() == MemberSession.State.BUSY
&& nowNanos - session.lastActivityAtNanos() >= workingSuspectAfterNanos;
// CB-643: the three message-layer facts CB-640 published. Read them here rather than
// leaving them false — that constant is what made 8 of the 9 fault states dead.
boolean queuedDelivery = messages.hasQueuedDelivery(session.terminalId());
boolean replyStranded = messages.hasStrandedReply(session.terminalId());
boolean orphanedDelegation = confirmOrphan(session.terminalId(),
messages.hasOrphanedDelegation(session.terminalId()));
HealthSnapshot snapshot = new HealthSnapshot(session.state(), status, accepted, queuedDelivery,
messages.hasInboxMessage(session.terminalId()), present, targetNotFound, controlLinkDown,
readinessGraceElapsed, orphanedDelegation, replyStranded, stalled);
HealthDecision decision = decide(snapshot, priors.getOrDefault(session.terminalId(), HealthPrior.NONE),
clock.getAsLong());
nowNanos);
priors.put(session.terminalId(), decision.prior());
reportTransition(session.terminalId(), decision.state());
}
priors.keySet().retainAll(current);
states.keySet().retainAll(current);
orphanStreaks.keySet().retainAll(current);
} catch (Throwable error) {
// A list failure is health evidence, and must never kill the monitor's only scheduler task.
// Any unclassified collection failure must never kill the monitor's only scheduler task.
log.warn("fleet health collection failed; will retry next tick", error);
} finally {
if (!scheduler.isShutdown()) {
@@ -99,6 +143,20 @@ public final class FleetHealthMonitor {
}
}
/**
* Debounce {@link MessageService#hasOrphanedDelegation} across ticks (CB-643). Returns true only
* once {@code observed} has held for {@link #ORPHAN_CONFIRM_TICKS} consecutive ticks; a single
* false reading resets the streak, so a transient race never reaches the classifier.
*/
private boolean confirmOrphan(String target, boolean observed) {
if (!observed) {
orphanStreaks.remove(target);
return false;
}
int streak = orphanStreaks.merge(target, 1, Integer::sum);
return streak >= ORPHAN_CONFIRM_TICKS;
}
void reportTransition(String target, HealthState next) {
HealthState previous = states.put(target, next);
if (previous == next) return;
@@ -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");
@@ -6,12 +6,14 @@ import dev.ltms.fleet.msg.TurnToken;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.time.Duration;
import java.util.List;
import java.util.Objects;
import java.util.Set;
import java.util.TreeSet;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.function.LongSupplier;
import java.util.regex.Pattern;
/**
@@ -61,6 +63,17 @@ public final class CompletionResolver implements TurnListener {
/** Cap the scraped tail so a long transcript can't return an unbounded blob. */
static final int MAX_SCRAPE_CHARS = 4000;
/**
* fleetd#164: the floor below which a {@code BUSY -> DONE} transition cannot be real work. A
* backend that rejects a turn outright (e.g. an HTTP 400 from the model, before the worker read
* a single file or produced a token) drives the exact same confirmed {@code working -> idle}
* transition a genuine completion does — just in about a second instead of the many seconds a
* real turn costs. {@link #onTurnComplete} cannot tell those two cases apart from the transition
* alone, so a turn that settles inside this floor is treated as a crash signature and resolved
* as a failure, never as a (possibly empty) success.
*/
public static final long MIN_TURN_NANOS = Duration.ofSeconds(2).toNanos();
private static final String CLIPPED_PANE_TAIL_MARKER =
"[Pane tail clipped: member did not call fleet_reply.]";
@@ -68,6 +81,7 @@ public final class CompletionResolver implements TurnListener {
private final Rendezvous rendezvous;
private final ExhaustedPatternLookup exhaustedPatterns;
private final ExhaustionSink exhaustionSink;
private final LongSupplier nowNanos;
/**
* Per-target record of the turn currently in flight: the exact {@link Rendezvous} waiter its
@@ -79,8 +93,21 @@ public final class CompletionResolver implements TurnListener {
* reference: a completion scrape equal to it means the worker produced no new output (the previous
* turn's wind-down sampled as this boundary), so it is suppressed. Overwritten on each delivery;
* cleared when the turn resolves. Package-private so tests can capture and replay a specific turn.
*
* <p>{@code deliveredAtNanos} (fleetd#164) is the {@link #nowNanos} reading taken at delivery —
* the other half of the {@link #MIN_TURN_NANOS} floor check, compared against a fresh reading at
* resolution time.
*/
record InFlight(CompletableFuture<Rendezvous.Resolution> waiter, String baseline) {
record InFlight(CompletableFuture<Rendezvous.Resolution> waiter, String baseline, long deliveredAtNanos) {
/**
* Convenience for tests exercising scrape/suppression logic that don't care about turn
* timing: back-dates the delivery far enough that {@link #MIN_TURN_NANOS} can never fire.
* Not used by production code — {@link #captureBaseline} always records a real reading.
*/
InFlight(CompletableFuture<Rendezvous.Resolution> waiter, String baseline) {
this(waiter, baseline, Long.MIN_VALUE / 2);
}
}
private final ConcurrentHashMap<String, InFlight> inFlight = new ConcurrentHashMap<>();
@@ -97,10 +124,25 @@ public final class CompletionResolver implements TurnListener {
*/
public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns,
ExhaustionSink exhaustionSink) {
this(agents, rendezvous, exhaustedPatterns, exhaustionSink, System::nanoTime);
}
/**
* Test constructor with an injectable clock (fleetd#164), matching the {@code LongSupplier}
* pattern {@link dev.ltms.fleet.session.SessionManager} and {@link dev.ltms.fleet.msg.MessageService}
* already use: lets a test place a turn's delivery and its resolution at an exact, controllable
* distance apart around the {@link #MIN_TURN_NANOS} floor, without a real sleep. Public (rather
* than package-private like those two) because callers that wire a full {@code MessageService}
* fixture — e.g. {@code MessageServiceTest} — construct this resolver directly from another
* package.
*/
public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns,
ExhaustionSink exhaustionSink, LongSupplier nowNanos) {
this.agents = agents;
this.rendezvous = rendezvous;
this.exhaustedPatterns = Objects.requireNonNull(exhaustedPatterns, "exhaustedPatterns");
this.exhaustionSink = Objects.requireNonNull(exhaustionSink, "exhaustionSink");
this.nowNanos = Objects.requireNonNull(nowNanos, "nowNanos");
}
@Override
@@ -130,7 +172,7 @@ public final class CompletionResolver implements TurnListener {
baseline = null; // fail open: no baseline ⇒ no suppression
log.debug("delivery baseline for {} failed: {}", target, e.getMessage());
}
inFlight.put(target, new InFlight(waiter, baseline));
inFlight.put(target, new InFlight(waiter, baseline, nowNanos.getAsLong()));
}
/** The turn currently baselined for {@code target}, or {@code null} — a test hook for the captureBaseline path. */
@@ -176,6 +218,15 @@ public final class CompletionResolver implements TurnListener {
inFlight.remove(target, turn);
return;
}
// fleetd#164: a BUSY -> DONE transition inside the floor cannot be real work — it's a crash
// signature (e.g. a backend HTTP 400 before the worker did anything), not a fast answer. Fail
// it before spending a scrape on the ordinary path; the reason still carries whatever is on
// screen, since that is usually the backend's own error.
long elapsedNanos = nowNanos.getAsLong() - turn.deliveredAtNanos();
if (elapsedNanos < MIN_TURN_NANOS) {
fail(target, turn, tooFastReason(target, elapsedNanos));
return;
}
String tail;
String assistantBlock = null;
int originalLength = 0;
@@ -187,20 +238,25 @@ public final class CompletionResolver implements TurnListener {
clipped = originalLength > MAX_SCRAPE_CHARS;
tail = clip(assistantBlock);
} catch (RuntimeException e) {
// The worker finished but we couldn't read its screen — still resolve the send so the
// caller unblocks; an empty tail beats hanging until the caller's timeout.
log.warn("completion scrape for {} failed; resolving with an empty tail: {}",
target, e.getMessage());
log.warn("completion scrape for {} failed: {}", target, e.getMessage());
tail = "";
scrapeFailed = true;
}
// fleetd#164: a scrape nobody could read, and a scrape that read cleanly but produced nothing,
// both used to resolve the send as a SUCCESS carrying "" — indistinguishable from a worker that
// genuinely finished with nothing to say. That is the defect: fail loudly instead, naming the
// member, so a caller (including a lead deciding whether to delegate again) can tell a lost
// turn from a real empty answer.
if (scrapeFailed || tail.isEmpty()) {
fail(target, turn, emptyScrapeReason(target, scrapeFailed));
return;
}
// Misattribution guard (CB-115): if the scrape is byte-identical to the pane content at
// delivery, this turn produced no new output — the boundary belongs to the previous turn's
// wind-down (common on rapid back-to-back sends). Suppress rather than resolve the send with
// a stale answer; the real fleet_reply (or a later genuine completion) resolves it instead.
// A scrape that failed to read is exempt — an empty tail there is "couldn't see", not "no change".
String baseline = turn.baseline();
if (!scrapeFailed && baseline != null && baseline.equals(tail)) {
if (baseline != null && baseline.equals(tail)) {
log.debug("suppressing misattributed completion for {} (no output change since delivery)",
target);
return; // keep the in-flight record: a later genuine completion still needs it
@@ -208,21 +264,19 @@ public final class CompletionResolver implements TurnListener {
// CB-578 stage A: a turn that ended with no fleet_reply AND whose scrape matches the
// backend's configured usage-limit pattern is a refusal, not an answer. Classify it as
// BACKEND_EXHAUSTED rather than handing the caller a scrape that reads like a real reply.
if (!scrapeFailed) {
Pattern exhausted = exhaustedPatterns.patternFor(target);
String matchedLine = exhausted == null ? null : firstMatchingLine(assistantBlock, exhausted);
if (matchedLine != null) {
String reason = "backend exhausted (usage limit): " + matchedLine;
if (rendezvous.resolveExhausted(waiter, reason)) {
inFlight.remove(target, turn);
log.warn("completion for {} classified BACKEND_EXHAUSTED (no fleet_reply; scrape "
+ "matched the profile's exhausted pattern): {}", target, reason);
// CB-578 stage B: only on the resolution that actually won the race — a late
// duplicate must never quarantine a credential twice for one refusal.
exhaustionSink.onExhausted(target, reason);
}
return;
Pattern exhausted = exhaustedPatterns.patternFor(target);
String matchedLine = exhausted == null ? null : firstMatchingLine(assistantBlock, exhausted);
if (matchedLine != null) {
String reason = "backend exhausted (usage limit): " + matchedLine;
if (rendezvous.resolveExhausted(waiter, reason)) {
inFlight.remove(target, turn);
log.warn("completion for {} classified BACKEND_EXHAUSTED (no fleet_reply; scrape "
+ "matched the profile's exhausted pattern): {}", target, reason);
// CB-578 stage B: only on the resolution that actually won the race — a late
// duplicate must never quarantine a credential twice for one refusal.
exhaustionSink.onExhausted(target, reason);
}
return;
}
String completion = clipped ? tail + "\n" + CLIPPED_PANE_TAIL_MARKER : tail;
if (rendezvous.resolveCompletion(waiter, completion)) {
@@ -272,6 +326,36 @@ public final class CompletionResolver implements TurnListener {
}
}
/**
* fleetd#164: the failure reason for a turn that settled inside {@link #MIN_TURN_NANOS} — names
* the member and both timings, and appends whatever the pane shows (usually the backend's own
* error) so the caller sees the cause, not just "it failed".
*/
private String tooFastReason(String target, long elapsedNanos) {
String scrape;
try {
scrape = clip(agents.read(target, SCRAPE_SOURCE));
} catch (RuntimeException e) {
scrape = "";
}
String reason = String.format(
"member %s went BUSY -> DONE in %dms (floor %dms) — too fast to be real work, most "
+ "likely a backend error before any work started",
target, elapsedNanos / 1_000_000, MIN_TURN_NANOS / 1_000_000);
return scrape.isBlank() ? reason : reason + ": " + scrape;
}
/**
* fleetd#164: the failure reason for a scrape that produced zero characters — names the member
* and says plainly that the turn produced nothing, so a caller (a lead deciding whether to
* delegate again included) never mistakes a lost turn for a genuinely empty reply.
*/
private static String emptyScrapeReason(String target, boolean scrapeFailed) {
return "member " + target + " turn completed with an empty scrape (0 chars) — "
+ (scrapeFailed ? "its pane could not be read; " : "")
+ "treating as a lost turn, not a real answer";
}
/**
* The first line of {@code text} matching {@code pattern}, stripped — the CB-578 stage A
* evidence carried in a {@code BACKEND_EXHAUSTED} reason so the operator sees the real refusal
@@ -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;
@@ -1033,34 +1033,40 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
* pass it through {@code tab.create}/{@code pane.split}. Returns the directory for teardown
* registration, or {@code null} when the policy does not apply.
*
* <p>The allow-list handed to the generator is the derived profile set ({@link
* MemberEnvAllowList#derive}) UNIONed with the exact keys of THIS launch's env map — names the
* daemon itself injects must survive its own control. {@code SSH_AUTH_SOCK} is added ONLY when
* the config explicitly allows it; by default it is absent, so the scrub blanks it like any
* other non-derived name.
* <p>The allow-list handed to the generator is the derived profile set UNIONed with the
* operator's own {@code memberCredentials.allow:} names ({@link MemberEnvAllowList#derive(
* Collection, Set)} — CB-633 follow-up) and with the exact keys of THIS launch's env map —
* names the daemon itself injects must survive its own control. {@code SSH_AUTH_SOCK} is added
* ONLY when the config explicitly allows it, EVEN IF the operator also listed it under
* {@code allow:}; by default it is absent, so the scrub blanks it like any other non-derived
* name. It stays a one-off decision because it is a live handle to the operator's ssh-agent, not
* a value — a member holding it can sign with every key the agent holds, so letting it ride in
* on the generic {@code allow:} list would hand that out for an unrelated reason.
*/
private Path applyEnvironmentAllowListPolicy(FleetConfig.Profile cfg, Launch launch) {
FleetConfig.MemberCredentials creds = memberCredentials == null ? null : memberCredentials.get();
if (creds == null || !creds.isAllowList()) {
return null;
}
Set<String> allowed = derivedAllowedNames(creds, launch);
String loginShell = resolveEnv("SHELL");
boolean zsh = loginShell != null && (loginShell.endsWith("/zsh") || loginShell.equals("zsh"));
if (!zsh) {
// A non-zsh login shell ignores ZDOTDIR entirely: NO scrub would run, so pretending
// otherwise would be worse than saying so. Warn loudly and fall back to the CB-596
// sentinel overlay over the enumerated known: names — weaker (a sourced file can undo
// it), but strictly better than nothing.
// it), but strictly better than nothing. Deliberately no "allowed N of M" line here: the
// scrub this count describes does not run on this path, so printing it would tell an
// operator that a fraction of names were blocked when the real number blocked is zero.
// logCredentialGap's WARN (below) is the only signal for this path.
warnNonZsh(loginShell);
overlayBlockedCredentials(launch.env(), creds);
logCredentialGap(creds);
return null;
}
Set<String> allowed = new java.util.TreeSet<>(MemberEnvAllowList.derive(profiles.values()));
if (creds.sshAuthSockAllowed()) {
allowed.add(SSH_AUTH_SOCK);
} // blocked by default: absent from the set ⇒ blanked by the scrub like any other name
allowed.addAll(launch.env().keySet());
// 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).
logAllowListCoverage(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 "
@@ -1069,8 +1075,42 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
return dir;
}
/**
* The full kept-name set for this spawn: the profile-derived names, unioned with {@code
* memberCredentials.allow:} (CB-633 follow-up — previously ignored by this whole policy), the
* ssh-agent handle when explicitly allowed, and the exact keys of THIS launch's own env map.
*/
private Set<String> derivedAllowedNames(FleetConfig.MemberCredentials creds, Launch launch) {
Set<String> allowed = new java.util.TreeSet<>(
MemberEnvAllowList.derive(profiles.values(), creds.allowSet()));
if (creds.sshAuthSockAllowed()) {
allowed.add(SSH_AUTH_SOCK);
} // blocked by default: absent from the set ⇒ blanked by the scrub like any other name
allowed.addAll(launch.env().keySet());
return allowed;
}
/**
* CB-633 follow-up: one INFO line per allow-list spawn WHOSE SCRUB ACTUALLY RUNS, so an operator
* can read a single log line and know the scrub ran and how much of the visible environment it
* will keep. Callable ONLY from the zsh branch of {@link #applyEnvironmentAllowListPolicy}, after
* the shell gate — logging it before that gate (or on the non-zsh fallback, where nothing is
* scrubbed) would tell an operator a fraction of names were blocked when the real number blocked
* is zero, which is worse than not logging at all. {@code M} is {@link #hostEnvNames}' size (the
* daemon's own environment — see that field's javadoc for why it stands in for the pane's, which
* the daemon has no channel to inspect at spawn time) and {@code N} is how many of those names
* survive {@code allowed} (including the {@code LC_*} prefix rule). Neither number is a constant:
* both come from the actual derived set and the actual environment this spawn sees. Never logs a
* variable NAME or VALUE — only the counts.
*/
private void logAllowListCoverage(Set<String> allowed) {
Set<String> hostNames = hostEnvNames.get();
long kept = hostNames.stream().filter(name -> MemberEnvAllowList.keeps(allowed, name)).count();
log.info("member credentials: allowed {} of {}", kept, hostNames.size());
}
/** The operator ssh-agent handle — kept ONLY by explicit config decision, never by default. */
private static final String SSH_AUTH_SOCK = "SSH_AUTH_SOCK";
private static final String SSH_AUTH_SOCK = MemberEnvAllowList.SSH_AUTH_SOCK;
/**
* CB-633: a non-zsh login shell means the allow-list control CANNOT run — say so once per
@@ -35,9 +35,27 @@ import java.util.TreeSet;
* <p>{@code SSH_AUTH_SOCK} is deliberately NOT here. It is a handle to the operator's ssh-agent — a
* member holding it can sign with the operator's keys — so keeping it is a config decision
* ({@code memberCredentials.sshAuthSock: allow}), not a derivation default.
*
* <p><b>CB-633 follow-up:</b> the union also includes {@code memberCredentials.allow:} — the
* operator's own explicit list. Before this, {@code policy: allow-list} silently ignored every name
* an operator wrote under {@code allow:} unless a profile happened to carry it too, which meant
* turning the policy on could blank credentials working members already depended on. {@code
* SSH_AUTH_SOCK} is the one exception: even when the operator lists it under {@code allow:}, it is
* excluded here and added back ONLY by the caller when {@code sshAuthSock: allow} is explicitly set
* (see {@link #SSH_AUTH_SOCK}'s javadoc) — it is a live handle to the operator's own ssh-agent, not
* a value, so treating it like any other allow-listed name would hand a member every key the
* operator's agent holds the moment they typed the name under {@code allow:} for an unrelated
* reason.
*/
public final class MemberEnvAllowList {
/**
* The operator's ssh-agent socket path. Deliberately excluded from {@link #derive}'s union of
* {@code memberCredentials.allow:} — see the class javadoc's CB-633 follow-up note. Governed
* ONLY by {@code memberCredentials.sshAuthSock}, never by appearing in {@code allow:}.
*/
public static final String SSH_AUTH_SOCK = "SSH_AUTH_SOCK";
/**
* Names that are not credentials and that a login shell or agent binary genuinely needs.
*
@@ -73,9 +91,21 @@ public final class MemberEnvAllowList {
/**
* Derive the allowed NAME set from the given profiles plus {@link #INFRASTRUCTURE_PASSTHROUGH}.
* Deterministic (sorted) so generated scrub files are diffable run-to-run.
* Equivalent to {@link #derive(Collection, Set)} with no operator-configured names — kept for
* callers (and existing tests) that only care about the profile-derived half.
*/
public static Set<String> derive(Collection<FleetConfig.Profile> profiles) {
return derive(profiles, Set.of());
}
/**
* Derive the allowed NAME set: the profile-derived union above, PLUS {@code configuredAllow} —
* the operator's own {@code memberCredentials.allow:} list (CB-633 follow-up). {@code
* SSH_AUTH_SOCK} is dropped from {@code configuredAllow} even if the operator listed it there;
* see the class javadoc for why. Deterministic (sorted) so generated scrub files are diffable
* run-to-run.
*/
public static Set<String> derive(Collection<FleetConfig.Profile> profiles, Set<String> configuredAllow) {
Set<String> derived = new TreeSet<>(INFRASTRUCTURE_PASSTHROUGH);
if (profiles != null) {
for (FleetConfig.Profile p : profiles) {
@@ -87,6 +117,13 @@ public final class MemberEnvAllowList {
}
}
}
if (configuredAllow != null) {
for (String name : configuredAllow) {
if (name != null && !name.isBlank() && !SSH_AUTH_SOCK.equals(name)) {
derived.add(name);
}
}
}
return Set.copyOf(derived);
}
@@ -29,8 +29,9 @@ import java.util.concurrent.TimeoutException;
*
* <p><strong>Mapping — consume-and-hold with deferred manual ack.</strong> Each target has a durable
* queue {@code agent.<target>.inbox}. The gateway that owns the target starts a manual-ack consumer
* ({@link #own}) that pulls persistent messages off that queue into an in-memory <em>held</em> map
* (keyed by {@code msgId}) but does <em>not</em> ack them. {@link #peek} returns that snapshot;
* ({@link #own}) that pulls persistent messages, up to its prefetch window, off that queue into an
* in-memory <em>held</em> map (keyed by {@code msgId}) but does <em>not</em> ack them.
* {@link #peek} returns that snapshot;
* {@link #ack} acks the broker delivery-tag and drops the entry. Because messages stay unacked until
* the owning gateway actually drains them, a crash (or a {@code java -jar} bounce) before caller-ack
* leaves them on the broker — it redelivers on reconnect. That is the durability the in-memory
@@ -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;
@@ -189,6 +190,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;
@@ -204,6 +206,25 @@ public final class MessageService {
new ConcurrentHashMap<>();
/** Async tickets paused on a specific {@code fleet_ask} turn. */
private final ConcurrentHashMap<String, Task> asyncTasksByTurn = new ConcurrentHashMap<>();
/**
* Targets whose most recent {@code fleet_reply} arrived with no send awaiting it (CB-640) —
* {@link Rendezvous#resolve} returned {@code false} and the reply was queued into the inbox
* instead (see {@link #reply}). The reply itself is not lost (it sits in the inbox for a
* later drain), but the stranding is a fact the health layer needs to see. Bounded by the
* target's own lifecycle rather than a TTL: an entry is cleared the next time this target's
* delivery is accepted ({@link #send}) or the target is torn down ({@link #abandon}), so the
* map holds at most one entry per session with an unresolved stranding right now.
*/
private final ConcurrentHashMap<String, Boolean> strandedReplies = new ConcurrentHashMap<>();
/**
* Targets whose last send timed out with {@link Outcome#TIMED_OUT_QUEUED} (CB-640) — the
* message never reached the {@link Injector} delivery window before the caller's deadline, so
* it is still sitting in the injector's own per-target queue. Set where {@link #send} already
* computes {@code wasDelivered} for that outcome; no new queue is kept here, only the fact.
* Cleared the same way as {@link #strandedReplies}: the next accepted delivery for the target
* ({@link #send} opening a fresh waiter) or a teardown ({@link #abandon}).
*/
private final ConcurrentHashMap<String, Boolean> queuedDeliveries = new ConcurrentHashMap<>();
private final AtomicLong ticketSeq = new AtomicLong();
private final ExecutorService asyncExecutor = Executors.newThreadPerTaskExecutor(
Thread.ofVirtual().name("bridge-async-", 0).factory());
@@ -238,6 +259,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;
@@ -246,6 +268,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);
@@ -258,7 +294,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. */
@@ -271,6 +307,55 @@ public final class MessageService {
return !inbox.peek(target).isEmpty();
}
/**
* Read-only delegation fact for fleet views (CB-640): {@code target}'s last send timed out
* before the {@link Injector} ever delivered it — the caller saw
* {@link Outcome#TIMED_OUT_QUEUED} (see the {@code TimeoutException} branch of {@link #send}),
* and the message is still sitting in the injector's per-target queue waiting for the worker
* to go idle. Distinct from {@link Outcome#TIMED_OUT_WORKING}, where delivery already happened
* and only the reply is outstanding. Cleared the next time this target's delivery is accepted
* or the target is abandoned — see {@link #queuedDeliveries}.
*/
public boolean hasQueuedDelivery(String target) {
return target != null && queuedDeliveries.containsKey(target);
}
/**
* Read-only delegation fact for fleet views (CB-640): {@code target}'s last {@code fleet_reply}
* arrived while no send was waiting for it, so {@link Rendezvous#resolve} returned
* {@code false} and {@link #reply} fell back to queueing it in the inbox (see the CB-307
* javadoc there and {@code docs/CB-307-Reliable-Delivery.md} §1). Cleared the next time this
* target's delivery is accepted or the target is abandoned — see {@link #strandedReplies}.
*/
public boolean hasStrandedReply(String target) {
return target != null && strandedReplies.containsKey(target);
}
/**
* Read-only delegation fact for fleet views (CB-640): an async ticket is still
* {@link Phase#PENDING} against {@code target}, yet nothing is actually in flight for it — no
* open rendezvous waiter ({@link #hasAcceptedDelivery}) and no message still sitting in the
* injector's queue ({@link #hasQueuedDelivery}). A healthy PENDING ticket can briefly look this
* way while its virtual thread has not yet been scheduled or is blocked on the session lock
* behind another send to the same target, so this is a snapshot fact for the health classifier
* to weigh across ticks, not proof on its own that the ticket is stuck. It also genuinely
* persists — not just as a passing race — once an async {@code fleet_ask} lapses unanswered:
* {@link #ask} clears the ticket's question and returns it to {@code PENDING}, but {@link #send}
* already closed the forward waiter the instant the question surfaced, so the target has
* neither an accepted nor a queued delivery left to show for it.
*/
public boolean hasOrphanedDelegation(String target) {
if (target == null || hasAcceptedDelivery(target) || hasQueuedDelivery(target)) {
return false;
}
for (Task task : tasks.values()) {
if (target.equals(task.target) && task.question == null && !task.future.isDone()) {
return true;
}
}
return false;
}
/**
* 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
@@ -288,6 +373,9 @@ public final class MessageService {
return true; // a live send took it — unchanged fast path
}
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.
strandedReplies.put(session, Boolean.TRUE);
// A rising inbox share is the signal CB-307 exists to make visible: the worker finished but
// nobody was waiting, so delivery now depends on the push loop and a drain.
count(FleetMetrics.REPLIES, "path", "inbox");
@@ -341,6 +429,9 @@ public final class MessageService {
* @return true if a live waiter was failed
*/
public boolean abandon(String target, String reason) {
// 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;
@@ -446,6 +537,11 @@ public final class MessageService {
// enqueue failure is safely closed by the finally below: nothing is left queued, and the
// failed send leaves no stale waiter behind.
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(target);
// CB-640: this send now owns target's delivery, so any earlier stranded-reply or
// still-queued fact no longer describes the live state — clear both rather than let
// them outlive the send that supersedes them.
strandedReplies.remove(target);
queuedDeliveries.remove(target);
try {
if (task != null) {
asyncTasksByWaiter.put(reply, task);
@@ -464,6 +560,11 @@ public final class MessageService {
} catch (TimeoutException e) {
boolean wasDelivered = delivered.isDone() && !delivered.isCompletedExceptionally();
log.debug("send to {} timed out (delivered={})", target, wasDelivered);
if (!wasDelivered) {
// CB-640: still sitting in the injector's queue, waiting for the member to
// go idle — record the fact for fleet health (see queuedDeliveries).
queuedDeliveries.put(target, Boolean.TRUE);
}
return recorded(new Reply(
wasDelivered ? Outcome.TIMED_OUT_WORKING : Outcome.TIMED_OUT_QUEUED, null));
} catch (ExecutionException e) {
@@ -704,7 +805,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";
}
@@ -31,6 +31,7 @@ import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.function.Function;
import java.util.function.Predicate;
import java.util.stream.Collectors;
/**
@@ -54,11 +55,12 @@ 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;
private final MemberPresence presence; // CB-113: which workers are MCP-connected (available)
private final Predicate<String> deliverable;
private final HttpServlet mcpServlet; // MCP Streamable-HTTP endpoint, mounted at /mcp (nullable)
private final CallerResolver auth; // CB-501: null → authz not enforced (legacy behaviour)
private final Metrics metrics; // CB-502: null → /metrics not exposed
@@ -70,9 +72,9 @@ public final class FleetApp {
* behaviour without each needing an auth fixture.
*/
public FleetApp(HerdrClient herdr, PeerLauncher workers, SessionManager sessions,
MessageService messages, MemberPresence presence,
HttpServlet mcpServlet) {
this(herdr, workers, sessions, messages, presence, mcpServlet, null, null);
MessageService messages, MemberPresence presence,
HttpServlet mcpServlet) {
this(herdr, workers, sessions, messages, presence, mcpServlet, null, null, presence::isPresent);
}
/**
@@ -82,13 +84,39 @@ public final class FleetApp {
* the endpoint
*/
public FleetApp(HerdrClient herdr, PeerLauncher workers, SessionManager sessions,
MessageService messages, MemberPresence presence,
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics) {
MessageService messages, MemberPresence presence,
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics) {
this(herdr, workers, sessions, messages, presence, mcpServlet, auth, metrics, presence::isPresent);
}
/**
* @param deliverable the injector's readiness gate, shared so status reports its real result
*/
public FleetApp(HerdrClient herdr, PeerLauncher workers, SessionManager sessions,
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;
this.presence = presence;
this.deliverable = deliverable;
this.mcpServlet = mcpServlet;
this.auth = auth;
this.metrics = metrics;
@@ -179,30 +207,64 @@ 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}.
*/
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;
}
if (memberHerdr != herdr) {
try {
memberHerdr.call("ping");
} catch (HerdrException e) {
ctx.status(503).json(Map.of(
"status", "degraded",
"herdr", "member unreachable",
"detail", e.getMessage()));
return;
}
}
ctx.status(200).json(Map.of(
"status", "ok",
"herdr", Map.of(
"version", pong.path("version").asText(""),
"protocol", pong.path("protocol").asInt())));
}
/** 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(""),
@@ -211,7 +273,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. */
@@ -507,9 +568,9 @@ public final class FleetApp {
/**
* Live lifecycle status of a worker (MCP `fleet_status` wraps this in CB-105), plus its
* <em>readiness</em> (CB-113): {@code ready} is true once the worker's Claude has connected the
* bridge MCP — the reliable "available to receive a task" signal, unlike bare {@code idle}, which
* is also true during boot.
* <em>readiness</em>: {@code ready} is true when the injector can deliver to the target. A
* spawned member must connect the bridge MCP first, while a known lead is ready without member
* presence. This differs from bare {@code idle}, which is also true during member boot.
*/
private void sessionStatus(Context ctx) {
String id = ctx.pathParam("id");
@@ -520,7 +581,7 @@ public final class FleetApp {
Map<String, Object> body = new LinkedHashMap<>();
body.put("sessionId", id);
body.put("status", messages.status(id).name().toLowerCase());
body.put("ready", presence.isPresent(id));
body.put("ready", deliverable.test(id));
// CB-582: a worker paused mid-turn in an async fleet_ask is otherwise invisible to a
// status poll — surface the open question and how to answer it, same as fleet_poll's
// Phase.ASKING view.
@@ -7,6 +7,8 @@ import java.io.BufferedReader;
import java.io.IOException;
import java.io.InputStreamReader;
import java.io.UncheckedIOException;
import java.net.URI;
import java.net.URISyntaxException;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
@@ -20,6 +22,7 @@ import java.util.Optional;
import java.util.Set;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.Consumer;
import java.util.stream.Collectors;
/**
@@ -64,6 +67,10 @@ public final class GitWorktrees implements Worktrees {
/** What {@link #isolateToolSurface} writes for {@code .autoenv}: a valid, empty env file. */
private static final String NEUTRAL_AUTOENV_CONFIG = "";
/** A credential helper command which reads only an environment variable at Git call time. */
private static final String ENVIRONMENT_CREDENTIAL_HELPER = "!f() { if [ \"$1\" = get ]; then "
+ "printf 'username=%s\\npassword=%s\\n\\n' git \"$WORKER_GITEA_TOKEN\"; fi; }; f";
/**
* A tracked project config that is hostile in a provisioned worktree, and what to replace it
* with. {@link #file} is the repo-relative path; {@link #stub} is a neutral but VALID payload for
@@ -81,6 +88,7 @@ public final class GitWorktrees implements Worktrees {
);
private final String configuredRoot;
private final Consumer<String> afterWorktreeAdded;
private final SecureRandom random = new SecureRandom();
private final AtomicLong seq = new AtomicLong();
@@ -91,7 +99,13 @@ public final class GitWorktrees implements Worktrees {
/** @param configuredRoot nullable absolute or relative path; null/blank derives a sibling of the repo root. */
public GitWorktrees(String configuredRoot) {
this(configuredRoot, _ -> {});
}
/** Test seam for changing a real worktree between its creation and its security check. */
GitWorktrees(String configuredRoot, Consumer<String> afterWorktreeAdded) {
this.configuredRoot = configuredRoot;
this.afterWorktreeAdded = afterWorktreeAdded == null ? _ -> {} : afterWorktreeAdded;
}
@Override
@@ -107,11 +121,148 @@ public final class GitWorktrees implements Worktrees {
}
String wt = path.toAbsolutePath().toString();
log.info("adding worktree branch={} path={} base={}", branch, wt, base);
removeUserInfoFromHttpsOrigin(repoRoot);
exec("git", "-C", repoRoot, "worktree", "add", wt, "-b", branch, base);
afterWorktreeAdded.accept(wt);
requireCredentialFreeHttpsOrigin(wt);
configureEnvironmentCredentialHelper(repoRoot, wt);
configureHttpsUrlRewriteForSshOrigin(repoRoot, wt);
isolateToolSurface(wt);
return wt;
}
/**
* A linked worktree shares its primary checkout's git config. Remove HTTPS user info before
* adding one, so a credential accidentally embedded in that config cannot reach the member.
*/
private void removeUserInfoFromHttpsOrigin(String repoRoot) {
if (exitCode("git", "-C", repoRoot, "config", "--get", "remote.origin.url") != 0) {
return;
}
String origin = exec("git", "-C", repoRoot, "config", "--get", "remote.origin.url").trim();
URI uri;
try {
uri = new URI(origin);
} catch (URISyntaxException e) {
throw new WorktreeException("origin URL is invalid; cannot provision a safe worktree", e);
}
if (!"https".equalsIgnoreCase(uri.getScheme()) || uri.getUserInfo() == null) {
return;
}
int schemeEnd = origin.indexOf("://") + 3;
int userInfoEnd = origin.indexOf('@', schemeEnd);
if (userInfoEnd < schemeEnd) {
throw new WorktreeException("origin URL has invalid HTTPS user info; cannot provision safely");
}
String cleanOrigin = origin.substring(0, schemeEnd) + origin.substring(userInfoEnd + 1);
exec("git", "-C", repoRoot, "remote", "set-url", "origin", cleanOrigin);
log.info("removed HTTPS user info from forge origin before provisioning worktree");
}
/** Refuse the worktree if Git still resolves any HTTPS origin URL with embedded credentials. */
private void requireCredentialFreeHttpsOrigin(String worktreePath) {
if (exitCode("git", "-C", worktreePath, "remote", "get-url", "--all", "origin") != 0) {
return;
}
String origins = exec("git", "-C", worktreePath, "remote", "get-url", "--all", "origin");
for (String origin : origins.split("\\R")) {
try {
URI uri = new URI(origin);
if ("https".equalsIgnoreCase(uri.getScheme()) && uri.getUserInfo() != null) {
throw new WorktreeException("worktree origin contains HTTPS user info; refusing provision");
}
} catch (URISyntaxException e) {
throw new WorktreeException("worktree origin URL is invalid; refusing provision", e);
}
}
}
/** 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");
// An empty helper resets values inherited from the system or global config. Without it Git
// asks the next helper after this one, which can expose an operator-level credential.
exec("git", "-C", worktreePath, "config", "--worktree", "--replace-all", "credential.helper", "");
exec("git", "-C", worktreePath, "config", "--worktree", "--add", "credential.helper",
ENVIRONMENT_CREDENTIAL_HELPER);
}
/**
* {@link #configureEnvironmentCredentialHelper} only ever fires for an HTTPS origin — Git never
* consults a {@code credential.helper} for an SSH transport. This repo's own origin is
* {@code ssh://git@git.ltms.dev:2224/fleet/fleetd.git}, so a member sitting on that origin never
* reaches the helper and the repo-scoped {@code WORKER_GITEA_TOKEN} is simply not used.
*
* <p>An earlier version of this javadoc justified the rewrite by claiming a member <em>cannot</em>
* push once {@code memberCredentials.policy: allow-list} blocks {@code SSH_AUTH_SOCK}, because
* "there is no private key file on this host, only an ssh-agent socket". That premise is false
* (fleetd #184): {@code ssh -G} resolves a readable, passphrase-free {@code IdentityFile} outside
* {@code ~/.ssh}, and a member — same OS user — pushes over SSH with the socket blanked. The
* rewrite is still worth having, but for the reason below rather than that one: it routes the
* member through its own scoped token instead of the operator's ssh identity, which is what makes
* a member's pushes attributable and revocable.
*
* <p>The fix is a <em>worktree-scoped</em> URL rewrite: {@code url.<https-base>.insteadOf
* <ssh-base>}, set with {@code --worktree} so it lands only in
* {@code <worktree>/.git/worktrees/<name>/config.worktree} (enabled by
* {@code extensions.worktreeConfig}, already turned on above) and never touches the shared
* repo-level config the primary checkout also reads. {@code insteadOf} — not
* {@code pushInsteadOf} — because a member may also need to fetch or rebase, and both should go
* through the member's own token for the same reason.
*
* <p>The host (and, for the rewrite's SSH-side match, the port) come from parsing the origin
* itself — never a hardcoded forge host, which is exactly what #177 removed. An origin that is
* already {@code https://} is left alone; the credential helper already covers it. An origin
* that is neither {@code ssh://} nor {@code https://} — including the scp-like shorthand
* ({@code git@host:path}, no scheme) — is left untouched deliberately: that shorthand's
* {@code host:path} split is defined by the user's ssh_config aliases, not by URI syntax, so
* guessing at it risks rewriting to the wrong place. A repo provisioned from that form keeps
* today's (broken, if the policy blocks the agent) SSH-only behaviour rather than a wrong rewrite.
*/
/**
* Blank the user-info of a remote URL before it reaches a log. A remote URL is not obviously a
* credential channel, which is exactly why one has leaked here three times ({@code git remote -v}
* printing a token inline, and fleetd #157 / #182). An {@code ssh://} authority normally carries
* only {@code git@}, so this usually changes nothing — it is here so that the one origin that
* does carry a secret cannot print it. Matches every {@code ://…@} pair, not just the first.
*/
private static String redactUserInfo(String url) {
return url == null ? null : url.replaceAll("://[^@/]*@", "://<redacted>@");
}
private void configureHttpsUrlRewriteForSshOrigin(String repoRoot, String worktreePath) {
if (exitCode("git", "-C", repoRoot, "config", "--get", "remote.origin.url") != 0) {
return;
}
String origin = exec("git", "-C", repoRoot, "config", "--get", "remote.origin.url").trim();
URI uri;
try {
uri = new URI(origin);
} catch (URISyntaxException e) {
log.warn("origin URL {} is not a valid URI; skipping worktree HTTPS rewrite", redactUserInfo(origin));
return;
}
String scheme = uri.getScheme();
if (!"ssh".equalsIgnoreCase(scheme)) {
// Already https:// (the credential helper covers it), or a scheme-less/scp-like origin
// left alone on purpose — see the javadoc above.
log.debug("origin scheme is not ssh ({}) — no worktree HTTPS rewrite needed", redactUserInfo(origin));
return;
}
String host = uri.getHost();
String authority = uri.getRawAuthority();
if (host == null || host.isBlank() || authority == null || authority.isBlank()) {
log.warn("ssh origin {} has no resolvable host; skipping worktree HTTPS rewrite", redactUserInfo(origin));
return;
}
String sshBase = "ssh://" + authority + "/";
String httpsBase = "https://" + host + "/";
exec("git", "-C", worktreePath, "config", "--worktree", "--replace-all",
"url." + httpsBase + ".insteadOf", sshBase);
log.info("worktree {} rewrites {} to {} (worktree-scoped; parent checkout untouched)",
worktreePath, redactUserInfo(sshBase), httpsBase);
}
/**
* Neutralize the worktree's worktree-hostile project configs so a worker inherits only the tools
* and environment its launcher mounts (the bridge via {@code --mcp-config}, the opencode config
@@ -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("));
}
}
@@ -1,14 +1,18 @@
package dev.ltms.fleet.health;
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;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.inject.Injector;
import dev.ltms.fleet.msg.InMemoryReplyInbox;
import dev.ltms.fleet.msg.MessageService;
import dev.ltms.fleet.msg.Rendezvous;
import dev.ltms.fleet.peer.MemberRole;
import dev.ltms.fleet.session.MemberSession;
import dev.ltms.fleet.session.SessionManager;
import dev.ltms.fleet.member.ClaudeCodeLauncher;
import dev.ltms.fleet.herdr.WorkspaceControl;
@@ -17,10 +21,15 @@ import dev.ltms.fleet.config.FleetConfig;
import org.junit.jupiter.api.Test;
import org.slf4j.LoggerFactory;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.BiConsumer;
import java.util.function.LongSupplier;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
@@ -40,7 +49,7 @@ class FleetHealthMonitorTest {
MessageService messages = new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox());
var scheduler = Executors.newSingleThreadScheduledExecutor();
FleetHealthMonitor monitor = new FleetHealthMonitor(agents, sessions::roster, messages, scheduler, () -> 1, 60,
(_, _) -> { });
600, (_, _) -> { });
monitor.tick();
monitor.stop();
assertEquals(1, herdr.calls.stream().filter(call -> call.method().equals("agent.list")).count());
@@ -52,7 +61,7 @@ class FleetHealthMonitorTest {
var scheduler = Executors.newSingleThreadScheduledExecutor();
FleetHealthMonitor monitor = new FleetHealthMonitor(agents, java.util.List::of,
new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()),
scheduler, () -> 1, 60, (_, _) -> { });
scheduler, () -> 1, 60, 600, (_, _) -> { });
monitor.tick();
herdr.healthy(true);
monitor.tick();
@@ -71,7 +80,7 @@ class FleetHealthMonitorTest {
var scheduler = Executors.newSingleThreadScheduledExecutor();
FleetHealthMonitor monitor = new FleetHealthMonitor(agents, java.util.List::of,
new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()),
scheduler, () -> 1, 60, (_, _) -> { });
scheduler, () -> 1, 60, 600, (_, _) -> { });
monitor.reportTransition("term_a", HealthState.TURN_BOUNDARY_LOST);
monitor.reportTransition("term_a", HealthState.TURN_BOUNDARY_LOST);
monitor.stop();
@@ -82,6 +91,148 @@ class FleetHealthMonitorTest {
}
}
@Test void failedListMarksEveryMemberControlLinkDownAndReschedules() {
Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
try {
FakeHerdr herdr = new FakeHerdr().healthy(false);
ScheduledThreadPoolExecutor scheduler = new ScheduledThreadPoolExecutor(1);
FleetHealthMonitor monitor = monitor(herdr, List.of(
member("term_one", MemberSession.State.READY, 0, 0),
member("term_two", MemberSession.State.BUSY, 0, 0)), scheduler, () -> 1, 600,
(_, _) -> { });
monitor.tick();
assertEquals(1, appender.list.stream().filter(event -> event.getFormattedMessage()
.contains("member=term_one state=CONTROL_LINK_DOWN")).count());
assertEquals(1, appender.list.stream().filter(event -> event.getFormattedMessage()
.contains("member=term_two state=CONTROL_LINK_DOWN")).count());
assertEquals(1, scheduler.getQueue().size());
monitor.stop();
} finally {
logger.detachAppender(appender);
}
}
@Test void missingRosterMemberGoesGoneAndFailsTargetOnceAcrossTicks() {
FakeHerdr herdr = new FakeHerdr();
RecordingFailTarget failTarget = new RecordingFailTarget();
var scheduler = Executors.newSingleThreadScheduledExecutor();
FleetHealthMonitor monitor = monitor(herdr,
List.of(member("term_missing", MemberSession.State.READY, 0, 0)), scheduler,
() -> 1, 600, failTarget);
monitor.tick();
monitor.tick();
monitor.stop();
assertEquals(1, failTarget.calls.size());
assertEquals("term_missing", failTarget.calls.get(0).target());
assertTrue(failTarget.calls.get(0).reason().contains("GONE"));
}
@Test void spawningMemberBecomesNeverReadyOnlyAfterReadinessGrace() {
Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
try {
AtomicLong clock = new AtomicLong(FleetHealthMonitor.READINESS_GRACE_NANOS - 1);
RecordingFailTarget failTarget = new RecordingFailTarget();
var scheduler = Executors.newSingleThreadScheduledExecutor();
FleetHealthMonitor monitor = monitor(new FakeHerdr(),
List.of(member("term_starting", MemberSession.State.SPAWNING, 0, 0)), scheduler,
clock::get, 600, failTarget);
monitor.tick();
assertEquals(0, failTarget.calls.size());
clock.set(FleetHealthMonitor.READINESS_GRACE_NANOS);
monitor.tick();
monitor.stop();
assertEquals(1, failTarget.calls.size());
assertTrue(failTarget.calls.get(0).reason().contains("NEVER_READY"));
assertTrue(appender.list.stream().anyMatch(event -> event.getFormattedMessage()
.contains("state=NEVER_READY previous=STARTING")));
} finally {
logger.detachAppender(appender);
}
}
@Test void workingSuspectConfigChangesTheStallBoundary() {
Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
try {
long nowNanos = TimeUnit.SECONDS.toNanos(600);
FakeHerdr herdr = new FakeHerdr()
.withAgent("short", "term_short", "pane_short", "tab_short")
.withAgent("long", "term_long", "pane_long", "tab_long");
FleetConfig.Health shortConfig = new FleetConfig.Health(true, 30, 300, null, null);
FleetConfig.Health longConfig = new FleetConfig.Health(true, 30, 601, null, null);
var shortScheduler = Executors.newSingleThreadScheduledExecutor();
var longScheduler = Executors.newSingleThreadScheduledExecutor();
FleetHealthMonitor shortMonitor = monitor(herdr,
List.of(member("term_short", MemberSession.State.BUSY, 0, 0)), shortScheduler,
() -> nowNanos, shortConfig.workingSuspectAfterOrDefault(), (_, _) -> { });
FleetHealthMonitor longMonitor = monitor(herdr,
List.of(member("term_long", MemberSession.State.BUSY, 0, 0)), longScheduler,
() -> nowNanos, longConfig.workingSuspectAfterOrDefault(), (_, _) -> { });
shortMonitor.tick();
longMonitor.tick();
shortMonitor.stop();
longMonitor.stop();
assertTrue(appender.list.stream().anyMatch(event -> event.getFormattedMessage()
.contains("member=term_short state=STALL_SUSPECTED")));
assertEquals(0, appender.list.stream().filter(event -> event.getFormattedMessage()
.contains("member=term_long state=STALL_SUSPECTED")).count());
} finally {
logger.detachAppender(appender);
}
}
@Test void workingSuspectConfigKeepsItsDefaultAndFloor() {
assertEquals(600, new FleetConfig.Health(true, null, null, null, null)
.workingSuspectAfterOrDefault());
assertEquals(300, new FleetConfig.Health(true, null, 1, null, null)
.workingSuspectAfterOrDefault());
}
@Test void goneMemberRecoveryLogsOnceWithoutRefiringTargetFailure() {
Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class);
Level previousLevel = logger.getLevel();
logger.setLevel(Level.INFO);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
try {
FakeHerdr herdr = new FakeHerdr();
RecordingFailTarget failTarget = new RecordingFailTarget();
var scheduler = Executors.newSingleThreadScheduledExecutor();
FleetHealthMonitor monitor = monitor(herdr,
List.of(member("term_recovered", MemberSession.State.READY, 0, 0)), scheduler,
() -> 1, 600, failTarget);
monitor.tick();
herdr.withAgent("recovered", "term_recovered", "pane_recovered", "tab_recovered");
monitor.tick();
monitor.stop();
assertEquals(1, failTarget.calls.size());
assertTrue(appender.list.stream().anyMatch(event -> event.getFormattedMessage()
.contains("member=term_recovered recovered state=IDLE previous=GONE")));
} finally {
logger.detachAppender(appender);
logger.setLevel(previousLevel);
}
}
// --- CB-580: a member that reaches GONE/NEVER_READY must fail its waiting tickets
private static FleetHealthMonitor monitorWith(BiConsumer<String, String> failTarget) {
@@ -90,7 +241,23 @@ class FleetHealthMonitorTest {
var scheduler = Executors.newSingleThreadScheduledExecutor();
return new FleetHealthMonitor(agents, java.util.List::of,
new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()),
scheduler, () -> 1, 60, failTarget);
scheduler, () -> 1, 60, 600, failTarget);
}
private static FleetHealthMonitor monitor(FakeHerdr herdr, List<MemberSession> roster,
java.util.concurrent.ScheduledExecutorService scheduler,
LongSupplier clock, long workingSuspectAfterSeconds,
BiConsumer<String, String> failTarget) {
AgentControl agents = new AgentControl(herdr);
return new FleetHealthMonitor(agents, () -> roster,
new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()),
scheduler, clock, 60, workingSuspectAfterSeconds, failTarget);
}
private static MemberSession member(String terminalId, MemberSession.State state,
long spawnedAtNanos, long lastActivityAtNanos) {
return new MemberSession("pane-" + terminalId, terminalId, "test", MemberRole.DEV, "/tmp", null,
spawnedAtNanos, lastActivityAtNanos, 0, state, null, null);
}
@Test void terminalTransitionFailsTheTargetOnce() {
@@ -167,4 +334,104 @@ class FleetHealthMonitorTest {
throw new RuntimeException("boom");
}
}
// --- CB-643: the two-tick gate on an orphaned delegation ------------------------------------
private static final String ORPHAN_WARNING = "member=term_a state=DELEGATION_ORPHANED";
/**
* Everything one orphan test needs: a live target whose async ticket really is orphaned. The
* recipe is the one MessageServiceTest proves for {@code hasOrphanedDelegation} — an unanswered
* {@code fleet_ask} lapses, so the ticket returns to PENDING while its forward waiter is already
* closed. {@code rendezvous} is exposed so a test can turn the fact off and on again.
*/
private record OrphanFleet(FakeHerdr herdr, Rendezvous rendezvous, MessageService messages,
FleetHealthMonitor monitor) {
}
private static OrphanFleet orphanedTarget(BiConsumer<String, String> failTarget) throws Exception {
FakeHerdr herdr = new FakeHerdr().withAgent("worker", "term_a", "pane-term_a", "tab_a")
.readText("$ prompt");
AgentControl agents = new AgentControl(herdr);
Rendezvous rendezvous = new Rendezvous();
Injector injector = new Injector(agents);
InMemoryReplyInbox inbox = new InMemoryReplyInbox();
inbox.own("term_a");
MessageService messages = new MessageService(agents, injector, rendezvous, inbox);
messages.sendAsync("term_a", "task that asks");
long deadline = System.currentTimeMillis() + 2000;
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting("term_a"), "the async send should have opened its waiter");
injector.onStatus("term_a", AgentStatus.IDLE); // deliver the task
injector.onStatus("term_a", AgentStatus.WORKING); // the worker picks it up
assertEquals(MessageService.AskOutcome.TIMED_OUT,
messages.ask("term_a", "which config?", 200).outcome());
assertTrue(messages.hasOrphanedDelegation("term_a"),
"the lapsed ask should leave a PENDING ticket with nothing in flight");
FleetHealthMonitor monitor = new FleetHealthMonitor(agents, () -> List.of(
member("term_a", MemberSession.State.READY, 0, 0)),
messages, new ScheduledThreadPoolExecutor(1), () -> 1, 60, 600, failTarget);
return new OrphanFleet(herdr, rendezvous, messages, monitor);
}
private static long orphanWarnings(ListAppender<ILoggingEvent> appender) {
return appender.list.stream()
.filter(event -> event.getFormattedMessage().contains(ORPHAN_WARNING)).count();
}
@Test void oneOrphanObservationIsNotYetReportedButTwoAre() throws Exception {
Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
OrphanFleet fleet = orphanedTarget((_, _) -> { });
try {
fleet.monitor().tick();
assertEquals(0, orphanWarnings(appender),
"one observation can be an ordinary race, so health must not report it yet");
fleet.monitor().tick();
assertEquals(1, orphanWarnings(appender),
"a second consecutive observation confirms the orphan and is reported once");
} finally {
fleet.messages().abandon("term_a", "test over");
fleet.monitor().stop();
logger.detachAppender(appender);
}
}
@Test void aSingleCleanTickResetsTheOrphanStreak() throws Exception {
Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
OrphanFleet fleet = orphanedTarget((_, _) -> { });
try {
fleet.monitor().tick(); // first observation: streak 1
// Something is in flight for the target again, so the orphan fact reads false.
var waiter = fleet.rendezvous().open("term_a");
fleet.monitor().tick();
assertEquals(0, orphanWarnings(appender), "a clean tick must clear the streak");
// Deregister that waiter — resolving it is not enough, the sender's close() is what
// removes it — so the ticket is orphaned again, from a streak of zero.
fleet.rendezvous().close("term_a", waiter);
assertTrue(fleet.messages().hasOrphanedDelegation("term_a"));
fleet.monitor().tick();
assertEquals(0, orphanWarnings(appender),
"the streak restarted, so this first observation is not reported either");
fleet.monitor().tick();
assertEquals(1, orphanWarnings(appender), "two consecutive observations report once");
} finally {
fleet.messages().abandon("term_a", "test over");
fleet.monitor().stop();
logger.detachAppender(appender);
}
}
}
@@ -40,6 +40,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
@@ -75,6 +76,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;
@@ -268,7 +278,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");
}
}
@@ -197,10 +197,15 @@ class CompletionResolverTest {
void resolvesSynchronouslyBeforePostTurnContextClearing() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ previous answer\n❯ ");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
// fleetd#164: an injectable clock, so this real (non-crash) turn lands outside MIN_TURN_NANOS
// — captureBaseline and resolveBeforePostAction below run back-to-back with no real delay.
long[] clock = {1_000_000_000L};
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none(), () -> clock[0]);
var waiter = rendezvous.open("term_a");
resolver.captureBaseline("term_a", new TurnToken("term_a", waiter));
herdr.readText("⏺ answer that /clear would erase\n❯ ");
clock[0] += CompletionResolver.MIN_TURN_NANOS + 1; // this turn took longer than the floor
resolver.resolveBeforePostAction("term_a");
@@ -219,7 +224,11 @@ class CompletionResolverTest {
String longBlock = "⏺ " + "x".repeat(CompletionResolver.MAX_SCRAPE_CHARS + 500) + "\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(longBlock);
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
// fleetd#164: an injectable clock so captureBaseline and resolve (back-to-back, no real
// delay) don't trip the too-fast-turn floor — this test is about the suppression guard, not timing.
long[] clock = {1_000_000_000L};
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none(), () -> clock[0]);
var waiter = rendezvous.open("term_a"); // a send is blocked on this turn
resolver.captureBaseline("term_a", new TurnToken("term_a", waiter)); // baseline is the clipped >cap block
@@ -227,6 +236,7 @@ class CompletionResolverTest {
assertEquals(CompletionResolver.MAX_SCRAPE_CHARS, turn.baseline().length(),
"the delivery baseline is clipped to the same cap resolve() applies to the tail");
clock[0] += CompletionResolver.MIN_TURN_NANOS + 1; // outside the floor
resolver.resolve("term_a", turn); // scrape unchanged → clipped tail == baseline → suppress
assertFalse(waiter.isDone(),
@@ -249,12 +259,13 @@ class CompletionResolverTest {
}
@Test
void resolvesWhenTheScrapeItselfFailsEvenWithABaselinePresent() {
// The most important branch of the CB-115 guard: a failed read means the resolver could not
// SEE the screen — "couldn't see", not "no change". It must still resolve the send (an empty
// tail beats hanging until the caller's timeout), even though a baseline was captured. The
// baseline here is "" (an empty pane at delivery), so without the !scrapeFailed clause the
// byte-identical guard would wrongly match the empty tail and suppress.
void aFailedScrapeResolvesAsAFailureEvenWithABaselinePresent() {
// fleetd#164: this test used to assert that a failed read resolved the send as a SUCCESS
// carrying an empty string ("an empty tail beats hanging until the caller's timeout") — that
// was the bug this ticket fixes: a lost turn and a genuine empty answer looked identical to
// every caller. This test encoded the bug and is changed here: a failed read must fail the
// send instead, naming the member, whether or not a baseline was captured (the baseline here
// is "", an empty pane at delivery — proof this isn't the CB-115 misattribution path either).
FakeHerdr herdr = new FakeHerdr().healthy(false); // agent.read throws HerdrException
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
@@ -265,8 +276,82 @@ class CompletionResolverTest {
assertTrue(waiter.isDone(),
"a failed scrape must still resolve the send, not hang until the caller's timeout");
assertEquals(Rendezvous.Kind.COMPLETION, waiter.getNow(null).kind());
assertEquals("", waiter.getNow(null).text(), "the tail is empty because the screen was unreadable");
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(),
"a failed read is a lost turn, not a successful empty reply");
assertTrue(waiter.getNow(null).text().contains("term_a"),
"the failure names the member: " + waiter.getNow(null).text());
assertTrue(waiter.getNow(null).text().contains("could not be read"),
"the failure explains the scrape could not be read: " + waiter.getNow(null).text());
}
// --- fleetd#164: an empty (but readable) scrape must never resolve as a success --------------
@Test
void anEmptyScrapeResolvesAsAFailureNamingTheMember() {
// The core defect: a scrape that read CLEANLY but produced zero characters used to resolve
// the send as a SUCCESS carrying "" — indistinguishable, to every caller, from a worker that
// genuinely finished with nothing to say. A lost turn must never look like a real empty reply.
FakeHerdr herdr = new FakeHerdr().readText("");
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(), "an empty scrape must still resolve the send, not hang");
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(),
"an empty scrape is a lost turn, not a successful empty reply");
assertTrue(waiter.getNow(null).text().contains("term_a"),
"the failure names the member: " + waiter.getNow(null).text());
assertTrue(waiter.getNow(null).text().toLowerCase().contains("empty"),
"the failure says the scrape was empty: " + waiter.getNow(null).text());
}
// --- fleetd#164: a BUSY -> DONE transition inside the floor is a crash, not a fast answer ------
@Test
void aBusyToDoneTransitionInsideTheFloorResolvesAsAFailure() {
// The exact fleetd#164 scenario: the backend returned an HTTP 400 before the worker did
// anything, and the member went BUSY -> DONE in ~1s. That transition alone is indistinguishable
// from a genuine (if unusually fast) completion, so the resolver leans on the floor to catch it.
FakeHerdr herdr = new FakeHerdr().readText("⏺ HTTP 400: invalid request\n❯ ");
Rendezvous rendezvous = new Rendezvous();
long[] clock = {10_000_000_000L};
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none(), () -> clock[0]);
var waiter = rendezvous.open("term_a");
var turn = new CompletionResolver.InFlight(waiter, null, clock[0]); // delivered "now"
clock[0] += CompletionResolver.MIN_TURN_NANOS - 1; // 1ns inside the floor
resolver.resolve("term_a", turn);
assertTrue(waiter.isDone(), "a suspiciously fast turn must still resolve (as a failure)");
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind());
assertTrue(waiter.getNow(null).text().contains("term_a"),
"the failure names the member: " + waiter.getNow(null).text());
assertTrue(waiter.getNow(null).text().contains("HTTP 400"),
"the failure carries whatever was on screen: " + waiter.getNow(null).text());
}
@Test
void aBusyToDoneTransitionJustOutsideTheFloorResolvesNormally() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ a real, if quick, answer\n❯ ");
Rendezvous rendezvous = new Rendezvous();
long[] clock = {10_000_000_000L};
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none(), () -> clock[0]);
var waiter = rendezvous.open("term_a");
var turn = new CompletionResolver.InFlight(waiter, null, clock[0]); // delivered "now"
clock[0] += CompletionResolver.MIN_TURN_NANOS + 1; // 1ns outside the floor
resolver.resolve("term_a", turn);
assertTrue(waiter.isDone());
assertEquals(Rendezvous.Kind.COMPLETION, waiter.getNow(null).kind(),
"a turn that took longer than the floor resolves normally");
assertEquals("a real, if quick, answer", waiter.getNow(null).text());
}
// --- CB-115/CB-116 fail guard: an already-done or absent waiter is left alone ---------
@@ -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,9 @@
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;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.FakeHerdr;
@@ -8,6 +12,7 @@ import dev.ltms.fleet.peer.Capability;
import dev.ltms.fleet.peer.MemberRole;
import dev.ltms.fleet.peer.SpawnRequest;
import org.junit.jupiter.api.Test;
import org.slf4j.LoggerFactory;
import java.nio.file.Files;
import java.nio.file.Path;
@@ -108,6 +113,122 @@ class HerdrPeerLauncherAllowListWiringTest {
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST, List.of(), List.of(), null);
}
/** Same as {@link #allowList()} but with an operator-configured {@code allow:} list. */
private static Supplier<FleetConfig.MemberCredentials> allowListWithAllow(List<String> allow) {
return () -> new FleetConfig.MemberCredentials(
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST, allow, List.of(), null);
}
/**
* CB-633 follow-up: a name that lives ONLY in {@code memberCredentials.allow:} — no profile
* mentions it — must survive the scrub the real spawn path generates. Calling {@code
* MemberEnvAllowList.derive} directly (as {@link MemberEnvAllowListTest} does) would pass even
* if {@code HerdrPeerLauncher} never threaded {@code allow:} into the derivation at all; this
* test goes through {@link HerdrPeerLauncher#spawn}, the method the daemon actually calls at
* spawn time, so it proves the union is wired in, not just correct in isolation.
*/
@Test
void spawningUnderAllowListPolicyIncludesAnOperatorConfiguredAllowName() {
FakeHerdr herdr = new FakeHerdr();
WiringLauncher launcher = new WiringLauncher(herdr,
allowListWithAllow(List.of("OPERATOR_ONLY_NAME")));
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
Path dir = Path.of(launcher.env.get("ZDOTDIR"));
assertTrue(readAll(dir.resolve(EnvAllowListScrub.SCRUB_FILE)).contains("OPERATOR_ONLY_NAME"),
"a name only in memberCredentials.allow: must reach the generated scrub through the "
+ "real launcher spawn path");
}
/**
* {@code SSH_AUTH_SOCK} is a live ssh-agent handle, not a value — it must stay blocked under
* {@code allow-list} even when the operator lists it under {@code allow:}, because {@code
* sshAuthSock} defaults to blocked. Governed ONLY by {@code memberCredentials.sshAuthSock}.
*/
@Test
void sshAuthSockStaysBlockedEvenWhenListedInMemberCredentialsAllow() {
FakeHerdr herdr = new FakeHerdr();
WiringLauncher launcher = new WiringLauncher(herdr,
allowListWithAllow(List.of("SSH_AUTH_SOCK")));
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
Path dir = Path.of(launcher.env.get("ZDOTDIR"));
String scrub = readAll(dir.resolve(EnvAllowListScrub.SCRUB_FILE));
assertFalse(scrub.contains("'SSH_AUTH_SOCK'"),
"SSH_AUTH_SOCK must not be on the derived allow-list just because the operator put "
+ "it under allow: — sshAuthSock is unset here, so it defaults to block");
}
/**
* CB-633 follow-up criterion 3: on every allow-list spawn the daemon logs one INFO line, shaped
* "member credentials: allowed N of M", with real counts — not constants. Real path: the count
* is asserted after a real {@link HerdrPeerLauncher#spawn} call, reading the log the production
* code actually emits.
*/
@Test
void logsAnAllowedCountLineAgainstTheHostEnvironmentOnEverySpawn() {
FakeHerdr herdr = new FakeHerdr();
// INJECTED is a key of this launch's own env map, so it always survives; the other two are
// neither derived from the profile nor configured anywhere, so they are blanked. Real
// N=1 (INJECTED), real M=3 (all three names) — neither number is hardcoded in the assertion
// by coincidence, they follow directly from this fixture.
Set<String> hostEnvNames = Set.of(INJECTED, "SOME_UNRELATED_NAME", "ANOTHER_UNRELATED_NAME");
WiringLauncher launcher = new WiringLauncher(herdr, allowList(), "/bin/zsh", () -> hostEnvNames);
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 {
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
} finally {
logger.detachAppender(appender);
logger.setLevel(original);
}
assertTrue(appender.list.stream()
.anyMatch(e -> "member credentials: allowed 1 of 3".equals(e.getFormattedMessage())),
"expected 'member credentials: allowed 1 of 3', got: "
+ appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
}
/**
* Lead-review fix: on a NON-zsh shell no scrub ever runs (bash ignores {@code ZDOTDIR}), so the
* "allowed N of M" line — which describes what the scrub does — must not be printed there either.
* Before this fix the line was logged BEFORE the zsh gate, so a non-zsh host printed e.g.
* "allowed 1 of 3" while blocking nothing at all, telling an operator a control ran when it did
* not. Real path: goes through {@link HerdrPeerLauncher#spawn}, same as the sibling test above,
* with the shell fixed to bash so the fallback branch is the one exercised.
*/
@Test
void noAllowedCountLineIsEmittedOnTheNonZshFallbackPath() {
FakeHerdr herdr = new FakeHerdr();
Set<String> hostEnvNames = Set.of(INJECTED, "SOME_UNRELATED_NAME", "ANOTHER_UNRELATED_NAME");
WiringLauncher launcher = new WiringLauncher(herdr, allowList(), "/bin/bash", () -> hostEnvNames);
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 {
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
} finally {
logger.detachAppender(appender);
logger.setLevel(original);
}
assertFalse(appender.list.stream()
.anyMatch(e -> e.getFormattedMessage().startsWith("member credentials: allowed ")),
"no scrub runs on a non-zsh shell, so no 'allowed N of M' count may be printed — got: "
+ appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
}
private static String readAll(Path p) {
try {
return Files.readString(p);
@@ -134,10 +255,16 @@ class HerdrPeerLauncherAllowListWiringTest {
}
WiringLauncher(FakeHerdr herdr, Supplier<FleetConfig.MemberCredentials> creds, String shell) {
this(herdr, creds, shell, null);
}
/** Plus an injectable {@code hostEnvNames} source, for the "allowed N of M" log line test. */
WiringLauncher(FakeHerdr herdr, Supplier<FleetConfig.MemberCredentials> creds, String shell,
Supplier<Set<String>> hostEnvNames) {
super("test", new AgentControl(herdr), new WorkspaceControl(herdr),
Map.of("test", profile()), "test",
name -> "SHELL".equals(name) ? shell : null,
0, () -> 0L, () -> { }, null, creds);
0, () -> 0L, () -> { }, null, creds, hostEnvNames);
}
@Override
@@ -88,6 +88,39 @@ class MemberEnvAllowListTest {
assertTrue(after.containsAll(Set.of("TOKEN_SECOND", "GIT_TOK", "SECOND_KEY")));
}
/**
* CB-633 follow-up: a name that appears ONLY in {@code memberCredentials.allow:} — no profile
* mentions it at all — must still survive the derivation. Before this fix {@code derive} never
* saw {@code allow:}, so setting {@code policy: allow-list} silently blanked exactly this name.
*/
@Test
void aNameOnlyInMemberCredentialsAllowSurvivesDerivation() {
FleetConfig.Profile p = profile("p", "TOKEN_A", null, null, Map.of("KEY_A", "v"));
Set<String> derived = MemberEnvAllowList.derive(List.of(p), Set.of("OPERATOR_ONLY_NAME"));
assertTrue(derived.contains("OPERATOR_ONLY_NAME"),
"memberCredentials.allow: must be unioned in, not ignored");
// and the profile-derived half must still be present — this is a union, not a replacement.
assertTrue(derived.contains("KEY_A"));
assertTrue(derived.contains("TOKEN_A"));
}
/**
* {@code SSH_AUTH_SOCK} is a live handle to the operator's ssh-agent, never a value — so it must
* stay excluded from the derived set even when the operator lists it under {@code allow:} for an
* unrelated reason. It is governed ONLY by {@code memberCredentials.sshAuthSock}, applied
* separately by the caller ({@code HerdrPeerLauncher}).
*/
@Test
void sshAuthSockInMemberCredentialsAllowIsStillExcluded() {
Set<String> derived = MemberEnvAllowList.derive(List.of(), Set.of("SSH_AUTH_SOCK", "OTHER_NAME"));
assertFalse(derived.contains("SSH_AUTH_SOCK"),
"SSH_AUTH_SOCK must never ride in on the generic allow: list");
assertTrue(derived.contains("OTHER_NAME"), "other allow: names are unaffected");
}
/** {@code LC_*} categories are infrastructure by prefix; everything else needs an exact match. */
@Test
void keepsMatchesExactlyPlusTheLocalePrefixRule() {
@@ -0,0 +1,149 @@
package dev.ltms.fleet.msg;
import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.DeliverCallback;
import com.rabbitmq.client.Delivery;
import com.rabbitmq.client.Envelope;
import org.junit.jupiter.api.Test;
import java.io.IOException;
import java.lang.reflect.InvocationHandler;
import java.lang.reflect.Proxy;
import java.nio.charset.StandardCharsets;
import java.util.ArrayDeque;
import java.util.LinkedHashMap;
import java.util.Map;
import java.util.concurrent.atomic.AtomicInteger;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
/**
* Pins the manual-ack prefetch behaviour without a broker. The fake channel models a broker that
* sends no more than its QoS window of unacked deliveries. If {@link AmqpReplyInbox} starts acking
* messages while it adds them to {@code held}, this test drains the whole fake queue instead.
*/
class AmqpReplyInboxPrefetchTest {
@Test
void unackedDeliveriesKeepTheHeldBacklogAtThePrefetchWindow() {
int prefetch = 3;
int published = 8;
PrefetchBroker broker = new PrefetchBroker();
for (int i = 0; i < published; i++) {
broker.publish("m" + i, "payload " + i);
}
try (AmqpReplyInbox inbox = new AmqpReplyInbox(connectionFor(broker.channel()), prefetch)) {
inbox.own("worker");
assertEquals(prefetch, inbox.peek("worker").size(),
"held messages must stop at the unacked prefetch window");
assertEquals(published - prefetch, broker.queuedCount(),
"messages beyond the window must remain on the broker");
assertEquals(0, broker.ackCount(), "receipt must not ack a held message");
inbox.ack("worker", "m0");
assertEquals(prefetch, inbox.peek("worker").size(),
"one caller ack frees exactly one slot for the broker");
assertEquals(published - prefetch - 1, broker.queuedCount(),
"only one queued message may enter after one caller ack");
assertEquals(1, broker.ackCount(), "only the caller ack may reach the broker");
}
}
private static Connection connectionFor(Channel consumeChannel) {
Channel publishChannel = (Channel) Proxy.newProxyInstance(
AmqpReplyInboxPrefetchTest.class.getClassLoader(), new Class<?>[] {Channel.class},
(proxy, method, args) -> defaultValue(method.getReturnType()));
AtomicInteger channelCalls = new AtomicInteger();
InvocationHandler handler = (proxy, method, args) -> {
if (method.getName().equals("createChannel") && (args == null || args.length == 0)) {
return channelCalls.getAndIncrement() == 0 ? consumeChannel : publishChannel;
}
return defaultValue(method.getReturnType());
};
return (Connection) Proxy.newProxyInstance(AmqpReplyInboxPrefetchTest.class.getClassLoader(),
new Class<?>[] {Connection.class}, handler);
}
private static final class PrefetchBroker implements InvocationHandler {
private final ArrayDeque<Delivery> queued = new ArrayDeque<>();
private final Map<Long, Delivery> unacked = new LinkedHashMap<>();
private DeliverCallback consumer;
private int prefetch;
private int acks;
private long nextTag = 1;
Channel channel() {
return (Channel) Proxy.newProxyInstance(AmqpReplyInboxPrefetchTest.class.getClassLoader(),
new Class<?>[] {Channel.class}, this);
}
void publish(String msgId, String content) {
queued.add(new Delivery(new Envelope(nextTag++, false, "", ""),
new AMQP.BasicProperties.Builder().messageId(msgId).build(),
content.getBytes(StandardCharsets.UTF_8)));
}
int queuedCount() {
return queued.size();
}
int ackCount() {
return acks;
}
@Override
public Object invoke(Object proxy, java.lang.reflect.Method method, Object[] args) throws IOException {
switch (method.getName()) {
case "basicQos" -> {
prefetch = (int) args[0];
return null;
}
case "basicConsume" -> {
assertFalse((boolean) args[1], "the inbox consumer must use manual acknowledgements");
consumer = (DeliverCallback) args[2];
deliverAvailable();
return "consumer";
}
case "basicAck" -> {
unacked.remove((long) args[0]);
acks++;
deliverAvailable();
return null;
}
default -> {
return defaultValue(method.getReturnType());
}
}
}
private void deliverAvailable() throws IOException {
while (consumer != null && unacked.size() < prefetch && !queued.isEmpty()) {
Delivery delivery = queued.removeFirst();
unacked.put(delivery.getEnvelope().getDeliveryTag(), delivery);
consumer.handle("consumer", delivery);
}
}
}
private static Object defaultValue(Class<?> type) {
if (!type.isPrimitive() || type == void.class) {
return null;
}
if (type == boolean.class) {
return false;
}
if (type == long.class) {
return 0L;
}
if (type == int.class) {
return 0;
}
return 0;
}
}
@@ -33,8 +33,17 @@ class MessageServiceTest {
private final FakeHerdr herdr = new FakeHerdr().readText("BUILD GREEN: 391 files");
private final AgentControl agents = new AgentControl(herdr);
private final Rendezvous rendezvous = new Rendezvous();
/**
* fleetd#164: this fixture drives a delivery and its completion back-to-back with no real time
* between them, so the real clock would trip {@link CompletionResolver#MIN_TURN_NANOS} on every
* completion-fallback test here. An ever-advancing fake clock stands in for the model-latency and
* herdr round-trips a real turn would spend, so each delivery-then-resolve pair still lands
* outside the floor.
*/
private final java.util.concurrent.atomic.AtomicLong resolverClock = new java.util.concurrent.atomic.AtomicLong();
private final CompletionResolver completion =
new CompletionResolver(agents, rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
new CompletionResolver(agents, rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none(),
() -> resolverClock.addAndGet(CompletionResolver.MIN_TURN_NANOS + 1));
private final Injector injector = new Injector(agents, completion);
private final InMemoryReplyInbox inbox = new InMemoryReplyInbox();
private final MessageService messages = new MessageService(agents, injector, rendezvous, inbox);
@@ -1081,4 +1090,147 @@ class MessageServiceTest {
assertEquals(phase, view.phase());
return view;
}
// --- CB-640: fleet health evidence accessors --------------------------------------------
@Test
void hasQueuedDeliveryIsFalseForAnUnknownTarget() {
assertFalse(messages.hasQueuedDelivery("nobody-ever-sent-here"));
}
@Test
void hasQueuedDeliveryIsFalseBeforeAnyTimeout() {
assertFalse(messages.hasQueuedDelivery(T));
}
@Test
void hasQueuedDeliveryIsTrueAfterAnUndeliveredSendTimesOut() {
// Nothing ever delivers the message (never goes IDLE/BLOCKED), so the send times out with
// TIMED_OUT_QUEUED — same setup as sendTimesOutBeforeDeliveryIsQueuedNotWorking above.
MessageService.Reply r = messages.send(T, "never delivered", 50);
assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, r.outcome());
assertTrue(messages.hasQueuedDelivery(T),
"a TIMED_OUT_QUEUED send leaves the message still queued in the injector");
}
@Test
void hasQueuedDeliveryClearsOnceTheTargetsNextDeliveryIsAccepted() throws Exception {
assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, messages.send(T, "first", 50).outcome());
assertTrue(messages.hasQueuedDelivery(T));
// A fresh send accepts delivery (opens its own waiter) — the stale queued fact is cleared.
CompletableFuture<MessageService.Reply> second = sendAsync();
awaitWaiting();
assertFalse(messages.hasQueuedDelivery(T),
"a fresh accepted delivery supersedes the earlier queued fact");
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
assertTrue(rendezvous.resolve(T, "second done"));
second.get(5, TimeUnit.SECONDS);
}
@Test
void hasQueuedDeliveryClearsOnAbandon() {
assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, messages.send(T, "first", 50).outcome());
assertTrue(messages.hasQueuedDelivery(T));
messages.abandon(T, "session released");
assertFalse(messages.hasQueuedDelivery(T), "a torn-down target has nothing left queued for it");
}
@Test
void hasStrandedReplyIsFalseForAnUnknownTarget() {
assertFalse(messages.hasStrandedReply("nobody-ever-sent-here"));
}
@Test
void hasStrandedReplyIsFalseWhenTheReplyResolvedALiveSend() throws Exception {
CompletableFuture<MessageService.Reply> send = sendAsync();
awaitWaiting();
assertTrue(messages.reply(T, "resolved-live"));
assertFalse(messages.hasStrandedReply(T), "a reply that resolved an open send is not stranded");
MessageService.Reply r = send.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.REPLIED, r.outcome());
}
@Test
void hasStrandedReplyIsTrueWhenNoSendWasWaiting() {
// No send is open for T — the reply queues into the inbox and is recorded as stranded.
assertTrue(messages.reply(T, "nobody was waiting"));
assertTrue(messages.hasStrandedReply(T),
"a reply with no open send strands, even though it is safely queued in the inbox");
}
@Test
void hasStrandedReplyClearsOnceTheTargetsNextDeliveryIsAccepted() throws Exception {
assertTrue(messages.reply(T, "stray"));
assertTrue(messages.hasStrandedReply(T));
// The next accepted delivery for T clears the stale stranding fact — the one case the
// ticket calls out as the one that matters.
CompletableFuture<MessageService.Reply> send = sendAsync();
awaitWaiting();
assertFalse(messages.hasStrandedReply(T),
"a stranded reply must clear once the target's delivery is accepted again");
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
assertTrue(rendezvous.resolve(T, "done"));
send.get(5, TimeUnit.SECONDS);
}
@Test
void hasStrandedReplyClearsOnAbandon() {
assertTrue(messages.reply(T, "stray"));
assertTrue(messages.hasStrandedReply(T));
messages.abandon(T, "session released");
assertFalse(messages.hasStrandedReply(T), "a torn-down target has nothing left to strand");
}
@Test
void hasOrphanedDelegationIsFalseForAnUnknownTarget() {
assertFalse(messages.hasOrphanedDelegation("nobody-ever-sent-here"));
}
@Test
void hasOrphanedDelegationIsFalseWhilePendingTicketsHaveAnAcceptedDelivery() throws Exception {
// Mirrors abandonFailsEveryPendingAsyncTicketForTheReleasedTarget above: "first" holds the
// session lock and its waiter is open, so the target genuinely has something in flight even
// though "second" and "third" are themselves parked (PENDING) behind the lock.
messages.sendAsync(T, "first task");
awaitWaiting(); // first task owns the target lock and rendezvous waiter
String second = messages.sendAsync(T, "second task");
assertEquals(MessageService.Phase.PENDING, messages.poll(second).phase());
assertFalse(messages.hasOrphanedDelegation(T),
"the target has an accepted delivery in flight (first), so nothing here is orphaned");
assertTrue(messages.abandon(T, "session released")); // release the lock and the parked tickets
}
@Test
void hasOrphanedDelegationIsTrueOnceAnUnansweredAskLapsesBackToPending() throws Exception {
// Same setup as unansweredAsyncQuestionReturnsTheTicketToPendingAndReleasesItsTarget above:
// once the ask lapses, the ticket goes back to PENDING but send() already closed the
// forward waiter the instant the question surfaced — nothing is left in flight for T.
String ticket = messages.sendAsync(T, "task that asks");
awaitWaiting();
injectDelivery();
assertEquals(MessageService.AskOutcome.TIMED_OUT,
messages.ask(T, "which config?", 200).outcome());
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket).phase());
assertFalse(messages.hasAcceptedDelivery(T), "the forward waiter closed when the question surfaced");
assertFalse(messages.hasQueuedDelivery(T), "this ticket never timed out as queued");
assertTrue(messages.hasOrphanedDelegation(T),
"a PENDING ticket with no accepted or queued delivery for its target is orphaned");
assertTrue(messages.abandon(T, "session released")); // clean up the still-open ticket
}
}
@@ -32,6 +32,7 @@ import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.UUID;
import java.util.function.Predicate;
import static org.junit.jupiter.api.Assertions.*;
@@ -64,6 +65,11 @@ class FleetAppTest {
}
private int start(FakeHerdr herdr, String workerBaseUrl, Set<String> allow, String placement, Worktrees worktrees) {
return start(herdr, workerBaseUrl, allow, placement, worktrees, ignored -> false);
}
private int start(FakeHerdr herdr, String workerBaseUrl, Set<String> allow, String placement,
Worktrees worktrees, Predicate<String> deliverable) {
FleetConfig.Profile wcfg = new FleetConfig.Profile(
"ltms-local", workerBaseUrl, "coder", null, "FLEETD_WORKER_TOKEN", null,
placement, "fleet", "worker: {profile} #{n}", null, null, null);
@@ -84,7 +90,8 @@ class FleetAppTest {
// it directly so the inbox contract holds for those endpoints.
inbox.own("term_a");
MessageService messages = new MessageService(agents, injector, rendezvous, inbox);
app = new FleetApp(herdr, workers, sessions, messages, this.presence, null)
app = new FleetApp(herdr, workers, sessions, messages, this.presence, null,
null, null, id -> this.presence.isPresent(id) || deliverable.test(id))
.build().start("127.0.0.1", 0);
return app.port();
}
@@ -484,6 +491,16 @@ class FleetAppTest {
assertTrue(mapper.readTree(req(port, "GET", "/sessions/term_a/status").body()).get("ready").asBoolean());
}
@Test
void sessionStatusReportsRegisteredLeadAsReady() throws Exception {
Map<String, String> leads = Map.of("term_lead", "terra");
int port = start(new FakeHerdr(), "http://gx00.gw:8000", Set.of("gx00.gw"), "tab",
new GitWorktrees(), leads::containsKey);
JsonNode body = mapper.readTree(req(port, "GET", "/sessions/term_lead/status").body());
assertTrue(body.get("ready").asBoolean());
}
/**
* CB-582: a lead polling {@code GET /sessions/{id}/status} on its normal cadence — not the
* ticket-scoped {@code /tasks/{ticket}} — must also see a worker's open async {@code fleet_ask}
@@ -0,0 +1,99 @@
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");
}
}
@@ -55,12 +55,21 @@ class GitWorktreesTest {
}
private static void git(Path cwd, String... args) throws Exception {
gitOutput(cwd, args);
}
private static String gitOutput(Path cwd, String... args) throws Exception {
List<String> cmd = new java.util.ArrayList<>(List.of("git"));
cmd.addAll(List.of(args));
Process p = new ProcessBuilder(cmd).directory(cwd.toFile()).redirectErrorStream(true).start();
ProcessBuilder pb = new ProcessBuilder(cmd).directory(cwd.toFile()).redirectErrorStream(true);
pb.environment().put("GIT_CONFIG_GLOBAL", "/dev/null");
pb.environment().put("GIT_CONFIG_SYSTEM", "/dev/null");
pb.environment().put("GIT_TERMINAL_PROMPT", "0");
Process p = pb.start();
String out = new String(p.getInputStream().readAllBytes());
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git timed out: " + String.join(" ", cmd));
assertEquals(0, p.exitValue(), "git " + String.join(" ", args) + " failed:\n" + out);
return out;
}
/** Pending changes to {@code file} in {@code cwd}, empty when git considers it unmodified. */
@@ -205,6 +214,144 @@ class GitWorktreesTest {
"expected an explicitly empty server map, got:\n" + body);
}
/**
* The worktree command shares the primary checkout's config, so this checks the URL git actually
* reads after {@link GitWorktrees#add}, rather than checking only a URL formatting helper.
*/
@Test
void aProvisionedWorktreeUsesACleanHttpsOrigin(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
git(repo, "remote", "add", "origin", "https://synthetic-test-token@git.ltms.dev/akb/kb.git");
String wt = new GitWorktrees(tmp.resolve("wts").toString())
.add(repo.toString(), "fleetd-157-safe-origin", "HEAD");
Path worktree = Path.of(wt);
String origin = gitOutput(worktree, "config", "--get", "remote.origin.url").trim();
assertEquals("https://git.ltms.dev/akb/kb.git", origin);
assertFalse(origin.contains("synthetic-test-token"), "provisioned worktree kept user info");
assertFalse(gitOutput(worktree, "remote", "-v").contains("synthetic-test-token"),
"git remote -v exposed user info");
assertFalse(gitOutput(worktree, "config", "--list").contains("synthetic-test-token"),
"git config --list exposed user info");
String helper = gitOutput(worktree, "config", "--worktree", "--get", "credential.helper");
assertTrue(helper.contains("WORKER_GITEA_TOKEN"), "credential helper does not read the member environment");
assertFalse(helper.contains("synthetic-test-token"), "credential helper stored user info");
}
@Test
void worktreeCredentialHelperCompletesWithoutUsingAnInheritedHelper(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
git(repo, "remote", "add", "origin", "https://git.ltms.dev/akb/kb.git");
String wt = new GitWorktrees(tmp.resolve("wts").toString())
.add(repo.toString(), "fleetd-157-helper", "HEAD");
Path globalConfig = tmp.resolve("global.gitconfig");
Files.writeString(globalConfig, """
[credential]
helper = !f() { printf 'username=%s\\npassword=%s\\n\\n' operator operator-secret; }; f
""");
ProcessBuilder pb = new ProcessBuilder("git", "credential", "fill")
.directory(Path.of(wt).toFile()).redirectErrorStream(true);
pb.environment().put("GIT_CONFIG_GLOBAL", globalConfig.toString());
pb.environment().put("GIT_CONFIG_SYSTEM", "/dev/null");
pb.environment().put("GIT_TERMINAL_PROMPT", "0");
pb.environment().put("WORKER_GITEA_TOKEN", "synthetic-worker-value");
Process p = pb.start();
p.getOutputStream().write("protocol=https\nhost=git.ltms.dev\n\n".getBytes(StandardCharsets.UTF_8));
p.getOutputStream().close();
String credential = new String(p.getInputStream().readAllBytes(), StandardCharsets.UTF_8);
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git credential fill timed out");
assertEquals(0, p.exitValue(), "git credential fill failed");
assertTrue(credential.contains("username=git"), "helper did not return its fixed username");
assertTrue(credential.contains("password=synthetic-worker-value"),
"helper did not return the worker token as the password");
assertFalse(credential.contains("operator-secret"), "Git used the inherited global helper");
}
/**
* fleetd #157 follow-up. An SSH origin never consults {@code credential.helper} — the environment
* credential helper set by {@link GitWorktrees#add} is therefore useless when a member's origin is
* SSH, which is exactly this repo's shape. The worktree must instead get a worktree-scoped
* {@code url.<https>.insteadOf <ssh>} rewrite so both fetch and push resolve to HTTPS, while the
* parent checkout — sharing the same repo-level origin config — must resolve the original SSH URL
* completely unchanged. The host/port here (a synthetic {@code forge.example.test:2222}, not
* {@code git.ltms.dev}) proves the rewrite is derived from the origin, not a hardcoded constant.
*/
@Test
void aProvisionedWorktreeRewritesAnSshOriginToHttpsWorktreeScoped(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
git(repo, "remote", "add", "origin", "ssh://git@forge.example.test:2222/acme/proj.git");
String wt = new GitWorktrees(tmp.resolve("wts").toString())
.add(repo.toString(), "fleetd-157-ssh-rewrite", "HEAD");
Path worktree = Path.of(wt);
assertEquals("https://forge.example.test/acme/proj.git",
gitOutput(worktree, "remote", "get-url", "origin").trim(),
"worktree fetch URL was not rewritten to HTTPS");
assertEquals("https://forge.example.test/acme/proj.git",
gitOutput(worktree, "remote", "get-url", "--push", "origin").trim(),
"worktree push URL was not rewritten to HTTPS");
// The raw config value is unchanged — only the resolved URL is rewritten, via insteadOf.
assertEquals("ssh://git@forge.example.test:2222/acme/proj.git",
gitOutput(worktree, "config", "--get", "remote.origin.url").trim());
assertEquals("ssh://git@forge.example.test:2222/acme/proj.git",
gitOutput(repo, "remote", "get-url", "origin").trim(),
"the parent checkout's fetch URL must be untouched");
assertEquals("ssh://git@forge.example.test:2222/acme/proj.git",
gitOutput(repo, "remote", "get-url", "--push", "origin").trim(),
"the parent checkout's push URL must be untouched");
}
/** An origin already on HTTPS is left alone — the environment credential helper already covers it. */
@Test
void aProvisionedWorktreeLeavesAnHttpsOriginAlone(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
git(repo, "remote", "add", "origin", "https://git.ltms.dev/akb/kb.git");
String wt = new GitWorktrees(tmp.resolve("wts").toString())
.add(repo.toString(), "fleetd-157-https-noop", "HEAD");
Path worktree = Path.of(wt);
assertEquals("https://git.ltms.dev/akb/kb.git",
gitOutput(worktree, "remote", "get-url", "origin").trim());
assertEquals(1, exitCode("git", "-C", wt, "config", "--worktree", "--get-regexp", "^url\\."),
"no url.*.insteadOf rewrite should be added for an already-HTTPS origin");
}
/** Test-local exit-code probe, mirroring {@link GitWorktrees#exitCode} for an assertion the
* production class does not expose. */
private static int exitCode(String... command) throws Exception {
ProcessBuilder pb = new ProcessBuilder(command).redirectErrorStream(true);
pb.environment().put("GIT_CONFIG_GLOBAL", "/dev/null");
pb.environment().put("GIT_CONFIG_SYSTEM", "/dev/null");
pb.environment().put("GIT_TERMINAL_PROMPT", "0");
Process p = pb.start();
p.getInputStream().readAllBytes();
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "command timed out: " + String.join(" ", command));
return p.exitValue();
}
@Test
void provisioningRefusesAWorktreeWhoseOriginStillHasHttpsUserInfo(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
git(repo, "remote", "add", "origin", "https://git.ltms.dev/akb/kb.git");
GitWorktrees worktrees = new GitWorktrees(tmp.resolve("wts").toString(), worktreePath -> {
try {
git(Path.of(worktreePath), "remote", "set-url", "origin",
"https://synthetic-test-token@git.ltms.dev/akb/kb.git");
} catch (Exception e) {
throw new RuntimeException(e);
}
});
WorktreeException error = assertThrows(WorktreeException.class,
() -> worktrees.add(repo.toString(), "fleetd-157-refuse-origin", "HEAD"));
assertEquals("worktree origin contains HTTPS user info; refusing provision", error.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 {