Compare commits

..

14 Commits

Author SHA1 Message Date
Dai Ha de026b8f8a #137 follow-up: document defect-1 and defect-2 reachability, drop unconstructible test
CI / contract (pull_request) Successful in 48s
CI / build (pull_request) Successful in 1m28s
Defect 1 (reply()'s askAnsweredAsyncTasks returning >1 candidate): confirmed
unreachable today. Documented why in three places — hasAsyncQuestion matches
any task with a stamped turnId (not just an open question), and answer()'s
clearAsyncQuestion(turnId, false) leaves that stamp in place until the
resumed turn's own future resolves — so a second task can never reach the
same eligible state while a first one holds it. Kept the defensive
inbox-fallback branch as defence in depth against that guarantee weakening,
per review instruction; no test seam added.

Defect 2 (abandon()'s broader `matching` filter applying a stranded reply to
more than one task): matching.size() >= 2 alone IS reachable (already
covered by abandonFailsEveryPendingAsyncTicketForTheReleasedTarget) and was
a real pre-fix bug (97f6c33's parent reused one drained reply for every
matching task). But hadStrandedReply == true together with matching.size()
>= 2, at the instant abandon() runs, is not constructible through the public
API: send() and answer() are the only two sites that ever hold a target's
session lock, and both open a Rendezvous waiter for that target as the first
thing they do while holding it — so "lock held" and "waiter open" are the
same fact throughout this class, and reply()'s fast path always resolves an
open waiter directly instead of stranding. A strand can only be created
while no task is accepted, and the moment the lock is next taken, that
acceptance clears the strand again before abandon() can observe both facts
together. Documented this in abandon()'s javadoc and removed the earlier
attempt at a deterministic test for the conjunction, whose apparent failure
was an invalid premise (the "accepted" task's own acceptance silently
cleared the strand it was meant to race against), not the fix being absent.
2026-08-31 23:17:50 +07:00
Dai Ha 97f6c33a45 #137: don't guess when a target has more than one open async task
CI / contract (pull_request) Successful in 1m4s
CI / build (pull_request) Successful in 1m17s
The CB-205-recovery fix in #205 assumed a target has at most one open
async task, with no guard. Fix three consequences:

- reply(): askAnsweredAsyncTask -> askAnsweredAsyncTasks (List). Exactly
  one candidate completes it (unchanged). Zero falls to the inbox
  (unchanged). More than one now ALSO falls to the inbox instead of
  picking an arbitrary ConcurrentHashMap iteration order, and logs a
  WARN naming the target and every candidate ticket.

- abandon(): a stranded reply now settles at most one matching task —
  the oldest by Task#createdNanos (a new field, the tiebreaker). Every
  other matching task keeps WORKER_FAILED, same as today.

- abandon(): if the chosen recovery task's complete() loses a race
  (another path resolved it first), the drained reply is republished
  to the inbox instead of being silently dropped.

Single-task behavior is unchanged; only the ambiguous case changes.
2026-08-31 22:46:23 +07:00
Dai Ha 7d4a4339c2 Merge branch 'worker/cb-137-ask-ticket-e7760c-2' of https://git.ltms.dev/fleet/fleetd into worker/cb137-ambiguous-task-4df3d8-4 2026-08-31 22:13:37 +07:00
Dai Ha a55079afbd #185: opt-in worktreeGroup, so a member running as another OS user can write its worktree
CI / contract (push) Successful in 50s
CI / build (push) Successful in 1m39s
Stage 3 of #185. A provisioned worktree and the repo's git store are made
group-writable when worktreeGroup names an OS group; absent, nothing changes.
The share pass runs after overlayParity, not inside add(), because
overlayParity copies more files in after add() returns.

This isolates credentials, not the repository: a member in the group can
still write the operator's git objects and refs.
2026-08-31 21:47:03 +07:00
Dai Ha e18ad4723b Merge remote-tracking branch 'refs/remotes/origin/cb206' 2026-08-31 21:47:03 +07:00
Dai Ha 6de8ac8972 #206: read opencode session ids from opencode.db, and pin the read-only open
opencode migrated its session store from a JSON file tree to SQLite in
January. OpenCodeSessionDiscovery still scanned the frozen tree, so it
returned null for every member: agentSessionId was never known and
resumeSessionId silently did nothing for every opencode profile, through
57 member spawns, with nothing logging that the search found nothing.
2026-08-31 21:46:32 +07:00
Dai Ha a052975420 #206: pin the read-only open with a test that actually fails without it
CI / build (pull_request) Successful in 1m18s
CI / contract (pull_request) Successful in 1m44s
setReadOnly(true) is the whole thing keeping fleetd out of the operator's
live 841MB opencode.db, and no test failed when it was removed.

The obvious test does not work. Making the database file unwritable and
checking the read still succeeds passes either way, because SQLite silently
downgrades a read-write open of an unwritable file to read-only. I wrote that
test, watched it pass with the flag removed, and threw it away.

What works: extract a package-private openReadOnly(), then ask that connection
to INSERT and require the refusal. Watched red with the flag removed, green
with it restored.

Also switches the test's INSERT helper to a PreparedStatement -- hand-escaped
SQL in a test is a pattern that gets copied into main code.
2026-08-31 21:46:07 +07:00
Dai Ha 6d82ca95a4 #185: skip absent git paths, and resolve the git dir instead of assuming .git
CI / contract (pull_request) Successful in 56s
CI / build (pull_request) Successful in 1m38s
Two defects in the stage-3 share pass, both of which would have failed EVERY
provisioning spawn once worktreeGroup was set, not only the two-user case.

.git/logs was handed to chgrp unguarded while packed-refs was guarded. It does
not exist with core.logAllRefUpdates=false, or before the first ref update, and
chgrp on a missing path exits non-zero -- surfacing as a WorktreeException that
blames a group which is in fact fine. Every path is now skipped when absent.

repoRoot + "/.git" was hardcoded. That is a FILE, not a directory, when the
checkout is itself a linked worktree -- the very thing this class creates for
every member. It now asks git: rev-parse --git-common-dir, resolved against
repoRoot because git answers relatively for an ordinary checkout.

Both new tests were watched failing with the fix removed before being kept.
2026-08-31 21:34:49 +07:00
Dai Ha 8067ee4ec4 fleetd #185: opt-in worktreeGroup config for group-shared worktrees
CI / contract (pull_request) Successful in 1m7s
CI / build (pull_request) Successful in 1m45s
Adds worktreeGroup (top-level FleetConfig key), Worktrees.shareWithGroup
(GitWorktrees impl: git config core.sharedRepository group + one-time
chgrp/chmod g+rwX/setgid fix-up over the worktree, .git/objects, refs,
logs, worktrees, and packed-refs when present), and wires SessionManager
to call it AFTER overlayParity so overlay files are covered too. Off by
default (byte-identical behaviour when unset). Documents the
credentials-not-repository caveat in the javadoc and example config.
2026-08-31 15:53:01 +07:00
Dai Ha 3743789e8d fleetd #206: read opencode session ids from opencode.db (SQLite), not the frozen JSON tree
CI / contract (pull_request) Successful in 1m5s
CI / build (pull_request) Successful in 1m44s
opencode migrated its session store to SQLite in January 2026; the JSON tree under
storage/session/<projectID>/ses_*.json stopped being written, so OpenCodeSessionDiscovery
returned null for every member forever, and fleet_spawn{resumeSessionId} was unreachable.

- Add org.xerial:sqlite-jdbc 3.53.4.0, opened read-only (SQLiteConfig.setReadOnly), so it
  never disturbs a live opencode process writing the WAL-mode database.
- Rewrite sessionIdForDirectory to run a parameterized SELECT ... WHERE directory = ?
  ORDER BY time_updated DESC LIMIT 1 against the session table. Still never throws: a
  missing database, a locked/corrupt one, or no matching row all return null.
- Log the silence that let this go unnoticed: WARN once per instance when opencode.db
  itself is missing (the layout moved again), DEBUG when it exists but no row matches
  yet (the normal interim answer right after a spawn).
- Replace the JSON-fixture tests with a synthetic-SQLite-db fixture; delete the tests
  that only proved the old JSON scan worked.
2026-08-31 15:52:06 +07:00
Dai Ha 388aba7632 #137: an answered turn's reply completes its own ticket, instead of a false failure
CI / contract (push) Successful in 1m12s
CI / build (push) Successful in 1m17s
A wait:false delegation whose worker used fleet_ask ended with fleet_poll{ticket}
reporting 'the worker session was released before it replied' -- naming a
worktree, a branch and a snapshot commit, so it read as lost work. The worker had
in fact replied in full.

The ticket guessed the ask rendezvous detached the turn. The real cause is
narrower: answer() (behind fleet_send{turnId}) waits only for the lead's own
bounded MCP call window. A resumed turn doing real work -- edits, a build, a
push, a PR -- routinely outlives it. On timeout answer()'s finally closed the
waiter, so the worker's later fleet_reply found none and fell to the session
inbox, leaving the ticket's future unresolved until fleet_stop forced it FAILED.

reply() now looks for the async task parked on this exact answered turn and
completes it with the real reply. That is safe against completing the wrong
ticket: answer() calls clearAsyncQuestion(turnId, false), so the task keeps its
turnId and stays in asyncTasksByTurn, and hasAsyncQuestion therefore still makes
send() return BUSY for a second async send to that target. At most one candidate
task can exist per target.

abandon() keeps an independent check: if a reply was stranded, a released session
reports REPLIED with that text rather than a failure -- so the recovery hint that
implies lost work never prints once a reply exists.

Verified before merging: build green unpiped, and both new tests drive the full
delegation path (send async, ask, answer, reply, poll the ticket) rather than
handing a Reply to a sink, which is the trap this ticket called out.

Co-authored-by: fleetd worker <worker@ltms.dev>
2026-08-31 14:24:54 +07:00
Dai Ha ad587eafa3 #172: keep the broker URI, password and all, out of every member pane
CI / contract (push) Successful in 1m10s
CI / build (push) Successful in 1m15s
broker.uriEnv names an environment variable holding amqp://user:password@host,
and it was reaching every member. Its name is not credential-shaped -- no TOKEN,
KEY or SECRET in it -- so every name-pattern heuristic missed it, and it sat on
neither credential list.

fleetd already knows the name: the operator wrote it in broker.uriEnv. So derive
the exclusion from the config rather than hoping an operator also remembers to
deny it. coordinator.uriEnv has the same shape and is excluded too; on this host
both resolve to the same variable.

Excluded even when the operator lists the name under memberCredentials.allow:,
following the SSH_AUTH_SOCK precedent. There is no override, because a member has
no legitimate use for the broker password.

Reviewer finding, recorded rather than overstated: this is only a hard guarantee
under policy: allow-list, where the ZDOTDIR scrub runs after the pane's shell has
sourced the operator's chain. Under deny-list the name is removed from the
pre-shell env only, and a login shell re-exports it. That is deny-list's existing
weakness rather than a regression here, but the javadoc now says so plainly
instead of implying a guarantee that path cannot give.

Co-authored-by: fleetd worker <worker@ltms.dev>
2026-08-31 14:21:48 +07:00
Dai Ha 08ce9aef11 #161: resolve a pane by process ancestry, closing a worker->primary escalation
CI / contract (push) Successful in 50s
CI / build (push) Successful in 1m52s
PaneLocator's javadoc always claimed it found 'the agent pane whose process
tree contains' a pid. It did not: paneOwnsPid matched only the pane's shell_pid
and its foreground_processes. A process a member spawned -- python3, curl, any
helper opening its own connection to 127.0.0.1:8765 -- matched no pane, so
CallerResolver fell through to loopback-trust and resolved it as the PRIMARY.
A member escalated to lead by shelling out.

terminalForPid now builds the caller's ancestor set once (bounded at 32
generations, with a cycle guard) and matches any ancestor against a pane's pids.
The set is reused across both herdr clients on the CB-185 two-daemon path.

This only ever ADDS matches, which is the safe direction: the failure mode of
the fix is a member correctly restricted, while the failure mode of the bug is a
member acting as the lead. The no-match case still returns null, so the lead --
which maps to a pane named by leaders: -- still resolves as primary.

Ancestry is walked through a new ParentResolver seam so the tests drive it from
a fake pid->parent map rather than spawning real processes.

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

- MessageService.reply(): before falling to the inbox, look for the async
  task this exact turn belongs to (already answered — question cleared,
  turnId still stamped — but not yet resolved) and complete it directly
  with the real reply, so fleet_poll{ticket} returns it.
- MessageService.abandon(): defense in depth, independent of the above —
  never write a false WORKER_FAILED once a reply reached the inbox for
  this target; recover and use its real content instead.
- Two new tests drive the full delegation path (async send -> ask ->
  answer with a short timeout -> reply -> poll/abandon), not a reply sink
  directly; both fail with the fix disabled and pass with it restored.
