Compare commits

...

9 Commits

Author SHA1 Message Date
Dai Ha 0373b6c41b #213: fix ZDOTDIR credential scrub gate + directory under memberHerdrSocket
CI / contract (pull_request) Successful in 1m10s
CI / build (pull_request) Successful in 2m17s
The memberCredentials.policy: allow-list ZDOTDIR scrub decided zsh-vs-not
using fleetd's own process $SHELL and wrote the generated scrub dir into
fleetd's own java.io.tmpdir. Under memberHerdrSocket: (member panes run as
a different OS user than fleetd's own process) this silently protects
nothing: the wrong shell decides the gate, and the directory can be
unreachable to the member.

- New FleetConfig.memberLoginShell: the member OS user's login shell,
  only ever read when memberHerdrSocket: is configured; fleetd's own
  $SHELL keeps deciding everything when memberHerdrSocket: is absent
  (byte-identical to before).
- HerdrPeerLauncher.applyEnvironmentAllowListPolicy: memberHerdrSocket +
  memberLoginShell not configured/non-zsh falls back to the CB-596
  sentinel overlay (warn loudly, never refuse to spawn). memberHerdrSocket
  + zsh memberLoginShell generates the ZDOTDIR under the configured
  worktreeRoot instead of java.io.tmpdir, shared with the existing
  worktreeGroup (reused, not a new key).
- EnvAllowListScrub: new generate(parentDir, allowedNames, group) overload
  shares the generated directory via pure-Java POSIX group ownership
  (rwxr-x--- dir, rw-r----- files) — no external process spawn.
- fleetd.example.yaml documents memberLoginShell: and worktreeGroup:'s
  reuse for the scrub directory (the live fleetd.yaml is gitignored).

4 new tests in HerdrPeerLauncherAllowListWiringTest cover the acceptance
criteria; 3 of the 4 were watched failing against the pre-fix code.
2026-09-01 11:07:11 +07:00
Dai Ha 966c58a3b8 #137: don't hand one stranded reply to every open async ticket (#215)
CI / contract (push) Successful in 1m31s
CI / build (push) Successful in 1m41s
abandon() drained the target's stranded reply once and then reused that same
Reply for every open task it walked past. Two open tickets on one target
therefore both came back REPLIED with the same text — one of them a reply the
worker never gave for that delegation.

One worker answer can settle at most one delegation. It now goes to the oldest
open task (lowest createdNanos) and every other open task keeps the ordinary
WORKER_FAILED path. If the chosen task turns out to be already resolved by
another path, the drained reply is published back to the inbox instead of being
dropped.

reply()'s matching side gets the same rule: more than one candidate means the
reply goes to the inbox rather than to a guess.

Reachability, checked rather than assumed:
- matching.size() >= 2 alone is reachable and was a real bug before this change.
- More than one candidate in reply() is not reachable today — hasAsyncQuestion
  matches any task with a stamped turnId, and answer()'s
  clearAsyncQuestion(turnId, false) leaves that stamp until the resumed turn
  resolves. Kept as defensive code, documented, no test seam added.
- matching.size() >= 2 together with a live strand is not reachable either:
  send() and answer() are the only two lock holders and both open a Rendezvous
  waiter inside the lock, so "lock held" and "waiter open" are one fact, and an
  acceptance always clears the strand first.

I checked the last point by building it in a scratch worktree: it can be forced
by closing an accepted send's waiter directly through Rendezvous, and it does go
red against the pre-fix loop — but that breaks the lock-and-waiter invariant
from outside the class, so no such test is added. A note in abandon()'s javadoc
says so, to save the next reader the same round trip.

Worker's REPORT-cb137.md left out of main.

mvn clean install: Tests run: 1071, Failures: 0, Errors: 0, Skipped: 0
2026-08-31 23:23:07 +07:00
Dai Ha 735c837604 #209: resolve agentSessionId lazily via the retained PeerHandle
CI / contract (push) Successful in 48s
CI / build (push) Successful in 2m18s
SessionManager called handle.agentSessionId() once at spawn and froze it in the
immutable MemberSession. For opencode that value is always null - the session
row does not exist yet when the pane is created - so fleet_list never reported
an agentSessionId and fleet_spawn{resumeSessionId} was unusable for that
backend. The handle's javadoc said 'the caller re-calls later'; no caller did,
and SessionManager did not even retain the handle.

Retain the PeerHandle per pane and re-resolve while the stored id is still
null, sticky once found, CAS-swapped into the registry.

roster() stays non-resolving - it is the roster supplier for LeadHeartbeatLoop
and FleetHealthMonitor, and resolving there would open opencode's 841MB SQLite
database on every tick, for every unresolved member, forever. rosterResolved()
carries the resolve and is used only by fleet_list and the REST roster, the two
surfaces that report the id. get(paneId) and release() resolve too, both
caller-driven.

A test pins the split: the plain roster() must never call agentSessionId()
again.
2026-08-31 22:21:16 +07:00
Dai Ha 457dc0330d #185 stage 2: the credential gap detector must not report on the wrong environment
CI / contract (push) Successful in 55s
CI / build (push) Successful in 1m59s
Under memberHerdrSocket: member panes run as a different OS user, so
hostEnvNames (fleetd's own environment) no longer describes what a member
inherits. logCredentialGap now reports "unknown, not clean" in that mode
instead of its usual UNBLOCKED/blanked conclusions, once per launcher, naming
the config key and scoping its count to fleetd's own process.

Byte-identical when memberHerdrSocket: is absent, pinned by a test.
2026-08-31 22:17:40 +07:00
Dai Ha a1a9015217 #199: /members returns rows under "members" 2026-08-31 22:17:35 +07:00
Dai Ha 5a811a3695 #199: GET /members returns its rows under "members", with "workers" kept as an alias
The endpoint became /members in the CB-634 rename but the body key stayed
"workers", so a caller that read "members" saw an empty fleet and reported
no members at all.

Emit the canonical "members" key. Keep "workers" as a deprecated alias so an
existing REST consumer keeps working - the out-of-band path a lead falls back
to when its MCP mount drops reads this endpoint.

The test now pins both keys and asserts they carry the same rows, so the alias
cannot silently drift. Watched failing without the fix:
"GET /members must return its rows under \"members\" ==> expected: not <null>"
2026-08-31 22:17:22 +07:00
Dai Ha 1fdaa74eb3 fleetd #209 follow-up: keep roster() off the resolve path
CI / build (pull_request) Successful in 1m20s
CI / contract (pull_request) Successful in 1m21s
roster() is the supplier for LeadHeartbeatLoop and FleetHealthMonitor
(both timer-driven) and for placement/exhaustion checks and the
metrics scrape — none of which read agentSessionId. Resolving there
meant every tick could open a lazy-resolving adapter's (opencode's)
on-disk session database once per member whose id was still unknown,
with no bound: a member whose id never appears would pay that cost for
the life of the process.

roster() goes back to its pre-#209 behavior (no resolve, no I/O). A
new rosterResolved() carries the resolve logic, and is used only by
the two surfaces that actually report agentSessionId to a caller:
fleet_list (FleetMcp.listFleet) and the REST roster
(FleetApp.listMembers). fleet_whoami's roster().stream() at
FleetMcp.java:768 does not surface the field, so it stays on the plain
roster(). get(paneId) (fleet_status) and the release() resolve are
caller-driven, not timers, and are unchanged.

Retargeted the roster-facing tests from #209 at rosterResolved(), and
added plainRosterDoesNotResolveAgentSessionId, which pins the split by
asserting the handle's agentSessionId() is not called again by
roster().
2026-08-31 22:15:16 +07:00
Dai Ha 847e8bd3fa #185 stage 2: stop the credential-gap detector reporting on the wrong environment
CI / contract (pull_request) Successful in 54s
CI / build (pull_request) Successful in 1m39s
When memberHerdrSocket: is configured, member panes run under a different OS
user than fleetd's own process, so hostEnvNames (fleetd's own environment)
no longer describes what a member pane inherits. logCredentialGap now checks
for that config key and, when set, logs a single WARN saying the gap is
UNKNOWN (not clean) and names the key, instead of printing the "inherits
them UNBLOCKED" / "the scrub blanks them" conclusions as fact. Behaviour is
byte-identical when memberHerdrSocket is absent (the default and only mode
this host runs).
2026-08-31 22:12:56 +07:00
Dai Ha f9fb387427 fleetd #209: re-poll agentSessionId against the retained PeerHandle
CI / contract (pull_request) Successful in 45s
CI / build (pull_request) Successful in 1m40s
SessionManager used to call PeerHandle.agentSessionId() exactly once at
spawn and freeze the answer into the immutable MemberSession. For
opencode that call always came back null, because opencode has not
written its on-disk session row yet when the pane is created, and no
caller ever re-asked the handle — it went out of scope at the end of
the spawn method. fleet_list therefore never reported agentSessionId
for an opencode member, and resumeSessionId was unusable for it.

