CB-518: weighted placement policy for worker spawns
This commit is contained in:
@@ -100,6 +100,8 @@ workers:
|
||||
mcpUrl: http://127.0.0.1:8765/mcp
|
||||
tokenEnv: BRIDGED_WORKER_TOKEN
|
||||
argv: ["ccs", "gx10"]
|
||||
weight: 0.5 # relative selection weight for placement: weighted
|
||||
maxLoad: 2 # max live workers on this profile (omit for unlimited)
|
||||
# gitTokenEnv: GITEA_TOKEN # opt-in: let this profile's workers open their own PR (CB-302)
|
||||
# gitHostEnv: GITEA_HOST # defaults to GITEA_HOST; injected only with gitTokenEnv
|
||||
# configDir: /Users/me/.ccs/instances/gx10 # CLAUDE_CONFIG_DIR — inherit that profile's skills/MCP
|
||||
@@ -112,6 +114,8 @@ workers:
|
||||
tabLabel: "worker: {profile} #{n}"
|
||||
mcpUrl: http://127.0.0.1:8765/mcp
|
||||
argv: ["ccs", "ollama"]
|
||||
weight: 0.5
|
||||
maxLoad: 2
|
||||
# CB-402: a second coding-agent kind, proving the PeerLauncher SPI is provider-neutral.
|
||||
# opencode is provider-agnostic and uses NONE of Claude's private seams: no ANTHROPIC_BASE_URL /
|
||||
# SubscriptionGuard (so it needs no `guard` host entry), no --mcp-config / --append-system-prompt.
|
||||
@@ -155,6 +159,9 @@ workers:
|
||||
# tabLabel: "opencode: {profile} #{n}"
|
||||
# mcpUrl: http://127.0.0.1:8765/mcp
|
||||
# argv: ["opencode"]
|
||||
# How an unqualified spawn chooses a profile: fixed (default, reproduces pre-CB-518 behaviour),
|
||||
# round-robin, or weighted. Omitting this key is a strict no-op for existing configs.
|
||||
placement: weighted
|
||||
defaultWorker: gx10
|
||||
|
||||
# Subscription boundary. A worker's base_url host MUST be one of these; the primary
|
||||
|
||||
@@ -33,6 +33,7 @@ import dev.ltms.bridged.session.SessionManager;
|
||||
import dev.ltms.bridged.peer.PeerLauncher;
|
||||
import dev.ltms.bridged.session.SessionReaper;
|
||||
import dev.ltms.bridged.worker.ClaudeCodeLauncher;
|
||||
import dev.ltms.bridged.placement.PlacementPolicies;
|
||||
import dev.ltms.bridged.worker.CompositePeerLauncher;
|
||||
import dev.ltms.bridged.worker.HerdrPeerLauncher;
|
||||
import dev.ltms.bridged.worker.OpenCodeLauncher;
|
||||
@@ -46,6 +47,8 @@ import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
import java.util.function.Function;
|
||||
|
||||
/**
|
||||
* {@code bridged} entry point. Wires the real herdr socket client to the REST app and
|
||||
@@ -111,7 +114,13 @@ public final class Bridged {
|
||||
opencodeProfiles, cfg.defaultProfile(), System::getenv,
|
||||
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs()));
|
||||
}
|
||||
PeerLauncher workers = new CompositePeerLauncher(adapters, cfg.defaultProfile());
|
||||
AtomicReference<Function<String, Integer>> liveCountRef = new AtomicReference<>(name -> 0);
|
||||
PeerLauncher workers = new CompositePeerLauncher(
|
||||
adapters,
|
||||
cfg.defaultProfile(),
|
||||
cfg.workerProfiles(),
|
||||
PlacementPolicies.fromName(cfg.placement()),
|
||||
profileName -> liveCountRef.get().apply(profileName));
|
||||
// CB-504: under supervision (launchd/systemd) bridged can start before herdr's socket
|
||||
// exists. The client itself is lazy — it connects per call — but the orphan reap below is
|
||||
// the first thing that actually talks to herdr, so without this wait a boot-order race
|
||||
@@ -136,6 +145,9 @@ public final class Bridged {
|
||||
contextCap = cfg.lifecycle().contextCap();
|
||||
}
|
||||
SessionManager sessions = new SessionManager(workers, new GitWorktrees(cfg.worktreeRoot()), contextCap);
|
||||
liveCountRef.set(profileName -> (int) sessions.roster().stream()
|
||||
.filter(s -> profileName.equals(s.profile()))
|
||||
.count());
|
||||
|
||||
// CB-303 part 1: idle-ttl reaper — only when configured, defaults to disabled.
|
||||
final SessionReaper reaper;
|
||||
|
||||
@@ -36,6 +36,8 @@ import java.util.Set;
|
||||
* @param primary optional pinned primary terminal config ({@code null} → derived from connection);
|
||||
* a non-blank {@code terminal} seeds {@code PrimaryRegistry} and prevents
|
||||
* connection-derived overrides, CB-307
|
||||
* @param placement how to choose a worker profile for an unqualified spawn:
|
||||
* {@code fixed} (default), {@code round-robin}, or {@code weighted}
|
||||
* @param auth API authentication mode ({@code null} → {@code loopback-trust}, the
|
||||
* historical behaviour), CB-501
|
||||
*/
|
||||
@@ -53,6 +55,7 @@ public record BridgedConfig(
|
||||
Integer spawnReadyPollMs,
|
||||
Broker broker,
|
||||
Primary primary,
|
||||
String placement,
|
||||
Auth auth) {
|
||||
|
||||
@JsonIgnoreProperties(ignoreUnknown = true)
|
||||
@@ -93,6 +96,12 @@ public record BridgedConfig(
|
||||
* (minimal-grant default — push over SSH stays free, PR-create is opt-in)
|
||||
* @param gitHostEnv name of the host env var holding the forge host (default {@code GITEA_HOST});
|
||||
* injected as {@code GITEA_HOST} <em>only</em> when {@code gitTokenEnv} is set
|
||||
* @param weight relative selection weight for {@code placement: weighted}. Absent or
|
||||
* non-positive ⇒ 1.0. Weights are normalised by the policy, so they need
|
||||
* not sum to 1.0.
|
||||
* @param maxLoad max live workers allowed on this profile at one time; absent or
|
||||
* non-positive ⇒ unlimited. Live means any session the registry still owns
|
||||
* (acquired and not yet released), in any state.
|
||||
* @param kind which peer launcher spawns this profile: {@code "claude-code"} (default —
|
||||
* the {@link dev.ltms.bridged.worker.ClaudeCodeLauncher}) or {@code "opencode"}.
|
||||
* The {@code CompositePeerLauncher} routes {@code spawn}/reap by this value, so
|
||||
@@ -108,7 +117,9 @@ public record BridgedConfig(
|
||||
List<String> parityOverlay,
|
||||
String gitTokenEnv, String gitHostEnv,
|
||||
String kind,
|
||||
Map<String, String> env) {
|
||||
Map<String, String> env,
|
||||
Float weight,
|
||||
Integer maxLoad) {
|
||||
|
||||
/** Peer kind spawned by {@link dev.ltms.bridged.worker.ClaudeCodeLauncher} (the default). */
|
||||
public static final String KIND_CLAUDE_CODE = "claude-code";
|
||||
@@ -135,6 +146,8 @@ public record BridgedConfig(
|
||||
// checkpoints need only set gitTokenEnv; it is injected only alongside a resolved token.
|
||||
gitHostEnv = (gitHostEnv == null || gitHostEnv.isBlank()) ? "GITEA_HOST" : gitHostEnv;
|
||||
env = (env == null) ? Map.of() : Map.copyOf(env);
|
||||
weight = (weight == null || weight <= 0.0f) ? 1.0f : weight;
|
||||
maxLoad = (maxLoad == null || maxLoad <= 0) ? null : maxLoad;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -147,7 +160,7 @@ public record BridgedConfig(
|
||||
String placement, String workspace, String tabLabel, String mcpUrl,
|
||||
String cwd, List<String> parityOverlay) {
|
||||
this(profile, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel,
|
||||
mcpUrl, cwd, parityOverlay, null, null, null);
|
||||
mcpUrl, cwd, parityOverlay, null, null, null, null, null, null);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -159,7 +172,7 @@ public record BridgedConfig(
|
||||
String placement, String workspace, String tabLabel, String mcpUrl,
|
||||
String cwd, List<String> parityOverlay, String gitTokenEnv, String gitHostEnv) {
|
||||
this(profile, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel,
|
||||
mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, null, null);
|
||||
mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, null, null, null, null);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -172,13 +185,13 @@ public record BridgedConfig(
|
||||
String cwd, List<String> parityOverlay, String gitTokenEnv, String gitHostEnv,
|
||||
String kind) {
|
||||
this(profile, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel,
|
||||
mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, null);
|
||||
mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, null, null, null);
|
||||
}
|
||||
|
||||
/** 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, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, env);
|
||||
mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, env, weight, maxLoad);
|
||||
}
|
||||
|
||||
/** True when this profile is served by the Claude Code adapter (the default kind). */
|
||||
@@ -385,9 +398,10 @@ public record BridgedConfig(
|
||||
Integer timeout = (spawnReadyTimeoutMs != null) ? spawnReadyTimeoutMs : 20000;
|
||||
Integer pollMs = (spawnReadyPollMs != null) ? spawnReadyPollMs : 300;
|
||||
Auth a = auth != null ? auth : new Auth(null, null);
|
||||
String placementOrDefault = (placement != null && !placement.isBlank()) ? placement : "fixed";
|
||||
// broker is left as-is: null (or an empty/blank uri) keeps the in-memory soft-state inbox.
|
||||
// primary is left as-is: null defaults to connection-derived identity.
|
||||
return new BridgedConfig(b, herdrSocket, worker, workers, defaultWorker, g, worktreeRoot, l, timeout, pollMs, broker, primary, a);
|
||||
return new BridgedConfig(b, herdrSocket, worker, workers, defaultWorker, g, worktreeRoot, l, timeout, pollMs, broker, primary, placementOrDefault, a);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -30,4 +30,15 @@ public interface PeerHandle {
|
||||
default String terminalId() {
|
||||
return null;
|
||||
}
|
||||
|
||||
/**
|
||||
* The worker profile that spawned this peer, if the launcher resolved one. A launcher that
|
||||
* performs dynamic profile selection (e.g. CB-518 weighted placement) sets this so the
|
||||
* session registry records the actual profile rather than the requested/default one.
|
||||
*
|
||||
* @return the profile name, or {@code null} when the launcher leaves it unspecified
|
||||
*/
|
||||
default String profile() {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,22 @@
|
||||
package dev.ltms.bridged.placement;
|
||||
|
||||
/**
|
||||
* Backward-compatible placement: an unqualified spawn always resolves to the configured default
|
||||
* profile, exactly as {@code CompositePeerLauncher} did before CB-518. This ignores caps and
|
||||
* reachability so that a pre-existing config behaves identically after upgrade.
|
||||
*/
|
||||
final class FixedPlacementPolicy implements PlacementPolicy {
|
||||
|
||||
@Override
|
||||
public PlacementCandidate select(PlacementContext ctx) {
|
||||
String d = ctx.defaultProfile();
|
||||
if (d != null && !d.isBlank()) {
|
||||
return new PlacementCandidate(d, null, 1.0f, null);
|
||||
}
|
||||
if (!ctx.candidates().isEmpty()) {
|
||||
PlacementCandidate first = ctx.candidates().getFirst();
|
||||
return new PlacementCandidate(first.profile(), null, first.weight(), first.maxLoad());
|
||||
}
|
||||
throw new PlacementException("no worker profiles configured");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,20 @@
|
||||
package dev.ltms.bridged.placement;
|
||||
|
||||
/**
|
||||
* A profile (and, in CB-308, a host) that can be chosen by a {@link PlacementPolicy}.
|
||||
*
|
||||
* <p>Keeping this as a small descriptor rather than a bare profile name lets CB-308 widen
|
||||
* selection to {@code (host, profile)} pairs without changing the policy interface.
|
||||
*/
|
||||
public record PlacementCandidate(String profile, String host, float weight, Integer maxLoad) {
|
||||
|
||||
/** A candidate with no explicit host (the single-host default) and the given weight/cap. */
|
||||
public static PlacementCandidate profile(String profile, float weight, Integer maxLoad) {
|
||||
return new PlacementCandidate(profile, null, weight, maxLoad);
|
||||
}
|
||||
|
||||
/** A candidate with no explicit host, unit weight, and no cap. */
|
||||
public static PlacementCandidate profile(String profile) {
|
||||
return new PlacementCandidate(profile, null, 1.0f, null);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,19 @@
|
||||
package dev.ltms.bridged.placement;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Set;
|
||||
import java.util.function.Function;
|
||||
|
||||
/**
|
||||
* Everything a {@link PlacementPolicy} needs to make one selection.
|
||||
*
|
||||
* @param defaultProfile profile a {@code fixed} policy should return (may be {@code null})
|
||||
* @param candidates every configured candidate; the policy filters out those at cap or unreachable
|
||||
* @param liveCount current live worker count per profile (from the session registry)
|
||||
* @param unreachable profiles already known to have failed in this spawn attempt
|
||||
*/
|
||||
public record PlacementContext(String defaultProfile,
|
||||
List<PlacementCandidate> candidates,
|
||||
Function<String, Integer> liveCount,
|
||||
Set<String> unreachable) {
|
||||
}
|
||||
@@ -0,0 +1,12 @@
|
||||
package dev.ltms.bridged.placement;
|
||||
|
||||
/**
|
||||
* Thrown when a {@link PlacementPolicy} has no candidate available. Kept as a distinct type so
|
||||
* callers can distinguish "no capacity" from a spawn-time transport failure.
|
||||
*/
|
||||
public final class PlacementException extends IllegalStateException {
|
||||
|
||||
public PlacementException(String message) {
|
||||
super(message);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,41 @@
|
||||
package dev.ltms.bridged.placement;
|
||||
|
||||
/**
|
||||
* Factory for the built-in placement policies.
|
||||
*/
|
||||
public final class PlacementPolicies {
|
||||
|
||||
private PlacementPolicies() {
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve a policy name from config. Absent/blank values and {@code "fixed"} return the
|
||||
* backward-compatible fixed policy; unknown names throw.
|
||||
*/
|
||||
public static PlacementPolicy fromName(String name) {
|
||||
String n = (name == null) ? "" : name.toLowerCase();
|
||||
if (n.isBlank() || "fixed".equals(n)) {
|
||||
return fixed();
|
||||
}
|
||||
if ("weighted".equals(n)) {
|
||||
return weighted();
|
||||
}
|
||||
if ("round-robin".equals(n)) {
|
||||
return roundRobin();
|
||||
}
|
||||
throw new IllegalArgumentException("unknown placement policy '" + name
|
||||
+ "' — must be one of: fixed, round-robin, weighted");
|
||||
}
|
||||
|
||||
public static PlacementPolicy fixed() {
|
||||
return new FixedPlacementPolicy();
|
||||
}
|
||||
|
||||
public static PlacementPolicy weighted() {
|
||||
return new WeightedRoundRobinPolicy();
|
||||
}
|
||||
|
||||
public static PlacementPolicy roundRobin() {
|
||||
return new RoundRobinPlacementPolicy();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,16 @@
|
||||
package dev.ltms.bridged.placement;
|
||||
|
||||
/**
|
||||
* How {@code bridged} chooses a worker profile when a spawn names none. Implementations are
|
||||
* deterministic and unit-testable; the caller (the composite launcher) handles failover retries.
|
||||
*/
|
||||
public interface PlacementPolicy {
|
||||
|
||||
/**
|
||||
* Pick one candidate from the configured set.
|
||||
*
|
||||
* @throws java.lang.IllegalStateException when no candidate is available, with a message naming
|
||||
* whether every profile is at capacity or unreachable
|
||||
*/
|
||||
PlacementCandidate select(PlacementContext ctx);
|
||||
}
|
||||
@@ -0,0 +1,65 @@
|
||||
package dev.ltms.bridged.placement;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
|
||||
/**
|
||||
* Shared filtering and empty-set reporting used by the built-in placement policies.
|
||||
*/
|
||||
final class PlacementPolicyUtil {
|
||||
|
||||
private PlacementPolicyUtil() {
|
||||
}
|
||||
|
||||
/**
|
||||
* Candidates that are not known-unreachable and have not reached their maxLoad.
|
||||
* A {@code null} maxLoad means unlimited.
|
||||
*/
|
||||
static List<PlacementCandidate> available(PlacementContext ctx) {
|
||||
List<PlacementCandidate> out = new ArrayList<>();
|
||||
for (PlacementCandidate c : ctx.candidates()) {
|
||||
if (ctx.unreachable().contains(c.profile())) {
|
||||
continue;
|
||||
}
|
||||
Integer cap = c.maxLoad();
|
||||
if (cap != null) {
|
||||
int live = ctx.liveCount().apply(c.profile());
|
||||
if (live >= cap) {
|
||||
continue;
|
||||
}
|
||||
}
|
||||
out.add(c);
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
/**
|
||||
* Build a clear exception describing why every candidate was dropped: all at capacity,
|
||||
* all unreachable, or a mix.
|
||||
*/
|
||||
static PlacementException emptyException(PlacementContext ctx) {
|
||||
int atCap = 0;
|
||||
int unreachable = 0;
|
||||
for (PlacementCandidate c : ctx.candidates()) {
|
||||
Integer cap = c.maxLoad();
|
||||
if (ctx.unreachable().contains(c.profile())) {
|
||||
unreachable++;
|
||||
} else if (cap != null && ctx.liveCount().apply(c.profile()) >= cap) {
|
||||
atCap++;
|
||||
}
|
||||
}
|
||||
|
||||
int total = ctx.candidates().size();
|
||||
if (total == 0) {
|
||||
return new PlacementException("no worker profiles configured");
|
||||
}
|
||||
if (atCap == total) {
|
||||
return new PlacementException("all worker profiles are at maxLoad");
|
||||
}
|
||||
if (unreachable == total) {
|
||||
return new PlacementException("all worker profiles are unreachable");
|
||||
}
|
||||
return new PlacementException("no worker profile available: " + atCap + " at maxLoad, "
|
||||
+ unreachable + " unreachable, " + (total - atCap - unreachable) + " remaining");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,23 @@
|
||||
package dev.ltms.bridged.placement;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
|
||||
/**
|
||||
* Deterministic round-robin over the profiles that still have capacity and are not known to be
|
||||
* unreachable. The index advances only on successful selections so the distribution stays even
|
||||
* across spawns.
|
||||
*/
|
||||
final class RoundRobinPlacementPolicy implements PlacementPolicy {
|
||||
|
||||
private final AtomicInteger index = new AtomicInteger(0);
|
||||
|
||||
@Override
|
||||
public synchronized PlacementCandidate select(PlacementContext ctx) {
|
||||
List<PlacementCandidate> available = PlacementPolicyUtil.available(ctx);
|
||||
if (available.isEmpty()) {
|
||||
throw PlacementPolicyUtil.emptyException(ctx);
|
||||
}
|
||||
return available.get(index.getAndIncrement() % available.size());
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,50 @@
|
||||
package dev.ltms.bridged.placement;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
|
||||
/**
|
||||
* Smooth weighted round-robin (nginx-style): for each selection, add the candidate's weight to
|
||||
* its current score, pick the highest score, then subtract the total weight of all available
|
||||
* candidates from the winner. Weights are not required to sum to 1.0; only their ratios matter.
|
||||
*
|
||||
* <p>The state is per-policy instance and protected by {@code synchronized} so concurrent spawns
|
||||
* see a consistent, deterministic sequence rather than interleaving updates.
|
||||
*/
|
||||
final class WeightedRoundRobinPolicy implements PlacementPolicy {
|
||||
|
||||
private final Map<String, Double> current = new ConcurrentHashMap<>();
|
||||
|
||||
@Override
|
||||
public synchronized PlacementCandidate select(PlacementContext ctx) {
|
||||
List<PlacementCandidate> available = PlacementPolicyUtil.available(ctx);
|
||||
if (available.isEmpty()) {
|
||||
throw PlacementPolicyUtil.emptyException(ctx);
|
||||
}
|
||||
|
||||
double total = 0.0;
|
||||
for (PlacementCandidate c : available) {
|
||||
total += c.weight();
|
||||
}
|
||||
if (total <= 0.0) {
|
||||
throw new PlacementException("all available profiles have non-positive weight");
|
||||
}
|
||||
|
||||
PlacementCandidate best = null;
|
||||
double bestScore = Double.NEGATIVE_INFINITY;
|
||||
for (PlacementCandidate c : available) {
|
||||
double score = current.merge(c.profile(), (double) c.weight(), (old, add) -> old + add);
|
||||
if (score > bestScore) {
|
||||
bestScore = score;
|
||||
best = c;
|
||||
}
|
||||
}
|
||||
if (best == null) {
|
||||
throw new PlacementException("no placement candidate could be selected");
|
||||
}
|
||||
|
||||
current.put(best.profile(), current.get(best.profile()) - total);
|
||||
return best;
|
||||
}
|
||||
}
|
||||
@@ -112,9 +112,8 @@ public final class SessionManager implements TurnListener {
|
||||
if (wt == null) {
|
||||
SpawnRequest req = new SpawnRequest(profile, requestedCwd, callerCwd);
|
||||
PeerHandle handle = launcher.spawn(req);
|
||||
String resolvedProfile = (profile == null || profile.isBlank())
|
||||
? launcher.defaultProfile() : profile;
|
||||
String cwd = launcher.effectiveCwd(req);
|
||||
String resolvedProfile = resolveProfile(handle, profile);
|
||||
String cwd = launcher.effectiveCwd(new SpawnRequest(resolvedProfile, requestedCwd, callerCwd));
|
||||
long now = nowNanos.getAsLong();
|
||||
WorkerSession session = new WorkerSession(
|
||||
handle.id(),
|
||||
@@ -210,7 +209,7 @@ public final class SessionManager implements TurnListener {
|
||||
|
||||
private WorkerSession acquireWithWorktree(String profile, String requestedCwd, String callerCwd,
|
||||
String ownerTerminal, WorktreeRequest wt) {
|
||||
String resolvedProfile = (profile == null || profile.isBlank())
|
||||
String preResolvedProfile = (profile == null || profile.isBlank())
|
||||
? launcher.defaultProfile() : profile;
|
||||
// CB-507: resolve through the launcher's CB-112 chain (requested → profile cwd → caller →
|
||||
// daemon cwd → "."), never the raw args. A plain REST spawn supplies neither a requested
|
||||
@@ -218,13 +217,13 @@ public final class SessionManager implements TurnListener {
|
||||
// `git -C null` on the command line — an NPE out of ProcessBuilder, surfacing as HTTP 500.
|
||||
// The non-worktree path always used this chain; only this branch was missed.
|
||||
String repoRoot = worktrees.repoRoot(
|
||||
launcher.effectiveCwd(new SpawnRequest(resolvedProfile, requestedCwd, callerCwd)));
|
||||
launcher.effectiveCwd(new SpawnRequest(preResolvedProfile, requestedCwd, callerCwd)));
|
||||
String branch = "worker/" + slug(wt.ticketSlug()) + "-" + nonce();
|
||||
String path = null;
|
||||
PeerHandle handle;
|
||||
try {
|
||||
path = worktrees.add(repoRoot, branch, wt.baseRef());
|
||||
worktrees.overlayParity(repoRoot, path, launcher.parityOverlay(resolvedProfile));
|
||||
worktrees.overlayParity(repoRoot, path, launcher.parityOverlay(preResolvedProfile));
|
||||
handle = launcher.spawn(new SpawnRequest(profile, path, callerCwd));
|
||||
} catch (RuntimeException e) {
|
||||
if (path != null) {
|
||||
@@ -236,12 +235,14 @@ public final class SessionManager implements TurnListener {
|
||||
}
|
||||
throw e;
|
||||
}
|
||||
String resolvedProfile = resolveProfile(handle, profile);
|
||||
String cwd = launcher.effectiveCwd(new SpawnRequest(resolvedProfile, path, callerCwd));
|
||||
long now = nowNanos.getAsLong();
|
||||
WorkerSession session = new WorkerSession(
|
||||
handle.id(),
|
||||
handle.terminalId(),
|
||||
resolvedProfile,
|
||||
resolveCwd(path, profile, callerCwd),
|
||||
cwd,
|
||||
ownerTerminal,
|
||||
now,
|
||||
now,
|
||||
@@ -264,6 +265,22 @@ public final class SessionManager implements TurnListener {
|
||||
return String.format("%06x", nonceRandom.nextInt(1 << 24)) + "-" + nonceSeq.incrementAndGet();
|
||||
}
|
||||
|
||||
/**
|
||||
* The profile to record for a session. A launcher that performed dynamic selection tells us
|
||||
* the actual profile via {@link PeerHandle#profile()}; otherwise fall back to what the caller
|
||||
* requested (or the launcher's default for a no-profile spawn).
|
||||
*/
|
||||
private String resolveProfile(PeerHandle handle, String requestedProfile) {
|
||||
String fromHandle = handle.profile();
|
||||
if (fromHandle != null && !fromHandle.isBlank()) {
|
||||
return fromHandle;
|
||||
}
|
||||
if (requestedProfile != null && !requestedProfile.isBlank()) {
|
||||
return requestedProfile;
|
||||
}
|
||||
return launcher.defaultProfile();
|
||||
}
|
||||
|
||||
/**
|
||||
* Release the old session and acquire a fresh one with the same profile and working directory.
|
||||
* The new session is guaranteed to have a pane id distinct from the old one (no-reuse invariant).
|
||||
@@ -467,10 +484,6 @@ public final class SessionManager implements TurnListener {
|
||||
return registry.replace(expected.paneId(), expected, updated);
|
||||
}
|
||||
|
||||
private String resolveCwd(String requestedCwd, String profileName, String callerCwd) {
|
||||
return launcher.effectiveCwd(new SpawnRequest(profileName, requestedCwd, callerCwd));
|
||||
}
|
||||
|
||||
/** WorkerPresence bridge that also drives the manager's READY transition. */
|
||||
private static final class PresenceBridge extends WorkerPresence {
|
||||
private final SessionManager sessions;
|
||||
|
||||
@@ -1,19 +1,28 @@
|
||||
package dev.ltms.bridged.worker;
|
||||
|
||||
import dev.ltms.bridged.config.BridgedConfig;
|
||||
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.PeerUnreachableException;
|
||||
import dev.ltms.bridged.peer.SpawnRequest;
|
||||
import dev.ltms.bridged.placement.PlacementCandidate;
|
||||
import dev.ltms.bridged.placement.PlacementContext;
|
||||
import dev.ltms.bridged.placement.PlacementPolicies;
|
||||
import dev.ltms.bridged.placement.PlacementPolicy;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.EnumSet;
|
||||
import java.util.HashSet;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.function.Function;
|
||||
|
||||
/**
|
||||
* The {@link PeerLauncher} the core actually holds when more than one adapter is configured — a thin
|
||||
@@ -35,6 +44,12 @@ import java.util.concurrent.ConcurrentHashMap;
|
||||
* 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>
|
||||
*
|
||||
* <p>CB-518: an unqualified spawn is routed through a {@link PlacementPolicy}. The default
|
||||
* {@code fixed} policy reproduces the historical default-profile behaviour; {@code weighted} uses
|
||||
* smooth weighted round-robin with {@code maxLoad} gating. If a chosen profile fails with
|
||||
* {@link PeerUnreachableException}, the composite advances to the next available candidate and
|
||||
* retries, bounded by the number of candidates.
|
||||
*/
|
||||
public final class CompositePeerLauncher implements PeerLauncher {
|
||||
|
||||
@@ -47,18 +62,46 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
/** paneId → the delegate that spawned it, so {@link #stop} tears down through the right adapter. */
|
||||
private final Map<String, HerdrPeerLauncher> spawnedBy = new ConcurrentHashMap<>();
|
||||
|
||||
private final Map<String, BridgedConfig.Worker> profileConfigs;
|
||||
private final PlacementPolicy placementPolicy;
|
||||
private final Function<String, Integer> liveCount;
|
||||
|
||||
/**
|
||||
* Backward-compatible constructor: fixed placement, no live-counting. Use this for tests and
|
||||
* simple wiring; it preserves the pre-CB-518 behaviour exactly.
|
||||
*
|
||||
* @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) {
|
||||
this(delegates, defaultProfile, Map.of(), PlacementPolicies.fixed(), name -> 0);
|
||||
}
|
||||
|
||||
/**
|
||||
* Production constructor with a placement policy and live-worker counter.
|
||||
*
|
||||
* @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 under {@code fixed} policy
|
||||
* @param profileConfigs all configured worker profiles (used for candidate weights/caps)
|
||||
* @param placementPolicy which policy governs unqualified spawns
|
||||
* @param liveCount live worker count per profile (must never return {@code null})
|
||||
*/
|
||||
public CompositePeerLauncher(List<HerdrPeerLauncher> delegates,
|
||||
String defaultProfile,
|
||||
Map<String, BridgedConfig.Worker> profileConfigs,
|
||||
PlacementPolicy placementPolicy,
|
||||
Function<String, Integer> liveCount) {
|
||||
if (delegates.isEmpty()) {
|
||||
throw new IllegalArgumentException("at least one peer adapter must be configured");
|
||||
}
|
||||
this.delegates = List.copyOf(delegates);
|
||||
this.defaultProfile = defaultProfile;
|
||||
this.profileConfigs = Map.copyOf(profileConfigs);
|
||||
this.placementPolicy = placementPolicy;
|
||||
this.liveCount = liveCount;
|
||||
Map<String, HerdrPeerLauncher> index = new LinkedHashMap<>();
|
||||
for (HerdrPeerLauncher d : this.delegates) {
|
||||
for (String profile : d.profiles()) {
|
||||
@@ -89,10 +132,64 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
|
||||
@Override
|
||||
public PeerHandle spawn(SpawnRequest req) {
|
||||
HerdrPeerLauncher d = route(req.profileName());
|
||||
PeerHandle handle = d.spawn(req);
|
||||
spawnedBy.put(handle.id(), d);
|
||||
return handle;
|
||||
String requestedProfile = req.profileName();
|
||||
if (requestedProfile != null && !requestedProfile.isBlank()) {
|
||||
// An explicit profile bypasses the policy entirely.
|
||||
HerdrPeerLauncher d = route(requestedProfile);
|
||||
PeerHandle handle = d.spawn(req);
|
||||
spawnedBy.put(handle.id(), d);
|
||||
return handle;
|
||||
}
|
||||
|
||||
List<PlacementCandidate> candidates = candidates();
|
||||
Set<String> unreachable = new HashSet<>();
|
||||
PlacementContext ctx = new PlacementContext(defaultProfile, candidates, liveCount, unreachable);
|
||||
|
||||
int maxAttempts = candidates.isEmpty() ? 1 : candidates.size();
|
||||
for (int attempt = 0; attempt < maxAttempts; attempt++) {
|
||||
PlacementCandidate chosen;
|
||||
try {
|
||||
chosen = placementPolicy.select(ctx);
|
||||
} catch (RuntimeException e) {
|
||||
// No candidate left (all at cap or all unreachable). The policy already threw a clear
|
||||
// message; do not wrap it in a generic PeerUnreachableException.
|
||||
throw e;
|
||||
}
|
||||
|
||||
HerdrPeerLauncher d = byProfile.get(chosen.profile());
|
||||
if (d == null) {
|
||||
// A configured profile with no adapter is a wiring bug; fail fast.
|
||||
unreachable.add(chosen.profile());
|
||||
continue;
|
||||
}
|
||||
|
||||
SpawnRequest routedReq = new SpawnRequest(chosen.profile(), req.requestedCwd(), req.callerCwd());
|
||||
try {
|
||||
PeerHandle handle = d.spawn(routedReq);
|
||||
spawnedBy.put(handle.id(), d);
|
||||
return handle;
|
||||
} catch (PeerUnreachableException e) {
|
||||
log.warn("spawn on profile {} unreachable, will retry next candidate if any: {}",
|
||||
chosen.profile(), e.getMessage());
|
||||
unreachable.add(chosen.profile());
|
||||
// Update the context for the next selection so the policy excludes this profile.
|
||||
ctx = new PlacementContext(defaultProfile, candidates, liveCount, unreachable);
|
||||
}
|
||||
}
|
||||
|
||||
throw new PeerUnreachableException(
|
||||
"no reachable worker profile available after trying " + unreachable.size()
|
||||
+ " candidate(s): " + String.join(", ", unreachable));
|
||||
}
|
||||
|
||||
/** Build the candidate list from the configured profiles, in definition order. */
|
||||
private List<PlacementCandidate> candidates() {
|
||||
List<PlacementCandidate> out = new ArrayList<>();
|
||||
for (Map.Entry<String, BridgedConfig.Worker> e : profileConfigs.entrySet()) {
|
||||
BridgedConfig.Worker w = e.getValue();
|
||||
out.add(new PlacementCandidate(e.getKey(), null, w.weight(), w.maxLoad()));
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
@@ -214,9 +214,10 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
if (spawnReadyTimeoutMs > 0) {
|
||||
waitUntilInjectableOrThrow(paneId);
|
||||
}
|
||||
// CB-519: the handle id is a host-unique UUID; the herdr pane it maps to stays internal.
|
||||
String id = UUID.randomUUID().toString();
|
||||
paneByAgentId.put(id, paneId);
|
||||
return new WorkerHandle(id, agent.terminalId());
|
||||
return new WorkerHandle(id, agent.terminalId(), requireProfile(req.profileName()).profile());
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -513,8 +514,8 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
+ spawnReadyTimeoutMs + "ms");
|
||||
}
|
||||
|
||||
/** A concrete {@link PeerHandle} wrapping herdr agent coordinates. */
|
||||
private record WorkerHandle(String id, String terminalId) implements PeerHandle {
|
||||
/** A concrete {@link PeerHandle} wrapping herdr agent coordinates and the profile that spawned it. */
|
||||
private record WorkerHandle(String id, String terminalId, String profile) implements PeerHandle {
|
||||
}
|
||||
|
||||
// --- shared helpers ------------------------------------------------------------------------
|
||||
|
||||
@@ -333,6 +333,9 @@ class BridgedConfigTest {
|
||||
parityOverlay: [".mcp.json", ".env"]
|
||||
gitTokenEnv: GITEA_TOKEN
|
||||
gitHostEnv: GITEA_HOST
|
||||
weight: 0.5
|
||||
maxLoad: 2
|
||||
placement: weighted
|
||||
lifecycle:
|
||||
idleTtlSeconds: 300
|
||||
contextCap: 10
|
||||
@@ -356,7 +359,10 @@ class BridgedConfigTest {
|
||||
assertEquals(java.util.List.of(".mcp.json", ".env"), w.parityOverlay());
|
||||
assertTrue(w.hasGitToken(), "gitTokenEnv binds and enables the CB-302 PR grant");
|
||||
assertEquals("GITEA_HOST", w.gitHostEnv());
|
||||
assertEquals(0.5f, w.weight(), 0.0001f, "weight binds as a float");
|
||||
assertEquals(2, w.maxLoad(), "maxLoad binds as an integer");
|
||||
|
||||
assertEquals("weighted", cfg.placement(), "placement binds at the top level");
|
||||
assertEquals(300, cfg.lifecycle().idleTtlSeconds());
|
||||
assertEquals(10, cfg.lifecycle().contextCap());
|
||||
assertEquals(5, cfg.lifecycle().drainTimeoutSeconds());
|
||||
@@ -365,4 +371,36 @@ class BridgedConfigTest {
|
||||
assertEquals(5, cfg.primary().remindersOrDefault());
|
||||
assertEquals(15000L, cfg.primary().backoffMsOrDefault());
|
||||
}
|
||||
|
||||
@Test
|
||||
void placementDefaultsToFixedForExistingConfigs(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("no-placement.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
workers:
|
||||
gx10:
|
||||
baseUrl: http://gx10.gw:8000
|
||||
""");
|
||||
|
||||
BridgedConfig cfg = BridgedConfig.load(f);
|
||||
assertEquals("fixed", cfg.placement(), "omitted placement must default to fixed");
|
||||
}
|
||||
|
||||
@Test
|
||||
void workerWeightAndMaxLoadDefaultSanely(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("no-weight.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
workers:
|
||||
gx10:
|
||||
baseUrl: http://gx10.gw:8000
|
||||
""");
|
||||
|
||||
BridgedConfig cfg = BridgedConfig.load(f);
|
||||
BridgedConfig.Worker w = cfg.workerProfiles().get("gx10");
|
||||
assertEquals(1.0f, w.weight(), 0.0001f, "absent weight defaults to 1.0");
|
||||
assertNull(w.maxLoad(), "absent maxLoad defaults to unlimited (null)");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,183 @@
|
||||
package dev.ltms.bridged.placement;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.HashSet;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.function.Function;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
/**
|
||||
* Unit tests for the placement policies. They run with no herdr and no launcher — pure selection
|
||||
* logic exercised through the descriptor type so CB-308 host expansion will not need to rewrite
|
||||
* these assertions.
|
||||
*/
|
||||
class PlacementPolicyTest {
|
||||
|
||||
private static Function<String, Integer> noSessions() {
|
||||
return name -> 0;
|
||||
}
|
||||
|
||||
private static PlacementContext ctx(List<PlacementCandidate> candidates,
|
||||
Function<String, Integer> liveCount,
|
||||
Set<String> unreachable) {
|
||||
return new PlacementContext("b", candidates, liveCount, unreachable);
|
||||
}
|
||||
|
||||
private static PlacementContext ctx(List<PlacementCandidate> candidates,
|
||||
Function<String, Integer> liveCount) {
|
||||
return ctx(candidates, liveCount, Set.of());
|
||||
}
|
||||
|
||||
@Test
|
||||
void fixedReturnsDefaultEvenIfOtherProfilesExist() {
|
||||
PlacementPolicy policy = PlacementPolicies.fixed();
|
||||
PlacementContext ctx = ctx(List.of(
|
||||
PlacementCandidate.profile("a"),
|
||||
PlacementCandidate.profile("b")), noSessions());
|
||||
assertEquals("b", policy.select(ctx).profile());
|
||||
}
|
||||
|
||||
@Test
|
||||
void fixedFallsBackToFirstCandidateWhenNoDefault() {
|
||||
PlacementPolicy policy = PlacementPolicies.fixed();
|
||||
PlacementContext ctx = new PlacementContext(null,
|
||||
List.of(PlacementCandidate.profile("a"), PlacementCandidate.profile("b")),
|
||||
noSessions(), Set.of());
|
||||
assertEquals("a", policy.select(ctx).profile());
|
||||
}
|
||||
|
||||
@Test
|
||||
void fixedThrowsWhenNoProfilesAndNoDefault() {
|
||||
PlacementPolicy policy = PlacementPolicies.fixed();
|
||||
PlacementContext ctx = new PlacementContext(null, List.of(), noSessions(), Set.of());
|
||||
assertThrows(PlacementException.class, () -> policy.select(ctx));
|
||||
}
|
||||
|
||||
@Test
|
||||
void roundRobinCyclesThroughAvailableProfiles() {
|
||||
PlacementPolicy policy = PlacementPolicies.roundRobin();
|
||||
List<PlacementCandidate> candidates = List.of(
|
||||
PlacementCandidate.profile("a"),
|
||||
PlacementCandidate.profile("b"),
|
||||
PlacementCandidate.profile("c"));
|
||||
assertEquals("a", policy.select(ctx(candidates, noSessions())).profile());
|
||||
assertEquals("b", policy.select(ctx(candidates, noSessions())).profile());
|
||||
assertEquals("c", policy.select(ctx(candidates, noSessions())).profile());
|
||||
assertEquals("a", policy.select(ctx(candidates, noSessions())).profile());
|
||||
}
|
||||
|
||||
@Test
|
||||
void roundRobinSkipsProfilesAtMaxLoad() {
|
||||
PlacementPolicy policy = PlacementPolicies.roundRobin();
|
||||
List<PlacementCandidate> candidates = List.of(
|
||||
PlacementCandidate.profile("a", 1.0f, 2),
|
||||
PlacementCandidate.profile("b", 1.0f, null));
|
||||
Function<String, Integer> liveCount = Map.of("a", 2)::get;
|
||||
for (int i = 0; i < 5; i++) {
|
||||
assertEquals("b", policy.select(ctx(candidates, liveCount)).profile());
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void roundRobinThrowsWhenAllAtMaxLoad() {
|
||||
PlacementPolicy policy = PlacementPolicies.roundRobin();
|
||||
List<PlacementCandidate> candidates = List.of(
|
||||
PlacementCandidate.profile("a", 1.0f, 1),
|
||||
PlacementCandidate.profile("b", 1.0f, 1));
|
||||
Function<String, Integer> liveCount = Map.of("a", 1, "b", 1)::get;
|
||||
PlacementException e = assertThrows(PlacementException.class,
|
||||
() -> policy.select(ctx(candidates, liveCount)));
|
||||
assertTrue(e.getMessage().contains("maxLoad"), e.getMessage());
|
||||
}
|
||||
|
||||
@Test
|
||||
void weightedAlternatesEvenlyWithEqualWeights() {
|
||||
PlacementPolicy policy = PlacementPolicies.weighted();
|
||||
List<PlacementCandidate> candidates = List.of(
|
||||
PlacementCandidate.profile("a", 0.5f, null),
|
||||
PlacementCandidate.profile("b", 0.5f, null));
|
||||
int a = 0, b = 0;
|
||||
for (int i = 0; i < 100; i++) {
|
||||
String p = policy.select(ctx(candidates, noSessions())).profile();
|
||||
if ("a".equals(p)) a++;
|
||||
else if ("b".equals(p)) b++;
|
||||
}
|
||||
assertEquals(50, a, "equal weights should split 50/50");
|
||||
assertEquals(50, b);
|
||||
}
|
||||
|
||||
@Test
|
||||
void weightedHoldsThreeToOneRatio() {
|
||||
PlacementPolicy policy = PlacementPolicies.weighted();
|
||||
List<PlacementCandidate> candidates = List.of(
|
||||
PlacementCandidate.profile("a", 0.75f, null),
|
||||
PlacementCandidate.profile("b", 0.25f, null));
|
||||
int a = 0, b = 0;
|
||||
for (int i = 0; i < 40; i++) {
|
||||
String p = policy.select(ctx(candidates, noSessions())).profile();
|
||||
if ("a".equals(p)) a++;
|
||||
else if ("b".equals(p)) b++;
|
||||
}
|
||||
assertEquals(30, a, "0.75/0.25 should yield a 3:1 ratio over a multiple of 4");
|
||||
assertEquals(10, b);
|
||||
}
|
||||
|
||||
@Test
|
||||
void weightedSkipsProfileAtMaxLoad() {
|
||||
PlacementPolicy policy = PlacementPolicies.weighted();
|
||||
List<PlacementCandidate> candidates = List.of(
|
||||
PlacementCandidate.profile("a", 1.0f, 1),
|
||||
PlacementCandidate.profile("b", 1.0f, null));
|
||||
Function<String, Integer> liveCount = name -> "a".equals(name) ? 1 : 0;
|
||||
for (int i = 0; i < 5; i++) {
|
||||
assertEquals("b", policy.select(ctx(candidates, liveCount)).profile());
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void weightedThrowsWhenAllAtMaxLoad() {
|
||||
PlacementPolicy policy = PlacementPolicies.weighted();
|
||||
List<PlacementCandidate> candidates = List.of(
|
||||
PlacementCandidate.profile("a", 1.0f, 1),
|
||||
PlacementCandidate.profile("b", 1.0f, 1));
|
||||
Function<String, Integer> liveCount = name -> 1;
|
||||
PlacementException e = assertThrows(PlacementException.class,
|
||||
() -> policy.select(ctx(candidates, liveCount)));
|
||||
assertTrue(e.getMessage().contains("maxLoad"), e.getMessage());
|
||||
}
|
||||
|
||||
@Test
|
||||
void weightedThrowsWhenAllUnreachable() {
|
||||
PlacementPolicy policy = PlacementPolicies.weighted();
|
||||
List<PlacementCandidate> candidates = List.of(
|
||||
PlacementCandidate.profile("a"),
|
||||
PlacementCandidate.profile("b"));
|
||||
PlacementException e = assertThrows(PlacementException.class,
|
||||
() -> policy.select(ctx(candidates, noSessions(), Set.of("a", "b"))));
|
||||
assertTrue(e.getMessage().contains("unreachable"), e.getMessage());
|
||||
}
|
||||
|
||||
@Test
|
||||
void mixedExclusionMessageNamesBothReasons() {
|
||||
PlacementPolicy policy = PlacementPolicies.weighted();
|
||||
List<PlacementCandidate> candidates = List.of(
|
||||
PlacementCandidate.profile("a", 1.0f, 1),
|
||||
PlacementCandidate.profile("b", 1.0f, null));
|
||||
Function<String, Integer> liveCount = name -> "a".equals(name) ? 1 : 0;
|
||||
Set<String> unreachable = new HashSet<>();
|
||||
unreachable.add("b");
|
||||
PlacementException e = assertThrows(PlacementException.class,
|
||||
() -> policy.select(ctx(candidates, liveCount, unreachable)));
|
||||
assertTrue(e.getMessage().contains("1 at maxLoad"), e.getMessage());
|
||||
assertTrue(e.getMessage().contains("1 unreachable"), e.getMessage());
|
||||
}
|
||||
|
||||
@Test
|
||||
void unknownPolicyNameThrows() {
|
||||
assertThrows(IllegalArgumentException.class, () -> PlacementPolicies.fromName("random"));
|
||||
}
|
||||
}
|
||||
@@ -512,7 +512,7 @@ class ClaudeCodeLauncherTest {
|
||||
BridgedConfig.Worker cfg = new BridgedConfig.Worker(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "BRIDGED_WORKER_TOKEN",
|
||||
List.of("claude"), "tab", "bridged-workers", "w #{n}", null, null, null, null, null,
|
||||
null, Map.of("JAVA_HOME", "/opt/jdk", "PATH", "/profile/bin"));
|
||||
null, Map.of("JAVA_HOME", "/opt/jdk", "PATH", "/profile/bin"), null, null);
|
||||
new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
|
||||
k -> "PATH".equals(k) ? "/daemon/bin" : null).spawn();
|
||||
@@ -533,7 +533,7 @@ class ClaudeCodeLauncherTest {
|
||||
BridgedConfig.Worker cfg = new BridgedConfig.Worker(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "BRIDGED_WORKER_TOKEN",
|
||||
List.of("claude"), "tab", "bridged-workers", "w #{n}", null, null, null, null, null,
|
||||
null, Map.of("ANTHROPIC_BASE_URL", "http://evil.example.com"));
|
||||
null, Map.of("ANTHROPIC_BASE_URL", "http://evil.example.com"), null, null);
|
||||
new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
|
||||
_ -> null).spawn();
|
||||
|
||||
@@ -2,17 +2,26 @@ package dev.ltms.bridged.worker;
|
||||
|
||||
import dev.ltms.bridged.config.BridgedConfig;
|
||||
import dev.ltms.bridged.guard.SubscriptionGuard;
|
||||
import dev.ltms.bridged.herdr.Agent;
|
||||
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.PeerHandle;
|
||||
import dev.ltms.bridged.peer.PeerLauncher;
|
||||
import dev.ltms.bridged.peer.PeerUnreachableException;
|
||||
import dev.ltms.bridged.peer.SpawnRequest;
|
||||
import dev.ltms.bridged.placement.PlacementException;
|
||||
import dev.ltms.bridged.placement.PlacementPolicies;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.EnumSet;
|
||||
import java.util.HashMap;
|
||||
import java.util.HashSet;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.function.Function;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
@@ -46,6 +55,73 @@ class CompositePeerLauncherTest {
|
||||
List.of(claudeAdapter(herdr), opencodeAdapter(herdr)), "claude");
|
||||
}
|
||||
|
||||
/**
|
||||
* A minimal concrete HerdrPeerLauncher for policy tests. It either returns a fake handle for the
|
||||
* requested profile or throws, depending on {@code failProfiles}. buildLaunch is a stub; only
|
||||
* spawn/stop/list/caps/reap are exercised by the composite.
|
||||
*/
|
||||
private static final class StubLauncher extends HerdrPeerLauncher {
|
||||
private final Set<String> failProfiles;
|
||||
private final Map<String, Integer> spawnCounts = new HashMap<>();
|
||||
|
||||
StubLauncher(String prefix, FakeHerdr herdr,
|
||||
Map<String, BridgedConfig.Worker> profiles, String defaultProfile,
|
||||
Set<String> failProfiles) {
|
||||
super(prefix, new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
profiles, defaultProfile, _ -> null, 0L, System::currentTimeMillis, () -> { });
|
||||
this.failProfiles = Set.copyOf(failProfiles);
|
||||
}
|
||||
|
||||
@Override
|
||||
protected Launch buildLaunch(BridgedConfig.Worker cfg) {
|
||||
return new Launch(Map.of(), List.of());
|
||||
}
|
||||
|
||||
@Override
|
||||
public PeerHandle spawn(SpawnRequest req) {
|
||||
String p = (req.profileName() == null || req.profileName().isBlank())
|
||||
? defaultProfile() : req.profileName();
|
||||
spawnCounts.merge(p, 1, Integer::sum);
|
||||
if (failProfiles.contains(p)) {
|
||||
throw new PeerUnreachableException(p + " is down");
|
||||
}
|
||||
return new PeerHandle() {
|
||||
@Override public String id() { return "pane-" + p; }
|
||||
@Override public String terminalId() { return "term-" + p; }
|
||||
@Override public String profile() { return p; }
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public void stop(String id) { }
|
||||
|
||||
@Override
|
||||
public List<Agent> list() { return List.of(); }
|
||||
|
||||
@Override
|
||||
public Set<Capability> capabilities() { return EnumSet.noneOf(Capability.class); }
|
||||
|
||||
@Override
|
||||
public int reapOrphanWorkers() { return 0; }
|
||||
|
||||
int spawnCount(String profile) {
|
||||
return spawnCounts.getOrDefault(profile, 0);
|
||||
}
|
||||
}
|
||||
|
||||
private static BridgedConfig.Worker stubWorker(String profile) {
|
||||
return new BridgedConfig.Worker(profile, "http://gx00.gw:8000", "coder",
|
||||
null, "BRIDGED_WORKER_TOKEN", List.of("claude"), "tab", "bridged-workers",
|
||||
"w #{n}", null, null, null, null, null, null, null, null, null);
|
||||
}
|
||||
|
||||
private static BridgedConfig.Worker stubWorker(String profile, float weight, Integer maxLoad) {
|
||||
return new BridgedConfig.Worker(profile, "http://gx00.gw:8000", "coder",
|
||||
null, "BRIDGED_WORKER_TOKEN", List.of("claude"), "tab", "bridged-workers",
|
||||
"w #{n}", null, null, null, null, null, null, null,
|
||||
weight, maxLoad);
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
private static String startedName(FakeHerdr herdr) {
|
||||
return (String) ((Map<String, Object>) herdr.lastCall("agent.start").params()).get("name");
|
||||
@@ -158,4 +234,99 @@ class CompositePeerLauncherTest {
|
||||
() -> new CompositePeerLauncher(List.of(), "claude"),
|
||||
"at least one adapter must be configured");
|
||||
}
|
||||
|
||||
@Test
|
||||
void fixedDefaultIsNoOpForUnqualifiedSpawns() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
CompositePeerLauncher composite = composite(herdr);
|
||||
|
||||
PeerHandle h = composite.spawn(new SpawnRequest(null, null, null));
|
||||
assertTrue(startedName(herdr).startsWith("claude-"),
|
||||
"fixed placement still routes an unqualified spawn to the default profile");
|
||||
assertEquals("claude", h.profile(), "the returned handle carries the resolved default profile");
|
||||
}
|
||||
|
||||
@Test
|
||||
void weightedPolicyGatesProfileAtMaxLoad() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Map<String, BridgedConfig.Worker> profiles = Map.of(
|
||||
"a", stubWorker("a", 1.0f, 1),
|
||||
"b", stubWorker("b", 1.0f, null));
|
||||
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "a", Set.of());
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(
|
||||
List.of(adapter), "a", profiles, PlacementPolicies.weighted(), name -> "a".equals(name) ? 1 : 0);
|
||||
|
||||
for (int i = 0; i < 5; i++) {
|
||||
PeerHandle h = composite.spawn(new SpawnRequest(null, null, null));
|
||||
assertEquals("b", h.profile(), "profile a is at maxLoad, so every spawn must land on b");
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void weightedPolicyDistributesAccordingToWeightRatio() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Map<String, BridgedConfig.Worker> profiles = Map.of(
|
||||
"a", stubWorker("a", 0.75f, null),
|
||||
"b", stubWorker("b", 0.25f, null));
|
||||
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "a", Set.of());
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(
|
||||
List.of(adapter), "a", profiles, PlacementPolicies.weighted(), name -> 0);
|
||||
|
||||
int a = 0, b = 0;
|
||||
for (int i = 0; i < 40; i++) {
|
||||
String p = composite.spawn(new SpawnRequest(null, null, null)).profile();
|
||||
if ("a".equals(p)) a++;
|
||||
else if ("b".equals(p)) b++;
|
||||
}
|
||||
assertEquals(30, a, "weighted distribution should hold the 3:1 ratio");
|
||||
assertEquals(10, b);
|
||||
}
|
||||
|
||||
@Test
|
||||
void failoverRetriesNextCandidateWhenProfileIsUnreachable() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Map<String, BridgedConfig.Worker> profiles = Map.of(
|
||||
"a", stubWorker("a"),
|
||||
"b", stubWorker("b"));
|
||||
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "a", Set.of("a"));
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(
|
||||
List.of(adapter), "a", profiles, PlacementPolicies.weighted(), name -> 0);
|
||||
|
||||
PeerHandle h = composite.spawn(new SpawnRequest(null, null, null));
|
||||
assertEquals("b", h.profile(), "the spawn must fail over from unreachable a to b");
|
||||
assertEquals(1, adapter.spawnCount("a"), "a was tried once and failed");
|
||||
assertEquals(1, adapter.spawnCount("b"), "b was tried once and succeeded");
|
||||
}
|
||||
|
||||
@Test
|
||||
void failoverBoundedByCandidateCount() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Map<String, BridgedConfig.Worker> profiles = Map.of(
|
||||
"a", stubWorker("a"),
|
||||
"b", stubWorker("b"));
|
||||
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "a", Set.of("a", "b"));
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(
|
||||
List.of(adapter), "a", profiles, PlacementPolicies.weighted(), name -> 0);
|
||||
|
||||
PeerUnreachableException e = assertThrows(PeerUnreachableException.class,
|
||||
() -> composite.spawn(new SpawnRequest(null, null, null)));
|
||||
assertTrue(e.getMessage().contains("no reachable worker profile"), e.getMessage());
|
||||
assertEquals(1, adapter.spawnCount("a"));
|
||||
assertEquals(1, adapter.spawnCount("b"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void emptyCandidateSetThrowsClearException() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Map<String, BridgedConfig.Worker> profiles = Map.of(
|
||||
"a", stubWorker("a", 1.0f, 1),
|
||||
"b", stubWorker("b", 1.0f, 1));
|
||||
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "a", Set.of());
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(
|
||||
List.of(adapter), "a", profiles, PlacementPolicies.weighted(), name -> 1);
|
||||
|
||||
PlacementException e = assertThrows(PlacementException.class,
|
||||
() -> composite.spawn(new SpawnRequest(null, null, null)));
|
||||
assertTrue(e.getMessage().contains("maxLoad"), e.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user