CB-622 Unit A: register fleet_* MCP tools, keep bridge_* working
fleet_* is now the documented tool name for all eleven MCP tools; each bridge_* twin is registered against the exact same handler (no logic duplication) and its description leads with a DEPRECATED notice. A bridge_* call logs one WARN naming the old and new name, once per name for the life of the process (a Set, not a numeric sentinel). REPLY_CHARTER in HerdrPeerLauncher now tells a spawned member to call fleet_reply — the one rule that must survive with no repo checkout. All other bridge_* string literals across mcp/, Javadoc, and tests were renamed to fleet_* for consistency with the new documented name.
This commit is contained in:
@@ -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,
|
||||
|
||||
Reference in New Issue
Block a user