Compare commits

...

27 Commits

Author SHA1 Message Date
Dai Ha cea1183f75 CB-569: pass composed charter to Claude
CI / contract (pull_request) Successful in 43s
CI / build (pull_request) Successful in 50s
2026-08-15 05:09:30 +02:00
Dai Ha e2af4c5ae4 CB-567 follow-up: realign the changed constructor signatures
CI / contract (push) Successful in 44s
CI / build (push) Successful in 1m26s
The supplier swap left three parameter lists indented one column past
their siblings and one javadoc line broken mid-sentence. No behaviour
change.
2026-08-15 05:06:03 +02:00
Dai Ha dfd5f82894 Merge CB-567: read the charter per spawn, compose it once (PR #38)
HerdrPeerLauncher now takes Supplier<BridgedConfig.Fleet> instead of
Supplier<String> tabLabelTemplate, and reads it once per spawn. A field
taken at construction would have made the charter deferred, and deferred
looks exactly like working — which is why the test uses a mutable
supplier and spawns twice, rather than ConfigRef.fixed().

Composition happens once in the base, not in each adapter: two copies
drift while both adapter-local tests keep passing. Role charter first,
reply charter last, because the final instruction is the one that must
not be overridden. The reply charter stays gated on hasMcp() — telling a
peer to call a tool it was not given is a bug — while the role charter is
not, being identity rather than a tool instruction.

Both REPLY_CHARTER copies collapse into one, and it now says 'spawned
member' rather than 'off-subscription worker'. The old text made the
launch prompt contradict bridge_whoami for an architect; both architects
read it in their own prompts and reported it.

The old buildLaunch overloads are removed rather than kept as defaults: a
surviving one is the same shape as a stale snapshot, a route that drops
role and charter while looking healthy. That removal also let the
OpenCode adapter drop its ThreadLocal resume-id hack, since LaunchSpec
now carries the value down the same path.
2026-08-15 05:04:08 +02:00
Dai Ha e81944cef6 CB-567: compose charters per spawn
CI / build (pull_request) Successful in 57s
CI / contract (pull_request) Successful in 1m3s
2026-08-15 05:00:48 +02:00
Dai Ha 7bdd39ab9a CB-566 follow-up: say why the five-arg Fleet constructor is kept
CI / contract (push) Successful in 43s
CI / build (push) Successful in 1m8s
Reviewing the merge I read the five-arg constructor as dead code and
removed it. That was wrong: the tests call it as BridgedConfig.Fleet,
which my grep for 'new Fleet(' did not match, and the build failed on
eight call sites. It is restored with a javadoc that says why keeping it
is safe here even though an overload that drops a new field is normally
the shape to avoid — nothing reads a charter through a constructor, and
Jackson binds the canonical one, so it cannot swallow an operator's YAML.

Also drop a redundant java.util.Arrays qualifier (the class is already
imported) and rewrap a javadoc line the change had left over-long.
2026-08-15 04:54:09 +02:00
Dai Ha f4b38f040e Merge CB-566: the charter config surface (PR #37)
fleet.charters is a validated Map<String,String>, not a record: Fleet is
@JsonIgnoreProperties(ignoreUnknown = true), so a record field named
architetc would be dropped in silence and the operator would never learn
of the typo. A map lets validateCharters see the bad key and refuse it.

The key sits under fleet: because ConfigRef already treats that block as
hot and changedDeferredKeys does not list it. A new top-level key would
inherit nothing, and forgetting to classify it means a reload prints
'config reloaded' and does nothing.

A blank value is refused while an absent one is fine: an absent key means
the operator configured no charter, a blank one means they tried and
failed. Refusing at both startup and reload is the point — wiring only
one of the two paths is the whole bug.
2026-08-15 04:52:48 +02:00
Dai Ha 799668e129 CB-566: add fleet charter config
CI / contract (pull_request) Successful in 43s
CI / build (pull_request) Successful in 55s
2026-08-15 04:50:23 +02:00
Dai Ha 890190263e CB-560/562/563: the shipped block tells the truth again
The canonical CLAUDE.md block is the instruction surface this daemon ships to
every agent that mounts it, so a code change that silently invalidates it is an
incomplete change. CB-548 added a third principal kind and the block was never
revisited. Four statements in it were simply false:

  * whoami was documented as returning only primary or worker;
  * "spawn/stop/send/drain are lead-only" — Authz permits SEND to an architect;
  * delivery was documented as idle/blocked — injectable() is IDLE|BLOCKED|DONE,
    and a spawned member must also have mounted the bridge MCP, which is the
    exact condition that made every architect undeliverable for a day;
  * the tool table said bridge_list returns `workers` — the JSON key is
    `members`.

Also corrected: the fallback ladder claimed each one-way signal identifies a
"worker", but an architect gets the same charter, the same mount and the same
env, so those signals identify a spawned member and only bridge_whoami
separates the two. The safe default stays "act as a worker" — it is the most
restricted member role.

The turn contract now covers both member kinds, and says why the completion
fallback is not a substitute for bridge_reply: it returns at most the last 4000
characters, so a long report reaches the lead with its end cut off. That is not
hypothetical — it happened twice today.

Found by a reviewer asked whether the block still matches the code. Verified
against injectable(), Authz, BridgeMcp.listFleet and LeadLauncher before
applying. The wiki template is updated in the same shape and re-checked
byte-identical.
2026-08-15 04:37:03 +02:00
Dai Ha f8522edacd Merge CB-564: give every silent failure path a voice
CI / build (push) Successful in 1m0s
CI / contract (push) Successful in 1m15s
The bridge cannot classify what it never emits. Fleet-health monitoring —
detect a wedged member, decide, escalate — is blocked on that, so this is the
foundation rather than the feature.

A survey of Injector, CompletionResolver, StatusPoller, SessionManager,
SessionReaper and MessageService found six conditions that ended a member's
usefulness while saying nothing useful:

  * Injector.drop()             — worker gone, queue cleared: SILENT
  * CompletionResolver.fail()   — "via turn-stall fallback" at DEBUG, no reason
  * SessionManager.onFailed()   — "session marked failed" at DEBUG, no stage
  * SessionManager.reapIdle()   — indistinguishable from any other release
  * SessionManager.acquire()    — spawn failure rethrown with no log at all
  * MessageService.abandon()    — failed a caller's request at DEBUG

The first two are the exact phrases that misled the CB-560 diagnosis: both name
a symptom and neither names a cause. They are now WARN and carry the reason,
the stage, and the counts.

Conditions already loud were left alone, and so were two by-design timeouts in
MessageService — an async model exists precisely for those, and promoting them
would turn healthy operation into noise.

Behaviour is unchanged: every edit is a log statement.

Merge note: the recycle test deleted by CB-565 conflicted with a test added
here. Resolved by keeping the new onTurnFailed assertion and dropping the
recycle test, which tests a method that no longer exists.
2026-08-15 04:36:09 +02:00
Dai Ha b958747855 Merge CB-565: remove the recycle trap (PR #35)
CI / contract (push) Successful in 58s
CI / build (push) Successful in 1m27s
SessionManager.recycle called the 4-argument acquire overload, which defaults
the role to DEV and requests no worktree. So it carried profile, cwd and owner
across and silently dropped two things: the member's role, and its worktree.

A recycled architect would have come back a plain worker, never rebound to its
slot, with nothing logged. Worse, a recycled member would have come back with
no worktree at all — and a member's uncommitted work exists in exactly one
place. Role loss is recoverable; that is not.

Nothing called it. The only reference outside its own javadoc was one test, and
the context-cap path calls release, not recycle. So this was a trap waiting for
its first caller, and that caller would not have noticed either loss.

Deleted rather than repaired, on the CB-561 precedent: an API that looks correct
and silently drops a property is worse than no API. The no-reuse invariant it
documented is still true and is now stated directly.

Found by a reviewer asked to hunt the rest of the CB-548 fallout. The worktree
half was found while verifying the report.
2026-08-15 04:31:46 +02:00
Dai Ha 078bde2c02 CB-564: give a voice to silent member-failure paths
Injector.drop, CompletionResolver.fail, SessionManager.onFailed/reapIdle/
acquire spawn failures, and MessageService.abandon used to fail a member
or a caller's request with no log, a bare DEBUG, or a log that named only
the symptom ("session marked failed", "failed send via turn-stall
fallback"). Each now logs at WARN and names the real cause and the
numbers involved. Observability only — no behaviour changed.
2026-08-15 04:31:42 +02:00
Dai Ha 0b10ea987b CB-565: remove unsafe session recycle
CI / contract (pull_request) Successful in 44s
CI / build (pull_request) Successful in 55s
2026-08-15 04:29:50 +02:00
Dai Ha 1553d38182 Merge CB-563: a clipped pane scrape says it was clipped (PR #34)
CI / contract (push) Successful in 45s
CI / build (push) Successful in 1m12s
When a member ends its turn without bridge_reply, CompletionResolver scrapes
the pane and resolves the waiting send with that text. The scrape is capped at
MAX_SCRAPE_CHARS (4000), and nothing told the caller when the cap had bitten.
A delegating lead could act on a report missing its end and believe it was
complete. That happened to me today: a member's full engineering report arrived
cut at exactly 4000 characters, and the only hint was a DEBUG line reading
"(4000 chars scraped)", which reads like a size and not like a warning.

The returned text now carries a marker when, and only when, it was clipped, and
the clip is logged at WARN with the original length and the cap.

The cap itself is unchanged. The problem was silence, not the number.

The CB-115 misattribution guard still compares the unmarked clipped tail to the
unmarked baseline, and the marker is appended only afterwards. Verified in the
code, not taken on report: clip() strips before truncating and the new length
check uses the same stripped length, so there is no off-by-one either.
2026-08-15 04:26:57 +02:00
Dai Ha d04b075996 CB-563: mark clipped completion scrapes
CI / contract (pull_request) Successful in 42s
CI / build (pull_request) Successful in 56s
2026-08-15 04:25:54 +02:00
Dai Ha 644927636d Merge CB-562: the readiness gate says why it gave up (PR #33)
CI / build (push) Successful in 1m1s
CI / contract (push) Successful in 1m14s
When the injector's readiness grace expired it cleared the queue and logged
nothing. The failure then surfaced elsewhere as a turn-stall, which names the
wrong cause. Diagnosing CB-560 cost two live spawns and a wrong first
hypothesis for exactly this reason: the logs said "session marked failed" and
"failed send via turn-stall fallback", and neither says the message was never
typed into the pane at all.

The expiry now logs the target, the number of messages being failed, the grace
in polls and seconds, and the real cause in plain words.

The grace in seconds is derived, not written down twice: Bridged's own
INJECT_POLL_MILLIS is deleted and Injector.POLL_INTERVAL_MILLIS is the single
source, passed to every StatusPoller. A cadence change can no longer leave a
log line confidently stating the wrong duration.

Behaviour is unchanged. This is the first structured health event in the
daemon, and the foundation the fleet-health work will build on.
2026-08-15 04:21:29 +02:00
Dai Ha 95e45007aa CB-562: single source for the injector poll cadence; tighten count assertion
CI / build (pull_request) Successful in 59s
CI / contract (pull_request) Successful in 1m17s
2026-08-15 04:20:29 +02:00
Dai Ha 7a583c4045 CB-562: log why the readiness gate gave up on a target
CI / build (pull_request) Successful in 57s
CI / contract (pull_request) Successful in 1m4s
2026-08-14 21:52:25 +02:00
Dai Ha 65bce058bb Merge CB-561: one public way to build a CallerResolver (PR #32)
CI / build (push) Successful in 1m17s
CI / contract (push) Successful in 1m16s
CallerResolver had 5 public constructors and 4 public factories, and only one
of them could ever produce an architect. The rest defaulted memberSlotRoles to
`_ -> null`, so every architect quietly fell through to Principal.worker().
Nothing logged, nothing threw — the role was simply off.

That is the same failure shape as CB-560, so the fix is structural rather than
a warning: withLeadsAndMembers(.., MemberRegistry) is now the only public
construction path. Two overloads with no caller at all are deleted; the rest
are package-private and marked test-only. No path remains that accepts
architect bindings without a slot-role lookup, so no runtime WARN is needed.

Also checked and closed: the suspected slot leak on shutdown drain is not
real. SessionManager.release calls memberLifecycle.released() for every
removed session, and drainAll routes every session through release,
SPAWNING included. Verified in code.
2026-08-14 21:50:04 +02:00
Dai Ha 48b437083b Merge CB-560: mark spawned members present (PR #31)
Merging CB-548 made an architect resolve as Role.ARCHITECT, and BridgeMcp
marked MCP presence only for Role.WORKER. So an architect was never present.
That one line broke two things, because SessionManager.asPresence() is the
same object the injector's readiness gate reads:

  * the session never left SPAWNING;
  * Bridged.deliverableTo was false, so the injector held every delivery,
    waited out READINESS_GRACE_POLLS (~60s) and failed the send.

Confirmed live twice, on claude-code and on opencode. Both logged
"session marked failed" then "failed send via turn-stall fallback", neither
of which names the real cause.

Principal.isSpawnedMember() now means "a spawned member with a pane" —
WORKER or ARCHITECT, never PRIMARY. The Bridged.deliverableTo javadoc, which
still stated the worker-only rule as fact, is corrected: it is the clearest
description of this invariant anywhere, and leaving it stale is how the bug
comes back.
2026-08-14 21:50:04 +02:00
Dai Ha fa0612859b CB-560: document spawned member presence
CI / build (pull_request) Successful in 1m0s
CI / contract (pull_request) Successful in 1m4s
2026-08-14 21:48:58 +02:00
Dai Ha 4ffbcd0b7d CB-561: require member registry for architects
CI / contract (pull_request) Successful in 43s
CI / build (pull_request) Successful in 50s
2026-08-14 21:47:53 +02:00
Dai Ha a0cd053fd9 CB-560: mark architect members present
CI / build (pull_request) Successful in 55s
CI / contract (pull_request) Successful in 1m4s
2026-08-14 21:45:28 +02:00
Dai Ha f9a5e066b5 Merge PR #30: bind a spawned architect to its slot (CB-548's missing half)
CI / contract (push) Successful in 48s
CI / build (push) Successful in 1m18s
MemberRegistry.bind() had 18 call sites, every one in its own unit test.
Nothing in production ever bound a terminal to a slot, so CallerResolver's
architect branch was unreachable, every spawned architect resolved as WORKER,
and Authz's architect row was dead code. Since SEND is granted to
isPrimary() || isArchitect(), no architect could message anyone.

The spawn lifecycle now binds on both acquire paths and unbinds on release,
through a MemberLifecycle seam with a no-op default, so all six SessionManager
constructors are unchanged.

The escalation guard is deliberately double-sided. MemberRegistry flattens
every pool, so dev:* and reviewer:* slots sit in the same map the architect
resolver reads. The lifecycle binds only ARCHITECT roles, AND CallerResolver
independently checks roleForSlot(slot) == ARCHITECT before granting. Either
alone would be enough today; together, a regression in one cannot escalate a
worker. A reviewer confirmed a third barrier already existed: Authz gates
SPAWN on isPrimary(), so only a lead can request role=ARCHITECT at all.

Verified by the primary: mvn clean install, 640 tests, 0 failures.

Two findings recorded, neither blocking:

1. Known race, low severity. registry.put() and memberLifecycle.acquired()
   are not atomic. A release landing between them unbinds nothing (no binding
   exists yet), then acquire binds a dead terminal that no later release will
   ever clear - the slot leaks. Only the shutdown drain can reach a SPAWNING
   session (the idle reaper skips it, and an explicit stop needs a paneId the
   spawn has not returned), so the leak dies with the process. Note that
   binding before put does NOT fix it: the drain then iterates a roster the
   session is not in yet. A real fix needs atomicity.

2. The legacy Supplier-based withLeadsAndMembers overload and the Map-form
   constructors now silently disable architect resolution: memberSlotRoles
   defaults to _ -> null, so the architect branch can never be taken there.
   Production uses the registry form, and the tests were migrated, so nothing
   fails - but a future caller gets workers with no error.
2026-08-14 21:25:45 +02:00
Dai Ha e9bc192160 CB-548: bind architect sessions to member slots
CI / contract (pull_request) Successful in 45s
CI / build (pull_request) Successful in 55s
2026-08-14 21:10:48 +02:00
Dai Ha e0a57988ad Merge PR #29: drop .claude/settings.local.json from the default parityOverlay
CI / build (push) Successful in 1m21s
CI / contract (push) Successful in 1m21s
A member no longer receives the primary's pre-approved tool grants by default.
The grants were inert — GitWorktrees.isolateToolSurface already strips each
worktree's .mcp.json to an empty server map — so this is defence in depth: two
independent guards instead of one. Same reasoning as CB-525, which removed the
sibling .mcp.json and left this file behind.

An explicit parityOverlay: in config is unaffected; only the default changes.

Verified by the primary: mvn clean install, 637 tests, 0 failures.
2026-08-14 21:02:25 +02:00
Dai Ha b414a74c26 CB-559: drop .claude/settings.local.json from the default parity overlay
CI / build (pull_request) Successful in 54s
CI / contract (pull_request) Successful in 1m7s
2026-08-14 20:56:16 +02:00
Dai Ha 2c3796d598 Merge: keep every credential in one store
CI / build (push) Successful in 55s
CI / contract (push) Successful in 1m13s
2026-08-14 20:38:36 +02:00
31 changed files with 957 additions and 342 deletions
+39 -31
View File
@@ -11,41 +11,45 @@
If no `bridge_*` MCP tools are mounted in this session, this section does not apply — skip it.
`bridged` is the **sole communication gateway** between agents here. The orchestrating session (the
**primary**) and every delegated peer (a **worker**) mount the *same* MCP server and talk only
**primary**) and every delegated peer (a **member**) mount the *same* MCP server and talk only
through its `bridge_*` tools. No session addresses a peer, a broker, or the network directly.
### Which role am I? — settle this before acting
**Both roles read this file.** A worker runs in a git worktree of this same repo, so it inherits
**Every role reads this file.** A member runs in a git worktree of this same repo, so it inherits
this `CLAUDE.md` verbatim, and every rule below is role-conditional.
**Call `bridge_whoami`.** It returns `{"role":"primary"}` or `{"role":"worker","sessionId":…,
"profile":…,"worktree":…,"branch":…}`, resolved by the daemon from your connection — unforgeable,
and the same resolution its authorization gate uses. Don't infer what you can ask.
**Call `bridge_whoami`.** It returns `primary`, `worker`, or `architect`, resolved by the daemon from
your connection — unforgeable, and the same resolution its authorization gate uses. A worker also
carries its `sessionId`, `profile`, `worktree` and `branch`; an architect carries the slot name it
was bound to. Don't infer what you can ask.
Only if that call is unavailable, fall back to these — each is one-way, so keep reading until one
fires: the reply charter in your system prompt (*"You are an off-subscription worker in the
claude-bridge fleet"*) ⇒ **worker**; bridge tools prefixed `mcp__bridge__*` ⇒ **worker** (the
launcher fixes that mount name; a primary's mount is named by whoever wrote its `.mcp.json`, so it
varies); `ANTHROPIC_BASE_URL` set ⇒ **worker** (Claude-model workers run on a clean env, so its
*absence* proves nothing). **Still unsure ⇒ act as a worker.** The two mistakes are not symmetric: a
primary acting as a worker is refused by the authorization gate — loud and self-correcting — while a
worker acting as the primary ends its turn with no `bridge_reply`, and the sender silently receives
nothing. Fail toward the recoverable error.
claude-bridge fleet"*) ⇒ **spawned member**; bridge tools prefixed `mcp__bridge__*` ⇒ **spawned
member** (the launcher fixes that mount name; a primary's mount is named by whoever wrote its
`.mcp.json`, so it varies); `ANTHROPIC_BASE_URL` set ⇒ **spawned member** (Claude-model members run
on a clean env, so its *absence* proves nothing). None of these separate a worker from an architect —
only `bridge_whoami` does. **Still unsure ⇒ act as a worker**, the most restricted member role. The
two mistakes are not symmetric: a primary acting as a worker is refused by the authorization gate —
loud and self-correcting — while a member acting as the primary ends its turn with no `bridge_reply`,
and the sender silently receives nothing. Fail toward the recoverable error.
### Invariants — both roles, no exceptions
1. **Never set, export, or forward `ANTHROPIC_BASE_URL`** (or `ANTHROPIC_AUTH_TOKEN`). The primary
stays on subscription; only the bridge puts a worker off it, at spawn. Mounting the bridge must
stays on subscription; only the bridge puts a member off it, at spawn. Mounting the bridge must
never move a session across that boundary.
2. **The bridge is the only channel.** Text you print in your terminal reaches nobody — the other
side cannot see your screen. An answer that isn't in a `bridge_*` call is silently discarded.
3. **Identity comes from the connection, never an argument.** Workers never pass a target; you
cannot act as another session. Spawn/stop/send/drain are lead-only; reply/ask are
only-as-itself — any peer may answer for its own pane, and for no other. A call outside your
role is refused, not queued.
cannot act as another session. Spawn/stop/drain are lead-only; **send is lead or architect**;
reply/ask are only-as-itself — any peer may answer for its own pane, and for no other. A call
outside your role is refused, not queued.
4. **Delivery is status-gated: one message per turn.** Don't busy-poll a peer's terminal and don't
re-send because a call looks slow — the bridge delivers when the peer is `idle`/`blocked`.
re-send because a call looks slow — the bridge delivers when the peer is `idle`, `blocked` or
`done`. A spawned member must **also** have mounted the bridge MCP: until it has, it is not
deliverable, and a send waits on that gate for ~60s and then fails without ever reaching its pane.
5. **Never drive the terminal multiplexer directly** (no `herdr` CLI, no socket). The bridge owns
policy; the multiplexer owns PTYs. Going around the bridge bypasses every rule above.
@@ -101,26 +105,26 @@ the merge — and merging on a reviewer's word is delegating it by proxy.
|---|---|
| Confirm your own role | `bridge_whoami` |
| See backends available | `bridge_profiles` |
| Start a worker | `bridge_spawn{profile?, cwd?, worktree?, ticket?}` → `sessionId` + `paneId` |
| See the fleet | `bridge_list` → `leads` (your peers) + `workers` · one peer's state: `bridge_status{sessionId}` |
| Start a member | `bridge_spawn{role?, profile?, cwd?, worktree?, ticket?}` → `sessionId` + `paneId` |
| See the fleet | `bridge_list` → `leads` (your peers) + `members` · one peer's state: `bridge_status{sessionId}` |
| Delegate (blocking) | `bridge_send{sessionId, content}` |
| Delegate (long task) | `bridge_send{sessionId, content, wait:false}` → ticket → `bridge_poll{ticket}` |
| Answer a worker's `bridge_ask` | `bridge_send{turnId, content}` — **not** `sessionId` |
| Answer a member's `bridge_ask` | `bridge_send{turnId, content}` — **not** `sessionId` |
| Message a **peer lead** | `bridge_send{sessionId: <their terminal>, content}` — `bridge_list` → `leads` reports it. Coordination only, **never** a task |
| Answer a peer lead that messaged you | `bridge_reply{content}` — the one case a lead replies |
| Collect a held reply | `bridge_poll{target}` · then `bridge_ack{target, msgId}` |
| Tear down | `bridge_stop{paneId}` |
| Tear down a member | `bridge_stop{paneId}` |
### Lead ↔ lead — coordinate, never delegate
`bridge_list` returns `leads` alongside `workers`; your own row carries `self: true`. Every other row
is a peer — an orchestrator with its own context, its own workers, and its own judgment. An empty
`workers` array means no workers are spawned; it says nothing about peers.
`bridge_list` returns `leads` alongside `members`; your own row carries `self: true`. Every other row
is a peer — an orchestrator with its own context, its own members, and its own judgment. An empty
`members` array means no members are spawned; it says nothing about peers.
**A lead never assigns a task to another lead.** Work goes to workers — only ever downward, never
sideways. Sending a peer a brief with acceptance criteria is a category error: a brief is a worker's
**A lead never assigns a task to another lead.** Work goes to members — only ever downward, never
sideways. Sending a peer a brief with acceptance criteria is a category error: a brief is a member's
artefact, and a peer is not yours to task. If a unit needs doing and it falls in your area, spawn a
worker and delegate it yourself; if it falls in the peer's area, say so and let the peer assign it.
member and delegate it yourself; if it falls in the peer's area, say so and let the peer assign it.
The traffic between leads is coordination and nothing else:
1. **Divide the map, not the work.** Agree who owns which area, then each of you assigns inside your
@@ -138,7 +142,7 @@ Being messaged by a peer does not make you its worker: answer with `bridge_reply
the substance if it is wrong. A peer that simply complies has thrown away the reason there are two of
you.
### Worker — the turn contract
### Member (worker or architect) — the turn contract
1. **Load the playbook skill the lead named** before doing anything else.
2. **Do the assigned scope only.** Note anything you spot outside it in one line; don't go hunt it.
@@ -147,6 +151,10 @@ you.
Don't ask what you could decide yourself.
4. **End the turn with exactly one `bridge_reply{content}`**, carrying your complete answer. This is
the whole handoff. No `bridge_reply` ⇒ the sender gets nothing and the exchange stalls.
Do **not** lean on the completion fallback to carry your answer for you: when you end a turn
without replying, the bridge scrapes your pane, and it can return only the last 4000 characters.
A clipped scrape is marked as partial, but the missing text is gone — your report reaches the
lead with its end cut off.
5. **Report honestly.** State only what you actually ran and its real output, including failures.
You mount **only** the bridge MCP — the primary's other servers (IDE, forge, docs) are not yours,
so never claim the result of a check you had no way to run.
@@ -157,9 +165,9 @@ you.
| Layer | Scope | Reaches |
|---|---|---|
| the launcher's reply charter | the one rule that must survive with no repo: *end every turn with `bridge_reply`* | every worker, at launch, every peer kind |
| **this section** | protocol + orchestration policy | primary **and** every Claude worker — tracked in git, so worktrees inherit it |
| role playbook skills | per-job procedure (commit/PR recipe, finding format) | a worker told to load one |
| the launcher's reply charter | the one rule that must survive with no repo: *end every turn with `bridge_reply`* | every spawned member, at launch, every peer kind — never a lead |
| **this section** | protocol + orchestration policy | primary **and** every member that reads the repo — tracked in git, so worktrees inherit it |
| role playbook skills | per-job procedure (commit/PR recipe, finding format) | a member told to load one |
| the bridge's own docs | design detail, flows, error model | on demand |
A rule belongs in **exactly one** layer — the outermost one that must obey it. Peers that don't read
+20 -2
View File
@@ -231,8 +231,8 @@ placement: weighted
#
# Not every key can move under a running daemon, and the difference is about what already exists
# when the reload happens — not about how important the key is:
# HOT → takes effect on the next spawn: the whole `fleet:` block (every role pool and
# `tabLabel`), `placement:`, and an existing profile's weight / maxLoad. Those are
# HOT → takes effect on the next spawn: the whole `fleet:` block (every role pool,
# `charters`, and `tabLabel`), `placement:`, and an existing profile's weight / maxLoad. Those are
# hot because the placement policy reads them through a supplier — being config is
# not by itself enough to make a key hot.
# DEFERRED → accepted into the new config, but the wiring built at startup keeps the old value
@@ -274,6 +274,24 @@ placement: weighted
# the candidates, in definition order. A dev and a reviewer staying anonymous is exactly compatible
# with being listed here; the entry key just names the entry.
fleet:
# Optional launch-charter text, keyed only by the singular role wire names: architect, dev,
# reviewer. Changes are HOT and reach the next spawn without a daemon restart. Do not put secrets
# here: a later launch step writes this text to a world-readable temp file, and ${ENV} interpolation
# is deliberately not supported.
charters:
architect: |-
You are an architect in this fleet. You refine work before anyone builds it:
scope, acceptance criteria, risks, and a unit split. You read the repo and
write analysis. You never commit production code and never open a PR.
A design task is worked by two architects. Design alone first, then exchange
and say plainly where you disagree. Do not concede just to agree.
dev: |-
You implement the one unit you were given, and nothing else. You test it,
commit it, and open your own pull request. You never merge.
reviewer: |-
You review the diff you were given. You report bugs, risks and missing tests.
You do not change code.
# Optional. Template for a member tab's label; {role}, {profile}, {model} and {n} are substituted.
# {n} counts per role+profile, so `dev: sonnet #2` really is the second sonnet dev. Because {role}
# comes from a closed enum, a generated label can never begin with a lead's tabPrefix.
@@ -70,9 +70,6 @@ public final class Bridged {
private static final Logger log = LoggerFactory.getLogger(Bridged.class);
/** How often the injector samples a busy worker's status while it has queued work. */
private static final long INJECT_POLL_MILLIS = 250;
/** CB-504: how long to wait at startup for herdr's socket before serving degraded. */
private static final long HERDR_WAIT_SECONDS = 30;
private static final long HERDR_WAIT_POLL_MILLIS = 500;
@@ -99,6 +96,7 @@ public final class Bridged {
// CB-542: a subscription:true profile whose env: reseats ANTHROPIC_BASE_URL/AUTH_TOKEN would
// reach an unguarded endpoint (the launcher skips SubscriptionGuard for it). Refuse at load.
cfg.validateSubscriptionProfiles();
cfg.validateCharters();
// CB-548: every architect slot must name a configured workers: profile — the strong-model
// backend the future spawn lifecycle would read. A stale reference dies here, not later.
cfg.validateMembers();
@@ -131,13 +129,13 @@ public final class Bridged {
adapters.add(new ClaudeCodeLauncher(agents, spaces, guard,
claudeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
() -> config.get().fleet().tabLabel()));
() -> config.get().fleet()));
}
if (!opencodeProfiles.isEmpty()) {
adapters.add(new OpenCodeLauncher(agents, spaces,
opencodeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
() -> config.get().fleet().tabLabel()));
() -> config.get().fleet()));
}
AtomicReference<Function<String, Integer>> liveCountRef = new AtomicReference<>(_ -> 0);
PeerLauncher workers = new CompositePeerLauncher(
@@ -243,6 +241,7 @@ public final class Bridged {
// what CallerResolver resolves against and what that lifecycle will read profiles from;
// nothing here spawns a slot.
MemberRegistry members = new MemberRegistry(cfg.fleet());
sessions.setMemberLifecycle(members);
if (!members.slots().isEmpty()) {
log.info("member slots: {} configured {} — none bound yet (a slot is idle until the "
+ "spawn lifecycle binds a live terminal to it)",
@@ -289,7 +288,7 @@ public final class Bridged {
};
Injector injector = new Injector(agents, turnListener, deliverableTo(presence, leads),
presence::forget);
StatusPoller poller = new StatusPoller(agents, injector, INJECT_POLL_MILLIS);
StatusPoller poller = new StatusPoller(agents, injector, Injector.POLL_INTERVAL_MILLIS);
poller.start();
// CB-307: reply inbox. A broker: block (with a uri) selects the AMQP-backed durable adapter;
@@ -375,13 +374,11 @@ public final class Bridged {
throw new IllegalStateException("auth.mode=token but env var " + cfg.auth().tokenEnv()
+ " is unset or empty — export it before starting bridged");
}
callers = CallerResolver.withLeadsAndMembers(identity, true, token, leads,
members::snapshot);
callers = CallerResolver.withLeadsAndMembers(identity, true, token, leads, members);
log.info("auth: token mode (bearer required for non-worker callers, env {})",
cfg.auth().tokenEnv());
} else {
callers = CallerResolver.withLeadsAndMembers(identity, false, null, leads,
members::snapshot);
callers = CallerResolver.withLeadsAndMembers(identity, false, null, leads, members);
log.info("auth: loopback-trust (any loopback non-worker caller is the primary)");
}
@@ -429,19 +426,20 @@ public final class Bridged {
}
/**
* The {@link Injector}'s readiness gate (CB-534): a target is deliverable if it is a worker whose
* agent has connected the bridge MCP, <em>or</em> a lead.
* The {@link Injector}'s readiness gate (CB-534): a target is deliverable if it is a spawned
* member whose agent has connected the bridge MCP, <em>or</em> a lead.
*
* <p>The gate exists for one reason — to hold a delivery out of a <em>spawned</em> worker's boot
* <p>The gate exists for one reason — to hold a delivery out of a <em>spawned</em> member's boot
* window, where herdr already reports {@code idle} but the TUI would drop an injected paste. That
* hazard is a property of spawning. A lead is never spawned: the operator started it and named it
* (or labelled its tab) only once it was up, so there is no boot window to guard.
*
* <p>A lead is also never enrolled in {@link MemberPresence} — {@code BridgeMcp} marks presence
* only for a worker, deliberately, since that map doubles as the worker roster's availability
* signal and a lead counted there would show up as an available worker. So without the second
* disjunct a lead is permanently un-deliverable: every lead→lead send sat on the gate for
* {@code READINESS_GRACE_POLLS} (~60s) and then failed having never been typed into the pane.
* for every spawned member (worker and architect), deliberately, since that map doubles as the
* member roster's availability signal and a lead counted there would show up as an available
* member. So without the second disjunct a lead is permanently un-deliverable: every
* lead→lead send sat on the gate for {@code READINESS_GRACE_POLLS} (~60s) and then failed
* having never been typed into the pane.
*
* <p>The lead set is read through the supplier on each call rather than snapshotted, so a lead
* discovered by {@code leadScan} after startup becomes deliverable without a restart.
@@ -1,10 +1,12 @@
package dev.ltms.bridged.auth;
import dev.ltms.bridged.mcp.ConnectionIdentity;
import dev.ltms.bridged.peer.MemberRole;
import java.nio.charset.StandardCharsets;
import java.security.MessageDigest;
import java.util.Map;
import java.util.function.Function;
import java.util.function.Supplier;
/**
@@ -60,14 +62,16 @@ public final class CallerResolver {
* in {@code Bridged} reads a constant from config, which is the degenerate live case.
*/
private final Supplier<Map<String, String>> architectTerminals;
private final Function<String, MemberRole> memberSlotRoles;
private final Function<String, String> memberSlotNames;
/** Loopback-trust resolver: no token required, historical behaviour. */
public CallerResolver(ConnectionIdentity identity) {
/** Loopback-trust resolver: no token required, historical behaviour. Test-only. */
CallerResolver(ConnectionIdentity identity) {
this(identity, false, null, Map.of());
}
/** As {@link #CallerResolver(ConnectionIdentity, boolean, String, Map)} with no leads pinned. */
public CallerResolver(ConnectionIdentity identity, boolean tokenMode, String token) {
/** As {@link #CallerResolver(ConnectionIdentity, boolean, String, Map)} with no leads pinned. Test-only. */
CallerResolver(ConnectionIdentity identity, boolean tokenMode, String token) {
this(identity, tokenMode, token, Map.of());
}
@@ -82,8 +86,8 @@ public final class CallerResolver {
* @param pinnedPrimaryTerminal the primary's own herdr {@code terminal_id}
* ({@code null}/blank = unpinned)
*/
public static CallerResolver pinnedTo(ConnectionIdentity identity, boolean tokenMode,
String token, String pinnedPrimaryTerminal) {
static CallerResolver pinnedTo(ConnectionIdentity identity, boolean tokenMode,
String token, String pinnedPrimaryTerminal) {
return new CallerResolver(identity, tokenMode, token,
pinnedPrimaryTerminal == null || pinnedPrimaryTerminal.isBlank()
? Map.of() : Map.of(pinnedPrimaryTerminal, "primary"));
@@ -98,21 +102,11 @@ public final class CallerResolver {
* {@link Role#PRIMARY} — rather than a worker. Empty = nothing pinned,
* so every pane resolves as a worker.
*/
public CallerResolver(ConnectionIdentity identity, boolean tokenMode, String token,
Map<String, String> leadTerminals) {
CallerResolver(ConnectionIdentity identity, boolean tokenMode, String token,
Map<String, String> leadTerminals) {
this(identity, tokenMode, token, fixed(leadTerminals), null);
}
/**
* Map-form of both registries (CB-548): lead terminals and the initial architect terminal
* bindings, each snapshotted at construction (a handed-over map is not offered as live state).
*/
public CallerResolver(ConnectionIdentity identity, boolean tokenMode, String token,
Map<String, String> leadTerminals,
Map<String, String> architectTerminals) {
this(identity, tokenMode, token, fixed(leadTerminals), fixed(architectTerminals));
}
/**
* Live-registry form: {@code leadTerminals} is consulted on every resolve, so leads discovered
* after startup (CB-531's tab scan) take effect without a restart.
@@ -121,26 +115,26 @@ public final class CallerResolver {
* {@link #pinnedTo}: {@code Map} and {@code Supplier} overloads are ambiguous for a literal
* {@code null}.
*/
public static CallerResolver withLeads(ConnectionIdentity identity, boolean tokenMode,
String token,
Supplier<Map<String, String>> leadTerminals) {
static CallerResolver withLeads(ConnectionIdentity identity, boolean tokenMode,
String token,
Supplier<Map<String, String>> leadTerminals) {
return new CallerResolver(identity, tokenMode, token, leadTerminals, null);
}
/**
* Live-registry form for both {@code leadTerminals} and the CB-548 architect registry: both
* are consulted on every resolve, so a slot binding injected after startup takes effect
* without a restart.
* Live registry form that can confirm a bound slot is an architect slot.
*
* <p>A static factory rather than a constructor overload, for the same reason as
* {@link #pinnedTo}: too many {@code Map}/{@code Supplier} combinations to make {@code null}
* unambiguous.
* <p>This is the only public construction path. It keeps terminal bindings and slot roles in
* the same {@link MemberRegistry}, so a configured architect can resolve as an architect.
*/
public static CallerResolver withLeadsAndMembers(ConnectionIdentity identity,
boolean tokenMode, String token,
Supplier<Map<String, String>> leadTerminals,
Supplier<Map<String, String>> architectTerminals) {
return new CallerResolver(identity, tokenMode, token, leadTerminals, architectTerminals);
boolean tokenMode, String token,
Supplier<Map<String, String>> leadTerminals,
MemberRegistry members) {
return new CallerResolver(identity, tokenMode, token, leadTerminals,
members == null ? null : members::snapshot,
members == null ? null : members::roleForSlot,
members == null ? null : members::nameForSlot);
}
private static Supplier<Map<String, String>> fixed(Map<String, String> leadTerminals) {
@@ -151,6 +145,21 @@ public final class CallerResolver {
private CallerResolver(ConnectionIdentity identity, boolean tokenMode, String token,
Supplier<Map<String, String>> leadTerminals,
Supplier<Map<String, String>> architectTerminals) {
this(identity, tokenMode, token, leadTerminals, architectTerminals, null);
}
private CallerResolver(ConnectionIdentity identity, boolean tokenMode, String token,
Supplier<Map<String, String>> leadTerminals,
Supplier<Map<String, String>> architectTerminals,
Function<String, MemberRole> memberSlotRoles) {
this(identity, tokenMode, token, leadTerminals, architectTerminals, memberSlotRoles, null);
}
private CallerResolver(ConnectionIdentity identity, boolean tokenMode, String token,
Supplier<Map<String, String>> leadTerminals,
Supplier<Map<String, String>> architectTerminals,
Function<String, MemberRole> memberSlotRoles,
Function<String, String> memberSlotNames) {
if (tokenMode && (token == null || token.isBlank())) {
throw new IllegalArgumentException(
"auth.mode=token requires a non-empty token; check that the env var named by "
@@ -161,6 +170,8 @@ public final class CallerResolver {
this.expectedToken = tokenMode ? token.getBytes(StandardCharsets.UTF_8) : null;
this.leadTerminals = leadTerminals == null ? Map::of : leadTerminals;
this.architectTerminals = architectTerminals == null ? Map::of : architectTerminals;
this.memberSlotRoles = memberSlotRoles == null ? _ -> null : memberSlotRoles;
this.memberSlotNames = memberSlotNames == null ? Function.identity() : memberSlotNames;
}
/**
@@ -206,11 +217,12 @@ public final class CallerResolver {
return Principal.leader(lead, c.terminal(), c.pid());
}
String slot = architectTerminals.get().get(c.terminal());
if (slot != null) {
if (slot != null && memberSlotRoles.apply(slot) == MemberRole.ARCHITECT) {
// The config/live binding names this pane as an architect slot's own. Same
// unforgeable pane mapping; the live binding, never a request argument, decides.
// Checked before the generic worker fallback, per the CB-548 precedence order.
return Principal.architect(slot, c.terminal(), c.pid());
// Check the slot role too: this defence in depth prevents a bad lifecycle bind from
// escalating a dev or reviewer into an architect. Checked before the worker fallback.
return Principal.architect(memberSlotNames.apply(slot), c.terminal(), c.pid());
}
return Principal.worker(c.terminal(), c.pid()); // unforgeable; never token-gated
}
@@ -0,0 +1,21 @@
package dev.ltms.bridged.auth;
import dev.ltms.bridged.peer.MemberRole;
/** Optional session lifecycle hook for live member-slot bindings. */
public interface MemberLifecycle {
MemberLifecycle NONE = new MemberLifecycle() {
@Override
public void acquired(MemberRole role, String profile, String terminal) {
}
@Override
public void released(String terminal) {
}
};
void acquired(MemberRole role, String profile, String terminal);
void released(String terminal);
}
@@ -2,11 +2,14 @@ package dev.ltms.bridged.auth;
import dev.ltms.bridged.config.BridgedConfig;
import dev.ltms.bridged.peer.MemberRole;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.Collections;
import java.util.HashMap;
import java.util.LinkedHashMap;
import java.util.Map;
import java.util.Objects;
/**
* The architect-slot registry (CB-548): every gateway-local architect name and the strong-model
@@ -29,7 +32,9 @@ import java.util.Map;
* exposes the map the resolver resolves against plus the profile lookup lifecycle will call.
* Nothing here creates or manages an architect session.
*/
public final class MemberRegistry {
public final class MemberRegistry implements MemberLifecycle {
private static final Logger log = LoggerFactory.getLogger(MemberRegistry.class);
/**
* One flattened {@code fleet:} entry.
@@ -125,6 +130,12 @@ public final class MemberRegistry {
return e == null ? null : e.role();
}
/** The unqualified configured name for a slot, or {@code null} if it is unknown. */
public String nameForSlot(String slotName) {
Entry e = slots.get(slotName);
return e == null ? null : e.name();
}
/** True when {@code slotName} is a configured architect slot. */
public boolean isSlot(String slotName) {
return slots.containsKey(slotName);
@@ -189,4 +200,33 @@ public final class MemberRegistry {
return true;
}
}
/**
* Bind only architect sessions to a free slot with the resolved profile.
*
* <p>The role check is lifecycle policy. {@link CallerResolver} repeats it when resolving a
* binding, so a later lifecycle regression cannot turn a worker into an architect.
*/
@Override
public void acquired(MemberRole role, String profile, String terminal) {
if (role != MemberRole.ARCHITECT || terminal == null || terminal.isBlank()) {
return;
}
// slotsFor preserves definition order, so duplicate-profile slots use the first free one.
for (Entry entry : slotsFor(MemberRole.ARCHITECT).values()) {
if (Objects.equals(profile, entry.profile()) && bind(entry.key(), terminal)) {
return;
}
}
log.info("member slot: no free architect slot for profile={}; session remains a worker", profile);
}
/** Unbind a released terminal using the compare-safe registry operation. */
@Override
public void released(String terminal) {
String slot = slotForTerminal(terminal);
if (slot != null) {
unbind(slot, terminal);
}
}
}
@@ -48,7 +48,7 @@ public record Principal(Role role, String terminal, long pid, String name) {
* made a lead unaddressable: {@link #ownsSession} could never be true for it, so
* {@code bridge_reply} was refused and one lead could send to another but never be answered.
* The terminal now means "which pane is this caller", the presence map keys on
* {@link #isWorker()} instead, and a lead is a peer that can both send and receive.
* {@link #isSpawnedMember()} instead, and a lead is a peer that can both send and receive.
*/
public static Principal leader(String name, String terminal, long pid) {
return new Principal(Role.PRIMARY, terminal, pid, name);
@@ -85,6 +85,16 @@ public record Principal(Role role, String terminal, long pid, String name) {
return role == Role.WORKER;
}
/**
* Whether this caller is a spawned member with its own pane.
*
* <p>Both workers and architects are spawned members. A lead is excluded because recording it
* as present would count it as an available member in the roster.
*/
public boolean isSpawnedMember() {
return role == Role.WORKER || role == Role.ARCHITECT;
}
public boolean isAnonymous() {
return role == Role.ANONYMOUS;
}
@@ -210,8 +210,14 @@ public record BridgedConfig(
// worker the primary's IDE servers, which are bound to the primary's checkout — so its
// navigation returned paths outside its own worktree. GitWorktrees now neutralizes that
// file instead; a worker's tools are whatever its launcher mounts.
// .claude/settings.local.json is the sibling that was left behind: it pre-approves tools
// (mcp__context7__*, mcp__jetbrains, mcp__intellij-index, Workflow(code-review)) a worker
// must never hold, and enables MCP servers by name. Its grants are currently INERT because
// GitWorktrees.isolateToolSurface strips every worktree's server map to empty — the named
// servers do not exist there to be enabled. This is defence in depth, not a live fix: the
// worker stays isolated only because this separate mechanism already removes the servers.
parityOverlay = (parityOverlay == null || parityOverlay.isEmpty())
? List.of(".claude/settings.local.json", ".env", ".envrc")
? List.of(".env", ".envrc")
: List.copyOf(parityOverlay);
// gitTokenEnv stays null when unset (opt-in). gitHostEnv defaults so operators enabling
// checkpoints need only set gitTokenEnv; it is injected only alongside a resolved token.
@@ -526,6 +532,7 @@ public record BridgedConfig(
* @param architects profiles the {@code architect} role may run on
* @param developers profiles the {@code dev} role may run on
* @param reviewers profiles the {@code reviewer} role may run on
* @param charters optional launch-charter text keyed by singular role wire name
* @param tabLabel template for a member tab's label; {@code {role}}, {@code {profile}},
* {@code {model}} and {@code {n}} (a per role+profile counter) are
* substituted. Default {@link #DEFAULT_TAB_LABEL}
@@ -535,6 +542,7 @@ public record BridgedConfig(
Map<String, Slot> architects,
Map<String, Slot> developers,
Map<String, Slot> reviewers,
Map<String, String> charters,
String tabLabel) {
/**
@@ -550,9 +558,24 @@ public record BridgedConfig(
architects = unmodifiableOrEmpty(architects);
developers = unmodifiableOrEmpty(developers);
reviewers = unmodifiableOrEmpty(reviewers);
charters = unmodifiableOrEmpty(charters);
tabLabel = (tabLabel == null || tabLabel.isBlank()) ? DEFAULT_TAB_LABEL : tabLabel;
}
/**
* A fleet with no configured launch charters — the shape every deployment had before
* CB-566, and what most tests want.
*
* <p>Kept deliberately, even though an overload that drops a new field is normally the
* shape to avoid. It is safe here because nothing <em>reads</em> a charter through a
* constructor: the launcher reads {@code fleet.charters()} from the live config. Jackson
* binds the canonical constructor, so this one cannot swallow an operator's YAML.
*/
public Fleet(Map<String, Leader> leaders, Map<String, Slot> architects,
Map<String, Slot> developers, Map<String, Slot> reviewers, String tabLabel) {
this(leaders, architects, developers, reviewers, null, tabLabel);
}
/**
* Deliberately not {@code Map.copyOf}: its iteration order is salted per JVM run, which
* would discard YAML definition order. The {@code fixed} placement policy answers with a
@@ -576,6 +599,11 @@ public record BridgedConfig(
};
}
/** The configured launch charter for {@code role}, or {@code null} when it is absent. */
public String charterFor(MemberRole role) {
return role == null ? null : charters.get(role.wireName());
}
/**
* The profile names {@code role} may run on, in definition order, without repeats.
*
@@ -1046,7 +1074,7 @@ public record BridgedConfig(
// fleet IS defaulted, unlike the leadScan: block it replaced, because an empty Fleet is not
// the same as an enabled one: every pool is empty, so no lead is scanned for or created and
// no role has a pool. Constructing it saves every reader a null check for no behaviour change.
Fleet f = (fleet != null) ? fleet : new Fleet(null, null, null, null, null);
Fleet f = (fleet != null) ? fleet : new Fleet(null, null, null, null, null, null);
// leadHeartbeat is left as-is (CB-551): null is "off", and LeadHeartbeat's own compact
// constructor defaults the fields of a block that IS present. Defaulting it here would
// switch the feature on for every config that never mentioned it.
@@ -1184,6 +1212,35 @@ public record BridgedConfig(
}
}
/**
* Reject configured charter entries that would remove a role's contract or never be read.
*
* <p>The map deliberately retains every key from {@code fleet.charters:}. A typed record would
* silently discard an unknown child because {@link Fleet} ignores unknown JSON properties, which
* would make a typo look like an accepted configuration.
*
* @throws IllegalStateException when a charter key is not a role wire name or its value is blank
*/
public void validateCharters() {
if (fleet == null || fleet.charters().isEmpty()) {
return;
}
List<String> valid = Arrays.stream(MemberRole.values())
.map(MemberRole::wireName)
.toList();
List<String> bad = new ArrayList<>();
fleet.charters().forEach((key, charter) -> {
if (!valid.contains(key)) {
bad.add("fleet.charters." + key + " is not a role wire name (valid: " + valid + ").");
} else if (charter == null || charter.isBlank()) {
bad.add("fleet.charters." + key + " is blank; a configured role needs charter text.");
}
});
if (!bad.isEmpty()) {
throw new IllegalStateException("refusing to start: " + String.join(" ", bad));
}
}
/**
* Reject a member slot whose {@code role} or {@code profile} does not resolve.
*
@@ -27,10 +27,10 @@ import java.util.function.Supplier;
*
* <ul>
* <li><strong>Hot</strong> — re-read per use, so a reload takes effect on the next spawn:
* {@code fleet:} (every role pool and {@code tabLabel}), {@code placement:}, and an existing
* profile's {@code weight} / {@code maxLoad}. Those three are read through a supplier on
* {@code CompositePeerLauncher}, which is what makes them hot — not the fact that they are
* config.</li>
* {@code fleet:} (every role pool, {@code charters}, and {@code tabLabel}),
* {@code placement:}, and an existing profile's {@code weight} / {@code maxLoad}. Those
* three are read through a supplier on {@code CompositePeerLauncher}, which is what makes
* them hot — not the fact that they are config.</li>
* <li><strong>Deferred</strong> — accepted into the new snapshot, but the wiring built at startup
* keeps the old value until a restart: {@code lifecycle:}, {@code leadHeartbeat:},
* {@code spawnReadyTimeoutMs} / {@code spawnReadyPollMs}, {@code guard:},
@@ -150,6 +150,7 @@ public final class ConfigRef implements Supplier<BridgedConfig> {
fresh.validateAuthExposure();
fresh.validateLeadTabPrefixes();
fresh.validateSubscriptionProfiles();
fresh.validateCharters();
fresh.validateMembers();
} catch (RuntimeException e) {
String msg = e.getMessage() == null ? e.toString() : e.getMessage();
@@ -55,6 +55,9 @@ public final class CompletionResolver implements TurnListener {
/** Cap the scraped tail so a long transcript can't return an unbounded blob. */
static final int MAX_SCRAPE_CHARS = 4000;
private static final String CLIPPED_PANE_TAIL_MARKER =
"[Pane tail clipped: member did not call bridge_reply.]";
private final AgentControl agents;
private final Rendezvous rendezvous;
@@ -147,9 +150,14 @@ public final class CompletionResolver implements TurnListener {
return;
}
String tail;
int originalLength = 0;
boolean clipped = false;
boolean scrapeFailed = false;
try {
tail = clip(lastAssistantBlock(agents.read(target, SCRAPE_SOURCE)));
String assistantBlock = lastAssistantBlock(agents.read(target, SCRAPE_SOURCE));
originalLength = assistantBlock.strip().length();
clipped = originalLength > MAX_SCRAPE_CHARS;
tail = clip(assistantBlock);
} catch (RuntimeException e) {
// The worker finished but we couldn't read its screen — still resolve the send so the
// caller unblocks; an empty tail beats hanging until the caller's timeout.
@@ -169,8 +177,14 @@ public final class CompletionResolver implements TurnListener {
target);
return; // keep the in-flight record: a later genuine completion still needs it
}
if (rendezvous.resolveCompletion(waiter, tail)) {
String completion = clipped ? tail + "\n" + CLIPPED_PANE_TAIL_MARKER : tail;
if (rendezvous.resolveCompletion(waiter, completion)) {
inFlight.remove(target, turn);
if (clipped) {
log.warn("completion scrape for {} clipped from {} chars to the {} char cap; "
+ "member did not call bridge_reply, so the pane tail is partial",
target, originalLength, MAX_SCRAPE_CHARS);
}
log.debug("resolved send to {} via turn-completion fallback ({} chars scraped)",
target, tail.length());
}
@@ -200,7 +214,7 @@ public final class CompletionResolver implements TurnListener {
}
if (rendezvous.resolveFailure(waiter, reason)) {
inFlight.remove(target, turn);
log.debug("failed send to {} via turn-stall fallback", target);
log.warn("failing send to {} via turn-stall fallback: {}", target, reason);
}
}
@@ -78,6 +78,15 @@ public final class Injector {
*/
private static final int READINESS_GRACE_POLLS = 240;
/**
* The single source for the injector poll cadence — how often the {@link StatusPoller} drives
* {@link #onStatus} at. {@code Bridged} passes this to every {@link StatusPoller} it constructs,
* and this class reads it to state the readiness grace in seconds on the CB-562 expiry log
* instead of hardcoding "60s". One constant, so a cadence change cannot silently desync a log
* that claims a grace duration.
*/
public static final long POLL_INTERVAL_MILLIS = 250;
private final AgentControl agents;
private final TurnListener turnListener;
private final Predicate<String> ready; // CB-113: a target is deliverable only when available
@@ -264,6 +273,11 @@ public final class Injector {
// fail every queued message and release the target (CB-114) instead of
// polling it indefinitely with the caller's future never completing.
notReady = new ArrayList<>(t.queue);
log.warn("readiness grace for {} expired after {} polls ({}s): target never "
+ "became deliverable, so failing {} queued message(s) that never "
+ "reached its pane",
target, READINESS_GRACE_POLLS,
READINESS_GRACE_POLLS * POLL_INTERVAL_MILLIS / 1000, notReady.size());
t.queue.clear();
t.notReadySincePoll = 0;
}
@@ -388,6 +402,10 @@ public final class Injector {
t.awaitingPostTurnPickup = false;
t.postTurnObserved = false;
}
log.warn("{} is gone, dropping its queue: {} message(s) failed{}; cause: {}", target,
pending.size(),
hadDeliveredTurn ? " (including one turn already in flight whose completion was never confirmed)" : "",
cause.getMessage());
forget.accept(target); // the worker is gone — clear its readiness/presence too (CB-114)
for (Pending p : pending) {
p.delivered().completeExceptionally(cause);
@@ -109,10 +109,10 @@ public final class BridgeMcp {
? callers.resolve(req.getRemoteAddr(), req.getRemotePort(),
req.getHeader("Authorization"))
: legacyPrincipal(identity, req.getRemoteAddr(), req.getRemotePort());
// CB-532: guard on the ROLE, not on the terminal being null. A named lead now
// carries its pane too, and enrolling a lead in the worker presence map would
// have it counted as an available worker.
if (p.isWorker()) presence.markPresent(p.terminal());
// CB-532: guard on the ROLE, not on the terminal being null. This excludes a
// lead, which carries its pane too, while including every spawned member role.
// Enrolling a lead would count it as an available member in the roster.
markSpawnedMemberPresent(p, presence);
return McpTransportContext.create(Map.of(
CALLER_TERMINAL, orEmpty(p.terminal()),
CALLER_PID, Long.toString(p.pid()),
@@ -365,6 +365,13 @@ public final class BridgeMcp {
return transport;
}
/** Mark a connected spawned member available for the injector readiness gate. */
static void markSpawnedMemberPresent(Principal caller, MemberPresence presence) {
if (caller.isSpawnedMember()) {
presence.markPresent(caller.terminal());
}
}
/** Graceful shutdown of the MCP server. */
public void close() {
server.closeGracefully();
@@ -44,23 +44,6 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
private final SubscriptionGuard guard;
/**
* Standing instruction appended to the worker's system prompt so it returns its result via
* {@code bridge_reply}. Injected as a launch flag, so nothing is written to the worker's
* profile — it is guidance, and a worker that never replies is caught by the send's timeout.
*/
static final String REPLY_CHARTER =
"You are an off-subscription worker in the claude-bridge fleet. Every message you "
+ "receive arrives through the bridge, and the ONLY channel back to the sender is the "
+ "bridge_reply MCP tool. Text you write in your terminal is NOT sent anywhere — the "
+ "sender cannot see your screen, so an in-terminal answer is silently discarded. "
+ "Therefore you MUST end EVERY turn by calling bridge_reply with `content` set to your "
+ "complete response. This holds for every message without exception — tasks, questions, "
+ "clarifications, acknowledgements, and ordinary back-and-forth conversation. Call "
+ "bridge_reply exactly once, as the final action of your turn, with your full answer in "
+ "`content`; never wait for confirmation first. If you end a turn without calling "
+ "bridge_reply, the sender receives nothing and the exchange stalls.";
/**
* Production constructor — disables the spawn-ready gate ({@code spawnReadyTimeoutMs == 0}) so
* existing deployments and tests keep the legacy non-blocking spawn semantics.
@@ -93,11 +76,11 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs, long spawnReadyPollMs,
Supplier<String> tabLabelTemplate) {
Supplier<BridgedConfig.Fleet> fleet) {
this(agents, spaces, guard, profiles, defaultProfile, env,
spawnReadyTimeoutMs,
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs),
tabLabelTemplate);
fleet);
}
/**
@@ -129,31 +112,19 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
/**
* Full testability constructor, plus the fleet-wide tab-label template (CB-557).
*
* @param tabLabelTemplate {@code fleet.tabLabel}; {@code null}/blank ⇒
* {@link BridgedConfig.Fleet#DEFAULT_TAB_LABEL}
* @param fleet live fleet config, read once for each spawn
*/
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs,
LongSupplier nowMillis, Runnable sleeper,
Supplier<String> tabLabelTemplate) {
Supplier<BridgedConfig.Fleet> fleet) {
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
spawnReadyTimeoutMs, nowMillis, sleeper, tabLabelTemplate);
spawnReadyTimeoutMs, nowMillis, sleeper, fleet);
this.guard = guard;
}
/**
* {@inheritDoc}
*
* <p>A legacy spawn with no session identity is a fresh, launcher-derived session — delegate to
* the session-aware form with no name and no resume id.
*/
@Override
protected Launch buildLaunch(BridgedConfig.Profile cfg) {
return buildLaunch(cfg, null, null);
}
/**
* {@inheritDoc}
*
@@ -164,7 +135,7 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
* applied here — see {@link #applySessionIdentity}.
*/
@Override
protected Launch buildLaunch(BridgedConfig.Profile cfg, String sessionName, String resumeSessionId) {
protected Launch buildLaunch(BridgedConfig.Profile cfg, LaunchSpec spec) {
// CB-539: a profile may deliberately opt into the subscription (subscription: true) when no
// off-subscription endpoint exists for it — e.g. `sonnet` on `ccs`. That profile gets no
// ANTHROPIC_BASE_URL/AUTH_TOKEN (there is nothing to point them at) and the guard's base_url
@@ -210,9 +181,9 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
// id via -r and passes no --session-id (the two conflict). Both are injected before the
// model flag so --model keeps outranking the operator's own argv.
// mutableArgv: argvWithBridge may hand back the profile's own (immutable) List.of when it
// has no MCP — session flags must be added into a list we own.
List<String> argv = mutableArgv(argvWithBridge(cfg));
String agentSessionId = applySessionIdentity(argv, sessionName, resumeSessionId);
// has neither MCP nor a charter — session flags must be added into a list we own.
List<String> argv = mutableArgv(argvWithBridge(cfg, spec.charter()));
String agentSessionId = applySessionIdentity(argv, spec.sessionName(), spec.resumeSessionId());
return new Launch(workerEnv, argvWithModel(argv, cfg), agentSessionId);
}
@@ -246,22 +217,26 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
}
/**
* The launch argv, plus — when {@code worker.mcpUrl} is set — inline {@code --mcp-config} for
* the bridge server and {@code --append-system-prompt} for the {@link #REPLY_CHARTER}. Neither
* touches the profile's config; both are pure command-line flags. This inline-flag mount is
* Claude Code specific — other adapters mount MCP and instructions their own way.
* The launch argv, plus an inline {@code --mcp-config} when {@code worker.mcpUrl} is set and
* {@code --append-system-prompt} when the base composed a charter. Neither touches the profile's
* config; both are pure command-line flags. This inline-flag mount is Claude Code specific —
* other adapters mount MCP and instructions their own way.
*/
private List<String> argvWithBridge(BridgedConfig.Profile cfg) {
if (!cfg.hasMcp()) {
private List<String> argvWithBridge(BridgedConfig.Profile cfg, String charter) {
if (!cfg.hasMcp() && charter == null) {
return cfg.argv();
}
String mcpJson = "{\"mcpServers\":{\"bridge\":{\"type\":\"http\",\"url\":\""
+ cfg.mcpUrl() + "\"}}}";
List<String> argv = mutableArgv(cfg.argv());
argv.add("--mcp-config");
argv.add(mcpJson);
argv.add("--append-system-prompt");
argv.add(REPLY_CHARTER);
if (cfg.hasMcp()) {
String mcpJson = "{\"mcpServers\":{\"bridge\":{\"type\":\"http\",\"url\":\""
+ cfg.mcpUrl() + "\"}}}";
argv.add("--mcp-config");
argv.add(mcpJson);
}
if (charter != null) {
argv.add("--append-system-prompt");
argv.add(charter);
}
return argv;
}
@@ -44,7 +44,7 @@ import java.util.regex.Pattern;
* <li>{@code namePrefix} (constructor arg) — the label prefix ({@code claude}, {@code opencode})
* that drives both unique naming and the orphan-reap pattern, so each adapter reaps only its
* own kind of pane and never another's.</li>
* <li>{@link #buildLaunch(BridgedConfig.Profile)} — the peer-specific env map + argv, including any
* <li>{@link #buildLaunch(BridgedConfig.Profile, LaunchSpec)} — the peer-specific env map + argv, including any
* subscription/guard check, MCP mount, and instruction injection. The base never sees how the
* peer is configured; it only places and starts the returned {@link Launch}.</li>
* </ul>
@@ -80,15 +80,25 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
private final AtomicLong nameSeq = new AtomicLong(); // per-peer counter (herdr agent names only)
/**
* The {@code fleet.tabLabel} template; a {@code null} supplier or a {@code null}/blank value ⇒
* {@link BridgedConfig.Fleet#DEFAULT_TAB_LABEL}. A profile's own {@code tabLabel} still
* overrides it.
* Live fleet config, read once per spawn. A null supplier or value leaves tab labels at their
* default and supplies no role charter. A profile's own {@code tabLabel} still overrides it.
*
* <p>CB-559: a supplier rather than a String, so a config reload renames the <em>next</em> tab
* without a restart. Existing tabs keep the label they were given — bridged does not rewrite a
* label it already wrote.
* <p>CB-559: a supplier rather than a snapshot, so a config reload affects the next launch
* without a restart. Existing tabs keep the label they were given.
*/
private final Supplier<String> tabLabelTemplate;
private final Supplier<BridgedConfig.Fleet> fleet;
/** The final instruction always requires a bridge reply when the bridge MCP is mounted. */
protected static final String REPLY_CHARTER =
"You are a spawned member in the claude-bridge fleet. Every message you receive arrives "
+ "through the bridge, and the ONLY channel back to the sender is the bridge_reply MCP tool. "
+ "Text you write in your terminal is NOT sent anywhere — the sender cannot see your screen, "
+ "so an in-terminal answer is silently discarded. Therefore you MUST end EVERY turn by calling "
+ "bridge_reply with `content` set to your complete response. This holds for every message without "
+ "exception — tasks, questions, clarifications, acknowledgements, and ordinary back-and-forth "
+ "conversation. Call bridge_reply exactly once, as the final action of your turn, with your full "
+ "answer in `content`; never wait for confirmation first. If you end a turn without calling "
+ "bridge_reply, the sender receives nothing and the exchange stalls.";
/**
* Tab numbers, counted per {@code role/profile} pair (CB-557).
@@ -142,21 +152,19 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
}
/**
* As above, plus the {@code fleet.tabLabel} template (CB-557).
* As above, plus the live {@code fleet} config (CB-557).
*
* @param tabLabelTemplate fleet-wide tab-label template, read per spawn (CB-559); {@code null},
* or a supplier yielding {@code null}/blank ⇒
* {@link BridgedConfig.Fleet#DEFAULT_TAB_LABEL}. A separate constructor
* rather than a new parameter on the one above, so every existing call
* site keeps the default without an edit.
* @param fleet live fleet config, read once per spawn; {@code null} ⇒ default tab label and no
* role charter. A separate constructor rather than a new parameter on the one
* above, so every existing call site keeps the default without an edit.
*/
protected HerdrPeerLauncher(String namePrefix, AgentControl agents, WorkspaceControl spaces,
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs,
LongSupplier nowMillis, Runnable sleeper,
Supplier<String> tabLabelTemplate) {
this.tabLabelTemplate = tabLabelTemplate;
Supplier<BridgedConfig.Fleet> fleet) {
this.fleet = fleet;
this.namePrefix = namePrefix;
this.agents = agents;
this.spaces = spaces;
@@ -175,22 +183,7 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
* Any subscription/guard check, MCP mount, and instruction injection happen here. The env map
* and argv are adapter-private; the base only places and starts what is returned.
*/
protected abstract Launch buildLaunch(BridgedConfig.Profile cfg);
/**
* Session-aware variant of {@link #buildLaunch(BridgedConfig.Profile)} (CB-547a). Default
* discards the session identity and delegates to the profile-only form, so an adapter that
* carries no durable peer session (opencode, say) inherits byte-identical behaviour and needs
* no change. An adapter that does (Claude Code) overrides this to mint/resume the id and to
* surface it on the returned {@link Launch#agentSessionId()}.
*
* @param cfg the resolved profile to spawn
* @param sessionName the bridge's logical session name, or null/blank for launcher-derived
* @param resumeSessionId the peer's own prior session id to resume, or null/blank for fresh
*/
protected Launch buildLaunch(BridgedConfig.Profile cfg, String sessionName, String resumeSessionId) {
return buildLaunch(cfg);
}
protected abstract Launch buildLaunch(BridgedConfig.Profile cfg, LaunchSpec spec);
/** Direct transport access for peer-specific, non-turn control operations. */
protected final AgentControl agents() {
@@ -225,6 +218,10 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
}
}
/** All per-spawn values adapters may need, including the base-composed effective charter. */
protected record LaunchSpec(String sessionName, String resumeSessionId, MemberRole role, String charter) {
}
// --- profile surface -----------------------------------------------------------------------
/** The configured peer profile names (what {@code spawn(profile)} accepts). */
@@ -295,18 +292,23 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
}
/**
* Spawn a peer with session identity (CB-547a). {@code sessionName} and {@code resumeSessionId}
* are threaded from the {@link SpawnRequest} into {@link #buildLaunch(BridgedConfig.Profile,
* String, String)}, and the launch's resolved agent-session id is returned alongside the agent
* so the caller can put it on the {@link PeerHandle}.
* Spawn a peer with session identity (CB-547a). The session values, role, and charter are
* threaded from the {@link SpawnRequest} into {@link #buildLaunch(BridgedConfig.Profile,
* LaunchSpec)}, and the launch's resolved agent-session id is returned alongside the agent so
* the caller can put it on the {@link PeerHandle}.
*/
protected Spawned spawnInternal(String profileName, String requestedCwd, String callerCwd,
String sessionName, String resumeSessionId, MemberRole role) {
BridgedConfig.Profile cfg = requireProfile(profileName);
Launch launch = buildLaunch(cfg, sessionName, resumeSessionId);
BridgedConfig.Fleet liveFleet = fleet == null ? null : fleet.get();
String roleCharter = liveFleet == null ? null : liveFleet.charterFor(role);
String replyCharter = cfg.hasMcp() ? REPLY_CHARTER : null;
String charter = roleCharter == null ? replyCharter
: replyCharter == null ? roleCharter : roleCharter + "\n\n" + replyCharter;
Launch launch = buildLaunch(cfg, new LaunchSpec(sessionName, resumeSessionId, role, charter));
String cwd = resolveCwd(requestedCwd, cfg, callerCwd);
Agent agent = cfg.tabPlacement()
? spawnInTab(cfg, launch.env(), launch.argv(), cwd, role)
? spawnInTab(cfg, launch.env(), launch.argv(), cwd, role, liveFleet)
: spawnAsPane(cfg, launch.env(), launch.argv(), cwd);
return new Spawned(agent, launch.agentSessionId());
}
@@ -382,7 +384,7 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
/** Dedicated worker space → own tab (carrying cwd+env) → start the peer into the seed pane. */
private Agent spawnInTab(BridgedConfig.Profile cfg, Map<String, String> workerEnv,
List<String> argv, String cwd, MemberRole role) {
List<String> argv, String cwd, MemberRole role, BridgedConfig.Fleet liveFleet) {
Workspace space = spaces.ensureWorkspace(cfg.workspace());
Tab.Created tab = spaces.createTab(space.workspaceId(), cwd, workerEnv);
log.info("spawning {} profile={} space={} tab={} cwd={}",
@@ -415,7 +417,7 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
tidy("label tab " + tab.tab().tabId(),
() -> spaces.renameTab(tab.tab().tabId(),
cfg.renderTabLabel(
tabLabelTemplate == null ? null : tabLabelTemplate.get(),
liveFleet == null ? null : liveFleet.tabLabel(),
role, nextLabelSeq(role, cfg.profile()))));
log.info("{} started pane={} tab={} terminal={}",
namePrefix, started.agent().paneId(), started.agent().tabId(), started.agent().terminalId());
@@ -54,38 +54,9 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
/** Writer for the generated {@code opencode.json}. */
private static final ObjectMapper JSON = new ObjectMapper();
/**
* Standing instruction written to the charter file and mounted via the config's
* {@code instructions} so the worker returns its result through {@code bridge_reply}. Kept on
* disk (not a launch flag) because opencode's {@code instructions} takes file paths, not inline
* text — the file is regenerated per spawn and never touches the worker's own profile.
*/
static final String REPLY_CHARTER =
"You are an off-subscription worker in the claude-bridge fleet, running under opencode. "
+ "Every message you receive arrives through the bridge, and the ONLY channel back to the "
+ "sender is the bridge_reply MCP tool. Text you write in your terminal is NOT sent "
+ "anywhere — the sender cannot see your screen, so an in-terminal answer is silently "
+ "discarded. Therefore you MUST end EVERY turn by calling bridge_reply with `content` set "
+ "to your complete response. This holds for every message without exception — tasks, "
+ "questions, clarifications, acknowledgements, and ordinary back-and-forth conversation. "
+ "Call bridge_reply exactly once, as the final action of your turn, with your full answer "
+ "in `content`; never wait for confirmation first. If you end a turn without calling "
+ "bridge_reply, the sender receives nothing and the exchange stalls.";
/** Root under which per-spawn opencode config dirs are created (injectable for tests). */
private final Path configRoot;
/**
* The current spawn's resume-target session id, threaded from {@link #spawn(SpawnRequest)} to
* {@link #buildLaunch} across the base's {@code spawn -> spawnInternal -> buildLaunch} chain,
* which carries no request. A plain field would race under concurrent spawns (the base supports
* them), so it is thread-local: each spawn captures its own request's id on its own thread, and
* {@code buildLaunch}, synchronous and same-thread, reads exactly that one. Set only around the
* {@code super.spawn} call and cleared in {@code finally}, so a paused/leftover value can never
* bleed into the next spawn.
*/
private final ThreadLocal<String> resumeSessionId = new ThreadLocal<>();
/**
* Session discovery against opencode's on-disk storage ({@link OpenCodeSessionDiscovery}) —
* the one seam that knows opencode's private session-file layout. Its root is injectable for
@@ -126,10 +97,10 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs, long spawnReadyPollMs,
Supplier<String> tabLabelTemplate) {
Supplier<BridgedConfig.Fleet> fleet) {
this(agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs,
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs),
defaultConfigRoot(), defaultDiscoveryRoot(), tabLabelTemplate);
defaultConfigRoot(), defaultDiscoveryRoot(), fleet);
}
/**
@@ -164,8 +135,7 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
/**
* Full testability constructor, plus the fleet-wide tab-label template (CB-557).
*
* @param tabLabelTemplate {@code fleet.tabLabel}; {@code null}/blank ⇒
* {@link BridgedConfig.Fleet#DEFAULT_TAB_LABEL}
* @param fleet live fleet config, read once for each spawn
*/
public OpenCodeLauncher(AgentControl agents, WorkspaceControl spaces,
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
@@ -173,9 +143,9 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
long spawnReadyTimeoutMs,
LongSupplier nowMillis, Runnable sleeper,
Path configRoot, Path discoveryRoot,
Supplier<String> tabLabelTemplate) {
Supplier<BridgedConfig.Fleet> fleet) {
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
spawnReadyTimeoutMs, nowMillis, sleeper, tabLabelTemplate);
spawnReadyTimeoutMs, nowMillis, sleeper, fleet);
this.configRoot = configRoot;
this.discovery = new OpenCodeSessionDiscovery(discoveryRoot);
}
@@ -199,14 +169,15 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
* with {@code -m}.
*/
@Override
protected Launch buildLaunch(BridgedConfig.Profile cfg) {
protected Launch buildLaunch(BridgedConfig.Profile cfg, LaunchSpec spec) {
Map<String, String> workerEnv = baseEnv(cfg);
// A config file is needed for the bridge MCP mount, for a pinned endpoint (CB-508), or both.
if (cfg.hasMcp() || hasCustomProvider(cfg)) {
workerEnv.put("OPENCODE_CONFIG", writeConfig(cfg).toString());
}
applyGitToken(workerEnv, cfg);
return new Launch(workerEnv, argvWithResume(argvWithModel(argvWithAuto(cfg), cfg)));
return new Launch(workerEnv,
argvWithResume(argvWithModel(argvWithAuto(cfg), cfg), spec.resumeSessionId()));
}
/**
@@ -243,11 +214,9 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
* The launch argv plus, on a resumed spawn, opencode's {@code -s <id>} flag to continue a prior
* conversation by its session id. {@code -s, --session <id>} resumes an existing session; on a
* fresh spawn (no resume target) no flag is added, letting opencode start a brand-new session.
* The id comes from the current spawn request's {@code resumeSessionId}, threaded per-thread by
* {@link #spawn(SpawnRequest)}.
* The id comes from the base launch spec.
*/
private List<String> argvWithResume(List<String> argv) {
String id = resumeSessionId.get();
private List<String> argvWithResume(List<String> argv, String id) {
if (id == null || id.isBlank()) {
return argv;
}
@@ -380,31 +349,11 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
return afterScheme.contains("/") ? trimmed : trimmed + "/v1";
}
/**
* {@inheritDoc}
*
* <p>adds this adapter's session-identity work around the base's spawn — as opencode cannot be
* told its session id at spawn (see {@link Capability#SESSION_RESUME} vs
* {@link Capability#SESSION_NAME}), identity is only ever adopted after the fact:
* <ul>
* <li>the request's {@code resumeSessionId} is remembered for {@link #buildLaunch} to turn
* into {@code -s <id>}; and</li>
* <li>the returned handle is wrapped so its
* {@link dev.ltms.bridged.peer.PeerHandle#agentSessionId()} performs lazy session
* discovery against opencode's storage (see {@link OpenCodeSessionDiscovery}) — always
* non-blocking, {@code null} until opencode has persisted the session record.</li>
* </ul>
*/
/** Add lazy on-disk session discovery to the base handle. */
@Override
public PeerHandle spawn(SpawnRequest req) {
resumeSessionId.set(req.resumeSessionId());
try {
PeerHandle inner = super.spawn(req);
return new SessionAwareHandle(inner, discovery, effectiveCwd(req));
} finally {
// Never let a paused/leftover resume id bleed into the next spawn on this thread.
resumeSessionId.remove();
}
PeerHandle inner = super.spawn(req);
return new SessionAwareHandle(inner, discovery, effectiveCwd(req));
}
/**
@@ -277,7 +277,7 @@ public final class MessageService {
}
boolean failed = rendezvous.resolveFailure(waiter, reason);
if (failed) {
log.debug("abandoned send to {}: {}", target, reason);
log.warn("abandoning the blocked send to {}: {}", target, reason);
}
return failed;
}
@@ -1,5 +1,6 @@
package dev.ltms.bridged.session;
import dev.ltms.bridged.auth.MemberLifecycle;
import dev.ltms.bridged.herdr.Agent;
import dev.ltms.bridged.inject.TurnListener;
import dev.ltms.bridged.inject.MemberPresence;
@@ -28,8 +29,7 @@ import java.util.function.LongSupplier;
* teardown on top.
*
* <p>The state machine is intentionally one-shot / no-reuse: every acquired worker is fresh,
* and a finished or released worker is torn down, never pooled. {@link #recycle} is a convenience
* for {@code release + acquire} with a new distinct pane id.
* and a finished or released worker is torn down, never pooled or reused.
*
* <p>The manager implements {@link TurnListener} so the injector's turn boundaries drive
* {@code READY → BUSY → DONE} (or {@code FAILED}). It exposes a {@link MemberPresence} view via
@@ -49,6 +49,7 @@ public final class SessionManager implements TurnListener {
private final LongSupplier nowNanos;
private final int contextCap;
private final boolean clearAfterTurn;
private volatile MemberLifecycle memberLifecycle = MemberLifecycle.NONE;
/** CB-520: notified with a terminalId on every acquire; no-op until wired. */
private final List<Consumer<String>> acquireListeners = new java.util.concurrent.CopyOnWriteArrayList<>();
@@ -139,7 +140,13 @@ public final class SessionManager implements TurnListener {
// to pick the profile out of that role's pool and to label the tab; a role kept only on
// the MemberSession is recorded after the spawn it was supposed to steer.
SpawnRequest req = new SpawnRequest(profile, requestedCwd, callerCwd, null, null, memberRole);
PeerHandle handle = launcher.spawn(req);
PeerHandle handle;
try {
handle = launcher.spawn(req);
} catch (RuntimeException e) {
log.warn("spawn failed for profile={} role={}: {}", profile, memberRole, e.getMessage());
throw e;
}
String resolvedProfile = resolveProfile(handle, profile);
String cwd = launcher.effectiveCwd(new SpawnRequest(resolvedProfile, requestedCwd, callerCwd));
long now = nowNanos.getAsLong();
@@ -157,6 +164,7 @@ public final class SessionManager implements TurnListener {
null,
null);
registry.put(handle.id(), session);
memberLifecycle.acquired(session.role(), session.profile(), session.terminalId());
log.debug("acquired session id={} terminal={} profile={} owner={}",
handle.id(), handle.terminalId(), session.profile(), session.ownerTerminal());
notifyAcquired(session.terminalId());
@@ -177,7 +185,7 @@ public final class SessionManager implements TurnListener {
* <p>CB-544: these are two concerns that used to be fused. Stopping the pane is correct on every
* teardown — the worker process must end. Removing the worktree is a destructive act that is only
* correct for a deliberately-finished teardown (an explicit stop of a completed session, the
* reaper releasing a genuinely idle one, a context-capped or recycled session). A shutdown drain
* reaper releasing a genuinely idle one, or a context-capped session). A shutdown drain
* must stop panes but preserve worktrees: a worker's uncommitted work exists in exactly one
* place — its worktree — so deleting it while the daemon simply goes down is silent data loss,
* with no copy and no error. Do NOT fuse these back together; the cost of an orphaned worktree
@@ -187,6 +195,7 @@ public final class SessionManager implements TurnListener {
MemberSession removed = registry.remove(paneId);
boolean preserveWorktree = cause == ReleaseCause.SHUTDOWN;
if (removed != null) {
memberLifecycle.released(removed.terminalId());
log.debug("releasing session pane={} terminal={} state={} cause={}",
removed.paneId(), removed.terminalId(), removed.state(), cause);
if (preserveWorktree && removed.worktree() != null) {
@@ -244,7 +253,7 @@ public final class SessionManager implements TurnListener {
/**
* Register a callback invoked with a session's {@code terminalId} whenever it is released
* (CB-516). Every teardown path funnels through {@link #release}, so one hook covers the REST
* and MCP stop tools, the idle-TTL reaper, {@code recycle}, and shutdown drain alike.
* and MCP stop tools, the idle-TTL reaper, and shutdown drain alike.
*
* <p>Added rather than injected because {@code MessageService} — one intended listener — is
* constructed after this manager (it needs the injector and rendezvous, which need the session
@@ -257,6 +266,11 @@ public final class SessionManager implements TurnListener {
}
}
/** Inject the optional member-slot lifecycle after construction without changing constructors. */
public void setMemberLifecycle(MemberLifecycle memberLifecycle) {
this.memberLifecycle = memberLifecycle == null ? MemberLifecycle.NONE : memberLifecycle;
}
/** A listener failure must never prevent the acquisition it is reacting to. */
private void notifyAcquired(String terminalId) {
if (terminalId == null) {
@@ -304,6 +318,8 @@ public final class SessionManager implements TurnListener {
worktrees.overlayParity(repoRoot, path, launcher.parityOverlay(preResolvedProfile));
handle = launcher.spawn(new SpawnRequest(profile, path, callerCwd, null, null, memberRole));
} catch (RuntimeException e) {
log.warn("spawn failed for profile={} role={} branch={} path={}: {}",
preResolvedProfile, memberRole, branch, path, e.getMessage());
if (path != null) {
try {
worktrees.remove(repoRoot, path);
@@ -330,6 +346,7 @@ public final class SessionManager implements TurnListener {
path,
branch);
registry.put(handle.id(), session);
memberLifecycle.acquired(session.role(), session.profile(), session.terminalId());
log.debug("acquired worktree session id={} terminal={} profile={} branch={} path={}",
handle.id(), handle.terminalId(), session.profile(), session.branch(), session.worktree());
notifyAcquired(session.terminalId());
@@ -360,19 +377,6 @@ public final class SessionManager implements TurnListener {
return launcher.defaultProfile();
}
/**
* Release the old session and acquire a fresh one with the same profile and working directory.
* The new session is guaranteed to have a pane id distinct from the old one (no-reuse invariant).
*/
public MemberSession recycle(String paneId) {
MemberSession old = registry.get(paneId);
if (old == null) {
throw new IllegalArgumentException("no session for paneId " + paneId);
}
release(paneId);
return acquire(old.profile(), old.cwd(), old.cwd(), old.ownerTerminal());
}
/** The session for {@code paneId}, if it is still registered and not released. */
public Optional<MemberSession> get(String paneId) {
return Optional.ofNullable(registry.get(paneId));
@@ -488,8 +492,10 @@ public final class SessionManager implements TurnListener {
MemberSession current = findByTerminal(target);
if (current == null) return;
if (current.state() == MemberSession.State.RELEASED) return;
MemberSession.State priorState = current.state();
if (replace(current, current.withState(MemberSession.State.FAILED))) {
log.debug("session marked failed terminal={} pane={}", target, current.paneId());
log.warn("member terminal={} pane={} can no longer be delegated to: its turn never resolved "
+ "(was {} when it failed)", target, current.paneId(), priorState);
}
}
@@ -505,7 +511,11 @@ public final class SessionManager implements TurnListener {
if (s.state() != MemberSession.State.READY && s.state() != MemberSession.State.DONE) {
continue;
}
if (now - s.lastActivityAtNanos() > idleTtlNanos) {
long idleNanos = now - s.lastActivityAtNanos();
if (idleNanos > idleTtlNanos) {
log.debug("reaping idle session terminal={} pane={}: idle {}s exceeds the {}s ttl",
s.terminalId(), s.paneId(), TimeUnit.NANOSECONDS.toSeconds(idleNanos),
TimeUnit.NANOSECONDS.toSeconds(idleTtlNanos));
release(s.paneId());
reaped++;
}
@@ -3,6 +3,8 @@ package dev.ltms.bridged.auth;
import dev.ltms.bridged.herdr.FakeHerdr;
import dev.ltms.bridged.herdr.PaneLocator;
import dev.ltms.bridged.mcp.ConnectionIdentity;
import dev.ltms.bridged.config.BridgedConfig;
import dev.ltms.bridged.peer.MemberRole;
import org.junit.jupiter.api.Test;
import java.util.Map;
@@ -32,6 +34,17 @@ class CallerResolverTest {
return identity(999_999);
}
private static MemberRegistry boundMembers(String slot, MemberRole role) {
Map<String, BridgedConfig.Slot> architects = role == MemberRole.ARCHITECT
? Map.of("lead-designer", new BridgedConfig.Slot("sonnet")) : Map.of();
Map<String, BridgedConfig.Slot> devs = role == MemberRole.DEV
? Map.of("builder", new BridgedConfig.Slot("sonnet")) : Map.of();
MemberRegistry members = new MemberRegistry(
new BridgedConfig.Fleet(Map.of(), architects, devs, Map.of(), null));
assertTrue(members.bind(slot, "term_a"));
return members;
}
@Test
void aLoopbackWorkerPaneResolvesToWorkerRegardlessOfAuthMode() {
Principal underTrust = new CallerResolver(workerIdentity()).resolve("127.0.0.1", 42, null);
@@ -302,8 +315,9 @@ class CallerResolverTest {
@Test
void aBoundArchitectPaneResolvesToArchitectBeforeTheWorkerFallback() {
// This is the production construction path used by Bridged.
Principal p = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
Map::of, () -> Map.of("term_a", "lead-designer"))
Map::of, boundMembers("architect:lead-designer", MemberRole.ARCHITECT))
.resolve("127.0.0.1", 42, null);
assertEquals(Role.ARCHITECT, p.role(),
@@ -316,7 +330,7 @@ class CallerResolverTest {
@Test
void anArchitectNeedsNoTokenEvenInTokenMode() {
Principal p = CallerResolver.withLeadsAndMembers(workerIdentity(), true, "s3cret",
Map::of, () -> Map.of("term_a", "lead-designer"))
Map::of, boundMembers("architect:lead-designer", MemberRole.ARCHITECT))
.resolve("127.0.0.1", 42, null);
assertEquals(Role.ARCHITECT, p.role(),
@@ -325,9 +339,10 @@ class CallerResolverTest {
@Test
void anUnboundPaneStillResolvesAsAWorker() {
Map<String, String> arch = Map.of("term_elsewhere", "reviewer");
MemberRegistry members = new MemberRegistry(new BridgedConfig.Fleet(Map.of(),
Map.of("lead-designer", new BridgedConfig.Slot("sonnet")), Map.of(), Map.of(), null));
Principal p = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
Map::of, () -> arch).resolve("127.0.0.1", 42, null);
Map::of, members).resolve("127.0.0.1", 42, null);
assertEquals(Role.WORKER, p.role());
assertNull(p.name());
@@ -337,7 +352,8 @@ class CallerResolverTest {
@Test
void aLeadWinsOverAnArchitectBindingForTheSamePane() {
Principal p = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
() -> Map.of("term_a", "opus-5.0"), () -> Map.of("term_a", "lead-designer"))
() -> Map.of("term_a", "opus-5.0"),
boundMembers("architect:lead-designer", MemberRole.ARCHITECT))
.resolve("127.0.0.1", 42, null);
assertEquals(Role.PRIMARY, p.role(),
@@ -349,26 +365,17 @@ class CallerResolverTest {
/** The registry is live, like leads: a binding injected after construction is honoured. */
@Test
void anArchitectBoundAfterConstructionIsHonouredWithoutRebuildingTheResolver() {
Map<String, String> live = new java.util.HashMap<>();
MemberRegistry members = new MemberRegistry(new BridgedConfig.Fleet(Map.of(),
Map.of("lead-designer", new BridgedConfig.Slot("sonnet")), Map.of(), Map.of(), null));
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
Map::of, () -> live);
Map::of, members);
assertEquals(Role.WORKER, r.resolve("127.0.0.1", 42, null).role());
live.put("term_a", "lead-designer"); // the later lifecycle binds the slot
assertTrue(members.bind("architect:lead-designer", "term_a")); // the later lifecycle binds the slot
assertEquals(Role.ARCHITECT, r.resolve("127.0.0.1", 42, null).role());
assertEquals("lead-designer", r.members().get("term_a"));
}
@Test
void theArchitectMapFormIsCopiedSoLaterMutationCannotGrantArchitect() {
Map<String, String> mutable = new java.util.LinkedHashMap<>();
CallerResolver r = new CallerResolver(workerIdentity(), false, null, Map.of(), mutable);
mutable.put("term_a", "sneaky");
assertEquals(Role.WORKER, r.resolve("127.0.0.1", 42, null).role());
assertEquals("architect:lead-designer", r.members().get("term_a"));
}
@Test
@@ -381,7 +388,8 @@ class CallerResolverTest {
@Test
void anArchitectOwnsItsOwnPaneAndNoOther() {
Principal arch = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
Map::of, () -> Map.of("term_a", "lead-designer")).resolve("127.0.0.1", 42, null);
Map::of, boundMembers("architect:lead-designer", MemberRole.ARCHITECT))
.resolve("127.0.0.1", 42, null);
assertTrue(arch.ownsSession("term_a"));
assertTrue(Authz.permits(arch, Authz.Action.REPLY, "term_a"));
@@ -389,6 +397,15 @@ class CallerResolverTest {
assertFalse(Authz.permits(arch, Authz.Action.REPLY, "term_b"));
}
@Test
void aBoundNonArchitectSlotStillResolvesAsAWorker() {
Principal p = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
Map::of, boundMembers("dev:builder", MemberRole.DEV))
.resolve("127.0.0.1", 42, null);
assertEquals(Role.WORKER, p.role(), "a dev binding must never grant architect rights");
}
@Test
void tokenModeRequiresANonEmptyConfiguredToken() {
ConnectionIdentity id = nonWorkerIdentity();
@@ -111,6 +111,28 @@ class BridgedConfigTest {
assertDoesNotThrow(() -> BridgedConfig.load(f));
}
@Test
void absentChartersRemainValidAndPresentChartersUseRoleWireNames(@TempDir Path dir) throws Exception {
Path absent = dir.resolve("absent.yaml");
Files.writeString(absent, "fleet: {}\n");
BridgedConfig withoutCharters = BridgedConfig.load(absent);
assertDoesNotThrow(withoutCharters::validateCharters);
assertNull(withoutCharters.fleet().charterFor(MemberRole.ARCHITECT));
Path blank = dir.resolve("blank.yaml");
Files.writeString(blank, "fleet:\n charters:\n architect: ' '\n");
IllegalStateException blankError = assertThrows(IllegalStateException.class,
() -> BridgedConfig.load(blank).validateCharters());
assertTrue(blankError.getMessage().contains("fleet.charters.architect is blank"));
Path unknown = dir.resolve("unknown.yaml");
Files.writeString(unknown, "fleet:\n charters:\n architetc: text\n");
IllegalStateException unknownError = assertThrows(IllegalStateException.class,
() -> BridgedConfig.load(unknown).validateCharters());
assertTrue(unknownError.getMessage().contains("architetc"));
assertTrue(unknownError.getMessage().contains("[architect, dev, reviewer]"));
}
/**
* CB-530. Unknown keys stay ignored — config must be allowed to run ahead of the code — but they
* must be NAMED at load. A whole block that parses, is dropped, and is never mentioned again is
@@ -1167,6 +1189,41 @@ class BridgedConfigTest {
"a profile without the key stays off-subscription (the default)");
}
@Test
void parityOverlayDefaultsToEnvFilesOnly(@TempDir Path dir) throws Exception {
Path f = dir.resolve("no-overlay.yaml");
Files.writeString(f, """
bind:
port: 8080
profiles:
gx10:
baseUrl: http://gx10.gw:8000
""");
BridgedConfig cfg = BridgedConfig.load(f);
assertEquals(List.of(".env", ".envrc"),
cfg.profiles().get("gx10").parityOverlay(),
"the default parity overlay is the env files; settings.local.json is no longer copied by default");
}
@Test
void parityOverlayExplicitListIsPreservedVerbatim(@TempDir Path dir) throws Exception {
Path f = dir.resolve("explicit-overlay.yaml");
Files.writeString(f, """
bind:
port: 8080
profiles:
gx10:
baseUrl: http://gx10.gw:8000
parityOverlay: [".claude/settings.local.json", ".env"]
""");
BridgedConfig cfg = BridgedConfig.load(f);
assertEquals(List.of(".claude/settings.local.json", ".env"),
cfg.profiles().get("gx10").parityOverlay(),
"an operator's explicit list survives verbatim — the default only changes when unset");
}
// ── CB-542: subscription:true must not smuggle an unguarded endpoint via env: ───────────────
@Test
@@ -66,6 +66,64 @@ class ConfigRefTest {
assertEquals("[{profile}] {role}", ref.get().fleet().tabLabel());
}
@Test
void aCharterChangeIsHotAndReachesTheLiveConfig(@TempDir Path dir) throws Exception {
Path f = dir.resolve("bridged.yaml");
Files.writeString(f, yaml("""
fleet:
charters:
architect: old charter
"""));
ConfigRef ref = refFor(f);
assertEquals("old charter", ref.get().fleet().charterFor(
dev.ltms.bridged.peer.MemberRole.ARCHITECT));
Files.writeString(f, yaml("""
fleet:
charters:
architect: new charter
"""));
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied());
assertTrue(out.deferred().isEmpty());
assertEquals("new charter", ref.get().fleet().charterFor(
dev.ltms.bridged.peer.MemberRole.ARCHITECT));
}
@Test
void invalidChartersRefuseReloadAndKeepTheRunningConfig(@TempDir Path dir) throws Exception {
Path f = dir.resolve("bridged.yaml");
Files.writeString(f, yaml("""
fleet:
charters:
architect: valid charter
"""));
ConfigRef ref = refFor(f);
BridgedConfig before = ref.get();
Files.writeString(f, yaml("""
fleet:
charters:
architect: " "
"""));
ConfigRef.Outcome blank = ref.reload();
assertFalse(blank.applied());
assertTrue(blank.error().contains("fleet.charters.architect is blank"));
assertSame(before, ref.get());
Files.writeString(f, yaml("""
fleet:
charters:
architetc: valid charter
"""));
ConfigRef.Outcome unknown = ref.reload();
assertFalse(unknown.applied());
assertTrue(unknown.error().contains("architetc"));
assertTrue(unknown.error().contains("architect"));
assertSame(before, ref.get());
}
/**
* The point of the whole class: a consumer holding the ref sees the new value without being
* rebuilt. A component that captured {@code get()} into a field would still show the old one.
@@ -1,9 +1,14 @@
package dev.ltms.bridged.inject;
import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.LoggerContext;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.FakeHerdr;
import dev.ltms.bridged.msg.Rendezvous;
import org.junit.jupiter.api.Test;
import org.slf4j.LoggerFactory;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
@@ -156,6 +161,33 @@ class CompletionResolverTest {
assertEquals("No, 391 = 17 × 23.", waiter.getNow(null).text());
}
@Test
void marksAClippedCompletionPaneTail() {
String block = "⏺ " + "x".repeat(CompletionResolver.MAX_SCRAPE_CHARS + 1) + "\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
assertEquals("x".repeat(CompletionResolver.MAX_SCRAPE_CHARS)
+ "\n[Pane tail clipped: member did not call bridge_reply.]",
waiter.getNow(null).text());
}
@Test
void leavesAnUnclippedCompletionPaneTailUnmarked() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ complete report\n❯ ");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
assertEquals("complete report", waiter.getNow(null).text());
}
@Test
void resolvesSynchronouslyBeforePostTurnContextClearing() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ previous answer\n❯ ");
@@ -177,7 +209,8 @@ class CompletionResolverTest {
// block while resolve compares against a clip()'d tail. For a block longer than MAX_SCRAPE_CHARS
// the two capped representations differ even when the pane never changed, so the CB-115
// byte-identical guard failed to fire and a stale completion could resolve the send. Both sides
// must clip identically; here an unchanged >cap block on rapid back-to-back turns stays suppressed.
// must clip identically. The returned-text marker is added only after this comparison, so an
// unchanged >cap block on rapid back-to-back turns still stays suppressed.
String longBlock = "⏺ " + "x".repeat(CompletionResolver.MAX_SCRAPE_CHARS + 500) + "\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(longBlock);
Rendezvous rendezvous = new Rendezvous();
@@ -270,6 +303,40 @@ class CompletionResolverTest {
assertEquals("stuck on an error screen", waiter.getNow(null).text());
}
@Test
void failIsLoggedAtWarnWithTheReason() {
// CB-564: this used to be a bare DEBUG "failed send to X via turn-stall fallback" — a symptom
// with no cause, and below the level anyone watching for member health would see. A fail that
// resolves a caller's blocked send is at least WARN and must carry the reason.
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
ch.qos.logback.classic.Logger resolverLog =
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(CompletionResolver.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.setContext(ctx);
appender.start();
resolverLog.addAppender(appender);
resolverLog.setLevel(Level.WARN);
try {
FakeHerdr herdr = new FakeHerdr().readText("stuck on an error screen");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
var waiter = rendezvous.open("term_a");
resolver.fail("term_a", null);
String warn = appender.list.stream()
.filter(e -> e.getLevel().equals(Level.WARN))
.map(ILoggingEvent::getFormattedMessage)
.findFirst()
.orElse("no turn-stall WARN logged");
assertTrue(warn.contains("term_a"), "the log names the target: " + warn);
assertTrue(warn.contains("stuck on an error screen"), "the log carries the reason: " + warn);
assertTrue(waiter.isDone());
} finally {
resolverLog.detachAppender(appender);
}
}
// --- CB-116 waiter identity: a late completion never crosses into the next turn ---------
@Test
@@ -1,10 +1,15 @@
package dev.ltms.bridged.inject;
import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.LoggerContext;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.AgentStatus;
import dev.ltms.bridged.herdr.FakeHerdr;
import dev.ltms.bridged.herdr.HerdrException;
import org.junit.jupiter.api.Test;
import org.slf4j.LoggerFactory;
import java.util.ArrayList;
import java.util.List;
@@ -369,6 +374,40 @@ class InjectorTest {
assertTrue(inj.activeTargets().isEmpty(), "the target is reclaimed, not polled forever");
}
@Test
void readinessGraceExpiryIsLogged() {
// CB-562: the grace-expiry path used to clear the queue silently, so a message that never
// reached the worker's pane surfaced elsewhere as an unrelated turn-stall failure. Assert the
// expiry now names the real cause. (ListAppender capture pattern mirrors AuditLogTest.)
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
ch.qos.logback.classic.Logger injectorLog =
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(Injector.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.setContext(ctx);
appender.start();
injectorLog.addAppender(appender);
injectorLog.setLevel(Level.WARN);
try {
Injector inj = new Injector(new AgentControl(herdr), TurnListener.NOOP, _ -> false, _ -> {
});
inj.enqueue(T, "task");
for (int i = 0; i < READINESS_SAMPLES; i++) inj.onStatus(T, AgentStatus.IDLE);
String warn = appender.list.stream()
.filter(e -> e.getLevel().equals(Level.WARN))
.map(ILoggingEvent::getFormattedMessage)
.findFirst()
.orElse("no grace-expiry WARN logged");
assertTrue(warn.contains(T), "the log names the target terminal: " + warn);
assertTrue(warn.contains("never reached"), "the log names the real cause: " + warn);
assertTrue(warn.contains("1 queued message"),
"the log carries the failed message count: " + warn);
} finally {
injectorLog.detachAppender(appender);
}
}
@Test
void aWorkerThatBecomesReadyWithinTheGraceIsDeliveredNormally() {
// The readiness grace must not fail a worker that is merely slow to boot: once it becomes
@@ -397,6 +436,38 @@ class InjectorTest {
assertEquals(List.of(T), forgotten, "drop clears the gone worker's presence");
}
@Test
void dropIsLogged() {
// CB-564: a vanished worker used to drop its queue with no log at all — the only trace was
// whatever failed downstream (e.g. a caller's send timing out with no clue why). Assert the
// drop itself now names the cause and the number of messages it failed.
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
ch.qos.logback.classic.Logger injectorLog =
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(Injector.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.setContext(ctx);
appender.start();
injectorLog.addAppender(appender);
injectorLog.setLevel(Level.WARN);
try {
Injector inj = new Injector(new AgentControl(herdr), TurnListener.NOOP, _ -> true, _ -> {
});
inj.enqueue(T, "orphan");
inj.drop(T, new HerdrException("worker gone", "pane_not_found", null));
String warn = appender.list.stream()
.filter(e -> e.getLevel().equals(Level.WARN))
.map(ILoggingEvent::getFormattedMessage)
.findFirst()
.orElse("no drop WARN logged");
assertTrue(warn.contains(T), "the log names the target terminal: " + warn);
assertTrue(warn.contains("1 message"), "the log carries the failed message count: " + warn);
assertTrue(warn.contains("worker gone"), "the log carries the real cause: " + warn);
} finally {
injectorLog.detachAppender(appender);
}
}
@Test
void pollerDeliversToAnIdleWorker() throws Exception {
// End-to-end through the poller: idle worker → message delivered without manual onStatus.
@@ -2,6 +2,7 @@ package dev.ltms.bridged.mcp;
import dev.ltms.bridged.auth.Authz;
import dev.ltms.bridged.auth.CallerResolver;
import dev.ltms.bridged.auth.MemberRegistry;
import dev.ltms.bridged.auth.Principal;
import dev.ltms.bridged.auth.Role;
import dev.ltms.bridged.config.BridgedConfig;
@@ -69,7 +70,8 @@ class BridgeMcpAuthzTest {
mcp = new BridgeMcp(messages, workers, sessions, identity, sessions.asPresence(),
new PrimaryRegistry(null),
enforce ? new CallerResolver(identity) : null,
enforce ? CallerResolver.withLeadsAndMembers(identity, false, null,
Map::of, new MemberRegistry(null)) : null,
metrics);
return mcp;
}
@@ -12,6 +12,7 @@ import dev.ltms.bridged.msg.Rendezvous;
import dev.ltms.bridged.session.FakeWorktrees;
import dev.ltms.bridged.session.SessionManager;
import dev.ltms.bridged.peer.MemberRole;
import dev.ltms.bridged.inject.MemberPresence;
import dev.ltms.bridged.session.MemberSession;
import dev.ltms.bridged.session.WorktreeRequest;
import dev.ltms.bridged.member.ClaudeCodeLauncher;
@@ -404,6 +405,30 @@ class BridgeMcpTest {
assertDoesNotThrow(() -> messages.ackReply("term_a", msgId));
}
@Test
void spawnedMembersAreMarkedPresent() {
Principal worker = Principal.worker("term_worker", 200);
Principal architect = Principal.architect("lead-designer", "term_design", 400);
MemberPresence presence = new MemberPresence();
BridgeMcp.markSpawnedMemberPresent(worker, presence);
BridgeMcp.markSpawnedMemberPresent(architect, presence);
assertTrue(presence.isPresent("term_worker"));
assertTrue(presence.isPresent("term_design"));
}
@Test
void nonMembersAreNotMarkedPresent() {
Principal lead = Principal.leader("opus", "term_lead", 100);
MemberPresence presence = new MemberPresence();
BridgeMcp.markSpawnedMemberPresent(lead, presence);
BridgeMcp.markSpawnedMemberPresent(Principal.anonymous(), presence);
assertFalse(presence.isPresent("term_lead"));
}
@Test
void statusReportsLiveAgentStatus() {
FakeHerdr blocked = new FakeHerdr().agentStatus("blocked");
@@ -102,6 +102,53 @@ class ClaudeCodeLauncherTest {
"the operator's own args are preserved, in order, ahead of the model flag");
}
@Test
void appendsTheBaseComposedRoleAndReplyCharter() {
FakeHerdr herdr = new FakeHerdr();
String roleCharter = "You review changes.";
BridgedConfig.Profile cfg = new BridgedConfig.Profile(
"sonnet", "http://gx00.gw:8000", null, null, "BRIDGED_WORKER_TOKEN",
List.of("claude"), "tab", "bridged-workers", "w #{n}",
"http://127.0.0.1:8765/mcp", null, null);
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null,
0, 0L, () -> fleet(Map.of("reviewer", roleCharter), null));
svc.spawn(new SpawnRequest("sonnet", null, null, null, null, MemberRole.REVIEWER));
List<String> args = spawnedArgs(herdr);
int flag = args.indexOf("--append-system-prompt");
assertEquals(1, args.stream().filter("--append-system-prompt"::equals).count(),
"the composed charter is passed once");
assertEquals(roleCharter + "\n\n" + HerdrPeerLauncher.REPLY_CHARTER, args.get(flag + 1),
"the role charter comes first and the reply rule comes last");
}
@Test
void profileWithoutMcpOrRoleCharterGetsNoSystemPrompt() {
FakeHerdr herdr = new FakeHerdr();
ClaudeCodeLauncher svc = labelService(herdr, () -> fleet(Map.of(), null));
svc.spawn(new SpawnRequest("sonnet", null, null, null, null, MemberRole.DEV));
assertFalse(spawnedArgs(herdr).contains("--append-system-prompt"));
}
@Test
void profileWithoutMcpStillGetsItsRoleCharter() {
FakeHerdr herdr = new FakeHerdr();
String roleCharter = "You design changes.";
ClaudeCodeLauncher svc = labelService(herdr, () -> fleet(Map.of("architect", roleCharter), null));
svc.spawn(new SpawnRequest("sonnet", null, null, null, null, MemberRole.ARCHITECT));
List<String> args = spawnedArgs(herdr);
int flag = args.indexOf("--append-system-prompt");
assertTrue(flag >= 0, "a role charter does not need an MCP mount");
assertEquals(roleCharter, args.get(flag + 1));
assertFalse(args.contains("--mcp-config"));
}
private ClaudeCodeLauncher multiProfile(FakeHerdr herdr) {
BridgedConfig.Profile gx10 = new BridgedConfig.Profile("gx10", "http://gx10.gw:8000", "coder",
null, "BRIDGED_WORKER_TOKEN", List.of("claude"), "tab", "bridged-workers", "w #{n}", null, null, null);
@@ -796,14 +843,18 @@ class ClaudeCodeLauncherTest {
.toList();
}
private static BridgedConfig.Fleet fleet(Map<String, String> charters, String tabLabel) {
return new BridgedConfig.Fleet(Map.of(), Map.of(), Map.of(), Map.of(), charters, tabLabel);
}
/** A profile with no {@code tabLabel:} of its own — the fleet template decides. */
private ClaudeCodeLauncher labelService(FakeHerdr herdr, Supplier<String> fleetTemplate) {
private ClaudeCodeLauncher labelService(FakeHerdr herdr, Supplier<BridgedConfig.Fleet> fleet) {
BridgedConfig.Profile cfg = new BridgedConfig.Profile(
"sonnet", "http://gx00.gw:8000", "sonnet", null, "BRIDGED_WORKER_TOKEN",
List.of("claude"), "tab", "bridged-workers", null, null, null, null);
return new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
_ -> null, 0, 0L, fleetTemplate);
_ -> null, 0, 0L, fleet);
}
/**
@@ -814,7 +865,7 @@ class ClaudeCodeLauncherTest {
@Test
void theFleetTemplateNamesTheRoleTheMemberWasSpawnedFor() {
FakeHerdr herdr = new FakeHerdr();
ClaudeCodeLauncher svc = labelService(herdr, () -> "{role}: {profile} #{n}");
ClaudeCodeLauncher svc = labelService(herdr, () -> fleet(null, "{role}: {profile} #{n}"));
svc.spawn(new SpawnRequest("sonnet", null, null, null, null, MemberRole.REVIEWER));
@@ -825,7 +876,7 @@ class ClaudeCodeLauncherTest {
@Test
void theCounterRunsPerRoleAndProfileNotPerFleet() {
FakeHerdr herdr = new FakeHerdr();
ClaudeCodeLauncher svc = labelService(herdr, () -> "{role}: {profile} #{n}");
ClaudeCodeLauncher svc = labelService(herdr, () -> fleet(null, "{role}: {profile} #{n}"));
svc.spawn(new SpawnRequest("sonnet", null, null, null, null, MemberRole.DEV));
svc.spawn(new SpawnRequest("sonnet", null, null, null, null, MemberRole.REVIEWER));
@@ -839,7 +890,7 @@ class ClaudeCodeLauncherTest {
@Test
void aBlankFleetTemplateFallsBackToTheRoleFirstDefault() {
FakeHerdr herdr = new FakeHerdr();
labelService(herdr, () -> null).spawn(
labelService(herdr, () -> fleet(null, null)).spawn(
new SpawnRequest("sonnet", null, null, null, null, MemberRole.ARCHITECT));
assertEquals(List.of("architect: sonnet #1"), tabLabels(herdr));
@@ -855,7 +906,7 @@ class ClaudeCodeLauncherTest {
List.of("claude"), "tab", "bridged-workers", "pinned {profile}", null, null, null);
new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
_ -> null, 0, 0L, () -> "{role}: {profile} #{n}")
_ -> null, 0, 0L, () -> fleet(null, "{role}: {profile} #{n}"))
.spawn(new SpawnRequest("sonnet", null, null, null, null, MemberRole.REVIEWER));
assertEquals(List.of("pinned sonnet"), tabLabels(herdr));
@@ -870,7 +921,7 @@ class ClaudeCodeLauncherTest {
void theTemplateIsReadOnEverySpawnSoAnEditTakesEffect() {
FakeHerdr herdr = new FakeHerdr();
AtomicReference<String> template = new AtomicReference<>("{role}: {profile} #{n}");
ClaudeCodeLauncher svc = labelService(herdr, template::get);
ClaudeCodeLauncher svc = labelService(herdr, () -> fleet(null, template.get()));
svc.spawn(new SpawnRequest("sonnet", null, null, null, null, MemberRole.DEV));
template.set("[{profile}] {role} {n}");
@@ -77,7 +77,7 @@ class CompositePeerLauncherTest {
}
@Override
protected Launch buildLaunch(BridgedConfig.Profile cfg) {
protected Launch buildLaunch(BridgedConfig.Profile cfg, LaunchSpec spec) {
return new Launch(Map.of(), List.of());
}
@@ -0,0 +1,71 @@
package dev.ltms.bridged.member;
import dev.ltms.bridged.config.BridgedConfig;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.FakeHerdr;
import dev.ltms.bridged.herdr.WorkspaceControl;
import dev.ltms.bridged.peer.Capability;
import dev.ltms.bridged.peer.MemberRole;
import dev.ltms.bridged.peer.SpawnRequest;
import org.junit.jupiter.api.Test;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.Supplier;
import static org.junit.jupiter.api.Assertions.assertEquals;
class HerdrPeerLauncherCharterTest {
@Test
void readsAndComposesTheFleetCharterForEachSpawn() {
AtomicReference<BridgedConfig.Fleet> fleet = new AtomicReference<>(fleet(Map.of()));
CapturingLauncher launcher = new CapturingLauncher(fleet::get);
launcher.spawn(new SpawnRequest("mcp", null, null, null, null, MemberRole.DEV));
fleet.set(fleet(Map.of("dev", "role charter")));
launcher.spawn(new SpawnRequest("mcp", null, null, null, null, MemberRole.DEV));
launcher.spawn(new SpawnRequest("no-mcp", null, null, null, null, MemberRole.DEV));
assertEquals(HerdrPeerLauncher.REPLY_CHARTER, launcher.specs.get(0).charter(),
"without a role charter, MCP profiles receive only the reply charter");
assertEquals("role charter\n\n" + HerdrPeerLauncher.REPLY_CHARTER, launcher.specs.get(1).charter(),
"the changed supplier value is read for the next spawn and the reply rule is last");
assertEquals("role charter", launcher.specs.get(2).charter(),
"a role charter does not depend on an MCP mount");
}
private static BridgedConfig.Fleet fleet(Map<String, String> charters) {
return new BridgedConfig.Fleet(Map.of(), Map.of(), Map.of(), Map.of(), charters, null);
}
private static final class CapturingLauncher extends HerdrPeerLauncher {
private final List<LaunchSpec> specs = new ArrayList<>();
CapturingLauncher(Supplier<BridgedConfig.Fleet> fleet) {
super("test", new AgentControl(new FakeHerdr()), new WorkspaceControl(new FakeHerdr()),
Map.of("mcp", profile("mcp", "http://bridge"),
"no-mcp", profile("no-mcp", null)),
"mcp", _ -> null, 0, () -> 0L, () -> { }, fleet);
}
@Override
protected Launch buildLaunch(BridgedConfig.Profile cfg, LaunchSpec spec) {
specs.add(spec);
return new Launch(Map.of(), List.of("test"));
}
@Override
public Set<Capability> capabilities() {
return Set.of();
}
private static BridgedConfig.Profile profile(String name, String mcpUrl) {
return new BridgedConfig.Profile(name, "http://gx00.gw:8000", null, null,
"BRIDGED_WORKER_TOKEN", List.of("test"), "pane", null, null, mcpUrl, null, null);
}
}
}
@@ -1,6 +1,7 @@
package dev.ltms.bridged.rest;
import dev.ltms.bridged.auth.CallerResolver;
import dev.ltms.bridged.auth.MemberRegistry;
import dev.ltms.bridged.config.BridgedConfig;
import dev.ltms.bridged.guard.SubscriptionGuard;
import dev.ltms.bridged.herdr.AgentControl;
@@ -65,9 +66,8 @@ class BridgedAppAuthTest {
MessageService messages = new MessageService(agents, injector, new Rendezvous());
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> pid);
CallerResolver callers = tokenMode
? new CallerResolver(identity, true, token)
: new CallerResolver(identity);
CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, tokenMode, token,
Map::of, new MemberRegistry(null));
metrics = BridgedMetrics.create(sessions, new dev.ltms.bridged.msg.InMemoryReplyInbox());
app = new BridgedApp(herdr, workers, sessions, messages, sessions.asPresence(), null,
@@ -1,5 +1,9 @@
package dev.ltms.bridged.session;
import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.LoggerContext;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import dev.ltms.bridged.config.BridgedConfig;
import dev.ltms.bridged.guard.SubscriptionGuard;
import dev.ltms.bridged.herdr.AgentControl;
@@ -8,6 +12,7 @@ import dev.ltms.bridged.herdr.WorkspaceControl;
import dev.ltms.bridged.member.ClaudeCodeLauncher;
import dev.ltms.bridged.peer.PeerUnreachableException;
import org.junit.jupiter.api.Test;
import org.slf4j.LoggerFactory;
import java.util.List;
import java.util.Map;
@@ -158,32 +163,38 @@ class SessionManagerTest {
}
@Test
void recycleProducesNewPaneIdAndOldOneIsGone() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
MemberSession oldSession = sessions.acquire("ltms-local", null, "/caller", "term_primary");
String oldPane = oldSession.paneId();
String oldTerminal = oldSession.terminalId();
void onTurnFailedIsLoggedAtWarnWithThePriorState() {
// CB-564: this transition used to be a bare DEBUG "session marked failed" — a symptom with no
// cause. A member that can no longer be delegated to must be at least WARN, and should name
// what stage it failed at (here: BUSY, i.e. a turn was in flight and never resolved).
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
ch.qos.logback.classic.Logger sessionLog =
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(SessionManager.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.setContext(ctx);
appender.start();
sessionLog.addAppender(appender);
sessionLog.setLevel(Level.WARN);
try {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
String terminal = session.terminalId();
sessions.asPresence().markPresent(terminal);
sessions.onDelivered(terminal);
MemberSession fresh = sessions.recycle(oldPane);
sessions.onTurnFailed(terminal);
assertNotEquals(oldPane, fresh.paneId(), "recycle yields a new pane id");
assertNotEquals(oldTerminal, fresh.terminalId(), "recycle yields a new terminal id");
assertEquals(oldSession.profile(), fresh.profile(), "profile is preserved");
assertEquals(oldSession.cwd(), fresh.cwd(), "cwd is preserved");
assertEquals(oldSession.ownerTerminal(), fresh.ownerTerminal(), "owner is preserved");
assertTrue(sessions.get(oldPane).isEmpty(), "old pane is deregistered");
assertEquals(1, sessions.roster().size(), "only the fresh session remains");
assertEquals(fresh.paneId(), sessions.roster().getFirst().paneId());
// The old session was the first spawn → pane w9:pRoot_1 (CB-519: the registry key is the
// uuid id, so teardown is asserted on the real pane coordinate).
long paneCloseCount = herdr.calls.stream()
.filter(c -> "pane.close".equals(c.method()))
.filter(c -> "w9:pRoot_1".equals(((Map<?, ?>) c.params()).get("pane_id")))
.count();
assertEquals(1, paneCloseCount, "the old worker was torn down");
String warn = appender.list.stream()
.filter(e -> e.getLevel().equals(Level.WARN))
.map(ILoggingEvent::getFormattedMessage)
.findFirst()
.orElse("no turn-failed WARN logged");
assertTrue(warn.contains(terminal), "the log names the member: " + warn);
assertTrue(warn.contains("BUSY"), "the log names the stage it failed at: " + warn);
} finally {
sessionLog.detachAppender(appender);
}
}
@Test
@@ -1,11 +1,13 @@
package dev.ltms.bridged.session;
import dev.ltms.bridged.auth.MemberRegistry;
import dev.ltms.bridged.config.BridgedConfig;
import dev.ltms.bridged.guard.SubscriptionGuard;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.FakeHerdr;
import dev.ltms.bridged.herdr.WorkspaceControl;
import dev.ltms.bridged.member.ClaudeCodeLauncher;
import dev.ltms.bridged.peer.MemberRole;
import org.junit.jupiter.api.Test;
import java.util.List;
@@ -22,6 +24,13 @@ import static org.junit.jupiter.api.Assertions.*;
*/
class WorktreeSessionManagerTest {
private static MemberRegistry members() {
return new MemberRegistry(new BridgedConfig.Fleet(Map.of(),
Map.of("architect", new BridgedConfig.Slot("ltms-local")),
Map.of("dev", new BridgedConfig.Slot("ltms-local")),
Map.of("reviewer", new BridgedConfig.Slot("ltms-local")), null));
}
private static ClaudeCodeLauncher workerService(FakeHerdr herdr) {
BridgedConfig.Profile cfg = new BridgedConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "BRIDGED_WORKER_TOKEN",
@@ -56,6 +65,36 @@ class WorktreeSessionManagerTest {
assertEquals("/caller/proj", startCwd(herdr), "spawn receives the caller's cwd");
}
@Test
void onlyArchitectsBindAndReleaseMakesTheirSlotReusable() {
FakeHerdr herdr = new FakeHerdr();
MemberRegistry members = members();
SessionManager sessions = new SessionManager(workerService(herdr), new FakeWorktrees());
sessions.setMemberLifecycle(members);
MemberSession architect = sessions.acquire("ltms-local", MemberRole.ARCHITECT,
null, "/caller/proj", null, null);
MemberSession dev = sessions.acquire("ltms-local", MemberRole.DEV,
null, "/caller/proj", null, null);
MemberSession reviewer = sessions.acquire("ltms-local", MemberRole.REVIEWER,
null, "/caller/proj", null, null);
assertEquals("architect:architect", members.slotForTerminal(architect.terminalId()));
assertNull(members.slotForTerminal(dev.terminalId()), "a dev must never receive architect rights");
assertNull(members.slotForTerminal(reviewer.terminalId()),
"a reviewer must never receive architect rights");
sessions.release(architect.paneId());
MemberSession replacement = sessions.acquire("ltms-local", MemberRole.ARCHITECT,
null, "/caller/proj", null, null);
assertEquals("architect:architect", members.slotForTerminal(replacement.terminalId()));
MemberSession overflow = sessions.acquire("ltms-local", MemberRole.ARCHITECT,
null, "/caller/proj", null, null);
assertNull(members.slotForTerminal(overflow.terminalId()),
"a full slot pool must not stop the architect spawn");
}
@Test
void worktreeAcquireProvisionsAndRecordsPathAndBranch() {
FakeHerdr herdr = new FakeHerdr();
@@ -79,6 +118,19 @@ class WorktreeSessionManagerTest {
assertEquals(expectedPath, s.cwd(), "session cwd is the worktree path");
}
@Test
void worktreeArchitectAcquireAlsoBindsItsSlot() {
FakeHerdr herdr = new FakeHerdr();
MemberRegistry members = members();
SessionManager sessions = new SessionManager(workerService(herdr), new FakeWorktrees());
sessions.setMemberLifecycle(members);
MemberSession architect = sessions.acquire("ltms-local", MemberRole.ARCHITECT,
null, "/caller/proj", null, new WorktreeRequest("cb-548", null));
assertEquals("architect:architect", members.slotForTerminal(architect.terminalId()));
}
@Test
void worktreeAcquireRunsParityOverlayWithProfileDefaults() {
FakeHerdr herdr = new FakeHerdr();
@@ -94,12 +146,12 @@ class WorktreeSessionManagerTest {
FakeWorktrees.OverlayCall overlay = worktrees.lastOverlay();
assertNotNull(overlay);
assertEquals("/repo", overlay.repoRoot());
assertEquals(List.of(".claude/settings.local.json", ".env", ".envrc"),
assertEquals(List.of(".env", ".envrc"),
overlay.requested(), "default parity overlay is used when unset");
assertFalse(overlay.requested().contains(".mcp.json"),
"CB-525: replicating the primary's MCP config gives a worker the primary's IDE "
+ "servers, which navigate its edits out of its own worktree");
assertEquals(List.of(".claude/settings.local.json", ".envrc"), overlay.copied(),
assertEquals(List.of(".envrc"), overlay.copied(),
"existing paths are copied; missing paths are skipped");
assertEquals(List.of(".envrc"), overlay.skipWorktree(),
"tracked copied paths are --skip-worktree'd");
+2 -9
View File
@@ -35,7 +35,7 @@ build on.
- **No checkpoint content.** Writing `STATE.md` + commit on teardown is CB-302; CB-301 only exposes
the release hook it will attach to.
"Recycle" under no-reuse is simply **release + fresh acquire** — a helper, not a pool operation.
Under no-reuse, a released session is terminal. A new `acquire` always creates a fresh session.
## Design
@@ -46,11 +46,6 @@ ownership on top.
**Package:** new `dev.ltms.bridged.session` — keeps the registry/lifecycle concern separate from
the `worker` spawn mechanics. Holds `SessionManager` + `WorkerSession`.
**`recycle` is IN SCOPE for CB-301** (decided): implement `recycle(paneId, …)` = `release` the old
session then `acquire` a fresh one, asserting a new distinct paneId (the no-reuse invariant). It is
a thin convenience over the two primitives, shipped now so the no-reuse teardown+respawn path is
covered by a test from day one.
### `WorkerSession` (record or small mutable holder)
| Field | Source | Notes |
@@ -88,7 +83,6 @@ SPAWNING|READY|BUSY|DONE --vanished/drop--> FAILED
final class SessionManager {
WorkerSession acquire(String profile, String requestedCwd, String callerCwd, String ownerTerminal);
void release(String paneId); // deterministic teardown + deregister
WorkerSession recycle(String paneId, ...); // release + acquire (no-reuse convenience)
Optional<WorkerSession> get(String paneId);
List<WorkerSession> roster(); // bridge-owned view (CB-304 consumes this)
// lifecycle hooks (package-private): onReady/onDelivered/onComplete/onFailed(target)
@@ -120,8 +114,7 @@ final class SessionManager {
3. `release` tears the worker down via `WorkerService.stop` and removes it from `roster()`;
a second `release` on the same paneId is a harmless no-op.
4. `onTurnFailed` / drop moves the session to `FAILED` and it is absent from the live roster.
5. `recycle` produces a new paneId and the old one is gone (no-reuse invariant).
6. `roster()` reflects exactly the sessions acquired-minus-released, joined with live status.
5. `roster()` reflects exactly the sessions acquired-minus-released, joined with live status.
## Seams left open (deliberately)