Compare commits
26 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| d9e45f8c2a | |||
| 890190263e | |||
| f8522edacd | |||
| b958747855 | |||
| 078bde2c02 | |||
| 0b10ea987b | |||
| 1553d38182 | |||
| d04b075996 | |||
| 644927636d | |||
| 95e45007aa | |||
| 7a583c4045 | |||
| 65bce058bb | |||
| 48b437083b | |||
| fa0612859b | |||
| 4ffbcd0b7d | |||
| a0cd053fd9 | |||
| f9a5e066b5 | |||
| e9bc192160 | |||
| e0a57988ad | |||
| b414a74c26 | |||
| 2c3796d598 | |||
| 7e0ff9ab06 | |||
| 6d1565d8b6 | |||
| e18e002d2f | |||
| c18572ea9d | |||
| e9d02bc5e8 |
@@ -92,24 +92,32 @@ route by who launches opencode:
|
||||
|
||||
| Launcher | Route |
|
||||
|---|---|
|
||||
| a human, from a terminal | `{file:…}` — see below |
|
||||
| a human, from a terminal | one central store, sourced by the login shell |
|
||||
| a spawner (bridge, CI, IDE) | `{env:…}`, with the spawner injecting the variable |
|
||||
|
||||
For the human case prefer **`{file:…}` with a workspace-relative path**, kept in a gitignored
|
||||
`.secrets/` directory beside `opencode.json`:
|
||||
**For the human case, keep every credential in one file the login shell sources.** Here that file
|
||||
is `${SHARED_ENV}/tools/secrets.sh`, sourced from `${SHARED_ENV}/.ltms`, kept at mode 600 and never
|
||||
committed. `opencode.json` then names variables and holds no values:
|
||||
|
||||
```json
|
||||
"headers": { "Authorization": "Bearer {file:.secrets/api-token}" }
|
||||
"headers": { "Authorization": "Bearer {env:CONTEXT7_TOKEN}" }
|
||||
```
|
||||
|
||||
Verified: opencode resolves relative `{file:}` paths against the project root, so this needs no
|
||||
shell setup at all — no rc export leaking the secret to every process, no direnv dependency.
|
||||
Both routes end at the same syntax, and that is the point. The file does not change when a human
|
||||
launches opencode instead of the bridge.
|
||||
|
||||
**The catch, and state it out loud:** `.secrets/` is gitignored, so a peer running in a git worktree
|
||||
does **not** get it — worktrees receive tracked files only, the same rule that makes `opencode.json`
|
||||
itself worth committing. Spawned peers must therefore be fed through `{env:…}` by whatever launches
|
||||
them. Check the variable *names* match: a spawner often injects under a different name than your
|
||||
shell uses, and the config has no fallback.
|
||||
**`{file:…}` also works, and this project moved away from it.** Opencode resolves a relative
|
||||
`{file:}` path against the project root, so a gitignored `.secrets/` beside `opencode.json` needs no
|
||||
shell setup at all. It did not fail; the problem is that it makes a second copy of the token. The
|
||||
same secret then lives in two places, and the copy you forget is the one that leaks or goes stale.
|
||||
One store with many references is easier to rotate and to audit.
|
||||
|
||||
**One catch survives either choice, so state it out loud:** what a human's shell exports does not
|
||||
reach a spawned peer, and neither does a gitignored `.secrets/` — a git worktree receives tracked
|
||||
files only. Spawned peers must be fed through `{env:…}` by whatever launches them. Check the
|
||||
variable *names* match: a spawner often injects under a different name than your shell uses, and
|
||||
the config has no fallback. In this repo the bridge goes further: it neutralizes a worktree's
|
||||
`opencode.json`, so a member cannot inherit the primary's credentials by accident.
|
||||
|
||||
## 4. Do not port machine-local MCP servers
|
||||
|
||||
|
||||
+3
-2
@@ -1,5 +1,6 @@
|
||||
# Workspace-scoped secrets for tools that read no settings cascade of their own
|
||||
# (opencode resolves these via {file:.secrets/...} in opencode.json).
|
||||
# No secret belongs in this repo any more: every credential lives in one shell-level store
|
||||
# (${SHARED_ENV}/tools/secrets.sh), and opencode.json reads it as {env:...}. This line stays as a
|
||||
# backstop, so a workspace-scoped copy that someone re-creates by habit still cannot be committed.
|
||||
.secrets/
|
||||
|
||||
# Settings backups inherit the env block — and secrets with it.
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -231,14 +231,18 @@ 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 / model /
|
||||
# tabLabel.
|
||||
# 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
|
||||
# until you restart: `lifecycle:`, `leadHeartbeat:`, `guard:`, `worktreeRoot:`,
|
||||
# `spawnReadyTimeoutMs` / `spawnReadyPollMs`, and ADDING or REMOVING a profile (a new
|
||||
# backend needs its own launcher, and launchers are built once). The reload logs
|
||||
# these by name rather than pretending they applied.
|
||||
# `spawnReadyTimeoutMs` / `spawnReadyPollMs`, ADDING or REMOVING a profile (a new
|
||||
# backend needs its own launcher, and launchers are built once), AND an existing
|
||||
# profile's launch settings — model, baseUrl, argv, env, configDir, mcpUrl, tabLabel.
|
||||
# The launcher takes a copy of `profiles:` at startup and resolves every spawn out of
|
||||
# that copy, so those never reach a launch until you restart. The reload logs them by
|
||||
# name rather than pretending they applied.
|
||||
# COLD → cannot change at all: `bind:`, `herdrSocket:`, `broker:` and `auth:`. The socket is
|
||||
# bound, the broker connection is open, and the auth mode decides who may reach the
|
||||
# port that is already listening.
|
||||
@@ -270,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();
|
||||
@@ -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,16 @@ public record BridgedConfig(
|
||||
architects = unmodifiableOrEmpty(architects);
|
||||
developers = unmodifiableOrEmpty(developers);
|
||||
reviewers = unmodifiableOrEmpty(reviewers);
|
||||
charters = unmodifiableOrEmpty(charters);
|
||||
tabLabel = (tabLabel == null || tabLabel.isBlank()) ? DEFAULT_TAB_LABEL : tabLabel;
|
||||
}
|
||||
|
||||
/** Convenience constructor for code that does not configure launch charters. */
|
||||
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 +591,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 +1066,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 +1204,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 = java.util.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.
|
||||
*
|
||||
|
||||
@@ -7,6 +7,7 @@ import java.nio.file.Path;
|
||||
import java.util.ArrayList;
|
||||
import java.util.LinkedHashSet;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
@@ -26,14 +27,19 @@ 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} / {@code model} / {@code tabLabel}.</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:},
|
||||
* {@code worktreeRoot:}, and <em>adding or removing</em> a profile (a new backend needs its
|
||||
* own launcher, which is constructed once). A reload logs these rather than pretending they
|
||||
* applied.</li>
|
||||
* {@code worktreeRoot:}, adding or removing a profile (a new backend needs its own launcher,
|
||||
* which is constructed once), <em>and an existing profile's launch settings</em> —
|
||||
* {@code model}, {@code baseUrl}, {@code argv}, {@code env}, {@code mcpUrl} and the rest.
|
||||
* {@code HerdrPeerLauncher} takes {@code Map.copyOf(profiles)} at construction and resolves
|
||||
* each spawn out of that copy, so those never reach a launch until the daemon restarts. A
|
||||
* reload logs these rather than pretending they applied.</li>
|
||||
* <li><strong>Cold</strong> — cannot change at all under a running daemon: {@code bind:},
|
||||
* {@code herdrSocket:}, {@code broker:} and {@code auth:}. The socket is bound, the broker
|
||||
* connection is open, and the auth mode decides who may reach the port that is already
|
||||
@@ -144,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();
|
||||
@@ -204,16 +211,60 @@ public final class ConfigRef implements Supplier<BridgedConfig> {
|
||||
|| !Objects.equals(old.spawnReadyPollMs(), fresh.spawnReadyPollMs())) {
|
||||
changed.add("spawnReady*");
|
||||
}
|
||||
// Only the profile SET is deferred: a new backend needs a launcher, and launchers are built
|
||||
// once at startup. An existing profile's fields are read per spawn and so are hot.
|
||||
Set<String> before = old.profiles() == null ? Set.of() : old.profiles().keySet();
|
||||
Set<String> after = fresh.profiles() == null ? Set.of() : fresh.profiles().keySet();
|
||||
if (!before.equals(after)) {
|
||||
Set<String> diff = new LinkedHashSet<>(before);
|
||||
diff.addAll(after);
|
||||
diff.removeIf(p -> before.contains(p) && after.contains(p));
|
||||
Map<String, BridgedConfig.Profile> before =
|
||||
old.profiles() == null ? Map.of() : old.profiles();
|
||||
Map<String, BridgedConfig.Profile> after =
|
||||
fresh.profiles() == null ? Map.of() : fresh.profiles();
|
||||
// Adding or removing a profile is deferred: a new backend needs its own launcher, and
|
||||
// launchers are built once at startup.
|
||||
if (!before.keySet().equals(after.keySet())) {
|
||||
Set<String> diff = new LinkedHashSet<>(before.keySet());
|
||||
diff.addAll(after.keySet());
|
||||
diff.removeIf(p -> before.containsKey(p) && after.containsKey(p));
|
||||
changed.add("profiles (added/removed: " + String.join(", ", diff) + ")");
|
||||
}
|
||||
// An EXISTING profile's launch settings are deferred too, and this is easy to get wrong:
|
||||
// `HerdrPeerLauncher` takes `Map.copyOf(profiles)` at construction and `spawn` resolves the
|
||||
// profile out of that snapshot, so a reloaded model/baseUrl/argv/env never reaches a launch.
|
||||
// Only weight and maxLoad are genuinely hot, because placement reads them through the
|
||||
// supplier on the composite rather than from the adapter's copy. Without this check a
|
||||
// changed model would report "config reloaded" and silently do nothing — the worst outcome
|
||||
// a reload can produce, because the operator has no reason to doubt it.
|
||||
List<String> relaunch = new ArrayList<>();
|
||||
before.forEach((name, was) -> {
|
||||
BridgedConfig.Profile now = after.get(name);
|
||||
if (now != null && !sameLaunchSettings(was, now)) {
|
||||
relaunch.add(name);
|
||||
}
|
||||
});
|
||||
if (!relaunch.isEmpty()) {
|
||||
changed.add("profiles." + String.join("/", relaunch) + " launch settings "
|
||||
+ "(model, baseUrl, argv, env, …) — the launcher holds a startup snapshot");
|
||||
}
|
||||
return changed;
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether two versions of a profile would launch a peer identically. Compares every component
|
||||
* the launcher reads at spawn; {@code weight} and {@code maxLoad} are excluded because those are
|
||||
* read live by the placement policy and really do take effect on the next spawn.
|
||||
*/
|
||||
private static boolean sameLaunchSettings(BridgedConfig.Profile a, BridgedConfig.Profile b) {
|
||||
return Objects.equals(a.baseUrl(), b.baseUrl())
|
||||
&& Objects.equals(a.model(), b.model())
|
||||
&& Objects.equals(a.configDir(), b.configDir())
|
||||
&& Objects.equals(a.tokenEnv(), b.tokenEnv())
|
||||
&& Objects.equals(a.argv(), b.argv())
|
||||
&& Objects.equals(a.placement(), b.placement())
|
||||
&& Objects.equals(a.workspace(), b.workspace())
|
||||
&& Objects.equals(a.tabLabel(), b.tabLabel())
|
||||
&& Objects.equals(a.mcpUrl(), b.mcpUrl())
|
||||
&& Objects.equals(a.cwd(), b.cwd())
|
||||
&& Objects.equals(a.parityOverlay(), b.parityOverlay())
|
||||
&& Objects.equals(a.gitTokenEnv(), b.gitTokenEnv())
|
||||
&& Objects.equals(a.gitHostEnv(), b.gitHostEnv())
|
||||
&& Objects.equals(a.kind(), b.kind())
|
||||
&& Objects.equals(a.env(), b.env())
|
||||
&& Objects.equals(a.subscription(), b.subscription());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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();
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -35,10 +35,17 @@ public final class GitWorktrees implements Worktrees {
|
||||
/** What {@link #isolateToolSurface} writes for {@code .mcp.json}: a valid, explicitly empty server map. */
|
||||
private static final String NEUTRAL_MCP_CONFIG = "{\n \"mcpServers\": {}\n}\n";
|
||||
|
||||
/** OpenCode's repo-level config. Tracked here, so it lands in every worktree; it carries
|
||||
* {@code {file:.secrets/...}} references to gitignored secrets that never reach a worktree, and
|
||||
* opencode refuses to start on a dangling reference — so it is neutralized and the worker gets
|
||||
* only the config its launcher writes via {@code OPENCODE_CONFIG}. */
|
||||
/** OpenCode's repo-level config. Tracked here, so it lands in every worktree, and it mounts the
|
||||
* primary's gitea and context7 servers with the primary's credentials. Neutralized so the worker
|
||||
* gets only the config its launcher writes via {@code OPENCODE_CONFIG}.
|
||||
*
|
||||
* <p>The reason has changed shape and is now stronger. It used to be a crash: the file carried
|
||||
* {@code {file:.secrets/...}} references to gitignored files that never reached a worktree, and
|
||||
* opencode refuses to start on a dangling reference (CB-543). Those credentials now live in one
|
||||
* shell-level store and the file reads them as {@code {env:...}}, so in a worktree the reference
|
||||
* resolves instead of failing. That is worse, not better: a member would silently inherit the
|
||||
* primary's admin-scoped {@code GITEA_ACCESS_TOKEN}. A loud crash became a quiet privilege leak,
|
||||
* so this entry protects a boundary now rather than papering over a startup error. */
|
||||
private static final String OPENCODE_CONFIG = "opencode.json";
|
||||
|
||||
/** What {@link #isolateToolSurface} writes for {@code opencode.json}: a valid, empty JSON object. */
|
||||
@@ -112,8 +119,8 @@ public final class GitWorktrees implements Worktrees {
|
||||
* paths outside its own worktree. That is not hypothetical: a CB-523 worker made all 59 of its
|
||||
* edits in the primary checkout while compiling its worktree, so every build it ran was of code
|
||||
* that did not contain its changes. {@code opencode.json} is the same trap one tool over — tracked,
|
||||
* so it lands in every worktree, referencing gitignored {@code .secrets/} files that never do, and
|
||||
* opencode refuses to start on the dangling reference. {@code .autoenv} extends the principle to a
|
||||
* so it lands in every worktree, and it mounts gitea and context7 with the primary's own
|
||||
* credentials, which a member must never hold. {@code .autoenv} extends the principle to a
|
||||
* config that is not tracked today: autoenv authorizes by path, so a fresh worktree path is always
|
||||
* unauthorized and its interactive prompt would block every spawn, so re-landing one must be safe.
|
||||
*
|
||||
|
||||
@@ -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.
|
||||
@@ -213,9 +271,12 @@ class ConfigRefTest {
|
||||
assertTrue(out.deferred().getFirst().contains("haiku"), out.deferred().toString());
|
||||
}
|
||||
|
||||
/** Changing an existing profile's fields is hot — no launcher has to be rebuilt for it. */
|
||||
/**
|
||||
* Changing an existing profile's weight or maxLoad IS hot: placement reads those live through
|
||||
* the supplier on the composite, so the next spawn already sees them.
|
||||
*/
|
||||
@Test
|
||||
void changingAnExistingProfilesFieldsIsHot(@TempDir Path dir) throws Exception {
|
||||
void changingAProfilesWeightOrMaxLoadIsHot(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("bridged.yaml");
|
||||
Files.writeString(f, yaml(""));
|
||||
ConfigRef ref = refFor(f);
|
||||
@@ -228,8 +289,9 @@ class ConfigRefTest {
|
||||
profiles:
|
||||
sonnet:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
model: sonnet-4-5
|
||||
model: sonnet
|
||||
maxLoad: 7
|
||||
weight: 3.0
|
||||
guard:
|
||||
offSubscriptionHosts:
|
||||
- gx00.gw
|
||||
@@ -238,7 +300,42 @@ class ConfigRefTest {
|
||||
|
||||
assertTrue(out.applied());
|
||||
assertTrue(out.deferred().isEmpty(), out.deferred().toString());
|
||||
assertEquals("sonnet-4-5", ref.get().profiles().get("sonnet").model());
|
||||
assertEquals(7, ref.get().profiles().get("sonnet").maxLoad());
|
||||
}
|
||||
|
||||
/**
|
||||
* Changing an existing profile's MODEL is deferred, not hot — and saying so is the whole point.
|
||||
* HerdrPeerLauncher takes Map.copyOf(profiles) at construction and resolves every spawn out of
|
||||
* that copy, so a reloaded model never reaches a launch. Reporting it as applied would be the
|
||||
* worst outcome a reload can produce: the operator has no reason to doubt a clean "reloaded".
|
||||
*/
|
||||
@Test
|
||||
void changingAProfilesLaunchSettingsIsReportedAsDeferred(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("bridged.yaml");
|
||||
Files.writeString(f, yaml(""));
|
||||
ConfigRef ref = refFor(f);
|
||||
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
herdrSocket: ~/.config/herdr/herdr.sock
|
||||
profiles:
|
||||
sonnet:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
model: deepseek-v4-flash
|
||||
guard:
|
||||
offSubscriptionHosts:
|
||||
- gx00.gw
|
||||
""");
|
||||
ConfigRef.Outcome out = ref.reload();
|
||||
|
||||
assertTrue(out.applied());
|
||||
assertEquals(1, out.deferred().size(), out.deferred().toString());
|
||||
assertTrue(out.deferred().getFirst().contains("sonnet"), out.deferred().toString());
|
||||
assertTrue(out.deferred().getFirst().contains("launch settings"), out.deferred().toString());
|
||||
// The snapshot still carries the new value — a restart is what makes it take effect.
|
||||
assertEquals("deepseek-v4-flash", ref.get().profiles().get("sonnet").model());
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
@@ -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");
|
||||
|
||||
@@ -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");
|
||||
|
||||
@@ -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)
|
||||
|
||||
|
||||
+3
-3
@@ -14,7 +14,7 @@
|
||||
"url": "https://ct7.ltms.dev/mcp",
|
||||
"enabled": true,
|
||||
"headers": {
|
||||
"Authorization": "Bearer {file:.secrets/context7-token}"
|
||||
"Authorization": "Bearer {env:CONTEXT7_TOKEN}"
|
||||
}
|
||||
},
|
||||
"gitea": {
|
||||
@@ -26,8 +26,8 @@
|
||||
],
|
||||
"enabled": true,
|
||||
"environment": {
|
||||
"GITEA_ACCESS_TOKEN": "{file:.secrets/gitea-token}",
|
||||
"GITEA_HOST": "{file:.secrets/gitea-host}"
|
||||
"GITEA_ACCESS_TOKEN": "{env:GITEA_ACCESS_TOKEN}",
|
||||
"GITEA_HOST": "{env:GITEA_HOST}"
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+1
-1
Submodule wiki updated: e5424f4665...7c50cce52e
Reference in New Issue
Block a user