From bc13b8e92cfbecfe5f533fee045b315fee86d615 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 13 Aug 2026 09:38:24 +0200 Subject: [PATCH] CB-530..536, CB-538 groundwork: leads as peers, and a fleet that can find itself MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Two leads now work as peers rather than one primary plus workers. The arc: CB-530/531 lead identity: `leaders:` names panes, `leadScan:` discovers them by tab label (LeadTabScanner, TTL-cached, worker spaces excluded). CB-532 leads can message each other AND be answered. Principal.leader now carries its terminal, so ownsSession() can be true for a lead; the "and you must be a worker" conjunct beside it protected nothing. Retires `primary:` — reply nudges follow the delegating lead, a binding recorded at bridge_send where both halves are known. CB-533 ClaudeCodeLauncher passes --model. argv is usually a wrapper (`ccs `) that re-exports its own model family, so ANTHROPIC_MODEL alone was silently overruled. CB-534 a lead is deliverable. The CB-113 readiness gate only opened for terminals in WorkerPresence, which only workers ever enter, so every lead->lead send waited out the ~60s grace and failed having never been typed. The gate guards a *spawned* peer's boot window; a lead is never spawned. CB-535 bridge_list returns `leads` alongside `workers`, with `self` on the caller's row. An empty worker roster no longer reads as "no peers". CB-536 CLAUDE.md: lead<->lead is coordinate-only, never sideways delegation. Propagated byte-identically to wiki/7-Use-Cases.md. MIXED PROVENANCE — recorded deliberately rather than hidden. This tree also carries in-progress CB-537 (context separation) authored by the peer lead gpt-sol-5.6 and its worker: Capability.CONTEXT_RESET, SessionManager.clearAfterTurn, and the Injector/TurnListener/CompletionResolver/launcher changes around it. That work was done in this shared working tree rather than a worktree, and is entangled with the above in BridgedConfig.java, Bridged.java and ClaudeCodeLauncher.java, so neither lead could stage its own half without sweeping in the other's. Committing the whole green state is the honest resolution; the peer branches from here. Note for whoever picks CB-537 up: the design in this commit is SUPERSEDED. Both leads agreed to replace the global `clearAfterTurn` boolean with per-delivery policy (inherit|fresh|thread) applied PRE-delivery, because a post-turn reset races by construction — Injector.onStatus clears awaitingCompletion and dequeues the next message in the same tick. `fresh` is also a correctness guarantee, so an adapter without a reset capability must refuse it rather than log a no-op. mvn clean install: Tests run: 464, Failures: 0, Errors: 0, Skipped: 0. BUILD SUCCESS. --- .claude/skills/port-to-opencode/SKILL.md | 172 ++++++++++++ .gitignore | 12 + CLAUDE.md | 38 ++- bridged/bridged.example.yaml | 42 +++ .../main/java/dev/ltms/bridged/Bridged.java | 91 +++++- .../java/dev/ltms/bridged/auth/Authz.java | 8 +- .../dev/ltms/bridged/auth/CallerResolver.java | 121 +++++--- .../java/dev/ltms/bridged/auth/Principal.java | 54 +++- .../ltms/bridged/config/BridgedConfig.java | 205 +++++++++++++- .../ltms/bridged/herdr/LeadTabScanner.java | 164 +++++++++++ .../bridged/inject/CompletionResolver.java | 9 + .../dev/ltms/bridged/inject/Injector.java | 52 +++- .../dev/ltms/bridged/inject/TurnListener.java | 19 ++ .../java/dev/ltms/bridged/mcp/BridgeMcp.java | 124 +++++++-- .../dev/ltms/bridged/mcp/PrimaryRegistry.java | 47 ++++ .../dev/ltms/bridged/msg/ReplyPushLoop.java | 46 +-- .../dev/ltms/bridged/peer/Capability.java | 6 + .../dev/ltms/bridged/peer/PeerLauncher.java | 9 + .../ltms/bridged/session/SessionManager.java | 45 ++- .../bridged/worker/ClaudeCodeLauncher.java | 42 ++- .../bridged/worker/CompositePeerLauncher.java | 10 + .../bridged/worker/HerdrPeerLauncher.java | 21 ++ .../bridged/BridgedDeliverabilityTest.java | 86 ++++++ .../ltms/bridged/auth/CallerResolverTest.java | 175 +++++++++++- .../bridged/config/BridgedConfigTest.java | 232 ++++++++++++++++ .../bridged/herdr/LeadTabScannerTest.java | 262 ++++++++++++++++++ .../inject/CompletionResolverTest.java | 15 + .../dev/ltms/bridged/inject/InjectorTest.java | 41 +++ .../dev/ltms/bridged/mcp/BridgeMcpTest.java | 55 +++- .../ltms/bridged/mcp/PrimaryRegistryTest.java | 64 +++++ .../bridged/session/SessionManagerTest.java | 58 +++- .../worker/ClaudeCodeLauncherTest.java | 78 +++++- .../worker/CompositePeerLauncherTest.java | 28 ++ opencode.json | 34 +++ 34 files changed, 2350 insertions(+), 115 deletions(-) create mode 100644 .claude/skills/port-to-opencode/SKILL.md create mode 100644 .gitignore create mode 100644 bridged/src/main/java/dev/ltms/bridged/herdr/LeadTabScanner.java create mode 100644 bridged/src/test/java/dev/ltms/bridged/BridgedDeliverabilityTest.java create mode 100644 bridged/src/test/java/dev/ltms/bridged/herdr/LeadTabScannerTest.java create mode 100644 opencode.json diff --git a/.claude/skills/port-to-opencode/SKILL.md b/.claude/skills/port-to-opencode/SKILL.md new file mode 100644 index 0000000..8bd2783 --- /dev/null +++ b/.claude/skills/port-to-opencode/SKILL.md @@ -0,0 +1,172 @@ +--- +name: port-to-opencode +description: Make an OpenCode session a first-class participant in a Claude Code workspace — instructions, MCP servers, and secrets — without duplicating config. Use when onboarding opencode to a project that already has CLAUDE.md and .mcp.json, or when an opencode peer needs the same tools and rules as the Claude session. +--- + +# Porting a Claude Code workspace to OpenCode + +**The headline: there is almost nothing to port.** OpenCode reads `CLAUDE.md` natively. The only +artifact you create is one `opencode.json` mapping MCP servers. Do not translate instructions, do +not generate a second rules file, and do not install a sync tool — every one of those makes the +workspace worse. + +Everything below was verified against `opencode 1.18.16` and the OpenCode docs. + +## 1. Know what you get for free + +OpenCode's instruction search order: + +``` +1. walking up from cwd: AGENTS.md , then CLAUDE.md +2. global: ~/.config/opencode/AGENTS.md +3. Claude Code global: ~/.claude/CLAUDE.md (unless disabled) +``` + +*"The first matching file wins in each category."* + +Consequences that decide the whole procedure: + +- **A project `CLAUDE.md` is already read.** No port needed. +- **Your user-level `~/.claude/CLAUDE.md` is already read too.** Global preferences carry over. +- **An `AGENTS.md` in the repo SHADOWS `CLAUDE.md`.** If one exists from a previous Codex port, + **delete it** — otherwise opencode reads the stale translated copy instead of the real rules. + This is the single most likely way to get this wrong. + +## 2. Create `opencode.json` for MCP servers only + +Project config lives at `opencode.json` in the repo root; the global one is +`~/.config/opencode/opencode.json`. **Configs merge, they do not replace** — so machine-local +servers belong in the global file and shared ones in the project file. + +Map each entry from `.mcp.json`: + +| `.mcp.json` | `opencode.json` | +|---|---| +| `"type": "http"` / `"sse"` | `"type": "remote"`, `"url"` | +| `"type": "stdio"` | `"type": "local"`, `"command": ["bin", "arg"]` | +| `"command"` + `"args"` | single `"command"` array | +| `"env"` | `"environment"` | +| `"headers"` | `"headers"` | + +```json +{ + "$schema": "https://opencode.ai/config.json", + "instructions": ["CLAUDE.md"], + "mcp": { + "bridged": { "type": "remote", "url": "http://127.0.0.1:8765/mcp", "enabled": true }, + "context7": { "type": "remote", "url": "https://example.dev/mcp", "enabled": true, + "headers": { "Authorization": "Bearer {env:CONTEXT7_TOKEN}" } }, + "gitea": { "type": "local", "command": ["gitea-mcp", "-t", "stdio"], "enabled": true, + "environment": { "GITEA_ACCESS_TOKEN": "{env:GITEA_ACCESS_TOKEN}" } } + } +} +``` + +Set `instructions` explicitly even though `CLAUDE.md` is found anyway — the fallback only applies +while no `AGENTS.md` exists, and being explicit survives someone adding one later. + +## 3. Reference secrets, never embed them + +OpenCode substitutes at load time, in both `headers` and `environment`: + +``` +{env:VARIABLE_NAME} value from the environment +{file:~/.secrets/token} value read from a file +``` + +**No credential ever belongs in `opencode.json`.** With `{env:…}` there is no reason to write one, +which is what makes this file safe to commit — and it must be committable, because a peer running +in a git worktree receives tracked files only. + +**But `{env:…}` reads OpenCode's *process* environment — and OpenCode has no env store of its +own.** A Claude Code `env` block in `~/.claude/settings.json` or `.claude/settings.local.json` does +**not** reach it: those files are Claude Code's, and the variables exist only inside processes +Claude Code spawned. Verify with a clean login shell, not the shell your agent hands you: + +```bash +env -u MY_TOKEN zsh -lc 'echo "${MY_TOKEN:-NOT IN PROFILE}"' +``` + +So `{env:…}` only works if something puts the variable in the environment first. Pick the supply +route by who launches opencode: + +| Launcher | Route | +|---|---| +| a human, from a terminal | `{file:…}` — see below | +| 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`: + +```json +"headers": { "Authorization": "Bearer {file:.secrets/api-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. + +**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. + +## 4. Do not port machine-local MCP servers + +IDE indexes, language servers, editor bridges — anything bound to *your* checkout — stay out of the +project file. Put them in `~/.config/opencode/opencode.json` if you want them personally. + +A committed project config reaches every worktree. A worker that mounts servers whose paths point +into the primary's checkout will edit the primary's files while building in its own — every build +passes, every change lands in the wrong tree. + +Port the servers the work needs. Leave the rest. + +## 5. Verify against the running agent, not the file + +A file on disk proves nothing about what the agent loaded. + +```bash +opencode run "In one line: state a rule from this project's instructions." +``` + +The answer must reflect the actual `CLAUDE.md`. If it answers generically, the instructions did not +reach the model and everything after this is built on sand. + +Then confirm the tools are mounted: + +```bash +opencode mcp list +``` + +**"connected" does not mean "working".** A stdio server with a missing credential still completes +the MCP handshake and reports green; only a real tool call reveals it. Verified: `gitea` showed +`✓ connected` with no token, then failed the first call with `token is required`. A remote server +is more honest (`⚠ needs authentication`), but do not rely on that difference — **exercise one +authenticated tool per server**: + +```bash +opencode run "Call 's . Report the result or the exact error. One line." +``` + +Check too that no server you deliberately withheld is present. + +## 6. Report + +State what you changed, which servers crossed and which you withheld and why, and quote the +verification answer verbatim. If any server failed to connect, say so plainly — a partially mounted +peer is worse than a missing one, because it looks configured. + +## Gotchas + +- **A leftover `AGENTS.md` silently wins over `CLAUDE.md`.** Check for one before anything else. +- **`opencode.json` is merged, not overridden** — a global entry and a project entry with the same + server name both matter; keep names distinct unless you intend to layer them. +- **`enabled: false`** turns a server off without deleting its config — prefer it over removal when + you may want the server back. +- **OAuth-based servers** store tokens in `~/.local/share/opencode/mcp-auth.json` after + `opencode mcp auth `; that is machine state, never config to commit. +- **Skills and subagents do not port.** OpenCode uses its own agent markdown under + `.opencode/agents/`; `.claude/skills/**` is not read. If a delegation brief tells a peer to load a + skill by name, that instruction has no effect on an opencode peer — spell the procedure out in the + brief, or author the equivalent agent file. diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..95bc6ee --- /dev/null +++ b/.gitignore @@ -0,0 +1,12 @@ +# Workspace-scoped secrets for tools that read no settings cascade of their own +# (opencode resolves these via {file:.secrets/...} in opencode.json). +.secrets/ + +# Settings backups inherit the env block — and secrets with it. +.claude/settings.local.json.bak* + +# Daemon runtime artefacts. bridged appends its log wherever it is launched from, so both the +# repo root and bridged/ collect one; neither belongs in git. +bridged.out +bridged/bridged.out +logs/ diff --git a/CLAUDE.md b/CLAUDE.md index 2b10caf..e2f7ca7 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -41,8 +41,9 @@ nothing. Fail toward the recoverable error. 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 primary-only; reply/ask are - worker-only-and-only-as-itself. A call outside your role is refused, not queued. + 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. 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`. 5. **Never drive the terminal multiplexer directly** (no `herdr` CLI, no socket). The bridge owns @@ -101,13 +102,42 @@ 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` · one worker's state: `bridge_status{sessionId}` | +| See the fleet | `bridge_list` → `leads` (your peers) + `workers` · 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` | +| Message a **peer lead** | `bridge_send{sessionId: , 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}` | +### 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. + +**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 +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. +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 + own. Split by **context ownership** — whoever already holds the context owns that area — and say + who takes what, in one message, before either of you starts. Two leads silently working the same + unit is the failure mode here, and neither notices until the merge. +2. **Share findings, hazards, and corrections.** What you have already discovered, what broke, what + the next person will trip on. This is the traffic that actually pays for the channel: it costs one + message and saves a peer a rediscovery. +3. **Verify a peer exactly as you verify yourself.** Peer status buys nothing: check the claim + against the code, and re-run the build. A peer's correction gets the same treatment — right or + wrong on the evidence, not on who said it. Neither of you merges the other's work unreviewed. + +Being messaged by a peer does not make you its worker: answer with `bridge_reply`, and push back on +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 1. **Load the playbook skill the lead named** before doing anything else. @@ -143,6 +173,8 @@ charter, not here. `mcp/ConnectionIdentity` (connection→role), and `worker/*Launcher` (`REPLY_CHARTER`). - **Skills available to delegate:** `implementer` (worktree → commit → push → own PR) and `reviewer` (scoped review → one structured finding). Name one in every delegation. +- **Primary-side skills** (not delegation playbooks — a worker cannot use them): + `port-to-opencode` (make an OpenCode session a participant in this workspace). - **Never commit** `.mcp.json` (the primary's local copy, flagged `--skip-worktree`) or `wiki/` (a submodule with its own remote). - **Flows and the error model** — rendezvous, `bridge_ask`, detached delivery, turn-done fallback — diff --git a/bridged/bridged.example.yaml b/bridged/bridged.example.yaml index c11ac06..a7e0927 100644 --- a/bridged/bridged.example.yaml +++ b/bridged/bridged.example.yaml @@ -42,6 +42,48 @@ bind: # pushReminders: 5 # max nudges before giving up (default 5) # pushBackoffMs: 15000 # delay between nudges (default 15000) +# CB-530: MORE THAN ONE LEAD. `primary:` above is singular by construction — every other pane +# resolves as a worker — which is right for one lead driving a fleet and wrong the moment two leads +# (say a Claude lead and an opencode lead) work as peers: the second is silently demoted and refused +# every orchestration call. List each lead's pane here and all of them resolve as leads. +# +# terminal → the ONLY field identity depends on; get it from that session's bridge_whoami +# kind/model → descriptive; they document what runs in the pane and are echoed by bridge_whoami +# +# A lead is never spawned — it pre-exists, which is exactly why it must be named rather than created. +# `bridge_whoami` reports `{"role":"primary","leader":""}`; role stays "primary" because a lead +# IS a primary for authorization, so nothing that keys on the role breaks. +# +# KEEP `primary:` when adding leads: it still addresses the CB-307 push loop, which needs a single +# destination for its nudges. If both name the same terminal, the `leaders:` entry wins. +# leaders: +# opus-5.0: +# terminal: term_0123456789abcd +# kind: claude +# gpt-sol-5.6: +# terminal: term_fedcba9876543 +# kind: opencode +# model: openai/gpt-5.6-terra + +# CB-531: FIND LEADS BY TAB NAME instead of pasting terminal ids. `leaders:` above needs an id that +# only exists once the session is running, so adding a lead is: open a tab, start the agent, ask it +# bridge_whoami, edit this file, restart the daemon. This block replaces all of that with a naming +# convention — label the tab `lead: ` when you open it and the pane is recognised on the next +# rescan, with no config edit and no restart. Reopen the tab later and the id changes; the label +# does not. +# +# bridged NEVER writes these labels. It renames worker tabs (see `tabLabel` below) but reads lead +# tabs read-only, so what is in the tab bar is always what you typed. Two things keep the convention +# from being a way to claim leadership: the configured worker spaces are excluded from the scan, so +# nothing bridged places can land in a matching tab; and startup REFUSES a `tabPrefix` that any +# worker `tabLabel` also matches, so the two namespaces cannot overlap by accident. +# +# Opt-in on purpose — this widens who resolves as a lead, so upgrading the daemon must never switch +# it on for you. Absent block = leads come only from `leaders:`/`primary:`, exactly as before. +# leadScan: +# tabPrefix: "lead:" # `lead: opus-5.0` ⇒ a lead named opus-5.0 (case-insensitive; default "lead:") +# intervalSeconds: 10 # rescan cadence, and the worst case before a new tab is recognised + # herdr Unix socket. Omit to use the client default # (${HERDR_SOCKET_PATH:-~/.config/herdr/herdr.sock}). herdrSocket: ~/.config/herdr/herdr.sock diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java index be9bfe1..38394b4 100644 --- a/bridged/src/main/java/dev/ltms/bridged/Bridged.java +++ b/bridged/src/main/java/dev/ltms/bridged/Bridged.java @@ -5,6 +5,7 @@ import dev.ltms.bridged.guard.SubscriptionGuard; import dev.ltms.bridged.herdr.AgentControl; import dev.ltms.bridged.herdr.HerdrClient; import dev.ltms.bridged.herdr.HerdrException; +import dev.ltms.bridged.herdr.LeadTabScanner; import dev.ltms.bridged.herdr.PaneLocator; import dev.ltms.bridged.herdr.UnixSocketHerdrClient; import dev.ltms.bridged.herdr.WorkspaceControl; @@ -46,9 +47,15 @@ import java.util.ArrayList; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; +import java.util.Objects; +import java.util.Set; import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicReference; import java.util.function.Function; +import java.util.function.Predicate; +import java.util.function.Supplier; +import java.util.stream.Collectors; /** * {@code bridged} entry point. Wires the real herdr socket client to the REST app and @@ -79,6 +86,7 @@ public final class Bridged { // refuses remote connections to a loopback socket. This throws rather than warns so the // dangerous configuration cannot be reached by ignoring a log line. cfg.validateAuthExposure(); + cfg.validateLeadScan(); Path socket = cfg.herdrSocket() != null && !cfg.herdrSocket().isBlank() ? Path.of(cfg.herdrSocket()) @@ -114,7 +122,7 @@ public final class Bridged { opencodeProfiles, cfg.defaultProfile(), System::getenv, cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs())); } - AtomicReference> liveCountRef = new AtomicReference<>(name -> 0); + AtomicReference> liveCountRef = new AtomicReference<>(_ -> 0); PeerLauncher workers = new CompositePeerLauncher( adapters, cfg.defaultProfile(), @@ -144,7 +152,9 @@ public final class Bridged { && cfg.lifecycle().contextCap() > 0) { contextCap = cfg.lifecycle().contextCap(); } - SessionManager sessions = new SessionManager(workers, new GitWorktrees(cfg.worktreeRoot()), contextCap); + boolean clearAfterTurn = cfg.lifecycle() != null && cfg.lifecycle().clearAfterTurn(); + SessionManager sessions = new SessionManager(workers, new GitWorktrees(cfg.worktreeRoot()), + System::nanoTime, contextCap, clearAfterTurn); liveCountRef.set(profileName -> (int) sessions.roster().stream() .filter(s -> profileName.equals(s.profile())) .count()); @@ -160,6 +170,32 @@ public final class Bridged { reaper = null; } + // CB-530: every pane the config names as a lead, merged from `leaders:` and the legacy + // singular pin. PrimaryRegistry below still tracks ONE terminal — it addresses the push + // loop's nudges, which need a single destination — so it keeps the legacy pin. + Map leadTerminals = cfg.leaderTerminals(); + if (leadTerminals.size() > 1) { + log.info("leads: {} panes recognised {}", leadTerminals.size(), leadTerminals.values()); + } + // CB-531: on top of the static registry, discover leads by the tab labels the operator + // writes. Opt-in, so a config with no `leadScan:` block resolves exactly as it did under + // CB-530 — the supplier is then a constant and never touches herdr. + final Supplier> leads; + if (cfg.leadScan() != null) { + var scan = cfg.leadScan(); + Set workerSpaces = cfg.workerProfiles().values().stream() + .map(BridgedConfig.Worker::workspace) + .filter(Objects::nonNull) + .collect(Collectors.toSet()); + leads = new LeadTabScanner(herdr, scan.tabPrefix(), workerSpaces, leadTerminals, + TimeUnit.SECONDS.toNanos(scan.intervalSeconds()), System::nanoTime); + log.info("lead scan: tabs labelled '{}…' host a lead (rescan every {}s, worker spaces {} " + + "excluded)", + scan.tabPrefix(), scan.intervalSeconds(), workerSpaces); + } else { + leads = () -> leadTerminals; + } + // Status-gated injector (CB-103): the single writer into workers, fed by a poller. // The blocking message endpoint (CB-104) is the producer; the poller is inert until then. // CB-106: a confirmed turn completion resolves a blocked send whose worker never replied. @@ -175,6 +211,17 @@ public final class Bridged { sessions.onTurnComplete(target); } + @Override + public boolean hasPostTurnAction(String target) { + return sessions.hasPostTurnAction(target); + } + + @Override + public boolean onTurnCompleteWithPostAction(String target) { + completion.resolveBeforePostAction(target); + return sessions.onTurnCompleteWithPostAction(target); + } + @Override public void onDelivered(String target) { completion.onDelivered(target); @@ -187,7 +234,8 @@ public final class Bridged { sessions.onTurnFailed(target); } }; - Injector injector = new Injector(agents, turnListener, presence::isPresent, presence::forget); + Injector injector = new Injector(agents, turnListener, deliverableTo(presence, leads), + presence::forget); StatusPoller poller = new StatusPoller(agents, injector, INJECT_POLL_MILLIS); poller.start(); @@ -207,7 +255,15 @@ public final class Bridged { // otherwise resolve as a worker and be refused every orchestration tool. String pinnedPrimaryTerminal = cfg.primary() != null ? cfg.primary().terminal() : null; PrimaryRegistry primaryRegistry = new PrimaryRegistry(pinnedPrimaryTerminal); - + // CB-532: `primary.terminal` is superseded and no longer needed for either of its jobs — + // identity comes from `leaders:`/`leadScan:`, and reply nudges now follow the delegating + // lead. Say so once at startup rather than leaving a redundant pin to look load-bearing. + if (pinnedPrimaryTerminal != null && !pinnedPrimaryTerminal.isBlank()) { + log.warn("primary.terminal is DEPRECATED (CB-532) and can be deleted: identity now comes " + + "from leaders:/leadScan:, and reply nudges follow the lead that delegated. " + + "It still works, and is still the fallback nudge destination when a restart " + + "has lost the delegation map. Its pushReminders/pushBackoffMs stay valid."); + } // CB-307: active push-to-primary loop — nudge the primary when replies land without an // open bridge_send. Uses its own lightweight scheduled executor, separate from the injector. int maxReminders = cfg.primary() != null ? cfg.primary().remindersOrDefault() : 5; @@ -232,6 +288,7 @@ public final class Bridged { sessions.onRelease(terminal -> { messages.abandon(terminal, "the worker session was released before it replied"); replyInbox.release(terminal); + primaryRegistry.forgetDelegation(terminal); // CB-532: don't leak the lead binding }); // MCP server face (CB-105): bridge_send/bridge_reply/bridge_status, mounted at /mcp. @@ -248,11 +305,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 = new CallerResolver(identity, true, token, pinnedPrimaryTerminal); + callers = CallerResolver.withLeads(identity, true, token, leads); log.info("auth: token mode (bearer required for non-worker callers, env {})", cfg.auth().tokenEnv()); } else { - callers = new CallerResolver(identity, false, null, pinnedPrimaryTerminal); + callers = CallerResolver.withLeads(identity, false, null, leads); log.info("auth: loopback-trust (any loopback non-worker caller is the primary)"); } @@ -287,6 +344,28 @@ public final class Bridged { cfg.bind().host(), cfg.bind().port(), socket); } + /** + * The {@link Injector}'s readiness gate (CB-534): a target is deliverable if it is a worker whose + * agent has connected the bridge MCP, or a lead. + * + *

