1040 lines
57 KiB
Java
1040 lines
57 KiB
Java
package dev.ltms.bridged.mcp;
|
|
|
|
import dev.ltms.bridged.auth.AuditLog;
|
|
import dev.ltms.bridged.auth.Authz;
|
|
import dev.ltms.bridged.auth.CallerResolver;
|
|
import dev.ltms.bridged.auth.Principal;
|
|
import dev.ltms.bridged.auth.Role;
|
|
import dev.ltms.bridged.guard.GuardException;
|
|
import dev.ltms.bridged.herdr.Agent;
|
|
import dev.ltms.bridged.metrics.BridgedMetrics;
|
|
import dev.ltms.bridged.metrics.Metrics;
|
|
import dev.ltms.bridged.inject.MemberPresence;
|
|
import dev.ltms.bridged.herdr.HerdrException;
|
|
import dev.ltms.bridged.msg.MessageService;
|
|
import dev.ltms.bridged.msg.Rendezvous;
|
|
import dev.ltms.bridged.peer.PeerUnreachableException;
|
|
import dev.ltms.bridged.session.SessionManager;
|
|
import dev.ltms.bridged.session.MemberSession;
|
|
import dev.ltms.bridged.session.WorktreeRequest;
|
|
import dev.ltms.bridged.peer.MemberRole;
|
|
import dev.ltms.bridged.peer.PeerLauncher;
|
|
import io.modelcontextprotocol.common.McpTransportContext;
|
|
import io.modelcontextprotocol.json.McpJsonMapper;
|
|
import io.modelcontextprotocol.json.jackson3.JacksonMcpJsonMapperSupplier;
|
|
import io.modelcontextprotocol.server.McpServer;
|
|
import io.modelcontextprotocol.server.McpSyncServer;
|
|
import io.modelcontextprotocol.server.McpSyncServerExchange;
|
|
import io.modelcontextprotocol.server.transport.HttpServletStreamableServerTransportProvider;
|
|
import io.modelcontextprotocol.spec.McpSchema;
|
|
import com.fasterxml.jackson.databind.ObjectMapper;
|
|
import jakarta.servlet.http.HttpServlet;
|
|
|
|
import java.util.LinkedHashMap;
|
|
import java.util.List;
|
|
import java.util.Map;
|
|
import java.util.Set;
|
|
import java.util.function.Function;
|
|
import java.util.function.LongSupplier;
|
|
import java.util.function.Supplier;
|
|
import java.util.stream.Collectors;
|
|
|
|
/**
|
|
* 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}.
|
|
*
|
|
* <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
|
|
* 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
|
|
* {@link McpSchema.CallToolResult}, so it is unit-testable without standing up the HTTP transport;
|
|
* the SDK owns the wire protocol. Mount {@link #servlet()} at {@code /mcp} on the daemon's Jetty.
|
|
*/
|
|
public final class BridgeMcp {
|
|
|
|
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
|
|
// 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
|
|
|
|
/** 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). */
|
|
static final String CALLER_PID = "callerPid";
|
|
/** Transport-context key under which the extractor stashes the resolved {@link Role} (CB-501). */
|
|
static final String CALLER_ROLE = "callerRole";
|
|
/** Transport-context key for the lead's configured name, when the caller is one (CB-530). */
|
|
static final String CALLER_NAME = "callerName";
|
|
|
|
private final HttpServletStreamableServerTransportProvider transport;
|
|
private final McpSyncServer server;
|
|
private final CallerResolver authz; // CB-501: null → authorization not enforced (legacy)
|
|
private final Metrics metrics; // CB-502: null → auth failures not counted
|
|
private final CapacitySource capacity;
|
|
private final HealthCoverageSource healthCoverage;
|
|
|
|
/** Capacity facts used by {@code bridge_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. */
|
|
public static CapacitySource none() { return new CapacitySource(_ -> 0, _ -> null, Set::of, System::nanoTime); }
|
|
boolean available() { return !configuredProfiles.get().isEmpty(); }
|
|
}
|
|
|
|
/** Coverage is supplied by the health wiring, not inferred from a missing dependency. */
|
|
public record HealthCoverageSource(Supplier<String> value) { }
|
|
|
|
/**
|
|
* @param callers resolves each call's {@link Principal}; {@code null} disables authorization.
|
|
* This surface needs its own enforcement: {@code /mcp} is a raw servlet on
|
|
* 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}
|
|
*/
|
|
public BridgeMcp(MessageService messages, PeerLauncher workers, SessionManager sessions,
|
|
ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry,
|
|
CallerResolver callers, Metrics metrics, CapacitySource capacity, HealthCoverageSource healthCoverage) {
|
|
this.capacity = capacity;
|
|
this.healthCoverage = healthCoverage;
|
|
McpJsonMapper json = new JacksonMcpJsonMapperSupplier().get();
|
|
this.transport = HttpServletStreamableServerTransportProvider.builder()
|
|
.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
|
|
// 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 -> {
|
|
// One resolution per call, shared with the REST surface via CallerResolver so
|
|
// the two paths cannot drift on who a caller is.
|
|
Principal p = callers != null
|
|
? callers.resolve(req.getRemoteAddr(), req.getRemotePort(),
|
|
req.getHeader("Authorization"))
|
|
: legacyPrincipal(identity, req.getRemoteAddr(), req.getRemotePort());
|
|
// CB-532: guard on the ROLE, not on the terminal being null. This excludes a
|
|
// lead, which carries its pane too, while including every spawned member role.
|
|
// Enrolling a lead would count it as an available member in the roster.
|
|
markSpawnedMemberPresent(p, presence);
|
|
return McpTransportContext.create(Map.of(
|
|
CALLER_TERMINAL, orEmpty(p.terminal()),
|
|
CALLER_PID, Long.toString(p.pid()),
|
|
CALLER_ROLE, p.role().name(),
|
|
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) -> {
|
|
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.SEND,
|
|
str(req.arguments(), "sessionId"));
|
|
if (denied != null) return denied;
|
|
String caller = callerTerminal(exchange);
|
|
// CB-548: only a PRIMARY caller may claim the legacy singleton "primary" fallback.
|
|
// An architect delegates as its own pane but must never become the fallback that
|
|
// no-delegation inbox nudges target as if it were the primary (the per-target
|
|
// delegation map does not cure the singleton).
|
|
recordPrimarySingleton(primaryRegistry, caller, principal(exchange));
|
|
Map<String, Object> a = req.arguments();
|
|
String target = str(a, "sessionId");
|
|
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
|
|
// 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));
|
|
}
|
|
// CB-548: delegator ownership (which lead's reply nudge this worker routes to,
|
|
// CB-532) is recorded only once the send is ACCEPTED — MessageService has won the
|
|
// session lock and queued delivery — via the accepted-delivery callback, never at
|
|
// request time. A concurrent sender that times out BUSY therefore cannot steal a
|
|
// live turn's reply routing without ever owning the turn.
|
|
Runnable onAccepted = () -> primaryRegistry.recordDelegation(target, caller);
|
|
// wait defaults to true (block for the reply); wait:false is fire-and-poll.
|
|
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) -> {
|
|
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) -> {
|
|
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) -> {
|
|
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) -> {
|
|
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) -> {
|
|
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) -> {
|
|
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.SPAWN, null);
|
|
if (denied != null) return denied;
|
|
String caller = callerTerminal(exchange);
|
|
// SPAWN is already auth-gated to PRIMARY (architects can never call it), but
|
|
// enforce the same invariant here: only a PRIMARY may claim the legacy singleton.
|
|
recordPrimarySingleton(primaryRegistry, caller, principal(exchange));
|
|
Map<String, Object> a = req.arguments();
|
|
// CB-112: worker inherits the primary's cwd unless the call pins one.
|
|
// CB-301: carry the caller's identity as the session owner (null for the primary).
|
|
// CB-301-ext: optional isolated worktree for parallel implementers.
|
|
String callerCwd = identity.cwdForPid(callerPid(exchange));
|
|
return spawn(sessions, str(a, "profile"), str(a, "role"), str(a, "cwd"), callerCwd,
|
|
callerTerminal(exchange), worktreeRequest(a));
|
|
})
|
|
.toolCall(listTool(), (exchange, _) -> {
|
|
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
|
|
if (denied != null) return denied;
|
|
return listFleet(workers, sessions, messages, capacity, healthCoverage,
|
|
callers == null ? Map.of() : callers.leads(),
|
|
callerTerminal(exchange));
|
|
})
|
|
.toolCall(stopTool(), (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, _) -> {
|
|
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
|
|
if (denied != null) return denied;
|
|
return profiles(workers);
|
|
})
|
|
.toolCall(whoamiTool(), (exchange, _) -> {
|
|
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
|
|
if (denied != null) return denied;
|
|
return whoami(principal(exchange), sessions);
|
|
})
|
|
.build();
|
|
this.authz = callers;
|
|
this.metrics = metrics;
|
|
}
|
|
|
|
/**
|
|
* Pre-CB-501 identity: worker if the connection maps to a pane, otherwise the primary. Used
|
|
* only by the legacy constructor, where authorization is not enforced anyway.
|
|
*/
|
|
private static Principal legacyPrincipal(ConnectionIdentity identity, String addr, int port) {
|
|
ConnectionIdentity.Caller c = identity.resolve(addr, port);
|
|
return c.terminal() != null
|
|
? Principal.worker(c.terminal(), c.pid())
|
|
: Principal.primary(c.pid());
|
|
}
|
|
|
|
/** The caller reconstructed from the transport context. */
|
|
private static Principal principal(McpSyncServerExchange exchange) {
|
|
return principalFrom(exchange.transportContext().get(CALLER_ROLE),
|
|
callerTerminal(exchange), callerPid(exchange), callerName(exchange));
|
|
}
|
|
|
|
/**
|
|
* Rebuild a {@link Principal} from the three values the context extractor stashed.
|
|
*
|
|
* <p>Split out from {@link #principal(McpSyncServerExchange)} so the identity rules are
|
|
* reachable without an {@code McpSyncServerExchange} — that is an SDK type this project has no
|
|
* mocking library to fabricate, which is why this logic had no test at all until CB-513.
|
|
*
|
|
* @param role the stashed {@link Role} name, or {@code null} on the legacy path
|
|
* @param terminal the worker terminal, or {@code null} for a non-worker
|
|
* @param pid the calling pid, or {@code -1}
|
|
*/
|
|
static Principal principalFrom(Object role, String terminal, long pid) {
|
|
return principalFrom(role, terminal, pid, null);
|
|
}
|
|
|
|
/** As {@link #principalFrom(Object, String, long)}, carrying a lead's name (CB-530). */
|
|
static Principal principalFrom(Object role, String terminal, long pid, String name) {
|
|
if (role == null) {
|
|
// No role stashed (legacy path): fall back to the historical interpretation.
|
|
return terminal != null ? Principal.worker(terminal, pid) : Principal.primary(pid);
|
|
}
|
|
return new Principal(Role.valueOf(role.toString()), terminal, pid, name);
|
|
}
|
|
|
|
/**
|
|
* Gate a tool call on the CB-505 table. Returns {@code null} when the call may proceed, or the
|
|
* error result to return when it may not.
|
|
*/
|
|
private McpSchema.CallToolResult deny(McpSyncServerExchange exchange, Authz.Action action,
|
|
String target) {
|
|
return denyFor(principal(exchange), action, target);
|
|
}
|
|
|
|
/**
|
|
* The policy half of {@link #deny}: everything except pulling the caller out of the MCP
|
|
* exchange. Kept separate so the authorization decision — the actual control — is unit-testable
|
|
* without fabricating an SDK {@code McpSyncServerExchange}.
|
|
*
|
|
* <p>This surface exists because the enforcement was previously unreachable from a test: no
|
|
* test constructs a {@code BridgeMcp}, so the whole MCP-side gate ran zero times in the suite
|
|
* while the REST-side equivalent had ten tests. A security control nothing exercises is a
|
|
* claim, not a control.
|
|
*
|
|
* @return {@code null} when the call may proceed, or the error result to return when it may not
|
|
*/
|
|
McpSchema.CallToolResult denyFor(Principal caller, Authz.Action action, String target) {
|
|
// The enforcement switch lives HERE rather than in the exchange-facing wrapper: any future
|
|
// tool that calls this directly must not be able to skip the gate by accident.
|
|
if (authz == null) {
|
|
return null; // legacy constructor: authorization not enforced
|
|
}
|
|
if (Authz.permits(caller, action, target)) {
|
|
if (action != Authz.Action.READ) {
|
|
AuditLog.allowed(caller, action, target); // reads would drown the trail
|
|
}
|
|
return null;
|
|
}
|
|
String reason = Authz.isUnauthenticated(caller) ? "unauthenticated" : "forbidden";
|
|
AuditLog.denied(caller, action, target, reason);
|
|
if (metrics != null) {
|
|
metrics.inc(BridgedMetrics.AUTH_FAILURES, "reason", reason);
|
|
}
|
|
return error(reason + ": " + caller.describe() + " may not " + action);
|
|
}
|
|
|
|
/**
|
|
* Update the legacy singleton "primary" fallback used for no-delegation inbox nudges (CB-548).
|
|
*
|
|
* <p>Only {@link Role#PRIMARY} callers — the unnamed primary and named leads alike — may claim
|
|
* it. An architect delegates as its own pane but must never become the fallback: the per-target
|
|
* delegation map ({@code PrimaryRegistry#recordDelegation}) does not cure the singleton, so an
|
|
* architect left here would draw nudges that belong to a primary. The decision uses the resolved
|
|
* role, never name/kind sniffing. A null {@code caller} (legacy/no-auth path) records nothing.
|
|
*
|
|
* <p>Split out of the tool handlers so the guard is unit-testable without fabricating an SDK
|
|
* {@code McpSyncServerExchange} (same pattern as {@link #denyFor}/{@link #principalFrom}).
|
|
*/
|
|
static void recordPrimarySingleton(PrimaryRegistry registry, String callerTerminal, Principal caller) {
|
|
if (caller != null && caller.isPrimary()) {
|
|
registry.record(callerTerminal);
|
|
}
|
|
}
|
|
|
|
/** The worker identity resolved from this call's connection, or {@code null} if the primary. */
|
|
private static String callerTerminal(McpSyncServerExchange exchange) {
|
|
Object v = exchange.transportContext().get(CALLER_TERMINAL);
|
|
String s = v == null ? null : v.toString();
|
|
return (s == null || s.isBlank()) ? null : s;
|
|
}
|
|
|
|
/** The lead name resolved from this call's connection, or {@code null} (CB-530). */
|
|
private static String callerName(McpSyncServerExchange exchange) {
|
|
Object v = exchange.transportContext().get(CALLER_NAME);
|
|
String s = v == null ? null : v.toString();
|
|
return (s == null || s.isBlank()) ? null : s;
|
|
}
|
|
|
|
/** The caller's PID resolved from this call's connection, or {@code -1} if unknown. */
|
|
private static long callerPid(McpSyncServerExchange exchange) {
|
|
Object v = exchange.transportContext().get(CALLER_PID);
|
|
try {
|
|
return v == null ? -1 : Long.parseLong(v.toString());
|
|
} catch (NumberFormatException e) {
|
|
return -1;
|
|
}
|
|
}
|
|
|
|
private static String orEmpty(String s) {
|
|
return s == null ? "" : s;
|
|
}
|
|
|
|
/** The Streamable-HTTP servlet to mount at {@code /mcp} on the daemon's Jetty. */
|
|
public HttpServlet servlet() {
|
|
return transport;
|
|
}
|
|
|
|
/** Mark a connected spawned member available for the injector readiness gate. */
|
|
static void markSpawnedMemberPresent(Principal caller, MemberPresence presence) {
|
|
if (caller.isSpawnedMember()) {
|
|
presence.markPresent(caller.terminal());
|
|
}
|
|
}
|
|
|
|
/** Graceful shutdown of the MCP server. */
|
|
public void close() {
|
|
server.closeGracefully();
|
|
}
|
|
|
|
// --- tool logic (thin adapters over the services; unit-testable) ---------------------------
|
|
|
|
/**
|
|
* {@code bridge_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
|
|
* a BUSY interloper never claims a turn it did not win. {@code null} disables recording.
|
|
*/
|
|
static McpSchema.CallToolResult send(MessageService messages, String sessionId, String content,
|
|
Long timeoutMs, Runnable onAccepted, Set<String> profiles) {
|
|
if (isBlank(sessionId) || isBlank(content)) {
|
|
return error("sessionId and content are required");
|
|
}
|
|
McpSchema.CallToolResult targetError = profileTargetError(sessionId, profiles);
|
|
if (targetError != null) {
|
|
return targetError;
|
|
}
|
|
long timeout = clamp(timeoutMs == null ? DEFAULT_TIMEOUT_MS : timeoutMs);
|
|
try {
|
|
return formatReply(messages.send(sessionId, content, timeout, onAccepted), timeout);
|
|
} catch (HerdrException e) {
|
|
return error("herdr error contacting session " + sessionId + ": " + e.getMessage());
|
|
}
|
|
}
|
|
|
|
/**
|
|
* {@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
|
|
* 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) {
|
|
if (isBlank(turnId) || isBlank(content)) {
|
|
return error("turnId and content are required to answer a worker's question");
|
|
}
|
|
long timeout = clamp(timeoutMs == null ? DEFAULT_TIMEOUT_MS : timeoutMs);
|
|
return formatReply(messages.answer(turnId, content, timeout), timeout);
|
|
}
|
|
|
|
/**
|
|
* {@code bridge_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 "
|
|
+ "from the connection");
|
|
}
|
|
if (isBlank(question)) {
|
|
return error("question is required");
|
|
}
|
|
long timeout = Math.clamp(timeoutMs == null ? ASK_DEFAULT_TIMEOUT_MS : timeoutMs, 1, ASK_MAX_TIMEOUT_MS);
|
|
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 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]");
|
|
};
|
|
}
|
|
|
|
/** Render a {@link MessageService.Reply} as a tool result — shared by {@link #send} and {@link #answer}. */
|
|
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
|
|
// 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());
|
|
// 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 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()
|
|
+ "\" 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)");
|
|
case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, BUSY -> text("[no reply within " + timeout + "ms — worker "
|
|
+ r.outcome().name().toLowerCase().replace("timed_out_", "") + "; retry or poll status]");
|
|
};
|
|
}
|
|
|
|
/**
|
|
* {@code bridge_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.
|
|
*
|
|
* This wires the accepted-delivery hook
|
|
* (CB-548) so an async flooding send records delegator ownership exactly once it is accepted.
|
|
*/
|
|
static McpSchema.CallToolResult sendAsync(MessageService messages, String sessionId, String content,
|
|
Runnable onAccepted, Set<String> profiles) {
|
|
if (isBlank(sessionId) || isBlank(content)) {
|
|
return error("sessionId and content are required");
|
|
}
|
|
McpSchema.CallToolResult targetError = profileTargetError(sessionId, profiles);
|
|
if (targetError != null) {
|
|
return targetError;
|
|
}
|
|
String ticket = messages.sendAsync(sessionId, content, onAccepted);
|
|
return text("accepted — task delegated. Poll bridge_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.");
|
|
}
|
|
return null;
|
|
}
|
|
|
|
/** {@code bridge_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);
|
|
if (replies.isEmpty()) {
|
|
return text("[]");
|
|
}
|
|
return text(json(replies));
|
|
}
|
|
if (isBlank(ticket)) {
|
|
return error("ticket (or target) is required");
|
|
}
|
|
MessageService.TaskView v = messages.poll(ticket);
|
|
if (v == null) {
|
|
return error("unknown ticket: " + ticket + " (never issued, or expired)");
|
|
}
|
|
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()
|
|
: 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()
|
|
+ "\" 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
|
|
* 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 "
|
|
+ "from the connection");
|
|
}
|
|
if (content == null) {
|
|
return error("content is required");
|
|
}
|
|
messages.reply(callerTerminal, content);
|
|
return text("delivered");
|
|
}
|
|
|
|
/** {@code bridge_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");
|
|
}
|
|
messages.ackReply(target, msgId);
|
|
return text("acknowledged " + msgId);
|
|
}
|
|
|
|
/** {@code bridge_status}: the live lifecycle status of a worker session. */
|
|
static McpSchema.CallToolResult status(MessageService messages, String sessionId) {
|
|
if (isBlank(sessionId)) {
|
|
return error("sessionId is required");
|
|
}
|
|
try {
|
|
return text(messages.status(sessionId).name().toLowerCase());
|
|
} catch (HerdrException e) {
|
|
return error("herdr error for session " + sessionId + ": " + e.getMessage());
|
|
}
|
|
}
|
|
|
|
/**
|
|
* {@code bridge_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
|
|
* from side channels the daemon does not control: a charter string in its system prompt, the
|
|
* 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
|
|
* 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
|
|
* registry has no record of — one that outlived a daemon restart — still gets its role and
|
|
* {@code sessionId}, which is the load-bearing part.
|
|
*/
|
|
static McpSchema.CallToolResult whoami(Principal caller, SessionManager sessions) {
|
|
Map<String, Object> m = new LinkedHashMap<>();
|
|
m.put("role", caller.role().name().toLowerCase());
|
|
if (caller.isArchitect()) {
|
|
// CB-548: the role reads "architect"; the name is the gateway-local slot the pane is
|
|
// bound to, and the pane itself so a peer knows where to reach it.
|
|
if (caller.name() != null) {
|
|
m.put("architect", caller.name());
|
|
}
|
|
if (caller.terminal() != null) {
|
|
m.put("sessionId", caller.terminal());
|
|
}
|
|
return text(json(m));
|
|
}
|
|
if (!caller.isWorker()) {
|
|
// CB-530: which lead, once more than one pane is configured as one. `role` deliberately
|
|
// still reads "primary" — the fallback ladder in CLAUDE.md keys on it, and a lead IS a
|
|
// primary for authorization; the name is additive so no existing reader breaks.
|
|
if (caller.name() != null) {
|
|
m.put("leader", caller.name());
|
|
}
|
|
// CB-532: a lead's own pane, so it can tell a peer where to reach it — and so an
|
|
// operator can read off which tab hosts which lead without going to herdr.
|
|
if (caller.terminal() != null) {
|
|
m.put("sessionId", caller.terminal());
|
|
}
|
|
return text(json(m));
|
|
}
|
|
m.put("sessionId", caller.terminal());
|
|
sessions.roster().stream()
|
|
.filter(s -> caller.terminal().equals(s.terminalId()))
|
|
.findFirst()
|
|
.ifPresent(s -> {
|
|
m.put("paneId", s.paneId());
|
|
m.put("profile", s.profile());
|
|
m.put("state", s.state().name().toLowerCase());
|
|
if (s.worktree() != null) {
|
|
m.put("worktree", s.worktree());
|
|
}
|
|
if (s.branch() != null) {
|
|
m.put("branch", s.branch());
|
|
}
|
|
if (s.ownerTerminal() != null) {
|
|
m.put("owner", s.ownerTerminal());
|
|
}
|
|
});
|
|
return text(json(m));
|
|
}
|
|
|
|
// --- fleet management logic (CB-108 / CB-301) --------------------------------------------
|
|
|
|
/** {@code bridge_spawn} without cwd/caller context (default resolution). */
|
|
static McpSchema.CallToolResult spawn(SessionManager sessions, String profile) {
|
|
return spawn(sessions, profile, null, null, null, null, null);
|
|
}
|
|
|
|
/**
|
|
* {@code bridge_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.
|
|
* CB-301: the session is registered with {@code ownerTerminal} as its owner.
|
|
* CB-301-ext: {@code worktreeRequest} non-null provisions an isolated git worktree.
|
|
*
|
|
* <p>{@code role} and {@code profile} are independent: the role picks the contract, the profile
|
|
* picks the backend. A reviewer on the same profile as the dev it reviews is a normal spawn.
|
|
*/
|
|
static McpSchema.CallToolResult spawn(SessionManager sessions, String profile, String role,
|
|
String requestedCwd, String callerCwd,
|
|
String ownerTerminal, WorktreeRequest worktreeRequest) {
|
|
MemberRole memberRole;
|
|
try {
|
|
memberRole = isBlank(role) ? MemberRole.DEV : MemberRole.parse(role);
|
|
} catch (IllegalArgumentException e) {
|
|
return error(e.getMessage());
|
|
}
|
|
try {
|
|
MemberSession member = sessions.acquire(isBlank(profile) ? null : profile, memberRole,
|
|
requestedCwd, callerCwd, ownerTerminal, worktreeRequest);
|
|
return text(json(memberView(member)));
|
|
} catch (GuardException e) {
|
|
return error("subscription boundary: " + e.getMessage());
|
|
} catch (IllegalArgumentException e) {
|
|
return error(e.getMessage()); // unknown / no-default profile
|
|
} catch (PeerUnreachableException e) {
|
|
return error("spawn timed out — worker pane never reached injectable state: " + e.getMessage());
|
|
} catch (HerdrException e) {
|
|
return error("herdr error spawning worker: " + e.getMessage());
|
|
}
|
|
}
|
|
|
|
/** Build a {@link WorktreeRequest} from {@code bridge_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)) {
|
|
return null;
|
|
}
|
|
String ticket = str(a, "ticket");
|
|
if (w instanceof String s) {
|
|
if (s.isBlank() || "false".equalsIgnoreCase(s)) {
|
|
return null;
|
|
}
|
|
if ("true".equalsIgnoreCase(s)) {
|
|
if (isBlank(ticket)) {
|
|
throw new IllegalArgumentException("worktree=true requires a ticket slug");
|
|
}
|
|
return new WorktreeRequest(ticket, null);
|
|
}
|
|
return new WorktreeRequest(s, null);
|
|
}
|
|
if (w instanceof Boolean b && b) {
|
|
if (isBlank(ticket)) {
|
|
throw new IllegalArgumentException("worktree=true requires a ticket slug");
|
|
}
|
|
return new WorktreeRequest(ticket, null);
|
|
}
|
|
return null;
|
|
}
|
|
|
|
/** {@code bridge_profiles}: the configured worker profiles and the default. */
|
|
static McpSchema.CallToolResult profiles(PeerLauncher workers) {
|
|
return text(json(Map.of(
|
|
"profiles", workers.profiles(),
|
|
"default", workers.defaultProfile() == null ? "" : workers.defaultProfile())));
|
|
}
|
|
|
|
/**
|
|
* {@code bridge_list}: the whole fleet — {@code leads} and {@code workers} — each merged with
|
|
* live herdr status. CB-519 decoupled the registry key (a host-unique id) from the herdr pane
|
|
* coordinate, so the join is on the terminal id, which both the session and the live agent carry.
|
|
*
|
|
* <p>CB-535 added the {@code leads} half. Until then this listed the worker roster alone, and a
|
|
* lead asking "who else is here?" got an empty array — which reads as <em>no peers</em> but
|
|
* actually means <em>no workers spawned</em>. There was no way at all for a lead to learn a
|
|
* peer's address; it had to be carried across by a human. Both halves are reported even when a
|
|
* half is empty, so an empty {@code workers} can no longer be mistaken for an empty fleet.
|
|
*
|
|
* <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.
|
|
*
|
|
* @param leads terminal_id → lead name, live from the resolver
|
|
* @param selfTerm the calling pane's terminal id, or blank for a caller with no pane
|
|
*/
|
|
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions,
|
|
Map<String, String> leads, String selfTerm) {
|
|
return listFleet(workers, sessions, null, CapacitySource.none(), new HealthCoverageSource(() -> "off"), leads, selfTerm);
|
|
}
|
|
|
|
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
|
|
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
|
Map<String, String> leads, String selfTerm) {
|
|
try {
|
|
Map<String, Agent> live = workers.list().stream()
|
|
.map(Agent.class::cast)
|
|
.filter(a -> a.terminalId() != null)
|
|
.collect(Collectors.toMap(Agent::terminalId, Function.identity(), (_, b) -> b));
|
|
List<Map<String, Object>> leadRows = leads.entrySet().stream()
|
|
.sorted(Map.Entry.comparingByValue())
|
|
.map(e -> leadView(e.getKey(), e.getValue(), live.get(e.getKey()), selfTerm))
|
|
.toList();
|
|
List<MemberSession> roster = sessions.roster();
|
|
List<Map<String, Object>> out = roster.stream()
|
|
.map(s -> memberCapacityView(s, live.get(s.terminalId()), messages, capacity.clock().getAsLong()))
|
|
.toList();
|
|
Set<String> profiles = new java.util.TreeSet<>(capacity.configuredProfiles().get());
|
|
roster.stream().map(MemberSession::profile).forEach(profiles::add);
|
|
Map<String, Object> result = new LinkedHashMap<>();
|
|
result.put("leads", leadRows); result.put("members", out);
|
|
result.put("healthCoverage", healthCoverage.value().get());
|
|
if (capacity.available()) result.put("capacity", profiles.stream()
|
|
.map(profile -> capacityView(profile, capacity.liveCount(), capacity.maxLoad(), roster, messages,
|
|
capacity.clock().getAsLong())).toList());
|
|
return text(json(result));
|
|
} catch (HerdrException e) {
|
|
return error("herdr error listing the fleet: " + e.getMessage());
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Capacity is advisory only. {@code reclaimable} says there is no bridge work, not that bridged
|
|
* may stop the member: the bridge has capacity facts but no work list, and choosing work needs
|
|
* authority it does not have. {@code idleForSeconds} is derived from monotonic nanoTime and has
|
|
* no wall-clock meaning across a daemon restart.
|
|
*/
|
|
private static Map<String, Object> memberCapacityView(MemberSession session, Agent live,
|
|
MessageService messages, long nowNanos) {
|
|
Map<String, Object> row = SessionManager.rosterView(session, live);
|
|
boolean open = messages != null && messages.hasAcceptedDelivery(session.terminalId());
|
|
boolean inbox = messages != null && messages.hasInboxMessage(session.terminalId());
|
|
boolean reclaimable = (session.state() == MemberSession.State.READY || session.state() == MemberSession.State.DONE)
|
|
&& !open && !inbox;
|
|
row.put("reclaimable", reclaimable);
|
|
row.put("idleForSeconds", reclaimable ? Math.max(0, (nowNanos - session.lastActivityAtNanos()) / 1_000_000_000L) : null);
|
|
return row;
|
|
}
|
|
|
|
private static Map<String, Object> capacityView(String profile, Function<String, Integer> liveCount,
|
|
Function<String, Integer> maxLoad, List<MemberSession> roster,
|
|
MessageService messages, long nowNanos) {
|
|
Integer cap = maxLoad.apply(profile);
|
|
int live = liveCount.apply(profile);
|
|
int reclaimable = (int) roster.stream().filter(s -> profile.equals(s.profile()))
|
|
.filter(s -> (s.state() == MemberSession.State.READY || s.state() == MemberSession.State.DONE))
|
|
.filter(s -> messages == null || (!messages.hasAcceptedDelivery(s.terminalId()) && !messages.hasInboxMessage(s.terminalId())))
|
|
.count();
|
|
Map<String, Object> row = new LinkedHashMap<>();
|
|
row.put("profile", profile); row.put("maxLoad", cap); row.put("live", live);
|
|
row.put("free", cap == null ? null : Math.max(0, cap - live)); row.put("reclaimable", reclaimable);
|
|
return row;
|
|
}
|
|
|
|
/**
|
|
* One lead's row: its address, its name, and whether it can be reached right now.
|
|
*
|
|
* <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,
|
|
* 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,
|
|
String selfTerm) {
|
|
Map<String, Object> m = new LinkedHashMap<>();
|
|
m.put("sessionId", terminal);
|
|
m.put("name", name);
|
|
m.put("status", live == null || live.status() == null
|
|
? "unknown" : live.status().name().toLowerCase());
|
|
if (terminal.equals(selfTerm)) {
|
|
m.put("self", true);
|
|
}
|
|
return m;
|
|
}
|
|
|
|
/** {@code bridge_stop}: tear a worker down by its pane id. */
|
|
static McpSchema.CallToolResult stop(SessionManager sessions, String paneId) {
|
|
if (isBlank(paneId)) {
|
|
return error("paneId is required");
|
|
}
|
|
try {
|
|
sessions.release(paneId);
|
|
return text("stopped " + paneId);
|
|
} catch (HerdrException e) {
|
|
return error("herdr error stopping " + paneId + ": " + e.getMessage());
|
|
}
|
|
}
|
|
|
|
/** CB-301 projection from the authoritative session registry. */
|
|
private static Map<String, Object> memberView(MemberSession s) {
|
|
Map<String, Object> m = new LinkedHashMap<>();
|
|
m.put("sessionId", s.terminalId());
|
|
m.put("paneId", s.paneId());
|
|
// Echo the role back so the caller can see what it actually got, not what it meant to ask
|
|
// for — a spawn that silently fell back to dev is otherwise invisible.
|
|
m.put("role", s.role() == null ? "dev" : s.role().wireName());
|
|
m.put("profile", s.profile());
|
|
m.put("status", s.state().name().toLowerCase());
|
|
if (s.worktree() != null) {
|
|
m.put("worktree", s.worktree());
|
|
}
|
|
if (s.branch() != null) {
|
|
m.put("branch", s.branch());
|
|
}
|
|
return m;
|
|
}
|
|
|
|
private static String json(Object o) {
|
|
try {
|
|
return MAPPER.writeValueAsString(o);
|
|
} catch (Exception e) {
|
|
return String.valueOf(o);
|
|
}
|
|
}
|
|
|
|
// --- tool schemas --------------------------------------------------------------------------
|
|
|
|
private static McpSchema.Tool sendTool() {
|
|
return tool("bridge_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.",
|
|
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 — "
|
|
+ "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",
|
|
"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 "
|
|
+ "primary; identity is resolved from your connection.",
|
|
objectSchema(Map.of(
|
|
"question", stringProp("The question to put to the primary"),
|
|
"timeoutMs", Map.of("type", "integer",
|
|
"description", "Max ms to wait for the primary's answer")),
|
|
List.of("question")));
|
|
}
|
|
|
|
private static McpSchema.Tool pollTool() {
|
|
return tool("bridge_poll",
|
|
"Check an async delegation (a bridge_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"),
|
|
"target", stringProp("Worker session id to drain pending replies from (optional)")),
|
|
List.of()));
|
|
}
|
|
|
|
private static McpSchema.Tool ackTool() {
|
|
return tool("bridge_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.",
|
|
objectSchema(Map.of(
|
|
"target", stringProp("Worker session id whose inbox to ack from"),
|
|
"msgId", stringProp("The message id to acknowledge")),
|
|
List.of("target", "msgId")));
|
|
}
|
|
|
|
private static McpSchema.Tool spawnTool() {
|
|
return tool("bridge_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 "
|
|
+ "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. Returns the member's "
|
|
+ "sessionId (use with bridge_send) and paneId (use with bridge_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)"),
|
|
"cwd", stringProp("Working directory for the member (omit to inherit yours)"),
|
|
"worktree", Map.of("type", "string", "description", "'true' or a ticket slug — requests an isolated git worktree"),
|
|
"ticket", stringProp("Ticket slug when worktree:true")),
|
|
List.of()));
|
|
}
|
|
|
|
private static McpSchema.Tool profilesTool() {
|
|
return tool("bridge_profiles",
|
|
"List the configured worker profiles (backends) and which one bridge_spawn uses by default.",
|
|
objectSchema(Map.of(), List.of()));
|
|
}
|
|
|
|
private static McpSchema.Tool listTool() {
|
|
return tool("bridge_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), "
|
|
+ "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, and live herdr status. An empty 'members' "
|
|
+ "means no members are spawned; it says nothing about peers.",
|
|
objectSchema(Map.of(), List.of()));
|
|
}
|
|
|
|
private static McpSchema.Tool stopTool() {
|
|
return tool("bridge_stop",
|
|
"Tear down a worker session by its paneId (from bridge_spawn or bridge_list).",
|
|
objectSchema(Map.of(
|
|
"paneId", stringProp("The worker's paneId to stop")),
|
|
List.of("paneId")));
|
|
}
|
|
|
|
private static McpSchema.Tool replyTool() {
|
|
// No session/target arg — the caller's identity is resolved from the connection.
|
|
return tool("bridge_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 "
|
|
+ "one of these. A lead uses it only to answer another lead that messaged "
|
|
+ "it — never to answer a worker, whose turn it is not.",
|
|
objectSchema(Map.of(
|
|
"content", stringProp("Your reply/answer")),
|
|
List.of("content")));
|
|
}
|
|
|
|
private static McpSchema.Tool statusTool() {
|
|
return tool("bridge_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")),
|
|
List.of("sessionId")));
|
|
}
|
|
|
|
private static McpSchema.Tool whoamiTool() {
|
|
return tool("bridge_whoami",
|
|
"Report who YOU are on the bridge — your role is resolved from your connection "
|
|
+ "(unforgeable), never from anything you claim. Returns role 'primary' (you "
|
|
+ "orchestrate: spawn/send/stop; 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 "
|
|
+ "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()));
|
|
}
|
|
|
|
// --- small helpers -------------------------------------------------------------------------
|
|
|
|
// The SDK 2.0.0 deprecates its own Tool builders without a stable replacement — isolate it here.
|
|
@SuppressWarnings("deprecation")
|
|
private static McpSchema.Tool tool(String name, String description, Map<String, Object> inputSchema) {
|
|
return McpSchema.Tool.builder(name).description(description).inputSchema(inputSchema).build();
|
|
}
|
|
|
|
private static Map<String, Object> objectSchema(Map<String, Object> properties, List<String> required) {
|
|
return Map.of("type", "object", "properties", properties, "required", required);
|
|
}
|
|
|
|
private static Map<String, Object> stringProp(String description) {
|
|
return Map.of("type", "string", "description", description);
|
|
}
|
|
|
|
private static McpSchema.CallToolResult text(String s) {
|
|
return McpSchema.CallToolResult.builder().addTextContent(s == null ? "" : s).build();
|
|
}
|
|
|
|
private static McpSchema.CallToolResult error(String s) {
|
|
return McpSchema.CallToolResult.builder().addTextContent(s).isError(true).build();
|
|
}
|
|
|
|
private static String str(Map<String, Object> args, String key) {
|
|
Object v = args.get(key);
|
|
return v == null ? null : v.toString();
|
|
}
|
|
|
|
private static Long timeoutMs(Map<String, Object> args) {
|
|
Object v = args.get("timeoutMs");
|
|
return v instanceof Number n ? n.longValue() : null;
|
|
}
|
|
|
|
private static long clamp(long ms) {
|
|
return Math.clamp(ms, 1, MAX_TIMEOUT_MS);
|
|
}
|
|
|
|
private static boolean isBlank(String s) {
|
|
return s == null || s.isBlank();
|
|
}
|
|
}
|