CB-622 Unit A: register fleet_* MCP tools, keep bridge_* working #135

Merged
ltms merged 1 commits from worker/cb-622a-165dff-1 into main 2026-08-22 21:58:03 +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.");
}
// CB-307: active push-to-primary loop — nudge the primary when replies land without an
// open bridge_send. Uses its own lightweight scheduled executor, separate from the injector.
// open fleet_send. Uses its own lightweight scheduled executor, separate from the injector.
int maxReminders = cfg.primary() != null ? cfg.primary().remindersOrDefault() : 5;
long backoffMs = cfg.primary() != null ? cfg.primary().backoffMsOrDefault() : 15_000L;
var pushScheduler = Executors.newSingleThreadScheduledExecutor(r ->
@@ -458,7 +458,7 @@ public final class Bridged {
primaryRegistry.forgetDelegation(detail.terminalId()); // CB-532: don't leak the lead binding
});
// MCP server face (CB-105): bridge_send/bridge_reply/bridge_status, mounted at /mcp.
// MCP server face (CB-105): fleet_send/fleet_reply/fleet_status, mounted at /mcp.
// Caller identity is resolved from the connection (peer PID → herdr pane), not arguments.
ConnectionIdentity identity = new ConnectionIdentity(
new PaneLocator(herdr), new LsofPeerPidLookup(), new LsofProcessCwdLookup());
@@ -94,7 +94,7 @@ public final class MemberRegistry implements MemberLifecycle {
* An immutable copy of the live {@code terminal_id → slot name} bindings.
*
* <p>Passed to {@link CallerResolver} as the source of architect identity, and what
* {@code bridge_whoami}/the roster will read to say which slot a pane hosts. Empty until the
* {@code fleet_whoami}/the roster will read to say which slot a pane hosts. Empty until the
* spawn lifecycle binds a slot.
*/
public Map<String, String> snapshot() {
@@ -39,14 +39,14 @@ public record Principal(Role role, String terminal, long pid, String name) {
*
* <p>Carries {@link Role#PRIMARY}: a lead <em>is</em> a primary as far as authorization goes,
* so every existing {@code isPrimary()} gate keeps working unchanged and the role table needed
* no new entry. The name is reporting only — it lets {@code bridge_whoami} say <em>which</em>
* no new entry. The name is reporting only — it lets {@code fleet_whoami} say <em>which</em>
* lead is asking once more than one is configured.
*
* <p><strong>CB-532: a lead now carries the terminal it was matched by.</strong> Under CB-530 it
* deliberately did not, because {@code terminal} meant "which worker pane" everywhere and a
* non-null one would have enrolled the lead in the worker presence map. That reading was what
* made a lead unaddressable: {@link #ownsSession} could never be true for it, so
* {@code bridge_reply} was refused and one lead could send to another but never be answered.
* {@code fleet_reply} was refused and one lead could send to another but never be answered.
* The terminal now means "which pane is this caller", the presence map keys on
* {@link #isSpawnedMember()} instead, and a lead is a peer that can both send and receive.
*/
@@ -63,7 +63,7 @@ public record Principal(Role role, String terminal, long pid, String name) {
* An architect (CB-548), identified by the slot it occupies and the pane bound to it.
*
* <p>Carries {@link Role#ARCHITECT}. {@code slotName} is reporting only — it lets
* {@code bridge_whoami} say <em>which</em> architect slot is asking, and it is the key the
* {@code fleet_whoami} say <em>which</em> architect slot is asking, and it is the key the
* (future) spawn lifecycle reads a profile back from. Identity is the {@code terminal}: like a
* worker's it comes from the connection and the live terminal→slot binding, so
* {@code ownsSession} works exactly as it does for a worker — an architect acts as its own
@@ -184,7 +184,7 @@ public record BridgedConfig(
* and {@code {n}} (per-worker number, to keep sibling tabs distinct)
* are substituted (default {@code "worker: {profile} #{n}"})
* @param mcpUrl bridge MCP URL to provision into the worker's {@code configDir} so it
* can call {@code bridge_reply} ({@code null}/blank → no provisioning; the
* can call {@code fleet_reply} ({@code null}/blank → no provisioning; the
* worker won't reply, only the fallback/timeout resolves the send)
* @param cwd fixed working directory for this profile's workers (CB-112 "told otherwise");
* {@code null}/blank → inherit the primary's cwd, else the daemon's
@@ -224,7 +224,7 @@ public record BridgedConfig(
* auto-select this profile": it is excluded from every automatic policy's
* pool the same way a quarantined candidate is (see
* {@code PlacementPolicyUtil}). This does not make the profile
* unreachable — an explicit {@code bridge_spawn{profile:"..."}} bypasses
* unreachable — an explicit {@code fleet_spawn{profile:"..."}} bypasses
* placement entirely and still resolves it. Weights among the remaining
* (non-excluded) candidates need not sum to 1.0; only their ratios matter.
* @param maxLoad max live workers allowed on this profile at one time. Absent
@@ -234,7 +234,7 @@ public record BridgedConfig(
* the same way a {@code weight <= 0} profile is (see
* {@code PlacementPolicyUtil.available()}, which already treats "at cap"
* and "excluded" alike), and an explicit
* {@code bridge_spawn{profile:"..."}} against it is refused too (see
* {@code fleet_spawn{profile:"..."}} against it is refused too (see
* {@code CompositePeerLauncher.enforceMaxLoad}) — a cap is a capacity
* statement that does not stop being true just because the profile was
* named directly. A negative value has no sane meaning (there is no
@@ -259,7 +259,7 @@ public record BridgedConfig(
* {@code ANTHROPIC_BASE_URL} or {@code ANTHROPIC_AUTH_TOKEN} is refused at
* config load (CB-542): on the subscription path no guard would vet it.
* @param exhaustedPattern regex matched against a completion-fallback scrape (CB-578 stage A) to
* classify a turn that ended with no {@code bridge_reply} as the backend
* classify a turn that ended with no {@code fleet_reply} as the backend
* having refused on a subscription usage limit, rather than a real answer.
* {@code null}/blank ⇒ the classification never fires for this profile and
* today's completion-fallback behaviour is unchanged. Every backend words
@@ -605,7 +605,7 @@ public record BridgedConfig(
* work as peers — the second is silently demoted and refused every orchestration call.
*
* <p>{@code kind} and {@code model} are descriptive only: they document what runs in the pane
* and are reported back by {@code bridge_whoami}.
* and are reported back by {@code fleet_whoami}.
*
* <p><b>A lead is now also creatable (CB-557).</b> Before, nothing spawned one — a lead
* pre-existed, which is why it had to be recognised by configuration rather than created. With
@@ -811,7 +811,7 @@ public record BridgedConfig(
/**
* Opt-in idle-lead heartbeat (CB-551): when the single lead has been continuously idle past
* {@code idleAfterSeconds} with no open {@code bridge_send} driving it, nudge it back to work.
* {@code idleAfterSeconds} with no open {@code fleet_send} driving it, nudge it back to work.
*
* <p>Deliberately opt-in ({@code null} ⇒ off, exactly like {@code leadScan:}). The heartbeat
* spends the operator's model subscription on its own initiative — it prompts the lead to start
@@ -6,7 +6,7 @@ import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
/**
* Counts turns that ended via the completion fallback instead of {@code bridge_reply}.
* Counts turns that ended via the completion fallback instead of {@code fleet_reply}.
* MUTE is an observation by target and profile, not a classifier state and never suppresses faults.
*/
public final class MuteCounter {
@@ -16,14 +16,14 @@ import java.util.regex.Pattern;
/**
* The CB-106 completion fallback: bridges the {@link Injector}'s turn-completion signal to the
* {@link Rendezvous} so a blocking {@code bridge_send} resolves even when the worker finishes its
* task without ever calling {@code bridge_reply} — the common case for a real delegated coding task.
* {@link Rendezvous} so a blocking {@code fleet_send} resolves even when the worker finishes its
* task without ever calling {@code fleet_reply} — the common case for a real delegated coding task.
*
* <p>On a confirmed {@code working → idle} boundary it scrapes the worker's recent transcript and
* resolves the awaiting send with that tail (a {@link Rendezvous.Kind#COMPLETION} resolution, so the
* caller can tell a scrape from a structured reply). It scrapes only when a send is actually waiting
* — a fleet worker's own turns, or a send that already timed out, cost no herdr traffic. An explicit
* {@code bridge_reply} that raced in first wins; {@link Rendezvous#resolveCompletion} is then a no-op.
* {@code fleet_reply} that raced in first wins; {@link Rendezvous#resolveCompletion} is then a no-op.
*
* <p>It also handles the CB-109 stall signal ({@link #onTurnFailed}): a worker that ran a turn then
* wedged in an {@code unknown} state resolves the send as a failure (with the error screen as
@@ -38,7 +38,7 @@ import java.util.regex.Pattern;
* <p><strong>Waiter-specific resolution (CB-116).</strong> On delivery we also capture the exact
* {@link Rendezvous} waiter this turn belongs to, and the completion/failure fallbacks resolve
* <em>that</em> waiter — never "whatever send is waiting now". A completion fallback runs on a virtual
* thread and can land after the worker's {@code bridge_reply} already resolved the turn and the
* thread and can land after the worker's {@code fleet_reply} already resolved the turn and the
* <em>next</em> send opened its own waiter on the same session; resolving the current waiter would
* then deliver turn N's stale scrape as turn N+1's answer. Targeting the captured waiter makes a late
* completion a harmless no-op (its waiter is already done) instead of a cross-turn stale reply.
@@ -62,7 +62,7 @@ public final class CompletionResolver implements TurnListener {
static final int MAX_SCRAPE_CHARS = 4000;
private static final String CLIPPED_PANE_TAIL_MARKER =
"[Pane tail clipped: member did not call bridge_reply.]";
"[Pane tail clipped: member did not call fleet_reply.]";
private final AgentControl agents;
private final Rendezvous rendezvous;
@@ -171,7 +171,7 @@ public final class CompletionResolver implements TurnListener {
void resolve(String target, InFlight turn) {
CompletableFuture<Rendezvous.Resolution> waiter = turn == null ? null : turn.waiter();
if (waiter == null || waiter.isDone()) {
// Nobody is blocked on THIS turn (it had no send, or its bridge_reply already won). Skip
// Nobody is blocked on THIS turn (it had no send, or its fleet_reply already won). Skip
// the scrape; resolving the current waiter here would be the CB-116 cross-turn stale reply.
inFlight.remove(target, turn);
return;
@@ -197,7 +197,7 @@ public final class CompletionResolver implements TurnListener {
// Misattribution guard (CB-115): if the scrape is byte-identical to the pane content at
// delivery, this turn produced no new output — the boundary belongs to the previous turn's
// wind-down (common on rapid back-to-back sends). Suppress rather than resolve the send with
// a stale answer; the real bridge_reply (or a later genuine completion) resolves it instead.
// a stale answer; the real fleet_reply (or a later genuine completion) resolves it instead.
// A scrape that failed to read is exempt — an empty tail there is "couldn't see", not "no change".
String baseline = turn.baseline();
if (!scrapeFailed && baseline != null && baseline.equals(tail)) {
@@ -205,7 +205,7 @@ public final class CompletionResolver implements TurnListener {
target);
return; // keep the in-flight record: a later genuine completion still needs it
}
// CB-578 stage A: a turn that ended with no bridge_reply AND whose scrape matches the
// CB-578 stage A: a turn that ended with no fleet_reply AND whose scrape matches the
// backend's configured usage-limit pattern is a refusal, not an answer. Classify it as
// BACKEND_EXHAUSTED rather than handing the caller a scrape that reads like a real reply.
if (!scrapeFailed) {
@@ -215,7 +215,7 @@ public final class CompletionResolver implements TurnListener {
String reason = "backend exhausted (usage limit): " + matchedLine;
if (rendezvous.resolveExhausted(waiter, reason)) {
inFlight.remove(target, turn);
log.warn("completion for {} classified BACKEND_EXHAUSTED (no bridge_reply; scrape "
log.warn("completion for {} classified BACKEND_EXHAUSTED (no fleet_reply; scrape "
+ "matched the profile's exhausted pattern): {}", target, reason);
// CB-578 stage B: only on the resolution that actually won the race — a late
// duplicate must never quarantine a credential twice for one refusal.
@@ -229,7 +229,7 @@ public final class CompletionResolver implements TurnListener {
inFlight.remove(target, turn);
if (clipped) {
log.warn("completion scrape for {} clipped from {} chars to the {} char cap; "
+ "member did not call bridge_reply, so the pane tail is partial",
+ "member did not call fleet_reply, so the pane tail is partial",
target, originalLength, MAX_SCRAPE_CHARS);
}
log.debug("resolved send to {} via turn-completion fallback ({} chars scraped)",
@@ -6,7 +6,7 @@ import dev.ltms.bridged.msg.TurnToken;
* Notified when a worker's delegated turn is observed to complete — a confirmed
* {@code WORKING → IDLE} transition after a delivery. This is the CB-106 completion signal the
* {@code CompletionResolver} uses to resolve a blocked send whose worker never called
* {@code bridge_reply}. Kept as a seam so the {@link Injector} needs no dependency on the message
* {@code fleet_reply}. Kept as a seam so the {@link Injector} needs no dependency on the message
* layer and stays unit-testable with a capturing fake.
*/
@FunctionalInterface
@@ -26,7 +26,7 @@ import java.util.stream.Collectors;
* must never receive:
* <ol>
* <li>it appends the <em>reply charter</em> — "you are an off-subscription worker … end every turn
* with {@code bridge_reply}". A lead is the orchestrator; telling it that it is a worker is
* with {@code fleet_reply}". A lead is the orchestrator; telling it that it is a worker is
* exactly backwards.</li>
* <li>it registers the session with {@code SessionManager}, which subjects it to the idle reaper,
* the context cap and the shutdown drain. An idle lead is the normal state of a lead, so the
@@ -37,19 +37,24 @@ import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.function.BiFunction;
import java.util.function.Function;
import java.util.function.LongSupplier;
import java.util.function.Supplier;
import java.util.stream.Collectors;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
/**
* The MCP SERVER face (CB-105): a Streamable-HTTP MCP server whose tools are <em>thin adapters</em>
* over the same {@link MessageService}/{@link Rendezvous} the REST routes use — so the two are
* validated by parity, not by re-implementing behaviour. The primary Opus calls {@code bridge_send}
* / {@code bridge_status}; the worker calls {@code bridge_reply}.
* validated by parity, not by re-implementing behaviour. The primary Opus calls {@code fleet_send}
* / {@code fleet_status}; the worker calls {@code fleet_reply}.
*
* <p>Beyond delegation the primary also manages the fleet here (CB-108): {@code bridge_spawn} /
* {@code bridge_list} / {@code bridge_stop} drive the {@link PeerLauncher} SPI so a worker's whole
* <p>Beyond delegation the primary also manages the fleet here (CB-108): {@code fleet_spawn} /
* {@code fleet_list} / {@code fleet_stop} drive the {@link PeerLauncher} SPI so a worker's whole
* lifecycle is managed through MCP, with each adapter's subscription boundary enforced inside it.
*
* <p>The tool <em>logic</em> lives in package-private static methods returning a
@@ -58,14 +63,24 @@ import java.util.stream.Collectors;
*/
public final class BridgeMcp {
private static final Logger log = LoggerFactory.getLogger(BridgeMcp.class);
private static final long DEFAULT_TIMEOUT_MS = 25_000;
private static final long MAX_TIMEOUT_MS = 120_000;
// bridge_ask blocks the WORKER's own MCP call, which its client caps near 60s — default under
// fleet_ask blocks the WORKER's own MCP call, which its client caps near 60s — default under
// that so the bridge returns a clean timeout before the client severs the call (CB-205).
private static final long ASK_DEFAULT_TIMEOUT_MS = 55_000;
private static final long ASK_MAX_TIMEOUT_MS = 115_000;
private static final ObjectMapper MAPPER = new ObjectMapper(); // worker-view JSON projections
/**
* CB-622: the product is renaming {@code bridge_*} tools to {@code fleet_*}. Both names reach
* the same handler (registered below); this set makes the "old name used" WARN fire once per
* old name for the life of the process, not once per call — a per-name flag, not a numeric
* sentinel, so it survives concurrent callers cleanly and reads unambiguously in a log.
*/
private static final Set<String> WARNED_DEPRECATED_NAMES = ConcurrentHashMap.newKeySet();
/** Transport-context key under which the extractor stashes the resolved caller identity. */
static final String CALLER_TERMINAL = "callerTerminal";
/** Transport-context key under which the extractor stashes the caller's PID (for cwd inherit). */
@@ -83,7 +98,7 @@ public final class BridgeMcp {
private final HealthCoverageSource healthCoverage;
private final QuarantineSource quarantine;
/** Capacity facts used by {@code bridge_list}; production must supply the placement live count. */
/** Capacity facts used by {@code fleet_list}; production must supply the placement live count. */
public record CapacitySource(Function<String, Integer> liveCount, Function<String, Integer> maxLoad,
Supplier<Set<String>> configuredProfiles, LongSupplier clock) {
/** Inert test-only source. It omits capacity rather than inventing zero live counts. */
@@ -95,7 +110,7 @@ public final class BridgeMcp {
public record HealthCoverageSource(Supplier<String> value) { }
/**
* CB-578 stage B quarantine facts used by {@code bridge_profiles}: a profile → credential id
* CB-578 stage B quarantine facts used by {@code fleet_profiles}: a profile → credential id
* lookup, plus the shared {@link BackendQuarantine} to read remaining cooldowns off.
*/
public record QuarantineSource(Function<String, String> credentialIdFor, BackendQuarantine quarantine) {
@@ -109,7 +124,7 @@ public final class BridgeMcp {
* Jetty's context handler and never passes through Javalin's {@code before}
* filter, so the REST guard does not cover it.
* @param metrics registry for auth-failure counting; may be {@code null}
* @param quarantine CB-578 stage B facts for {@code bridge_profiles}; required — pass
* @param quarantine CB-578 stage B facts for {@code fleet_profiles}; required — pass
* {@link QuarantineSource#none()} for a caller that does not want the feature
*/
public BridgeMcp(MessageService messages, PeerLauncher workers, SessionManager sessions,
@@ -124,7 +139,7 @@ public final class BridgeMcp {
.jsonMapper(json)
.mcpEndpoint("/mcp")
// Resolve the caller from the connection (peer PID → herdr pane) in one lookup: the
// worker terminal for bridge_reply (no spoofable arg), and the PID so bridge_spawn can
// worker terminal for fleet_reply (no spoofable arg), and the PID so fleet_spawn can
// inherit the primary's cwd (CB-112). Any contact from a worker marks it available
// (CB-113) — its MCP initialize is the reliable "the agent is up" signal.
.contextExtractor(req -> {
@@ -145,10 +160,10 @@ public final class BridgeMcp {
CALLER_NAME, orEmpty(p.name())));
})
.build();
this.server = McpServer.sync(transport)
.serverInfo("bridge", "0.1.0")
.capabilities(McpSchema.ServerCapabilities.builder().tools(true).build())
.toolCall(sendTool(), (exchange, req) -> {
// CB-622: each handler is built once and reused for BOTH its fleet_* tool and its
// deprecated bridge_* twin (registered below), so the two names can never drift apart.
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> sendHandler =
(exchange, req) -> {
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.SEND,
str(req.arguments(), "sessionId"));
if (denied != null) return denied;
@@ -163,7 +178,7 @@ public final class BridgeMcp {
String content = str(a, "content");
String turnId = str(a, "turnId");
if (turnId != null && !turnId.isBlank()) {
// Answering a worker's bridge_ask (CB-205): resolve its blocked question and
// Answering a worker's fleet_ask (CB-205): resolve its blocked question and
// block for the worker's reply as it resumes the same turn. This is the same
// delegation, so ownership is left untouched (CB-548) — never re-recorded.
return answer(messages, turnId, content, timeoutMs(a));
@@ -178,43 +193,49 @@ public final class BridgeMcp {
return Boolean.FALSE.equals(a.get("wait"))
? sendAsync(messages, target, content, onAccepted, workers.profiles())
: send(messages, target, content, timeoutMs(a), onAccepted, workers.profiles());
})
// bridge_reply's identity is the CONNECTION, never an argument — so the authz check
// is "is this caller a worker at all", and it can only ever reply as itself.
.toolCall(replyTool(), (exchange, req) -> {
};
// fleet_reply's identity is the CONNECTION, never an argument — so the authz check
// is "is this caller a worker at all", and it can only ever reply as itself.
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> replyHandler =
(exchange, req) -> {
String self = callerTerminal(exchange);
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.REPLY, self);
if (denied != null) return denied;
return reply(messages, self, str(req.arguments(), "content"));
})
// bridge_ask (CB-205): a worker's mid-turn question — identity from the CONNECTION.
.toolCall(askTool(), (exchange, req) -> {
};
// fleet_ask (CB-205): a worker's mid-turn question — identity from the CONNECTION.
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> askHandler =
(exchange, req) -> {
String self = callerTerminal(exchange);
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.ASK, self);
if (denied != null) return denied;
return ask(messages, self, str(req.arguments(), "question"), timeoutMs(req.arguments()));
})
.toolCall(statusTool(), (exchange, req) -> {
};
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> statusHandler =
(exchange, req) -> {
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
if (denied != null) return denied;
return status(messages, str(req.arguments(), "sessionId"));
})
.toolCall(pollTool(), (exchange, req) -> {
};
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> pollHandler =
(exchange, req) -> {
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
if (denied != null) return denied;
Map<String, Object> a = req.arguments();
return poll(messages, str(a, "ticket"), str(a, "target"));
})
// CB-307 Increment 3: per-msgId ack (not needed in v1 but supported by the inbox).
// Acking removes a reply from the inbox, so it is a drain, not a read.
.toolCall(ackTool(), (exchange, req) -> {
};
// CB-307 Increment 3: per-msgId ack (not needed in v1 but supported by the inbox).
// Acking removes a reply from the inbox, so it is a drain, not a read.
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> ackHandler =
(exchange, req) -> {
Map<String, Object> a = req.arguments();
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.DRAIN, str(a, "target"));
if (denied != null) return denied;
return ack(messages, str(a, "target"), str(a, "msgId"));
})
// Fleet management (CB-108): spawn/list/stop over ClaudeCodeLauncher.
.toolCall(spawnTool(), (exchange, req) -> {
};
// Fleet management (CB-108): spawn/list/stop over ClaudeCodeLauncher.
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> spawnHandler =
(exchange, req) -> {
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.SPAWN, null);
if (denied != null) return denied;
String caller = callerTerminal(exchange);
@@ -229,30 +250,72 @@ public final class BridgeMcp {
return spawn(sessions, str(a, "profile"), str(a, "role"), str(a, "cwd"), callerCwd,
callerTerminal(exchange), worktreeRequest(a),
str(a, "sessionName"), str(a, "resumeSessionId"));
})
.toolCall(listTool(), (exchange, _) -> {
};
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> listHandler =
(exchange, _) -> {
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
if (denied != null) return denied;
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine,
callers == null ? Map.of() : callers.leads(),
callerTerminal(exchange));
})
.toolCall(stopTool(), (exchange, req) -> {
};
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> stopHandler =
(exchange, req) -> {
String paneId = str(req.arguments(), "paneId");
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.STOP, paneId);
if (denied != null) return denied;
return stop(sessions, paneId);
})
.toolCall(profilesTool(), (exchange, _) -> {
};
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> profilesHandler =
(exchange, _) -> {
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
if (denied != null) return denied;
return profiles(workers, quarantine);
})
.toolCall(whoamiTool(), (exchange, _) -> {
};
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> whoamiHandler =
(exchange, _) -> {
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
if (denied != null) return denied;
return whoami(principal(exchange), sessions);
})
};
McpSchema.Tool fleetSend = sendTool();
McpSchema.Tool fleetReply = replyTool();
McpSchema.Tool fleetAsk = askTool();
McpSchema.Tool fleetStatus = statusTool();
McpSchema.Tool fleetPoll = pollTool();
McpSchema.Tool fleetAck = ackTool();
McpSchema.Tool fleetSpawn = spawnTool();
McpSchema.Tool fleetList = listTool();
McpSchema.Tool fleetStop = stopTool();
McpSchema.Tool fleetProfiles = profilesTool();
McpSchema.Tool fleetWhoami = whoamiTool();
this.server = McpServer.sync(transport)
.serverInfo("bridge", "0.1.0")
.capabilities(McpSchema.ServerCapabilities.builder().tools(true).build())
.toolCall(fleetSend, sendHandler)
.toolCall(deprecatedTwin(fleetSend, "bridge_send"), deprecatedHandler(fleetSend, "bridge_send", sendHandler))
.toolCall(fleetReply, replyHandler)
.toolCall(deprecatedTwin(fleetReply, "bridge_reply"), deprecatedHandler(fleetReply, "bridge_reply", replyHandler))
.toolCall(fleetAsk, askHandler)
.toolCall(deprecatedTwin(fleetAsk, "bridge_ask"), deprecatedHandler(fleetAsk, "bridge_ask", askHandler))
.toolCall(fleetStatus, statusHandler)
.toolCall(deprecatedTwin(fleetStatus, "bridge_status"), deprecatedHandler(fleetStatus, "bridge_status", statusHandler))
.toolCall(fleetPoll, pollHandler)
.toolCall(deprecatedTwin(fleetPoll, "bridge_poll"), deprecatedHandler(fleetPoll, "bridge_poll", pollHandler))
.toolCall(fleetAck, ackHandler)
.toolCall(deprecatedTwin(fleetAck, "bridge_ack"), deprecatedHandler(fleetAck, "bridge_ack", ackHandler))
.toolCall(fleetSpawn, spawnHandler)
.toolCall(deprecatedTwin(fleetSpawn, "bridge_spawn"), deprecatedHandler(fleetSpawn, "bridge_spawn", spawnHandler))
.toolCall(fleetList, listHandler)
.toolCall(deprecatedTwin(fleetList, "bridge_list"), deprecatedHandler(fleetList, "bridge_list", listHandler))
.toolCall(fleetStop, stopHandler)
.toolCall(deprecatedTwin(fleetStop, "bridge_stop"), deprecatedHandler(fleetStop, "bridge_stop", stopHandler))
.toolCall(fleetProfiles, profilesHandler)
.toolCall(deprecatedTwin(fleetProfiles, "bridge_profiles"), deprecatedHandler(fleetProfiles, "bridge_profiles", profilesHandler))
.toolCall(fleetWhoami, whoamiHandler)
.toolCall(deprecatedTwin(fleetWhoami, "bridge_whoami"), deprecatedHandler(fleetWhoami, "bridge_whoami", whoamiHandler))
.build();
this.authz = callers;
this.metrics = metrics;
@@ -406,7 +469,7 @@ public final class BridgeMcp {
// --- tool logic (thin adapters over the services; unit-testable) ---------------------------
/**
* {@code bridge_send}: delegate {@code content} to a worker session and block for its reply.
* {@code fleet_send}: delegate {@code content} to a worker session and block for its reply.
* The configured profiles are required so a profile name can never bypass target validation.
*
* (CB-548): {@code onAccepted} records delegator ownership the instant the send is accepted, so
@@ -430,8 +493,8 @@ public final class BridgeMcp {
}
/**
* {@code bridge_send} carrying a {@code turnId}: the primary's answer to a worker's
* {@code bridge_ask} (CB-205). Resolves the worker's blocked question and blocks for its reply as
* {@code fleet_send} carrying a {@code turnId}: the primary's answer to a worker's
* {@code fleet_ask} (CB-205). Resolves the worker's blocked question and blocks for its reply as
* it resumes the same turn — surfaced to the primary identically to a normal send.
*/
static McpSchema.CallToolResult answer(MessageService messages, String turnId, String content, Long timeoutMs) {
@@ -443,13 +506,13 @@ public final class BridgeMcp {
}
/**
* {@code bridge_ask} (CB-205): a worker pauses its delegated turn to ask the primary, blocking
* {@code fleet_ask} (CB-205): a worker pauses its delegated turn to ask the primary, blocking
* until the primary answers. The worker is identified by its connection ({@code callerTerminal}),
* never an argument — a {@code null} means the caller is not a known worker.
*/
static McpSchema.CallToolResult ask(MessageService messages, String callerTerminal, String question, Long timeoutMs) {
if (callerTerminal == null) {
return error("bridge_ask is for workers only — could not identify the calling worker "
return error("fleet_ask is for workers only — could not identify the calling worker "
+ "from the connection");
}
if (isBlank(question)) {
@@ -459,10 +522,10 @@ public final class BridgeMcp {
MessageService.AskResult r = messages.ask(callerTerminal, question, timeout);
return switch (r.outcome()) {
case ANSWERED -> text(r.answer());
case NO_WAITER -> error("no primary is awaiting this turn — bridge_ask only works while a "
+ "bridge_send delegation is open to answer it");
case NO_WAITER -> error("no primary is awaiting this turn — fleet_ask only works while a "
+ "fleet_send delegation is open to answer it");
case TIMED_OUT -> text("[no answer within " + timeout + "ms — the primary did not respond; "
+ "proceed on your best judgement, then call bridge_reply to end the turn]");
+ "proceed on your best judgement, then call fleet_reply to end the turn]");
};
}
@@ -470,10 +533,10 @@ public final class BridgeMcp {
private static McpSchema.CallToolResult formatReply(MessageService.Reply r, long timeout) {
return switch (r.outcome()) {
case REPLIED -> text(r.text());
// The worker's turn finished but it never called bridge_reply — hand back the scraped
// The worker's turn finished but it never called fleet_reply — hand back the scraped
// transcript tail, flagged so the primary knows it isn't a structured reply.
case COMPLETED_UNREPLIED -> text(
"[worker finished without a structured bridge_reply — transcript tail follows]\n" + r.text());
"[worker finished without a structured fleet_reply — transcript tail follows]\n" + r.text());
// The worker ran the turn then wedged (CB-109) — surface the error context.
case WORKER_FAILED -> text("[worker failed — turn ended in an unrecoverable state]\n" + r.text());
// The backend refused on a subscription usage limit (CB-578 stage A) — the worker's
@@ -483,7 +546,7 @@ public final class BridgeMcp {
+ "usage limit]\n" + r.text());
// The worker paused mid-turn to ask (CB-205) — tell the primary how to answer in-turn.
case QUESTION -> text("[question] the worker paused to ask before it can finish:\n" + r.text()
+ "\n\nAnswer it by calling bridge_send again with turnId=\"" + r.turnId()
+ "\n\nAnswer it by calling fleet_send again with turnId=\"" + r.turnId()
+ "\" and content set to your answer; the worker resumes the same turn.");
case STALE_TURN -> error("that question is no longer open — it timed out or was already "
+ "answered (turnId stale)");
@@ -493,7 +556,7 @@ public final class BridgeMcp {
}
/**
* {@code bridge_send} with {@code wait:false}: delegate {@code content} and return a ticket
* {@code fleet_send} with {@code wait:false}: delegate {@code content} and return a ticket
* immediately (fire-and-poll), so a long task isn't cut off by the caller's MCP call timeout.
* The configured profiles are required so a profile name can never bypass target validation.
*
@@ -510,19 +573,19 @@ public final class BridgeMcp {
return targetError;
}
String ticket = messages.sendAsync(sessionId, content, onAccepted);
return text("accepted — task delegated. Poll bridge_poll with ticket=" + ticket);
return text("accepted — task delegated. Poll fleet_poll with ticket=" + ticket);
}
/** A configured profile is never a send target; other unknown values may be herdr-owned panes. */
private static McpSchema.CallToolResult profileTargetError(String sessionId, Set<String> profiles) {
if (profiles.contains(sessionId)) {
return error("unknown send target \"" + sessionId + "\": it is a configured profile name, not a "
+ "session id. Call bridge_list to find a member or lead sessionId.");
+ "session id. Call fleet_list to find a member or lead sessionId.");
}
return null;
}
/** {@code bridge_poll}: check an async delegation by ticket, or drain a worker's inbox by target. */
/** {@code fleet_poll}: check an async delegation by ticket, or drain a worker's inbox by target. */
static McpSchema.CallToolResult poll(MessageService messages, String ticket, String target) {
if (!isBlank(target)) {
var replies = messages.drainReplies(target);
@@ -540,25 +603,25 @@ public final class BridgeMcp {
}
return switch (v.phase()) {
case DONE -> text(v.replySource() != null && v.replySource().equals("transcript")
? "[done — worker finished without a structured bridge_reply; transcript tail follows]\n" + v.reply()
? "[done — worker finished without a structured fleet_reply; transcript tail follows]\n" + v.reply()
: v.reply());
case PENDING -> text("[pending — " + v.detail() + "]");
case ASKING -> text("[question — worker is waiting for your answer]\n" + v.reply()
+ "\n\nAnswer it by calling bridge_send again with turnId=\"" + v.turnId()
+ "\n\nAnswer it by calling fleet_send again with turnId=\"" + v.turnId()
+ "\" and content set to your answer; the worker resumes the same turn.");
case FAILED -> text("[failed — " + v.detail() + "]");
};
}
/**
* {@code bridge_reply}: the worker returns its structured answer, resolving the awaiting send
* {@code fleet_reply}: the worker returns its structured answer, resolving the awaiting send
* or — when no send is open — queueing the reply in the inbox for later drain (CB-307).
* {@code callerTerminal} is resolved from the connection (never an argument); a {@code null}
* means the caller is not a known worker (e.g. the primary called it by mistake).
*/
static McpSchema.CallToolResult reply(MessageService messages, String callerTerminal, String content) {
if (callerTerminal == null) {
return error("bridge_reply is for workers only — could not identify the calling worker "
return error("fleet_reply is for workers only — could not identify the calling worker "
+ "from the connection");
}
if (content == null) {
@@ -568,7 +631,7 @@ public final class BridgeMcp {
return text("delivered");
}
/** {@code bridge_ack}: acknowledge (remove) a specific reply from the inbox. */
/** {@code fleet_ack}: acknowledge (remove) a specific reply from the inbox. */
static McpSchema.CallToolResult ack(MessageService messages, String target, String msgId) {
if (isBlank(target) || isBlank(msgId)) {
return error("target and msgId are required");
@@ -578,8 +641,8 @@ public final class BridgeMcp {
}
/**
* {@code bridge_status}: the live lifecycle status of a worker session, plus — when the worker
* is paused mid-turn in an async {@code bridge_ask} (CB-582) — the open question and how to
* {@code fleet_status}: the live lifecycle status of a worker session, plus — when the worker
* is paused mid-turn in an async {@code fleet_ask} (CB-582) — the open question and how to
* answer it, so a lead on its normal poll cadence does not need the ticket to notice.
*/
static McpSchema.CallToolResult status(MessageService messages, String sessionId) {
@@ -593,7 +656,7 @@ public final class BridgeMcp {
return text(base);
}
return text(base + "\n\n[question — worker is waiting for your answer]\n" + ask.question()
+ "\n\nAnswer it by calling bridge_send again with turnId=\"" + ask.turnId()
+ "\n\nAnswer it by calling fleet_send again with turnId=\"" + ask.turnId()
+ "\" and content set to your answer; the worker resumes the same turn."
+ " (ticket " + ask.ticket() + ")");
} catch (HerdrException e) {
@@ -602,7 +665,7 @@ public final class BridgeMcp {
}
/**
* {@code bridge_whoami}: the caller's own identity, as the daemon already resolved it.
* {@code fleet_whoami}: the caller's own identity, as the daemon already resolved it.
*
* <p>Every other tool <em>consumes</em> this identity — the authorization gate, the reply
* rendezvous, the cwd inherit — but none reported it, so an agent had to infer its own role
@@ -610,7 +673,7 @@ public final class BridgeMcp {
* name its MCP mount happens to carry, or {@code ANTHROPIC_BASE_URL} (which Claude-model
* workers do not set). The failure mode of guessing is asymmetric and silent: a primary that
* mistakes itself for a worker is refused by {@link Authz} and learns immediately, while a
* worker that mistakes itself for the primary ends its turn without {@code bridge_reply} and
* worker that mistakes itself for the primary ends its turn without {@code fleet_reply} and
* the sender simply receives nothing. This tool removes the guess.
*
* <p>For a worker the session registry adds what it knows about that session. A worker the
@@ -668,13 +731,13 @@ public final class BridgeMcp {
// --- fleet management logic (CB-108 / CB-301) --------------------------------------------
/** {@code bridge_spawn} without cwd/caller context (default resolution). */
/** {@code fleet_spawn} without cwd/caller context (default resolution). */
static McpSchema.CallToolResult spawn(SessionManager sessions, String profile) {
return spawn(sessions, profile, null, null, null, null, null, null, null);
}
/**
* {@code bridge_spawn}: launch a guard-checked member for {@code profile} (blank → the default
* {@code fleet_spawn}: launch a guard-checked member for {@code profile} (blank → the default
* profile) under {@code role} (blank → {@code dev}), and return its session id + pane id. The
* member's cwd is {@code requestedCwd} if given, else the profile's config, else
* {@code callerCwd} (the primary's directory), else the daemon's.
@@ -718,7 +781,7 @@ public final class BridgeMcp {
}
}
/** Build a {@link WorktreeRequest} from {@code bridge_spawn}'s optional {@code worktree}/{@code ticket} args. */
/** Build a {@link WorktreeRequest} from {@code fleet_spawn}'s optional {@code worktree}/{@code ticket} args. */
private static WorktreeRequest worktreeRequest(Map<String, Object> a) {
Object w = a.get("worktree");
if (w == null || Boolean.FALSE.equals(w)) {
@@ -747,7 +810,7 @@ public final class BridgeMcp {
}
/**
* {@code bridge_profiles}: the configured worker profiles, the default, and — CB-578 stage B —
* {@code fleet_profiles}: the configured worker profiles, the default, and — CB-578 stage B —
* which of them are currently quarantined (backend exhausted) and for how much longer. The
* {@code quarantined} key is present only when at least one profile is, so a fleet where nothing
* has ever been quarantined gets exactly the pre-stage-B shape.
@@ -776,7 +839,7 @@ public final class BridgeMcp {
}
/**
* {@code bridge_list}: the whole fleet — {@code leads} and {@code workers} — each merged with
* {@code fleet_list}: the whole fleet — {@code leads} and {@code workers} — each merged with
* live herdr status. CB-519 decoupled the registry key (a host-unique id) from the herdr pane
* coordinate, so the join is on the terminal id, which both the session and the live agent carry.
*
@@ -789,10 +852,10 @@ public final class BridgeMcp {
* <p>Leads are drawn from the resolver rather than from a second registry, so an address listed
* here is one that would actually resolve as a lead — see {@link CallerResolver#leads()}. The
* caller's own row is flagged {@code "self": true}: a peer needs to tell its own pane apart from
* a peer's, and the alternative is every lead calling {@code bridge_whoami} to subtract itself.
* a peer's, and the alternative is every lead calling {@code fleet_whoami} to subtract itself.
*
* <p>CB-583: the {@code capacity} rows reuse {@code quarantine} (the same {@link QuarantineSource}
* {@code bridge_profiles} reads) so the two surfaces cannot disagree about which profile is
* {@code fleet_profiles} reads) so the two surfaces cannot disagree about which profile is
* quarantined — see {@link #capacityView}.
*
* @param leads terminal_id → lead name, live from the resolver
@@ -856,7 +919,7 @@ public final class BridgeMcp {
* CB-583: {@code free} alone cannot tell a lead "busy, will free up" from "refusing, and
* nothing changes for N seconds" — those need different decisions. So a quarantined profile
* forces {@code free} to 0, whatever its {@code maxLoad}/{@code live} say, and the row carries
* the same {@code credentialId}/{@code quarantinedForSeconds} facts {@code bridge_profiles}
* the same {@code credentialId}/{@code quarantinedForSeconds} facts {@code fleet_profiles}
* reports, reusing {@link QuarantineSource} rather than a second lookup. Both new keys are
* added only when the profile is actually quarantined, so an ordinary fleet's rows are
* byte-identical to before this change.
@@ -889,7 +952,7 @@ public final class BridgeMcp {
*
* <p>{@code status} is herdr's live view, and {@code unknown} when herdr is not tracking that
* pane as an agent — the honest answer, and the one that matters: a lead whose pane herdr cannot
* see is a lead a {@code bridge_send} cannot be typed into. It is reported rather than hidden,
* see is a lead a {@code fleet_send} cannot be typed into. It is reported rather than hidden,
* because a peer that has gone unreachable is exactly what the sender needs to know.
*/
private static Map<String, Object> leadView(String terminal, String name, Agent live,
@@ -905,7 +968,7 @@ public final class BridgeMcp {
return m;
}
/** {@code bridge_stop}: tear a worker down by its pane id. */
/** {@code fleet_stop}: tear a worker down by its pane id. */
static McpSchema.CallToolResult stop(SessionManager sessions, String paneId) {
if (isBlank(paneId)) {
return error("paneId is required");
@@ -948,25 +1011,25 @@ public final class BridgeMcp {
// --- tool schemas --------------------------------------------------------------------------
private static McpSchema.Tool sendTool() {
return tool("bridge_send",
return tool("fleet_send",
"Delegate a task to a worker session. By default blocks until the worker replies and "
+ "returns its reply (or a 'still working / queued' note on timeout). Pass wait:false "
+ "for a long task to return a ticket immediately, then poll it with bridge_poll. To "
+ "answer a worker's bridge_ask, pass its turnId (with content) instead of sessionId.",
+ "for a long task to return a ticket immediately, then poll it with fleet_poll. To "
+ "answer a worker's fleet_ask, pass its turnId (with content) instead of sessionId.",
objectSchema(Map.of(
"sessionId", stringProp("The worker session id (herdr terminal_id) to delegate to"),
"content", stringProp("The task/message to send to the worker (or your answer, with turnId)"),
"timeoutMs", Map.of("type", "integer", "description", "Max ms to wait for a reply (blocking mode)"),
"wait", Map.of("type", "boolean",
"description", "Block for the reply (default true); false returns a ticket to poll"),
"turnId", stringProp("When answering a worker's bridge_ask, its question turnId — "
"turnId", stringProp("When answering a worker's fleet_ask, its question turnId — "
+ "routes your answer back into the same turn (omit for a normal delegation)")),
List.of("content")));
}
private static McpSchema.Tool askTool() {
// No target/session arg — the worker's identity is resolved from the connection.
return tool("bridge_ask",
return tool("fleet_ask",
"Pause your current delegated turn to ask the primary a question, blocking until it "
+ "answers — then resume the same turn with the answer. Use this when only the "
+ "primary has a decision or detail you need to continue. You do not address the "
@@ -979,19 +1042,19 @@ public final class BridgeMcp {
}
private static McpSchema.Tool pollTool() {
return tool("bridge_poll",
"Check an async delegation (a bridge_send with wait:false) by its ticket: "
return tool("fleet_poll",
"Check an async delegation (a fleet_send with wait:false) by its ticket: "
+ "pending, done (with the worker's reply), or failed. When target (a worker "
+ "session id) is present instead of ticket, drain that worker's inbox of "
+ "replies delivered when no send was open.",
objectSchema(Map.of(
"ticket", stringProp("The ticket returned by bridge_send wait:false"),
"ticket", stringProp("The ticket returned by fleet_send wait:false"),
"target", stringProp("Worker session id to drain pending replies from (optional)")),
List.of()));
}
private static McpSchema.Tool ackTool() {
return tool("bridge_ack",
return tool("fleet_ack",
"Acknowledge (remove) a specific reply from a worker's inbox. Use when the primary "
+ "has processed a reply and wants to confirm it, leaving other pending replies "
+ "in the inbox for later drain.",
@@ -1002,22 +1065,22 @@ public final class BridgeMcp {
}
private static McpSchema.Tool spawnTool() {
return tool("bridge_spawn",
return tool("fleet_spawn",
"Spawn a new off-subscription member session. A member has two independent attributes: "
+ "role (what it is for) and profile (which backend it runs on). Pass role to pick "
+ "the contract — 'dev' implements a unit and opens its own PR, 'reviewer' reviews a "
+ "diff it did not write, 'architect' refines a ticket before anyone builds it; omit "
+ "it for 'dev'. Pass profile (from bridge_profiles) to pick the backend, or omit it "
+ "it for 'dev'. Pass profile (from fleet_profiles) to pick the backend, or omit it "
+ "for the default. The two are independent: a reviewer may run on the same profile "
+ "as the dev it reviews. The member opens your current directory by default; pass "
+ "cwd to pin a different one. Pass worktree:true (with ticket) or "
+ "worktree:<ticket-slug> to provision an isolated git worktree. Pass resumeSessionId "
+ "to relaunch onto a prior conversation instead of starting cold — this requires an "
+ "explicit profile whose backend supports it (bridge_list shows agentSessionId for "
+ "explicit profile whose backend supports it (fleet_list shows agentSessionId for "
+ "resumable members), and is refused otherwise rather than silently starting fresh. "
+ "sessionName gives the member a display name in its own UI when the backend supports "
+ "one. Returns the member's sessionId (use with bridge_send) and paneId (use with "
+ "bridge_stop).",
+ "one. Returns the member's sessionId (use with fleet_send) and paneId (use with "
+ "fleet_stop).",
objectSchema(Map.of(
"role", stringProp("What the member is for: architect, dev or reviewer (default dev)"),
"profile", stringProp("Which backend to run it on (omit for the default profile)"),
@@ -1025,33 +1088,33 @@ public final class BridgeMcp {
"worktree", Map.of("type", "string", "description", "'true' or a ticket slug — requests an isolated git worktree"),
"ticket", stringProp("Ticket slug when worktree:true"),
"sessionName", stringProp("Logical display name for the member's own session, when its backend supports one"),
"resumeSessionId", stringProp("A prior member's agentSessionId (from bridge_list) to resume — requires an explicit profile that supports it")),
"resumeSessionId", stringProp("A prior member's agentSessionId (from fleet_list) to resume — requires an explicit profile that supports it")),
List.of()));
}
private static McpSchema.Tool profilesTool() {
return tool("bridge_profiles",
"List the configured worker profiles (backends) and which one bridge_spawn uses by "
return tool("fleet_profiles",
"List the configured worker profiles (backends) and which one fleet_spawn uses by "
+ "default. A 'quarantined' map is present when a backend-exhausted refusal put "
+ "a profile's credential on cooldown — bridge_spawn onto it is refused until "
+ "a profile's credential on cooldown — fleet_spawn onto it is refused until "
+ "quarantinedForSeconds elapses; a profile sharing that credential is listed too.",
objectSchema(Map.of(), List.of()));
}
private static McpSchema.Tool listTool() {
return tool("bridge_list",
return tool("fleet_list",
"List the whole fleet the bridge tracks, in two parts. 'leads' are your PEERS — other "
+ "orchestrators, each with its sessionId (the address to bridge_send to), "
+ "orchestrators, each with its sessionId (the address to fleet_send to), "
+ "name, live status, and 'self': true on your own row; this is how you "
+ "discover a peer lead without being told its address. 'members' are the "
+ "sessions delegated to — each with sessionId, paneId, role (architect/dev/"
+ "reviewer), profile (the backend it runs on), state, optional "
+ "worktree/branch/owner/agentSessionId (the id to pass as bridge_spawn's "
+ "worktree/branch/owner/agentSessionId (the id to pass as fleet_spawn's "
+ "resumeSessionId to relaunch onto that same conversation, when the backend "
+ "supports it), and live herdr status. An empty 'members' "
+ "means no members are spawned; it says nothing about peers. When capacity "
+ "facts are configured, a 'capacity' row per profile also reports free: 0 for "
+ "a quarantined profile's credential (see bridge_profiles), whatever its "
+ "a quarantined profile's credential (see fleet_profiles), whatever its "
+ "maxLoad/live — with credentialId and quarantinedForSeconds naming the "
+ "quarantine, so 'free: 0, busy' can be told apart from 'free: 0, refusing "
+ "for N seconds'.",
@@ -1059,8 +1122,8 @@ public final class BridgeMcp {
}
private static McpSchema.Tool stopTool() {
return tool("bridge_stop",
"Tear down a worker session by its paneId (from bridge_spawn or bridge_list).",
return tool("fleet_stop",
"Tear down a worker session by its paneId (from fleet_spawn or fleet_list).",
objectSchema(Map.of(
"paneId", stringProp("The worker's paneId to stop")),
List.of("paneId")));
@@ -1068,9 +1131,9 @@ public final class BridgeMcp {
private static McpSchema.Tool replyTool() {
// No session/target arg — the caller's identity is resolved from the connection.
return tool("bridge_reply",
return tool("fleet_reply",
"Return your structured answer for a message you were sent, resolving the sender's "
+ "blocked bridge_send. A worker MUST end every delegated turn with exactly "
+ "blocked fleet_send. A worker MUST end every delegated turn with exactly "
+ "one of these. A lead uses it only to answer another lead that messaged "
+ "it — never to answer a worker, whose turn it is not.",
objectSchema(Map.of(
@@ -1079,7 +1142,7 @@ public final class BridgeMcp {
}
private static McpSchema.Tool statusTool() {
return tool("bridge_status",
return tool("fleet_status",
"Get the live lifecycle status (idle/working/blocked/unknown) of a worker session.",
objectSchema(Map.of(
"sessionId", stringProp("The worker session id to query")),
@@ -1087,20 +1150,58 @@ public final class BridgeMcp {
}
private static McpSchema.Tool whoamiTool() {
return tool("bridge_whoami",
return tool("fleet_whoami",
"Report who YOU are on the bridge — your role is resolved from your connection "
+ "(unforgeable), never from anything you claim. Returns role 'primary' (you "
+ "orchestrate: spawn/send/stop; reply ONLY to answer a peer lead that "
+ "messaged you, never to answer a worker), 'architect' (you delegate turns "
+ "and reply/ask as your own pane, but cannot spawn/stop/drain), or 'worker' "
+ "(you were delegated to: you must end every turn with exactly one "
+ "bridge_reply, and cannot spawn or send), plus 'leader'/'architect' naming "
+ "fleet_reply, and cannot spawn or send), plus 'leader'/'architect' naming "
+ "which one you are, your own sessionId, and profile/worktree/branch when "
+ "you are a worker. Call this first when following role-conditional "
+ "instructions rather than guessing.",
objectSchema(Map.of(), List.of()));
}
// --- CB-622: bridge_* -> fleet_* rename, kept working under both names -------------------
/**
* The deprecated {@code bridge_*} twin of {@code fleetTool}: same name-minus-prefix schema,
* with a description that leads with the deprecation notice so a client listing tools sees it
* immediately. Reuses {@code fleetTool}'s input schema rather than restating it, so the two
* can never drift on parameters.
*/
static McpSchema.Tool deprecatedTwin(McpSchema.Tool fleetTool, String oldName) {
return tool(oldName, "DEPRECATED: use " + fleetTool.name() + " instead. " + fleetTool.description(),
fleetTool.inputSchema());
}
/**
* Wrap {@code handler} so a call under the deprecated {@code oldName} logs one WARN naming
* the old and new name, then runs the exact SAME handler {@code fleetTool}'s name uses — no
* logic is duplicated between the two registrations.
*/
static BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> deprecatedHandler(
McpSchema.Tool fleetTool, String oldName,
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> handler) {
return (exchange, req) -> {
warnDeprecatedOnce(oldName, fleetTool.name());
return handler.apply(exchange, req);
};
}
/**
* Log one WARN naming {@code oldName} and {@code newName} — once per {@code oldName} for the
* life of the process, not once per call. {@link #WARNED_DEPRECATED_NAMES} is a per-name flag
* (a {@link Set}), not a call counter, so this never conflates "warned" with any numeric state.
*/
static void warnDeprecatedOnce(String oldName, String newName) {
if (WARNED_DEPRECATED_NAMES.add(oldName)) {
log.warn("{} is deprecated; use {} instead", oldName, newName);
}
}
// --- small helpers -------------------------------------------------------------------------
// The SDK 2.0.0 deprecates its own Tool builders without a stable replacement — isolate it here.
@@ -11,7 +11,7 @@ import java.util.concurrent.atomic.AtomicReference;
* Single-slot, thread-safe registry for the primary's herdr {@code terminal_id}.
*
* <p>Populated from the caller terminal of orchestration-side MCP tools
* ({@code bridge_send}, {@code bridge_spawn}) — tools that only the primary calls.
* ({@code fleet_send}, {@code fleet_spawn}) — tools that only the primary calls.
* A pinned terminal (from config) seeds the registry at construction and makes
* subsequent {@link #record(String)} calls no-ops.
*
@@ -29,7 +29,7 @@ public final class PrimaryRegistry {
/**
* CB-532: worker terminal → the lead that delegated to it. The single slot above answers "who is
* THE primary", a question with no correct answer once two leads orchestrate the same fleet:
* whichever called {@code bridge_send} first captured every nudge, including nudges for the
* whichever called {@code fleet_send} first captured every nudge, including nudges for the
* other lead's delegations. This map answers the question that actually matters — "who is
* waiting on THIS worker" — and is what lets {@code primary.terminal} be retired.
*/
@@ -71,7 +71,7 @@ public final class PrimaryRegistry {
*
* <p>Called from the {@code MessageService} accepted-delivery hook — only after a send has won
* the session's send lock and queued delivery — where both halves are known (CB-548). It is
* deliberately <em>not</em> called at {@code bridge_send} request time: a concurrent sender that
* deliberately <em>not</em> called at {@code fleet_send} request time: a concurrent sender that
* times out {@code BUSY} must not steal a live delegation's reply routing without ever owning
* the turn. Last writer wins — if a second lead's later send is accepted, replies follow the
* lead that most recently delegated to it, which is the one waiting.
@@ -214,7 +214,7 @@ public final class CompositePeerLauncher implements PeerLauncher {
}
}
}
// Order-preserving for the same reason, and because profiles() is user-visible (bridge_profiles).
// Order-preserving for the same reason, and because profiles() is user-visible (fleet_profiles).
this.byProfile = Collections.unmodifiableMap(index);
}
@@ -267,7 +267,7 @@ public final class CompositePeerLauncher implements PeerLauncher {
// CB-557: an unqualified spawn is placed inside the pool of the role it asked for, not across
// the whole profile list. An EXPLICIT profile (above) is left alone on purpose — it is the
// operator overriding, and refusing it would break `bridge_spawn{profile:"opus"}`, which
// operator overriding, and refusing it would break `fleet_spawn{profile:"opus"}`, which
// carries no role and so would be judged against the dev pool it was never meant for.
List<PlacementCandidate> candidates = candidates(req.role());
String roleDefault = defaultProfileFor(req.role());
@@ -117,14 +117,14 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
/** The final instruction always requires a bridge reply when the bridge MCP is mounted. */
protected static final String REPLY_CHARTER =
"You are a spawned member in the claude-bridge fleet. Every message you receive arrives "
+ "through the bridge, and the ONLY channel back to the sender is the bridge_reply MCP tool. "
+ "through the bridge, and the ONLY channel back to the sender is the fleet_reply MCP tool. "
+ "Text you write in your terminal is NOT sent anywhere — the sender cannot see your screen, "
+ "so an in-terminal answer is silently discarded. Therefore you MUST end EVERY turn by calling "
+ "bridge_reply with `content` set to your complete response. This holds for every message without "
+ "fleet_reply with `content` set to your complete response. This holds for every message without "
+ "exception — tasks, questions, clarifications, acknowledgements, and ordinary back-and-forth "
+ "conversation. Call bridge_reply exactly once, as the final action of your turn, with your full "
+ "conversation. Call fleet_reply exactly once, as the final action of your turn, with your full "
+ "answer in `content`; never wait for confirmation first. If you end a turn without calling "
+ "bridge_reply, the sender receives nothing and the exchange stalls.";
+ "fleet_reply, the sender receives nothing and the exchange stalls.";
/**
* Tab numbers, counted per {@code role/profile} pair (CB-557).
@@ -413,7 +413,7 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
: replyCharter == null ? roleCharter : roleCharter + "\n\n" + replyCharter;
// CB-571: fingerprint the exact composed charter bytes once, here in the base, before the
// string leaves for an adapter — so Claude and OpenCode derive the same digest. A failed
// start has no bridge_spawn result and no roster row, so the failure log below is the only
// start has no fleet_spawn result and no roster row, so the failure log below is the only
// surface the byte count can appear on. The charter text itself is never logged.
CharterReceipt receipt = CharterReceipt.compose(role, cfg.profile(), roleCharter, charter);
String cwd = resolveCwd(requestedCwd, cfg, callerCwd);
@@ -247,7 +247,7 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
* permissions opencode does not explicitly deny. It is unconditional, not a preference: a
* spawned peer has no human at its pane — the bridge spawned it — so one that stops at an
* approval prompt is a wedged agent, indistinguishable from a legitimate mid-turn wait and
* unable to end its turn with {@code bridge_reply}. opencode's own help calls this
* unable to end its turn with {@code fleet_reply}. opencode's own help calls this
* "dangerous!", but the blast radius here is already bounded by design: a worker runs in its
* own git worktree on its own branch, is off-subscription, and cannot merge — the lead is the
* gate.
@@ -18,7 +18,7 @@ import java.util.function.Supplier;
/**
* CB-551: an opt-in heartbeat that nudges the single idle lead back to work once it has been
* continuously idle past a quiet period with no open {@code bridge_send} driving it.
* continuously idle past a quiet period with no open {@code fleet_send} driving it.
*
* <p>Why this exists: the fleet is ONE lead + architects + workers, so an idle, stalled lead is a
* single point of failure for the fleet's progress. {@link ReplyPushLoop} nudges the lead only when
@@ -290,7 +290,7 @@ public final class LeadHeartbeatLoop {
/** The nudge body, phrased for the two cases the heartbeat distinguishes. */
String nudgeText() {
StringBuilder sb = new StringBuilder(
"Heartbeat: you are idle and no bridge_send is waiting on you.");
"Heartbeat: you are idle and no fleet_send is waiting on you.");
if (hasPending()) {
sb.append(" The fleet has state to collect: ").append(pendingDetail());
} else {
@@ -309,10 +309,10 @@ public final class LeadHeartbeatLoop {
.append(" pending collection");
if (!replyTargets.isEmpty()) {
// Render each as the exact command so the lead can act without parsing: the nearest
// analogue to ReplyPushLoop's bridge_poll(target=...) nudge.
// analogue to ReplyPushLoop's fleet_poll(target=...) nudge.
sb.append(" (")
.append(String.join(", ",
replyTargets.stream().map(t -> "bridge_poll(target=" + t + ")").toList()))
replyTargets.stream().map(t -> "fleet_poll(target=" + t + ")").toList()))
.append(")");
}
sb.append(", ").append(doneSessions).append(" DONE session").append(doneSessions == 1 ? "" : "s")
@@ -24,7 +24,7 @@ import java.util.function.LongSupplier;
/**
* The blocking delegation feature (CB-104): deliver {@code content} into a worker and block until
* the worker returns a <em>structured reply</em> via {@code bridge_reply} (the {@link Rendezvous}),
* the worker returns a <em>structured reply</em> via {@code fleet_reply} (the {@link Rendezvous}),
* then hand that reply back. Delivery is the {@link Injector}'s job (the background poller sends it
* when the worker is injectable); this service never drives the injector or scrapes the terminal —
* completion is the worker's explicit reply, not a guess about {@code agent_status}.
@@ -62,10 +62,10 @@ public final class MessageService {
/** Outcome of a blocking send. */
public enum Outcome {
/** The worker called {@code bridge_reply}; {@code text} holds the structured answer. */
/** The worker called {@code fleet_reply}; {@code text} holds the structured answer. */
REPLIED,
/**
* The worker's delegated turn finished without a {@code bridge_reply} (CB-106 fallback);
* The worker's delegated turn finished without a {@code fleet_reply} (CB-106 fallback);
* {@code text} is the scraped transcript tail rather than a structured answer.
*/
COMPLETED_UNREPLIED,
@@ -75,7 +75,7 @@ public final class MessageService {
*/
WORKER_FAILED,
/**
* The turn finished without a {@code bridge_reply} and the scrape matched the backend's
* The turn finished without a {@code fleet_reply} and the scrape matched the backend's
* configured usage-limit refusal pattern (CB-578 stage A); {@code text} is the reason,
* carrying the matched line. The worker's pane is healthy — only its account is refusing —
* so this is never reported as a completed reply, and is kept distinct from
@@ -96,14 +96,14 @@ public final class MessageService {
BUSY,
/**
* An answer ({@link #answer(String, String, long)}) referenced a {@code turnId} that is no
* longer open — the worker's {@code bridge_ask} already timed out or was answered.
* longer open — the worker's {@code fleet_ask} already timed out or was answered.
*/
STALE_TURN
}
/**
* @param outcome how the send ended (or paused)
* @param text the worker's answer when {@link #completed()} (a structured {@code bridge_reply}
* @param text the worker's answer when {@link #completed()} (a structured {@code fleet_reply}
* for {@link Outcome#REPLIED}, a scraped transcript tail for
* {@link Outcome#COMPLETED_UNREPLIED}), or the question for {@link Outcome#QUESTION},
* else {@code null}
@@ -122,7 +122,7 @@ public final class MessageService {
}
}
/** How a worker's {@code bridge_ask} (CB-205) resolved. */
/** How a worker's {@code fleet_ask} (CB-205) resolved. */
public enum AskOutcome {
/** The primary answered; {@link AskResult#answer} carries it. */
ANSWERED,
@@ -132,7 +132,7 @@ public final class MessageService {
TIMED_OUT
}
/** The outcome of a worker's {@code bridge_ask}: how it resolved and (if answered) the answer. */
/** The outcome of a worker's {@code fleet_ask}: how it resolved and (if answered) the answer. */
public record AskResult(AskOutcome outcome, String answer) {
}
@@ -140,7 +140,7 @@ public final class MessageService {
public enum Phase {
/** Delegated and in flight — queued for the worker or being worked. */
PENDING,
/** The worker is paused in {@code bridge_ask}; {@link TaskView#reply} and {@link TaskView#turnId} identify it. */
/** The worker is paused in {@code fleet_ask}; {@link TaskView#reply} and {@link TaskView#turnId} identify it. */
ASKING,
/** The worker's turn finished; {@link TaskView#reply} holds the answer. */
DONE,
@@ -153,7 +153,7 @@ public final class MessageService {
*
* @param reply the answer when {@link #phase} is {@link Phase#DONE}, or the question when
* {@link #phase} is {@link Phase#ASKING}; otherwise {@code null}
* @param replySource {@code "reply"} (structured {@code bridge_reply}) or {@code "transcript"}
* @param replySource {@code "reply"} (structured {@code fleet_reply}) or {@code "transcript"}
* (completion scrape) when {@link Phase#DONE}, else {@code null}
* @param detail a human note (live worker status while pending, ask state, or failure reason)
* @param turnId correlation id for an {@link Phase#ASKING} ticket, else {@code null}
@@ -179,11 +179,11 @@ public final class MessageService {
}
/**
* A worker session's currently-open {@code bridge_ask} question, surfaced so {@code bridge_status}
* A worker session's currently-open {@code fleet_ask} question, surfaced so {@code fleet_status}
* can show it without the caller needing the ticket first (CB-582). Only covers async
* (fire-and-poll) delegations, which track the question on their {@link Task}; a blocking
* ({@code wait:true}) send already hands the question straight back to its own caller, so there is
* nothing hidden left for {@code bridge_status} to surface in that case.
* nothing hidden left for {@code fleet_status} to surface in that case.
*/
public record PendingAsk(String ticket, String question, String turnId) {
}
@@ -202,7 +202,7 @@ public final class MessageService {
/** Async task that owns each exact forward rendezvous waiter. */
private final ConcurrentHashMap<CompletableFuture<Rendezvous.Resolution>, Task> asyncTasksByWaiter =
new ConcurrentHashMap<>();
/** Async tickets paused on a specific {@code bridge_ask} turn. */
/** Async tickets paused on a specific {@code fleet_ask} turn. */
private final ConcurrentHashMap<String, Task> asyncTasksByTurn = new ConcurrentHashMap<>();
private final AtomicLong ticketSeq = new AtomicLong();
private final ExecutorService asyncExecutor = Executors.newThreadPerTaskExecutor(
@@ -216,7 +216,7 @@ public final class MessageService {
* (CB-588) whenever an async ticket started by {@link #sendAsync} reaches a
* terminal phase, whenever {@link #poll} hands a terminal ticket to its caller,
* and (CB-582) whenever an async ticket's worker pauses mid-turn in
* {@code bridge_ask} or that pause ends (answered or lapsed)
* {@code fleet_ask} or that pause ends (answered or lapsed)
*/
public MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous,
ReplyInbox inbox, ReplyPushLoop pushLoop) {
@@ -272,11 +272,11 @@ public final class MessageService {
}
/**
* Route a worker's explicit {@code bridge_reply}: resolve an open send, or queue it in the
* Route a worker's explicit {@code fleet_reply}: resolve an open send, or queue it in the
* inbox if no send is currently open. Unlike the bare {@link Rendezvous#resolve}, a no-waiter
* result is <em>not</em> a failure — the reply is held for later drain.
*
* <p><strong>Do NOT use this for mid-turn questions.</strong> {@code bridge_ask} /
* <p><strong>Do NOT use this for mid-turn questions.</strong> {@code fleet_ask} /
* {@link Rendezvous#resolveQuestion} must keep today's {@code NO_WAITER} behaviour — questions
* are interactive and must never be queued.
*
@@ -329,7 +329,7 @@ public final class MessageService {
* Abandon any send still waiting on {@code target} because its session has gone away (CB-516).
*
* <p>Without this, tearing a worker down left its rendezvous waiter open: a blocking
* {@code bridge_send} kept blocking, and an async one kept reporting {@code PENDING} until
* {@code fleet_send} kept blocking, and an async one kept reporting {@code PENDING} until
* {@link #ASYNC_TIMEOUT_MS} — thirty minutes — even though the worker provably no longer
* existed and the delegation could never complete. Worse, {@code poll} already had the evidence
* (it calls {@code liveStatus} to build its detail string and gets back {@code "unknown"}) and
@@ -484,7 +484,7 @@ public final class MessageService {
/**
* A worker's mid-turn question (CB-205 reverse rendezvous): surface {@code question} to the
* primary by resolving its open blocking {@code bridge_send}, then block this (worker) call until
* primary by resolving its open blocking {@code fleet_send}, then block this (worker) call until
* the primary answers via {@link #answer} or {@code timeoutMillis} elapses. Identity is the
* worker's own session — it does not address the primary.
*
@@ -509,7 +509,7 @@ public final class MessageService {
rendezvous.closeAsk(ticket.turnId());
return new AskResult(AskOutcome.NO_WAITER, null); // no primary is blocked on this worker
}
// CB-582: the question just became visible via bridge_poll (Phase.ASKING) for an async
// CB-582: the question just became visible via fleet_poll (Phase.ASKING) for an async
// (wait:false) delegation — nudge the lead's own pane the same way a terminal ticket does
// (CB-588), since the lead's normal poll cadence is minutes away and the reverse-rendezvous
// window (~55s, see BridgeMcp/BridgedApp) is far shorter. A blocking (wait:true) send has
@@ -523,7 +523,7 @@ public final class MessageService {
String answer = ticket.answer().get(timeoutMillis, TimeUnit.MILLISECONDS);
return new AskResult(AskOutcome.ANSWERED, answer);
} catch (TimeoutException e) {
log.debug("bridge_ask from {} went unanswered in {}ms", workerSession, timeoutMillis);
log.debug("fleet_ask from {} went unanswered in {}ms", workerSession, timeoutMillis);
clearAsyncQuestion(ticket.turnId(), true);
return new AskResult(AskOutcome.TIMED_OUT, null);
} catch (ExecutionException e) {
@@ -549,13 +549,13 @@ public final class MessageService {
}
/**
* The primary's answer to a worker's {@code bridge_ask} (CB-205): resolve the worker's blocked
* The primary's answer to a worker's {@code fleet_ask} (CB-205): resolve the worker's blocked
* question identified by {@code turnId}, then — like a fresh {@link #send} — block for the worker's
* eventual {@code bridge_reply} as it finishes the resumed turn. The worker session is derived from
* eventual {@code fleet_reply} as it finishes the resumed turn. The worker session is derived from
* {@code turnId}, never a caller argument.
*
* <p>Unlike {@link #send} this does not re-inject through the {@link Injector}: the worker is
* mid-turn (already picked up), so the answer flows back through its own open {@code bridge_ask}
* mid-turn (already picked up), so the answer flows back through its own open {@code fleet_ask}
* call, not a new status-gated delivery. The forward waiter is opened <em>before</em> the worker
* is unblocked so a reply that lands the instant it resumes is not lost.
*/
@@ -624,9 +624,9 @@ public final class MessageService {
tasks.put(ticket, task);
if (pushLoop != null) {
// CB-588: task.future only ever completes on a terminal phase (DONE or a failure) — a
// worker paused in bridge_ask leaves it running, per finishAsyncTask's own contract — so
// worker paused in fleet_ask leaves it running, per finishAsyncTask's own contract — so
// this fires exactly once, from whichever path completes it: finishAsyncTask(task, result)
// below on any non-QUESTION outcome of send() — a worker's bridge_reply, the CB-106
// below on any non-QUESTION outcome of send() — a worker's fleet_reply, the CB-106
// completion fallback, a CB-109 wedge, TIMED_OUT, BUSY, or BACKEND_EXHAUSTED — the same
// finishAsyncTask reached via answer()'s finishAsyncTask(turnId, result) once a QUESTION
// is resolved, completeExceptionally(t) just below when send() itself throws, or a CB-516
@@ -719,7 +719,7 @@ public final class MessageService {
* too, its own {@code pendingTickets} entry would outlive the ticket it names: an unpolled ticket
* (or one the reminder cap already gave up on) is pruned here but never collected there, so it
* lingers in {@code pendingTickets} forever and rides along on every later nudge to the same lead
* — naming a ticket {@code bridge_poll} can no longer find (CB-588 follow-up).
* — naming a ticket {@code fleet_poll} can no longer find (CB-588 follow-up).
*/
private void pruneTerminalTickets() {
long cutoff = nowNanos.getAsLong() - TICKET_TTL_NANOS;
@@ -784,10 +784,10 @@ public final class MessageService {
}
/**
* The question {@code workerSession} is currently paused on via {@code bridge_ask}, if any
* (CB-582) — {@code bridge_status} uses this to show a pending question without the caller
* The question {@code workerSession} is currently paused on via {@code fleet_ask}, if any
* (CB-582) — {@code fleet_status} uses this to show a pending question without the caller
* needing the ticket. {@code null} when the session has no open async question (including a
* session mid a <em>blocking</em> {@code bridge_ask}, which has no {@link Task} to look up — see
* session mid a <em>blocking</em> {@code fleet_ask}, which has no {@link Task} to look up — see
* {@link PendingAsk}).
*/
public PendingAsk pendingAsk(String workerSession) {
@@ -5,9 +5,9 @@ import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicLong;
/**
* The reply rendezvous: where a blocking {@code bridge_send} awaits how the worker's delegated turn
* The reply rendezvous: where a blocking {@code fleet_send} awaits how the worker's delegated turn
* ends. The sending (primary) request thread {@link #open}s a waiter; it is resolved either by the
* worker's explicit {@code bridge_reply} ({@link #resolve}, arriving on a different thread via
* worker's explicit {@code fleet_reply} ({@link #resolve}, arriving on a different thread via
* {@code POST /sessions/{id}/reply}) or — the CB-106 fallback — by the injector observing the
* worker's delegated turn return to idle without a reply ({@link #resolveCompletion}).
*
@@ -27,14 +27,14 @@ public final class Rendezvous {
/** How a delegated turn ended (or paused). */
public enum Kind {
/** The worker called {@code bridge_reply} with a structured answer. */
/** The worker called {@code fleet_reply} with a structured answer. */
REPLY,
/** The worker's turn finished without a {@code bridge_reply}; {@code text} is a scrape. */
/** The worker's turn finished without a {@code fleet_reply}; {@code text} is a scrape. */
COMPLETION,
/** The worker ran the turn then wedged (CB-109); {@code text} is the failure context. */
FAILED,
/**
* The turn finished without a {@code bridge_reply}, and the scrape matched the backend's
* The turn finished without a {@code fleet_reply}, and the scrape matched the backend's
* configured usage-limit refusal pattern (CB-578 stage A); {@code text} is the reason,
* carrying the matched line. The pane is healthy — only the account is refusing — so this
* is kept separate from a session simply going {@code GONE}.
@@ -43,7 +43,7 @@ public final class Rendezvous {
/**
* The worker paused mid-turn to ask the primary a question (CB-205 reverse rendezvous);
* {@code text} is the question and {@code turnId} correlates the primary's answer back to
* the worker's blocked {@code bridge_ask}. Not terminal — the turn resumes after the answer.
* the worker's blocked {@code fleet_ask}. Not terminal — the turn resumes after the answer.
*/
QUESTION
}
@@ -75,7 +75,7 @@ public final class Rendezvous {
/** Reverse rendezvous (CB-205): worker questions awaiting the primary's answer, keyed by {@code turnId}. */
private final ConcurrentHashMap<String, AskWaiter> asks = new ConcurrentHashMap<>();
private final AtomicLong askSeq = new AtomicLong();
/** Per-session index of the currently-open ask, so duplicate bridge_ask calls coalesce onto one turn. */
/** Per-session index of the currently-open ask, so duplicate fleet_ask calls coalesce onto one turn. */
private final ConcurrentHashMap<String, String> openAsksBySession = new ConcurrentHashMap<>();
/**
@@ -132,7 +132,7 @@ public final class Rendezvous {
return complete(session, new Resolution(Kind.REPLY, content));
}
// --- reverse rendezvous (CB-205 bridge_ask) ------------------------------------------------
// --- reverse rendezvous (CB-205 fleet_ask) ------------------------------------------------
/**
* Open a reverse-rendezvous waiter for a worker's mid-turn question. If {@code session} already has
@@ -166,7 +166,7 @@ public final class Rendezvous {
}
/**
* Surface a worker's mid-turn {@code question} by resolving the primary's open {@code bridge_send}
* Surface a worker's mid-turn {@code question} by resolving the primary's open {@code fleet_send}
* with a {@link Kind#QUESTION} carrying {@code turnId}. Same session-keyed semantics as
* {@link #resolve}: the one outstanding send for {@code session} unblocks with the question.
*
@@ -184,7 +184,7 @@ public final class Rendezvous {
}
/**
* Resolve a worker's blocked {@code bridge_ask} with the primary's {@code answer}, unblocking it
* Resolve a worker's blocked {@code fleet_ask} with the primary's {@code answer}, unblocking it
* to resume its turn.
*
* @return {@code true} if the ask was still open and got the answer; {@code false} if the
@@ -195,7 +195,7 @@ public final class Rendezvous {
return w != null && w.answer().complete(answer);
}
/** Drop a reverse-rendezvous turn once its {@code bridge_ask} has resolved (answered or lapsed). */
/** Drop a reverse-rendezvous turn once its {@code fleet_ask} has resolved (answered or lapsed). */
public void closeAsk(String turnId) {
AskWaiter w = asks.get(turnId);
if (w == null) {
@@ -209,9 +209,9 @@ public final class Rendezvous {
/**
* Resolve a specific captured {@code waiter} as a completion (the delegated turn finished with no
* {@code bridge_reply}); {@code text} is the scraped transcript tail. The waiter is the one
* {@code fleet_reply}); {@code text} is the scraped transcript tail. The waiter is the one
* captured when this turn was delivered, so a late completion for turn N cannot land on turn N+1's
* send (CB-116). A no-op if that waiter was already resolved — a raced {@code bridge_reply} wins.
* send (CB-116). A no-op if that waiter was already resolved — a raced {@code fleet_reply} wins.
*
* @return {@code true} if this call resolved the waiter, {@code false} if it was null or already resolved
*/
@@ -233,7 +233,7 @@ public final class Rendezvous {
/**
* Resolve a specific captured {@code waiter} as {@link Kind#BACKEND_EXHAUSTED} (CB-578 stage A):
* the turn finished with no {@code bridge_reply} and the scrape matched the backend's configured
* the turn finished with no {@code fleet_reply} and the scrape matched the backend's configured
* usage-limit pattern; {@code reason} carries the matched line. Like
* {@link #resolveCompletion(CompletableFuture, String)} it targets the exact captured send
* (CB-116). A no-op if that waiter was already resolved — first resolution wins.
@@ -19,9 +19,9 @@ import java.util.stream.Collectors;
/**
* A status-gated push loop that nudges a lead's own herdr pane when it has uncollected work
* waiting: a worker reply queued with no live {@code bridge_send} to resolve it (CB-307), an
* async delegation ticket ({@code bridge_send(wait:false)}) that reached a terminal phase
* (CB-588), or an async ticket's worker pausing mid-turn in {@code bridge_ask} to await an answer
* waiting: a worker reply queued with no live {@code fleet_send} to resolve it (CB-307), an
* async delegation ticket ({@code fleet_send(wait:false)}) that reached a terminal phase
* (CB-588), or an async ticket's worker pausing mid-turn in {@code fleet_ask} to await an answer
* (CB-582).
*
* <p><strong>CB-590: one schedule per lead.</strong> All three kinds of work are triggered
@@ -50,24 +50,24 @@ import java.util.stream.Collectors;
public final class ReplyPushLoop {
private static final Logger log = LoggerFactory.getLogger(ReplyPushLoop.class);
static final String NUDGE_FORMAT = "Worker %s returned a reply — run bridge_poll(target=%s) to collect it";
static final String NUDGE_FORMAT = "Worker %s returned a reply — run fleet_poll(target=%s) to collect it";
/** Coalesced form, several uncollected replies for the same lead. */
static final String REPLIES_NUDGE_FORMAT =
"%d workers returned replies — run bridge_poll(target=...) for each to collect them: %s";
"%d workers returned replies — run fleet_poll(target=...) for each to collect them: %s";
/** Singular form, one uncollected ticket. */
static final String TICKET_NUDGE_FORMAT =
"Ticket %s finished%s — run bridge_poll(ticket=%s) to collect it";
"Ticket %s finished%s — run fleet_poll(ticket=%s) to collect it";
/** Coalesced form, several uncollected tickets for the same lead. */
static final String TICKETS_NUDGE_FORMAT =
"%d tickets finished%s — run bridge_poll(ticket=...) for each to collect them: %s";
/** Singular form, one worker paused mid-turn in bridge_ask (CB-582) — names the answer call directly. */
"%d tickets finished%s — run fleet_poll(ticket=...) for each to collect them: %s";
/** Singular form, one worker paused mid-turn in fleet_ask (CB-582) — names the answer call directly. */
static final String QUESTION_NUDGE_FORMAT =
"Worker %s asked a question (ticket %s) — answer it with bridge_send(turnId=\"%s\", "
"Worker %s asked a question (ticket %s) — answer it with fleet_send(turnId=\"%s\", "
+ "content=...) to resume its turn:\n%s";
/** Coalesced form, several open questions for the same lead. */
static final String QUESTIONS_NUDGE_FORMAT =
"%d workers are paused on a question — run bridge_poll(ticket=...) for each, then answer "
+ "with bridge_send(turnId=..., content=...): %s";
"%d workers are paused on a question — run fleet_poll(ticket=...) for each, then answer "
+ "with fleet_send(turnId=..., content=...): %s";
private final PrimaryRegistry primaryRegistry;
private final AgentControl agents;
@@ -87,7 +87,7 @@ public final class ReplyPushLoop {
/** Tickets that have gone terminal but not yet been polled, keyed by ticket. */
private final ConcurrentHashMap<String, PendingTicket> pendingTickets = new ConcurrentHashMap<>();
/**
* Open {@code bridge_ask} questions not yet answered or lapsed, keyed by {@code turnId}
* Open {@code fleet_ask} questions not yet answered or lapsed, keyed by {@code turnId}
* (CB-582). A question's own nudge count is tracked the same per-item way as
* {@link #pendingTickets} (CB-598): a fresh question keeps its source eligible regardless of
* how depleted an older, still-open question's count is.
@@ -130,7 +130,7 @@ public final class ReplyPushLoop {
/**
* Reply targets still pending for {@code lead} — registered via {@link #onReplyQueued} and
* whose inbox still holds an unacked message. A target whose inbox has since drained (acked,
* or collected via a live {@code bridge_send} rendezvous instead) is dropped from
* or collected via a live {@code fleet_send} rendezvous instead) is dropped from
* {@link #pendingReplies} here rather than lingering forever; there is no explicit "reply
* collected" callback the way {@link #ticketCollected} exists for tickets, so the inbox itself
* is the only signal.
@@ -323,7 +323,7 @@ public final class ReplyPushLoop {
}
/**
* Called when an async delegation ticket ({@code bridge_send(wait:false)}, CB-107) reaches a
* Called when an async delegation ticket ({@code fleet_send(wait:false)}, CB-107) reaches a
* terminal phase — DONE or a failure. Unlike {@link #onReplyQueued}, which nudges about the
* durable-inbox no-waiter path, this covers the path {@code MessageService.reply} takes when a
* fire-and-poll send's own rendezvous waiter resolves the reply directly: that path returns
@@ -352,7 +352,7 @@ public final class ReplyPushLoop {
}
/**
* Called when a ticket's terminal state has been collected via {@code bridge_poll}. Removes it
* Called when a ticket's terminal state has been collected via {@code fleet_poll}. Removes it
* from the pending set so a scheduled tick — and any nudge it sends — never names a ticket the
* lead already has. A ticket that was never pending (unknown ticket, or one nudged with no push
* loop configured) is a no-op.
@@ -362,16 +362,16 @@ public final class ReplyPushLoop {
}
/**
* Called when an async ticket's worker pauses mid-turn in {@code bridge_ask} (CB-582): the
* question is now visible via {@code bridge_poll} (Phase.ASKING), but the reverse-rendezvous
* Called when an async ticket's worker pauses mid-turn in {@code fleet_ask} (CB-582): the
* question is now visible via {@code fleet_poll} (Phase.ASKING), but the reverse-rendezvous
* window it opened with (~55s default, see {@code BridgeMcp}/{@code BridgedApp}) is far shorter
* than a lead's normal minutes-long poll cadence — exactly the gap this closes. Resolves the
* delegating lead the same way {@link #onTicketTerminal} does and coalesces onto the same
* per-lead schedule (CB-590).
*
* @param ticket the async ticket the question belongs to (for {@code bridge_poll})
* @param ticket the async ticket the question belongs to (for {@code fleet_poll})
* @param target the worker session that asked
* @param turnId correlation id the lead answers with ({@code bridge_send turnId=...})
* @param turnId correlation id the lead answers with ({@code fleet_send turnId=...})
* @param question the question text
*/
public void onQuestionOpened(String ticket, String target, String turnId, String question) {
@@ -386,7 +386,7 @@ public final class ReplyPushLoop {
}
/**
* Called when a worker's {@code bridge_ask} resolves — answered or lapsed unanswered — so a
* Called when a worker's {@code fleet_ask} resolves — answered or lapsed unanswered — so a
* scheduled tick never nudges about a question the lead already handled. A {@code turnId} that
* was never pending (never nudged, or already closed) is a no-op.
*/
@@ -9,7 +9,7 @@ package dev.ltms.bridged.peer;
public enum Capability {
/**
* The peer supports {@code bridge_ask} rendezvous — pausing its delegated turn to ask
* The peer supports {@code fleet_ask} rendezvous — pausing its delegated turn to ask
* the primary a question, then resuming once answered. All Claude Code peers support this.
*/
MID_TURN_ASK,
@@ -11,7 +11,7 @@ package dev.ltms.bridged.placement;
* a usage limit, not a transient capacity or reachability concern.
* <li>Weight 0 (CB-554): {@code fixed} is still automatic selection, so a profile the operator
* marked "never auto-select me" ({@code weight <= 0}) must be skipped here exactly as
* {@code weighted}/{@code round-robin} skip it — an explicit {@code bridge_spawn} naming
* {@code weighted}/{@code round-robin} skip it — an explicit {@code fleet_spawn} naming
* the profile is unaffected, only this automatic fallback walk.
* </ul>
* A fleet where nothing is ever quarantined or weight-0 never exercises either path, so today's
@@ -26,7 +26,7 @@ public record PlacementCandidate(String profile, String host, float weight, Inte
* not need to distinguish "explicit 0" from "absent" itself.
*
* <p>Exclusion is about <em>automatic</em> selection only — an explicit
* {@code bridge_spawn{profile:"..."}} bypasses placement entirely and is unaffected.
* {@code fleet_spawn{profile:"..."}} bypasses placement entirely and is unaffected.
*/
public boolean excluded() {
return weight <= 0.0f;
@@ -46,7 +46,7 @@ public final class BridgedApp {
/** Default blocking window for a message; kept under typical HTTP idle timeouts. */
private static final long DEFAULT_MESSAGE_TIMEOUT_MS = 25_000;
private static final long MAX_MESSAGE_TIMEOUT_MS = 120_000;
/** Blocking window for a worker's bridge_ask (CB-205); the worker's MCP client caps its own call. */
/** Blocking window for a worker's fleet_ask (CB-205); the worker's MCP client caps its own call. */
private static final long DEFAULT_ASK_TIMEOUT_MS = 55_000;
private static final long MAX_ASK_TIMEOUT_MS = 115_000;
@@ -121,11 +121,11 @@ public final class BridgedApp {
app.get("/profiles", this::profiles); // configured backend profiles
app.post("/members", this::spawnMember); // optional ?role=&profile= or {"role":…,"profile":…}
app.delete("/members/{paneId}", this::stopMember);
app.post("/sessions/{id}/message", this::sendMessage); // bridge_send (primary; blocking, wait:false, or answer via turnId)
app.post("/sessions/{id}/reply", this::replyMessage); // bridge_reply (worker)
app.post("/sessions/{id}/message", this::sendMessage); // fleet_send (primary; blocking, wait:false, or answer via turnId)
app.post("/sessions/{id}/reply", this::replyMessage); // fleet_reply (worker)
app.get("/sessions/{id}/replies", this::drainReplies); // drain reply inbox (CB-307)
app.post("/sessions/{id}/ask", this::askMessage); // bridge_ask (worker → primary, CB-205)
app.get("/sessions/{id}/status", this::sessionStatus); // bridge_status
app.post("/sessions/{id}/ask", this::askMessage); // fleet_ask (worker → primary, CB-205)
app.get("/sessions/{id}/status", this::sessionStatus); // fleet_status
app.get("/tasks/{ticket}", this::taskStatus); // poll an async (wait:false) send
return app;
}
@@ -341,7 +341,7 @@ public final class BridgedApp {
/**
* The blocking delegation call (CB-104): inject {@code content} into the worker via the
* status-gated injector and block until the worker returns a structured {@code bridge_reply}.
* status-gated injector and block until the worker returns a structured {@code fleet_reply}.
* Times out with a typed 202 (working / queued / busy) rather than an error — the message may
* still land.
*/
@@ -370,7 +370,7 @@ public final class BridgedApp {
}
timeout = Math.clamp(timeout, 1, MAX_MESSAGE_TIMEOUT_MS);
// Answering a worker's bridge_ask (CB-205): always blocks, and derives the worker from turnId.
// Answering a worker's fleet_ask (CB-205): always blocks, and derives the worker from turnId.
if (turnId != null && !turnId.isBlank()) {
writeReply(ctx, id, messages.answer(turnId, content, timeout), timeout);
return;
@@ -392,7 +392,7 @@ public final class BridgedApp {
/**
* Render a {@link MessageService.Reply} onto the response — shared by a normal send and a
* bridge_ask answer. A structured/scraped completion is 200; a worker's mid-turn question a 202
* fleet_ask answer. A structured/scraped completion is 200; a worker's mid-turn question a 202
* (with its {@code turnId}); a stale answer a 409; every other non-terminal outcome a typed 202.
*/
private void writeReply(Context ctx, String id, MessageService.Reply reply, long timeout) {
@@ -404,7 +404,7 @@ public final class BridgedApp {
"sessionId", id, "error", "stale_turn",
"detail", "that question is no longer open (timed out or already answered)"));
case REPLIED, COMPLETED_UNREPLIED -> {
// replySource distinguishes a structured bridge_reply from the CB-106 completion
// replySource distinguishes a structured fleet_reply from the CB-106 completion
// fallback (a scrape of the worker's transcript when it finished without replying).
String source = reply.outcome() == MessageService.Outcome.REPLIED ? "reply" : "transcript";
ctx.status(200).json(Map.of("sessionId", id, "reply", reply.text(), "replySource", source));
@@ -428,7 +428,7 @@ public final class BridgedApp {
}
/**
* A worker's mid-turn question ({@code bridge_ask}, CB-205) — surfaces to the primary's open
* A worker's mid-turn question ({@code fleet_ask}, CB-205) — surfaces to the primary's open
* blocking send and blocks until it answers. 200 with the answer, 409 if no delegation is open,
* 202 if the primary stayed silent.
*/
@@ -465,7 +465,7 @@ public final class BridgedApp {
}
/**
* The worker's structured reply ({@code bridge_reply}) — resolves the blocking send awaiting
* The worker's structured reply ({@code fleet_reply}) — resolves the blocking send awaiting
* on this session, or queues the reply in the inbox when no send is open (CB-307).
*/
private void replyMessage(Context ctx) {
@@ -505,7 +505,7 @@ public final class BridgedApp {
}
/**
* Live lifecycle status of a worker (MCP `bridge_status` wraps this in CB-105), plus its
* Live lifecycle status of a worker (MCP `fleet_status` wraps this in CB-105), plus its
* <em>readiness</em> (CB-113): {@code ready} is true once the worker's Claude has connected the
* bridge MCP — the reliable "available to receive a task" signal, unlike bare {@code idle}, which
* is also true during boot.
@@ -520,8 +520,8 @@ public final class BridgedApp {
body.put("sessionId", id);
body.put("status", messages.status(id).name().toLowerCase());
body.put("ready", presence.isPresent(id));
// CB-582: a worker paused mid-turn in an async bridge_ask is otherwise invisible to a
// status poll — surface the open question and how to answer it, same as bridge_poll's
// CB-582: a worker paused mid-turn in an async fleet_ask is otherwise invisible to a
// status poll — surface the open question and how to answer it, same as fleet_poll's
// Phase.ASKING view.
MessageService.PendingAsk ask = messages.pendingAsk(id);
if (ask != null) {
@@ -55,7 +55,7 @@ class CallerResolverTest {
assertEquals("term_a", underTrust.terminal());
assertEquals(Role.WORKER, underToken.role(),
"worker identity is unforgeable and must never be token-gated — otherwise enabling "
+ "auth would lock the whole fleet out of bridge_reply");
+ "auth would lock the whole fleet out of fleet_reply");
assertEquals("term_a", underToken.terminal());
}
@@ -177,7 +177,7 @@ class CompletionResolverTest {
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
assertEquals("x".repeat(CompletionResolver.MAX_SCRAPE_CHARS)
+ "\n[Pane tail clipped: member did not call bridge_reply.]",
+ "\n[Pane tail clipped: member did not call fleet_reply.]",
waiter.getNow(null).text());
}
@@ -347,7 +347,7 @@ class CompletionResolverTest {
@Test
void aLateCompletionForOneTurnNeverResolvesTheNextTurnsWaiter() {
// The cross-turn stale reply the conversation test surfaced: turn N's completion fallback
// fires AFTER turn N was resolved by an explicit bridge_reply and turn N+1 has opened its own
// fires AFTER turn N was resolved by an explicit fleet_reply and turn N+1 has opened its own
// waiter on the same session. Resolving "whatever is waiting now" would hand turn N's stale
// scrape to turn N+1; targeting turn N's captured waiter makes the late completion a no-op.
FakeHerdr herdr = new FakeHerdr().readText("⏺ turn N answer\n❯ ");
@@ -163,7 +163,7 @@ class LeadLauncherTest {
/**
* The single most important assertion here. The worker charter tells its reader it is an
* off-subscription worker that must end every turn with bridge_reply — the opposite of what an
* off-subscription worker that must end every turn with fleet_reply — the opposite of what an
* orchestrator is. A lead must never receive it.
*/
@Test
@@ -174,7 +174,7 @@ class LeadLauncherTest {
List<String> args = startedArgs(herdr);
assertFalse(args.contains("--append-system-prompt"),
"the reply charter is a worker contract and must not be injected into a lead");
assertTrue(args.stream().noneMatch(a -> a.contains("bridge_reply")), args.toString());
assertTrue(args.stream().noneMatch(a -> a.contains("fleet_reply")), args.toString());
}
/** It still mounts the bridge — a lead that cannot orchestrate is pointless. */
@@ -19,16 +19,24 @@ import dev.ltms.bridged.member.ClaudeCodeLauncher;
import dev.ltms.bridged.member.CompositePeerLauncher;
import dev.ltms.bridged.placement.BackendQuarantine;
import dev.ltms.bridged.placement.PlacementPolicies;
import io.modelcontextprotocol.server.McpSyncServerExchange;
import io.modelcontextprotocol.spec.McpSchema;
import dev.ltms.bridged.msg.InMemoryReplyInbox;
import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.slf4j.LoggerFactory;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
import java.util.function.BiFunction;
import static org.junit.jupiter.api.Assertions.*;
@@ -83,7 +91,7 @@ class BridgeMcpTest {
@Test
void sendThenReplyRoundTrips() throws Exception {
// bridge_send blocks; bridge_reply resolves it with the worker's structured answer.
// fleet_send blocks; fleet_reply resolves it with the worker's structured answer.
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
() -> BridgeMcp.send(messages, "term_a", "review this", 4000L, null, Set.of()));
@@ -106,7 +114,7 @@ class BridgeMcpTest {
@Test
void asyncSendReturnsATicketThenPollReportsTheReply() throws Exception {
// wait:false parity — a ticket is issued, resolved by a reply, and surfaced by bridge_poll.
// wait:false parity — a ticket is issued, resolved by a reply, and surfaced by fleet_poll.
McpSchema.CallToolResult accepted = BridgeMcp.sendAsync(messages, "term_a", "do it", null, Set.of());
assertNotEquals(Boolean.TRUE, accepted.isError());
String out = textOf(accepted);
@@ -290,7 +298,7 @@ class BridgeMcpTest {
assertTrue(async.isError());
assertTrue(textOf(blocking).contains("sol"));
assertTrue(textOf(blocking).contains("configured profile name"));
assertTrue(textOf(blocking).contains("bridge_list"));
assertTrue(textOf(blocking).contains("fleet_list"));
assertFalse(textOf(async).contains("ticket="));
}
@@ -324,7 +332,7 @@ class BridgeMcpTest {
// A reply with no open send queues it in the inbox.
BridgeMcp.reply(messages, "term_a", "queued-msg");
// bridge_poll with target drains the inbox.
// fleet_poll with target drains the inbox.
McpSchema.CallToolResult res = BridgeMcp.poll(messages, null, "term_a");
assertNotEquals(Boolean.TRUE, res.isError());
String text = textOf(res);
@@ -359,7 +367,7 @@ class BridgeMcpTest {
String afterMarker = qt.substring(qt.indexOf("turnId=\"") + "turnId=\"".length());
String turnId = afterMarker.substring(0, afterMarker.indexOf('"'));
// The primary answers via bridge_send(turnId); this blocks again for the worker's reply.
// The primary answers via fleet_send(turnId); this blocks again for the worker's reply.
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
() -> BridgeMcp.answer(messages, turnId, "config.yaml", 5000L));
@@ -655,7 +663,7 @@ class BridgeMcpTest {
assertTrue(out.contains("\"name\":\"gpt-sol-5.6\""), out);
assertTrue(out.contains("\"sessionId\":\"term_peer\""), out);
// The caller's own row is flagged, and only the caller's — a peer must be distinguishable
// from self without a second bridge_whoami call.
// from self without a second fleet_whoami call.
assertEquals(1, out.split("\"self\":true", -1).length - 1, out);
assertTrue(out.indexOf("term_me") < out.indexOf("\"self\":true"), out);
}
@@ -730,7 +738,7 @@ class BridgeMcpTest {
assertEquals(1, before.size(), "one reply in the inbox");
String msgId = before.getFirst().msgId();
// Publish the same reply again and ack it via bridge_ack surface.
// Publish the same reply again and ack it via fleet_ack surface.
BridgeMcp.reply(messages, "term_a", "orphan-again");
var peeked = messages.drainReplies("term_a");
assertEquals(1, peeked.size(), "one fresh reply in the inbox");
@@ -774,8 +782,8 @@ class BridgeMcpTest {
}
/**
* CB-582: a lead polling {@code bridge_status} on its normal cadence — not {@code bridge_poll}
* — must also see a worker's open async {@code bridge_ask} question, since the reverse-rendezvous
* CB-582: a lead polling {@code fleet_status} on its normal cadence — not {@code fleet_poll}
* — must also see a worker's open async {@code fleet_ask} question, since the reverse-rendezvous
* window it opened with is far shorter than that cadence.
*/
@Test
@@ -819,7 +827,7 @@ class BridgeMcpTest {
answer.get(5, TimeUnit.SECONDS);
}
// --- bridge_whoami: the caller's own identity, so an agent never has to guess its role -------
// --- fleet_whoami: the caller's own identity, so an agent never has to guess its role -------
@Test
void whoamiReportsThePrimaryAsPrimaryAndNothingElse() {
@@ -994,7 +1002,7 @@ class BridgeMcpTest {
assertTrue(textOf(res).contains("architect, dev, reviewer"), textOf(res));
}
// ── CB-584: bridge_spawn accepts sessionName/resumeSessionId; roster shows agentSessionId ──
// ── CB-584: fleet_spawn accepts sessionName/resumeSessionId; roster shows agentSessionId ──
@Test
void spawnWithResumeSessionIdPutsTheIdOnTheRoster() {
@@ -1020,4 +1028,103 @@ class BridgeMcpTest {
assertEquals(Boolean.TRUE, res.isError());
assertTrue(textOf(res).contains("explicit profile"), textOf(res));
}
// ── CB-622: bridge_* -> fleet_* rename, both names answer through the SAME handler ────────
//
// These are unit tests of the wiring helpers (deprecatedTwin/deprecatedHandler/
// warnDeprecatedOnce), not a live MCP client call — proving the tool is REGISTERED and
// ROUTES correctly. Whether a real MCP client can actually invoke a tool by either name is
// the live check the lead runs; see CB-618's lesson that a green build here is not proof of
// that.
private static McpSchema.Tool fakeTool(String name) {
return McpSchema.Tool.builder(name)
.description("does the thing")
.inputSchema(Map.of("type", "object", "properties", Map.of(), "required", List.of()))
.build();
}
@Test
void deprecatedTwinNamesTheOldToolAndDefersToTheFleetDescription() {
McpSchema.Tool fleetTool = fakeTool("fleet_cb622_twin");
McpSchema.Tool twin = BridgeMcp.deprecatedTwin(fleetTool, "bridge_cb622_twin");
assertEquals("bridge_cb622_twin", twin.name());
assertTrue(twin.description().startsWith("DEPRECATED: use fleet_cb622_twin instead."),
twin.description());
assertTrue(twin.description().contains(fleetTool.description()), twin.description());
assertEquals(fleetTool.inputSchema(), twin.inputSchema(), "the twin must not restate the schema");
}
@Test
void oldNameAndNewNameReachTheExactSameHandler() {
McpSchema.Tool fleetTool = fakeTool("fleet_cb622_dual");
List<String> handlerCalls = new ArrayList<>();
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> handler =
(exchange, req) -> {
handlerCalls.add(req.name());
return McpSchema.CallToolResult.builder().addTextContent("handled:" + req.name()).build();
};
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> deprecated =
BridgeMcp.deprecatedHandler(fleetTool, "bridge_cb622_dual", handler);
// Calling the fleet_* name directly and calling the wrapped bridge_* name both end up
// running the SAME `handler` instance — not a copy of its logic.
McpSchema.CallToolResult viaFleet = handler.apply(null, new McpSchema.CallToolRequest("fleet_cb622_dual", Map.of()));
McpSchema.CallToolResult viaBridge = deprecated.apply(null, new McpSchema.CallToolRequest("bridge_cb622_dual", Map.of()));
assertEquals(textOf(viaFleet), "handled:fleet_cb622_dual");
assertEquals(textOf(viaBridge), "handled:bridge_cb622_dual");
assertEquals(2, handlerCalls.size(), "the same handler ran for both calls");
}
@Test
void deprecatedNameWarnsOnceForTheProcessNotOncePerCall() {
Logger logger = (Logger) LoggerFactory.getLogger(BridgeMcp.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
// A name unique to this test run so another test's use of the mechanism (or a re-run in
// the same JVM) cannot leave the "already warned" flag set before this assertion.
String oldName = "bridge_cb622_warnonce_" + System.identityHashCode(appender);
try {
BridgeMcp.warnDeprecatedOnce(oldName, "fleet_cb622_warnonce");
BridgeMcp.warnDeprecatedOnce(oldName, "fleet_cb622_warnonce");
BridgeMcp.warnDeprecatedOnce(oldName, "fleet_cb622_warnonce");
List<ILoggingEvent> matching = appender.list.stream()
.filter(e -> e.getFormattedMessage().contains(oldName))
.toList();
assertEquals(1, matching.size(), "three calls with the same old name must warn exactly once");
assertEquals(Level.WARN, matching.get(0).getLevel());
assertTrue(matching.get(0).getFormattedMessage().contains("fleet_cb622_warnonce"),
"the warning must name the new tool too");
} finally {
logger.detachAppender(appender);
}
}
@Test
void deprecatedNameWarnsAgainForADifferentOldName() {
Logger logger = (Logger) LoggerFactory.getLogger(BridgeMcp.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
String suffix = System.identityHashCode(appender) + "";
String nameA = "bridge_cb622_multi_a_" + suffix;
String nameB = "bridge_cb622_multi_b_" + suffix;
try {
BridgeMcp.warnDeprecatedOnce(nameA, "fleet_cb622_multi_a");
BridgeMcp.warnDeprecatedOnce(nameB, "fleet_cb622_multi_b");
long distinctNamesWarned = appender.list.stream()
.map(ILoggingEvent::getFormattedMessage)
.filter(m -> m.contains(nameA) || m.contains(nameB))
.count();
assertEquals(2, distinctNamesWarned, "each distinct old name gets its own warning");
} finally {
logger.detachAppender(appender);
}
}
}
@@ -110,7 +110,7 @@ class PrimaryRegistryTest {
/**
* The bug that made `primary.terminal` unretirable: with one slot, whichever lead called
* bridge_send first captured every nudge — including nudges for the other lead's delegations.
* fleet_send first captured every nudge — including nudges for the other lead's delegations.
*/
@Test
void aNudgeGoesToTheLeadThatDelegatedToThatWorker() {
@@ -57,7 +57,7 @@ class ClaudeCodeLauncherTest {
assertTrue(args.stream().anyMatch(a -> a.contains("\"bridge\"") && a.contains("http://127.0.0.1:8765/mcp")),
"inline bridge MCP config present");
assertTrue(args.contains("--append-system-prompt"));
assertTrue(args.stream().anyMatch(a -> a.contains("bridge_reply")), "reply charter present");
assertTrue(args.stream().anyMatch(a -> a.contains("fleet_reply")), "reply charter present");
}
@Test
@@ -358,7 +358,7 @@ class CompositePeerLauncherTest {
@Test
void explicitSpawnStillSucceedsOnWeightZeroProfile() {
// CB-554: weight: 0 excludes a profile from AUTOMATIC selection only — an explicit
// bridge_spawn{profile:"a"} must still work exactly as today (e.g. `opus` on the
// fleet_spawn{profile:"a"} must still work exactly as today (e.g. `opus` on the
// operator's own subscription, kept weight-0 so it is never picked automatically).
FakeHerdr herdr = new FakeHerdr();
Map<String, BridgedConfig.Profile> profiles = ordered(
@@ -593,7 +593,7 @@ class CompositePeerLauncherTest {
/**
* An explicit profile is the operator overriding and is NOT judged against the pool. It must
* stay that way: an unrolled `bridge_spawn{profile:"opus"}` carries no role, so it defaults to
* stay that way: an unrolled `fleet_spawn{profile:"opus"}` carries no role, so it defaults to
* DEV, and enforcing the pool here would refuse a spawn the operator asked for by name.
*/
@Test
@@ -169,8 +169,8 @@ class LeadHeartbeatLoopTest {
String t = pendingFleet().nudgeText();
assertTrue(t.startsWith("Heartbeat:"), "the nudge identifies itself as a heartbeat");
assertTrue(t.contains("2 worker replies pending collection"), t);
assertTrue(t.contains("bridge_poll(target=term_a)"), t);
assertTrue(t.contains("bridge_poll(target=term_b)"), t);
assertTrue(t.contains("fleet_poll(target=term_a)"), t);
assertTrue(t.contains("fleet_poll(target=term_b)"), t);
assertTrue(t.contains("1 DONE session awaiting teardown"), t);
assertTrue(t.contains("3 workers live"), t);
}
@@ -68,7 +68,7 @@ class MessageServiceTest {
injector.onStatus(T, AgentStatus.IDLE); // deliver the task (baselines the pre-turn content)
injector.onStatus(T, AgentStatus.WORKING); // worker picks it up and works
herdr.readText("BUILD GREEN: 391 files"); // the worker's turn produced new output
injector.onStatus(T, AgentStatus.IDLE); // working → idle: turn complete, no bridge_reply
injector.onStatus(T, AgentStatus.IDLE); // working → idle: turn complete, no fleet_reply
MessageService.Reply reply = send.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.COMPLETED_UNREPLIED, reply.outcome(),
@@ -153,7 +153,7 @@ class MessageServiceTest {
+ "(worker unreachable or stuck)", resolution.text());
}
// --- bridge_ask reverse rendezvous (CB-205) ------------------------------------------------
// --- fleet_ask reverse rendezvous (CB-205) ------------------------------------------------
@Test
void askSurfacesAsAQuestionAndTheAnswerResumesTheSameTurn() throws Exception {
@@ -172,7 +172,7 @@ class MessageServiceTest {
assertEquals("which config file?", q.text());
assertNotNull(q.turnId(), "a question carries a turnId to answer on");
// The primary answers via bridge_send(turnId); this blocks again for the worker's reply.
// The primary answers via fleet_send(turnId); this blocks again for the worker's reply.
CompletableFuture<MessageService.Reply> answer =
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000));
@@ -196,7 +196,7 @@ class MessageServiceTest {
injector.onStatus(T, AgentStatus.IDLE); // deliver
injector.onStatus(T, AgentStatus.WORKING); // worker picks it up, then pauses to ask
// A transport retry: two concurrent bridge_ask calls from the same worker session.
// A transport retry: two concurrent fleet_ask calls from the same worker session.
CompletableFuture<MessageService.AskResult> ask1 =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
CompletableFuture<MessageService.AskResult> ask2 =
@@ -297,7 +297,7 @@ class MessageServiceTest {
assertNotNull(q.turnId());
// The primary answers, unblocking the worker; but the worker never sends the follow-up
// bridge_reply, so the answering send rides out its short window as still-working.
// fleet_reply, so the answering send rides out its short window as still-working.
MessageService.Reply answer = messages.answer(q.turnId(), "config.yaml", 200);
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answer.outcome(),
"an answered worker that never replies times out as still working");
@@ -332,7 +332,7 @@ class MessageServiceTest {
assertNotNull(view, "a resolved async send must become DONE");
assertEquals(MessageService.Phase.DONE, view.phase());
assertEquals("async result", view.reply(), "the completed ticket reports the reply");
assertEquals("reply", view.replySource(), "a structured bridge_reply is sourced from 'reply'");
assertEquals("reply", view.replySource(), "a structured fleet_reply is sourced from 'reply'");
}
@Test
@@ -363,7 +363,7 @@ class MessageServiceTest {
/**
* The bug CB-548 fixes: L holds worker W, then architect A attempts W and times out BUSY. With
* delegator ownership recorded at {@code bridge_send} <em>request</em> time, A's rejected call
* delegator ownership recorded at {@code fleet_send} <em>request</em> time, A's rejected call
* would overwrite L — and W's late no-waiter reply would be pushed to A, who never owned the
* turn. The accepted-delivery hook must not fire for a BUSY send, so L stays the delegator.
*/
@@ -423,7 +423,7 @@ class MessageServiceTest {
}
/**
* CB-548 requirement: answering an existing {@code bridge_ask} is the SAME delegation, so it must
* CB-548 requirement: answering an existing {@code fleet_ask} is the SAME delegation, so it must
* not rewrite ownership. L accepted the send (owned), the worker paused to ask, and L answers via
* turnId — ownership stays L throughout; the answer path never touches the registry.
*/
@@ -533,12 +533,12 @@ class MessageServiceTest {
@Test
void aQuestionIsNeverQueuedInTheInbox() {
// No send is open — bridge_ask with no delegation returns NO_WAITER,
// No send is open — fleet_ask with no delegation returns NO_WAITER,
// and the question text MUST NOT appear in the reply inbox.
// The inbox is only fed by MessageService.reply(), not by bridge_ask.
// The inbox is only fed by MessageService.reply(), not by fleet_ask.
MessageService.AskResult r = messages.ask(T, "anyone there?", 500);
assertEquals(MessageService.AskOutcome.NO_WAITER, r.outcome(),
"bridge_ask with no open delegation must return NO_WAITER, never queued");
"fleet_ask with no open delegation must return NO_WAITER, never queued");
assertTrue(messages.drainReplies(T).isEmpty(), "questions must never be queued");
}
@@ -550,7 +550,7 @@ class MessageServiceTest {
awaitUninterruptibly(T);
injectDelivery();
// The worker never sends bridge_reply, but the turn completes.
// The worker never sends fleet_reply, but the turn completes.
herdr.readText("done-scraped");
completion.onTurnComplete(T); // The fallback arms and resolves the captured waiter.
@@ -714,7 +714,7 @@ class MessageServiceTest {
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
}
// --- CB-582: bridge_status pendingAsk() ------------------------------------------------------
// --- CB-582: fleet_status pendingAsk() ------------------------------------------------------
@Test
void pendingAskReturnsNullWhenNoQuestionIsOpen() throws Exception {
@@ -741,7 +741,7 @@ class MessageServiceTest {
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
MessageService.PendingAsk pending = messages.pendingAsk(T);
assertNotNull(pending, "bridge_status should see the open question");
assertNotNull(pending, "fleet_status should see the open question");
assertEquals(ticket, pending.ticket());
assertEquals("which config file?", pending.question());
assertEquals(asking.turnId(), pending.turnId());
@@ -760,7 +760,7 @@ class MessageServiceTest {
// MessageService.reply's rendezvous fast path is exactly what an async ticket always takes
// (sendAsync registers a rendezvous waiter — see asyncTasksByWaiter), so it never reached
// ReplyPushLoop.onReplyQueued. These prove the ticket reaches ReplyPushLoop through the new
// onTicketTerminal entry point instead, with no bridge_poll from the lead first.
// onTicketTerminal entry point instead, with no fleet_poll from the lead first.
private static final String LEAD = "term_lead";
@@ -808,7 +808,7 @@ class MessageServiceTest {
awaitNudge(wiring.leadHerdr());
String nudge = wiring.leadHerdr().lastCall("agent.prompt").params().toString();
assertTrue(nudge.contains(ticket), "the nudge should name the ticket: " + nudge);
assertTrue(nudge.contains("bridge_poll(ticket="),
assertTrue(nudge.contains("fleet_poll(ticket="),
"the nudge should name the exact ticket-collecting call: " + nudge);
assertFalse(nudge.toUpperCase().contains("FAILED"),
"a successfully-replied ticket's nudge must not say it failed: " + nudge);
@@ -874,7 +874,7 @@ class MessageServiceTest {
}
}
// --- CB-582: bridge_ask question-open nudges --------------------------------------------------
// --- CB-582: fleet_ask question-open nudges --------------------------------------------------
@Test
void anAsyncTicketThatPausesOnAQuestionNudgesTheLeadWithNoPriorPollCall() throws Exception {
@@ -891,7 +891,7 @@ class MessageServiceTest {
String nudge = wiring.leadHerdr().lastCall("agent.prompt").params().toString();
assertTrue(nudge.contains(ticket), "the nudge should name the ticket: " + nudge);
assertTrue(nudge.contains(asking.turnId()), "the nudge should name the turnId: " + nudge);
assertTrue(nudge.contains("bridge_send(turnId="),
assertTrue(nudge.contains("fleet_send(turnId="),
"the nudge should name the exact answer call: " + nudge);
assertTrue(nudge.contains("which config file?"), "the nudge should include the question: " + nudge);
@@ -997,7 +997,7 @@ class MessageServiceTest {
* {@code pruneTerminalTickets} drops entries from it once {@link MessageService#TICKET_TTL_NANOS}
* elapses. Before this test, that prune never told {@code ReplyPushLoop} — its own
* {@code pendingTickets} entry for a pruned, never-collected ticket had no remover at all, so it
* rode along on every later nudge to the same lead, naming a ticket {@code bridge_poll} could no
* rode along on every later nudge to the same lead, naming a ticket {@code fleet_poll} could no
* longer find. Uses the injectable clock (mirroring {@code SessionManager}'s {@code nowNanos} seam
* for its idle reaper) to cross the 10-minute TTL without a real wait.
*/
@@ -1040,10 +1040,10 @@ class MessageServiceTest {
String latestNudge = wiring.leadHerdr().lastCall("agent.prompt").params().toString();
assertTrue(latestNudge.contains(fresh), "the fresh ticket's nudge must still arrive: " + latestNudge);
assertFalse(latestNudge.contains(stale),
"a pruned ticket must never be named in a later nudge — it is gone and bridge_poll "
"a pruned ticket must never be named in a later nudge — it is gone and fleet_poll "
+ "on it would return nothing: " + latestNudge);
// And bridge_poll(ticket=stale) really does return nothing now — the nudge would have lied.
// And fleet_poll(ticket=stale) really does return nothing now — the nudge would have lied.
assertNull(wiring.service().poll(stale), "the pruned ticket must actually be gone, not just unmentioned");
}
}
@@ -13,7 +13,7 @@ import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* The reverse rendezvous (CB-205): the {@code bridge_ask} registry that lets a worker pause mid-turn
* The reverse rendezvous (CB-205): the {@code fleet_ask} registry that lets a worker pause mid-turn
* to ask the primary. Unit-level — the message-layer round-trip is covered in {@link MessageServiceTest}.
*/
class RendezvousTest {
@@ -168,8 +168,8 @@ class ReplyPushLoopTest {
// Exactly one nudge = exactly 1 agent.prompt call (it submits itself)
assertEquals(1, rec.sendCount());
assertTrue(rec.sentParams().stream()
.anyMatch(e -> e.getValue().toString().contains("bridge_poll")),
"nudge text should contain bridge_poll");
.anyMatch(e -> e.getValue().toString().contains("fleet_poll")),
"nudge text should contain fleet_poll");
}
@Test
@@ -226,14 +226,14 @@ class ReplyPushLoopTest {
void nudgeFormatIsCorrect() {
String nudge = ReplyPushLoop.NUDGE_FORMAT.formatted(WORKER, WORKER);
assertTrue(nudge.contains("Worker term_worker"));
assertTrue(nudge.contains("bridge_poll(target=term_worker)"));
assertTrue(nudge.contains("fleet_poll(target=term_worker)"));
}
@Test
void repliesNudgeFormatIsCorrect() {
String multi = ReplyPushLoop.REPLIES_NUDGE_FORMAT.formatted(2, "term_worker1, term_worker2");
assertTrue(multi.contains("2 workers"));
assertTrue(multi.contains("bridge_poll(target=...)"));
assertTrue(multi.contains("fleet_poll(target=...)"));
}
// --- CB-588: async ticket terminal nudges — decide() logic on tickets -----------------------
@@ -289,8 +289,8 @@ class ReplyPushLoopTest {
assertEquals(1, rec.sendCount());
String nudge = rec.sentParams().getFirst().getValue().toString();
assertTrue(nudge.contains("task-1"), "nudge should name the ticket");
assertTrue(nudge.contains("bridge_poll(ticket="), "nudge should name the exact ticket-poll call");
assertFalse(nudge.contains("bridge_poll(target="), "a ticket-only nudge must not tell the lead to run the target-poll call");
assertTrue(nudge.contains("fleet_poll(ticket="), "nudge should name the exact ticket-poll call");
assertFalse(nudge.contains("fleet_poll(target="), "a ticket-only nudge must not tell the lead to run the target-poll call");
}
@Test
@@ -468,11 +468,11 @@ class ReplyPushLoopTest {
void ticketNudgeFormatIsCorrect() {
String single = ReplyPushLoop.TICKET_NUDGE_FORMAT.formatted("task-1", "", "task-1");
assertTrue(single.contains("Ticket task-1"));
assertTrue(single.contains("bridge_poll(ticket=task-1)"));
assertTrue(single.contains("fleet_poll(ticket=task-1)"));
String multi = ReplyPushLoop.TICKETS_NUDGE_FORMAT.formatted(2, "", "task-1, task-2");
assertTrue(multi.contains("2 tickets"));
assertTrue(multi.contains("bridge_poll(ticket=...)"));
assertTrue(multi.contains("fleet_poll(ticket=...)"));
}
// --- CB-307 nudge path is unchanged (regression) --------------------------------------------
@@ -480,7 +480,7 @@ class ReplyPushLoopTest {
@Test
void inboxNudgeStillUsesTheOriginalTargetPollCall() {
String nudge = ReplyPushLoop.NUDGE_FORMAT.formatted(WORKER, WORKER);
assertTrue(nudge.contains("bridge_poll(target=" + WORKER + ")"),
assertTrue(nudge.contains("fleet_poll(target=" + WORKER + ")"),
"CB-588/CB-590 must not change the CB-307 inbox nudge's call shape");
}
@@ -502,7 +502,7 @@ class ReplyPushLoopTest {
"a reply and a ticket for the same lead must coalesce onto ONE schedule — "
+ "two nudge injections into the same lead pane must never overlap");
String nudge = rec.sentParams().getFirst().getValue().toString();
assertTrue(nudge.contains("bridge_poll(target=" + WORKER + ")"),
assertTrue(nudge.contains("fleet_poll(target=" + WORKER + ")"),
"the combined nudge must still mention the reply: " + nudge);
assertTrue(nudge.contains("task-1"), "the combined nudge must still mention the ticket: " + nudge);
}
@@ -525,7 +525,7 @@ class ReplyPushLoopTest {
Thread.sleep(200);
assertEquals(1, rec.sendCount(), "exactly one nudge once injectable — reply and ticket coalesced");
String nudge = rec.sentParams().getFirst().getValue().toString();
assertTrue(nudge.contains("bridge_poll(target=" + WORKER + ")"), "the reply must not be dropped: " + nudge);
assertTrue(nudge.contains("fleet_poll(target=" + WORKER + ")"), "the reply must not be dropped: " + nudge);
assertTrue(nudge.contains("task-1"), "the ticket must not be dropped: " + nudge);
}
@@ -657,7 +657,7 @@ class ReplyPushLoopTest {
assertTrue(loop.isActive(), "the schedule must stay active after nudging the fresh ticket");
}
// --- CB-582: bridge_ask question-open nudges -------------------------------------------------
// --- CB-582: fleet_ask question-open nudges -------------------------------------------------
@Test
void onQuestionOpenedWithNoKnownLeadNeverStartsASchedule() throws Exception {
@@ -706,7 +706,7 @@ class ReplyPushLoopTest {
String nudge = rec.sentParams().getFirst().getValue().toString();
assertTrue(nudge.contains("task-1"), "nudge should name the ticket: " + nudge);
assertTrue(nudge.contains("term_worker#1"), "nudge should name the turnId: " + nudge);
assertTrue(nudge.contains("bridge_send(turnId="), "nudge should name the exact answer call: " + nudge);
assertTrue(nudge.contains("fleet_send(turnId="), "nudge should name the exact answer call: " + nudge);
assertTrue(nudge.contains("which config file?"), "nudge should include the question text: " + nudge);
}
@@ -775,12 +775,12 @@ class ReplyPushLoopTest {
String single = ReplyPushLoop.QUESTION_NUDGE_FORMAT.formatted(
WORKER, "task-1", "term_worker#1", "which config?");
assertTrue(single.contains("Worker term_worker"));
assertTrue(single.contains("bridge_send(turnId=\"term_worker#1\""));
assertTrue(single.contains("fleet_send(turnId=\"term_worker#1\""));
assertTrue(single.contains("which config?"));
String multi = ReplyPushLoop.QUESTIONS_NUDGE_FORMAT.formatted(2, "task-1 (turnId=t1), task-2 (turnId=t2)");
assertTrue(multi.contains("2 workers"));
assertTrue(multi.contains("bridge_poll(ticket=...)"));
assertTrue(multi.contains("fleet_poll(ticket=...)"));
}
// --- metrics (CB-512) ----------------------------------------------------------------------
@@ -355,7 +355,7 @@ class BridgedAppTest {
@Test
void messageReturnsTheWorkersStructuredReply() throws Exception {
// CB-104 (option C): the blocking send resolves on the worker's bridge_reply, not a scrape.
// CB-104 (option C): the blocking send resolves on the worker's fleet_reply, not a scrape.
FakeHerdr herdr = new FakeHerdr().agentStatus("idle"); // poller delivers the injection
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
@@ -486,7 +486,7 @@ class BridgedAppTest {
/**
* CB-582: a lead polling {@code GET /sessions/{id}/status} on its normal cadence — not the
* ticket-scoped {@code /tasks/{ticket}} — must also see a worker's open async {@code bridge_ask}
* ticket-scoped {@code /tasks/{ticket}} — must also see a worker's open async {@code fleet_ask}
* question, since the reverse-rendezvous window it opened with is far shorter than that cadence.
*/
@Test
@@ -528,7 +528,7 @@ class BridgedAppTest {
assertEquals(turnId, body.get("turnId").asText());
assertEquals(ticket, body.get("ticket").asText());
// Answer it — via the same /message route bridge_send uses, keyed by turnId — so the
// Answer it — via the same /message route fleet_send uses, keyed by turnId — so the
// background ask thread does not linger past the test.
var answer = java.util.concurrent.CompletableFuture.supplyAsync(() -> {
try {
@@ -206,7 +206,7 @@ class SessionManagerTest {
@Test
void rosterViewExposesTheCharterReceiptButNeverTheCharterText() {
// The roster (bridge_list and GET /members both render through rosterView) must let a lead
// The roster (fleet_list and GET /members both render through rosterView) must let a lead
// see which charter a member got, without ever carrying the charter prose itself (CB-571).
MemberSession s = new MemberSession("p1", "term1", "prof", MemberRole.DEV, "/cwd", null,
0, 0, 0, MemberSession.State.READY, null, null,