Retain each spawn's PeerHandle in SessionManager, keyed by paneId, and
re-resolve a still-null agentSessionId against it from roster(), get(),
and release() (so a released member's detail also carries a
late-resolved id). Resolution is bounded: only sessions with a still-
null id do any work, a resolved id is never looked up again, and a
throwing handle degrades to "unresolved" rather than breaking the
caller. MemberSession gains a withAgentSessionId wither in the same
style as withState/withActivity.
2026-08-31 22:04:20 +07:00
13 changed files with 1062 additions and 184 deletions
-166
View File
@@ -1,166 +0,0 @@
# CB-137 / fleetd issue #137 — report
## Real root cause (not the hypothesis in the ticket)
I read `MessageService.java` and `Rendezvous.java` before changing anything. The mechanism is real,
but the exact place it happens is `MessageService.answer()`, not "the reply goes to the inbox on
purpose" in general.
1. A lead delegates with `fleet_send{wait:false}` → `sendAsync()` creates a `Task` and runs `send()`
on a background virtual thread with a 30-minute internal budget (`ASYNC_TIMEOUT_MS`).
2. The worker calls `fleet_ask`. That resolves the open rendezvous waiter with `Kind.QUESTION`, so
`send()` returns immediately and the `Task` is left open (its `future` stays unresolved — see the
comment in `sendAsync`'s lambda: "Keep the accepted owner until answer() finishes it").
3. The lead answers with `fleet_send{turnId, content}`. This calls `FleetMcp.answer()` →
`MessageService.answer(turnId, content, timeout)`. The `timeout` here is **not** the generous
30-minute async budget — it is the MCP tool's own bounded wait: `DEFAULT_TIMEOUT_MS = 25_000`,
clamped to at most `MAX_TIMEOUT_MS = 120_000` (`FleetMcp.java:71-72,495,512`). This is the same
~60–120s window every blocking `fleet_send` call is capped at (documented elsewhere as "the
caller's own MCP client call timeout").
4. `answer()` opens a **fresh** rendezvous waiter for the worker session and blocks on it for at most
that window. If the worker's resumed turn takes longer than that to actually finish (very
plausible — the resumed turn can mean more edits, a build, a commit, a push, opening a PR), the
wait times out. On timeout, `answer()`'s `finally` block unconditionally calls
`rendezvous.close(workerSession, reply)`, **removing the waiter from the map**, and returns
`Outcome.TIMED_OUT_WORKING` to the lead.
5. The worker keeps working, unaware anything happened, and eventually calls `fleet_reply`. That
reaches `MessageService.reply(session, content)`, which tries `rendezvous.resolve(session,
content)` — but the waiter was already closed in step 4, so `resolve` returns `false`. `reply()`
then falls back to `inbox.publish(...)` and marks `strandedReplies.put(session, true)`
(CB-640 bookkeeping) — the reply is safely held, but **the async `Task`'s `future` is never
completed**.
6. `fleet_poll{ticket}` keeps returning `PENDING` forever (the `Task` never resolves) — until the
lead eventually calls `fleet_stop`. That fires `sessions.onRelease` → `messages.abandon(target,
reason)` (`Fleetd.java:481-497`), where `reason` is built with the exact text from the bug report
("the worker session was released before it replied; worktree=... branch=... snapshot=...",
`Fleetd.java:484-487`). `abandon()`'s loop finds the still-open `Task` (`question == null`, future
not done) and completes it as `WORKER_FAILED` with that misleading reason — even though the
worker's real reply is sitting, intact, in the inbox the whole time.
So: the reported behaviour is correct, and the specific trigger is `answer()`'s own bounded wait
being shorter than the worker's real resumed-turn time — not anything to do with the ~55s
`fleet_ask` window itself (that part, issue #61, is untouched).
## Fix
Two changes in `fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java`, both scoped to the
ticket/reply routing and the terminal-state text — `fleet_ask`'s own window and mechanics are
untouched.
**1. `reply()` — priority 1 (the ticket resolves with the real reply).**
Before falling back to the inbox, `reply()` now looks for an async `Task` that is specifically in the
"already answered but not yet resolved" state (`question == null`, `turnId != null` — set once
`answer()` has cleared the question but before anything completed the future, `!future.isDone()`).
If one exists for this `target`, the worker's reply completes that `Task`'s future directly as
`Outcome.REPLIED` with the real content, and the reply never touches the inbox at all. A task that
was never asked has `turnId == null` and can never match, so ordinary (no-`fleet_ask`) delegations
are unaffected — they already resolve through the pre-existing rendezvous fast path.
I chose this over leaving `answer()`'s own timeout behaviour untouched and instead keeping its
rendezvous waiter open in the background: that alternative works but reopens the "at most one
waiter per session" invariant (`Rendezvous.open` throws on a double-open) to a new class of races
with a fresh send arriving mid-window. The `send()` path already guards against sending into an
answered-but-still-resolving worker via `hasAsyncQuestion(target)` (checks `asyncTasksByTurn`,
which still holds the task until it resolves), so routing through `reply()` gets the same protection
without touching `answer()`'s waiter lifecycle at all — the smaller, safer diff.
**2. `abandon()` — priority 3 (required independently, "even if you fix (1)").**
Before marking any of a released target's still-open tasks `WORKER_FAILED`, `abandon()` now checks
`hasStrandedReply(target)` (the existing CB-640 fact — true whenever the *last* `reply()` for this
target fell through to the inbox). If true, it drains the inbox (`recoverStrandedReply`) and — if it
actually finds a message — completes the task as `REPLIED` with that real content instead of writing
the failure. This is deliberately a **separate** check from fix 1: fix 1 already prevents the
inbox-stranding from happening in the exact scenario this ticket describes, so by the time
`abandon()` runs the task is normally already resolved and `abandon()`'s `complete()` call is a
harmless no-op. This second check exists so that if some *other* future path ever strands a reply
in the inbox without resolving its ticket, `abandon()` still refuses to report a false failure —
"if a reply reached any sink for that turn, the terminal state is done," per the ticket. I verified
both are required by disabling each independently and confirming the two new tests fail (see below).
**Priority 4 (the snapshot/worktree hint).** Handled as a consequence of both fixes rather than a
separate branch: once a task resolves as `REPLIED` (via either fix), `abandon()` never calls
`new Reply(Outcome.WORKER_FAILED, reason)` for that task at all, so the "the worker session was
released before it replied; worktree=... branch=... snapshot=..." text is never constructed or
attached to that ticket's outcome. It still appears, correctly, for a task that never got a reply
(the existing `abandonFailsEveryPendingAsyncTicketForTheReleasedTarget` /
`anAbandonedAsyncTaskPollsAsFailedNotPending` tests still pass unchanged).
**Priority 2** was not needed — fix 1 makes `fleet_poll{ticket}` return the actual reply (the
higher-priority option), so I did not fall back to "the ticket merely resolves as done with no
content."
## Tests — driven through the real delegation path, not the reply sink directly
Both new tests in `fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java` go through
`sendAsync` → `injectDelivery` → `ask` → `answer` (with a short timeout, so it genuinely times out,
mirroring the ~25–120s real MCP-call bound vs. a longer resumed turn) → `reply` → `poll`/`abandon`.
No test constructs a `Reply` and hands it to a sink directly.
- `aReplyAfterAnswerTimesOutStillCompletesTheAsyncTicket` — asserts `fleet_poll{ticket}` (via
`messages.poll`) reaches `Phase.DONE` with the worker's actual reply text and
`replySource() == "reply"`, and that `hasStrandedReply(T)` stays `false` (proves the reply never
touched the inbox at all — fix 1 caught it).
- `fleetStopAfterAnOrphanedReplyDoesNotFailTheTicket` — same setup, then calls `abandon(T, "the
worker session was released before it replied")` (what `fleet_stop` triggers) and asserts it
returns `false` (no failure recorded) and the ticket still polls `DONE` with the real reply.
**Proof both fail without the change.** I temporarily short-circuited both new private methods
(`askAnsweredAsyncTask` → always `null`, `recoverStrandedReply` → always `null`) — i.e. disabled
both fixes — and ran just these two tests:
```
[ERROR] Tests run: 2, Failures: 2, Errors: 0, Skipped: 0
dev.ltms.fleet.msg.MessageServiceTest.aReplyAfterAnswerTimesOutStillCompletesTheAsyncTicket
org.opentest4j.AssertionFailedError: expected: <DONE> but was: <PENDING>
dev.ltms.fleet.msg.MessageServiceTest.fleetStopAfterAnOrphanedReplyDoesNotFailTheTicket
org.opentest4j.AssertionFailedError: a reply already arrived, so nothing here is a genuine failure
==> expected: <false> but was: <true>
```
This is the exact bug: the ticket stays `PENDING` forever, and `abandon()` reports `true` (a
failure) even though a reply had already arrived. I then restored both fixes (verified with
`grep -n "TEMP #137-proof"` finding nothing) and re-ran — both pass.
## Build
Ran from `fleetd/`, unpiped, full output read (not `| tail`):
```
mvn clean install
...
[INFO] Tests run: 1039, Failures: 0, Errors: 0, Skipped: 0
[INFO] BUILD SUCCESS
[INFO] Total time: 36.315 s
```
Main was at 1037 tests; this branch adds the 2 new tests above → 1039, all green, `exit=0`.
## What I could NOT check
- No IDE tooling is mounted for me (worker), so no `ide_diagnostics`/IntelliJ inspection pass — only
`mvn clean install` (compiler + full test suite), as the worker procedure allows.
- I cannot restart the daemon or dogfood this live — I have no forge/daemon control. This is
unverified against a real herdr pane, a real MCP client's ~60s call cap, or a real worker session;
everything above is verified only through the JUnit fixture's simulated timing
(`FakeHerdr`/`injector.onStatus`/direct `messages.answer(...,150)` calls), not a live fleet.
A primary should still consider a short live dogfood (an async delegation that asks, gets answered,
and takes longer than ~2 minutes to reply) before calling this closed.
- I did not touch, and did not re-verify, the `fleet_ask` ~55s window itself (issue #61) — out of
scope per the brief.
## Scope note (not investigated further)
`answer()`'s nested/double-`fleet_ask` case (the worker asks a second question before ever
replying to the first answer) has some pre-existing behaviour around which `turnId` a `QUESTION`
resolution gets attributed to that I did not fully untangle — it predates this change, my fix does
not touch it, and it is unrelated to the reported defect. Flagging only; not investigated further.
## Handoff
- Branch: `worker/cb-137-ask-ticket-e7760c-2`
- Worktree root: `/Users/dai.ha/LTMS/.bridged-worktrees/734324-2`
- Files changed:
- `fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java`
- `fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java`
- `REPORT-cb137.md` (this file)
- Build: `Tests run: 1039, Failures: 0, Errors: 0, Skipped: 0` / `BUILD SUCCESS` (verbatim above)
+16
View File
@@ -117,6 +117,15 @@ 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
# fleetd #213: the login shell the member OS user (memberHerdrSocket above) actually runs. ONLY
# read when memberHerdrSocket is set — fleetd's own $SHELL says nothing about a pane running
# under a different OS user, and there is no channel to ask herdr for that user's shell, so this
# must be told rather than guessed. Absent, blank, or anything not ending in "zsh" is treated the
# same as "not zsh": the memberCredentials.policy: allow-list ZDOTDIR scrub (see worktreeGroup
# below) is skipped in favour of the weaker CB-596 sentinel overlay — a degraded control, never a
# refusal to spawn. When memberHerdrSocket is absent this key is never consulted at all.
# memberLoginShell: /bin/zsh
# 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
@@ -643,6 +652,13 @@ guard:
# CAUTION: this isolates credentials, not the repository — a member in the group can still
# write the operator's git objects and refs in the shared repo. The operator running fleetd
# must already be a member of the named group, or every provisioning spawn fails loudly.
#
# fleetd #213: this is also the ONE group the memberCredentials.policy: allow-list ZDOTDIR scrub
# reuses when memberHerdrSocket is set — deliberately not a second config key. Under
# memberHerdrSocket, the scrub directory is generated under worktreeRoot (never java.io.tmpdir,
# which the member OS user cannot reach) and shared read-only with this group. If worktreeGroup
# is unset while memberHerdrSocket is set, the scrub cannot be guaranteed reachable by the member,
# so fleetd falls back to the weaker CB-596 sentinel overlay instead (a WARN names the gap).
# worktreeGroup: fleet-workers
# Session lifecycle limits (CB-303). All knobs are opt-in; omit or set to null to keep
@@ -84,6 +84,18 @@ import java.util.Set;
* <strong>This isolates credentials, not the repository</strong>: a member in
* the group can still write the operator's git objects and refs in the shared
* repo. See {@link dev.ltms.fleet.session.Worktrees#shareWithGroup}.
* @param memberLoginShell fleetd #213: the login shell the member's OS user actually runs, ONLY
* meaningful (and only ever read) when {@code memberHerdrSocket} is
* configured — that mode spawns member panes under a different OS user than
* fleetd's own process, so fleetd's own {@code $SHELL} says nothing about what
* that pane runs. There is no channel to ask herdr for another user's shell, so
* this must be told, never guessed. {@code null}/blank (or a value not ending
* in {@code zsh}) is treated the same as "not zsh": the {@code
* memberCredentials.policy: allow-list} ZDOTDIR scrub is skipped in favour of
* the CB-596 sentinel overlay — a degraded control, never a refusal to spawn.
* When {@code memberHerdrSocket} is NOT configured this field is never
* consulted at all; fleetd keeps reading its own {@code $SHELL}, exactly as
* before this field existed.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record FleetConfig(
@@ -107,7 +119,20 @@ public record FleetConfig(
Integer quarantineCooldownSeconds,
MemberCredentials memberCredentials,
Coordinator coordinator,
String worktreeGroup) {
String worktreeGroup,
String memberLoginShell) {
/** Back-compat form before the {@code memberLoginShell} key was added. */
public FleetConfig(Bind bind, String herdrSocket, String memberHerdrSocket, Map<String, Profile> profiles,
Guard guard, String worktreeRoot, Lifecycle lifecycle, Integer spawnReadyTimeoutMs,
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
LeadHeartbeat leadHeartbeat, Health health, String placement, Auth auth,
ConfigReload configReload, Integer quarantineCooldownSeconds,
MemberCredentials memberCredentials, Coordinator coordinator, String worktreeGroup) {
this(bind, herdrSocket, memberHerdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth,
configReload, quarantineCooldownSeconds, memberCredentials, coordinator, worktreeGroup, null);
}
/** Back-compat form before the {@code worktreeGroup} key was added. */
public FleetConfig(Bind bind, String herdrSocket, String memberHerdrSocket, Map<String, Profile> profiles,
@@ -118,7 +143,7 @@ public record FleetConfig(
MemberCredentials memberCredentials, Coordinator coordinator) {
this(bind, herdrSocket, memberHerdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth,
configReload, quarantineCooldownSeconds, memberCredentials, coordinator, null);
configReload, quarantineCooldownSeconds, memberCredentials, coordinator, null, null);
}
/** Back-compat form before the {@code coordinator:} block was added. */
@@ -1346,7 +1371,7 @@ public record FleetConfig(
"bind", "herdrSocket", "memberHerdrSocket", "profiles", "guard", "worktreeRoot",
"lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs", "broker", "primary", "fleet",
"leadHeartbeat", "health", "placement", "auth", "configReload", "quarantineCooldownSeconds",
"memberCredentials", "coordinator", "worktreeGroup");
"memberCredentials", "coordinator", "worktreeGroup", "memberLoginShell");
/** Load and validate config from {@code path}. */
public static FleetConfig load(Path path) {
@@ -1966,9 +1991,12 @@ public record FleetConfig(
// and this ticket's Coordinator is config-only anyway (nothing yet reads it at startup).
// worktreeGroup is left as-is (fleetd #185 stage 3): null/blank is "off", and there is no
// sane non-null default — an OS group name is operator-specific.
// memberLoginShell is left as-is (fleetd #213), like worktreeGroup: null/blank is "not
// configured", and there is no sane non-null default — a member's login shell is
// operator-specific and only meaningful when memberHerdrSocket is also set.
return new FleetConfig(b, herdrSocket, memberHerdrSocket, profiles, g, worktreeRoot, l, timeout, pollMs,
broker, primary, f, leadHeartbeat, health, placementOrDefault, a, configReload,
quarantineCooldown, mc, coordinator, worktreeGroup);
quarantineCooldown, mc, coordinator, worktreeGroup, memberLoginShell);
}
/**
@@ -951,7 +951,10 @@ public final class FleetMcp {
.sorted(Map.Entry.comparingByValue())
.map(e -> leadView(e.getKey(), e.getValue(), live.get(e.getKey()), selfTerm))
.toList();
List<MemberSession> roster = sessions.roster();
// fleetd #209: this is the caller-driven fleet_list read that actually reports
// agentSessionId (via memberCapacityView -> SessionManager.rosterView), so it uses the
// resolving roster; the heartbeat/health/metrics timers stay on the plain sessions.roster().
List<MemberSession> roster = sessions.rosterResolved();
List<Map<String, Object>> out = roster.stream()
.map(s -> memberCapacityView(s, live.get(s.terminalId()), messages, capacity.clock().getAsLong()))
.toList();
@@ -7,6 +7,9 @@ import java.io.IOException;
import java.io.UncheckedIOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.attribute.GroupPrincipal;
import java.nio.file.attribute.PosixFileAttributeView;
import java.nio.file.attribute.PosixFilePermissions;
import java.time.Duration;
import java.time.Instant;
import java.util.ArrayList;
@@ -116,6 +119,70 @@ public final class EnvAllowListScrub {
}
}
/**
* fleetd #213: as {@link #generate(Path, Set)}, plus share the generated directory with
* {@code group} — the member OS user's group (the operator's existing {@code worktreeGroup:}
* name, reused rather than inventing a second one) — so a member running under a different OS
* user than fleetd's own process can still read what it needs from a directory placed outside
* {@code java.io.tmpdir}. {@code group} null/blank ⇒ identical to {@link #generate(Path, Set)};
* this is the single-daemon (no {@code memberHerdrSocket}) shape, where the pane is fleetd's own
* uid and no group sharing is needed.
*
* @throws UncheckedIOException also when {@code group} does not resolve on this host, or a
* group-ownership/permission call is refused — the same "fail
* loudly rather than start unprotected" contract as above: a scrub
* the configured member user cannot even read is not a working
* control.
*/
public static Path generate(Path parentDir, Set<String> allowedNames, String group) {
Path dir = generate(parentDir, allowedNames);
if (group != null && !group.isBlank()) {
shareWithGroup(dir, group);
}
return dir;
}
/**
* chgrp/chmod-equivalent over the freshly generated directory and the startup files already
* written into it: owner keeps full access, {@code group} gets traverse+read on the directory
* ({@code rwxr-x---}, so a login shell under that group can find and source the files) and
* read-only on each file ({@code rw-r-----}) — deliberately no group WRITE anywhere, since a
* member never needs to add or change fleetd's own generated scrub. (The scrub script's own
* report write inside the pane consequently fails closed rather than open — see {@code
* scrub.zsh}'s trailing {@code 2>/dev/null} — which {@link
* dev.ltms.fleet.member.HerdrPeerLauncher#releaseZdotdir} already treats as "cannot be
* confirmed to have run" rather than success.)
*/
private static void shareWithGroup(Path dir, String group) {
try {
GroupPrincipal principal = dir.getFileSystem().getUserPrincipalLookupService()
.lookupPrincipalByGroupName(group);
setGroupAndPermissions(dir, principal, "rwxr-x---");
try (Stream<Path> entries = Files.list(dir)) {
for (Path file : entries.toList()) {
setGroupAndPermissions(file, principal, "rw-r-----");
}
}
} catch (IOException e) {
throw new UncheckedIOException("cannot share generated ZDOTDIR " + dir + " with group '"
+ group + "' — the group must exist, and the fleetd operator ("
+ System.getProperty("user.name") + ") must be a member of it", e);
} catch (UnsupportedOperationException e) {
throw new UncheckedIOException("cannot share generated ZDOTDIR " + dir + " with group '"
+ group + "' — this filesystem does not support POSIX group ownership",
new IOException(e));
}
}
private static void setGroupAndPermissions(Path path, GroupPrincipal group, String perms) throws IOException {
PosixFileAttributeView view = Files.getFileAttributeView(path, PosixFileAttributeView.class);
if (view == null) {
throw new IOException("POSIX file attributes are not supported for " + path);
}
view.setGroup(group);
Files.setPosixFilePermissions(path, PosixFilePermissions.fromString(perms));
}
/** One operator-sourcing startup file: source the {@code $HOME} counterpart, change nothing else. */
private static String homeSourcingFile(String name) {
return """
@@ -112,6 +112,12 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
* daemon's own process is started the same way (a login shell sourcing the same secret store —
* see CB-592's investigation of {@code secrets.sh}), so on a single-host deployment its env
* mirrors what the pane's login shell is about to export.
*
* <p>fleetd #185 stage 2: that mirroring assumption holds only while the member pane runs under
* the SAME OS user as the daemon. When {@code memberHerdrSocket:} is configured, member panes
* run on a second herdr owned by a different user — different {@code $HOME}, different {@code
* secrets.sh}, different environment entirely — so this field's data no longer describes what a
* member pane inherits. See {@link #logCredentialGap} for how that mode is handled.
*/
private final Supplier<Set<String>> hostEnvNames;
@@ -1071,6 +1077,29 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
* 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.
*
* <p>fleetd #213: which shell decides the zsh gate, and where the generated directory lives,
* both depend on whether {@code memberHerdrSocket:} is configured — see {@link
* #memberHerdrSocketConfigured()}'s javadoc for why fleetd's own {@code $SHELL} and {@code
* java.io.tmpdir} describe the wrong process once member panes run under a different OS user.
* <ul>
* <li>{@code memberHerdrSocket} ABSENT (today's only mode): byte-identical to before this
* fix — fleetd's own {@code $SHELL} decides zsh, and the directory is generated under
* {@code java.io.tmpdir}.</li>
* <li>{@code memberHerdrSocket} PRESENT: the configured {@code memberLoginShell:} decides
* zsh instead — fleetd's own {@code $SHELL} is never consulted, since it names a
* different user's shell, not the member's. Absent/non-zsh falls back exactly like the
* non-zsh case below. When it IS zsh, the directory still cannot go under {@code
* java.io.tmpdir} (mode 0700, unreadable by another uid — the exact gap fleetd #213
* exists to close), so it is generated under {@code worktreeRoot} instead and shared
* read-only with {@code worktreeGroup} — the same group {@link
* dev.ltms.fleet.session.Worktrees#shareWithGroup} already uses, reused rather than
* inventing a second group key. Either one missing means the scrub cannot be guaranteed
* reachable by the member, which is the same "cannot guarantee the scrub runs" case as a
* non-zsh shell, so it gets the identical fallback.</li>
* </ul>
* In every branch: never refuse to spawn. A degraded credential control must not become an
* outage for an opt-in feature.
*/
private Path applyEnvironmentAllowListPolicy(FleetConfig.Profile cfg, Launch launch) {
FleetConfig.MemberCredentials creds = memberCredentials == null ? null : memberCredentials.get();
@@ -1078,8 +1107,14 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
return null;
}
Set<String> allowed = derivedAllowedNames(creds, launch);
String loginShell = resolveEnv("SHELL");
boolean zsh = loginShell != null && (loginShell.endsWith("/zsh") || loginShell.equals("zsh"));
boolean memberHerdrSocket = memberHerdrSocketConfigured();
// fleetd #213 defect 1: under memberHerdrSocket the member pane runs as a DIFFERENT OS
// user, so fleetd's own $SHELL says nothing about what that pane runs — resolveEnv("SHELL")
// must not even be called on this path, only the explicit memberLoginShell: config can
// answer it. With memberHerdrSocket absent, nothing here changes: fleetd's own $SHELL is
// still the input, exactly as before this fix.
String loginShell = memberHerdrSocket ? configuredMemberLoginShell() : resolveEnv("SHELL");
boolean zsh = isZshShell(loginShell);
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
@@ -1093,13 +1128,35 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
logCredentialGap(creds, null);
return null;
}
Path parentDir;
String group = null;
if (memberHerdrSocket) {
// fleetd #213 defect 2: java.io.tmpdir is fleetd's own per-user temp dir (mode 0700 on
// macOS) — a member running as a different uid cannot even traverse it, let alone read
// the generated files. worktreeRoot is the only configured location a different-uid
// member can be given access to, and only WITH worktreeGroup to grant that access —
// absent either, the scrub cannot be guaranteed reachable, so this falls back exactly
// like the non-zsh case above rather than generating a directory nothing can read.
parentDir = memberScrubParentDir();
group = memberGroup();
if (parentDir == null || group == null) {
warnCannotShareScrubDirectory();
overlayBlockedCredentials(launch.env(), creds);
logCredentialGap(creds, null);
return null;
}
} else {
parentDir = Path.of(System.getProperty("java.io.tmpdir"));
}
// 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). Same
// reasoning gates logCredentialGap's wording: passing the derived `allowed` set (non-null)
// here, and ONLY here, is what tells it the scrub will really blank an unkept name — #192.
logAllowListCoverage(allowed);
logCredentialGap(creds, allowed);
Path dir = EnvAllowListScrub.generate(Path.of(System.getProperty("java.io.tmpdir")), allowed);
Path dir = memberHerdrSocket
? EnvAllowListScrub.generate(parentDir, allowed, group)
: EnvAllowListScrub.generate(parentDir, allowed);
launch.env().put("ZDOTDIR", dir.toAbsolutePath().toString());
log.info("memberCredentials policy=allow-list: profile={} generated ZDOTDIR {} — derived "
+ "allow-list holds {} name(s); the pane reports allowed N of M at release",
@@ -1107,6 +1164,55 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
return dir;
}
/** True when {@code shell} is a zsh login shell path or bare name — the ZDOTDIR gate. */
private static boolean isZshShell(String shell) {
return shell != null && (shell.endsWith("/zsh") || shell.equals("zsh"));
}
/**
* fleetd #213: the configured {@code memberLoginShell:}, or {@code null} when unconfigured.
* Called ONLY from the {@code memberHerdrSocket}-configured branch of {@link
* #applyEnvironmentAllowListPolicy} — fleetd's own {@code $SHELL} is never read on that path.
* {@link #config} being {@code null} (an older test call site, or a launcher that never
* threaded the full config through) is treated the same as "not configured".
*/
private String configuredMemberLoginShell() {
FleetConfig cfg = config == null ? null : config.get();
return cfg == null ? null : cfg.memberLoginShell();
}
/**
* fleetd #213: {@code worktreeRoot}, as the ZDOTDIR scrub's parent directory under {@code
* memberHerdrSocket}, or {@code null} when unconfigured — the same "cannot guarantee the scrub
* runs" gap as {@link #memberGroup()} being unset (see {@link
* #applyEnvironmentAllowListPolicy}). Deliberately no sibling-of-repo-root default here, unlike
* {@code GitWorktrees}' own {@code worktreeRoot} resolution: that default is a convenience for
* provisioning a worktree that will exist regardless, whereas an unconfigured value here means
* fleetd has no operator-endorsed location to put a credential-bearing directory a different OS
* user must reach, so falling back to the overlay is the honest answer, not a guess.
*/
private Path memberScrubParentDir() {
FleetConfig cfg = config == null ? null : config.get();
if (cfg == null || cfg.worktreeRoot() == null || cfg.worktreeRoot().isBlank()) {
return null;
}
return Path.of(cfg.worktreeRoot());
}
/**
* fleetd #213: the configured {@code worktreeGroup:}, or {@code null} when unset/blank. Reuses
* the group {@link dev.ltms.fleet.session.Worktrees#shareWithGroup} already establishes for
* provisioned worktrees, rather than a second group key — see {@link
* #applyEnvironmentAllowListPolicy}.
*/
private String memberGroup() {
FleetConfig cfg = config == null ? null : config.get();
if (cfg == null || cfg.worktreeGroup() == null || cfg.worktreeGroup().isBlank()) {
return null;
}
return cfg.worktreeGroup();
}
/**
* 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
@@ -1164,6 +1270,30 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
}
}
/** Guards {@link #warnCannotShareScrubDirectory} to one WARN per launcher instance. */
private final AtomicBoolean cannotShareScrubDirWarned = new AtomicBoolean();
/**
* fleetd #213: {@code memberHerdrSocket} is configured and the member login shell IS zsh, but
* {@code worktreeRoot} and/or {@code worktreeGroup} is missing, so the generated ZDOTDIR cannot
* be placed anywhere the member's OS user can reach — {@code java.io.tmpdir} is fleetd's own
* 0700 temp dir, unreadable by another uid, which is the exact gap this ticket exists to close.
* Say so once per launcher instance, instead of either generating a directory nothing can read
* (protection theatre) or refusing to spawn (turning a degraded credential control into an
* outage for an opt-in feature).
*/
private void warnCannotShareScrubDirectory() {
if (cannotShareScrubDirWarned.compareAndSet(false, true)) {
log.warn("memberCredentials policy=allow-list: memberHerdrSocket is configured and the "
+ "member login shell is zsh, but worktreeRoot and/or worktreeGroup is not "
+ "configured — the generated ZDOTDIR cannot be placed where the member's OS "
+ "user can read it (java.io.tmpdir is fleetd's own, unreadable by another uid), "
+ "so the scrub cannot be guaranteed to run. Falling back to the CB-596 sentinel "
+ "overlay. Configure both worktreeRoot and worktreeGroup to enable the "
+ "allow-list scrub under memberHerdrSocket.");
}
}
/**
* CB-633 teardown half: read the pane's scrub report (the denominator report the generated
* scrub wrote) and delete the directory. Called from {@link #stop}, which is the one funnel
@@ -1230,6 +1360,67 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
*/
private final AtomicBoolean allowListGapLogged = new AtomicBoolean();
/**
* fleetd #185 stage 2: guards {@link #warnUnknownMemberEnvironment} to one WARN per launcher
* instance, not one per spawn — the same one-per-instance shape as {@link #unprotectedGapLogged}
* and {@link #allowListGapLogged}, kept as its own flag for the same reason those two are split:
* this mode is orthogonal to which of the other two branches would otherwise have fired.
*/
private final AtomicBoolean unknownMemberEnvironmentWarned = new AtomicBoolean();
/**
* fleetd #185 stage 2: whether {@code memberHerdrSocket:} is configured, i.e. member panes run
* on a second herdr owned by a different OS user than the daemon's own process. Re-read from the
* live config on every call (same hot-reload shape as {@link #memberCredentials}), never cached,
* so a config reload takes effect on the next spawn without a restart.
*
* <p>{@link #config} is {@code null} on any call site that never threaded the full config
* through (every production {@code HerdrPeerLauncher} does; a handful of older tests do not) —
* treated the same as "not configured", which is the correct, permissive default: it is exactly
* today's single-daemon behaviour.
*/
private boolean memberHerdrSocketConfigured() {
if (config == null) {
return false;
}
FleetConfig cfg = config.get();
return cfg != null && cfg.memberHerdrSocket() != null && !cfg.memberHerdrSocket().isBlank();
}
/**
* fleetd #185 stage 2: the single replacement WARN for {@link #logCredentialGap}'s usual
* conclusions when {@code memberHerdrSocket:} is configured. {@link #hostEnvNames} (and
* everything derived from it — {@code known}/{@code allow} coverage, the allow-list scrub's
* derived set) describes the DAEMON's own environment; under this config key member panes run as
* a different OS user with a different environment entirely, so neither "every member pane
* inherits them UNBLOCKED" nor "the scrub blanks them" is evidence-backed here — both would be
* reporting on the wrong process. Logged once, names the config key, and states the honest
* conclusion: the gap for member panes is UNKNOWN, not clean, so {@code memberCredentials} cannot
* be verified from this daemon. The one count it does report is scoped explicitly to fleetd's own
* environment, never presented as if it said anything about the member's — see {@link
* #logCredentialGap}'s javadoc for why this branch exists.
*/
private void warnUnknownMemberEnvironment(FleetConfig.MemberCredentials creds) {
if (!unknownMemberEnvironmentWarned.compareAndSet(false, true)) {
return;
}
Set<String> covered = new HashSet<>(creds.known());
covered.addAll(creds.allow());
Set<String> hostNames = hostEnvNames.get();
long gapInFleetdsOwnEnv = hostNames.stream()
.filter(name -> CREDENTIAL_SHAPED_NAME.matcher(name).matches())
.filter(name -> !covered.contains(name))
.count();
log.warn("memberCredentials gap: memberHerdrSocket is configured, so member panes run under "
+ "a different OS user than fleetd's own process, with a different environment "
+ "entirely — fleetd has no channel to read that user's environment. {} of the "
+ "{} names in fleetd's OWN environment are credential-shaped and not on "
+ "known:/allow:, but that count describes fleetd's process, not the member "
+ "herdr's. The credential gap for member panes is UNKNOWN, not clean, and "
+ "memberCredentials cannot be verified from here.",
gapInFleetdsOwnEnv, hostNames.size());
}
/**
* CB-596 criterion 4: a credential-shaped host env var name on neither {@code known} nor
* {@code allow} is not silently allowed — it is reported. {@link #hostEnvNames} enumerates the
@@ -1258,8 +1449,22 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
* (same severity, and same guard, as the deny-by-default case — a name genuinely reaching a
* member unprotected is equally serious whichever path put it there), and the names it says are
* blanked keep the INFO.
*
* <p>fleetd #185 stage 2: everything above assumes the member pane runs under the same OS user
* as the daemon, so {@link #hostEnvNames} mirrors what the pane inherits — see that field's
* javadoc. When {@code memberHerdrSocket:} is configured that assumption is false: the member
* pane runs on a second herdr owned by a <em>different</em> user, and neither conclusion below
* ("inherits them UNBLOCKED" / "the scrub blanks them") is backed by evidence about that user's
* environment. So this method checks that first and, when configured, reports the honest
* "unknown, not clean" conclusion instead — see {@link #warnUnknownMemberEnvironment}. When
* {@code memberHerdrSocket:} is absent (the default, and the only mode this host runs) this
* branch is never taken and every line below is unchanged.
*/
private void logCredentialGap(FleetConfig.MemberCredentials creds, Set<String> effectiveAllowed) {
if (memberHerdrSocketConfigured()) {
warnUnknownMemberEnvironment(creds);
return;
}
Set<String> covered = new HashSet<>(creds.known());
covered.addAll(creds.allow());
List<String> gap = hostEnvNames.get().stream()
@@ -568,6 +568,14 @@ public final class MessageService {
* conjunction. {@code matching.size() >= 2} alone, without a strand, is exactly what
* {@code abandonFailsEveryPendingAsyncTicketForTheReleasedTarget} already covers.
*
* <p><strong>Do not "fix" that gap with a test that reaches past this class.</strong> A test can
* build both facts by calling {@link Rendezvous#close} itself on the waiter an accepted
* {@link #send} is still blocked on: the send keeps the lock, no waiter is registered any more, a
* second parked task stays open, and the next {@link #reply} then strands. That was checked, and
* such a test does go red against the pre-fix loop. But it only goes red because it broke the
* lock-and-waiter invariant above from outside — no caller of this class ever does that — so it
* pins a state production cannot reach, and would read to the next person as if it could.
*
* @return true if a live waiter or an async task was failed (never true for one recovered as a
* reply — see the note above)
*/
@@ -315,10 +315,18 @@ public final class FleetApp {
.map(Agent.class::cast)
.filter(a -> a.terminalId() != null)
.collect(Collectors.toMap(Agent::terminalId, Function.identity(), (_, b) -> b));
List<Map<String, Object>> out = sessions.roster().stream()
// fleetd #209: this REST roster reports agentSessionId via SessionManager.rosterView, so it
// uses the resolving roster read (caller-driven, not a timer) rather than the plain one.
List<Map<String, Object>> out = sessions.rosterResolved().stream()
.map(s -> SessionManager.rosterView(s, live.get(s.terminalId())))
.toList();
Map<String, Object> body = new LinkedHashMap<>();
// fleetd #199: the endpoint became /members in the CB-634 rename but the body key stayed
// "workers", so a caller that read "members" saw an empty fleet and reported no members at
// all. "members" is the canonical key; "workers" stays as a deprecated alias so an existing
// REST consumer keeps working — the out-of-band path a lead falls back to when its MCP mount
// drops reads this endpoint. Drop the alias once nothing reads it.
body.put("members", out);
body.put("workers", out);
// CB-586: operator visibility for the refs/wip snapshot store without shelling into the
// repo — how many snapshot refs exist and roughly what they cost. Present only once a
@@ -89,4 +89,15 @@ public record MemberSession(
return new MemberSession(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
nowNanos, turnCount + 1, state, worktree, branch, charterReceipt, agentSessionId);
}
/**
* Return a copy with {@code agentSessionId} resolved to a non-null value (fleetd #209). Some
* adapters (opencode) cannot answer {@link dev.ltms.fleet.peer.PeerHandle#agentSessionId()} at
* spawn time — the peer has not persisted its session record yet — so the id is discovered on
* a later poll and swapped into the otherwise-immutable session via this wither.
*/
public MemberSession withAgentSessionId(String agentSessionId) {
return new MemberSession(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
lastActivityAtNanos, turnCount, state, worktree, branch, charterReceipt, agentSessionId);
}
}
@@ -46,6 +46,15 @@ public final class SessionManager implements TurnListener {
private final PeerLauncher launcher;
private final Worktrees worktrees;
private final ConcurrentHashMap<String /*paneId*/, MemberSession> registry = new ConcurrentHashMap<>();
/**
* fleetd #209: the live {@link PeerHandle} for every registered pane, retained solely so
* {@link #resolveAgentSessionId} can re-poll {@link PeerHandle#agentSessionId()} after spawn.
* The handle used to go out of scope at the end of the spawn method, so a launcher that answers
* the id lazily (opencode — the on-disk session row is written after the pane is created) could
* never be re-asked, and {@code fleet_list}/{@code fleet_spawn resumeSessionId} never saw it.
* Populated on every spawn path, removed on {@link #release}.
*/
private final ConcurrentHashMap<String /*paneId*/, PeerHandle> handles = new ConcurrentHashMap<>();
private final MemberPresence presence;
private final SecureRandom nonceRandom = new SecureRandom();
private final AtomicLong nonceSeq = new AtomicLong();
@@ -205,6 +214,7 @@ public final class SessionManager implements TurnListener {
handle.charterReceipt(),
handle.agentSessionId());
registry.put(handle.id(), session);
handles.put(handle.id(), handle);
memberLifecycle.acquired(session.role(), session.profile(), session.terminalId());
log.debug("acquired session id={} terminal={} profile={} owner={}",
handle.id(), handle.terminalId(), session.profile(), session.ownerTerminal());
@@ -257,6 +267,10 @@ public final class SessionManager implements TurnListener {
*/
private void release(String paneId, ReleaseCause cause) {
MemberSession removed = registry.remove(paneId);
// fleetd #209: remove right alongside the registry entry so a released session's handle is
// never leaked — but keep the local reference below, so the id can still be resolved for
// the ReleaseDetail this teardown notifies with.
PeerHandle removedHandle = handles.remove(paneId);
boolean preserveWorktree = cause == ReleaseCause.SHUTDOWN;
String snapshotRef = null;
if (removed != null) {
@@ -303,8 +317,11 @@ public final class SessionManager implements TurnListener {
// too, so a failed ticket's detail can point a lead at the same tree to re-dispatch.
// CB-584 (issue #65 criterion 5): carry agentSessionId alongside them, so a lead can
// also resume the member's conversation, not just re-dispatch onto its files.
notifyReleased(new ReleaseDetail(removed.terminalId(), removed.worktree(),
removed.branch(), snapshotRef, removed.agentSessionId()));
// fleetd #209: a late-resolving adapter (opencode) may only now have an id — resolve
// one last time so a released member's detail carries the id it now has.
MemberSession resolved = resolveAgentSessionId(removed, removedHandle);
notifyReleased(new ReleaseDetail(resolved.terminalId(), resolved.worktree(),
resolved.branch(), snapshotRef, resolved.agentSessionId()));
}
}
// CB-581: the pane must always stop, even if the dirty check above threw. A session removed
@@ -508,6 +525,7 @@ public final class SessionManager implements TurnListener {
handle.charterReceipt(),
handle.agentSessionId());
registry.put(handle.id(), session);
handles.put(handle.id(), handle);
memberLifecycle.acquired(session.role(), session.profile(), session.terminalId());
log.debug("acquired worktree session id={} terminal={} profile={} branch={} path={}",
handle.id(), handle.terminalId(), session.profile(), session.branch(), session.worktree());
@@ -539,16 +557,73 @@ public final class SessionManager implements TurnListener {
return launcher.defaultProfile();
}
/** The session for {@code paneId}, if it is still registered and not released. */
/**
* The session for {@code paneId}, if it is still registered and not released. fleetd #209:
* resolves a still-unknown {@code agentSessionId} against the retained handle before returning,
* so {@code fleet_status} sees an id a lazy-resolving adapter has since written.
*/
public Optional<MemberSession> get(String paneId) {
return Optional.ofNullable(registry.get(paneId));
return Optional.ofNullable(registry.get(paneId)).map(this::resolveAgentSessionId);
}
/** Fleet-owned roster: all registered sessions (acquired minus released). */
/**
* Fleet-owned roster: all registered sessions (acquired minus released). Deliberately does
* <strong>not</strong> resolve {@code agentSessionId} (fleetd #209 follow-up) — this is the
* roster supplier on the heartbeat and health-tick timers ({@code LeadHeartbeatLoop},
* {@code FleetHealthMonitor} in {@code Fleetd}), on the placement/exhaustion paths, and on the
* metrics scrape ({@code FleetMetrics}), all called far more often than any caller actually
* reads {@code agentSessionId}. Resolving here would mean every tick opens a lazy-resolving
* adapter's (opencode's) on-disk session store once per member whose id is still unknown — and
* for a member whose id never appears, that cost never stops, for the life of the process. Use
* {@link #rosterResolved()} instead wherever the id must be current.
*/
public List<MemberSession> roster() {
return List.copyOf(registry.values());
}
/**
* {@link #roster()}, with each session's still-unknown {@code agentSessionId} re-resolved
* against its retained handle (fleetd #209) — so a caller that actually reports the id (
* {@code fleet_list}, the REST roster) sees one a lazy-resolving adapter (opencode) has since
* written, rather than the null frozen in at spawn time. Reserved for caller-driven reads, not
* timers: see {@link #roster()}'s javadoc for why the plain roster must stay non-resolving.
*/
public List<MemberSession> rosterResolved() {
return registry.values().stream().map(this::resolveAgentSessionId).toList();
}
/**
* Resolve {@code session}'s {@code agentSessionId} if still unknown, re-polling the retained
* {@link PeerHandle} for this pane (fleetd #209). A no-op — returning {@code session} unchanged
* — once the id is already known, once no handle is retained for this pane (never spawned, or
* already released), or if the handle throws while answering. A resolved id is best-effort
* CAS-swapped into the registry via {@link #replace}; a lost race just means another caller
* already applied the same update, so the resolved value is returned either way.
*/
private MemberSession resolveAgentSessionId(MemberSession session) {
return resolveAgentSessionId(session, handles.get(session.paneId()));
}
private MemberSession resolveAgentSessionId(MemberSession session, PeerHandle handle) {
if (session.agentSessionId() != null || handle == null) {
return session;
}
String resolved;
try {
resolved = handle.agentSessionId();
} catch (RuntimeException e) {
log.debug("agentSessionId lookup failed for pane={} terminal={}: {}",
session.paneId(), session.terminalId(), e.toString());
return session;
}
if (resolved == null) {
return session;
}
MemberSession updated = session.withAgentSessionId(resolved);
replace(session, updated); // best-effort; a lost CAS just means the resolved value stands anyway
return updated;
}
/**
* CB-304 merged roster+live view. The registry is authoritative for worktree, branch,
* profile, owner, and state; the optional live agent supplies the herdr-reported status.
@@ -12,20 +12,26 @@ 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.junit.jupiter.api.io.TempDir;
import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.attribute.PosixFileAttributeView;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.function.Function;
import java.util.function.Supplier;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.junit.jupiter.api.Assumptions.assumeTrue;
/**
* CB-633: proves the allow-list scrub is actually WIRED INTO the spawn path — not merely that its
@@ -119,6 +125,18 @@ class HerdrPeerLauncherAllowListWiringTest {
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST, allow, List.of(), null);
}
/**
* Same as {@link #allowList()} but with an operator-configured {@code known:} list, so the
* fallback overlay (the non-zsh / no-memberLoginShell path) has something visible to shadow —
* {@link FleetConfig.MemberCredentials#blockedSet()} is {@code known - allow}, so an empty
* {@code known} (what {@link #allowList()} uses) blocks nothing and a fallback test would have
* no sentinel entry to assert on.
*/
private static Supplier<FleetConfig.MemberCredentials> allowListWithKnown(List<String> known) {
return () -> new FleetConfig.MemberCredentials(
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST, List.of(), known, 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
@@ -257,6 +275,268 @@ class HerdrPeerLauncherAllowListWiringTest {
+ appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
}
/**
* fleetd #185 stage 2 pin: with {@code memberHerdrSocket:} absent (today's only mode, and the
* default — this host runs no other), the gap detector's WARN/INFO conclusions read exactly as
* they did before this fix. Real path: {@code FLEETD_WORKER_TOKEN} (the test profile's own
* {@code tokenEnv}) is a name the derived allow-list keeps, so it gets the "UNBLOCKED" WARN;
* {@code SOME_UNKNOWN_SECRET_TOKEN} is not derived from anywhere, so it gets the "scrub blanks
* them" INFO. This is the exact shape #185 stage 2 must not touch on this path.
*/
@Test
void gapConclusionsAreByteIdenticalWhenMemberHerdrSocketIsAbsent() {
FakeHerdr herdr = new FakeHerdr();
Set<String> hostEnvNames = Set.of("FLEETD_WORKER_TOKEN", "SOME_UNKNOWN_SECRET_TOKEN");
WiringLauncher launcher = new WiringLauncher(herdr, allowList(), "/bin/zsh", () -> hostEnvNames,
() -> config(null));
List<String> messages = spawnAndCaptureLogs(launcher);
assertTrue(messages.contains("memberCredentials gap: 1 credential-shaped env var name(s) are on "
+ "neither known: nor allow: — the derived allow-list keeps them anyway (a "
+ "profile's gitTokenEnv/gitHostEnv/tokenEnv/env: names one, or this spawn "
+ "injects it), so every member pane inherits them UNBLOCKED — [FLEETD_WORKER_TOKEN]. "
+ "Add each to memberCredentials.known (or .allow if a member legitimately needs "
+ "it), or remove it from whatever profile setting derives it in."),
"expected the pre-existing UNBLOCKED WARN unchanged, got: " + messages);
assertTrue(messages.contains("memberCredentials gap: 1 credential-shaped env var name(s) are on "
+ "neither known: nor allow: — [SOME_UNKNOWN_SECRET_TOKEN]. The allow-list scrub "
+ "blanks them anyway (they are not on the derived allow-list), so no member pane "
+ "keeps them; add each to memberCredentials.known or .allow to make that explicit."),
"expected the pre-existing 'scrub blanks them' INFO unchanged, got: " + messages);
}
/**
* fleetd #185 stage 2: with {@code memberHerdrSocket:} configured, member panes run under a
* different OS user — {@link HerdrPeerLauncher#hostEnvNames} describes fleetd's own process, not
* that user's. Neither "inherits them UNBLOCKED" nor "scrub blanks them" is evidence-backed
* there, so neither may print; the single unknown-environment WARN must, naming the config key.
*/
@Test
void gapDetectorReportsUnknownInsteadOfAConclusionWhenMemberHerdrSocketIsConfigured() {
FakeHerdr herdr = new FakeHerdr();
Set<String> hostEnvNames = Set.of("FLEETD_WORKER_TOKEN", "SOME_UNKNOWN_SECRET_TOKEN");
WiringLauncher launcher = new WiringLauncher(herdr, allowList(), "/bin/zsh", () -> hostEnvNames,
() -> configWithMemberHerdrSocket("/tmp/other-user-herdr.sock"));
List<String> messages = spawnAndCaptureLogs(launcher);
assertTrue(messages.stream().anyMatch(m -> m.contains("memberHerdrSocket")
&& m.contains("UNKNOWN") && m.contains("cannot be verified")),
"expected the unknown-member-environment WARN naming memberHerdrSocket, got: " + messages);
assertFalse(messages.stream().anyMatch(m -> m.contains("UNBLOCKED")),
"the 'inherits them UNBLOCKED' conclusion must not print once the evidence is about "
+ "the wrong (daemon's own) environment — got: " + messages);
assertFalse(messages.stream().anyMatch(m -> m.contains("scrub blanks them")),
"the 'scrub blanks them' conclusion must not print once the evidence is about the "
+ "wrong (daemon's own) environment — got: " + messages);
}
/**
* fleetd #185 stage 2: the unknown-environment WARN is a standing fact about this launcher's
* configuration, not per-spawn news — it must fire once per launcher instance, the same shape as
* every other one-time WARN in this class (e.g. {@code warnNonZsh}).
*/
@Test
void theUnknownEnvironmentWarnFiresOnceNotOncePerSpawn() {
FakeHerdr herdr = new FakeHerdr();
Set<String> hostEnvNames = Set.of("FLEETD_WORKER_TOKEN", "SOME_UNKNOWN_SECRET_TOKEN");
WiringLauncher launcher = new WiringLauncher(herdr, allowList(), "/bin/zsh", () -> hostEnvNames,
() -> configWithMemberHerdrSocket("/tmp/other-user-herdr.sock"));
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
Level original = logger.getLevel();
logger.setLevel(Level.WARN);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
try {
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
} finally {
logger.detachAppender(appender);
logger.setLevel(original);
}
long count = appender.list.stream()
.filter(e -> e.getFormattedMessage().contains("cannot be verified from here"))
.count();
assertEquals(1, count, "the unknown-member-environment WARN must fire once per launcher "
+ "instance, not once per spawn — got " + count + " occurrence(s) among: "
+ appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
}
/**
* Hard constraint: the gap detector must never log an env var VALUE, only its NAME. {@code
* SOME_UNKNOWN_SECRET_TOKEN} resolves to a distinctive canary value through the same {@code env}
* lookup the launcher uses elsewhere (SHELL, PATH, token resolution) — proving the value IS
* resolvable does not mean the detector reads it, since {@link HerdrPeerLauncher#hostEnvNames}
* (names only) is its data source, never {@code env.apply(name)} for those names.
*/
@Test
void theGapDetectorNeverLogsAnEnvVarValueOnlyItsName() {
FakeHerdr herdr = new FakeHerdr();
String canary = "sekrit-value-CANARY-9f3a1b7c";
Set<String> hostEnvNames = Set.of("FLEETD_WORKER_TOKEN", "SOME_UNKNOWN_SECRET_TOKEN");
WiringLauncher launcher = new WiringLauncher(herdr, allowList(), "/bin/zsh", () -> hostEnvNames,
() -> config(null), Map.of("SOME_UNKNOWN_SECRET_TOKEN", canary));
List<String> messages = spawnAndCaptureLogs(launcher);
assertTrue(messages.stream().anyMatch(m -> m.contains("SOME_UNKNOWN_SECRET_TOKEN")),
"expected the credential-shaped NAME to appear in the log, got: " + messages);
assertFalse(messages.stream().anyMatch(m -> m.contains(canary)),
"the log must never contain an env var VALUE, only its NAME — got: " + messages);
}
/**
* fleetd #213 defect 1, acceptance criterion 1: {@code memberHerdrSocket} configured and {@code
* memberLoginShell} configured as non-zsh must fall back to the sentinel overlay exactly like a
* non-zsh {@code $SHELL} does today — and the "generated ZDOTDIR" INFO must not appear, since no
* scrub actually runs. The WiringLauncher's own {@code env("SHELL")} is deliberately set to
* {@code /bin/zsh} — the OPPOSITE of what {@code memberLoginShell} says — so a launcher that
* (incorrectly) fell back to fleetd's own {@code $SHELL} here would wrongly pass the gate and
* fail this test.
*/
@Test
void memberHerdrSocketWithNonZshMemberLoginShellFallsBackToTheOverlay() {
FakeHerdr herdr = new FakeHerdr();
WiringLauncher launcher = new WiringLauncher(herdr, allowListWithKnown(List.of("SOME_TOKEN")),
"/bin/zsh", null,
() -> configWithMemberHerdrSocket("/tmp/other-user-herdr.sock", "/bin/bash"));
List<String> messages = spawnAndCaptureLogs(launcher);
assertEquals("blocked-by-fleetd-cb596-see-gitea-issue-82", launcher.env.get("SOME_TOKEN"),
"a non-zsh memberLoginShell must fall back to the CB-596 sentinel overlay, exactly "
+ "like a non-zsh $SHELL does when memberHerdrSocket is absent");
assertFalse(launcher.env.containsKey("ZDOTDIR"),
"no scrub directory may be generated when the configured member login shell is not zsh");
assertFalse(messages.stream().anyMatch(m -> m.contains("generated ZDOTDIR")),
"the 'generated ZDOTDIR' INFO must not appear when the scrub never runs — got: " + messages);
}
/**
* fleetd #213 defect 1, acceptance criterion 2: {@code memberHerdrSocket} configured and NO
* {@code memberLoginShell} configured must fall back exactly like criterion 1 above — AND
* fleetd's own {@code $SHELL} must never even be consulted (not merely "not decisive"). The
* fixture's {@code env} function reports {@code /bin/zsh} for {@code SHELL} — a value that would
* WRONGLY pass the zsh gate if the fix regressed to reading it — while flagging whether it was
* ever asked for at all, so this test fails loudly on either kind of regression.
*/
@Test
void memberHerdrSocketWithNoMemberLoginShellFallsBackAndNeverConsultsFleetdsOwnShell() {
FakeHerdr herdr = new FakeHerdr();
AtomicBoolean shellQueried = new AtomicBoolean(false);
Function<String, String> env = name -> {
if ("SHELL".equals(name)) {
shellQueried.set(true);
return "/bin/zsh"; // would wrongly pass the zsh gate if this ever leaked through
}
return null;
};
WiringLauncher launcher = new WiringLauncher(herdr, allowListWithKnown(List.of("SOME_TOKEN")), env,
() -> configWithMemberHerdrSocket("/tmp/other-user-herdr.sock", null));
List<String> messages = spawnAndCaptureLogs(launcher);
assertFalse(shellQueried.get(), "fleetd's own $SHELL must never be consulted once "
+ "memberHerdrSocket is configured — only memberLoginShell: may decide the gate");
assertEquals("blocked-by-fleetd-cb596-see-gitea-issue-82", launcher.env.get("SOME_TOKEN"),
"no memberLoginShell configured must fall back to the sentinel overlay, same as a "
+ "configured non-zsh shell");
assertFalse(messages.stream().anyMatch(m -> m.contains("generated ZDOTDIR")),
"no scrub may run without a configured memberLoginShell — got: " + messages);
}
/**
* fleetd #213 defect 2, acceptance criterion 3: with {@code memberHerdrSocket} configured, a
* zsh {@code memberLoginShell}, and {@code worktreeRoot}/{@code worktreeGroup} both configured,
* the generated scrub directory must live under {@code worktreeRoot} — NEVER under {@code
* java.io.tmpdir}, which is fleetd's own 0700 temp dir and unreadable by the member's different
* OS user. {@code worktreeGroup} is set to the CURRENT process's own primary group so {@code
* EnvAllowListScrub}'s group-sharing step resolves on whatever host runs this test, rather than
* hardcoding a group name that may not exist here.
*/
@Test
void memberHerdrSocketWithZshMemberLoginShellPutsTheScrubOutsideJavaIoTmpdir(@TempDir Path worktreeRoot)
throws IOException {
String group = currentUserGroup();
FakeHerdr herdr = new FakeHerdr();
WiringLauncher launcher = new WiringLauncher(herdr, allowList(), "/bin/bash-should-be-ignored", null,
() -> configWithMemberHerdrSocketRootAndGroup("/tmp/other-user-herdr.sock", "/bin/zsh",
worktreeRoot.toString(), group));
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
String zdotdir = launcher.env.get("ZDOTDIR");
assertNotNull(zdotdir, "a zsh memberLoginShell with worktreeRoot+worktreeGroup configured "
+ "must still generate a ZDOTDIR");
Path dir = Path.of(zdotdir);
// The immediate PARENT is asserted (not merely startsWith(java.io.tmpdir)), because a
// JUnit @TempDir is itself carved out of the JVM's java.io.tmpdir — startsWith alone would
// pass by coincidence of the test fixture, not because the launcher used worktreeRoot.
assertEquals(worktreeRoot.toAbsolutePath().normalize(), dir.getParent(),
"the generated scrub directory's parent must be the configured worktreeRoot, not "
+ "System.getProperty(\"java.io.tmpdir\") — got parent " + dir.getParent());
}
/**
* fleetd #213, acceptance criterion 4: with {@code memberHerdrSocket} absent (today's only
* mode), behaviour must be byte-identical to before this fix — fleetd's own {@code $SHELL}
* still decides the gate (proven here, not merely assumed, by flagging the lookup), and the
* scrub still lands under {@code java.io.tmpdir}.
*/
@Test
void memberHerdrSocketAbsentStillConsultsFleetdsOwnShellAndBehavesAsBefore() {
FakeHerdr herdr = new FakeHerdr();
AtomicBoolean shellQueried = new AtomicBoolean(false);
Function<String, String> env = name -> {
if ("SHELL".equals(name)) {
shellQueried.set(true);
return "/bin/zsh";
}
return null;
};
WiringLauncher launcher = new WiringLauncher(herdr, allowList(), env, null); // memberHerdrSocket absent
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
assertTrue(shellQueried.get(), "with memberHerdrSocket absent, fleetd's own $SHELL must "
+ "still decide the zsh gate, unchanged from before this fix");
String zdotdir = launcher.env.get("ZDOTDIR");
assertNotNull(zdotdir, "SHELL=/bin/zsh with memberHerdrSocket absent must still generate a "
+ "ZDOTDIR, as before this fix");
Path dir = Path.of(zdotdir);
assertTrue(dir.startsWith(Path.of(System.getProperty("java.io.tmpdir"))),
"with memberHerdrSocket absent the scrub directory must still be generated under "
+ "java.io.tmpdir, unchanged from before this fix: " + dir);
}
/** The current process's own primary group — resolvable on whatever host runs this test. */
private static String currentUserGroup() throws IOException {
PosixFileAttributeView view = Files.getFileAttributeView(Path.of("."), PosixFileAttributeView.class);
assumeTrue(view != null, "this host's filesystem does not support POSIX group ownership");
return view.readAttributes().group().getName();
}
/** Spawn once through the real launcher path, capturing every INFO+ line this class logs. */
private static List<String> spawnAndCaptureLogs(HerdrPeerLauncher launcher) {
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);
}
return appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList();
}
private static String readAll(Path p) {
try {
return Files.readString(p);
@@ -294,12 +574,40 @@ class HerdrPeerLauncherAllowListWiringTest {
WiringLauncher(FakeHerdr herdr, Supplier<FleetConfig.MemberCredentials> creds, String shell,
Supplier<Set<String>> hostEnvNames, Supplier<FleetConfig> config) {
this(herdr, creds, shell, hostEnvNames, config, Map.of());
}
/**
* Plus a host-env value map (name → value), resolved through the same {@code env} lookup
* every adapter uses for {@code SHELL}/{@code PATH}/token resolution — fleetd #185 stage 2's
* "the gap detector never logs a value" tests use this to prove a value that IS resolvable
* for a credential-shaped name never reaches the log, since the detector only ever reads
* {@code hostEnvNames} (names), never {@code env.apply(name)} (values), for those names.
*/
WiringLauncher(FakeHerdr herdr, Supplier<FleetConfig.MemberCredentials> creds, String shell,
Supplier<Set<String>> hostEnvNames, Supplier<FleetConfig> config,
Map<String, String> extraEnvValues) {
super("test", new AgentControl(herdr), new WorkspaceControl(herdr),
Map.of("test", profile()), "test",
name -> "SHELL".equals(name) ? shell : null,
name -> "SHELL".equals(name) ? shell
: (extraEnvValues != null && extraEnvValues.containsKey(name))
? extraEnvValues.get(name) : null,
0, () -> 0L, () -> { }, null, creds, hostEnvNames, config);
}
/**
* Full control over the {@code env} lookup, bypassing the {@code shell}/{@code
* extraEnvValues} convenience above entirely — fleetd #213's "fleetd's own $SHELL must
* never be consulted" tests need to OBSERVE whether {@code SHELL} was ever looked up, which
* a plain value substitution cannot do.
*/
WiringLauncher(FakeHerdr herdr, Supplier<FleetConfig.MemberCredentials> creds,
Function<String, String> env, Supplier<FleetConfig> config) {
super("test", new AgentControl(herdr), new WorkspaceControl(herdr),
Map.of("test", profile()), "test",
env, 0, () -> 0L, () -> { }, null, creds, null, config);
}
@Override
protected Launch buildLaunch(FleetConfig.Profile cfg, LaunchSpec spec) {
Map<String, String> launchEnv = baseEnv(cfg);
@@ -321,6 +629,82 @@ class HerdrPeerLauncherAllowListWiringTest {
null, null, null, null, null).withDefaults();
}
/**
* fleetd #185 stage 2: a config with {@code memberHerdrSocket:} set — member panes run on a
* second herdr owned by a different OS user, so {@link HerdrPeerLauncher#hostEnvNames} no
* longer describes what a member pane inherits.
*/
private static FleetConfig configWithMemberHerdrSocket(String memberHerdrSocket) {
return new FleetConfig(null, null, memberHerdrSocket, Map.of(), null, null, null, null, null,
null, null, null, null, null, null, null, null, null, null, null).withDefaults();
}
/**
* fleetd #213: as {@link #configWithMemberHerdrSocket(String)}, plus the {@code
* memberLoginShell:} the member's OS user actually runs — the config key {@link
* HerdrPeerLauncher#applyEnvironmentAllowListPolicy} must consult instead of fleetd's own
* {@code $SHELL} once {@code memberHerdrSocket} is configured.
*/
private static FleetConfig configWithMemberHerdrSocket(String memberHerdrSocket, String memberLoginShell) {
return new FleetConfig(
null, // bind
null, // herdrSocket
memberHerdrSocket, // memberHerdrSocket
Map.of(), // profiles
null, // guard
null, // worktreeRoot
null, // lifecycle
null, // spawnReadyTimeoutMs
null, // spawnReadyPollMs
null, // broker
null, // primary
null, // fleet
null, // leadHeartbeat
null, // health
null, // placement
null, // auth
null, // configReload
null, // quarantineCooldownSeconds
null, // memberCredentials
null, // coordinator
null, // worktreeGroup
memberLoginShell // memberLoginShell
).withDefaults();
}
/**
* fleetd #213: as above, plus {@code worktreeRoot:}/{@code worktreeGroup:} — both required for
* the ZDOTDIR scrub to run at all once {@code memberHerdrSocket} is configured; either missing
* falls back to the sentinel overlay, same as a non-zsh {@code memberLoginShell}.
*/
private static FleetConfig configWithMemberHerdrSocketRootAndGroup(String memberHerdrSocket,
String memberLoginShell, String worktreeRoot, String worktreeGroup) {
return new FleetConfig(
null, // bind
null, // herdrSocket
memberHerdrSocket, // memberHerdrSocket
Map.of(), // profiles
null, // guard
worktreeRoot, // worktreeRoot
null, // lifecycle
null, // spawnReadyTimeoutMs
null, // spawnReadyPollMs
null, // broker
null, // primary
null, // fleet
null, // leadHeartbeat
null, // health
null, // placement
null, // auth
null, // configReload
null, // quarantineCooldownSeconds
null, // memberCredentials
null, // coordinator
worktreeGroup, // worktreeGroup
memberLoginShell // memberLoginShell
).withDefaults();
}
/** The generated directory is a temp directory; make sure the test does not leave a pile. */
@Test
void theGeneratedDirectoryIsRemovedWhenThePaneIsStopped() {
@@ -218,9 +218,17 @@ class FleetAppTest {
HttpResponse<String> res = req(port, "GET", "/members");
assertEquals(200, res.statusCode());
JsonNode workers = mapper.readTree(res.body()).get("workers");
assertEquals(1, workers.size());
JsonNode w = workers.get(0);
JsonNode body = mapper.readTree(res.body());
// fleetd #199: "members" is the canonical key. The endpoint is /members, so a caller that
// reads "members" must not see an empty fleet. "workers" is kept only as a deprecated alias
// and must carry the same rows — assert both, or the alias can silently drift.
JsonNode members = body.get("members");
assertNotNull(members, "GET /members must return its rows under \"members\"");
assertEquals(1, members.size());
JsonNode workers = body.get("workers");
assertNotNull(workers, "the deprecated \"workers\" alias is still emitted");
assertEquals(members, workers, "the alias must carry the same rows as \"members\"");
JsonNode w = members.get(0);
assertEquals(spawned.get("terminalId").asText(), w.get("sessionId").asText());
assertEquals(paneId, w.get("paneId").asText());
assertEquals("ltms-local", w.get("profile").asText());
@@ -931,4 +931,235 @@ class SessionManagerTest {
assertTrue(e.getMessage().contains("SESSION_RESUME"), e.getMessage());
assertTrue(e.getMessage().contains("stub-profile"), e.getMessage());
}
// ── fleetd #209: agentSessionId resolved lazily against the retained handle ─────────────────
//
// The opencode adapter cannot answer PeerHandle.agentSessionId() at spawn time — the on-disk
// session row is written only after the pane is live — so the id must be re-polled on a LATER
// call, against the SAME handle instance the launcher returned at spawn. SessionManager used
// to let that handle go out of scope at the end of the spawn method, so no caller ever re-asked
// it and fleet_list/fleet_status never saw the id. LazyIdHandle below reproduces exactly that
// shape: null on the first N calls (the spawn-time call included), a real id after.
/**
* A {@link PeerHandle} whose {@link #agentSessionId()} answers {@code null} for its first
* {@code nullCalls} invocations, then either a fixed id or a configured throw on every call
* after that — the shape of the opencode bug (fleetd #209): the session row is not written
* until after the pane is live, so early polls come back empty and a later one finds it.
*/
private static final class LazyIdHandle implements PeerHandle {
private final String id;
private final String terminalId;
private final int nullCalls;
private final String resolvedId;
private final java.util.concurrent.atomic.AtomicInteger calls =
new java.util.concurrent.atomic.AtomicInteger();
private volatile RuntimeException throwAfter;
LazyIdHandle(String id, String terminalId, int nullCalls, String resolvedId) {
this.id = id;
this.terminalId = terminalId;
this.nullCalls = nullCalls;
this.resolvedId = resolvedId;
}
/** After the null calls are exhausted, throw instead of answering the resolved id. */
LazyIdHandle throwing(RuntimeException e) {
this.throwAfter = e;
return this;
}
@Override
public String id() {
return id;
}
@Override
public String terminalId() {
return terminalId;
}
@Override
public String agentSessionId() {
int n = calls.incrementAndGet();
if (n <= nullCalls) {
return null;
}
if (throwAfter != null) {
throw throwAfter;
}
return resolvedId;
}
@Override
public CharterReceipt charterReceipt() {
return null;
}
int callCount() {
return calls.get();
}
}
/** A minimal {@link PeerLauncher} that hands out pre-built {@link LazyIdHandle}s, one per spawn. */
private static final class LazyIdLauncher implements PeerLauncher {
private final java.util.Deque<LazyIdHandle> queued = new java.util.ArrayDeque<>();
LazyIdLauncher queue(LazyIdHandle handle) {
queued.add(handle);
return this;
}
@Override
public Set<Capability> capabilities() {
return Set.of(Capability.WORKTREE, Capability.SESSION_RESUME);
}
@Override
public Set<Capability> capabilitiesFor(String profileName) {
return capabilities();
}
@Override
public PeerHandle spawn(SpawnRequest req) {
LazyIdHandle handle = queued.poll();
if (handle == null) {
throw new IllegalStateException("no queued LazyIdHandle for this spawn");
}
return handle;
}
@Override
public Set<String> profiles() {
return Set.of("lazy");
}
@Override
public String defaultProfile() {
return "lazy";
}
@Override
public String effectiveCwd(SpawnRequest req) {
return "/cwd";
}
@Override
public List<String> parityOverlay(String profileName) {
return List.of();
}
@Override
public List<?> list() {
return List.of();
}
@Override
public int reapOrphanWorkers() {
return 0;
}
@Override
public void stop(String id) {
}
@Override
public boolean clearContext(String id) {
return false;
}
}
@Test
void plainRosterDoesNotResolveAgentSessionId() {
// fleetd #209 follow-up: roster() sits on the heartbeat/health-tick timers (and the metrics
// scrape), so it must never trigger the resolve lookup — for opencode that lookup opens an
// on-disk session database, and a member whose id never appears would pay that cost forever.
// rosterResolved() is the one to use when a caller actually reports the id.
LazyIdHandle handle = new LazyIdHandle("p0", "t0", 1, "oc-session-0");
SessionManager sessions = new SessionManager(new LazyIdLauncher().queue(handle));
MemberSession acquired = sessions.acquire("lazy", "/cwd", "/caller", null);
assertEquals(1, handle.callCount(), "sanity: only the spawn-time call happened so far");
List<MemberSession> roster = sessions.roster();
assertEquals(1, roster.size());
assertEquals(acquired.paneId(), roster.getFirst().paneId());
assertNull(roster.getFirst().agentSessionId(), "the plain roster must not resolve the id");
assertEquals(1, handle.callCount(),
"roster() must never call agentSessionId() again — it sits on the heartbeat/health timers");
}
@Test
void rosterResolvedResolvesALateAgentSessionIdFromTheRetainedHandle() {
LazyIdHandle handle = new LazyIdHandle("p1", "t1", 1, "oc-session-1");
SessionManager sessions = new SessionManager(new LazyIdLauncher().queue(handle));
MemberSession acquired = sessions.acquire("lazy", "/cwd", "/caller", null);
assertNull(acquired.agentSessionId(),
"opencode has not written its session row yet at spawn time");
List<MemberSession> roster = sessions.rosterResolved();
assertEquals(1, roster.size());
assertEquals("oc-session-1", roster.getFirst().agentSessionId(),
"fleet_list must see the id once the adapter can answer it");
Map<String, Object> view = SessionManager.rosterView(roster.getFirst(), null);
assertEquals("oc-session-1", view.get("agentSessionId"),
"rosterView renders whatever rosterResolved() resolved");
}
@Test
void getResolvesALateAgentSessionIdFromTheRetainedHandle() {
LazyIdHandle handle = new LazyIdHandle("p2", "t2", 1, "oc-session-2");
SessionManager sessions = new SessionManager(new LazyIdLauncher().queue(handle));
MemberSession acquired = sessions.acquire("lazy", "/cwd", "/caller", null);
MemberSession resolved = sessions.get(acquired.paneId()).orElseThrow();
assertEquals("oc-session-2", resolved.agentSessionId(),
"fleet_status (single-session lookup) must also see the late-resolved id");
}
@Test
void releaseCarriesALateResolvedAgentSessionIdIntoTheReleaseDetail() {
LazyIdHandle handle = new LazyIdHandle("p3", "t3", 1, "oc-session-3");
SessionManager sessions = new SessionManager(new LazyIdLauncher().queue(handle));
MemberSession acquired = sessions.acquire("lazy", "/cwd", "/caller", null);
java.util.List<String> released = new java.util.concurrent.CopyOnWriteArrayList<>();
sessions.onRelease(detail -> released.add(detail.agentSessionId()));
sessions.release(acquired.paneId());
assertEquals(java.util.List.of("oc-session-3"), released,
"a released member's detail carries the id it has since resolved, not the null "
+ "frozen in at spawn time");
}
@Test
void aThrowingHandleDoesNotBreakRosterResolved() {
LazyIdHandle handle = new LazyIdHandle("p4", "t4", 1, "oc-session-4")
.throwing(new RuntimeException("sqlite locked"));
SessionManager sessions = new SessionManager(new LazyIdLauncher().queue(handle));
sessions.acquire("lazy", "/cwd", "/caller", null);
List<MemberSession> roster = assertDoesNotThrow(sessions::rosterResolved,
"a handle that throws resolving its id must not break the roster read");
assertEquals(1, roster.size());
assertNull(roster.getFirst().agentSessionId(), "the id stays unresolved when the lookup throws");
}
@Test
void aResolvedAgentSessionIdIsNotLookedUpAgain() {
LazyIdHandle handle = new LazyIdHandle("p5", "t5", 1, "oc-session-5");
SessionManager sessions = new SessionManager(new LazyIdLauncher().queue(handle));
sessions.acquire("lazy", "/cwd", "/caller", null);
assertEquals(1, handle.callCount(), "sanity: only the spawn-time call happened so far");
sessions.rosterResolved();
assertEquals(2, handle.callCount(), "the first rosterResolved() read resolves the id");
sessions.rosterResolved();
assertEquals(2, handle.callCount(), "once resolved, the id must not be looked up again");
}
}