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}" + } + } + } +}