From a7f0211e2f13b982abfcc56df8003adcd54bcca5 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Tue, 4 Aug 2026 17:50:04 +0200 Subject: [PATCH] CB-518: weighted placement policy for worker spawns --- bridged/bridged.example.yaml | 7 + .../main/java/dev/ltms/bridged/Bridged.java | 14 +- .../ltms/bridged/config/BridgedConfig.java | 26 ++- .../dev/ltms/bridged/peer/PeerHandle.java | 11 ++ .../placement/FixedPlacementPolicy.java | 22 +++ .../bridged/placement/PlacementCandidate.java | 20 ++ .../bridged/placement/PlacementContext.java | 19 ++ .../bridged/placement/PlacementException.java | 12 ++ .../bridged/placement/PlacementPolicies.java | 41 ++++ .../bridged/placement/PlacementPolicy.java | 16 ++ .../placement/PlacementPolicyUtil.java | 65 +++++++ .../placement/RoundRobinPlacementPolicy.java | 23 +++ .../placement/WeightedRoundRobinPolicy.java | 50 +++++ .../ltms/bridged/session/SessionManager.java | 35 ++-- .../bridged/worker/CompositePeerLauncher.java | 105 +++++++++- .../bridged/worker/HerdrPeerLauncher.java | 7 +- .../bridged/config/BridgedConfigTest.java | 38 ++++ .../placement/PlacementPolicyTest.java | 183 ++++++++++++++++++ .../worker/ClaudeCodeLauncherTest.java | 4 +- .../worker/CompositePeerLauncherTest.java | 171 ++++++++++++++++ 20 files changed, 842 insertions(+), 27 deletions(-) create mode 100644 bridged/src/main/java/dev/ltms/bridged/placement/FixedPlacementPolicy.java create mode 100644 bridged/src/main/java/dev/ltms/bridged/placement/PlacementCandidate.java create mode 100644 bridged/src/main/java/dev/ltms/bridged/placement/PlacementContext.java create mode 100644 bridged/src/main/java/dev/ltms/bridged/placement/PlacementException.java create mode 100644 bridged/src/main/java/dev/ltms/bridged/placement/PlacementPolicies.java create mode 100644 bridged/src/main/java/dev/ltms/bridged/placement/PlacementPolicy.java create mode 100644 bridged/src/main/java/dev/ltms/bridged/placement/PlacementPolicyUtil.java create mode 100644 bridged/src/main/java/dev/ltms/bridged/placement/RoundRobinPlacementPolicy.java create mode 100644 bridged/src/main/java/dev/ltms/bridged/placement/WeightedRoundRobinPolicy.java create mode 100644 bridged/src/test/java/dev/ltms/bridged/placement/PlacementPolicyTest.java diff --git a/bridged/bridged.example.yaml b/bridged/bridged.example.yaml index 080812c..6fc69bb 100644 --- a/bridged/bridged.example.yaml +++ b/bridged/bridged.example.yaml @@ -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 diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java index 6720ae5..be9bfe1 100644 --- a/bridged/src/main/java/dev/ltms/bridged/Bridged.java +++ b/bridged/src/main/java/dev/ltms/bridged/Bridged.java @@ -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> 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; diff --git a/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java b/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java index 83c6526..fce933e 100644 --- a/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java +++ b/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java @@ -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} only 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 parityOverlay, String gitTokenEnv, String gitHostEnv, String kind, - Map env) { + Map 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 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 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 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); } /** diff --git a/bridged/src/main/java/dev/ltms/bridged/peer/PeerHandle.java b/bridged/src/main/java/dev/ltms/bridged/peer/PeerHandle.java index bf46bd9..e290451 100644 --- a/bridged/src/main/java/dev/ltms/bridged/peer/PeerHandle.java +++ b/bridged/src/main/java/dev/ltms/bridged/peer/PeerHandle.java @@ -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; + } } diff --git a/bridged/src/main/java/dev/ltms/bridged/placement/FixedPlacementPolicy.java b/bridged/src/main/java/dev/ltms/bridged/placement/FixedPlacementPolicy.java new file mode 100644 index 0000000..e0896ce --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/placement/FixedPlacementPolicy.java @@ -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"); + } +} diff --git a/bridged/src/main/java/dev/ltms/bridged/placement/PlacementCandidate.java b/bridged/src/main/java/dev/ltms/bridged/placement/PlacementCandidate.java new file mode 100644 index 0000000..91232d0 --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/placement/PlacementCandidate.java @@ -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}. + * + *

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); + } +} diff --git a/bridged/src/main/java/dev/ltms/bridged/placement/PlacementContext.java b/bridged/src/main/java/dev/ltms/bridged/placement/PlacementContext.java new file mode 100644 index 0000000..068f526 --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/placement/PlacementContext.java @@ -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 candidates, + Function liveCount, + Set unreachable) { +} diff --git a/bridged/src/main/java/dev/ltms/bridged/placement/PlacementException.java b/bridged/src/main/java/dev/ltms/bridged/placement/PlacementException.java new file mode 100644 index 0000000..9afdff3 --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/placement/PlacementException.java @@ -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); + } +} diff --git a/bridged/src/main/java/dev/ltms/bridged/placement/PlacementPolicies.java b/bridged/src/main/java/dev/ltms/bridged/placement/PlacementPolicies.java new file mode 100644 index 0000000..a339602 --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/placement/PlacementPolicies.java @@ -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(); + } +} diff --git a/bridged/src/main/java/dev/ltms/bridged/placement/PlacementPolicy.java b/bridged/src/main/java/dev/ltms/bridged/placement/PlacementPolicy.java new file mode 100644 index 0000000..8d8de1d --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/placement/PlacementPolicy.java @@ -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); +} diff --git a/bridged/src/main/java/dev/ltms/bridged/placement/PlacementPolicyUtil.java b/bridged/src/main/java/dev/ltms/bridged/placement/PlacementPolicyUtil.java new file mode 100644 index 0000000..1e32f8b --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/placement/PlacementPolicyUtil.java @@ -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 available(PlacementContext ctx) { + List 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"); + } +} diff --git a/bridged/src/main/java/dev/ltms/bridged/placement/RoundRobinPlacementPolicy.java b/bridged/src/main/java/dev/ltms/bridged/placement/RoundRobinPlacementPolicy.java new file mode 100644 index 0000000..c76115f --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/placement/RoundRobinPlacementPolicy.java @@ -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 available = PlacementPolicyUtil.available(ctx); + if (available.isEmpty()) { + throw PlacementPolicyUtil.emptyException(ctx); + } + return available.get(index.getAndIncrement() % available.size()); + } +} diff --git a/bridged/src/main/java/dev/ltms/bridged/placement/WeightedRoundRobinPolicy.java b/bridged/src/main/java/dev/ltms/bridged/placement/WeightedRoundRobinPolicy.java new file mode 100644 index 0000000..c97ca7b --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/placement/WeightedRoundRobinPolicy.java @@ -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. + * + *

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 current = new ConcurrentHashMap<>(); + + @Override + public synchronized PlacementCandidate select(PlacementContext ctx) { + List 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; + } +} diff --git a/bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java b/bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java index 9af1940..82ae833 100644 --- a/bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java +++ b/bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java @@ -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; diff --git a/bridged/src/main/java/dev/ltms/bridged/worker/CompositePeerLauncher.java b/bridged/src/main/java/dev/ltms/bridged/worker/CompositePeerLauncher.java index 45922a9..40ad67e 100644 --- a/bridged/src/main/java/dev/ltms/bridged/worker/CompositePeerLauncher.java +++ b/bridged/src/main/java/dev/ltms/bridged/worker/CompositePeerLauncher.java @@ -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. * + * + *

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 spawnedBy = new ConcurrentHashMap<>(); + private final Map profileConfigs; + private final PlacementPolicy placementPolicy; + private final Function 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 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 delegates, + String defaultProfile, + Map profileConfigs, + PlacementPolicy placementPolicy, + Function 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 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 candidates = candidates(); + Set 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 candidates() { + List out = new ArrayList<>(); + for (Map.Entry e : profileConfigs.entrySet()) { + BridgedConfig.Worker w = e.getValue(); + out.add(new PlacementCandidate(e.getKey(), null, w.weight(), w.maxLoad())); + } + return out; } @Override diff --git a/bridged/src/main/java/dev/ltms/bridged/worker/HerdrPeerLauncher.java b/bridged/src/main/java/dev/ltms/bridged/worker/HerdrPeerLauncher.java index f5d4ccc..0958bef 100644 --- a/bridged/src/main/java/dev/ltms/bridged/worker/HerdrPeerLauncher.java +++ b/bridged/src/main/java/dev/ltms/bridged/worker/HerdrPeerLauncher.java @@ -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 ------------------------------------------------------------------------ diff --git a/bridged/src/test/java/dev/ltms/bridged/config/BridgedConfigTest.java b/bridged/src/test/java/dev/ltms/bridged/config/BridgedConfigTest.java index 5e396f2..1590ca3 100644 --- a/bridged/src/test/java/dev/ltms/bridged/config/BridgedConfigTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/config/BridgedConfigTest.java @@ -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)"); + } } diff --git a/bridged/src/test/java/dev/ltms/bridged/placement/PlacementPolicyTest.java b/bridged/src/test/java/dev/ltms/bridged/placement/PlacementPolicyTest.java new file mode 100644 index 0000000..105056b --- /dev/null +++ b/bridged/src/test/java/dev/ltms/bridged/placement/PlacementPolicyTest.java @@ -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 noSessions() { + return name -> 0; + } + + private static PlacementContext ctx(List candidates, + Function liveCount, + Set unreachable) { + return new PlacementContext("b", candidates, liveCount, unreachable); + } + + private static PlacementContext ctx(List candidates, + Function 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 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 candidates = List.of( + PlacementCandidate.profile("a", 1.0f, 2), + PlacementCandidate.profile("b", 1.0f, null)); + Function 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 candidates = List.of( + PlacementCandidate.profile("a", 1.0f, 1), + PlacementCandidate.profile("b", 1.0f, 1)); + Function 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 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 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 candidates = List.of( + PlacementCandidate.profile("a", 1.0f, 1), + PlacementCandidate.profile("b", 1.0f, null)); + Function 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 candidates = List.of( + PlacementCandidate.profile("a", 1.0f, 1), + PlacementCandidate.profile("b", 1.0f, 1)); + Function 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 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 candidates = List.of( + PlacementCandidate.profile("a", 1.0f, 1), + PlacementCandidate.profile("b", 1.0f, null)); + Function liveCount = name -> "a".equals(name) ? 1 : 0; + Set 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")); + } +} diff --git a/bridged/src/test/java/dev/ltms/bridged/worker/ClaudeCodeLauncherTest.java b/bridged/src/test/java/dev/ltms/bridged/worker/ClaudeCodeLauncherTest.java index b0e68f7..48769b9 100644 --- a/bridged/src/test/java/dev/ltms/bridged/worker/ClaudeCodeLauncherTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/worker/ClaudeCodeLauncherTest.java @@ -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(); diff --git a/bridged/src/test/java/dev/ltms/bridged/worker/CompositePeerLauncherTest.java b/bridged/src/test/java/dev/ltms/bridged/worker/CompositePeerLauncherTest.java index c888fc1..08f006f 100644 --- a/bridged/src/test/java/dev/ltms/bridged/worker/CompositePeerLauncherTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/worker/CompositePeerLauncherTest.java @@ -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 failProfiles; + private final Map spawnCounts = new HashMap<>(); + + StubLauncher(String prefix, FakeHerdr herdr, + Map profiles, String defaultProfile, + Set 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 list() { return List.of(); } + + @Override + public Set 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) 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 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 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 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 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 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()); + } }