CB-402 Increment 4: CompositePeerLauncher — route the fleet by kind
Introduce the router the core holds when more than one adapter is configured: one HerdrPeerLauncher per peer kind, dispatched by profile (spawn/effectiveCwd/parityOverlay), by pane id (stop, via a spawn-time owner map), and fanned out + combined for the fleet-wide queries (list dedup by pane id, reap/caps union, profiles union). The ctor rejects an empty adapter list and a profile two adapters both claim. Wire it in Bridged.main: partition workerProfiles() by kind (claude-code is the always-present default adapter; opencode is added when any profile opts in) and front both with the composite. This lets BridgeMcp and BridgedApp finally take the PeerLauncher SPI instead of a concrete ClaudeCodeLauncher — the two (ClaudeCodeLauncher) casts in Bridged are gone. list() elements are cast to herdr Agent at the point of the herdr-specific roster view, where that assumption actually lives. Add BridgedConfig.Worker.isClaudeCode()/isOpenCode() kind predicates (the wiring uses isOpenCode; both are unit-tested). Drop the long-dead 'rendezvous' constructor param threaded into BridgeMcp and BridgedApp. 10 CompositePeerLauncherTest cases over two real adapters on one FakeHerdr: profile routing (observed via the started agent's claude-/opencode- name prefix), default resolution, unknown-profile and duplicate-profile rejection, caps union, list dedup, reap sum, and stop teardown. 266 tests green.
This commit is contained in:
@@ -28,11 +28,18 @@ import dev.ltms.bridged.session.SessionManager;
|
|||||||
import dev.ltms.bridged.peer.PeerLauncher;
|
import dev.ltms.bridged.peer.PeerLauncher;
|
||||||
import dev.ltms.bridged.session.SessionReaper;
|
import dev.ltms.bridged.session.SessionReaper;
|
||||||
import dev.ltms.bridged.worker.ClaudeCodeLauncher;
|
import dev.ltms.bridged.worker.ClaudeCodeLauncher;
|
||||||
|
import dev.ltms.bridged.worker.CompositePeerLauncher;
|
||||||
|
import dev.ltms.bridged.worker.HerdrPeerLauncher;
|
||||||
|
import dev.ltms.bridged.worker.OpenCodeLauncher;
|
||||||
import io.javalin.Javalin;
|
import io.javalin.Javalin;
|
||||||
import org.slf4j.Logger;
|
import org.slf4j.Logger;
|
||||||
import org.slf4j.LoggerFactory;
|
import org.slf4j.LoggerFactory;
|
||||||
|
|
||||||
import java.nio.file.Path;
|
import java.nio.file.Path;
|
||||||
|
import java.util.ArrayList;
|
||||||
|
import java.util.LinkedHashMap;
|
||||||
|
import java.util.List;
|
||||||
|
import java.util.Map;
|
||||||
import java.util.concurrent.Executors;
|
import java.util.concurrent.Executors;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -63,9 +70,33 @@ public final class Bridged {
|
|||||||
|
|
||||||
AgentControl agents = new AgentControl(herdr);
|
AgentControl agents = new AgentControl(herdr);
|
||||||
WorkspaceControl spaces = new WorkspaceControl(herdr);
|
WorkspaceControl spaces = new WorkspaceControl(herdr);
|
||||||
PeerLauncher workers = new ClaudeCodeLauncher(agents, spaces, guard,
|
// CB-402: one adapter per configured peer kind, fronted by a composite router. A profile's
|
||||||
cfg.workerProfiles(), cfg.defaultProfile(), System::getenv,
|
// `kind:` selects its adapter — claude-code (the default) and opencode partition the profile
|
||||||
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs());
|
// set — and the composite dispatches each SPI call to the adapter that owns the profile/pane.
|
||||||
|
Map<String, BridgedConfig.Worker> claudeProfiles = new LinkedHashMap<>();
|
||||||
|
Map<String, BridgedConfig.Worker> opencodeProfiles = new LinkedHashMap<>();
|
||||||
|
cfg.workerProfiles().forEach((name, w) -> {
|
||||||
|
if (w.isOpenCode()) {
|
||||||
|
opencodeProfiles.put(name, w);
|
||||||
|
} else {
|
||||||
|
claudeProfiles.put(name, w);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
List<HerdrPeerLauncher> adapters = new ArrayList<>();
|
||||||
|
// The claude-code adapter is the always-present default; keep it even with no profiles (so a
|
||||||
|
// bridge configured with no workers, or opencode-only, still has a well-defined base adapter)
|
||||||
|
// unless opencode is the only kind configured.
|
||||||
|
if (!claudeProfiles.isEmpty() || opencodeProfiles.isEmpty()) {
|
||||||
|
adapters.add(new ClaudeCodeLauncher(agents, spaces, guard,
|
||||||
|
claudeProfiles, cfg.defaultProfile(), System::getenv,
|
||||||
|
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs()));
|
||||||
|
}
|
||||||
|
if (!opencodeProfiles.isEmpty()) {
|
||||||
|
adapters.add(new OpenCodeLauncher(agents, spaces,
|
||||||
|
opencodeProfiles, cfg.defaultProfile(), System::getenv,
|
||||||
|
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs()));
|
||||||
|
}
|
||||||
|
PeerLauncher workers = new CompositePeerLauncher(adapters, cfg.defaultProfile());
|
||||||
// CB-117: herdr keeps worker panes alive across a daemon restart, and their ids died with
|
// CB-117: herdr keeps worker panes alive across a daemon restart, and their ids died with
|
||||||
// the previous process — reap those leaked orphans now, before we start serving.
|
// the previous process — reap those leaked orphans now, before we start serving.
|
||||||
workers.reapOrphanWorkers();
|
workers.reapOrphanWorkers();
|
||||||
@@ -151,8 +182,7 @@ public final class Bridged {
|
|||||||
// Caller identity is resolved from the connection (peer PID → herdr pane), not arguments.
|
// Caller identity is resolved from the connection (peer PID → herdr pane), not arguments.
|
||||||
ConnectionIdentity identity = new ConnectionIdentity(
|
ConnectionIdentity identity = new ConnectionIdentity(
|
||||||
new PaneLocator(herdr), new LsofPeerPidLookup(), new LsofProcessCwdLookup());
|
new PaneLocator(herdr), new LsofPeerPidLookup(), new LsofProcessCwdLookup());
|
||||||
// Cast to ClaudeCodeLauncher: BridgeMcp is not yet migrated to PeerLauncher (Stage A scope).
|
BridgeMcp mcp = new BridgeMcp(messages, workers, sessions, identity, presence,
|
||||||
BridgeMcp mcp = new BridgeMcp(messages, rendezvous, (ClaudeCodeLauncher) workers, sessions, identity, presence,
|
|
||||||
primaryRegistry);
|
primaryRegistry);
|
||||||
|
|
||||||
// CB-303 part 3: single ordered shutdown hook. Drain sessions first while herdr is still
|
// CB-303 part 3: single ordered shutdown hook. Drain sessions first while herdr is still
|
||||||
@@ -176,7 +206,7 @@ public final class Bridged {
|
|||||||
herdr.close();
|
herdr.close();
|
||||||
}));
|
}));
|
||||||
|
|
||||||
Javalin app = new BridgedApp(herdr, (ClaudeCodeLauncher) workers, sessions, messages, rendezvous, presence, mcp.servlet()).build();
|
Javalin app = new BridgedApp(herdr, workers, sessions, messages, presence, mcp.servlet()).build();
|
||||||
app.start(cfg.bind().host(), cfg.bind().port());
|
app.start(cfg.bind().host(), cfg.bind().port());
|
||||||
log.info("bridged listening on {}:{}, herdr socket {}",
|
log.info("bridged listening on {}:{}, herdr socket {}",
|
||||||
cfg.bind().host(), cfg.bind().port(), socket);
|
cfg.bind().host(), cfg.bind().port(), socket);
|
||||||
|
|||||||
@@ -163,6 +163,16 @@ public record BridgedConfig(
|
|||||||
mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind);
|
mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** True when this profile is served by the Claude Code adapter (the default kind). */
|
||||||
|
public boolean isClaudeCode() {
|
||||||
|
return KIND_CLAUDE_CODE.equals(kind);
|
||||||
|
}
|
||||||
|
|
||||||
|
/** True when this profile is served by the opencode adapter (CB-402). */
|
||||||
|
public boolean isOpenCode() {
|
||||||
|
return KIND_OPENCODE.equals(kind);
|
||||||
|
}
|
||||||
|
|
||||||
/** True when this profile's workers are granted a forge token to open their own PR (CB-302). */
|
/** True when this profile's workers are granted a forge token to open their own PR (CB-302). */
|
||||||
public boolean hasGitToken() {
|
public boolean hasGitToken() {
|
||||||
return gitTokenEnv != null && !gitTokenEnv.isBlank();
|
return gitTokenEnv != null && !gitTokenEnv.isBlank();
|
||||||
|
|||||||
@@ -10,7 +10,7 @@ import dev.ltms.bridged.peer.PeerUnreachableException;
|
|||||||
import dev.ltms.bridged.session.SessionManager;
|
import dev.ltms.bridged.session.SessionManager;
|
||||||
import dev.ltms.bridged.session.WorkerSession;
|
import dev.ltms.bridged.session.WorkerSession;
|
||||||
import dev.ltms.bridged.session.WorktreeRequest;
|
import dev.ltms.bridged.session.WorktreeRequest;
|
||||||
import dev.ltms.bridged.worker.ClaudeCodeLauncher;
|
import dev.ltms.bridged.peer.PeerLauncher;
|
||||||
import io.modelcontextprotocol.common.McpTransportContext;
|
import io.modelcontextprotocol.common.McpTransportContext;
|
||||||
import io.modelcontextprotocol.json.McpJsonMapper;
|
import io.modelcontextprotocol.json.McpJsonMapper;
|
||||||
import io.modelcontextprotocol.json.jackson3.JacksonMcpJsonMapperSupplier;
|
import io.modelcontextprotocol.json.jackson3.JacksonMcpJsonMapperSupplier;
|
||||||
@@ -35,8 +35,8 @@ import java.util.stream.Collectors;
|
|||||||
* / {@code bridge_status}; the worker calls {@code bridge_reply}.
|
* / {@code bridge_status}; the worker calls {@code bridge_reply}.
|
||||||
*
|
*
|
||||||
* <p>Beyond delegation the primary also manages the fleet here (CB-108): {@code bridge_spawn} /
|
* <p>Beyond delegation the primary also manages the fleet here (CB-108): {@code bridge_spawn} /
|
||||||
* {@code bridge_list} / {@code bridge_stop} adapt {@link ClaudeCodeLauncher} so a worker's whole lifecycle
|
* {@code bridge_list} / {@code bridge_stop} drive the {@link PeerLauncher} SPI so a worker's whole
|
||||||
* is driven through MCP, with the subscription boundary still enforced inside {@code ClaudeCodeLauncher}.
|
* lifecycle is managed through MCP, with each adapter's subscription boundary enforced inside it.
|
||||||
*
|
*
|
||||||
* <p>The tool <em>logic</em> lives in package-private static methods returning a
|
* <p>The tool <em>logic</em> lives in package-private static methods returning a
|
||||||
* {@link McpSchema.CallToolResult}, so it is unit-testable without standing up the HTTP transport;
|
* {@link McpSchema.CallToolResult}, so it is unit-testable without standing up the HTTP transport;
|
||||||
@@ -60,7 +60,7 @@ public final class BridgeMcp {
|
|||||||
private final HttpServletStreamableServerTransportProvider transport;
|
private final HttpServletStreamableServerTransportProvider transport;
|
||||||
private final McpSyncServer server;
|
private final McpSyncServer server;
|
||||||
|
|
||||||
public BridgeMcp(MessageService messages, Rendezvous rendezvous, ClaudeCodeLauncher workers,
|
public BridgeMcp(MessageService messages, PeerLauncher workers,
|
||||||
SessionManager sessions, ConnectionIdentity identity, WorkerPresence presence,
|
SessionManager sessions, ConnectionIdentity identity, WorkerPresence presence,
|
||||||
PrimaryRegistry primaryRegistry) {
|
PrimaryRegistry primaryRegistry) {
|
||||||
McpJsonMapper json = new JacksonMcpJsonMapperSupplier().get();
|
McpJsonMapper json = new JacksonMcpJsonMapperSupplier().get();
|
||||||
@@ -373,16 +373,17 @@ public final class BridgeMcp {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/** {@code bridge_profiles}: the configured worker profiles and the default. */
|
/** {@code bridge_profiles}: the configured worker profiles and the default. */
|
||||||
static McpSchema.CallToolResult profiles(ClaudeCodeLauncher workers) {
|
static McpSchema.CallToolResult profiles(PeerLauncher workers) {
|
||||||
return text(json(Map.of(
|
return text(json(Map.of(
|
||||||
"profiles", workers.profiles(),
|
"profiles", workers.profiles(),
|
||||||
"default", workers.defaultProfile() == null ? "" : workers.defaultProfile())));
|
"default", workers.defaultProfile() == null ? "" : workers.defaultProfile())));
|
||||||
}
|
}
|
||||||
|
|
||||||
/** {@code bridge_list}: bridge-owned roster merged with live herdr status by paneId. */
|
/** {@code bridge_list}: bridge-owned roster merged with live herdr status by paneId. */
|
||||||
static McpSchema.CallToolResult listWorkers(ClaudeCodeLauncher workers, SessionManager sessions) {
|
static McpSchema.CallToolResult listWorkers(PeerLauncher workers, SessionManager sessions) {
|
||||||
try {
|
try {
|
||||||
Map<String, Agent> live = workers.list().stream()
|
Map<String, Agent> live = workers.list().stream()
|
||||||
|
.map(Agent.class::cast)
|
||||||
.filter(a -> a.paneId() != null)
|
.filter(a -> a.paneId() != null)
|
||||||
.collect(Collectors.toMap(Agent::paneId, Function.identity(), (_, b) -> b));
|
.collect(Collectors.toMap(Agent::paneId, Function.identity(), (_, b) -> b));
|
||||||
List<Map<String, Object>> out = sessions.roster().stream()
|
List<Map<String, Object>> out = sessions.roster().stream()
|
||||||
|
|||||||
@@ -9,11 +9,10 @@ import dev.ltms.bridged.herdr.HerdrException;
|
|||||||
import dev.ltms.bridged.inject.WorkerPresence;
|
import dev.ltms.bridged.inject.WorkerPresence;
|
||||||
import dev.ltms.bridged.peer.PeerUnreachableException;
|
import dev.ltms.bridged.peer.PeerUnreachableException;
|
||||||
import dev.ltms.bridged.msg.MessageService;
|
import dev.ltms.bridged.msg.MessageService;
|
||||||
import dev.ltms.bridged.msg.Rendezvous;
|
|
||||||
import dev.ltms.bridged.session.SessionManager;
|
import dev.ltms.bridged.session.SessionManager;
|
||||||
import dev.ltms.bridged.session.WorkerSession;
|
import dev.ltms.bridged.session.WorkerSession;
|
||||||
import dev.ltms.bridged.session.WorktreeRequest;
|
import dev.ltms.bridged.session.WorktreeRequest;
|
||||||
import dev.ltms.bridged.worker.ClaudeCodeLauncher;
|
import dev.ltms.bridged.peer.PeerLauncher;
|
||||||
import io.javalin.Javalin;
|
import io.javalin.Javalin;
|
||||||
import io.javalin.http.Context;
|
import io.javalin.http.Context;
|
||||||
import jakarta.servlet.http.HttpServlet;
|
import jakarta.servlet.http.HttpServlet;
|
||||||
@@ -45,22 +44,20 @@ public final class BridgedApp {
|
|||||||
private static final long MAX_ASK_TIMEOUT_MS = 115_000;
|
private static final long MAX_ASK_TIMEOUT_MS = 115_000;
|
||||||
|
|
||||||
private final HerdrClient herdr;
|
private final HerdrClient herdr;
|
||||||
private final ClaudeCodeLauncher workers;
|
private final PeerLauncher workers;
|
||||||
private final SessionManager sessions; // CB-301: authoritative session registry
|
private final SessionManager sessions; // CB-301: authoritative session registry
|
||||||
private final MessageService messages;
|
private final MessageService messages;
|
||||||
private final Rendezvous rendezvous;
|
|
||||||
private final WorkerPresence presence; // CB-113: which workers are MCP-connected (available)
|
private final WorkerPresence presence; // CB-113: which workers are MCP-connected (available)
|
||||||
private final HttpServlet mcpServlet; // MCP Streamable-HTTP endpoint, mounted at /mcp (nullable)
|
private final HttpServlet mcpServlet; // MCP Streamable-HTTP endpoint, mounted at /mcp (nullable)
|
||||||
private final ObjectMapper mapper = new ObjectMapper();
|
private final ObjectMapper mapper = new ObjectMapper();
|
||||||
|
|
||||||
public BridgedApp(HerdrClient herdr, ClaudeCodeLauncher workers, SessionManager sessions,
|
public BridgedApp(HerdrClient herdr, PeerLauncher workers, SessionManager sessions,
|
||||||
MessageService messages, Rendezvous rendezvous, WorkerPresence presence,
|
MessageService messages, WorkerPresence presence,
|
||||||
HttpServlet mcpServlet) {
|
HttpServlet mcpServlet) {
|
||||||
this.herdr = herdr;
|
this.herdr = herdr;
|
||||||
this.workers = workers;
|
this.workers = workers;
|
||||||
this.sessions = sessions;
|
this.sessions = sessions;
|
||||||
this.messages = messages;
|
this.messages = messages;
|
||||||
this.rendezvous = rendezvous;
|
|
||||||
this.presence = presence;
|
this.presence = presence;
|
||||||
this.mcpServlet = mcpServlet;
|
this.mcpServlet = mcpServlet;
|
||||||
}
|
}
|
||||||
@@ -125,12 +122,14 @@ public final class BridgedApp {
|
|||||||
|
|
||||||
/** Discovery: every agent herdr tracks, keyed by its Claude session UUID. */
|
/** Discovery: every agent herdr tracks, keyed by its Claude session UUID. */
|
||||||
private void agents(Context ctx) {
|
private void agents(Context ctx) {
|
||||||
ctx.status(200).json(Map.of("agents", workers.list().stream().map(BridgedApp::view).toList()));
|
ctx.status(200).json(Map.of("agents",
|
||||||
|
workers.list().stream().map(Agent.class::cast).map(BridgedApp::view).toList()));
|
||||||
}
|
}
|
||||||
|
|
||||||
/** CB-304: bridge-owned roster merged with live herdr status by paneId. */
|
/** CB-304: bridge-owned roster merged with live herdr status by paneId. */
|
||||||
private void listWorkers(Context ctx) {
|
private void listWorkers(Context ctx) {
|
||||||
Map<String, Agent> live = workers.list().stream()
|
Map<String, Agent> live = workers.list().stream()
|
||||||
|
.map(Agent.class::cast)
|
||||||
.filter(a -> a.paneId() != null)
|
.filter(a -> a.paneId() != null)
|
||||||
.collect(Collectors.toMap(Agent::paneId, Function.identity(), (_, b) -> b));
|
.collect(Collectors.toMap(Agent::paneId, Function.identity(), (_, b) -> b));
|
||||||
List<Map<String, Object>> out = sessions.roster().stream()
|
List<Map<String, Object>> out = sessions.roster().stream()
|
||||||
|
|||||||
@@ -0,0 +1,160 @@
|
|||||||
|
package dev.ltms.bridged.worker;
|
||||||
|
|
||||||
|
import dev.ltms.bridged.herdr.Agent;
|
||||||
|
import dev.ltms.bridged.peer.Capability;
|
||||||
|
import dev.ltms.bridged.peer.PeerHandle;
|
||||||
|
import dev.ltms.bridged.peer.PeerLauncher;
|
||||||
|
import dev.ltms.bridged.peer.SpawnRequest;
|
||||||
|
import org.slf4j.Logger;
|
||||||
|
import org.slf4j.LoggerFactory;
|
||||||
|
|
||||||
|
import java.util.EnumSet;
|
||||||
|
import java.util.LinkedHashMap;
|
||||||
|
import java.util.List;
|
||||||
|
import java.util.Map;
|
||||||
|
import java.util.Set;
|
||||||
|
import java.util.concurrent.ConcurrentHashMap;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* The {@link PeerLauncher} the core actually holds when more than one adapter is configured — a thin
|
||||||
|
* router in front of one {@link HerdrPeerLauncher} per peer {@code kind} (Claude Code, opencode, …).
|
||||||
|
* It owns no transport of its own; it dispatches each SPI call to the delegate that owns the profile
|
||||||
|
* involved, and fans the fleet-wide queries (list/reap/caps/profiles) across all delegates.
|
||||||
|
*
|
||||||
|
* <p>Routing rules:
|
||||||
|
* <ul>
|
||||||
|
* <li><strong>By profile</strong> — {@link #spawn}, {@link #effectiveCwd}, {@link #parityOverlay}
|
||||||
|
* resolve the profile (a null/blank name → the global {@link #defaultProfile}) and delegate to
|
||||||
|
* the single adapter that declares it. Profiles partition cleanly across adapters: the
|
||||||
|
* constructor rejects a name claimed by two.</li>
|
||||||
|
* <li><strong>By pane id</strong> — {@link #stop} routes to the adapter that spawned that pane
|
||||||
|
* (recorded at spawn time). A pane the composite never spawned (only real for a caller that
|
||||||
|
* hand-rolls an id) falls back to the first delegate; teardown is pane-id addressed and
|
||||||
|
* tab cleanup is single-occupant guarded, so it is safe either way.</li>
|
||||||
|
* <li><strong>Fleet-wide</strong> — {@link #reapOrphanWorkers} and {@link #capabilities} fan out
|
||||||
|
* and combine. {@link #list} is deduplicated by pane id because every herdr-backed delegate
|
||||||
|
* shares one herdr connection and so reports the same global agent set.</li>
|
||||||
|
* </ul>
|
||||||
|
*/
|
||||||
|
public final class CompositePeerLauncher implements PeerLauncher {
|
||||||
|
|
||||||
|
private static final Logger log = LoggerFactory.getLogger(CompositePeerLauncher.class);
|
||||||
|
|
||||||
|
private final List<HerdrPeerLauncher> delegates;
|
||||||
|
private final Map<String, HerdrPeerLauncher> byProfile;
|
||||||
|
private final String defaultProfile;
|
||||||
|
|
||||||
|
/** paneId → the delegate that spawned it, so {@link #stop} tears down through the right adapter. */
|
||||||
|
private final Map<String, HerdrPeerLauncher> spawnedBy = new ConcurrentHashMap<>();
|
||||||
|
|
||||||
|
/**
|
||||||
|
* @param delegates one adapter per configured peer kind; must be non-empty and declare
|
||||||
|
* disjoint profile-name sets
|
||||||
|
* @param defaultProfile the profile a no-argument spawn resolves to (may be null)
|
||||||
|
* @throws IllegalArgumentException if {@code delegates} is empty or two adapters claim one profile
|
||||||
|
*/
|
||||||
|
public CompositePeerLauncher(List<HerdrPeerLauncher> delegates, String defaultProfile) {
|
||||||
|
if (delegates.isEmpty()) {
|
||||||
|
throw new IllegalArgumentException("at least one peer adapter must be configured");
|
||||||
|
}
|
||||||
|
this.delegates = List.copyOf(delegates);
|
||||||
|
this.defaultProfile = defaultProfile;
|
||||||
|
Map<String, HerdrPeerLauncher> index = new LinkedHashMap<>();
|
||||||
|
for (HerdrPeerLauncher d : this.delegates) {
|
||||||
|
for (String profile : d.profiles()) {
|
||||||
|
HerdrPeerLauncher prev = index.putIfAbsent(profile, d);
|
||||||
|
if (prev != null) {
|
||||||
|
throw new IllegalArgumentException(
|
||||||
|
"worker profile '" + profile + "' is claimed by two peer adapters");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
this.byProfile = Map.copyOf(index);
|
||||||
|
}
|
||||||
|
|
||||||
|
/** The adapter owning {@code profileName} (null/blank → the default). Throws on an unknown profile. */
|
||||||
|
private HerdrPeerLauncher route(String profileName) {
|
||||||
|
String resolved = (profileName == null || profileName.isBlank()) ? defaultProfile : profileName;
|
||||||
|
if (resolved == null) {
|
||||||
|
// No profile and no default configured — hand to the first delegate so it raises the
|
||||||
|
// same "no default" error it would on its own; keeps the SPI contract single-sourced.
|
||||||
|
return delegates.getFirst();
|
||||||
|
}
|
||||||
|
HerdrPeerLauncher d = byProfile.get(resolved);
|
||||||
|
if (d == null) {
|
||||||
|
throw new IllegalArgumentException("unknown worker profile: " + resolved);
|
||||||
|
}
|
||||||
|
return d;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public PeerHandle spawn(SpawnRequest req) {
|
||||||
|
HerdrPeerLauncher d = route(req.profileName());
|
||||||
|
PeerHandle handle = d.spawn(req);
|
||||||
|
spawnedBy.put(handle.id(), d);
|
||||||
|
return handle;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public String effectiveCwd(SpawnRequest req) {
|
||||||
|
return route(req.profileName()).effectiveCwd(req);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public List<String> parityOverlay(String profileName) {
|
||||||
|
return route(profileName).parityOverlay(profileName);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public void stop(String id) {
|
||||||
|
HerdrPeerLauncher d = spawnedBy.remove(id);
|
||||||
|
if (d == null) {
|
||||||
|
log.debug("stop({}) — no recorded owner, routing to the first adapter (pane-addressed)", id);
|
||||||
|
d = delegates.getFirst();
|
||||||
|
}
|
||||||
|
d.stop(id);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public Set<String> profiles() {
|
||||||
|
return byProfile.keySet();
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public String defaultProfile() {
|
||||||
|
return defaultProfile;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Every herdr agent, deduplicated by pane id (all delegates share one herdr and list globally). */
|
||||||
|
@Override
|
||||||
|
public List<Agent> list() {
|
||||||
|
Map<String, Agent> byPane = new LinkedHashMap<>();
|
||||||
|
for (HerdrPeerLauncher d : delegates) {
|
||||||
|
for (Agent a : d.list()) {
|
||||||
|
if (a.paneId() != null) {
|
||||||
|
byPane.putIfAbsent(a.paneId(), a);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return List.copyOf(byPane.values());
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public int reapOrphanWorkers() {
|
||||||
|
int reaped = 0;
|
||||||
|
for (HerdrPeerLauncher d : delegates) {
|
||||||
|
reaped += d.reapOrphanWorkers();
|
||||||
|
}
|
||||||
|
return reaped;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** The union of every adapter's capabilities — a capability any adapter offers, the fleet offers. */
|
||||||
|
@Override
|
||||||
|
public Set<Capability> capabilities() {
|
||||||
|
EnumSet<Capability> caps = EnumSet.noneOf(Capability.class);
|
||||||
|
for (HerdrPeerLauncher d : delegates) {
|
||||||
|
caps.addAll(d.capabilities());
|
||||||
|
}
|
||||||
|
return Set.copyOf(caps);
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -194,6 +194,27 @@ class BridgedConfigTest {
|
|||||||
"kind is normalised to lower-case so YAML casing does not matter");
|
"kind is normalised to lower-case so YAML casing does not matter");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void kindPredicatesReflectTheResolvedKind(@TempDir Path dir) throws Exception {
|
||||||
|
Path f = dir.resolve("kind-predicates.yaml");
|
||||||
|
Files.writeString(f, """
|
||||||
|
workers:
|
||||||
|
claude:
|
||||||
|
baseUrl: http://gx10.gw:8000
|
||||||
|
gemini:
|
||||||
|
kind: opencode
|
||||||
|
model: google/gemini-2.5-pro
|
||||||
|
""");
|
||||||
|
|
||||||
|
BridgedConfig cfg = BridgedConfig.load(f);
|
||||||
|
BridgedConfig.Worker claude = cfg.workerProfiles().get("claude");
|
||||||
|
BridgedConfig.Worker gemini = cfg.workerProfiles().get("gemini");
|
||||||
|
assertTrue(claude.isClaudeCode(), "the default-kind worker is claude-code");
|
||||||
|
assertFalse(claude.isOpenCode(), "a claude-code worker is not opencode");
|
||||||
|
assertTrue(gemini.isOpenCode(), "the kind: opencode worker is opencode");
|
||||||
|
assertFalse(gemini.isClaudeCode(), "an opencode worker is not claude-code");
|
||||||
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
void argvDefaultsToTheKindBinaryWhenUnset(@TempDir Path dir) throws Exception {
|
void argvDefaultsToTheKindBinaryWhenUnset(@TempDir Path dir) throws Exception {
|
||||||
Path f = dir.resolve("kind-argv.yaml");
|
Path f = dir.resolve("kind-argv.yaml");
|
||||||
|
|||||||
@@ -75,7 +75,7 @@ class BridgedAppTest {
|
|||||||
poller.start();
|
poller.start();
|
||||||
Rendezvous rendezvous = new Rendezvous();
|
Rendezvous rendezvous = new Rendezvous();
|
||||||
MessageService messages = new MessageService(agents, injector, rendezvous);
|
MessageService messages = new MessageService(agents, injector, rendezvous);
|
||||||
app = new BridgedApp(herdr, workers, sessions, messages, rendezvous, this.presence, null)
|
app = new BridgedApp(herdr, workers, sessions, messages, this.presence, null)
|
||||||
.build().start("127.0.0.1", 0);
|
.build().start("127.0.0.1", 0);
|
||||||
return app.port();
|
return app.port();
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,158 @@
|
|||||||
|
package dev.ltms.bridged.worker;
|
||||||
|
|
||||||
|
import dev.ltms.bridged.config.BridgedConfig;
|
||||||
|
import dev.ltms.bridged.guard.SubscriptionGuard;
|
||||||
|
import dev.ltms.bridged.herdr.AgentControl;
|
||||||
|
import dev.ltms.bridged.herdr.FakeHerdr;
|
||||||
|
import dev.ltms.bridged.herdr.WorkspaceControl;
|
||||||
|
import dev.ltms.bridged.peer.PeerHandle;
|
||||||
|
import dev.ltms.bridged.peer.PeerLauncher;
|
||||||
|
import dev.ltms.bridged.peer.SpawnRequest;
|
||||||
|
import org.junit.jupiter.api.Test;
|
||||||
|
|
||||||
|
import java.util.List;
|
||||||
|
import java.util.Map;
|
||||||
|
import java.util.Set;
|
||||||
|
|
||||||
|
import static org.junit.jupiter.api.Assertions.*;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* The composite router: profile → owning adapter for spawn/cwd/parity, pane id → owner for stop,
|
||||||
|
* and fleet-wide union/dedup for list/reap/caps/profiles. Exercised through two real adapters —
|
||||||
|
* claude-code + opencode — over one FakeHerdr, so each call is observed reaching the right adapter
|
||||||
|
* (the started herdr agent name carries that adapter's {@code claude-}/{@code opencode-} prefix).
|
||||||
|
*/
|
||||||
|
class CompositePeerLauncherTest {
|
||||||
|
|
||||||
|
private ClaudeCodeLauncher claudeAdapter(FakeHerdr herdr) {
|
||||||
|
// 12-arg back-compat Worker ctor → kind defaults to claude-code.
|
||||||
|
BridgedConfig.Worker claude = new BridgedConfig.Worker("claude", "http://gx00.gw:8000", "coder",
|
||||||
|
null, "BRIDGED_WORKER_TOKEN", List.of("claude"), "tab", "bridged-workers", "w #{n}",
|
||||||
|
null, null, null);
|
||||||
|
return new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||||
|
new SubscriptionGuard(Set.of("gx00.gw")), Map.of("claude", claude), "claude", _ -> null);
|
||||||
|
}
|
||||||
|
|
||||||
|
private OpenCodeLauncher opencodeAdapter(FakeHerdr herdr) {
|
||||||
|
BridgedConfig.Worker gemini = new BridgedConfig.Worker("gemini", null, "google/gemini-2.5-pro",
|
||||||
|
null, "BRIDGED_WORKER_TOKEN", List.of("opencode"), "tab", "bridged-workers", "w #{n}",
|
||||||
|
null, null, null, "GITEA_ACCESS_TOKEN", null, BridgedConfig.Worker.KIND_OPENCODE);
|
||||||
|
return new OpenCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||||
|
Map.of("gemini", gemini), "gemini", _ -> "tok");
|
||||||
|
}
|
||||||
|
|
||||||
|
private CompositePeerLauncher composite(FakeHerdr herdr) {
|
||||||
|
return new CompositePeerLauncher(
|
||||||
|
List.of(claudeAdapter(herdr), opencodeAdapter(herdr)), "claude");
|
||||||
|
}
|
||||||
|
|
||||||
|
@SuppressWarnings("unchecked")
|
||||||
|
private static String startedName(FakeHerdr herdr) {
|
||||||
|
return (String) ((Map<String, Object>) herdr.lastCall("agent.start").params()).get("name");
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void spawnRoutesEachProfileToItsOwningAdapter() {
|
||||||
|
FakeHerdr herdr = new FakeHerdr();
|
||||||
|
PeerLauncher composite = composite(herdr);
|
||||||
|
|
||||||
|
composite.spawn(new SpawnRequest("gemini", null, null));
|
||||||
|
assertTrue(startedName(herdr).startsWith("opencode-"),
|
||||||
|
"the gemini profile is spawned by the opencode adapter: " + startedName(herdr));
|
||||||
|
|
||||||
|
composite.spawn(new SpawnRequest("claude", null, null));
|
||||||
|
assertTrue(startedName(herdr).startsWith("claude-"),
|
||||||
|
"the claude profile is spawned by the claude-code adapter: " + startedName(herdr));
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void nullProfileResolvesTheDefaultAndRoutesToItsOwner() {
|
||||||
|
FakeHerdr herdr = new FakeHerdr();
|
||||||
|
composite(herdr).spawn(new SpawnRequest(null, null, null));
|
||||||
|
assertTrue(startedName(herdr).startsWith("claude-"),
|
||||||
|
"a no-profile spawn resolves the default (claude) and routes to its adapter");
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void unknownProfileIsRejected() {
|
||||||
|
FakeHerdr herdr = new FakeHerdr();
|
||||||
|
PeerLauncher composite = composite(herdr);
|
||||||
|
assertThrows(IllegalArgumentException.class,
|
||||||
|
() -> composite.spawn(new SpawnRequest("nope", null, null)),
|
||||||
|
"a profile no adapter declares is an error");
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void profilesAndDefaultAreExposedAcrossAdapters() {
|
||||||
|
FakeHerdr herdr = new FakeHerdr();
|
||||||
|
PeerLauncher composite = composite(herdr);
|
||||||
|
assertEquals(Set.of("claude", "gemini"), composite.profiles(),
|
||||||
|
"profiles are the union of every adapter's profiles");
|
||||||
|
assertEquals("claude", composite.defaultProfile());
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void capabilitiesAreTheUnionOfEveryAdapter() {
|
||||||
|
FakeHerdr herdr = new FakeHerdr();
|
||||||
|
ClaudeCodeLauncher claude = claudeAdapter(herdr);
|
||||||
|
OpenCodeLauncher opencode = opencodeAdapter(herdr);
|
||||||
|
PeerLauncher composite = new CompositePeerLauncher(List.of(claude, opencode), "claude");
|
||||||
|
|
||||||
|
assertTrue(composite.capabilities().containsAll(claude.capabilities()),
|
||||||
|
"the fleet offers every claude-code capability");
|
||||||
|
assertTrue(composite.capabilities().containsAll(opencode.capabilities()),
|
||||||
|
"the fleet offers every opencode capability (incl. SELF_PR from its git-token profile)");
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void listIsDeduplicatedByPaneIdAcrossAdaptersSharingHerdr() {
|
||||||
|
FakeHerdr herdr = new FakeHerdr();
|
||||||
|
PeerLauncher composite = composite(herdr);
|
||||||
|
// Both adapters wrap the same herdr, so each list() returns the same global agent set;
|
||||||
|
// the composite must return each pane once, not once per adapter.
|
||||||
|
assertEquals(1, composite.list().size(),
|
||||||
|
"the single herdr-tracked pane appears once, not duplicated per adapter");
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void reapSumsAcrossAdaptersAndEachAdapterReapsOnlyItsOwnPrefix() {
|
||||||
|
// One foreign opencode orphan + one foreign claude orphan, from a prior daemon (different nonce).
|
||||||
|
FakeHerdr herdr = new FakeHerdr()
|
||||||
|
.withAgent("opencode-gemini-ffffff-1", "term_o", "wQ:pO", "wQ:tO")
|
||||||
|
.withAgent("claude-claude-eeeeee-1", "term_c", "wQ:pC", "wQ:tC");
|
||||||
|
PeerLauncher composite = composite(herdr);
|
||||||
|
assertEquals(2, composite.reapOrphanWorkers(),
|
||||||
|
"both orphans are reaped — one by each adapter, summed by the composite");
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void stopTearsDownAPaneSpawnedThroughTheComposite() {
|
||||||
|
FakeHerdr herdr = new FakeHerdr();
|
||||||
|
PeerLauncher composite = composite(herdr);
|
||||||
|
PeerHandle handle = composite.spawn(new SpawnRequest("gemini", null, null));
|
||||||
|
|
||||||
|
composite.stop(handle.id());
|
||||||
|
assertTrue(herdr.calls.stream()
|
||||||
|
.anyMatch(c -> c.method().equals("pane.close")
|
||||||
|
&& handle.id().equals(((Map<?, ?>) c.params()).get("pane_id"))),
|
||||||
|
"stop routes to the spawning adapter and closes that worker's pane");
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void constructorRejectsAProfileClaimedByTwoAdapters() {
|
||||||
|
FakeHerdr herdr = new FakeHerdr();
|
||||||
|
// Two opencode adapters both declaring "gemini" — a profile-name collision.
|
||||||
|
OpenCodeLauncher a = opencodeAdapter(herdr);
|
||||||
|
OpenCodeLauncher b = opencodeAdapter(herdr);
|
||||||
|
assertThrows(IllegalArgumentException.class,
|
||||||
|
() -> new CompositePeerLauncher(List.of(a, b), "gemini"),
|
||||||
|
"a profile two adapters both claim is a configuration error");
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void constructorRejectsAnEmptyAdapterList() {
|
||||||
|
assertThrows(IllegalArgumentException.class,
|
||||||
|
() -> new CompositePeerLauncher(List.of(), "claude"),
|
||||||
|
"at least one adapter must be configured");
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user