Compare commits

..

1 Commits

Author SHA1 Message Date
Dai Ha d9e45f8c2a CB-566: add fleet charter config
CI / build (pull_request) Successful in 1m13s
CI / contract (pull_request) Successful in 1m15s
2026-08-15 04:49:53 +02:00
13 changed files with 228 additions and 634 deletions
@@ -129,13 +129,13 @@ public final class Bridged {
adapters.add(new ClaudeCodeLauncher(agents, spaces, guard,
claudeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
() -> config.get().fleet()));
() -> config.get().fleet().tabLabel()));
}
if (!opencodeProfiles.isEmpty()) {
adapters.add(new OpenCodeLauncher(agents, spaces,
opencodeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
() -> config.get().fleet()));
() -> config.get().fleet().tabLabel()));
}
AtomicReference<Function<String, Integer>> liveCountRef = new AtomicReference<>(_ -> 0);
PeerLauncher workers = new CompositePeerLauncher(
@@ -562,15 +562,7 @@ public record BridgedConfig(
tabLabel = (tabLabel == null || tabLabel.isBlank()) ? DEFAULT_TAB_LABEL : tabLabel;
}
/**
* A fleet with no configured launch charters — the shape every deployment had before
* CB-566, and what most tests want.
*
* <p>Kept deliberately, even though an overload that drops a new field is normally the
* shape to avoid. It is safe here because nothing <em>reads</em> a charter through a
* constructor: the launcher reads {@code fleet.charters()} from the live config. Jackson
* binds the canonical constructor, so this one cannot swallow an operator's YAML.
*/
/** Convenience constructor for code that does not configure launch charters. */
public Fleet(Map<String, Leader> leaders, Map<String, Slot> architects,
Map<String, Slot> developers, Map<String, Slot> reviewers, String tabLabel) {
this(leaders, architects, developers, reviewers, null, tabLabel);
@@ -1225,7 +1217,7 @@ public record BridgedConfig(
if (fleet == null || fleet.charters().isEmpty()) {
return;
}
List<String> valid = Arrays.stream(MemberRole.values())
List<String> valid = java.util.Arrays.stream(MemberRole.values())
.map(MemberRole::wireName)
.toList();
List<String> bad = new ArrayList<>();
@@ -28,9 +28,9 @@ import java.util.function.Supplier;
* <ul>
* <li><strong>Hot</strong> — re-read per use, so a reload takes effect on the next spawn:
* {@code fleet:} (every role pool, {@code charters}, and {@code tabLabel}),
* {@code placement:}, and an existing profile's {@code weight} / {@code maxLoad}. Those
* three are read through a supplier on {@code CompositePeerLauncher}, which is what makes
* them hot — not the fact that they are config.</li>
* {@code placement:}, and an existing profile's {@code weight} / {@code maxLoad}. Those three are read through a supplier on
* {@code CompositePeerLauncher}, which is what makes them hot — not the fact that they are
* config.</li>
* <li><strong>Deferred</strong> — accepted into the new snapshot, but the wiring built at startup
* keeps the old value until a restart: {@code lifecycle:}, {@code leadHeartbeat:},
* {@code spawnReadyTimeoutMs} / {@code spawnReadyPollMs}, {@code guard:},
@@ -33,7 +33,6 @@ 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.stream.Collectors;
@@ -152,8 +151,8 @@ public final class BridgeMcp {
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());
? sendAsync(messages, target, content, onAccepted)
: send(messages, target, content, timeoutMs(a), onAccepted);
})
// 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.
@@ -380,22 +379,21 @@ 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. */
static McpSchema.CallToolResult send(MessageService messages, String sessionId, String content, Long timeoutMs) {
return send(messages, sessionId, content, timeoutMs, null);
}
/**
* {@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.
*
* As {@link #send(MessageService, String, String, Long)}, wiring an accepted-delivery hook
* (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) {
Long timeoutMs, Runnable onAccepted) {
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);
@@ -465,33 +463,24 @@ public final class BridgeMcp {
/**
* {@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
*/
static McpSchema.CallToolResult sendAsync(MessageService messages, String sessionId, String content) {
return sendAsync(messages, sessionId, content, null);
}
/**
* As {@link #sendAsync(MessageService, String, String)}, wiring 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) {
Runnable onAccepted) {
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)) {
@@ -513,9 +502,6 @@ public final class BridgeMcp {
? "[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() + "]");
};
}
@@ -44,6 +44,23 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
private final SubscriptionGuard guard;
/**
* Standing instruction appended to the worker's system prompt so it returns its result via
* {@code bridge_reply}. Injected as a launch flag, so nothing is written to the worker's
* profile — it is guidance, and a worker that never replies is caught by the send's timeout.
*/
static final String REPLY_CHARTER =
"You are an off-subscription worker 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. 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 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 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.";
/**
* Production constructor — disables the spawn-ready gate ({@code spawnReadyTimeoutMs == 0}) so
* existing deployments and tests keep the legacy non-blocking spawn semantics.
@@ -76,11 +93,11 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs, long spawnReadyPollMs,
Supplier<BridgedConfig.Fleet> fleet) {
Supplier<String> tabLabelTemplate) {
this(agents, spaces, guard, profiles, defaultProfile, env,
spawnReadyTimeoutMs,
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs),
fleet);
tabLabelTemplate);
}
/**
@@ -112,19 +129,31 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
/**
* Full testability constructor, plus the fleet-wide tab-label template (CB-557).
*
* @param fleet live fleet config, read once for each spawn
* @param tabLabelTemplate {@code fleet.tabLabel}; {@code null}/blank ⇒
* {@link BridgedConfig.Fleet#DEFAULT_TAB_LABEL}
*/
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs,
LongSupplier nowMillis, Runnable sleeper,
Supplier<BridgedConfig.Fleet> fleet) {
Supplier<String> tabLabelTemplate) {
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
spawnReadyTimeoutMs, nowMillis, sleeper, fleet);
spawnReadyTimeoutMs, nowMillis, sleeper, tabLabelTemplate);
this.guard = guard;
}
/**
* {@inheritDoc}
*
* <p>A legacy spawn with no session identity is a fresh, launcher-derived session — delegate to
* the session-aware form with no name and no resume id.
*/
@Override
protected Launch buildLaunch(BridgedConfig.Profile cfg) {
return buildLaunch(cfg, null, null);
}
/**
* {@inheritDoc}
*
@@ -135,7 +164,7 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
* applied here — see {@link #applySessionIdentity}.
*/
@Override
protected Launch buildLaunch(BridgedConfig.Profile cfg, LaunchSpec spec) {
protected Launch buildLaunch(BridgedConfig.Profile cfg, String sessionName, String resumeSessionId) {
// CB-539: a profile may deliberately opt into the subscription (subscription: true) when no
// off-subscription endpoint exists for it — e.g. `sonnet` on `ccs`. That profile gets no
// ANTHROPIC_BASE_URL/AUTH_TOKEN (there is nothing to point them at) and the guard's base_url
@@ -181,9 +210,9 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
// id via -r and passes no --session-id (the two conflict). Both are injected before the
// model flag so --model keeps outranking the operator's own argv.
// mutableArgv: argvWithBridge may hand back the profile's own (immutable) List.of when it
// has neither MCP nor a charter — session flags must be added into a list we own.
List<String> argv = mutableArgv(argvWithBridge(cfg, spec.charter()));
String agentSessionId = applySessionIdentity(argv, spec.sessionName(), spec.resumeSessionId());
// has no MCP — session flags must be added into a list we own.
List<String> argv = mutableArgv(argvWithBridge(cfg));
String agentSessionId = applySessionIdentity(argv, sessionName, resumeSessionId);
return new Launch(workerEnv, argvWithModel(argv, cfg), agentSessionId);
}
@@ -217,26 +246,22 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
}
/**
* The launch argv, plus an inline {@code --mcp-config} when {@code worker.mcpUrl} is set and
* {@code --append-system-prompt} when the base composed a charter. Neither touches the profile's
* config; both are pure command-line flags. This inline-flag mount is Claude Code specific —
* other adapters mount MCP and instructions their own way.
* The launch argv, plus — when {@code worker.mcpUrl} is set — inline {@code --mcp-config} for
* the bridge server and {@code --append-system-prompt} for the {@link #REPLY_CHARTER}. Neither
* touches the profile's config; both are pure command-line flags. This inline-flag mount is
* Claude Code specific — other adapters mount MCP and instructions their own way.
*/
private List<String> argvWithBridge(BridgedConfig.Profile cfg, String charter) {
if (!cfg.hasMcp() && charter == null) {
private List<String> argvWithBridge(BridgedConfig.Profile cfg) {
if (!cfg.hasMcp()) {
return cfg.argv();
}
String mcpJson = "{\"mcpServers\":{\"bridge\":{\"type\":\"http\",\"url\":\""
+ cfg.mcpUrl() + "\"}}}";
List<String> argv = mutableArgv(cfg.argv());
if (cfg.hasMcp()) {
String mcpJson = "{\"mcpServers\":{\"bridge\":{\"type\":\"http\",\"url\":\""
+ cfg.mcpUrl() + "\"}}}";
argv.add("--mcp-config");
argv.add(mcpJson);
}
if (charter != null) {
argv.add("--append-system-prompt");
argv.add(charter);
}
argv.add("--mcp-config");
argv.add(mcpJson);
argv.add("--append-system-prompt");
argv.add(REPLY_CHARTER);
return argv;
}
@@ -44,7 +44,7 @@ import java.util.regex.Pattern;
* <li>{@code namePrefix} (constructor arg) — the label prefix ({@code claude}, {@code opencode})
* that drives both unique naming and the orphan-reap pattern, so each adapter reaps only its
* own kind of pane and never another's.</li>
* <li>{@link #buildLaunch(BridgedConfig.Profile, LaunchSpec)} — the peer-specific env map + argv, including any
* <li>{@link #buildLaunch(BridgedConfig.Profile)} — the peer-specific env map + argv, including any
* subscription/guard check, MCP mount, and instruction injection. The base never sees how the
* peer is configured; it only places and starts the returned {@link Launch}.</li>
* </ul>
@@ -80,25 +80,15 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
private final AtomicLong nameSeq = new AtomicLong(); // per-peer counter (herdr agent names only)
/**
* Live fleet config, read once per spawn. A null supplier or value leaves tab labels at their
* default and supplies no role charter. A profile's own {@code tabLabel} still overrides it.
* The {@code fleet.tabLabel} template; a {@code null} supplier or a {@code null}/blank value ⇒
* {@link BridgedConfig.Fleet#DEFAULT_TAB_LABEL}. A profile's own {@code tabLabel} still
* overrides it.
*
* <p>CB-559: a supplier rather than a snapshot, so a config reload affects the next launch
* without a restart. Existing tabs keep the label they were given.
* <p>CB-559: a supplier rather than a String, so a config reload renames the <em>next</em> tab
* without a restart. Existing tabs keep the label they were given — bridged does not rewrite a
* label it already wrote.
*/
private final Supplier<BridgedConfig.Fleet> fleet;
/** 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. "
+ "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 "
+ "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 "
+ "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.";
private final Supplier<String> tabLabelTemplate;
/**
* Tab numbers, counted per {@code role/profile} pair (CB-557).
@@ -152,19 +142,21 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
}
/**
* As above, plus the live {@code fleet} config (CB-557).
* As above, plus the {@code fleet.tabLabel} template (CB-557).
*
* @param fleet live fleet config, read once per spawn; {@code null} ⇒ default tab label and no
* role charter. A separate constructor rather than a new parameter on the one
* above, so every existing call site keeps the default without an edit.
* @param tabLabelTemplate fleet-wide tab-label template, read per spawn (CB-559); {@code null},
* or a supplier yielding {@code null}/blank ⇒
* {@link BridgedConfig.Fleet#DEFAULT_TAB_LABEL}. A separate constructor
* rather than a new parameter on the one above, so every existing call
* site keeps the default without an edit.
*/
protected HerdrPeerLauncher(String namePrefix, AgentControl agents, WorkspaceControl spaces,
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs,
LongSupplier nowMillis, Runnable sleeper,
Supplier<BridgedConfig.Fleet> fleet) {
this.fleet = fleet;
Supplier<String> tabLabelTemplate) {
this.tabLabelTemplate = tabLabelTemplate;
this.namePrefix = namePrefix;
this.agents = agents;
this.spaces = spaces;
@@ -183,7 +175,22 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
* Any subscription/guard check, MCP mount, and instruction injection happen here. The env map
* and argv are adapter-private; the base only places and starts what is returned.
*/
protected abstract Launch buildLaunch(BridgedConfig.Profile cfg, LaunchSpec spec);
protected abstract Launch buildLaunch(BridgedConfig.Profile cfg);
/**
* Session-aware variant of {@link #buildLaunch(BridgedConfig.Profile)} (CB-547a). Default
* discards the session identity and delegates to the profile-only form, so an adapter that
* carries no durable peer session (opencode, say) inherits byte-identical behaviour and needs
* no change. An adapter that does (Claude Code) overrides this to mint/resume the id and to
* surface it on the returned {@link Launch#agentSessionId()}.
*
* @param cfg the resolved profile to spawn
* @param sessionName the bridge's logical session name, or null/blank for launcher-derived
* @param resumeSessionId the peer's own prior session id to resume, or null/blank for fresh
*/
protected Launch buildLaunch(BridgedConfig.Profile cfg, String sessionName, String resumeSessionId) {
return buildLaunch(cfg);
}
/** Direct transport access for peer-specific, non-turn control operations. */
protected final AgentControl agents() {
@@ -218,10 +225,6 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
}
}
/** All per-spawn values adapters may need, including the base-composed effective charter. */
protected record LaunchSpec(String sessionName, String resumeSessionId, MemberRole role, String charter) {
}
// --- profile surface -----------------------------------------------------------------------
/** The configured peer profile names (what {@code spawn(profile)} accepts). */
@@ -292,23 +295,18 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
}
/**
* Spawn a peer with session identity (CB-547a). The session values, role, and charter are
* threaded from the {@link SpawnRequest} into {@link #buildLaunch(BridgedConfig.Profile,
* LaunchSpec)}, and the launch's resolved agent-session id is returned alongside the agent so
* the caller can put it on the {@link PeerHandle}.
* Spawn a peer with session identity (CB-547a). {@code sessionName} and {@code resumeSessionId}
* are threaded from the {@link SpawnRequest} into {@link #buildLaunch(BridgedConfig.Profile,
* String, String)}, and the launch's resolved agent-session id is returned alongside the agent
* so the caller can put it on the {@link PeerHandle}.
*/
protected Spawned spawnInternal(String profileName, String requestedCwd, String callerCwd,
String sessionName, String resumeSessionId, MemberRole role) {
BridgedConfig.Profile cfg = requireProfile(profileName);
BridgedConfig.Fleet liveFleet = fleet == null ? null : fleet.get();
String roleCharter = liveFleet == null ? null : liveFleet.charterFor(role);
String replyCharter = cfg.hasMcp() ? REPLY_CHARTER : null;
String charter = roleCharter == null ? replyCharter
: replyCharter == null ? roleCharter : roleCharter + "\n\n" + replyCharter;
Launch launch = buildLaunch(cfg, new LaunchSpec(sessionName, resumeSessionId, role, charter));
Launch launch = buildLaunch(cfg, sessionName, resumeSessionId);
String cwd = resolveCwd(requestedCwd, cfg, callerCwd);
Agent agent = cfg.tabPlacement()
? spawnInTab(cfg, launch.env(), launch.argv(), cwd, role, liveFleet)
? spawnInTab(cfg, launch.env(), launch.argv(), cwd, role)
: spawnAsPane(cfg, launch.env(), launch.argv(), cwd);
return new Spawned(agent, launch.agentSessionId());
}
@@ -384,7 +382,7 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
/** Dedicated worker space → own tab (carrying cwd+env) → start the peer into the seed pane. */
private Agent spawnInTab(BridgedConfig.Profile cfg, Map<String, String> workerEnv,
List<String> argv, String cwd, MemberRole role, BridgedConfig.Fleet liveFleet) {
List<String> argv, String cwd, MemberRole role) {
Workspace space = spaces.ensureWorkspace(cfg.workspace());
Tab.Created tab = spaces.createTab(space.workspaceId(), cwd, workerEnv);
log.info("spawning {} profile={} space={} tab={} cwd={}",
@@ -417,7 +415,7 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
tidy("label tab " + tab.tab().tabId(),
() -> spaces.renameTab(tab.tab().tabId(),
cfg.renderTabLabel(
liveFleet == null ? null : liveFleet.tabLabel(),
tabLabelTemplate == null ? null : tabLabelTemplate.get(),
role, nextLabelSeq(role, cfg.profile()))));
log.info("{} started pane={} tab={} terminal={}",
namePrefix, started.agent().paneId(), started.agent().tabId(), started.agent().terminalId());
@@ -38,7 +38,7 @@ import java.util.function.Supplier;
* <li><strong>File-based MCP mount + instructions.</strong> opencode has no inline
* {@code --mcp-config}/{@code --append-system-prompt}. Instead the bridge writes an ephemeral
* {@code opencode.json} that declares the bridge as a {@code remote} MCP server and lists a
* member-charter file under {@code instructions}, then points the worker at it with
* reply-charter file under {@code instructions}, then points the worker at it with
* {@code OPENCODE_CONFIG}. This is the one place the launcher touches disk — Claude never did.</li>
* <li><strong>Model as a flag.</strong> the {@code provider/model} selector is passed as
* {@code -m}, not an env var.</li>
@@ -54,9 +54,38 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
/** Writer for the generated {@code opencode.json}. */
private static final ObjectMapper JSON = new ObjectMapper();
/**
* Standing instruction written to the charter file and mounted via the config's
* {@code instructions} so the worker returns its result through {@code bridge_reply}. Kept on
* disk (not a launch flag) because opencode's {@code instructions} takes file paths, not inline
* text — the file is regenerated per spawn and never touches the worker's own profile.
*/
static final String REPLY_CHARTER =
"You are an off-subscription worker in the claude-bridge fleet, running under opencode. "
+ "Every message you receive arrives through the bridge, and the ONLY channel back to the "
+ "sender is the bridge_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 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 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.";
/** Root under which per-spawn opencode config dirs are created (injectable for tests). */
private final Path configRoot;
/**
* The current spawn's resume-target session id, threaded from {@link #spawn(SpawnRequest)} to
* {@link #buildLaunch} across the base's {@code spawn -> spawnInternal -> buildLaunch} chain,
* which carries no request. A plain field would race under concurrent spawns (the base supports
* them), so it is thread-local: each spawn captures its own request's id on its own thread, and
* {@code buildLaunch}, synchronous and same-thread, reads exactly that one. Set only around the
* {@code super.spawn} call and cleared in {@code finally}, so a paused/leftover value can never
* bleed into the next spawn.
*/
private final ThreadLocal<String> resumeSessionId = new ThreadLocal<>();
/**
* Session discovery against opencode's on-disk storage ({@link OpenCodeSessionDiscovery}) —
* the one seam that knows opencode's private session-file layout. Its root is injectable for
@@ -97,10 +126,10 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs, long spawnReadyPollMs,
Supplier<BridgedConfig.Fleet> fleet) {
Supplier<String> tabLabelTemplate) {
this(agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs,
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs),
defaultConfigRoot(), defaultDiscoveryRoot(), fleet);
defaultConfigRoot(), defaultDiscoveryRoot(), tabLabelTemplate);
}
/**
@@ -135,7 +164,8 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
/**
* Full testability constructor, plus the fleet-wide tab-label template (CB-557).
*
* @param fleet live fleet config, read once for each spawn
* @param tabLabelTemplate {@code fleet.tabLabel}; {@code null}/blank ⇒
* {@link BridgedConfig.Fleet#DEFAULT_TAB_LABEL}
*/
public OpenCodeLauncher(AgentControl agents, WorkspaceControl spaces,
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
@@ -143,9 +173,9 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
long spawnReadyTimeoutMs,
LongSupplier nowMillis, Runnable sleeper,
Path configRoot, Path discoveryRoot,
Supplier<BridgedConfig.Fleet> fleet) {
Supplier<String> tabLabelTemplate) {
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
spawnReadyTimeoutMs, nowMillis, sleeper, fleet);
spawnReadyTimeoutMs, nowMillis, sleeper, tabLabelTemplate);
this.configRoot = configRoot;
this.discovery = new OpenCodeSessionDiscovery(discoveryRoot);
}
@@ -163,21 +193,20 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
* {@inheritDoc}
*
* <p>Builds the opencode launch: no {@code ANTHROPIC_*} and no guard (opencode reads its own
* provider credentials); when the profile mounts the bridge MCP or has a member charter,
* generate an ephemeral {@code opencode.json} (remote MCP server + member-charter instructions)
* and point the worker at it via {@code OPENCODE_CONFIG}; carry the parity-neutral git-forge
* grant; and select the model with {@code -m}.
* provider credentials); when the profile mounts the bridge MCP, generate an ephemeral
* {@code opencode.json} (remote MCP server + reply-charter instructions) and point the worker at
* it via {@code OPENCODE_CONFIG}; carry the parity-neutral git-forge grant; and select the model
* with {@code -m}.
*/
@Override
protected Launch buildLaunch(BridgedConfig.Profile cfg, LaunchSpec spec) {
protected Launch buildLaunch(BridgedConfig.Profile cfg) {
Map<String, String> workerEnv = baseEnv(cfg);
// A config file is needed for the bridge MCP mount, a member charter, or a pinned endpoint (CB-508).
if (cfg.hasMcp() || spec.charter() != null || hasCustomProvider(cfg)) {
workerEnv.put("OPENCODE_CONFIG", writeConfig(cfg, spec.charter()).toString());
// A config file is needed for the bridge MCP mount, for a pinned endpoint (CB-508), or both.
if (cfg.hasMcp() || hasCustomProvider(cfg)) {
workerEnv.put("OPENCODE_CONFIG", writeConfig(cfg).toString());
}
applyGitToken(workerEnv, cfg);
return new Launch(workerEnv,
argvWithResume(argvWithModel(argvWithAuto(cfg), cfg), spec.resumeSessionId()));
return new Launch(workerEnv, argvWithResume(argvWithModel(argvWithAuto(cfg), cfg)));
}
/**
@@ -214,9 +243,11 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
* The launch argv plus, on a resumed spawn, opencode's {@code -s <id>} flag to continue a prior
* conversation by its session id. {@code -s, --session <id>} resumes an existing session; on a
* fresh spawn (no resume target) no flag is added, letting opencode start a brand-new session.
* The id comes from the base launch spec.
* The id comes from the current spawn request's {@code resumeSessionId}, threaded per-thread by
* {@link #spawn(SpawnRequest)}.
*/
private List<String> argvWithResume(List<String> argv, String id) {
private List<String> argvWithResume(List<String> argv) {
String id = resumeSessionId.get();
if (id == null || id.isBlank()) {
return argv;
}
@@ -236,12 +267,12 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
}
/**
* Write an ephemeral {@code opencode.json} (and the member-charter file it references) into a
* Write an ephemeral {@code opencode.json} (and the reply-charter file it references) into a
* fresh per-spawn directory under {@link #configRoot}, and return the config file's path for
* {@code OPENCODE_CONFIG}. The dir is unique per spawn so concurrent workers never race on it;
* it is best-effort cleaned on JVM exit (worker config is disposable — regenerated every spawn).
*/
private Path writeConfig(BridgedConfig.Profile cfg, String charterText) {
private Path writeConfig(BridgedConfig.Profile cfg) {
try {
Path dir = Files.createTempDirectory(configRoot, "bridged-opencode-");
dir.toFile().deleteOnExit();
@@ -264,19 +295,16 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
// if per-profile control is ever wanted, add a profile knob rather than dropping this.
root.putObject("compaction").put("auto", true);
if (charterText != null) {
Path charter = dir.resolve("member-charter.md");
Files.writeString(charter, charterText);
if (cfg.hasMcp()) {
Path charter = dir.resolve("reply-charter.md");
Files.writeString(charter, REPLY_CHARTER);
charter.toFile().deleteOnExit();
root.putArray("instructions").add(charter.toAbsolutePath().toString());
}
if (cfg.hasMcp()) {
ObjectNode bridge = root.putObject("mcp").putObject("bridge");
bridge.put("type", "remote");
bridge.put("url", cfg.mcpUrl());
bridge.put("enabled", true);
root.putArray("instructions").add(charter.toAbsolutePath().toString());
}
if (hasCustomProvider(cfg)) {
addCustomProvider(root, cfg);
@@ -352,11 +380,31 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
return afterScheme.contains("/") ? trimmed : trimmed + "/v1";
}
/** Add lazy on-disk session discovery to the base handle. */
/**
* {@inheritDoc}
*
* <p>adds this adapter's session-identity work around the base's spawn — as opencode cannot be
* told its session id at spawn (see {@link Capability#SESSION_RESUME} vs
* {@link Capability#SESSION_NAME}), identity is only ever adopted after the fact:
* <ul>
* <li>the request's {@code resumeSessionId} is remembered for {@link #buildLaunch} to turn
* into {@code -s <id>}; and</li>
* <li>the returned handle is wrapped so its
* {@link dev.ltms.bridged.peer.PeerHandle#agentSessionId()} performs lazy session
* discovery against opencode's storage (see {@link OpenCodeSessionDiscovery}) — always
* non-blocking, {@code null} until opencode has persisted the session record.</li>
* </ul>
*/
@Override
public PeerHandle spawn(SpawnRequest req) {
PeerHandle inner = super.spawn(req);
return new SessionAwareHandle(inner, discovery, effectiveCwd(req));
resumeSessionId.set(req.resumeSessionId());
try {
PeerHandle inner = super.spawn(req);
return new SessionAwareHandle(inner, discovery, effectiveCwd(req));
} finally {
// Never let a paused/leftover resume id bleed into the next spawn on this thread.
resumeSessionId.remove();
}
}
/**
@@ -127,8 +127,6 @@ 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. */
ASKING,
/** The worker's turn finished; {@link TaskView#reply} holds the answer. */
DONE,
/** The delegation could not complete (timed out, worker gone, or busy). */
@@ -138,28 +136,16 @@ public final class MessageService {
/**
* A poll snapshot of an async delegation.
*
* @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 reply the answer when {@link #phase} is {@link Phase#DONE}, else {@code null}
* @param replySource {@code "reply"} (structured {@code bridge_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}
* @param detail a human note (live worker status while pending, or the failure reason)
*/
public record TaskView(String ticket, Phase phase, String reply, String replySource, String detail,
String turnId) {
public record TaskView(String ticket, Phase phase, String reply, String replySource, String detail) {
}
/** An in-flight or finished async delegation, keyed by its ticket. */
private static final class Task {
private final String target;
private final CompletableFuture<Reply> future = new CompletableFuture<>();
private final long createdNanos = System.nanoTime();
private volatile Reply question;
private volatile String turnId;
private Task(String target) {
this.target = target;
}
private record Task(String target, CompletableFuture<Reply> future, long createdNanos) {
}
private final AgentControl agents;
@@ -170,10 +156,6 @@ public final class MessageService {
private final Metrics metrics; // CB-502: nullable — no registry in unit tests
private final ConcurrentHashMap<String, ReentrantLock> sessionLocks = new ConcurrentHashMap<>();
private final ConcurrentHashMap<String, Task> tasks = new ConcurrentHashMap<>();
/** The async task that currently owns a target's send lock. */
private final ConcurrentHashMap<String, Task> asyncTasksByTarget = new ConcurrentHashMap<>();
/** Async tickets paused on a specific {@code bridge_ask} turn. */
private final ConcurrentHashMap<String, Task> asyncTasksByTurn = new ConcurrentHashMap<>();
private final AtomicLong ticketSeq = new AtomicLong();
private final ExecutorService asyncExecutor = Executors.newThreadPerTaskExecutor(
Thread.ofVirtual().name("bridge-async-", 0).factory());
@@ -351,9 +333,6 @@ public final class MessageService {
return new Reply(Outcome.BUSY, null); // another send held the session the whole window
}
try {
if (hasAsyncQuestion(target)) {
return new Reply(Outcome.BUSY, null); // the worker's current turn is paused for its lead
}
// Open the waiter BEFORE queueing delivery (CB-548). A fast reply — the worker already
// injectable the instant we enqueue — otherwise arrives before the waiter is registered
// and orphans into the inbox while this send blocks to the timeout (the enqueue-before-
@@ -414,14 +393,12 @@ public final class MessageService {
rendezvous.closeAsk(ticket.turnId());
return new AskResult(AskOutcome.NO_WAITER, null); // no primary is blocked on this worker
}
markAsyncQuestion(workerSession, question, ticket.turnId());
}
try {
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);
clearAsyncQuestion(ticket.turnId(), true);
return new AskResult(AskOutcome.TIMED_OUT, null);
} catch (ExecutionException e) {
Throwable cause = e.getCause();
@@ -464,12 +441,9 @@ public final class MessageService {
rendezvous.close(workerSession, reply);
return new Reply(Outcome.STALE_TURN, null); // lapsed between the lookup and the unblock
}
clearAsyncQuestion(turnId, false);
try {
Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS);
Reply result = new Reply(outcomeOf(r.kind()), r.text(), r.turnId());
finishAsyncTask(turnId, result);
return result;
return new Reply(outcomeOf(r.kind()), r.text(), r.turnId());
} catch (TimeoutException e) {
// The worker resumed but hasn't replied yet — no completion fallback arms an answered
// turn (it never re-entered the injector), so a silent worker rides out the window.
@@ -509,27 +483,9 @@ public final class MessageService {
*/
public String sendAsync(String target, String content, Runnable onAccepted) {
String ticket = "task-" + ticketSeq.incrementAndGet();
Task task = new Task(target);
tasks.put(ticket, task);
asyncExecutor.submit(() -> {
try {
Runnable trackingAccepted = () -> {
if (onAccepted != null) {
onAccepted.run();
}
asyncTasksByTarget.put(target, task);
};
Reply result = send(target, content, ASYNC_TIMEOUT_MS, trackingAccepted);
if (result.outcome() == Outcome.QUESTION) {
asyncTasksByTarget.remove(target, task);
} else {
finishAsyncTask(task, result);
}
} catch (Throwable t) {
task.future.completeExceptionally(t);
asyncTasksByTarget.remove(target, task);
}
});
CompletableFuture<Reply> future = CompletableFuture.supplyAsync(
() -> send(target, content, ASYNC_TIMEOUT_MS, onAccepted), asyncExecutor);
tasks.put(ticket, new Task(target, future, System.nanoTime()));
pruneTerminalTickets();
log.debug("async send {} -> {}", ticket, target);
return ticket;
@@ -545,32 +501,27 @@ public final class MessageService {
if (task == null) {
return null;
}
CompletableFuture<Reply> f = task.future;
CompletableFuture<Reply> f = task.future();
if (!f.isDone()) {
Reply question = task.question;
if (question != null) {
return new TaskView(ticket, Phase.ASKING, question.text(), null,
"worker is waiting for your answer", question.turnId());
}
return new TaskView(ticket, Phase.PENDING, null, null, "worker " + liveStatus(task.target), null);
return new TaskView(ticket, Phase.PENDING, null, null, "worker " + liveStatus(task.target()));
}
Reply r;
try {
r = f.getNow(null);
} catch (CompletionException | java.util.concurrent.CancellationException e) {
Throwable cause = (e instanceof CompletionException ce && ce.getCause() != null) ? ce.getCause() : e;
return new TaskView(ticket, Phase.FAILED, null, null, cause.getMessage(), null);
return new TaskView(ticket, Phase.FAILED, null, null, cause.getMessage());
}
if (r.completed()) {
String source = r.outcome() == Outcome.REPLIED ? "reply" : "transcript";
return new TaskView(ticket, Phase.DONE, r.text(), source, null, null);
return new TaskView(ticket, Phase.DONE, r.text(), source, null);
}
// A wedged worker (CB-109) carries the error context as its reason; the timeout/busy
// outcomes carry none, so fall back to the outcome name.
String detail = r.outcome() == Outcome.WORKER_FAILED && r.text() != null
? r.text()
: "no reply — " + r.outcome().name().toLowerCase();
return new TaskView(ticket, Phase.FAILED, null, null, detail, null);
return new TaskView(ticket, Phase.FAILED, null, null, detail);
}
/** Best-effort live worker status for a pending poll; never throws (a lookup error is just noise). */
@@ -585,51 +536,7 @@ public final class MessageService {
/** Drop finished tickets older than the TTL so the registry cannot grow without bound. */
private void pruneTerminalTickets() {
long cutoff = System.nanoTime() - TICKET_TTL_NANOS;
tasks.values().removeIf(t -> t.future.isDone() && t.createdNanos < cutoff);
}
/** Record the active question for an async ticket; blocking sends have no entry and stay unchanged. */
private void markAsyncQuestion(String target, String text, String turnId) {
Task task = asyncTasksByTarget.get(target);
if (task != null) {
task.question = new Reply(Outcome.QUESTION, text, turnId);
task.turnId = turnId;
asyncTasksByTurn.put(turnId, task);
}
}
/** Clear an answered or lapsed question, but only when it matches the ticket's current turn. */
private void clearAsyncQuestion(String turnId, boolean forgetTurn) {
Task task = asyncTasksByTurn.get(turnId);
if (task != null && turnId.equals(task.turnId)) {
task.question = null;
if (forgetTurn) {
asyncTasksByTurn.remove(turnId, task);
task.turnId = null;
}
}
}
/** Complete and detach an async ticket after its worker's actual terminal reply. */
private void finishAsyncTask(Task task, Reply result) {
task.future.complete(result);
asyncTasksByTarget.remove(task.target, task);
if (task.turnId != null) {
asyncTasksByTurn.remove(task.turnId, task);
}
}
/** Complete the async ticket correlated to a specific answered turn. */
private void finishAsyncTask(String turnId, Reply result) {
Task task = asyncTasksByTurn.get(turnId);
if (task != null) {
finishAsyncTask(task, result);
}
}
/** A new send must not open a waiter while an async ticket owns this worker's paused turn. */
private boolean hasAsyncQuestion(String target) {
return asyncTasksByTurn.values().stream().anyMatch(task -> target.equals(task.target));
tasks.values().removeIf(t -> t.future().isDone() && t.createdNanos() < cutoff);
}
/** Release the async executor. */
@@ -53,18 +53,6 @@ class BridgeMcpTest {
return ((McpSchema.TextContent) r.content().getFirst()).text();
}
private void assertSendRoundTrips(String target, Set<String> profiles) throws Exception {
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
() -> BridgeMcp.send(messages, target, "hi", 4000L, null, profiles));
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(target) && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting(target), "send should be accepted for " + target);
BridgeMcp.reply(messages, target, "received");
assertEquals("received", textOf(send.get(6, TimeUnit.SECONDS)));
}
private static ClaudeCodeLauncher workerService(FakeHerdr h, String baseUrl, Set<String> allow) {
BridgedConfig.Profile cfg = new BridgedConfig.Profile(
"ltms-local", baseUrl, "coder", null, "BRIDGED_WORKER_TOKEN", null,
@@ -81,7 +69,7 @@ class BridgeMcpTest {
void sendThenReplyRoundTrips() throws Exception {
// bridge_send blocks; bridge_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()));
() -> BridgeMcp.send(messages, "term_a", "review this", 4000L));
// Wait until the send has opened its waiter so the reply resolves it (CB-307: reply now
// queues in the inbox if no waiter is open, which would break the round-trip).
@@ -103,7 +91,7 @@ class BridgeMcpTest {
@Test
void asyncSendReturnsATicketThenPollReportsTheReply() throws Exception {
// wait:false parity — a ticket is issued, resolved by a reply, and surfaced by bridge_poll.
McpSchema.CallToolResult accepted = BridgeMcp.sendAsync(messages, "term_a", "do it", null, Set.of());
McpSchema.CallToolResult accepted = BridgeMcp.sendAsync(messages, "term_a", "do it");
assertNotEquals(Boolean.TRUE, accepted.isError());
String out = textOf(accepted);
assertTrue(out.contains("ticket="), out);
@@ -132,131 +120,6 @@ class BridgeMcpTest {
assertEquals("async LGTM", textOf(polled));
}
@Test
void asyncSendSurfacesAnAskThenKeepsTheTicketForTheFinalReply() throws Exception {
McpSchema.CallToolResult accepted = BridgeMcp.sendAsync(messages, "term_a", "do it", null, Set.of());
String ticket = textOf(accepted).substring(textOf(accepted).indexOf("ticket=") + "ticket=".length()).trim();
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting("term_a"));
CompletableFuture<McpSchema.CallToolResult> ask = CompletableFuture.supplyAsync(
() -> BridgeMcp.ask(messages, "term_a", "which config?", 5000L));
McpSchema.CallToolResult question = BridgeMcp.poll(messages, ticket, null);
deadline = System.currentTimeMillis() + 3000;
while (!textOf(question).contains("[question") && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
question = BridgeMcp.poll(messages, ticket, null);
}
assertTrue(textOf(question).contains("which config?"), textOf(question));
String questionText = textOf(question);
String afterTurnId = questionText.substring(questionText.indexOf("turnId=\"") + "turnId=\"".length());
String turnId = afterTurnId.substring(0, afterTurnId.indexOf('"'));
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
() -> BridgeMcp.answer(messages, turnId, "config.yaml", 5000L));
assertEquals("config.yaml", textOf(ask.get(6, TimeUnit.SECONDS)));
deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting("term_a"));
BridgeMcp.reply(messages, "term_a", "done");
assertEquals("done", textOf(answer.get(6, TimeUnit.SECONDS)));
McpSchema.CallToolResult done = BridgeMcp.poll(messages, ticket, null);
deadline = System.currentTimeMillis() + 3000;
while (!"done".equals(textOf(done)) && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
done = BridgeMcp.poll(messages, ticket, null);
}
assertEquals("done", textOf(done));
}
@Test
void unansweredAsyncAskReturnsTheTicketToPending() throws Exception {
McpSchema.CallToolResult accepted = BridgeMcp.sendAsync(messages, "term_a", "do it", null, Set.of());
String ticket = textOf(accepted).substring(textOf(accepted).indexOf("ticket=") + "ticket=".length()).trim();
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting("term_a"));
McpSchema.CallToolResult ask = BridgeMcp.ask(messages, "term_a", "still there?", 50L);
assertTrue(textOf(ask).contains("no answer"), textOf(ask));
assertTrue(textOf(BridgeMcp.poll(messages, ticket, null)).startsWith("[pending"));
BridgeMcp.reply(messages, "term_a", "finished after timeout");
assertEquals("finished after timeout", messages.drainReplies("term_a").getFirst().content());
}
@Test
void asyncSendFailureDoesNotLeaveItsTicketPending() throws Exception {
String ticket = messages.sendAsync("term_a", "do it", () -> {
throw new IllegalStateException("accept failed");
});
long deadline = System.currentTimeMillis() + 3000;
MessageService.TaskView view = messages.poll(ticket);
while (view.phase() == MessageService.Phase.PENDING && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
view = messages.poll(ticket);
}
assertEquals(MessageService.Phase.FAILED, view.phase());
assertEquals("accept failed", view.detail());
}
@Test
void anotherAsyncTicketCannotCaptureAReplyWhileTheFirstTicketIsAsking() throws Exception {
McpSchema.CallToolResult firstAccepted = BridgeMcp.sendAsync(messages, "term_a", "first", null, Set.of());
String firstTicket = textOf(firstAccepted).substring(textOf(firstAccepted).indexOf("ticket=") + "ticket=".length()).trim();
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting("term_a"));
CompletableFuture<McpSchema.CallToolResult> ask = CompletableFuture.supplyAsync(
() -> BridgeMcp.ask(messages, "term_a", "which config?", 5000L));
MessageService.TaskView first = messages.poll(firstTicket);
deadline = System.currentTimeMillis() + 3000;
while (first.phase() != MessageService.Phase.ASKING && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
first = messages.poll(firstTicket);
}
assertEquals(MessageService.Phase.ASKING, first.phase());
String firstTurnId = first.turnId();
McpSchema.CallToolResult secondAccepted = BridgeMcp.sendAsync(messages, "term_a", "second", null, Set.of());
String secondTicket = textOf(secondAccepted).substring(textOf(secondAccepted).indexOf("ticket=") + "ticket=".length()).trim();
MessageService.TaskView second = messages.poll(secondTicket);
deadline = System.currentTimeMillis() + 3000;
while (second.phase() == MessageService.Phase.PENDING && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
second = messages.poll(secondTicket);
}
assertEquals(MessageService.Phase.FAILED, second.phase());
BridgeMcp.reply(messages, "term_a", "late reply");
assertEquals("late reply", messages.drainReplies("term_a").getFirst().content());
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
() -> BridgeMcp.answer(messages, firstTurnId, "config.yaml", 5000L));
assertEquals("config.yaml", textOf(ask.get(6, TimeUnit.SECONDS)));
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
}
BridgeMcp.reply(messages, "term_a", "done");
assertEquals("done", textOf(answer.get(6, TimeUnit.SECONDS)));
}
@Test
void pollUnknownTicketIsAnError() {
McpSchema.CallToolResult res = BridgeMcp.poll(messages, "task-999", null);
@@ -266,40 +129,15 @@ class BridgeMcpTest {
@Test
void sendTimesOutWithAWorkingNote() {
McpSchema.CallToolResult res = BridgeMcp.send(messages, "term_a", "hi", 120L, null, Set.of());
McpSchema.CallToolResult res = BridgeMcp.send(messages, "term_a", "hi", 120L);
assertNotEquals(Boolean.TRUE, res.isError(), "a timeout is informational, not a tool error");
assertTrue(textOf(res).contains("no reply"), "got: " + textOf(res));
}
@Test
void sendRejectsMissingArgs() {
assertTrue(BridgeMcp.send(messages, null, "hi", null, null, Set.of()).isError());
assertTrue(BridgeMcp.send(messages, "term_a", " ", null, null, Set.of()).isError());
}
@Test
void sendRejectsAConfiguredProfileNameBeforeAcceptingIt() {
McpSchema.CallToolResult blocking = BridgeMcp.send(messages, "sol", "hi", 100L, null, Set.of("sol"));
McpSchema.CallToolResult async = BridgeMcp.sendAsync(messages, "sol", "hi", null, Set.of("sol"));
assertTrue(blocking.isError());
assertTrue(async.isError());
assertTrue(textOf(blocking).contains("sol"));
assertTrue(textOf(blocking).contains("configured profile name"));
assertTrue(textOf(blocking).contains("bridge_list"));
assertFalse(textOf(async).contains("ticket="));
}
@Test
void sendAllowsPeerLeadMemberAndUnclassifiedTargets() throws Exception {
Set<String> profiles = Set.of("sol");
assertSendRoundTrips("term_peer_lead", profiles);
assertSendRoundTrips("term_live_member", profiles);
// A herdr-owned pane outside the bridge roster cannot be classified at accept time.
McpSchema.CallToolResult result = BridgeMcp.send(messages, "external-pane", "hi", 10L, null, profiles);
assertFalse(result.isError(), "an unclassified target must not be rejected at acceptance time");
assertTrue(BridgeMcp.send(messages, null, "hi", null).isError());
assertTrue(BridgeMcp.send(messages, "term_a", " ", null).isError());
}
@Test
@@ -335,7 +173,7 @@ class BridgeMcpTest {
void askThenAnswerRoundTrips() throws Exception {
// The primary delegates and blocks; wait until its waiter is open before the worker asks.
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
() -> BridgeMcp.send(messages, "term_a", "do X", 5000L, null, Set.of()));
() -> BridgeMcp.send(messages, "term_a", "do X", 5000L));
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
//noinspection BusyWait
@@ -102,53 +102,6 @@ class ClaudeCodeLauncherTest {
"the operator's own args are preserved, in order, ahead of the model flag");
}
@Test
void appendsTheBaseComposedRoleAndReplyCharter() {
FakeHerdr herdr = new FakeHerdr();
String roleCharter = "You review changes.";
BridgedConfig.Profile cfg = new BridgedConfig.Profile(
"sonnet", "http://gx00.gw:8000", null, null, "BRIDGED_WORKER_TOKEN",
List.of("claude"), "tab", "bridged-workers", "w #{n}",
"http://127.0.0.1:8765/mcp", null, null);
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null,
0, 0L, () -> fleet(Map.of("reviewer", roleCharter), null));
svc.spawn(new SpawnRequest("sonnet", null, null, null, null, MemberRole.REVIEWER));
List<String> args = spawnedArgs(herdr);
int flag = args.indexOf("--append-system-prompt");
assertEquals(1, args.stream().filter("--append-system-prompt"::equals).count(),
"the composed charter is passed once");
assertEquals(roleCharter + "\n\n" + HerdrPeerLauncher.REPLY_CHARTER, args.get(flag + 1),
"the role charter comes first and the reply rule comes last");
}
@Test
void profileWithoutMcpOrRoleCharterGetsNoSystemPrompt() {
FakeHerdr herdr = new FakeHerdr();
ClaudeCodeLauncher svc = labelService(herdr, () -> fleet(Map.of(), null));
svc.spawn(new SpawnRequest("sonnet", null, null, null, null, MemberRole.DEV));
assertFalse(spawnedArgs(herdr).contains("--append-system-prompt"));
}
@Test
void profileWithoutMcpStillGetsItsRoleCharter() {
FakeHerdr herdr = new FakeHerdr();
String roleCharter = "You design changes.";
ClaudeCodeLauncher svc = labelService(herdr, () -> fleet(Map.of("architect", roleCharter), null));
svc.spawn(new SpawnRequest("sonnet", null, null, null, null, MemberRole.ARCHITECT));
List<String> args = spawnedArgs(herdr);
int flag = args.indexOf("--append-system-prompt");
assertTrue(flag >= 0, "a role charter does not need an MCP mount");
assertEquals(roleCharter, args.get(flag + 1));
assertFalse(args.contains("--mcp-config"));
}
private ClaudeCodeLauncher multiProfile(FakeHerdr herdr) {
BridgedConfig.Profile gx10 = new BridgedConfig.Profile("gx10", "http://gx10.gw:8000", "coder",
null, "BRIDGED_WORKER_TOKEN", List.of("claude"), "tab", "bridged-workers", "w #{n}", null, null, null);
@@ -843,18 +796,14 @@ class ClaudeCodeLauncherTest {
.toList();
}
private static BridgedConfig.Fleet fleet(Map<String, String> charters, String tabLabel) {
return new BridgedConfig.Fleet(Map.of(), Map.of(), Map.of(), Map.of(), charters, tabLabel);
}
/** A profile with no {@code tabLabel:} of its own — the fleet template decides. */
private ClaudeCodeLauncher labelService(FakeHerdr herdr, Supplier<BridgedConfig.Fleet> fleet) {
private ClaudeCodeLauncher labelService(FakeHerdr herdr, Supplier<String> fleetTemplate) {
BridgedConfig.Profile cfg = new BridgedConfig.Profile(
"sonnet", "http://gx00.gw:8000", "sonnet", null, "BRIDGED_WORKER_TOKEN",
List.of("claude"), "tab", "bridged-workers", null, null, null, null);
return new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
_ -> null, 0, 0L, fleet);
_ -> null, 0, 0L, fleetTemplate);
}
/**
@@ -865,7 +814,7 @@ class ClaudeCodeLauncherTest {
@Test
void theFleetTemplateNamesTheRoleTheMemberWasSpawnedFor() {
FakeHerdr herdr = new FakeHerdr();
ClaudeCodeLauncher svc = labelService(herdr, () -> fleet(null, "{role}: {profile} #{n}"));
ClaudeCodeLauncher svc = labelService(herdr, () -> "{role}: {profile} #{n}");
svc.spawn(new SpawnRequest("sonnet", null, null, null, null, MemberRole.REVIEWER));
@@ -876,7 +825,7 @@ class ClaudeCodeLauncherTest {
@Test
void theCounterRunsPerRoleAndProfileNotPerFleet() {
FakeHerdr herdr = new FakeHerdr();
ClaudeCodeLauncher svc = labelService(herdr, () -> fleet(null, "{role}: {profile} #{n}"));
ClaudeCodeLauncher svc = labelService(herdr, () -> "{role}: {profile} #{n}");
svc.spawn(new SpawnRequest("sonnet", null, null, null, null, MemberRole.DEV));
svc.spawn(new SpawnRequest("sonnet", null, null, null, null, MemberRole.REVIEWER));
@@ -890,7 +839,7 @@ class ClaudeCodeLauncherTest {
@Test
void aBlankFleetTemplateFallsBackToTheRoleFirstDefault() {
FakeHerdr herdr = new FakeHerdr();
labelService(herdr, () -> fleet(null, null)).spawn(
labelService(herdr, () -> null).spawn(
new SpawnRequest("sonnet", null, null, null, null, MemberRole.ARCHITECT));
assertEquals(List.of("architect: sonnet #1"), tabLabels(herdr));
@@ -906,7 +855,7 @@ class ClaudeCodeLauncherTest {
List.of("claude"), "tab", "bridged-workers", "pinned {profile}", null, null, null);
new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
_ -> null, 0, 0L, () -> fleet(null, "{role}: {profile} #{n}"))
_ -> null, 0, 0L, () -> "{role}: {profile} #{n}")
.spawn(new SpawnRequest("sonnet", null, null, null, null, MemberRole.REVIEWER));
assertEquals(List.of("pinned sonnet"), tabLabels(herdr));
@@ -921,7 +870,7 @@ class ClaudeCodeLauncherTest {
void theTemplateIsReadOnEverySpawnSoAnEditTakesEffect() {
FakeHerdr herdr = new FakeHerdr();
AtomicReference<String> template = new AtomicReference<>("{role}: {profile} #{n}");
ClaudeCodeLauncher svc = labelService(herdr, () -> fleet(null, template.get()));
ClaudeCodeLauncher svc = labelService(herdr, template::get);
svc.spawn(new SpawnRequest("sonnet", null, null, null, null, MemberRole.DEV));
template.set("[{profile}] {role} {n}");
@@ -77,7 +77,7 @@ class CompositePeerLauncherTest {
}
@Override
protected Launch buildLaunch(BridgedConfig.Profile cfg, LaunchSpec spec) {
protected Launch buildLaunch(BridgedConfig.Profile cfg) {
return new Launch(Map.of(), List.of());
}
@@ -1,71 +0,0 @@
package dev.ltms.bridged.member;
import dev.ltms.bridged.config.BridgedConfig;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.FakeHerdr;
import dev.ltms.bridged.herdr.WorkspaceControl;
import dev.ltms.bridged.peer.Capability;
import dev.ltms.bridged.peer.MemberRole;
import dev.ltms.bridged.peer.SpawnRequest;
import org.junit.jupiter.api.Test;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.Supplier;
import static org.junit.jupiter.api.Assertions.assertEquals;
class HerdrPeerLauncherCharterTest {
@Test
void readsAndComposesTheFleetCharterForEachSpawn() {
AtomicReference<BridgedConfig.Fleet> fleet = new AtomicReference<>(fleet(Map.of()));
CapturingLauncher launcher = new CapturingLauncher(fleet::get);
launcher.spawn(new SpawnRequest("mcp", null, null, null, null, MemberRole.DEV));
fleet.set(fleet(Map.of("dev", "role charter")));
launcher.spawn(new SpawnRequest("mcp", null, null, null, null, MemberRole.DEV));
launcher.spawn(new SpawnRequest("no-mcp", null, null, null, null, MemberRole.DEV));
assertEquals(HerdrPeerLauncher.REPLY_CHARTER, launcher.specs.get(0).charter(),
"without a role charter, MCP profiles receive only the reply charter");
assertEquals("role charter\n\n" + HerdrPeerLauncher.REPLY_CHARTER, launcher.specs.get(1).charter(),
"the changed supplier value is read for the next spawn and the reply rule is last");
assertEquals("role charter", launcher.specs.get(2).charter(),
"a role charter does not depend on an MCP mount");
}
private static BridgedConfig.Fleet fleet(Map<String, String> charters) {
return new BridgedConfig.Fleet(Map.of(), Map.of(), Map.of(), Map.of(), charters, null);
}
private static final class CapturingLauncher extends HerdrPeerLauncher {
private final List<LaunchSpec> specs = new ArrayList<>();
CapturingLauncher(Supplier<BridgedConfig.Fleet> fleet) {
super("test", new AgentControl(new FakeHerdr()), new WorkspaceControl(new FakeHerdr()),
Map.of("mcp", profile("mcp", "http://bridge"),
"no-mcp", profile("no-mcp", null)),
"mcp", _ -> null, 0, () -> 0L, () -> { }, fleet);
}
@Override
protected Launch buildLaunch(BridgedConfig.Profile cfg, LaunchSpec spec) {
specs.add(spec);
return new Launch(Map.of(), List.of("test"));
}
@Override
public Set<Capability> capabilities() {
return Set.of();
}
private static BridgedConfig.Profile profile(String name, String mcpUrl) {
return new BridgedConfig.Profile(name, "http://gx00.gw:8000", null, null,
"BRIDGED_WORKER_TOKEN", List.of("test"), "pane", null, null, mcpUrl, null, null);
}
}
}
@@ -17,10 +17,6 @@ import java.nio.file.Files;
import java.nio.file.Path;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.function.Supplier;
import static org.junit.jupiter.api.Assertions.*;
@@ -38,19 +34,12 @@ class OpenCodeLauncherTest {
}
/** Gate-disabled launcher whose per-spawn config dirs land under an inspectable temp root. */
private static OpenCodeLauncher service(FakeHerdr herdr, Path configRoot, BridgedConfig.Profile cfg) {
private OpenCodeLauncher service(FakeHerdr herdr, Path configRoot, BridgedConfig.Profile cfg) {
return new OpenCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
Map.of(cfg.profile(), cfg), cfg.profile(), k -> "GITEA_ACCESS_TOKEN".equals(k) ? "tok" : null,
0, System::currentTimeMillis, () -> { }, configRoot, configRoot);
}
private static OpenCodeLauncher service(FakeHerdr herdr, Path configRoot, BridgedConfig.Profile cfg,
Supplier<BridgedConfig.Fleet> fleet) {
return new OpenCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null,
0, System::currentTimeMillis, () -> { }, configRoot, configRoot, fleet);
}
@SuppressWarnings("unchecked")
private static Map<String, Object> lastStart(FakeHerdr herdr) {
return (Map<String, Object>) herdr.lastCall("agent.start").params();
@@ -73,10 +62,8 @@ class OpenCodeLauncherTest {
@Test
void writesRemoteMcpConfigAndCharterInstructionsWhenMcpUrlSet(@TempDir Path root) throws Exception {
FakeHerdr herdr = new FakeHerdr();
BridgedConfig.Fleet fleet = new BridgedConfig.Fleet(Map.of(), Map.of(), Map.of(), Map.of(),
Map.of("dev", "role rule"), null);
service(herdr, root, opencodeCfg("google/gemini-2.5-pro", "http://127.0.0.1:8765/mcp", null),
() -> fleet).spawn();
service(herdr, root, opencodeCfg("google/gemini-2.5-pro", "http://127.0.0.1:8765/mcp", null))
.spawn();
Map<String, String> env = startEnv(herdr);
assertNull(env.get("ANTHROPIC_BASE_URL"), "opencode carries no ANTHROPIC_* / subscription boundary");
@@ -95,15 +82,13 @@ class OpenCodeLauncherTest {
"the profile's bridge MCP url is present");
assertTrue(bridge.path("enabled").asBoolean(), "the bridge server is enabled");
assertTrue(json.path("instructions").isArray() && !json.path("instructions").isEmpty(),
"the member charter is mounted via instructions");
"the reply charter is mounted via instructions");
// The instructions entry is a real file path holding the composed member charter.
Path charter = Path.of(cfgPath).resolveSibling("member-charter.md");
// The instructions entry is a real file path holding the reply charter.
Path charter = Path.of(cfgPath).resolveSibling("reply-charter.md");
assertTrue(Files.exists(charter), "the charter file the config references was written");
assertEquals("role rule\n\n" + HerdrPeerLauncher.REPLY_CHARTER, Files.readString(charter),
"the composed charter keeps the role rule first and the reply rule last");
assertEquals(charter.toAbsolutePath().toString(), json.path("instructions").get(0).asText(),
"instructions names the charter file by its absolute path");
assertTrue(Files.readString(charter).contains("bridge_reply"),
"the charter instructs the worker to answer via bridge_reply");
}
@Test
@@ -115,69 +100,6 @@ class OpenCodeLauncherTest {
"no bridge MCP url → no config file and no OPENCODE_CONFIG");
}
@Test
void roleCharterWithoutMcpOrCustomProviderStillWritesAConfig(@TempDir Path root) throws Exception {
FakeHerdr herdr = new FakeHerdr();
BridgedConfig.Fleet fleet = new BridgedConfig.Fleet(Map.of(), Map.of(), Map.of(), Map.of(),
Map.of("dev", "role rule"), null);
Path configRoot = Files.createDirectory(root.resolve("configs"));
Path checkout = Files.createDirectory(root.resolve("checkout"));
service(herdr, configRoot, opencodeCfg("google/gemini-2.5-pro", null, null), () -> fleet).spawn();
String cfgPath = startEnv(herdr).get("OPENCODE_CONFIG");
assertNotNull(cfgPath, "a role charter needs a config even without MCP or custom provider");
JsonNode json = new ObjectMapper().readTree(Path.of(cfgPath).toFile());
Path charter = Path.of(json.path("instructions").get(0).asText());
assertEquals("role rule", Files.readString(charter), "the base-composed role charter is unchanged");
assertTrue(json.path("mcp").isMissingNode(), "a charter does not add an MCP mount");
assertTrue(charter.startsWith(configRoot), "the charter is written under the temp config root");
try (var files = Files.walk(checkout)) {
assertFalse(files.anyMatch(path -> path.getFileName().toString().equals("member-charter.md")),
"the worker checkout receives no charter file");
}
}
@Test
void nullCharterWritesNoCharterFileOrInstructions(@TempDir Path root) throws Exception {
FakeHerdr herdr = new FakeHerdr();
service(herdr, root, pinnedCfg("local-vllm/model", "http://127.0.0.1:8000", null),
() -> new BridgedConfig.Fleet(Map.of(), Map.of(), Map.of(), Map.of(), Map.of(), null)).spawn();
String config = startEnv(herdr).get("OPENCODE_CONFIG");
assertNotNull(config, "the custom provider still needs a config");
JsonNode json = new ObjectMapper().readTree(Path.of(config).toFile());
assertTrue(json.path("instructions").isMissingNode(), "a null charter adds no instructions entry");
try (var files = Files.walk(root)) {
assertFalse(files.anyMatch(path -> path.getFileName().toString().equals("member-charter.md")),
"a null charter creates no charter file");
}
}
@Test
void concurrentSpawnsWriteSeparateCharterDirectories(@TempDir Path root) throws Exception {
BridgedConfig.Profile cfg = opencodeCfg("google/gemini-2.5-pro", "http://127.0.0.1:8765/mcp", null);
ExecutorService executor = Executors.newFixedThreadPool(2);
try {
Future<String> first = executor.submit(() -> spawnConfigPath(root, cfg));
Future<String> second = executor.submit(() -> spawnConfigPath(root, cfg));
Path firstCharter = Path.of(first.get()).resolveSibling("member-charter.md");
Path secondCharter = Path.of(second.get()).resolveSibling("member-charter.md");
assertNotEquals(firstCharter.getParent(), secondCharter.getParent(),
"each concurrent spawn owns a separate config directory");
assertTrue(Files.exists(firstCharter));
assertTrue(Files.exists(secondCharter));
} finally {
executor.shutdownNow();
}
}
private static String spawnConfigPath(Path root, BridgedConfig.Profile cfg) {
FakeHerdr herdr = new FakeHerdr();
service(herdr, root, cfg).spawn();
return startEnv(herdr).get("OPENCODE_CONFIG");
}
@Test
void passesTheModelAsDashMFlagAlongsideAutoApprove(@TempDir Path root) {
FakeHerdr herdr = new FakeHerdr();