Compare commits

..

9 Commits

Author SHA1 Message Date
Dai Ha 445a45f6e1 fleetd#211: classify BACKEND_EXHAUSTED/BACKEND_ERROR from the raw scrape as a fallback
CI / contract (pull_request) Successful in 50s
CI / build (pull_request) Successful in 1m17s
CompletionResolver.resolve() returned an empty-scrape failure before the
BACKEND_EXHAUSTED / BACKEND_ERROR classification ever ran, whenever
lastAssistantBlock() found no usable text — most commonly a pane with no ⏺
marker at all, whose boundary scan starts at the top of the raw screen and
breaks immediately on the first line of TUI chrome. Since BACKEND_EXHAUSTED
is the only caller of exhaustionSink, this meant an exhausted backend was
recorded as "produced nothing" instead of being quarantined.

Fix: run the same two classifications against the raw (untrimmed) scrape as
a fallback, only inside the empty-scrape failure branch. A pane that already
yields a usable assistant block never reaches this branch, so the existing
narrow match is unchanged. lastAssistantBlock stays the sole source of the
reply text; only classification ever consults the raw scrape.
2026-09-01 10:50:10 +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
12 changed files with 744 additions and 178 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)
@@ -237,11 +237,13 @@ public final class CompletionResolver implements TurnListener {
}
String tail;
String assistantBlock = null;
String rawScrape = null;
int originalLength = 0;
boolean clipped = false;
boolean scrapeFailed = false;
try {
assistantBlock = lastAssistantBlock(agents.read(target, SCRAPE_SOURCE));
rawScrape = agents.read(target, SCRAPE_SOURCE);
assistantBlock = lastAssistantBlock(rawScrape);
originalLength = assistantBlock.strip().length();
clipped = originalLength > MAX_SCRAPE_CHARS;
tail = clip(assistantBlock);
@@ -256,6 +258,20 @@ public final class CompletionResolver implements TurnListener {
// 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()) {
// fleetd#211: lastAssistantBlock() found nothing usable — most often a pane with no ⏺
// marker at all, whose boundary scan then starts at the top of the raw screen and breaks
// immediately on the first line of TUI chrome (╭, │, ❯, …). Before giving up as a lost
// turn, run the same exhaustion/backend-error classification against the RAW scrape as a
// fallback, ONLY here. A pane that already yielded a usable block never reaches this
// branch, so the narrow (trimmed) match on the normal path below is completely unchanged
// — zero new false positives there. Every pane this fallback examines was already headed
// for the empty-scrape failure, so a wrong label here is strictly less bad than silently
// losing an exhaustion signal: the alternative outcome is already a failure, just one that
// never quarantines the credential. lastAssistantBlock stays the source of the reply
// TEXT everywhere else; only classification ever consults the raw scrape, and only here.
if (rawScrape != null && classifyRawScrapeFallback(target, turn, waiter, rawScrape)) {
return;
}
fail(target, turn, emptyScrapeReason(target, scrapeFailed));
return;
}
@@ -313,6 +329,38 @@ public final class CompletionResolver implements TurnListener {
}
}
/**
* fleetd#211: the raw-scrape fallback classification, run only when {@link #lastAssistantBlock}
* found nothing usable (see the call site in {@link #resolve}). Mirrors the two classifications
* the normal path already applies to the trimmed assistant block — exhaustion first, then the
* narrow {@link #BACKEND_ERROR} pattern — against {@code raw} instead, and reports whether one of
* them handled the turn (resolved the waiter or failed it) so the caller skips the empty-scrape
* failure. Never runs on the normal (non-empty-block) path, and never touches the reply text.
*/
private boolean classifyRawScrapeFallback(String target, InFlight turn,
CompletableFuture<Rendezvous.Resolution> waiter, String raw) {
Pattern exhausted = exhaustedPatterns.patternFor(target);
String matchedLine = exhausted == null ? null : firstMatchingLine(raw, 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 from the raw scrape (no "
+ "usable assistant block; no fleet_reply): {}", 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 true;
}
String backendError = firstMatchingLine(raw, BACKEND_ERROR);
if (backendError != null) {
fail(target, turn, "member " + target + " ended on a backend error: " + backendError);
return true;
}
return false;
}
/** Synchronous fail (the unit-testable core of {@link #onTurnFailed}). */
void fail(String target, InFlight turn) {
fail(target, turn, null);
@@ -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();
@@ -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;
@@ -1230,6 +1236,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 +1325,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.
@@ -633,6 +633,109 @@ class CompletionResolverTest {
"the failure carries the rest of the pane, not only the matched line: " + reason);
}
// --- fleetd#211: raw-scrape fallback classification when there is no usable assistant block ---
@Test
void anExhaustionLineWithNoMarkerAndLeadingChromeIsClassifiedFromTheRawScrapeAndNotifiesTheSink() {
// No ⏺ anywhere, and the first visible line is TUI chrome (╭). lastAssistantBlock's boundary
// scan starts at the top of the raw screen and breaks immediately, so the trimmed block is "".
// The fix: fall back to matching the RAW scrape so this doesn't get lost as an empty scrape.
String block = """
╭──────────────────────────────────────╮
The usage limit has been reached. Try again later.
""";
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached");
java.util.List<String> notified = new java.util.ArrayList<>();
ExhaustionSink sink = (target, reason) -> notified.add(target + ": " + reason);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, sink);
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
assertTrue(waiter.isDone(), "a raw-scrape match still resolves the blocked send");
assertEquals(Rendezvous.Kind.BACKEND_EXHAUSTED, waiter.getNow(null).kind(),
"classified from the raw scrape even though the trimmed block was empty");
assertEquals(1, notified.size(),
"the sink is the whole point of this ticket — it must be notified: " + notified);
assertTrue(notified.get(0).startsWith("term_a: "), "the sink is told which target exhausted");
assertTrue(notified.get(0).contains("The usage limit has been reached"),
"the sink is told the matched reason: " + notified.get(0));
}
@Test
void aBackendErrorLineWithNoMarkerAndLeadingChromeIsClassifiedFromTheRawScrape() {
String block = """
╭──────────────────────────────────────╮
API Error: 400 invalid request body
""";
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
assertTrue(waiter.isDone(), "a raw-scrape backend-error match still resolves the blocked send");
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(),
"classified BACKEND_ERROR from the raw scrape even though the trimmed block was empty");
assertTrue(waiter.getNow(null).text().contains("API Error: 400 invalid request body"),
"the failure carries the matched line: " + waiter.getNow(null).text());
}
@Test
void anOrdinaryPaneWithANormalAssistantBlockIsUnaffectedByTheRawScrapeFallback() {
// Pin: on a pane that already yields a usable block, the fallback branch is never reached —
// same outcome, same text, sink not called — even though the raw screen around the marker
// would itself match the configured exhausted pattern.
String block = "The usage limit has been reached, but this is a leading TUI line above the "
+ "marker.\n⏺ complete report\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached");
java.util.List<String> notified = new java.util.ArrayList<>();
ExhaustionSink sink = (target, reason) -> notified.add(target + ": " + reason);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, sink);
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
assertTrue(waiter.isDone());
assertEquals(Rendezvous.Kind.COMPLETION, waiter.getNow(null).kind(),
"unchanged: a usable assistant block never reaches the raw-scrape fallback");
assertEquals("complete report", waiter.getNow(null).text(), "the reply text is unaffected");
assertTrue(notified.isEmpty(), "the fallback never runs, so the sink is never called");
}
@Test
void aGenuinelyEmptyScrapeStillFailsAsEmptyAndNeverNotifiesTheSink() {
// The false-positive pin: no exhaustion or backend-error text anywhere on the pane (just
// chrome, no marker) — the raw-scrape fallback must not manufacture a classification, and
// the sink must stay untouched.
String block = """
╭──────────────────────────────────────╮
│ > │
╰──────────────────────────────────────╯
""";
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached");
java.util.List<String> notified = new java.util.ArrayList<>();
ExhaustionSink sink = (target, reason) -> notified.add(target + ": " + reason);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, sink);
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(),
"no exhaustion or backend-error text anywhere ⇒ this stays the ordinary empty-scrape failure");
assertTrue(waiter.getNow(null).text().toLowerCase().contains("empty"),
"the failure still says the scrape was empty: " + waiter.getNow(null).text());
assertTrue(notified.isEmpty(), "a genuinely empty pane must never quarantine a credential");
}
@Test
void coverageIsOffWhenNoProfileHasAPatternConfigured() {
assertEquals("off (no profile has an exhaustedPattern configured; profiles: [terra])",
@@ -257,6 +257,137 @@ 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);
}
/** 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,9 +425,24 @@ 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);
}
@@ -321,6 +467,16 @@ 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();
}
/** 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");
}
}