Compare commits

...

1 Commits

Author SHA1 Message Date
Dai Ha 41a3114d03 CB-622 Unit A: register fleet_* MCP tools, keep bridge_* working
CI / build (pull_request) Successful in 1m2s
CI / contract (pull_request) Successful in 1m8s
fleet_* is now the documented tool name for all eleven MCP tools; each
bridge_* twin is registered against the exact same handler (no logic
duplication) and its description leads with a DEPRECATED notice. A
bridge_* call logs one WARN naming the old and new name, once per name
for the life of the process (a Set, not a numeric sentinel).

REPLY_CHARTER in HerdrPeerLauncher now tells a spawned member to call
fleet_reply — the one rule that must survive with no repo checkout.

All other bridge_* string literals across mcp/, Javadoc, and tests were
renamed to fleet_* for consistency with the new documented name.
2026-08-22 21:54:32 +02:00
34 changed files with 496 additions and 288 deletions
@@ -381,7 +381,7 @@ public final class Bridged {
+ "has lost the delegation map. Its pushReminders/pushBackoffMs stay valid."); + "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 // 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; int maxReminders = cfg.primary() != null ? cfg.primary().remindersOrDefault() : 5;
long backoffMs = cfg.primary() != null ? cfg.primary().backoffMsOrDefault() : 15_000L; long backoffMs = cfg.primary() != null ? cfg.primary().backoffMsOrDefault() : 15_000L;
var pushScheduler = Executors.newSingleThreadScheduledExecutor(r -> 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 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. // Caller identity is resolved from the connection (peer PID → herdr pane), not arguments.
ConnectionIdentity identity = new ConnectionIdentity( ConnectionIdentity identity = new ConnectionIdentity(
new PaneLocator(herdr), new LsofPeerPidLookup(), new LsofProcessCwdLookup()); 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. * 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 * <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. * spawn lifecycle binds a slot.
*/ */
public Map<String, String> snapshot() { 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, * <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 * 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. * 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 * <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 * 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 * 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 * 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 * 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. * {@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. * 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 * <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 * (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 * 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 * {@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) * and {@code {n}} (per-worker number, to keep sibling tabs distinct)
* are substituted (default {@code "worker: {profile} #{n}"}) * are substituted (default {@code "worker: {profile} #{n}"})
* @param mcpUrl bridge MCP URL to provision into the worker's {@code configDir} so it * @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) * 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"); * @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 * {@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 * auto-select this profile": it is excluded from every automatic policy's
* pool the same way a quarantined candidate is (see * pool the same way a quarantined candidate is (see
* {@code PlacementPolicyUtil}). This does not make the profile * {@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 * placement entirely and still resolves it. Weights among the remaining
* (non-excluded) candidates need not sum to 1.0; only their ratios matter. * (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 * @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 * the same way a {@code weight <= 0} profile is (see
* {@code PlacementPolicyUtil.available()}, which already treats "at cap" * {@code PlacementPolicyUtil.available()}, which already treats "at cap"
* and "excluded" alike), and an explicit * 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 * {@code CompositePeerLauncher.enforceMaxLoad}) — a cap is a capacity
* statement that does not stop being true just because the profile was * statement that does not stop being true just because the profile was
* named directly. A negative value has no sane meaning (there is no * 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 * {@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. * 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 * @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. * having refused on a subscription usage limit, rather than a real answer.
* {@code null}/blank ⇒ the classification never fires for this profile and * {@code null}/blank ⇒ the classification never fires for this profile and
* today's completion-fallback behaviour is unchanged. Every backend words * 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. * 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 * <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 * <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 * 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 * 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 * <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 * 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; 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. * MUTE is an observation by target and profile, not a classifier state and never suppresses faults.
*/ */
public final class MuteCounter { 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 * 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 * {@link Rendezvous} so a blocking {@code fleet_send} resolves even when the worker finishes its
* task without ever calling {@code bridge_reply} — the common case for a real delegated coding task. * 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 * <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 * 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 * 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 * — 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 * <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 * 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 * <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 * {@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 * <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 * <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 * 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. * 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; static final int MAX_SCRAPE_CHARS = 4000;
private static final String CLIPPED_PANE_TAIL_MARKER = 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 AgentControl agents;
private final Rendezvous rendezvous; private final Rendezvous rendezvous;
@@ -171,7 +171,7 @@ public final class CompletionResolver implements TurnListener {
void resolve(String target, InFlight turn) { void resolve(String target, InFlight turn) {
CompletableFuture<Rendezvous.Resolution> waiter = turn == null ? null : turn.waiter(); CompletableFuture<Rendezvous.Resolution> waiter = turn == null ? null : turn.waiter();
if (waiter == null || waiter.isDone()) { 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. // the scrape; resolving the current waiter here would be the CB-116 cross-turn stale reply.
inFlight.remove(target, turn); inFlight.remove(target, turn);
return; 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 // 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 // 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 // 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". // A scrape that failed to read is exempt — an empty tail there is "couldn't see", not "no change".
String baseline = turn.baseline(); String baseline = turn.baseline();
if (!scrapeFailed && baseline != null && baseline.equals(tail)) { if (!scrapeFailed && baseline != null && baseline.equals(tail)) {
@@ -205,7 +205,7 @@ public final class CompletionResolver implements TurnListener {
target); target);
return; // keep the in-flight record: a later genuine completion still needs it 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'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. // BACKEND_EXHAUSTED rather than handing the caller a scrape that reads like a real reply.
if (!scrapeFailed) { if (!scrapeFailed) {
@@ -215,7 +215,7 @@ public final class CompletionResolver implements TurnListener {
String reason = "backend exhausted (usage limit): " + matchedLine; String reason = "backend exhausted (usage limit): " + matchedLine;
if (rendezvous.resolveExhausted(waiter, reason)) { if (rendezvous.resolveExhausted(waiter, reason)) {
inFlight.remove(target, turn); 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); + "matched the profile's exhausted pattern): {}", target, reason);
// CB-578 stage B: only on the resolution that actually won the race — a late // 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. // duplicate must never quarantine a credential twice for one refusal.
@@ -229,7 +229,7 @@ public final class CompletionResolver implements TurnListener {
inFlight.remove(target, turn); inFlight.remove(target, turn);
if (clipped) { if (clipped) {
log.warn("completion scrape for {} clipped from {} chars to the {} char cap; " 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); target, originalLength, MAX_SCRAPE_CHARS);
} }
log.debug("resolved send to {} via turn-completion fallback ({} chars scraped)", 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 * 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 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 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. * layer and stays unit-testable with a capturing fake.
*/ */
@FunctionalInterface @FunctionalInterface
@@ -26,7 +26,7 @@ import java.util.stream.Collectors;
* must never receive: * must never receive:
* <ol> * <ol>
* <li>it appends the <em>reply charter</em> — "you are an off-subscription worker … end every turn * <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> * exactly backwards.</li>
* <li>it registers the session with {@code SessionManager}, which subjects it to the idle reaper, * <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 * 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.Map;
import java.util.Objects; import java.util.Objects;
import java.util.Set; import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.function.BiFunction;
import java.util.function.Function; import java.util.function.Function;
import java.util.function.LongSupplier; import java.util.function.LongSupplier;
import java.util.function.Supplier; import java.util.function.Supplier;
import java.util.stream.Collectors; 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> * 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 * 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} * validated by parity, not by re-implementing behaviour. The primary Opus calls {@code fleet_send}
* / {@code bridge_status}; the worker calls {@code bridge_reply}. * / {@code fleet_status}; the worker calls {@code fleet_reply}.
* *
* <p>Beyond delegation the primary also manages the fleet here (CB-108): {@code bridge_spawn} / * <p>Beyond delegation the primary also manages the fleet here (CB-108): {@code fleet_spawn} /
* {@code bridge_list} / {@code bridge_stop} drive the {@link PeerLauncher} SPI so a worker's whole * {@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. * 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 * <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 { 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 DEFAULT_TIMEOUT_MS = 25_000;
private static final long MAX_TIMEOUT_MS = 120_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). // 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_DEFAULT_TIMEOUT_MS = 55_000;
private static final long ASK_MAX_TIMEOUT_MS = 115_000; private static final long ASK_MAX_TIMEOUT_MS = 115_000;
private static final ObjectMapper MAPPER = new ObjectMapper(); // worker-view JSON projections 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. */ /** Transport-context key under which the extractor stashes the resolved caller identity. */
static final String CALLER_TERMINAL = "callerTerminal"; static final String CALLER_TERMINAL = "callerTerminal";
/** Transport-context key under which the extractor stashes the caller's PID (for cwd inherit). */ /** 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 HealthCoverageSource healthCoverage;
private final QuarantineSource quarantine; 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, public record CapacitySource(Function<String, Integer> liveCount, Function<String, Integer> maxLoad,
Supplier<Set<String>> configuredProfiles, LongSupplier clock) { Supplier<Set<String>> configuredProfiles, LongSupplier clock) {
/** Inert test-only source. It omits capacity rather than inventing zero live counts. */ /** 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) { } 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. * lookup, plus the shared {@link BackendQuarantine} to read remaining cooldowns off.
*/ */
public record QuarantineSource(Function<String, String> credentialIdFor, BackendQuarantine quarantine) { 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} * Jetty's context handler and never passes through Javalin's {@code before}
* filter, so the REST guard does not cover it. * filter, so the REST guard does not cover it.
* @param metrics registry for auth-failure counting; may be {@code null} * @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 * {@link QuarantineSource#none()} for a caller that does not want the feature
*/ */
public BridgeMcp(MessageService messages, PeerLauncher workers, SessionManager sessions, public BridgeMcp(MessageService messages, PeerLauncher workers, SessionManager sessions,
@@ -124,7 +139,7 @@ public final class BridgeMcp {
.jsonMapper(json) .jsonMapper(json)
.mcpEndpoint("/mcp") .mcpEndpoint("/mcp")
// Resolve the caller from the connection (peer PID → herdr pane) in one lookup: the // 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 // 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. // (CB-113) — its MCP initialize is the reliable "the agent is up" signal.
.contextExtractor(req -> { .contextExtractor(req -> {
@@ -145,10 +160,10 @@ public final class BridgeMcp {
CALLER_NAME, orEmpty(p.name()))); CALLER_NAME, orEmpty(p.name())));
}) })
.build(); .build();
this.server = McpServer.sync(transport) // CB-622: each handler is built once and reused for BOTH its fleet_* tool and its
.serverInfo("bridge", "0.1.0") // deprecated bridge_* twin (registered below), so the two names can never drift apart.
.capabilities(McpSchema.ServerCapabilities.builder().tools(true).build()) BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> sendHandler =
.toolCall(sendTool(), (exchange, req) -> { (exchange, req) -> {
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.SEND, McpSchema.CallToolResult denied = deny(exchange, Authz.Action.SEND,
str(req.arguments(), "sessionId")); str(req.arguments(), "sessionId"));
if (denied != null) return denied; if (denied != null) return denied;
@@ -163,7 +178,7 @@ public final class BridgeMcp {
String content = str(a, "content"); String content = str(a, "content");
String turnId = str(a, "turnId"); String turnId = str(a, "turnId");
if (turnId != null && !turnId.isBlank()) { 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 // 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. // delegation, so ownership is left untouched (CB-548) — never re-recorded.
return answer(messages, turnId, content, timeoutMs(a)); return answer(messages, turnId, content, timeoutMs(a));
@@ -178,43 +193,49 @@ public final class BridgeMcp {
return Boolean.FALSE.equals(a.get("wait")) return Boolean.FALSE.equals(a.get("wait"))
? sendAsync(messages, target, content, onAccepted, workers.profiles()) ? sendAsync(messages, target, content, onAccepted, workers.profiles())
: send(messages, target, content, timeoutMs(a), 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 // 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. // is "is this caller a worker at all", and it can only ever reply as itself.
.toolCall(replyTool(), (exchange, req) -> { BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> replyHandler =
(exchange, req) -> {
String self = callerTerminal(exchange); String self = callerTerminal(exchange);
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.REPLY, self); McpSchema.CallToolResult denied = deny(exchange, Authz.Action.REPLY, self);
if (denied != null) return denied; if (denied != null) return denied;
return reply(messages, self, str(req.arguments(), "content")); return reply(messages, self, str(req.arguments(), "content"));
}) };
// bridge_ask (CB-205): a worker's mid-turn question — identity from the CONNECTION. // fleet_ask (CB-205): a worker's mid-turn question — identity from the CONNECTION.
.toolCall(askTool(), (exchange, req) -> { BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> askHandler =
(exchange, req) -> {
String self = callerTerminal(exchange); String self = callerTerminal(exchange);
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.ASK, self); McpSchema.CallToolResult denied = deny(exchange, Authz.Action.ASK, self);
if (denied != null) return denied; if (denied != null) return denied;
return ask(messages, self, str(req.arguments(), "question"), timeoutMs(req.arguments())); 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); McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
if (denied != null) return denied; if (denied != null) return denied;
return status(messages, str(req.arguments(), "sessionId")); 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); McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
if (denied != null) return denied; if (denied != null) return denied;
Map<String, Object> a = req.arguments(); Map<String, Object> a = req.arguments();
return poll(messages, str(a, "ticket"), str(a, "target")); 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). // 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. // Acking removes a reply from the inbox, so it is a drain, not a read.
.toolCall(ackTool(), (exchange, req) -> { BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> ackHandler =
(exchange, req) -> {
Map<String, Object> a = req.arguments(); Map<String, Object> a = req.arguments();
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.DRAIN, str(a, "target")); McpSchema.CallToolResult denied = deny(exchange, Authz.Action.DRAIN, str(a, "target"));
if (denied != null) return denied; if (denied != null) return denied;
return ack(messages, str(a, "target"), str(a, "msgId")); return ack(messages, str(a, "target"), str(a, "msgId"));
}) };
// Fleet management (CB-108): spawn/list/stop over ClaudeCodeLauncher. // Fleet management (CB-108): spawn/list/stop over ClaudeCodeLauncher.
.toolCall(spawnTool(), (exchange, req) -> { BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> spawnHandler =
(exchange, req) -> {
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.SPAWN, null); McpSchema.CallToolResult denied = deny(exchange, Authz.Action.SPAWN, null);
if (denied != null) return denied; if (denied != null) return denied;
String caller = callerTerminal(exchange); 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, return spawn(sessions, str(a, "profile"), str(a, "role"), str(a, "cwd"), callerCwd,
callerTerminal(exchange), worktreeRequest(a), callerTerminal(exchange), worktreeRequest(a),
str(a, "sessionName"), str(a, "resumeSessionId")); 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); McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
if (denied != null) return denied; if (denied != null) return denied;
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine,
callers == null ? Map.of() : callers.leads(), callers == null ? Map.of() : callers.leads(),
callerTerminal(exchange)); callerTerminal(exchange));
}) };
.toolCall(stopTool(), (exchange, req) -> { BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> stopHandler =
(exchange, req) -> {
String paneId = str(req.arguments(), "paneId"); String paneId = str(req.arguments(), "paneId");
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.STOP, paneId); McpSchema.CallToolResult denied = deny(exchange, Authz.Action.STOP, paneId);
if (denied != null) return denied; if (denied != null) return denied;
return stop(sessions, paneId); return stop(sessions, paneId);
}) };
.toolCall(profilesTool(), (exchange, _) -> { BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> profilesHandler =
(exchange, _) -> {
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null); McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
if (denied != null) return denied; if (denied != null) return denied;
return profiles(workers, quarantine); return profiles(workers, quarantine);
}) };
.toolCall(whoamiTool(), (exchange, _) -> { BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> whoamiHandler =
(exchange, _) -> {
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null); McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
if (denied != null) return denied; if (denied != null) return denied;
return whoami(principal(exchange), sessions); 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(); .build();
this.authz = callers; this.authz = callers;
this.metrics = metrics; this.metrics = metrics;
@@ -406,7 +469,7 @@ public final class BridgeMcp {
// --- tool logic (thin adapters over the services; unit-testable) --------------------------- // --- 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. * 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 * (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 fleet_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_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. * 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) { 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}), * 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. * 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) { static McpSchema.CallToolResult ask(MessageService messages, String callerTerminal, String question, Long timeoutMs) {
if (callerTerminal == null) { 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"); + "from the connection");
} }
if (isBlank(question)) { if (isBlank(question)) {
@@ -459,10 +522,10 @@ public final class BridgeMcp {
MessageService.AskResult r = messages.ask(callerTerminal, question, timeout); MessageService.AskResult r = messages.ask(callerTerminal, question, timeout);
return switch (r.outcome()) { return switch (r.outcome()) {
case ANSWERED -> text(r.answer()); case ANSWERED -> text(r.answer());
case NO_WAITER -> error("no primary is awaiting this turn — bridge_ask only works while a " case NO_WAITER -> error("no primary is awaiting this turn — fleet_ask only works while a "
+ "bridge_send delegation is open to answer it"); + "fleet_send delegation is open to answer it");
case TIMED_OUT -> text("[no answer within " + timeout + "ms — the primary did not respond; " 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) { private static McpSchema.CallToolResult formatReply(MessageService.Reply r, long timeout) {
return switch (r.outcome()) { return switch (r.outcome()) {
case REPLIED -> text(r.text()); 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. // transcript tail, flagged so the primary knows it isn't a structured reply.
case COMPLETED_UNREPLIED -> text( 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. // 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()); 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 // 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()); + "usage limit]\n" + r.text());
// The worker paused mid-turn to ask (CB-205) — tell the primary how to answer in-turn. // 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() 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."); + "\" 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 " case STALE_TURN -> error("that question is no longer open — it timed out or was already "
+ "answered (turnId stale)"); + "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. * 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. * The configured profiles are required so a profile name can never bypass target validation.
* *
@@ -510,19 +573,19 @@ public final class BridgeMcp {
return targetError; return targetError;
} }
String ticket = messages.sendAsync(sessionId, content, onAccepted); 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. */ /** 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) { private static McpSchema.CallToolResult profileTargetError(String sessionId, Set<String> profiles) {
if (profiles.contains(sessionId)) { if (profiles.contains(sessionId)) {
return error("unknown send target \"" + sessionId + "\": it is a configured profile name, not a " 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; 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) { static McpSchema.CallToolResult poll(MessageService messages, String ticket, String target) {
if (!isBlank(target)) { if (!isBlank(target)) {
var replies = messages.drainReplies(target); var replies = messages.drainReplies(target);
@@ -540,25 +603,25 @@ public final class BridgeMcp {
} }
return switch (v.phase()) { return switch (v.phase()) {
case DONE -> text(v.replySource() != null && v.replySource().equals("transcript") 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()); : v.reply());
case PENDING -> text("[pending — " + v.detail() + "]"); case PENDING -> text("[pending — " + v.detail() + "]");
case ASKING -> text("[question — worker is waiting for your answer]\n" + v.reply() 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."); + "\" and content set to your answer; the worker resumes the same turn.");
case FAILED -> text("[failed — " + v.detail() + "]"); 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). * 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} * {@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). * 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) { static McpSchema.CallToolResult reply(MessageService messages, String callerTerminal, String content) {
if (callerTerminal == null) { 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"); + "from the connection");
} }
if (content == null) { if (content == null) {
@@ -568,7 +631,7 @@ public final class BridgeMcp {
return text("delivered"); 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) { static McpSchema.CallToolResult ack(MessageService messages, String target, String msgId) {
if (isBlank(target) || isBlank(msgId)) { if (isBlank(target) || isBlank(msgId)) {
return error("target and msgId are required"); 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 * {@code fleet_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 * 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. * 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) { static McpSchema.CallToolResult status(MessageService messages, String sessionId) {
@@ -593,7 +656,7 @@ public final class BridgeMcp {
return text(base); return text(base);
} }
return text(base + "\n\n[question — worker is waiting for your answer]\n" + ask.question() 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." + "\" and content set to your answer; the worker resumes the same turn."
+ " (ticket " + ask.ticket() + ")"); + " (ticket " + ask.ticket() + ")");
} catch (HerdrException e) { } 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 * <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 * 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 * 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 * 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 * 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. * 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 * <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) -------------------------------------------- // --- 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) { static McpSchema.CallToolResult spawn(SessionManager sessions, String profile) {
return spawn(sessions, profile, null, null, null, null, null, null, null); 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 * 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 * member's cwd is {@code requestedCwd} if given, else the profile's config, else
* {@code callerCwd} (the primary's directory), else the daemon's. * {@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) { private static WorktreeRequest worktreeRequest(Map<String, Object> a) {
Object w = a.get("worktree"); Object w = a.get("worktree");
if (w == null || Boolean.FALSE.equals(w)) { 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 * 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 * {@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. * 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 * 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. * 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 * <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 * 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 * 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} * <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}. * quarantined — see {@link #capacityView}.
* *
* @param leads terminal_id → lead name, live from the resolver * @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 * 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 * 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 * 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 * 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 * added only when the profile is actually quarantined, so an ordinary fleet's rows are
* byte-identical to before this change. * 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 * <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 * 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. * 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, private static Map<String, Object> leadView(String terminal, String name, Agent live,
@@ -905,7 +968,7 @@ public final class BridgeMcp {
return m; 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) { static McpSchema.CallToolResult stop(SessionManager sessions, String paneId) {
if (isBlank(paneId)) { if (isBlank(paneId)) {
return error("paneId is required"); return error("paneId is required");
@@ -948,25 +1011,25 @@ public final class BridgeMcp {
// --- tool schemas -------------------------------------------------------------------------- // --- tool schemas --------------------------------------------------------------------------
private static McpSchema.Tool sendTool() { 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 " "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 " + "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 " + "for a long task to return a ticket immediately, then poll it with fleet_poll. To "
+ "answer a worker's bridge_ask, pass its turnId (with content) instead of sessionId.", + "answer a worker's fleet_ask, pass its turnId (with content) instead of sessionId.",
objectSchema(Map.of( objectSchema(Map.of(
"sessionId", stringProp("The worker session id (herdr terminal_id) to delegate to"), "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)"), "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)"), "timeoutMs", Map.of("type", "integer", "description", "Max ms to wait for a reply (blocking mode)"),
"wait", Map.of("type", "boolean", "wait", Map.of("type", "boolean",
"description", "Block for the reply (default true); false returns a ticket to poll"), "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)")), + "routes your answer back into the same turn (omit for a normal delegation)")),
List.of("content"))); List.of("content")));
} }
private static McpSchema.Tool askTool() { private static McpSchema.Tool askTool() {
// No target/session arg — the worker's identity is resolved from the connection. // 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 " "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 " + "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 " + "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() { private static McpSchema.Tool pollTool() {
return tool("bridge_poll", return tool("fleet_poll",
"Check an async delegation (a bridge_send with wait:false) by its ticket: " "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 " + "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 " + "session id) is present instead of ticket, drain that worker's inbox of "
+ "replies delivered when no send was open.", + "replies delivered when no send was open.",
objectSchema(Map.of( 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)")), "target", stringProp("Worker session id to drain pending replies from (optional)")),
List.of())); List.of()));
} }
private static McpSchema.Tool ackTool() { 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 " "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 " + "has processed a reply and wants to confirm it, leaving other pending replies "
+ "in the inbox for later drain.", + "in the inbox for later drain.",
@@ -1002,22 +1065,22 @@ public final class BridgeMcp {
} }
private static McpSchema.Tool spawnTool() { 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: " "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 " + "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 " + "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 " + "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 " + "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 " + "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 " + "cwd to pin a different one. Pass worktree:true (with ticket) or "
+ "worktree:<ticket-slug> to provision an isolated git worktree. Pass resumeSessionId " + "worktree:<ticket-slug> to provision an isolated git worktree. Pass resumeSessionId "
+ "to relaunch onto a prior conversation instead of starting cold — this requires an " + "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. " + "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 " + "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 " + "one. Returns the member's sessionId (use with fleet_send) and paneId (use with "
+ "bridge_stop).", + "fleet_stop).",
objectSchema(Map.of( objectSchema(Map.of(
"role", stringProp("What the member is for: architect, dev or reviewer (default dev)"), "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)"), "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"), "worktree", Map.of("type", "string", "description", "'true' or a ticket slug — requests an isolated git worktree"),
"ticket", stringProp("Ticket slug when worktree:true"), "ticket", stringProp("Ticket slug when worktree:true"),
"sessionName", stringProp("Logical display name for the member's own session, when its backend supports one"), "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())); List.of()));
} }
private static McpSchema.Tool profilesTool() { private static McpSchema.Tool profilesTool() {
return tool("bridge_profiles", return tool("fleet_profiles",
"List the configured worker profiles (backends) and which one bridge_spawn uses by " "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 " + "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.", + "quarantinedForSeconds elapses; a profile sharing that credential is listed too.",
objectSchema(Map.of(), List.of())); objectSchema(Map.of(), List.of()));
} }
private static McpSchema.Tool listTool() { 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 " "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 " + "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 " + "discover a peer lead without being told its address. 'members' are the "
+ "sessions delegated to — each with sessionId, paneId, role (architect/dev/" + "sessions delegated to — each with sessionId, paneId, role (architect/dev/"
+ "reviewer), profile (the backend it runs on), state, optional " + "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 " + "resumeSessionId to relaunch onto that same conversation, when the backend "
+ "supports it), and live herdr status. An empty 'members' " + "supports it), and live herdr status. An empty 'members' "
+ "means no members are spawned; it says nothing about peers. When capacity " + "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 " + "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 " + "maxLoad/live — with credentialId and quarantinedForSeconds naming the "
+ "quarantine, so 'free: 0, busy' can be told apart from 'free: 0, refusing " + "quarantine, so 'free: 0, busy' can be told apart from 'free: 0, refusing "
+ "for N seconds'.", + "for N seconds'.",
@@ -1059,8 +1122,8 @@ public final class BridgeMcp {
} }
private static McpSchema.Tool stopTool() { private static McpSchema.Tool stopTool() {
return tool("bridge_stop", return tool("fleet_stop",
"Tear down a worker session by its paneId (from bridge_spawn or bridge_list).", "Tear down a worker session by its paneId (from fleet_spawn or fleet_list).",
objectSchema(Map.of( objectSchema(Map.of(
"paneId", stringProp("The worker's paneId to stop")), "paneId", stringProp("The worker's paneId to stop")),
List.of("paneId"))); List.of("paneId")));
@@ -1068,9 +1131,9 @@ public final class BridgeMcp {
private static McpSchema.Tool replyTool() { private static McpSchema.Tool replyTool() {
// No session/target arg — the caller's identity is resolved from the connection. // No session/target arg — the caller's identity is resolved from the connection.
return tool("bridge_reply", return tool("fleet_reply",
"Return your structured answer for a message you were sent, resolving the sender's " "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 " + "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.", + "it — never to answer a worker, whose turn it is not.",
objectSchema(Map.of( objectSchema(Map.of(
@@ -1079,7 +1142,7 @@ public final class BridgeMcp {
} }
private static McpSchema.Tool statusTool() { 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.", "Get the live lifecycle status (idle/working/blocked/unknown) of a worker session.",
objectSchema(Map.of( objectSchema(Map.of(
"sessionId", stringProp("The worker session id to query")), "sessionId", stringProp("The worker session id to query")),
@@ -1087,20 +1150,58 @@ public final class BridgeMcp {
} }
private static McpSchema.Tool whoamiTool() { 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 " "Report who YOU are on the bridge — your role is resolved from your connection "
+ "(unforgeable), never from anything you claim. Returns role 'primary' (you " + "(unforgeable), never from anything you claim. Returns role 'primary' (you "
+ "orchestrate: spawn/send/stop; reply ONLY to answer a peer lead that " + "orchestrate: spawn/send/stop; reply ONLY to answer a peer lead that "
+ "messaged you, never to answer a worker), 'architect' (you delegate turns " + "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' " + "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 " + "(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 " + "which one you are, your own sessionId, and profile/worktree/branch when "
+ "you are a worker. Call this first when following role-conditional " + "you are a worker. Call this first when following role-conditional "
+ "instructions rather than guessing.", + "instructions rather than guessing.",
objectSchema(Map.of(), List.of())); 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 ------------------------------------------------------------------------- // --- small helpers -------------------------------------------------------------------------
// The SDK 2.0.0 deprecates its own Tool builders without a stable replacement — isolate it here. // 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}. * Single-slot, thread-safe registry for the primary's herdr {@code terminal_id}.
* *
* <p>Populated from the caller terminal of orchestration-side MCP tools * <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 * A pinned terminal (from config) seeds the registry at construction and makes
* subsequent {@link #record(String)} calls no-ops. * 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 * 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: * 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 * 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. * 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 * <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 * 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 * 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 * 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. * 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); 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 // 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 // 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. // carries no role and so would be judged against the dev pool it was never meant for.
List<PlacementCandidate> candidates = candidates(req.role()); List<PlacementCandidate> candidates = candidates(req.role());
String roleDefault = defaultProfileFor(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. */ /** The final instruction always requires a bridge reply when the bridge MCP is mounted. */
protected static final String REPLY_CHARTER = protected static final String REPLY_CHARTER =
"You are a spawned member in the claude-bridge fleet. Every message you receive arrives " "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, " + "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 " + "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 " + "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 " + "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). * 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; : replyCharter == null ? roleCharter : roleCharter + "\n\n" + replyCharter;
// CB-571: fingerprint the exact composed charter bytes once, here in the base, before the // 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 // 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. // surface the byte count can appear on. The charter text itself is never logged.
CharterReceipt receipt = CharterReceipt.compose(role, cfg.profile(), roleCharter, charter); CharterReceipt receipt = CharterReceipt.compose(role, cfg.profile(), roleCharter, charter);
String cwd = resolveCwd(requestedCwd, cfg, callerCwd); 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 * 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 * 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 * 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 * "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 * own git worktree on its own branch, is off-subscription, and cannot merge — the lead is the
* gate. * 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 * 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 * <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 * 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. */ /** The nudge body, phrased for the two cases the heartbeat distinguishes. */
String nudgeText() { String nudgeText() {
StringBuilder sb = new StringBuilder( 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()) { if (hasPending()) {
sb.append(" The fleet has state to collect: ").append(pendingDetail()); sb.append(" The fleet has state to collect: ").append(pendingDetail());
} else { } else {
@@ -309,10 +309,10 @@ public final class LeadHeartbeatLoop {
.append(" pending collection"); .append(" pending collection");
if (!replyTargets.isEmpty()) { if (!replyTargets.isEmpty()) {
// Render each as the exact command so the lead can act without parsing: the nearest // 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(" (") sb.append(" (")
.append(String.join(", ", .append(String.join(", ",
replyTargets.stream().map(t -> "bridge_poll(target=" + t + ")").toList())) replyTargets.stream().map(t -> "fleet_poll(target=" + t + ")").toList()))
.append(")"); .append(")");
} }
sb.append(", ").append(doneSessions).append(" DONE session").append(doneSessions == 1 ? "" : "s") 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 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 * 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 — * 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}. * 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. */ /** Outcome of a blocking send. */
public enum Outcome { 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, 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. * {@code text} is the scraped transcript tail rather than a structured answer.
*/ */
COMPLETED_UNREPLIED, COMPLETED_UNREPLIED,
@@ -75,7 +75,7 @@ public final class MessageService {
*/ */
WORKER_FAILED, 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, * 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 — * 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 * so this is never reported as a completed reply, and is kept distinct from
@@ -96,14 +96,14 @@ public final class MessageService {
BUSY, BUSY,
/** /**
* An answer ({@link #answer(String, String, long)}) referenced a {@code turnId} that is no * 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 STALE_TURN
} }
/** /**
* @param outcome how the send ended (or paused) * @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 * for {@link Outcome#REPLIED}, a scraped transcript tail for
* {@link Outcome#COMPLETED_UNREPLIED}), or the question for {@link Outcome#QUESTION}, * {@link Outcome#COMPLETED_UNREPLIED}), or the question for {@link Outcome#QUESTION},
* else {@code null} * 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 { public enum AskOutcome {
/** The primary answered; {@link AskResult#answer} carries it. */ /** The primary answered; {@link AskResult#answer} carries it. */
ANSWERED, ANSWERED,
@@ -132,7 +132,7 @@ public final class MessageService {
TIMED_OUT 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) { public record AskResult(AskOutcome outcome, String answer) {
} }
@@ -140,7 +140,7 @@ public final class MessageService {
public enum Phase { public enum Phase {
/** Delegated and in flight — queued for the worker or being worked. */ /** Delegated and in flight — queued for the worker or being worked. */
PENDING, 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, ASKING,
/** The worker's turn finished; {@link TaskView#reply} holds the answer. */ /** The worker's turn finished; {@link TaskView#reply} holds the answer. */
DONE, DONE,
@@ -153,7 +153,7 @@ public final class MessageService {
* *
* @param reply the answer when {@link #phase} is {@link Phase#DONE}, or the question when * @param reply the answer when {@link #phase} is {@link Phase#DONE}, or the question when
* {@link #phase} is {@link Phase#ASKING}; otherwise {@code null} * {@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} * (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 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} * @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 * 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 * (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 * ({@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) { 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. */ /** Async task that owns each exact forward rendezvous waiter. */
private final ConcurrentHashMap<CompletableFuture<Rendezvous.Resolution>, Task> asyncTasksByWaiter = private final ConcurrentHashMap<CompletableFuture<Rendezvous.Resolution>, Task> asyncTasksByWaiter =
new ConcurrentHashMap<>(); 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 ConcurrentHashMap<String, Task> asyncTasksByTurn = new ConcurrentHashMap<>();
private final AtomicLong ticketSeq = new AtomicLong(); private final AtomicLong ticketSeq = new AtomicLong();
private final ExecutorService asyncExecutor = Executors.newThreadPerTaskExecutor( 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 * (CB-588) whenever an async ticket started by {@link #sendAsync} reaches a
* terminal phase, whenever {@link #poll} hands a terminal ticket to its caller, * 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 * 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, public MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous,
ReplyInbox inbox, ReplyPushLoop pushLoop) { 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 * 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. * 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 * {@link Rendezvous#resolveQuestion} must keep today's {@code NO_WAITER} behaviour — questions
* are interactive and must never be queued. * 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). * 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 * <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 * {@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 * 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 * (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 * 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 * the primary answers via {@link #answer} or {@code timeoutMillis} elapses. Identity is the
* worker's own session — it does not address the primary. * worker's own session — it does not address the primary.
* *
@@ -509,7 +509,7 @@ public final class MessageService {
rendezvous.closeAsk(ticket.turnId()); rendezvous.closeAsk(ticket.turnId());
return new AskResult(AskOutcome.NO_WAITER, null); // no primary is blocked on this worker 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 // (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 // (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 // 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); String answer = ticket.answer().get(timeoutMillis, TimeUnit.MILLISECONDS);
return new AskResult(AskOutcome.ANSWERED, answer); return new AskResult(AskOutcome.ANSWERED, answer);
} catch (TimeoutException e) { } 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); clearAsyncQuestion(ticket.turnId(), true);
return new AskResult(AskOutcome.TIMED_OUT, null); return new AskResult(AskOutcome.TIMED_OUT, null);
} catch (ExecutionException e) { } 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 * 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. * {@code turnId}, never a caller argument.
* *
* <p>Unlike {@link #send} this does not re-inject through the {@link Injector}: the worker is * <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 * 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. * 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); tasks.put(ticket, task);
if (pushLoop != null) { if (pushLoop != null) {
// CB-588: task.future only ever completes on a terminal phase (DONE or a failure) — a // 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) // 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 // 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 // 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 // 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 * 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 * (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 * 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() { private void pruneTerminalTickets() {
long cutoff = nowNanos.getAsLong() - TICKET_TTL_NANOS; 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 * The question {@code workerSession} is currently paused on via {@code fleet_ask}, if any
* (CB-582) — {@code bridge_status} uses this to show a pending question without the caller * (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 * 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}). * {@link PendingAsk}).
*/ */
public PendingAsk pendingAsk(String workerSession) { public PendingAsk pendingAsk(String workerSession) {
@@ -5,9 +5,9 @@ import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicLong; 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 * 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 * {@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}). * 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). */ /** How a delegated turn ended (or paused). */
public enum Kind { public enum Kind {
/** The worker called {@code bridge_reply} with a structured answer. */ /** The worker called {@code fleet_reply} with a structured answer. */
REPLY, 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, COMPLETION,
/** The worker ran the turn then wedged (CB-109); {@code text} is the failure context. */ /** The worker ran the turn then wedged (CB-109); {@code text} is the failure context. */
FAILED, 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, * 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 * 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}. * 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); * 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 * {@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 QUESTION
} }
@@ -75,7 +75,7 @@ public final class Rendezvous {
/** Reverse rendezvous (CB-205): worker questions awaiting the primary's answer, keyed by {@code turnId}. */ /** 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 ConcurrentHashMap<String, AskWaiter> asks = new ConcurrentHashMap<>();
private final AtomicLong askSeq = new AtomicLong(); 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<>(); private final ConcurrentHashMap<String, String> openAsksBySession = new ConcurrentHashMap<>();
/** /**
@@ -132,7 +132,7 @@ public final class Rendezvous {
return complete(session, new Resolution(Kind.REPLY, content)); 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 * 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 * 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. * {@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. * to resume its turn.
* *
* @return {@code true} if the ask was still open and got the answer; {@code false} if the * @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); 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) { public void closeAsk(String turnId) {
AskWaiter w = asks.get(turnId); AskWaiter w = asks.get(turnId);
if (w == null) { 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 * 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 * 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 * @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): * 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 * usage-limit pattern; {@code reason} carries the matched line. Like
* {@link #resolveCompletion(CompletableFuture, String)} it targets the exact captured send * {@link #resolveCompletion(CompletableFuture, String)} it targets the exact captured send
* (CB-116). A no-op if that waiter was already resolved — first resolution wins. * (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 * 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 * waiting: a worker reply queued with no live {@code fleet_send} to resolve it (CB-307), an
* async delegation ticket ({@code bridge_send(wait:false)}) that reached a terminal phase * 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 bridge_ask} to await an answer * (CB-588), or an async ticket's worker pausing mid-turn in {@code fleet_ask} to await an answer
* (CB-582). * (CB-582).
* *
* <p><strong>CB-590: one schedule per lead.</strong> All three kinds of work are triggered * <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 { public final class ReplyPushLoop {
private static final Logger log = LoggerFactory.getLogger(ReplyPushLoop.class); 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. */ /** Coalesced form, several uncollected replies for the same lead. */
static final String REPLIES_NUDGE_FORMAT = 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. */ /** Singular form, one uncollected ticket. */
static final String TICKET_NUDGE_FORMAT = 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. */ /** Coalesced form, several uncollected tickets for the same lead. */
static final String TICKETS_NUDGE_FORMAT = static final String TICKETS_NUDGE_FORMAT =
"%d tickets finished%s — run bridge_poll(ticket=...) for each to collect them: %s"; "%d tickets finished%s — run fleet_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. */ /** Singular form, one worker paused mid-turn in fleet_ask (CB-582) — names the answer call directly. */
static final String QUESTION_NUDGE_FORMAT = 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"; + "content=...) to resume its turn:\n%s";
/** Coalesced form, several open questions for the same lead. */ /** Coalesced form, several open questions for the same lead. */
static final String QUESTIONS_NUDGE_FORMAT = static final String QUESTIONS_NUDGE_FORMAT =
"%d workers are paused on a question — run bridge_poll(ticket=...) for each, then answer " "%d workers are paused on a question — run fleet_poll(ticket=...) for each, then answer "
+ "with bridge_send(turnId=..., content=...): %s"; + "with fleet_send(turnId=..., content=...): %s";
private final PrimaryRegistry primaryRegistry; private final PrimaryRegistry primaryRegistry;
private final AgentControl agents; 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. */ /** Tickets that have gone terminal but not yet been polled, keyed by ticket. */
private final ConcurrentHashMap<String, PendingTicket> pendingTickets = new ConcurrentHashMap<>(); 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 * (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 * {@link #pendingTickets} (CB-598): a fresh question keeps its source eligible regardless of
* how depleted an older, still-open question's count is. * 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 * 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, * 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 * {@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 * collected" callback the way {@link #ticketCollected} exists for tickets, so the inbox itself
* is the only signal. * 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 * 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 * 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 * 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 * 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 * lead already has. A ticket that was never pending (unknown ticket, or one nudged with no push
* loop configured) is a no-op. * 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 * Called when an async ticket's worker pauses mid-turn in {@code fleet_ask} (CB-582): the
* question is now visible via {@code bridge_poll} (Phase.ASKING), but the reverse-rendezvous * 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 * 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 * 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 * delegating lead the same way {@link #onTicketTerminal} does and coalesces onto the same
* per-lead schedule (CB-590). * 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 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 * @param question the question text
*/ */
public void onQuestionOpened(String ticket, String target, String turnId, String question) { 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 * 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. * was never pending (never nudged, or already closed) is a no-op.
*/ */
@@ -9,7 +9,7 @@ package dev.ltms.bridged.peer;
public enum Capability { 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. * the primary a question, then resuming once answered. All Claude Code peers support this.
*/ */
MID_TURN_ASK, MID_TURN_ASK,
@@ -11,7 +11,7 @@ package dev.ltms.bridged.placement;
* a usage limit, not a transient capacity or reachability concern. * 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 * <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 * 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. * the profile is unaffected, only this automatic fallback walk.
* </ul> * </ul>
* A fleet where nothing is ever quarantined or weight-0 never exercises either path, so today's * 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. * not need to distinguish "explicit 0" from "absent" itself.
* *
* <p>Exclusion is about <em>automatic</em> selection only — an explicit * <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() { public boolean excluded() {
return weight <= 0.0f; return weight <= 0.0f;
@@ -46,7 +46,7 @@ public final class BridgedApp {
/** Default blocking window for a message; kept under typical HTTP idle timeouts. */ /** 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 DEFAULT_MESSAGE_TIMEOUT_MS = 25_000;
private static final long MAX_MESSAGE_TIMEOUT_MS = 120_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 DEFAULT_ASK_TIMEOUT_MS = 55_000;
private static final long MAX_ASK_TIMEOUT_MS = 115_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.get("/profiles", this::profiles); // configured backend profiles
app.post("/members", this::spawnMember); // optional ?role=&profile= or {"role":…,"profile":…} app.post("/members", this::spawnMember); // optional ?role=&profile= or {"role":…,"profile":…}
app.delete("/members/{paneId}", this::stopMember); 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}/message", this::sendMessage); // fleet_send (primary; blocking, wait:false, or answer via turnId)
app.post("/sessions/{id}/reply", this::replyMessage); // bridge_reply (worker) app.post("/sessions/{id}/reply", this::replyMessage); // fleet_reply (worker)
app.get("/sessions/{id}/replies", this::drainReplies); // drain reply inbox (CB-307) 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.post("/sessions/{id}/ask", this::askMessage); // fleet_ask (worker → primary, CB-205)
app.get("/sessions/{id}/status", this::sessionStatus); // bridge_status app.get("/sessions/{id}/status", this::sessionStatus); // fleet_status
app.get("/tasks/{ticket}", this::taskStatus); // poll an async (wait:false) send app.get("/tasks/{ticket}", this::taskStatus); // poll an async (wait:false) send
return app; return app;
} }
@@ -341,7 +341,7 @@ public final class BridgedApp {
/** /**
* The blocking delegation call (CB-104): inject {@code content} into the worker via the * 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 * Times out with a typed 202 (working / queued / busy) rather than an error — the message may
* still land. * still land.
*/ */
@@ -370,7 +370,7 @@ public final class BridgedApp {
} }
timeout = Math.clamp(timeout, 1, MAX_MESSAGE_TIMEOUT_MS); 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()) { if (turnId != null && !turnId.isBlank()) {
writeReply(ctx, id, messages.answer(turnId, content, timeout), timeout); writeReply(ctx, id, messages.answer(turnId, content, timeout), timeout);
return; return;
@@ -392,7 +392,7 @@ public final class BridgedApp {
/** /**
* Render a {@link MessageService.Reply} onto the response — shared by a normal send and a * 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. * (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) { 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", "sessionId", id, "error", "stale_turn",
"detail", "that question is no longer open (timed out or already answered)")); "detail", "that question is no longer open (timed out or already answered)"));
case REPLIED, COMPLETED_UNREPLIED -> { 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). // fallback (a scrape of the worker's transcript when it finished without replying).
String source = reply.outcome() == MessageService.Outcome.REPLIED ? "reply" : "transcript"; String source = reply.outcome() == MessageService.Outcome.REPLIED ? "reply" : "transcript";
ctx.status(200).json(Map.of("sessionId", id, "reply", reply.text(), "replySource", source)); 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, * blocking send and blocks until it answers. 200 with the answer, 409 if no delegation is open,
* 202 if the primary stayed silent. * 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). * on this session, or queues the reply in the inbox when no send is open (CB-307).
*/ */
private void replyMessage(Context ctx) { 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 * <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 * bridge MCP — the reliable "available to receive a task" signal, unlike bare {@code idle}, which
* is also true during boot. * is also true during boot.
@@ -520,8 +520,8 @@ public final class BridgedApp {
body.put("sessionId", id); body.put("sessionId", id);
body.put("status", messages.status(id).name().toLowerCase()); body.put("status", messages.status(id).name().toLowerCase());
body.put("ready", presence.isPresent(id)); body.put("ready", presence.isPresent(id));
// CB-582: a worker paused mid-turn in an async bridge_ask is otherwise invisible to a // 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 bridge_poll's // status poll — surface the open question and how to answer it, same as fleet_poll's
// Phase.ASKING view. // Phase.ASKING view.
MessageService.PendingAsk ask = messages.pendingAsk(id); MessageService.PendingAsk ask = messages.pendingAsk(id);
if (ask != null) { if (ask != null) {
@@ -55,7 +55,7 @@ class CallerResolverTest {
assertEquals("term_a", underTrust.terminal()); assertEquals("term_a", underTrust.terminal());
assertEquals(Role.WORKER, underToken.role(), assertEquals(Role.WORKER, underToken.role(),
"worker identity is unforgeable and must never be token-gated — otherwise enabling " "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()); assertEquals("term_a", underToken.terminal());
} }
@@ -177,7 +177,7 @@ class CompletionResolverTest {
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null)); resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
assertEquals("x".repeat(CompletionResolver.MAX_SCRAPE_CHARS) 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()); waiter.getNow(null).text());
} }
@@ -347,7 +347,7 @@ class CompletionResolverTest {
@Test @Test
void aLateCompletionForOneTurnNeverResolvesTheNextTurnsWaiter() { void aLateCompletionForOneTurnNeverResolvesTheNextTurnsWaiter() {
// The cross-turn stale reply the conversation test surfaced: turn N's completion fallback // 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 // 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. // 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❯ "); 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 * 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. * orchestrator is. A lead must never receive it.
*/ */
@Test @Test
@@ -174,7 +174,7 @@ class LeadLauncherTest {
List<String> args = startedArgs(herdr); List<String> args = startedArgs(herdr);
assertFalse(args.contains("--append-system-prompt"), assertFalse(args.contains("--append-system-prompt"),
"the reply charter is a worker contract and must not be injected into a lead"); "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. */ /** 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.member.CompositePeerLauncher;
import dev.ltms.bridged.placement.BackendQuarantine; import dev.ltms.bridged.placement.BackendQuarantine;
import dev.ltms.bridged.placement.PlacementPolicies; import dev.ltms.bridged.placement.PlacementPolicies;
import io.modelcontextprotocol.server.McpSyncServerExchange;
import io.modelcontextprotocol.spec.McpSchema; import io.modelcontextprotocol.spec.McpSchema;
import dev.ltms.bridged.msg.InMemoryReplyInbox; 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.BeforeEach;
import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Test;
import org.slf4j.LoggerFactory;
import java.util.ArrayList;
import java.util.List; import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.Set; import java.util.Set;
import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import java.util.function.BiFunction;
import static org.junit.jupiter.api.Assertions.*; import static org.junit.jupiter.api.Assertions.*;
@@ -83,7 +91,7 @@ class BridgeMcpTest {
@Test @Test
void sendThenReplyRoundTrips() throws Exception { 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( CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
() -> BridgeMcp.send(messages, "term_a", "review this", 4000L, null, Set.of())); () -> BridgeMcp.send(messages, "term_a", "review this", 4000L, null, Set.of()));
@@ -106,7 +114,7 @@ class BridgeMcpTest {
@Test @Test
void asyncSendReturnsATicketThenPollReportsTheReply() throws Exception { 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()); McpSchema.CallToolResult accepted = BridgeMcp.sendAsync(messages, "term_a", "do it", null, Set.of());
assertNotEquals(Boolean.TRUE, accepted.isError()); assertNotEquals(Boolean.TRUE, accepted.isError());
String out = textOf(accepted); String out = textOf(accepted);
@@ -290,7 +298,7 @@ class BridgeMcpTest {
assertTrue(async.isError()); assertTrue(async.isError());
assertTrue(textOf(blocking).contains("sol")); assertTrue(textOf(blocking).contains("sol"));
assertTrue(textOf(blocking).contains("configured profile name")); assertTrue(textOf(blocking).contains("configured profile name"));
assertTrue(textOf(blocking).contains("bridge_list")); assertTrue(textOf(blocking).contains("fleet_list"));
assertFalse(textOf(async).contains("ticket=")); assertFalse(textOf(async).contains("ticket="));
} }
@@ -324,7 +332,7 @@ class BridgeMcpTest {
// A reply with no open send queues it in the inbox. // A reply with no open send queues it in the inbox.
BridgeMcp.reply(messages, "term_a", "queued-msg"); 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"); McpSchema.CallToolResult res = BridgeMcp.poll(messages, null, "term_a");
assertNotEquals(Boolean.TRUE, res.isError()); assertNotEquals(Boolean.TRUE, res.isError());
String text = textOf(res); String text = textOf(res);
@@ -359,7 +367,7 @@ class BridgeMcpTest {
String afterMarker = qt.substring(qt.indexOf("turnId=\"") + "turnId=\"".length()); String afterMarker = qt.substring(qt.indexOf("turnId=\"") + "turnId=\"".length());
String turnId = afterMarker.substring(0, afterMarker.indexOf('"')); 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( CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
() -> BridgeMcp.answer(messages, turnId, "config.yaml", 5000L)); () -> 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("\"name\":\"gpt-sol-5.6\""), out);
assertTrue(out.contains("\"sessionId\":\"term_peer\""), 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 // 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); assertEquals(1, out.split("\"self\":true", -1).length - 1, out);
assertTrue(out.indexOf("term_me") < out.indexOf("\"self\":true"), 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"); assertEquals(1, before.size(), "one reply in the inbox");
String msgId = before.getFirst().msgId(); 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"); BridgeMcp.reply(messages, "term_a", "orphan-again");
var peeked = messages.drainReplies("term_a"); var peeked = messages.drainReplies("term_a");
assertEquals(1, peeked.size(), "one fresh reply in the inbox"); 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} * 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 bridge_ask} question, since the reverse-rendezvous * — 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. * window it opened with is far shorter than that cadence.
*/ */
@Test @Test
@@ -819,7 +827,7 @@ class BridgeMcpTest {
answer.get(5, TimeUnit.SECONDS); 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 @Test
void whoamiReportsThePrimaryAsPrimaryAndNothingElse() { void whoamiReportsThePrimaryAsPrimaryAndNothingElse() {
@@ -994,7 +1002,7 @@ class BridgeMcpTest {
assertTrue(textOf(res).contains("architect, dev, reviewer"), textOf(res)); 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 @Test
void spawnWithResumeSessionIdPutsTheIdOnTheRoster() { void spawnWithResumeSessionIdPutsTheIdOnTheRoster() {
@@ -1020,4 +1028,103 @@ class BridgeMcpTest {
assertEquals(Boolean.TRUE, res.isError()); assertEquals(Boolean.TRUE, res.isError());
assertTrue(textOf(res).contains("explicit profile"), textOf(res)); 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 * 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 @Test
void aNudgeGoesToTheLeadThatDelegatedToThatWorker() { 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")), assertTrue(args.stream().anyMatch(a -> a.contains("\"bridge\"") && a.contains("http://127.0.0.1:8765/mcp")),
"inline bridge MCP config present"); "inline bridge MCP config present");
assertTrue(args.contains("--append-system-prompt")); 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 @Test
@@ -358,7 +358,7 @@ class CompositePeerLauncherTest {
@Test @Test
void explicitSpawnStillSucceedsOnWeightZeroProfile() { void explicitSpawnStillSucceedsOnWeightZeroProfile() {
// CB-554: weight: 0 excludes a profile from AUTOMATIC selection only — an explicit // 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). // operator's own subscription, kept weight-0 so it is never picked automatically).
FakeHerdr herdr = new FakeHerdr(); FakeHerdr herdr = new FakeHerdr();
Map<String, BridgedConfig.Profile> profiles = ordered( 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 * 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. * DEV, and enforcing the pool here would refuse a spawn the operator asked for by name.
*/ */
@Test @Test
@@ -169,8 +169,8 @@ class LeadHeartbeatLoopTest {
String t = pendingFleet().nudgeText(); String t = pendingFleet().nudgeText();
assertTrue(t.startsWith("Heartbeat:"), "the nudge identifies itself as a heartbeat"); assertTrue(t.startsWith("Heartbeat:"), "the nudge identifies itself as a heartbeat");
assertTrue(t.contains("2 worker replies pending collection"), t); assertTrue(t.contains("2 worker replies pending collection"), t);
assertTrue(t.contains("bridge_poll(target=term_a)"), t); assertTrue(t.contains("fleet_poll(target=term_a)"), t);
assertTrue(t.contains("bridge_poll(target=term_b)"), t); assertTrue(t.contains("fleet_poll(target=term_b)"), t);
assertTrue(t.contains("1 DONE session awaiting teardown"), t); assertTrue(t.contains("1 DONE session awaiting teardown"), t);
assertTrue(t.contains("3 workers live"), 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.IDLE); // deliver the task (baselines the pre-turn content)
injector.onStatus(T, AgentStatus.WORKING); // worker picks it up and works injector.onStatus(T, AgentStatus.WORKING); // worker picks it up and works
herdr.readText("BUILD GREEN: 391 files"); // the worker's turn produced new output 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); MessageService.Reply reply = send.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.COMPLETED_UNREPLIED, reply.outcome(), assertEquals(MessageService.Outcome.COMPLETED_UNREPLIED, reply.outcome(),
@@ -153,7 +153,7 @@ class MessageServiceTest {
+ "(worker unreachable or stuck)", resolution.text()); + "(worker unreachable or stuck)", resolution.text());
} }
// --- bridge_ask reverse rendezvous (CB-205) ------------------------------------------------ // --- fleet_ask reverse rendezvous (CB-205) ------------------------------------------------
@Test @Test
void askSurfacesAsAQuestionAndTheAnswerResumesTheSameTurn() throws Exception { void askSurfacesAsAQuestionAndTheAnswerResumesTheSameTurn() throws Exception {
@@ -172,7 +172,7 @@ class MessageServiceTest {
assertEquals("which config file?", q.text()); assertEquals("which config file?", q.text());
assertNotNull(q.turnId(), "a question carries a turnId to answer on"); 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<MessageService.Reply> answer =
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000)); 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.IDLE); // deliver
injector.onStatus(T, AgentStatus.WORKING); // worker picks it up, then pauses to ask 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<MessageService.AskResult> ask1 =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000)); CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
CompletableFuture<MessageService.AskResult> ask2 = CompletableFuture<MessageService.AskResult> ask2 =
@@ -297,7 +297,7 @@ class MessageServiceTest {
assertNotNull(q.turnId()); assertNotNull(q.turnId());
// The primary answers, unblocking the worker; but the worker never sends the follow-up // 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); MessageService.Reply answer = messages.answer(q.turnId(), "config.yaml", 200);
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answer.outcome(), assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answer.outcome(),
"an answered worker that never replies times out as still working"); "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"); assertNotNull(view, "a resolved async send must become DONE");
assertEquals(MessageService.Phase.DONE, view.phase()); assertEquals(MessageService.Phase.DONE, view.phase());
assertEquals("async result", view.reply(), "the completed ticket reports the reply"); 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 @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 * 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 * 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. * 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 * 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. * turnId — ownership stays L throughout; the answer path never touches the registry.
*/ */
@@ -533,12 +533,12 @@ class MessageServiceTest {
@Test @Test
void aQuestionIsNeverQueuedInTheInbox() { 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. // 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); MessageService.AskResult r = messages.ask(T, "anyone there?", 500);
assertEquals(MessageService.AskOutcome.NO_WAITER, r.outcome(), 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"); assertTrue(messages.drainReplies(T).isEmpty(), "questions must never be queued");
} }
@@ -550,7 +550,7 @@ class MessageServiceTest {
awaitUninterruptibly(T); awaitUninterruptibly(T);
injectDelivery(); 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"); herdr.readText("done-scraped");
completion.onTurnComplete(T); // The fallback arms and resolves the captured waiter. 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()); assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
} }
// --- CB-582: bridge_status pendingAsk() ------------------------------------------------------ // --- CB-582: fleet_status pendingAsk() ------------------------------------------------------
@Test @Test
void pendingAskReturnsNullWhenNoQuestionIsOpen() throws Exception { void pendingAskReturnsNullWhenNoQuestionIsOpen() throws Exception {
@@ -741,7 +741,7 @@ class MessageServiceTest {
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING); MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
MessageService.PendingAsk pending = messages.pendingAsk(T); 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(ticket, pending.ticket());
assertEquals("which config file?", pending.question()); assertEquals("which config file?", pending.question());
assertEquals(asking.turnId(), pending.turnId()); 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 // 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 // (sendAsync registers a rendezvous waiter — see asyncTasksByWaiter), so it never reached
// ReplyPushLoop.onReplyQueued. These prove the ticket reaches ReplyPushLoop through the new // 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"; private static final String LEAD = "term_lead";
@@ -808,7 +808,7 @@ class MessageServiceTest {
awaitNudge(wiring.leadHerdr()); awaitNudge(wiring.leadHerdr());
String nudge = wiring.leadHerdr().lastCall("agent.prompt").params().toString(); String nudge = wiring.leadHerdr().lastCall("agent.prompt").params().toString();
assertTrue(nudge.contains(ticket), "the nudge should name the ticket: " + nudge); 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); "the nudge should name the exact ticket-collecting call: " + nudge);
assertFalse(nudge.toUpperCase().contains("FAILED"), assertFalse(nudge.toUpperCase().contains("FAILED"),
"a successfully-replied ticket's nudge must not say it failed: " + nudge); "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 @Test
void anAsyncTicketThatPausesOnAQuestionNudgesTheLeadWithNoPriorPollCall() throws Exception { void anAsyncTicketThatPausesOnAQuestionNudgesTheLeadWithNoPriorPollCall() throws Exception {
@@ -891,7 +891,7 @@ class MessageServiceTest {
String nudge = wiring.leadHerdr().lastCall("agent.prompt").params().toString(); String nudge = wiring.leadHerdr().lastCall("agent.prompt").params().toString();
assertTrue(nudge.contains(ticket), "the nudge should name the ticket: " + nudge); 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(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); "the nudge should name the exact answer call: " + nudge);
assertTrue(nudge.contains("which config file?"), "the nudge should include the question: " + 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} * {@code pruneTerminalTickets} drops entries from it once {@link MessageService#TICKET_TTL_NANOS}
* elapses. Before this test, that prune never told {@code ReplyPushLoop} — its own * 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 * {@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 * 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. * 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(); String latestNudge = wiring.leadHerdr().lastCall("agent.prompt").params().toString();
assertTrue(latestNudge.contains(fresh), "the fresh ticket's nudge must still arrive: " + latestNudge); assertTrue(latestNudge.contains(fresh), "the fresh ticket's nudge must still arrive: " + latestNudge);
assertFalse(latestNudge.contains(stale), 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); + "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"); 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; 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}. * to ask the primary. Unit-level — the message-layer round-trip is covered in {@link MessageServiceTest}.
*/ */
class RendezvousTest { class RendezvousTest {
@@ -168,8 +168,8 @@ class ReplyPushLoopTest {
// Exactly one nudge = exactly 1 agent.prompt call (it submits itself) // Exactly one nudge = exactly 1 agent.prompt call (it submits itself)
assertEquals(1, rec.sendCount()); assertEquals(1, rec.sendCount());
assertTrue(rec.sentParams().stream() assertTrue(rec.sentParams().stream()
.anyMatch(e -> e.getValue().toString().contains("bridge_poll")), .anyMatch(e -> e.getValue().toString().contains("fleet_poll")),
"nudge text should contain bridge_poll"); "nudge text should contain fleet_poll");
} }
@Test @Test
@@ -226,14 +226,14 @@ class ReplyPushLoopTest {
void nudgeFormatIsCorrect() { void nudgeFormatIsCorrect() {
String nudge = ReplyPushLoop.NUDGE_FORMAT.formatted(WORKER, WORKER); String nudge = ReplyPushLoop.NUDGE_FORMAT.formatted(WORKER, WORKER);
assertTrue(nudge.contains("Worker term_worker")); assertTrue(nudge.contains("Worker term_worker"));
assertTrue(nudge.contains("bridge_poll(target=term_worker)")); assertTrue(nudge.contains("fleet_poll(target=term_worker)"));
} }
@Test @Test
void repliesNudgeFormatIsCorrect() { void repliesNudgeFormatIsCorrect() {
String multi = ReplyPushLoop.REPLIES_NUDGE_FORMAT.formatted(2, "term_worker1, term_worker2"); String multi = ReplyPushLoop.REPLIES_NUDGE_FORMAT.formatted(2, "term_worker1, term_worker2");
assertTrue(multi.contains("2 workers")); 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 ----------------------- // --- CB-588: async ticket terminal nudges — decide() logic on tickets -----------------------
@@ -289,8 +289,8 @@ class ReplyPushLoopTest {
assertEquals(1, rec.sendCount()); assertEquals(1, rec.sendCount());
String nudge = rec.sentParams().getFirst().getValue().toString(); String nudge = rec.sentParams().getFirst().getValue().toString();
assertTrue(nudge.contains("task-1"), "nudge should name the ticket"); assertTrue(nudge.contains("task-1"), "nudge should name the ticket");
assertTrue(nudge.contains("bridge_poll(ticket="), "nudge should name the exact ticket-poll call"); assertTrue(nudge.contains("fleet_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"); assertFalse(nudge.contains("fleet_poll(target="), "a ticket-only nudge must not tell the lead to run the target-poll call");
} }
@Test @Test
@@ -468,11 +468,11 @@ class ReplyPushLoopTest {
void ticketNudgeFormatIsCorrect() { void ticketNudgeFormatIsCorrect() {
String single = ReplyPushLoop.TICKET_NUDGE_FORMAT.formatted("task-1", "", "task-1"); String single = ReplyPushLoop.TICKET_NUDGE_FORMAT.formatted("task-1", "", "task-1");
assertTrue(single.contains("Ticket 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"); String multi = ReplyPushLoop.TICKETS_NUDGE_FORMAT.formatted(2, "", "task-1, task-2");
assertTrue(multi.contains("2 tickets")); assertTrue(multi.contains("2 tickets"));
assertTrue(multi.contains("bridge_poll(ticket=...)")); assertTrue(multi.contains("fleet_poll(ticket=...)"));
} }
// --- CB-307 nudge path is unchanged (regression) -------------------------------------------- // --- CB-307 nudge path is unchanged (regression) --------------------------------------------
@@ -480,7 +480,7 @@ class ReplyPushLoopTest {
@Test @Test
void inboxNudgeStillUsesTheOriginalTargetPollCall() { void inboxNudgeStillUsesTheOriginalTargetPollCall() {
String nudge = ReplyPushLoop.NUDGE_FORMAT.formatted(WORKER, WORKER); 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"); "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 — " "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"); + "two nudge injections into the same lead pane must never overlap");
String nudge = rec.sentParams().getFirst().getValue().toString(); 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); "the combined nudge must still mention the reply: " + nudge);
assertTrue(nudge.contains("task-1"), "the combined nudge must still mention the ticket: " + nudge); assertTrue(nudge.contains("task-1"), "the combined nudge must still mention the ticket: " + nudge);
} }
@@ -525,7 +525,7 @@ class ReplyPushLoopTest {
Thread.sleep(200); Thread.sleep(200);
assertEquals(1, rec.sendCount(), "exactly one nudge once injectable — reply and ticket coalesced"); assertEquals(1, rec.sendCount(), "exactly one nudge once injectable — reply and ticket coalesced");
String nudge = rec.sentParams().getFirst().getValue().toString(); 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); 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"); 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 @Test
void onQuestionOpenedWithNoKnownLeadNeverStartsASchedule() throws Exception { void onQuestionOpenedWithNoKnownLeadNeverStartsASchedule() throws Exception {
@@ -706,7 +706,7 @@ class ReplyPushLoopTest {
String nudge = rec.sentParams().getFirst().getValue().toString(); String nudge = rec.sentParams().getFirst().getValue().toString();
assertTrue(nudge.contains("task-1"), "nudge should name the ticket: " + nudge); 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("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); 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( String single = ReplyPushLoop.QUESTION_NUDGE_FORMAT.formatted(
WORKER, "task-1", "term_worker#1", "which config?"); WORKER, "task-1", "term_worker#1", "which config?");
assertTrue(single.contains("Worker term_worker")); 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?")); assertTrue(single.contains("which config?"));
String multi = ReplyPushLoop.QUESTIONS_NUDGE_FORMAT.formatted(2, "task-1 (turnId=t1), task-2 (turnId=t2)"); String multi = ReplyPushLoop.QUESTIONS_NUDGE_FORMAT.formatted(2, "task-1 (turnId=t1), task-2 (turnId=t2)");
assertTrue(multi.contains("2 workers")); assertTrue(multi.contains("2 workers"));
assertTrue(multi.contains("bridge_poll(ticket=...)")); assertTrue(multi.contains("fleet_poll(ticket=...)"));
} }
// --- metrics (CB-512) ---------------------------------------------------------------------- // --- metrics (CB-512) ----------------------------------------------------------------------
@@ -355,7 +355,7 @@ class BridgedAppTest {
@Test @Test
void messageReturnsTheWorkersStructuredReply() throws Exception { 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 FakeHerdr herdr = new FakeHerdr().agentStatus("idle"); // poller delivers the injection
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw")); 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 * 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. * question, since the reverse-rendezvous window it opened with is far shorter than that cadence.
*/ */
@Test @Test
@@ -528,7 +528,7 @@ class BridgedAppTest {
assertEquals(turnId, body.get("turnId").asText()); assertEquals(turnId, body.get("turnId").asText());
assertEquals(ticket, body.get("ticket").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. // background ask thread does not linger past the test.
var answer = java.util.concurrent.CompletableFuture.supplyAsync(() -> { var answer = java.util.concurrent.CompletableFuture.supplyAsync(() -> {
try { try {
@@ -206,7 +206,7 @@ class SessionManagerTest {
@Test @Test
void rosterViewExposesTheCharterReceiptButNeverTheCharterText() { 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). // 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, MemberSession s = new MemberSession("p1", "term1", "prof", MemberRole.DEV, "/cwd", null,
0, 0, 0, MemberSession.State.READY, null, null, 0, 0, 0, MemberSession.State.READY, null, null,