Compare commits
17 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| d9168de43e | |||
| c3fa1136d4 | |||
| cabcd87b66 | |||
| bf0e09b1a2 | |||
| e1eb50ce65 | |||
| 049e7d9d54 | |||
| ff3b49cd1e | |||
| 0373b6c41b | |||
| 445a45f6e1 | |||
| 966c58a3b8 | |||
| 735c837604 | |||
| 457dc0330d | |||
| a1a9015217 | |||
| 5a811a3695 | |||
| 1fdaa74eb3 | |||
| 847e8bd3fa | |||
| f9fb387427 |
-166
@@ -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)
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -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,43 @@ 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) {
|
||||
// Carry the pane, not just the matched line — the same fleetd#164 rule the normal path
|
||||
// above applies. Here it matters more, not less: the trimmed block was empty, so the raw
|
||||
// scrape is the ONLY copy of whatever the member managed to say. Clipped to the same cap
|
||||
// the normal path uses, since a raw screen has no boundary trimming to bound it.
|
||||
fail(target, turn, "member " + target + " ended on a backend error: " + backendError
|
||||
+ "\n--- pane tail ---\n" + clip(raw));
|
||||
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();
|
||||
|
||||
@@ -278,18 +278,22 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
|
||||
|
||||
/**
|
||||
* Add the Claude-specific session-identity flags to {@code argv} and return the peer's OWN
|
||||
* session id — the resume handle. A resume request passes the prior id via {@code -r} and
|
||||
* returns that id; a fresh named session mints a new UUID, passes it via {@code --session-id},
|
||||
* and returns the mint. The bridge's logical name rides along as {@code -n} when present. When
|
||||
* <em>no</em> identity is requested (sessionName and resumeSessionId both blank) this adds
|
||||
* nothing and returns {@code null}, keeping the legacy no-identity launch byte-identical.
|
||||
* session id — the resume handle. A resume passes the prior id via {@code -r} and returns
|
||||
* that id; every other spawn mints a new UUID, passes it via {@code --session-id}, and
|
||||
* returns the mint. The bridge's logical name rides along as {@code -n} when present.
|
||||
*
|
||||
* <p>fleetd #214: the mint is unconditional. A plain {@code fleet_spawn} passes neither
|
||||
* sessionName nor resumeSessionId, yet the member must still be resumable, and this id is
|
||||
* the only resume handle a claude-code member has — unlike opencode, nothing resolves it
|
||||
* after the launch. Checked against the real binary (claude 2.1.252): the flag is safe on
|
||||
* every spawn. The binary takes only a valid UUID — it refuses any other value at argument
|
||||
* parsing ("Invalid session ID. Must be a valid UUID.") — so the {@code UUID.randomUUID()}
|
||||
* mint is required, not incidental. The only flag interaction the binary documents is with
|
||||
* {@code -r} (both claim the session id), and the resume branch above never combines the two.
|
||||
*/
|
||||
private static String applySessionIdentity(List<String> argv, String sessionName, String resumeSessionId) {
|
||||
boolean resuming = resumeSessionId != null && !resumeSessionId.isBlank();
|
||||
boolean named = sessionName != null && !sessionName.isBlank();
|
||||
if (!resuming && !named) {
|
||||
return null; // no identity requested — keep the legacy launch byte-identical
|
||||
}
|
||||
if (named) {
|
||||
argv.add("-n");
|
||||
argv.add(sessionName);
|
||||
@@ -320,11 +324,19 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
|
||||
*
|
||||
* <p>CB-618: Claude Code refuses to start when BOTH {@code --append-system-prompt} and
|
||||
* {@code --append-system-prompt-file} are on the command line ("Cannot use both ... Please use
|
||||
* only one"), so the two charters can never travel on separate flags. When both are present they
|
||||
* are concatenated into the one file, role charter first and reply charter last — last is where
|
||||
* the reply rule must sit, because it is the rule that must survive. When only the reply charter
|
||||
* is present it keeps its proven inline {@code --append-system-prompt} delivery, which is also
|
||||
* the only form that reaches a member with no repo checkout.
|
||||
* only one"), so the two charters can never travel on separate flags. They are concatenated
|
||||
* into the one file, role charter first and reply charter last — last is where the reply rule
|
||||
* must sit, because it is the rule that must survive.
|
||||
*
|
||||
* <p>fleetd #220: a lone reply charter used to ride inline on {@code --append-system-prompt},
|
||||
* which put ~800 bytes of prose on the command line herdr types into the pane. That line is
|
||||
* capped at {@value HerdrPeerLauncher#PANE_COMMAND_BYTE_LIMIT} bytes by the pty itself, and
|
||||
* everything past the cap is dropped with no error from any layer. The charter alone left about
|
||||
* 50 bytes of headroom, so adding one flag ({@code --session-id}, fleetd #214) truncated the
|
||||
* LAST argument instead — {@code --autocompact 250000} arrived as {@code --autocompact 25},
|
||||
* claude rejected it, and every claude-code spawn died as an unexplained readiness timeout.
|
||||
* The charter now always travels as a file, which takes the prose off the command line for
|
||||
* good; {@link HerdrPeerLauncher#checkPaneCommandFits} is the backstop for whatever grows next.
|
||||
*/
|
||||
private List<String> argvWithFleet(FleetConfig.Profile cfg, LaunchSpec spec) {
|
||||
String roleCharter = nonBlank(spec.roleCharter());
|
||||
@@ -341,16 +353,13 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
|
||||
argv.add("--mcp-config");
|
||||
argv.add(mcpConfigJson(cfg));
|
||||
}
|
||||
// Combine the charters in order role -> reply, dropping any that are absent. When two
|
||||
// or more survive they must ride one --append-system-prompt-file (CB-618 forbids the inline
|
||||
// flag and the file flag together). A lone reply charter keeps its proven inline delivery.
|
||||
// Combine the charters in order role -> reply, dropping any that are absent. They ride one
|
||||
// --append-system-prompt-file (CB-618 forbids the inline flag and the file flag together),
|
||||
// always — fleetd #220: charter prose on the command line overruns the pane's byte cap.
|
||||
List<String> charters = new java.util.ArrayList<>(2);
|
||||
if (roleCharter != null) charters.add(roleCharter);
|
||||
if (replyCharter != null) charters.add(replyCharter);
|
||||
if (charters.size() == 1 && replyCharter != null && roleCharter == null) {
|
||||
argv.add("--append-system-prompt");
|
||||
argv.add(replyCharter);
|
||||
} else if (!charters.isEmpty()) {
|
||||
if (!charters.isEmpty()) {
|
||||
argv.add("--append-system-prompt-file");
|
||||
argv.add(writeCharterFile(String.join("\n\n", charters)).toString());
|
||||
}
|
||||
|
||||
@@ -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,77 @@ 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 a freshly generated directory and the flat files already written
|
||||
* into it: owner keeps full access, {@code group} gets traverse+read on the directory ({@code
|
||||
* rwxr-x---}, so a member process — a login shell reading it via {@code ZDOTDIR}, or another
|
||||
* process simply opening a file under it — running under that group can find and read 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 what fleetd generated. (For the ZDOTDIR scrub
|
||||
* specifically, this also means the scrub script's own report write inside the pane 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.)
|
||||
*
|
||||
* <p>Package-private and named generically on purpose: fleetd #213 built this for the ZDOTDIR
|
||||
* scrub directory, and fleetd #219 reuses it verbatim for {@link
|
||||
* dev.ltms.fleet.member.OpenCodeLauncher}'s ephemeral {@code opencode.json} directory — both are
|
||||
* "a fleetd-generated directory of flat files that a different-uid member process must read but
|
||||
* never write," so the sharing mechanism is shared rather than copied a second time.
|
||||
*/
|
||||
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 directory " + 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 directory " + 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;
|
||||
|
||||
@@ -676,6 +682,7 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
// Protocol 19 resolves the executable from the agent kind (== namePrefix here), so
|
||||
// argv[0] — the configured executable — is dropped and only the extra args are passed.
|
||||
List<String> args = argv.isEmpty() ? argv : argv.subList(1, argv.size());
|
||||
checkPaneCommandFits(cfg, argv);
|
||||
HerdrException last = null;
|
||||
for (int attempt = 0; attempt < NAME_RETRIES; attempt++) {
|
||||
long seq = nameSeq.incrementAndGet();
|
||||
@@ -691,6 +698,59 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
throw last;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #220: herdr does not exec the launch command — it TYPES it into the pane as one line,
|
||||
* and a pty line buffer holds only {@value #PANE_COMMAND_BYTE_LIMIT} bytes (BSD/macOS {@code
|
||||
* MAX_CANON}). Everything past that byte is dropped. Nothing reports it: herdr answers "agent
|
||||
* started", the backend exits on the mangled argument it was handed, the pane closes, and the
|
||||
* only symptom is {@link #waitUntilInjectableOrThrow} timing out 20 seconds later with no
|
||||
* reason. That is exactly how #214 broke every claude-code spawn — one 50-byte flag pushed a
|
||||
* 978-byte command to 1028, and the tail that got cut was {@code --autocompact 250000}.
|
||||
*
|
||||
* <p>So measure it here and refuse, loudly and immediately, rather than spawn something that
|
||||
* cannot work. The estimate is deliberately conservative: fleetd cannot see herdr's quoting, so
|
||||
* every argument is charged its own bytes plus a separator and a quote pair. An over-estimate
|
||||
* costs a clear error at a length that was already unsafe; an under-estimate would let the
|
||||
* silent truncation back in.
|
||||
*
|
||||
* @throws PeerUnreachableException when the command cannot fit — the same failure the spawn
|
||||
* would have hit anyway, named at the point it is still
|
||||
* explainable
|
||||
*/
|
||||
private void checkPaneCommandFits(FleetConfig.Profile cfg, List<String> argv) {
|
||||
int bytes = 0;
|
||||
String longest = null;
|
||||
int longestBytes = 0;
|
||||
for (String arg : argv) {
|
||||
int argBytes = arg == null ? 0 : arg.getBytes(java.nio.charset.StandardCharsets.UTF_8).length;
|
||||
bytes += argBytes + QUOTING_OVERHEAD_PER_ARG;
|
||||
if (argBytes > longestBytes) {
|
||||
longestBytes = argBytes;
|
||||
longest = arg;
|
||||
}
|
||||
}
|
||||
if (bytes <= PANE_COMMAND_BYTE_LIMIT) {
|
||||
return;
|
||||
}
|
||||
String culprit = longest == null ? "<none>"
|
||||
: longest.substring(0, Math.min(longest.length(), 60)) + (longest.length() > 60 ? "…" : "");
|
||||
throw new PeerUnreachableException(
|
||||
"launch command for profile " + cfg.profile() + " is about " + bytes + " bytes, over the "
|
||||
+ PANE_COMMAND_BYTE_LIMIT + "-byte limit of the pane line herdr types it into. "
|
||||
+ "The pty would drop the tail silently and the backend would exit on a mangled "
|
||||
+ "argument. Longest argument is " + longestBytes + " bytes: " + culprit
|
||||
+ " — move it off the command line (a file flag) or shorten it.");
|
||||
}
|
||||
|
||||
/**
|
||||
* The pty line buffer herdr types a launch command into: BSD/macOS {@code MAX_CANON}. Not a
|
||||
* fleetd choice and not configurable — see {@link #checkPaneCommandFits}.
|
||||
*/
|
||||
static final int PANE_COMMAND_BYTE_LIMIT = 1024;
|
||||
|
||||
/** Per-argument allowance for the separating space and a shell quote pair fleetd cannot see. */
|
||||
private static final int QUOTING_OVERHEAD_PER_ARG = 3;
|
||||
|
||||
/** Start the agent into {@code paneId}, waiting out the seed shell's boot with the sleeper. */
|
||||
private Agent startAwaitingShellPrompt(String name, List<String> args, String paneId) {
|
||||
HerdrException busy = null;
|
||||
@@ -854,20 +914,51 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
*/
|
||||
private void waitUntilInjectableOrThrow(String paneId) {
|
||||
long deadline = nowMillis.getAsLong() + spawnReadyTimeoutMs;
|
||||
Object lastStatus = null;
|
||||
while (nowMillis.getAsLong() < deadline) {
|
||||
if (agents.status(paneId).injectable()) {
|
||||
var status = agents.status(paneId);
|
||||
lastStatus = status;
|
||||
if (status.injectable()) {
|
||||
log.debug("peer pane={} reached injectable state", paneId);
|
||||
return;
|
||||
}
|
||||
sleeper.run();
|
||||
}
|
||||
log.warn("peer pane={} did not become injectable within {}ms — closing", paneId, spawnReadyTimeoutMs);
|
||||
// fleetd #220: read the pane BEFORE stop() closes it. Without this the gate says only that
|
||||
// it timed out, which is true of every cause — a backend that never launched, a binary that
|
||||
// rejected an argument and exited, a trust prompt, a login shell that hung. The pane holds
|
||||
// the one copy of that answer and it is destroyed a line later.
|
||||
log.warn("peer pane={} did not become injectable within {}ms (last status {}) — closing. "
|
||||
+ "Pane tail:\n{}",
|
||||
paneId, spawnReadyTimeoutMs, lastStatus, readPaneQuietly(paneId));
|
||||
stop(paneId);
|
||||
throw new PeerUnreachableException(
|
||||
"worker pane " + paneId + " did not reach injectable state within "
|
||||
+ spawnReadyTimeoutMs + "ms");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #220: the pane's recent output, clipped, for the readiness-gate timeout log — or a
|
||||
* short note when it cannot be read. Best-effort by construction: this runs on a path that is
|
||||
* already failing, so it must never replace the real error with one of its own.
|
||||
*/
|
||||
private String readPaneQuietly(String paneId) {
|
||||
try {
|
||||
String pane = agents.read(paneId, "recent");
|
||||
if (pane == null || pane.isBlank()) {
|
||||
return "<pane read returned nothing>";
|
||||
}
|
||||
return pane.length() <= SPAWN_FAILURE_PANE_CHARS
|
||||
? pane
|
||||
: pane.substring(pane.length() - SPAWN_FAILURE_PANE_CHARS);
|
||||
} catch (RuntimeException e) {
|
||||
return "<pane could not be read: " + e.getMessage() + ">";
|
||||
}
|
||||
}
|
||||
|
||||
/** How much of a failed spawn's pane the timeout log carries. */
|
||||
private static final int SPAWN_FAILURE_PANE_CHARS = 4000;
|
||||
|
||||
/**
|
||||
* A concrete {@link PeerHandle} wrapping herdr agent coordinates, the profile that spawned it,
|
||||
* the session identity the launch resolved (CB-547a): the bridge's logical name and the peer's
|
||||
@@ -1071,6 +1162,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 +1192,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 +1213,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 +1249,63 @@ 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.
|
||||
*
|
||||
* <p>Package-private (fleetd #219) so {@link OpenCodeLauncher} can reuse the exact same
|
||||
* "different OS user, put it under worktreeRoot instead of java.io.tmpdir" resolution for its
|
||||
* own ephemeral {@code opencode.json} directory, rather than re-reading {@code config} a second
|
||||
* time with a second copy of this null/blank handling.
|
||||
*/
|
||||
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}.
|
||||
*
|
||||
* <p>Package-private (fleetd #219) — reused by {@link OpenCodeLauncher} alongside {@link
|
||||
* #memberScrubParentDir()}; see that method's javadoc.
|
||||
*/
|
||||
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 +1363,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 +1453,71 @@ 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.
|
||||
*
|
||||
* <p>Package-private (fleetd #219) — {@link OpenCodeLauncher} reuses this same gate to decide
|
||||
* where its own ephemeral {@code opencode.json} directory (site 1) and its opencode session
|
||||
* discovery (site 2) may run, rather than re-deriving "is this a multi-uid fleet" a second way.
|
||||
*/
|
||||
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 +1546,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()
|
||||
|
||||
@@ -20,6 +20,8 @@ import java.util.EnumSet;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.function.BooleanSupplier;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.LongSupplier;
|
||||
import java.util.function.Supplier;
|
||||
@@ -215,7 +217,43 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
return Path.of(System.getProperty("java.io.tmpdir"));
|
||||
}
|
||||
|
||||
/** The default opencode storage root: {@code ~/.local/share/opencode} (the XDG data dir). */
|
||||
/**
|
||||
* The default opencode storage root: {@code ~/.local/share/opencode} (the XDG data dir) —
|
||||
* always FLEETD's OWN {@code user.home}, whichever OS user runs the daemon.
|
||||
*
|
||||
* <p><b>fleetd #219 site 2 — a decision, not a patch.</b> Under {@code memberHerdrSocket:} the
|
||||
* member pane runs as a <em>different</em> OS user, and opencode writes {@code opencode.db}
|
||||
* under <em>that</em> user's {@code $HOME}, not fleetd's. Scanning fleetd's own {@code
|
||||
* user.home} is therefore looking in the wrong place — a wrong-LOCATION failure, not a
|
||||
* wrong-PERMISSION one like site 1, and it fails quietly: {@link
|
||||
* SessionAwareHandle#agentSessionId()} would keep returning {@code null} forever, which reads
|
||||
* as "opencode does not support resume" rather than "fleetd looked in the wrong home." fleetd
|
||||
* #209 is the reason that silence is unacceptable.
|
||||
*
|
||||
* <p>Three ways to close the gap were weighed:
|
||||
* <ol>
|
||||
* <li><b>Make the member's home configurable.</b> Correct in principle, but this ticket's
|
||||
* scope is the two existing call sites, not a new config key — {@code memberHerdrSocket}
|
||||
* already carries the second herdr's socket path, not its user's home, and inventing a
|
||||
* parallel key here without also wiring it through discovery's actual callers is a
|
||||
* half-shipped feature (the exact shape CB-596/CB-611 warn against).</li>
|
||||
* <li><b>Derive it</b> (e.g. from {@code worktreeRoot}'s owner, or {@code getent passwd}).
|
||||
* Rejected: nothing in this codebase resolves a Unix username to a home directory today,
|
||||
* and guessing wrong would silently point discovery at a THIRD wrong location — worse
|
||||
* than the current gap, because it would look like it should work.</li>
|
||||
* <li><b>Declare discovery unavailable</b> under {@code memberHerdrSocket}, and say so once,
|
||||
* loudly, instead of scanning a directory that structurally cannot hold the answer.</li>
|
||||
* </ol>
|
||||
*
|
||||
* <p>Option 3 is taken — the one this ticket says to default to when unsure. {@link
|
||||
* OpenCodeLauncher#spawn} routes {@link SessionAwareHandle#agentSessionId()} through {@link
|
||||
* HerdrPeerLauncher#memberHerdrSocketConfigured()} before ever calling {@link
|
||||
* OpenCodeSessionDiscovery#sessionIdForDirectory}, so under {@code memberHerdrSocket} the
|
||||
* database at this root is never even opened, and one WARN per launcher instance names the gap
|
||||
* instead of the {@code null} return reading as "unsupported." Capability advertising is
|
||||
* unaffected: {@link #capabilities()} always includes {@code SESSION_RESUME}, since {@code
|
||||
* memberHerdrSocket} absent (today's only live mode) is unchanged by this decision.
|
||||
*/
|
||||
private static Path defaultDiscoveryRoot() {
|
||||
return Path.of(System.getProperty("user.home"), ".local", "share", "opencode");
|
||||
}
|
||||
@@ -331,7 +369,7 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
*/
|
||||
private Path writeConfig(FleetConfig.Profile cfg, String charterText, String cwd) {
|
||||
try {
|
||||
Path dir = Files.createTempDirectory(configRoot, "fleetd-opencode-");
|
||||
Path dir = Files.createTempDirectory(configParentDir(), "fleetd-opencode-");
|
||||
dir.toFile().deleteOnExit();
|
||||
|
||||
ObjectNode root = JSON.createObjectNode();
|
||||
@@ -402,6 +440,14 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
// carries operator-supplied values (URL, model id, api key), so escaping must be real.
|
||||
Files.writeString(cfgFile, JSON.writerWithDefaultPrettyPrinter().writeValueAsString(root));
|
||||
cfgFile.toFile().deleteOnExit();
|
||||
if (memberHerdrSocketConfigured()) {
|
||||
// fleetd #219: the same "different OS user" gap fleetd #213 closed for the ZDOTDIR
|
||||
// scrub — share read-only with worktreeGroup rather than leaving the directory under
|
||||
// fleetd's own 0700 java.io.tmpdir, where the member's OS user could not even
|
||||
// traverse it. memberGroup() cannot be null here: configParentDir() above already
|
||||
// refused this spawn if either worktreeRoot or worktreeGroup was missing.
|
||||
EnvAllowListScrub.shareWithGroup(dir, memberGroup());
|
||||
}
|
||||
return cfgFile;
|
||||
} catch (IOException e) {
|
||||
throw new UncheckedIOException(
|
||||
@@ -409,6 +455,57 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #219 site 1: where {@link #writeConfig} creates its per-spawn directory.
|
||||
*
|
||||
* <ul>
|
||||
* <li>{@code memberHerdrSocket} ABSENT (today's only mode): byte-identical to before this
|
||||
* fix — always {@link #configRoot} (defaults to {@code java.io.tmpdir}, fleetd's own
|
||||
* process).</li>
|
||||
* <li>{@code memberHerdrSocket} PRESENT: {@code java.io.tmpdir} is fleetd's own per-user temp
|
||||
* dir (mode {@code 0700} on macOS) — the member pane runs as a DIFFERENT OS user under
|
||||
* this config key and cannot even traverse it, so the directory holding {@code
|
||||
* opencode.json} (which tells the member where the bridge MCP is) and the member charter
|
||||
* would be unreadable to the very process it is written for. The directory instead goes
|
||||
* under {@code worktreeRoot}, shared read-only with {@code worktreeGroup} via {@link
|
||||
* EnvAllowListScrub#shareWithGroup} — the SAME mechanism fleetd #213 built for the ZDOTDIR
|
||||
* scrub, reused here rather than duplicated (see {@link
|
||||
* HerdrPeerLauncher#memberScrubParentDir()}).</li>
|
||||
* </ul>
|
||||
*
|
||||
* <p><b>Unlike the ZDOTDIR scrub, a missing {@code worktreeRoot}/{@code worktreeGroup} here
|
||||
* REFUSES the spawn instead of degrading.</b> The ZDOTDIR scrub is a credential CONTROL: a
|
||||
* degraded control (CB-596's sentinel overlay) is still worth having. This config file is not a
|
||||
* control — it is the ONLY way the member learns where the bridge MCP lives. Writing it
|
||||
* somewhere the member cannot read would not degrade anything; it would spawn a member that
|
||||
* occupies a pane and never becomes deliverable, since {@code fleet_send} waits ~60s on the
|
||||
* readiness gate and then fails with nothing pointing at a temp directory as the cause. Refusing
|
||||
* up front, with a message that names the missing config key, is the honest failure — an
|
||||
* undeliverable member is not a working spawn either way, so nothing is lost by refusing loudly
|
||||
* instead of failing silently later.
|
||||
*
|
||||
* @throws IllegalStateException when {@code memberHerdrSocket} is configured but {@code
|
||||
* worktreeRoot} and/or {@code worktreeGroup} is not
|
||||
*/
|
||||
private Path configParentDir() {
|
||||
if (!memberHerdrSocketConfigured()) {
|
||||
return configRoot;
|
||||
}
|
||||
Path root = memberScrubParentDir();
|
||||
String group = memberGroup();
|
||||
if (root == null || group == null) {
|
||||
throw new IllegalStateException("memberHerdrSocket is configured, so opencode's config "
|
||||
+ "directory (opencode.json + member charter) must be placed where the member's "
|
||||
+ "OS user can read it — worktreeRoot, shared via worktreeGroup — but "
|
||||
+ (root == null ? "worktreeRoot" : "worktreeGroup") + " is not configured. "
|
||||
+ "Refusing to spawn rather than write a config the member cannot read: that "
|
||||
+ "member would occupy a pane and never become deliverable, with nothing "
|
||||
+ "pointing at the real cause. Configure both worktreeRoot and worktreeGroup to "
|
||||
+ "enable opencode member spawns under memberHerdrSocket.");
|
||||
}
|
||||
return root;
|
||||
}
|
||||
|
||||
/**
|
||||
* Declare a custom OpenAI-compatible provider so the worker talks to a pinned endpoint (a local
|
||||
* vLLM, say) instead of opencode's default gateway (CB-508).
|
||||
@@ -517,11 +614,16 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
return afterScheme.contains("/") ? trimmed : trimmed + "/v1";
|
||||
}
|
||||
|
||||
/** One WARN per launcher instance for the fleetd #219 site-2 discovery-unavailable gap. */
|
||||
private final AtomicBoolean discoveryUnavailableWarned =
|
||||
new AtomicBoolean();
|
||||
|
||||
/** Add lazy on-disk session discovery to the base handle. */
|
||||
@Override
|
||||
public PeerHandle spawn(SpawnRequest req) {
|
||||
PeerHandle inner = super.spawn(req);
|
||||
return new SessionAwareHandle(inner, discovery, effectiveCwd(req));
|
||||
return new SessionAwareHandle(inner, discovery, effectiveCwd(req),
|
||||
this::memberHerdrSocketConfigured, discoveryUnavailableWarned);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -536,11 +638,17 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
private final PeerHandle delegate;
|
||||
private final OpenCodeSessionDiscovery discovery;
|
||||
private final String cwd;
|
||||
private final BooleanSupplier discoveryUnavailable;
|
||||
private final AtomicBoolean discoveryUnavailableWarned;
|
||||
|
||||
SessionAwareHandle(PeerHandle delegate, OpenCodeSessionDiscovery discovery, String cwd) {
|
||||
SessionAwareHandle(PeerHandle delegate, OpenCodeSessionDiscovery discovery, String cwd,
|
||||
BooleanSupplier discoveryUnavailable,
|
||||
AtomicBoolean discoveryUnavailableWarned) {
|
||||
this.delegate = delegate;
|
||||
this.discovery = discovery;
|
||||
this.cwd = cwd;
|
||||
this.discoveryUnavailable = discoveryUnavailable;
|
||||
this.discoveryUnavailableWarned = discoveryUnavailableWarned;
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -565,6 +673,22 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
|
||||
@Override
|
||||
public String agentSessionId() {
|
||||
// fleetd #219 site 2: under memberHerdrSocket the member pane runs as a different OS
|
||||
// user, so opencode.db lives under THAT user's $HOME, not the one discoveryRoot was
|
||||
// built from (see OpenCodeLauncher#defaultDiscoveryRoot's javadoc for the full
|
||||
// reasoning). Scanning fleetd's own $HOME under that config would only ever find "no
|
||||
// row" and read as "resume unsupported" — declare it unavailable instead, once, loudly.
|
||||
if (discoveryUnavailable.getAsBoolean()) {
|
||||
if (discoveryUnavailableWarned.compareAndSet(false, true)) {
|
||||
log.warn("opencode session discovery unavailable: memberHerdrSocket is "
|
||||
+ "configured, so opencode's on-disk session database lives under the "
|
||||
+ "MEMBER's own $HOME, not fleetd's ({}) — agentSessionId will stay null "
|
||||
+ "for every opencode member under this config, and SESSION_RESUME "
|
||||
+ "cannot be honored (fleetd #209/#219).",
|
||||
System.getProperty("user.home"));
|
||||
}
|
||||
return null;
|
||||
}
|
||||
// Lazy + retried, never a spawn-time blocker: opencode writes the session record only
|
||||
// when the session is first persisted, so null here is the correct interim answer and
|
||||
// the caller re-calls later (each call re-scans, picking up a record that has since
|
||||
|
||||
@@ -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,115 @@ 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());
|
||||
assertTrue(waiter.getNow(null).text().contains("--- pane tail ---"),
|
||||
"fleetd#164: the failure must carry the pane, not only the matched line — with an "
|
||||
+ "empty trimmed block the raw scrape is the only copy of what the member said: "
|
||||
+ waiter.getNow(null).text());
|
||||
assertTrue(waiter.getNow(null).text().contains("╭"),
|
||||
"the carried pane is the raw scrape, chrome included: " + 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])",
|
||||
|
||||
@@ -61,8 +61,75 @@ class ClaudeCodeLauncherTest {
|
||||
assertTrue(args.stream().noneMatch(a -> a.contains("\"bridge\"")),
|
||||
"the mount is named fleet since CB-632 — a member addresses its tools as "
|
||||
+ "mcp__fleet__*, and CLAUDE.md's role-detection ladder names that prefix");
|
||||
assertTrue(args.contains("--append-system-prompt"));
|
||||
assertTrue(args.stream().anyMatch(a -> a.contains("fleet_reply")), "reply charter present");
|
||||
// fleetd #220: the charter travels as a FILE, never inline — charter prose on the command
|
||||
// line overruns the byte cap of the pane line herdr types it into.
|
||||
assertTrue(args.contains("--append-system-prompt-file"));
|
||||
assertFalse(args.contains("--append-system-prompt"),
|
||||
"the inline flag would put ~800 bytes of prose on the pane command line");
|
||||
String charterFile = args.get(args.indexOf("--append-system-prompt-file") + 1);
|
||||
assertTrue(readFile(charterFile).contains("fleet_reply"), "reply charter present in the file");
|
||||
}
|
||||
|
||||
/** Read a charter file the launcher wrote, failing the test rather than the build on an IO error. */
|
||||
private static String readFile(String path) {
|
||||
try {
|
||||
return java.nio.file.Files.readString(java.nio.file.Path.of(path));
|
||||
} catch (java.io.IOException e) {
|
||||
throw new AssertionError("charter file " + path + " is not readable", e);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #220 regression. herdr TYPES the launch command into the pane, and the pty line buffer
|
||||
* holds 1024 bytes — past that the tail is dropped with no error anywhere, so the backend exits
|
||||
* on a mangled argument and the spawn dies as an unexplained readiness timeout. That is what
|
||||
* happened when #214 added --session-id to a command already 978 bytes long: --autocompact
|
||||
* 250000 arrived as --autocompact 25. This asserts the whole assembled command still fits, with
|
||||
* the flags a real spawn carries (MCP mount, charter, model, autocompact, session id).
|
||||
*/
|
||||
@Test
|
||||
void theAssembledLaunchCommandFitsThePaneLineLimit() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
// Shaped like the live sonnet profile, because the bug is a SUM: the charter alone fits,
|
||||
// and so does every flag alone. Only model + autocompact + session id on top of the charter
|
||||
// crossed the cap, which is why nothing caught it until a member failed to spawn.
|
||||
FleetConfig.Profile cfg = new FleetConfig.Profile(
|
||||
"sonnet", null, "claude-sonnet-5", null, null,
|
||||
List.of("claude"), "tab", "fleetd-workers", "w #{n}", "http://127.0.0.1:8765/mcp",
|
||||
null, null, null, null, null, null, null, null, true, null, null, null, null, null,
|
||||
250000);
|
||||
new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of()), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null)
|
||||
.spawn();
|
||||
|
||||
List<String> args = spawnedArgs(herdr);
|
||||
int bytes = "claude".length();
|
||||
for (String arg : args) {
|
||||
bytes += arg.getBytes(java.nio.charset.StandardCharsets.UTF_8).length + 3;
|
||||
}
|
||||
assertTrue(bytes <= HerdrPeerLauncher.PANE_COMMAND_BYTE_LIMIT,
|
||||
"the launch command must fit the pane line: " + bytes + " bytes vs limit "
|
||||
+ HerdrPeerLauncher.PANE_COMMAND_BYTE_LIMIT + " — args " + args);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #220: the guard refuses a command that cannot fit, instead of letting the pty drop the
|
||||
* tail. The refusal must name the size and the argument to blame — a spawn that fails with
|
||||
* "did not reach injectable state" tells the operator nothing, which is the whole reason this
|
||||
* bug took a live pane scrape to find.
|
||||
*/
|
||||
@Test
|
||||
void anOverlongLaunchCommandIsRefusedWithTheSizeAndTheCulprit() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
String huge = "x".repeat(1500);
|
||||
ClaudeCodeLauncher launcher = service(herdr, List.of("claude", huge), null);
|
||||
|
||||
PeerUnreachableException refused = assertThrows(PeerUnreachableException.class, launcher::spawn);
|
||||
|
||||
assertTrue(refused.getMessage().contains("1024"), "names the limit: " + refused.getMessage());
|
||||
assertTrue(refused.getMessage().contains("1500"), "names the culprit's size: " + refused.getMessage());
|
||||
assertFalse(herdr.called("agent.start"),
|
||||
"nothing may be started — a truncated command is worse than no spawn");
|
||||
}
|
||||
|
||||
// CB-634: a profile with ideMcpUrl set mounts the IDE Index MCP as a second server and pins
|
||||
@@ -262,10 +329,16 @@ class ClaudeCodeLauncherTest {
|
||||
|
||||
Map<?, ?> start = (Map<?, ?>) herdr.lastCall("agent.start").params();
|
||||
assertEquals("claude", start.get("kind"), "herdr launches the canonical executable by kind");
|
||||
// CB-533: the shared fixture pins model "coder", so the model flag is the whole args list.
|
||||
// What this test guards is that argv[0] is NOT repeated — herdr supplies it from `kind`.
|
||||
assertEquals(List.of("--model", "coder"), start.get("args"),
|
||||
"the configured executable is not repeated in args");
|
||||
// fleetd #214: a plain spawn now also carries fleetd's own minted session id.
|
||||
List<String> args = spawnedArgs(herdr);
|
||||
assertFalse(args.contains("claude"), "the configured executable is not repeated in args");
|
||||
int flag = args.indexOf("--session-id");
|
||||
assertTrue(flag >= 0, "a plain spawn mints a session id: " + args);
|
||||
assertDoesNotThrow(() -> UUID.fromString(args.get(flag + 1)), "the minted id is a valid UUID");
|
||||
assertEquals(List.of("--model", "coder"),
|
||||
List.of(args.get(args.size() - 2), args.get(args.size() - 1)),
|
||||
"the CB-533 model flag still trails the launch flags: " + args);
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -278,7 +351,15 @@ class ClaudeCodeLauncherTest {
|
||||
assertFalse(args.contains("--append-system-prompt"), "no reply charter without mcpUrl");
|
||||
// CB-533: the model flag is independent of the MCP mount — pinning the model is not part of
|
||||
// "mount the bridge", so an unmounted worker still runs the model its profile names.
|
||||
assertEquals(List.of("--verbose", "--model", "coder"), args,
|
||||
// fleetd #214: a plain spawn now carries fleetd's own minted session id; drop the two
|
||||
// mint elements when checking the rest of the argv.
|
||||
int flag = args.indexOf("--session-id");
|
||||
assertTrue(flag >= 0, "a plain spawn mints a session id even without a bridge mount: " + args);
|
||||
assertDoesNotThrow(() -> UUID.fromString(args.get(flag + 1)), "the minted id is a valid UUID");
|
||||
List<String> rest = new java.util.ArrayList<>(args);
|
||||
rest.remove(flag + 1);
|
||||
rest.remove(flag);
|
||||
assertEquals(List.of("--verbose", "--model", "coder"), rest,
|
||||
"the operator's own args are preserved, in order, ahead of the model flag");
|
||||
}
|
||||
|
||||
@@ -649,7 +730,7 @@ class ClaudeCodeLauncherTest {
|
||||
assertEquals("/work/proj", cwd, "effectiveCwd via SpawnRequest must match the three-arg resolution");
|
||||
}
|
||||
|
||||
// --- CB-547a: durable session identity (mint / resume / no-identity legacy) -----------------
|
||||
// --- CB-547a / fleetd #214: durable session identity (always mint / resume) -----------------
|
||||
|
||||
@Test
|
||||
void freshSpawnMintsASessionIdAndPassesTheName() {
|
||||
@@ -684,18 +765,28 @@ class ClaudeCodeLauncherTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
void noIdentitySpawnKeepsTheLegacyArgvAndCarriesNoSessionHandle() {
|
||||
void plainSpawnMintsASessionIdSoEveryMemberIsResumable() {
|
||||
// fleetd #214: a plain spawn passes no sessionName and no resumeSessionId, yet the member
|
||||
// must still be resumable — the id is minted unconditionally, and it is the ONLY resume
|
||||
// handle a claude-code member has (unlike opencode, nothing resolves it after the launch).
|
||||
// The binary requires a valid UUID (checked against claude 2.1.252: a non-UUID is refused
|
||||
// at argument parsing with "Invalid session ID. Must be a valid UUID.").
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
ClaudeCodeLauncher svc = service(herdr, List.of("ccs", "ltms-local"), null);
|
||||
|
||||
PeerHandle handle = svc.spawn(new SpawnRequest("ltms-local", null, null));
|
||||
|
||||
List<String> args = spawnedArgs(herdr);
|
||||
assertFalse(args.contains("--session-id"), "no identity → no --session-id");
|
||||
assertFalse(args.contains("-n"), "no identity → no -n");
|
||||
assertFalse(args.contains("-r"), "no identity → no -r");
|
||||
assertNull(handle.agentSessionId(), "no identity → no resume handle");
|
||||
assertNull(handle.sessionName(), "no identity → no logical name");
|
||||
int flag = args.indexOf("--session-id");
|
||||
assertTrue(flag >= 0 && flag + 1 < args.size(),
|
||||
"--session-id is minted even when no identity is requested: " + args);
|
||||
String minted = args.get(flag + 1);
|
||||
assertDoesNotThrow(() -> UUID.fromString(minted), "--session-id is a valid UUID: " + minted);
|
||||
assertEquals(minted, handle.agentSessionId(),
|
||||
"the resume handle is the minted id, so every member is resumable from fleet_list");
|
||||
assertFalse(args.contains("-n"), "no sessionName was requested → no -n");
|
||||
assertFalse(args.contains("-r"), "no resume was requested → no -r");
|
||||
assertNull(handle.sessionName(), "no sessionName was requested → no logical name");
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
+385
-1
@@ -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() {
|
||||
|
||||
@@ -1,5 +1,9 @@
|
||||
package dev.ltms.fleet.member;
|
||||
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.Logger;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
@@ -14,9 +18,13 @@ import dev.ltms.fleet.peer.PeerUnreachableException;
|
||||
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.nio.file.attribute.PosixFilePermissions;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.ExecutorService;
|
||||
@@ -25,6 +33,7 @@ import java.util.concurrent.Future;
|
||||
import java.util.function.Supplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
import static org.junit.jupiter.api.Assumptions.assumeTrue;
|
||||
|
||||
/**
|
||||
* The opencode adapter's launch build: a file-based MCP mount + reply-charter instructions (no
|
||||
@@ -610,4 +619,196 @@ class OpenCodeLauncherTest {
|
||||
assertTrue(json.path("mcp").path("intellij").isMissingNode(),
|
||||
"no IDE server when ideMcpUrl is unset");
|
||||
}
|
||||
|
||||
// --- fleetd #219: config root + discovery root under memberHerdrSocket ------------------------
|
||||
|
||||
/** A config with {@code memberHerdrSocket:} set, and optionally {@code worktreeRoot:}/{@code worktreeGroup:}. */
|
||||
private static FleetConfig configWithMemberHerdrSocket(String worktreeRoot, String worktreeGroup) {
|
||||
return new FleetConfig(
|
||||
null, // bind
|
||||
null, // herdrSocket
|
||||
"/tmp/other-user.sock", // 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
|
||||
null // memberLoginShell
|
||||
).withDefaults();
|
||||
}
|
||||
|
||||
private static OpenCodeLauncher serviceWithConfig(FakeHerdr herdr, Path configRoot, Path discoveryRoot,
|
||||
FleetConfig.Profile cfg, Supplier<FleetConfig> config) {
|
||||
return new OpenCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null,
|
||||
0, System::currentTimeMillis, () -> { }, configRoot, discoveryRoot, null, null, config);
|
||||
}
|
||||
|
||||
/** 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();
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #219 site 1, acceptance criterion 1: with {@code memberHerdrSocket} configured and both
|
||||
* {@code worktreeRoot}/{@code worktreeGroup} set, the generated {@code opencode.json} directory
|
||||
* lives under {@code worktreeRoot} — NEVER under the injected {@code configRoot} (standing in for
|
||||
* {@code java.io.tmpdir}, fleetd's own 0700 temp dir, unreadable by the member's different OS
|
||||
* user) — and is shared read-only with the group via the SAME mechanism (fleetd #213's {@link
|
||||
* EnvAllowListScrub#shareWithGroup}) the ZDOTDIR scrub uses.
|
||||
*/
|
||||
@Test
|
||||
void memberHerdrSocketWithWorktreeRootAndGroupPutsConfigDirUnderWorktreeRootAndSharesIt(
|
||||
@TempDir Path configRoot, @TempDir Path worktreeRoot) throws Exception {
|
||||
String group = currentUserGroup();
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
serviceWithConfig(herdr, configRoot, configRoot,
|
||||
opencodeCfg("google/gemini-2.5-pro", "http://127.0.0.1:8765/mcp", null),
|
||||
() -> configWithMemberHerdrSocket(worktreeRoot.toString(), group)).spawn();
|
||||
|
||||
String cfgPath = startEnv(herdr).get("OPENCODE_CONFIG");
|
||||
assertNotNull(cfgPath, "the profile still needs a config file");
|
||||
Path cfgFile = Path.of(cfgPath);
|
||||
Path dir = cfgFile.getParent();
|
||||
assertEquals(worktreeRoot.toAbsolutePath().normalize(), dir.getParent(),
|
||||
"the generated directory's parent must be worktreeRoot, not the injected configRoot "
|
||||
+ "standing in for java.io.tmpdir — got parent " + dir.getParent());
|
||||
assertFalse(dir.startsWith(configRoot),
|
||||
"the generated directory must NOT be created under configRoot when memberHerdrSocket "
|
||||
+ "is configured: " + dir);
|
||||
|
||||
assertEquals("rwxr-x---", PosixFilePermissions.toString(Files.getPosixFilePermissions(dir)),
|
||||
"the directory must be group-traversable+readable, owner-only writable");
|
||||
assertEquals("rw-r-----", PosixFilePermissions.toString(Files.getPosixFilePermissions(cfgFile)),
|
||||
"opencode.json must be group-readable, never group-writable");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #219 site 1, acceptance criterion 2: with {@code memberHerdrSocket} configured but
|
||||
* NEITHER {@code worktreeRoot} nor {@code worktreeGroup} set, the launcher must refuse the spawn
|
||||
* rather than write a config under {@code java.io.tmpdir} the member cannot read — that member
|
||||
* would occupy a pane and never become deliverable, with nothing pointing at the real cause.
|
||||
*/
|
||||
@Test
|
||||
void memberHerdrSocketWithoutWorktreeRootOrGroupRefusesTheSpawn(@TempDir Path configRoot) {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
OpenCodeLauncher launcher = serviceWithConfig(herdr, configRoot, configRoot,
|
||||
opencodeCfg("google/gemini-2.5-pro", "http://127.0.0.1:8765/mcp", null),
|
||||
() -> configWithMemberHerdrSocket(null, null));
|
||||
|
||||
IllegalStateException ex = assertThrows(IllegalStateException.class,
|
||||
() -> launcher.spawn(new SpawnRequest(null, null, null)),
|
||||
"a missing worktreeRoot/worktreeGroup must refuse the spawn, not write an unreadable config");
|
||||
assertTrue(ex.getMessage().contains("worktreeRoot"),
|
||||
"the refusal must name the missing config key — got: " + ex.getMessage());
|
||||
assertFalse(herdr.called("tab.create"),
|
||||
"the spawn must be refused BEFORE any pane is created — got calls: " + herdr.calls);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #219 site 1, acceptance criterion 2 (the other missing half): {@code worktreeRoot} set
|
||||
* but {@code worktreeGroup} missing must ALSO refuse — either one alone is not enough to
|
||||
* guarantee the member's OS user can read the directory.
|
||||
*/
|
||||
@Test
|
||||
void memberHerdrSocketWithWorktreeRootButNoGroupRefusesTheSpawn(
|
||||
@TempDir Path configRoot, @TempDir Path worktreeRoot) {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
OpenCodeLauncher launcher = serviceWithConfig(herdr, configRoot, configRoot,
|
||||
opencodeCfg("google/gemini-2.5-pro", "http://127.0.0.1:8765/mcp", null),
|
||||
() -> configWithMemberHerdrSocket(worktreeRoot.toString(), null));
|
||||
|
||||
IllegalStateException ex = assertThrows(IllegalStateException.class,
|
||||
() -> launcher.spawn(new SpawnRequest(null, null, null)));
|
||||
assertTrue(ex.getMessage().contains("worktreeGroup"),
|
||||
"worktreeRoot alone is not enough — got: " + ex.getMessage());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #219 site 1, acceptance criterion 3: with {@code memberHerdrSocket} ABSENT — even when
|
||||
* a live, non-null {@code config} supplier is threaded through (not merely {@code config == null},
|
||||
* which every other test in this file already exercises) — the generated directory must still
|
||||
* land directly under the injected {@code configRoot}, byte-identical to before this fix.
|
||||
*/
|
||||
@Test
|
||||
void memberHerdrSocketAbsentStaysUnderConfigRootEvenWithALiveConfigSupplier(@TempDir Path configRoot)
|
||||
throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FleetConfig config = new FleetConfig(null, null, null, Map.of(), null, null, null, null, null,
|
||||
null, null, null, null, null, null, null, null, null, null, null, null, null).withDefaults();
|
||||
serviceWithConfig(herdr, configRoot, configRoot,
|
||||
opencodeCfg("google/gemini-2.5-pro", "http://127.0.0.1:8765/mcp", null), () -> config)
|
||||
.spawn();
|
||||
|
||||
String cfgPath = startEnv(herdr).get("OPENCODE_CONFIG");
|
||||
assertNotNull(cfgPath);
|
||||
assertTrue(Path.of(cfgPath).startsWith(configRoot),
|
||||
"with memberHerdrSocket absent, the config directory must still be created directly "
|
||||
+ "under configRoot, unchanged from before this fix");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #219 site 2: under {@code memberHerdrSocket}, opencode session discovery must be
|
||||
* declared unavailable rather than silently scanning fleetd's own {@code discoveryRoot} — which,
|
||||
* under this config key, is NOT where the member's opencode actually writes its session
|
||||
* database. This test proves the gate is real, not merely "no record yet": a matching record IS
|
||||
* written to {@code discoveryRoot} (the exact fixture {@link
|
||||
* #theHandleDiscoversTheSessionIdForTheWorkersCwdOnlyAfterItAppears} proves discovery would
|
||||
* otherwise find), and {@code agentSessionId()} must still return {@code null} — proving the
|
||||
* gate, not a coincidental absence of data, is what produced the null. One WARN is also logged,
|
||||
* exactly once even across repeated calls.
|
||||
*/
|
||||
@Test
|
||||
void discoveryIsUnavailableUnderMemberHerdrSocketEvenWhenARecordExists(
|
||||
@TempDir Path configRoot, @TempDir Path worktreeRoot, @TempDir Path discRoot) throws Exception {
|
||||
String group = currentUserGroup();
|
||||
OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "ses_should_be_hidden", "/work/dir", 1000L);
|
||||
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
OpenCodeLauncher launcher = new OpenCodeLauncher(new AgentControl(herdr),
|
||||
new WorkspaceControl(herdr), Map.of("gemini", opencodeCfg(null, null, null)),
|
||||
"gemini", _ -> null, 0, System::currentTimeMillis, () -> { }, configRoot, discRoot,
|
||||
null, null, () -> configWithMemberHerdrSocket(worktreeRoot.toString(), group));
|
||||
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(OpenCodeLauncher.class);
|
||||
Level original = logger.getLevel();
|
||||
logger.setLevel(Level.INFO);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
PeerHandle handle;
|
||||
try {
|
||||
handle = launcher.spawn(new SpawnRequest(null, "/work/dir", null));
|
||||
assertNull(handle.agentSessionId(),
|
||||
"memberHerdrSocket configured: discovery must stay unavailable even though a "
|
||||
+ "matching record exists in discoveryRoot");
|
||||
assertNull(handle.agentSessionId(), "the gate must hold on a second call too");
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
logger.setLevel(original);
|
||||
}
|
||||
List<String> warnings = appender.list.stream()
|
||||
.filter(e -> e.getLevel() == Level.WARN)
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.toList();
|
||||
assertEquals(1, warnings.size(),
|
||||
"exactly one WARN across two agentSessionId() calls — got: " + warnings);
|
||||
assertTrue(warnings.get(0).contains("memberHerdrSocket"),
|
||||
"the WARN must name memberHerdrSocket as the reason — got: " + warnings.get(0));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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());
|
||||
|
||||
@@ -887,15 +887,18 @@ class SessionManagerTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
void acquireWithNeitherSessionFieldLeavesAgentSessionIdNull() {
|
||||
void acquireWithNeitherSessionFieldStillMintsAnAgentSessionId() {
|
||||
// fleetd #214: the claude-code launcher mints a session id for EVERY spawn, so a member is
|
||||
// resumable even when the spawn asked for no session identity.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
SessionManager sessions = sessionManager(herdr);
|
||||
|
||||
MemberSession s = sessions.acquire("ltms-local", null, null, null);
|
||||
|
||||
assertNull(s.agentSessionId(), "no identity requested — unchanged from before CB-584");
|
||||
assertFalse(SessionManager.rosterView(s, null).containsKey("agentSessionId"),
|
||||
"a null id is omitted from the roster, like charterSha256 for a receipt-less session");
|
||||
assertNotNull(s.agentSessionId(),
|
||||
"fleetd #214: a plain spawn mints a session id, so every member is resumable");
|
||||
assertTrue(SessionManager.rosterView(s, null).containsKey("agentSessionId"),
|
||||
"the minted id is in the roster, so fleet_list advertises every member's resume handle");
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -931,4 +934,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");
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user