CB-111: multi-profile workers — named backends selectable at spawn

The bridge was hard-wired to one worker profile. Config now takes a 'workers'
map keyed by profile name plus 'defaultWorker'; WorkerService holds the map and
gains spawn(profile) (spawn() uses the default). Selection threads through the
surfaces: REST POST /workers ?profile= / {"profile":…} + GET /profiles; MCP
bridge_spawn {profile?} + new bridge_profiles. Each profile's base_url is
guard-checked independently, so gx10 and ollama can run side by side and you
address each worker by its returned sessionId. Backward-compatible: the legacy
singular 'worker:' block still loads as a one-entry profile map.
This commit is contained in:
Dai Ha
2026-07-15 15:54:31 +02:00
parent ab771ea24e
commit 7f1b6b3a0d
10 changed files with 303 additions and 52 deletions
@@ -51,7 +51,8 @@ public final class Bridged {
AgentControl agents = new AgentControl(herdr);
WorkspaceControl spaces = new WorkspaceControl(herdr);
WorkerService workers = new WorkerService(agents, spaces, guard, cfg.worker(), System::getenv);
WorkerService workers = new WorkerService(agents, spaces, guard,
cfg.workerProfiles(), cfg.defaultProfile(), System::getenv);
// Status-gated injector (CB-103): the single writer into workers, fed by a poller.
// The blocking message endpoint (CB-104) is the producer; the poller is inert until then.
@@ -8,7 +8,9 @@ import java.io.IOException;
import java.io.UncheckedIOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
/**
@@ -16,16 +18,21 @@ import java.util.Set;
* {@code bridged.example.yaml}). Unknown keys are ignored so config can grow ahead
* of the code.
*
* @param bind REST/MCP listen host:port
* @param herdrSocket path to herdr's Unix socket ({@code null} → client default)
* @param worker worker-spawn settings
* @param guard subscription-boundary allowlist
* @param bind REST/MCP listen host:port
* @param herdrSocket path to herdr's Unix socket ({@code null} → client default)
* @param worker single worker profile (legacy; superseded by {@code workers})
* @param workers named worker profiles, keyed by profile name (multi-backend fleet)
* @param defaultWorker which {@code workers} key a no-argument spawn uses ({@code null} → the
* single {@code worker}, or the sole/first profile)
* @param guard subscription-boundary allowlist
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record BridgedConfig(
Bind bind,
String herdrSocket,
Worker worker,
Map<String, Worker> workers,
String defaultWorker,
Guard guard) {
@JsonIgnoreProperties(ignoreUnknown = true)
@@ -69,6 +76,11 @@ public record BridgedConfig(
tabLabel = (tabLabel == null || tabLabel.isBlank()) ? "worker: {profile} #{n}" : tabLabel;
}
/** A copy with {@code profile} set — used to default a profile to its {@code workers} key. */
public Worker withProfile(String p) {
return new Worker(p, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel, mcpUrl);
}
/** True when workers should land in their own tab in the worker space. */
public boolean tabPlacement() {
return "tab".equals(placement);
@@ -108,6 +120,40 @@ public record BridgedConfig(
}
}
/**
* The effective worker profiles, keyed by profile name. Prefers the {@code workers} map (each
* value's {@code profile} defaulted to its key); falls back to the legacy singular {@code worker}
* (keyed by its own profile). Empty if neither is configured.
*/
public Map<String, Worker> workerProfiles() {
if (workers != null && !workers.isEmpty()) {
Map<String, Worker> out = new LinkedHashMap<>();
workers.forEach((name, w) -> out.put(name,
(w.profile() == null || w.profile().isBlank()) ? w.withProfile(name) : w));
return Map.copyOf(out);
}
if (worker != null) {
String name = (worker.profile() == null || worker.profile().isBlank()) ? "default" : worker.profile();
return Map.of(name, worker);
}
return Map.of();
}
/**
* The profile a no-argument spawn uses: {@code defaultWorker} if set, else the legacy single
* {@code worker}'s profile, else the sole/first configured profile, else {@code null}.
*/
public String defaultProfile() {
if (defaultWorker != null && !defaultWorker.isBlank()) {
return defaultWorker;
}
if (worker != null && worker.profile() != null && !worker.profile().isBlank()) {
return worker.profile();
}
Map<String, Worker> p = workerProfiles();
return p.isEmpty() ? null : p.keySet().iterator().next();
}
private static final ObjectMapper YAML = new ObjectMapper(new YAMLFactory());
/** Load and validate config from {@code path}. */
@@ -124,6 +170,6 @@ public record BridgedConfig(
public BridgedConfig withDefaults() {
Bind b = bind != null ? bind : new Bind(null, 0);
Guard g = guard != null ? guard : new Guard(List.of());
return new BridgedConfig(b, herdrSocket, worker, g);
return new BridgedConfig(b, herdrSocket, worker, workers, defaultWorker, g);
}
}
@@ -76,9 +76,10 @@ public final class BridgeMcp {
.toolCall(pollTool(), (_, req) ->
poll(messages, str(req.arguments(), "ticket")))
// Fleet management (CB-108): spawn/list/stop over WorkerService.
.toolCall(spawnTool(), (_, _) -> spawn(workers))
.toolCall(spawnTool(), (_, req) -> spawn(workers, str(req.arguments(), "profile")))
.toolCall(listTool(), (_, _) -> listWorkers(workers))
.toolCall(stopTool(), (_, req) -> stop(workers, str(req.arguments(), "paneId")))
.toolCall(profilesTool(), (_, _) -> profiles(workers))
.build();
}
@@ -195,17 +196,30 @@ public final class BridgeMcp {
// --- fleet management logic (CB-108) -------------------------------------------------------
/** {@code bridge_spawn}: launch a guard-checked worker and return its session id + pane id. */
static McpSchema.CallToolResult spawn(WorkerService workers) {
/**
* {@code bridge_spawn}: launch a guard-checked worker for {@code profile} (blank → the default
* profile) and return its session id + pane id.
*/
static McpSchema.CallToolResult spawn(WorkerService workers, String profile) {
try {
return text(json(workerView(workers.spawn())));
Agent worker = isBlank(profile) ? workers.spawn() : workers.spawn(profile);
return text(json(workerView(worker)));
} catch (GuardException e) {
return error("subscription boundary: " + e.getMessage());
} catch (IllegalArgumentException e) {
return error(e.getMessage()); // unknown / no-default profile
} catch (HerdrException e) {
return error("herdr error spawning worker: " + e.getMessage());
}
}
/** {@code bridge_profiles}: the configured worker profiles and the default. */
static McpSchema.CallToolResult profiles(WorkerService workers) {
return text(json(Map.of(
"profiles", workers.profiles(),
"default", workers.defaultProfile() == null ? "" : workers.defaultProfile())));
}
/** {@code bridge_list}: every worker herdr tracks (session id, pane, status). */
static McpSchema.CallToolResult listWorkers(WorkerService workers) {
try {
@@ -273,8 +287,17 @@ public final class BridgeMcp {
private static McpSchema.Tool spawnTool() {
return tool("bridge_spawn",
"Spawn a new off-subscription worker session for the configured profile. Returns its "
+ "sessionId (use with bridge_send) and paneId (use with bridge_stop).",
"Spawn a new off-subscription worker session. Pass a profile (from bridge_profiles) to "
+ "pick the backend, or omit it for the default. Returns the worker's sessionId "
+ "(use with bridge_send) and paneId (use with bridge_stop).",
objectSchema(Map.of(
"profile", stringProp("Worker profile to spawn (omit for the default profile)")),
List.of()));
}
private static McpSchema.Tool profilesTool() {
return tool("bridge_profiles",
"List the configured worker profiles (backends) and which one bridge_spawn uses by default.",
objectSchema(Map.of(), List.of()));
}
@@ -63,7 +63,8 @@ public final class BridgedApp {
app.get("/healthz", this::healthz);
app.get("/sessions", this::sessions);
app.get("/agents", this::agents);
app.post("/workers", this::spawnWorker);
app.get("/profiles", this::profiles); // configured worker profiles
app.post("/workers", this::spawnWorker); // optional ?profile= or {"profile":…}
app.delete("/workers/{paneId}", this::stopWorker);
app.post("/sessions/{id}/message", this::sendMessage); // bridge_send (primary; blocking or wait:false)
app.post("/sessions/{id}/reply", this::replyMessage); // bridge_reply (worker)
@@ -109,13 +110,37 @@ public final class BridgedApp {
ctx.status(200).json(Map.of("agents", workers.list().stream().map(BridgedApp::view).toList()));
}
/** Spawn a guard-checked worker. 403 if the base_url would breach the subscription boundary. */
/** The configured worker profiles and which one a no-argument spawn uses. */
private void profiles(Context ctx) {
ctx.status(200).json(Map.of(
"profiles", workers.profiles(),
"default", workers.defaultProfile() == null ? "" : workers.defaultProfile()));
}
/**
* Spawn a guard-checked worker. An optional {@code profile} (query param or {@code {"profile":…}}
* body) picks which configured profile; omitted → the default. 403 if the base_url would breach
* the subscription boundary, 400 for an unknown profile.
*/
private void spawnWorker(Context ctx) {
String profile = ctx.queryParam("profile");
if (profile == null || profile.isBlank()) {
try {
String body = ctx.body();
if (!body.isBlank()) {
profile = mapper.readTree(body).path("profile").asText(null);
}
} catch (Exception ignored) {
// A malformed/empty body just means "no profile" → fall through to the default.
}
}
try {
Agent worker = workers.spawn();
Agent worker = (profile == null || profile.isBlank()) ? workers.spawn() : workers.spawn(profile);
ctx.status(201).json(view(worker));
} catch (GuardException e) {
ctx.status(403).json(Map.of("error", "subscription_boundary", "detail", e.getMessage()));
} catch (IllegalArgumentException e) {
ctx.status(400).json(Map.of("error", "unknown_profile", "detail", e.getMessage()));
}
}
@@ -16,6 +16,7 @@ import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.Function;
@@ -44,7 +45,8 @@ public final class WorkerService {
private final AgentControl agents;
private final WorkspaceControl spaces;
private final SubscriptionGuard guard;
private final BridgedConfig.Worker cfg;
private final Map<String, BridgedConfig.Worker> profiles; // profile name → spawn settings
private final String defaultProfile; // profile a no-arg spawn uses (nullable)
private final Function<String, String> env; // host env lookup (injectable for tests)
private final AtomicLong nameSeq = new AtomicLong(); // per-worker counter (also the tab #)
@@ -64,16 +66,42 @@ public final class WorkerService {
private final String nameNonce = String.format("%06x", new SecureRandom().nextInt(1 << 24));
public WorkerService(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
BridgedConfig.Worker cfg, Function<String, String> env) {
Map<String, BridgedConfig.Worker> profiles, String defaultProfile,
Function<String, String> env) {
this.agents = agents;
this.spaces = spaces;
this.guard = guard;
this.cfg = cfg;
this.profiles = Map.copyOf(profiles);
this.defaultProfile = defaultProfile;
this.env = env;
}
/** Spawn a worker for the configured profile. Guard runs before any herdr call. */
/** The configured worker profile names (what {@code spawn(profile)} accepts). */
public Set<String> profiles() {
return profiles.keySet();
}
/** The profile a no-argument {@link #spawn()} uses, or {@code null} if none is configured. */
public String defaultProfile() {
return defaultProfile;
}
/** Spawn a worker for the default profile. Guard runs before any herdr call. */
public Agent spawn() {
if (defaultProfile == null || defaultProfile.isBlank()) {
throw new IllegalArgumentException("no default worker profile is configured — "
+ "pass a profile; configured: " + profiles.keySet());
}
return spawn(defaultProfile);
}
/** Spawn a worker for a named profile. Guard runs before any herdr call. */
public Agent spawn(String profileName) {
BridgedConfig.Worker cfg = profiles.get(profileName);
if (cfg == null) {
throw new IllegalArgumentException("unknown worker profile '" + profileName
+ "' — configured: " + profiles.keySet());
}
String baseUrl = cfg.baseUrl();
guard.assertWorker(baseUrl); // hard stop before we spawn anything
@@ -86,9 +114,9 @@ public final class WorkerService {
// Mount the bridge MCP + reply charter as launch FLAGS (non-invasive: nothing written to
// the worker's profile/config dir). Identity is connection-based, so the mount is shared.
List<String> argv = argvWithBridge();
List<String> argv = argvWithBridge(cfg);
return cfg.tabPlacement() ? spawnInTab(workerEnv, argv) : spawnAsPane(workerEnv, argv);
return cfg.tabPlacement() ? spawnInTab(cfg, workerEnv, argv) : spawnAsPane(cfg, workerEnv, argv);
}
/**
@@ -96,7 +124,7 @@ public final class WorkerService {
* 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.
*/
private List<String> argvWithBridge() {
private List<String> argvWithBridge(BridgedConfig.Worker cfg) {
if (!cfg.hasMcp()) {
return cfg.argv();
}
@@ -111,7 +139,7 @@ public final class WorkerService {
}
/** Dedicated worker space → own tab → drop the placeholder shell so only the worker remains. */
private Agent spawnInTab(Map<String, String> workerEnv, List<String> argv) {
private Agent spawnInTab(BridgedConfig.Worker cfg, Map<String, String> workerEnv, List<String> argv) {
Workspace space = spaces.ensureWorkspace(cfg.workspace());
Tab.Created tab = spaces.createTab(space.workspaceId());
log.info("spawning worker profile={} base_url={} space={} tab={}",
@@ -119,7 +147,7 @@ public final class WorkerService {
Started started;
try {
started = startUniquelyNamed(workerEnv, argv, tab.tab().tabId());
started = startUniquelyNamed(cfg, workerEnv, argv, tab.tab().tabId());
} catch (RuntimeException e) {
// The worker never started — don't leave the tab we just created orphaned.
// Best-effort cleanup; never let it mask the real spawn failure.
@@ -159,10 +187,10 @@ public final class WorkerService {
}
/** Legacy placement: herdr splits the currently-focused tab. */
private Agent spawnAsPane(Map<String, String> workerEnv, List<String> argv) {
private Agent spawnAsPane(BridgedConfig.Worker cfg, Map<String, String> workerEnv, List<String> argv) {
log.info("spawning worker (pane placement) profile={} base_url={} argv={}",
cfg.profile(), cfg.baseUrl(), argv);
Agent worker = startUniquelyNamed(workerEnv, argv, null).agent();
Agent worker = startUniquelyNamed(cfg, workerEnv, argv, null).agent();
log.info("worker started pane={} terminal={}", worker.paneId(), worker.terminalId());
return worker;
}
@@ -181,7 +209,8 @@ public final class WorkerService {
* retry is a belt-and-braces backstop for the astronomically unlikely nonce+seq clash;
* the name is a label only — herdr detects kind and status from terminal output, not it.
*/
private Started startUniquelyNamed(Map<String, String> workerEnv, List<String> argv, String tabId) {
private Started startUniquelyNamed(BridgedConfig.Worker cfg, Map<String, String> workerEnv,
List<String> argv, String tabId) {
HerdrException last = null;
for (int attempt = 0; attempt < NAME_RETRIES; attempt++) {
long seq = nameSeq.incrementAndGet();
@@ -214,7 +243,10 @@ public final class WorkerService {
* a genuinely failed teardown is not reported as done.
*/
public void stop(String paneId) {
WorkspaceControl.PaneLocation loc = cfg.tabPlacement() ? spaces.locatePane(paneId) : null;
// Teardown knows only the paneId, not which profile spawned it. Attempt tab cleanup when any
// profile uses tab placement (so the bridge may have created a dedicated worker tab); the
// single-occupant check below is what actually protects the user's shared tabs.
WorkspaceControl.PaneLocation loc = usesTabPlacement() ? spaces.locatePane(paneId) : null;
try {
agents.close(paneId);
} catch (HerdrException e) {
@@ -229,6 +261,11 @@ public final class WorkerService {
}
}
/** Whether any configured profile places workers in their own tab (so tabs may need cleanup). */
private boolean usesTabPlacement() {
return profiles.values().stream().anyMatch(BridgedConfig.Worker::tabPlacement);
}
/** True when a herdr error means the target is already gone (safe to treat as done). */
private static boolean isAlreadyGone(HerdrException e) {
return e.code() != null && e.code().endsWith("_not_found");