The gate exists for one reason — to hold a delivery out of a spawned worker'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. + * + *

A lead is also never enrolled in {@link WorkerPresence} — {@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. + * + *

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. + */ + static Predicate deliverableTo(WorkerPresence presence, Supplier> leads) { + return target -> presence.isPresent(target) || leads.get().containsKey(target); + } + /** * Poll herdr's {@code ping} until it answers or {@link #HERDR_WAIT_SECONDS} elapses (CB-504). * diff --git a/bridged/src/main/java/dev/ltms/bridged/auth/Authz.java b/bridged/src/main/java/dev/ltms/bridged/auth/Authz.java index f23cd69..0f96310 100644 --- a/bridged/src/main/java/dev/ltms/bridged/auth/Authz.java +++ b/bridged/src/main/java/dev/ltms/bridged/auth/Authz.java @@ -50,9 +50,11 @@ public final class Authz { // worker escalating into the orchestrator role. case SPAWN, STOP, SEND, DRAIN -> caller.isPrimary(); - // The load-bearing rule: a worker acts only as itself. The primary is deliberately - // excluded — a reply/ask is a worker's own turn output, and letting the primary forge - // one would corrupt the rendezvous correlation it is itself waiting on. + // The load-bearing rule: a caller acts only as the pane it occupies. CB-532 widened who + // that can be — a lead answering another lead is replying for its OWN terminal, which + // this already permits — while the rule itself is unchanged, and is what stops anyone + // forging a reply for a rendezvous someone else is waiting on. An unnamed primary + // (token/loopback, no pane) owns nothing and is still excluded. case REPLY, ASK -> caller.ownsSession(targetSession); // Observation is open to both authenticated roles: a worker legitimately polls its own diff --git a/bridged/src/main/java/dev/ltms/bridged/auth/CallerResolver.java b/bridged/src/main/java/dev/ltms/bridged/auth/CallerResolver.java index e986c4f..5f10789 100644 --- a/bridged/src/main/java/dev/ltms/bridged/auth/CallerResolver.java +++ b/bridged/src/main/java/dev/ltms/bridged/auth/CallerResolver.java @@ -4,6 +4,8 @@ import dev.ltms.bridged.mcp.ConnectionIdentity; import java.nio.charset.StandardCharsets; import java.security.MessageDigest; +import java.util.Map; +import java.util.function.Supplier; /** * Resolves every caller to a {@link Principal}, for both entry paths into the core (CB-501). @@ -15,10 +17,13 @@ import java.security.MessageDigest; * *

Resolution order — connection identity first, token second, nothing third: *

    - *
  1. A loopback peer PID that maps to the pinned {@code primary.terminal} pane (CB-307) ⇒ - * {@link Role#PRIMARY}. The pane mapping is as unforgeable as a worker's, and the config - * explicitly names that pane as the primary's own — without this rule a primary running - * inside a herdr pane is misread as a worker and locked out of orchestration.
  2. + *
  3. A loopback peer PID that maps to a pane named by {@code leaders:}, by the legacy + * {@code primary.terminal} pin, or by an operator-labelled lead tab (CB-307, CB-530, CB-531) + * ⇒ {@link Role#PRIMARY}, carrying that lead's + * name. The pane mapping is as unforgeable as a worker's, and the config explicitly names + * that pane as a lead's own — without this rule a lead running inside a herdr pane + * is misread as a worker and locked out of orchestration. More than one pane may be named, + * so two leads can work as peers rather than one being demoted.
  4. *
  5. A loopback peer PID that maps to any other herdr pane ⇒ {@link Role#WORKER}. This is * unforgeable (the OS reports the PID, herdr owns the PID→pane map) and is honoured * regardless of auth mode, so enabling auth never breaks the fleet.
  6. @@ -33,29 +38,79 @@ public final class CallerResolver { private final ConnectionIdentity identity; private final boolean tokenMode; private final byte[] expectedToken; // null unless tokenMode - private final String pinnedPrimaryTerminal; // null unless primary.terminal is configured + /** + * terminal_id → lead name; empty when nothing is pinned. CB-530. + * + *

    A supplier rather than a map because the registry is no longer fixed at startup: CB-531 + * discovers leads by scanning herdr for operator-labelled tabs, so a lead that opens its tab + * after the daemon booted must still be recognised. Consulted per resolve; the scanner behind + * it is TTL-cached, so this is a map lookup in the common case. + */ + private final Supplier> leadTerminals; /** Loopback-trust resolver: no token required, historical behaviour. */ public CallerResolver(ConnectionIdentity identity) { - this(identity, false, null, null); + this(identity, false, null, Map.of()); } - /** As {@link #CallerResolver(ConnectionIdentity, boolean, String, String)} with no pin. */ + /** As {@link #CallerResolver(ConnectionIdentity, boolean, String, Map)} with no leads pinned. */ public CallerResolver(ConnectionIdentity identity, boolean tokenMode, String token) { - this(identity, tokenMode, token, null); + this(identity, tokenMode, token, Map.of()); } /** - * @param identity connection-based worker identification - * @param tokenMode when true, a non-worker caller must present a valid bearer token - * @param token the expected bearer token; required (non-blank) when - * {@code tokenMode} - * @param pinnedPrimaryTerminal the primary's own herdr {@code terminal_id} from - * {@code primary.terminal} ({@code null}/blank = unpinned); a - * caller resolving to this pane is the primary, not a worker + * Single-pin form for the legacy {@code primary.terminal}-only configuration — one lead, named + * {@code primary}. + * + *

    A static factory rather than a fourth constructor overload on purpose: {@code String} and + * {@code Map} overloads are ambiguous for a literal {@code null} argument, which is a compile + * error at the call site and exactly the shape "unpinned" is written in. + * + * @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) { + return new CallerResolver(identity, tokenMode, token, + pinnedPrimaryTerminal == null || pinnedPrimaryTerminal.isBlank() + ? Map.of() : Map.of(pinnedPrimaryTerminal, "primary")); + } + + /** + * @param identity connection-based worker identification + * @param tokenMode when true, a non-worker caller must present a valid bearer token + * @param token the expected bearer token; required (non-blank) when {@code tokenMode} + * @param leadTerminals herdr {@code terminal_id} → lead name for every configured lead + * (CB-530). A caller resolving to one of these panes is that lead — a + * {@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, - String pinnedPrimaryTerminal) { + Map leadTerminals) { + this(identity, tokenMode, token, fixed(leadTerminals)); + } + + /** + * 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. + * + *

    A static factory rather than a fourth constructor overload, for the same reason as + * {@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> leadTerminals) { + return new CallerResolver(identity, tokenMode, token, leadTerminals); + } + + private static Supplier> fixed(Map leadTerminals) { + Map snapshot = leadTerminals == null ? Map.of() : Map.copyOf(leadTerminals); + return () -> snapshot; + } + + private CallerResolver(ConnectionIdentity identity, boolean tokenMode, String token, + Supplier> leadTerminals) { if (tokenMode && (token == null || token.isBlank())) { throw new IllegalArgumentException( "auth.mode=token requires a non-empty token; check that the env var named by " @@ -64,9 +119,20 @@ public final class CallerResolver { this.identity = identity; this.tokenMode = tokenMode; this.expectedToken = tokenMode ? token.getBytes(StandardCharsets.UTF_8) : null; - this.pinnedPrimaryTerminal = - pinnedPrimaryTerminal == null || pinnedPrimaryTerminal.isBlank() - ? null : pinnedPrimaryTerminal; + this.leadTerminals = leadTerminals == null ? Map::of : leadTerminals; + } + + /** + * The currently-recognised leads, {@code terminal_id → name} (CB-535). + * + *

    Deliberately read from the same supplier {@link #resolve} consults, rather than from a + * second copy handed to the roster: a lead that is listed but would not resolve + * (or the reverse) is an address a peer cannot actually reach, and the two answers drifting apart + * is precisely the confusion this exists to end. Live, so a lead discovered by the tab scan after + * startup appears without a restart. + */ + public Map leads() { + return leadTerminals.get(); } /** @@ -79,10 +145,11 @@ public final class CallerResolver { public Principal resolve(String remoteAddr, int remotePort, String authorizationHeader) { ConnectionIdentity.Caller c = identity.resolve(remoteAddr, remotePort); if (c.terminal() != null) { - if (c.terminal().equals(pinnedPrimaryTerminal)) { - // The config names this pane as the primary's own. The pane mapping is exactly as + String lead = leadTerminals.get().get(c.terminal()); + if (lead != null) { + // The config names this pane as a lead's own. The pane mapping is exactly as // unforgeable as a worker's, so it outranks the token path — no credential needed. - return Principal.primary(c.pid()); + return Principal.leader(lead, c.terminal(), c.pid()); } return Principal.worker(c.terminal(), c.pid()); // unforgeable; never token-gated } @@ -99,16 +166,6 @@ public final class CallerResolver { return isLoopback(remoteAddr) ? Principal.primary(c.pid()) : Principal.anonymous(); } - /** The working directory of the calling process (CB-112 spawn cwd inheritance), or {@code null}. */ - public String cwdForPid(long pid) { - return identity.cwdForPid(pid); - } - - /** True when auth requires a bearer token of non-worker callers. */ - public boolean tokenMode() { - return tokenMode; - } - private boolean presentedTokenMatches(String authorizationHeader) { String presented = bearerValue(authorizationHeader); if (presented == null) { diff --git a/bridged/src/main/java/dev/ltms/bridged/auth/Principal.java b/bridged/src/main/java/dev/ltms/bridged/auth/Principal.java index 5d739b0..e4803bf 100644 --- a/bridged/src/main/java/dev/ltms/bridged/auth/Principal.java +++ b/bridged/src/main/java/dev/ltms/bridged/auth/Principal.java @@ -5,21 +5,54 @@ package dev.ltms.bridged.auth; * identifies which worker it is (CB-501). * * @param role what this caller is authorized to act as - * @param terminal the worker's herdr terminal id; {@code null} for {@code PRIMARY}/{@code ANONYMOUS} + * @param terminal the herdr {@code terminal_id} of the pane this caller occupies — a worker's, or + * (since CB-532) a named lead's; {@code null} for an unnamed primary resolved off a + * token or loopback trust, and for {@code ANONYMOUS} * @param pid the connecting process id, or {@code -1} when not resolvable (audit context) + * @param name for a lead resolved from the CB-530 {@code leaders:} registry, which lead it is; + * {@code null} for every other caller, including an unnamed primary */ -public record Principal(Role role, String terminal, long pid) { +public record Principal(Role role, String terminal, long pid, String name) { + + /** + * Three-arg form for the callers that have no name to carry (workers, anonymous, and the + * token/loopback primary paths). Kept so adding CB-530's {@code name} did not churn every + * construction site — and so a reconstructed principal without a stashed name still works. + */ + public Principal(Role role, String terminal, long pid) { + this(role, terminal, pid, null); + } /** A caller authenticated as nothing — the default when no check establishes anything else. */ public static Principal anonymous() { return new Principal(Role.ANONYMOUS, null, -1); } - /** The orchestrating session. */ + /** The orchestrating session, unnamed (token or loopback-trust path). */ public static Principal primary(long pid) { return new Principal(Role.PRIMARY, null, pid); } + /** + * A named lead from the {@code leaders:} registry (CB-530). + * + *

    Carries {@link Role#PRIMARY}: a lead is a primary as far as authorization goes, + * so every existing {@code isPrimary()} gate keeps working unchanged and the role table needed + * no new entry. The name is reporting only — it lets {@code bridge_whoami} say which + * lead is asking once more than one is configured. + * + *

    CB-532: a lead now carries the terminal it was matched by. Under CB-530 it + * deliberately did not, because {@code terminal} meant "which worker pane" everywhere and a + * non-null one would have enrolled the lead in the worker presence map. That reading was what + * 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. + */ + public static Principal leader(String name, String terminal, long pid) { + return new Principal(Role.PRIMARY, terminal, pid, name); + } + /** A worker peer, identified by its herdr pane. */ public static Principal worker(String terminal, long pid) { return new Principal(Role.WORKER, terminal, pid); @@ -39,18 +72,25 @@ public record Principal(Role role, String terminal, long pid) { /** * Whether this caller may act as {@code sessionId} — the "own session only" rule that - * keeps one worker from replying or asking on another's behalf. Only a worker can own a - * session, and only its own. + * keeps one peer from replying or asking on another's behalf. + * + *

    The rule is about identity, not rank: a caller may act as the pane it demonstrably + * occupies, and as no other. CB-532 dropped the extra {@code isWorker()} conjunct that used to + * be here. It was not what enforced the rule — {@code terminal.equals(sessionId)} is, and that + * terminal comes from the connection, so it cannot be forged either way. All the conjunct did + * was make a lead permanently unable to answer anyone, since a lead's terminal was null and a + * lead is not a worker. A caller with no terminal at all (an off-host or token-authenticated + * primary) still owns nothing, which is the case the null check covers. */ public boolean ownsSession(String sessionId) { - return isWorker() && terminal != null && terminal.equals(sessionId); + return terminal != null && terminal.equals(sessionId); } /** Short, non-sensitive description for audit lines and error details. */ public String describe() { return switch (role) { case WORKER -> "worker:" + terminal; - case PRIMARY -> "primary"; + case PRIMARY -> name == null ? "primary" : "leader:" + name; case ANONYMOUS -> "anonymous"; }; } diff --git a/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java b/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java index c1816df..a732e93 100644 --- a/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java +++ b/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java @@ -3,6 +3,8 @@ package dev.ltms.bridged.config; import com.fasterxml.jackson.annotation.JsonIgnoreProperties; import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.dataformat.yaml.YAMLFactory; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import java.io.IOException; import java.io.UncheckedIOException; @@ -17,7 +19,8 @@ import java.util.Set; /** * {@code bridged} configuration, loaded from a YAML file (see * {@code bridged.example.yaml}). Unknown keys are ignored so config can grow ahead - * of the code. + * of the code — but an unknown top-level key is logged as a WARN at load (CB-530), because + * silently dropping a whole block is indistinguishable from honouring it. * * @param bind REST/MCP listen host:port * @param herdrSocket path to herdr's Unix socket ({@code null} → client default) @@ -37,6 +40,11 @@ import java.util.Set; * @param primary optional pinned primary terminal config ({@code null} → derived from connection); * a non-blank {@code terminal} seeds {@code PrimaryRegistry} and prevents * connection-derived overrides, CB-307 + * @param leaders named panes that orchestrate rather than are orchestrated (CB-530), keyed by + * lead name; supersedes the singular {@code primary} pin, which stays honoured. + * See {@link #leaderTerminals()} for how the two merge + * @param leadScan opt-in discovery of leads by tab label (CB-531); {@code null} ⇒ no scanning, + * and only {@code leaders:}/{@code primary:} name a lead * @param placement how to choose a worker profile for an unqualified spawn: * {@code fixed} (default), {@code round-robin}, or {@code weighted} * @param auth API authentication mode ({@code null} → {@code loopback-trust}, the @@ -56,6 +64,8 @@ public record BridgedConfig( Integer spawnReadyPollMs, Broker broker, Primary primary, + Map leaders, + LeadScan leadScan, String placement, Auth auth) { @@ -246,9 +256,12 @@ public record BridgedConfig( * ({@code null} → disabled) * @param drainTimeoutSeconds seconds to wait for {@code BUSY} sessions to finish before * forced teardown on shutdown (default 5 when unset) + * @param clearAfterTurn whether a reusable worker discards its conversation context after + * every completed delegated turn (default false) */ @JsonIgnoreProperties(ignoreUnknown = true) - public record Lifecycle(Integer idleTtlSeconds, Integer contextCap, Integer drainTimeoutSeconds) { + public record Lifecycle(Integer idleTtlSeconds, Integer contextCap, Integer drainTimeoutSeconds, + boolean clearAfterTurn) { } /** @@ -295,6 +308,84 @@ public record BridgedConfig( } } + /** + * One entry of the CB-530 {@code leaders:} registry — a pane that orchestrates rather than one + * that is orchestrated. + * + *

    Why a registry and not a second {@code primary:}: {@code primary.terminal} is singular by + * construction, so a session in any other pane resolves as a worker. That is correct while one + * lead drives a fleet, and wrong the moment two leads (say an Opus lead and an opencode lead) + * work as peers — the second is silently demoted and refused every orchestration call. + * + *

    {@code kind} and {@code model} are descriptive only at this stage: they document what runs + * in the pane and are reported back by {@code bridge_whoami}. Nothing spawns a lead — a lead + * pre-exists, which is precisely why it must be recognised by configuration rather than created. + * + * @param terminal the lead's herdr {@code terminal_id}; the only field identity depends on + * @param kind which agent runs there ({@code claude}, {@code opencode}, …); descriptive + * @param model the model or selector it runs, for operators reading the roster; descriptive + */ + @JsonIgnoreProperties(ignoreUnknown = true) + public record Leader(String terminal, String kind, String model) { + } + + /** + * Discover leads by tab label instead of by pasted {@code terminal_id} (CB-531). + * + *

    Why: a lead is not spawned, so its {@code terminal_id} exists only once a human has opened + * the tab and started the agent — which makes {@code leaders:} a three-step ritual (start it, + * ask it its id, edit config, restart) repeated per lead. Naming the tab is one step, done at + * the moment the operator is already there. The convention also survives what an id does not: + * close the tab and reopen it and the id changes, while the label is retyped as-is. + * + *

    Deliberately opt-in ({@code null} ⇒ off). Turning it on widens who resolves as + * {@link dev.ltms.bridged.auth.Role#PRIMARY}, and a config that never asked for it must not + * acquire that by upgrading the daemon. + * + *

    bridged never writes these labels — see {@link dev.ltms.bridged.herdr.LeadTabScanner} for + * why that one-way direction is what keeps the convention trustworthy. + * + * @param tabPrefix label prefix marking a lead's tab, matched case-insensitively; the + * remainder is the lead's name ({@code "lead: opus-5.0"} → {@code + * opus-5.0}). Default {@code "lead:"} + * @param intervalSeconds how long a scan is cached before herdr is asked again; also the worst + * case before a newly-labelled tab is recognised. Default 10 + */ + @JsonIgnoreProperties(ignoreUnknown = true) + public record LeadScan(String tabPrefix, Integer intervalSeconds) { + public LeadScan { + tabPrefix = (tabPrefix == null || tabPrefix.isBlank()) ? "lead:" : tabPrefix.strip(); + intervalSeconds = (intervalSeconds == null || intervalSeconds <= 0) ? 10 : intervalSeconds; + } + } + + /** + * The terminal → lead-name map that {@link dev.ltms.bridged.auth.CallerResolver} resolves + * against, merging the {@code leaders:} registry with the legacy singular {@code primary:} pin. + * + *

    Precedence: an explicit {@code leaders:} entry wins over the {@code primary:} pin for the + * same terminal. The pin is the older, less expressive spelling of the same fact, so when both + * name a pane the named entry is the one an operator meant. The pin is still honoured on its + * own — a config carrying only {@code primary:} behaves exactly as it did before CB-530. + * + * @return an unmodifiable map, empty when neither block is configured (nothing is pinned, and + * every pane therefore resolves as a worker — the pre-CB-307 behaviour) + */ + public Map leaderTerminals() { + Map byTerminal = new LinkedHashMap<>(); + if (leaders != null) { + leaders.forEach((name, leader) -> { + if (leader != null && leader.terminal() != null && !leader.terminal().isBlank()) { + byTerminal.put(leader.terminal(), name); + } + }); + } + if (primary != null && primary.terminal() != null && !primary.terminal().isBlank()) { + byTerminal.putIfAbsent(primary.terminal(), "primary"); + } + return Collections.unmodifiableMap(byTerminal); + } + /** * API authentication (CB-501). Governs how a caller that is not an on-host worker * pane proves it is the primary. @@ -387,30 +478,94 @@ public record BridgedConfig( return p.isEmpty() ? null : p.keySet().iterator().next(); } + private static final Logger log = LoggerFactory.getLogger(BridgedConfig.class); + private static final ObjectMapper YAML = new ObjectMapper(new YAMLFactory()); + /** + * Top-level keys this version understands. Used only to warn about the rest — see + * {@link #warnUnknownTopLevelKeys}. Keep in step with the record components. + */ + private static final Set KNOWN_TOP_LEVEL_KEYS = Set.of( + "bind", "herdrSocket", "worker", "workers", "defaultWorker", "guard", "worktreeRoot", + "lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs", "broker", "primary", "leaders", + "leadScan", "placement", "auth"); + /** Load and validate config from {@code path}. */ public static BridgedConfig load(Path path) { try { - BridgedConfig cfg = YAML.readValue(Files.readString(path), BridgedConfig.class); + String yaml = Files.readString(path); + warnUnknownTopLevelKeys(yaml, path); + BridgedConfig cfg = YAML.readValue(yaml, BridgedConfig.class); return cfg.withDefaults(); } catch (IOException e) { throw new UncheckedIOException("cannot read bridged config at " + path, e); } } + /** + * Log a WARN naming any top-level key this version does not understand (CB-530). + * + *

    Why this exists: every record here is {@code @JsonIgnoreProperties(ignoreUnknown = true)}, + * which is deliberate — config must be allowed to grow ahead of the code, and a rolled-back + * daemon must still start. The cost is that a whole block can be written, parsed, dropped, and + * never mentioned again. That is exactly how a hand-written {@code leaders:} registry came to + * look configured while being inert: the daemon started, nothing complained, and the only way + * to discover it was reading the config class. + * + *

    A warning rather than a failure, on purpose. Failing closed would turn "the config names + * something this build has not learned yet" into a daemon that will not boot — which is the + * forward-compatibility this annotation was chosen to preserve. Loud, not fatal. + */ + private static void warnUnknownTopLevelKeys(String yaml, Path path) { + List unknown = unknownTopLevelKeys(yaml); + if (!unknown.isEmpty()) { + log.warn("{}: ignoring unknown top-level config key(s) {} — this build does not " + + "understand them, so they have NO effect. Check for a typo, or a " + + "feature not in this version.", + path, unknown); + } + } + + /** + * The top-level keys in {@code yaml} that this build does not understand, sorted. Package-private + * so the guardrail is asserted directly rather than through a log appender. + * + * @return empty when everything is known, or when {@code yaml} is not a mapping at all (a + * malformed file is {@code readValue}'s error to report, not this method's) + */ + static List unknownTopLevelKeys(String yaml) { + Map raw; + try { + raw = YAML.readValue(yaml, Map.class); + } catch (IOException | IllegalArgumentException e) { + return List.of(); + } + if (raw == null) { + return List.of(); + } + return raw.keySet().stream() + .map(String::valueOf) + .filter(k -> !KNOWN_TOP_LEVEL_KEYS.contains(k)) + .sorted() + .toList(); + } + /** Fill in nested defaults so callers never see nulls for structural fields. */ public BridgedConfig withDefaults() { Bind b = bind != null ? bind : new Bind(null, 0); Guard g = guard != null ? guard : new Guard(List.of()); - Lifecycle l = lifecycle != null ? lifecycle : new Lifecycle(null, null, null); + Lifecycle l = lifecycle != null ? lifecycle : new Lifecycle(null, null, null, false); Integer timeout = (spawnReadyTimeoutMs != null) ? spawnReadyTimeoutMs : 20000; Integer pollMs = (spawnReadyPollMs != null) ? spawnReadyPollMs : 300; Auth a = auth != null ? auth : new Auth(null, null); String placementOrDefault = (placement != null && !placement.isBlank()) ? placement : "fixed"; // broker is left as-is: null (or an empty/blank uri) keeps the in-memory soft-state inbox. // primary is left as-is: null defaults to connection-derived identity. - return new BridgedConfig(b, herdrSocket, worker, workers, defaultWorker, g, worktreeRoot, l, timeout, pollMs, broker, primary, placementOrDefault, a); + // leadScan is left as-is: null is "off", and LeadScan'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. + return new BridgedConfig(b, herdrSocket, worker, workers, defaultWorker, g, worktreeRoot, l, timeout, pollMs, broker, primary, leaders, leadScan, placementOrDefault, a); } /** @@ -437,6 +592,46 @@ public record BridgedConfig( + "bind to 127.0.0.1 and put a reverse proxy in front."); } + /** + * Reject a lead-scan convention that a worker tab would also satisfy (CB-531). + * + *

    The scan reads a tab label and concludes "a lead lives here". bridged also writes + * tab labels — every worker gets {@code tabLabel} rendered into its tab. Choose a + * {@code leadScan.tabPrefix} that a worker template matches and the daemon starts labelling its + * own workers as leads, promoting the entire fleet to {@link dev.ltms.bridged.auth.Role#PRIMARY} + * with no message and no diff. The worker-space exclusion in + * {@link dev.ltms.bridged.herdr.LeadTabScanner} already blocks the realistic path, but defence + * that depends on one workspace label holding is not defence enough for a privilege boundary. + * + *

    Fatal rather than a warning, unlike {@link #warnUnknownTopLevelKeys}: an unknown key means + * a feature does nothing, while this means a feature does the opposite of what it says. + * + * @throws IllegalStateException when any worker profile's {@code tabLabel} starts with the + * configured lead prefix + */ + public void validateLeadScan() { + if (leadScan == null) { + return; + } + String prefix = leadScan.tabPrefix(); + List clashing = workerProfiles().entrySet().stream() + .filter(e -> e.getValue().tabLabel() != null + && e.getValue().tabLabel().strip() + .regionMatches(true, 0, prefix, 0, prefix.length())) + .map(Map.Entry::getKey) + .sorted() + .toList(); + if (clashing.isEmpty()) { + return; + } + throw new IllegalStateException( + "refusing to start: leadScan.tabPrefix=\"" + prefix + "\" also matches the tabLabel " + + "of worker profile(s) " + clashing + ". Every worker spawned under them " + + "would be read back as a lead and granted spawn/stop/send on the whole " + + "fleet. Change one of the two so worker tabs and lead tabs cannot be " + + "confused."); + } + /** True for the loopback addresses and the unspecified-but-local forms we treat as same-host. */ private static boolean isLoopbackBind(String host) { if (host == null || host.isBlank()) { diff --git a/bridged/src/main/java/dev/ltms/bridged/herdr/LeadTabScanner.java b/bridged/src/main/java/dev/ltms/bridged/herdr/LeadTabScanner.java new file mode 100644 index 0000000..92deac6 --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/herdr/LeadTabScanner.java @@ -0,0 +1,164 @@ +package dev.ltms.bridged.herdr; + +import com.fasterxml.jackson.databind.JsonNode; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.util.Collections; +import java.util.LinkedHashMap; +import java.util.Map; +import java.util.Set; +import java.util.function.LongSupplier; +import java.util.function.Supplier; + +/** + * Discovers which panes host a lead by scanning herdr for tabs the operator labelled by convention + * (CB-531), and hands {@link dev.ltms.bridged.auth.CallerResolver} the resulting + * {@code terminal_id → lead name} map. + * + *

    Why scan at all. A lead is never spawned — a human opens a tab and starts an + * agent in it — so the daemon cannot learn a lead's {@code terminal_id} at creation time the way it + * does a worker's. CB-530 solved that by having the operator paste each id into {@code leaders:}, + * which works but costs a config edit and a daemon restart per lead, and the id is only obtainable + * by first starting the session and asking it. Scanning closes that loop: label the tab, and the + * pane is recognised on the next resolve. + * + *

    Direction of trust. The label names the lead; it never grants + * anything a pane could take for itself. Three properties keep that honest: + *

      + *
    1. bridged never renames a lead tab. The operator's label is read-only input, so what is in + * the tab bar is always what the human wrote — no round-trip where the daemon's own rename + * becomes the evidence for its next decision.
    2. + *
    3. Worker spaces are excluded wholesale ({@code excludedWorkspaceLabels}), so a worker cannot + * become a lead by being placed — as a split, say — inside a matching tab.
    4. + *
    5. A worker cannot rename a tab: {@code tab.rename} is reachable only through + * {@link WorkspaceControl}, which no {@code bridge_*} tool exposes. The label is writable by + * the human at the terminal and by nobody the bridge is defending against.
    6. + *
    + * The remaining hazard is an operator one — a worker {@code tabLabel} template that + * happens to start with the same prefix would promote the whole fleet — and that is refused at + * startup by {@code BridgedConfig.validateLeadScan} rather than documented here. + * + *

    Caching. {@link #get()} is on the request path (every resolve), so the scan + * is TTL-cached and a stale-but-valid map is preferred to a herdr round-trip. A failed scan keeps + * the previous answer instead of emptying it — a herdr hiccup must not silently demote a live lead + * mid-session. + */ +public final class LeadTabScanner implements Supplier> { + + private static final Logger log = LoggerFactory.getLogger(LeadTabScanner.class); + + private final HerdrClient herdr; + private final String tabPrefix; + private final Set excludedWorkspaceLabels; + private final Map configuredLeads; + private final long ttlNanos; + private final LongSupplier clock; + + private Map cached; + private long scannedAtNanos; + private boolean everScanned; + + /** + * @param herdr the herdr client to query ({@code workspace.list}, + * {@code tab.list}, {@code pane.list} — all read-only) + * @param tabPrefix a tab whose label starts with this (case-insensitively) hosts a + * lead; the rest of the label, trimmed, is the lead's name + * @param excludedWorkspaceLabels workspaces never scanned — the configured worker spaces + * @param configuredLeads the static {@code leaders:}/{@code primary:} registry, merged + * over every scan result. Explicit config outranks the + * convention, and survives a scan that cannot run at all + * @param ttlNanos how long a scan result is reused before the next one + * @param clock nanosecond time source ({@code System::nanoTime} in production) + */ + public LeadTabScanner(HerdrClient herdr, String tabPrefix, Set excludedWorkspaceLabels, + Map configuredLeads, long ttlNanos, LongSupplier clock) { + this.herdr = herdr; + this.tabPrefix = tabPrefix == null || tabPrefix.isBlank() ? "lead:" : tabPrefix.strip(); + this.excludedWorkspaceLabels = excludedWorkspaceLabels == null + ? Set.of() : Set.copyOf(excludedWorkspaceLabels); + this.configuredLeads = configuredLeads == null ? Map.of() : Map.copyOf(configuredLeads); + this.ttlNanos = ttlNanos; + this.clock = clock; + this.cached = this.configuredLeads; + } + + /** + * The current {@code terminal_id → lead name} map, rescanning when the cache has expired. + * + *

    Synchronized so a burst of concurrent calls produces one scan rather than one each; a scan + * is a handful of RPCs over a Unix socket and is rate-limited to one per TTL. + */ + @Override + public synchronized Map get() { + long now = clock.getAsLong(); + if (everScanned && now - scannedAtNanos < ttlNanos) { + return cached; + } + // Stamp before scanning, not after: a herdr that is down must cost one attempt per TTL, not + // one per request. + scannedAtNanos = now; + everScanned = true; + try { + Map fresh = scan(); + if (!fresh.equals(cached)) { + log.info("lead panes: {}", fresh); + } + cached = fresh; + } catch (HerdrException e) { + log.warn("lead-tab scan failed, keeping the {} lead(s) already known: {}", + cached.size(), e.getMessage()); + } + return cached; + } + + /** One full pass: labelled tabs → their panes → those panes' terminals. */ + private Map scan() { + Map nameByTab = new LinkedHashMap<>(); + for (JsonNode w : herdr.call("workspace.list").path("workspaces")) { + Workspace ws = Workspace.from(w); + if (ws.workspaceId() == null || excludedWorkspaceLabels.contains(ws.label())) { + continue; + } + for (JsonNode t : herdr.call("tab.list", Map.of("workspace_id", ws.workspaceId())).path("tabs")) { + Tab tab = Tab.from(t); + String name = leadNameOf(tab.label()); + if (name != null && tab.tabId() != null) { + nameByTab.put(tab.tabId(), name); + } + } + } + + Map byTerminal = new LinkedHashMap<>(); + if (!nameByTab.isEmpty()) { + // One pane.list for every tab: panes carry tab_id, so the join is local. + for (JsonNode p : herdr.call("pane.list", Map.of()).path("panes")) { + String name = nameByTab.get(p.path("tab_id").asText(null)); + String terminal = p.path("terminal_id").asText(null); + if (name != null && terminal != null && !terminal.isBlank()) { + byTerminal.put(terminal, name); + } + } + } + byTerminal.putAll(configuredLeads); // an explicit pin outranks a label + return Collections.unmodifiableMap(byTerminal); + } + + /** + * The lead name a tab label declares, or {@code null} if it declares none. + * + *

    {@code "lead: opus-5.0"} → {@code "opus-5.0"}. A bare {@code "lead:"} names nobody and is + * rejected: an unnamed lead would resolve as {@code PRIMARY} with nothing to attribute it to. + */ + private String leadNameOf(String label) { + if (label == null) { + return null; + } + String l = label.strip(); + if (!l.regionMatches(true, 0, tabPrefix, 0, tabPrefix.length())) { + return null; + } + String name = l.substring(tabPrefix.length()).strip(); + return name.isEmpty() ? null : name; + } +} diff --git a/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java b/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java index 2c815b4..87fcf6c 100644 --- a/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java +++ b/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java @@ -122,6 +122,15 @@ public final class CompletionResolver implements TurnListener { Thread.ofVirtual().name("completion-" + target).start(() -> resolve(target, turn)); } + /** + * Resolve the completed turn before adapter housekeeping can erase its rendered output. This is + * intentionally synchronous and used only when a post-turn context reset is enabled; the normal + * path remains off-loaded so polling is not blocked by a scrape. + */ + public void resolveBeforePostAction(String target) { + resolve(target, inFlight.get(target)); + } + @Override public void onTurnFailed(String target) { InFlight turn = inFlight.get(target); diff --git a/bridged/src/main/java/dev/ltms/bridged/inject/Injector.java b/bridged/src/main/java/dev/ltms/bridged/inject/Injector.java index 20e61c6..f72b902 100644 --- a/bridged/src/main/java/dev/ltms/bridged/inject/Injector.java +++ b/bridged/src/main/java/dev/ltms/bridged/inject/Injector.java @@ -132,6 +132,10 @@ public final class Injector { boolean turnObserved; // saw a real `working` sample since that delivery (turn ran) int unknownSinceTurn; // consecutive `unknown` samples while a delegation is outstanding (CB-109) int notReadySincePoll; // consecutive injectable samples a queued message waited on the readiness gate (CB-114) + boolean postTurnPending; // completion observed; adapter housekeeping has not started yet + boolean awaitingPostTurnPickup; + boolean postTurnObserved; + int injectableSincePostTurnPickup; synchronized void add(Pending p) { queue.add(p); @@ -173,9 +177,15 @@ public final class Injector { boolean turnCompleted = false; boolean turnFailed = false; boolean resubmit = false; + boolean startPostTurn = false; List notReady = null; // queued messages failed because the worker never became ready synchronized (t) { if (status == AgentStatus.WORKING) { + if (t.awaitingPostTurnPickup) { + t.awaitingPostTurnPickup = false; + t.injectableSincePostTurnPickup = 0; + t.postTurnObserved = true; + } // Definitive pickup: the worker is busy on our last message, and (if a delivery is // outstanding) a real turn is now confirmed to be running. t.awaitingPickup = false; @@ -185,6 +195,16 @@ public final class Injector { if (t.awaitingCompletion) t.turnObserved = true; } else if (status.injectable()) { // IDLE or BLOCKED t.unknownSinceTurn = 0; + if (t.awaitingPostTurnPickup) { + if (++t.injectableSincePostTurnPickup >= PICKUP_GRACE_POLLS) { + t.awaitingPostTurnPickup = false; + t.injectableSincePostTurnPickup = 0; + } else { + resubmit = true; + } + } else if (t.postTurnObserved) { + t.postTurnObserved = false; + } if (t.awaitingPickup) { if (++t.injectableSincePickup >= PICKUP_GRACE_POLLS) { // Pickup edge was never sampled (turn faster than the poll, or status lag). @@ -209,11 +229,16 @@ public final class Injector { t.awaitingCompletion = false; t.turnObserved = false; turnCompleted = true; + if (turnListener.hasPostTurnAction(target)) { + t.postTurnPending = true; + startPostTurn = true; + } } // Deliver the next queued message only once the prior turn is fully settled, so a // completion is never confused with the pickup of the following message — and only // once the worker is available (CB-113), so we never paste into its boot window. - if (!t.awaitingCompletion) { + if (!t.awaitingCompletion && !t.postTurnPending + && !t.awaitingPostTurnPickup && !t.postTurnObserved) { Pending p = t.queue.peek(); if (p != null && ready.test(target)) { t.notReadySincePoll = 0; @@ -263,7 +288,8 @@ public final class Injector { // Reclaim the entry once the worker is fully quiescent (nothing queued, no pickup or // completion awaited), so the map cannot grow without bound across short-lived workers. - if (t.queue.isEmpty() && !t.awaitingPickup && !t.awaitingCompletion) { + if (t.queue.isEmpty() && !t.awaitingPickup && !t.awaitingCompletion + && !t.postTurnPending && !t.awaitingPostTurnPickup && !t.postTurnObserved) { targets.remove(target, t); } } @@ -290,7 +316,21 @@ public final class Injector { turnListener.onTurnFailed(target); } if (turnCompleted) { - turnListener.onTurnComplete(target); + if (startPostTurn) { + boolean started = turnListener.onTurnCompleteWithPostAction(target); + synchronized (t) { + t.postTurnPending = false; + if (started) { + t.awaitingPostTurnPickup = true; + t.injectableSincePostTurnPickup = 0; + } + if (t.queue.isEmpty() && !t.awaitingPostTurnPickup) { + targets.remove(target, t); + } + } + } else { + turnListener.onTurnComplete(target); + } } if (turnFailed) { turnListener.onTurnFailed(target); @@ -317,7 +357,8 @@ public final class Injector { .filter(e -> { synchronized (e.getValue()) { Target t = e.getValue(); - return !t.queue.isEmpty() || t.awaitingPickup || t.awaitingCompletion; + return !t.queue.isEmpty() || t.awaitingPickup || t.awaitingCompletion + || t.postTurnPending || t.awaitingPostTurnPickup || t.postTurnObserved; } }) .map(java.util.Map.Entry::getKey) @@ -343,6 +384,9 @@ public final class Injector { hadDeliveredTurn = t.awaitingCompletion; t.awaitingCompletion = false; t.awaitingPickup = false; + t.postTurnPending = false; + t.awaitingPostTurnPickup = false; + t.postTurnObserved = false; } forget.accept(target); // the worker is gone — clear its readiness/presence too (CB-114) for (Pending p : pending) { diff --git a/bridged/src/main/java/dev/ltms/bridged/inject/TurnListener.java b/bridged/src/main/java/dev/ltms/bridged/inject/TurnListener.java index b1e752b..6ac5062 100644 --- a/bridged/src/main/java/dev/ltms/bridged/inject/TurnListener.java +++ b/bridged/src/main/java/dev/ltms/bridged/inject/TurnListener.java @@ -13,6 +13,25 @@ public interface TurnListener { /** A worker's delegated turn finished (worker returned to idle after visibly working). */ void onTurnComplete(String target); + /** + * Whether completion must pause ordinary delivery while adapter-specific housekeeping starts. + * This is queried before the injector considers the next queued message, closing the same-tick + * ordering gap. Implementations must not mutate state here. + */ + default boolean hasPostTurnAction(String target) { + return false; + } + + /** + * Complete the delegated turn and start its non-turn housekeeping operation. + * + * @return true when the injector must observe that operation settle before the next delivery + */ + default boolean onTurnCompleteWithPostAction(String target) { + onTurnComplete(target); + return false; + } + /** * A worker that visibly ran a delegated turn then wedged in a non-idle, non-working state * (CB-109) — e.g. an error screen herdr classifies as {@code unknown} — so no diff --git a/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java index 07f0a59..223f65e 100644 --- a/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java +++ b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java @@ -65,6 +65,8 @@ public final class BridgeMcp { static final String CALLER_PID = "callerPid"; /** Transport-context key under which the extractor stashes the resolved {@link Role} (CB-501). */ static final String CALLER_ROLE = "callerRole"; + /** Transport-context key for the lead's configured name, when the caller is one (CB-530). */ + static final String CALLER_NAME = "callerName"; private final HttpServletStreamableServerTransportProvider transport; private final McpSyncServer server; @@ -106,11 +108,15 @@ public final class BridgeMcp { ? callers.resolve(req.getRemoteAddr(), req.getRemotePort(), req.getHeader("Authorization")) : legacyPrincipal(identity, req.getRemoteAddr(), req.getRemotePort()); - presence.markPresent(p.terminal()); // no-op for the primary (null terminal) + // 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()); return McpTransportContext.create(Map.of( CALLER_TERMINAL, orEmpty(p.terminal()), CALLER_PID, Long.toString(p.pid()), - CALLER_ROLE, p.role().name())); + CALLER_ROLE, p.role().name(), + CALLER_NAME, orEmpty(p.name()))); }) .build(); this.server = McpServer.sync(transport) @@ -123,6 +129,9 @@ public final class BridgeMcp { String caller = callerTerminal(exchange); if (caller != null) primaryRegistry.record(caller); Map a = req.arguments(); + // CB-532: remember WHICH lead is waiting on this worker, so its reply nudge goes + // back to that lead rather than to whichever one happened to send first. + primaryRegistry.recordDelegation(str(a, "sessionId"), caller); String turnId = str(a, "turnId"); if (turnId != null && !turnId.isBlank()) { // Answering a worker's bridge_ask (CB-205): resolve its blocked question and @@ -185,7 +194,9 @@ public final class BridgeMcp { .toolCall(listTool(), (exchange, _) -> { McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null); if (denied != null) return denied; - return listWorkers(workers, sessions); + return listFleet(workers, sessions, + callers == null ? Map.of() : callers.leads(), + callerTerminal(exchange)); }) .toolCall(stopTool(), (exchange, req) -> { String paneId = str(req.arguments(), "paneId"); @@ -222,7 +233,7 @@ public final class BridgeMcp { /** The caller reconstructed from the transport context. */ private static Principal principal(McpSyncServerExchange exchange) { return principalFrom(exchange.transportContext().get(CALLER_ROLE), - callerTerminal(exchange), callerPid(exchange)); + callerTerminal(exchange), callerPid(exchange), callerName(exchange)); } /** @@ -237,11 +248,16 @@ public final class BridgeMcp { * @param pid the calling pid, or {@code -1} */ static Principal principalFrom(Object role, String terminal, long pid) { + return principalFrom(role, terminal, pid, null); + } + + /** As {@link #principalFrom(Object, String, long)}, carrying a lead's name (CB-530). */ + static Principal principalFrom(Object role, String terminal, long pid, String name) { if (role == null) { // No role stashed (legacy path): fall back to the historical interpretation. return terminal != null ? Principal.worker(terminal, pid) : Principal.primary(pid); } - return new Principal(Role.valueOf(role.toString()), terminal, pid); + return new Principal(Role.valueOf(role.toString()), terminal, pid, name); } /** @@ -292,6 +308,13 @@ public final class BridgeMcp { return (s == null || s.isBlank()) ? null : s; } + /** The lead name resolved from this call's connection, or {@code null} (CB-530). */ + private static String callerName(McpSyncServerExchange exchange) { + Object v = exchange.transportContext().get(CALLER_NAME); + String s = v == null ? null : v.toString(); + return (s == null || s.isBlank()) ? null : s; + } + /** The caller's PID resolved from this call's connection, or {@code -1} if unknown. */ private static long callerPid(McpSyncServerExchange exchange) { Object v = exchange.transportContext().get(CALLER_PID); @@ -485,6 +508,17 @@ public final class BridgeMcp { Map m = new LinkedHashMap<>(); m.put("role", caller.role().name().toLowerCase()); if (!caller.isWorker()) { + // CB-530: which lead, once more than one pane is configured as one. `role` deliberately + // still reads "primary" — the fallback ladder in CLAUDE.md keys on it, and a lead IS a + // primary for authorization; the name is additive so no existing reader breaks. + if (caller.name() != null) { + m.put("leader", caller.name()); + } + // CB-532: a lead's own pane, so it can tell a peer where to reach it — and so an + // operator can read off which tab hosts which lead without going to herdr. + if (caller.terminal() != null) { + m.put("sessionId", caller.terminal()); + } return text(json(m)); } m.put("sessionId", caller.terminal()); @@ -576,25 +610,65 @@ public final class BridgeMcp { } /** - * {@code bridge_list}: bridge-owned roster merged with live herdr status. CB-519 decoupled the - * registry key (a host-unique id) from the herdr pane coordinate, so the join is on the - * terminal id, which both the session and the live agent carry. + * {@code bridge_list}: the whole fleet — {@code leads} and {@code workers} — each merged with + * live herdr status. CB-519 decoupled the registry key (a host-unique id) from the herdr pane + * coordinate, so the join is on the terminal id, which both the session and the live agent carry. + * + *

    CB-535 added the {@code leads} half. Until then this listed the worker roster alone, and a + * lead asking "who else is here?" got an empty array — which reads as no peers but + * actually means no workers spawned. There was no way at all for a lead to learn a + * peer's address; it had to be carried across by a human. Both halves are reported even when a + * half is empty, so an empty {@code workers} can no longer be mistaken for an empty fleet. + * + *

    Leads are drawn from the resolver rather than from a second registry, so an address listed + * here is one that would actually resolve as a lead — see {@link CallerResolver#leads()}. The + * caller's own row is flagged {@code "self": true}: a peer needs to tell its own pane apart from + * a peer's, and the alternative is every lead calling {@code bridge_whoami} to subtract itself. + * + * @param leads terminal_id → lead name, live from the resolver + * @param selfTerm the calling pane's terminal id, or blank for a caller with no pane */ - static McpSchema.CallToolResult listWorkers(PeerLauncher workers, SessionManager sessions) { + static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, + Map leads, String selfTerm) { try { Map live = workers.list().stream() .map(Agent.class::cast) .filter(a -> a.terminalId() != null) .collect(Collectors.toMap(Agent::terminalId, Function.identity(), (_, b) -> b)); + List> leadRows = leads.entrySet().stream() + .sorted(Map.Entry.comparingByValue()) + .map(e -> leadView(e.getKey(), e.getValue(), live.get(e.getKey()), selfTerm)) + .toList(); List> out = sessions.roster().stream() .map(s -> SessionManager.rosterView(s, live.get(s.terminalId()))) .toList(); - return text(json(Map.of("workers", out))); + return text(json(Map.of("leads", leadRows, "workers", out))); } catch (HerdrException e) { - return error("herdr error listing workers: " + e.getMessage()); + return error("herdr error listing the fleet: " + e.getMessage()); } } + /** + * One lead's row: its address, its name, and whether it can be reached right now. + * + *

    {@code status} is herdr's live view, and {@code unknown} when herdr is not tracking that + * pane as an agent — the honest answer, and the one that matters: a lead whose pane herdr cannot + * see is a lead a {@code bridge_send} cannot be typed into. It is reported rather than hidden, + * because a peer that has gone unreachable is exactly what the sender needs to know. + */ + private static Map leadView(String terminal, String name, Agent live, + String selfTerm) { + Map m = new LinkedHashMap<>(); + m.put("sessionId", terminal); + m.put("name", name); + m.put("status", live == null || live.status() == null + ? "unknown" : live.status().name().toLowerCase()); + if (terminal.equals(selfTerm)) { + m.put("self", true); + } + return m; + } + /** {@code bridge_stop}: tear a worker down by its pane id. */ static McpSchema.CallToolResult stop(SessionManager sessions, String paneId) { if (isBlank(paneId)) { @@ -710,8 +784,13 @@ public final class BridgeMcp { private static McpSchema.Tool listTool() { return tool("bridge_list", - "List the worker sessions the bridge tracks — each with its sessionId, paneId, profile, " - + "state, optional worktree/branch/owner, and live herdr status.", + "List the whole fleet the bridge tracks, in two parts. 'leads' are your PEERS — other " + + "orchestrators, each with its sessionId (the address to bridge_send to), " + + "name, live status, and 'self': true on your own row; this is how you " + + "discover a peer lead without being told its address. 'workers' are the " + + "sessions delegated to — each with sessionId, paneId, profile, state, " + + "optional worktree/branch/owner, and live herdr status. An empty 'workers' " + + "means no workers are spawned; it says nothing about peers.", objectSchema(Map.of(), List.of())); } @@ -724,10 +803,12 @@ public final class BridgeMcp { } private static McpSchema.Tool replyTool() { - // No session/target arg — the worker's identity is resolved from the connection. + // No session/target arg — the caller's identity is resolved from the connection. return tool("bridge_reply", - "Return your structured answer for the task you were delegated, " - + "resolving the caller's blocked bridge_send.", + "Return your structured answer for a message you were sent, resolving the sender's " + + "blocked bridge_send. A worker MUST end every delegated turn with exactly " + + "one of these. A lead uses it only to answer another lead that messaged " + + "it — never to answer a worker, whose turn it is not.", objectSchema(Map.of( "content", stringProp("Your reply/answer")), List.of("content"))); @@ -745,11 +826,12 @@ public final class BridgeMcp { return tool("bridge_whoami", "Report who YOU are on the bridge — your role is resolved from your connection " + "(unforgeable), never from anything you claim. Returns role 'primary' (you " - + "orchestrate: spawn/send/stop, and you must never call bridge_reply) or " - + "'worker' (you were delegated to: you must end every turn with exactly one " - + "bridge_reply, and cannot spawn or send), plus your own sessionId, profile, " - + "worktree and branch when you are a worker. Call this first when following " - + "role-conditional instructions rather than guessing your role.", + + "orchestrate: spawn/send/stop; reply ONLY to answer a peer lead that " + + "messaged you, never to answer a worker) or 'worker' (you were delegated " + + "to: you must end every turn with exactly one bridge_reply, and cannot " + + "spawn or send), plus 'leader' naming which lead you are, your own " + + "sessionId, and profile/worktree/branch when you are a worker. Call this " + + "first when following role-conditional instructions rather than guessing.", objectSchema(Map.of(), List.of())); } diff --git a/bridged/src/main/java/dev/ltms/bridged/mcp/PrimaryRegistry.java b/bridged/src/main/java/dev/ltms/bridged/mcp/PrimaryRegistry.java index a04d86b..64b156a 100644 --- a/bridged/src/main/java/dev/ltms/bridged/mcp/PrimaryRegistry.java +++ b/bridged/src/main/java/dev/ltms/bridged/mcp/PrimaryRegistry.java @@ -4,6 +4,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.util.Optional; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicReference; /** @@ -25,6 +26,15 @@ public final class PrimaryRegistry { private final AtomicReference terminal = new AtomicReference<>(); private final boolean pinned; + /** + * CB-532: worker terminal → the lead that delegated to it. The single slot above answers "who is + * THE primary", a question with no correct answer once two leads orchestrate the same fleet: + * whichever called {@code bridge_send} first captured every nudge, including nudges for the + * other lead's delegations. This map answers the question that actually matters — "who is + * waiting on THIS worker" — and is what lets {@code primary.terminal} be retired. + */ + private final ConcurrentHashMap leadByTarget = new ConcurrentHashMap<>(); + /** * @param pinnedTerminal an optional pinned terminal from config ({@code null}/blank = unpinned) */ @@ -56,6 +66,43 @@ public final class PrimaryRegistry { } } + /** + * Record that {@code leadTerminal} delegated to worker {@code target} (CB-532). + * + *

    Called at {@code bridge_send} time, where both halves are known: the target is the tool's + * argument and the lead is resolved from the connection. Last writer wins — if a second lead + * takes over a worker, replies follow the lead that most recently delegated to it, which is the + * one waiting. + */ + public void recordDelegation(String target, String leadTerminal) { + if (target == null || target.isBlank() || leadTerminal == null || leadTerminal.isBlank()) { + return; + } + leadByTarget.put(target, leadTerminal); + } + + /** Forget a worker's delegating lead — call on release, so a torn-down session leaks nothing. */ + public void forgetDelegation(String target) { + if (target != null) { + leadByTarget.remove(target); + } + } + + /** + * Where a nudge about {@code target}'s reply should go: the lead that delegated to it, falling + * back to the single known primary. + * + *

    The fallback matters after a daemon restart, which loses the map while the durable inbox + * keeps the reply. With one lead the fallback is unambiguous and correct. With several and no + * recorded delegation there is no right answer, so this returns empty rather than guessing — + * delivery degrades to pull, which is exactly what the durable inbox is for, instead of + * interrupting the wrong lead with someone else's result. + */ + public Optional nudgeTargetFor(String target) { + String lead = target == null ? null : leadByTarget.get(target); + return lead != null ? Optional.of(lead) : Optional.ofNullable(terminal.get()); + } + /** The known primary terminal, or empty if not yet learned (and not pinned). */ public Optional primaryTerminal() { return Optional.ofNullable(terminal.get()); diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java b/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java index 8fa0eb7..2ff0712 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/ReplyPushLoop.java @@ -59,10 +59,10 @@ public final class ReplyPushLoop { this.metrics = metrics; } - /** Record a counter sample when a registry is wired; a no-op in unit tests. */ - private void count(String name, String... labels) { + /** Count one nudge outcome when a registry is wired; a no-op in unit tests. */ + private void countNudge(String outcome) { if (metrics != null) { - metrics.inc(name, labels); + metrics.inc(BridgedMetrics.PUSH_NUDGES, "outcome", outcome); } } @@ -79,8 +79,13 @@ public final class ReplyPushLoop { * @return the action the caller should take */ Action decide(String target, int reminderCount) { - if (!primaryRegistry.isKnown()) { - log.debug("push: primary unknown, stopping reminder for {}", target); + // CB-532: the destination is per-delegation — the lead that sent this worker its work, not + // "the primary". With two leads orchestrating one fleet the singular question has no right + // answer, and answering it anyway interrupted whichever lead happened to call bridge_send + // first with results it never asked for. + var nudgeTarget = primaryRegistry.nudgeTargetFor(target); + if (nudgeTarget.isEmpty()) { + log.debug("push: no lead is known to be waiting on {}, stopping reminder", target); return Action.STOP; } if (inbox.peek(target).isEmpty()) { @@ -89,21 +94,21 @@ public final class ReplyPushLoop { } if (reminderCount >= maxReminders) { log.debug("push: reminder cap ({}) reached for {}, stopping", maxReminders, target); - count(BridgedMetrics.PUSH_NUDGES, "outcome", "exhausted"); + countNudge("exhausted"); return Action.STOP; } - var primaryTerminal = primaryRegistry.primaryTerminal().orElseThrow(); + String leadTerminal = nudgeTarget.get(); AgentStatus status; try { - status = agents.status(primaryTerminal); + status = agents.status(leadTerminal); } catch (RuntimeException e) { - log.debug("push: status check failed for primary {}, will retry", primaryTerminal, e); + log.debug("push: status check failed for lead {}, will retry", leadTerminal, e); return Action.WAIT_BUSY; } if (status.injectable()) { return Action.INJECT; } - log.debug("push: primary {} is {} (not injectable), waiting", primaryTerminal, status); + log.debug("push: lead {} is {} (not injectable), waiting", leadTerminal, status); return Action.WAIT_BUSY; } @@ -143,16 +148,23 @@ public final class ReplyPushLoop { /** Send the nudge and log the event. */ private void injectNudge(String target, int reminderCount) { - var primaryTerminal = primaryRegistry.primaryTerminal().orElseThrow(); + // Re-read rather than threading it down from decide(): the delegating lead can change + // between the decision and the injection, and the nudge should follow the current one. + var lead = primaryRegistry.nudgeTargetFor(target); + if (lead.isEmpty()) { + log.debug("push: lead for {} disappeared before the nudge could be sent", target); + return; + } + String leadTerminal = lead.get(); String nudge = NUDGE_FORMAT.formatted(target, target); try { - agents.send(primaryTerminal, nudge); - log.debug("push: nudge {}/{} sent to primary {} for target {}", - reminderCount + 1, maxReminders, primaryTerminal, target); - count(BridgedMetrics.PUSH_NUDGES, "outcome", "delivered"); + agents.send(leadTerminal, nudge); + log.debug("push: nudge {}/{} sent to lead {} for target {}", + reminderCount + 1, maxReminders, leadTerminal, target); + countNudge("delivered"); } catch (RuntimeException e) { - log.warn("push: failed to nudge primary {} for target {} (reminder {}/{}): {}", - primaryTerminal, target, reminderCount + 1, maxReminders, e.toString()); + log.warn("push: failed to nudge lead {} for target {} (reminder {}/{}): {}", + leadTerminal, target, reminderCount + 1, maxReminders, e.toString()); } } diff --git a/bridged/src/main/java/dev/ltms/bridged/peer/Capability.java b/bridged/src/main/java/dev/ltms/bridged/peer/Capability.java index 14adfbc..15974a6 100644 --- a/bridged/src/main/java/dev/ltms/bridged/peer/Capability.java +++ b/bridged/src/main/java/dev/ltms/bridged/peer/Capability.java @@ -26,6 +26,12 @@ public enum Capability { */ WORKTREE, + /** + * The peer can discard its current conversation context without starting a delegated bridge + * turn. The command or API used to do that is adapter-specific. + */ + CONTEXT_RESET, + /** * The spawner can reconcile orphaned peers on boot — workers that outlived a prior daemon * process and whose pane ids died with it (CB-117). Claude Code over herdr supports this diff --git a/bridged/src/main/java/dev/ltms/bridged/peer/PeerLauncher.java b/bridged/src/main/java/dev/ltms/bridged/peer/PeerLauncher.java index 5e86005..5ebd958 100644 --- a/bridged/src/main/java/dev/ltms/bridged/peer/PeerLauncher.java +++ b/bridged/src/main/java/dev/ltms/bridged/peer/PeerLauncher.java @@ -82,4 +82,13 @@ public interface PeerLauncher { * when safe to do so. */ void stop(String id); + + /** + * Discard the context of the peer identified by {@code id}. Implementations must bypass normal + * bridge delivery/turn accounting. Unsupported peer kinds return {@code false} without sending + * a guessed command. + * + * @return {@code true} when a reset was sent and its status transition must settle before reuse + */ + boolean clearContext(String id); } diff --git a/bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java b/bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java index 82ae833..b560bf7 100644 --- a/bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java +++ b/bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java @@ -47,6 +47,7 @@ public final class SessionManager implements TurnListener { private final AtomicLong nonceSeq = new AtomicLong(); private final LongSupplier nowNanos; private final int contextCap; + private final boolean clearAfterTurn; /** CB-520: notified with a terminalId on every acquire; no-op until wired. */ private final List> acquireListeners = new java.util.concurrent.CopyOnWriteArrayList<>(); @@ -55,31 +56,37 @@ public final class SessionManager implements TurnListener { /** Backward-compatible constructor: shared-tree sessions, production git seam. */ public SessionManager(PeerLauncher launcher) { - this(launcher, new GitWorktrees(), System::nanoTime, 0); + this(launcher, new GitWorktrees(), System::nanoTime, 0, false); } /** Backward-compatible constructor with an injectable worktree seam. */ public SessionManager(PeerLauncher launcher, Worktrees worktrees) { - this(launcher, worktrees, System::nanoTime, 0); + this(launcher, worktrees, System::nanoTime, 0, false); } /** Test constructor with an injectable clock. */ public SessionManager(PeerLauncher launcher, Worktrees worktrees, LongSupplier nowNanos) { - this(launcher, worktrees, nowNanos, 0); + this(launcher, worktrees, nowNanos, 0, false); } /** Production constructor with a configured context turn cap. */ public SessionManager(PeerLauncher launcher, Worktrees worktrees, int contextCap) { - this(launcher, worktrees, System::nanoTime, contextCap); + this(launcher, worktrees, System::nanoTime, contextCap, false); } public SessionManager(PeerLauncher launcher, Worktrees worktrees, LongSupplier nowNanos, - int contextCap) { + int contextCap) { + this(launcher, worktrees, nowNanos, contextCap, false); + } + + public SessionManager(PeerLauncher launcher, Worktrees worktrees, LongSupplier nowNanos, + int contextCap, boolean clearAfterTurn) { this.launcher = launcher; this.worktrees = worktrees; this.presence = new PresenceBridge(this); this.nowNanos = nowNanos; this.contextCap = contextCap; + this.clearAfterTurn = clearAfterTurn; } /** @@ -355,8 +362,25 @@ public final class SessionManager implements TurnListener { /** Lifecycle hook: the worker's delegated turn completed successfully. */ @Override public void onTurnComplete(String target) { + completeTurn(target, false); + } + + @Override + public boolean hasPostTurnAction(String target) { + if (!clearAfterTurn) return false; WorkerSession current = findByTerminal(target); - if (current == null || current.state() != WorkerSession.State.BUSY) return; + return current != null && current.state() == WorkerSession.State.BUSY + && (contextCap <= 0 || current.turnCount() < contextCap); + } + + @Override + public boolean onTurnCompleteWithPostAction(String target) { + return completeTurn(target, true); + } + + private boolean completeTurn(String target, boolean startContextReset) { + WorkerSession current = findByTerminal(target); + if (current == null || current.state() != WorkerSession.State.BUSY) return false; long now = nowNanos.getAsLong(); WorkerSession updated = current.withState(WorkerSession.State.DONE).withActivity(now); if (replace(current, updated)) { @@ -365,6 +389,15 @@ public final class SessionManager implements TurnListener { } if (contextCap > 0 && updated.turnCount() >= contextCap) { release(current.paneId()); + return false; + } + if (!startContextReset || !clearAfterTurn) return false; + try { + return launcher.clearContext(current.paneId()); + } catch (RuntimeException e) { + log.warn("context reset failed for terminal={} pane={}; continuing without reset: {}", + target, current.paneId(), e.getMessage()); + return false; } } diff --git a/bridged/src/main/java/dev/ltms/bridged/worker/ClaudeCodeLauncher.java b/bridged/src/main/java/dev/ltms/bridged/worker/ClaudeCodeLauncher.java index 93c2cdb..6a721ac 100644 --- a/bridged/src/main/java/dev/ltms/bridged/worker/ClaudeCodeLauncher.java +++ b/bridged/src/main/java/dev/ltms/bridged/worker/ClaudeCodeLauncher.java @@ -126,7 +126,7 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher { putIfPresent(workerEnv, "ANTHROPIC_AUTH_TOKEN", env.apply(cfg.tokenEnv())); applyGitToken(workerEnv, cfg); - return new Launch(workerEnv, argvWithBridge(cfg)); + return new Launch(workerEnv, argvWithModel(argvWithBridge(cfg), cfg)); } /** @@ -149,6 +149,32 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher { return argv; } + /** + * Pin the model on the command line as well as in {@code ANTHROPIC_MODEL} (CB-533). + * + *

    The env var alone is not a reliable pin for this adapter, because the argv is usually a + * launcher rather than {@code claude} itself — {@code ["ccs", ""]} — and {@code ccs} + * exports its profile's own model family ({@code ANTHROPIC_MODEL}, {@code DEFAULT_OPUS/SONNET/ + * HAIKU}, {@code CLAUDE_CODE_SUBAGENT_MODEL}) over whatever it inherited. A worker profile that + * set {@code model:} therefore got silently overruled by its own launcher. Claude Code's + * {@code --model} flag outranks the environment, and {@code ccs [claude-args...]} + * passes trailing arguments through, so the flag survives the wrapper. + * + *

    Appended last so it also outranks anything in the operator's own {@code argv}. Profiles + * that deliberately leave {@code model:} unset (letting {@code ccs} own model selection, as + * {@code gx10} does) are untouched — this adds nothing when there is nothing to add. This is + * the {@code kind: claude} counterpart of the opencode adapter's {@code -m provider/model}. + */ + private static List argvWithModel(List argv, BridgedConfig.Worker cfg) { + if (cfg.model() == null || cfg.model().isBlank()) { + return argv; + } + List withModel = mutableArgv(argv); + withModel.add("--model"); + withModel.add(cfg.model()); + return withModel; + } + // --- Agent-returning convenience spawns (used by callers/tests that want the herdr Agent) --- /** Spawn a worker for the default profile in the resolved default cwd. */ @@ -170,13 +196,25 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher { @Override public Set capabilities() { - Set caps = EnumSet.of(Capability.MID_TURN_ASK, Capability.WORKTREE, Capability.ORPHAN_REAP); + Set caps = EnumSet.of(Capability.MID_TURN_ASK, Capability.WORKTREE, + Capability.CONTEXT_RESET, Capability.ORPHAN_REAP); if (hasGitTokenProfile()) { caps.add(Capability.SELF_PR); } return Set.copyOf(caps); } + @Override + public boolean clearContext(String id) { + String target = agentTarget(id); + if (target == null) { + return false; + } + // This deliberately bypasses Injector: /clear is housekeeping, not a delegated turn. + agents().send(target, "/clear"); + return true; + } + /** Whether any configured profile opts into a git-forge token (required for {@link Capability#SELF_PR}). */ private boolean hasGitTokenProfile() { return profileConfigs().stream().anyMatch(BridgedConfig.Worker::hasGitToken); diff --git a/bridged/src/main/java/dev/ltms/bridged/worker/CompositePeerLauncher.java b/bridged/src/main/java/dev/ltms/bridged/worker/CompositePeerLauncher.java index 3df94b7..5e35025 100644 --- a/bridged/src/main/java/dev/ltms/bridged/worker/CompositePeerLauncher.java +++ b/bridged/src/main/java/dev/ltms/bridged/worker/CompositePeerLauncher.java @@ -217,6 +217,16 @@ public final class CompositePeerLauncher implements PeerLauncher { d.stop(id); } + @Override + public boolean clearContext(String id) { + HerdrPeerLauncher delegate = spawnedBy.get(id); + if (delegate == null) { + log.debug("clearContext({}) ignored — no recorded owning adapter", id); + return false; + } + return delegate.clearContext(id); + } + @Override public Set profiles() { return byProfile.keySet(); diff --git a/bridged/src/main/java/dev/ltms/bridged/worker/HerdrPeerLauncher.java b/bridged/src/main/java/dev/ltms/bridged/worker/HerdrPeerLauncher.java index 0958bef..a3a934e 100644 --- a/bridged/src/main/java/dev/ltms/bridged/worker/HerdrPeerLauncher.java +++ b/bridged/src/main/java/dev/ltms/bridged/worker/HerdrPeerLauncher.java @@ -24,6 +24,7 @@ import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicLong; import java.util.function.Function; import java.util.function.LongSupplier; @@ -90,6 +91,7 @@ public abstract class HerdrPeerLauncher implements PeerLauncher { // exact pane it must tear down. The pane id is launcher-private (never the routing key) — see // PeerHandle.id(). private final ConcurrentMap paneByAgentId = new ConcurrentHashMap<>(); + private final AtomicBoolean resetUnsupportedLogged = new AtomicBoolean(); /** * @param namePrefix label prefix for this peer kind (drives naming and reap) @@ -128,6 +130,25 @@ public abstract class HerdrPeerLauncher implements PeerLauncher { */ protected abstract Launch buildLaunch(BridgedConfig.Worker cfg); + /** Direct transport access for peer-specific, non-turn control operations. */ + protected final AgentControl agents() { + return agents; + } + + /** Resolve the public peer id to the launcher's private herdr target. */ + protected final String agentTarget(String id) { + return paneByAgentId.get(id); + } + + @Override + public boolean clearContext(String id) { + if (resetUnsupportedLogged.compareAndSet(false, true)) { + log.warn("context reset is unsupported for peer kind {}; clearAfterTurn is a no-op", + namePrefix); + } + return false; + } + /** A peer-specific launch: the herdr {@code env} map and {@code argv}. */ protected record Launch(Map env, List argv) { } diff --git a/bridged/src/test/java/dev/ltms/bridged/BridgedDeliverabilityTest.java b/bridged/src/test/java/dev/ltms/bridged/BridgedDeliverabilityTest.java new file mode 100644 index 0000000..fcdeaab --- /dev/null +++ b/bridged/src/test/java/dev/ltms/bridged/BridgedDeliverabilityTest.java @@ -0,0 +1,86 @@ +package dev.ltms.bridged; + +import dev.ltms.bridged.inject.WorkerPresence; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; + +import java.util.HashMap; +import java.util.Map; +import java.util.function.Predicate; +import java.util.function.Supplier; + +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * CB-534: the injector's readiness gate must open for a lead as well as for a present worker. + * + *

    The bug these cover was silent and slow: a lead was never marked present (only workers are), so + * every lead→lead delivery sat on the gate for the full readiness grace and failed ~60s later without + * a keystroke ever reaching the pane. + */ +class BridgedDeliverabilityTest { + + private static Supplier> leads(Map m) { + return () -> m; + } + + @Test + @DisplayName("a worker that has connected its MCP is deliverable") + void presentWorkerIsDeliverable() { + WorkerPresence presence = new WorkerPresence(); + presence.markPresent("term_worker"); + + assertTrue(Bridged.deliverableTo(presence, leads(Map.of())).test("term_worker")); + } + + @Test + @DisplayName("a worker still in its boot window is held back") + void absentWorkerIsNotDeliverable() { + assertFalse(Bridged.deliverableTo(new WorkerPresence(), leads(Map.of())).test("term_booting")); + } + + @Test + @DisplayName("a lead is deliverable without ever being marked present") + void leadIsDeliverableWithoutPresence() { + WorkerPresence presence = new WorkerPresence(); + Predicate deliverable = + Bridged.deliverableTo(presence, leads(Map.of("term_lead", "opus-5.0"))); + + assertFalse(presence.isPresent("term_lead"), "a lead is never enrolled in worker presence"); + assertTrue(deliverable.test("term_lead"), "…and must be deliverable anyway"); + } + + @Test + @DisplayName("an unknown terminal is deliverable to neither") + void strangerIsNotDeliverable() { + WorkerPresence presence = new WorkerPresence(); + presence.markPresent("term_worker"); + + assertFalse(Bridged.deliverableTo(presence, leads(Map.of("term_lead", "opus-5.0"))) + .test("term_stranger")); + } + + @Test + @DisplayName("a lead discovered after startup becomes deliverable with no restart") + void leadSetIsReadThroughOnEveryCall() { + Map discovered = new HashMap<>(); + Predicate deliverable = Bridged.deliverableTo(new WorkerPresence(), leads(discovered)); + + assertFalse(deliverable.test("term_late")); + discovered.put("term_late", "gpt-sol-5.6"); // leadScan picks up a newly labelled tab + assertTrue(deliverable.test("term_late"), "the supplier must be re-read, not snapshotted"); + } + + @Test + @DisplayName("forgetting a torn-down worker does not strip a lead of its deliverability") + void forgetDoesNotDisarmALead() { + WorkerPresence presence = new WorkerPresence(); + Predicate deliverable = + Bridged.deliverableTo(presence, leads(Map.of("term_lead", "opus-5.0"))); + + presence.forget("term_lead"); // the injector's cleanup path runs against every target + + assertTrue(deliverable.test("term_lead")); + } +} diff --git a/bridged/src/test/java/dev/ltms/bridged/auth/CallerResolverTest.java b/bridged/src/test/java/dev/ltms/bridged/auth/CallerResolverTest.java index 425ade3..bea0417 100644 --- a/bridged/src/test/java/dev/ltms/bridged/auth/CallerResolverTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/auth/CallerResolverTest.java @@ -5,6 +5,8 @@ import dev.ltms.bridged.herdr.PaneLocator; import dev.ltms.bridged.mcp.ConnectionIdentity; import org.junit.jupiter.api.Test; +import java.util.Map; + import static org.junit.jupiter.api.Assertions.*; /** @@ -48,7 +50,7 @@ class CallerResolverTest { void aPinnedPrimaryTerminalResolvesToPrimaryNotWorker() { // The primary's own session lives in a herdr pane (term_a here). Without the pin the pane // match wins and the primary is locked out of spawn/send/stop as a misread worker. - Principal p = new CallerResolver(workerIdentity(), false, null, "term_a") + Principal p = CallerResolver.pinnedTo(workerIdentity(), false, null, "term_a") .resolve("127.0.0.1", 42, null); assertEquals(Role.PRIMARY, p.role()); @@ -56,7 +58,7 @@ class CallerResolverTest { @Test void aPinnedPrimaryTerminalNeedsNoTokenEvenInTokenMode() { - Principal p = new CallerResolver(workerIdentity(), true, "s3cret", "term_a") + Principal p = CallerResolver.pinnedTo(workerIdentity(), true, "s3cret", "term_a") .resolve("127.0.0.1", 42, null); assertEquals(Role.PRIMARY, p.role(), @@ -65,7 +67,7 @@ class CallerResolverTest { @Test void otherPanesRemainWorkersWhenAPinIsSet() { - Principal p = new CallerResolver(workerIdentity(), false, null, "term_someone_else") + Principal p = CallerResolver.pinnedTo(workerIdentity(), false, null, "term_someone_else") .resolve("127.0.0.1", 42, null); assertEquals(Role.WORKER, p.role()); @@ -76,9 +78,9 @@ class CallerResolverTest { @Test void aBlankPinLeavesWorkerResolutionUntouched() { assertEquals(Role.WORKER, - new CallerResolver(workerIdentity(), false, null, " ").resolve("127.0.0.1", 42, null).role()); + CallerResolver.pinnedTo(workerIdentity(), false, null, " ").resolve("127.0.0.1", 42, null).role()); assertEquals(Role.WORKER, - new CallerResolver(workerIdentity(), false, null, null).resolve("127.0.0.1", 42, null).role()); + CallerResolver.pinnedTo(workerIdentity(), false, null, null).resolve("127.0.0.1", 42, null).role()); } @Test @@ -133,6 +135,169 @@ class CallerResolverTest { assertEquals(Role.ANONYMOUS, p.role()); } + // ── CB-530: the leaders registry ──────────────────────────────────────────────────────────── + // The fake resolves exactly one pane (term_a) from a PID, so "two leads both resolve" is + // asserted at the config layer (BridgedConfigTest#leaderTerminals…). What matters here is that + // resolution is a REGISTRY LOOKUP rather than a single equality test against one pin. + + @Test + void aRegisteredLeadPaneResolvesToPrimaryCarryingItsName() { + Principal p = new CallerResolver(workerIdentity(), false, null, Map.of("term_a", "opus-5.0")) + .resolve("127.0.0.1", 42, null); + + assertEquals(Role.PRIMARY, p.role()); + assertEquals("opus-5.0", p.name(), "whoami must be able to say WHICH lead is asking"); + assertEquals("term_a", p.terminal(), + "CB-532: a lead carries the pane it was matched by. Without it ownsSession() can " + + "never be true for a lead, so it can send to a peer but never answer one"); + } + + @Test + void aSecondLeadIsRecognisedRatherThanSilentlyDemoted() { + // The regression this feature exists for: with a singular pin, whichever lead was not the + // pin resolved as a worker and was refused every orchestration call. + Principal p = new CallerResolver(workerIdentity(), false, null, + Map.of("term_elsewhere", "gpt-sol-5.6", "term_a", "opus-5.0")) + .resolve("127.0.0.1", 42, null); + + assertEquals(Role.PRIMARY, p.role()); + assertEquals("opus-5.0", p.name()); + } + + @Test + void aPaneAbsentFromTheRegistryIsStillAWorker() { + Principal p = new CallerResolver(workerIdentity(), false, null, + Map.of("term_elsewhere", "gpt-sol-5.6")) + .resolve("127.0.0.1", 42, null); + + assertEquals(Role.WORKER, p.role()); + assertEquals("term_a", p.terminal()); + assertNull(p.name()); + } + + @Test + void aRegisteredLeadNeedsNoTokenEvenInTokenMode() { + Principal p = new CallerResolver(workerIdentity(), true, "s3cret", Map.of("term_a", "opus")) + .resolve("127.0.0.1", 42, null); + + assertEquals(Role.PRIMARY, p.role(), + "the pane mapping is as unforgeable as a worker's — it outranks the token path"); + assertEquals("opus", p.name()); + } + + /** The pre-CB-530 spelling must keep working, exactly, including for configs that never migrate. */ + @Test + void theLegacySinglePinBehavesAsALeadNamedPrimary() { + Principal p = CallerResolver.pinnedTo(workerIdentity(), false, null, "term_a") + .resolve("127.0.0.1", 42, null); + + assertEquals(Role.PRIMARY, p.role()); + assertEquals("primary", p.name()); + } + + @Test + void anEmptyRegistryLeavesEveryPaneAWorker() { + Map noLeads = null; + assertEquals(Role.WORKER, + new CallerResolver(workerIdentity(), false, null, Map.of()) + .resolve("127.0.0.1", 42, null).role()); + assertEquals(Role.WORKER, + new CallerResolver(workerIdentity(), false, null, noLeads) + .resolve("127.0.0.1", 42, null).role()); + // CB-531: and the same for the live-registry form, whose supplier may also be absent. + assertEquals(Role.WORKER, + CallerResolver.withLeads(workerIdentity(), false, null, null) + .resolve("127.0.0.1", 42, null).role()); + } + + /** The audit line must distinguish leads once several exist, or a log says nothing useful. */ + @Test + void describeNamesTheLeadButStillReadsPrimaryWhenUnnamed() { + assertEquals("leader:opus-5.0", Principal.leader("opus-5.0", "term_a", 1).describe()); + assertEquals("primary", Principal.primary(1).describe()); + assertEquals("worker:term_a", Principal.worker("term_a", 1).describe()); + } + + // ── CB-532: a lead is an addressable peer, not only a sender ──────────────────────────────── + + /** + * The regression this ticket exists for: two leads could both be recognised (CB-530/531) and + * still not converse, because REPLY is gated on ownsSession() and a lead owned nothing. + */ + @Test + void aLeadOwnsItsOwnPaneSoItMayAnswerAPeer() { + Principal lead = new CallerResolver(workerIdentity(), false, null, Map.of("term_a", "opus-5.0")) + .resolve("127.0.0.1", 42, null); + + assertTrue(lead.ownsSession("term_a")); + assertTrue(Authz.permits(lead, Authz.Action.REPLY, "term_a"), + "a lead answering a peer replies for its OWN terminal — the rendezvous the sender " + + "opened is keyed on exactly that"); + assertTrue(Authz.permits(lead, Authz.Action.ASK, "term_a")); + } + + @Test + void aLeadStillCannotActAsAnyoneElse() { + Principal lead = new CallerResolver(workerIdentity(), false, null, Map.of("term_a", "opus-5.0")) + .resolve("127.0.0.1", 42, null); + + assertFalse(lead.ownsSession("term_someone_else")); + assertFalse(Authz.permits(lead, Authz.Action.REPLY, "term_someone_else"), + "widening WHO may reply must not widen WHAT they may reply as"); + } + + /** A primary with no pane — token mode, or off-host — owns nothing and must stay a sender only. */ + @Test + void anUnnamedPrimaryWithNoPaneOwnsNothing() { + Principal p = new CallerResolver(nonWorkerIdentity(), true, "s3cret") + .resolve("127.0.0.1", 99, "Bearer s3cret"); + + assertEquals(Role.PRIMARY, p.role()); + assertNull(p.terminal()); + assertFalse(p.ownsSession(null), "a null terminal must never match a null session id"); + assertFalse(Authz.permits(p, Authz.Action.REPLY, null)); + } + + @Test + void aLeadKeepsEveryOrchestrationRightItAlreadyHad() { + Principal lead = new CallerResolver(workerIdentity(), false, null, Map.of("term_a", "opus-5.0")) + .resolve("127.0.0.1", 42, null); + + assertTrue(Authz.permits(lead, Authz.Action.SPAWN, null)); + assertTrue(Authz.permits(lead, Authz.Action.SEND, "term_worker")); + assertTrue(Authz.permits(lead, Authz.Action.STOP, null)); + assertTrue(Authz.permits(lead, Authz.Action.DRAIN, null)); + } + + /** + * CB-531: the registry is read per resolve, not snapshotted at construction — a lead that + * labels its tab after the daemon booted is recognised without a restart. + */ + @Test + void aLeadRegisteredAfterConstructionIsHonouredWithoutRebuildingTheResolver() { + Map live = new java.util.HashMap<>(); + CallerResolver r = CallerResolver.withLeads(workerIdentity(), false, null, () -> live); + + assertEquals(Role.WORKER, r.resolve("127.0.0.1", 42, null).role()); + + live.put("term_a", "gpt-sol-5.6"); // the scanner sees a newly-labelled tab + + Principal p = r.resolve("127.0.0.1", 42, null); + assertEquals(Role.PRIMARY, p.role()); + assertEquals("gpt-sol-5.6", p.name()); + } + + /** The map form must stay a snapshot: a caller handing over a map is not offering live state. */ + @Test + void theMapFormIsCopiedSoLaterMutationCannotGrantLeadership() { + Map mutable = new java.util.HashMap<>(); + CallerResolver r = new CallerResolver(workerIdentity(), false, null, mutable); + + mutable.put("term_a", "sneaky"); + + assertEquals(Role.WORKER, r.resolve("127.0.0.1", 42, null).role()); + } + @Test void tokenModeRequiresANonEmptyConfiguredToken() { ConnectionIdentity id = nonWorkerIdentity(); diff --git a/bridged/src/test/java/dev/ltms/bridged/config/BridgedConfigTest.java b/bridged/src/test/java/dev/ltms/bridged/config/BridgedConfigTest.java index 2960b70..8ed6587 100644 --- a/bridged/src/test/java/dev/ltms/bridged/config/BridgedConfigTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/config/BridgedConfigTest.java @@ -5,6 +5,8 @@ import org.junit.jupiter.api.io.TempDir; import java.nio.file.Files; import java.nio.file.Path; +import java.util.List; +import java.util.Map; import java.util.Set; import static org.junit.jupiter.api.Assertions.*; @@ -45,6 +47,7 @@ class BridgedConfigTest { assertEquals(9000, cfg.bind().port()); assertNotNull(cfg.guard(), "guard must default to empty, never null"); assertTrue(cfg.guard().offSubscriptionHosts().isEmpty()); + assertFalse(cfg.lifecycle().clearAfterTurn(), "context clearing is opt-in"); } @Test @@ -96,6 +99,233 @@ class BridgedConfigTest { assertDoesNotThrow(() -> BridgedConfig.load(f)); } + /** + * 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 + * indistinguishable from one that works: that is exactly how a hand-written `leaders:` registry + * came to look configured while being inert. + */ + @Test + void unknownTopLevelKeysAreNamedSoADroppedBlockCannotLookLikeAWorkingOne() { + assertEquals(List.of("futureFeature", "leedars"), + BridgedConfig.unknownTopLevelKeys( + "bind:\n port: 8080\nleedars:\n a: b\nfutureFeature: true\n"), + "a typo'd key is the common case and must be reported by name"); + } + + @Test + void everyKeyThisBuildUnderstandsIsAbsentFromTheUnknownList() { + assertTrue(BridgedConfig.unknownTopLevelKeys(""" + bind: + port: 8080 + herdrSocket: /tmp/s + workers: {} + defaultWorker: a + guard: {} + worktreeRoot: /tmp + lifecycle: {} + spawnReadyTimeoutMs: 1 + spawnReadyPollMs: 1 + broker: {} + primary: {} + leaders: {} + leadScan: {} + placement: fixed + auth: {} + """).isEmpty(), "the known-key set must not drift from the record components"); + } + + @Test + void aMalformedOrEmptyDocumentIsNotReportedAsUnknownKeys() { + assertTrue(BridgedConfig.unknownTopLevelKeys("").isEmpty()); + assertTrue(BridgedConfig.unknownTopLevelKeys("just a scalar").isEmpty()); + } + + // ── CB-531: lead discovery by tab label ───────────────────────────────────────────────────── + + @Test + void leadScanIsOffUnlessTheBlockIsPresent(@TempDir Path dir) throws Exception { + Path f = dir.resolve("no-scan.yaml"); + Files.writeString(f, "bind:\n port: 8080\n"); + + assertNull(BridgedConfig.load(f).leadScan(), + "turning this on widens who resolves as PRIMARY — upgrading the daemon must not do that"); + } + + @Test + void leadScanDefaultsItsFieldsWhenTheBlockIsPresentButBare(@TempDir Path dir) throws Exception { + Path f = dir.resolve("bare-scan.yaml"); + Files.writeString(f, "bind:\n port: 8080\nleadScan: {}\n"); + + BridgedConfig.LeadScan scan = BridgedConfig.load(f).leadScan(); + assertEquals("lead:", scan.tabPrefix()); + assertEquals(10, scan.intervalSeconds()); + } + + @Test + void leadScanReadsAnExplicitPrefixAndInterval(@TempDir Path dir) throws Exception { + Path f = dir.resolve("scan.yaml"); + Files.writeString(f, """ + bind: + port: 8080 + leadScan: + tabPrefix: "drive:" + intervalSeconds: 30 + """); + + BridgedConfig.LeadScan scan = BridgedConfig.load(f).leadScan(); + assertEquals("drive:", scan.tabPrefix()); + assertEquals(30, scan.intervalSeconds()); + } + + /** + * The hazard the guard exists for: bridged writes worker tab labels and reads lead tab labels. + * Overlap the two and every worker it spawns is read back as a lead. + */ + @Test + void aLeadPrefixThatAWorkerTabLabelAlsoMatchesRefusesToStart(@TempDir Path dir) throws Exception { + Path f = dir.resolve("collide.yaml"); + Files.writeString(f, """ + bind: + port: 8080 + workers: + gx10: + tabLabel: "lead: {profile} #{n}" + leadScan: + tabPrefix: "lead:" + """); + BridgedConfig cfg = BridgedConfig.load(f); + + IllegalStateException e = assertThrows(IllegalStateException.class, cfg::validateLeadScan); + assertTrue(e.getMessage().contains("gx10"), "the message must name the offending profile"); + } + + @Test + void theDefaultWorkerTabLabelDoesNotCollideWithTheDefaultLeadPrefix(@TempDir Path dir) throws Exception { + Path f = dir.resolve("ok.yaml"); + Files.writeString(f, """ + bind: + port: 8080 + workers: + gx10: + baseUrl: http://gx00.gw:8000 + leadScan: {} + """); + + assertDoesNotThrow(() -> BridgedConfig.load(f).validateLeadScan()); + } + + @Test + void theCollisionGuardIsANoOpWhenScanningIsOff(@TempDir Path dir) throws Exception { + Path f = dir.resolve("off.yaml"); + Files.writeString(f, """ + bind: + port: 8080 + workers: + gx10: + tabLabel: "lead: {profile}" + """); + + assertDoesNotThrow(() -> BridgedConfig.load(f).validateLeadScan(), + "a label that collides with a convention nobody reads is not a problem"); + } + + // ── CB-530: the leaders registry ──────────────────────────────────────────────────────────── + + @Test + void leadersBlockRegistersEveryPaneByName(@TempDir Path dir) throws Exception { + Path f = dir.resolve("leaders.yaml"); + Files.writeString(f, """ + bind: + port: 8080 + leaders: + opus-5.0: + terminal: term_opus + kind: claude + gpt-sol-5.6: + terminal: term_sol + kind: opencode + model: openai/gpt-5.6-terra + """); + + BridgedConfig cfg = BridgedConfig.load(f); + + assertEquals(Set.of("opus-5.0", "gpt-sol-5.6"), cfg.leaders().keySet()); + assertEquals("opencode", cfg.leaders().get("gpt-sol-5.6").kind()); + assertEquals("openai/gpt-5.6-terra", cfg.leaders().get("gpt-sol-5.6").model()); + // The whole point: BOTH panes resolve as leads, so neither is demoted to worker. + assertEquals(Map.of("term_opus", "opus-5.0", "term_sol", "gpt-sol-5.6"), + cfg.leaderTerminals()); + } + + @Test + void aLegacyPrimaryPinAloneStillRegistersAsALeadNamedPrimary(@TempDir Path dir) throws Exception { + Path f = dir.resolve("legacy-pin.yaml"); + Files.writeString(f, "bind:\n port: 8080\nprimary:\n terminal: term_fixed\n"); + + assertEquals(Map.of("term_fixed", "primary"), BridgedConfig.load(f).leaderTerminals(), + "configs that never migrate must behave exactly as they did before CB-530"); + } + + @Test + void anExplicitLeadersEntryWinsOverThePinForTheSameTerminal(@TempDir Path dir) throws Exception { + Path f = dir.resolve("both.yaml"); + Files.writeString(f, """ + bind: + port: 8080 + primary: + terminal: term_shared + leaders: + opus-5.0: + terminal: term_shared + """); + + assertEquals(Map.of("term_shared", "opus-5.0"), BridgedConfig.load(f).leaderTerminals(), + "the pin is the older spelling of the same fact; the named entry is what was meant"); + } + + @Test + void bothBlocksTogetherRegisterTheUnionOfTheirTerminals(@TempDir Path dir) throws Exception { + Path f = dir.resolve("union.yaml"); + Files.writeString(f, """ + bind: + port: 8080 + primary: + terminal: term_pinned + leaders: + gpt-sol-5.6: + terminal: term_sol + """); + + assertEquals(Map.of("term_pinned", "primary", "term_sol", "gpt-sol-5.6"), + BridgedConfig.load(f).leaderTerminals()); + } + + @Test + void neitherBlockLeavesNothingRegistered(@TempDir Path dir) throws Exception { + Path f = dir.resolve("none.yaml"); + Files.writeString(f, "bind:\n port: 8080\n"); + + assertTrue(BridgedConfig.load(f).leaderTerminals().isEmpty()); + } + + /** A lead entry with no terminal identifies nothing — it must not register a null key. */ + @Test + void aLeadWithoutATerminalIsNotRegistered(@TempDir Path dir) throws Exception { + Path f = dir.resolve("no-terminal.yaml"); + Files.writeString(f, """ + bind: + port: 8080 + leaders: + sketch: + kind: opencode + real: + terminal: term_real + """); + + assertEquals(Map.of("term_real", "real"), BridgedConfig.load(f).leaderTerminals()); + } + @Test void absentBrokerBlockLeavesInboxSoftState(@TempDir Path dir) throws Exception { Path f = dir.resolve("no-broker.yaml"); @@ -345,6 +575,7 @@ class BridgedConfigTest { idleTtlSeconds: 300 contextCap: 10 drainTimeoutSeconds: 5 + clearAfterTurn: true broker: uri: amqp://guest:guest@127.0.0.1:5672 primary: @@ -371,6 +602,7 @@ class BridgedConfigTest { assertEquals(300, cfg.lifecycle().idleTtlSeconds()); assertEquals(10, cfg.lifecycle().contextCap()); assertEquals(5, cfg.lifecycle().drainTimeoutSeconds()); + assertTrue(cfg.lifecycle().clearAfterTurn()); assertEquals("amqp://guest:guest@127.0.0.1:5672", cfg.broker().uri()); assertEquals("term_abc123", cfg.primary().terminal()); assertEquals(5, cfg.primary().remindersOrDefault()); diff --git a/bridged/src/test/java/dev/ltms/bridged/herdr/LeadTabScannerTest.java b/bridged/src/test/java/dev/ltms/bridged/herdr/LeadTabScannerTest.java new file mode 100644 index 0000000..8ccb2d6 --- /dev/null +++ b/bridged/src/test/java/dev/ltms/bridged/herdr/LeadTabScannerTest.java @@ -0,0 +1,262 @@ +package dev.ltms.bridged.herdr; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import org.junit.jupiter.api.Test; + +import java.util.ArrayList; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; + +import static org.junit.jupiter.api.Assertions.*; + +/** + * CB-531. A lead is never spawned, so the daemon has to find it: these assert that an + * operator-labelled tab is what makes a pane a lead, and — just as importantly — what does not. + */ +class LeadTabScannerTest { + + private static final ObjectMapper MAPPER = new ObjectMapper(); + private static final long TTL = TimeUnit.SECONDS.toNanos(10); + + /** + * A herdr whose workspace/tab/pane topology is declared per test. Counts calls so the caching + * contract can be asserted, and can be made to fail on demand. + */ + private static final class TopologyHerdr implements HerdrClient { + /** workspace_id → label. */ + final Map workspaces = new LinkedHashMap<>(); + /** tab_id → [workspace_id, label]. */ + final Map tabs = new LinkedHashMap<>(); + /** pane_id → [tab_id, terminal_id]. */ + final Map panes = new LinkedHashMap<>(); + int calls; + boolean failing; + + TopologyHerdr workspace(String id, String label) { + workspaces.put(id, label); + return this; + } + + TopologyHerdr tab(String tabId, String workspaceId, String label) { + tabs.put(tabId, new String[]{workspaceId, label}); + return this; + } + + TopologyHerdr pane(String paneId, String tabId, String terminalId) { + panes.put(paneId, new String[]{tabId, terminalId}); + return this; + } + + @Override + public JsonNode call(String method, Object params) { + calls++; + if (failing) { + throw new HerdrException("socket closed"); + } + List items = new ArrayList<>(); + switch (method) { + case "workspace.list" -> { + workspaces.forEach((id, label) -> items.add( + "{\"workspace_id\":\"%s\",\"label\":\"%s\"}".formatted(id, label))); + return read("{\"workspaces\":[%s]}".formatted(String.join(",", items))); + } + case "tab.list" -> { + String ws = String.valueOf(((Map) params).get("workspace_id")); + tabs.forEach((id, t) -> { + if (ws.equals(t[0])) { + items.add(("{\"tab_id\":\"%s\",\"workspace_id\":\"%s\",\"label\":%s," + + "\"pane_count\":1}").formatted(id, t[0], + t[1] == null ? "null" : "\"" + t[1] + "\"")); + } + }); + return read("{\"tabs\":[%s]}".formatted(String.join(",", items))); + } + case "pane.list" -> { + panes.forEach((id, p) -> items.add( + "{\"pane_id\":\"%s\",\"tab_id\":\"%s\",\"terminal_id\":\"%s\"}" + .formatted(id, p[0], p[1]))); + return read("{\"panes\":[%s]}".formatted(String.join(",", items))); + } + default -> throw new AssertionError("unexpected herdr call: " + method); + } + } + + private static JsonNode read(String json) { + try { + return MAPPER.readTree(json); + } catch (Exception e) { + throw new AssertionError(e); + } + } + + @Override + public void close() { + } + } + + /** + * The usual shape: one user space with lead tabs, one worker space bridged owns. + * + *

    Not closed: {@code close()} is a no-op on this fake, and every test needs the handle after + * the scanner is built (to mutate the topology or read {@code calls}). + */ + @SuppressWarnings("resource") + private TopologyHerdr twoLeads() { + return new TopologyHerdr() + .workspace("w1", "main") + .workspace("w9", "bridged-workers") + .tab("w1:t1", "w1", "lead: opus-5.0") + .tab("w1:t2", "w1", "lead: gpt-sol-5.6") + .tab("w1:t3", "w1", "notes") + .tab("w9:t1", "w9", "worker: gx10 #1") + .pane("w1:p1", "w1:t1", "term_opus") + .pane("w1:p2", "w1:t2", "term_gpt") + .pane("w1:p3", "w1:t3", "term_notes") + .pane("w9:p1", "w9:t1", "term_worker"); + } + + private LeadTabScanner scanner(TopologyHerdr herdr, Map configured, + AtomicLong clock) { + return new LeadTabScanner(herdr, "lead:", Set.of("bridged-workers"), configured, TTL, + clock::get); + } + + @Test + void everyLabelledTabBecomesALeadNamedByItsLabel() { + Map leads = scanner(twoLeads(), Map.of(), new AtomicLong()).get(); + + assertEquals(Map.of("term_opus", "opus-5.0", "term_gpt", "gpt-sol-5.6"), leads, + "two leads discovered from labels alone — no terminal_id was ever configured"); + } + + @Test + void anUnlabelledTabContributesNothing() { + assertFalse(scanner(twoLeads(), Map.of(), new AtomicLong()).get().containsKey("term_notes")); + } + + /** + * The guard that matters: bridged labels its own worker tabs, so if a worker space were scanned + * a naming accident would promote the fleet. The exclusion is by workspace, not by hoping the + * worker template never collides. + */ + @Test + void aTabInAWorkerSpaceIsNeverALeadEvenWhenItsLabelMatches() { + TopologyHerdr herdr = twoLeads().tab("w9:t2", "w9", "lead: impostor") + .pane("w9:p2", "w9:t2", "term_impostor"); + + assertFalse(scanner(herdr, Map.of(), new AtomicLong()).get().containsKey("term_impostor")); + } + + @Test + void aBarePrefixNamesNobodyAndIsRejected() { + TopologyHerdr herdr = new TopologyHerdr().workspace("w1", "main") + .tab("w1:t1", "w1", "lead:").pane("w1:p1", "w1:t1", "term_a"); + + assertEquals(Map.of(), scanner(herdr, Map.of(), new AtomicLong()).get(), + "a lead with no name would resolve as PRIMARY with nothing to attribute it to"); + } + + @Test + void thePrefixMatchesCaseInsensitivelyAndTheNameIsTrimmed() { + TopologyHerdr herdr = new TopologyHerdr().workspace("w1", "main") + .tab("w1:t1", "w1", " LEAD: opus-5.0 ").pane("w1:p1", "w1:t1", "term_a"); + + assertEquals(Map.of("term_a", "opus-5.0"), scanner(herdr, Map.of(), new AtomicLong()).get()); + } + + @Test + void everyPaneInALeadTabResolvesAsThatLead() { + // A human may split their own lead tab. Both panes are theirs, so both are that lead — + // nothing bridged placed can land here (see the worker-space test above). + TopologyHerdr herdr = twoLeads().pane("w1:p1b", "w1:t1", "term_opus_split"); + + assertEquals("opus-5.0", scanner(herdr, Map.of(), new AtomicLong()).get().get("term_opus_split")); + } + + @Test + void anExplicitlyConfiguredLeadIsMergedInAndOutranksALabel() { + Map configured = Map.of("term_opus", "pinned-name", "term_extra", "from-config"); + + Map leads = scanner(twoLeads(), configured, new AtomicLong()).get(); + + assertEquals("pinned-name", leads.get("term_opus"), "an explicit pin is the operator's last word"); + assertEquals("from-config", leads.get("term_extra"), "a configured lead needs no tab at all"); + assertEquals("gpt-sol-5.6", leads.get("term_gpt")); + } + + // ── caching ───────────────────────────────────────────────────────────────────────────────── + + @Test + void aSecondLookupWithinTheTtlDoesNotTouchHerdr() { + TopologyHerdr herdr = twoLeads(); + AtomicLong clock = new AtomicLong(); + LeadTabScanner s = scanner(herdr, Map.of(), clock); + + s.get(); + int afterFirst = herdr.calls; + clock.addAndGet(TTL - 1); + s.get(); + + assertEquals(afterFirst, herdr.calls, + "resolve() runs on every request — an un-cached scan would put herdr on that path"); + } + + @Test + void aTabLabelledAfterStartupIsPickedUpOnceTheTtlExpires() { + TopologyHerdr herdr = twoLeads(); + AtomicLong clock = new AtomicLong(); + LeadTabScanner s = scanner(herdr, Map.of(), clock); + assertFalse(s.get().containsKey("term_notes")); + + herdr.tab("w1:t3", "w1", "lead: late-arrival"); // the operator renames their tab + clock.addAndGet(TTL); + + assertEquals("late-arrival", s.get().get("term_notes"), + "the whole point over `leaders:`: no config edit, no restart"); + } + + @Test + void aFailedScanKeepsTheLeadsAlreadyKnownRatherThanDemotingThem() { + TopologyHerdr herdr = twoLeads(); + AtomicLong clock = new AtomicLong(); + LeadTabScanner s = scanner(herdr, Map.of(), clock); + Map before = s.get(); + + herdr.failing = true; + clock.addAndGet(TTL); + + assertEquals(before, s.get(), + "a herdr hiccup must not silently demote a live lead to a worker mid-session"); + } + + @Test + void aFailedFirstScanStillHonoursTheConfiguredLeads() { + TopologyHerdr herdr = twoLeads(); + herdr.failing = true; + + Map leads = scanner(herdr, Map.of("term_x", "opus-5.0"), new AtomicLong()).get(); + + assertEquals(Map.of("term_x", "opus-5.0"), leads, + "config-named leads must not depend on herdr answering at all"); + } + + @Test + void aDownHerdrIsRetriedOncePerTtlNotOncePerRequest() { + TopologyHerdr herdr = twoLeads(); + herdr.failing = true; + AtomicLong clock = new AtomicLong(); + LeadTabScanner s = scanner(herdr, Map.of(), clock); + + s.get(); + int afterFirst = herdr.calls; + s.get(); + s.get(); + + assertEquals(afterFirst, herdr.calls, "the failure path must be rate-limited too"); + } +} diff --git a/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java b/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java index c46d5d6..4b67084 100644 --- a/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java @@ -156,6 +156,21 @@ class CompletionResolverTest { assertEquals("No, 391 = 17 × 23.", waiter.getNow(null).text()); } + @Test + void resolvesSynchronouslyBeforePostTurnContextClearing() { + FakeHerdr herdr = new FakeHerdr().readText("⏺ previous answer\n❯ "); + Rendezvous rendezvous = new Rendezvous(); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous); + var waiter = rendezvous.open("term_a"); + resolver.captureBaseline("term_a"); + herdr.readText("⏺ answer that /clear would erase\n❯ "); + + resolver.resolveBeforePostAction("term_a"); + + assertTrue(waiter.isDone(), "the answer is captured before the adapter sends /clear"); + assertEquals("answer that /clear would erase", waiter.getNow(null).text()); + } + @Test void suppressesAnUnchangedCompletionEvenWhenTheBlockExceedsTheScrapeCap() { // The fan-out issue-hunt finding: captureBaseline once stored the RAW (unclipped) assistant diff --git a/bridged/src/test/java/dev/ltms/bridged/inject/InjectorTest.java b/bridged/src/test/java/dev/ltms/bridged/inject/InjectorTest.java index 0d7cffa..86534ca 100644 --- a/bridged/src/test/java/dev/ltms/bridged/inject/InjectorTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/inject/InjectorTest.java @@ -196,6 +196,47 @@ class InjectorTest { assertEquals(List.of(T), completed, "a confirmed working→idle fires exactly one completion"); } + @Test + void postTurnResetSettlesBeforeNextDelegationWithoutBecomingATurn() { + AgentControl agents = new AgentControl(herdr); + class ResetListener implements TurnListener { + int completed; + + @Override + public void onTurnComplete(String target) { + completed++; + } + + @Override + public boolean hasPostTurnAction(String target) { + return true; + } + + @Override + public boolean onTurnCompleteWithPostAction(String target) { + completed++; + agents.send(target, "/clear"); // direct housekeeping, never Injector.enqueue + return true; + } + } + ResetListener listener = new ResetListener(); + Injector inj = new Injector(agents, listener); + inj.enqueue(T, "first"); + inj.enqueue(T, "second"); + + inj.onStatus(T, AgentStatus.IDLE); // first delegation + inj.onStatus(T, AgentStatus.WORKING); + inj.onStatus(T, AgentStatus.IDLE); // first complete; reset dispatched + assertEquals(List.of("first", "/clear"), sent(), + "same-tick completion must not let the queued delegation overtake reset"); + + inj.onStatus(T, AgentStatus.WORKING); // reset picked up, but this is not a bridge turn + inj.onStatus(T, AgentStatus.IDLE); // reset settled; second may now deliver + assertEquals(List.of("first", "/clear", "second"), sent()); + assertEquals(1, listener.completed, + "reset settlement must not recursively emit another turn completion"); + } + @Test void doesNotSynthesizeCompletionFromAnUnconfirmedTurn() { List completed = new ArrayList<>(); diff --git a/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java b/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java index a5ad4a9..6e65fa5 100644 --- a/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java @@ -288,7 +288,8 @@ class BridgeMcpTest { WorkerSession s = sessions.acquire("ltms-local", null, "/caller/proj", "term_primary", new WorktreeRequest("cb-304", null)); - McpSchema.CallToolResult res = BridgeMcp.listWorkers(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions); + McpSchema.CallToolResult res = BridgeMcp.listFleet( + workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, Map.of(), ""); assertNotEquals(Boolean.TRUE, res.isError()); String out = textOf(res); @@ -302,6 +303,58 @@ class BridgeMcpTest { assertTrue(out.contains("\"liveStatus\":\"unknown\""), out); } + @Test + void listReportsLeadsAndFlagsTheCallersOwnRow() { + FakeHerdr h = new FakeHerdr(); + SessionManager sessions = new SessionManager( + workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw"))); + + McpSchema.CallToolResult res = BridgeMcp.listFleet( + workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, + Map.of("term_me", "opus-5.0", "term_peer", "gpt-sol-5.6"), "term_me"); + + assertNotEquals(Boolean.TRUE, res.isError()); + String out = textOf(res); + assertTrue(out.contains("\"name\":\"opus-5.0\""), out); + assertTrue(out.contains("\"name\":\"gpt-sol-5.6\""), out); + assertTrue(out.contains("\"sessionId\":\"term_peer\""), out); + // The caller's own row is flagged, and only the caller's — a peer must be distinguishable + // from self without a second bridge_whoami call. + assertEquals(1, out.split("\"self\":true", -1).length - 1, out); + assertTrue(out.indexOf("term_me") < out.indexOf("\"self\":true"), out); + } + + @Test + void listReportsBothHalvesEvenWhenEmpty() { + FakeHerdr h = new FakeHerdr(); + SessionManager sessions = new SessionManager( + workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw"))); + + McpSchema.CallToolResult res = BridgeMcp.listFleet( + workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, Map.of(), ""); + + // An absent "leads" key is what made an empty worker roster read as "no peers" (CB-535). + String out = textOf(res); + assertTrue(out.contains("\"leads\":[]"), out); + assertTrue(out.contains("\"workers\":[]"), out); + } + + @Test + void listReportsALeadHerdrCannotSeeAsUnknown() { + FakeHerdr h = new FakeHerdr(); + SessionManager sessions = new SessionManager( + workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw"))); + + McpSchema.CallToolResult res = BridgeMcp.listFleet( + workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, + Map.of("term_ghost", "gone-away"), "term_me"); + + // Reported, not hidden: an unreachable peer is exactly what a would-be sender needs to see. + String out = textOf(res); + assertTrue(out.contains("\"name\":\"gone-away\""), out); + assertTrue(out.contains("\"status\":\"unknown\""), out); + } + @Test void stopTearsDownAWorkerByPane() { FakeHerdr h = new FakeHerdr(); diff --git a/bridged/src/test/java/dev/ltms/bridged/mcp/PrimaryRegistryTest.java b/bridged/src/test/java/dev/ltms/bridged/mcp/PrimaryRegistryTest.java index 1a26485..0960d68 100644 --- a/bridged/src/test/java/dev/ltms/bridged/mcp/PrimaryRegistryTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/mcp/PrimaryRegistryTest.java @@ -105,4 +105,68 @@ class PrimaryRegistryTest { reg.record("term_found"); assertEquals("term_found", reg.primaryTerminal().get()); } + + // ── CB-532: nudges follow the delegating lead, not "the primary" ──────────────────────────── + + /** + * The bug that made `primary.terminal` unretirable: with one slot, whichever lead called + * bridge_send first captured every nudge — including nudges for the other lead's delegations. + */ + @Test + void aNudgeGoesToTheLeadThatDelegatedToThatWorker() { + var reg = new PrimaryRegistry(null); + reg.recordDelegation("term_worker_a", "term_lead_opus"); + reg.recordDelegation("term_worker_b", "term_lead_sol"); + + assertEquals("term_lead_opus", reg.nudgeTargetFor("term_worker_a").orElseThrow()); + assertEquals("term_lead_sol", reg.nudgeTargetFor("term_worker_b").orElseThrow()); + } + + @Test + void theMostRecentDelegatorWinsWhenAWorkerChangesHands() { + var reg = new PrimaryRegistry(null); + reg.recordDelegation("term_worker", "term_lead_opus"); + reg.recordDelegation("term_worker", "term_lead_sol"); + + assertEquals("term_lead_sol", reg.nudgeTargetFor("term_worker").orElseThrow(), + "the lead waiting on the reply is the one that sent the work most recently"); + } + + /** After a restart the map is empty while the durable inbox still holds the reply. */ + @Test + void anUnknownDelegationFallsBackToThePinnedPrimary() { + var reg = new PrimaryRegistry("term_pinned"); + + assertEquals("term_pinned", reg.nudgeTargetFor("term_never_seen").orElseThrow()); + } + + @Test + void withNoPinAndNoDelegationNobodyIsNudged() { + var reg = new PrimaryRegistry(null); + + assertTrue(reg.nudgeTargetFor("term_worker").isEmpty(), + "guessing would interrupt the wrong lead with someone else's result; the durable " + + "inbox makes pull the correct degradation"); + } + + @Test + void releasingASessionForgetsItsLead() { + var reg = new PrimaryRegistry(null); + reg.recordDelegation("term_worker", "term_lead"); + + reg.forgetDelegation("term_worker"); + + assertTrue(reg.nudgeTargetFor("term_worker").isEmpty()); + } + + @Test + void aBlankOrNullDelegationIsIgnoredRatherThanStored() { + var reg = new PrimaryRegistry(null); + reg.recordDelegation("term_worker", " "); + reg.recordDelegation(null, "term_lead"); + reg.recordDelegation(" ", "term_lead"); + + assertTrue(reg.nudgeTargetFor("term_worker").isEmpty()); + assertTrue(reg.nudgeTargetFor(null).isEmpty()); + } } diff --git a/bridged/src/test/java/dev/ltms/bridged/session/SessionManagerTest.java b/bridged/src/test/java/dev/ltms/bridged/session/SessionManagerTest.java index 640bfa4..f921eee 100644 --- a/bridged/src/test/java/dev/ltms/bridged/session/SessionManagerTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/session/SessionManagerTest.java @@ -39,13 +39,18 @@ class SessionManagerTest { } private SessionManager sessionManager(FakeHerdr herdr, LongSupplier clock, int contextCap) { + return sessionManager(herdr, clock, contextCap, false); + } + + private SessionManager sessionManager(FakeHerdr herdr, LongSupplier clock, int contextCap, + boolean clearAfterTurn) { BridgedConfig.Worker cfg = new BridgedConfig.Worker( "ltms-local", "http://gx00.gw:8000", "coder", null, "BRIDGED_WORKER_TOKEN", List.of("ccs", "ltms-local"), "tab", "bridged-workers", "worker: {profile} #{n}", null, null, null); ClaudeCodeLauncher workers = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr), new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null); - return new SessionManager(workers, new GitWorktrees(), clock, contextCap); + return new SessionManager(workers, new GitWorktrees(), clock, contextCap, clearAfterTurn); } @Test @@ -336,6 +341,50 @@ class SessionManagerTest { "forced release tears the worker pane down exactly once"); } + @Test + void clearAfterTurnResetsContextWithoutDoubleCountingTheTurn() { + FakeHerdr herdr = new FakeHerdr(); + SessionManager sessions = sessionManager(herdr, () -> 0L, 0, true); + WorkerSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary"); + sessions.asPresence().markPresent(session.terminalId()); + + sessions.onDelivered(session.terminalId()); + assertTrue(sessions.onTurnCompleteWithPostAction(session.terminalId())); + + WorkerSession updated = sessions.get(session.paneId()).orElseThrow(); + assertEquals(1, updated.turnCount(), "the reset is housekeeping, not a second delegation"); + assertEquals(List.of("/clear"), promptTexts(herdr)); + } + + @Test + void contextCapReleaseWinsOverClearAfterTurn() { + FakeHerdr herdr = new FakeHerdr(); + SessionManager sessions = sessionManager(herdr, () -> 0L, 1, true); + WorkerSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary"); + sessions.asPresence().markPresent(session.terminalId()); + sessions.onDelivered(session.terminalId()); + + assertFalse(sessions.hasPostTurnAction(session.terminalId()), + "a session at its cap will be released, not reset for reuse"); + assertFalse(sessions.onTurnCompleteWithPostAction(session.terminalId())); + assertTrue(sessions.get(session.paneId()).isEmpty()); + assertTrue(promptTexts(herdr).isEmpty(), "never send /clear into a worker being torn down"); + } + + @Test + void clearAfterTurnFalsePreservesCompletionWithoutAControlPrompt() { + FakeHerdr herdr = new FakeHerdr(); + SessionManager sessions = sessionManager(herdr, () -> 0L, 0, false); + WorkerSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary"); + sessions.asPresence().markPresent(session.terminalId()); + sessions.onDelivered(session.terminalId()); + + sessions.onTurnComplete(session.terminalId()); + + assertEquals(WorkerSession.State.DONE, sessions.get(session.paneId()).orElseThrow().state()); + assertTrue(promptTexts(herdr).isEmpty()); + } + @Test void drainAllReleasesBusyAndReadySessionsAndWaitsForBusy() { long[] clock = {0}; @@ -367,6 +416,13 @@ class SessionManagerTest { .count(); } + private static List promptTexts(FakeHerdr herdr) { + return herdr.calls.stream() + .filter(c -> "agent.prompt".equals(c.method())) + .map(c -> String.valueOf(((Map) c.params()).get("text"))) + .toList(); + } + // --- CB-306 spawn-readiness gate: no half-registered session on timeout ---------------- @Test diff --git a/bridged/src/test/java/dev/ltms/bridged/worker/ClaudeCodeLauncherTest.java b/bridged/src/test/java/dev/ltms/bridged/worker/ClaudeCodeLauncherTest.java index a6c6ab4..f9c9563 100644 --- a/bridged/src/test/java/dev/ltms/bridged/worker/ClaudeCodeLauncherTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/worker/ClaudeCodeLauncherTest.java @@ -78,14 +78,24 @@ class ClaudeCodeLauncherTest { Map start = (Map) herdr.lastCall("agent.start").params(); assertEquals("claude", start.get("kind"), "herdr launches the canonical executable by kind"); - assertEquals(List.of(), start.get("args"), "the configured executable is not repeated in args"); + // CB-533: the shared fixture pins model "coder", so the model flag is the whole args list. + // What this test guards is that argv[0] is NOT repeated — herdr supplies it from `kind`. + assertEquals(List.of("--model", "coder"), start.get("args"), + "the configured executable is not repeated in args"); } @Test void noBridgeFlagsWhenMcpUrlAbsent() { FakeHerdr herdr = new FakeHerdr(); service(herdr, List.of("claude", "--verbose"), null).spawn(); - assertEquals(List.of("--verbose"), spawnedArgs(herdr), "extra args untouched without mcpUrl"); + + List args = spawnedArgs(herdr); + assertFalse(args.contains("--mcp-config"), "no bridge mount without mcpUrl"); + assertFalse(args.contains("--append-system-prompt"), "no reply charter without mcpUrl"); + // CB-533: the model flag is independent of the MCP mount — pinning the model is not part of + // "mount the bridge", so an unmounted worker still runs the model its profile names. + assertEquals(List.of("--verbose", "--model", "coder"), args, + "the operator's own args are preserved, in order, ahead of the model flag"); } private ClaudeCodeLauncher multiProfile(FakeHerdr herdr) { @@ -293,6 +303,20 @@ class ClaudeCodeLauncherTest { assertTrue(caps.contains(Capability.MID_TURN_ASK), "every Claude Code peer supports mid-turn ask"); assertTrue(caps.contains(Capability.WORKTREE), "every CLI peer supports worktree cwd"); assertTrue(caps.contains(Capability.ORPHAN_REAP), "every herdr launcher supports orphan reap"); + assertTrue(caps.contains(Capability.CONTEXT_RESET), "Claude Code supports /clear"); + } + + @Test + void clearContextUsesTheClaudeCommandThroughTheOwningHandle() { + FakeHerdr herdr = new FakeHerdr(); + ClaudeCodeLauncher svc = service(herdr, List.of("claude"), null); + PeerHandle handle = svc.spawn(new SpawnRequest(null, null, null)); + + assertTrue(svc.clearContext(handle.id())); + + Map prompt = (Map) herdr.lastCall("agent.prompt").params(); + assertEquals("/clear", prompt.get("text")); + assertEquals("w9:pRoot_1", prompt.get("target")); } @Test @@ -541,4 +565,54 @@ class ClaudeCodeLauncherTest { assertEquals("http://gx00.gw:8000", startEnv(herdr).get("ANTHROPIC_BASE_URL"), "the guard-checked baseUrl must win over any env: entry, or the boundary is bypassable"); } + + // ── CB-533: the model is pinned on the command line, not only in the environment ──────────── + + /** A launcher for a profile identical but for its {@code model:} — the only variable here. */ + private ClaudeCodeLauncher serviceWithModel(FakeHerdr herdr, String model) { + BridgedConfig.Worker cfg = profileWithModel(model); + return new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr), + new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), + _ -> null); + } + + private static BridgedConfig.Worker profileWithModel(String model) { + return new BridgedConfig.Worker("sonnet", "http://gx00.gw:8000", model, null, + "BRIDGED_WORKER_TOKEN", List.of("ccs", "sonnet"), "tab", "bridged-workers", + "w #{n}", "http://127.0.0.1:8765/mcp", null, null); + } + + @Test + void aConfiguredModelIsPassedAsAModelFlagAsWellAsTheEnvVar() { + // ANTHROPIC_MODEL alone loses to `ccs`, which exports its own model family over whatever it + // inherited — so a profile that set model: was silently overruled by its own launcher. + FakeHerdr herdr = new FakeHerdr(); + serviceWithModel(herdr, "claude-sonnet-5").spawn("sonnet", null, null); + + assertEquals("claude-sonnet-5", startEnv(herdr).get("ANTHROPIC_MODEL")); + List args = spawnedArgs(herdr); + int flag = args.indexOf("--model"); + assertTrue(flag >= 0, "the flag is what survives a wrapper argv like [ccs, sonnet]"); + assertEquals("claude-sonnet-5", args.get(flag + 1)); + } + + @Test + void theModelFlagComesLastSoItOutranksTheOperatorsOwnArgv() { + FakeHerdr herdr = new FakeHerdr(); + serviceWithModel(herdr, "claude-sonnet-5").spawn("sonnet", null, null); + + List args = spawnedArgs(herdr); + assertEquals(args.size() - 2, args.indexOf("--model")); + } + + @Test + void aProfileWithNoModelGetsNoModelFlag() { + // gx10 deliberately leaves model: unset so ccs owns selection; adding a flag would make + // this file a second source of truth for exactly the thing it declines to decide. + FakeHerdr herdr = new FakeHerdr(); + serviceWithModel(herdr, null).spawn("sonnet", null, null); + + assertFalse(spawnedArgs(herdr).contains("--model")); + assertNull(startEnv(herdr).get("ANTHROPIC_MODEL")); + } } diff --git a/bridged/src/test/java/dev/ltms/bridged/worker/CompositePeerLauncherTest.java b/bridged/src/test/java/dev/ltms/bridged/worker/CompositePeerLauncherTest.java index 87ec958..518e789 100644 --- a/bridged/src/test/java/dev/ltms/bridged/worker/CompositePeerLauncherTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/worker/CompositePeerLauncherTest.java @@ -1,5 +1,8 @@ package dev.ltms.bridged.worker; +import ch.qos.logback.classic.Logger; +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.Agent; @@ -14,6 +17,7 @@ import dev.ltms.bridged.peer.SpawnRequest; import dev.ltms.bridged.placement.PlacementException; import dev.ltms.bridged.placement.PlacementPolicies; import org.junit.jupiter.api.Test; +import org.slf4j.LoggerFactory; import java.util.EnumSet; import java.util.HashMap; @@ -232,6 +236,30 @@ class CompositePeerLauncherTest { "stop routes to the spawning adapter and closes exactly that worker's pane"); } + @Test + void opencodeContextResetIsANoOpAndWarnsOnlyOnce() { + FakeHerdr herdr = new FakeHerdr(); + CompositePeerLauncher composite = composite(herdr); + PeerHandle handle = composite.spawn(new SpawnRequest("gemini", null, null)); + Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class); + ListAppender appender = new ListAppender<>(); + appender.start(); + logger.addAppender(appender); + try { + assertFalse(composite.clearContext(handle.id())); + assertFalse(composite.clearContext(handle.id())); + } finally { + logger.detachAppender(appender); + } + + assertFalse(opencodeAdapter(herdr).capabilities().contains(Capability.CONTEXT_RESET)); + assertTrue(herdr.calls.stream().noneMatch(c -> "agent.prompt".equals(c.method())), + "never type Claude's /clear into an opencode prompt"); + assertEquals(1, appender.list.stream() + .filter(e -> e.getFormattedMessage().contains("context reset is unsupported")) + .count(), "unsupported reset is logged once per adapter, not once per turn"); + } + @Test void constructorRejectsAProfileClaimedByTwoAdapters() { FakeHerdr herdr = new FakeHerdr(); diff --git a/opencode.json b/opencode.json new file mode 100644 index 0000000..de4c029 --- /dev/null +++ b/opencode.json @@ -0,0 +1,34 @@ +{ + "$schema": "https://opencode.ai/config.json", + "instructions": [ + "CLAUDE.md" + ], + "mcp": { + "bridged": { + "type": "remote", + "url": "http://127.0.0.1:8765/mcp", + "enabled": true + }, + "context7": { + "type": "remote", + "url": "https://ct7.ltms.dev/mcp", + "enabled": true, + "headers": { + "Authorization": "Bearer {file:.secrets/context7-token}" + } + }, + "gitea": { + "type": "local", + "command": [ + "gitea-mcp", + "-t", + "stdio" + ], + "enabled": true, + "environment": { + "GITEA_ACCESS_TOKEN": "{file:.secrets/gitea-token}", + "GITEA_HOST": "{file:.secrets/gitea-host}" + } + } + } +}