Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 41a3114d03 |
@@ -13,7 +13,7 @@ A design task is worked by two architects. Design alone first, then exchange and
|
||||
say plainly where you disagree. Do not concede just to agree.
|
||||
|
||||
Do only the assigned scope. Note anything outside that scope in one line and do not
|
||||
investigate it further. Use `fleet_ask{question}` only when a decision belongs to
|
||||
investigate it further. Use `bridge_ask{question}` only when a decision belongs to
|
||||
the lead, such as an unclear requirement or two defensible fixes. Do not ask about
|
||||
something you can decide by reading more code.
|
||||
|
||||
|
||||
@@ -10,7 +10,7 @@ git worktree and branch. Never check out, rebase onto, or push to `main`. Confir
|
||||
the worktree root and branch before you edit. Use only paths under that root.
|
||||
|
||||
Do only the assigned scope. Note anything outside that scope in one line and do not
|
||||
investigate it further. Use `fleet_ask{question}` only when a decision belongs to
|
||||
investigate it further. Use `bridge_ask{question}` only when a decision belongs to
|
||||
the lead, such as an unclear requirement or two defensible fixes. Do not ask about
|
||||
something you can decide by reading more code.
|
||||
|
||||
|
||||
@@ -12,7 +12,7 @@ Read the whole assigned scope before judging it. Review only that scope. If you
|
||||
something outside it, note it in one line and do not investigate it further. Do not
|
||||
run the build. The owner makes changes and runs checks.
|
||||
|
||||
Use `fleet_ask{question}` only when a decision belongs to the lead, such as an
|
||||
Use `bridge_ask{question}` only when a decision belongs to the lead, such as an
|
||||
unclear requirement or two defensible fixes. Do not ask about something you can
|
||||
decide by reading more code.
|
||||
|
||||
|
||||
@@ -5,7 +5,7 @@ description: Implementer-role procedure for a bridged worker — verify your wor
|
||||
|
||||
# Implementer worker — procedure
|
||||
|
||||
The turn contract (one `fleet_reply`, `fleet_ask` for the lead's decisions, honest reporting,
|
||||
The turn contract (one `bridge_reply`, `bridge_ask` for the lead's decisions, honest reporting,
|
||||
never merge, never commit `.mcp.json` or `wiki/`) is in **`CLAUDE.md` → Bridge communication →
|
||||
Worker** and already applies. This skill is only the *implement-and-hand-off procedure*.
|
||||
|
||||
@@ -103,7 +103,7 @@ fix it if the cause is yours (e.g. branch not pushed yet), and report the failur
|
||||
inventing a URL. If `GITEA_TOKEN` is unset your profile was not granted PR-create: push the branch
|
||||
and report its name so the lead opens the PR.
|
||||
|
||||
## 6. Hand off — what goes in `fleet_reply`
|
||||
## 6. Hand off — what goes in `bridge_reply`
|
||||
|
||||
The reply is the entire handoff; the lead cannot see your terminal.
|
||||
|
||||
@@ -130,7 +130,7 @@ sequenceDiagram
|
||||
I->>G: git push -u origin HEAD
|
||||
I->>G: POST /pulls (GITEA_TOKEN) — open PR to main
|
||||
G-->>I: html_url
|
||||
I->>L: fleet_reply(PR url, branch, files, tests)
|
||||
I->>L: bridge_reply(PR url, branch, files, tests)
|
||||
Note over L,G: the lead reviews the PR and merges on green — you never merge
|
||||
```
|
||||
|
||||
|
||||
@@ -5,7 +5,7 @@ description: Reviewer-role procedure for a bridged worker — how to work a revi
|
||||
|
||||
# Reviewer worker — procedure
|
||||
|
||||
The turn contract (one `fleet_reply`, `fleet_ask` for the lead's decisions, honest reporting,
|
||||
The turn contract (one `bridge_reply`, `bridge_ask` for the lead's decisions, honest reporting,
|
||||
never merge) is in **`CLAUDE.md` → Bridge communication → Worker** and already applies. This
|
||||
skill is only the *review procedure*: how to work the scope, and the exact shape of what you
|
||||
send back.
|
||||
@@ -24,13 +24,13 @@ wrong.
|
||||
covers what it was given.
|
||||
- Do **not** edit files or run the build. You review; the owner acts.
|
||||
|
||||
## 3. Reach for `fleet_ask` only for a genuine fork
|
||||
## 3. Reach for `bridge_ask` only for a genuine fork
|
||||
|
||||
Ambiguous requirement, a missing acceptance criterion, "intended or a bug?", or two defensible
|
||||
fixes with different consequences — those are the lead's call, and guessing produces a
|
||||
confident-but-wrong finding. Anything you could settle by reading more code is yours to settle.
|
||||
|
||||
## 4. The finding — what goes in `fleet_reply`
|
||||
## 4. The finding — what goes in `bridge_reply`
|
||||
|
||||
Report the **single most important** real issue in the scope, in these four lines, under
|
||||
~90 words:
|
||||
|
||||
@@ -13,7 +13,7 @@ A design task is worked by two architects. Design alone first, then exchange and
|
||||
say plainly where you disagree. Do not concede just to agree.
|
||||
|
||||
Do only the assigned scope. Note anything outside that scope in one line and do not
|
||||
investigate it further. Use `fleet_ask{question}` only when a decision belongs to
|
||||
investigate it further. Use `bridge_ask{question}` only when a decision belongs to
|
||||
the lead, such as an unclear requirement or two defensible fixes. Do not ask about
|
||||
something you can decide by reading more code.
|
||||
|
||||
|
||||
@@ -10,7 +10,7 @@ git worktree and branch. Never check out, rebase onto, or push to `main`. Confir
|
||||
the worktree root and branch before you edit. Use only paths under that root.
|
||||
|
||||
Do only the assigned scope. Note anything outside that scope in one line and do not
|
||||
investigate it further. Use `fleet_ask{question}` only when a decision belongs to
|
||||
investigate it further. Use `bridge_ask{question}` only when a decision belongs to
|
||||
the lead, such as an unclear requirement or two defensible fixes. Do not ask about
|
||||
something you can decide by reading more code.
|
||||
|
||||
|
||||
@@ -12,7 +12,7 @@ Read the whole assigned scope before judging it. Review only that scope. If you
|
||||
something outside it, note it in one line and do not investigate it further. Do not
|
||||
run the build. The owner makes changes and runs checks.
|
||||
|
||||
Use `fleet_ask{question}` only when a decision belongs to the lead, such as an
|
||||
Use `bridge_ask{question}` only when a decision belongs to the lead, such as an
|
||||
unclear requirement or two defensible fixes. Do not ask about something you can
|
||||
decide by reading more code.
|
||||
|
||||
|
||||
@@ -34,8 +34,8 @@ flowchart LR
|
||||
W["worker claude pane<br/>ANTHROPIC_BASE_URL set<br/>MCP client"]
|
||||
M["llm.ltms.dev<br/>(the one gateway)"]
|
||||
|
||||
OPUS -->|"MCP fleet_send (blocks)"| SRV
|
||||
W -.->|"MCP fleet_reply"| SRV
|
||||
OPUS -->|"MCP bridge_send (blocks)"| SRV
|
||||
W -.->|"MCP bridge_reply"| SRV
|
||||
CLI -->|"Unix socket<br/>send_text · events.subscribe"| HERDR
|
||||
HERDR -->|"drives PTY"| W
|
||||
W -->|"inference"| M
|
||||
@@ -51,16 +51,13 @@ flowchart LR
|
||||
plain daemon (no Anthropic quota), so it may poll/subscribe freely.
|
||||
- **One gateway (unified MCP setup):** `bridged` is the **sole communication path** for every
|
||||
Claude session. Primary and workers each mount it as an MCP server (one `claude mcp add`
|
||||
line, same on both) and talk over MCP tools — `fleet_send` / `fleet_reply` /
|
||||
`fleet_status` (with `fleet_ask` planned for the blocked-worker path). **No Claude session
|
||||
ever addresses a broker, a peer, or the network
|
||||
directly**; any queue is `bridged`-internal. MCP tool I/O never sets `ANTHROPIC_BASE_URL`, so
|
||||
mounting the bridge is subscription-safe by construction.
|
||||
**Tool naming:** the tools were renamed from `bridge_*` to `fleet_*` (CB-622). The daemon
|
||||
still answers the old `bridge_*` names for one release, but they are deprecated — use the
|
||||
`fleet_*` names.
|
||||
- **How the primary consumes a reply:** a single **blocking MCP call** (`fleet_send`);
|
||||
`bridged` holds it open until the worker calls `fleet_reply` or its turn hits
|
||||
line, same on both) and talk over MCP tools — `bridge_send` / `bridge_reply` /
|
||||
`bridge_status` (with `bridge_ask` planned for the blocked-worker path). **No Claude session
|
||||
ever addresses a broker, a peer, or the network
|
||||
directly**; any queue is `bridged`-internal. MCP tool I/O never sets `ANTHROPIC_BASE_URL`, so
|
||||
mounting the bridge is subscription-safe by construction.
|
||||
- **How the primary consumes a reply:** a single **blocking MCP call** (`bridge_send`);
|
||||
`bridged` holds it open until the worker calls `bridge_reply` or its turn hits
|
||||
`agent_status=done`, then returns the reply as the tool result. No cross-turn busy-poll, so
|
||||
no quota burn. SSE is an optional side-channel for humans/dashboards watching status.
|
||||
- **Worker → primary** rides `bridged`'s **MCP rendezvous** — the reply resolves the primary's
|
||||
@@ -101,9 +98,9 @@ tests run separately via `mvn test -Pcontract`):
|
||||
|
||||
- **Core gateway** — herdr socket client (contract-tested vs live 0.7.0); guard-checked worker
|
||||
spawn with `ANTHROPIC_BASE_URL` injected only into the worker's env; status-gated injector;
|
||||
blocking `fleet_send` with reply rendezvous; MCP server as a thin adapter over the REST core.
|
||||
- **MCP tools** — `fleet_send` / `fleet_reply` / `fleet_status` (messaging) and `fleet_spawn`
|
||||
/ `fleet_list` / `fleet_stop` / `fleet_profiles` / `fleet_poll` (fleet). Caller identity is
|
||||
blocking `bridge_send` with reply rendezvous; MCP server as a thin adapter over the REST core.
|
||||
- **MCP tools** — `bridge_send` / `bridge_reply` / `bridge_status` (messaging) and `bridge_spawn`
|
||||
/ `bridge_list` / `bridge_stop` / `bridge_profiles` / `bridge_poll` (fleet). Caller identity is
|
||||
connection-based (loopback peer PID → herdr pane), so the same mount serves primary and workers.
|
||||
- **Delivery reliability** — completion fallback (a confirmed `working→idle` turn resolves a
|
||||
send); async fire-and-poll (beats the caller's MCP call timeout for long tasks); and failure
|
||||
@@ -111,7 +108,7 @@ tests run separately via `mvn test -Pcontract`):
|
||||
- **Fleet** — multiple worker profiles, each with an independent base_url guard check; workers
|
||||
inherit the primary's working directory (never `$HOME`); a readiness gate holds delivery until
|
||||
a worker's Claude has connected the bridge MCP (no paste lost into its boot window).
|
||||
- **Blocked-worker path** — `fleet_ask` reverse rendezvous: a worker pauses its delegated turn to
|
||||
- **Blocked-worker path** — `bridge_ask` reverse rendezvous: a worker pauses its delegated turn to
|
||||
ask the primary and resumes the *same* turn with the answer (CB-205).
|
||||
- **Session lifecycle** — session manager with spawn/reuse/recycle, `idle_ttl` reaper, `context_cap`,
|
||||
and graceful drain on shutdown (CB-301/CB-303); per-worker git worktrees on their own branch with
|
||||
|
||||
@@ -5,10 +5,10 @@
|
||||
|
||||
## Why this exists
|
||||
|
||||
Stage 2 gave a worker→primary reply a **durable place to wait** when no `fleet_send` is
|
||||
Stage 2 gave a worker→primary reply a **durable place to wait** when no `bridge_send` is
|
||||
open: it lands in `agent.<target>.inbox` on the broker and survives a daemon bounce. But
|
||||
delivery is still **pull** — the primary only sees the reply if it happens to call
|
||||
`fleet_poll(target)` / `GET /sessions/{id}/replies`. A reply can sit indefinitely while
|
||||
`bridge_poll(target)` / `GET /sessions/{id}/replies`. A reply can sit indefinitely while
|
||||
the primary works on something else.
|
||||
|
||||
This layer makes delivery **active**: the bridge *pushes* a nudge to the primary the moment
|
||||
@@ -26,11 +26,11 @@ pointed at the primary's pane instead.
|
||||
|
||||
```mermaid
|
||||
flowchart LR
|
||||
W["worker"] -->|"fleet_reply (no open send)"| MS["MessageService.reply"]
|
||||
W["worker"] -->|"bridge_reply (no open send)"| MS["MessageService.reply"]
|
||||
MS -->|"inbox.publish"| INBOX[("agent.<target>.inbox<br/>(durable, LavinMQ)")]
|
||||
MS -->|"notify"| LOOP["ReplyPushLoop"]
|
||||
LOOP -->|"status-gated inject"| PANE["primary's herdr pane"]
|
||||
PANE -->|"primary drains"| DRAIN["fleet_poll(target)<br/>= peek + ack"]
|
||||
PANE -->|"primary drains"| DRAIN["bridge_poll(target)<br/>= peek + ack"]
|
||||
DRAIN -->|"inbox now empty"| LOOP
|
||||
LOOP -.->|"still non-empty →<br/>re-inject on backoff"| PANE
|
||||
classDef store fill:#2c5282,stroke:#1a365d,color:#ffffff;
|
||||
@@ -51,11 +51,11 @@ the caller runs in a herdr pane on this host. Today it's discarded for the prima
|
||||
(`presence.markPresent` is a no-op on it).
|
||||
|
||||
**Plan:** a single-slot `PrimaryRegistry` (thread-safe) holding the primary's `terminal_id`.
|
||||
Populate it from the **orchestration-side** MCP tools — `fleet_send`, `fleet_spawn`,
|
||||
`fleet_poll`, `fleet_list`, `fleet_status`, `fleet_profiles` — capturing
|
||||
Populate it from the **orchestration-side** MCP tools — `bridge_send`, `bridge_spawn`,
|
||||
`bridge_poll`, `bridge_list`, `bridge_status`, `bridge_profiles` — capturing
|
||||
`callerTerminal(exchange)` when it is (a) non-null and (b) **not** a registered worker
|
||||
session in `SessionManager`. That caller is, by construction, the primary. Worker-side tools
|
||||
(`fleet_reply`, `fleet_ask`) never set it.
|
||||
(`bridge_reply`, `bridge_ask`) never set it.
|
||||
|
||||
- **Config override / pin:** a `primary: { terminal: "<id>" }` block in `BridgedConfig`
|
||||
(nested record, same shape as `Broker`). Lets an operator pin it, or supply it when
|
||||
@@ -71,13 +71,13 @@ A `ReplyPushLoop` component, notified at the single no-waiter call site
|
||||
(`MessageService.reply` → the `inbox.publish` branch, `MessageService.java:192`).
|
||||
|
||||
- **Inject a nudge, not the payload.** The injected turn tells the primary *to drain*
|
||||
(e.g. "Worker `<target>` returned a reply — run `fleet_poll(target=<target>)` to collect
|
||||
(e.g. "Worker `<target>` returned a reply — run `bridge_poll(target=<target>)` to collect
|
||||
it"), it does **not** carry the reply text. Rationale: replies can be large/multiline and
|
||||
terminal injection would mangle them; the drain response is the clean transport. Keeps the
|
||||
push idempotent — re-nudging is harmless.
|
||||
- **Ack = drain.** The primary draining (`drainReplies` = peek + ack) is the acknowledgement.
|
||||
The loop's **stop condition is `inbox.peek(target).isEmpty()`** — the reply is gone from the
|
||||
inbox because it was acked. No new `fleet_ack` tool needed for v1 (see Increment 3).
|
||||
inbox because it was acked. No new `bridge_ack` tool needed for v1 (see Increment 3).
|
||||
- **Status-gated injection (mechanism (b), chosen).** A dedicated lightweight scheduled loop,
|
||||
**not** the worker `Injector`. It injects via `AgentControl.send(primaryTerminal, nudge)`
|
||||
(the same herdr `agent.send` = `pane send-text` + submit that delivers to workers) only when
|
||||
@@ -93,10 +93,10 @@ A `ReplyPushLoop` component, notified at the single no-waiter call site
|
||||
the reply remains in the durable inbox and the next natural poll (or a later worker reply's
|
||||
nudge) still surfaces it. Bounded so the bridge never spams the primary.
|
||||
|
||||
### Increment 3 — optional per-`msgId` `fleet_ack` tool (deferred)
|
||||
### Increment 3 — optional per-`msgId` `bridge_ack` tool (deferred)
|
||||
|
||||
Drain-as-ack is coarse: it clears *all* pending replies for a target at once. If finer
|
||||
control is ever needed (ack one reply, leave others held), add a `fleet_ack(msgId)` tool
|
||||
control is ever needed (ack one reply, leave others held), add a `bridge_ack(msgId)` tool
|
||||
mapping to `inbox.ack(target, msgId)` — the port already supports per-`msgId` ack. Not built
|
||||
in v1; the stop-on-empty loop is sufficient.
|
||||
|
||||
@@ -132,7 +132,7 @@ boundary**. The bridge is signalling the primary that it has mail — not drivin
|
||||
non-null terminal AND not a registered session" predicate; the loop's stop-on-empty and
|
||||
bounded-reminder logic with an injected clock + a fake injector (no real herdr).
|
||||
- **Live dogfood (primary-side):** with the daemon on the broker jar + a real worker,
|
||||
delegate a task, let the worker reply after the `fleet_send` window closes, and observe the
|
||||
delegate a task, let the worker reply after the `bridge_send` window closes, and observe the
|
||||
bridge inject a drain nudge into *this* primary pane; confirm draining stops the reminders;
|
||||
confirm an unreachable primary (registry empty) degrades to pull with no loss.
|
||||
|
||||
|
||||
@@ -381,7 +381,7 @@ public final class Bridged {
|
||||
+ "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.
|
||||
// open fleet_send. Uses its own lightweight scheduled executor, separate from the injector.
|
||||
int maxReminders = cfg.primary() != null ? cfg.primary().remindersOrDefault() : 5;
|
||||
long backoffMs = cfg.primary() != null ? cfg.primary().backoffMsOrDefault() : 15_000L;
|
||||
var pushScheduler = Executors.newSingleThreadScheduledExecutor(r ->
|
||||
@@ -458,7 +458,7 @@ public final class Bridged {
|
||||
primaryRegistry.forgetDelegation(detail.terminalId()); // CB-532: don't leak the lead binding
|
||||
});
|
||||
|
||||
// MCP server face (CB-105): bridge_send/bridge_reply/bridge_status, mounted at /mcp.
|
||||
// MCP server face (CB-105): fleet_send/fleet_reply/fleet_status, mounted at /mcp.
|
||||
// Caller identity is resolved from the connection (peer PID → herdr pane), not arguments.
|
||||
ConnectionIdentity identity = new ConnectionIdentity(
|
||||
new PaneLocator(herdr), new LsofPeerPidLookup(), new LsofProcessCwdLookup());
|
||||
|
||||
@@ -94,7 +94,7 @@ public final class MemberRegistry implements MemberLifecycle {
|
||||
* An immutable copy of the live {@code terminal_id → slot name} bindings.
|
||||
*
|
||||
* <p>Passed to {@link CallerResolver} as the source of architect identity, and what
|
||||
* {@code bridge_whoami}/the roster will read to say which slot a pane hosts. Empty until the
|
||||
* {@code fleet_whoami}/the roster will read to say which slot a pane hosts. Empty until the
|
||||
* spawn lifecycle binds a slot.
|
||||
*/
|
||||
public Map<String, String> snapshot() {
|
||||
|
||||
@@ -39,14 +39,14 @@ public record Principal(Role role, String terminal, long pid, String name) {
|
||||
*
|
||||
* <p>Carries {@link Role#PRIMARY}: a lead <em>is</em> 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 <em>which</em>
|
||||
* no new entry. The name is reporting only — it lets {@code fleet_whoami} say <em>which</em>
|
||||
* lead is asking once more than one is configured.
|
||||
*
|
||||
* <p><strong>CB-532: a lead now carries the terminal it was matched by.</strong> 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.
|
||||
* {@code fleet_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 #isSpawnedMember()} instead, and a lead is a peer that can both send and receive.
|
||||
*/
|
||||
@@ -63,7 +63,7 @@ public record Principal(Role role, String terminal, long pid, String name) {
|
||||
* An architect (CB-548), identified by the slot it occupies and the pane bound to it.
|
||||
*
|
||||
* <p>Carries {@link Role#ARCHITECT}. {@code slotName} is reporting only — it lets
|
||||
* {@code bridge_whoami} say <em>which</em> architect slot is asking, and it is the key the
|
||||
* {@code fleet_whoami} say <em>which</em> architect slot is asking, and it is the key the
|
||||
* (future) spawn lifecycle reads a profile back from. Identity is the {@code terminal}: like a
|
||||
* worker's it comes from the connection and the live terminal→slot binding, so
|
||||
* {@code ownsSession} works exactly as it does for a worker — an architect acts as its own
|
||||
|
||||
@@ -184,7 +184,7 @@ public record BridgedConfig(
|
||||
* and {@code {n}} (per-worker number, to keep sibling tabs distinct)
|
||||
* are substituted (default {@code "worker: {profile} #{n}"})
|
||||
* @param mcpUrl bridge MCP URL to provision into the worker's {@code configDir} so it
|
||||
* can call {@code bridge_reply} ({@code null}/blank → no provisioning; the
|
||||
* can call {@code fleet_reply} ({@code null}/blank → no provisioning; the
|
||||
* worker won't reply, only the fallback/timeout resolves the send)
|
||||
* @param cwd fixed working directory for this profile's workers (CB-112 "told otherwise");
|
||||
* {@code null}/blank → inherit the primary's cwd, else the daemon's
|
||||
@@ -224,7 +224,7 @@ public record BridgedConfig(
|
||||
* auto-select this profile": it is excluded from every automatic policy's
|
||||
* pool the same way a quarantined candidate is (see
|
||||
* {@code PlacementPolicyUtil}). This does not make the profile
|
||||
* unreachable — an explicit {@code bridge_spawn{profile:"..."}} bypasses
|
||||
* unreachable — an explicit {@code fleet_spawn{profile:"..."}} bypasses
|
||||
* placement entirely and still resolves it. Weights among the remaining
|
||||
* (non-excluded) candidates need not sum to 1.0; only their ratios matter.
|
||||
* @param maxLoad max live workers allowed on this profile at one time. Absent
|
||||
@@ -234,7 +234,7 @@ public record BridgedConfig(
|
||||
* the same way a {@code weight <= 0} profile is (see
|
||||
* {@code PlacementPolicyUtil.available()}, which already treats "at cap"
|
||||
* and "excluded" alike), and an explicit
|
||||
* {@code bridge_spawn{profile:"..."}} against it is refused too (see
|
||||
* {@code fleet_spawn{profile:"..."}} against it is refused too (see
|
||||
* {@code CompositePeerLauncher.enforceMaxLoad}) — a cap is a capacity
|
||||
* statement that does not stop being true just because the profile was
|
||||
* named directly. A negative value has no sane meaning (there is no
|
||||
@@ -259,7 +259,7 @@ public record BridgedConfig(
|
||||
* {@code ANTHROPIC_BASE_URL} or {@code ANTHROPIC_AUTH_TOKEN} is refused at
|
||||
* config load (CB-542): on the subscription path no guard would vet it.
|
||||
* @param exhaustedPattern regex matched against a completion-fallback scrape (CB-578 stage A) to
|
||||
* classify a turn that ended with no {@code bridge_reply} as the backend
|
||||
* classify a turn that ended with no {@code fleet_reply} as the backend
|
||||
* having refused on a subscription usage limit, rather than a real answer.
|
||||
* {@code null}/blank ⇒ the classification never fires for this profile and
|
||||
* today's completion-fallback behaviour is unchanged. Every backend words
|
||||
@@ -605,7 +605,7 @@ public record BridgedConfig(
|
||||
* work as peers — the second is silently demoted and refused every orchestration call.
|
||||
*
|
||||
* <p>{@code kind} and {@code model} are descriptive only: they document what runs in the pane
|
||||
* and are reported back by {@code bridge_whoami}.
|
||||
* and are reported back by {@code fleet_whoami}.
|
||||
*
|
||||
* <p><b>A lead is now also creatable (CB-557).</b> Before, nothing spawned one — a lead
|
||||
* pre-existed, which is why it had to be recognised by configuration rather than created. With
|
||||
@@ -811,7 +811,7 @@ public record BridgedConfig(
|
||||
|
||||
/**
|
||||
* Opt-in idle-lead heartbeat (CB-551): when the single lead has been continuously idle past
|
||||
* {@code idleAfterSeconds} with no open {@code bridge_send} driving it, nudge it back to work.
|
||||
* {@code idleAfterSeconds} with no open {@code fleet_send} driving it, nudge it back to work.
|
||||
*
|
||||
* <p>Deliberately opt-in ({@code null} ⇒ off, exactly like {@code leadScan:}). The heartbeat
|
||||
* spends the operator's model subscription on its own initiative — it prompts the lead to start
|
||||
|
||||
@@ -6,7 +6,7 @@ import java.util.Map;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
|
||||
/**
|
||||
* Counts turns that ended via the completion fallback instead of {@code bridge_reply}.
|
||||
* Counts turns that ended via the completion fallback instead of {@code fleet_reply}.
|
||||
* MUTE is an observation by target and profile, not a classifier state and never suppresses faults.
|
||||
*/
|
||||
public final class MuteCounter {
|
||||
|
||||
@@ -16,14 +16,14 @@ import java.util.regex.Pattern;
|
||||
|
||||
/**
|
||||
* The CB-106 completion fallback: bridges the {@link Injector}'s turn-completion signal to the
|
||||
* {@link Rendezvous} so a blocking {@code bridge_send} resolves even when the worker finishes its
|
||||
* task without ever calling {@code bridge_reply} — the common case for a real delegated coding task.
|
||||
* {@link Rendezvous} so a blocking {@code fleet_send} resolves even when the worker finishes its
|
||||
* task without ever calling {@code fleet_reply} — the common case for a real delegated coding task.
|
||||
*
|
||||
* <p>On a confirmed {@code working → idle} boundary it scrapes the worker's recent transcript and
|
||||
* resolves the awaiting send with that tail (a {@link Rendezvous.Kind#COMPLETION} resolution, so the
|
||||
* caller can tell a scrape from a structured reply). It scrapes only when a send is actually waiting
|
||||
* — a fleet worker's own turns, or a send that already timed out, cost no herdr traffic. An explicit
|
||||
* {@code bridge_reply} that raced in first wins; {@link Rendezvous#resolveCompletion} is then a no-op.
|
||||
* {@code fleet_reply} that raced in first wins; {@link Rendezvous#resolveCompletion} is then a no-op.
|
||||
*
|
||||
* <p>It also handles the CB-109 stall signal ({@link #onTurnFailed}): a worker that ran a turn then
|
||||
* wedged in an {@code unknown} state resolves the send as a failure (with the error screen as
|
||||
@@ -38,7 +38,7 @@ import java.util.regex.Pattern;
|
||||
* <p><strong>Waiter-specific resolution (CB-116).</strong> On delivery we also capture the exact
|
||||
* {@link Rendezvous} waiter this turn belongs to, and the completion/failure fallbacks resolve
|
||||
* <em>that</em> waiter — never "whatever send is waiting now". A completion fallback runs on a virtual
|
||||
* thread and can land after the worker's {@code bridge_reply} already resolved the turn and the
|
||||
* thread and can land after the worker's {@code fleet_reply} already resolved the turn and the
|
||||
* <em>next</em> send opened its own waiter on the same session; resolving the current waiter would
|
||||
* then deliver turn N's stale scrape as turn N+1's answer. Targeting the captured waiter makes a late
|
||||
* completion a harmless no-op (its waiter is already done) instead of a cross-turn stale reply.
|
||||
@@ -62,7 +62,7 @@ public final class CompletionResolver implements TurnListener {
|
||||
static final int MAX_SCRAPE_CHARS = 4000;
|
||||
|
||||
private static final String CLIPPED_PANE_TAIL_MARKER =
|
||||
"[Pane tail clipped: member did not call bridge_reply.]";
|
||||
"[Pane tail clipped: member did not call fleet_reply.]";
|
||||
|
||||
private final AgentControl agents;
|
||||
private final Rendezvous rendezvous;
|
||||
@@ -171,7 +171,7 @@ public final class CompletionResolver implements TurnListener {
|
||||
void resolve(String target, InFlight turn) {
|
||||
CompletableFuture<Rendezvous.Resolution> waiter = turn == null ? null : turn.waiter();
|
||||
if (waiter == null || waiter.isDone()) {
|
||||
// Nobody is blocked on THIS turn (it had no send, or its bridge_reply already won). Skip
|
||||
// Nobody is blocked on THIS turn (it had no send, or its fleet_reply already won). Skip
|
||||
// the scrape; resolving the current waiter here would be the CB-116 cross-turn stale reply.
|
||||
inFlight.remove(target, turn);
|
||||
return;
|
||||
@@ -197,7 +197,7 @@ public final class CompletionResolver implements TurnListener {
|
||||
// Misattribution guard (CB-115): if the scrape is byte-identical to the pane content at
|
||||
// delivery, this turn produced no new output — the boundary belongs to the previous turn's
|
||||
// wind-down (common on rapid back-to-back sends). Suppress rather than resolve the send with
|
||||
// a stale answer; the real bridge_reply (or a later genuine completion) resolves it instead.
|
||||
// a stale answer; the real fleet_reply (or a later genuine completion) resolves it instead.
|
||||
// A scrape that failed to read is exempt — an empty tail there is "couldn't see", not "no change".
|
||||
String baseline = turn.baseline();
|
||||
if (!scrapeFailed && baseline != null && baseline.equals(tail)) {
|
||||
@@ -205,7 +205,7 @@ public final class CompletionResolver implements TurnListener {
|
||||
target);
|
||||
return; // keep the in-flight record: a later genuine completion still needs it
|
||||
}
|
||||
// CB-578 stage A: a turn that ended with no bridge_reply AND whose scrape matches the
|
||||
// CB-578 stage A: a turn that ended with no fleet_reply AND whose scrape matches the
|
||||
// backend's configured usage-limit pattern is a refusal, not an answer. Classify it as
|
||||
// BACKEND_EXHAUSTED rather than handing the caller a scrape that reads like a real reply.
|
||||
if (!scrapeFailed) {
|
||||
@@ -215,7 +215,7 @@ public final class CompletionResolver implements TurnListener {
|
||||
String reason = "backend exhausted (usage limit): " + matchedLine;
|
||||
if (rendezvous.resolveExhausted(waiter, reason)) {
|
||||
inFlight.remove(target, turn);
|
||||
log.warn("completion for {} classified BACKEND_EXHAUSTED (no bridge_reply; scrape "
|
||||
log.warn("completion for {} classified BACKEND_EXHAUSTED (no fleet_reply; scrape "
|
||||
+ "matched the profile's exhausted pattern): {}", target, reason);
|
||||
// CB-578 stage B: only on the resolution that actually won the race — a late
|
||||
// duplicate must never quarantine a credential twice for one refusal.
|
||||
@@ -229,7 +229,7 @@ public final class CompletionResolver implements TurnListener {
|
||||
inFlight.remove(target, turn);
|
||||
if (clipped) {
|
||||
log.warn("completion scrape for {} clipped from {} chars to the {} char cap; "
|
||||
+ "member did not call bridge_reply, so the pane tail is partial",
|
||||
+ "member did not call fleet_reply, so the pane tail is partial",
|
||||
target, originalLength, MAX_SCRAPE_CHARS);
|
||||
}
|
||||
log.debug("resolved send to {} via turn-completion fallback ({} chars scraped)",
|
||||
|
||||
@@ -6,7 +6,7 @@ import dev.ltms.bridged.msg.TurnToken;
|
||||
* Notified when a worker's delegated turn is observed to complete — a confirmed
|
||||
* {@code WORKING → IDLE} transition after a delivery. This is the CB-106 completion signal the
|
||||
* {@code CompletionResolver} uses to resolve a blocked send whose worker never called
|
||||
* {@code bridge_reply}. Kept as a seam so the {@link Injector} needs no dependency on the message
|
||||
* {@code fleet_reply}. Kept as a seam so the {@link Injector} needs no dependency on the message
|
||||
* layer and stays unit-testable with a capturing fake.
|
||||
*/
|
||||
@FunctionalInterface
|
||||
|
||||
@@ -26,7 +26,7 @@ import java.util.stream.Collectors;
|
||||
* must never receive:
|
||||
* <ol>
|
||||
* <li>it appends the <em>reply charter</em> — "you are an off-subscription worker … end every turn
|
||||
* with {@code bridge_reply}". A lead is the orchestrator; telling it that it is a worker is
|
||||
* with {@code fleet_reply}". A lead is the orchestrator; telling it that it is a worker is
|
||||
* exactly backwards.</li>
|
||||
* <li>it registers the session with {@code SessionManager}, which subjects it to the idle reaper,
|
||||
* the context cap and the shutdown drain. An idle lead is the normal state of a lead, so the
|
||||
|
||||
@@ -37,19 +37,24 @@ import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.function.BiFunction;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.LongSupplier;
|
||||
import java.util.function.Supplier;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
/**
|
||||
* The MCP SERVER face (CB-105): a Streamable-HTTP MCP server whose tools are <em>thin adapters</em>
|
||||
* over the same {@link MessageService}/{@link Rendezvous} the REST routes use — so the two are
|
||||
* validated by parity, not by re-implementing behaviour. The primary Opus calls {@code bridge_send}
|
||||
* / {@code bridge_status}; the worker calls {@code bridge_reply}.
|
||||
* validated by parity, not by re-implementing behaviour. The primary Opus calls {@code fleet_send}
|
||||
* / {@code fleet_status}; the worker calls {@code fleet_reply}.
|
||||
*
|
||||
* <p>Beyond delegation the primary also manages the fleet here (CB-108): {@code bridge_spawn} /
|
||||
* {@code bridge_list} / {@code bridge_stop} drive the {@link PeerLauncher} SPI so a worker's whole
|
||||
* <p>Beyond delegation the primary also manages the fleet here (CB-108): {@code fleet_spawn} /
|
||||
* {@code fleet_list} / {@code fleet_stop} drive the {@link PeerLauncher} SPI so a worker's whole
|
||||
* lifecycle is managed through MCP, with each adapter's subscription boundary enforced inside it.
|
||||
*
|
||||
* <p>The tool <em>logic</em> lives in package-private static methods returning a
|
||||
@@ -58,14 +63,24 @@ import java.util.stream.Collectors;
|
||||
*/
|
||||
public final class BridgeMcp {
|
||||
|
||||
private static final Logger log = LoggerFactory.getLogger(BridgeMcp.class);
|
||||
|
||||
private static final long DEFAULT_TIMEOUT_MS = 25_000;
|
||||
private static final long MAX_TIMEOUT_MS = 120_000;
|
||||
// bridge_ask blocks the WORKER's own MCP call, which its client caps near 60s — default under
|
||||
// fleet_ask blocks the WORKER's own MCP call, which its client caps near 60s — default under
|
||||
// that so the bridge returns a clean timeout before the client severs the call (CB-205).
|
||||
private static final long ASK_DEFAULT_TIMEOUT_MS = 55_000;
|
||||
private static final long ASK_MAX_TIMEOUT_MS = 115_000;
|
||||
private static final ObjectMapper MAPPER = new ObjectMapper(); // worker-view JSON projections
|
||||
|
||||
/**
|
||||
* CB-622: the product is renaming {@code bridge_*} tools to {@code fleet_*}. Both names reach
|
||||
* the same handler (registered below); this set makes the "old name used" WARN fire once per
|
||||
* old name for the life of the process, not once per call — a per-name flag, not a numeric
|
||||
* sentinel, so it survives concurrent callers cleanly and reads unambiguously in a log.
|
||||
*/
|
||||
private static final Set<String> WARNED_DEPRECATED_NAMES = ConcurrentHashMap.newKeySet();
|
||||
|
||||
/** Transport-context key under which the extractor stashes the resolved caller identity. */
|
||||
static final String CALLER_TERMINAL = "callerTerminal";
|
||||
/** Transport-context key under which the extractor stashes the caller's PID (for cwd inherit). */
|
||||
@@ -83,7 +98,7 @@ public final class BridgeMcp {
|
||||
private final HealthCoverageSource healthCoverage;
|
||||
private final QuarantineSource quarantine;
|
||||
|
||||
/** Capacity facts used by {@code bridge_list}; production must supply the placement live count. */
|
||||
/** Capacity facts used by {@code fleet_list}; production must supply the placement live count. */
|
||||
public record CapacitySource(Function<String, Integer> liveCount, Function<String, Integer> maxLoad,
|
||||
Supplier<Set<String>> configuredProfiles, LongSupplier clock) {
|
||||
/** Inert test-only source. It omits capacity rather than inventing zero live counts. */
|
||||
@@ -95,7 +110,7 @@ public final class BridgeMcp {
|
||||
public record HealthCoverageSource(Supplier<String> value) { }
|
||||
|
||||
/**
|
||||
* CB-578 stage B quarantine facts used by {@code bridge_profiles}: a profile → credential id
|
||||
* CB-578 stage B quarantine facts used by {@code fleet_profiles}: a profile → credential id
|
||||
* lookup, plus the shared {@link BackendQuarantine} to read remaining cooldowns off.
|
||||
*/
|
||||
public record QuarantineSource(Function<String, String> credentialIdFor, BackendQuarantine quarantine) {
|
||||
@@ -109,7 +124,7 @@ public final class BridgeMcp {
|
||||
* Jetty's context handler and never passes through Javalin's {@code before}
|
||||
* filter, so the REST guard does not cover it.
|
||||
* @param metrics registry for auth-failure counting; may be {@code null}
|
||||
* @param quarantine CB-578 stage B facts for {@code bridge_profiles}; required — pass
|
||||
* @param quarantine CB-578 stage B facts for {@code fleet_profiles}; required — pass
|
||||
* {@link QuarantineSource#none()} for a caller that does not want the feature
|
||||
*/
|
||||
public BridgeMcp(MessageService messages, PeerLauncher workers, SessionManager sessions,
|
||||
@@ -124,7 +139,7 @@ public final class BridgeMcp {
|
||||
.jsonMapper(json)
|
||||
.mcpEndpoint("/mcp")
|
||||
// Resolve the caller from the connection (peer PID → herdr pane) in one lookup: the
|
||||
// worker terminal for bridge_reply (no spoofable arg), and the PID so bridge_spawn can
|
||||
// worker terminal for fleet_reply (no spoofable arg), and the PID so fleet_spawn can
|
||||
// inherit the primary's cwd (CB-112). Any contact from a worker marks it available
|
||||
// (CB-113) — its MCP initialize is the reliable "the agent is up" signal.
|
||||
.contextExtractor(req -> {
|
||||
@@ -145,10 +160,10 @@ public final class BridgeMcp {
|
||||
CALLER_NAME, orEmpty(p.name())));
|
||||
})
|
||||
.build();
|
||||
this.server = McpServer.sync(transport)
|
||||
.serverInfo("bridge", "0.1.0")
|
||||
.capabilities(McpSchema.ServerCapabilities.builder().tools(true).build())
|
||||
.toolCall(sendTool(), (exchange, req) -> {
|
||||
// CB-622: each handler is built once and reused for BOTH its fleet_* tool and its
|
||||
// deprecated bridge_* twin (registered below), so the two names can never drift apart.
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> sendHandler =
|
||||
(exchange, req) -> {
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.SEND,
|
||||
str(req.arguments(), "sessionId"));
|
||||
if (denied != null) return denied;
|
||||
@@ -163,7 +178,7 @@ public final class BridgeMcp {
|
||||
String content = str(a, "content");
|
||||
String turnId = str(a, "turnId");
|
||||
if (turnId != null && !turnId.isBlank()) {
|
||||
// Answering a worker's bridge_ask (CB-205): resolve its blocked question and
|
||||
// Answering a worker's fleet_ask (CB-205): resolve its blocked question and
|
||||
// block for the worker's reply as it resumes the same turn. This is the same
|
||||
// delegation, so ownership is left untouched (CB-548) — never re-recorded.
|
||||
return answer(messages, turnId, content, timeoutMs(a));
|
||||
@@ -178,43 +193,49 @@ public final class BridgeMcp {
|
||||
return Boolean.FALSE.equals(a.get("wait"))
|
||||
? sendAsync(messages, target, content, onAccepted, workers.profiles())
|
||||
: send(messages, target, content, timeoutMs(a), onAccepted, workers.profiles());
|
||||
})
|
||||
// bridge_reply's identity is the CONNECTION, never an argument — so the authz check
|
||||
// is "is this caller a worker at all", and it can only ever reply as itself.
|
||||
.toolCall(replyTool(), (exchange, req) -> {
|
||||
};
|
||||
// fleet_reply's identity is the CONNECTION, never an argument — so the authz check
|
||||
// is "is this caller a worker at all", and it can only ever reply as itself.
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> replyHandler =
|
||||
(exchange, req) -> {
|
||||
String self = callerTerminal(exchange);
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.REPLY, self);
|
||||
if (denied != null) return denied;
|
||||
return reply(messages, self, str(req.arguments(), "content"));
|
||||
})
|
||||
// bridge_ask (CB-205): a worker's mid-turn question — identity from the CONNECTION.
|
||||
.toolCall(askTool(), (exchange, req) -> {
|
||||
};
|
||||
// fleet_ask (CB-205): a worker's mid-turn question — identity from the CONNECTION.
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> askHandler =
|
||||
(exchange, req) -> {
|
||||
String self = callerTerminal(exchange);
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.ASK, self);
|
||||
if (denied != null) return denied;
|
||||
return ask(messages, self, str(req.arguments(), "question"), timeoutMs(req.arguments()));
|
||||
})
|
||||
.toolCall(statusTool(), (exchange, req) -> {
|
||||
};
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> statusHandler =
|
||||
(exchange, req) -> {
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
|
||||
if (denied != null) return denied;
|
||||
return status(messages, str(req.arguments(), "sessionId"));
|
||||
})
|
||||
.toolCall(pollTool(), (exchange, req) -> {
|
||||
};
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> pollHandler =
|
||||
(exchange, req) -> {
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
|
||||
if (denied != null) return denied;
|
||||
Map<String, Object> a = req.arguments();
|
||||
return poll(messages, str(a, "ticket"), str(a, "target"));
|
||||
})
|
||||
// CB-307 Increment 3: per-msgId ack (not needed in v1 but supported by the inbox).
|
||||
// Acking removes a reply from the inbox, so it is a drain, not a read.
|
||||
.toolCall(ackTool(), (exchange, req) -> {
|
||||
};
|
||||
// CB-307 Increment 3: per-msgId ack (not needed in v1 but supported by the inbox).
|
||||
// Acking removes a reply from the inbox, so it is a drain, not a read.
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> ackHandler =
|
||||
(exchange, req) -> {
|
||||
Map<String, Object> a = req.arguments();
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.DRAIN, str(a, "target"));
|
||||
if (denied != null) return denied;
|
||||
return ack(messages, str(a, "target"), str(a, "msgId"));
|
||||
})
|
||||
// Fleet management (CB-108): spawn/list/stop over ClaudeCodeLauncher.
|
||||
.toolCall(spawnTool(), (exchange, req) -> {
|
||||
};
|
||||
// Fleet management (CB-108): spawn/list/stop over ClaudeCodeLauncher.
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> spawnHandler =
|
||||
(exchange, req) -> {
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.SPAWN, null);
|
||||
if (denied != null) return denied;
|
||||
String caller = callerTerminal(exchange);
|
||||
@@ -229,30 +250,72 @@ public final class BridgeMcp {
|
||||
return spawn(sessions, str(a, "profile"), str(a, "role"), str(a, "cwd"), callerCwd,
|
||||
callerTerminal(exchange), worktreeRequest(a),
|
||||
str(a, "sessionName"), str(a, "resumeSessionId"));
|
||||
})
|
||||
.toolCall(listTool(), (exchange, _) -> {
|
||||
};
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> listHandler =
|
||||
(exchange, _) -> {
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
|
||||
if (denied != null) return denied;
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine,
|
||||
callers == null ? Map.of() : callers.leads(),
|
||||
callerTerminal(exchange));
|
||||
})
|
||||
.toolCall(stopTool(), (exchange, req) -> {
|
||||
};
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> stopHandler =
|
||||
(exchange, req) -> {
|
||||
String paneId = str(req.arguments(), "paneId");
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.STOP, paneId);
|
||||
if (denied != null) return denied;
|
||||
return stop(sessions, paneId);
|
||||
})
|
||||
.toolCall(profilesTool(), (exchange, _) -> {
|
||||
};
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> profilesHandler =
|
||||
(exchange, _) -> {
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
|
||||
if (denied != null) return denied;
|
||||
return profiles(workers, quarantine);
|
||||
})
|
||||
.toolCall(whoamiTool(), (exchange, _) -> {
|
||||
};
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> whoamiHandler =
|
||||
(exchange, _) -> {
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
|
||||
if (denied != null) return denied;
|
||||
return whoami(principal(exchange), sessions);
|
||||
})
|
||||
};
|
||||
|
||||
McpSchema.Tool fleetSend = sendTool();
|
||||
McpSchema.Tool fleetReply = replyTool();
|
||||
McpSchema.Tool fleetAsk = askTool();
|
||||
McpSchema.Tool fleetStatus = statusTool();
|
||||
McpSchema.Tool fleetPoll = pollTool();
|
||||
McpSchema.Tool fleetAck = ackTool();
|
||||
McpSchema.Tool fleetSpawn = spawnTool();
|
||||
McpSchema.Tool fleetList = listTool();
|
||||
McpSchema.Tool fleetStop = stopTool();
|
||||
McpSchema.Tool fleetProfiles = profilesTool();
|
||||
McpSchema.Tool fleetWhoami = whoamiTool();
|
||||
|
||||
this.server = McpServer.sync(transport)
|
||||
.serverInfo("bridge", "0.1.0")
|
||||
.capabilities(McpSchema.ServerCapabilities.builder().tools(true).build())
|
||||
.toolCall(fleetSend, sendHandler)
|
||||
.toolCall(deprecatedTwin(fleetSend, "bridge_send"), deprecatedHandler(fleetSend, "bridge_send", sendHandler))
|
||||
.toolCall(fleetReply, replyHandler)
|
||||
.toolCall(deprecatedTwin(fleetReply, "bridge_reply"), deprecatedHandler(fleetReply, "bridge_reply", replyHandler))
|
||||
.toolCall(fleetAsk, askHandler)
|
||||
.toolCall(deprecatedTwin(fleetAsk, "bridge_ask"), deprecatedHandler(fleetAsk, "bridge_ask", askHandler))
|
||||
.toolCall(fleetStatus, statusHandler)
|
||||
.toolCall(deprecatedTwin(fleetStatus, "bridge_status"), deprecatedHandler(fleetStatus, "bridge_status", statusHandler))
|
||||
.toolCall(fleetPoll, pollHandler)
|
||||
.toolCall(deprecatedTwin(fleetPoll, "bridge_poll"), deprecatedHandler(fleetPoll, "bridge_poll", pollHandler))
|
||||
.toolCall(fleetAck, ackHandler)
|
||||
.toolCall(deprecatedTwin(fleetAck, "bridge_ack"), deprecatedHandler(fleetAck, "bridge_ack", ackHandler))
|
||||
.toolCall(fleetSpawn, spawnHandler)
|
||||
.toolCall(deprecatedTwin(fleetSpawn, "bridge_spawn"), deprecatedHandler(fleetSpawn, "bridge_spawn", spawnHandler))
|
||||
.toolCall(fleetList, listHandler)
|
||||
.toolCall(deprecatedTwin(fleetList, "bridge_list"), deprecatedHandler(fleetList, "bridge_list", listHandler))
|
||||
.toolCall(fleetStop, stopHandler)
|
||||
.toolCall(deprecatedTwin(fleetStop, "bridge_stop"), deprecatedHandler(fleetStop, "bridge_stop", stopHandler))
|
||||
.toolCall(fleetProfiles, profilesHandler)
|
||||
.toolCall(deprecatedTwin(fleetProfiles, "bridge_profiles"), deprecatedHandler(fleetProfiles, "bridge_profiles", profilesHandler))
|
||||
.toolCall(fleetWhoami, whoamiHandler)
|
||||
.toolCall(deprecatedTwin(fleetWhoami, "bridge_whoami"), deprecatedHandler(fleetWhoami, "bridge_whoami", whoamiHandler))
|
||||
.build();
|
||||
this.authz = callers;
|
||||
this.metrics = metrics;
|
||||
@@ -406,7 +469,7 @@ public final class BridgeMcp {
|
||||
// --- tool logic (thin adapters over the services; unit-testable) ---------------------------
|
||||
|
||||
/**
|
||||
* {@code bridge_send}: delegate {@code content} to a worker session and block for its reply.
|
||||
* {@code fleet_send}: delegate {@code content} to a worker session and block for its reply.
|
||||
* The configured profiles are required so a profile name can never bypass target validation.
|
||||
*
|
||||
* (CB-548): {@code onAccepted} records delegator ownership the instant the send is accepted, so
|
||||
@@ -430,8 +493,8 @@ public final class BridgeMcp {
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code bridge_send} carrying a {@code turnId}: the primary's answer to a worker's
|
||||
* {@code bridge_ask} (CB-205). Resolves the worker's blocked question and blocks for its reply as
|
||||
* {@code fleet_send} carrying a {@code turnId}: the primary's answer to a worker's
|
||||
* {@code fleet_ask} (CB-205). Resolves the worker's blocked question and blocks for its reply as
|
||||
* it resumes the same turn — surfaced to the primary identically to a normal send.
|
||||
*/
|
||||
static McpSchema.CallToolResult answer(MessageService messages, String turnId, String content, Long timeoutMs) {
|
||||
@@ -443,13 +506,13 @@ public final class BridgeMcp {
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code bridge_ask} (CB-205): a worker pauses its delegated turn to ask the primary, blocking
|
||||
* {@code fleet_ask} (CB-205): a worker pauses its delegated turn to ask the primary, blocking
|
||||
* until the primary answers. The worker is identified by its connection ({@code callerTerminal}),
|
||||
* never an argument — a {@code null} means the caller is not a known worker.
|
||||
*/
|
||||
static McpSchema.CallToolResult ask(MessageService messages, String callerTerminal, String question, Long timeoutMs) {
|
||||
if (callerTerminal == null) {
|
||||
return error("bridge_ask is for workers only — could not identify the calling worker "
|
||||
return error("fleet_ask is for workers only — could not identify the calling worker "
|
||||
+ "from the connection");
|
||||
}
|
||||
if (isBlank(question)) {
|
||||
@@ -459,10 +522,10 @@ public final class BridgeMcp {
|
||||
MessageService.AskResult r = messages.ask(callerTerminal, question, timeout);
|
||||
return switch (r.outcome()) {
|
||||
case ANSWERED -> text(r.answer());
|
||||
case NO_WAITER -> error("no primary is awaiting this turn — bridge_ask only works while a "
|
||||
+ "bridge_send delegation is open to answer it");
|
||||
case NO_WAITER -> error("no primary is awaiting this turn — fleet_ask only works while a "
|
||||
+ "fleet_send delegation is open to answer it");
|
||||
case TIMED_OUT -> text("[no answer within " + timeout + "ms — the primary did not respond; "
|
||||
+ "proceed on your best judgement, then call bridge_reply to end the turn]");
|
||||
+ "proceed on your best judgement, then call fleet_reply to end the turn]");
|
||||
};
|
||||
}
|
||||
|
||||
@@ -470,10 +533,10 @@ public final class BridgeMcp {
|
||||
private static McpSchema.CallToolResult formatReply(MessageService.Reply r, long timeout) {
|
||||
return switch (r.outcome()) {
|
||||
case REPLIED -> text(r.text());
|
||||
// The worker's turn finished but it never called bridge_reply — hand back the scraped
|
||||
// The worker's turn finished but it never called fleet_reply — hand back the scraped
|
||||
// transcript tail, flagged so the primary knows it isn't a structured reply.
|
||||
case COMPLETED_UNREPLIED -> text(
|
||||
"[worker finished without a structured bridge_reply — transcript tail follows]\n" + r.text());
|
||||
"[worker finished without a structured fleet_reply — transcript tail follows]\n" + r.text());
|
||||
// The worker ran the turn then wedged (CB-109) — surface the error context.
|
||||
case WORKER_FAILED -> text("[worker failed — turn ended in an unrecoverable state]\n" + r.text());
|
||||
// The backend refused on a subscription usage limit (CB-578 stage A) — the worker's
|
||||
@@ -483,7 +546,7 @@ public final class BridgeMcp {
|
||||
+ "usage limit]\n" + r.text());
|
||||
// The worker paused mid-turn to ask (CB-205) — tell the primary how to answer in-turn.
|
||||
case QUESTION -> text("[question] the worker paused to ask before it can finish:\n" + r.text()
|
||||
+ "\n\nAnswer it by calling bridge_send again with turnId=\"" + r.turnId()
|
||||
+ "\n\nAnswer it by calling fleet_send again with turnId=\"" + r.turnId()
|
||||
+ "\" and content set to your answer; the worker resumes the same turn.");
|
||||
case STALE_TURN -> error("that question is no longer open — it timed out or was already "
|
||||
+ "answered (turnId stale)");
|
||||
@@ -493,7 +556,7 @@ public final class BridgeMcp {
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code bridge_send} with {@code wait:false}: delegate {@code content} and return a ticket
|
||||
* {@code fleet_send} with {@code wait:false}: delegate {@code content} and return a ticket
|
||||
* immediately (fire-and-poll), so a long task isn't cut off by the caller's MCP call timeout.
|
||||
* The configured profiles are required so a profile name can never bypass target validation.
|
||||
*
|
||||
@@ -510,19 +573,19 @@ public final class BridgeMcp {
|
||||
return targetError;
|
||||
}
|
||||
String ticket = messages.sendAsync(sessionId, content, onAccepted);
|
||||
return text("accepted — task delegated. Poll bridge_poll with ticket=" + ticket);
|
||||
return text("accepted — task delegated. Poll fleet_poll with ticket=" + ticket);
|
||||
}
|
||||
|
||||
/** A configured profile is never a send target; other unknown values may be herdr-owned panes. */
|
||||
private static McpSchema.CallToolResult profileTargetError(String sessionId, Set<String> profiles) {
|
||||
if (profiles.contains(sessionId)) {
|
||||
return error("unknown send target \"" + sessionId + "\": it is a configured profile name, not a "
|
||||
+ "session id. Call bridge_list to find a member or lead sessionId.");
|
||||
+ "session id. Call fleet_list to find a member or lead sessionId.");
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
/** {@code bridge_poll}: check an async delegation by ticket, or drain a worker's inbox by target. */
|
||||
/** {@code fleet_poll}: check an async delegation by ticket, or drain a worker's inbox by target. */
|
||||
static McpSchema.CallToolResult poll(MessageService messages, String ticket, String target) {
|
||||
if (!isBlank(target)) {
|
||||
var replies = messages.drainReplies(target);
|
||||
@@ -540,25 +603,25 @@ public final class BridgeMcp {
|
||||
}
|
||||
return switch (v.phase()) {
|
||||
case DONE -> text(v.replySource() != null && v.replySource().equals("transcript")
|
||||
? "[done — worker finished without a structured bridge_reply; transcript tail follows]\n" + v.reply()
|
||||
? "[done — worker finished without a structured fleet_reply; transcript tail follows]\n" + v.reply()
|
||||
: v.reply());
|
||||
case PENDING -> text("[pending — " + v.detail() + "]");
|
||||
case ASKING -> text("[question — worker is waiting for your answer]\n" + v.reply()
|
||||
+ "\n\nAnswer it by calling bridge_send again with turnId=\"" + v.turnId()
|
||||
+ "\n\nAnswer it by calling fleet_send again with turnId=\"" + v.turnId()
|
||||
+ "\" and content set to your answer; the worker resumes the same turn.");
|
||||
case FAILED -> text("[failed — " + v.detail() + "]");
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code bridge_reply}: the worker returns its structured answer, resolving the awaiting send
|
||||
* {@code fleet_reply}: the worker returns its structured answer, resolving the awaiting send
|
||||
* or — when no send is open — queueing the reply in the inbox for later drain (CB-307).
|
||||
* {@code callerTerminal} is resolved from the connection (never an argument); a {@code null}
|
||||
* means the caller is not a known worker (e.g. the primary called it by mistake).
|
||||
*/
|
||||
static McpSchema.CallToolResult reply(MessageService messages, String callerTerminal, String content) {
|
||||
if (callerTerminal == null) {
|
||||
return error("bridge_reply is for workers only — could not identify the calling worker "
|
||||
return error("fleet_reply is for workers only — could not identify the calling worker "
|
||||
+ "from the connection");
|
||||
}
|
||||
if (content == null) {
|
||||
@@ -568,7 +631,7 @@ public final class BridgeMcp {
|
||||
return text("delivered");
|
||||
}
|
||||
|
||||
/** {@code bridge_ack}: acknowledge (remove) a specific reply from the inbox. */
|
||||
/** {@code fleet_ack}: acknowledge (remove) a specific reply from the inbox. */
|
||||
static McpSchema.CallToolResult ack(MessageService messages, String target, String msgId) {
|
||||
if (isBlank(target) || isBlank(msgId)) {
|
||||
return error("target and msgId are required");
|
||||
@@ -578,8 +641,8 @@ public final class BridgeMcp {
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code bridge_status}: the live lifecycle status of a worker session, plus — when the worker
|
||||
* is paused mid-turn in an async {@code bridge_ask} (CB-582) — the open question and how to
|
||||
* {@code fleet_status}: the live lifecycle status of a worker session, plus — when the worker
|
||||
* is paused mid-turn in an async {@code fleet_ask} (CB-582) — the open question and how to
|
||||
* answer it, so a lead on its normal poll cadence does not need the ticket to notice.
|
||||
*/
|
||||
static McpSchema.CallToolResult status(MessageService messages, String sessionId) {
|
||||
@@ -593,7 +656,7 @@ public final class BridgeMcp {
|
||||
return text(base);
|
||||
}
|
||||
return text(base + "\n\n[question — worker is waiting for your answer]\n" + ask.question()
|
||||
+ "\n\nAnswer it by calling bridge_send again with turnId=\"" + ask.turnId()
|
||||
+ "\n\nAnswer it by calling fleet_send again with turnId=\"" + ask.turnId()
|
||||
+ "\" and content set to your answer; the worker resumes the same turn."
|
||||
+ " (ticket " + ask.ticket() + ")");
|
||||
} catch (HerdrException e) {
|
||||
@@ -602,7 +665,7 @@ public final class BridgeMcp {
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code bridge_whoami}: the caller's own identity, as the daemon already resolved it.
|
||||
* {@code fleet_whoami}: the caller's own identity, as the daemon already resolved it.
|
||||
*
|
||||
* <p>Every other tool <em>consumes</em> this identity — the authorization gate, the reply
|
||||
* rendezvous, the cwd inherit — but none reported it, so an agent had to infer its own role
|
||||
@@ -610,7 +673,7 @@ public final class BridgeMcp {
|
||||
* name its MCP mount happens to carry, or {@code ANTHROPIC_BASE_URL} (which Claude-model
|
||||
* workers do not set). The failure mode of guessing is asymmetric and silent: a primary that
|
||||
* mistakes itself for a worker is refused by {@link Authz} and learns immediately, while a
|
||||
* worker that mistakes itself for the primary ends its turn without {@code bridge_reply} and
|
||||
* worker that mistakes itself for the primary ends its turn without {@code fleet_reply} and
|
||||
* the sender simply receives nothing. This tool removes the guess.
|
||||
*
|
||||
* <p>For a worker the session registry adds what it knows about that session. A worker the
|
||||
@@ -668,13 +731,13 @@ public final class BridgeMcp {
|
||||
|
||||
// --- fleet management logic (CB-108 / CB-301) --------------------------------------------
|
||||
|
||||
/** {@code bridge_spawn} without cwd/caller context (default resolution). */
|
||||
/** {@code fleet_spawn} without cwd/caller context (default resolution). */
|
||||
static McpSchema.CallToolResult spawn(SessionManager sessions, String profile) {
|
||||
return spawn(sessions, profile, null, null, null, null, null, null, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code bridge_spawn}: launch a guard-checked member for {@code profile} (blank → the default
|
||||
* {@code fleet_spawn}: launch a guard-checked member for {@code profile} (blank → the default
|
||||
* profile) under {@code role} (blank → {@code dev}), and return its session id + pane id. The
|
||||
* member's cwd is {@code requestedCwd} if given, else the profile's config, else
|
||||
* {@code callerCwd} (the primary's directory), else the daemon's.
|
||||
@@ -718,7 +781,7 @@ public final class BridgeMcp {
|
||||
}
|
||||
}
|
||||
|
||||
/** Build a {@link WorktreeRequest} from {@code bridge_spawn}'s optional {@code worktree}/{@code ticket} args. */
|
||||
/** Build a {@link WorktreeRequest} from {@code fleet_spawn}'s optional {@code worktree}/{@code ticket} args. */
|
||||
private static WorktreeRequest worktreeRequest(Map<String, Object> a) {
|
||||
Object w = a.get("worktree");
|
||||
if (w == null || Boolean.FALSE.equals(w)) {
|
||||
@@ -747,7 +810,7 @@ public final class BridgeMcp {
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code bridge_profiles}: the configured worker profiles, the default, and — CB-578 stage B —
|
||||
* {@code fleet_profiles}: the configured worker profiles, the default, and — CB-578 stage B —
|
||||
* which of them are currently quarantined (backend exhausted) and for how much longer. The
|
||||
* {@code quarantined} key is present only when at least one profile is, so a fleet where nothing
|
||||
* has ever been quarantined gets exactly the pre-stage-B shape.
|
||||
@@ -776,7 +839,7 @@ public final class BridgeMcp {
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code bridge_list}: the whole fleet — {@code leads} and {@code workers} — each merged with
|
||||
* {@code fleet_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.
|
||||
*
|
||||
@@ -789,10 +852,10 @@ public final class BridgeMcp {
|
||||
* <p>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.
|
||||
* a peer's, and the alternative is every lead calling {@code fleet_whoami} to subtract itself.
|
||||
*
|
||||
* <p>CB-583: the {@code capacity} rows reuse {@code quarantine} (the same {@link QuarantineSource}
|
||||
* {@code bridge_profiles} reads) so the two surfaces cannot disagree about which profile is
|
||||
* {@code fleet_profiles} reads) so the two surfaces cannot disagree about which profile is
|
||||
* quarantined — see {@link #capacityView}.
|
||||
*
|
||||
* @param leads terminal_id → lead name, live from the resolver
|
||||
@@ -856,7 +919,7 @@ public final class BridgeMcp {
|
||||
* CB-583: {@code free} alone cannot tell a lead "busy, will free up" from "refusing, and
|
||||
* nothing changes for N seconds" — those need different decisions. So a quarantined profile
|
||||
* forces {@code free} to 0, whatever its {@code maxLoad}/{@code live} say, and the row carries
|
||||
* the same {@code credentialId}/{@code quarantinedForSeconds} facts {@code bridge_profiles}
|
||||
* the same {@code credentialId}/{@code quarantinedForSeconds} facts {@code fleet_profiles}
|
||||
* reports, reusing {@link QuarantineSource} rather than a second lookup. Both new keys are
|
||||
* added only when the profile is actually quarantined, so an ordinary fleet's rows are
|
||||
* byte-identical to before this change.
|
||||
@@ -889,7 +952,7 @@ public final class BridgeMcp {
|
||||
*
|
||||
* <p>{@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,
|
||||
* see is a lead a {@code fleet_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<String, Object> leadView(String terminal, String name, Agent live,
|
||||
@@ -905,7 +968,7 @@ public final class BridgeMcp {
|
||||
return m;
|
||||
}
|
||||
|
||||
/** {@code bridge_stop}: tear a worker down by its pane id. */
|
||||
/** {@code fleet_stop}: tear a worker down by its pane id. */
|
||||
static McpSchema.CallToolResult stop(SessionManager sessions, String paneId) {
|
||||
if (isBlank(paneId)) {
|
||||
return error("paneId is required");
|
||||
@@ -948,25 +1011,25 @@ public final class BridgeMcp {
|
||||
// --- tool schemas --------------------------------------------------------------------------
|
||||
|
||||
private static McpSchema.Tool sendTool() {
|
||||
return tool("bridge_send",
|
||||
return tool("fleet_send",
|
||||
"Delegate a task to a worker session. By default blocks until the worker replies and "
|
||||
+ "returns its reply (or a 'still working / queued' note on timeout). Pass wait:false "
|
||||
+ "for a long task to return a ticket immediately, then poll it with bridge_poll. To "
|
||||
+ "answer a worker's bridge_ask, pass its turnId (with content) instead of sessionId.",
|
||||
+ "for a long task to return a ticket immediately, then poll it with fleet_poll. To "
|
||||
+ "answer a worker's fleet_ask, pass its turnId (with content) instead of sessionId.",
|
||||
objectSchema(Map.of(
|
||||
"sessionId", stringProp("The worker session id (herdr terminal_id) to delegate to"),
|
||||
"content", stringProp("The task/message to send to the worker (or your answer, with turnId)"),
|
||||
"timeoutMs", Map.of("type", "integer", "description", "Max ms to wait for a reply (blocking mode)"),
|
||||
"wait", Map.of("type", "boolean",
|
||||
"description", "Block for the reply (default true); false returns a ticket to poll"),
|
||||
"turnId", stringProp("When answering a worker's bridge_ask, its question turnId — "
|
||||
"turnId", stringProp("When answering a worker's fleet_ask, its question turnId — "
|
||||
+ "routes your answer back into the same turn (omit for a normal delegation)")),
|
||||
List.of("content")));
|
||||
}
|
||||
|
||||
private static McpSchema.Tool askTool() {
|
||||
// No target/session arg — the worker's identity is resolved from the connection.
|
||||
return tool("bridge_ask",
|
||||
return tool("fleet_ask",
|
||||
"Pause your current delegated turn to ask the primary a question, blocking until it "
|
||||
+ "answers — then resume the same turn with the answer. Use this when only the "
|
||||
+ "primary has a decision or detail you need to continue. You do not address the "
|
||||
@@ -979,19 +1042,19 @@ public final class BridgeMcp {
|
||||
}
|
||||
|
||||
private static McpSchema.Tool pollTool() {
|
||||
return tool("bridge_poll",
|
||||
"Check an async delegation (a bridge_send with wait:false) by its ticket: "
|
||||
return tool("fleet_poll",
|
||||
"Check an async delegation (a fleet_send with wait:false) by its ticket: "
|
||||
+ "pending, done (with the worker's reply), or failed. When target (a worker "
|
||||
+ "session id) is present instead of ticket, drain that worker's inbox of "
|
||||
+ "replies delivered when no send was open.",
|
||||
objectSchema(Map.of(
|
||||
"ticket", stringProp("The ticket returned by bridge_send wait:false"),
|
||||
"ticket", stringProp("The ticket returned by fleet_send wait:false"),
|
||||
"target", stringProp("Worker session id to drain pending replies from (optional)")),
|
||||
List.of()));
|
||||
}
|
||||
|
||||
private static McpSchema.Tool ackTool() {
|
||||
return tool("bridge_ack",
|
||||
return tool("fleet_ack",
|
||||
"Acknowledge (remove) a specific reply from a worker's inbox. Use when the primary "
|
||||
+ "has processed a reply and wants to confirm it, leaving other pending replies "
|
||||
+ "in the inbox for later drain.",
|
||||
@@ -1002,22 +1065,22 @@ public final class BridgeMcp {
|
||||
}
|
||||
|
||||
private static McpSchema.Tool spawnTool() {
|
||||
return tool("bridge_spawn",
|
||||
return tool("fleet_spawn",
|
||||
"Spawn a new off-subscription member session. A member has two independent attributes: "
|
||||
+ "role (what it is for) and profile (which backend it runs on). Pass role to pick "
|
||||
+ "the contract — 'dev' implements a unit and opens its own PR, 'reviewer' reviews a "
|
||||
+ "diff it did not write, 'architect' refines a ticket before anyone builds it; omit "
|
||||
+ "it for 'dev'. Pass profile (from bridge_profiles) to pick the backend, or omit it "
|
||||
+ "it for 'dev'. Pass profile (from fleet_profiles) to pick the backend, or omit it "
|
||||
+ "for the default. The two are independent: a reviewer may run on the same profile "
|
||||
+ "as the dev it reviews. The member opens your current directory by default; pass "
|
||||
+ "cwd to pin a different one. Pass worktree:true (with ticket) or "
|
||||
+ "worktree:<ticket-slug> to provision an isolated git worktree. Pass resumeSessionId "
|
||||
+ "to relaunch onto a prior conversation instead of starting cold — this requires an "
|
||||
+ "explicit profile whose backend supports it (bridge_list shows agentSessionId for "
|
||||
+ "explicit profile whose backend supports it (fleet_list shows agentSessionId for "
|
||||
+ "resumable members), and is refused otherwise rather than silently starting fresh. "
|
||||
+ "sessionName gives the member a display name in its own UI when the backend supports "
|
||||
+ "one. Returns the member's sessionId (use with bridge_send) and paneId (use with "
|
||||
+ "bridge_stop).",
|
||||
+ "one. Returns the member's sessionId (use with fleet_send) and paneId (use with "
|
||||
+ "fleet_stop).",
|
||||
objectSchema(Map.of(
|
||||
"role", stringProp("What the member is for: architect, dev or reviewer (default dev)"),
|
||||
"profile", stringProp("Which backend to run it on (omit for the default profile)"),
|
||||
@@ -1025,33 +1088,33 @@ public final class BridgeMcp {
|
||||
"worktree", Map.of("type", "string", "description", "'true' or a ticket slug — requests an isolated git worktree"),
|
||||
"ticket", stringProp("Ticket slug when worktree:true"),
|
||||
"sessionName", stringProp("Logical display name for the member's own session, when its backend supports one"),
|
||||
"resumeSessionId", stringProp("A prior member's agentSessionId (from bridge_list) to resume — requires an explicit profile that supports it")),
|
||||
"resumeSessionId", stringProp("A prior member's agentSessionId (from fleet_list) to resume — requires an explicit profile that supports it")),
|
||||
List.of()));
|
||||
}
|
||||
|
||||
private static McpSchema.Tool profilesTool() {
|
||||
return tool("bridge_profiles",
|
||||
"List the configured worker profiles (backends) and which one bridge_spawn uses by "
|
||||
return tool("fleet_profiles",
|
||||
"List the configured worker profiles (backends) and which one fleet_spawn uses by "
|
||||
+ "default. A 'quarantined' map is present when a backend-exhausted refusal put "
|
||||
+ "a profile's credential on cooldown — bridge_spawn onto it is refused until "
|
||||
+ "a profile's credential on cooldown — fleet_spawn onto it is refused until "
|
||||
+ "quarantinedForSeconds elapses; a profile sharing that credential is listed too.",
|
||||
objectSchema(Map.of(), List.of()));
|
||||
}
|
||||
|
||||
private static McpSchema.Tool listTool() {
|
||||
return tool("bridge_list",
|
||||
return tool("fleet_list",
|
||||
"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), "
|
||||
+ "orchestrators, each with its sessionId (the address to fleet_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. 'members' are the "
|
||||
+ "sessions delegated to — each with sessionId, paneId, role (architect/dev/"
|
||||
+ "reviewer), profile (the backend it runs on), state, optional "
|
||||
+ "worktree/branch/owner/agentSessionId (the id to pass as bridge_spawn's "
|
||||
+ "worktree/branch/owner/agentSessionId (the id to pass as fleet_spawn's "
|
||||
+ "resumeSessionId to relaunch onto that same conversation, when the backend "
|
||||
+ "supports it), and live herdr status. An empty 'members' "
|
||||
+ "means no members are spawned; it says nothing about peers. When capacity "
|
||||
+ "facts are configured, a 'capacity' row per profile also reports free: 0 for "
|
||||
+ "a quarantined profile's credential (see bridge_profiles), whatever its "
|
||||
+ "a quarantined profile's credential (see fleet_profiles), whatever its "
|
||||
+ "maxLoad/live — with credentialId and quarantinedForSeconds naming the "
|
||||
+ "quarantine, so 'free: 0, busy' can be told apart from 'free: 0, refusing "
|
||||
+ "for N seconds'.",
|
||||
@@ -1059,8 +1122,8 @@ public final class BridgeMcp {
|
||||
}
|
||||
|
||||
private static McpSchema.Tool stopTool() {
|
||||
return tool("bridge_stop",
|
||||
"Tear down a worker session by its paneId (from bridge_spawn or bridge_list).",
|
||||
return tool("fleet_stop",
|
||||
"Tear down a worker session by its paneId (from fleet_spawn or fleet_list).",
|
||||
objectSchema(Map.of(
|
||||
"paneId", stringProp("The worker's paneId to stop")),
|
||||
List.of("paneId")));
|
||||
@@ -1068,9 +1131,9 @@ public final class BridgeMcp {
|
||||
|
||||
private static McpSchema.Tool replyTool() {
|
||||
// No session/target arg — the caller's identity is resolved from the connection.
|
||||
return tool("bridge_reply",
|
||||
return tool("fleet_reply",
|
||||
"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 "
|
||||
+ "blocked fleet_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(
|
||||
@@ -1079,7 +1142,7 @@ public final class BridgeMcp {
|
||||
}
|
||||
|
||||
private static McpSchema.Tool statusTool() {
|
||||
return tool("bridge_status",
|
||||
return tool("fleet_status",
|
||||
"Get the live lifecycle status (idle/working/blocked/unknown) of a worker session.",
|
||||
objectSchema(Map.of(
|
||||
"sessionId", stringProp("The worker session id to query")),
|
||||
@@ -1087,20 +1150,58 @@ public final class BridgeMcp {
|
||||
}
|
||||
|
||||
private static McpSchema.Tool whoamiTool() {
|
||||
return tool("bridge_whoami",
|
||||
return tool("fleet_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; reply ONLY to answer a peer lead that "
|
||||
+ "messaged you, never to answer a worker), 'architect' (you delegate turns "
|
||||
+ "and reply/ask as your own pane, but cannot spawn/stop/drain), or 'worker' "
|
||||
+ "(you were delegated to: you must end every turn with exactly one "
|
||||
+ "bridge_reply, and cannot spawn or send), plus 'leader'/'architect' naming "
|
||||
+ "fleet_reply, and cannot spawn or send), plus 'leader'/'architect' naming "
|
||||
+ "which one 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()));
|
||||
}
|
||||
|
||||
// --- CB-622: bridge_* -> fleet_* rename, kept working under both names -------------------
|
||||
|
||||
/**
|
||||
* The deprecated {@code bridge_*} twin of {@code fleetTool}: same name-minus-prefix schema,
|
||||
* with a description that leads with the deprecation notice so a client listing tools sees it
|
||||
* immediately. Reuses {@code fleetTool}'s input schema rather than restating it, so the two
|
||||
* can never drift on parameters.
|
||||
*/
|
||||
static McpSchema.Tool deprecatedTwin(McpSchema.Tool fleetTool, String oldName) {
|
||||
return tool(oldName, "DEPRECATED: use " + fleetTool.name() + " instead. " + fleetTool.description(),
|
||||
fleetTool.inputSchema());
|
||||
}
|
||||
|
||||
/**
|
||||
* Wrap {@code handler} so a call under the deprecated {@code oldName} logs one WARN naming
|
||||
* the old and new name, then runs the exact SAME handler {@code fleetTool}'s name uses — no
|
||||
* logic is duplicated between the two registrations.
|
||||
*/
|
||||
static BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> deprecatedHandler(
|
||||
McpSchema.Tool fleetTool, String oldName,
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> handler) {
|
||||
return (exchange, req) -> {
|
||||
warnDeprecatedOnce(oldName, fleetTool.name());
|
||||
return handler.apply(exchange, req);
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Log one WARN naming {@code oldName} and {@code newName} — once per {@code oldName} for the
|
||||
* life of the process, not once per call. {@link #WARNED_DEPRECATED_NAMES} is a per-name flag
|
||||
* (a {@link Set}), not a call counter, so this never conflates "warned" with any numeric state.
|
||||
*/
|
||||
static void warnDeprecatedOnce(String oldName, String newName) {
|
||||
if (WARNED_DEPRECATED_NAMES.add(oldName)) {
|
||||
log.warn("{} is deprecated; use {} instead", oldName, newName);
|
||||
}
|
||||
}
|
||||
|
||||
// --- small helpers -------------------------------------------------------------------------
|
||||
|
||||
// The SDK 2.0.0 deprecates its own Tool builders without a stable replacement — isolate it here.
|
||||
|
||||
@@ -11,7 +11,7 @@ import java.util.concurrent.atomic.AtomicReference;
|
||||
* Single-slot, thread-safe registry for the primary's herdr {@code terminal_id}.
|
||||
*
|
||||
* <p>Populated from the caller terminal of orchestration-side MCP tools
|
||||
* ({@code bridge_send}, {@code bridge_spawn}) — tools that only the primary calls.
|
||||
* ({@code fleet_send}, {@code fleet_spawn}) — tools that only the primary calls.
|
||||
* A pinned terminal (from config) seeds the registry at construction and makes
|
||||
* subsequent {@link #record(String)} calls no-ops.
|
||||
*
|
||||
@@ -29,7 +29,7 @@ public final class PrimaryRegistry {
|
||||
/**
|
||||
* 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
|
||||
* whichever called {@code fleet_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.
|
||||
*/
|
||||
@@ -71,7 +71,7 @@ public final class PrimaryRegistry {
|
||||
*
|
||||
* <p>Called from the {@code MessageService} accepted-delivery hook — only after a send has won
|
||||
* the session's send lock and queued delivery — where both halves are known (CB-548). It is
|
||||
* deliberately <em>not</em> called at {@code bridge_send} request time: a concurrent sender that
|
||||
* deliberately <em>not</em> called at {@code fleet_send} request time: a concurrent sender that
|
||||
* times out {@code BUSY} must not steal a live delegation's reply routing without ever owning
|
||||
* the turn. Last writer wins — if a second lead's later send is accepted, replies follow the
|
||||
* lead that most recently delegated to it, which is the one waiting.
|
||||
|
||||
@@ -214,7 +214,7 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
}
|
||||
}
|
||||
}
|
||||
// Order-preserving for the same reason, and because profiles() is user-visible (bridge_profiles).
|
||||
// Order-preserving for the same reason, and because profiles() is user-visible (fleet_profiles).
|
||||
this.byProfile = Collections.unmodifiableMap(index);
|
||||
}
|
||||
|
||||
@@ -267,7 +267,7 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
|
||||
// CB-557: an unqualified spawn is placed inside the pool of the role it asked for, not across
|
||||
// the whole profile list. An EXPLICIT profile (above) is left alone on purpose — it is the
|
||||
// operator overriding, and refusing it would break `bridge_spawn{profile:"opus"}`, which
|
||||
// operator overriding, and refusing it would break `fleet_spawn{profile:"opus"}`, which
|
||||
// carries no role and so would be judged against the dev pool it was never meant for.
|
||||
List<PlacementCandidate> candidates = candidates(req.role());
|
||||
String roleDefault = defaultProfileFor(req.role());
|
||||
|
||||
@@ -117,14 +117,14 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
/** The final instruction always requires a bridge reply when the bridge MCP is mounted. */
|
||||
protected static final String REPLY_CHARTER =
|
||||
"You are a spawned member in the claude-bridge fleet. Every message you receive arrives "
|
||||
+ "through the bridge, and the ONLY channel back to the sender is the bridge_reply MCP tool. "
|
||||
+ "through the bridge, and the ONLY channel back to the sender is the fleet_reply MCP tool. "
|
||||
+ "Text you write in your terminal is NOT sent anywhere — the sender cannot see your screen, "
|
||||
+ "so an in-terminal answer is silently discarded. Therefore you MUST end EVERY turn by calling "
|
||||
+ "bridge_reply with `content` set to your complete response. This holds for every message without "
|
||||
+ "fleet_reply with `content` set to your complete response. This holds for every message without "
|
||||
+ "exception — tasks, questions, clarifications, acknowledgements, and ordinary back-and-forth "
|
||||
+ "conversation. Call bridge_reply exactly once, as the final action of your turn, with your full "
|
||||
+ "conversation. Call fleet_reply exactly once, as the final action of your turn, with your full "
|
||||
+ "answer in `content`; never wait for confirmation first. If you end a turn without calling "
|
||||
+ "bridge_reply, the sender receives nothing and the exchange stalls.";
|
||||
+ "fleet_reply, the sender receives nothing and the exchange stalls.";
|
||||
|
||||
/**
|
||||
* Tab numbers, counted per {@code role/profile} pair (CB-557).
|
||||
@@ -413,7 +413,7 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
: replyCharter == null ? roleCharter : roleCharter + "\n\n" + replyCharter;
|
||||
// CB-571: fingerprint the exact composed charter bytes once, here in the base, before the
|
||||
// string leaves for an adapter — so Claude and OpenCode derive the same digest. A failed
|
||||
// start has no bridge_spawn result and no roster row, so the failure log below is the only
|
||||
// start has no fleet_spawn result and no roster row, so the failure log below is the only
|
||||
// surface the byte count can appear on. The charter text itself is never logged.
|
||||
CharterReceipt receipt = CharterReceipt.compose(role, cfg.profile(), roleCharter, charter);
|
||||
String cwd = resolveCwd(requestedCwd, cfg, callerCwd);
|
||||
|
||||
@@ -247,7 +247,7 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
* permissions opencode does not explicitly deny. It is unconditional, not a preference: a
|
||||
* spawned peer has no human at its pane — the bridge spawned it — so one that stops at an
|
||||
* approval prompt is a wedged agent, indistinguishable from a legitimate mid-turn wait and
|
||||
* unable to end its turn with {@code bridge_reply}. opencode's own help calls this
|
||||
* unable to end its turn with {@code fleet_reply}. opencode's own help calls this
|
||||
* "dangerous!", but the blast radius here is already bounded by design: a worker runs in its
|
||||
* own git worktree on its own branch, is off-subscription, and cannot merge — the lead is the
|
||||
* gate.
|
||||
|
||||
@@ -18,7 +18,7 @@ import java.util.function.Supplier;
|
||||
|
||||
/**
|
||||
* CB-551: an opt-in heartbeat that nudges the single idle lead back to work once it has been
|
||||
* continuously idle past a quiet period with no open {@code bridge_send} driving it.
|
||||
* continuously idle past a quiet period with no open {@code fleet_send} driving it.
|
||||
*
|
||||
* <p>Why this exists: the fleet is ONE lead + architects + workers, so an idle, stalled lead is a
|
||||
* single point of failure for the fleet's progress. {@link ReplyPushLoop} nudges the lead only when
|
||||
@@ -290,7 +290,7 @@ public final class LeadHeartbeatLoop {
|
||||
/** The nudge body, phrased for the two cases the heartbeat distinguishes. */
|
||||
String nudgeText() {
|
||||
StringBuilder sb = new StringBuilder(
|
||||
"Heartbeat: you are idle and no bridge_send is waiting on you.");
|
||||
"Heartbeat: you are idle and no fleet_send is waiting on you.");
|
||||
if (hasPending()) {
|
||||
sb.append(" The fleet has state to collect: ").append(pendingDetail());
|
||||
} else {
|
||||
@@ -309,10 +309,10 @@ public final class LeadHeartbeatLoop {
|
||||
.append(" pending collection");
|
||||
if (!replyTargets.isEmpty()) {
|
||||
// Render each as the exact command so the lead can act without parsing: the nearest
|
||||
// analogue to ReplyPushLoop's bridge_poll(target=...) nudge.
|
||||
// analogue to ReplyPushLoop's fleet_poll(target=...) nudge.
|
||||
sb.append(" (")
|
||||
.append(String.join(", ",
|
||||
replyTargets.stream().map(t -> "bridge_poll(target=" + t + ")").toList()))
|
||||
replyTargets.stream().map(t -> "fleet_poll(target=" + t + ")").toList()))
|
||||
.append(")");
|
||||
}
|
||||
sb.append(", ").append(doneSessions).append(" DONE session").append(doneSessions == 1 ? "" : "s")
|
||||
|
||||
@@ -24,7 +24,7 @@ import java.util.function.LongSupplier;
|
||||
|
||||
/**
|
||||
* The blocking delegation feature (CB-104): deliver {@code content} into a worker and block until
|
||||
* the worker returns a <em>structured reply</em> via {@code bridge_reply} (the {@link Rendezvous}),
|
||||
* the worker returns a <em>structured reply</em> via {@code fleet_reply} (the {@link Rendezvous}),
|
||||
* then hand that reply back. Delivery is the {@link Injector}'s job (the background poller sends it
|
||||
* when the worker is injectable); this service never drives the injector or scrapes the terminal —
|
||||
* completion is the worker's explicit reply, not a guess about {@code agent_status}.
|
||||
@@ -62,10 +62,10 @@ public final class MessageService {
|
||||
|
||||
/** Outcome of a blocking send. */
|
||||
public enum Outcome {
|
||||
/** The worker called {@code bridge_reply}; {@code text} holds the structured answer. */
|
||||
/** The worker called {@code fleet_reply}; {@code text} holds the structured answer. */
|
||||
REPLIED,
|
||||
/**
|
||||
* The worker's delegated turn finished without a {@code bridge_reply} (CB-106 fallback);
|
||||
* The worker's delegated turn finished without a {@code fleet_reply} (CB-106 fallback);
|
||||
* {@code text} is the scraped transcript tail rather than a structured answer.
|
||||
*/
|
||||
COMPLETED_UNREPLIED,
|
||||
@@ -75,7 +75,7 @@ public final class MessageService {
|
||||
*/
|
||||
WORKER_FAILED,
|
||||
/**
|
||||
* The turn finished without a {@code bridge_reply} and the scrape matched the backend's
|
||||
* The turn finished without a {@code fleet_reply} and the scrape matched the backend's
|
||||
* configured usage-limit refusal pattern (CB-578 stage A); {@code text} is the reason,
|
||||
* carrying the matched line. The worker's pane is healthy — only its account is refusing —
|
||||
* so this is never reported as a completed reply, and is kept distinct from
|
||||
@@ -96,14 +96,14 @@ public final class MessageService {
|
||||
BUSY,
|
||||
/**
|
||||
* An answer ({@link #answer(String, String, long)}) referenced a {@code turnId} that is no
|
||||
* longer open — the worker's {@code bridge_ask} already timed out or was answered.
|
||||
* longer open — the worker's {@code fleet_ask} already timed out or was answered.
|
||||
*/
|
||||
STALE_TURN
|
||||
}
|
||||
|
||||
/**
|
||||
* @param outcome how the send ended (or paused)
|
||||
* @param text the worker's answer when {@link #completed()} (a structured {@code bridge_reply}
|
||||
* @param text the worker's answer when {@link #completed()} (a structured {@code fleet_reply}
|
||||
* for {@link Outcome#REPLIED}, a scraped transcript tail for
|
||||
* {@link Outcome#COMPLETED_UNREPLIED}), or the question for {@link Outcome#QUESTION},
|
||||
* else {@code null}
|
||||
@@ -122,7 +122,7 @@ public final class MessageService {
|
||||
}
|
||||
}
|
||||
|
||||
/** How a worker's {@code bridge_ask} (CB-205) resolved. */
|
||||
/** How a worker's {@code fleet_ask} (CB-205) resolved. */
|
||||
public enum AskOutcome {
|
||||
/** The primary answered; {@link AskResult#answer} carries it. */
|
||||
ANSWERED,
|
||||
@@ -132,7 +132,7 @@ public final class MessageService {
|
||||
TIMED_OUT
|
||||
}
|
||||
|
||||
/** The outcome of a worker's {@code bridge_ask}: how it resolved and (if answered) the answer. */
|
||||
/** The outcome of a worker's {@code fleet_ask}: how it resolved and (if answered) the answer. */
|
||||
public record AskResult(AskOutcome outcome, String answer) {
|
||||
}
|
||||
|
||||
@@ -140,7 +140,7 @@ public final class MessageService {
|
||||
public enum Phase {
|
||||
/** Delegated and in flight — queued for the worker or being worked. */
|
||||
PENDING,
|
||||
/** The worker is paused in {@code bridge_ask}; {@link TaskView#reply} and {@link TaskView#turnId} identify it. */
|
||||
/** The worker is paused in {@code fleet_ask}; {@link TaskView#reply} and {@link TaskView#turnId} identify it. */
|
||||
ASKING,
|
||||
/** The worker's turn finished; {@link TaskView#reply} holds the answer. */
|
||||
DONE,
|
||||
@@ -153,7 +153,7 @@ public final class MessageService {
|
||||
*
|
||||
* @param reply the answer when {@link #phase} is {@link Phase#DONE}, or the question when
|
||||
* {@link #phase} is {@link Phase#ASKING}; otherwise {@code null}
|
||||
* @param replySource {@code "reply"} (structured {@code bridge_reply}) or {@code "transcript"}
|
||||
* @param replySource {@code "reply"} (structured {@code fleet_reply}) or {@code "transcript"}
|
||||
* (completion scrape) when {@link Phase#DONE}, else {@code null}
|
||||
* @param detail a human note (live worker status while pending, ask state, or failure reason)
|
||||
* @param turnId correlation id for an {@link Phase#ASKING} ticket, else {@code null}
|
||||
@@ -179,11 +179,11 @@ public final class MessageService {
|
||||
}
|
||||
|
||||
/**
|
||||
* A worker session's currently-open {@code bridge_ask} question, surfaced so {@code bridge_status}
|
||||
* A worker session's currently-open {@code fleet_ask} question, surfaced so {@code fleet_status}
|
||||
* can show it without the caller needing the ticket first (CB-582). Only covers async
|
||||
* (fire-and-poll) delegations, which track the question on their {@link Task}; a blocking
|
||||
* ({@code wait:true}) send already hands the question straight back to its own caller, so there is
|
||||
* nothing hidden left for {@code bridge_status} to surface in that case.
|
||||
* nothing hidden left for {@code fleet_status} to surface in that case.
|
||||
*/
|
||||
public record PendingAsk(String ticket, String question, String turnId) {
|
||||
}
|
||||
@@ -202,7 +202,7 @@ public final class MessageService {
|
||||
/** Async task that owns each exact forward rendezvous waiter. */
|
||||
private final ConcurrentHashMap<CompletableFuture<Rendezvous.Resolution>, Task> asyncTasksByWaiter =
|
||||
new ConcurrentHashMap<>();
|
||||
/** Async tickets paused on a specific {@code bridge_ask} turn. */
|
||||
/** Async tickets paused on a specific {@code fleet_ask} turn. */
|
||||
private final ConcurrentHashMap<String, Task> asyncTasksByTurn = new ConcurrentHashMap<>();
|
||||
private final AtomicLong ticketSeq = new AtomicLong();
|
||||
private final ExecutorService asyncExecutor = Executors.newThreadPerTaskExecutor(
|
||||
@@ -216,7 +216,7 @@ public final class MessageService {
|
||||
* (CB-588) whenever an async ticket started by {@link #sendAsync} reaches a
|
||||
* terminal phase, whenever {@link #poll} hands a terminal ticket to its caller,
|
||||
* and (CB-582) whenever an async ticket's worker pauses mid-turn in
|
||||
* {@code bridge_ask} or that pause ends (answered or lapsed)
|
||||
* {@code fleet_ask} or that pause ends (answered or lapsed)
|
||||
*/
|
||||
public MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous,
|
||||
ReplyInbox inbox, ReplyPushLoop pushLoop) {
|
||||
@@ -272,11 +272,11 @@ public final class MessageService {
|
||||
}
|
||||
|
||||
/**
|
||||
* Route a worker's explicit {@code bridge_reply}: resolve an open send, or queue it in the
|
||||
* Route a worker's explicit {@code fleet_reply}: resolve an open send, or queue it in the
|
||||
* inbox if no send is currently open. Unlike the bare {@link Rendezvous#resolve}, a no-waiter
|
||||
* result is <em>not</em> a failure — the reply is held for later drain.
|
||||
*
|
||||
* <p><strong>Do NOT use this for mid-turn questions.</strong> {@code bridge_ask} /
|
||||
* <p><strong>Do NOT use this for mid-turn questions.</strong> {@code fleet_ask} /
|
||||
* {@link Rendezvous#resolveQuestion} must keep today's {@code NO_WAITER} behaviour — questions
|
||||
* are interactive and must never be queued.
|
||||
*
|
||||
@@ -329,7 +329,7 @@ public final class MessageService {
|
||||
* Abandon any send still waiting on {@code target} because its session has gone away (CB-516).
|
||||
*
|
||||
* <p>Without this, tearing a worker down left its rendezvous waiter open: a blocking
|
||||
* {@code bridge_send} kept blocking, and an async one kept reporting {@code PENDING} until
|
||||
* {@code fleet_send} kept blocking, and an async one kept reporting {@code PENDING} until
|
||||
* {@link #ASYNC_TIMEOUT_MS} — thirty minutes — even though the worker provably no longer
|
||||
* existed and the delegation could never complete. Worse, {@code poll} already had the evidence
|
||||
* (it calls {@code liveStatus} to build its detail string and gets back {@code "unknown"}) and
|
||||
@@ -484,7 +484,7 @@ public final class MessageService {
|
||||
|
||||
/**
|
||||
* A worker's mid-turn question (CB-205 reverse rendezvous): surface {@code question} to the
|
||||
* primary by resolving its open blocking {@code bridge_send}, then block this (worker) call until
|
||||
* primary by resolving its open blocking {@code fleet_send}, then block this (worker) call until
|
||||
* the primary answers via {@link #answer} or {@code timeoutMillis} elapses. Identity is the
|
||||
* worker's own session — it does not address the primary.
|
||||
*
|
||||
@@ -509,7 +509,7 @@ public final class MessageService {
|
||||
rendezvous.closeAsk(ticket.turnId());
|
||||
return new AskResult(AskOutcome.NO_WAITER, null); // no primary is blocked on this worker
|
||||
}
|
||||
// CB-582: the question just became visible via bridge_poll (Phase.ASKING) for an async
|
||||
// CB-582: the question just became visible via fleet_poll (Phase.ASKING) for an async
|
||||
// (wait:false) delegation — nudge the lead's own pane the same way a terminal ticket does
|
||||
// (CB-588), since the lead's normal poll cadence is minutes away and the reverse-rendezvous
|
||||
// window (~55s, see BridgeMcp/BridgedApp) is far shorter. A blocking (wait:true) send has
|
||||
@@ -523,7 +523,7 @@ public final class MessageService {
|
||||
String answer = ticket.answer().get(timeoutMillis, TimeUnit.MILLISECONDS);
|
||||
return new AskResult(AskOutcome.ANSWERED, answer);
|
||||
} catch (TimeoutException e) {
|
||||
log.debug("bridge_ask from {} went unanswered in {}ms", workerSession, timeoutMillis);
|
||||
log.debug("fleet_ask from {} went unanswered in {}ms", workerSession, timeoutMillis);
|
||||
clearAsyncQuestion(ticket.turnId(), true);
|
||||
return new AskResult(AskOutcome.TIMED_OUT, null);
|
||||
} catch (ExecutionException e) {
|
||||
@@ -549,13 +549,13 @@ public final class MessageService {
|
||||
}
|
||||
|
||||
/**
|
||||
* The primary's answer to a worker's {@code bridge_ask} (CB-205): resolve the worker's blocked
|
||||
* The primary's answer to a worker's {@code fleet_ask} (CB-205): resolve the worker's blocked
|
||||
* question identified by {@code turnId}, then — like a fresh {@link #send} — block for the worker's
|
||||
* eventual {@code bridge_reply} as it finishes the resumed turn. The worker session is derived from
|
||||
* eventual {@code fleet_reply} as it finishes the resumed turn. The worker session is derived from
|
||||
* {@code turnId}, never a caller argument.
|
||||
*
|
||||
* <p>Unlike {@link #send} this does not re-inject through the {@link Injector}: the worker is
|
||||
* mid-turn (already picked up), so the answer flows back through its own open {@code bridge_ask}
|
||||
* mid-turn (already picked up), so the answer flows back through its own open {@code fleet_ask}
|
||||
* call, not a new status-gated delivery. The forward waiter is opened <em>before</em> the worker
|
||||
* is unblocked so a reply that lands the instant it resumes is not lost.
|
||||
*/
|
||||
@@ -624,9 +624,9 @@ public final class MessageService {
|
||||
tasks.put(ticket, task);
|
||||
if (pushLoop != null) {
|
||||
// CB-588: task.future only ever completes on a terminal phase (DONE or a failure) — a
|
||||
// worker paused in bridge_ask leaves it running, per finishAsyncTask's own contract — so
|
||||
// worker paused in fleet_ask leaves it running, per finishAsyncTask's own contract — so
|
||||
// this fires exactly once, from whichever path completes it: finishAsyncTask(task, result)
|
||||
// below on any non-QUESTION outcome of send() — a worker's bridge_reply, the CB-106
|
||||
// below on any non-QUESTION outcome of send() — a worker's fleet_reply, the CB-106
|
||||
// completion fallback, a CB-109 wedge, TIMED_OUT, BUSY, or BACKEND_EXHAUSTED — the same
|
||||
// finishAsyncTask reached via answer()'s finishAsyncTask(turnId, result) once a QUESTION
|
||||
// is resolved, completeExceptionally(t) just below when send() itself throws, or a CB-516
|
||||
@@ -719,7 +719,7 @@ public final class MessageService {
|
||||
* too, its own {@code pendingTickets} entry would outlive the ticket it names: an unpolled ticket
|
||||
* (or one the reminder cap already gave up on) is pruned here but never collected there, so it
|
||||
* lingers in {@code pendingTickets} forever and rides along on every later nudge to the same lead
|
||||
* — naming a ticket {@code bridge_poll} can no longer find (CB-588 follow-up).
|
||||
* — naming a ticket {@code fleet_poll} can no longer find (CB-588 follow-up).
|
||||
*/
|
||||
private void pruneTerminalTickets() {
|
||||
long cutoff = nowNanos.getAsLong() - TICKET_TTL_NANOS;
|
||||
@@ -784,10 +784,10 @@ public final class MessageService {
|
||||
}
|
||||
|
||||
/**
|
||||
* The question {@code workerSession} is currently paused on via {@code bridge_ask}, if any
|
||||
* (CB-582) — {@code bridge_status} uses this to show a pending question without the caller
|
||||
* The question {@code workerSession} is currently paused on via {@code fleet_ask}, if any
|
||||
* (CB-582) — {@code fleet_status} uses this to show a pending question without the caller
|
||||
* needing the ticket. {@code null} when the session has no open async question (including a
|
||||
* session mid a <em>blocking</em> {@code bridge_ask}, which has no {@link Task} to look up — see
|
||||
* session mid a <em>blocking</em> {@code fleet_ask}, which has no {@link Task} to look up — see
|
||||
* {@link PendingAsk}).
|
||||
*/
|
||||
public PendingAsk pendingAsk(String workerSession) {
|
||||
|
||||
@@ -5,9 +5,9 @@ import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
|
||||
/**
|
||||
* The reply rendezvous: where a blocking {@code bridge_send} awaits how the worker's delegated turn
|
||||
* The reply rendezvous: where a blocking {@code fleet_send} awaits how the worker's delegated turn
|
||||
* ends. The sending (primary) request thread {@link #open}s a waiter; it is resolved either by the
|
||||
* worker's explicit {@code bridge_reply} ({@link #resolve}, arriving on a different thread via
|
||||
* worker's explicit {@code fleet_reply} ({@link #resolve}, arriving on a different thread via
|
||||
* {@code POST /sessions/{id}/reply}) or — the CB-106 fallback — by the injector observing the
|
||||
* worker's delegated turn return to idle without a reply ({@link #resolveCompletion}).
|
||||
*
|
||||
@@ -27,14 +27,14 @@ public final class Rendezvous {
|
||||
|
||||
/** How a delegated turn ended (or paused). */
|
||||
public enum Kind {
|
||||
/** The worker called {@code bridge_reply} with a structured answer. */
|
||||
/** The worker called {@code fleet_reply} with a structured answer. */
|
||||
REPLY,
|
||||
/** The worker's turn finished without a {@code bridge_reply}; {@code text} is a scrape. */
|
||||
/** The worker's turn finished without a {@code fleet_reply}; {@code text} is a scrape. */
|
||||
COMPLETION,
|
||||
/** The worker ran the turn then wedged (CB-109); {@code text} is the failure context. */
|
||||
FAILED,
|
||||
/**
|
||||
* The turn finished without a {@code bridge_reply}, and the scrape matched the backend's
|
||||
* The turn finished without a {@code fleet_reply}, and the scrape matched the backend's
|
||||
* configured usage-limit refusal pattern (CB-578 stage A); {@code text} is the reason,
|
||||
* carrying the matched line. The pane is healthy — only the account is refusing — so this
|
||||
* is kept separate from a session simply going {@code GONE}.
|
||||
@@ -43,7 +43,7 @@ public final class Rendezvous {
|
||||
/**
|
||||
* The worker paused mid-turn to ask the primary a question (CB-205 reverse rendezvous);
|
||||
* {@code text} is the question and {@code turnId} correlates the primary's answer back to
|
||||
* the worker's blocked {@code bridge_ask}. Not terminal — the turn resumes after the answer.
|
||||
* the worker's blocked {@code fleet_ask}. Not terminal — the turn resumes after the answer.
|
||||
*/
|
||||
QUESTION
|
||||
}
|
||||
@@ -75,7 +75,7 @@ public final class Rendezvous {
|
||||
/** Reverse rendezvous (CB-205): worker questions awaiting the primary's answer, keyed by {@code turnId}. */
|
||||
private final ConcurrentHashMap<String, AskWaiter> asks = new ConcurrentHashMap<>();
|
||||
private final AtomicLong askSeq = new AtomicLong();
|
||||
/** Per-session index of the currently-open ask, so duplicate bridge_ask calls coalesce onto one turn. */
|
||||
/** Per-session index of the currently-open ask, so duplicate fleet_ask calls coalesce onto one turn. */
|
||||
private final ConcurrentHashMap<String, String> openAsksBySession = new ConcurrentHashMap<>();
|
||||
|
||||
/**
|
||||
@@ -132,7 +132,7 @@ public final class Rendezvous {
|
||||
return complete(session, new Resolution(Kind.REPLY, content));
|
||||
}
|
||||
|
||||
// --- reverse rendezvous (CB-205 bridge_ask) ------------------------------------------------
|
||||
// --- reverse rendezvous (CB-205 fleet_ask) ------------------------------------------------
|
||||
|
||||
/**
|
||||
* Open a reverse-rendezvous waiter for a worker's mid-turn question. If {@code session} already has
|
||||
@@ -166,7 +166,7 @@ public final class Rendezvous {
|
||||
}
|
||||
|
||||
/**
|
||||
* Surface a worker's mid-turn {@code question} by resolving the primary's open {@code bridge_send}
|
||||
* Surface a worker's mid-turn {@code question} by resolving the primary's open {@code fleet_send}
|
||||
* with a {@link Kind#QUESTION} carrying {@code turnId}. Same session-keyed semantics as
|
||||
* {@link #resolve}: the one outstanding send for {@code session} unblocks with the question.
|
||||
*
|
||||
@@ -184,7 +184,7 @@ public final class Rendezvous {
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve a worker's blocked {@code bridge_ask} with the primary's {@code answer}, unblocking it
|
||||
* Resolve a worker's blocked {@code fleet_ask} with the primary's {@code answer}, unblocking it
|
||||
* to resume its turn.
|
||||
*
|
||||
* @return {@code true} if the ask was still open and got the answer; {@code false} if the
|
||||
@@ -195,7 +195,7 @@ public final class Rendezvous {
|
||||
return w != null && w.answer().complete(answer);
|
||||
}
|
||||
|
||||
/** Drop a reverse-rendezvous turn once its {@code bridge_ask} has resolved (answered or lapsed). */
|
||||
/** Drop a reverse-rendezvous turn once its {@code fleet_ask} has resolved (answered or lapsed). */
|
||||
public void closeAsk(String turnId) {
|
||||
AskWaiter w = asks.get(turnId);
|
||||
if (w == null) {
|
||||
@@ -209,9 +209,9 @@ public final class Rendezvous {
|
||||
|
||||
/**
|
||||
* Resolve a specific captured {@code waiter} as a completion (the delegated turn finished with no
|
||||
* {@code bridge_reply}); {@code text} is the scraped transcript tail. The waiter is the one
|
||||
* {@code fleet_reply}); {@code text} is the scraped transcript tail. The waiter is the one
|
||||
* captured when this turn was delivered, so a late completion for turn N cannot land on turn N+1's
|
||||
* send (CB-116). A no-op if that waiter was already resolved — a raced {@code bridge_reply} wins.
|
||||
* send (CB-116). A no-op if that waiter was already resolved — a raced {@code fleet_reply} wins.
|
||||
*
|
||||
* @return {@code true} if this call resolved the waiter, {@code false} if it was null or already resolved
|
||||
*/
|
||||
@@ -233,7 +233,7 @@ public final class Rendezvous {
|
||||
|
||||
/**
|
||||
* Resolve a specific captured {@code waiter} as {@link Kind#BACKEND_EXHAUSTED} (CB-578 stage A):
|
||||
* the turn finished with no {@code bridge_reply} and the scrape matched the backend's configured
|
||||
* the turn finished with no {@code fleet_reply} and the scrape matched the backend's configured
|
||||
* usage-limit pattern; {@code reason} carries the matched line. Like
|
||||
* {@link #resolveCompletion(CompletableFuture, String)} it targets the exact captured send
|
||||
* (CB-116). A no-op if that waiter was already resolved — first resolution wins.
|
||||
|
||||
@@ -19,9 +19,9 @@ import java.util.stream.Collectors;
|
||||
|
||||
/**
|
||||
* A status-gated push loop that nudges a lead's own herdr pane when it has uncollected work
|
||||
* waiting: a worker reply queued with no live {@code bridge_send} to resolve it (CB-307), an
|
||||
* async delegation ticket ({@code bridge_send(wait:false)}) that reached a terminal phase
|
||||
* (CB-588), or an async ticket's worker pausing mid-turn in {@code bridge_ask} to await an answer
|
||||
* waiting: a worker reply queued with no live {@code fleet_send} to resolve it (CB-307), an
|
||||
* async delegation ticket ({@code fleet_send(wait:false)}) that reached a terminal phase
|
||||
* (CB-588), or an async ticket's worker pausing mid-turn in {@code fleet_ask} to await an answer
|
||||
* (CB-582).
|
||||
*
|
||||
* <p><strong>CB-590: one schedule per lead.</strong> All three kinds of work are triggered
|
||||
@@ -50,24 +50,24 @@ import java.util.stream.Collectors;
|
||||
public final class ReplyPushLoop {
|
||||
|
||||
private static final Logger log = LoggerFactory.getLogger(ReplyPushLoop.class);
|
||||
static final String NUDGE_FORMAT = "Worker %s returned a reply — run bridge_poll(target=%s) to collect it";
|
||||
static final String NUDGE_FORMAT = "Worker %s returned a reply — run fleet_poll(target=%s) to collect it";
|
||||
/** Coalesced form, several uncollected replies for the same lead. */
|
||||
static final String REPLIES_NUDGE_FORMAT =
|
||||
"%d workers returned replies — run bridge_poll(target=...) for each to collect them: %s";
|
||||
"%d workers returned replies — run fleet_poll(target=...) for each to collect them: %s";
|
||||
/** Singular form, one uncollected ticket. */
|
||||
static final String TICKET_NUDGE_FORMAT =
|
||||
"Ticket %s finished%s — run bridge_poll(ticket=%s) to collect it";
|
||||
"Ticket %s finished%s — run fleet_poll(ticket=%s) to collect it";
|
||||
/** Coalesced form, several uncollected tickets for the same lead. */
|
||||
static final String TICKETS_NUDGE_FORMAT =
|
||||
"%d tickets finished%s — run bridge_poll(ticket=...) for each to collect them: %s";
|
||||
/** Singular form, one worker paused mid-turn in bridge_ask (CB-582) — names the answer call directly. */
|
||||
"%d tickets finished%s — run fleet_poll(ticket=...) for each to collect them: %s";
|
||||
/** Singular form, one worker paused mid-turn in fleet_ask (CB-582) — names the answer call directly. */
|
||||
static final String QUESTION_NUDGE_FORMAT =
|
||||
"Worker %s asked a question (ticket %s) — answer it with bridge_send(turnId=\"%s\", "
|
||||
"Worker %s asked a question (ticket %s) — answer it with fleet_send(turnId=\"%s\", "
|
||||
+ "content=...) to resume its turn:\n%s";
|
||||
/** Coalesced form, several open questions for the same lead. */
|
||||
static final String QUESTIONS_NUDGE_FORMAT =
|
||||
"%d workers are paused on a question — run bridge_poll(ticket=...) for each, then answer "
|
||||
+ "with bridge_send(turnId=..., content=...): %s";
|
||||
"%d workers are paused on a question — run fleet_poll(ticket=...) for each, then answer "
|
||||
+ "with fleet_send(turnId=..., content=...): %s";
|
||||
|
||||
private final PrimaryRegistry primaryRegistry;
|
||||
private final AgentControl agents;
|
||||
@@ -87,7 +87,7 @@ public final class ReplyPushLoop {
|
||||
/** Tickets that have gone terminal but not yet been polled, keyed by ticket. */
|
||||
private final ConcurrentHashMap<String, PendingTicket> pendingTickets = new ConcurrentHashMap<>();
|
||||
/**
|
||||
* Open {@code bridge_ask} questions not yet answered or lapsed, keyed by {@code turnId}
|
||||
* Open {@code fleet_ask} questions not yet answered or lapsed, keyed by {@code turnId}
|
||||
* (CB-582). A question's own nudge count is tracked the same per-item way as
|
||||
* {@link #pendingTickets} (CB-598): a fresh question keeps its source eligible regardless of
|
||||
* how depleted an older, still-open question's count is.
|
||||
@@ -130,7 +130,7 @@ public final class ReplyPushLoop {
|
||||
/**
|
||||
* Reply targets still pending for {@code lead} — registered via {@link #onReplyQueued} and
|
||||
* whose inbox still holds an unacked message. A target whose inbox has since drained (acked,
|
||||
* or collected via a live {@code bridge_send} rendezvous instead) is dropped from
|
||||
* or collected via a live {@code fleet_send} rendezvous instead) is dropped from
|
||||
* {@link #pendingReplies} here rather than lingering forever; there is no explicit "reply
|
||||
* collected" callback the way {@link #ticketCollected} exists for tickets, so the inbox itself
|
||||
* is the only signal.
|
||||
@@ -323,7 +323,7 @@ public final class ReplyPushLoop {
|
||||
}
|
||||
|
||||
/**
|
||||
* Called when an async delegation ticket ({@code bridge_send(wait:false)}, CB-107) reaches a
|
||||
* Called when an async delegation ticket ({@code fleet_send(wait:false)}, CB-107) reaches a
|
||||
* terminal phase — DONE or a failure. Unlike {@link #onReplyQueued}, which nudges about the
|
||||
* durable-inbox no-waiter path, this covers the path {@code MessageService.reply} takes when a
|
||||
* fire-and-poll send's own rendezvous waiter resolves the reply directly: that path returns
|
||||
@@ -352,7 +352,7 @@ public final class ReplyPushLoop {
|
||||
}
|
||||
|
||||
/**
|
||||
* Called when a ticket's terminal state has been collected via {@code bridge_poll}. Removes it
|
||||
* Called when a ticket's terminal state has been collected via {@code fleet_poll}. Removes it
|
||||
* from the pending set so a scheduled tick — and any nudge it sends — never names a ticket the
|
||||
* lead already has. A ticket that was never pending (unknown ticket, or one nudged with no push
|
||||
* loop configured) is a no-op.
|
||||
@@ -362,16 +362,16 @@ public final class ReplyPushLoop {
|
||||
}
|
||||
|
||||
/**
|
||||
* Called when an async ticket's worker pauses mid-turn in {@code bridge_ask} (CB-582): the
|
||||
* question is now visible via {@code bridge_poll} (Phase.ASKING), but the reverse-rendezvous
|
||||
* Called when an async ticket's worker pauses mid-turn in {@code fleet_ask} (CB-582): the
|
||||
* question is now visible via {@code fleet_poll} (Phase.ASKING), but the reverse-rendezvous
|
||||
* window it opened with (~55s default, see {@code BridgeMcp}/{@code BridgedApp}) is far shorter
|
||||
* than a lead's normal minutes-long poll cadence — exactly the gap this closes. Resolves the
|
||||
* delegating lead the same way {@link #onTicketTerminal} does and coalesces onto the same
|
||||
* per-lead schedule (CB-590).
|
||||
*
|
||||
* @param ticket the async ticket the question belongs to (for {@code bridge_poll})
|
||||
* @param ticket the async ticket the question belongs to (for {@code fleet_poll})
|
||||
* @param target the worker session that asked
|
||||
* @param turnId correlation id the lead answers with ({@code bridge_send turnId=...})
|
||||
* @param turnId correlation id the lead answers with ({@code fleet_send turnId=...})
|
||||
* @param question the question text
|
||||
*/
|
||||
public void onQuestionOpened(String ticket, String target, String turnId, String question) {
|
||||
@@ -386,7 +386,7 @@ public final class ReplyPushLoop {
|
||||
}
|
||||
|
||||
/**
|
||||
* Called when a worker's {@code bridge_ask} resolves — answered or lapsed unanswered — so a
|
||||
* Called when a worker's {@code fleet_ask} resolves — answered or lapsed unanswered — so a
|
||||
* scheduled tick never nudges about a question the lead already handled. A {@code turnId} that
|
||||
* was never pending (never nudged, or already closed) is a no-op.
|
||||
*/
|
||||
|
||||
@@ -9,7 +9,7 @@ package dev.ltms.bridged.peer;
|
||||
public enum Capability {
|
||||
|
||||
/**
|
||||
* The peer supports {@code bridge_ask} rendezvous — pausing its delegated turn to ask
|
||||
* The peer supports {@code fleet_ask} rendezvous — pausing its delegated turn to ask
|
||||
* the primary a question, then resuming once answered. All Claude Code peers support this.
|
||||
*/
|
||||
MID_TURN_ASK,
|
||||
|
||||
@@ -11,7 +11,7 @@ package dev.ltms.bridged.placement;
|
||||
* a usage limit, not a transient capacity or reachability concern.
|
||||
* <li>Weight 0 (CB-554): {@code fixed} is still automatic selection, so a profile the operator
|
||||
* marked "never auto-select me" ({@code weight <= 0}) must be skipped here exactly as
|
||||
* {@code weighted}/{@code round-robin} skip it — an explicit {@code bridge_spawn} naming
|
||||
* {@code weighted}/{@code round-robin} skip it — an explicit {@code fleet_spawn} naming
|
||||
* the profile is unaffected, only this automatic fallback walk.
|
||||
* </ul>
|
||||
* A fleet where nothing is ever quarantined or weight-0 never exercises either path, so today's
|
||||
|
||||
@@ -26,7 +26,7 @@ public record PlacementCandidate(String profile, String host, float weight, Inte
|
||||
* not need to distinguish "explicit 0" from "absent" itself.
|
||||
*
|
||||
* <p>Exclusion is about <em>automatic</em> selection only — an explicit
|
||||
* {@code bridge_spawn{profile:"..."}} bypasses placement entirely and is unaffected.
|
||||
* {@code fleet_spawn{profile:"..."}} bypasses placement entirely and is unaffected.
|
||||
*/
|
||||
public boolean excluded() {
|
||||
return weight <= 0.0f;
|
||||
|
||||
@@ -46,7 +46,7 @@ public final class BridgedApp {
|
||||
/** Default blocking window for a message; kept under typical HTTP idle timeouts. */
|
||||
private static final long DEFAULT_MESSAGE_TIMEOUT_MS = 25_000;
|
||||
private static final long MAX_MESSAGE_TIMEOUT_MS = 120_000;
|
||||
/** Blocking window for a worker's bridge_ask (CB-205); the worker's MCP client caps its own call. */
|
||||
/** Blocking window for a worker's fleet_ask (CB-205); the worker's MCP client caps its own call. */
|
||||
private static final long DEFAULT_ASK_TIMEOUT_MS = 55_000;
|
||||
private static final long MAX_ASK_TIMEOUT_MS = 115_000;
|
||||
|
||||
@@ -121,11 +121,11 @@ public final class BridgedApp {
|
||||
app.get("/profiles", this::profiles); // configured backend profiles
|
||||
app.post("/members", this::spawnMember); // optional ?role=&profile= or {"role":…,"profile":…}
|
||||
app.delete("/members/{paneId}", this::stopMember);
|
||||
app.post("/sessions/{id}/message", this::sendMessage); // bridge_send (primary; blocking, wait:false, or answer via turnId)
|
||||
app.post("/sessions/{id}/reply", this::replyMessage); // bridge_reply (worker)
|
||||
app.post("/sessions/{id}/message", this::sendMessage); // fleet_send (primary; blocking, wait:false, or answer via turnId)
|
||||
app.post("/sessions/{id}/reply", this::replyMessage); // fleet_reply (worker)
|
||||
app.get("/sessions/{id}/replies", this::drainReplies); // drain reply inbox (CB-307)
|
||||
app.post("/sessions/{id}/ask", this::askMessage); // bridge_ask (worker → primary, CB-205)
|
||||
app.get("/sessions/{id}/status", this::sessionStatus); // bridge_status
|
||||
app.post("/sessions/{id}/ask", this::askMessage); // fleet_ask (worker → primary, CB-205)
|
||||
app.get("/sessions/{id}/status", this::sessionStatus); // fleet_status
|
||||
app.get("/tasks/{ticket}", this::taskStatus); // poll an async (wait:false) send
|
||||
return app;
|
||||
}
|
||||
@@ -341,7 +341,7 @@ public final class BridgedApp {
|
||||
|
||||
/**
|
||||
* The blocking delegation call (CB-104): inject {@code content} into the worker via the
|
||||
* status-gated injector and block until the worker returns a structured {@code bridge_reply}.
|
||||
* status-gated injector and block until the worker returns a structured {@code fleet_reply}.
|
||||
* Times out with a typed 202 (working / queued / busy) rather than an error — the message may
|
||||
* still land.
|
||||
*/
|
||||
@@ -370,7 +370,7 @@ public final class BridgedApp {
|
||||
}
|
||||
timeout = Math.clamp(timeout, 1, MAX_MESSAGE_TIMEOUT_MS);
|
||||
|
||||
// Answering a worker's bridge_ask (CB-205): always blocks, and derives the worker from turnId.
|
||||
// Answering a worker's fleet_ask (CB-205): always blocks, and derives the worker from turnId.
|
||||
if (turnId != null && !turnId.isBlank()) {
|
||||
writeReply(ctx, id, messages.answer(turnId, content, timeout), timeout);
|
||||
return;
|
||||
@@ -392,7 +392,7 @@ public final class BridgedApp {
|
||||
|
||||
/**
|
||||
* Render a {@link MessageService.Reply} onto the response — shared by a normal send and a
|
||||
* bridge_ask answer. A structured/scraped completion is 200; a worker's mid-turn question a 202
|
||||
* fleet_ask answer. A structured/scraped completion is 200; a worker's mid-turn question a 202
|
||||
* (with its {@code turnId}); a stale answer a 409; every other non-terminal outcome a typed 202.
|
||||
*/
|
||||
private void writeReply(Context ctx, String id, MessageService.Reply reply, long timeout) {
|
||||
@@ -404,7 +404,7 @@ public final class BridgedApp {
|
||||
"sessionId", id, "error", "stale_turn",
|
||||
"detail", "that question is no longer open (timed out or already answered)"));
|
||||
case REPLIED, COMPLETED_UNREPLIED -> {
|
||||
// replySource distinguishes a structured bridge_reply from the CB-106 completion
|
||||
// replySource distinguishes a structured fleet_reply from the CB-106 completion
|
||||
// fallback (a scrape of the worker's transcript when it finished without replying).
|
||||
String source = reply.outcome() == MessageService.Outcome.REPLIED ? "reply" : "transcript";
|
||||
ctx.status(200).json(Map.of("sessionId", id, "reply", reply.text(), "replySource", source));
|
||||
@@ -428,7 +428,7 @@ public final class BridgedApp {
|
||||
}
|
||||
|
||||
/**
|
||||
* A worker's mid-turn question ({@code bridge_ask}, CB-205) — surfaces to the primary's open
|
||||
* A worker's mid-turn question ({@code fleet_ask}, CB-205) — surfaces to the primary's open
|
||||
* blocking send and blocks until it answers. 200 with the answer, 409 if no delegation is open,
|
||||
* 202 if the primary stayed silent.
|
||||
*/
|
||||
@@ -465,7 +465,7 @@ public final class BridgedApp {
|
||||
}
|
||||
|
||||
/**
|
||||
* The worker's structured reply ({@code bridge_reply}) — resolves the blocking send awaiting
|
||||
* The worker's structured reply ({@code fleet_reply}) — resolves the blocking send awaiting
|
||||
* on this session, or queues the reply in the inbox when no send is open (CB-307).
|
||||
*/
|
||||
private void replyMessage(Context ctx) {
|
||||
@@ -505,7 +505,7 @@ public final class BridgedApp {
|
||||
}
|
||||
|
||||
/**
|
||||
* Live lifecycle status of a worker (MCP `bridge_status` wraps this in CB-105), plus its
|
||||
* Live lifecycle status of a worker (MCP `fleet_status` wraps this in CB-105), plus its
|
||||
* <em>readiness</em> (CB-113): {@code ready} is true once the worker's Claude has connected the
|
||||
* bridge MCP — the reliable "available to receive a task" signal, unlike bare {@code idle}, which
|
||||
* is also true during boot.
|
||||
@@ -520,8 +520,8 @@ public final class BridgedApp {
|
||||
body.put("sessionId", id);
|
||||
body.put("status", messages.status(id).name().toLowerCase());
|
||||
body.put("ready", presence.isPresent(id));
|
||||
// CB-582: a worker paused mid-turn in an async bridge_ask is otherwise invisible to a
|
||||
// status poll — surface the open question and how to answer it, same as bridge_poll's
|
||||
// CB-582: a worker paused mid-turn in an async fleet_ask is otherwise invisible to a
|
||||
// status poll — surface the open question and how to answer it, same as fleet_poll's
|
||||
// Phase.ASKING view.
|
||||
MessageService.PendingAsk ask = messages.pendingAsk(id);
|
||||
if (ask != null) {
|
||||
|
||||
@@ -55,7 +55,7 @@ class CallerResolverTest {
|
||||
assertEquals("term_a", underTrust.terminal());
|
||||
assertEquals(Role.WORKER, underToken.role(),
|
||||
"worker identity is unforgeable and must never be token-gated — otherwise enabling "
|
||||
+ "auth would lock the whole fleet out of bridge_reply");
|
||||
+ "auth would lock the whole fleet out of fleet_reply");
|
||||
assertEquals("term_a", underToken.terminal());
|
||||
}
|
||||
|
||||
|
||||
@@ -177,7 +177,7 @@ class CompletionResolverTest {
|
||||
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
|
||||
|
||||
assertEquals("x".repeat(CompletionResolver.MAX_SCRAPE_CHARS)
|
||||
+ "\n[Pane tail clipped: member did not call bridge_reply.]",
|
||||
+ "\n[Pane tail clipped: member did not call fleet_reply.]",
|
||||
waiter.getNow(null).text());
|
||||
}
|
||||
|
||||
@@ -347,7 +347,7 @@ class CompletionResolverTest {
|
||||
@Test
|
||||
void aLateCompletionForOneTurnNeverResolvesTheNextTurnsWaiter() {
|
||||
// The cross-turn stale reply the conversation test surfaced: turn N's completion fallback
|
||||
// fires AFTER turn N was resolved by an explicit bridge_reply and turn N+1 has opened its own
|
||||
// fires AFTER turn N was resolved by an explicit fleet_reply and turn N+1 has opened its own
|
||||
// waiter on the same session. Resolving "whatever is waiting now" would hand turn N's stale
|
||||
// scrape to turn N+1; targeting turn N's captured waiter makes the late completion a no-op.
|
||||
FakeHerdr herdr = new FakeHerdr().readText("⏺ turn N answer\n❯ ");
|
||||
|
||||
@@ -163,7 +163,7 @@ class LeadLauncherTest {
|
||||
|
||||
/**
|
||||
* The single most important assertion here. The worker charter tells its reader it is an
|
||||
* off-subscription worker that must end every turn with bridge_reply — the opposite of what an
|
||||
* off-subscription worker that must end every turn with fleet_reply — the opposite of what an
|
||||
* orchestrator is. A lead must never receive it.
|
||||
*/
|
||||
@Test
|
||||
@@ -174,7 +174,7 @@ class LeadLauncherTest {
|
||||
List<String> args = startedArgs(herdr);
|
||||
assertFalse(args.contains("--append-system-prompt"),
|
||||
"the reply charter is a worker contract and must not be injected into a lead");
|
||||
assertTrue(args.stream().noneMatch(a -> a.contains("bridge_reply")), args.toString());
|
||||
assertTrue(args.stream().noneMatch(a -> a.contains("fleet_reply")), args.toString());
|
||||
}
|
||||
|
||||
/** It still mounts the bridge — a lead that cannot orchestrate is pointless. */
|
||||
|
||||
@@ -19,16 +19,24 @@ import dev.ltms.bridged.member.ClaudeCodeLauncher;
|
||||
import dev.ltms.bridged.member.CompositePeerLauncher;
|
||||
import dev.ltms.bridged.placement.BackendQuarantine;
|
||||
import dev.ltms.bridged.placement.PlacementPolicies;
|
||||
import io.modelcontextprotocol.server.McpSyncServerExchange;
|
||||
import io.modelcontextprotocol.spec.McpSchema;
|
||||
import dev.ltms.bridged.msg.InMemoryReplyInbox;
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.Logger;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import org.junit.jupiter.api.BeforeEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.function.BiFunction;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
@@ -83,7 +91,7 @@ class BridgeMcpTest {
|
||||
|
||||
@Test
|
||||
void sendThenReplyRoundTrips() throws Exception {
|
||||
// bridge_send blocks; bridge_reply resolves it with the worker's structured answer.
|
||||
// fleet_send blocks; fleet_reply resolves it with the worker's structured answer.
|
||||
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
|
||||
() -> BridgeMcp.send(messages, "term_a", "review this", 4000L, null, Set.of()));
|
||||
|
||||
@@ -106,7 +114,7 @@ class BridgeMcpTest {
|
||||
|
||||
@Test
|
||||
void asyncSendReturnsATicketThenPollReportsTheReply() throws Exception {
|
||||
// wait:false parity — a ticket is issued, resolved by a reply, and surfaced by bridge_poll.
|
||||
// wait:false parity — a ticket is issued, resolved by a reply, and surfaced by fleet_poll.
|
||||
McpSchema.CallToolResult accepted = BridgeMcp.sendAsync(messages, "term_a", "do it", null, Set.of());
|
||||
assertNotEquals(Boolean.TRUE, accepted.isError());
|
||||
String out = textOf(accepted);
|
||||
@@ -290,7 +298,7 @@ class BridgeMcpTest {
|
||||
assertTrue(async.isError());
|
||||
assertTrue(textOf(blocking).contains("sol"));
|
||||
assertTrue(textOf(blocking).contains("configured profile name"));
|
||||
assertTrue(textOf(blocking).contains("bridge_list"));
|
||||
assertTrue(textOf(blocking).contains("fleet_list"));
|
||||
assertFalse(textOf(async).contains("ticket="));
|
||||
}
|
||||
|
||||
@@ -324,7 +332,7 @@ class BridgeMcpTest {
|
||||
// A reply with no open send queues it in the inbox.
|
||||
BridgeMcp.reply(messages, "term_a", "queued-msg");
|
||||
|
||||
// bridge_poll with target drains the inbox.
|
||||
// fleet_poll with target drains the inbox.
|
||||
McpSchema.CallToolResult res = BridgeMcp.poll(messages, null, "term_a");
|
||||
assertNotEquals(Boolean.TRUE, res.isError());
|
||||
String text = textOf(res);
|
||||
@@ -359,7 +367,7 @@ class BridgeMcpTest {
|
||||
String afterMarker = qt.substring(qt.indexOf("turnId=\"") + "turnId=\"".length());
|
||||
String turnId = afterMarker.substring(0, afterMarker.indexOf('"'));
|
||||
|
||||
// The primary answers via bridge_send(turnId); this blocks again for the worker's reply.
|
||||
// The primary answers via fleet_send(turnId); this blocks again for the worker's reply.
|
||||
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
|
||||
() -> BridgeMcp.answer(messages, turnId, "config.yaml", 5000L));
|
||||
|
||||
@@ -655,7 +663,7 @@ class BridgeMcpTest {
|
||||
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.
|
||||
// from self without a second fleet_whoami call.
|
||||
assertEquals(1, out.split("\"self\":true", -1).length - 1, out);
|
||||
assertTrue(out.indexOf("term_me") < out.indexOf("\"self\":true"), out);
|
||||
}
|
||||
@@ -730,7 +738,7 @@ class BridgeMcpTest {
|
||||
assertEquals(1, before.size(), "one reply in the inbox");
|
||||
String msgId = before.getFirst().msgId();
|
||||
|
||||
// Publish the same reply again and ack it via bridge_ack surface.
|
||||
// Publish the same reply again and ack it via fleet_ack surface.
|
||||
BridgeMcp.reply(messages, "term_a", "orphan-again");
|
||||
var peeked = messages.drainReplies("term_a");
|
||||
assertEquals(1, peeked.size(), "one fresh reply in the inbox");
|
||||
@@ -774,8 +782,8 @@ class BridgeMcpTest {
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-582: a lead polling {@code bridge_status} on its normal cadence — not {@code bridge_poll}
|
||||
* — must also see a worker's open async {@code bridge_ask} question, since the reverse-rendezvous
|
||||
* CB-582: a lead polling {@code fleet_status} on its normal cadence — not {@code fleet_poll}
|
||||
* — must also see a worker's open async {@code fleet_ask} question, since the reverse-rendezvous
|
||||
* window it opened with is far shorter than that cadence.
|
||||
*/
|
||||
@Test
|
||||
@@ -819,7 +827,7 @@ class BridgeMcpTest {
|
||||
answer.get(5, TimeUnit.SECONDS);
|
||||
}
|
||||
|
||||
// --- bridge_whoami: the caller's own identity, so an agent never has to guess its role -------
|
||||
// --- fleet_whoami: the caller's own identity, so an agent never has to guess its role -------
|
||||
|
||||
@Test
|
||||
void whoamiReportsThePrimaryAsPrimaryAndNothingElse() {
|
||||
@@ -994,7 +1002,7 @@ class BridgeMcpTest {
|
||||
assertTrue(textOf(res).contains("architect, dev, reviewer"), textOf(res));
|
||||
}
|
||||
|
||||
// ── CB-584: bridge_spawn accepts sessionName/resumeSessionId; roster shows agentSessionId ──
|
||||
// ── CB-584: fleet_spawn accepts sessionName/resumeSessionId; roster shows agentSessionId ──
|
||||
|
||||
@Test
|
||||
void spawnWithResumeSessionIdPutsTheIdOnTheRoster() {
|
||||
@@ -1020,4 +1028,103 @@ class BridgeMcpTest {
|
||||
assertEquals(Boolean.TRUE, res.isError());
|
||||
assertTrue(textOf(res).contains("explicit profile"), textOf(res));
|
||||
}
|
||||
|
||||
// ── CB-622: bridge_* -> fleet_* rename, both names answer through the SAME handler ────────
|
||||
//
|
||||
// These are unit tests of the wiring helpers (deprecatedTwin/deprecatedHandler/
|
||||
// warnDeprecatedOnce), not a live MCP client call — proving the tool is REGISTERED and
|
||||
// ROUTES correctly. Whether a real MCP client can actually invoke a tool by either name is
|
||||
// the live check the lead runs; see CB-618's lesson that a green build here is not proof of
|
||||
// that.
|
||||
|
||||
private static McpSchema.Tool fakeTool(String name) {
|
||||
return McpSchema.Tool.builder(name)
|
||||
.description("does the thing")
|
||||
.inputSchema(Map.of("type", "object", "properties", Map.of(), "required", List.of()))
|
||||
.build();
|
||||
}
|
||||
|
||||
@Test
|
||||
void deprecatedTwinNamesTheOldToolAndDefersToTheFleetDescription() {
|
||||
McpSchema.Tool fleetTool = fakeTool("fleet_cb622_twin");
|
||||
McpSchema.Tool twin = BridgeMcp.deprecatedTwin(fleetTool, "bridge_cb622_twin");
|
||||
|
||||
assertEquals("bridge_cb622_twin", twin.name());
|
||||
assertTrue(twin.description().startsWith("DEPRECATED: use fleet_cb622_twin instead."),
|
||||
twin.description());
|
||||
assertTrue(twin.description().contains(fleetTool.description()), twin.description());
|
||||
assertEquals(fleetTool.inputSchema(), twin.inputSchema(), "the twin must not restate the schema");
|
||||
}
|
||||
|
||||
@Test
|
||||
void oldNameAndNewNameReachTheExactSameHandler() {
|
||||
McpSchema.Tool fleetTool = fakeTool("fleet_cb622_dual");
|
||||
List<String> handlerCalls = new ArrayList<>();
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> handler =
|
||||
(exchange, req) -> {
|
||||
handlerCalls.add(req.name());
|
||||
return McpSchema.CallToolResult.builder().addTextContent("handled:" + req.name()).build();
|
||||
};
|
||||
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> deprecated =
|
||||
BridgeMcp.deprecatedHandler(fleetTool, "bridge_cb622_dual", handler);
|
||||
|
||||
// Calling the fleet_* name directly and calling the wrapped bridge_* name both end up
|
||||
// running the SAME `handler` instance — not a copy of its logic.
|
||||
McpSchema.CallToolResult viaFleet = handler.apply(null, new McpSchema.CallToolRequest("fleet_cb622_dual", Map.of()));
|
||||
McpSchema.CallToolResult viaBridge = deprecated.apply(null, new McpSchema.CallToolRequest("bridge_cb622_dual", Map.of()));
|
||||
|
||||
assertEquals(textOf(viaFleet), "handled:fleet_cb622_dual");
|
||||
assertEquals(textOf(viaBridge), "handled:bridge_cb622_dual");
|
||||
assertEquals(2, handlerCalls.size(), "the same handler ran for both calls");
|
||||
}
|
||||
|
||||
@Test
|
||||
void deprecatedNameWarnsOnceForTheProcessNotOncePerCall() {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(BridgeMcp.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
// A name unique to this test run so another test's use of the mechanism (or a re-run in
|
||||
// the same JVM) cannot leave the "already warned" flag set before this assertion.
|
||||
String oldName = "bridge_cb622_warnonce_" + System.identityHashCode(appender);
|
||||
try {
|
||||
BridgeMcp.warnDeprecatedOnce(oldName, "fleet_cb622_warnonce");
|
||||
BridgeMcp.warnDeprecatedOnce(oldName, "fleet_cb622_warnonce");
|
||||
BridgeMcp.warnDeprecatedOnce(oldName, "fleet_cb622_warnonce");
|
||||
|
||||
List<ILoggingEvent> matching = appender.list.stream()
|
||||
.filter(e -> e.getFormattedMessage().contains(oldName))
|
||||
.toList();
|
||||
assertEquals(1, matching.size(), "three calls with the same old name must warn exactly once");
|
||||
assertEquals(Level.WARN, matching.get(0).getLevel());
|
||||
assertTrue(matching.get(0).getFormattedMessage().contains("fleet_cb622_warnonce"),
|
||||
"the warning must name the new tool too");
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void deprecatedNameWarnsAgainForADifferentOldName() {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(BridgeMcp.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
String suffix = System.identityHashCode(appender) + "";
|
||||
String nameA = "bridge_cb622_multi_a_" + suffix;
|
||||
String nameB = "bridge_cb622_multi_b_" + suffix;
|
||||
try {
|
||||
BridgeMcp.warnDeprecatedOnce(nameA, "fleet_cb622_multi_a");
|
||||
BridgeMcp.warnDeprecatedOnce(nameB, "fleet_cb622_multi_b");
|
||||
|
||||
long distinctNamesWarned = appender.list.stream()
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.filter(m -> m.contains(nameA) || m.contains(nameB))
|
||||
.count();
|
||||
assertEquals(2, distinctNamesWarned, "each distinct old name gets its own warning");
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -110,7 +110,7 @@ class PrimaryRegistryTest {
|
||||
|
||||
/**
|
||||
* 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.
|
||||
* fleet_send first captured every nudge — including nudges for the other lead's delegations.
|
||||
*/
|
||||
@Test
|
||||
void aNudgeGoesToTheLeadThatDelegatedToThatWorker() {
|
||||
|
||||
@@ -57,7 +57,7 @@ class ClaudeCodeLauncherTest {
|
||||
assertTrue(args.stream().anyMatch(a -> a.contains("\"bridge\"") && a.contains("http://127.0.0.1:8765/mcp")),
|
||||
"inline bridge MCP config present");
|
||||
assertTrue(args.contains("--append-system-prompt"));
|
||||
assertTrue(args.stream().anyMatch(a -> a.contains("bridge_reply")), "reply charter present");
|
||||
assertTrue(args.stream().anyMatch(a -> a.contains("fleet_reply")), "reply charter present");
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
@@ -358,7 +358,7 @@ class CompositePeerLauncherTest {
|
||||
@Test
|
||||
void explicitSpawnStillSucceedsOnWeightZeroProfile() {
|
||||
// CB-554: weight: 0 excludes a profile from AUTOMATIC selection only — an explicit
|
||||
// bridge_spawn{profile:"a"} must still work exactly as today (e.g. `opus` on the
|
||||
// fleet_spawn{profile:"a"} must still work exactly as today (e.g. `opus` on the
|
||||
// operator's own subscription, kept weight-0 so it is never picked automatically).
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Map<String, BridgedConfig.Profile> profiles = ordered(
|
||||
@@ -593,7 +593,7 @@ class CompositePeerLauncherTest {
|
||||
|
||||
/**
|
||||
* An explicit profile is the operator overriding and is NOT judged against the pool. It must
|
||||
* stay that way: an unrolled `bridge_spawn{profile:"opus"}` carries no role, so it defaults to
|
||||
* stay that way: an unrolled `fleet_spawn{profile:"opus"}` carries no role, so it defaults to
|
||||
* DEV, and enforcing the pool here would refuse a spawn the operator asked for by name.
|
||||
*/
|
||||
@Test
|
||||
|
||||
@@ -169,8 +169,8 @@ class LeadHeartbeatLoopTest {
|
||||
String t = pendingFleet().nudgeText();
|
||||
assertTrue(t.startsWith("Heartbeat:"), "the nudge identifies itself as a heartbeat");
|
||||
assertTrue(t.contains("2 worker replies pending collection"), t);
|
||||
assertTrue(t.contains("bridge_poll(target=term_a)"), t);
|
||||
assertTrue(t.contains("bridge_poll(target=term_b)"), t);
|
||||
assertTrue(t.contains("fleet_poll(target=term_a)"), t);
|
||||
assertTrue(t.contains("fleet_poll(target=term_b)"), t);
|
||||
assertTrue(t.contains("1 DONE session awaiting teardown"), t);
|
||||
assertTrue(t.contains("3 workers live"), t);
|
||||
}
|
||||
|
||||
@@ -68,7 +68,7 @@ class MessageServiceTest {
|
||||
injector.onStatus(T, AgentStatus.IDLE); // deliver the task (baselines the pre-turn content)
|
||||
injector.onStatus(T, AgentStatus.WORKING); // worker picks it up and works
|
||||
herdr.readText("BUILD GREEN: 391 files"); // the worker's turn produced new output
|
||||
injector.onStatus(T, AgentStatus.IDLE); // working → idle: turn complete, no bridge_reply
|
||||
injector.onStatus(T, AgentStatus.IDLE); // working → idle: turn complete, no fleet_reply
|
||||
|
||||
MessageService.Reply reply = send.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(MessageService.Outcome.COMPLETED_UNREPLIED, reply.outcome(),
|
||||
@@ -153,7 +153,7 @@ class MessageServiceTest {
|
||||
+ "(worker unreachable or stuck)", resolution.text());
|
||||
}
|
||||
|
||||
// --- bridge_ask reverse rendezvous (CB-205) ------------------------------------------------
|
||||
// --- fleet_ask reverse rendezvous (CB-205) ------------------------------------------------
|
||||
|
||||
@Test
|
||||
void askSurfacesAsAQuestionAndTheAnswerResumesTheSameTurn() throws Exception {
|
||||
@@ -172,7 +172,7 @@ class MessageServiceTest {
|
||||
assertEquals("which config file?", q.text());
|
||||
assertNotNull(q.turnId(), "a question carries a turnId to answer on");
|
||||
|
||||
// The primary answers via bridge_send(turnId); this blocks again for the worker's reply.
|
||||
// The primary answers via fleet_send(turnId); this blocks again for the worker's reply.
|
||||
CompletableFuture<MessageService.Reply> answer =
|
||||
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000));
|
||||
|
||||
@@ -196,7 +196,7 @@ class MessageServiceTest {
|
||||
injector.onStatus(T, AgentStatus.IDLE); // deliver
|
||||
injector.onStatus(T, AgentStatus.WORKING); // worker picks it up, then pauses to ask
|
||||
|
||||
// A transport retry: two concurrent bridge_ask calls from the same worker session.
|
||||
// A transport retry: two concurrent fleet_ask calls from the same worker session.
|
||||
CompletableFuture<MessageService.AskResult> ask1 =
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
|
||||
CompletableFuture<MessageService.AskResult> ask2 =
|
||||
@@ -297,7 +297,7 @@ class MessageServiceTest {
|
||||
assertNotNull(q.turnId());
|
||||
|
||||
// The primary answers, unblocking the worker; but the worker never sends the follow-up
|
||||
// bridge_reply, so the answering send rides out its short window as still-working.
|
||||
// fleet_reply, so the answering send rides out its short window as still-working.
|
||||
MessageService.Reply answer = messages.answer(q.turnId(), "config.yaml", 200);
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answer.outcome(),
|
||||
"an answered worker that never replies times out as still working");
|
||||
@@ -332,7 +332,7 @@ class MessageServiceTest {
|
||||
assertNotNull(view, "a resolved async send must become DONE");
|
||||
assertEquals(MessageService.Phase.DONE, view.phase());
|
||||
assertEquals("async result", view.reply(), "the completed ticket reports the reply");
|
||||
assertEquals("reply", view.replySource(), "a structured bridge_reply is sourced from 'reply'");
|
||||
assertEquals("reply", view.replySource(), "a structured fleet_reply is sourced from 'reply'");
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -363,7 +363,7 @@ class MessageServiceTest {
|
||||
|
||||
/**
|
||||
* The bug CB-548 fixes: L holds worker W, then architect A attempts W and times out BUSY. With
|
||||
* delegator ownership recorded at {@code bridge_send} <em>request</em> time, A's rejected call
|
||||
* delegator ownership recorded at {@code fleet_send} <em>request</em> time, A's rejected call
|
||||
* would overwrite L — and W's late no-waiter reply would be pushed to A, who never owned the
|
||||
* turn. The accepted-delivery hook must not fire for a BUSY send, so L stays the delegator.
|
||||
*/
|
||||
@@ -423,7 +423,7 @@ class MessageServiceTest {
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-548 requirement: answering an existing {@code bridge_ask} is the SAME delegation, so it must
|
||||
* CB-548 requirement: answering an existing {@code fleet_ask} is the SAME delegation, so it must
|
||||
* not rewrite ownership. L accepted the send (owned), the worker paused to ask, and L answers via
|
||||
* turnId — ownership stays L throughout; the answer path never touches the registry.
|
||||
*/
|
||||
@@ -533,12 +533,12 @@ class MessageServiceTest {
|
||||
|
||||
@Test
|
||||
void aQuestionIsNeverQueuedInTheInbox() {
|
||||
// No send is open — bridge_ask with no delegation returns NO_WAITER,
|
||||
// No send is open — fleet_ask with no delegation returns NO_WAITER,
|
||||
// and the question text MUST NOT appear in the reply inbox.
|
||||
// The inbox is only fed by MessageService.reply(), not by bridge_ask.
|
||||
// The inbox is only fed by MessageService.reply(), not by fleet_ask.
|
||||
MessageService.AskResult r = messages.ask(T, "anyone there?", 500);
|
||||
assertEquals(MessageService.AskOutcome.NO_WAITER, r.outcome(),
|
||||
"bridge_ask with no open delegation must return NO_WAITER, never queued");
|
||||
"fleet_ask with no open delegation must return NO_WAITER, never queued");
|
||||
|
||||
assertTrue(messages.drainReplies(T).isEmpty(), "questions must never be queued");
|
||||
}
|
||||
@@ -550,7 +550,7 @@ class MessageServiceTest {
|
||||
awaitUninterruptibly(T);
|
||||
injectDelivery();
|
||||
|
||||
// The worker never sends bridge_reply, but the turn completes.
|
||||
// The worker never sends fleet_reply, but the turn completes.
|
||||
herdr.readText("done-scraped");
|
||||
completion.onTurnComplete(T); // The fallback arms and resolves the captured waiter.
|
||||
|
||||
@@ -714,7 +714,7 @@ class MessageServiceTest {
|
||||
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
|
||||
}
|
||||
|
||||
// --- CB-582: bridge_status pendingAsk() ------------------------------------------------------
|
||||
// --- CB-582: fleet_status pendingAsk() ------------------------------------------------------
|
||||
|
||||
@Test
|
||||
void pendingAskReturnsNullWhenNoQuestionIsOpen() throws Exception {
|
||||
@@ -741,7 +741,7 @@ class MessageServiceTest {
|
||||
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
|
||||
|
||||
MessageService.PendingAsk pending = messages.pendingAsk(T);
|
||||
assertNotNull(pending, "bridge_status should see the open question");
|
||||
assertNotNull(pending, "fleet_status should see the open question");
|
||||
assertEquals(ticket, pending.ticket());
|
||||
assertEquals("which config file?", pending.question());
|
||||
assertEquals(asking.turnId(), pending.turnId());
|
||||
@@ -760,7 +760,7 @@ class MessageServiceTest {
|
||||
// MessageService.reply's rendezvous fast path is exactly what an async ticket always takes
|
||||
// (sendAsync registers a rendezvous waiter — see asyncTasksByWaiter), so it never reached
|
||||
// ReplyPushLoop.onReplyQueued. These prove the ticket reaches ReplyPushLoop through the new
|
||||
// onTicketTerminal entry point instead, with no bridge_poll from the lead first.
|
||||
// onTicketTerminal entry point instead, with no fleet_poll from the lead first.
|
||||
|
||||
private static final String LEAD = "term_lead";
|
||||
|
||||
@@ -808,7 +808,7 @@ class MessageServiceTest {
|
||||
awaitNudge(wiring.leadHerdr());
|
||||
String nudge = wiring.leadHerdr().lastCall("agent.prompt").params().toString();
|
||||
assertTrue(nudge.contains(ticket), "the nudge should name the ticket: " + nudge);
|
||||
assertTrue(nudge.contains("bridge_poll(ticket="),
|
||||
assertTrue(nudge.contains("fleet_poll(ticket="),
|
||||
"the nudge should name the exact ticket-collecting call: " + nudge);
|
||||
assertFalse(nudge.toUpperCase().contains("FAILED"),
|
||||
"a successfully-replied ticket's nudge must not say it failed: " + nudge);
|
||||
@@ -874,7 +874,7 @@ class MessageServiceTest {
|
||||
}
|
||||
}
|
||||
|
||||
// --- CB-582: bridge_ask question-open nudges --------------------------------------------------
|
||||
// --- CB-582: fleet_ask question-open nudges --------------------------------------------------
|
||||
|
||||
@Test
|
||||
void anAsyncTicketThatPausesOnAQuestionNudgesTheLeadWithNoPriorPollCall() throws Exception {
|
||||
@@ -891,7 +891,7 @@ class MessageServiceTest {
|
||||
String nudge = wiring.leadHerdr().lastCall("agent.prompt").params().toString();
|
||||
assertTrue(nudge.contains(ticket), "the nudge should name the ticket: " + nudge);
|
||||
assertTrue(nudge.contains(asking.turnId()), "the nudge should name the turnId: " + nudge);
|
||||
assertTrue(nudge.contains("bridge_send(turnId="),
|
||||
assertTrue(nudge.contains("fleet_send(turnId="),
|
||||
"the nudge should name the exact answer call: " + nudge);
|
||||
assertTrue(nudge.contains("which config file?"), "the nudge should include the question: " + nudge);
|
||||
|
||||
@@ -997,7 +997,7 @@ class MessageServiceTest {
|
||||
* {@code pruneTerminalTickets} drops entries from it once {@link MessageService#TICKET_TTL_NANOS}
|
||||
* elapses. Before this test, that prune never told {@code ReplyPushLoop} — its own
|
||||
* {@code pendingTickets} entry for a pruned, never-collected ticket had no remover at all, so it
|
||||
* rode along on every later nudge to the same lead, naming a ticket {@code bridge_poll} could no
|
||||
* rode along on every later nudge to the same lead, naming a ticket {@code fleet_poll} could no
|
||||
* longer find. Uses the injectable clock (mirroring {@code SessionManager}'s {@code nowNanos} seam
|
||||
* for its idle reaper) to cross the 10-minute TTL without a real wait.
|
||||
*/
|
||||
@@ -1040,10 +1040,10 @@ class MessageServiceTest {
|
||||
String latestNudge = wiring.leadHerdr().lastCall("agent.prompt").params().toString();
|
||||
assertTrue(latestNudge.contains(fresh), "the fresh ticket's nudge must still arrive: " + latestNudge);
|
||||
assertFalse(latestNudge.contains(stale),
|
||||
"a pruned ticket must never be named in a later nudge — it is gone and bridge_poll "
|
||||
"a pruned ticket must never be named in a later nudge — it is gone and fleet_poll "
|
||||
+ "on it would return nothing: " + latestNudge);
|
||||
|
||||
// And bridge_poll(ticket=stale) really does return nothing now — the nudge would have lied.
|
||||
// And fleet_poll(ticket=stale) really does return nothing now — the nudge would have lied.
|
||||
assertNull(wiring.service().poll(stale), "the pruned ticket must actually be gone, not just unmentioned");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -13,7 +13,7 @@ import static org.junit.jupiter.api.Assertions.assertThrows;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* The reverse rendezvous (CB-205): the {@code bridge_ask} registry that lets a worker pause mid-turn
|
||||
* The reverse rendezvous (CB-205): the {@code fleet_ask} registry that lets a worker pause mid-turn
|
||||
* to ask the primary. Unit-level — the message-layer round-trip is covered in {@link MessageServiceTest}.
|
||||
*/
|
||||
class RendezvousTest {
|
||||
|
||||
@@ -168,8 +168,8 @@ class ReplyPushLoopTest {
|
||||
// Exactly one nudge = exactly 1 agent.prompt call (it submits itself)
|
||||
assertEquals(1, rec.sendCount());
|
||||
assertTrue(rec.sentParams().stream()
|
||||
.anyMatch(e -> e.getValue().toString().contains("bridge_poll")),
|
||||
"nudge text should contain bridge_poll");
|
||||
.anyMatch(e -> e.getValue().toString().contains("fleet_poll")),
|
||||
"nudge text should contain fleet_poll");
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -226,14 +226,14 @@ class ReplyPushLoopTest {
|
||||
void nudgeFormatIsCorrect() {
|
||||
String nudge = ReplyPushLoop.NUDGE_FORMAT.formatted(WORKER, WORKER);
|
||||
assertTrue(nudge.contains("Worker term_worker"));
|
||||
assertTrue(nudge.contains("bridge_poll(target=term_worker)"));
|
||||
assertTrue(nudge.contains("fleet_poll(target=term_worker)"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void repliesNudgeFormatIsCorrect() {
|
||||
String multi = ReplyPushLoop.REPLIES_NUDGE_FORMAT.formatted(2, "term_worker1, term_worker2");
|
||||
assertTrue(multi.contains("2 workers"));
|
||||
assertTrue(multi.contains("bridge_poll(target=...)"));
|
||||
assertTrue(multi.contains("fleet_poll(target=...)"));
|
||||
}
|
||||
|
||||
// --- CB-588: async ticket terminal nudges — decide() logic on tickets -----------------------
|
||||
@@ -289,8 +289,8 @@ class ReplyPushLoopTest {
|
||||
assertEquals(1, rec.sendCount());
|
||||
String nudge = rec.sentParams().getFirst().getValue().toString();
|
||||
assertTrue(nudge.contains("task-1"), "nudge should name the ticket");
|
||||
assertTrue(nudge.contains("bridge_poll(ticket="), "nudge should name the exact ticket-poll call");
|
||||
assertFalse(nudge.contains("bridge_poll(target="), "a ticket-only nudge must not tell the lead to run the target-poll call");
|
||||
assertTrue(nudge.contains("fleet_poll(ticket="), "nudge should name the exact ticket-poll call");
|
||||
assertFalse(nudge.contains("fleet_poll(target="), "a ticket-only nudge must not tell the lead to run the target-poll call");
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -468,11 +468,11 @@ class ReplyPushLoopTest {
|
||||
void ticketNudgeFormatIsCorrect() {
|
||||
String single = ReplyPushLoop.TICKET_NUDGE_FORMAT.formatted("task-1", "", "task-1");
|
||||
assertTrue(single.contains("Ticket task-1"));
|
||||
assertTrue(single.contains("bridge_poll(ticket=task-1)"));
|
||||
assertTrue(single.contains("fleet_poll(ticket=task-1)"));
|
||||
|
||||
String multi = ReplyPushLoop.TICKETS_NUDGE_FORMAT.formatted(2, "", "task-1, task-2");
|
||||
assertTrue(multi.contains("2 tickets"));
|
||||
assertTrue(multi.contains("bridge_poll(ticket=...)"));
|
||||
assertTrue(multi.contains("fleet_poll(ticket=...)"));
|
||||
}
|
||||
|
||||
// --- CB-307 nudge path is unchanged (regression) --------------------------------------------
|
||||
@@ -480,7 +480,7 @@ class ReplyPushLoopTest {
|
||||
@Test
|
||||
void inboxNudgeStillUsesTheOriginalTargetPollCall() {
|
||||
String nudge = ReplyPushLoop.NUDGE_FORMAT.formatted(WORKER, WORKER);
|
||||
assertTrue(nudge.contains("bridge_poll(target=" + WORKER + ")"),
|
||||
assertTrue(nudge.contains("fleet_poll(target=" + WORKER + ")"),
|
||||
"CB-588/CB-590 must not change the CB-307 inbox nudge's call shape");
|
||||
}
|
||||
|
||||
@@ -502,7 +502,7 @@ class ReplyPushLoopTest {
|
||||
"a reply and a ticket for the same lead must coalesce onto ONE schedule — "
|
||||
+ "two nudge injections into the same lead pane must never overlap");
|
||||
String nudge = rec.sentParams().getFirst().getValue().toString();
|
||||
assertTrue(nudge.contains("bridge_poll(target=" + WORKER + ")"),
|
||||
assertTrue(nudge.contains("fleet_poll(target=" + WORKER + ")"),
|
||||
"the combined nudge must still mention the reply: " + nudge);
|
||||
assertTrue(nudge.contains("task-1"), "the combined nudge must still mention the ticket: " + nudge);
|
||||
}
|
||||
@@ -525,7 +525,7 @@ class ReplyPushLoopTest {
|
||||
Thread.sleep(200);
|
||||
assertEquals(1, rec.sendCount(), "exactly one nudge once injectable — reply and ticket coalesced");
|
||||
String nudge = rec.sentParams().getFirst().getValue().toString();
|
||||
assertTrue(nudge.contains("bridge_poll(target=" + WORKER + ")"), "the reply must not be dropped: " + nudge);
|
||||
assertTrue(nudge.contains("fleet_poll(target=" + WORKER + ")"), "the reply must not be dropped: " + nudge);
|
||||
assertTrue(nudge.contains("task-1"), "the ticket must not be dropped: " + nudge);
|
||||
}
|
||||
|
||||
@@ -657,7 +657,7 @@ class ReplyPushLoopTest {
|
||||
assertTrue(loop.isActive(), "the schedule must stay active after nudging the fresh ticket");
|
||||
}
|
||||
|
||||
// --- CB-582: bridge_ask question-open nudges -------------------------------------------------
|
||||
// --- CB-582: fleet_ask question-open nudges -------------------------------------------------
|
||||
|
||||
@Test
|
||||
void onQuestionOpenedWithNoKnownLeadNeverStartsASchedule() throws Exception {
|
||||
@@ -706,7 +706,7 @@ class ReplyPushLoopTest {
|
||||
String nudge = rec.sentParams().getFirst().getValue().toString();
|
||||
assertTrue(nudge.contains("task-1"), "nudge should name the ticket: " + nudge);
|
||||
assertTrue(nudge.contains("term_worker#1"), "nudge should name the turnId: " + nudge);
|
||||
assertTrue(nudge.contains("bridge_send(turnId="), "nudge should name the exact answer call: " + nudge);
|
||||
assertTrue(nudge.contains("fleet_send(turnId="), "nudge should name the exact answer call: " + nudge);
|
||||
assertTrue(nudge.contains("which config file?"), "nudge should include the question text: " + nudge);
|
||||
}
|
||||
|
||||
@@ -775,12 +775,12 @@ class ReplyPushLoopTest {
|
||||
String single = ReplyPushLoop.QUESTION_NUDGE_FORMAT.formatted(
|
||||
WORKER, "task-1", "term_worker#1", "which config?");
|
||||
assertTrue(single.contains("Worker term_worker"));
|
||||
assertTrue(single.contains("bridge_send(turnId=\"term_worker#1\""));
|
||||
assertTrue(single.contains("fleet_send(turnId=\"term_worker#1\""));
|
||||
assertTrue(single.contains("which config?"));
|
||||
|
||||
String multi = ReplyPushLoop.QUESTIONS_NUDGE_FORMAT.formatted(2, "task-1 (turnId=t1), task-2 (turnId=t2)");
|
||||
assertTrue(multi.contains("2 workers"));
|
||||
assertTrue(multi.contains("bridge_poll(ticket=...)"));
|
||||
assertTrue(multi.contains("fleet_poll(ticket=...)"));
|
||||
}
|
||||
|
||||
// --- metrics (CB-512) ----------------------------------------------------------------------
|
||||
|
||||
@@ -355,7 +355,7 @@ class BridgedAppTest {
|
||||
|
||||
@Test
|
||||
void messageReturnsTheWorkersStructuredReply() throws Exception {
|
||||
// CB-104 (option C): the blocking send resolves on the worker's bridge_reply, not a scrape.
|
||||
// CB-104 (option C): the blocking send resolves on the worker's fleet_reply, not a scrape.
|
||||
FakeHerdr herdr = new FakeHerdr().agentStatus("idle"); // poller delivers the injection
|
||||
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
|
||||
|
||||
@@ -486,7 +486,7 @@ class BridgedAppTest {
|
||||
|
||||
/**
|
||||
* CB-582: a lead polling {@code GET /sessions/{id}/status} on its normal cadence — not the
|
||||
* ticket-scoped {@code /tasks/{ticket}} — must also see a worker's open async {@code bridge_ask}
|
||||
* ticket-scoped {@code /tasks/{ticket}} — must also see a worker's open async {@code fleet_ask}
|
||||
* question, since the reverse-rendezvous window it opened with is far shorter than that cadence.
|
||||
*/
|
||||
@Test
|
||||
@@ -528,7 +528,7 @@ class BridgedAppTest {
|
||||
assertEquals(turnId, body.get("turnId").asText());
|
||||
assertEquals(ticket, body.get("ticket").asText());
|
||||
|
||||
// Answer it — via the same /message route bridge_send uses, keyed by turnId — so the
|
||||
// Answer it — via the same /message route fleet_send uses, keyed by turnId — so the
|
||||
// background ask thread does not linger past the test.
|
||||
var answer = java.util.concurrent.CompletableFuture.supplyAsync(() -> {
|
||||
try {
|
||||
|
||||
@@ -206,7 +206,7 @@ class SessionManagerTest {
|
||||
|
||||
@Test
|
||||
void rosterViewExposesTheCharterReceiptButNeverTheCharterText() {
|
||||
// The roster (bridge_list and GET /members both render through rosterView) must let a lead
|
||||
// The roster (fleet_list and GET /members both render through rosterView) must let a lead
|
||||
// see which charter a member got, without ever carrying the charter prose itself (CB-571).
|
||||
MemberSession s = new MemberSession("p1", "term1", "prof", MemberRole.DEV, "/cwd", null,
|
||||
0, 0, 0, MemberSession.State.READY, null, null,
|
||||
|
||||
@@ -15,7 +15,7 @@ Consequences today:
|
||||
since when?" without shelling to herdr for a raw agent list (no state, no ownership, no age).
|
||||
- Cleanup of a worker that outlived its owning process depends entirely on the boot-time
|
||||
name-nonce **reaper** (CB-117) — there is no live, authoritative roster during a run.
|
||||
- `fleet_list` (CB-304) can only surface herdr's view, not a bridge-owned roster.
|
||||
- `bridge_list` (CB-304) can only surface herdr's view, not a bridge-owned roster.
|
||||
- There is no seam for per-session policy (checkpoint on teardown → CB-302; idle_ttl /
|
||||
context_cap / drain → CB-303).
|
||||
|
||||
@@ -100,9 +100,9 @@ final class SessionManager {
|
||||
`TurnListener` alongside `CompletionResolver` so it sees turn boundaries, and give it the
|
||||
`WorkerPresence` signal for `READY`.
|
||||
- **`BridgeMcp.spawn` / `BridgedApp.spawnWorker`** — route spawn through `SessionManager.acquire`
|
||||
(carry `callerTerminal` as `ownerTerminal`). **`fleet_stop` / `DELETE /workers/{paneId}`** →
|
||||
(carry `callerTerminal` as `ownerTerminal`). **`bridge_stop` / `DELETE /workers/{paneId}`** →
|
||||
`SessionManager.release`.
|
||||
- **`fleet_list` / `GET /sessions` (CB-304 later)** — read `SessionManager.roster()`.
|
||||
- **`bridge_list` / `GET /sessions` (CB-304 later)** — read `SessionManager.roster()`.
|
||||
- **`MessageService`** — no change required for one-shot; a later CB-303 auto-release hook can call
|
||||
`release` from `onTurnComplete` under policy.
|
||||
|
||||
@@ -121,4 +121,4 @@ final class SessionManager {
|
||||
- **CB-302** — attach a checkpoint step (`STATE.md` + commit) to the `release` path.
|
||||
- **CB-303** — a policy loop over `roster()` using `spawnedAtNanos`/state to auto-`release` on
|
||||
`idle_ttl`, or drain on `context_cap`.
|
||||
- **CB-304** — `fleet_list` reads `roster()` for a bridge-owned roster + live join.
|
||||
- **CB-304** — `bridge_list` reads `roster()` for a bridge-owned roster + live join.
|
||||
|
||||
@@ -3,7 +3,7 @@
|
||||
**Status:** ✅ shipped — implemented at commit `97ecc71` (per-worker git worktree + config-parity
|
||||
overlay). As-built: `session/GitWorktrees.java` behind the `Worktrees` port, wired in
|
||||
`Bridged.main` and configurable via `worktreeRoot` / per-profile `parityOverlay`
|
||||
(see `bridged.example.yaml`). Branch/worktree surface in `fleet_list` landed with CB-304
|
||||
(see `bridged.example.yaml`). Branch/worktree surface in `bridge_list` landed with CB-304
|
||||
(`9fe04bf`); the worker-opened-PR checkpoint landed as CB-302 (`64e70ef`).
|
||||
**Extends:** [CB-301 Session Manager](CB-301-Session-Manager.md) (shipped, commit `54d907c`).
|
||||
**Realizes:** the config-parity requirement in [Worker Git Workflow](Worker-Git-Workflow.md).
|
||||
@@ -122,7 +122,7 @@ public void release(String paneId) {
|
||||
|
||||
### Surface: MCP + REST
|
||||
|
||||
- `fleet_spawn` gains an optional `worktree` arg: `true`, or a ticket slug string. Truthy ⇒ build a
|
||||
- `bridge_spawn` gains an optional `worktree` arg: `true`, or a ticket slug string. Truthy ⇒ build a
|
||||
`WorktreeRequest(slug, null)` and call the 5-arg `acquire`.
|
||||
- `POST /workers` gains `worktree` (+ optional `ticket`) in the body/query, same mapping.
|
||||
- `workerView`/`view(WorkerSession)` include `worktree` and `branch` **when non-null** (omit for
|
||||
|
||||
@@ -6,11 +6,11 @@
|
||||
|
||||
## 1. Problem
|
||||
|
||||
`fleet_spawn` today returns a session the instant the herdr pane is started. The pane is not
|
||||
`bridge_spawn` today returns a session the instant the herdr pane is started. The pane is not
|
||||
yet a usable Claude REPL — it may still be sitting at the folder-trust prompt, or the CLI may
|
||||
never come up at all. Nothing blocks or times out on that. Consequences:
|
||||
|
||||
- A `fleet_send` to a not-yet-ready worker surfaces as a **~60 s MCP-client timeout** (the send
|
||||
- A `bridge_send` to a not-yet-ready worker surfaces as a **~60 s MCP-client timeout** (the send
|
||||
blocks waiting for a turn that can't start) instead of a fast, explicit spawn failure.
|
||||
- A worker stuck at the folder-trust prompt lingers in `SPAWNING` forever; nothing fails it.
|
||||
|
||||
@@ -99,11 +99,11 @@ Jackson ignores unknown keys, so omitting them in existing YAML is safe; pick sa
|
||||
`spawn` throws — verify the new exception flows through it (worktree removed, nothing registered).
|
||||
- The **non-worktree** path registers the session only *after* `spawn` returns, so a throw means no
|
||||
half-live `SPAWNING` session is ever registered — confirm this and add a test.
|
||||
- `fleet_spawn` (MCP verb) must return an **error result** carrying the exception message, not a
|
||||
- `bridge_spawn` (MCP verb) must return an **error result** carrying the exception message, not a
|
||||
success with a dead session. Trace `BridgeMcp`/`BridgedApp` spawn handlers and make sure the
|
||||
exception becomes a clean tool error, not an uncaught 500 with a stack trace.
|
||||
|
||||
**Out of scope (do NOT do here):** gating `fleet_send` on session `READY` (existing status-gate +
|
||||
**Out of scope (do NOT do here):** gating `bridge_send` on session `READY` (existing status-gate +
|
||||
this spawn gate already close the window), MCP-handshake-as-readiness signal, the CB-307 broker,
|
||||
any config `kind:` discriminator, any second adapter.
|
||||
|
||||
|
||||
@@ -13,7 +13,7 @@ Stage 2 (the AMQP/LavinMQ adapter behind the same port) is explicitly **out of s
|
||||
## 1. The bug this fixes (grounded in current code)
|
||||
|
||||
The reverse (worker→primary) path is `Rendezvous` — a `ConcurrentHashMap<session, CompletableFuture<Resolution>>`
|
||||
of **live blocking waiters only**. No queue, no store. When a worker calls `fleet_reply` and **no send
|
||||
of **live blocking waiters only**. No queue, no store. When a worker calls `bridge_reply` and **no send
|
||||
is currently open** for that worker:
|
||||
|
||||
- `Rendezvous.resolve(session, content)` → `complete(...)` → `waiters.get(session) == null` →
|
||||
@@ -22,7 +22,7 @@ is currently open** for that worker:
|
||||
`BridgeMcp.reply` returns `error("no send is awaiting a reply for this worker")` (`mcp/BridgeMcp.java:270-272`);
|
||||
REST returns `409 no_pending_send` (`rest/BridgedApp.java:339-345`).
|
||||
|
||||
This is the observed "communication break": a worker that finishes just after its `fleet_send` timed
|
||||
This is the observed "communication break": a worker that finishes just after its `bridge_send` timed
|
||||
out (the ~60s sync window) replies into the void. There is **no message-id, dedup, or ack** anywhere in
|
||||
the message path today.
|
||||
|
||||
@@ -79,7 +79,7 @@ in `MessageService`, which already owns the `Rendezvous` and will own the `Reply
|
||||
|
||||
- Add `MessageService.reply(String session, String content)`:
|
||||
```java
|
||||
/** Route a worker's explicit fleet_reply: resolve an open send, or queue it in the inbox if none. */
|
||||
/** Route a worker's explicit bridge_reply: resolve an open send, or queue it in the inbox if none. */
|
||||
public boolean reply(String session, String content) {
|
||||
if (rendezvous.resolve(session, content)) {
|
||||
return true; // a live send took it — unchanged fast path
|
||||
@@ -94,7 +94,7 @@ in `MessageService`, which already owns the `Rendezvous` and will own the `Reply
|
||||
- `BridgedApp.replyMessage` (`rest/BridgedApp.java:330-346`) — return `200` (queued) instead of
|
||||
`409 no_pending_send`.
|
||||
|
||||
**DO NOT touch the QUESTION path.** `fleet_ask` / `rendezvous.resolveQuestion` must keep today's
|
||||
**DO NOT touch the QUESTION path.** `bridge_ask` / `rendezvous.resolveQuestion` must keep today's
|
||||
`NO_WAITER` behaviour — a mid-turn question is **interactive** (the worker blocks synchronously and cannot
|
||||
consume a late answer), so it must **never** be queued. Only terminal `REPLY`s go to the inbox.
|
||||
|
||||
@@ -110,11 +110,11 @@ The primary re-checks a worker it delegated to. Expose a drain keyed by **worker
|
||||
- Add `MessageService.drainReplies(String target)`: `peek` the inbox, `ack` each returned `msgId`, hand
|
||||
back the `List<InboxMessage>` (or just the contents). At-least-once: peek→deliver→ack (ack only after
|
||||
the caller has them, so an in-flight failure re-surfaces them).
|
||||
- **DECISION (required default): expose via the existing poll verb, keyed by target.** Extend `fleet_poll`
|
||||
- **DECISION (required default): expose via the existing poll verb, keyed by target.** Extend `bridge_poll`
|
||||
to accept an optional `target` (worker session) and, when present, return that worker's drained replies —
|
||||
alongside a matching REST route `GET /sessions/{id}/replies`. Do **not** change `send`/`answer` semantics
|
||||
(do not drain inside `send` — that conflates "deliver to worker" with "collect its mail"). Keep the
|
||||
existing ticket-based `fleet_poll(ticket)` path working unchanged. If you see a cleaner surface, still
|
||||
existing ticket-based `bridge_poll(ticket)` path working unchanged. If you see a cleaner surface, still
|
||||
ship this default and note the alternative for review.
|
||||
|
||||
## 3. Config
|
||||
@@ -126,11 +126,11 @@ Stage-2 AMQP adapter: absent → in-memory, present → AMQP).
|
||||
## 4. Acceptance criteria (what the primary will verify)
|
||||
|
||||
1. New `ReplyInbox` + `InboxMessage` + `InMemoryReplyInbox` in `dev.ltms.bridged.msg`.
|
||||
2. `fleet_reply` with **no open send** now **succeeds and queues** (no more `error` / `409`); the reply is
|
||||
2. `bridge_reply` with **no open send** now **succeeds and queues** (no more `error` / `409`); the reply is
|
||||
later retrievable and identical.
|
||||
3. The queued reply is drainable by the primary keyed by target; draining **acks** it (a second drain
|
||||
returns nothing); dedup by `msgId` (re-publishing the same id does not double-queue).
|
||||
4. **QUESTION path unchanged** — `fleet_ask` with no open send still returns `NO_WAITER` (add/keep a test
|
||||
4. **QUESTION path unchanged** — `bridge_ask` with no open send still returns `NO_WAITER` (add/keep a test
|
||||
proving a question is never queued).
|
||||
5. Completion/failure fallbacks unchanged.
|
||||
6. Unit tests covering: `InMemoryReplyInbox` publish/peek/ack/dedup/FIFO/concurrency; `MessageService.reply`
|
||||
|
||||
@@ -111,7 +111,7 @@ host's terminals.*
|
||||
|
||||
```mermaid
|
||||
flowchart LR
|
||||
send["fleet_send(globalId, msg)"] --> lookup{"directory:<br/>is globalId local?"}
|
||||
send["bridge_send(globalId, msg)"] --> lookup{"directory:<br/>is globalId local?"}
|
||||
lookup -->|"yes"| local["inject via local herdr<br/>(today's Injector path)"]
|
||||
lookup -->|"no"| pub["publish agent.<id>.inbox<br/>(broker routes to owning gateway)"]
|
||||
pub --> consume["owning gateway consumes<br/>→ injects into its local herdr"]
|
||||
@@ -158,17 +158,17 @@ sequenceDiagram
|
||||
participant BR as broker
|
||||
participant GA as gateway A
|
||||
participant P as primary (host A, MCP client)
|
||||
W->>GB: fleet_reply
|
||||
W->>GB: bridge_reply
|
||||
GB->>BR: publish primary-bound (durable, msg id)
|
||||
BR->>GA: route to A's primary inbox
|
||||
Note over GA: held durably until the primary pulls
|
||||
P->>GA: blocking fleet_send resolves / fleet_poll
|
||||
P->>GA: blocking bridge_send resolves / bridge_poll
|
||||
GA-->>P: reply (then ACK to broker)
|
||||
```
|
||||
|
||||
*Figure 4 — the broker makes the middle hop lossless, ordered, and idempotent; the **final** hop
|
||||
into the primary is still a **pull** (gateway A holds the message until the primary's blocking
|
||||
`fleet_send` or `fleet_poll`). Cross-host neither improves nor worsens this — it just spans hosts.
|
||||
`bridge_send` or `bridge_poll`). Cross-host neither improves nor worsens this — it just spans hosts.
|
||||
This is precisely the gap CB-307 closes on one host and CB-308 stretches across hosts.*
|
||||
|
||||
## 6. Staging & dependencies
|
||||
@@ -207,7 +207,7 @@ they overlap (notably: the envelope is no longer optional, and dedup is split by
|
||||
gateway's host. This extends the single-host invariant — *identity comes from the connection,
|
||||
never an argument* — across the broker: cross-host, identity comes from the key. Complements
|
||||
(not replaces) per-gateway broker logins over TLS.
|
||||
2. **Profiles are owned by the worker's host.** `fleet_spawn(profile, host)` resolves the name in
|
||||
2. **Profiles are owned by the worker's host.** `bridge_spawn(profile, host)` resolves the name in
|
||||
the *target* gateway's `bridged.yaml`. Gateways advertise their profile names in presence
|
||||
heartbeats, so a leader sees what each host offers before spawning; an unknown name is a clear
|
||||
error from the target. Secrets (base URLs, tokens) never leave the host that uses them.
|
||||
@@ -259,7 +259,7 @@ they overlap (notably: the envelope is no longer optional, and dedup is split by
|
||||
serialize acks fleet-wide. Ordering caveat: a *return* (unroutable) arrives **before** the
|
||||
confirm, so "confirmed" ≠ "routed"; the sender checks the returned-set at confirm time.
|
||||
`mandatory` is false only for `BROADCAST`, where an empty group is legal silence.
|
||||
10. **Queue lifecycle is session lifecycle.** `fleet_stop`/reap deletes the worker's inbox queue
|
||||
10. **Queue lifecycle is session lifecycle.** `bridge_stop`/reap deletes the worker's inbox queue
|
||||
(its `broadcast.*` bindings die with it — no broadcasts to the dead); `x-expires` collects
|
||||
queues orphaned by a crashed gateway (long for main/orchestrator inboxes, short for workers).
|
||||
Queue names carry a version suffix (`.v2`): AMQP refuses to redeclare an existing durable
|
||||
|
||||
@@ -7,7 +7,7 @@ the core learning any one peer's environment.
|
||||
## 1. Why
|
||||
|
||||
`claude-bridge` is a **communication bus between heterogeneous AI agents** — its stable surface is
|
||||
the protocol (`fleet_spawn / send / poll / reply / ask / list / stop / status`), and that surface
|
||||
the protocol (`bridge_spawn / send / poll / reply / ask / list / stop / status`), and that surface
|
||||
should stay provider-neutral. Today the daemon can only materialize one kind of peer: an
|
||||
off-subscription Claude Code CLI over herdr. Everything specific to *how that peer is set up*
|
||||
(`ANTHROPIC_BASE_URL`, the subscription guard, `--mcp-config`/system-prompt flags, `claude-*`
|
||||
@@ -126,7 +126,7 @@ existing configs. (Jackson already ignores unknown keys, so adding `kind` is bac
|
||||
|
||||
```mermaid
|
||||
sequenceDiagram
|
||||
participant MCP as fleet_spawn (MCP/REST)
|
||||
participant MCP as bridge_spawn (MCP/REST)
|
||||
participant SM as SessionManager
|
||||
participant L as PeerLauncher (by profile.kind)
|
||||
participant T as transport (herdr)
|
||||
@@ -150,7 +150,7 @@ degrades gracefully when a launcher lacks one:
|
||||
|
||||
| Capability | Meaning | Claude Code | Codex (likely) | Human |
|
||||
|---|---|---|---|---|
|
||||
| `MID_TURN_ASK` | supports `fleet_ask` rendezvous | ✓ | ? | ✗ |
|
||||
| `MID_TURN_ASK` | supports `bridge_ask` rendezvous | ✓ | ? | ✗ |
|
||||
| `SELF_PR` | can open its own PR at checkpoint (CB-302) | ✓ (opt-in token) | ? | ✗ |
|
||||
| `WORKTREE` | can run in a provisioned git worktree | ✓ | ✓ | ✗ |
|
||||
| `ORPHAN_REAP` | spawner can reconcile orphaned peers on boot | ✓ | ? | ✗ |
|
||||
|
||||
@@ -184,7 +184,7 @@ sequenceDiagram
|
||||
participant O as OpenCodeLauncher
|
||||
participant B as HerdrPeerLauncher (base)
|
||||
participant H as herdr
|
||||
P->>M: fleet_spawn(profile="oc-impl")
|
||||
P->>M: bridge_spawn(profile="oc-impl")
|
||||
M->>C: spawn(SpawnRequest)
|
||||
C->>C: kind(profile)=="opencode"
|
||||
C->>O: spawn(req)
|
||||
@@ -226,11 +226,11 @@ until it is "major" (Stage-B whole), per the CB-401 bar.
|
||||
- **opencode TUI ⇄ herdr injection.** herdr drives a pane by typing into a TUI. Must confirm
|
||||
opencode's TUI accepts injected keystrokes/submit the way `claude` does, and reaches an
|
||||
`injectable` status the CB-306 gate recognizes. *Validation:* spawn one opencode worker,
|
||||
watch the readiness gate pass, `fleet_send` a trivial task.
|
||||
watch the readiness gate pass, `bridge_send` a trivial task.
|
||||
- **Bridge MCP visibility in opencode.** Confirm `OPENCODE_CONFIG` (or `opencode mcp add`)
|
||||
actually surfaces the `fleet_*` tools inside the opencode session, and that `fleet_reply`
|
||||
actually surfaces the `bridge_*` tools inside the opencode session, and that `bridge_reply`
|
||||
is callable — the reply-charter is worthless if the tool isn't mounted. *Validation:* the
|
||||
worker completes a task by calling `fleet_reply`; the reply lands via the CB-307 path.
|
||||
worker completes a task by calling `bridge_reply`; the reply lands via the CB-307 path.
|
||||
- **opencode MCP/config schema drift.** opencode is fast-moving (1.1.31 today). Pin the config
|
||||
schema we generate against the installed version; treat the exact keys (`type: "remote"` vs
|
||||
`"http"`, `instructions` shape) as a dogfood-verified fact, not an assumption.
|
||||
@@ -281,7 +281,7 @@ as a dogfood-verified fact for 1.18.5.
|
||||
| §5 risk | Result |
|
||||
|---|---|
|
||||
| opencode TUI ⇄ herdr injection; CB-306 gate | ✅ `peer pane=wD:p3 reached injectable state` ~0.6s after `agent.start` |
|
||||
| Bridge MCP visible + `fleet_reply` callable | ✅ MCP `initialize` from `Implementation[name=opencode, version=1.18.5]`; worker replied through the tool |
|
||||
| Bridge MCP visible + `bridge_reply` callable | ✅ MCP `initialize` from `Implementation[name=opencode, version=1.18.5]`; worker replied through the tool |
|
||||
| Config schema drift (1.1.31 → 1.18.5) | ✅ unchanged, see above |
|
||||
| Provider credentials | ✅ free tier, zero credentials |
|
||||
|
||||
@@ -291,7 +291,7 @@ Full lifecycle exercised through the REST surface:
|
||||
to `OpenCodeLauncher` (`spawning opencode profile=opencode-free`), pane `wD:p3`.
|
||||
2. Readiness: `{"ready":true,"status":"idle"}`, roster state `ready`.
|
||||
3. `POST /sessions/{id}/message` → **`{"replySource":"reply","reply":"391"}`** — a *structured*
|
||||
`fleet_reply`, not the CB-115 completion-fallback transcript scrape. The clean path.
|
||||
`bridge_reply`, not the CB-115 completion-fallback transcript scrape. The clean path.
|
||||
4. `DELETE /workers/wD:p3` → `204`, roster empty, tolerant teardown (`tab_not_found` ignored —
|
||||
opencode had already closed its own tab).
|
||||
|
||||
|
||||
@@ -45,7 +45,7 @@ flowchart TB
|
||||
w1["worker pane (gx00 vLLM)"]
|
||||
w2["worker pane (ollama)"]
|
||||
human --> primary
|
||||
primary -->|"fleet_send / spawn / ask"| daemon
|
||||
primary -->|"bridge_send / spawn / ask"| daemon
|
||||
daemon --> comp
|
||||
comp --> cc
|
||||
comp --> oc
|
||||
@@ -150,7 +150,7 @@ sequenceDiagram
|
||||
participant SL as SandboxLauncher
|
||||
participant SB as sandbox (peer-owned)
|
||||
participant A as agent in sandbox
|
||||
M->>D: fleet_spawn(profile=backend, role=backend)
|
||||
M->>D: bridge_spawn(profile=backend, role=backend)
|
||||
D->>SL: spawn(SpawnRequest)
|
||||
SL->>SB: start entrypoint (image = backend role)
|
||||
Note over SL,SB: bridge injects guarded ANTHROPIC_BASE_URL,<br/>mounts bridge MCP url + reply charter
|
||||
@@ -213,12 +213,12 @@ sequenceDiagram
|
||||
participant MA as main A (Opus)
|
||||
participant BR as bridged / broker
|
||||
participant MB as main B (cloud)
|
||||
MA->>BR: fleet_send(to = main B, msg)
|
||||
MA->>BR: bridge_send(to = main B, msg)
|
||||
BR->>BR: publish agent.B.inbox (durable, msg id)
|
||||
Note over BR: held until B pulls (B is a client too)
|
||||
MB->>BR: blocking fleet_send / poll resolves
|
||||
MB->>BR: blocking bridge_send / poll resolves
|
||||
BR-->>MB: msg (then ACK)
|
||||
MB->>BR: fleet_reply(to = main A)
|
||||
MB->>BR: bridge_reply(to = main A)
|
||||
BR->>BR: publish agent.A.inbox
|
||||
MA->>BR: poll resolves
|
||||
BR-->>MA: reply
|
||||
@@ -287,7 +287,7 @@ sequenceDiagram
|
||||
H->>O: high-level goal (large context)
|
||||
O->>O: name/resume main A session
|
||||
O->>MA: task + SCOPED context slice (not the whole history)
|
||||
MA->>W: fleet_send(delegation, carrying only the relevant slice)
|
||||
MA->>W: bridge_send(delegation, carrying only the relevant slice)
|
||||
W-->>MA: result
|
||||
MA-->>O: rollup
|
||||
O->>O: fold into orchestrator context, pick next main/turn
|
||||
@@ -440,13 +440,13 @@ sequenceDiagram
|
||||
participant BR as broker
|
||||
participant GB as gateway B
|
||||
participant SB as sandbox agent (host B, container)
|
||||
MA->>GA: fleet_send(globalId on B, msg)
|
||||
MA->>GA: bridge_send(globalId on B, msg)
|
||||
GA->>GA: directory lookup - is globalId local? NO
|
||||
GA->>BR: publish agent.ID.inbox (durable)
|
||||
BR->>GB: route to the owning gateway
|
||||
GB->>SB: inject via B's LOCAL herdr (keystrokes)
|
||||
Note over GB,SB: SandboxLauncher already spawned the container -<br/>its PTY is in B's herdr, CB-306 readiness passed
|
||||
SB-->>GB: fleet_reply (to B's LOCAL MCP endpoint)
|
||||
SB-->>GB: bridge_reply (to B's LOCAL MCP endpoint)
|
||||
GB->>BR: publish primary-bound (durable, msg id)
|
||||
BR->>GA: route back to A
|
||||
Note over GA: held until the main pulls (the main is a client)
|
||||
|
||||
@@ -196,7 +196,7 @@ every member spawn. The wiki keeps the direct route open precisely because "if t
|
||||
nothing that matters is blocked."
|
||||
|
||||
So keep it. Add a second profile `local-direct` pointing at `http://gx00.gw:8000` with **`weight: 0`**
|
||||
— never auto-selected, still spawnable with an explicit `fleet_spawn{profile: "local-direct"}`.
|
||||
— never auto-selected, still spawnable with an explicit `bridge_spawn{profile: "local-direct"}`.
|
||||
That is exactly what CB-554 made `weight: 0` mean, and it turns the escape hatch into something the
|
||||
lead can actually reach during an incident.
|
||||
|
||||
@@ -266,7 +266,7 @@ Each of these cost someone real debugging time upstream. They apply to us.
|
||||
1. **Rotating a token restarts the auth proxy, which drops in-flight streaming responses.** For us
|
||||
that means rotating `AI_GATEWAY_TOKEN` kills every live member mid-turn, and an async ticket's
|
||||
report goes with it. This is the same rule as a daemon redeploy: **drain the fleet first**
|
||||
(`fleet_list` → `fleet_poll` anything wanted → `fleet_stop`), then rotate.
|
||||
(`bridge_list` → `bridge_poll` anything wanted → `bridge_stop`), then rotate.
|
||||
2. **The gateway's own `SecurityPolicy` fails open.** Standalone `aigw run` accepts it and silently
|
||||
ignores it — an unauthenticated request returned **200**. Auth is the Caddy proxy in front, and
|
||||
nothing else. Never reason as if the gateway authenticates.
|
||||
@@ -283,10 +283,10 @@ Each of these cost someone real debugging time upstream. They apply to us.
|
||||
|
||||
Merging config is not proving it. The checks, in order:
|
||||
|
||||
1. `fleet_spawn{profile: "gx"}` succeeds and the member completes a real turn ending in
|
||||
`fleet_reply`. This is the first proof of the token, the URL and the model name, and it risks
|
||||
1. `bridge_spawn{profile: "gx"}` succeeds and the member completes a real turn ending in
|
||||
`bridge_reply`. This is the first proof of the token, the URL and the model name, and it risks
|
||||
nothing the fleet depends on.
|
||||
2. `fleet_spawn{profile: "local"}` succeeds. If the guard allowlist was missed, this **throws** — a
|
||||
2. `bridge_spawn{profile: "local"}` succeeds. If the guard allowlist was missed, this **throws** — a
|
||||
loud, self-correcting failure, which is the good kind. If the restart was missed, it also throws,
|
||||
for the same reason.
|
||||
3. A `local` member completes a turn. That exercises streaming through two TLS edges, the auth proxy
|
||||
@@ -297,9 +297,9 @@ Merging config is not proving it. The checks, in order:
|
||||
deltas. For `gx` on `/v1`, this answers the open question in §3 rather than assuming it.
|
||||
5. The cockpit at `auth.ltms.dev` shows requests counted against the `claude-bridge` consumer, not
|
||||
`legacy`. That is the whole point of taking our own token.
|
||||
6. `fleet_spawn{profile: "local-direct"}` still works, so the escape hatch is real rather than
|
||||
6. `bridge_spawn{profile: "local-direct"}` still works, so the escape hatch is real rather than
|
||||
theoretical.
|
||||
7. `fleet_list` shows `gx` carrying no `credentialId`, so a `sol`/`terra` exhaustion cannot
|
||||
7. `bridge_list` shows `gx` carrying no `credentialId`, so a `sol`/`terra` exhaustion cannot
|
||||
quarantine it. This is the single-point-of-failure claim in §3b, checked rather than asserted.
|
||||
|
||||
---
|
||||
|
||||
@@ -172,7 +172,7 @@ The current split is the mirror image of what you'd expect:
|
||||
| REST routes | ❌ **none at all** — the session id is taken from the URL path and trusted | ❌ none |
|
||||
|
||||
So REST is the *more* exposed surface: `POST /sessions/{id}/reply` accepts any `{id}` from the
|
||||
path, whereas the MCP `fleet_reply` derives the worker from the connection and refuses to read it
|
||||
path, whereas the MCP `bridge_reply` derives the worker from the connection and refuses to read it
|
||||
from an argument. Loopback-only bind is what makes this safe today.
|
||||
|
||||
**Therefore CB-505 must enforce on both paths against one shared resolver** — not at a single
|
||||
|
||||
+19
-19
@@ -1,7 +1,7 @@
|
||||
# M4 - Fleet health, recovery, routing, and capacity
|
||||
|
||||
**Status:** Design accepted on 2026-08-15. CB-573 part 1 has shipped the classification model and
|
||||
the `fleet_list` capacity view; the remaining M4 units are not yet shipped. See
|
||||
the `bridge_list` capacity view; the remaining M4 units are not yet shipped. See
|
||||
[Unit 2 — what has landed so far](#unit-2---what-has-landed-so-far) before planning Unit 2 work:
|
||||
some of its criteria were met by separate CB tickets, and one of them contradicts the unit text.
|
||||
**Scope:** Fleet evidence, safe mechanical repair, lead routing, capacity reporting, and optional
|
||||
@@ -145,7 +145,7 @@ Idle may drive configured resource cleanup. It never opens an incident and never
|
||||
| `TURN_BOUNDARY_LOST` | Same session turn stays `BUSY`; same accepted task stays open; two raw snapshots show `IDLE` or `DONE` | Strong disagreement. Strict reconciliation may repair it. |
|
||||
| `ERROR_ON_SCREEN` | Suspicious non-working state survives grace; `detection` matches a tested adapter-specific fatal signature | Certain only for the matched signature. A bare word such as `Exception` is not enough. |
|
||||
| `STALL_SUSPECTED` | Open turn is older than the configured threshold; two normalised `recent_unwrapped` digests are unchanged; no boundary or reply occurs | Not certain. A long valid API call can look the same. Lead decides. |
|
||||
| `MUTE` | Turn resolves through completion fallback instead of `fleet_reply` | Certain that no structured reply won. It does not prove an MCP failure. A single event is a metric, not an incident. |
|
||||
| `MUTE` | Turn resolves through completion fallback instead of `bridge_reply` | Certain that no structured reply won. It does not prove an MCP failure. A single event is a metric, not an incident. |
|
||||
| `REPLY_STRANDED` | Typed reply or health message remains after owning-lead push reaches its cap | Collection failed. This does not explain whether the lead is busy, dead, or ignoring the nudge. |
|
||||
| `DELEGATION_ORPHANED` | Target is gone, failed, or released, but one or more tasks remain `PENDING` after reconciliation grace | Certain bridge invariant failure. This is not an inbox-drain fault. |
|
||||
| `WORK_PRODUCT_AT_RISK` | Provisioned branch has commits after its recorded base; member is `DONE`, `FAILED`, or preserved after release; no turn or inbox item remains; long-idle threshold passed | A warning, not proof of loss. Work may already have an open pull request or a squash merge. |
|
||||
@@ -230,8 +230,8 @@ Only the lead may:
|
||||
- choose how to use partial work in a worktree;
|
||||
- restart herdr or change network, model, credentials, backend, or configuration.
|
||||
|
||||
Reports include literal safe tool calls such as `fleet_status(sessionId="...")`,
|
||||
`fleet_poll(ticket="...")`, `fleet_list()`, and optional `fleet_stop(paneId="...")`. A judgement
|
||||
Reports include literal safe tool calls such as `bridge_status(sessionId="...")`,
|
||||
`bridge_poll(ticket="...")`, `bridge_list()`, and optional `bridge_stop(paneId="...")`. A judgement
|
||||
state never presents stop as the only action.
|
||||
|
||||
### 5.3 Release causes and worktree safety
|
||||
@@ -251,7 +251,7 @@ preserving cause.
|
||||
|
||||
Before abnormal release removes the live session, M4 writes an atomic manifest under the worktree
|
||||
root. It records session identity, owner, role, profile, repository, path, branch, base commit,
|
||||
release cause, release time, state, and pending task ids. `fleet_list.preservedWorktrees` loads these
|
||||
release cause, release time, state, and pending task ids. `bridge_list.preservedWorktrees` loads these
|
||||
manifests after restart. Stop output and WARN logs also name the path and cause. M4 never
|
||||
auto-deletes a preserved worktree.
|
||||
|
||||
@@ -348,7 +348,7 @@ Invalid known-format data is copied byte-for-byte to durable queue
|
||||
before the original is acknowledged. A failed quarantine handoff leaves the original unacknowledged.
|
||||
The raw body never enters logs.
|
||||
|
||||
Decode failure creates a redacted WARN, metric, `fleet_list` summary, and routed health incident.
|
||||
Decode failure creates a redacted WARN, metric, `bridge_list` summary, and routed health incident.
|
||||
One bad entry never escapes the consumer callback and never stops later valid messages.
|
||||
|
||||
Safe downgrade is not supported. The previous build ignores `content_type` and would show typed JSON
|
||||
@@ -399,13 +399,13 @@ Choose the candidate with the fewest assigned foreign incidents. Break ties by s
|
||||
then terminal id. Pin the recipient. Reassign only if that peer becomes unhealthy or retires. A
|
||||
routing generation marks a reassignment, and old pending assignments become superseded.
|
||||
|
||||
`fleet_list` lead rows show health, health age, assigned foreign incident count, and a bounded list
|
||||
`bridge_list` lead rows show health, health age, assigned foreign incident count, and a bounded list
|
||||
of incident id, subject, state, severity, age, and routing generation. The top-level view also shows
|
||||
owner, recipient, and routing reason.
|
||||
|
||||
A peer incident is published under the recipient lead's inbox key, not the failed subject's key. Its
|
||||
status-gated nudge names the failed lead and gives the exact
|
||||
`fleet_poll(target="<recipient-terminal>")` call.
|
||||
`bridge_poll(target="<recipient-terminal>")` call.
|
||||
|
||||
### 7.5 Lead inbox ownership
|
||||
|
||||
@@ -427,7 +427,7 @@ An idle, reachable lead may receive the existing bounded nudge. An unreachable o
|
||||
has no safe in-loop recovery. The bridge must not restart or replace it. A new lead would not have the
|
||||
failed lead's plan or context, and an uncertain relaunch could create two orchestrators.
|
||||
|
||||
With no webhook, only `fleet_list`, `/healthz`, metrics, WARN logs, and the incident journal remain.
|
||||
With no webhook, only `bridge_list`, `/healthz`, metrics, WARN logs, and the incident journal remain.
|
||||
These are passive surfaces. They are not a human notification.
|
||||
|
||||
## 8. Lost-boundary reconciliation
|
||||
@@ -521,7 +521,7 @@ task poll source, and metrics. The lead sees:
|
||||
|
||||
```text
|
||||
[repaired completion - bridged detected a lost turn boundary. The member did not call
|
||||
fleet_reply; pane-derived text follows and may be partial]
|
||||
bridge_reply; pane-derived text follows and may be partial]
|
||||
```
|
||||
|
||||
Clipped text also keeps the existing clipped-tail marker.
|
||||
@@ -551,7 +551,7 @@ failure operation once. It never recreates the task.
|
||||
|
||||
Capacity is a view, not a health state.
|
||||
|
||||
`fleet_list` adds one block per profile:
|
||||
`bridge_list` adds one block per profile:
|
||||
|
||||
```text
|
||||
profile, maxLoad, live, free, reclaimable
|
||||
@@ -612,7 +612,7 @@ incident feeds lead-health evidence. If the lead then becomes unhealthy, peer or
|
||||
Webhook mode requires a resolved environment variable. Turning notification off stops outbound
|
||||
attempts but keeps incidents. Turning it back on resumes still-open human incidents.
|
||||
|
||||
Without a sink, `fleet_list.healthCoverage` states that human escalation is unavailable. `/healthz`
|
||||
Without a sink, `bridge_list.healthCoverage` states that human escalation is unavailable. `/healthz`
|
||||
keeps its existing HTTP liveness result and adds a nested `fleetHealth.status=partial` component.
|
||||
Metrics and one startup or reload WARN expose the same limit.
|
||||
|
||||
@@ -792,7 +792,7 @@ Acceptance criteria:
|
||||
them.
|
||||
16. Explicit stop is state-aware. Any pending task or non-terminal state preserves the worktree.
|
||||
17. Atomic preserved-worktree manifests reload after restart and appear in lead-only
|
||||
`fleet_list.preservedWorktrees`.
|
||||
`bridge_list.preservedWorktrees`.
|
||||
18. Manifest failure preserves the worktree and opens an operator-visible health failure.
|
||||
19. Provision records the base commit. Terminal, long-idle worktrees report
|
||||
`WORK_PRODUCT_AT_RISK` only under the evidence in Section 4.2 and never auto-delete work.
|
||||
@@ -820,7 +820,7 @@ met, and I did not run one.
|
||||
| 14 | `DELEGATION_ORPHANED` | 3 files | **Partial.** The health state exists. The teardown-invariant check that creates it, and the retry rule, do not. |
|
||||
| 15 | `SPAWN_ROLLBACK` | 0 files | **Contradicted — see below.** |
|
||||
| 16 | — | — | Partial at best. CB-576 made release preserve a dirty worktree; whether explicit stop is state-aware is not checked. |
|
||||
| 17, 18 | `preservedWorktrees` | 0 files | Not started. No manifest, and no lead-only `fleet_list` field. |
|
||||
| 17, 18 | `preservedWorktrees` | 0 files | Not started. No manifest, and no lead-only `bridge_list` field. |
|
||||
| 19 | `WORK_PRODUCT_AT_RISK` | 0 files | Not started. |
|
||||
|
||||
**Criterion 15 no longer matches the code, and the code is right.** It says "normal `COMPLETED`
|
||||
@@ -837,7 +837,7 @@ preserve it. `SPAWN_ROLLBACK` itself does not exist yet.
|
||||
### Unit 3 - Typed inbox and member routing
|
||||
|
||||
Scope: semantic record, AMQP migration, both adapters, member routing, polling, and member health in
|
||||
`fleet_list`.
|
||||
`bridge_list`.
|
||||
|
||||
Acceptance criteria:
|
||||
|
||||
@@ -852,7 +852,7 @@ Acceptance criteria:
|
||||
6. Invalid known data never escapes the callback, appears as a reply, or blocks later valid messages.
|
||||
7. Invalid data reaches durable per-target quarantine before original ack. Failed handoff leaves the
|
||||
original unacked.
|
||||
8. Decode failures create redacted WARN, metric, `fleet_list` summary, and routed incident without
|
||||
8. Decode failures create redacted WARN, metric, `bridge_list` summary, and routed incident without
|
||||
raw content.
|
||||
9. Both adapters pass one semantic contract for fields, FIFO, dedup, ownership, ack, and release.
|
||||
10. Lead keys require explicit ownership. Publication never claims a queue.
|
||||
@@ -865,7 +865,7 @@ Acceptance criteria:
|
||||
15. RabbitMQ contract tests pass with `mvn test -Pcontract`. The same cases run once on production
|
||||
LavinMQ, or the release states that LavinMQ was not checked.
|
||||
16. Member incidents route to the exact delegating lead and never resolve a task rendezvous.
|
||||
17. `fleet_list` shows compact member health and capacity without pane content. Member
|
||||
17. `bridge_list` shows compact member health and capacity without pane content. Member
|
||||
`idleForSeconds` is present only when no accepted turn or inbox item exists.
|
||||
|
||||
### Unit 4 - Lead health and peer routing
|
||||
@@ -887,7 +887,7 @@ Acceptance criteria:
|
||||
8. A selected working peer is not interrupted. Its push waits for an injectable window.
|
||||
9. Recipient assignment stays pinned. Reassignment increments generation and supersedes old pending
|
||||
assignment.
|
||||
10. `fleet_list` shows bounded foreign assignments, recipient, reason, and generation without pane
|
||||
10. `bridge_list` shows bounded foreign assignments, recipient, reason, and generation without pane
|
||||
content.
|
||||
11. `LeadInboxRegistry` owns configured and discovered lead keys before publication.
|
||||
12. Missing leads keep ownership. Retirement needs an empty queue and handled incidents.
|
||||
@@ -911,7 +911,7 @@ Acceptance criteria:
|
||||
3. Webhook mode requires a resolved environment value. Bad notification config does not disable an
|
||||
already valid detector.
|
||||
4. Mode changes keep open incidents. Re-enable resumes eligible incidents.
|
||||
5. `fleet_list`, `/healthz`, metrics, and one WARN show partial coverage without a sink. HTTP
|
||||
5. `bridge_list`, `/healthz`, metrics, and one WARN show partial coverage without a sink. HTTP
|
||||
liveness behavior stays unchanged.
|
||||
6. One-lead, no-sink coverage states that lead failure has no active notification or recovery.
|
||||
7. Incident and outbound dedupe use the stable keys in Section 10.3.
|
||||
|
||||
+13
-13
@@ -220,17 +220,17 @@ sequenceDiagram
|
||||
participant H as herdr
|
||||
participant W as Worker (Claude)
|
||||
|
||||
P->>B: fleet_send("do X", target=w) — blocks
|
||||
P->>B: bridge_send("do X", target=w) — blocks
|
||||
B->>B: register waiter(w)
|
||||
B->>H: agent.send(w, "do X") (idle window)
|
||||
H-->>W: prompt injected
|
||||
W->>W: works the turn
|
||||
W->>B: fleet_reply("result")
|
||||
W->>B: bridge_reply("result")
|
||||
B->>B: resolve waiter(w)
|
||||
B-->>P: { outcome:"reply", text:"result" }
|
||||
```
|
||||
|
||||
### 6.2 Clarification — reverse rendezvous (`fleet_ask`)
|
||||
### 6.2 Clarification — reverse rendezvous (`bridge_ask`)
|
||||
|
||||
The worker pauses mid-turn to ask; the primary answers; the worker resumes in the same turn.
|
||||
|
||||
@@ -240,14 +240,14 @@ sequenceDiagram
|
||||
participant B as bridged
|
||||
participant W as Worker
|
||||
|
||||
P->>B: fleet_send("do X", target=w) — blocks
|
||||
P->>B: bridge_send("do X", target=w) — blocks
|
||||
B-->>W: "do X" (injected)
|
||||
W->>B: fleet_ask("which config?") — worker blocks
|
||||
W->>B: bridge_ask("which config?") — worker blocks
|
||||
B-->>P: { outcome:"question", text:"which config?", turn_id }
|
||||
P->>B: fleet_send("config.yaml", target=w, turn_id) — blocks again
|
||||
B-->>W: resolve fleet_ask → { answer:"config.yaml" }
|
||||
P->>B: bridge_send("config.yaml", target=w, turn_id) — blocks again
|
||||
B-->>W: resolve bridge_ask → { answer:"config.yaml" }
|
||||
W->>W: resumes same turn
|
||||
W->>B: fleet_reply("done")
|
||||
W->>B: bridge_reply("done")
|
||||
B-->>P: { outcome:"reply", text:"done" }
|
||||
```
|
||||
|
||||
@@ -261,10 +261,10 @@ sequenceDiagram
|
||||
participant B as bridged
|
||||
participant W as Worker
|
||||
|
||||
P->>B: fleet_send("do X", target=w, block=false)
|
||||
P->>B: bridge_send("do X", target=w, block=false)
|
||||
B-->>P: { outcome:"dispatched", dispatch_id }
|
||||
P->>P: continues its own work
|
||||
W->>B: fleet_reply("result")
|
||||
W->>B: bridge_reply("result")
|
||||
Note over B: no waiter → detached path
|
||||
B->>B: Injector.enqueue(primary_pane, "result")
|
||||
B-->>P: injected into idle pane (status-gated)
|
||||
@@ -272,7 +272,7 @@ sequenceDiagram
|
||||
|
||||
### 6.4 Uncooperative worker — turn-done fallback
|
||||
|
||||
A worker that never calls `fleet_reply` still returns a result: `bridged` reads its terminal
|
||||
A worker that never calls `bridge_reply` still returns a result: `bridged` reads its terminal
|
||||
tail when the turn completes.
|
||||
|
||||
```mermaid
|
||||
@@ -281,9 +281,9 @@ sequenceDiagram
|
||||
participant B as bridged
|
||||
participant W as Worker
|
||||
|
||||
P->>B: fleet_send("do X", target=w) — blocks
|
||||
P->>B: bridge_send("do X", target=w) — blocks
|
||||
B-->>W: "do X" (injected)
|
||||
W->>W: works, never calls fleet_reply
|
||||
W->>W: works, never calls bridge_reply
|
||||
B->>B: StatusPoller sees agent_status → idle/done
|
||||
B->>B: AgentControl.read(w, "recent")
|
||||
B-->>P: { outcome:"turn_done", text:<terminal tail> }
|
||||
|
||||
@@ -19,12 +19,12 @@ of this one.
|
||||
## One gateway for all messages
|
||||
|
||||
- **Everyone talks through the same door.** The leader and every worker connect to the same
|
||||
MCP server and use only its tools: `fleet_whoami` · `fleet_profiles` · `fleet_spawn` ·
|
||||
`fleet_list` · `fleet_status` · `fleet_send` · `fleet_reply` · `fleet_ask` ·
|
||||
`fleet_poll` · `fleet_ack` · `fleet_stop`.
|
||||
MCP server and use only its tools: `bridge_whoami` · `bridge_profiles` · `bridge_spawn` ·
|
||||
`bridge_list` · `bridge_status` · `bridge_send` · `bridge_reply` · `bridge_ask` ·
|
||||
`bridge_poll` · `bridge_ack` · `bridge_stop`.
|
||||
- **You are who your connection says you are.** The bridge finds out who is calling from the
|
||||
connection itself, never from a name the caller sends. So a worker cannot pretend to be
|
||||
someone else, and `fleet_whoami` tells each agent its own role — no guessing.
|
||||
someone else, and `bridge_whoami` tells each agent its own role — no guessing.
|
||||
- **The subscription line cannot be crossed.** Only a spawned worker gets
|
||||
`ANTHROPIC_BASE_URL`; the leader never does. Each worker profile has a list of allowed
|
||||
model hosts, checked before anything starts.
|
||||
@@ -59,7 +59,7 @@ had to ask for its replies. That gap is now closed on a single machine:
|
||||
- **The leader gets a tap on the shoulder.** When a reply lands, the bridge nudges the
|
||||
leader's own pane — only when the leader is free, and only a few times. If the leader is on
|
||||
another machine, this quietly falls back to pick-up mode; the reply still waits.
|
||||
- **Workers can ask questions.** With `fleet_ask`, a worker can pause mid-task, ask the
|
||||
- **Workers can ask questions.** With `bridge_ask`, a worker can pause mid-task, ask the
|
||||
leader something, and continue the *same* task with the answer.
|
||||
|
||||
## More than one kind of worker
|
||||
|
||||
@@ -20,7 +20,7 @@ requirement, not a nice-to-have.**
|
||||
- **Worktree provisioned by the daemon** — `SessionManager` creates a dedicated git worktree +
|
||||
branch per session, **hydrates it to full config parity** (below), and tears it down on release.
|
||||
- **Worker opens its own PR** — the worker commits, pushes its branch, and opens the PR/MR itself,
|
||||
returning the PR URL in its `fleet_reply`.
|
||||
returning the PR URL in its `bridge_reply`.
|
||||
|
||||
## Why worktrees (the hazard being fixed)
|
||||
|
||||
@@ -103,7 +103,7 @@ sequenceDiagram
|
||||
Note over W: implement in the isolated worktree
|
||||
W->>G: git commit + git push (SSH, same user)
|
||||
W->>G: open PR (branch to main)
|
||||
W-->>P: fleet_reply (prUrl, branch, summary, tests)
|
||||
W-->>P: bridge_reply (prUrl, branch, summary, tests)
|
||||
P->>SM: release(paneId)
|
||||
SM->>G: git worktree remove wt
|
||||
Note over G: branch + PR persist for review/merge
|
||||
@@ -170,7 +170,7 @@ A worker-facing playbook (sibling to the existing `reviewer` skill):
|
||||
3. **Push** your branch (`git push -u origin HEAD`).
|
||||
4. **Open a PR** to `main` (option A `curl`, or the decided mechanism) with a title/body describing
|
||||
the change and referencing the ticket.
|
||||
5. **Reply** via `fleet_reply` with the **PR URL**, branch name, files changed, and test names —
|
||||
5. **Reply** via `bridge_reply` with the **PR URL**, branch name, files changed, and test names —
|
||||
that reply is the whole handoff.
|
||||
6. Do **not** merge; do **not** touch `.mcp.json` or `wiki/`.
|
||||
|
||||
|
||||
@@ -14,7 +14,7 @@ rule that **a worker inherits the primary's directory** (never `$HOME`), and how
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
A["fleet_spawn / POST /workers"] --> B{"explicit cwd?<br/>(profile cwd or spawn arg)"}
|
||||
A["bridge_spawn / POST /workers"] --> B{"explicit cwd?<br/>(profile cwd or spawn arg)"}
|
||||
B -->|"yes — told otherwise"| C["use that cwd"]
|
||||
B -->|"no"| D{"caller PID resolvable?<br/>(MCP peer PID)"}
|
||||
D -->|"yes"| E["cwd = the primary's cwd<br/>lsof -a -p PID -d cwd"]
|
||||
@@ -51,7 +51,7 @@ only affect the seed shell, which the bridge closes).
|
||||
| # | Source | When |
|
||||
|---|--------|------|
|
||||
| 1 | Explicit `cwd` — a per-profile `cwd:` in config, or a spawn argument | "told otherwise" — pin a fixed workdir |
|
||||
| 2 | The **primary's cwd**, auto-detected from the `fleet_spawn` caller | normal MCP spawn from the primary |
|
||||
| 2 | The **primary's cwd**, auto-detected from the `bridge_spawn` caller | normal MCP spawn from the primary |
|
||||
| 3 | The `bridged` daemon's own cwd | REST spawn / off-host caller — **never `$HOME`** |
|
||||
|
||||
The primary's cwd (source 2) is discoverable with no new plumbing: `bridged` already resolves the MCP
|
||||
@@ -65,7 +65,7 @@ sequenceDiagram
|
||||
participant B as "bridged"
|
||||
participant O as "OS (lsof)"
|
||||
participant H as "herdr"
|
||||
P->>B: "fleet_spawn {profile} (no cwd)"
|
||||
P->>B: "bridge_spawn {profile} (no cwd)"
|
||||
B->>O: "peer PID for this connection's port"
|
||||
O-->>B: "pid"
|
||||
B->>O: "cwd of pid (lsof -d cwd)"
|
||||
@@ -79,7 +79,7 @@ sequenceDiagram
|
||||
|
||||
> **Status:** implemented (CB-112). `bridged` threads the resolved `cwd` onto **`agent.start {cwd}`**
|
||||
> (verified: the worker process is rooted there), keeping the single shared worker space. On an MCP
|
||||
> `fleet_spawn` the primary's cwd is auto-detected from the caller's PID; over REST (no MCP caller)
|
||||
> `bridge_spawn` the primary's cwd is auto-detected from the caller's PID; over REST (no MCP caller)
|
||||
> it is the explicit `cwd` param else the daemon's cwd. Both placements (`tab` and legacy `pane`)
|
||||
> carry it, since it rides `agent.start`.
|
||||
|
||||
@@ -133,5 +133,5 @@ unattended.
|
||||
|
||||
## See also
|
||||
|
||||
- `docs/MCP-Contract.md` — the tool surface (`fleet_spawn`, `fleet_profiles`, …).
|
||||
- `docs/MCP-Contract.md` — the tool surface (`bridge_spawn`, `bridge_profiles`, …).
|
||||
- `wiki/2-Message-Server.md` — the herdr `agent.*` / `workspace.*` schema (`workspace.create {cwd}`).
|
||||
|
||||
+3
-3
@@ -6,7 +6,7 @@ running `bridged` daemon**, captures the full transcript, and grades the channel
|
||||
|
||||
This is the committed form of the ad-hoc channel test that discovered the CB-115 gaps
|
||||
(herdr `unknown` misclassification wedging delivery, dirty completion scrapes, and workers
|
||||
never calling `fleet_reply` in conversation). Run it after any change to the injector,
|
||||
never calling `bridge_reply` in conversation). Run it after any change to the injector,
|
||||
status handling, completion/failure paths, or the worker reply charter.
|
||||
|
||||
## What it exercises
|
||||
@@ -28,7 +28,7 @@ sequenceDiagram
|
||||
T->>B: POST /sessions/{id}/message {wait:false}
|
||||
B-->>T: ticket
|
||||
B->>W: inject prompt (status-gated)
|
||||
W-->>B: fleet_reply
|
||||
W-->>B: bridge_reply
|
||||
T->>B: GET /tasks/{ticket} (poll)
|
||||
B-->>T: done + reply
|
||||
end
|
||||
@@ -68,7 +68,7 @@ Per-turn grade:
|
||||
|
||||
| Grade | Meaning |
|
||||
|------------|---------------------------------------------------------------------|
|
||||
| `OK` | delivered and resolved by an explicit `fleet_reply` (`source=reply`) |
|
||||
| `OK` | delivered and resolved by an explicit `bridge_reply` (`source=reply`) |
|
||||
| `DEGRADED` | delivered and answered, but resolved via completion-scrape fallback |
|
||||
| `EMPTY` | turn completed but the reply was empty |
|
||||
| `FAILED` | the worker's turn ended in failure (`phase=failed`) |
|
||||
|
||||
Reference in New Issue
Block a user