2026-08-31 14:12:23 +07:00
30 changed files with 1572 additions and 388 deletions
+166
View File
@@ -0,0 +1,166 @@
# CB-137 / fleetd issue #137 — report
## Real root cause (not the hypothesis in the ticket)
I read `MessageService.java` and `Rendezvous.java` before changing anything. The mechanism is real,
but the exact place it happens is `MessageService.answer()`, not "the reply goes to the inbox on
purpose" in general.
1. A lead delegates with `fleet_send{wait:false}` → `sendAsync()` creates a `Task` and runs `send()`
on a background virtual thread with a 30-minute internal budget (`ASYNC_TIMEOUT_MS`).
2. The worker calls `fleet_ask`. That resolves the open rendezvous waiter with `Kind.QUESTION`, so
`send()` returns immediately and the `Task` is left open (its `future` stays unresolved — see the
comment in `sendAsync`'s lambda: "Keep the accepted owner until answer() finishes it").
3. The lead answers with `fleet_send{turnId, content}`. This calls `FleetMcp.answer()` →
`MessageService.answer(turnId, content, timeout)`. The `timeout` here is **not** the generous
30-minute async budget — it is the MCP tool's own bounded wait: `DEFAULT_TIMEOUT_MS = 25_000`,
clamped to at most `MAX_TIMEOUT_MS = 120_000` (`FleetMcp.java:71-72,495,512`). This is the same
~60–120s window every blocking `fleet_send` call is capped at (documented elsewhere as "the
caller's own MCP client call timeout").
4. `answer()` opens a **fresh** rendezvous waiter for the worker session and blocks on it for at most
that window. If the worker's resumed turn takes longer than that to actually finish (very
plausible — the resumed turn can mean more edits, a build, a commit, a push, opening a PR), the
wait times out. On timeout, `answer()`'s `finally` block unconditionally calls
`rendezvous.close(workerSession, reply)`, **removing the waiter from the map**, and returns
`Outcome.TIMED_OUT_WORKING` to the lead.
5. The worker keeps working, unaware anything happened, and eventually calls `fleet_reply`. That
reaches `MessageService.reply(session, content)`, which tries `rendezvous.resolve(session,
content)` — but the waiter was already closed in step 4, so `resolve` returns `false`. `reply()`
then falls back to `inbox.publish(...)` and marks `strandedReplies.put(session, true)`
(CB-640 bookkeeping) — the reply is safely held, but **the async `Task`'s `future` is never
completed**.
6. `fleet_poll{ticket}` keeps returning `PENDING` forever (the `Task` never resolves) — until the
lead eventually calls `fleet_stop`. That fires `sessions.onRelease` → `messages.abandon(target,
reason)` (`Fleetd.java:481-497`), where `reason` is built with the exact text from the bug report
("the worker session was released before it replied; worktree=... branch=... snapshot=...",
`Fleetd.java:484-487`). `abandon()`'s loop finds the still-open `Task` (`question == null`, future
not done) and completes it as `WORKER_FAILED` with that misleading reason — even though the
worker's real reply is sitting, intact, in the inbox the whole time.
So: the reported behaviour is correct, and the specific trigger is `answer()`'s own bounded wait
being shorter than the worker's real resumed-turn time — not anything to do with the ~55s
`fleet_ask` window itself (that part, issue #61, is untouched).
## Fix
Two changes in `fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java`, both scoped to the
ticket/reply routing and the terminal-state text — `fleet_ask`'s own window and mechanics are
untouched.
**1. `reply()` — priority 1 (the ticket resolves with the real reply).**
Before falling back to the inbox, `reply()` now looks for an async `Task` that is specifically in the
"already answered but not yet resolved" state (`question == null`, `turnId != null` — set once
`answer()` has cleared the question but before anything completed the future, `!future.isDone()`).
If one exists for this `target`, the worker's reply completes that `Task`'s future directly as
`Outcome.REPLIED` with the real content, and the reply never touches the inbox at all. A task that
was never asked has `turnId == null` and can never match, so ordinary (no-`fleet_ask`) delegations
are unaffected — they already resolve through the pre-existing rendezvous fast path.
I chose this over leaving `answer()`'s own timeout behaviour untouched and instead keeping its
rendezvous waiter open in the background: that alternative works but reopens the "at most one
waiter per session" invariant (`Rendezvous.open` throws on a double-open) to a new class of races
with a fresh send arriving mid-window. The `send()` path already guards against sending into an
answered-but-still-resolving worker via `hasAsyncQuestion(target)` (checks `asyncTasksByTurn`,
which still holds the task until it resolves), so routing through `reply()` gets the same protection
without touching `answer()`'s waiter lifecycle at all — the smaller, safer diff.
**2. `abandon()` — priority 3 (required independently, "even if you fix (1)").**
Before marking any of a released target's still-open tasks `WORKER_FAILED`, `abandon()` now checks
`hasStrandedReply(target)` (the existing CB-640 fact — true whenever the *last* `reply()` for this
target fell through to the inbox). If true, it drains the inbox (`recoverStrandedReply`) and — if it
actually finds a message — completes the task as `REPLIED` with that real content instead of writing
the failure. This is deliberately a **separate** check from fix 1: fix 1 already prevents the
inbox-stranding from happening in the exact scenario this ticket describes, so by the time
`abandon()` runs the task is normally already resolved and `abandon()`'s `complete()` call is a
harmless no-op. This second check exists so that if some *other* future path ever strands a reply
in the inbox without resolving its ticket, `abandon()` still refuses to report a false failure —
"if a reply reached any sink for that turn, the terminal state is done," per the ticket. I verified
both are required by disabling each independently and confirming the two new tests fail (see below).
**Priority 4 (the snapshot/worktree hint).** Handled as a consequence of both fixes rather than a
separate branch: once a task resolves as `REPLIED` (via either fix), `abandon()` never calls
`new Reply(Outcome.WORKER_FAILED, reason)` for that task at all, so the "the worker session was
released before it replied; worktree=... branch=... snapshot=..." text is never constructed or
attached to that ticket's outcome. It still appears, correctly, for a task that never got a reply
(the existing `abandonFailsEveryPendingAsyncTicketForTheReleasedTarget` /
`anAbandonedAsyncTaskPollsAsFailedNotPending` tests still pass unchanged).
**Priority 2** was not needed — fix 1 makes `fleet_poll{ticket}` return the actual reply (the
higher-priority option), so I did not fall back to "the ticket merely resolves as done with no
content."
## Tests — driven through the real delegation path, not the reply sink directly
Both new tests in `fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java` go through
`sendAsync` → `injectDelivery` → `ask` → `answer` (with a short timeout, so it genuinely times out,
mirroring the ~25–120s real MCP-call bound vs. a longer resumed turn) → `reply` → `poll`/`abandon`.
No test constructs a `Reply` and hands it to a sink directly.
- `aReplyAfterAnswerTimesOutStillCompletesTheAsyncTicket` — asserts `fleet_poll{ticket}` (via
`messages.poll`) reaches `Phase.DONE` with the worker's actual reply text and
`replySource() == "reply"`, and that `hasStrandedReply(T)` stays `false` (proves the reply never
touched the inbox at all — fix 1 caught it).
- `fleetStopAfterAnOrphanedReplyDoesNotFailTheTicket` — same setup, then calls `abandon(T, "the
worker session was released before it replied")` (what `fleet_stop` triggers) and asserts it
returns `false` (no failure recorded) and the ticket still polls `DONE` with the real reply.
**Proof both fail without the change.** I temporarily short-circuited both new private methods
(`askAnsweredAsyncTask` → always `null`, `recoverStrandedReply` → always `null`) — i.e. disabled
both fixes — and ran just these two tests:
```
[ERROR] Tests run: 2, Failures: 2, Errors: 0, Skipped: 0
dev.ltms.fleet.msg.MessageServiceTest.aReplyAfterAnswerTimesOutStillCompletesTheAsyncTicket
org.opentest4j.AssertionFailedError: expected: <DONE> but was: <PENDING>
dev.ltms.fleet.msg.MessageServiceTest.fleetStopAfterAnOrphanedReplyDoesNotFailTheTicket
org.opentest4j.AssertionFailedError: a reply already arrived, so nothing here is a genuine failure
==> expected: <false> but was: <true>
```
This is the exact bug: the ticket stays `PENDING` forever, and `abandon()` reports `true` (a
failure) even though a reply had already arrived. I then restored both fixes (verified with
`grep -n "TEMP #137-proof"` finding nothing) and re-ran — both pass.
## Build
Ran from `fleetd/`, unpiped, full output read (not `| tail`):
```
mvn clean install
...
[INFO] Tests run: 1039, Failures: 0, Errors: 0, Skipped: 0
[INFO] BUILD SUCCESS
[INFO] Total time: 36.315 s
```
Main was at 1037 tests; this branch adds the 2 new tests above → 1039, all green, `exit=0`.
## What I could NOT check
- No IDE tooling is mounted for me (worker), so no `ide_diagnostics`/IntelliJ inspection pass — only
`mvn clean install` (compiler + full test suite), as the worker procedure allows.
- I cannot restart the daemon or dogfood this live — I have no forge/daemon control. This is
unverified against a real herdr pane, a real MCP client's ~60s call cap, or a real worker session;
everything above is verified only through the JUnit fixture's simulated timing
(`FakeHerdr`/`injector.onStatus`/direct `messages.answer(...,150)` calls), not a live fleet.
A primary should still consider a short live dogfood (an async delegation that asks, gets answered,
and takes longer than ~2 minutes to reply) before calling this closed.
- I did not touch, and did not re-verify, the `fleet_ask` ~55s window itself (issue #61) — out of
scope per the brief.
## Scope note (not investigated further)
`answer()`'s nested/double-`fleet_ask` case (the worker asks a second question before ever
replying to the first answer) has some pre-existing behaviour around which `turnId` a `QUESTION`
resolution gets attributed to that I did not fully untangle — it predates this change, my fix does
not touch it, and it is unrelated to the reported defect. Flagging only; not investigated further.
## Handoff
- Branch: `worker/cb-137-ask-ticket-e7760c-2`
- Worktree root: `/Users/dai.ha/LTMS/.bridged-worktrees/734324-2`
- Files changed:
- `fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java`
- `fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java`
- `REPORT-cb137.md` (this file)
- Build: `Tests run: 1039, Failures: 0, Errors: 0, Skipped: 0` / `BUILD SUCCESS` (verbatim above)
-52
View File
@@ -1,52 +0,0 @@
# CB-175 report
## Change
`OpenCodeLauncher` reads the newest opencode session record for the worker cwd after spawn readiness.
It compares the requested `provider/model` selector with `model.providerID/model.id` from the record.
`variant` is not compared because a profile selector has no variant part.
An absent record, malformed record, or incomplete model object is unknown evidence. It does not log
an error or quarantine the profile.
On a real mismatch, fleetd logs an ERROR with the requested and resolved selectors. The mismatch goes
through `ExhaustionSink` into the existing `BackendQuarantine` and uses the profile's
`effectiveCredentialId()`.
I chose a permanent, process-lifetime quarantine. A withdrawn selector cannot become correct after a
cooldown. A timed retry could silently use the paid fallback again. `fleet_list` will show the usual
quarantine state, with a very large remaining time, until fleetd restarts after an operator fixes the
profile.
## Tests
Added tests for an exact match, a mismatch, and missing or unreadable session storage.
I proved the mismatch test fails without the quarantine call. I commented out the call and ran:
```text
mvn -Dtest=OpenCodeLauncherTest#differentResolvedModelPermanentlyQuarantinesTheProfile test
```
The result was:
```text
[ERROR] Tests run: 1, Failures: 1, Errors: 0, Skipped: 0
org.opentest4j.AssertionFailedError: a fallback model must block later spawns ==> expected: <true> but was: <false>
[INFO] BUILD FAILURE
```
I restored the call. I then ran `mvn clean install` in `fleetd/` without a pipe. Its result was:
```text
[INFO] Tests run: 1040, Failures: 0, Errors: 0, Skipped: 0
[INFO] BUILD SUCCESS
```
## Limits and scope
I could not spawn a real opencode member or restart fleetd. I did not test this end to end against a
live opencode session database.
I confirmed `ClaudeCodeLauncher` passes `--model` but does not read back the resolved model. I did
not change it because it is outside this ticket's scope.
+11
View File
@@ -634,6 +634,17 @@ guard:
# to a sibling directory of the repo root.
# worktreeRoot: /Users/me/src/.bridged-worktrees
# Worktree group sharing (fleetd #185 stage 3). OPTIONAL, off by default. Names an OS group
# that a provisioned worktree's repo is made group-writable for (git config
# core.sharedRepository group, plus a one-time chgrp/chmod/setgid fix-up), so a member spawned
# under a DIFFERENT OS user (see memberHerdrSocket) can write its own worktree, its
# per-worktree git metadata, and its own commit objects — without it, every file GitWorktrees
# creates is owned by fleetd's own uid and unwritable by another user.
# 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.
# worktreeGroup: fleet-workers
# Session lifecycle limits (CB-303). All knobs are opt-in; omit or set to null to keep
# the feature disabled. By default the daemon never reaps, caps, or drains sessions.
# idleTtlSeconds → reap READY/DONE sessions idle longer than this (never BUSY/SPAWNING)
+18
View File
@@ -28,6 +28,7 @@
<testcontainers.version>1.20.4</testcontainers.version>
<commons-compress.version>1.27.1</commons-compress.version>
<commons-lang3.version>3.18.0</commons-lang3.version>
<sqlite-jdbc.version>3.53.4.0</sqlite-jdbc.version>
</properties>
<!--
@@ -44,6 +45,12 @@
3.0-rc5; bumping Jackson 3 to the patched 3.2.x breaks the SDK (annotation mismatch).
Only the loopback /mcp endpoint parses this JSON, from trusted local Claude clients.
The 11.0.23 -> 11.0.25 bump did clear jetty CVE-2024-8184 (5.9) and CVE-2024-6763.
fleetd #206: org.xerial:sqlite-jdbc 3.53.4.0 (added for OpenCodeSessionDiscovery) — the
only known advisory against this artifact is CVE-2023-32697 (RCE via an attacker-controlled
JDBC URL), fixed in 3.41.2.2; 3.53.4.0 is well past that fix and OSV.dev reports no open
advisory against it. Checked via the OSV.dev API (no Mend.io/JetBrains IDE MCP mount
available from this worktree) on 2026-08-31.
-->
<!-- Force the latest patched Jetty 11.x across all Javalin-pulled Jetty modules (no version
@@ -123,6 +130,17 @@
<version>${amqp.version}</version>
</dependency>
<!-- fleetd #206: opencode moved its session store from a JSON tree to SQLite
(opencode.db). This is the JDBC driver OpenCodeSessionDiscovery uses to read it
read-only. Ships bundled native libraries (linux/mac/windows, several archs), so it
is a heavier jar than most deps here — see the pom's dependency-security note below
for the size/CVE tradeoff actually measured. -->
<dependency>
<groupId>org.xerial</groupId>
<artifactId>sqlite-jdbc</artifactId>
<version>${sqlite-jdbc.version}</version>
</dependency>
<!-- Logging -->
<dependency>
<groupId>org.slf4j</groupId>
@@ -168,17 +168,6 @@ public final class Fleetd {
claudeProfiles.put(name, w);
}
});
// One tracker covers timed backend exhaustion and permanent model-selector mismatches. The
// latter cannot heal on a retry, so OpenCodeLauncher uses quarantinePermanently through the
// sink below rather than letting a cooldown reopen a paid fallback.
BackendQuarantine quarantine = new BackendQuarantine(System::nanoTime,
TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds()));
ExhaustionSink modelMismatchSink = (profileName, reason) -> {
FleetConfig.Profile profile = config.get().profiles().get(profileName);
if (profile != null) {
quarantine.quarantinePermanently(profile.effectiveCredentialId());
}
};
List<HerdrPeerLauncher> adapters = new ArrayList<>();
// The claude-code adapter is the always-present default; keep it even with no profiles (so a
// bridge configured with no workers, or opencode-only, still has a well-defined base adapter)
@@ -188,20 +177,22 @@ public final class Fleetd {
claudeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
() -> config.get().fleet(),
() -> config.get().memberCredentials()));
() -> config.get().memberCredentials(), null, config::get));
}
if (!opencodeProfiles.isEmpty()) {
adapters.add(new OpenCodeLauncher(router.memberAgents(), router.memberSpaces(),
opencodeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
() -> config.get().fleet(),
() -> config.get().memberCredentials(), modelMismatchSink));
() -> config.get().memberCredentials(), config::get));
}
AtomicReference<Function<String, Integer>> liveCountRef = new AtomicReference<>(_ -> 0);
// CB-578 stage B: one quarantine tracker for the whole daemon, shared between the launcher
// (checked at spawn) and the exhaustion sink wired in below (written on BACKEND_EXHAUSTED).
// The cooldown is deferred (see FleetConfig#quarantineCooldownSeconds): it is read once
// here, at startup, and a config reload only changes it for a daemon restart.
BackendQuarantine quarantine = new BackendQuarantine(System::nanoTime,
TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds()));
PeerLauncher workers = new CompositePeerLauncher(
adapters,
cfg.effectiveDefaultProfile(),
@@ -233,7 +224,7 @@ public final class Fleetd {
contextCap = cfg.lifecycle().contextCap();
}
boolean clearAfterTurn = cfg.lifecycle() != null && cfg.lifecycle().clearAfterTurn();
SessionManager sessions = new SessionManager(workers, new GitWorktrees(cfg.worktreeRoot()),
SessionManager sessions = new SessionManager(workers, new GitWorktrees(cfg.worktreeRoot(), cfg.worktreeGroup()),
System::nanoTime, contextCap, clearAfterTurn);
liveCountRef.set(profileName -> (int) sessions.roster().stream()
.filter(s -> profileName.equals(s.profile()))
@@ -76,6 +76,14 @@ import java.util.Set;
* stay on {@code broker}'s vhost). {@code null} → no lead mailbox is opened.
* Config parsing + accessors only — nothing here wires it into a live
* {@code LeadMailbox}; that is a separate ticket. See {@link Coordinator}.
* @param worktreeGroup optional OS group name (fleetd #185 stage 3) that makes a provisioned
* worktree's repo group-shared, so a member running as a different OS user
* (see {@code memberHerdrSocket}) can write its own worktree, its per-worktree
* git metadata, and its own commit objects. {@code null}/blank/empty ⇒ off,
* today's behaviour unchanged (every file stays owned by fleetd's own uid).
* <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}.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record FleetConfig(
@@ -98,7 +106,20 @@ public record FleetConfig(
ConfigReload configReload,
Integer quarantineCooldownSeconds,
MemberCredentials memberCredentials,
Coordinator coordinator) {
Coordinator coordinator,
String worktreeGroup) {
/** Back-compat form before the {@code worktreeGroup} 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) {
this(bind, herdrSocket, memberHerdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth,
configReload, quarantineCooldownSeconds, memberCredentials, coordinator, null);
}
/** Back-compat form before the {@code coordinator:} block was added. */
public FleetConfig(Bind bind, String herdrSocket, Map<String, Profile> profiles, Guard guard,
@@ -109,7 +130,7 @@ public record FleetConfig(
MemberCredentials memberCredentials) {
this(bind, herdrSocket, null, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth,
configReload, quarantineCooldownSeconds, memberCredentials, null);
configReload, quarantineCooldownSeconds, memberCredentials, null, null);
}
/** Back-compat form before the CB-596 {@code memberCredentials:} block was added. */
@@ -1325,7 +1346,7 @@ public record FleetConfig(
"bind", "herdrSocket", "memberHerdrSocket", "profiles", "guard", "worktreeRoot",
"lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs", "broker", "primary", "fleet",
"leadHeartbeat", "health", "placement", "auth", "configReload", "quarantineCooldownSeconds",
"memberCredentials", "coordinator");
"memberCredentials", "coordinator", "worktreeGroup");
/** Load and validate config from {@code path}. */
public static FleetConfig load(Path path) {
@@ -1943,9 +1964,11 @@ public record FleetConfig(
: new MemberCredentials(null, List.of(), List.of());
// coordinator is left as-is, like broker/primary above: null keeps no LeadMailbox opened,
// and this ticket's Coordinator is config-only anyway (nothing yet reads it at startup).
// 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.
return new FleetConfig(b, herdrSocket, memberHerdrSocket, profiles, g, worktreeRoot, l, timeout, pollMs,
broker, primary, f, leadHeartbeat, health, placementOrDefault, a, configReload,
quarantineCooldown, mc, coordinator);
quarantineCooldown, mc, coordinator, worktreeGroup);
}
/**
@@ -2,8 +2,11 @@ package dev.ltms.fleet.herdr;
import com.fasterxml.jackson.databind.JsonNode;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Map;
import java.util.OptionalLong;
import java.util.Set;
/**
* Resolves which herdr pane a process belongs to — the herdr half of connection-based MCP
@@ -12,23 +15,46 @@ import java.util.Map;
* is calling without the worker sending anything spoofable.
*
* <p>herdr owns the PID→pane truth: {@code pane.process_info} reports each pane's {@code shell_pid}
* and foreground process PIDs. This scans agent panes; a spawn-time {@code pid→terminal} cache is
* the obvious optimization once wired into {@code ClaudeCodeLauncher}.
* and foreground process PIDs. A pid that is neither of those directly — e.g. a grandchild a
* worker spawned, such as a {@code python3} or {@code curl} helper that opens its own MCP
* connection — is resolved by walking its ancestry (via {@link ParentResolver}) up to the root and
* matching any ancestor against a pane's {@code shell_pid} or foreground pids (CB-161). Without
* this walk such a pid matches no pane, and the caller falls through to loopback-trust and is
* resolved as the primary — a worker→primary privilege escalation.
*
* <p>This scans agent panes; a spawn-time {@code pid→terminal} cache is the obvious optimization
* once wired into {@code ClaudeCodeLauncher}.
*
* <p>CB-185 split the fleet across two herdr daemons — lead operations on one, members on the
* other ({@code memberHerdrSocket}). A caller's pane can live on <em>either</em> daemon (a lead's
* MCP connection resolves against the lead daemon; a member's against the member daemon), so this
* must be able to search more than one client. {@link #PaneLocator(HerdrClient, HerdrClient)}
* searches the lead client first, then the member client, and collapses to a single scan when the
* two are the same object (the historical single-daemon deployment).
* two are the same object (the historical single-daemon deployment). The caller's ancestor set is
* computed once per {@link #terminalForPid} call and reused across every client searched — it
* does not depend on which daemon a pane happens to live on.
*/
public final class PaneLocator {
/**
* Bound on how many ancestor generations {@link #ancestorsOf} walks. This runs on every MCP
* call, so a cycle or a pathologically deep process tree must not hang identity resolution;
* 32 generations is far more than any real worker→helper process tree needs.
*/
private static final int MAX_ANCESTRY_DEPTH = 32;
private final List<HerdrClient> herdrs;
private final ParentResolver parentResolver;
/** Search only this client — the single-daemon deployment. */
public PaneLocator(HerdrClient herdr) {
this(herdr, ParentResolver.PROCESS_HANDLE);
}
/** Search only this client, resolving ancestry through {@code parentResolver} — for tests. */
public PaneLocator(HerdrClient herdr, ParentResolver parentResolver) {
this.herdrs = List.of(herdr);
this.parentResolver = parentResolver;
}
/**
@@ -37,7 +63,13 @@ public final class PaneLocator {
* collapses to one client and one scan, exactly {@link #PaneLocator(HerdrClient)}'s behaviour.
*/
public PaneLocator(HerdrClient lead, HerdrClient member) {
this(lead, member, ParentResolver.PROCESS_HANDLE);
}
/** Two-daemon deployment, resolving ancestry through {@code parentResolver} — for tests. */
public PaneLocator(HerdrClient lead, HerdrClient member, ParentResolver parentResolver) {
this.herdrs = lead == member ? List.of(lead) : List.of(lead, member);
this.parentResolver = parentResolver;
}
/**
@@ -49,8 +81,9 @@ public final class PaneLocator {
if (pid <= 0) {
return null;
}
Set<Long> ancestry = ancestorsOf(pid);
for (HerdrClient herdr : herdrs) {
String terminal = terminalForPid(herdr, pid);
String terminal = terminalForPid(herdr, ancestry);
if (terminal != null) {
return terminal;
}
@@ -58,28 +91,54 @@ public final class PaneLocator {
return null;
}
private static String terminalForPid(HerdrClient herdr, long pid) {
/**
* {@code pid} itself plus its ancestor chain, walked through {@link #parentResolver} up to
* {@link #MAX_ANCESTRY_DEPTH} generations or pid 1, whichever comes first. A vanished ancestor
* ({@link ParentResolver#parentOf} returning empty) ends the walk without error — it just means
* the chain is shorter than the bound. A cycle in a fake resolver is caught by the "already
* seen" check and also ends the walk, so this can never loop.
*/
private Set<Long> ancestorsOf(long pid) {
Set<Long> ancestry = new LinkedHashSet<>();
long current = pid;
for (int depth = 0; depth < MAX_ANCESTRY_DEPTH; depth++) {
if (current <= 0 || !ancestry.add(current)) {
break; // vanished/invalid pid, or a cycle back to a pid already recorded
}
if (current == 1) {
break; // reached the root of the process tree
}
OptionalLong parent = parentResolver.parentOf(current);
if (parent.isEmpty()) {
break; // vanished ancestor — not an error, just the end of the chain
}
current = parent.getAsLong();
}
return ancestry;
}
private static String terminalForPid(HerdrClient herdr, Set<Long> ancestry) {
for (JsonNode pane : herdr.call("pane.list", Map.of()).path("panes")) {
String paneId = pane.path("pane_id").asText(null);
if (paneId != null && paneOwnsPid(herdr, paneId, pid)) {
if (paneId != null && paneOwnsAnyOf(herdr, paneId, ancestry)) {
return pane.path("terminal_id").asText(null);
}
}
return null;
}
private static boolean paneOwnsPid(HerdrClient herdr, String paneId, long pid) {
private static boolean paneOwnsAnyOf(HerdrClient herdr, String paneId, Set<Long> ancestry) {
JsonNode info;
try {
info = herdr.call("pane.process_info", Map.of("pane_id", paneId)).path("process_info");
} catch (HerdrException e) {
return false; // pane vanished mid-scan — just skip it
}
if (info.path("shell_pid").asLong(-1) == pid) {
if (ancestry.contains(info.path("shell_pid").asLong(-1))) {
return true;
}
for (JsonNode p : info.path("foreground_processes")) {
if (p.path("pid").asLong(-1) == pid) {
if (ancestry.contains(p.path("pid").asLong(-1))) {
return true;
}
}
@@ -0,0 +1,21 @@
package dev.ltms.fleet.herdr;
import java.util.OptionalLong;
/**
* Resolves a pid's parent pid — the seam {@link PaneLocator} walks a process's ancestry through,
* so its tests can drive the walk from a fake pid→parent map instead of spawning real processes.
*
* <p>{@link #PROCESS_HANDLE} is the production implementation, backed by {@link ProcessHandle}.
*/
public interface ParentResolver {
/** The parent pid of {@code pid}, or empty if {@code pid} is gone or has no known parent. */
OptionalLong parentOf(long pid);
/** Production resolver: asks the JVM's {@link ProcessHandle} view of the OS process tree. */
ParentResolver PROCESS_HANDLE = pid -> ProcessHandle.of(pid)
.flatMap(ProcessHandle::parent)
.map(parent -> OptionalLong.of(parent.pid()))
.orElse(OptionalLong.empty());
}
@@ -94,13 +94,26 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs, long spawnReadyPollMs,
Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
long spawnReadyTimeoutMs, long spawnReadyPollMs,
Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
this(agents, spaces, guard, profiles, defaultProfile, env, spawnReadyTimeoutMs, spawnReadyPollMs,
fleet, memberCredentials, null, null);
}
/** Production constructor, plus the live config for URI environment exclusions. */
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs, long spawnReadyPollMs,
Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials,
Supplier<Set<String>> hostEnvNames,
Supplier<FleetConfig> config) {
this(agents, spaces, guard, profiles, defaultProfile, env,
spawnReadyTimeoutMs,
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs),
fleet, memberCredentials);
fleet, memberCredentials, hostEnvNames, config);
}
/**
@@ -153,10 +166,22 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
Function<String, String> env,
long spawnReadyTimeoutMs,
LongSupplier nowMillis, Runnable sleeper,
Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
this(agents, spaces, guard, profiles, defaultProfile, env, spawnReadyTimeoutMs, nowMillis, sleeper,
fleet, memberCredentials, null, null);
}
/** Full testability constructor, plus the live config for URI environment exclusions. */
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env, long spawnReadyTimeoutMs,
LongSupplier nowMillis, Runnable sleeper, Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials,
Supplier<Set<String>> hostEnvNames,
Supplier<FleetConfig> config) {
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials);
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials, hostEnvNames, config);
this.guard = guard;
}
@@ -174,7 +199,7 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
Supplier<FleetConfig.MemberCredentials> memberCredentials,
Supplier<Set<String>> hostEnvNames) {
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials, hostEnvNames);
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials, hostEnvNames, null);
this.guard = guard;
}
@@ -170,6 +170,8 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
/** Guards {@link #warnNonZsh} to one WARN per launcher instance, not one per spawn. */
private final AtomicBoolean nonZshShellWarned = new AtomicBoolean();
/** Live config provides URI environment names that must never enter member panes. */
private final Supplier<FleetConfig> config;
/**
* @param namePrefix label prefix for this peer kind (drives naming and reap)
@@ -238,8 +240,19 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
long spawnReadyTimeoutMs,
LongSupplier nowMillis, Runnable sleeper,
Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials,
Supplier<Set<String>> hostEnvNames) {
this(namePrefix, agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs, nowMillis,
sleeper, fleet, memberCredentials, hostEnvNames, null);
}
/** As above, plus the live full config for secret-bearing URI environment names. */
protected HerdrPeerLauncher(String namePrefix, AgentControl agents, WorkspaceControl spaces,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env, long spawnReadyTimeoutMs,
LongSupplier nowMillis, Runnable sleeper, Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials,
Supplier<Set<String>> hostEnvNames) {
Supplier<Set<String>> hostEnvNames, Supplier<FleetConfig> config) {
this.fleet = fleet;
this.namePrefix = namePrefix;
this.agents = agents;
@@ -252,6 +265,7 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
this.sleeper = sleeper;
this.memberCredentials = memberCredentials;
this.hostEnvNames = hostEnvNames != null ? hostEnvNames : () -> System.getenv().keySet();
this.config = config;
}
// --- adapter seams -------------------------------------------------------------------------
@@ -1028,9 +1042,11 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
}
/** Put {@link #BLOCKED_CREDENTIAL_SENTINEL} over every blocked name in the pane-creation env map. */
private static void overlayBlockedCredentials(Map<String, String> workerEnv,
FleetConfig.MemberCredentials creds) {
for (String name : creds.blockedSet()) {
private void overlayBlockedCredentials(Map<String, String> workerEnv,
FleetConfig.MemberCredentials creds) {
Set<String> blocked = new java.util.TreeSet<>(creds.blockedSet());
blocked.addAll(brokerUriEnvNames());
for (String name : blocked) {
workerEnv.put(name, BLOCKED_CREDENTIAL_SENTINEL);
}
}
@@ -1097,15 +1113,21 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
* ssh-agent handle when explicitly allowed, and the exact keys of THIS launch's own env map.
*/
private Set<String> derivedAllowedNames(FleetConfig.MemberCredentials creds, Launch launch) {
Set<String> brokerUriEnvNames = brokerUriEnvNames();
Set<String> allowed = new java.util.TreeSet<>(
MemberEnvAllowList.derive(profiles.values(), creds.allowSet()));
MemberEnvAllowList.derive(profiles.values(), creds.allowSet(), brokerUriEnvNames));
if (creds.sshAuthSockAllowed()) {
allowed.add(SSH_AUTH_SOCK);
} // blocked by default: absent from the set ⇒ blanked by the scrub like any other name
allowed.addAll(launch.env().keySet());
allowed.removeAll(brokerUriEnvNames);
return allowed;
}
private Set<String> brokerUriEnvNames() {
return MemberEnvAllowList.brokerUriEnvNames(config == null ? null : config.get());
}
/**
* CB-633 follow-up: one INFO line per allow-list spawn WHOSE SCRUB ACTUALLY RUNS, so an operator
* can read a single log line and know the scrub ran and how much of the visible environment it
@@ -40,8 +40,9 @@ import java.util.TreeSet;
* operator's own explicit list. Before this, {@code policy: allow-list} silently ignored every name
* an operator wrote under {@code allow:} unless a profile happened to carry it too, which meant
* turning the policy on could blank credentials working members already depended on. {@code
* SSH_AUTH_SOCK} is the one exception: even when the operator lists it under {@code allow:}, it is
* excluded here and added back ONLY by the caller when {@code sshAuthSock: allow} is explicitly set
* SSH_AUTH_SOCK} and configured broker URI environment names are exceptions: even when the operator
* lists them under {@code allow:}, they are excluded here. {@code SSH_AUTH_SOCK} is added back ONLY
* by the caller when {@code sshAuthSock: allow} is explicitly set
* (see {@link #SSH_AUTH_SOCK}'s javadoc) — it is a live handle to the operator's own ssh-agent, not
* a value, so treating it like any other allow-listed name would hand a member every key the
* operator's agent holds the moment they typed the name under {@code allow:} for an unrelated
@@ -106,6 +107,16 @@ public final class MemberEnvAllowList {
* run-to-run.
*/
public static Set<String> derive(Collection<FleetConfig.Profile> profiles, Set<String> configuredAllow) {
return derive(profiles, configuredAllow, Set.of());
}
/**
* As {@link #derive(Collection, Set)}, while excluding names that fleetd knows carry credentials.
* A configured broker URI contains its AMQP password inline, so it must never reach a member,
* even when an operator put its variable name in {@code memberCredentials.allow:}.
*/
public static Set<String> derive(Collection<FleetConfig.Profile> profiles, Set<String> configuredAllow,
Set<String> excludedNames) {
Set<String> derived = new TreeSet<>(INFRASTRUCTURE_PASSTHROUGH);
if (profiles != null) {
for (FleetConfig.Profile p : profiles) {
@@ -124,9 +135,37 @@ public final class MemberEnvAllowList {
}
}
}
if (excludedNames != null) {
derived.removeAll(excludedNames);
}
return Set.copyOf(derived);
}
/**
* The host environment names whose values are AMQP URIs with inline passwords. Both broker
* connections belong to fleetd, never to a member pane. Blank and absent configuration changes
* nothing.
*
* <p><b>How strong this exclusion is depends on the policy, and the difference matters.</b>
* Under {@code policy: allow-list} it is enforced by the generated ZDOTDIR scrub, which runs
* AFTER the pane's shell has sourced the operator's chain — so a login shell that re-exports the
* name is still blanked. Under the deny-list policy there is no scrub: the name is only removed
* from the pre-shell env map, and a login shell that sources the operator's secret store
* re-exports it. That is the long-standing weakness of deny-list (a sourced file can undo it),
* not something this exclusion introduces, but it means deny-list deployments do NOT get this
* guarantee. The same caveat applies to the non-zsh path, which has no scrub at all — see
* {@code HerdrPeerLauncher#applyEnvironmentAllowListPolicy}.
*/
public static Set<String> brokerUriEnvNames(FleetConfig config) {
if (config == null) {
return Set.of();
}
Set<String> names = new TreeSet<>();
addUriEnvIfPresent(names, config.broker());
addUriEnvIfPresent(names, config.coordinator());
return Set.copyOf(names);
}
/**
* Whether {@code name} survives the scrub when {@code allowedNames} is the derived set: an exact
* match, or an infrastructure-prefixed name ({@code LC_*}). Prefix rules live ONLY here and in
@@ -146,4 +185,16 @@ public final class MemberEnvAllowList {
into.add(name);
}
}
private static void addUriEnvIfPresent(Set<String> into, FleetConfig.Broker broker) {
if (broker != null && broker.hasUriEnv()) {
into.add(broker.uriEnv());
}
}
private static void addUriEnvIfPresent(Set<String> into, FleetConfig.Coordinator coordinator) {
if (coordinator != null && coordinator.hasUriEnv()) {
into.add(coordinator.uriEnv());
}
}
}
@@ -6,7 +6,6 @@ import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.herdr.Agent;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.inject.ExhaustionSink;
import dev.ltms.fleet.peer.Capability;
import dev.ltms.fleet.peer.CharterReceipt;
import dev.ltms.fleet.peer.PeerHandle;
@@ -72,9 +71,6 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
*/
private final OpenCodeSessionDiscovery discovery;
/** Receives the profile when a resolved model differs from its requested selector. */
private final ExhaustionSink modelMismatchSink;
/**
* Production constructor — disables the spawn-ready gate ({@code spawnReadyTimeoutMs == 0}) so it
* matches the legacy non-blocking spawn semantics. Config dirs are created under the JVM temp dir.
@@ -120,25 +116,23 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
public OpenCodeLauncher(AgentControl agents, WorkspaceControl spaces,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs, long spawnReadyPollMs,
Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
this(agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs,
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs),
defaultConfigRoot(), defaultDiscoveryRoot(), fleet, memberCredentials, ExhaustionSink.none());
}
/** Production constructor with permanent-quarantine wiring for a model mismatch. */
public OpenCodeLauncher(AgentControl agents, WorkspaceControl spaces,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs, long spawnReadyPollMs,
Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
this(agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs, spawnReadyPollMs,
fleet, memberCredentials, null);
}
/** Production constructor, plus the live config for URI environment exclusions. */
public OpenCodeLauncher(AgentControl agents, WorkspaceControl spaces,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env, long spawnReadyTimeoutMs, long spawnReadyPollMs,
Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials,
ExhaustionSink modelMismatchSink) {
Supplier<FleetConfig> config) {
this(agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs,
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs),
defaultConfigRoot(), defaultDiscoveryRoot(), fleet, memberCredentials, modelMismatchSink);
defaultConfigRoot(), defaultDiscoveryRoot(), fleet, memberCredentials, config);
}
/**
@@ -180,10 +174,12 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
Function<String, String> env,
long spawnReadyTimeoutMs,
LongSupplier nowMillis, Runnable sleeper,
Path configRoot, Path discoveryRoot,
Supplier<FleetConfig.Fleet> fleet) {
this(agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs, nowMillis, sleeper,
configRoot, discoveryRoot, fleet, null, ExhaustionSink.none());
Path configRoot, Path discoveryRoot,
Supplier<FleetConfig.Fleet> fleet) {
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
spawnReadyTimeoutMs, nowMillis, sleeper, fleet);
this.configRoot = configRoot;
this.discovery = new OpenCodeSessionDiscovery(discoveryRoot);
}
/**
@@ -194,32 +190,25 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
Function<String, String> env,
long spawnReadyTimeoutMs,
LongSupplier nowMillis, Runnable sleeper,
Path configRoot, Path discoveryRoot,
Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
Path configRoot, Path discoveryRoot,
Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
this(agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs, nowMillis, sleeper,
configRoot, discoveryRoot, fleet, memberCredentials, ExhaustionSink.none());
configRoot, discoveryRoot, fleet, memberCredentials, null);
}
/**
* Full constructor with the model-mismatch quarantine callback. The callback is an
* {@link ExhaustionSink} so model mismatches use the existing quarantine path rather than a
* second state tracker.
*/
/** Full testability constructor, plus the live config for URI environment exclusions. */
public OpenCodeLauncher(AgentControl agents, WorkspaceControl spaces,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs,
LongSupplier nowMillis, Runnable sleeper,
Path configRoot, Path discoveryRoot,
Function<String, String> env, long spawnReadyTimeoutMs,
LongSupplier nowMillis, Runnable sleeper, Path configRoot, Path discoveryRoot,
Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials,
ExhaustionSink modelMismatchSink) {
Supplier<FleetConfig> config) {
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials);
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials, null, config);
this.configRoot = configRoot;
this.discovery = new OpenCodeSessionDiscovery(discoveryRoot);
this.modelMismatchSink = modelMismatchSink;
}
private static Path defaultConfigRoot() {
@@ -532,34 +521,9 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
@Override
public PeerHandle spawn(SpawnRequest req) {
PeerHandle inner = super.spawn(req);
verifyResolvedModel(requireProfile(req.profileName()), effectiveCwd(req));
return new SessionAwareHandle(inner, discovery, effectiveCwd(req));
}
/**
* Read the session record once after the spawn-readiness gate. opencode writes the model it
* actually selected there. No record, incomplete model object, or a selector without one slash
* is unknown evidence, so it must not quarantine a working profile.
*/
private void verifyResolvedModel(FleetConfig.Profile cfg, String cwd) {
String[] requested = splitProviderModel(cfg.model());
if (requested == null) {
return;
}
OpenCodeSessionDiscovery.SessionRecord record = discovery.sessionForDirectory(cwd);
if (record == null || record.providerId() == null || record.modelId() == null) {
return;
}
String actual = record.providerId() + "/" + record.modelId();
if (cfg.model().equals(actual)) {
return;
}
log.error("opencode model mismatch for profile '{}': requested '{}' but resolved '{}'",
cfg.profile(), cfg.model(), actual);
modelMismatchSink.onExhausted(cfg.profile(), "opencode model mismatch: requested "
+ cfg.model() + ", resolved " + actual);
}
/**
* A {@link PeerHandle} that delegates everything to the base's worker handle but resolves
* {@link #agentSessionId()} lazily through opencode session discovery. Delegate-only, so the
@@ -1,143 +1,120 @@
package dev.ltms.fleet.member;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.sqlite.SQLiteConfig;
import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.stream.Stream;
import java.sql.Connection;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.util.concurrent.atomic.AtomicBoolean;
/**
* Resolves the opencode session id for a fleetd worker from opencode's on-disk storage — the
* only place this adapter touches opencode's private layout, and deliberately the <em>only</em>
* class that does.
*
* <p><strong>Why this is isolated behind one seam.</strong> The layout is version-coupled and not a
* stable contract: opencode writes one JSON file per session under
* {@code <storageRoot>/session/<projectID>/<ses_*.json>}, and each record carries a
* {@code "version"} field (e.g. {@code "1.1.31"}), so the exact directory shape, file naming, and
* field names can move between opencode releases. opencode also ships a headless HTTP server that
* may supersede file scanning entirely. Everything this adapter knows about that private storage —
* its shape, naming, and field names — lives here, so a layout change, or a switch to the HTTP
* server, changes exactly one class and nothing in {@link OpenCodeLauncher}.
* <p><strong>Why this is isolated behind one seam.</strong> The layout is version-coupled and not
* a stable contract: opencode persists its session state in a SQLite database at
* {@code <storageRoot>/opencode.db} (a {@code session} table, one row per session, keyed by id and
* carrying a {@code directory} column). That schema can move between opencode releases exactly
* like the JSON-file layout it replaced did (opencode migrated off a one-JSON-file-per-session
* tree under {@code <storageRoot>/storage/session/<projectID>/ses_*.json} in January 2026 — that
* tree is now a frozen migration artefact nothing writes, which is why this class no longer reads
* it). opencode also ships a headless HTTP server that may supersede both of these entirely.
* Everything this adapter knows about that private storage — its shape and column names — lives
* here, so a layout change, or a switch to the HTTP server, changes exactly one class and nothing
* in {@link OpenCodeLauncher}.
*
* <p>The determinism that makes this useful is structural, not a guess: every fleetd worker runs
* in its own unique git worktree, so the record's {@code directory} (its project root) equals the
* in its own unique git worktree, so the row's {@code directory} (its project root) equals the
* worker's cwd identifies <em>its</em> session unambiguously. We match on {@code directory} rather
* than diffing {@code opencode session list} before/after — that races under concurrent spawns, and
* the CLI listing does not even show the directory.
*
* <p>All reads are best-effort and never throw: a missing or unreadable storage root, a record that
* fails to parse, or a directory with no record yet all yield {@code null}, and the caller (the
* session handle) treats that as "identity not resolved yet" and retries later.
* <p>All reads are best-effort and never throw: a missing or unreadable database, a query that
* fails, or a directory with no row yet all yield {@code null}, and the caller (the session
* handle) treats that as "identity not resolved yet" and retries later. The database is opened
* read-only and never written to: opencode itself may be running and writing it concurrently (WAL
* mode), and this class must never disturb that.
*/
final class OpenCodeSessionDiscovery {
private static final Logger log = LoggerFactory.getLogger(OpenCodeSessionDiscovery.class);
private final Path storageRoot; // e.g. ~/.local/share/opencode (injectable for tests)
private final ObjectMapper json;
private final Path databasePath;
private final AtomicBoolean warnedMissingDatabase = new AtomicBoolean(false);
OpenCodeSessionDiscovery(Path storageRoot) {
this.storageRoot = storageRoot;
this.json = new ObjectMapper();
this.databasePath = storageRoot.resolve("opencode.db");
}
/**
* The opencode session id whose record references {@code directory} (the worker's cwd), or
* {@code null} when no record matches yet. When several records share the directory — e.g.
* repeated spawns into the same worktree — the <em>most recently modified</em> one wins: it is
* A connection to {@link #databasePath} opened with SQLite's {@code SQLITE_OPEN_READONLY}
* flag: it never creates the file, never writes, and never touches WAL or journal mode.
* opencode may be running and writing this database concurrently, and this class must never
* disturb it.
*
* <p>Package-private so a test can hold the connection and prove it refuses a write. That is
* the only way to pin this property: making the file unwritable does <em>not</em> work,
* because SQLite silently downgrades a read-write open of an unwritable file to read-only, so
* such a test passes whether or not the flag is set.
*/
Connection openReadOnly() throws SQLException {
SQLiteConfig config = new SQLiteConfig();
config.setReadOnly(true);
return config.createConnection("jdbc:sqlite:" + databasePath);
}
/**
* The opencode session id whose row references {@code directory} (the worker's cwd), or
* {@code null} when no row matches yet. When several rows share the directory — e.g. repeated
* spawns into the same worktree — the row with the highest {@code time_updated} wins: it is
* the session the pane most likely corresponds to.
*
* <p>Never throws: a missing {@code storageRoot}, an unreadable/malformed record, or a
* directory that has not been persisted yet all resolve to {@code null} rather than failing a
* spawn. A fleetd worker's session record is written lazily (when the session is first
* persisted), so {@code null} here is the normal answer right after the pane is ready, and the
* caller retries later.
* <p>Never throws: a missing {@code opencode.db}, a locked/unreadable database, a query
* failure, or a directory that has not been persisted yet all resolve to {@code null} rather
* than failing a spawn. A fleetd worker's session row is written lazily (when the session is
* first persisted), so {@code null} here is the normal answer right after the pane is ready,
* and the caller retries later.
*
* @param directory the worker's cwd, as resolved for this spawn
* @return the matching session id, or {@code null} if none is known yet
*/
String sessionIdForDirectory(String directory) {
SessionRecord record = sessionForDirectory(directory);
return record == null ? null : record.id();
}
/**
* The newest session record for {@code directory}, or {@code null} when opencode has not written
* one yet. This is the single storage seam for both session identity and resolved-model checks.
*/
SessionRecord sessionForDirectory(String directory) {
if (directory == null || directory.isBlank()) {
return null;
}
Path sessionRoot = storageRoot.resolve("session");
if (!Files.isDirectory(sessionRoot)) {
if (!Files.isRegularFile(databasePath)) {
if (warnedMissingDatabase.compareAndSet(false, true)) {
log.warn("opencode session database not found at {} — opencode's on-disk layout "
+ "may have moved again; session discovery will keep returning null",
databasePath);
}
return null;
}
SessionRecord best = null;
long bestMtime = Long.MIN_VALUE;
try (Stream<Path> projectDirs = Files.list(sessionRoot)) {
for (Path projectDir : projectDirs.filter(Files::isDirectory).toList()) {
try (Stream<Path> records = Files.list(projectDir)) {
for (Path record : records.toList()) {
SessionRecord matched = matchRecord(record, directory);
if (matched == null) {
continue;
}
long mtime = lastModifiedEpochMillis(record);
if (mtime > bestMtime) {
bestMtime = mtime;
best = matched;
}
}
} catch (IOException ignored) {
// one project dir unreadable — skip it; another may still match
String sql = "SELECT id FROM session WHERE directory = ? ORDER BY time_updated DESC LIMIT 1";
try (Connection connection = openReadOnly();
PreparedStatement statement = connection.prepareStatement(sql)) {
statement.setString(1, directory);
try (ResultSet rows = statement.executeQuery()) {
if (rows.next()) {
return rows.getString("id");
}
}
} catch (IOException ignored) {
// storage root vanished or became unreadable — "no session known yet"
} catch (SQLException e) {
// Locked, corrupt, or otherwise unreadable — never fatal to a spawn. Not the
// "database moved" signal (the file exists), so this stays below WARN.
log.debug("opencode session database unreadable at {}: {}", databasePath, e.toString());
return null;
}
return best;
}
/**
* The record's session id when it references {@code directory}, else {@code null}. A record
* that is not JSON, lacks {@code id}/{@code directory}, or points at a different directory is
* simply not our session; a malformed one is skipped, never fatal.
*/
private SessionRecord matchRecord(Path record, String directory) {
try {
JsonNode node = json.readTree(record.toFile());
JsonNode id = node == null ? null : node.get("id");
JsonNode dir = node == null ? null : node.get("directory");
if (id == null || dir == null || !directory.equals(dir.asText())) {
return null;
}
JsonNode model = node.path("model");
String providerId = text(model, "providerID");
String modelId = text(model, "id");
return new SessionRecord(id.asText(), providerId, modelId);
} catch (IOException e) {
return null;
}
}
private static String text(JsonNode node, String name) {
JsonNode value = node.get(name);
return value == null || value.isNull() || value.asText().isBlank() ? null : value.asText();
}
/** The opencode fields fleetd reads from one session record. Null model fields mean unknown. */
record SessionRecord(String id, String providerId, String modelId) {
}
/** The record's last-modified epoch ms, or {@code Long.MIN_VALUE} if unreadable (never wins). */
private static long lastModifiedEpochMillis(Path record) {
try {
return Files.getLastModifiedTime(record).toMillis();
} catch (IOException e) {
return Long.MIN_VALUE;
}
log.debug("no opencode session row for directory (root={}, directory={})",
storageRoot, directory);
return null;
}
}
@@ -9,6 +9,7 @@ import dev.ltms.fleet.metrics.Metrics;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.ArrayList;
import java.util.List;
import java.util.UUID;
import java.util.concurrent.CompletableFuture;
@@ -167,6 +168,12 @@ public final class MessageService {
private static final class Task {
private final String ticket;
private final String target;
/**
* When this task was created (#137 fix): the tiebreaker for which of several open tasks on
* one target gets a recovered reply in {@link #abandon} — the oldest, since it is the one
* that has been waiting longest.
*/
private final long createdNanos;
private final CompletableFuture<Reply> future = new CompletableFuture<>();
/**
* When {@link #future} resolved, or {@code null} while it is still pending — the clock
@@ -184,6 +191,7 @@ public final class MessageService {
private Task(String ticket, String target, LongSupplier nowNanos) {
this.ticket = ticket;
this.target = target;
this.createdNanos = nowNanos.getAsLong();
future.whenComplete((reply, ex) -> completedNanos = nowNanos.getAsLong());
}
}
@@ -366,21 +374,61 @@ public final class MessageService {
}
/**
* Route a worker's explicit {@code fleet_reply}: resolve an open send, or queue it in the
* inbox if no send is currently open. Unlike the bare {@link Rendezvous#resolve}, a no-waiter
* result is <em>not</em> a failure — the reply is held for later drain.
* Route a worker's explicit {@code fleet_reply}: resolve an open send, complete an async ticket
* still parked waiting on this exact turn's answer, or — only once neither applies — queue it in
* the inbox. Unlike the bare {@link Rendezvous#resolve}, a no-waiter result is <em>not</em> a
* failure — the reply is held for later drain.
*
* <p><strong>Do NOT use this for mid-turn questions.</strong> {@code fleet_ask} /
* {@link Rendezvous#resolveQuestion} must keep today's {@code NO_WAITER} behaviour — questions
* are interactive and must never be queued.
*
* @return always {@code true} — the reply either resolved a live send or was queued
* <p><strong>Ambiguous match also falls to the inbox.</strong> {@link #askAnsweredAsyncTasks}
* cannot actually return more than one entry today (see its own javadoc for why — in short,
* {@link #hasAsyncQuestion} keeps a target BUSY, so no second task can reach this state, for as
* long as an earlier one's {@code turnId} is still stamped). That is an emergent guarantee from
* two other facts, not one this method enforces, so this branch stays in as defence in depth
* rather than being removed as dead code: if it ever weakens, returning whichever candidate a
* {@code ConcurrentHashMap} iteration reaches first would let a genuine reply complete the
* <em>wrong</em> ticket — silently handing the lead something that reads like a correct answer to
* a delegation the worker never touched, which is worse than a failure because the lead acts on
* it. When more than one candidate exists, guessing is not safe: fall back to the inbox exactly
* as the zero-candidate case does, and let {@link #abandon} apply the eventual recovery
* deterministically instead.
*
* @return always {@code true} — the reply resolved a live send, completed a parked ticket, or
* was queued
*/
public boolean reply(String session, String content) {
if (rendezvous.resolve(session, content)) {
count(FleetMetrics.REPLIES, "path", "rendezvous");
return true; // a live send took it — unchanged fast path
}
// #137: no live rendezvous waiter, but this may be the worker's real fleet_reply resuming a
// turn that {@link #answer} already gave up waiting on. answer()'s own bounded wait (the
// primary's fleet_send{turnId} call, capped well under a minute) can time out and close its
// waiter long before the worker — now actually resuming real work — finishes and replies. That
// reply used to have nowhere to land but the session inbox, leaving the async ticket's future
// unresolved forever: fleet_poll{ticket} stayed PENDING until fleet_stop's abandon() forced it
// FAILED with a misleading "session released before it replied" reason, even though the reply
// had, in fact, arrived. Completing the matching ticket directly here means fleet_poll{ticket}
// sees the real reply instead.
List<Task> candidates = askAnsweredAsyncTasks(session);
if (candidates.size() == 1) {
Task orphan = candidates.get(0);
if (orphan.future.complete(new Reply(Outcome.REPLIED, content))) {
if (orphan.turnId != null) {
asyncTasksByTurn.remove(orphan.turnId, orphan);
}
count(FleetMetrics.REPLIES, "path", "async-recovered");
return true; // the ticket itself took it — no inbox stranding at all
}
} else if (candidates.size() > 1) {
List<String> tickets = candidates.stream().map(t -> t.ticket).toList();
log.warn("reply from {} matches {} open async tickets {} — cannot tell which one it "
+ "answers, queuing to the inbox instead of guessing", session, candidates.size(),
tickets);
}
inbox.publish(session, UUID.randomUUID().toString(), content);
// CB-640: record the stranding itself (not just the reply text) so fleet health can see a
// worker whose replies keep missing their waiter, not only the queue depth this leaves behind.
@@ -394,6 +442,41 @@ public final class MessageService {
return true; // held, not lost
}
/**
* Every still-open async task on {@code target} whose {@code fleet_ask} was already answered —
* its {@link Task#turnId} is stamped but its {@link Task#question} was cleared by {@link #answer}
* — yet whose future is not resolved yet (#137). Empty if no such task exists, including the
* common case where {@code target}'s worker never used {@code fleet_ask} at all (a task that was
* never asked has {@code turnId == null}, so it can never match here and only ever completes
* through the ordinary rendezvous fast path in {@link #reply}).
*
* <p><strong>Returns at most one entry today — verified, not assumed.</strong> {@link #send}
* refuses to open a waiter on {@code target} while {@link #hasAsyncQuestion} is true, and that
* check matches ANY task whose {@code turnId} is still stamped in {@code asyncTasksByTurn} —
* not only while its question is still open. {@link #answer} deliberately leaves that stamp in
* place ({@code clearAsyncQuestion(turnId, false)}) until the resumed turn's own future actually
* resolves, at which point {@link #finishAsyncTask} both removes the stamp AND completes that
* task's future in the same call. So a second task can never reach "{@code turnId} stamped, future
* still open" — the exact pair this method matches on — while a first one already holds it: by
* the time the stamp is gone, so is the eligibility. This is an emergent property of those two
* facts holding together, not something this method (or its callers) enforces on its own — flip
* {@code forgetTurn} to {@code true} in that one {@link #answer} call and it silently stops being
* true, with nothing left to fail loudly. The callers below still handle "more than one" as
* defence in depth against exactly that, not because they exercise it today: {@link #reply}
* treats it as unresolvable and falls back to the inbox; {@link #abandon} would pick the oldest
* deterministically (its own {@code matching} list has no such guarantee — see its javadoc).
*/
private List<Task> askAnsweredAsyncTasks(String target) {
List<Task> candidates = new ArrayList<>();
for (Task task : tasks.values()) {
if (target.equals(task.target) && task.question == null && task.turnId != null
&& !task.future.isDone()) {
candidates.add(task);
}
}
return candidates;
}
/** Record a counter sample when a registry is wired; a no-op in unit tests. */
private void count(String name, String... labels) {
if (metrics != null) {
@@ -435,19 +518,100 @@ public final class MessageService {
* <p>Resolving the waiter as a failure — rather than letting it time out — also means the
* outcome is counted, so a torn-down delegation stops being invisible to {@code /metrics}.
*
* @return true if a live waiter was failed
* <p><strong>#137 defence in depth.</strong> {@link #reply} already hands a worker's real
* {@code fleet_reply} straight to the async ticket it belongs to whenever exactly one is still
* parked waiting for it (see {@link #askAnsweredAsyncTasks}), so by the time a session is
* released its tasks are normally already resolved — this loop's {@code complete} calls are then
* harmless no-ops (a {@link CompletableFuture} can only resolve once). But should some other path
* someday strand a reply in the inbox without completing its ticket, checking
* {@link #hasStrandedReply(String)} here — before ever writing a failure — means a torn-down
* session whose worker in fact replied is still reported {@code REPLIED} with that reply's own
* text, never the misleading "the worker session was released before it replied" (which also
* means the snapshot/worktree recovery hint that follows it never prints once a reply exists).
*
* <p><strong>At most one task gets the recovered reply — and here, unlike {@link #reply}'s
* {@link #askAnsweredAsyncTasks}, {@code matching.size() >= 2} alone is reachable today.</strong>
* This method's {@code matching} filter has no {@code turnId != null} requirement, so it matches
* any plain (never-asked) open task too — and {@link #sendAsync} does not limit a target to one
* of those: a second {@code fleet_send{wait:false}} at a target that is still busy returns its own
* ticket immediately and simply parks its {@link #send} behind the target's session lock for up
* to {@link #ASYNC_TIMEOUT_MS}, exactly as {@code abandonFailsEveryPendingAsyncTicketForTheReleasedTarget}
* already proves. Before this fix, the loop below drained the strand once and then reused that
* same {@code Reply} for <em>every</em> task it walked past — so two open tasks really did both
* complete {@code REPLIED} with the same text (see the pre-fix loop in commit 97f6c33's parent).
* A stranded reply is one worker answer, so it can settle at most one open task on this target —
* never every open task, and never a guess. When more than one task is still open here, the
* recovered reply goes to the <em>oldest</em> (lowest {@link Task#createdNanos}) — it has been
* waiting longest, so it is the one most likely to be what the reply actually answers. Every
* other open task keeps the ordinary {@code WORKER_FAILED} path it would take without a stranded
* reply at all.
*
* <p><strong>{@code matching.size() >= 2} together with {@code hadStrandedReply} is a different
* question, and today it is defence in depth rather than a path this codebase's public API can
* drive.</strong> This class has exactly two sites that ever acquire a target's entry in
* {@code sessionLocks} — {@link #send} and {@link #answer} — and both open a {@link Rendezvous}
* waiter for that same target as the very first thing they do after acquiring the lock, then hold
* lock and waiter together for the rest of their critical section ({@link #send} also clears
* {@link #strandedReplies} right there, the instant it opens its waiter — before it ever enqueues
* delivery). So "the session lock is held" and "a live waiter is open for it" are the same fact
* throughout this class, and {@link #reply}'s fast path always resolves a currently-open waiter
* directly rather than stranding. The two facts this method wants therefore cannot be produced
* side by side: while the lock is held, a real reply resolves the open waiter directly and never
* reaches {@link #strandedReplies}; the instant the lock is free, any parked matching task's own
* {@link #send} that is scheduled next wins it and, by opening its waiter, clears the strand again
* before this method ever runs. There is no way to hold that lock open-but-unaccepted from outside
* {@link #send}/{@link #answer} to freeze a window in between. Constructing both facts at once
* through {@code sendAsync}/{@code reply}/{@code ask}/{@code answer} would need a race against
* virtual-thread scheduling, not a deterministic sequence — so the oldest-wins code below stays as
* defence in depth against a regression to that mechanism (e.g. clearing {@link #strandedReplies}
* on a narrower condition than "any acceptance"), not because today's test suite exercises the
* conjunction. {@code matching.size() >= 2} alone, without a strand, is exactly what
* {@code abandonFailsEveryPendingAsyncTicketForTheReleasedTarget} already covers.
*
* @return true if a live waiter or an async task was failed (never true for one recovered as a
* reply — see the note above)
*/
public boolean abandon(String target, String reason) {
boolean hadStrandedReply = hasStrandedReply(target);
// CB-640: the session is gone — nothing will ever accept or deliver into it now.
strandedReplies.remove(target);
queuedDeliveries.remove(target);
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(target);
boolean failed = waiter != null && !waiter.isDone() && rendezvous.resolveFailure(waiter, reason);
boolean asyncFailed = false;
List<Task> matching = new ArrayList<>();
for (Task task : tasks.values()) {
if (target.equals(task.target) && task.question == null
&& task.future.complete(new Reply(Outcome.WORKER_FAILED, reason))) {
asyncFailed = true;
if (target.equals(task.target) && task.question == null && !task.future.isDone()) {
matching.add(task);
}
}
Task recoveryTask = null;
if (hadStrandedReply && !matching.isEmpty()) {
recoveryTask = matching.get(0);
for (Task candidate : matching) {
if (candidate.createdNanos < recoveryTask.createdNanos) {
recoveryTask = candidate;
}
}
}
Reply recovered = recoveryTask != null ? recoverStrandedReply(target) : null;
for (Task task : matching) {
boolean isRecovery = task == recoveryTask && recovered != null;
Reply outcome = isRecovery ? recovered : new Reply(Outcome.WORKER_FAILED, reason);
if (task.future.complete(outcome)) {
if (outcome.outcome() == Outcome.WORKER_FAILED) {
asyncFailed = true;
} else if (task.turnId != null) {
asyncTasksByTurn.remove(task.turnId, task);
}
} else if (isRecovery) {
// The recovered reply was already drained out of the inbox, but this task resolved
// through another path (e.g. a concurrent reply() or a second abandon() racing this
// one) between us choosing it and completing it here. Put the reply back rather than
// lose it silently — it may still belong to some other still-open task, or the next
// caller that drains this target's inbox.
inbox.publish(target, UUID.randomUUID().toString(), recovered.text());
}
}
if (failed) {
@@ -456,6 +620,25 @@ public final class MessageService {
return failed || asyncFailed;
}
/**
* Drain {@code target}'s inbox and hand its content back as a {@link Outcome#REPLIED} result
* (#137 defence in depth for {@link #abandon}) — {@code null} if it turned out empty (the
* stranding fact raced away, e.g. a lead's own {@code fleet_poll} on the raw session already
* drained it first). When more than one message is queued, only the newest is the worker's actual
* final answer ({@link #drainReplies} returns them oldest-first).
*
* <p>This does drain (removes the messages from the inbox) before the caller knows whether the
* task it is recovering for will actually accept them — {@link #abandon} is the one that puts a
* reply back if its {@code complete} call turns out to lose the race.
*/
private Reply recoverStrandedReply(String target) {
var messages = drainReplies(target);
if (messages.isEmpty()) {
return null;
}
return new Reply(Outcome.REPLIED, messages.get(messages.size() - 1).content());
}
/**
* Acknowledge a specific reply by {@code msgId} for {@code target}. Removes it from the inbox
* so that a subsequent drain or peek no longer returns it.
@@ -76,20 +76,6 @@ public final class BackendQuarantine {
quarantinedUntilNanos.put(credentialId, nowNanos.getAsLong() + cooldownNanos);
}
/**
* Quarantine {@code credentialId} until the daemon restarts. This is for a configuration error
* that cannot heal with time, unlike an exhausted backend. A model selector that opencode silently
* resolves to another model stays wrong until an operator changes the profile, so a cooldown would
* start the same unsafe work again.
*/
public void quarantinePermanently(String credentialId) {
Objects.requireNonNull(credentialId, "credentialId");
if (inert) {
return;
}
quarantinedUntilNanos.put(credentialId, Long.MAX_VALUE);
}
/** Whether {@code credentialId} is quarantined right now. */
public boolean isQuarantined(String credentialId) {
return remainingNanos(credentialId) > 0;
@@ -124,7 +110,6 @@ public final class BackendQuarantine {
}
private static long toSecondsRoundedUp(long nanos) {
long seconds = nanos / 1_000_000_000L;
return seconds + (nanos % 1_000_000_000L == 0 ? 0 : 1);
return (nanos + 999_999_999L) / 1_000_000_000L;
}
}
@@ -23,6 +23,7 @@ import java.util.Set;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.Consumer;
import java.util.function.Function;
import java.util.stream.Collectors;
/**
@@ -88,24 +89,64 @@ public final class GitWorktrees implements Worktrees {
);
private final String configuredRoot;
/** OS group name for {@link #shareWithGroup} (fleetd #185 stage 3); {@code null} ⇒ feature off. */
private final String group;
private final Consumer<String> afterWorktreeAdded;
/** How {@link #shareWithGroup}'s processes (git config / chgrp / chmod / find) actually run.
* Defaults to the real {@link #exec(String...)}. Package-private test seam so a unit test can
* prove "no group configured ⇒ zero processes spawned" and inspect exactly what a configured
* group runs, without a real second OS user or OS group on this host. */
private final Function<String[], String> shareGroupRunner;
private final SecureRandom random = new SecureRandom();
private final AtomicLong seq = new AtomicLong();
/** Default constructor: worktree root is derived per-repo as {@code <repoRoot>/../.bridged-worktrees}. */
public GitWorktrees() {
this(null);
this(null, (String) null);
}
/** @param configuredRoot nullable absolute or relative path; null/blank derives a sibling of the repo root. */
/**
* @param configuredRoot nullable absolute or relative path; null/blank derives a sibling of
* the repo root. No {@code worktreeGroup} configured — {@link #shareWithGroup}
* is a no-op.
*/
public GitWorktrees(String configuredRoot) {
this(configuredRoot, _ -> {});
this(configuredRoot, (String) null);
}
/**
* @param configuredRoot nullable absolute or relative path; null/blank derives a sibling of
* the repo root.
* @param group optional OS group name (fleetd #185 stage 3, {@code worktreeGroup:} in
* config); null/blank ⇒ {@link #shareWithGroup} is a no-op.
*/
public GitWorktrees(String configuredRoot, String group) {
this(configuredRoot, group, _ -> {});
}
/** Test seam for changing a real worktree between its creation and its security check. */
GitWorktrees(String configuredRoot, Consumer<String> afterWorktreeAdded) {
this(configuredRoot, null, afterWorktreeAdded);
}
/** Test seam combining a configurable {@code group} with {@link #afterWorktreeAdded}. */
GitWorktrees(String configuredRoot, String group, Consumer<String> afterWorktreeAdded) {
this(configuredRoot, group, afterWorktreeAdded, null);
}
/**
* Full test seam: also overrides how {@link #shareWithGroup}'s processes run (fleetd #185
* stage 3), so a unit test can prove "no group configured ⇒ no process spawned" and inspect
* exactly what commands a configured group runs, without a real second OS user/group.
*
* @param shareGroupRunner {@code null} ⇒ the real {@link #exec(String...)}.
*/
GitWorktrees(String configuredRoot, String group, Consumer<String> afterWorktreeAdded,
Function<String[], String> shareGroupRunner) {
this.configuredRoot = configuredRoot;
this.group = (group == null || group.isBlank()) ? null : group;
this.afterWorktreeAdded = afterWorktreeAdded == null ? _ -> {} : afterWorktreeAdded;
this.shareGroupRunner = shareGroupRunner != null ? shareGroupRunner : this::exec;
}
@Override
@@ -607,6 +648,126 @@ public final class GitWorktrees implements Worktrees {
return deleted;
}
/**
* {@inheritDoc}
*
* <p>fleetd #185 stage 3. No-op — no process spawned, nothing logged — when {@link #group} is
* null/blank. Otherwise:
* <ol>
* <li>{@code git -C repoRoot config core.sharedRepository group} so every future write by
* either uid stays group-writable;</li>
* <li>a one-time {@code chgrp}/{@code chmod g+rwX} fix-up over the worktree directory and,
* under the repo's <em>common</em> git directory, {@code objects}, {@code refs},
* {@code logs}, {@code worktrees} and {@code packed-refs} — with setgid
* ({@code chmod g+s}) applied only to the directories among them, so files created later
* inherit the group;</li>
* <li>one INFO line naming the group and the paths touched.</li>
* </ol>
*
* <p><b>Every path is skipped when it does not exist.</b> {@code .git/logs} is absent in a repo
* with {@code core.logAllRefUpdates=false} or one that has had no ref update yet, and
* {@code packed-refs} is absent until refs are packed. Passing a missing path to {@code chgrp}
* exits non-zero, which would fail <em>every</em> provisioning spawn with a message blaming a
* group that is in fact fine.
*
* <p><b>The git directory is resolved, not assumed.</b> {@code <repoRoot>/.git} is a
* <em>file</em>, not a directory, when the checkout is itself a linked worktree — the very
* thing this class creates for every member. {@code git rev-parse --git-common-dir} gives the
* real shared store, and it may answer relatively, so it is resolved against {@code repoRoot}.
*
* <p><b>The fix-up re-runs on every spawn, by design.</b> {@code core.sharedRepository=group}
* governs only what git writes <em>after</em> it is set; the walk is what covers everything
* already on disk. It is not redundant work to optimise away — dropping it silently leaves
* pre-existing objects unreadable to the member. It costs three walks of the object store per
* spawn (about 3000 files in this repo, well under a second, but it grows with the repo).
*
* <p>This only fixes up file ownership/permissions on the operator's shared repo so a
* different-uid member can write to it — it isolates credentials, not the repository. A member
* in the group can still write the operator's git objects and refs.
*
* <p>Fails loudly: a missing group, or a {@code chgrp}/{@code chmod} refused because the
* operator is not a member of it, becomes a {@link WorktreeException} naming the group — never
* a silent skip that leaves a member unable to work with nothing in the log to explain why.
*/
@Override
public void shareWithGroup(String repoRoot, String worktreePath) {
if (group == null) {
return;
}
List<String> touched = new ArrayList<>();
try {
shareGroupRunner.apply(new String[]{"git", "-C", repoRoot, "config", "core.sharedRepository", "group"});
String commonDir = gitCommonDir(repoRoot);
shareGroupPathIfPresent(worktreePath, true, touched);
for (String name : List.of("objects", "refs", "logs", "worktrees")) {
shareGroupPathIfPresent(commonDir + "/" + name, true, touched);
}
shareGroupPathIfPresent(commonDir + "/packed-refs", false, touched);
} catch (WorktreeException e) {
throw new WorktreeException("cannot share worktree with group '" + group + "': "
+ e.getMessage() + " — the group must exist, and the fleetd operator ("
+ System.getProperty("user.name") + ") must be a member of it", e);
}
log.info("worktreeGroup={} shared repoRoot={} worktreePath={} paths={}",
group, repoRoot, worktreePath, touched);
}
/**
* The repo's <em>common</em> git directory as an absolute path — where {@code objects},
* {@code refs} and {@code worktrees} actually live. {@code git rev-parse --git-common-dir}
* answers relative to {@code repoRoot} in the ordinary case ({@code .git}) and absolutely for a
* linked worktree, so the answer is resolved against {@code repoRoot} either way. Never
* hardcode {@code repoRoot + "/.git"}: that is a FILE when the checkout is itself a linked
* worktree.
*/
private String gitCommonDir(String repoRoot) {
String answer = shareGroupRunner.apply(
new String[]{"git", "-C", repoRoot, "rev-parse", "--git-common-dir"});
String trimmed = answer == null ? "" : answer.trim();
if (trimmed.isEmpty()) {
trimmed = ".git";
}
return Path.of(repoRoot).resolve(trimmed).normalize().toString();
}
/**
* {@link #shareGroupPath} when {@code path} exists, recording it in {@code touched}; otherwise
* nothing at all. A missing path is normal, not an error — see {@link #shareWithGroup}'s
* javadoc for which ones are routinely absent and why passing them to {@code chgrp} would fail
* every spawn.
*/
private void shareGroupPathIfPresent(String path, boolean recursive, List<String> touched) {
if (!Files.exists(Path.of(path))) {
return;
}
shareGroupPath(path, recursive);
touched.add(path);
}
/**
* {@code chgrp}/{@code chmod g+rwX} {@code path} to {@link #group}. When {@code recursive},
* also walks the directories under {@code path} (including {@code path} itself, when it is a
* directory) and sets setgid on each — directories only, per the javadoc on
* {@link #shareWithGroup}.
*/
private void shareGroupPath(String path, boolean recursive) {
List<String> chgrp = new ArrayList<>(List.of("chgrp"));
if (recursive) chgrp.add("-R");
chgrp.add(group);
chgrp.add(path);
shareGroupRunner.apply(chgrp.toArray(new String[0]));
List<String> chmod = new ArrayList<>(List.of("chmod"));
if (recursive) chmod.add("-R");
chmod.add("g+rwX");
chmod.add(path);
shareGroupRunner.apply(chmod.toArray(new String[0]));
if (recursive) {
shareGroupRunner.apply(new String[]{"find", path, "-type", "d", "-exec", "chmod", "g+s", "{}", "+"});
}
}
/**
* Every {@code refs/wip/*} ref (see {@link WipRef}). The committer date is read as a unix
* count of seconds and converted to millis. {@code %00} (NUL) separates the fields because a
@@ -472,6 +472,10 @@ public final class SessionManager implements TurnListener {
try {
path = worktrees.add(repoRoot, branch, wt.baseRef());
worktrees.overlayParity(repoRoot, path, launcher.parityOverlay(preResolvedProfile));
// fleetd #185 stage 3: MUST run after overlayParity, not folded into add() — overlayParity
// copies more files into the worktree after add() returns, so sharing the group any earlier
// leaves those overlay files operator-owned and read-only for a different-uid member.
worktrees.shareWithGroup(repoRoot, path);
handle = launcher.spawn(new SpawnRequest(profile, path, callerCwd, sessionName, resumeSessionId, memberRole));
} catch (RuntimeException e) {
log.warn("spawn failed for profile={} role={} branch={} path={}: {}",
@@ -98,4 +98,21 @@ public interface Worktrees {
/** CB-586: the operator-visible census of {@code refs/wip/*} in one repository. */
record WipRefStats(int count, long costBytes) {
}
/**
* Make {@code repoRoot}'s git store and {@code worktreePath} writable by the configured group
* (fleetd #185 stage 3), so a member spawned as a different OS user (see
* {@code memberHerdrSocket}) can write its own worktree, its per-worktree git metadata, and
* its own commit objects. No-op when no group is configured.
*
* <p><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 — this only fixes file
* ownership/permissions so a different-uid member can work at all, it grants no narrower access
* than that.
*
* @param repoRoot the repository whose git store ({@code .git/objects}, {@code refs},
* {@code logs}, {@code worktrees}, {@code packed-refs}) needs sharing
* @param worktreePath the linked worktree's own directory
*/
void shareWithGroup(String repoRoot, String worktreePath);
}
@@ -1036,6 +1036,35 @@ class FleetConfigTest {
assertEquals(LeadMailbox.DEFAULT_PREFETCH, noEnv.prefetchOrDefault());
}
@Test
void absentWorktreeGroupLeavesItNull(@TempDir Path dir) throws Exception {
Path f = dir.resolve("no-worktree-group.yaml");
Files.writeString(f, "bind:\n port: 8080\n");
FleetConfig cfg = FleetConfig.load(f);
assertNull(cfg.worktreeGroup(), "no worktreeGroup: key → null → GitWorktrees.shareWithGroup is a no-op");
}
@Test
void worktreeGroupKeyParses(@TempDir Path dir) throws Exception {
Path f = dir.resolve("worktree-group.yaml");
Files.writeString(f, "bind:\n port: 8080\nworktreeGroup: fleet-workers\n");
FleetConfig cfg = FleetConfig.load(f);
assertEquals("fleet-workers", cfg.worktreeGroup());
}
@Test
void absentWorktreeGroupSurvivesTheBackCompatConstructorChain() {
// fleetd #185 stage 3: withDefaults() (and every pre-existing call site) must not silently
// drop a live worktreeGroup by routing through a back-compat constructor that defaults it
// to null.
FleetConfig cfg = new FleetConfig(null, null, null, Map.of(), null, null, null, null, null,
null, null, null, null, null, null, null, null, null, null, null, "fleet-workers");
assertEquals("fleet-workers", cfg.withDefaults().worktreeGroup(),
"withDefaults() must carry a configured worktreeGroup through unchanged");
}
@Test
void absentPrimaryBlockLeavesPrimaryNull(@TempDir Path dir) throws Exception {
Path f = dir.resolve("no-primary.yaml");
@@ -0,0 +1,27 @@
package dev.ltms.fleet.herdr;
import java.util.HashMap;
import java.util.Map;
import java.util.OptionalLong;
/**
* Fake {@link ParentResolver} backed by an explicit pid→parent map — lets {@link PaneLocatorTest}
* drive {@link PaneLocator}'s ancestry walk (grandchild pids, cycles) without spawning real OS
* processes.
*/
final class FakeParentResolver implements ParentResolver {
private final Map<Long, Long> parents = new HashMap<>();
/** {@code pid}'s parent is {@code parentPid}. A pid with no entry here has no known parent. */
FakeParentResolver parent(long pid, long parentPid) {
parents.put(pid, parentPid);
return this;
}
@Override
public OptionalLong parentOf(long pid) {
Long parent = parents.get(pid);
return parent == null ? OptionalLong.empty() : OptionalLong.of(parent);
}
}
@@ -1,7 +1,11 @@
package dev.ltms.fleet.herdr;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.junit.jupiter.api.Test;
import java.util.concurrent.atomic.AtomicInteger;
import static org.junit.jupiter.api.Assertions.*;
/** Unit tests for PID → pane resolution (the herdr half of connection-based MCP identity). */
@@ -63,4 +67,123 @@ class PaneLocatorTest {
long paneListCalls = shared.calls.stream().filter(c -> c.method().equals("pane.list")).count();
assertEquals(1, paneListCalls, "same-object lead/member must scan exactly once, not twice");
}
// --- ancestry walk (CB-161: grandchild pids matched no pane, resolving as primary) --------
@Test
void stillResolvesAPidThatIsExactlyThePaneShellPid() {
// Regression: a pid with no parent chain at all — no ancestry walk is needed to match it.
OnePaneHerdr pane = new OnePaneHerdr("term_x", "pX", 5000, 6000);
PaneLocator loc = new PaneLocator(pane, new FakeParentResolver());
assertEquals("term_x", loc.terminalForPid(5000));
}
@Test
void stillResolvesAPidThatIsExactlyAForegroundPid() {
// Regression: same as above, but matching via the foreground-processes list.
OnePaneHerdr pane = new OnePaneHerdr("term_x", "pX", 5000, 6000);
PaneLocator loc = new PaneLocator(pane, new FakeParentResolver());
assertEquals("term_x", loc.terminalForPid(6000));
}
@Test
void resolvesAGrandchildPidTwoLevelsBelowTheShellPid() {
// The bug: a helper process a worker spawns (python3, curl, ...) is a grandchild of the
// pane's shell — not the shell_pid and not a foreground pid directly. Before the fix,
// paneOwnsPid only checked direct pid equality, so this pid matched no pane and the
// caller fell through to loopback-trust as the primary.
OnePaneHerdr pane = new OnePaneHerdr("term_x", "pX", 5000, 6000);
FakeParentResolver parents = new FakeParentResolver()
.parent(7002, 7001) // grandchild -> child
.parent(7001, 5000); // child -> shell (the pane's shell_pid)
PaneLocator loc = new PaneLocator(pane, parents);
assertEquals("term_x", loc.terminalForPid(7002));
}
@Test
void nullForAPidWhoseAncestryMatchesNoPane() {
// Must not break the other direction: a pid that truly belongs to nothing here (e.g. the
// real primary) must still resolve to null. Resolving everything to a worker would demote
// the actual lead and refuse every orchestration call.
OnePaneHerdr pane = new OnePaneHerdr("term_x", "pX", 5000, 6000);
FakeParentResolver parents = new FakeParentResolver()
.parent(9002, 9001)
.parent(9001, 9000); // chain never reaches 5000 or 6000
PaneLocator loc = new PaneLocator(pane, parents);
assertNull(loc.terminalForPid(9002));
}
@Test
void ancestryWalkTerminatesOnACycleInsteadOfHanging() {
// A fake (or corrupted) parent map that cycles must not hang identity resolution, which
// runs on every MCP call. The walk must still terminate and correctly resolve to null.
OnePaneHerdr pane = new OnePaneHerdr("term_x", "pX", 5000, 6000);
FakeParentResolver parents = new FakeParentResolver()
.parent(100, 101)
.parent(101, 100); // cycle, never reaches the pane's pids
PaneLocator loc = new PaneLocator(pane, parents);
assertNull(loc.terminalForPid(100));
}
@Test
void ancestrySetIsComputedOnceAcrossBothClientsInTheTwoDaemonConstructor() {
// CB-185: the two-daemon constructor searches lead then member. The ancestor set is
// per-caller, not per-client — it must be walked once and reused, not recomputed for
// each client searched.
OnePaneHerdr pane = new OnePaneHerdr("term_x", "pX", 5000, 6000);
FakeParentResolver parents = new FakeParentResolver()
.parent(7002, 7001)
.parent(7001, 5000);
AtomicInteger calls = new AtomicInteger();
ParentResolver counting = pid -> {
calls.incrementAndGet();
return parents.parentOf(pid);
};
HerdrClient noPanes = new FakeHerdr().withNoPanes();
PaneLocator two = new PaneLocator(noPanes, pane, counting);
assertEquals("term_x", two.terminalForPid(7002));
assertEquals(3, calls.get(), "ancestry must be walked once (3 lookups: 7002, 7001, 5000), "
+ "not re-walked per herdr client");
}
/** Minimal single-pane {@link HerdrClient} fake, purpose-built for the ancestry tests above. */
private static final class OnePaneHerdr implements HerdrClient {
private final ObjectMapper mapper = new ObjectMapper();
private final String terminalId;
private final String paneId;
private final long shellPid;
private final long foregroundPid;
OnePaneHerdr(String terminalId, String paneId, long shellPid, long foregroundPid) {
this.terminalId = terminalId;
this.paneId = paneId;
this.shellPid = shellPid;
this.foregroundPid = foregroundPid;
}
@Override
public JsonNode call(String method, Object params) {
try {
return switch (method) {
case "pane.list" -> mapper.readTree(("""
{"type":"pane_list","panes":[
{"pane_id":"%s","terminal_id":"%s","workspace_id":"w1","tab_id":"w1:t1","agent":"claude"}]}""")
.formatted(paneId, terminalId));
case "pane.process_info" -> mapper.readTree(("""
{"type":"pane_process_info","process_info":{"pane_id":"%s","shell_pid":%d,
"foreground_processes":[{"pid":%d,"name":"node","argv0":"claude"}]}}""")
.formatted(paneId, shellPid, foregroundPid));
default -> throw new HerdrException("OnePaneHerdr has no canned response for " + method);
};
} catch (HerdrException e) {
throw e;
} catch (Exception e) {
throw new HerdrException("OnePaneHerdr decode failed for " + method, e);
}
}
@Override
public void close() {
}
}
}
@@ -161,6 +161,34 @@ class HerdrPeerLauncherAllowListWiringTest {
+ "it under allow: — sshAuthSock is unset here, so it defaults to block");
}
@Test
void brokerUriEnvStaysBlockedWhenListedInMemberCredentialsAllow() {
FakeHerdr herdr = new FakeHerdr();
WiringLauncher launcher = new WiringLauncher(herdr,
allowListWithAllow(List.of("BROKER_CONNECTION_URI")), "/bin/zsh", null,
() -> config("BROKER_CONNECTION_URI"));
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
Path dir = Path.of(launcher.env.get("ZDOTDIR"));
assertFalse(readAll(dir.resolve(EnvAllowListScrub.SCRUB_FILE)).contains("'BROKER_CONNECTION_URI'"),
"broker.uriEnv must not reach a member even when listed in memberCredentials.allow:");
}
@Test
void brokerUriEnvIsDeniedUnderTheDenyListPolicyEvenWhenAllowed() {
FakeHerdr herdr = new FakeHerdr();
WiringLauncher launcher = new WiringLauncher(herdr,
() -> new FleetConfig.MemberCredentials(null, List.of("BROKER_CONNECTION_URI"), List.of(), null),
"/bin/bash", null,
() -> config("BROKER_CONNECTION_URI"));
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
assertEquals("blocked-by-fleetd-cb596-see-gitea-issue-82", launcher.env.get("BROKER_CONNECTION_URI"),
"the deny-list overlay must deny broker.uriEnv even when allow: names it");
}
/**
* CB-633 follow-up criterion 3: on every allow-list spawn the daemon logs one INFO line, shaped
* "member credentials: allowed N of M", with real counts — not constants. Real path: the count
@@ -260,15 +288,24 @@ class HerdrPeerLauncherAllowListWiringTest {
/** Plus an injectable {@code hostEnvNames} source, for the "allowed N of M" log line test. */
WiringLauncher(FakeHerdr herdr, Supplier<FleetConfig.MemberCredentials> creds, String shell,
Supplier<Set<String>> hostEnvNames) {
Supplier<Set<String>> hostEnvNames) {
this(herdr, creds, shell, hostEnvNames, null);
}
WiringLauncher(FakeHerdr herdr, Supplier<FleetConfig.MemberCredentials> creds, String shell,
Supplier<Set<String>> hostEnvNames, Supplier<FleetConfig> config) {
super("test", new AgentControl(herdr), new WorkspaceControl(herdr),
Map.of("test", profile()), "test",
name -> "SHELL".equals(name) ? shell : null,
0, () -> 0L, () -> { }, null, creds, hostEnvNames);
0, () -> 0L, () -> { }, null, creds, hostEnvNames, config);
}
@Override
protected Launch buildLaunch(FleetConfig.Profile cfg, LaunchSpec spec) {
Map<String, String> launchEnv = baseEnv(cfg);
launchEnv.putAll(env);
env.clear();
env.putAll(launchEnv);
return new Launch(env, List.of("test"));
}
@@ -278,6 +315,12 @@ class HerdrPeerLauncherAllowListWiringTest {
}
}
private static FleetConfig config(String brokerUriEnv) {
return new FleetConfig(null, null, null, Map.of(), null, null, null, null, null,
new FleetConfig.Broker(null, brokerUriEnv, null), null, null, null, null, null,
null, null, null, null, null).withDefaults();
}
/** The generated directory is a temp directory; make sure the test does not leave a pile. */
@Test
void theGeneratedDirectoryIsRemovedWhenThePaneIsStopped() {
@@ -121,6 +121,34 @@ class MemberEnvAllowListTest {
assertTrue(derived.contains("OTHER_NAME"), "other allow: names are unaffected");
}
@Test
void configuredBrokerAndCoordinatorUriEnvNamesAreExcludedEvenWhenAllowed() {
FleetConfig config = config("BROKER_CONNECTION_URI", "COORDINATOR_CONNECTION_URI");
Set<String> excluded = MemberEnvAllowList.brokerUriEnvNames(config);
Set<String> derived = MemberEnvAllowList.derive(List.of(),
Set.of("BROKER_CONNECTION_URI", "COORDINATOR_CONNECTION_URI", "OTHER_NAME"), excluded);
assertFalse(derived.contains("BROKER_CONNECTION_URI"),
"broker.uriEnv is secret-bearing and must not ride in on allow:");
assertFalse(derived.contains("COORDINATOR_CONNECTION_URI"),
"coordinator.uriEnv has the same inline-password shape");
assertTrue(derived.contains("OTHER_NAME"), "unrelated allow: entries are unaffected");
}
@Test
void absentOrBlankBrokerUriEnvAddsNoExclusions() {
assertTrue(MemberEnvAllowList.brokerUriEnvNames(config(null, null)).isEmpty());
assertTrue(MemberEnvAllowList.brokerUriEnvNames(config(" ", "")).isEmpty());
}
private static FleetConfig config(String brokerUriEnv, String coordinatorUriEnv) {
return new FleetConfig(null, null, null, Map.of(), null, null, null, null, null,
new FleetConfig.Broker(null, brokerUriEnv, null), null, null, null, null, null,
null, null, null, null,
new FleetConfig.Coordinator(null, coordinatorUriEnv, null, null)).withDefaults();
}
/** {@code LC_*} categories are infrastructure by prefix; everything else needs an exact match. */
@Test
void keepsMatchesExactlyPlusTheLocalePrefixRule() {
@@ -12,7 +12,6 @@ import dev.ltms.fleet.peer.MemberRole;
import dev.ltms.fleet.peer.PeerHandle;
import dev.ltms.fleet.peer.PeerUnreachableException;
import dev.ltms.fleet.peer.SpawnRequest;
import dev.ltms.fleet.placement.BackendQuarantine;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
@@ -23,7 +22,6 @@ import java.util.Map;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.function.Supplier;
import static org.junit.jupiter.api.Assertions.*;
@@ -63,23 +61,6 @@ class OpenCodeLauncherTest {
0, System::currentTimeMillis, () -> { }, configRoot, configRoot, null, () -> creds);
}
private static OpenCodeLauncher serviceWithModelMismatchSink(FakeHerdr herdr, Path root,
FleetConfig.Profile cfg,
BackendQuarantine quarantine) {
return new OpenCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null,
0, System::currentTimeMillis, () -> { }, root, root, null, null,
(profile, _) -> quarantine.quarantinePermanently(profile));
}
private static void writeResolvedModelRecord(Path root, String directory, String providerId,
String modelId) throws Exception {
Path record = Files.createDirectories(root.resolve("session").resolve("p1")).resolve("ses_a.json");
Files.writeString(record, "{\"id\":\"ses_a\",\"directory\":\"" + directory
+ "\",\"model\":{\"id\":\"" + modelId + "\",\"providerID\":\""
+ providerId + "\",\"variant\":\"high\"}}");
}
@SuppressWarnings("unchecked")
private static Map<String, Object> lastStart(FakeHerdr herdr) {
return (Map<String, Object>) herdr.lastCall("agent.start").params();
@@ -230,50 +211,6 @@ class OpenCodeLauncherTest {
"--auto is unconditional: a model-less worker still must never block on approval");
}
// --- CB-175: verify opencode's recorded resolved model after spawn readiness ----------------
@Test
void matchingResolvedModelDoesNotQuarantineTheProfile(@TempDir Path root) throws Exception {
FleetConfig.Profile cfg = opencodeCfg("opencode/x-preview-f-free", null, null);
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(1));
writeResolvedModelRecord(root, root.toString(), "opencode", "x-preview-f-free");
serviceWithModelMismatchSink(new FakeHerdr(), root, cfg, quarantine)
.spawn(new SpawnRequest(null, root.toString(), null, null, null, MemberRole.DEV));
assertFalse(quarantine.isQuarantined("gemini"), "the exact provider/model match is safe");
}
@Test
void differentResolvedModelPermanentlyQuarantinesTheProfile(@TempDir Path root) throws Exception {
FleetConfig.Profile cfg = opencodeCfg("opencode/x-preview-f-free", null, null);
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(1));
writeResolvedModelRecord(root, root.toString(), "openai", "gpt-5.6-sol");
serviceWithModelMismatchSink(new FakeHerdr(), root, cfg, quarantine)
.spawn(new SpawnRequest(null, root.toString(), null, null, null, MemberRole.DEV));
assertTrue(quarantine.isQuarantined("gemini"), "a fallback model must block later spawns");
assertEquals(Long.MAX_VALUE / 1_000_000_000L + 1, quarantine.remainingSeconds("gemini").orElseThrow(),
"a withdrawn selector cannot become safe after the normal cooldown");
}
@Test
void missingOrUnreadableSessionDatabaseDoesNotQuarantine(@TempDir Path root) throws Exception {
FleetConfig.Profile cfg = opencodeCfg("opencode/x-preview-f-free", null, null);
BackendQuarantine missing = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(1));
serviceWithModelMismatchSink(new FakeHerdr(), root, cfg, missing)
.spawn(new SpawnRequest(null, root.toString(), null, null, null, MemberRole.DEV));
assertFalse(missing.isQuarantined("gemini"), "a missing database is unknown evidence");
Path broken = Files.createDirectories(root.resolve("session").resolve("p1")).resolve("ses_a.json");
Files.writeString(broken, "not JSON");
BackendQuarantine unreadable = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(1));
serviceWithModelMismatchSink(new FakeHerdr(), root, cfg, unreadable)
.spawn(new SpawnRequest(null, root.toString(), null, null, null, MemberRole.DEV));
assertFalse(unreadable.isQuarantined("gemini"), "an unreadable database is unknown evidence");
}
// --- CB-617: --agent <role> when the role has an agent-definition file --------------------
@Test
@@ -397,8 +334,7 @@ class OpenCodeLauncherTest {
assertNull(handle.agentSessionId(), "no record yet → null, not a spawn-time block");
// Once the record appears (here: same cwd), lazy discovery resolves it — the handle's
// session id matches its own worktree, not another's.
OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "p1", "ses_a.json",
"ses_resolved", "/work/dir", 1000L);
OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "ses_resolved", "/work/dir", 1000L);
assertEquals("ses_resolved", handle.agentSessionId(),
"agentSessionId() re-scans and picks up a record that has since been written");
}
@@ -5,68 +5,89 @@ import org.junit.jupiter.api.io.TempDir;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.attribute.FileTime;
import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.PreparedStatement;
import java.sql.SQLException;
import java.sql.Statement;
import static org.junit.jupiter.api.Assertions.*;
/**
* {@link OpenCodeSessionDiscovery} matches an opencode session record by the worker's cwd (its
* {@code directory}) against opencode's on-disk storage. These tests populate a TEMP storage root
* themselves — never the operator's real {@code ~/.local/share/opencode}.
* {@link OpenCodeSessionDiscovery} matches an opencode session row by the worker's cwd (its
* {@code directory}) against opencode's {@code opencode.db} SQLite database. These tests build a
* SYNTHETIC database themselves, in a JUnit temp directory — never the operator's real
* {@code ~/.local/share/opencode/opencode.db}, which a live opencode process may be writing.
*/
class OpenCodeSessionDiscoveryTest {
/**
* Write a session record {@code {"id":..., "directory":...}} under
* {@code <root>/session/<projectID>/<fileName>} and stamp it with a known last-modified time,
* so "most recently modified wins" is deterministic. Static so the launcher test can reuse it.
* Create {@code <root>/opencode.db} with a minimal {@code session} table (just the columns
* {@link OpenCodeSessionDiscovery} reads: {@code id}, {@code directory}, {@code time_updated})
* and insert one row. Static so {@link OpenCodeLauncherTest} can reuse it.
*/
static void writeRecord(Path root, String projectId, String fileName, String id,
String directory, long lastModifiedEpochMillis) throws Exception {
Path dir = root.resolve("session").resolve(projectId);
Files.createDirectories(dir);
Path file = dir.resolve(fileName);
Files.writeString(file, "{\"id\":\"" + id + "\",\"directory\":\"" + directory
+ "\",\"projectID\":\"" + projectId + "\",\"version\":\"1.1.31\"}");
Files.setLastModifiedTime(file, FileTime.fromMillis(lastModifiedEpochMillis));
static void writeRecord(Path root, String id, String directory, long timeUpdated) throws Exception {
Path db = root.resolve("opencode.db");
try (Connection connection = DriverManager.getConnection("jdbc:sqlite:" + db)) {
try (Statement statement = connection.createStatement()) {
statement.execute("CREATE TABLE IF NOT EXISTS session ("
+ "id TEXT PRIMARY KEY, directory TEXT, time_updated INTEGER)");
}
// Bound parameters, not string interpolation: the class under test uses a
// PreparedStatement, and a hand-escaped INSERT here is a pattern someone copies out.
try (PreparedStatement insert = connection.prepareStatement(
"INSERT INTO session (id, directory, time_updated) VALUES (?, ?, ?)")) {
insert.setString(1, id);
insert.setString(2, directory);
insert.setLong(3, timeUpdated);
insert.executeUpdate();
}
}
}
@Test
void findsTheRecordWhoseDirectoryEqualsTheCwd(@TempDir Path root) throws Exception {
writeRecord(root, "p1", "ses_a.json", "ses_aaa", "/w/a", 1000L);
writeRecord(root, "p2", "ses_b.json", "ses_bbb", "/w/b", 2000L);
void findsTheRowWhoseDirectoryEqualsTheCwd(@TempDir Path root) throws Exception {
writeRecord(root, "ses_aaa", "/w/a", 1000L);
writeRecord(root, "ses_bbb", "/w/b", 2000L);
assertEquals("ses_bbb", new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/b"),
"the record whose directory equals the cwd is the one found");
"the row whose directory equals the cwd is the one found");
assertEquals("ses_aaa", new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"));
}
@Test
void aNonMatchingDirectoryYieldsNullRatherThanAMismatch(@TempDir Path root) throws Exception {
writeRecord(root, "p1", "ses_a.json", "ses_aaa", "/w/a", 1000L);
writeRecord(root, "ses_aaa", "/w/a", 1000L);
assertNull(new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/other"),
"no record for this cwd yet → null, not a wrong session");
"no row for this cwd yet → null, not a wrong session");
}
@Test
void prefersTheMostRecentlyModifiedRecordWhenSeveralMatch(@TempDir Path root) throws Exception {
writeRecord(root, "p1", "old.json", "ses_old", "/w/a", 1000L);
writeRecord(root, "p2", "new.json", "ses_new", "/w/a", 5000L);
void prefersTheMostRecentlyUpdatedRowWhenSeveralMatch(@TempDir Path root) throws Exception {
writeRecord(root, "ses_old", "/w/a", 1000L);
writeRecord(root, "ses_new", "/w/a", 5000L);
assertEquals("ses_new", new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"),
"the freshest record for the cwd wins");
"the row with the highest time_updated for the cwd wins");
}
@Test
void aMissingOrEmptyStorageRootYieldsNullWithoutThrowing(@TempDir Path root) throws Exception {
// Missing: no session dir at all under the root.
void aMissingDatabaseYieldsNullWithoutThrowing(@TempDir Path root) {
// No opencode.db at all under the root.
assertNull(new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"));
}
// Present but empty: a session dir with nothing in it produces no match, not a throw.
Path emptyRoot = root.resolve("empty");
Files.createDirectories(emptyRoot.resolve("session"));
assertNull(new OpenCodeSessionDiscovery(emptyRoot).sessionIdForDirectory("/w/a"));
@Test
void anEmptyDatabaseYieldsNullWithoutThrowing(@TempDir Path root) throws Exception {
Path db = root.resolve("opencode.db");
try (Connection connection = DriverManager.getConnection("jdbc:sqlite:" + db);
Statement statement = connection.createStatement()) {
statement.execute("CREATE TABLE session (id TEXT PRIMARY KEY, directory TEXT, "
+ "time_updated INTEGER)");
}
assertNull(new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"));
}
@Test
@@ -76,15 +97,41 @@ class OpenCodeSessionDiscoveryTest {
assertNull(discovery.sessionIdForDirectory(" "));
}
/**
* The one line standing between fleetd and writing the operator's live {@code opencode.db} —
* 841MB, with a running opencode writing it — is {@code config.setReadOnly(true)} in
* {@link OpenCodeSessionDiscovery#openReadOnly()}. Delete it and every other test in this class
* still passes, so this is the test that guards it.
*
* <p>It asks the connection to write, and requires a refusal. The obvious alternative — make
* the database file unwritable and check the read still works — proves nothing: SQLite silently
* downgrades a read-write open of an unwritable file to read-only, so that test passes either
* way. It was tried and watched pass with the flag removed.
*/
@Test
void aMalformedRecordIsSkippedRatherThanFatal(@TempDir Path root) throws Exception {
// A record that fails to parse must not abort the scan of its siblings.
Path dir = root.resolve("session").resolve("p1");
Files.createDirectories(dir);
Files.writeString(dir.resolve("broken.json"), "{not valid json");
writeRecord(root, "p1", "good.json", "ses_good", "/w/a", 1000L);
void theDatabaseIsOpenedReadOnly(@TempDir Path root) throws Exception {
writeRecord(root, "ses_aaa", "/w/a", 1000L);
assertEquals("ses_good", new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"),
"an unreadable record is skipped; a later valid one still matches");
try (Connection connection = new OpenCodeSessionDiscovery(root).openReadOnly();
Statement statement = connection.createStatement()) {
SQLException refused = assertThrows(SQLException.class,
() -> statement.executeUpdate("INSERT INTO session (id, directory, time_updated) "
+ "VALUES ('ses_zzz', '/w/z', 1)"),
"the connection must REFUSE a write — opencode is writing this database live");
assertTrue(refused.getMessage().toLowerCase().contains("readonly")
|| refused.getMessage().toLowerCase().contains("read-only"),
"the refusal must be about read-only, not some other error: " + refused.getMessage());
}
}
@Test
void aCorruptDatabaseFileYieldsNullWithoutThrowing(@TempDir Path root) throws Exception {
// A file at opencode.db that is not a SQLite database at all — the open/query must fail
// safe, never fatal to a spawn.
Path db = root.resolve("opencode.db");
Files.writeString(db, "this is not a sqlite database");
assertNull(new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"),
"an unreadable database resolves to null, not an exception");
}
}
@@ -688,6 +688,32 @@ class MessageServiceTest {
assertFailedTicket(third, "agent target term_a not found");
}
// --- #137 follow-up: abandon() must not guess when more than one task is open ---------------
//
// A test combining a genuine stranded reply (hasStrandedReply(T)==true) with two simultaneously
// open matching tasks was attempted here and removed after investigation showed the combination
// is not reachable through the public API today, not merely hard to time right:
//
// This class has exactly two call sites that ever hold a target's entry in the session-lock map
// (send() and answer()), and both open a Rendezvous waiter for that same target as the first thing
// they do after acquiring the lock, holding lock and waiter together for their whole critical
// section. So "the lock is held" and "a live waiter is open" are the same fact throughout this
// class. reply()'s fast path always resolves a currently-open waiter directly instead of
// stranding — so a strand can only be created while NO task is accepted (lock free), and the
// instant the lock is next taken (by any parked matching task's own send(), the moment it is
// scheduled), that acceptance clears strandedReplies again (see send()'s CB-640 comment) before
// abandon() can ever observe both facts together. Confirmed empirically too: an earlier version of
// this test stranded a reply, then created an "accepted" task (awaitWaiting()) followed by a
// "parked" one — and the accepted task's own acceptance silently cleared the strand it was
// supposed to be racing against, so the parked task came back WORKER_FAILED instead of DONE, not
// because the fix was missing but because the test's premise could not be constructed.
//
// The reachable half — matching.size() >= 2 alone, no strand — is exactly what
// abandonFailsEveryPendingAsyncTicketForTheReleasedTarget already covers (all fail, none guess).
// The oldest-wins code in abandon() stays as defence in depth (see its own javadoc) against a
// regression that would make the conjunction reachable, e.g. clearing strandedReplies on a
// narrower condition than "any acceptance" — not because this suite exercises it today.
@Test
void abandonDoesNotFailAnAsyncTicketWaitingForAnAnswer() throws Exception {
String ticket = messages.sendAsync(T, "task that asks");
@@ -710,6 +736,68 @@ class MessageServiceTest {
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
}
// --- #137: a fleet_ask round-trip must not orphan the ticket's own reply -------------------
//
// The primary's fleet_send{turnId} answer call is itself bounded (a real MCP call, capped well
// under a minute) — far shorter than a resumed turn can genuinely take to finish real work. These
// drive the exact real delegation path (async send -> worker asks -> primary answers -> primary's
// own wait gives up -> worker's real fleet_reply arrives afterwards) rather than calling a reply
// sink directly, since the bug is specifically about which sink the resumed turn's reply reaches.
@Test
void aReplyAfterAnswerTimesOutStillCompletesTheAsyncTicket() throws Exception {
String ticket = messages.sendAsync(T, "task that asks");
awaitWaiting();
injectDelivery();
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
// The primary answers, but its own bounded wait for the worker's resumed turn is short and
// expires before the worker (still genuinely working) gets back to it.
MessageService.Reply answerReply = messages.answer(asking.turnId(), "config.yaml", 150);
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome(),
"the primary's own bounded wait gives up before the worker finishes resuming");
// The worker keeps working past that window and only now calls fleet_reply.
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE);
assertEquals("PR opened: https://example/pulls/42", done.reply(),
"fleet_poll{ticket} must return the worker's real reply, not stay pending forever");
assertEquals("reply", done.replySource());
assertFalse(messages.hasStrandedReply(T),
"the reply completed its own ticket directly and never touched the inbox");
}
@Test
void fleetStopAfterAnOrphanedReplyDoesNotFailTheTicket() throws Exception {
String ticket = messages.sendAsync(T, "task that asks");
awaitWaiting();
injectDelivery();
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
MessageService.Reply answerReply = messages.answer(asking.turnId(), "config.yaml", 150);
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome());
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
// fleet_stop tears the worker's session down right after the reply landed — this must never
// report the misleading "the worker session was released before it replied": a reply is
// exactly what happened.
assertFalse(messages.abandon(T, "the worker session was released before it replied"),
"a reply already arrived, so nothing here is a genuine failure");
MessageService.TaskView view = awaitTicketPhase(ticket, MessageService.Phase.DONE);
assertEquals("PR opened: https://example/pulls/42", view.reply());
}
@Test
void unansweredAsyncQuestionReturnsTheTicketToPendingAndReleasesItsTarget() throws Exception {
String ticket = messages.sendAsync(T, "task that asks");
@@ -30,12 +30,19 @@ public final class FakeWorktrees implements Worktrees {
public record PruneCall(String repoRoot, long minAgeMillis) {
}
public record ShareCall(String repoRoot, String worktreePath) {
}
private final List<AddCall> addCalls = new CopyOnWriteArrayList<>();
private final List<RemoveCall> removeCalls = new CopyOnWriteArrayList<>();
private final List<OverlayCall> overlayCalls = new CopyOnWriteArrayList<>();
private final List<RepoRootCall> repoRootCalls = new CopyOnWriteArrayList<>();
private final List<SnapshotCall> snapshotCalls = new CopyOnWriteArrayList<>();
private final List<PruneCall> pruneCalls = new CopyOnWriteArrayList<>();
private final List<ShareCall> shareCalls = new CopyOnWriteArrayList<>();
/** Tags every {@code overlayParity}/{@code shareWithGroup} call in call order, so a test can
* pin that sharing runs after the overlay copy (fleetd #185 stage 3). */
private final List<String> overlayShareOrder = new CopyOnWriteArrayList<>();
private final Set<String> existingPaths = ConcurrentHashMap.newKeySet();
private final Set<String> trackedPaths = ConcurrentHashMap.newKeySet();
private final AtomicLong snapshotSeq = new AtomicLong();
@@ -136,6 +143,13 @@ public final class FakeWorktrees implements Worktrees {
}
overlayCalls.add(new OverlayCall(repoRoot, worktreePath, List.copyOf(overlay),
List.copyOf(copied), List.copyOf(skipped)));
overlayShareOrder.add("overlay:" + worktreePath);
}
@Override
public void shareWithGroup(String repoRoot, String worktreePath) {
shareCalls.add(new ShareCall(repoRoot, worktreePath));
overlayShareOrder.add("share:" + worktreePath);
}
@Override
@@ -203,4 +217,17 @@ public final class FakeWorktrees implements Worktrees {
public SnapshotCall lastSnapshot() {
return snapshotCalls.isEmpty() ? null : snapshotCalls.getLast();
}
public List<ShareCall> shareCalls() {
return List.copyOf(shareCalls);
}
public ShareCall lastShare() {
return shareCalls.isEmpty() ? null : shareCalls.getLast();
}
/** Call-order tags ({@code "overlay:<path>"}/{@code "share:<path>"}) — see field javadoc. */
public List<String> overlayShareOrder() {
return List.copyOf(overlayShareOrder);
}
}
@@ -924,4 +924,179 @@ class GitWorktreesTest {
assertEquals(2, stats.count(), "two snapshot refs are reported");
assertTrue(stats.costBytes() > 0, "the cost of the snapshots is a positive byte count");
}
/**
* fleetd #185 stage 3: a recording {@link java.util.function.Function} test seam stands in for
* every process {@link GitWorktrees#shareWithGroup} would run — no real second OS user/group
* exists on this host, so these are unit tests against that seam, not a live-group integration
* test (out of scope per the ticket).
*/
private static List<String> joined(String[] command) {
return List.of(command);
}
/** Remove {@code path} and anything under it. Tolerates an already-absent path. */
private static void deleteRecursively(Path path) throws Exception {
if (!Files.exists(path)) {
return;
}
if (Files.isDirectory(path)) {
try (java.util.stream.Stream<Path> children = Files.list(path)) {
for (Path child : children.toList()) {
deleteRecursively(child);
}
}
}
Files.delete(path);
}
/** {@code worktreeGroup} absent ⇒ zero processes spawned and no git config written. */
@Test
void shareWithGroupIsNoopWhenNoGroupConfigured(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
List<List<String>> recorded = new java.util.ArrayList<>();
java.util.function.Function<String[], String> recordingRunner = cmd -> {
recorded.add(joined(cmd));
return "";
};
GitWorktrees gitWorktrees = new GitWorktrees(tmp.resolve("wts").toString(), null, _ -> {}, recordingRunner);
gitWorktrees.shareWithGroup(repo.toString(), repo.resolve("some-worktree").toString());
assertTrue(recorded.isEmpty(), "no group configured must spawn no process at all: " + recorded);
}
/** A configured group runs {@code git config core.sharedRepository group} first, then
* chgrp/chmod/setgid over every path {@link GitWorktrees#shareWithGroup} documents. */
@Test
void shareWithGroupRunsConfigThenChgrpChmodSetgidPerPath(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
String repoRoot = repo.toString();
Path worktree = Files.createDirectories(repo.resolve("some-worktree"));
String worktreePath = worktree.toString();
Files.createDirectories(repo.resolve(".git/worktrees"));
List<List<String>> recorded = new java.util.ArrayList<>();
java.util.function.Function<String[], String> recordingRunner = cmd -> {
recorded.add(joined(cmd));
// What real git answers for an ordinary (non-linked) checkout: relative to repoRoot.
return List.of(cmd).contains("--git-common-dir") ? ".git\n" : "";
};
GitWorktrees gitWorktrees =
new GitWorktrees(tmp.resolve("wts").toString(), "devteam", _ -> {}, recordingRunner);
gitWorktrees.shareWithGroup(repoRoot, worktreePath);
assertEquals(List.of("git", "-C", repoRoot, "config", "core.sharedRepository", "group"), recorded.get(0),
"core.sharedRepository must be set first, so it keeps working after the one-time fix-up");
assertTrue(recorded.contains(List.of("git", "-C", repoRoot, "rev-parse", "--git-common-dir")),
"the git dir must be asked for, never hardcoded as <repoRoot>/.git — that is a FILE "
+ "when the checkout is itself a linked worktree: " + recorded);
for (String dir : List.of(worktreePath, repoRoot + "/.git/objects", repoRoot + "/.git/refs",
repoRoot + "/.git/logs", repoRoot + "/.git/worktrees")) {
assertTrue(recorded.contains(List.of("chgrp", "-R", "devteam", dir)), "missing chgrp -R for " + dir);
assertTrue(recorded.contains(List.of("chmod", "-R", "g+rwX", dir)), "missing chmod -R for " + dir);
assertTrue(recorded.contains(List.of("find", dir, "-type", "d", "-exec", "chmod", "g+s", "{}", "+")),
"missing setgid find pass for " + dir);
}
// packed-refs does not exist in a freshly-init'd repo (only git gc / pack-refs creates it) —
// tolerated absence, so it must not appear at all: no recursive/-R treatment for a plain file.
String packedRefs = repoRoot + "/.git/packed-refs";
assertTrue(recorded.stream().noneMatch(c -> c.contains(packedRefs)),
"packed-refs is absent here and must be skipped, not chgrp'd: " + recorded);
}
/**
* A path that does not exist is skipped, never handed to {@code chgrp}. {@code .git/logs} is
* absent whenever {@code core.logAllRefUpdates} is false or no ref has been updated yet, and
* {@code chgrp} on a missing path exits non-zero — which would fail EVERY provisioning spawn
* with a message blaming a group that is in fact fine.
*/
@Test
void shareWithGroupSkipsPathsThatDoNotExist(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
String repoRoot = repo.toString();
deleteRecursively(repo.resolve(".git/logs"));
assertFalse(Files.exists(repo.resolve(".git/logs")), "fixture: .git/logs must be gone");
List<List<String>> recorded = new java.util.ArrayList<>();
java.util.function.Function<String[], String> recordingRunner = cmd -> {
recorded.add(joined(cmd));
return List.of(cmd).contains("--git-common-dir") ? ".git\n" : "";
};
GitWorktrees gitWorktrees =
new GitWorktrees(tmp.resolve("wts").toString(), "devteam", _ -> {}, recordingRunner);
gitWorktrees.shareWithGroup(repoRoot, repo.resolve("no-such-worktree").toString());
String logs = repoRoot + "/.git/logs";
assertTrue(recorded.stream().noneMatch(c -> c.contains(logs)),
"a missing .git/logs must be skipped, not chgrp'd: " + recorded);
assertTrue(recorded.stream().noneMatch(c -> c.contains(repo.resolve("no-such-worktree").toString())),
"a missing worktree path must be skipped too: " + recorded);
assertTrue(recorded.contains(List.of("chgrp", "-R", "devteam", repoRoot + "/.git/objects")),
"paths that DO exist are still shared: " + recorded);
}
/**
* The git store is located by {@code rev-parse --git-common-dir}, not by appending
* {@code /.git}. When git answers with an absolute path — what it does for a linked worktree,
* where {@code <repoRoot>/.git} is a file — every shared path must follow that answer.
*/
@Test
void shareWithGroupFollowsAnAbsoluteGitCommonDir(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
Path realGitDir = repo.resolve(".git");
List<List<String>> recorded = new java.util.ArrayList<>();
java.util.function.Function<String[], String> recordingRunner = cmd -> {
recorded.add(joined(cmd));
return List.of(cmd).contains("--git-common-dir") ? realGitDir + "\n" : "";
};
GitWorktrees gitWorktrees =
new GitWorktrees(tmp.resolve("wts").toString(), "devteam", _ -> {}, recordingRunner);
gitWorktrees.shareWithGroup(tmp.resolve("some/linked/worktree").toString(),
repo.resolve("wt").toString());
assertTrue(recorded.contains(List.of("chgrp", "-R", "devteam", realGitDir + "/objects")),
"objects must be taken from the reported common dir, not <repoRoot>/.git: " + recorded);
}
/** {@code packed-refs}, when present, is chgrp/chmod'd but never setgid'd (it is a file, not a dir). */
@Test
void shareWithGroupIncludesPackedRefsWhenPresent(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
String repoRoot = repo.toString();
Path packedRefsPath = repo.resolve(".git/packed-refs");
Files.writeString(packedRefsPath, "");
List<List<String>> recorded = new java.util.ArrayList<>();
java.util.function.Function<String[], String> recordingRunner = cmd -> {
recorded.add(joined(cmd));
return "";
};
GitWorktrees gitWorktrees =
new GitWorktrees(tmp.resolve("wts").toString(), "devteam", _ -> {}, recordingRunner);
gitWorktrees.shareWithGroup(repoRoot, repo.resolve("some-worktree").toString());
String packedRefs = packedRefsPath.toString();
assertTrue(recorded.contains(List.of("chgrp", "devteam", packedRefs)),
"packed-refs must be chgrp'd non-recursively when present: " + recorded);
assertTrue(recorded.contains(List.of("chmod", "g+rwX", packedRefs)),
"packed-refs must be chmod'd non-recursively when present: " + recorded);
assertTrue(recorded.stream().noneMatch(c -> c.contains("find") && c.contains(packedRefs)),
"packed-refs (a file) must never get the recursive setgid pass: " + recorded);
}
/** A group that does not exist (or that the operator is not a member of) fails loudly, naming it. */
@Test
void shareWithGroupThrowsNamingTheGroupWhenChgrpFails(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
GitWorktrees gitWorktrees = new GitWorktrees(tmp.resolve("wts").toString(), "cb185-nonexistent-group-zz");
String wt = new GitWorktrees(tmp.resolve("wts").toString()).add(repo.toString(), "cb-185-share", "HEAD");
WorktreeException e = assertThrows(WorktreeException.class,
() -> gitWorktrees.shareWithGroup(repo.toString(), wt));
assertTrue(e.getMessage().contains("cb185-nonexistent-group-zz"),
"exception must name the missing/refused group: " + e.getMessage());
}
}
@@ -148,6 +148,10 @@ class SessionManagerTest {
return 0;
}
@Override
public void shareWithGroup(String repoRoot, String worktreePath) {
}
List<String> removeCalls() {
return List.copyOf(removeCalls);
}
@@ -163,6 +163,37 @@ class WorktreeSessionManagerTest {
"tracked copied paths are --skip-worktree'd");
}
/**
* fleetd #185 stage 3, THE TRAP: {@code overlayParity} copies more files into the worktree
* AFTER {@code add} returns, so {@code shareWithGroup} must run after it, not folded into
* {@code add()} — otherwise every overlay file lands operator-owned and unwritable for a
* different-uid member, with a green test suite hiding it.
*/
@Test
void shareWithGroupRunsAfterOverlayParityNotBeforeIt() {
FakeHerdr herdr = new FakeHerdr();
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt")
.track(".envrc");
SessionManager sessions = new SessionManager(workerService(herdr), worktrees);
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-185", null));
assertEquals(1, worktrees.overlayCalls().size(), "overlayParity ran exactly once");
assertEquals(1, worktrees.shareCalls().size(), "shareWithGroup ran exactly once");
FakeWorktrees.OverlayCall overlay = worktrees.lastOverlay();
FakeWorktrees.ShareCall share = worktrees.lastShare();
assertEquals(s.worktree(), overlay.worktreePath());
assertEquals(s.worktree(), share.worktreePath());
List<String> order = worktrees.overlayShareOrder();
int overlayIndex = order.indexOf("overlay:" + s.worktree());
int shareIndex = order.indexOf("share:" + s.worktree());
assertTrue(overlayIndex >= 0 && shareIndex >= 0, "both calls must be recorded: " + order);
assertTrue(overlayIndex < shareIndex,
"shareWithGroup MUST run after overlayParity, not before/inside add(): " + order);
}
@Test
void releaseRemovesWorktreeButDoesNotDeleteBranch() {
FakeHerdr herdr = new FakeHerdr();