Compare commits
14 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| e501d39988 | |||
| c00a86b32c | |||
| d3ae0350a2 | |||
| 2c2196a1f1 | |||
| 35ade14630 | |||
| 541df87272 | |||
| 337b6ccd6e | |||
| 42f46dfe9a | |||
| 9088d2b2c5 | |||
| 500bfa2c33 | |||
| 9ca9c43dfa | |||
| 1966c69994 | |||
| 976eff8ad1 | |||
| 0af902ec43 |
@@ -47,15 +47,18 @@ bind:
|
||||
# (say a Claude lead and an opencode lead) work as peers: the second is silently demoted and refused
|
||||
# every orchestration call. List each lead's pane here and all of them resolve as leads.
|
||||
#
|
||||
# terminal → the ONLY field identity depends on; get it from that session's bridge_whoami
|
||||
# tab → the ONLY field identity depends on (CB-579); the exact label of the tab hosting the lead.
|
||||
# Label the tab yourself, or let bridged label one it launches — see `fleet.leaders:` below.
|
||||
# kind/model → descriptive; they document what runs in the pane and are echoed by bridge_whoami
|
||||
#
|
||||
# A lead is never spawned — it pre-exists, which is exactly why it must be named rather than created.
|
||||
# A lead's tab must already carry its label (or be launched by bridged, which labels it) — there is
|
||||
# no terminal id to paste in and nothing to re-pin when the session restarts: the tab survives, so
|
||||
# the same label resolves the same lead again on the next scan.
|
||||
# `bridge_whoami` reports `{"role":"primary","leader":"<name>"}`; role stays "primary" because a lead
|
||||
# IS a primary for authorization, so nothing that keys on the role breaks.
|
||||
#
|
||||
# KEEP `primary:` when adding leads: it still addresses the CB-307 push loop, which needs a single
|
||||
# destination for its nudges. If both name the same terminal, the `fleet.leaders:` entry wins.
|
||||
# destination for its nudges, and is a separate mechanism from lead identity — see `fleet.leaders:`.
|
||||
#
|
||||
# Leads are configured under `fleet.leaders:` — see THE FLEET further down.
|
||||
#
|
||||
@@ -135,6 +138,23 @@ herdrSocket: ~/.config/herdr/herdr.sock
|
||||
# SSH is unaffected). The token value itself is never stored in this file.
|
||||
# gitHostEnv → host env var holding the forge host (default GITEA_HOST). Injected as
|
||||
# GITEA_HOST *only* alongside a resolved gitTokenEnv.
|
||||
# exhaustedPattern → regex matched against a completion-fallback scrape (CB-578 stage A) to
|
||||
# classify a turn that ended with no bridge_reply as the backend having
|
||||
# refused on a subscription usage limit, rather than a real answer. Opt-in —
|
||||
# omit and this profile's completion fallback behaves exactly as before.
|
||||
# Every backend words its refusal differently, so this is config, never a
|
||||
# vendor string baked into bridged itself.
|
||||
# DEFERRED: compiled once into a startup pattern map — editing it needs a
|
||||
# daemon restart, same as this profile's model/baseUrl/argv.
|
||||
# credentialId → CB-578 stage B: the credential this profile quarantines WITH when a
|
||||
# BACKEND_EXHAUSTED classification fires. Two profiles that set the SAME
|
||||
# credentialId share one quarantine — the case this exists for is two models
|
||||
# on one account (e.g. sol and terra both billing one OpenAI credential): an
|
||||
# exhaustion on either one must lock out both, or the fleet just walks onto
|
||||
# the same dead account under the sibling's name. Opt-in — omit and this
|
||||
# profile quarantines alone, under its own name, exactly as if the field did
|
||||
# not exist. Cooldown length is the top-level quarantineCooldownSeconds below.
|
||||
# HOT: read live at every spawn/exhaustion check — no restart needed.
|
||||
# env → extra environment for this profile's workers, as a literal key/value map
|
||||
# (CB-511). Use it to give workers a toolchain.
|
||||
#
|
||||
@@ -168,6 +188,8 @@ profiles:
|
||||
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
|
||||
# exhaustedPattern: "usage limit has been reached" # opt-in: classify a usage-limit refusal (CB-578)
|
||||
# credentialId: shared-openai # opt-in: quarantine together with every other profile sharing this id (CB-578)
|
||||
# configDir: /Users/me/.ccs/instances/gx10 # CLAUDE_CONFIG_DIR — inherit that profile's skills/MCP
|
||||
# cwd: /Users/me/src/myrepo # pin the working dir; omit to inherit the primary's
|
||||
# parityOverlay: [".claude/settings.local.json", ".env", ".envrc"] # never add .mcp.json — see above
|
||||
@@ -233,6 +255,15 @@ profiles:
|
||||
# round-robin, or weighted. Omitting this key is a strict no-op for existing configs.
|
||||
placement: weighted
|
||||
|
||||
# How long a credential sits out after a BACKEND_EXHAUSTED classification (CB-578 stage B), in
|
||||
# seconds, before a spawn may land on it again. Applies to every profile's effective credential
|
||||
# (its own name, or its credentialId if set above) — there is no per-profile override. Default
|
||||
# 1800 (30 minutes) when omitted or non-positive.
|
||||
# DEFERRED: baked once into the BackendQuarantine built at startup — a running quarantine keeps
|
||||
# its original cooldown regardless; a new value only applies to a quarantine that starts after a
|
||||
# restart. Editing this needs a daemon restart to take effect.
|
||||
# quarantineCooldownSeconds: 1800
|
||||
|
||||
# Re-read this file without restarting the daemon (CB-559). Off unless you add this block, so an
|
||||
# upgraded bridged keeps the old behaviour: the file is read once at boot and never again.
|
||||
# enabled → turn the watch on. bridged checks the file's modified time on a timer and
|
||||
@@ -242,17 +273,20 @@ placement: weighted
|
||||
# Not every key can move under a running daemon, and the difference is about what already exists
|
||||
# when the reload happens — not about how important the key is:
|
||||
# HOT → takes effect on the next spawn: the whole `fleet:` block (every role pool,
|
||||
# `charters`, and `tabLabel`), `placement:`, and an existing profile's weight / maxLoad. Those are
|
||||
# hot because the placement policy reads them through a supplier — being config is
|
||||
# `charters`, and `tabLabel`), `placement:`, and an existing profile's weight / maxLoad
|
||||
# / credentialId. Those are hot because the placement policy (and, for credentialId,
|
||||
# the CB-578 stage B quarantine check) reads them through a supplier — being config is
|
||||
# not by itself enough to make a key hot.
|
||||
# DEFERRED → accepted into the new config, but the wiring built at startup keeps the old value
|
||||
# until you restart: `lifecycle:`, `leadHeartbeat:`, `guard:`, `worktreeRoot:`,
|
||||
# `spawnReadyTimeoutMs` / `spawnReadyPollMs`, ADDING or REMOVING a profile (a new
|
||||
# backend needs its own launcher, and launchers are built once), AND an existing
|
||||
# profile's launch settings — model, baseUrl, argv, env, configDir, mcpUrl, tabLabel.
|
||||
# The launcher takes a copy of `profiles:` at startup and resolves every spawn out of
|
||||
# that copy, so those never reach a launch until you restart. The reload logs them by
|
||||
# name rather than pretending they applied.
|
||||
# `spawnReadyTimeoutMs` / `spawnReadyPollMs`, `quarantineCooldownSeconds` (CB-578
|
||||
# stage B — baked once into the quarantine tracker built at startup), ADDING or
|
||||
# REMOVING a profile (a new backend needs its own launcher, and launchers are built
|
||||
# once), AND an existing profile's launch settings — model, baseUrl, argv, env,
|
||||
# configDir, mcpUrl, tabLabel, exhaustedPattern. The launcher takes a copy of
|
||||
# `profiles:` at startup and resolves every spawn out of that copy, so those never
|
||||
# reach a launch until you restart. The reload logs them by name rather than
|
||||
# pretending they applied.
|
||||
# COLD → cannot change at all: `bind:`, `herdrSocket:`, `broker:` and `auth:`. The socket is
|
||||
# bound, the broker connection is open, and the auth mode decides who may reach the
|
||||
# port that is already listening.
|
||||
@@ -309,15 +343,18 @@ fleet:
|
||||
|
||||
# Panes that orchestrate rather than are orchestrated. A lead may now be CREATED as well as
|
||||
# recognised: give it a `profile:` and the daemon launches the shortfall when fewer than
|
||||
# `instances` are live. Give it only a `terminal:` and it is recognise-only, as before.
|
||||
# `instances` are live. Omit `profile:` and it is recognise-only, as before.
|
||||
#
|
||||
# `tabPrefix` is the naming convention that finds a lead without pasting a terminal id: label the
|
||||
# tab `lead: <name>` when you open it and the pane is recognised on the next rescan. Reopen the
|
||||
# tab later and the id changes; the label does not.
|
||||
# `tab:` (CB-579) is REQUIRED and is the only field identity depends on — the exact label of the
|
||||
# tab hosting the lead, matched case-insensitively. Label the tab yourself and put that same
|
||||
# string here, and the pane is recognised on the next rescan. Reopen the tab later, or the session
|
||||
# inside it restarts — the terminal id changes; the tab, and its label, do not, so no config edit
|
||||
# follows a restart.
|
||||
#
|
||||
# A lead the daemon launches is labelled BY the daemon, using the same convention, so it is found
|
||||
# A lead the daemon launches is labelled BY the daemon with this same `tab:` value, so it is found
|
||||
# by the same scan. A lead counts as live only when herdr also reports a running agent in that
|
||||
# tab — a label left behind by a session that died does not block the relaunch.
|
||||
# tab — a label left behind by a session that died does not block the relaunch, and a tab that is
|
||||
# gone entirely drops out of the next scan rather than being remembered forever.
|
||||
#
|
||||
# An auto-launched lead is NOT a member: it gets no worker reply charter, is never registered with
|
||||
# the session lifecycle (the idle reaper would kill your orchestrator), and stays on the
|
||||
@@ -326,10 +363,9 @@ fleet:
|
||||
# opus-5.0:
|
||||
# profile: opus # omit to never create this lead, only recognise it
|
||||
# instances: 1 # desired live count; only the shortfall is launched. 0 = off
|
||||
# terminal: term_0123456789abcd # optional hand-pin; usually found by tabPrefix instead.
|
||||
# # A running agent on this terminal also counts as live, so a
|
||||
# # lead you opened by hand is not relaunched under you.
|
||||
# tabPrefix: "lead:" # `lead: opus-5.0` ⇒ a lead named opus-5.0 (case-insensitive)
|
||||
# tab: "lead: opus-5.0" # REQUIRED — the exact tab label this lead lives in
|
||||
# tabPrefix: "lead:" # only used to guard against a worker tabLabel colliding with
|
||||
# # this convention at startup; plays no part in matching a lead
|
||||
# scanIntervalSeconds: 10 # rescan cadence, and the worst case before a new tab is seen
|
||||
# workspace: leads # where a launched lead's tab is created (default "leads").
|
||||
# # MUST NOT be a member workspace — those are excluded from the
|
||||
@@ -337,7 +373,7 @@ fleet:
|
||||
# cwd: /path/to/repo # the launched lead's working directory (default: bridged's own)
|
||||
# kind: claude # descriptive; reported by bridge_whoami
|
||||
# gpt-sol-5.6:
|
||||
# terminal: term_fedcba9876543
|
||||
# tab: "lead: gpt-sol-5.6"
|
||||
# kind: opencode
|
||||
# model: openai/gpt-5.6-terra
|
||||
|
||||
|
||||
@@ -13,6 +13,8 @@ import dev.ltms.bridged.herdr.PaneLocator;
|
||||
import dev.ltms.bridged.herdr.UnixSocketHerdrClient;
|
||||
import dev.ltms.bridged.herdr.WorkspaceControl;
|
||||
import dev.ltms.bridged.inject.CompletionResolver;
|
||||
import dev.ltms.bridged.inject.ExhaustedPatternLookup;
|
||||
import dev.ltms.bridged.inject.ExhaustionSink;
|
||||
import dev.ltms.bridged.inject.Injector;
|
||||
import dev.ltms.bridged.inject.StatusPoller;
|
||||
import dev.ltms.bridged.inject.TurnListener;
|
||||
@@ -36,6 +38,7 @@ import dev.ltms.bridged.msg.LeadHeartbeatLoop;
|
||||
import dev.ltms.bridged.msg.ReplyPushLoop;
|
||||
import dev.ltms.bridged.rest.BridgedApp;
|
||||
import dev.ltms.bridged.session.GitWorktrees;
|
||||
import dev.ltms.bridged.session.MemberSession;
|
||||
import dev.ltms.bridged.session.SessionManager;
|
||||
import dev.ltms.bridged.peer.PeerLauncher;
|
||||
import dev.ltms.bridged.session.SessionReaper;
|
||||
@@ -43,6 +46,7 @@ import dev.ltms.bridged.member.ClaudeCodeLauncher;
|
||||
import dev.ltms.bridged.member.CompositePeerLauncher;
|
||||
import dev.ltms.bridged.member.HerdrPeerLauncher;
|
||||
import dev.ltms.bridged.member.OpenCodeLauncher;
|
||||
import dev.ltms.bridged.placement.BackendQuarantine;
|
||||
import io.javalin.Javalin;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
@@ -60,6 +64,7 @@ import java.util.concurrent.atomic.AtomicReference;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.Predicate;
|
||||
import java.util.function.Supplier;
|
||||
import java.util.regex.Pattern;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
/**
|
||||
@@ -139,11 +144,18 @@ public final class Bridged {
|
||||
() -> config.get().fleet()));
|
||||
}
|
||||
AtomicReference<Function<String, Integer>> liveCountRef = new AtomicReference<>(_ -> 0);
|
||||
// CB-578 stage B: one quarantine tracker for the whole daemon, shared between the launcher
|
||||
// (checked at spawn) and the exhaustion sink wired in below (written on BACKEND_EXHAUSTED).
|
||||
// The cooldown is deferred (see BridgedConfig#quarantineCooldownSeconds): it is read once
|
||||
// here, at startup, and a config reload only changes it for a daemon restart.
|
||||
BackendQuarantine quarantine = new BackendQuarantine(System::nanoTime,
|
||||
TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds()));
|
||||
PeerLauncher workers = new CompositePeerLauncher(
|
||||
adapters,
|
||||
cfg.effectiveDefaultProfile(),
|
||||
config,
|
||||
profileName -> liveCountRef.get().apply(profileName));
|
||||
profileName -> liveCountRef.get().apply(profileName),
|
||||
quarantine);
|
||||
// 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
|
||||
@@ -193,34 +205,34 @@ public final class Bridged {
|
||||
if (leadTerminals.size() > 1) {
|
||||
log.info("leads: {} panes recognised {}", leadTerminals.size(), leadTerminals.values());
|
||||
}
|
||||
// CB-531: on top of the static registry, discover leads by the tab labels the operator
|
||||
// writes. CB-557 moved the settings onto the lead they describe, so scanning is on whenever
|
||||
// a `fleet.leaders:` entry exists — with no leads configured the supplier is a constant and
|
||||
// never touches herdr, exactly as a missing `leadScan:` block used to behave.
|
||||
// CB-531: on top of the legacy primary.terminal pin, discover leads by the tab labels the
|
||||
// operator writes. CB-557 moved the settings onto the lead they describe, so scanning is on
|
||||
// whenever a `fleet.leaders:` entry exists — with no leads configured the supplier is a
|
||||
// constant and never touches herdr, exactly as a missing `leadScan:` block used to behave.
|
||||
// CB-579: each lead now names its own exact `tab:` label, so one scanner discovers every
|
||||
// configured lead regardless of how differently their tabs are labelled — the old
|
||||
// single-shared-tabPrefix limitation (and its warning) is gone.
|
||||
final Supplier<Map<String, String>> leads;
|
||||
var leaders = cfg.fleet().leaders();
|
||||
if (!leaders.isEmpty()) {
|
||||
// One scanner, so one prefix and one interval. Distinct per-lead prefixes would need a
|
||||
// scanner each; until a config actually wants that, take the first entry's settings and
|
||||
// say so, rather than silently honouring one lead's prefix and dropping another's.
|
||||
var scan = leaders.values().iterator().next();
|
||||
Set<String> memberSpaces = cfg.profiles().values().stream()
|
||||
.map(BridgedConfig.Profile::workspace)
|
||||
.filter(Objects::nonNull)
|
||||
.collect(Collectors.toSet());
|
||||
leads = new LeadTabScanner(herdr, scan.tabPrefix(), memberSpaces, leadTerminals,
|
||||
TimeUnit.SECONDS.toNanos(scan.scanIntervalSeconds()), System::nanoTime);
|
||||
log.info("lead scan: tabs labelled '{}…' host a lead (rescan every {}s, member spaces {} "
|
||||
+ "excluded)",
|
||||
scan.tabPrefix(), scan.scanIntervalSeconds(), memberSpaces);
|
||||
long distinctPrefixes = leaders.values().stream()
|
||||
.map(BridgedConfig.Leader::tabPrefix).distinct().count();
|
||||
if (distinctPrefixes > 1) {
|
||||
log.warn("fleet.leaders declares {} different tabPrefix values; only '{}' is scanned "
|
||||
+ "for. Give every lead the same tabPrefix, or leads under the others "
|
||||
+ "will not be discovered.",
|
||||
distinctPrefixes, scan.tabPrefix());
|
||||
}
|
||||
Map<String, String> tabToName = new LinkedHashMap<>();
|
||||
leaders.forEach((name, leader) -> {
|
||||
if (leader != null && leader.tab() != null && !leader.tab().isBlank()) {
|
||||
tabToName.put(leader.tab(), name);
|
||||
}
|
||||
});
|
||||
// One shared rescan cadence: still taken from the first entry, as before — it is an
|
||||
// operational cadence, not identity, so there is no correctness reason to give every
|
||||
// lead its own scanner.
|
||||
int scanIntervalSeconds = leaders.values().iterator().next().scanIntervalSeconds();
|
||||
leads = new LeadTabScanner(herdr, tabToName, memberSpaces,
|
||||
TimeUnit.SECONDS.toNanos(scanIntervalSeconds), System::nanoTime);
|
||||
log.info("lead scan: tabs {} host a lead (rescan every {}s, member spaces {} excluded)",
|
||||
tabToName.keySet(), scanIntervalSeconds, memberSpaces);
|
||||
} else {
|
||||
leads = () -> leadTerminals;
|
||||
}
|
||||
@@ -253,7 +265,40 @@ public final class Bridged {
|
||||
// The blocking message endpoint (CB-104) is the producer; the poller is inert until then.
|
||||
// CB-106: a confirmed turn completion resolves a blocked send whose worker never replied.
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver completion = new CompletionResolver(agents, rendezvous);
|
||||
// CB-578 stage A: classify a completion-fallback scrape that matches a profile's configured
|
||||
// usage-limit refusal as BACKEND_EXHAUSTED rather than handing it back as a real answer.
|
||||
// Compiled once at startup, keyed by profile name; a profile with no exhaustedPattern is
|
||||
// simply absent here, so its workers keep today's completion-fallback behaviour unchanged.
|
||||
Map<String, Pattern> exhaustedPatternsByProfile = new LinkedHashMap<>();
|
||||
cfg.profiles().forEach((name, profile) -> {
|
||||
if (profile.hasExhaustedPattern()) {
|
||||
exhaustedPatternsByProfile.put(name, Pattern.compile(profile.exhaustedPattern()));
|
||||
}
|
||||
});
|
||||
ExhaustedPatternLookup exhaustedPatterns = target -> sessions.roster().stream()
|
||||
.filter(session -> target.equals(session.terminalId()))
|
||||
.findFirst()
|
||||
.map(session -> exhaustedPatternsByProfile.get(session.profile()))
|
||||
.orElse(null);
|
||||
log.info("backend-exhausted classification (CB-578 stage A): {}",
|
||||
CompletionResolver.coverage(cfg.profiles().keySet(), exhaustedPatternsByProfile.keySet()));
|
||||
// CB-578 stage B: on a classification that actually wins, quarantine the exhausted profile's
|
||||
// CREDENTIAL — not the profile name — so a profile sharing that credential (e.g. two models
|
||||
// on one OpenAI account) is refused too, not just the one that happened to report it. Reads
|
||||
// the profile config live off `config`, so a credentialId edit is hot: no restart needed.
|
||||
ExhaustionSink exhaustionSink = (target, reason) -> sessions.roster().stream()
|
||||
.filter(session -> target.equals(session.terminalId()))
|
||||
.findFirst()
|
||||
.map(MemberSession::profile)
|
||||
.map(profileName -> config.get().profiles().get(profileName))
|
||||
.ifPresent(profile -> {
|
||||
String credentialId = profile.effectiveCredentialId();
|
||||
quarantine.quarantine(credentialId);
|
||||
log.warn("credential '{}' quarantined for {}s (profile '{}' classified "
|
||||
+ "BACKEND_EXHAUSTED): {}", credentialId,
|
||||
cfg.quarantineCooldownSeconds(), profile.profile(), reason);
|
||||
});
|
||||
CompletionResolver completion = new CompletionResolver(agents, rendezvous, exhaustedPatterns, exhaustionSink);
|
||||
// CB-113: deliver only to an available worker (its MCP is connected), never its boot window.
|
||||
// CB-301: the manager's presence bridge records availability and drives SPAWNING → READY.
|
||||
MemberPresence presence = sessions.asPresence();
|
||||
@@ -360,8 +405,10 @@ public final class Bridged {
|
||||
var healthScheduler = Executors.newSingleThreadScheduledExecutor(r ->
|
||||
Thread.ofVirtual().name("bridge-health-").unstarted(r));
|
||||
if (cfg.health() != null && cfg.health().isEnabled()) {
|
||||
// CB-580: a member found GONE/NEVER_READY must fail whatever ticket is waiting on it,
|
||||
// through the same idempotent target-wide operation CB-516 already uses on release.
|
||||
healthMonitor = new FleetHealthMonitor(agents, sessions::roster, messages, healthScheduler,
|
||||
System::nanoTime, cfg.health().intervalOrDefault());
|
||||
System::nanoTime, cfg.health().intervalOrDefault(), messages::abandon);
|
||||
String coverage = FleetHealthMonitor.coverage(true,
|
||||
cfg.health().notifications() != null && cfg.health().notifications().configured());
|
||||
if ("detection-only".equals(coverage)) {
|
||||
@@ -419,7 +466,11 @@ public final class Bridged {
|
||||
var health = config.get().health();
|
||||
return FleetHealthMonitor.coverage(health != null && health.isEnabled(),
|
||||
health != null && health.notifications() != null && health.notifications().configured());
|
||||
}));
|
||||
}),
|
||||
new BridgeMcp.QuarantineSource(profile -> {
|
||||
var configured = config.get().profiles().get(profile);
|
||||
return configured == null ? null : configured.effectiveCredentialId();
|
||||
}, quarantine));
|
||||
|
||||
// CB-559: opt-in config reload. With no `configReload:` block nothing is constructed, so an
|
||||
// upgraded daemon behaves exactly as before — the file is read once at boot and never again.
|
||||
|
||||
@@ -58,6 +58,12 @@ import java.util.Set;
|
||||
* {@code fixed} (default), {@code round-robin}, or {@code weighted}
|
||||
* @param auth API authentication mode ({@code null} → {@code loopback-trust}, the
|
||||
* historical behaviour), CB-501
|
||||
* @param quarantineCooldownSeconds how long a credential stays quarantined after a
|
||||
* {@code BACKEND_EXHAUSTED} classification (CB-578 stage B); {@code null}/{@code
|
||||
* <=0} → {@link #DEFAULT_QUARANTINE_COOLDOWN_SECONDS}. Baked once into the
|
||||
* {@code BackendQuarantine} built at startup, so it is DEFERRED: changing it
|
||||
* needs a restart, and a quarantine already running keeps whatever cooldown was
|
||||
* live when it started.
|
||||
*/
|
||||
@JsonIgnoreProperties(ignoreUnknown = true)
|
||||
public record BridgedConfig(
|
||||
@@ -76,7 +82,11 @@ public record BridgedConfig(
|
||||
Health health,
|
||||
String placement,
|
||||
Auth auth,
|
||||
ConfigReload configReload) {
|
||||
ConfigReload configReload,
|
||||
Integer quarantineCooldownSeconds) {
|
||||
|
||||
/** Default cooldown (CB-578 stage B) when {@code quarantineCooldownSeconds} is absent/non-positive. */
|
||||
public static final int DEFAULT_QUARANTINE_COOLDOWN_SECONDS = 1800;
|
||||
|
||||
/** Back-compat 14-arg form — no {@code configReload:} block, so file watching stays off. */
|
||||
public BridgedConfig(Bind bind, String herdrSocket, Map<String, Profile> profiles, Guard guard,
|
||||
@@ -84,7 +94,7 @@ public record BridgedConfig(
|
||||
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
|
||||
LeadHeartbeat leadHeartbeat, String placement, Auth auth) {
|
||||
this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
|
||||
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, null, placement, auth, null);
|
||||
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, null, placement, auth, null, null);
|
||||
}
|
||||
|
||||
/** Back-compat form before the optional {@code health:} block was added. */
|
||||
@@ -93,7 +103,18 @@ public record BridgedConfig(
|
||||
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
|
||||
LeadHeartbeat leadHeartbeat, String placement, Auth auth, ConfigReload configReload) {
|
||||
this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
|
||||
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, null, placement, auth, configReload);
|
||||
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, null, placement, auth, configReload, null);
|
||||
}
|
||||
|
||||
/** Back-compat form before the CB-578 stage B {@code quarantineCooldownSeconds} field was added. */
|
||||
public BridgedConfig(Bind bind, String herdrSocket, Map<String, Profile> profiles, Guard guard,
|
||||
String worktreeRoot, Lifecycle lifecycle, Integer spawnReadyTimeoutMs,
|
||||
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
|
||||
LeadHeartbeat leadHeartbeat, Health health, String placement, Auth auth,
|
||||
ConfigReload configReload) {
|
||||
this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
|
||||
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth,
|
||||
configReload, null);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -181,6 +202,23 @@ public record BridgedConfig(
|
||||
* opposite intents). For the same reason, an {@code env:} entry naming
|
||||
* {@code ANTHROPIC_BASE_URL} or {@code ANTHROPIC_AUTH_TOKEN} is refused at
|
||||
* config load (CB-542): on the subscription path no guard would vet it.
|
||||
* @param exhaustedPattern regex matched against a completion-fallback scrape (CB-578 stage A) to
|
||||
* classify a turn that ended with no {@code bridge_reply} as the backend
|
||||
* having refused on a subscription usage limit, rather than a real answer.
|
||||
* {@code null}/blank ⇒ the classification never fires for this profile and
|
||||
* today's completion-fallback behaviour is unchanged. Every backend words
|
||||
* its refusal differently, so this is config, never a vendor string in code.
|
||||
* @param credentialId the shared account this profile authenticates as (CB-578 stage B). Two
|
||||
* or more profiles setting the <em>same</em> non-blank value are quarantined
|
||||
* together by one {@code exhaustedPattern} classification on any one of
|
||||
* them — the case a single OpenAI (or any other) credential backing several
|
||||
* profiles ({@code sol}, {@code terra}, ...) needs, so a spawn does not walk
|
||||
* straight onto the other profile sharing the same exhausted account.
|
||||
* {@code null}/blank ⇒ this profile's {@link #effectiveCredentialId()} is
|
||||
* its own name, so it quarantines alone — today's behaviour for every
|
||||
* profile that does not opt in. Read live off the current config, so it is
|
||||
* HOT: a change takes effect on the next exhaustion classification / spawn,
|
||||
* no restart needed.
|
||||
*/
|
||||
@JsonIgnoreProperties(ignoreUnknown = true)
|
||||
public record Profile(String profile, String baseUrl, String model,
|
||||
@@ -193,7 +231,9 @@ public record BridgedConfig(
|
||||
Map<String, String> env,
|
||||
Float weight,
|
||||
Integer maxLoad,
|
||||
Boolean subscription) {
|
||||
Boolean subscription,
|
||||
String exhaustedPattern,
|
||||
String credentialId) {
|
||||
|
||||
/** Peer kind spawned by {@link dev.ltms.bridged.member.ClaudeCodeLauncher} (the default). */
|
||||
public static final String KIND_CLAUDE_CODE = "claude-code";
|
||||
@@ -236,6 +276,13 @@ public record BridgedConfig(
|
||||
weight = (weight == null || weight <= 0.0f) ? 1.0f : weight;
|
||||
maxLoad = (maxLoad == null || maxLoad <= 0) ? null : maxLoad;
|
||||
subscription = (subscription != null && subscription) ? Boolean.TRUE : Boolean.FALSE;
|
||||
// exhaustedPattern stays null when unset/blank (opt-in) — no defaulting, no vendor
|
||||
// wording: an unconfigured profile keeps today's completion-fallback behaviour exactly.
|
||||
exhaustedPattern = (exhaustedPattern == null || exhaustedPattern.isBlank()) ? null : exhaustedPattern;
|
||||
// credentialId stays null when unset/blank — effectiveCredentialId() is where the
|
||||
// "quarantines alone" fallback actually lives, so today's behaviour needs no defaulting
|
||||
// here at all.
|
||||
credentialId = (credentialId == null || credentialId.isBlank()) ? null : credentialId;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -248,7 +295,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, null, null, null, null);
|
||||
mcpUrl, cwd, parityOverlay, null, null, null, null, null, null, null, null);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -260,7 +307,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, null, null, null);
|
||||
mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, null, null, null, null, null, null);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -273,13 +320,14 @@ 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, null, null, null);
|
||||
mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, null, null, null, null, null);
|
||||
}
|
||||
|
||||
/** A copy with {@code profile} set — used to default a profile to its {@code workers} key. */
|
||||
public Profile withProfile(String p) {
|
||||
return new Profile(p, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel,
|
||||
mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, env, weight, maxLoad, subscription);
|
||||
mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, env, weight, maxLoad, subscription,
|
||||
exhaustedPattern, credentialId);
|
||||
}
|
||||
|
||||
/** True when this profile is served by the Claude Code adapter (the default kind). */
|
||||
@@ -312,7 +360,39 @@ public record BridgedConfig(
|
||||
String cwd, List<String> parityOverlay, String gitTokenEnv, String gitHostEnv,
|
||||
String kind, Map<String, String> env, Float weight, Integer maxLoad) {
|
||||
this(profile, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel,
|
||||
mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, env, weight, maxLoad, null);
|
||||
mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, env, weight, maxLoad, null, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* Backward-compatible constructor without the CB-578 stage B {@code credentialId} — the
|
||||
* profile quarantines alone under its own name (see {@link #effectiveCredentialId()}). Keeps
|
||||
* pre-stage-B call sites (and any YAML that omits the key) compiling and behaving identically.
|
||||
*/
|
||||
public Profile(String profile, String baseUrl, String model,
|
||||
String configDir, String tokenEnv, List<String> argv,
|
||||
String placement, String workspace, String tabLabel, String mcpUrl,
|
||||
String cwd, List<String> parityOverlay, String gitTokenEnv, String gitHostEnv,
|
||||
String kind, Map<String, String> env, Float weight, Integer maxLoad,
|
||||
Boolean subscription, String exhaustedPattern) {
|
||||
this(profile, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel,
|
||||
mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, env, weight, maxLoad,
|
||||
subscription, exhaustedPattern, null);
|
||||
}
|
||||
|
||||
/** True when this profile's CB-578 stage A backend-exhausted classification is configured. */
|
||||
public boolean hasExhaustedPattern() {
|
||||
return exhaustedPattern != null;
|
||||
}
|
||||
|
||||
/**
|
||||
* The credential group this profile quarantines with (CB-578 stage B): the configured
|
||||
* {@link #credentialId} when set, else this profile's own name — so an unconfigured profile
|
||||
* quarantines alone, exactly as it did before this field existed. Two profiles that set the
|
||||
* same non-blank {@code credentialId} share one quarantine: a {@code BACKEND_EXHAUSTED}
|
||||
* classification on either one quarantines both.
|
||||
*/
|
||||
public String effectiveCredentialId() {
|
||||
return (credentialId == null || credentialId.isBlank()) ? profile : credentialId;
|
||||
}
|
||||
|
||||
/** True when this profile's workers are granted a forge token to open their own PR (CB-302). */
|
||||
@@ -458,24 +538,32 @@ public record BridgedConfig(
|
||||
* pre-existed, which is why it had to be recognised by configuration rather than created. With
|
||||
* {@code profile} and {@code instances} the daemon may stand one up when none is live, so the
|
||||
* pane no longer has to exist before the daemon does. Recognition still comes first: a lead
|
||||
* already running under {@code tabPrefix} is adopted, and only the shortfall is launched.
|
||||
* already running in its configured {@code tab} is adopted, and only the shortfall is launched.
|
||||
*
|
||||
* <p><b>{@code tab} replaced {@code terminal} (CB-579).</b> A herdr {@code terminal_id} changes
|
||||
* every time the lead's session restarts, so pinning one cost a config edit and a daemon restart
|
||||
* per restart. A tab is stable: a human opens it once, it holds exactly one pane, and its label
|
||||
* survives restarts of the agent inside it — so identity is now the tab label alone.
|
||||
*
|
||||
* @param profile the {@code profiles:} entry to launch this lead on when one must
|
||||
* be created; {@code null} ⇒ recognise-only, never create
|
||||
* @param terminal the lead's herdr {@code terminal_id} when pinned by hand; the only
|
||||
* field identity depends on. {@code null} ⇒ found by {@code tabPrefix}
|
||||
* @param tab the exact tab label hosting this lead, matched case-insensitively;
|
||||
* the only field identity depends on. Required — a lead with no
|
||||
* {@code tab} can never be discovered, launched or not
|
||||
* @param instances how many of this lead should be live (default 1). The daemon
|
||||
* launches only the shortfall, so a restart adopts rather than doubles
|
||||
* @param tabPrefix label prefix marking this lead's tab, matched case-insensitively;
|
||||
* the remainder is the lead's name ({@code "lead: opus"} →
|
||||
* {@code opus}). Default {@code "lead:"}
|
||||
* @param tabPrefix no longer used to find a lead's tab — {@code tab} is matched
|
||||
* exactly. Its only remaining job is the startup collision guard
|
||||
* ({@link #validateLeadTabPrefixes()}), which still uses it to refuse
|
||||
* a worker {@code tabLabel} template that could be misread as a lead.
|
||||
* Default {@code "lead:"}
|
||||
* @param scanIntervalSeconds how long a tab scan is cached before herdr is asked again; also the
|
||||
* worst case before a newly-labelled tab is recognised. Default 10
|
||||
* @param kind which agent runs there ({@code claude}, {@code opencode}, …)
|
||||
* @param model the model or selector it runs, for operators reading the roster
|
||||
*/
|
||||
@JsonIgnoreProperties(ignoreUnknown = true)
|
||||
public record Leader(String profile, String terminal, Integer instances, String tabPrefix,
|
||||
public record Leader(String profile, String tab, Integer instances, String tabPrefix,
|
||||
Integer scanIntervalSeconds, String kind, String model,
|
||||
String workspace, String cwd) {
|
||||
|
||||
@@ -493,12 +581,13 @@ public record BridgedConfig(
|
||||
(scanIntervalSeconds == null || scanIntervalSeconds <= 0) ? 10 : scanIntervalSeconds;
|
||||
workspace = (workspace == null || workspace.isBlank())
|
||||
? DEFAULT_WORKSPACE : workspace.strip();
|
||||
tab = (tab == null || tab.isBlank()) ? null : tab.strip();
|
||||
}
|
||||
|
||||
/** Back-compat 7-arg form — no workspace or cwd, so both take their defaults. */
|
||||
public Leader(String profile, String terminal, Integer instances, String tabPrefix,
|
||||
public Leader(String profile, String tab, Integer instances, String tabPrefix,
|
||||
Integer scanIntervalSeconds, String kind, String model) {
|
||||
this(profile, terminal, instances, tabPrefix, scanIntervalSeconds, kind, model, null, null);
|
||||
this(profile, tab, instances, tabPrefix, scanIntervalSeconds, kind, model, null, null);
|
||||
}
|
||||
|
||||
/** True when this lead may be launched by the daemon rather than only recognised. */
|
||||
@@ -506,9 +595,9 @@ public record BridgedConfig(
|
||||
return profile != null && !profile.isBlank() && instances > 0;
|
||||
}
|
||||
|
||||
/** The tab label an auto-launched instance of this lead gets — what the scanner reads back. */
|
||||
public String tabLabel(String name) {
|
||||
return tabPrefix + " " + name;
|
||||
/** The tab label an auto-launched instance of this lead gets — its configured {@code tab}. */
|
||||
public String tabLabel() {
|
||||
return tab;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -707,28 +796,20 @@ public record BridgedConfig(
|
||||
}
|
||||
|
||||
/**
|
||||
* The terminal → lead-name map that {@link dev.ltms.bridged.auth.CallerResolver} resolves
|
||||
* against, merging the {@code leaders:} registry with the legacy singular {@code primary:} pin.
|
||||
* The terminal → lead-name map seeded from the legacy singular {@code primary:} pin (CB-530).
|
||||
*
|
||||
* <p>Precedence: an explicit {@code leaders:} entry wins over the {@code primary:} pin for the
|
||||
* same terminal. The pin is the older, less expressive spelling of the same fact, so when both
|
||||
* name a pane the named entry is the one an operator meant. The pin is still honoured on its
|
||||
* own — a config carrying only {@code primary:} behaves exactly as it did before CB-530.
|
||||
* <p>{@code fleet.leaders} no longer carries a per-entry terminal pin (CB-579): a lead's identity
|
||||
* comes from its {@code tab} alone, resolved live by {@code LeadTabScanner}. This method now
|
||||
* exists only for the {@code primary.terminal} fallback — a config that never migrated off it
|
||||
* still resolves that one pane as a lead named {@code "primary"}, exactly as before CB-530.
|
||||
*
|
||||
* @return an unmodifiable map, empty when neither block is configured (nothing is pinned, and
|
||||
* every pane therefore resolves as a worker — the pre-CB-307 behaviour)
|
||||
* @return an unmodifiable map, empty when {@code primary.terminal} is not configured (nothing is
|
||||
* pinned, and every pane therefore resolves as a worker — the pre-CB-307 behaviour)
|
||||
*/
|
||||
public Map<String, String> leaderTerminals() {
|
||||
Map<String, String> byTerminal = new LinkedHashMap<>();
|
||||
if (fleet != null) {
|
||||
fleet.leaders().forEach((name, leader) -> {
|
||||
if (leader != null && leader.terminal() != null && !leader.terminal().isBlank()) {
|
||||
byTerminal.put(leader.terminal(), name);
|
||||
}
|
||||
});
|
||||
}
|
||||
if (primary != null && primary.terminal() != null && !primary.terminal().isBlank()) {
|
||||
byTerminal.putIfAbsent(primary.terminal(), "primary");
|
||||
byTerminal.put(primary.terminal(), "primary");
|
||||
}
|
||||
return Collections.unmodifiableMap(byTerminal);
|
||||
}
|
||||
@@ -833,13 +914,14 @@ public record BridgedConfig(
|
||||
private static final Set<String> KNOWN_TOP_LEVEL_KEYS = Set.of(
|
||||
"bind", "herdrSocket", "profiles", "guard", "worktreeRoot",
|
||||
"lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs", "broker", "primary", "fleet",
|
||||
"leadHeartbeat", "health", "placement", "auth", "configReload");
|
||||
"leadHeartbeat", "health", "placement", "auth", "configReload", "quarantineCooldownSeconds");
|
||||
|
||||
/** Load and validate config from {@code path}. */
|
||||
public static BridgedConfig load(Path path) {
|
||||
try {
|
||||
String yaml = Files.readString(path);
|
||||
rejectRenamedTopLevelKeys(yaml);
|
||||
rejectLeaderTerminalKey(yaml);
|
||||
warnUnknownTopLevelKeys(yaml, path);
|
||||
rejectDuplicateMemberSlots(yaml);
|
||||
BridgedConfig cfg = YAML.readValue(yaml, BridgedConfig.class);
|
||||
@@ -1064,6 +1146,43 @@ public record BridgedConfig(
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Reject a config whose {@code fleet.leaders.<name>} still carries the retired {@code terminal:}
|
||||
* pin (CB-579), naming {@code tab:} as its replacement.
|
||||
*
|
||||
* <p>{@code Leader} is {@code @JsonIgnoreProperties(ignoreUnknown = true)}, so simply dropping
|
||||
* the record component would make a leftover {@code terminal:} key silently no-op — the daemon
|
||||
* would start, the pin would never take effect, and nothing would say why. Fatal and specific
|
||||
* instead, exactly like {@link #rejectRenamedTopLevelKeys}, which this mirrors for a key one
|
||||
* level deeper than the ones that method covers.
|
||||
*
|
||||
* @param yaml the raw config text
|
||||
* @throws IllegalStateException when any {@code fleet.leaders.<name>.terminal} key is present
|
||||
*/
|
||||
static void rejectLeaderTerminalKey(String yaml) {
|
||||
Map<?, ?> raw;
|
||||
try {
|
||||
raw = YAML.readValue(yaml, Map.class);
|
||||
} catch (IOException | IllegalArgumentException e) {
|
||||
return; // a malformed file is reported by the real parse, not here
|
||||
}
|
||||
if (raw == null || !(raw.get("fleet") instanceof Map<?, ?> fleet)
|
||||
|| !(fleet.get("leaders") instanceof Map<?, ?> leaders)) {
|
||||
return;
|
||||
}
|
||||
List<String> bad = leaders.entrySet().stream()
|
||||
.filter(e -> e.getValue() instanceof Map<?, ?> leader && leader.containsKey("terminal"))
|
||||
.map(e -> String.valueOf(e.getKey()))
|
||||
.sorted()
|
||||
.toList();
|
||||
if (!bad.isEmpty()) {
|
||||
throw new IllegalStateException("refusing to start: fleet.leaders entries ["
|
||||
+ String.join(", ", bad) + "] still use the retired 'terminal:' key — replace it "
|
||||
+ "with 'tab:', the exact tab label hosting the lead. A terminal_id changes on "
|
||||
+ "every restart of the lead's session; a tab label does not.");
|
||||
}
|
||||
}
|
||||
|
||||
static List<String> unknownTopLevelKeys(String yaml) {
|
||||
Map<?, ?> raw;
|
||||
try {
|
||||
@@ -1102,8 +1221,15 @@ public record BridgedConfig(
|
||||
// configReload is left as-is: null is "off", and ConfigReload's own compact constructor
|
||||
// defaults the fields of a block that IS present. Defaulting it here would start watching
|
||||
// the file for every config that never asked to be watched.
|
||||
// quarantineCooldownSeconds IS defaulted, unlike leadHeartbeat/configReload above: it has no
|
||||
// separate on/off switch of its own (BackendQuarantine only ever quarantines a credential
|
||||
// after a BACKEND_EXHAUSTED classification, which stays opt-in via exhaustedPattern), so a
|
||||
// config that never mentions it should still get a sane cooldown rather than a null one.
|
||||
Integer quarantineCooldown = (quarantineCooldownSeconds != null && quarantineCooldownSeconds > 0)
|
||||
? quarantineCooldownSeconds : DEFAULT_QUARANTINE_COOLDOWN_SECONDS;
|
||||
return new BridgedConfig(b, herdrSocket, profiles, g, worktreeRoot, l, timeout, pollMs,
|
||||
broker, primary, f, leadHeartbeat, health, placementOrDefault, a, configReload);
|
||||
broker, primary, f, leadHeartbeat, health, placementOrDefault, a, configReload,
|
||||
quarantineCooldown);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1313,10 +1439,9 @@ public record BridgedConfig(
|
||||
+ "', which is not a configured profiles: entry (have: " + profiles.keySet()
|
||||
+ ").");
|
||||
}
|
||||
if (!leader.isCreatable() && (leader.terminal() == null || leader.terminal().isBlank())) {
|
||||
bad.add("fleet.leaders." + name + " can neither be found nor created — it pins no "
|
||||
+ "terminal: and names no profile: to launch one on. Give it one or the "
|
||||
+ "other, or drop the entry.");
|
||||
if (leader.tab() == null || leader.tab().isBlank()) {
|
||||
bad.add("fleet.leaders." + name + " has no tab: — a lead is now found (and, if "
|
||||
+ "auto-launched, labelled) purely by its tab, so every entry must name one.");
|
||||
}
|
||||
});
|
||||
if (!bad.isEmpty()) {
|
||||
|
||||
@@ -30,13 +30,22 @@ import java.util.function.Supplier;
|
||||
* {@code fleet:} (every role pool, {@code charters}, and {@code tabLabel}),
|
||||
* {@code placement:}, and an existing profile's {@code weight} / {@code maxLoad}. Those
|
||||
* three are read through a supplier on {@code CompositePeerLauncher}, which is what makes
|
||||
* them hot — not the fact that they are config.</li>
|
||||
* them hot — not the fact that they are config. <strong>This does NOT include
|
||||
* {@code fleet.leaders}</strong>: {@code Bridged.main} reads {@code cfg.fleet().leaders()}
|
||||
* once at startup to build the {@code LeadTabScanner} and the {@code LeadLauncher}, and
|
||||
* neither is reconstructed on reload — so a lead added, removed, or re-{@code tab}'d under
|
||||
* {@code fleet.leaders} needs a restart, the same as any deferred key below.</li>
|
||||
* <li><strong>Deferred</strong> — accepted into the new snapshot, but the wiring built at startup
|
||||
* keeps the old value until a restart: {@code lifecycle:}, {@code leadHeartbeat:},
|
||||
* {@code spawnReadyTimeoutMs} / {@code spawnReadyPollMs}, {@code guard:},
|
||||
* {@code worktreeRoot:}, adding or removing a profile (a new backend needs its own launcher,
|
||||
* {@code spawnReadyTimeoutMs} / {@code spawnReadyPollMs}, {@code quarantineCooldownSeconds}
|
||||
* (CB-578 stage B — baked once into the {@code BackendQuarantine} built at startup),
|
||||
* {@code guard:}, {@code worktreeRoot:}, adding or removing a profile (a new backend needs its own launcher,
|
||||
* which is constructed once), <em>and an existing profile's launch settings</em> —
|
||||
* {@code model}, {@code baseUrl}, {@code argv}, {@code env}, {@code mcpUrl} and the rest.
|
||||
* {@code model}, {@code baseUrl}, {@code argv}, {@code env}, {@code mcpUrl},
|
||||
* {@code exhaustedPattern} (CB-578 stage A — compiled once into {@code Bridged.main}'s
|
||||
* pattern map at startup), and the rest. {@code credentialId} (CB-578 stage B) is NOT on
|
||||
* this list — it is read live off the config supplier at every quarantine check and
|
||||
* exhaustion event, exactly like {@code weight} / {@code maxLoad}, so it is hot instead.
|
||||
* {@code HerdrPeerLauncher} takes {@code Map.copyOf(profiles)} at construction and resolves
|
||||
* each spawn out of that copy, so those never reach a launch until the daemon restarts. A
|
||||
* reload logs these rather than pretending they applied.</li>
|
||||
@@ -211,6 +220,12 @@ public final class ConfigRef implements Supplier<BridgedConfig> {
|
||||
|| !Objects.equals(old.spawnReadyPollMs(), fresh.spawnReadyPollMs())) {
|
||||
changed.add("spawnReady*");
|
||||
}
|
||||
// CB-578 stage B: baked once into the BackendQuarantine built at startup — a running
|
||||
// quarantine keeps its original cooldown regardless, and a new cooldown only applies to a
|
||||
// quarantine that starts after a restart.
|
||||
if (!Objects.equals(old.quarantineCooldownSeconds(), fresh.quarantineCooldownSeconds())) {
|
||||
changed.add("quarantineCooldownSeconds");
|
||||
}
|
||||
Map<String, BridgedConfig.Profile> before =
|
||||
old.profiles() == null ? Map.of() : old.profiles();
|
||||
Map<String, BridgedConfig.Profile> after =
|
||||
@@ -246,8 +261,10 @@ public final class ConfigRef implements Supplier<BridgedConfig> {
|
||||
|
||||
/**
|
||||
* Whether two versions of a profile would launch a peer identically. Compares every component
|
||||
* the launcher reads at spawn; {@code weight} and {@code maxLoad} are excluded because those are
|
||||
* read live by the placement policy and really do take effect on the next spawn.
|
||||
* the launcher reads at spawn; {@code weight}, {@code maxLoad} and {@code credentialId} are
|
||||
* excluded because those are read live (by the placement policy and, for credentialId, by
|
||||
* {@code CompositePeerLauncher}/the CB-578 stage B exhaustion sink) and really do take effect on
|
||||
* the next spawn.
|
||||
*/
|
||||
private static boolean sameLaunchSettings(BridgedConfig.Profile a, BridgedConfig.Profile b) {
|
||||
return Objects.equals(a.baseUrl(), b.baseUrl())
|
||||
@@ -265,6 +282,10 @@ public final class ConfigRef implements Supplier<BridgedConfig> {
|
||||
&& Objects.equals(a.gitHostEnv(), b.gitHostEnv())
|
||||
&& Objects.equals(a.kind(), b.kind())
|
||||
&& Objects.equals(a.env(), b.env())
|
||||
&& Objects.equals(a.subscription(), b.subscription());
|
||||
&& Objects.equals(a.subscription(), b.subscription())
|
||||
// CB-578 stage B: exhaustedPattern is compiled once into Bridged.main's pattern map
|
||||
// at startup (see ExhaustedPatternLookup wiring) — a reload never re-reads it, so a
|
||||
// changed pattern must be reported as deferred, exactly like model/baseUrl/argv.
|
||||
&& Objects.equals(a.exhaustedPattern(), b.exhaustedPattern());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -12,34 +12,50 @@ import java.util.HashMap;
|
||||
import java.util.HashSet;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.function.BiConsumer;
|
||||
import java.util.function.LongSupplier;
|
||||
import java.util.function.Supplier;
|
||||
|
||||
/** Slow whole-fleet evidence collection. It is deliberately separate from the delivery poller. */
|
||||
public final class FleetHealthMonitor {
|
||||
private static final Logger log = LoggerFactory.getLogger(FleetHealthMonitor.class);
|
||||
|
||||
/** Bounded attempts to run {@link #failTarget} for one transition. Never retried tick-to-tick (CB-580). */
|
||||
static final int MAX_FAIL_TARGET_ATTEMPTS = 3;
|
||||
|
||||
private final AgentControl agents;
|
||||
private final Supplier<List<MemberSession>> roster;
|
||||
private final MessageService messages;
|
||||
private final ScheduledExecutorService scheduler;
|
||||
private final LongSupplier clock;
|
||||
private final long intervalSeconds;
|
||||
private final BiConsumer<String, String> failTarget;
|
||||
private final Map<String, HealthPrior> priors = new HashMap<>();
|
||||
private final Map<String, HealthState> states = new HashMap<>();
|
||||
|
||||
// These facts need the evidence publishers introduced by later M4 units. They are not negatives.
|
||||
private static final boolean NOT_YET_OBSERVED = false;
|
||||
|
||||
/**
|
||||
* @param failTarget CB-568's idempotent target-wide failure operation (e.g. {@code messages::abandon}),
|
||||
* invoked once when a member transitions into a terminal health state. Required —
|
||||
* there is deliberately no defaulting overload; a caller that does not want the
|
||||
* fail-tickets-on-terminal-health behavior must pass an explicit inert value (see
|
||||
* {@code TestTurnTokens.inert} / {@code BridgeMcp.CapacitySource.none()} for the pattern).
|
||||
*/
|
||||
public FleetHealthMonitor(AgentControl agents, Supplier<List<MemberSession>> roster, MessageService messages,
|
||||
ScheduledExecutorService scheduler, LongSupplier clock, long intervalSeconds) {
|
||||
ScheduledExecutorService scheduler, LongSupplier clock, long intervalSeconds,
|
||||
BiConsumer<String, String> failTarget) {
|
||||
this.agents = agents;
|
||||
this.roster = roster;
|
||||
this.messages = messages;
|
||||
this.scheduler = scheduler;
|
||||
this.clock = clock;
|
||||
this.intervalSeconds = intervalSeconds;
|
||||
this.failTarget = Objects.requireNonNull(failTarget, "failTarget");
|
||||
}
|
||||
|
||||
/** Pure per-member decision seam. */
|
||||
@@ -91,6 +107,34 @@ public final class FleetHealthMonitor {
|
||||
} else if (previous != null && fault(previous)) {
|
||||
log.info("fleet health member={} recovered state={} previous={}", target, next, previous);
|
||||
}
|
||||
// CB-580: a member entering GONE/NEVER_READY must not leave its waiting tickets pending
|
||||
// forever. Fire exactly once per transition — never on a tick where the state is unchanged,
|
||||
// which is what made the rejected commit call abandon() once per tick for as long as a
|
||||
// member stayed terminal.
|
||||
if (terminal(next)) {
|
||||
failTerminalTarget(target, next);
|
||||
}
|
||||
}
|
||||
|
||||
private void failTerminalTarget(String target, HealthState state) {
|
||||
String reason = "fleet health: member reached terminal state " + state.name();
|
||||
RuntimeException last = null;
|
||||
for (int attempt = 1; attempt <= MAX_FAIL_TARGET_ATTEMPTS; attempt++) {
|
||||
try {
|
||||
failTarget.accept(target, reason);
|
||||
return;
|
||||
} catch (RuntimeException error) {
|
||||
last = error;
|
||||
log.warn("fleet health: failTarget attempt {}/{} failed for member={} state={}",
|
||||
attempt, MAX_FAIL_TARGET_ATTEMPTS, target, state, error);
|
||||
}
|
||||
}
|
||||
log.warn("fleet health: giving up on failTarget for member={} state={} after {} attempts",
|
||||
target, state, MAX_FAIL_TARGET_ATTEMPTS, last);
|
||||
}
|
||||
|
||||
private static boolean terminal(HealthState state) {
|
||||
return state == HealthState.GONE || state == HealthState.NEVER_READY;
|
||||
}
|
||||
|
||||
private static boolean fault(HealthState state) {
|
||||
|
||||
@@ -6,6 +6,7 @@ import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.Collections;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.Locale;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.function.LongSupplier;
|
||||
@@ -23,6 +24,15 @@ import java.util.function.Supplier;
|
||||
* by first starting the session and asking it. Scanning closes that loop: label the tab, and the
|
||||
* pane is recognised on the next resolve.
|
||||
*
|
||||
* <p><strong>CB-579 — matched by name, not prefix.</strong> This used to strip one shared
|
||||
* {@code tabPrefix} off a label to derive the lead's name, and merged a config-supplied
|
||||
* {@code terminal_id} pin over every scan result so the pin could never expire. Both are gone: each
|
||||
* lead now configures its own exact {@code tab} label ({@code fleet.leaders.<name>.tab}), so this
|
||||
* class is handed a {@code tab → name} map up front and matches labels against it exactly
|
||||
* (case-insensitively). There is no merge step — a scan result is the whole answer. That is the
|
||||
* fix for the bug this replaces: a {@code terminal_id} pin surviving in config after the pane it
|
||||
* named was gone, so the daemon kept treating a dead session as a live lead forever.
|
||||
*
|
||||
* <p><strong>Direction of trust.</strong> The label names the lead; it never <em>grants</em>
|
||||
* anything a pane could take for itself. Three properties keep that honest:
|
||||
* <ol>
|
||||
@@ -48,7 +58,7 @@ import java.util.function.Supplier;
|
||||
* ever make a decision that <em>removes</em> something based on this map, add the same check.
|
||||
* The remaining hazard is an <em>operator</em> one — a worker {@code tabLabel} template that
|
||||
* happens to start with the same prefix would promote the whole fleet — and that is refused at
|
||||
* startup by {@code BridgedConfig.validateLeadScan} rather than documented here.
|
||||
* startup by {@code BridgedConfig.validateLeadTabPrefixes} rather than documented here.
|
||||
*
|
||||
* <p><strong>Caching.</strong> {@link #get()} is on the request path (every resolve), so the scan
|
||||
* is TTL-cached and a stale-but-valid map is preferred to a herdr round-trip. A failed scan keeps
|
||||
@@ -60,38 +70,46 @@ public final class LeadTabScanner implements Supplier<Map<String, String>> {
|
||||
private static final Logger log = LoggerFactory.getLogger(LeadTabScanner.class);
|
||||
|
||||
private final HerdrClient herdr;
|
||||
private final String tabPrefix;
|
||||
private final Map<String, String> tabToName;
|
||||
private final Set<String> excludedWorkspaceLabels;
|
||||
private final Map<String, String> configuredLeads;
|
||||
private final long ttlNanos;
|
||||
private final LongSupplier clock;
|
||||
|
||||
private Map<String, String> cached;
|
||||
private Map<String, String> cached = Map.of();
|
||||
private long scannedAtNanos;
|
||||
private boolean everScanned;
|
||||
|
||||
/**
|
||||
* @param herdr the herdr client to query ({@code workspace.list},
|
||||
* {@code tab.list}, {@code pane.list} — all read-only)
|
||||
* @param tabPrefix a tab whose label starts with this (case-insensitively) hosts a
|
||||
* lead; the rest of the label, trimmed, is the lead's name
|
||||
* @param tabToName every configured lead's exact tab label → its name
|
||||
* ({@code fleet.leaders.<name>.tab}), matched case-insensitively
|
||||
* @param excludedWorkspaceLabels workspaces never scanned — the configured worker spaces
|
||||
* @param configuredLeads the static {@code leaders:}/{@code primary:} registry, merged
|
||||
* over every scan result. Explicit config outranks the
|
||||
* convention, and survives a scan that cannot run at all
|
||||
* @param ttlNanos how long a scan result is reused before the next one
|
||||
* @param clock nanosecond time source ({@code System::nanoTime} in production)
|
||||
*/
|
||||
public LeadTabScanner(HerdrClient herdr, String tabPrefix, Set<String> excludedWorkspaceLabels,
|
||||
Map<String, String> configuredLeads, long ttlNanos, LongSupplier clock) {
|
||||
public LeadTabScanner(HerdrClient herdr, Map<String, String> tabToName,
|
||||
Set<String> excludedWorkspaceLabels, long ttlNanos, LongSupplier clock) {
|
||||
this.herdr = herdr;
|
||||
this.tabPrefix = tabPrefix == null || tabPrefix.isBlank() ? "lead:" : tabPrefix.strip();
|
||||
this.tabToName = normalize(tabToName);
|
||||
this.excludedWorkspaceLabels = excludedWorkspaceLabels == null
|
||||
? Set.of() : Set.copyOf(excludedWorkspaceLabels);
|
||||
this.configuredLeads = configuredLeads == null ? Map.of() : Map.copyOf(configuredLeads);
|
||||
this.ttlNanos = ttlNanos;
|
||||
this.clock = clock;
|
||||
this.cached = this.configuredLeads;
|
||||
}
|
||||
|
||||
/** Keys stripped and lower-cased once, so every lookup is a plain map hit. */
|
||||
private static Map<String, String> normalize(Map<String, String> tabToName) {
|
||||
if (tabToName == null || tabToName.isEmpty()) {
|
||||
return Map.of();
|
||||
}
|
||||
Map<String, String> out = new LinkedHashMap<>();
|
||||
tabToName.forEach((tab, name) -> {
|
||||
if (tab != null && !tab.isBlank() && name != null && !name.isBlank()) {
|
||||
out.put(tab.strip().toLowerCase(Locale.ROOT), name);
|
||||
}
|
||||
});
|
||||
return Collections.unmodifiableMap(out);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -151,25 +169,20 @@ public final class LeadTabScanner implements Supplier<Map<String, String>> {
|
||||
}
|
||||
}
|
||||
}
|
||||
byTerminal.putAll(configuredLeads); // an explicit pin outranks a label
|
||||
return Collections.unmodifiableMap(byTerminal);
|
||||
}
|
||||
|
||||
/**
|
||||
* The lead name a tab label declares, or {@code null} if it declares none.
|
||||
* The lead name a tab label declares, or {@code null} if it names none of the configured leads.
|
||||
*
|
||||
* <p>{@code "lead: opus-5.0"} → {@code "opus-5.0"}. A bare {@code "lead:"} names nobody and is
|
||||
* rejected: an unnamed lead would resolve as {@code PRIMARY} with nothing to attribute it to.
|
||||
* <p>Exact match (case-insensitive, ends stripped) against {@link #tabToName} — no prefix
|
||||
* stripping, so an operator's {@code "lead: something-else"} tab is never mistaken for a
|
||||
* configured lead just because it shares a prefix.
|
||||
*/
|
||||
private String leadNameOf(String label) {
|
||||
if (label == null) {
|
||||
return null;
|
||||
}
|
||||
String l = label.strip();
|
||||
if (!l.regionMatches(true, 0, tabPrefix, 0, tabPrefix.length())) {
|
||||
return null;
|
||||
}
|
||||
String name = l.substring(tabPrefix.length()).strip();
|
||||
return name.isEmpty() ? null : name;
|
||||
return tabToName.get(label.strip().toLowerCase(Locale.ROOT));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -6,8 +6,13 @@ import dev.ltms.bridged.msg.TurnToken;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Objects;
|
||||
import java.util.Set;
|
||||
import java.util.TreeSet;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.regex.Pattern;
|
||||
|
||||
/**
|
||||
* The CB-106 completion fallback: bridges the {@link Injector}'s turn-completion signal to the
|
||||
@@ -61,6 +66,8 @@ public final class CompletionResolver implements TurnListener {
|
||||
|
||||
private final AgentControl agents;
|
||||
private final Rendezvous rendezvous;
|
||||
private final ExhaustedPatternLookup exhaustedPatterns;
|
||||
private final ExhaustionSink exhaustionSink;
|
||||
|
||||
/**
|
||||
* Per-target record of the turn currently in flight: the exact {@link Rendezvous} waiter its
|
||||
@@ -78,9 +85,22 @@ public final class CompletionResolver implements TurnListener {
|
||||
|
||||
private final ConcurrentHashMap<String, InFlight> inFlight = new ConcurrentHashMap<>();
|
||||
|
||||
public CompletionResolver(AgentControl agents, Rendezvous rendezvous) {
|
||||
/**
|
||||
* @param exhaustedPatterns CB-578 stage A: per-target lookup for a profile's configured
|
||||
* usage-limit refusal pattern. Required — there is deliberately no
|
||||
* defaulting overload; a caller that does not want the classification
|
||||
* must pass an explicit inert value ({@link ExhaustedPatternLookup#none()}).
|
||||
* @param exhaustionSink CB-578 stage B: notified when a {@code BACKEND_EXHAUSTED}
|
||||
* classification actually resolves a waiter. Required for the same
|
||||
* reason as {@code exhaustedPatterns} — pass {@link ExhaustionSink#none()}
|
||||
* to opt out.
|
||||
*/
|
||||
public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns,
|
||||
ExhaustionSink exhaustionSink) {
|
||||
this.agents = agents;
|
||||
this.rendezvous = rendezvous;
|
||||
this.exhaustedPatterns = Objects.requireNonNull(exhaustedPatterns, "exhaustedPatterns");
|
||||
this.exhaustionSink = Objects.requireNonNull(exhaustionSink, "exhaustionSink");
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -157,11 +177,12 @@ public final class CompletionResolver implements TurnListener {
|
||||
return;
|
||||
}
|
||||
String tail;
|
||||
String assistantBlock = null;
|
||||
int originalLength = 0;
|
||||
boolean clipped = false;
|
||||
boolean scrapeFailed = false;
|
||||
try {
|
||||
String assistantBlock = lastAssistantBlock(agents.read(target, SCRAPE_SOURCE));
|
||||
assistantBlock = lastAssistantBlock(agents.read(target, SCRAPE_SOURCE));
|
||||
originalLength = assistantBlock.strip().length();
|
||||
clipped = originalLength > MAX_SCRAPE_CHARS;
|
||||
tail = clip(assistantBlock);
|
||||
@@ -184,6 +205,25 @@ public final class CompletionResolver implements TurnListener {
|
||||
target);
|
||||
return; // keep the in-flight record: a later genuine completion still needs it
|
||||
}
|
||||
// CB-578 stage A: a turn that ended with no bridge_reply AND whose scrape matches the
|
||||
// backend's configured usage-limit pattern is a refusal, not an answer. Classify it as
|
||||
// BACKEND_EXHAUSTED rather than handing the caller a scrape that reads like a real reply.
|
||||
if (!scrapeFailed) {
|
||||
Pattern exhausted = exhaustedPatterns.patternFor(target);
|
||||
String matchedLine = exhausted == null ? null : firstMatchingLine(assistantBlock, exhausted);
|
||||
if (matchedLine != null) {
|
||||
String reason = "backend exhausted (usage limit): " + matchedLine;
|
||||
if (rendezvous.resolveExhausted(waiter, reason)) {
|
||||
inFlight.remove(target, turn);
|
||||
log.warn("completion for {} classified BACKEND_EXHAUSTED (no bridge_reply; scrape "
|
||||
+ "matched the profile's exhausted pattern): {}", target, reason);
|
||||
// CB-578 stage B: only on the resolution that actually won the race — a late
|
||||
// duplicate must never quarantine a credential twice for one refusal.
|
||||
exhaustionSink.onExhausted(target, reason);
|
||||
}
|
||||
return;
|
||||
}
|
||||
}
|
||||
String completion = clipped ? tail + "\n" + CLIPPED_PANE_TAIL_MARKER : tail;
|
||||
if (rendezvous.resolveCompletion(waiter, completion)) {
|
||||
inFlight.remove(target, turn);
|
||||
@@ -232,6 +272,45 @@ public final class CompletionResolver implements TurnListener {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* The first line of {@code text} matching {@code pattern}, stripped — the CB-578 stage A
|
||||
* evidence carried in a {@code BACKEND_EXHAUSTED} reason so the operator sees the real refusal
|
||||
* text, never a generic label. {@code null} if no line matches.
|
||||
*/
|
||||
static String firstMatchingLine(String text, Pattern pattern) {
|
||||
if (text == null || text.isEmpty()) return null;
|
||||
for (String line : text.split("\n", -1)) {
|
||||
if (pattern.matcher(line).find()) {
|
||||
return line.strip();
|
||||
}
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
/**
|
||||
* Coverage summary for the CB-578 stage A exhausted-pattern classification, logged at startup
|
||||
* the way {@link dev.ltms.bridged.health.FleetHealthMonitor#coverage} is — so an operator can
|
||||
* see whether the classification is on, and for which profiles, without reading every
|
||||
* profile's config by hand.
|
||||
*
|
||||
* @param allProfiles every configured profile name
|
||||
* @param configuredProfiles the subset of {@code allProfiles} that carry an exhausted pattern
|
||||
*/
|
||||
public static String coverage(Set<String> allProfiles, Set<String> configuredProfiles) {
|
||||
if (configuredProfiles.isEmpty()) {
|
||||
return "off (no profile has an exhaustedPattern configured; profiles: " + sorted(allProfiles) + ")";
|
||||
}
|
||||
Set<String> unconfigured = new TreeSet<>(allProfiles);
|
||||
unconfigured.removeAll(configuredProfiles);
|
||||
return unconfigured.isEmpty()
|
||||
? "full (all profiles configured: " + sorted(allProfiles) + ")"
|
||||
: "partial (configured: " + sorted(configuredProfiles) + "; not configured: " + sorted(unconfigured) + ")";
|
||||
}
|
||||
|
||||
private static List<String> sorted(Set<String> names) {
|
||||
return names.stream().sorted().toList();
|
||||
}
|
||||
|
||||
private static String clip(String s) {
|
||||
if (s == null) return "";
|
||||
String trimmed = s.strip();
|
||||
|
||||
@@ -0,0 +1,27 @@
|
||||
package dev.ltms.bridged.inject;
|
||||
|
||||
import java.util.regex.Pattern;
|
||||
|
||||
/**
|
||||
* Per-target lookup for a profile's configured usage-limit refusal pattern (CB-578 stage A): how
|
||||
* {@link CompletionResolver} tells a backend that refused on a subscription usage limit — the
|
||||
* worker's pane stays healthy, but the account is exhausted — apart from a genuine completion.
|
||||
*
|
||||
* <p>The pattern is always profile config, never a vendor string in Java source: every backend
|
||||
* words its refusal differently, so a hardcoded sentence would only ever match one of them.
|
||||
*/
|
||||
@FunctionalInterface
|
||||
public interface ExhaustedPatternLookup {
|
||||
|
||||
/** The compiled pattern configured for {@code target}'s profile, or {@code null} if none. */
|
||||
Pattern patternFor(String target);
|
||||
|
||||
/**
|
||||
* Inert lookup — no profile has a pattern configured, so the classification never fires and
|
||||
* the completion fallback behaves exactly as before CB-578 stage A. The explicit stand-in a
|
||||
* caller (or a test not exercising this feature) passes instead of a defaulting overload.
|
||||
*/
|
||||
static ExhaustedPatternLookup none() {
|
||||
return target -> null;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,31 @@
|
||||
package dev.ltms.bridged.inject;
|
||||
|
||||
/**
|
||||
* Notified when {@link CompletionResolver} actually delivers a {@code BACKEND_EXHAUSTED}
|
||||
* classification to a waiting send (CB-578 stage B) — never on a race that lost (see
|
||||
* {@link CompletionResolver#resolve}, which only calls this after
|
||||
* {@code Rendezvous.resolveExhausted} returns {@code true}).
|
||||
*
|
||||
* <p>{@link CompletionResolver} knows only {@code target} (a herdr terminal id); it has no notion of
|
||||
* profiles or credentials, so mapping {@code target} to whatever should be quarantined is entirely
|
||||
* the sink's job — see {@code Bridged.main}'s wiring, which resolves target → session → profile →
|
||||
* {@code effectiveCredentialId()} and calls {@code BackendQuarantine.quarantine} on it.
|
||||
*/
|
||||
@FunctionalInterface
|
||||
public interface ExhaustionSink {
|
||||
|
||||
/**
|
||||
* @param target the herdr terminal id whose turn was classified {@code BACKEND_EXHAUSTED}
|
||||
* @param reason the matched-line reason carried by the classification
|
||||
*/
|
||||
void onExhausted(String target, String reason);
|
||||
|
||||
/**
|
||||
* Inert sink — nothing happens on exhaustion. The explicit stand-in a caller (or a test not
|
||||
* exercising this feature) passes instead of a defaulting overload, exactly like
|
||||
* {@link ExhaustedPatternLookup#none()}.
|
||||
*/
|
||||
static ExhaustionSink none() {
|
||||
return (target, reason) -> { };
|
||||
}
|
||||
}
|
||||
@@ -59,7 +59,7 @@ public final class LeadLauncher {
|
||||
/**
|
||||
* @param agents herdr agent control (start, list)
|
||||
* @param spaces workspace / tab control (ensure, create, label, list)
|
||||
* @param cfg the loaded config — {@code fleet.leaders}, {@code profiles} and the lead pins
|
||||
* @param cfg the loaded config — {@code fleet.leaders}, {@code profiles} and each lead's tab
|
||||
*/
|
||||
public LeadLauncher(AgentControl agents, WorkspaceControl spaces, BridgedConfig cfg) {
|
||||
this.agents = agents;
|
||||
@@ -102,8 +102,8 @@ public final class LeadLauncher {
|
||||
continue;
|
||||
}
|
||||
if (!lead.isCreatable()) {
|
||||
// A lead with a `terminal:` pin and no `profile:` is recognise-only by design: the
|
||||
// operator opens it by hand. Say so once rather than looking like a silent failure.
|
||||
// A lead with a `tab:` but no `profile:` is recognise-only by design: the operator
|
||||
// opens it by hand. Say so once rather than looking like a silent failure.
|
||||
log.info("lead '{}' is not live, and names no profile — it can be recognised but not "
|
||||
+ "launched. Add `profile:` under fleet.leaders.{} to have bridged start it.",
|
||||
name, name);
|
||||
@@ -127,18 +127,16 @@ public final class LeadLauncher {
|
||||
}
|
||||
|
||||
/**
|
||||
* How many live leads exist per configured name.
|
||||
* How many live leads exist per configured name: a running agent in a tab labelled with that
|
||||
* lead's exact {@code tab} (CB-579). Member workspaces are excluded, exactly as the scanner
|
||||
* excludes them: a member must not be counted as a lead because it happens to sit in a matching
|
||||
* tab.
|
||||
*
|
||||
* <p>Two independent pieces of evidence, because either alone double-spawns:
|
||||
* <ul>
|
||||
* <li>a running agent in a tab labelled {@code "<tabPrefix> <name>"} — how an auto-launched
|
||||
* lead, or an operator following the labelling convention, is found;</li>
|
||||
* <li>a running agent on a terminal the config pins in {@code fleet.leaders.<name>.terminal} —
|
||||
* how a lead the operator opened and pinned by hand is found. Without this, a pinned lead
|
||||
* whose tab carries no matching label would be relaunched on every boot.</li>
|
||||
* </ul>
|
||||
* Member workspaces are excluded, exactly as the scanner excludes them: a member must not be
|
||||
* counted as a lead because it happens to sit in a matching tab.
|
||||
* <p>There used to be a second path here — a running agent on the terminal a
|
||||
* {@code fleet.leaders.<name>.terminal} pin named, for a lead opened and pinned by hand. That
|
||||
* pin is retired: {@code tab} is now the only field identity depends on, and {@link Agent}
|
||||
* already carries {@link Agent#tabId()} directly, so a hand-opened lead is found the same way an
|
||||
* auto-launched one is — by labelling its tab to match.
|
||||
*/
|
||||
private Map<String, Integer> liveLeads(Map<String, BridgedConfig.Leader> leaders) {
|
||||
Set<String> memberSpaces = cfg.profiles().values().stream()
|
||||
@@ -160,20 +158,9 @@ public final class LeadLauncher {
|
||||
}
|
||||
}
|
||||
|
||||
// terminalId → the lead name the config pins it to.
|
||||
Map<String, String> nameByPinnedTerminal = new LinkedHashMap<>();
|
||||
leaders.forEach((name, lead) -> {
|
||||
if (lead.terminal() != null && !lead.terminal().isBlank()) {
|
||||
nameByPinnedTerminal.put(lead.terminal().strip(), name);
|
||||
}
|
||||
});
|
||||
|
||||
Map<String, Integer> counts = new LinkedHashMap<>();
|
||||
for (Agent a : agents.list()) {
|
||||
String name = nameByTab.get(a.tabId());
|
||||
if (name == null) {
|
||||
name = nameByPinnedTerminal.get(a.terminalId());
|
||||
}
|
||||
if (name != null) {
|
||||
counts.merge(name, 1, Integer::sum);
|
||||
}
|
||||
@@ -184,7 +171,7 @@ public final class LeadLauncher {
|
||||
/**
|
||||
* The configured lead a tab label names, or {@code null} for a label that names none.
|
||||
*
|
||||
* <p>Matched against the declared lead names rather than by splitting on the prefix, so an
|
||||
* <p>Matched exactly (case-insensitively) against each lead's configured {@code tab}, so an
|
||||
* operator's {@code "lead: something-else"} tab is not mistaken for a configured lead.
|
||||
*/
|
||||
private String leadNameOf(String label, Map<String, BridgedConfig.Leader> leaders) {
|
||||
@@ -193,7 +180,8 @@ public final class LeadLauncher {
|
||||
}
|
||||
String l = label.strip();
|
||||
for (Map.Entry<String, BridgedConfig.Leader> e : leaders.entrySet()) {
|
||||
if (l.equalsIgnoreCase(e.getValue().tabLabel(e.getKey()).strip())) {
|
||||
String tab = e.getValue().tabLabel();
|
||||
if (tab != null && l.equalsIgnoreCase(tab.strip())) {
|
||||
return e.getKey();
|
||||
}
|
||||
}
|
||||
@@ -202,7 +190,7 @@ public final class LeadLauncher {
|
||||
|
||||
/** Start one lead. Returns false (having logged) rather than throwing on any failure. */
|
||||
private boolean launch(String name, BridgedConfig.Leader lead, BridgedConfig.Profile profile) {
|
||||
String label = lead.tabLabel(name);
|
||||
String label = lead.tabLabel();
|
||||
String cwd = (lead.cwd() == null || lead.cwd().isBlank())
|
||||
? System.getProperty("user.dir") : lead.cwd();
|
||||
|
||||
|
||||
@@ -14,6 +14,7 @@ import dev.ltms.bridged.herdr.HerdrException;
|
||||
import dev.ltms.bridged.msg.MessageService;
|
||||
import dev.ltms.bridged.msg.Rendezvous;
|
||||
import dev.ltms.bridged.peer.PeerUnreachableException;
|
||||
import dev.ltms.bridged.placement.BackendQuarantine;
|
||||
import dev.ltms.bridged.session.SessionManager;
|
||||
import dev.ltms.bridged.session.MemberSession;
|
||||
import dev.ltms.bridged.session.WorktreeRequest;
|
||||
@@ -33,6 +34,7 @@ import jakarta.servlet.http.HttpServlet;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
import java.util.Set;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.LongSupplier;
|
||||
@@ -78,6 +80,7 @@ public final class BridgeMcp {
|
||||
private final Metrics metrics; // CB-502: null → auth failures not counted
|
||||
private final CapacitySource capacity;
|
||||
private final HealthCoverageSource healthCoverage;
|
||||
private final QuarantineSource quarantine;
|
||||
|
||||
/** Capacity facts used by {@code bridge_list}; production must supply the placement live count. */
|
||||
public record CapacitySource(Function<String, Integer> liveCount, Function<String, Integer> maxLoad,
|
||||
@@ -90,17 +93,30 @@ public final class BridgeMcp {
|
||||
/** Coverage is supplied by the health wiring, not inferred from a missing dependency. */
|
||||
public record HealthCoverageSource(Supplier<String> value) { }
|
||||
|
||||
/**
|
||||
* CB-578 stage B quarantine facts used by {@code bridge_profiles}: a profile → credential id
|
||||
* lookup, plus the shared {@link BackendQuarantine} to read remaining cooldowns off.
|
||||
*/
|
||||
public record QuarantineSource(Function<String, String> credentialIdFor, BackendQuarantine quarantine) {
|
||||
/** Inert source — no profile is ever reported quarantined. Explicit stand-in, not a default. */
|
||||
public static QuarantineSource none() { return new QuarantineSource(_ -> null, BackendQuarantine.none()); }
|
||||
}
|
||||
|
||||
/**
|
||||
* @param callers resolves each call's {@link Principal}; {@code null} disables authorization.
|
||||
* This surface needs its own enforcement: {@code /mcp} is a raw servlet on
|
||||
* Jetty's context handler and never passes through Javalin's {@code before}
|
||||
* filter, so the REST guard does not cover it.
|
||||
* @param metrics registry for auth-failure counting; may be {@code null}
|
||||
* @param quarantine CB-578 stage B facts for {@code bridge_profiles}; required — pass
|
||||
* {@link QuarantineSource#none()} for a caller that does not want the feature
|
||||
*/
|
||||
public BridgeMcp(MessageService messages, PeerLauncher workers, SessionManager sessions,
|
||||
ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry,
|
||||
CallerResolver callers, Metrics metrics, CapacitySource capacity, HealthCoverageSource healthCoverage) {
|
||||
CallerResolver callers, Metrics metrics, CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
QuarantineSource quarantine) {
|
||||
this.capacity = capacity;
|
||||
this.quarantine = Objects.requireNonNull(quarantine, "quarantine");
|
||||
this.healthCoverage = healthCoverage;
|
||||
McpJsonMapper json = new JacksonMcpJsonMapperSupplier().get();
|
||||
this.transport = HttpServletStreamableServerTransportProvider.builder()
|
||||
@@ -228,7 +244,7 @@ public final class BridgeMcp {
|
||||
.toolCall(profilesTool(), (exchange, _) -> {
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
|
||||
if (denied != null) return denied;
|
||||
return profiles(workers);
|
||||
return profiles(workers, quarantine);
|
||||
})
|
||||
.toolCall(whoamiTool(), (exchange, _) -> {
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
|
||||
@@ -458,6 +474,11 @@ public final class BridgeMcp {
|
||||
"[worker finished without a structured bridge_reply — transcript tail follows]\n" + r.text());
|
||||
// The worker ran the turn then wedged (CB-109) — surface the error context.
|
||||
case WORKER_FAILED -> text("[worker failed — turn ended in an unrecoverable state]\n" + r.text());
|
||||
// The backend refused on a subscription usage limit (CB-578 stage A) — the worker's
|
||||
// pane stayed healthy, but its account is exhausted. Distinct from WORKER_FAILED so the
|
||||
// primary gets the real cause, not a generic wedge.
|
||||
case BACKEND_EXHAUSTED -> text("[backend exhausted — the worker's account refused on a "
|
||||
+ "usage limit]\n" + r.text());
|
||||
// The worker paused mid-turn to ask (CB-205) — tell the primary how to answer in-turn.
|
||||
case QUESTION -> text("[question] the worker paused to ask before it can finish:\n" + r.text()
|
||||
+ "\n\nAnswer it by calling bridge_send again with turnId=\"" + r.turnId()
|
||||
@@ -701,11 +722,33 @@ public final class BridgeMcp {
|
||||
return null;
|
||||
}
|
||||
|
||||
/** {@code bridge_profiles}: the configured worker profiles and the default. */
|
||||
static McpSchema.CallToolResult profiles(PeerLauncher workers) {
|
||||
return text(json(Map.of(
|
||||
"profiles", workers.profiles(),
|
||||
"default", workers.defaultProfile() == null ? "" : workers.defaultProfile())));
|
||||
/**
|
||||
* {@code bridge_profiles}: the configured worker profiles, the default, and — CB-578 stage B —
|
||||
* which of them are currently quarantined (backend exhausted) and for how much longer. The
|
||||
* {@code quarantined} key is present only when at least one profile is, so a fleet where nothing
|
||||
* has ever been quarantined gets exactly the pre-stage-B shape.
|
||||
*/
|
||||
static McpSchema.CallToolResult profiles(PeerLauncher workers, QuarantineSource quarantine) {
|
||||
Map<String, Object> result = new LinkedHashMap<>();
|
||||
result.put("profiles", workers.profiles());
|
||||
result.put("default", workers.defaultProfile() == null ? "" : workers.defaultProfile());
|
||||
Map<String, Object> quarantined = new LinkedHashMap<>();
|
||||
for (String profile : workers.profiles()) {
|
||||
String credentialId = quarantine.credentialIdFor().apply(profile);
|
||||
if (credentialId == null) {
|
||||
continue;
|
||||
}
|
||||
quarantine.quarantine().remainingSeconds(credentialId).ifPresent(remaining -> {
|
||||
Map<String, Object> row = new LinkedHashMap<>();
|
||||
row.put("credentialId", credentialId);
|
||||
row.put("quarantinedForSeconds", remaining);
|
||||
quarantined.put(profile, row);
|
||||
});
|
||||
}
|
||||
if (!quarantined.isEmpty()) {
|
||||
result.put("quarantined", quarantined);
|
||||
}
|
||||
return text(json(result));
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -935,7 +978,10 @@ public final class BridgeMcp {
|
||||
|
||||
private static McpSchema.Tool profilesTool() {
|
||||
return tool("bridge_profiles",
|
||||
"List the configured worker profiles (backends) and which one bridge_spawn uses by default.",
|
||||
"List the configured worker profiles (backends) and which one bridge_spawn uses by "
|
||||
+ "default. A 'quarantined' map is present when a backend-exhausted refusal put "
|
||||
+ "a profile's credential on cooldown — bridge_spawn onto it is refused until "
|
||||
+ "quarantinedForSeconds elapses; a profile sharing that credential is listed too.",
|
||||
objectSchema(Map.of(), List.of()));
|
||||
}
|
||||
|
||||
|
||||
@@ -8,6 +8,7 @@ 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.BackendQuarantine;
|
||||
import dev.ltms.bridged.placement.PlacementCandidate;
|
||||
import dev.ltms.bridged.placement.PlacementContext;
|
||||
import dev.ltms.bridged.placement.PlacementException;
|
||||
@@ -23,10 +24,12 @@ import java.util.HashSet;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.Supplier;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
/**
|
||||
* The {@link PeerLauncher} the core actually holds when more than one adapter is configured — a thin
|
||||
@@ -78,6 +81,9 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
private final Supplier<Map<String, BridgedConfig.Profile>> profileConfigs;
|
||||
private final Supplier<PlacementPolicy> placementPolicy;
|
||||
|
||||
/** CB-578 stage B: credential cooldown, checked before an explicit spawn and filtered into placement. */
|
||||
private final BackendQuarantine quarantine;
|
||||
|
||||
/**
|
||||
* CB-557: the role pools an unqualified spawn draws its candidates from. A supplier that yields
|
||||
* {@code null}, and an empty pool for a role, both fall back to every configured profile — the
|
||||
@@ -99,7 +105,10 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
}
|
||||
|
||||
/**
|
||||
* Production constructor with a placement policy and live-worker counter.
|
||||
* Production constructor with a placement policy and live-worker counter. Quarantine (CB-578
|
||||
* stage B) is off for this constructor — {@link BackendQuarantine#none()} — since it predates
|
||||
* the feature and existing callers of this exact overload never exercised it; use the 7-arg
|
||||
* overload below to wire a real {@link BackendQuarantine}.
|
||||
*
|
||||
* @param delegates one adapter per configured peer kind; must be non-empty and declare
|
||||
* disjoint profile-name sets
|
||||
@@ -113,13 +122,14 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
Map<String, BridgedConfig.Profile> profileConfigs,
|
||||
PlacementPolicy placementPolicy,
|
||||
Function<String, Integer> liveCount) {
|
||||
this(delegates, defaultProfile, profileConfigs, placementPolicy, liveCount, null);
|
||||
this(delegates, defaultProfile, profileConfigs, placementPolicy, liveCount, null, BackendQuarantine.none());
|
||||
}
|
||||
|
||||
/**
|
||||
* Production constructor with role pools (CB-557). An unqualified spawn draws its candidates from
|
||||
* {@code fleet.<role>} instead of from every configured profile, so a reviewer is placed on a
|
||||
* reviewer backend and never on, say, the architect-only one.
|
||||
* reviewer backend and never on, say, the architect-only one. Quarantine is off for this
|
||||
* constructor too, for the same reason as the 5-arg overload above.
|
||||
*
|
||||
* @param fleet the configured role pools; {@code null} ⇒ every profile is a candidate for every
|
||||
* role, which is the pre-CB-557 behaviour
|
||||
@@ -130,29 +140,50 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
PlacementPolicy placementPolicy,
|
||||
Function<String, Integer> liveCount,
|
||||
BridgedConfig.Fleet fleet) {
|
||||
this(delegates, defaultProfile, profileConfigs, placementPolicy, liveCount, fleet, BackendQuarantine.none());
|
||||
}
|
||||
|
||||
/**
|
||||
* Production constructor with role pools and quarantine (CB-578 stage B). The full-featured
|
||||
* non-reloading form; {@link #CompositePeerLauncher(List, String, Supplier, Function, BackendQuarantine)}
|
||||
* is what {@code Bridged.main} actually wires up.
|
||||
*
|
||||
* @param quarantine required — pass {@link BackendQuarantine#none()} for a caller that does not
|
||||
* want the feature, never a defaulting overload (CB-578 stage B's own rule).
|
||||
*/
|
||||
public CompositePeerLauncher(List<HerdrPeerLauncher> delegates,
|
||||
String defaultProfile,
|
||||
Map<String, BridgedConfig.Profile> profileConfigs,
|
||||
PlacementPolicy placementPolicy,
|
||||
Function<String, Integer> liveCount,
|
||||
BridgedConfig.Fleet fleet,
|
||||
BackendQuarantine quarantine) {
|
||||
// LinkedHashMap, not Map.copyOf: candidates() promises definition order and the weighted
|
||||
// policy breaks exact-weight ties on it, so a salted iteration order would make placement
|
||||
// differ from one JVM run to the next.
|
||||
this(delegates, defaultProfile,
|
||||
constant(Collections.unmodifiableMap(new LinkedHashMap<>(profileConfigs))),
|
||||
constant(placementPolicy), liveCount, constant(fleet));
|
||||
constant(placementPolicy), liveCount, constant(fleet), quarantine);
|
||||
}
|
||||
|
||||
/**
|
||||
* Production constructor that re-reads its placement inputs per spawn (CB-559), so a config
|
||||
* reload retargets the next member without a restart.
|
||||
*
|
||||
* @param config the live configuration — read at every spawn, never captured
|
||||
* @param config the live configuration — read at every spawn, never captured
|
||||
* @param quarantine required — CB-578 stage B; pass {@link BackendQuarantine#none()} to opt out
|
||||
*/
|
||||
public CompositePeerLauncher(List<HerdrPeerLauncher> delegates,
|
||||
String defaultProfile,
|
||||
Supplier<BridgedConfig> config,
|
||||
Function<String, Integer> liveCount) {
|
||||
Function<String, Integer> liveCount,
|
||||
BackendQuarantine quarantine) {
|
||||
this(delegates, defaultProfile,
|
||||
() -> config.get().profiles(),
|
||||
() -> PlacementPolicies.fromName(config.get().placement()),
|
||||
liveCount,
|
||||
() -> config.get().fleet());
|
||||
() -> config.get().fleet(),
|
||||
quarantine);
|
||||
}
|
||||
|
||||
/** The all-suppliers form every other constructor funnels into. */
|
||||
@@ -161,8 +192,10 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
Supplier<Map<String, BridgedConfig.Profile>> profileConfigs,
|
||||
Supplier<PlacementPolicy> placementPolicy,
|
||||
Function<String, Integer> liveCount,
|
||||
Supplier<BridgedConfig.Fleet> fleet) {
|
||||
Supplier<BridgedConfig.Fleet> fleet,
|
||||
BackendQuarantine quarantine) {
|
||||
this.fleet = fleet;
|
||||
this.quarantine = Objects.requireNonNull(quarantine, "quarantine");
|
||||
if (delegates.isEmpty()) {
|
||||
throw new IllegalArgumentException("at least one peer adapter must be configured");
|
||||
}
|
||||
@@ -225,6 +258,7 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
// the charter makes explicit-profile spawns the normal path — so skipping the check
|
||||
// here would leave the cap dead config in real operation.
|
||||
HerdrPeerLauncher d = route(requestedProfile);
|
||||
enforceNotQuarantined(requestedProfile);
|
||||
enforceMaxLoad(requestedProfile);
|
||||
PeerHandle handle = d.spawn(req);
|
||||
spawnedBy.put(handle.id(), d);
|
||||
@@ -238,7 +272,10 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
List<PlacementCandidate> candidates = candidates(req.role());
|
||||
String roleDefault = defaultProfileFor(req.role());
|
||||
Set<String> unreachable = new HashSet<>();
|
||||
PlacementContext ctx = new PlacementContext(roleDefault, candidates, liveCount, unreachable);
|
||||
// CB-578 stage B: computed once up front — a quarantine's expiry cannot pass within one spawn
|
||||
// call, so re-deriving it per retry would only cost work, never change the answer.
|
||||
Set<String> quarantined = quarantinedProfiles(candidates);
|
||||
PlacementContext ctx = new PlacementContext(roleDefault, candidates, liveCount, unreachable, quarantined);
|
||||
|
||||
int maxAttempts = candidates.isEmpty() ? 1 : candidates.size();
|
||||
for (int attempt = 0; attempt < maxAttempts; attempt++) {
|
||||
@@ -268,7 +305,7 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
chosen.profile(), e.getMessage());
|
||||
unreachable.add(chosen.profile());
|
||||
// Update the context for the next selection so the policy excludes this profile.
|
||||
ctx = new PlacementContext(roleDefault, candidates, liveCount, unreachable);
|
||||
ctx = new PlacementContext(roleDefault, candidates, liveCount, unreachable, quarantined);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -298,6 +335,42 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
* @param profile the profile the caller explicitly named
|
||||
* @throws PlacementException when the profile is at capacity
|
||||
*/
|
||||
/**
|
||||
* Refuse an explicit-profile spawn whose credential is quarantined (CB-578 stage B): a prior
|
||||
* {@code BACKEND_EXHAUSTED} classification on this profile, or on another profile sharing its
|
||||
* {@code credentialId}, is still on cooldown.
|
||||
*
|
||||
* <p>Checked before {@link #enforceMaxLoad}, and for the same reason that check exists: an
|
||||
* explicit profile bypasses the placement policy's filtering entirely, so without this the cap
|
||||
* (there, quarantine here) would be dead config on the very path the charter calls normal.
|
||||
* Deliberately no fallback to another profile, matching {@link #enforceMaxLoad}'s own reasoning —
|
||||
* the caller named this profile for a cost/model reason.
|
||||
*
|
||||
* @throws PlacementException naming the profile, its credential, and the remaining cooldown
|
||||
*/
|
||||
private void enforceNotQuarantined(String profile) {
|
||||
String credentialId = credentialIdFor(profile);
|
||||
quarantine.remainingSeconds(credentialId).ifPresent(remaining -> {
|
||||
throw new PlacementException("worker profile '" + profile + "' is quarantined "
|
||||
+ "(credential '" + credentialId + "' exhausted; ~" + remaining
|
||||
+ "s remaining) — refusing spawn");
|
||||
});
|
||||
}
|
||||
|
||||
/** {@code profile}'s credential group (CB-578 stage B), or the profile's own name if unconfigured. */
|
||||
private String credentialIdFor(String profile) {
|
||||
BridgedConfig.Profile cfg = profiles0().get(profile);
|
||||
return cfg == null ? profile : cfg.effectiveCredentialId();
|
||||
}
|
||||
|
||||
/** The subset of {@code candidates} whose credential is currently quarantined (CB-578 stage B). */
|
||||
private Set<String> quarantinedProfiles(List<PlacementCandidate> candidates) {
|
||||
return candidates.stream()
|
||||
.map(PlacementCandidate::profile)
|
||||
.filter(p -> quarantine.isQuarantined(credentialIdFor(p)))
|
||||
.collect(Collectors.toSet());
|
||||
}
|
||||
|
||||
private void enforceMaxLoad(String profile) {
|
||||
// Absent config, or a config whose maxLoad normalized to null (non-positive ⇒ unlimited at
|
||||
// load), means no cap — never cap what wasn't configured.
|
||||
|
||||
@@ -7,6 +7,7 @@ import dev.ltms.bridged.herdr.HerdrException;
|
||||
import dev.ltms.bridged.herdr.Tab;
|
||||
import dev.ltms.bridged.herdr.Workspace;
|
||||
import dev.ltms.bridged.herdr.WorkspaceControl;
|
||||
import dev.ltms.bridged.peer.CharterReceipt;
|
||||
import dev.ltms.bridged.peer.MemberRole;
|
||||
import dev.ltms.bridged.peer.PeerHandle;
|
||||
import dev.ltms.bridged.peer.PeerLauncher;
|
||||
@@ -269,8 +270,11 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
|
||||
// --- spawn ---------------------------------------------------------------------------------
|
||||
|
||||
/** A started peer plus the launch's agent-session id (the resume handle, or null). */
|
||||
private record Spawned(Agent agent, String agentSessionId) {
|
||||
/**
|
||||
* A started peer plus the launch's agent-session id (the resume handle, or null) and the
|
||||
* charter receipt (CB-571) the base composed for it.
|
||||
*/
|
||||
private record Spawned(Agent agent, String agentSessionId, CharterReceipt receipt) {
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -305,12 +309,41 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
String replyCharter = cfg.hasMcp() ? REPLY_CHARTER : null;
|
||||
String charter = roleCharter == null ? replyCharter
|
||||
: replyCharter == null ? roleCharter : roleCharter + "\n\n" + replyCharter;
|
||||
Launch launch = buildLaunch(cfg, new LaunchSpec(sessionName, resumeSessionId, role, charter));
|
||||
String cwd = resolveCwd(requestedCwd, cfg, callerCwd);
|
||||
Agent agent = cfg.tabPlacement()
|
||||
? spawnInTab(cfg, launch.env(), launch.argv(), cwd, role, liveFleet)
|
||||
: spawnAsPane(cfg, launch.env(), launch.argv(), cwd);
|
||||
return new Spawned(agent, launch.agentSessionId());
|
||||
// CB-571: fingerprint the exact composed charter bytes once, here in the base, before the
|
||||
// string leaves for an adapter — so Claude and OpenCode derive the same digest. A failed
|
||||
// start has no bridge_spawn result and no roster row, so the failure log below is the only
|
||||
// surface the byte count can appear on. The charter text itself is never logged.
|
||||
CharterReceipt receipt = CharterReceipt.compose(role, cfg.profile(), roleCharter, charter);
|
||||
try {
|
||||
Launch launch = buildLaunch(cfg, new LaunchSpec(sessionName, resumeSessionId, role, charter));
|
||||
String cwd = resolveCwd(requestedCwd, cfg, callerCwd);
|
||||
Agent agent = cfg.tabPlacement()
|
||||
? spawnInTab(cfg, launch.env(), launch.argv(), cwd, role, liveFleet)
|
||||
: spawnAsPane(cfg, launch.env(), launch.argv(), cwd, charter);
|
||||
logCharterReceipt(receipt, true);
|
||||
return new Spawned(agent, launch.agentSessionId(), receipt);
|
||||
} catch (RuntimeException e) {
|
||||
logCharterReceipt(receipt, false);
|
||||
throw e;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* The one place the charter's size and digest appear in the logs. {@code success} true after a
|
||||
* start, false from the failure path of {@link #spawnInternal} where no handle or roster row
|
||||
* exists to carry the receipt. Always metadata only — never the charter text.
|
||||
*/
|
||||
private static void logCharterReceipt(CharterReceipt receipt, boolean success) {
|
||||
String role = receipt.role() == null ? "" : receipt.role().wireName();
|
||||
if (success) {
|
||||
log.info("spawned role={} profile={} charterSource={} charterSha256={} charterBytes={}",
|
||||
role, receipt.profile(), receipt.charterSource(),
|
||||
receipt.charterSha256(), receipt.charterBytes());
|
||||
} else {
|
||||
log.warn("spawn failed; charter role={} profile={} charterSource={} charterSha256={} charterBytes={}",
|
||||
role, receipt.profile(), receipt.charterSource(),
|
||||
receipt.charterSha256(), receipt.charterBytes());
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -349,7 +382,7 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
String id = UUID.randomUUID().toString();
|
||||
paneByAgentId.put(id, paneId);
|
||||
return new WorkerHandle(id, agent.terminalId(), requireProfile(req.profileName()).profile(),
|
||||
req.sessionName(), spawned.agentSessionId());
|
||||
req.sessionName(), spawned.agentSessionId(), spawned.receipt());
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -433,11 +466,19 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
}
|
||||
}
|
||||
|
||||
/** Legacy placement: split the currently-focused tab; the peer still starts in {@code cwd}. */
|
||||
/**
|
||||
* Legacy placement: split the currently-focused tab; the peer still starts in {@code cwd}.
|
||||
*
|
||||
* <p>CB-571: this is the one legacy log that printed the full argv, and the charter travels
|
||||
* inside argv — so the charter text went to the daemon log on every pane-placement spawn. The
|
||||
* {@code spawnInTab} path never logs argv, so only this site is fixed. {@code charter} is the
|
||||
* composed charter, if any; its argv element is replaced by its digest so the log still shows
|
||||
* which args were passed without exposing the charter prose.
|
||||
*/
|
||||
private Agent spawnAsPane(BridgedConfig.Profile cfg, Map<String, String> workerEnv,
|
||||
List<String> argv, String cwd) {
|
||||
List<String> argv, String cwd, String charter) {
|
||||
log.info("spawning {} (pane placement) profile={} cwd={} argv={}",
|
||||
namePrefix, cfg.profile(), cwd, argv);
|
||||
namePrefix, cfg.profile(), cwd, redactCharter(argv, charter));
|
||||
String paneId = spaces.splitPane(cwd, workerEnv);
|
||||
if (paneId == null) {
|
||||
throw new IllegalStateException("pane.split returned no pane — cannot start a peer");
|
||||
@@ -447,6 +488,21 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
return peer;
|
||||
}
|
||||
|
||||
/**
|
||||
* A copy of {@code argv} with an element equal to {@code charter} replaced by its digest, so
|
||||
* the pane log never prints the charter prose. The charter is handed to an adapter as one argv
|
||||
* element, so exact-equality is the right match; every other argument passes through unchanged.
|
||||
*/
|
||||
private static List<String> redactCharter(List<String> argv, String charter) {
|
||||
if (charter == null || charter.isBlank() || argv == null || argv.isEmpty()) {
|
||||
return argv;
|
||||
}
|
||||
String digest = CharterReceipt.digestOf(charter);
|
||||
return argv.stream()
|
||||
.map(a -> a.equals(charter) ? "<charter sha256=" + digest + ">" : a)
|
||||
.toList();
|
||||
}
|
||||
|
||||
/** A started peer together with the sequence its unique name/label used. */
|
||||
private record Started(Agent agent, long seq) {
|
||||
}
|
||||
@@ -651,11 +707,18 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
|
||||
/**
|
||||
* A concrete {@link PeerHandle} wrapping herdr agent coordinates, the profile that spawned it,
|
||||
* and the session identity the launch resolved (CB-547a): the bridge's logical name and the
|
||||
* peer's own session id, both null when the spawn carried no identity.
|
||||
* the session identity the launch resolved (CB-547a): the bridge's logical name and the peer's
|
||||
* own session id, both null when the spawn carried no identity — and the charter receipt
|
||||
* (CB-571) the base computed for this launch.
|
||||
*/
|
||||
private record WorkerHandle(String id, String terminalId, String profile,
|
||||
String sessionName, String agentSessionId) implements PeerHandle {
|
||||
String sessionName, String agentSessionId,
|
||||
CharterReceipt receipt) implements PeerHandle {
|
||||
|
||||
@Override
|
||||
public CharterReceipt charterReceipt() {
|
||||
return receipt;
|
||||
}
|
||||
}
|
||||
|
||||
// --- shared helpers ------------------------------------------------------------------------
|
||||
|
||||
@@ -7,6 +7,7 @@ import dev.ltms.bridged.herdr.Agent;
|
||||
import dev.ltms.bridged.herdr.AgentControl;
|
||||
import dev.ltms.bridged.herdr.WorkspaceControl;
|
||||
import dev.ltms.bridged.peer.Capability;
|
||||
import dev.ltms.bridged.peer.CharterReceipt;
|
||||
import dev.ltms.bridged.peer.PeerHandle;
|
||||
import dev.ltms.bridged.peer.SpawnRequest;
|
||||
|
||||
@@ -406,6 +407,11 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
// appeared).
|
||||
return discovery.sessionIdForDirectory(cwd);
|
||||
}
|
||||
|
||||
@Override
|
||||
public CharterReceipt charterReceipt() {
|
||||
return delegate.charterReceipt();
|
||||
}
|
||||
}
|
||||
|
||||
// --- Agent-returning convenience spawns (used by callers/tests that want the herdr Agent) ---
|
||||
|
||||
@@ -69,6 +69,14 @@ public final class MessageService {
|
||||
* failure context (e.g. the error screen). Terminal, but not a successful completion.
|
||||
*/
|
||||
WORKER_FAILED,
|
||||
/**
|
||||
* The turn finished without a {@code bridge_reply} and the scrape matched the backend's
|
||||
* configured usage-limit refusal pattern (CB-578 stage A); {@code text} is the reason,
|
||||
* carrying the matched line. The worker's pane is healthy — only its account is refusing —
|
||||
* so this is never reported as a completed reply, and is kept distinct from
|
||||
* {@link #WORKER_FAILED} (a wedged worker) and a session simply going {@code GONE}.
|
||||
*/
|
||||
BACKEND_EXHAUSTED,
|
||||
/**
|
||||
* The worker paused mid-turn to ask the primary a question (CB-205); {@code text} is the
|
||||
* question and {@code turnId} correlates the answer. Not terminal — the primary answers with
|
||||
@@ -280,6 +288,7 @@ public final class MessageService {
|
||||
case COMPLETED_UNREPLIED -> "completion_fallback";
|
||||
case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, BUSY -> "timeout";
|
||||
case WORKER_FAILED -> "failed";
|
||||
case BACKEND_EXHAUSTED -> "backend_exhausted";
|
||||
case STALE_TURN, QUESTION -> null; // not a completed delegation
|
||||
};
|
||||
}
|
||||
@@ -591,9 +600,11 @@ public final class MessageService {
|
||||
String source = r.outcome() == Outcome.REPLIED ? "reply" : "transcript";
|
||||
return new TaskView(ticket, Phase.DONE, r.text(), source, null, null);
|
||||
}
|
||||
// A wedged worker (CB-109) carries the error context as its reason; the timeout/busy
|
||||
// outcomes carry none, so fall back to the outcome name.
|
||||
String detail = r.outcome() == Outcome.WORKER_FAILED && r.text() != null
|
||||
// A wedged worker (CB-109) or a backend-exhausted classification (CB-578 stage A) carries
|
||||
// the real cause as its reason; the timeout/busy outcomes carry none, so fall back to the
|
||||
// outcome name.
|
||||
boolean carriesReason = r.outcome() == Outcome.WORKER_FAILED || r.outcome() == Outcome.BACKEND_EXHAUSTED;
|
||||
String detail = carriesReason && r.text() != null
|
||||
? r.text()
|
||||
: "no reply — " + r.outcome().name().toLowerCase();
|
||||
return new TaskView(ticket, Phase.FAILED, null, null, detail, null);
|
||||
@@ -669,6 +680,7 @@ public final class MessageService {
|
||||
case REPLY -> Outcome.REPLIED;
|
||||
case COMPLETION -> Outcome.COMPLETED_UNREPLIED;
|
||||
case FAILED -> Outcome.WORKER_FAILED;
|
||||
case BACKEND_EXHAUSTED -> Outcome.BACKEND_EXHAUSTED;
|
||||
case QUESTION -> Outcome.QUESTION;
|
||||
};
|
||||
}
|
||||
|
||||
@@ -33,6 +33,13 @@ public final class Rendezvous {
|
||||
COMPLETION,
|
||||
/** The worker ran the turn then wedged (CB-109); {@code text} is the failure context. */
|
||||
FAILED,
|
||||
/**
|
||||
* The turn finished without a {@code bridge_reply}, and the scrape matched the backend's
|
||||
* configured usage-limit refusal pattern (CB-578 stage A); {@code text} is the reason,
|
||||
* carrying the matched line. The pane is healthy — only the account is refusing — so this
|
||||
* is kept separate from a session simply going {@code GONE}.
|
||||
*/
|
||||
BACKEND_EXHAUSTED,
|
||||
/**
|
||||
* The worker paused mid-turn to ask the primary a question (CB-205 reverse rendezvous);
|
||||
* {@code text} is the question and {@code turnId} correlates the primary's answer back to
|
||||
@@ -224,6 +231,19 @@ public final class Rendezvous {
|
||||
return waiter != null && waiter.complete(new Resolution(Kind.FAILED, reason));
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve a specific captured {@code waiter} as {@link Kind#BACKEND_EXHAUSTED} (CB-578 stage A):
|
||||
* the turn finished with no {@code bridge_reply} and the scrape matched the backend's configured
|
||||
* usage-limit pattern; {@code reason} carries the matched line. Like
|
||||
* {@link #resolveCompletion(CompletableFuture, String)} it targets the exact captured send
|
||||
* (CB-116). A no-op if that waiter was already resolved — first resolution wins.
|
||||
*
|
||||
* @return {@code true} if this call resolved the waiter, {@code false} if it was null or already resolved
|
||||
*/
|
||||
public boolean resolveExhausted(CompletableFuture<Resolution> waiter, String reason) {
|
||||
return waiter != null && waiter.complete(new Resolution(Kind.BACKEND_EXHAUSTED, reason));
|
||||
}
|
||||
|
||||
private boolean complete(String session, Resolution resolution) {
|
||||
CompletableFuture<Resolution> waiter = waiters.get(session);
|
||||
return waiter != null && waiter.complete(resolution);
|
||||
|
||||
@@ -0,0 +1,84 @@
|
||||
package dev.ltms.bridged.peer;
|
||||
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.security.MessageDigest;
|
||||
import java.security.NoSuchAlgorithmException;
|
||||
import java.util.HexFormat;
|
||||
|
||||
/**
|
||||
* CB-571: a fingerprint of the exact charter bytes handed to a spawned member.
|
||||
*
|
||||
* <p>Lets an operator prove <em>which</em> charter a member actually got, without ever logging the
|
||||
* charter text. The digest covers the exact composed UTF-8 string {@code HerdrPeerLauncher} passes
|
||||
* to its adapter as {@code LaunchSpec.charter()}, so every adapter that receives the same string —
|
||||
* Claude inlining it, OpenCode writing it to a file — produces the same digest for the same config.
|
||||
* Two spawns of the same role from the same config agree; editing the charter changes the digest.
|
||||
*
|
||||
* <p>Deliberately places no charter prose. A charter is operator-authored text that may name
|
||||
* internal projects or unreleased plans, and logs get tailed, shipped, and pasted into tickets.
|
||||
* The {@code charterSource} key is what the operator wants to confirm, and it carries no content.
|
||||
*/
|
||||
public record CharterReceipt(
|
||||
MemberRole role,
|
||||
String profile,
|
||||
String charterSource,
|
||||
String charterSha256,
|
||||
int charterBytes) {
|
||||
|
||||
/** Source reported when the role has no configured charter, so the field is never omitted. */
|
||||
public static final String NO_SOURCE = "none";
|
||||
|
||||
/**
|
||||
* The config key that supplied the role's charter text, e.g. {@code fleet.charters.architect}.
|
||||
*/
|
||||
public static String sourceKey(MemberRole role) {
|
||||
return "fleet.charters." + (role == null ? "?" : role.wireName());
|
||||
}
|
||||
|
||||
/**
|
||||
* Fingerprint the composed charter for {@code role} on {@code profile}. {@code configured} is
|
||||
* the role's charter text as read from config ({@code null} when none is configured);
|
||||
* {@code composed} is the exact string the launcher will pass to the adapter — the reply
|
||||
* charter may be appended to {@code configured}, or stand alone when no role charter exists.
|
||||
*
|
||||
* <p>No composed charter at all is reported as an explicit absence — a {@code null} digest and
|
||||
* a zero byte count — never a digest of the empty string, which would hide the fact that no
|
||||
* text was supplied. {@code configured} being {@code null} while {@code composed} is the reply
|
||||
* charter alone is a normal case, and the source says so.
|
||||
*/
|
||||
public static CharterReceipt compose(MemberRole role, String profile,
|
||||
String configured, String composed) {
|
||||
String source = (configured == null || configured.isBlank())
|
||||
? NO_SOURCE : sourceKey(role);
|
||||
if (composed == null) {
|
||||
return new CharterReceipt(role, profile, source, null, 0);
|
||||
}
|
||||
byte[] bytes = composed.getBytes(StandardCharsets.UTF_8);
|
||||
return new CharterReceipt(role, profile, source, digestOf(composed), bytes.length);
|
||||
}
|
||||
|
||||
/** Whether the composed charter was absent (no text was given to the member). */
|
||||
public boolean absent() {
|
||||
return charterSha256 == null;
|
||||
}
|
||||
|
||||
/**
|
||||
* The stable SHA-256 hex digest of {@code text}, or {@code null} for null/blank text. Used both
|
||||
* for the receipt's fingerprint and to redact a charter argument in a spawn log.
|
||||
*/
|
||||
public static String digestOf(String text) {
|
||||
if (text == null || text.isBlank()) {
|
||||
return null;
|
||||
}
|
||||
return sha256Hex(text.getBytes(StandardCharsets.UTF_8));
|
||||
}
|
||||
|
||||
private static String sha256Hex(byte[] bytes) {
|
||||
try {
|
||||
MessageDigest md = MessageDigest.getInstance("SHA-256");
|
||||
return HexFormat.of().formatHex(md.digest(bytes));
|
||||
} catch (NoSuchAlgorithmException e) {
|
||||
throw new IllegalStateException("SHA-256 is unavailable", e);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -67,4 +67,18 @@ public interface PeerHandle {
|
||||
default String agentSessionId() {
|
||||
return null;
|
||||
}
|
||||
|
||||
/**
|
||||
* The charter receipt (CB-571) for this peer's launch — the fingerprint of the exact charter
|
||||
* bytes it was started with. {@code null} when the launcher records none (a non-instrumented
|
||||
* adapter, or a launcher before this field); the session registry stores it so the spawn result
|
||||
* and the roster row can show an operator which charter a member actually got.
|
||||
*
|
||||
* <p>Deliberately not a {@code default}: a decorator that forgets to override this silently
|
||||
* answers {@code null} for a question it has no basis to answer, and the gap surfaces only as
|
||||
* a missing roster field, not a compile error. Every implementation must answer explicitly.
|
||||
*
|
||||
* @return the fingerprint, or {@code null} when the launcher carries none
|
||||
*/
|
||||
CharterReceipt charterReceipt();
|
||||
}
|
||||
|
||||
@@ -0,0 +1,99 @@
|
||||
package dev.ltms.bridged.placement;
|
||||
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
import java.util.OptionalLong;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
/**
|
||||
* Where a credential (not a profile — see {@code BridgedConfig.Profile#effectiveCredentialId()})
|
||||
* sits out a cooldown after a {@code BACKEND_EXHAUSTED} classification (CB-578 stage B), so a fresh
|
||||
* spawn does not walk straight back onto the account that just refused on a usage limit.
|
||||
*
|
||||
* <p>Keyed by credential id, never by profile name: two profiles sharing one credential (e.g. two
|
||||
* models on the same OpenAI account) share one quarantine — {@link #quarantine} one credential id
|
||||
* and every profile whose {@code effectiveCredentialId()} equals it is quarantined too, without this
|
||||
* class knowing anything about profiles at all. That mapping is the caller's job (see
|
||||
* {@code CompositePeerLauncher} and {@code dev.ltms.bridged.inject.ExhaustionSink}).
|
||||
*
|
||||
* <p>The clock is injected ({@link LongSupplier}, conventionally {@code System::nanoTime} like
|
||||
* {@code FleetHealthMonitor}), never read inline, so a quarantine's expiry is testable without a
|
||||
* real sleep.
|
||||
*/
|
||||
public final class BackendQuarantine {
|
||||
|
||||
private final ConcurrentHashMap<String, Long> quarantinedUntilNanos = new ConcurrentHashMap<>();
|
||||
private final LongSupplier nowNanos;
|
||||
private final long cooldownNanos;
|
||||
|
||||
/**
|
||||
* @param nowNanos monotonic clock, injected for testability
|
||||
* @param cooldownNanos how long a fresh {@link #quarantine} call blocks the credential for;
|
||||
* must be positive
|
||||
*/
|
||||
public BackendQuarantine(LongSupplier nowNanos, long cooldownNanos) {
|
||||
this.nowNanos = Objects.requireNonNull(nowNanos, "nowNanos");
|
||||
if (cooldownNanos <= 0) {
|
||||
throw new IllegalArgumentException("cooldownNanos must be positive: " + cooldownNanos);
|
||||
}
|
||||
this.cooldownNanos = cooldownNanos;
|
||||
}
|
||||
|
||||
/**
|
||||
* Inert quarantine — nothing is ever quarantined unless {@link #quarantine} is actually called on
|
||||
* this instance. The explicit stand-in a caller (or a test not exercising this feature) passes
|
||||
* instead of a defaulting overload, exactly like {@code ExhaustedPatternLookup.none()}.
|
||||
*/
|
||||
public static BackendQuarantine none() {
|
||||
return new BackendQuarantine(() -> 0L, 1);
|
||||
}
|
||||
|
||||
/**
|
||||
* Quarantine {@code credentialId} for the configured cooldown, starting now. A repeat call while
|
||||
* already quarantined restarts the cooldown at full length — a fresh refusal is fresh evidence the
|
||||
* account is still exhausted, not a reason to let an earlier, shorter wait stand.
|
||||
*/
|
||||
public void quarantine(String credentialId) {
|
||||
Objects.requireNonNull(credentialId, "credentialId");
|
||||
quarantinedUntilNanos.put(credentialId, nowNanos.getAsLong() + cooldownNanos);
|
||||
}
|
||||
|
||||
/** Whether {@code credentialId} is quarantined right now. */
|
||||
public boolean isQuarantined(String credentialId) {
|
||||
return remainingNanos(credentialId) > 0;
|
||||
}
|
||||
|
||||
/** Seconds left on {@code credentialId}'s quarantine, or empty when it is not quarantined. */
|
||||
public OptionalLong remainingSeconds(String credentialId) {
|
||||
long remaining = remainingNanos(credentialId);
|
||||
return remaining > 0 ? OptionalLong.of(toSecondsRoundedUp(remaining)) : OptionalLong.empty();
|
||||
}
|
||||
|
||||
/**
|
||||
* Every currently-quarantined credential id and its remaining seconds (CB-578 stage B fleet
|
||||
* reporting) — expired entries are never included. Not pruned from the backing map here: it stays
|
||||
* small (bounded by the number of distinct credentials ever exhausted) and a lazily-stale entry is
|
||||
* harmless, since every read already checks the deadline.
|
||||
*/
|
||||
public Map<String, Long> activeRemainingSeconds() {
|
||||
Map<String, Long> out = new LinkedHashMap<>();
|
||||
quarantinedUntilNanos.forEach((credentialId, deadline) -> {
|
||||
long remaining = deadline - nowNanos.getAsLong();
|
||||
if (remaining > 0) {
|
||||
out.put(credentialId, toSecondsRoundedUp(remaining));
|
||||
}
|
||||
});
|
||||
return out;
|
||||
}
|
||||
|
||||
private long remainingNanos(String credentialId) {
|
||||
Long deadline = quarantinedUntilNanos.get(credentialId);
|
||||
return deadline == null ? 0L : deadline - nowNanos.getAsLong();
|
||||
}
|
||||
|
||||
private static long toSecondsRoundedUp(long nanos) {
|
||||
return (nanos + 999_999_999L) / 1_000_000_000L;
|
||||
}
|
||||
}
|
||||
@@ -4,18 +4,28 @@ 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.
|
||||
*
|
||||
* <p>Quarantine (CB-578 stage B) is the one exception: a quarantined default is a credential that
|
||||
* just refused on a usage limit, not a transient capacity or reachability concern, so {@code fixed}
|
||||
* steps to the first non-quarantined candidate instead of walking straight back onto it. A fleet
|
||||
* where nothing is ever quarantined never exercises this path, so today's behaviour is unchanged.
|
||||
*/
|
||||
final class FixedPlacementPolicy implements PlacementPolicy {
|
||||
|
||||
@Override
|
||||
public PlacementCandidate select(PlacementContext ctx) {
|
||||
String d = ctx.defaultProfile();
|
||||
if (d != null && !d.isBlank()) {
|
||||
if (d != null && !d.isBlank() && !ctx.quarantined().contains(d)) {
|
||||
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());
|
||||
for (PlacementCandidate c : ctx.candidates()) {
|
||||
if (!ctx.quarantined().contains(c.profile())) {
|
||||
return new PlacementCandidate(c.profile(), null, c.weight(), c.maxLoad());
|
||||
}
|
||||
}
|
||||
if (d != null && !d.isBlank()) {
|
||||
throw new PlacementException("worker profile '" + d + "' is quarantined (backend "
|
||||
+ "exhausted) and no un-quarantined candidate is available");
|
||||
}
|
||||
throw new PlacementException("no worker profiles configured");
|
||||
}
|
||||
|
||||
@@ -11,9 +11,13 @@ import java.util.function.Function;
|
||||
* @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
|
||||
* @param quarantined profiles whose credential is currently quarantined (CB-578 stage B) — a
|
||||
* {@code BACKEND_EXHAUSTED} classification put it, or a profile it shares a
|
||||
* credential with, on cooldown. Filtered the same way as {@code unreachable}.
|
||||
*/
|
||||
public record PlacementContext(String defaultProfile,
|
||||
List<PlacementCandidate> candidates,
|
||||
Function<String, Integer> liveCount,
|
||||
Set<String> unreachable) {
|
||||
Set<String> unreachable,
|
||||
Set<String> quarantined) {
|
||||
}
|
||||
|
||||
@@ -12,13 +12,13 @@ final class PlacementPolicyUtil {
|
||||
}
|
||||
|
||||
/**
|
||||
* Candidates that are not known-unreachable and have not reached their maxLoad.
|
||||
* A {@code null} maxLoad means unlimited.
|
||||
* Candidates that are not known-unreachable, not quarantined (CB-578 stage B), 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())) {
|
||||
if (ctx.unreachable().contains(c.profile()) || ctx.quarantined().contains(c.profile())) {
|
||||
continue;
|
||||
}
|
||||
Integer cap = c.maxLoad();
|
||||
@@ -35,14 +35,17 @@ final class PlacementPolicyUtil {
|
||||
|
||||
/**
|
||||
* Build a clear exception describing why every candidate was dropped: all at capacity,
|
||||
* all unreachable, or a mix.
|
||||
* all unreachable, all quarantined, or a mix.
|
||||
*/
|
||||
static PlacementException emptyException(PlacementContext ctx) {
|
||||
int atCap = 0;
|
||||
int unreachable = 0;
|
||||
int quarantined = 0;
|
||||
for (PlacementCandidate c : ctx.candidates()) {
|
||||
Integer cap = c.maxLoad();
|
||||
if (ctx.unreachable().contains(c.profile())) {
|
||||
if (ctx.quarantined().contains(c.profile())) {
|
||||
quarantined++;
|
||||
} else if (ctx.unreachable().contains(c.profile())) {
|
||||
unreachable++;
|
||||
} else if (cap != null && ctx.liveCount().apply(c.profile()) >= cap) {
|
||||
atCap++;
|
||||
@@ -53,6 +56,9 @@ final class PlacementPolicyUtil {
|
||||
if (total == 0) {
|
||||
return new PlacementException("no worker profiles configured");
|
||||
}
|
||||
if (quarantined == total) {
|
||||
return new PlacementException("all worker profiles are quarantined (backend exhausted)");
|
||||
}
|
||||
if (atCap == total) {
|
||||
return new PlacementException("all worker profiles are at maxLoad");
|
||||
}
|
||||
@@ -60,6 +66,7 @@ final class PlacementPolicyUtil {
|
||||
return new PlacementException("all worker profiles are unreachable");
|
||||
}
|
||||
return new PlacementException("no worker profile available: " + atCap + " at maxLoad, "
|
||||
+ unreachable + " unreachable, " + (total - atCap - unreachable) + " remaining");
|
||||
+ unreachable + " unreachable, " + quarantined + " quarantined, "
|
||||
+ (total - atCap - unreachable - quarantined) + " remaining");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -403,9 +403,12 @@ public final class BridgedApp {
|
||||
case TIMED_OUT_QUEUED -> "queued";
|
||||
case BUSY -> "busy";
|
||||
case WORKER_FAILED -> "failed";
|
||||
case BACKEND_EXHAUSTED -> "backend_exhausted";
|
||||
default -> "done"; // unreachable (terminal outcomes handled above)
|
||||
},
|
||||
"detail", reply.outcome() == MessageService.Outcome.WORKER_FAILED && reply.text() != null
|
||||
"detail", (reply.outcome() == MessageService.Outcome.WORKER_FAILED
|
||||
|| reply.outcome() == MessageService.Outcome.BACKEND_EXHAUSTED)
|
||||
&& reply.text() != null
|
||||
? reply.text()
|
||||
: "no reply within " + timeout + "ms; poll status or retry"));
|
||||
}
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
package dev.ltms.bridged.session;
|
||||
|
||||
import dev.ltms.bridged.peer.CharterReceipt;
|
||||
import dev.ltms.bridged.peer.MemberRole;
|
||||
|
||||
/**
|
||||
@@ -23,6 +24,8 @@ import dev.ltms.bridged.peer.MemberRole;
|
||||
* @param lastActivityAtNanos {@link System#nanoTime()} of the most recent lifecycle event
|
||||
* @param turnCount number of delegated turns that have been delivered to this session
|
||||
* @param state current lifecycle state in the one-shot FSM
|
||||
* @param charterReceipt the fingerprint (CB-571) of the charter bytes this member was started
|
||||
* with; {@code null} for a session whose launcher recorded none
|
||||
*/
|
||||
public record MemberSession(
|
||||
String paneId,
|
||||
@@ -36,7 +39,8 @@ public record MemberSession(
|
||||
int turnCount,
|
||||
State state,
|
||||
String worktree,
|
||||
String branch) {
|
||||
String branch,
|
||||
CharterReceipt charterReceipt) {
|
||||
|
||||
/** One-shot worker lifecycle states. */
|
||||
public enum State {
|
||||
@@ -48,21 +52,34 @@ public record MemberSession(
|
||||
RELEASED
|
||||
}
|
||||
|
||||
/**
|
||||
* Backward-compatible shape: a session with no charter receipt (a test or a launcher before
|
||||
* CB-571). A separate constructor rather than a new parameter on the canonical one, so existing
|
||||
* call sites that have nothing to record keep compiling unchanged.
|
||||
*/
|
||||
public MemberSession(String paneId, String terminalId, String profile, MemberRole role,
|
||||
String cwd, String ownerTerminal, long spawnedAtNanos,
|
||||
long lastActivityAtNanos, int turnCount, State state,
|
||||
String worktree, String branch) {
|
||||
this(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
|
||||
lastActivityAtNanos, turnCount, state, worktree, branch, null);
|
||||
}
|
||||
|
||||
/** Return a copy of this session in {@code state}. */
|
||||
public MemberSession withState(State state) {
|
||||
return new MemberSession(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
|
||||
lastActivityAtNanos, turnCount, state, worktree, branch);
|
||||
lastActivityAtNanos, turnCount, state, worktree, branch, charterReceipt);
|
||||
}
|
||||
|
||||
/** Return a copy with {@code lastActivityAtNanos} updated to {@code nowNanos}. */
|
||||
public MemberSession withActivity(long nowNanos) {
|
||||
return new MemberSession(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
|
||||
nowNanos, turnCount, state, worktree, branch);
|
||||
nowNanos, turnCount, state, worktree, branch, charterReceipt);
|
||||
}
|
||||
|
||||
/** Return a copy with the turn count incremented and activity timestamped at {@code nowNanos}. */
|
||||
public MemberSession bumpTurn(long nowNanos) {
|
||||
return new MemberSession(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
|
||||
nowNanos, turnCount + 1, state, worktree, branch);
|
||||
nowNanos, turnCount + 1, state, worktree, branch, charterReceipt);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -163,7 +163,8 @@ public final class SessionManager implements TurnListener {
|
||||
0,
|
||||
MemberSession.State.SPAWNING,
|
||||
null,
|
||||
null);
|
||||
null,
|
||||
handle.charterReceipt());
|
||||
registry.put(handle.id(), session);
|
||||
memberLifecycle.acquired(session.role(), session.profile(), session.terminalId());
|
||||
log.debug("acquired session id={} terminal={} profile={} owner={}",
|
||||
@@ -355,7 +356,8 @@ public final class SessionManager implements TurnListener {
|
||||
0,
|
||||
MemberSession.State.SPAWNING,
|
||||
path,
|
||||
branch);
|
||||
branch,
|
||||
handle.charterReceipt());
|
||||
registry.put(handle.id(), session);
|
||||
memberLifecycle.acquired(session.role(), session.profile(), session.terminalId());
|
||||
log.debug("acquired worktree session id={} terminal={} profile={} branch={} path={}",
|
||||
@@ -421,6 +423,16 @@ public final class SessionManager implements TurnListener {
|
||||
if (session.ownerTerminal() != null) {
|
||||
m.put("owner", session.ownerTerminal());
|
||||
}
|
||||
// CB-571: which charter this member was started with — never the charter text itself. The
|
||||
// digest lets a lead tell at a glance whether all members got the same charter; the source
|
||||
// records whether a role charter was configured ("fleet.charters.<role>") or only the reply
|
||||
// charter was composed ("none").
|
||||
if (session.charterReceipt() != null) {
|
||||
m.put("charterSource", session.charterReceipt().charterSource());
|
||||
if (session.charterReceipt().charterSha256() != null) {
|
||||
m.put("charterSha256", session.charterReceipt().charterSha256());
|
||||
}
|
||||
}
|
||||
m.put("liveStatus", live == null ? "unknown" : live.status().name().toLowerCase());
|
||||
return m;
|
||||
}
|
||||
|
||||
@@ -194,7 +194,7 @@ class BridgedConfigTest {
|
||||
fleet:
|
||||
leaders:
|
||||
opus:
|
||||
terminal: term_opus
|
||||
tab: "lead: opus"
|
||||
""");
|
||||
|
||||
BridgedConfig.Leader lead = BridgedConfig.load(f).fleet().leaders().get("opus");
|
||||
@@ -212,7 +212,7 @@ class BridgedConfigTest {
|
||||
fleet:
|
||||
leaders:
|
||||
opus:
|
||||
terminal: term_opus
|
||||
tab: "drive: opus"
|
||||
tabPrefix: "drive:"
|
||||
scanIntervalSeconds: 30
|
||||
""");
|
||||
@@ -224,7 +224,9 @@ class BridgedConfigTest {
|
||||
|
||||
/**
|
||||
* The pane no longer has to exist before the daemon does (CB-557): a lead naming a profile may
|
||||
* be launched, while one that names only a terminal is recognised and never created.
|
||||
* be launched, while one that names no profile is recognised and never created. Either way it
|
||||
* still needs its own {@code tab:} (CB-579) — that part is unconditional, see
|
||||
* {@link #aLeadWithNoTabRefusesToStart}.
|
||||
*/
|
||||
@Test
|
||||
void aLeadIsCreatableOnlyWhenItNamesAProfile(@TempDir Path dir) throws Exception {
|
||||
@@ -239,8 +241,9 @@ class BridgedConfigTest {
|
||||
leaders:
|
||||
launched:
|
||||
profile: opus
|
||||
tab: "lead: launched"
|
||||
pinned:
|
||||
terminal: term_opus
|
||||
tab: "lead: pinned"
|
||||
""");
|
||||
|
||||
var leaders = BridgedConfig.load(f).fleet().leaders();
|
||||
@@ -249,8 +252,12 @@ class BridgedConfigTest {
|
||||
"no profile to launch on ⇒ recognise-only, the pre-CB-557 behaviour");
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-579: {@code tab} is the only field a lead's identity depends on now, so it is required
|
||||
* whether the entry is creatable or recognise-only — without it the entry can never be found.
|
||||
*/
|
||||
@Test
|
||||
void aLeadThatCanBeNeitherFoundNorCreatedRefusesToStart(@TempDir Path dir) throws Exception {
|
||||
void aLeadWithNoTabRefusesToStart(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("useless-lead.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
@@ -264,6 +271,80 @@ class BridgedConfigTest {
|
||||
|
||||
IllegalStateException e = assertThrows(IllegalStateException.class, cfg::validateMembers);
|
||||
assertTrue(e.getMessage().contains("ghost"), "the message must name the useless entry");
|
||||
assertTrue(e.getMessage().contains("tab:"), "the message must say what is missing");
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-579 acceptance (2): a config still spelling {@code fleet.leaders.<name>.terminal} must fail
|
||||
* loudly at load, not be silently dropped by {@code Leader}'s {@code @JsonIgnoreProperties}.
|
||||
*/
|
||||
@Test
|
||||
void aLeaderTerminalKeyFailsLoadAndNamesTabAsTheReplacement(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("stale-terminal.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
fleet:
|
||||
leaders:
|
||||
opus:
|
||||
terminal: term_opus
|
||||
""");
|
||||
|
||||
IllegalStateException e =
|
||||
assertThrows(IllegalStateException.class, () -> BridgedConfig.load(f));
|
||||
assertTrue(e.getMessage().contains("opus"), "the message must name the offending entry");
|
||||
assertTrue(e.getMessage().contains("tab:"), "the message must name the replacement key");
|
||||
assertTrue(e.getMessage().contains("terminal"), "the message must name the retired key");
|
||||
}
|
||||
|
||||
/** The same refusal, and it must name every offending entry, not just the first. */
|
||||
@Test
|
||||
void everyLeaderStillUsingTerminalIsReportedAtOnce(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("stale-terminals.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
fleet:
|
||||
leaders:
|
||||
opus:
|
||||
terminal: term_opus
|
||||
sol:
|
||||
terminal: term_sol
|
||||
""");
|
||||
|
||||
IllegalStateException e =
|
||||
assertThrows(IllegalStateException.class, () -> BridgedConfig.load(f));
|
||||
assertTrue(e.getMessage().contains("opus"));
|
||||
assertTrue(e.getMessage().contains("sol"));
|
||||
}
|
||||
|
||||
/** A {@code terminal:} anywhere else in the document (not under a leader entry) is unaffected. */
|
||||
@Test
|
||||
void aTerminalKeyOutsideFleetLeadersIsNotRejected(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("primary-terminal-ok.yaml");
|
||||
Files.writeString(f, "bind:\n port: 8080\nprimary:\n terminal: term_fixed\n");
|
||||
|
||||
assertDoesNotThrow(() -> BridgedConfig.load(f));
|
||||
}
|
||||
|
||||
/** CB-579 acceptance (3): distinct `tab:` labels need no shared prefix — one scanner finds both. */
|
||||
@Test
|
||||
void twoLeadersWithDifferentTabsAreBothConfigured(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("two-tabs.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
fleet:
|
||||
leaders:
|
||||
opus:
|
||||
tab: "lead: opus"
|
||||
sol:
|
||||
tab: "captain: sol"
|
||||
""");
|
||||
|
||||
var leaders = BridgedConfig.load(f).fleet().leaders();
|
||||
assertEquals("lead: opus", leaders.get("opus").tab());
|
||||
assertEquals("captain: sol", leaders.get("sol").tab());
|
||||
}
|
||||
|
||||
// ── CB-551: the idle-lead heartbeat ─────────────────────────────────────────────────────────
|
||||
@@ -328,7 +409,7 @@ class BridgedConfigTest {
|
||||
fleet:
|
||||
leaders:
|
||||
opus:
|
||||
terminal: term_opus
|
||||
tab: "lead: opus"
|
||||
tabPrefix: "lead:"
|
||||
""");
|
||||
BridgedConfig cfg = BridgedConfig.load(f);
|
||||
@@ -349,7 +430,7 @@ class BridgedConfigTest {
|
||||
tabLabel: "lead: {role} {profile}"
|
||||
leaders:
|
||||
opus:
|
||||
terminal: term_opus
|
||||
tab: "lead: opus"
|
||||
""");
|
||||
BridgedConfig cfg = BridgedConfig.load(f);
|
||||
|
||||
@@ -374,7 +455,7 @@ class BridgedConfigTest {
|
||||
fleet:
|
||||
leaders:
|
||||
opus:
|
||||
terminal: term_opus
|
||||
tab: "lead: opus"
|
||||
""");
|
||||
|
||||
assertDoesNotThrow(() -> BridgedConfig.load(f).validateLeadTabPrefixes());
|
||||
@@ -400,7 +481,7 @@ class BridgedConfigTest {
|
||||
"a label that collides with a convention nobody reads is not a problem");
|
||||
}
|
||||
|
||||
// ── CB-530: the leaders registry ────────────────────────────────────────────────────────────
|
||||
// ── CB-530/CB-579: the leaders registry ─────────────────────────────────────────────────────
|
||||
|
||||
@Test
|
||||
void leadersBlockRegistersEveryPaneByName(@TempDir Path dir) throws Exception {
|
||||
@@ -411,10 +492,10 @@ class BridgedConfigTest {
|
||||
fleet:
|
||||
leaders:
|
||||
opus-5.0:
|
||||
terminal: term_opus
|
||||
tab: "lead: opus-5.0"
|
||||
kind: claude
|
||||
gpt-sol-5.6:
|
||||
terminal: term_sol
|
||||
tab: "lead: gpt-sol-5.6"
|
||||
kind: opencode
|
||||
model: openai/gpt-5.6-terra
|
||||
""");
|
||||
@@ -425,9 +506,9 @@ class BridgedConfigTest {
|
||||
assertEquals(Set.of("opus-5.0", "gpt-sol-5.6"), leaders.keySet());
|
||||
assertEquals("opencode", leaders.get("gpt-sol-5.6").kind());
|
||||
assertEquals("openai/gpt-5.6-terra", leaders.get("gpt-sol-5.6").model());
|
||||
// The whole point: BOTH panes resolve as leads, so neither is demoted to worker.
|
||||
assertEquals(Map.of("term_opus", "opus-5.0", "term_sol", "gpt-sol-5.6"),
|
||||
cfg.leaderTerminals());
|
||||
// Identity is the tab now (CB-579) — both entries carry their own, distinct label.
|
||||
assertEquals("lead: opus-5.0", leaders.get("opus-5.0").tab());
|
||||
assertEquals("lead: gpt-sol-5.6", leaders.get("gpt-sol-5.6").tab());
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -439,42 +520,6 @@ class BridgedConfigTest {
|
||||
"configs that never migrate must behave exactly as they did before CB-530");
|
||||
}
|
||||
|
||||
@Test
|
||||
void anExplicitLeadersEntryWinsOverThePinForTheSameTerminal(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("both.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
primary:
|
||||
terminal: term_shared
|
||||
fleet:
|
||||
leaders:
|
||||
opus-5.0:
|
||||
terminal: term_shared
|
||||
""");
|
||||
|
||||
assertEquals(Map.of("term_shared", "opus-5.0"), BridgedConfig.load(f).leaderTerminals(),
|
||||
"the pin is the older spelling of the same fact; the named entry is what was meant");
|
||||
}
|
||||
|
||||
@Test
|
||||
void bothBlocksTogetherRegisterTheUnionOfTheirTerminals(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("union.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
primary:
|
||||
terminal: term_pinned
|
||||
fleet:
|
||||
leaders:
|
||||
gpt-sol-5.6:
|
||||
terminal: term_sol
|
||||
""");
|
||||
|
||||
assertEquals(Map.of("term_pinned", "primary", "term_sol", "gpt-sol-5.6"),
|
||||
BridgedConfig.load(f).leaderTerminals());
|
||||
}
|
||||
|
||||
@Test
|
||||
void neitherBlockLeavesNothingRegistered(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("none.yaml");
|
||||
@@ -483,22 +528,25 @@ class BridgedConfigTest {
|
||||
assertTrue(BridgedConfig.load(f).leaderTerminals().isEmpty());
|
||||
}
|
||||
|
||||
/** A lead entry with no terminal identifies nothing — it must not register a null key. */
|
||||
/**
|
||||
* CB-579: {@code fleet.leaders} no longer feeds {@code leaderTerminals()} at all — a lead's
|
||||
* identity comes from the live tab scan, not a config-held terminal map. This method now exists
|
||||
* only for the {@code primary.terminal} fallback.
|
||||
*/
|
||||
@Test
|
||||
void aLeadWithoutATerminalIsNotRegistered(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("no-terminal.yaml");
|
||||
void fleetLeadersNeverContributesToLeaderTerminals(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("leaders-only.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
fleet:
|
||||
leaders:
|
||||
sketch:
|
||||
kind: opencode
|
||||
real:
|
||||
terminal: term_real
|
||||
opus-5.0:
|
||||
tab: "lead: opus-5.0"
|
||||
""");
|
||||
|
||||
assertEquals(Map.of("term_real", "real"), BridgedConfig.load(f).leaderTerminals());
|
||||
assertTrue(BridgedConfig.load(f).leaderTerminals().isEmpty(),
|
||||
"no primary.terminal pin ⇒ nothing registered, even with fleet.leaders configured");
|
||||
}
|
||||
|
||||
// ── CB-548: the architects registry ────────────────────────────────────────────────────────
|
||||
@@ -1092,6 +1140,7 @@ class BridgedConfigTest {
|
||||
gitHostEnv: GITEA_HOST
|
||||
weight: 0.5
|
||||
maxLoad: 2
|
||||
exhaustedPattern: "usage limit has been reached"
|
||||
placement: weighted
|
||||
lifecycle:
|
||||
idleTtlSeconds: 300
|
||||
@@ -1123,6 +1172,8 @@ class BridgedConfigTest {
|
||||
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");
|
||||
assertTrue(w.hasExhaustedPattern(), "exhaustedPattern binds and enables the CB-578 stage A classification");
|
||||
assertEquals("usage limit has been reached", w.exhaustedPattern());
|
||||
|
||||
assertEquals("weighted", cfg.placement(), "placement binds at the top level");
|
||||
assertEquals(300, cfg.lifecycle().idleTtlSeconds());
|
||||
|
||||
@@ -338,6 +338,54 @@ class ConfigRefTest {
|
||||
assertEquals("deepseek-v4-flash", ref.get().profiles().get("sonnet").model());
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-578 stage B: exhaustedPattern is compiled once into Bridged.main's pattern map at startup
|
||||
* (see ExhaustedPatternLookup), so a reload never re-reads it — a changed pattern must be
|
||||
* reported deferred exactly like model/baseUrl, not silently claimed as applied.
|
||||
*/
|
||||
@Test
|
||||
void changingAProfilesExhaustedPatternIsReportedAsDeferred(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("bridged.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
herdrSocket: ~/.config/herdr/herdr.sock
|
||||
profiles:
|
||||
sonnet:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
model: sonnet
|
||||
exhaustedPattern: "usage limit has been reached"
|
||||
guard:
|
||||
offSubscriptionHosts:
|
||||
- gx00.gw
|
||||
""");
|
||||
ConfigRef ref = refFor(f);
|
||||
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
herdrSocket: ~/.config/herdr/herdr.sock
|
||||
profiles:
|
||||
sonnet:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
model: sonnet
|
||||
exhaustedPattern: "rate limit exceeded"
|
||||
guard:
|
||||
offSubscriptionHosts:
|
||||
- gx00.gw
|
||||
""");
|
||||
ConfigRef.Outcome out = ref.reload();
|
||||
|
||||
assertTrue(out.applied());
|
||||
assertEquals(1, out.deferred().size(), out.deferred().toString());
|
||||
assertTrue(out.deferred().getFirst().contains("sonnet"), out.deferred().toString());
|
||||
assertTrue(out.deferred().getFirst().contains("launch settings"), out.deferred().toString());
|
||||
// The snapshot still carries the new value — a restart is what makes it take effect.
|
||||
assertEquals("rate limit exceeded", ref.get().profiles().get("sonnet").exhaustedPattern());
|
||||
}
|
||||
|
||||
@Test
|
||||
void aFixedRefHasNoFileAndRefusesToReload() {
|
||||
BridgedConfig cfg = new BridgedConfig(null, null, null, null, null, null,
|
||||
|
||||
@@ -20,8 +20,10 @@ import org.slf4j.LoggerFactory;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.function.BiConsumer;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
class FleetHealthMonitorTest {
|
||||
@Test void oneTickUsesOneFleetListForAnyRosterSize() {
|
||||
@@ -37,7 +39,8 @@ class FleetHealthMonitorTest {
|
||||
herdr.calls.clear();
|
||||
MessageService messages = new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox());
|
||||
var scheduler = Executors.newSingleThreadScheduledExecutor();
|
||||
FleetHealthMonitor monitor = new FleetHealthMonitor(agents, sessions::roster, messages, scheduler, () -> 1, 60);
|
||||
FleetHealthMonitor monitor = new FleetHealthMonitor(agents, sessions::roster, messages, scheduler, () -> 1, 60,
|
||||
(_, _) -> { });
|
||||
monitor.tick();
|
||||
monitor.stop();
|
||||
assertEquals(1, herdr.calls.stream().filter(call -> call.method().equals("agent.list")).count());
|
||||
@@ -49,7 +52,7 @@ class FleetHealthMonitorTest {
|
||||
var scheduler = Executors.newSingleThreadScheduledExecutor();
|
||||
FleetHealthMonitor monitor = new FleetHealthMonitor(agents, java.util.List::of,
|
||||
new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()),
|
||||
scheduler, () -> 1, 60);
|
||||
scheduler, () -> 1, 60, (_, _) -> { });
|
||||
monitor.tick();
|
||||
herdr.healthy(true);
|
||||
monitor.tick();
|
||||
@@ -68,7 +71,7 @@ class FleetHealthMonitorTest {
|
||||
var scheduler = Executors.newSingleThreadScheduledExecutor();
|
||||
FleetHealthMonitor monitor = new FleetHealthMonitor(agents, java.util.List::of,
|
||||
new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()),
|
||||
scheduler, () -> 1, 60);
|
||||
scheduler, () -> 1, 60, (_, _) -> { });
|
||||
monitor.reportTransition("term_a", HealthState.TURN_BOUNDARY_LOST);
|
||||
monitor.reportTransition("term_a", HealthState.TURN_BOUNDARY_LOST);
|
||||
monitor.stop();
|
||||
@@ -78,4 +81,90 @@ class FleetHealthMonitorTest {
|
||||
logger.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
|
||||
// --- CB-580: a member that reaches GONE/NEVER_READY must fail its waiting tickets
|
||||
|
||||
private static FleetHealthMonitor monitorWith(BiConsumer<String, String> failTarget) {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
var scheduler = Executors.newSingleThreadScheduledExecutor();
|
||||
return new FleetHealthMonitor(agents, java.util.List::of,
|
||||
new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()),
|
||||
scheduler, () -> 1, 60, failTarget);
|
||||
}
|
||||
|
||||
@Test void terminalTransitionFailsTheTargetOnce() {
|
||||
RecordingFailTarget failTarget = new RecordingFailTarget();
|
||||
FleetHealthMonitor monitor = monitorWith(failTarget);
|
||||
monitor.reportTransition("term_a", HealthState.GONE);
|
||||
monitor.stop();
|
||||
assertEquals(1, failTarget.calls.size());
|
||||
assertEquals("term_a", failTarget.calls.get(0).target());
|
||||
assertTrue(failTarget.calls.get(0).reason().contains("GONE"));
|
||||
}
|
||||
|
||||
@Test void neverReadyNamesItselfAsTheReason() {
|
||||
RecordingFailTarget failTarget = new RecordingFailTarget();
|
||||
FleetHealthMonitor monitor = monitorWith(failTarget);
|
||||
monitor.reportTransition("term_a", HealthState.NEVER_READY);
|
||||
monitor.stop();
|
||||
assertEquals(1, failTarget.calls.size());
|
||||
assertTrue(failTarget.calls.get(0).reason().contains("NEVER_READY"));
|
||||
}
|
||||
|
||||
@Test void stayingInATerminalStateProducesOneFailureNotN() {
|
||||
RecordingFailTarget failTarget = new RecordingFailTarget();
|
||||
FleetHealthMonitor monitor = monitorWith(failTarget);
|
||||
monitor.reportTransition("term_a", HealthState.GONE);
|
||||
monitor.reportTransition("term_a", HealthState.GONE);
|
||||
monitor.reportTransition("term_a", HealthState.GONE);
|
||||
monitor.reportTransition("term_a", HealthState.GONE);
|
||||
monitor.stop();
|
||||
assertEquals(1, failTarget.calls.size());
|
||||
}
|
||||
|
||||
@Test void aNonTerminalFaultStateDoesNotFailTheTarget() {
|
||||
RecordingFailTarget failTarget = new RecordingFailTarget();
|
||||
FleetHealthMonitor monitor = monitorWith(failTarget);
|
||||
monitor.reportTransition("term_a", HealthState.TURN_BOUNDARY_LOST);
|
||||
monitor.stop();
|
||||
assertEquals(0, failTarget.calls.size());
|
||||
}
|
||||
|
||||
@Test void failTargetRetryIsBounded() {
|
||||
AlwaysThrowingFailTarget failTarget = new AlwaysThrowingFailTarget();
|
||||
FleetHealthMonitor monitor = monitorWith(failTarget);
|
||||
monitor.reportTransition("term_a", HealthState.GONE);
|
||||
monitor.stop();
|
||||
assertEquals(FleetHealthMonitor.MAX_FAIL_TARGET_ATTEMPTS, failTarget.calls);
|
||||
}
|
||||
|
||||
@Test void exhaustedRetryStillDoesNotRefireOnAnUnchangedTick() {
|
||||
AlwaysThrowingFailTarget failTarget = new AlwaysThrowingFailTarget();
|
||||
FleetHealthMonitor monitor = monitorWith(failTarget);
|
||||
monitor.reportTransition("term_a", HealthState.GONE);
|
||||
int afterFirstTransition = failTarget.calls;
|
||||
monitor.reportTransition("term_a", HealthState.GONE);
|
||||
monitor.stop();
|
||||
assertEquals(afterFirstTransition, failTarget.calls);
|
||||
}
|
||||
|
||||
private record RecordedCall(String target, String reason) { }
|
||||
|
||||
private static final class RecordingFailTarget implements BiConsumer<String, String> {
|
||||
final java.util.List<RecordedCall> calls = new java.util.ArrayList<>();
|
||||
|
||||
@Override public void accept(String target, String reason) {
|
||||
calls.add(new RecordedCall(target, reason));
|
||||
}
|
||||
}
|
||||
|
||||
private static final class AlwaysThrowingFailTarget implements BiConsumer<String, String> {
|
||||
int calls = 0;
|
||||
|
||||
@Override public void accept(String target, String reason) {
|
||||
calls++;
|
||||
throw new RuntimeException("boom");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -15,8 +15,9 @@ import java.util.concurrent.atomic.AtomicLong;
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
/**
|
||||
* CB-531. A lead is never spawned, so the daemon has to <em>find</em> it: these assert that an
|
||||
* operator-labelled tab is what makes a pane a lead, and — just as importantly — what does not.
|
||||
* CB-531/CB-579. A lead is never spawned, so the daemon has to <em>find</em> it: these assert that
|
||||
* an operator-labelled tab matching a configured {@code tab:} is what makes a pane a lead, and —
|
||||
* just as importantly — what does not, and that a stale entry does not linger forever.
|
||||
*/
|
||||
class LeadTabScannerTest {
|
||||
|
||||
@@ -120,23 +121,48 @@ class LeadTabScannerTest {
|
||||
.pane("w9:p1", "w9:t1", "term_worker");
|
||||
}
|
||||
|
||||
private LeadTabScanner scanner(TopologyHerdr herdr, Map<String, String> configured,
|
||||
/** The {@code tab:} → name map {@code twoLeads()}'s two lead tabs are configured under. */
|
||||
private static Map<String, String> twoLeadsConfigured() {
|
||||
return Map.of("lead: opus-5.0", "opus-5.0", "lead: gpt-sol-5.6", "gpt-sol-5.6");
|
||||
}
|
||||
|
||||
private LeadTabScanner scanner(TopologyHerdr herdr, Map<String, String> tabToName,
|
||||
AtomicLong clock) {
|
||||
return new LeadTabScanner(herdr, "lead:", Set.of("bridged-workers"), configured, TTL,
|
||||
clock::get);
|
||||
return new LeadTabScanner(herdr, tabToName, Set.of("bridged-workers"), TTL, clock::get);
|
||||
}
|
||||
|
||||
@Test
|
||||
void everyLabelledTabBecomesALeadNamedByItsLabel() {
|
||||
Map<String, String> leads = scanner(twoLeads(), Map.of(), new AtomicLong()).get();
|
||||
void everyConfiguredTabBecomesALeadNamedByItsEntry() {
|
||||
Map<String, String> leads = scanner(twoLeads(), twoLeadsConfigured(), new AtomicLong()).get();
|
||||
|
||||
assertEquals(Map.of("term_opus", "opus-5.0", "term_gpt", "gpt-sol-5.6"), leads,
|
||||
"two leads discovered from labels alone — no terminal_id was ever configured");
|
||||
"two leads discovered by their configured tab — no terminal_id was ever configured");
|
||||
}
|
||||
|
||||
@Test
|
||||
void anUnlabelledTabContributesNothing() {
|
||||
assertFalse(scanner(twoLeads(), Map.of(), new AtomicLong()).get().containsKey("term_notes"));
|
||||
void anUnconfiguredTabContributesNothing() {
|
||||
assertFalse(scanner(twoLeads(), twoLeadsConfigured(), new AtomicLong())
|
||||
.get().containsKey("term_notes"));
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-579: matching is exact against the configured map now, not a shared prefix — two leads with
|
||||
* completely different labels are both discovered by one scanner, no convention required.
|
||||
*/
|
||||
@Test
|
||||
void twoLeadsWithCompletelyDifferentLabelsAreBothDiscovered() {
|
||||
TopologyHerdr herdr = new TopologyHerdr()
|
||||
.workspace("w1", "main")
|
||||
.tab("w1:t1", "w1", "orchestrator: opus")
|
||||
.tab("w1:t2", "w1", "captain: sol")
|
||||
.pane("w1:p1", "w1:t1", "term_opus")
|
||||
.pane("w1:p2", "w1:t2", "term_sol");
|
||||
Map<String, String> tabToName = Map.of("orchestrator: opus", "opus", "captain: sol", "sol");
|
||||
|
||||
Map<String, String> leads = scanner(herdr, tabToName, new AtomicLong()).get();
|
||||
|
||||
assertEquals(Map.of("term_opus", "opus", "term_sol", "sol"), leads,
|
||||
"no shared prefix needed — each lead is matched by its own configured tab");
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -148,25 +174,28 @@ class LeadTabScannerTest {
|
||||
void aTabInAWorkerSpaceIsNeverALeadEvenWhenItsLabelMatches() {
|
||||
TopologyHerdr herdr = twoLeads().tab("w9:t2", "w9", "lead: impostor")
|
||||
.pane("w9:p2", "w9:t2", "term_impostor");
|
||||
Map<String, String> tabToName = new LinkedHashMap<>(twoLeadsConfigured());
|
||||
tabToName.put("lead: impostor", "impostor");
|
||||
|
||||
assertFalse(scanner(herdr, Map.of(), new AtomicLong()).get().containsKey("term_impostor"));
|
||||
assertFalse(scanner(herdr, tabToName, new AtomicLong()).get().containsKey("term_impostor"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void aBarePrefixNamesNobodyAndIsRejected() {
|
||||
void aLabelWithNoConfiguredEntryIsIgnored() {
|
||||
TopologyHerdr herdr = new TopologyHerdr().workspace("w1", "main")
|
||||
.tab("w1:t1", "w1", "lead:").pane("w1:p1", "w1:t1", "term_a");
|
||||
.tab("w1:t1", "w1", "lead: nobody-configured").pane("w1:p1", "w1:t1", "term_a");
|
||||
|
||||
assertEquals(Map.of(), scanner(herdr, Map.of(), new AtomicLong()).get(),
|
||||
"a lead with no name would resolve as PRIMARY with nothing to attribute it to");
|
||||
assertEquals(Map.of(), scanner(herdr, twoLeadsConfigured(), new AtomicLong()).get(),
|
||||
"a label that names no configured lead resolves nobody");
|
||||
}
|
||||
|
||||
@Test
|
||||
void thePrefixMatchesCaseInsensitivelyAndTheNameIsTrimmed() {
|
||||
void matchingIsCaseInsensitiveAndToleratesSurroundingWhitespace() {
|
||||
TopologyHerdr herdr = new TopologyHerdr().workspace("w1", "main")
|
||||
.tab("w1:t1", "w1", " LEAD: opus-5.0 ").pane("w1:p1", "w1:t1", "term_a");
|
||||
.tab("w1:t1", "w1", " LEAD: Opus-5.0 ").pane("w1:p1", "w1:t1", "term_a");
|
||||
|
||||
assertEquals(Map.of("term_a", "opus-5.0"), scanner(herdr, Map.of(), new AtomicLong()).get());
|
||||
assertEquals(Map.of("term_a", "opus-5.0"),
|
||||
scanner(herdr, Map.of("lead: Opus-5.0", "opus-5.0"), new AtomicLong()).get());
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -175,18 +204,51 @@ class LeadTabScannerTest {
|
||||
// nothing bridged placed can land here (see the worker-space test above).
|
||||
TopologyHerdr herdr = twoLeads().pane("w1:p1b", "w1:t1", "term_opus_split");
|
||||
|
||||
assertEquals("opus-5.0", scanner(herdr, Map.of(), new AtomicLong()).get().get("term_opus_split"));
|
||||
assertEquals("opus-5.0",
|
||||
scanner(herdr, twoLeadsConfigured(), new AtomicLong()).get().get("term_opus_split"));
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-579 acceptance (6): this is the bug the ticket closes. A stale pin used to be merged back
|
||||
* over every scan and never expire; now a scan is the whole answer, so a lead whose tab is gone
|
||||
* drops out on the very next scan.
|
||||
*/
|
||||
@Test
|
||||
void anExplicitlyConfiguredLeadIsMergedInAndOutranksALabel() {
|
||||
Map<String, String> configured = Map.of("term_opus", "pinned-name", "term_extra", "from-config");
|
||||
void aTabNoLongerPresentDropsTheLeadOnTheNextScan() {
|
||||
TopologyHerdr herdr = twoLeads();
|
||||
AtomicLong clock = new AtomicLong();
|
||||
LeadTabScanner s = scanner(herdr, twoLeadsConfigured(), clock);
|
||||
assertTrue(s.get().containsKey("term_opus"));
|
||||
|
||||
Map<String, String> leads = scanner(twoLeads(), configured, new AtomicLong()).get();
|
||||
// The session behind term_opus restarted — herdr no longer reports that tab or pane at all.
|
||||
herdr.tabs.remove("w1:t1");
|
||||
herdr.panes.remove("w1:p1");
|
||||
clock.addAndGet(TTL);
|
||||
|
||||
assertEquals("pinned-name", leads.get("term_opus"), "an explicit pin is the operator's last word");
|
||||
assertEquals("from-config", leads.get("term_extra"), "a configured lead needs no tab at all");
|
||||
assertEquals("gpt-sol-5.6", leads.get("term_gpt"));
|
||||
assertFalse(s.get().containsKey("term_opus"),
|
||||
"a stale entry must expire once the tab it named is gone, not be merged back forever");
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-579 acceptance (5): the whole point of matching by tab instead of {@code terminal_id} — a
|
||||
* restart changes the terminal, not the tab, so the lead resolves under the same name with no
|
||||
* config edit.
|
||||
*/
|
||||
@Test
|
||||
void aLeadRestartingInTheSameTabResolvesUnderTheSameName() {
|
||||
TopologyHerdr herdr = twoLeads();
|
||||
AtomicLong clock = new AtomicLong();
|
||||
LeadTabScanner s = scanner(herdr, twoLeadsConfigured(), clock);
|
||||
assertEquals("opus-5.0", s.get().get("term_opus"));
|
||||
|
||||
// The session restarts: herdr assigns the pane a new terminal_id, same tab (w1:t1).
|
||||
herdr.panes.remove("w1:p1");
|
||||
herdr.pane("w1:p1", "w1:t1", "term_opus_v2");
|
||||
clock.addAndGet(TTL);
|
||||
|
||||
Map<String, String> leads = s.get();
|
||||
assertEquals("opus-5.0", leads.get("term_opus_v2"), "the new terminal resolves immediately");
|
||||
assertFalse(leads.containsKey("term_opus"), "the old terminal_id is simply gone, not carried");
|
||||
}
|
||||
|
||||
// ── caching ─────────────────────────────────────────────────────────────────────────────────
|
||||
@@ -195,7 +257,7 @@ class LeadTabScannerTest {
|
||||
void aSecondLookupWithinTheTtlDoesNotTouchHerdr() {
|
||||
TopologyHerdr herdr = twoLeads();
|
||||
AtomicLong clock = new AtomicLong();
|
||||
LeadTabScanner s = scanner(herdr, Map.of(), clock);
|
||||
LeadTabScanner s = scanner(herdr, twoLeadsConfigured(), clock);
|
||||
|
||||
s.get();
|
||||
int afterFirst = herdr.calls;
|
||||
@@ -210,21 +272,23 @@ class LeadTabScannerTest {
|
||||
void aTabLabelledAfterStartupIsPickedUpOnceTheTtlExpires() {
|
||||
TopologyHerdr herdr = twoLeads();
|
||||
AtomicLong clock = new AtomicLong();
|
||||
LeadTabScanner s = scanner(herdr, Map.of(), clock);
|
||||
Map<String, String> tabToName = new LinkedHashMap<>(twoLeadsConfigured());
|
||||
tabToName.put("lead: late-arrival", "late-arrival");
|
||||
LeadTabScanner s = scanner(herdr, tabToName, clock);
|
||||
assertFalse(s.get().containsKey("term_notes"));
|
||||
|
||||
herdr.tab("w1:t3", "w1", "lead: late-arrival"); // the operator renames their tab
|
||||
clock.addAndGet(TTL);
|
||||
|
||||
assertEquals("late-arrival", s.get().get("term_notes"),
|
||||
"the whole point over `leaders:`: no config edit, no restart");
|
||||
"the whole point over a config-held terminal_id: no config edit, no restart");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aFailedScanKeepsTheLeadsAlreadyKnownRatherThanDemotingThem() {
|
||||
TopologyHerdr herdr = twoLeads();
|
||||
AtomicLong clock = new AtomicLong();
|
||||
LeadTabScanner s = scanner(herdr, Map.of(), clock);
|
||||
LeadTabScanner s = scanner(herdr, twoLeadsConfigured(), clock);
|
||||
Map<String, String> before = s.get();
|
||||
|
||||
herdr.failing = true;
|
||||
@@ -235,14 +299,14 @@ class LeadTabScannerTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
void aFailedFirstScanStillHonoursTheConfiguredLeads() {
|
||||
void aFailedFirstScanReturnsEmptyRatherThanThrowing() {
|
||||
TopologyHerdr herdr = twoLeads();
|
||||
herdr.failing = true;
|
||||
|
||||
Map<String, String> leads = scanner(herdr, Map.of("term_x", "opus-5.0"), new AtomicLong()).get();
|
||||
Map<String, String> leads = scanner(herdr, twoLeadsConfigured(), new AtomicLong()).get();
|
||||
|
||||
assertEquals(Map.of("term_x", "opus-5.0"), leads,
|
||||
"config-named leads must not depend on herdr answering at all");
|
||||
assertEquals(Map.of(), leads,
|
||||
"with nothing scanned yet and no override to fall back on, the map is simply empty");
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -250,7 +314,7 @@ class LeadTabScannerTest {
|
||||
TopologyHerdr herdr = twoLeads();
|
||||
herdr.failing = true;
|
||||
AtomicLong clock = new AtomicLong();
|
||||
LeadTabScanner s = scanner(herdr, Map.of(), clock);
|
||||
LeadTabScanner s = scanner(herdr, twoLeadsConfigured(), clock);
|
||||
|
||||
s.get();
|
||||
int afterFirst = herdr.calls;
|
||||
|
||||
@@ -12,6 +12,9 @@ import dev.ltms.bridged.msg.TurnToken;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.Set;
|
||||
import java.util.regex.Pattern;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
@@ -23,7 +26,7 @@ class CompletionResolverTest {
|
||||
void skipsTheScrapeWhenNoSendIsWaiting() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Rendezvous rendezvous = new Rendezvous(); // no waiter opened
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
|
||||
resolver.resolve("term_a", null); // no in-flight turn captured for this target
|
||||
|
||||
@@ -35,7 +38,7 @@ class CompletionResolverTest {
|
||||
void failSkipsTheScrapeWhenNoSendIsWaiting() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Rendezvous rendezvous = new Rendezvous(); // no waiter opened
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
|
||||
resolver.fail("term_a", null); // no in-flight turn, and no registered waiter to fall back to
|
||||
|
||||
@@ -47,7 +50,7 @@ class CompletionResolverTest {
|
||||
void captureBaselineSkipsTheReadWhenNoSendIsWaiting() {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("⏺ X\n❯ ");
|
||||
Rendezvous rendezvous = new Rendezvous(); // no waiter opened
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
|
||||
resolver.captureBaseline("term_a", TestTurnTokens.inert("term_a")); // no send to attribute a later completion to
|
||||
|
||||
@@ -136,7 +139,7 @@ class CompletionResolverTest {
|
||||
// send must NOT be resolved with the stale answer.
|
||||
FakeHerdr herdr = new FakeHerdr().readText("⏺ 391\n❯ ");
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
|
||||
var waiter = rendezvous.open("term_a"); // a send is blocked on this turn
|
||||
// The turn as captured at delivery: its waiter, and the previous turn's answer still on screen.
|
||||
@@ -151,7 +154,7 @@ class CompletionResolverTest {
|
||||
void resolvesACompletionWhoseScrapeChangedSinceDelivery() {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("⏺ No, 391 = 17 × 23.\n❯ "); // the worker's real answer
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
// Delivery baseline was the previous turn's "391"; the scrape now differs → resolve.
|
||||
@@ -168,7 +171,7 @@ class CompletionResolverTest {
|
||||
String block = "⏺ " + "x".repeat(CompletionResolver.MAX_SCRAPE_CHARS + 1) + "\n❯ ";
|
||||
FakeHerdr herdr = new FakeHerdr().readText(block);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
|
||||
@@ -182,7 +185,7 @@ class CompletionResolverTest {
|
||||
void leavesAnUnclippedCompletionPaneTailUnmarked() {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("⏺ complete report\n❯ ");
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
|
||||
@@ -194,7 +197,7 @@ class CompletionResolverTest {
|
||||
void resolvesSynchronouslyBeforePostTurnContextClearing() {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("⏺ previous answer\n❯ ");
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
var waiter = rendezvous.open("term_a");
|
||||
resolver.captureBaseline("term_a", new TurnToken("term_a", waiter));
|
||||
herdr.readText("⏺ answer that /clear would erase\n❯ ");
|
||||
@@ -216,7 +219,7 @@ class CompletionResolverTest {
|
||||
String longBlock = "⏺ " + "x".repeat(CompletionResolver.MAX_SCRAPE_CHARS + 500) + "\n❯ ";
|
||||
FakeHerdr herdr = new FakeHerdr().readText(longBlock);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
|
||||
var waiter = rendezvous.open("term_a"); // a send is blocked on this turn
|
||||
resolver.captureBaseline("term_a", new TurnToken("term_a", waiter)); // baseline is the clipped >cap block
|
||||
@@ -236,7 +239,7 @@ class CompletionResolverTest {
|
||||
// No delivery baseline (e.g. the pre-turn read failed) ⇒ never suppress; the completion resolves.
|
||||
FakeHerdr herdr = new FakeHerdr().readText("⏺ hello\n❯ ");
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
|
||||
@@ -254,7 +257,7 @@ class CompletionResolverTest {
|
||||
// byte-identical guard would wrongly match the empty tail and suppress.
|
||||
FakeHerdr herdr = new FakeHerdr().healthy(false); // agent.read throws HerdrException
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
var turn = new CompletionResolver.InFlight(waiter, ""); // empty pane baselined at delivery
|
||||
@@ -274,7 +277,7 @@ class CompletionResolverTest {
|
||||
// fail must not overwrite that value, and must not even scrape the worker — nobody needs it.
|
||||
FakeHerdr herdr = new FakeHerdr().readText("an error screen");
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
var turn = new CompletionResolver.InFlight(waiter, null);
|
||||
@@ -295,7 +298,7 @@ class CompletionResolverTest {
|
||||
// fail falls back to the waiter currently registered on the Rendezvous and fails it.
|
||||
FakeHerdr herdr = new FakeHerdr().readText("stuck on an error screen");
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
|
||||
var waiter = rendezvous.open("term_a"); // send registered, but no captureBaseline ever ran
|
||||
resolver.fail("term_a", null); // no in-flight turn → fall back to the registered waiter
|
||||
@@ -321,7 +324,7 @@ class CompletionResolverTest {
|
||||
try {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("stuck on an error screen");
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
var waiter = rendezvous.open("term_a");
|
||||
|
||||
resolver.fail("term_a", null);
|
||||
@@ -349,7 +352,7 @@ class CompletionResolverTest {
|
||||
// scrape to turn N+1; targeting turn N's captured waiter makes the late completion a no-op.
|
||||
FakeHerdr herdr = new FakeHerdr().readText("⏺ turn N answer\n❯ ");
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
|
||||
var waiterN = rendezvous.open("term_a"); // turn N's send
|
||||
// The turn as the injector captured it at delivery (waiter + pre-turn baseline).
|
||||
@@ -370,4 +373,129 @@ class CompletionResolverTest {
|
||||
"turn N stays resolved by its own reply");
|
||||
assertTrue(rendezvous.isWaiting("term_a"), "turn N+1 is still awaiting its own resolution");
|
||||
}
|
||||
|
||||
// --- CB-578 stage A: backend-exhausted classification ---------------------------------
|
||||
|
||||
@Test
|
||||
void classifiesAMatchingScrapeAsBackendExhaustedInsteadOfACompletedReply() {
|
||||
String block = "⏺ Working on it...\nThe usage limit has been reached. Try again later.\n❯ ";
|
||||
FakeHerdr herdr = new FakeHerdr().readText(block);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached");
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, ExhaustionSink.none());
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
|
||||
|
||||
assertTrue(waiter.isDone(), "a matching scrape still resolves the blocked send");
|
||||
assertEquals(Rendezvous.Kind.BACKEND_EXHAUSTED, waiter.getNow(null).kind(),
|
||||
"not reported as a completed reply — the classification is distinct");
|
||||
}
|
||||
|
||||
@Test
|
||||
void theExhaustedReasonCarriesTheMatchedLine() {
|
||||
String block = "⏺ Working on it...\nThe usage limit has been reached. Try again later.\n❯ ";
|
||||
FakeHerdr herdr = new FakeHerdr().readText(block);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached");
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, ExhaustionSink.none());
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
|
||||
|
||||
assertEquals("backend exhausted (usage limit): The usage limit has been reached. Try again later.",
|
||||
waiter.getNow(null).text(), "the reason names the real cause and carries the matched line");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aWinningBackendExhaustedClassificationNotifiesTheExhaustionSink() {
|
||||
String block = "⏺ Working on it...\nThe usage limit has been reached. Try again later.\n❯ ";
|
||||
FakeHerdr herdr = new FakeHerdr().readText(block);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached");
|
||||
java.util.List<String> notified = new java.util.ArrayList<>();
|
||||
ExhaustionSink sink = (target, reason) -> notified.add(target + ": " + reason);
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, sink);
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
|
||||
|
||||
assertEquals(1, notified.size(), "the sink is notified exactly once for the winning classification");
|
||||
assertTrue(notified.get(0).startsWith("term_a: "), "the sink is told which target exhausted");
|
||||
assertTrue(notified.get(0).contains("The usage limit has been reached"),
|
||||
"the sink is told the matched reason: " + notified.get(0));
|
||||
}
|
||||
|
||||
@Test
|
||||
void aLosingBackendExhaustedClassificationNeverNotifiesTheExhaustionSink() {
|
||||
// The waiter was already resolved (e.g. by the worker's own reply) before this scrape landed —
|
||||
// resolveExhausted loses the race and must return false, so the sink must not fire either.
|
||||
String block = "⏺ Working on it...\nThe usage limit has been reached. Try again later.\n❯ ";
|
||||
FakeHerdr herdr = new FakeHerdr().readText(block);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached");
|
||||
java.util.List<String> notified = new java.util.ArrayList<>();
|
||||
ExhaustionSink sink = (target, reason) -> notified.add(target + ": " + reason);
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, sink);
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
var turn = new CompletionResolver.InFlight(waiter, null);
|
||||
assertTrue(rendezvous.resolveCompletion(waiter, "already replied"));
|
||||
|
||||
resolver.resolve("term_a", turn);
|
||||
|
||||
assertTrue(notified.isEmpty(), "a classification that loses the race must not quarantine anything");
|
||||
assertEquals("already replied", waiter.getNow(null).text(), "the earlier resolution stands untouched");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aNonMatchingScrapeResolvesAsAnOrdinaryCompletion() {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("⏺ complete report\n❯ ");
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached");
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, ExhaustionSink.none());
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
|
||||
|
||||
assertEquals(Rendezvous.Kind.COMPLETION, waiter.getNow(null).kind(),
|
||||
"a scrape that does not match the pattern is an ordinary completion");
|
||||
assertEquals("complete report", waiter.getNow(null).text());
|
||||
}
|
||||
|
||||
@Test
|
||||
void aProfileWithNoConfiguredPatternKeepsTodaysCompletionFallbackUnchanged() {
|
||||
// Even a scrape that WOULD have matched some other profile's pattern must resolve as a
|
||||
// plain completion when this target's own profile has none configured (CB-578 criterion 4).
|
||||
String block = "⏺ The usage limit has been reached.\n❯ ";
|
||||
FakeHerdr herdr = new FakeHerdr().readText(block);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver =
|
||||
new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
|
||||
|
||||
assertEquals(Rendezvous.Kind.COMPLETION, waiter.getNow(null).kind(),
|
||||
"no pattern configured for this target's profile ⇒ unchanged completion-fallback behaviour");
|
||||
assertEquals("The usage limit has been reached.", waiter.getNow(null).text());
|
||||
}
|
||||
|
||||
@Test
|
||||
void coverageIsOffWhenNoProfileHasAPatternConfigured() {
|
||||
assertEquals("off (no profile has an exhaustedPattern configured; profiles: [terra])",
|
||||
CompletionResolver.coverage(Set.of("terra"), Set.of()));
|
||||
}
|
||||
|
||||
@Test
|
||||
void coverageIsFullWhenEveryProfileHasAPatternConfigured() {
|
||||
assertEquals("full (all profiles configured: [gx10, terra])",
|
||||
CompletionResolver.coverage(Set.of("terra", "gx10"), Set.of("terra", "gx10")));
|
||||
}
|
||||
|
||||
@Test
|
||||
void coverageIsPartialAndNamesWhichProfilesAreConfigured() {
|
||||
assertEquals("partial (configured: [terra]; not configured: [gx10])",
|
||||
CompletionResolver.coverage(Set.of("terra", "gx10"), Set.of("terra")));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -28,7 +28,7 @@ class LeadLauncherTest {
|
||||
List.of("ccs", "ltms"), "tab", "bridged-workers", null,
|
||||
"http://127.0.0.1:8765/mcp", null, null,
|
||||
null, null, null,
|
||||
Map.of("CLAUDE_CODE_AUTO_COMPACT_WINDOW", "300000"), null, null, true);
|
||||
Map.of("CLAUDE_CODE_AUTO_COMPACT_WINDOW", "300000"), null, null, true, null);
|
||||
}
|
||||
|
||||
private static BridgedConfig configWith(BridgedConfig.Leader lead) {
|
||||
@@ -41,8 +41,8 @@ class LeadLauncherTest {
|
||||
null, null, fleet, null, "fixed", null).withDefaults();
|
||||
}
|
||||
|
||||
private static BridgedConfig.Leader lead(String profile, String terminal, int instances) {
|
||||
return new BridgedConfig.Leader(profile, terminal, instances, "lead:", 10, null, null,
|
||||
private static BridgedConfig.Leader lead(String profile, String tab, int instances) {
|
||||
return new BridgedConfig.Leader(profile, tab, instances, "lead:", 10, null, null,
|
||||
"leads", "/repo");
|
||||
}
|
||||
|
||||
@@ -70,16 +70,16 @@ class LeadLauncherTest {
|
||||
void startsTheDeclaredLeadWhenNoneIsRunning() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
|
||||
assertEquals(1, launcher(herdr, configWith(lead("opus", null, 1))).ensureLeads());
|
||||
assertEquals(1, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads());
|
||||
assertTrue(herdr.called("agent.start"), "a lead must actually be started");
|
||||
assertEquals("lead-opus", startedName(herdr));
|
||||
}
|
||||
|
||||
/** The tab is labelled so the scanner finds the lead on the next resolve. */
|
||||
/** The tab is labelled with the configured `tab:` so the scanner finds the lead on the next resolve. */
|
||||
@Test
|
||||
void labelsTheTabWithThePrefixTheScannerReadsBack() {
|
||||
void labelsTheTabWithTheConfiguredTabValue() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
launcher(herdr, configWith(lead("opus", null, 1))).ensureLeads();
|
||||
launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads();
|
||||
|
||||
assertEquals("lead: opus",
|
||||
((Map<?, ?>) herdr.lastCall("tab.rename").params()).get("label"));
|
||||
@@ -90,7 +90,7 @@ class LeadLauncherTest {
|
||||
void startsAsManyInstancesAsAreDeclared() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
|
||||
assertEquals(2, launcher(herdr, configWith(lead("opus", null, 2))).ensureLeads());
|
||||
assertEquals(2, launcher(herdr, configWith(lead("opus", "lead: opus", 2))).ensureLeads());
|
||||
assertEquals(2, herdr.calls.stream().filter(c -> c.method().equals("agent.start")).count());
|
||||
}
|
||||
|
||||
@@ -104,7 +104,7 @@ class LeadLauncherTest {
|
||||
.withTab("wL", "wL:t1", "lead: opus")
|
||||
.withAgent("lead-opus", "term_lead", "wL:p1", "wL:t1");
|
||||
|
||||
assertEquals(0, launcher(herdr, configWith(lead("opus", null, 1))).ensureLeads());
|
||||
assertEquals(0, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads());
|
||||
assertFalse(herdr.called("agent.start"), "the live lead must not be duplicated");
|
||||
}
|
||||
|
||||
@@ -118,20 +118,23 @@ class LeadLauncherTest {
|
||||
.withWorkspace("wL", "leads")
|
||||
.withTab("wL", "wL:t1", "lead: opus"); // label only — nothing running in it
|
||||
|
||||
assertEquals(1, launcher(herdr, configWith(lead("opus", null, 1))).ensureLeads(),
|
||||
assertEquals(1, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads(),
|
||||
"a stale label is not a lead; the lead must be relaunched");
|
||||
}
|
||||
|
||||
/**
|
||||
* A lead the operator opened by hand and pinned with `terminal:` is live even though its tab
|
||||
* carries no matching label. Counting labels alone would relaunch it on every boot.
|
||||
* A lead the operator opened by hand is live once its tab carries the configured `tab:` label —
|
||||
* CB-579 retired the `terminal:` pin, so a hand-opened lead is found the same way an
|
||||
* auto-launched one is, by its tab, not by a terminal id nobody wrote down in advance.
|
||||
*/
|
||||
@Test
|
||||
void aPinnedTerminalWithARunningAgentCountsAsLive() {
|
||||
void aHandOpenedLeadWithTheConfiguredTabLabelCountsAsLive() {
|
||||
FakeHerdr herdr = new FakeHerdr()
|
||||
.withAgent("hand-opened", "term_pinned", "wX:p1", "wX:t1");
|
||||
.withWorkspace("wX", "main")
|
||||
.withTab("wX", "wX:t1", "lead: opus")
|
||||
.withAgent("hand-opened", "term_hand", "wX:p1", "wX:t1");
|
||||
|
||||
assertEquals(0, launcher(herdr, configWith(lead("opus", "term_pinned", 1))).ensureLeads());
|
||||
assertEquals(0, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads());
|
||||
assertFalse(herdr.called("agent.start"));
|
||||
}
|
||||
|
||||
@@ -143,7 +146,7 @@ class LeadLauncherTest {
|
||||
.withTab("wM", "wM:t1", "lead: opus") // a member tab that looks like a lead
|
||||
.withAgent("claude-opus-x", "term_m", "wM:p1", "wM:t1");
|
||||
|
||||
assertEquals(1, launcher(herdr, configWith(lead("opus", null, 1))).ensureLeads(),
|
||||
assertEquals(1, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads(),
|
||||
"a member in a lead-labelled tab is not a lead, so the real lead is still missing");
|
||||
}
|
||||
|
||||
@@ -152,7 +155,7 @@ class LeadLauncherTest {
|
||||
void anUncountableHerdrStartsNothing() {
|
||||
FakeHerdr herdr = new FakeHerdr().healthy(false);
|
||||
|
||||
assertEquals(0, launcher(herdr, configWith(lead("opus", null, 1))).ensureLeads());
|
||||
assertEquals(0, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads());
|
||||
assertFalse(herdr.called("agent.start"));
|
||||
}
|
||||
|
||||
@@ -166,7 +169,7 @@ class LeadLauncherTest {
|
||||
@Test
|
||||
void theLeadNeverReceivesTheWorkerReplyCharter() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
launcher(herdr, configWith(lead("opus", null, 1))).ensureLeads();
|
||||
launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads();
|
||||
|
||||
List<String> args = startedArgs(herdr);
|
||||
assertFalse(args.contains("--append-system-prompt"),
|
||||
@@ -178,7 +181,7 @@ class LeadLauncherTest {
|
||||
@Test
|
||||
void theLeadMountsTheBridgeMcpAndPinsItsModel() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
launcher(herdr, configWith(lead("opus", null, 1))).ensureLeads();
|
||||
launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads();
|
||||
|
||||
List<String> args = startedArgs(herdr);
|
||||
assertTrue(args.contains("--mcp-config"));
|
||||
@@ -192,7 +195,7 @@ class LeadLauncherTest {
|
||||
@Test
|
||||
void theLeadEnvCarriesNoAnthropicBinding() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
launcher(herdr, configWith(lead("opus", null, 1))).ensureLeads();
|
||||
launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads();
|
||||
|
||||
Map<String, String> env = tabEnv(herdr);
|
||||
assertNull(env.get("ANTHROPIC_BASE_URL"));
|
||||
@@ -205,7 +208,7 @@ class LeadLauncherTest {
|
||||
@Test
|
||||
void theLeadTabIsCreatedOutsideEveryMemberWorkspace() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
launcher(herdr, configWith(lead("opus", null, 1))).ensureLeads();
|
||||
launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads();
|
||||
|
||||
String label = (String) ((Map<?, ?>) herdr.lastCall("workspace.create").params()).get("label");
|
||||
assertEquals("leads", label);
|
||||
@@ -214,12 +217,12 @@ class LeadLauncherTest {
|
||||
|
||||
// ── recognise-only and misconfiguration ───────────────────────────────────────────────────
|
||||
|
||||
/** A lead with a pin but no profile is recognise-only by design — not an error, not a launch. */
|
||||
/** A lead with a tab but no profile is recognise-only by design — not an error, not a launch. */
|
||||
@Test
|
||||
void aLeadThatNamesNoProfileIsRecognisedButNeverLaunched() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
|
||||
assertEquals(0, launcher(herdr, configWith(lead(null, "term_dead", 1))).ensureLeads());
|
||||
assertEquals(0, launcher(herdr, configWith(lead(null, "lead: dead", 1))).ensureLeads());
|
||||
assertFalse(herdr.called("agent.start"));
|
||||
}
|
||||
|
||||
@@ -228,7 +231,7 @@ class LeadLauncherTest {
|
||||
void zeroInstancesLaunchesNothing() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
|
||||
assertEquals(0, launcher(herdr, configWith(lead("opus", null, 0))).ensureLeads());
|
||||
assertEquals(0, launcher(herdr, configWith(lead("opus", "lead: opus", 0))).ensureLeads());
|
||||
assertFalse(herdr.called("agent.start"));
|
||||
}
|
||||
|
||||
@@ -237,7 +240,7 @@ class LeadLauncherTest {
|
||||
void anUnknownProfileIsSkippedRatherThanThrown() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
|
||||
assertEquals(0, launcher(herdr, configWith(lead("nope", null, 1))).ensureLeads());
|
||||
assertEquals(0, launcher(herdr, configWith(lead("nope", "lead: opus", 1))).ensureLeads());
|
||||
assertFalse(herdr.called("agent.start"));
|
||||
}
|
||||
|
||||
|
||||
@@ -72,7 +72,8 @@ class BridgeMcpAuthzTest {
|
||||
new PrimaryRegistry(null),
|
||||
enforce ? CallerResolver.withLeadsAndMembers(identity, false, null,
|
||||
Map::of, new MemberRegistry(null)) : null,
|
||||
metrics, BridgeMcp.CapacitySource.none(), new BridgeMcp.HealthCoverageSource(() -> "off"));
|
||||
metrics, BridgeMcp.CapacitySource.none(), new BridgeMcp.HealthCoverageSource(() -> "off"),
|
||||
BridgeMcp.QuarantineSource.none());
|
||||
return mcp;
|
||||
}
|
||||
|
||||
|
||||
@@ -16,6 +16,7 @@ import dev.ltms.bridged.inject.MemberPresence;
|
||||
import dev.ltms.bridged.session.MemberSession;
|
||||
import dev.ltms.bridged.session.WorktreeRequest;
|
||||
import dev.ltms.bridged.member.ClaudeCodeLauncher;
|
||||
import dev.ltms.bridged.placement.BackendQuarantine;
|
||||
import io.modelcontextprotocol.spec.McpSchema;
|
||||
import dev.ltms.bridged.msg.InMemoryReplyInbox;
|
||||
import org.junit.jupiter.api.BeforeEach;
|
||||
@@ -437,11 +438,29 @@ class BridgeMcpTest {
|
||||
@Test
|
||||
void profilesListsConfiguredProfilesAndDefault() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
McpSchema.CallToolResult res = BridgeMcp.profiles(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
|
||||
McpSchema.CallToolResult res = BridgeMcp.profiles(
|
||||
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), BridgeMcp.QuarantineSource.none());
|
||||
assertNotEquals(Boolean.TRUE, res.isError());
|
||||
String out = textOf(res);
|
||||
assertTrue(out.contains("ltms-local"), out);
|
||||
assertTrue(out.contains("\"default\":\"ltms-local\""), out);
|
||||
assertFalse(out.contains("quarantined"), "no profile is quarantined, so the key is omitted: " + out);
|
||||
}
|
||||
|
||||
@Test
|
||||
void profilesReportsAQuarantinedCredential() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
|
||||
quarantine.quarantine("shared-openai");
|
||||
BridgeMcp.QuarantineSource source = new BridgeMcp.QuarantineSource(
|
||||
profile -> "ltms-local".equals(profile) ? "shared-openai" : null, quarantine);
|
||||
McpSchema.CallToolResult res = BridgeMcp.profiles(
|
||||
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), source);
|
||||
assertNotEquals(Boolean.TRUE, res.isError());
|
||||
String out = textOf(res);
|
||||
assertTrue(out.contains("\"quarantined\""), out);
|
||||
assertTrue(out.contains("shared-openai"), out);
|
||||
assertTrue(out.contains("\"quarantinedForSeconds\":1800"), out);
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
@@ -733,7 +733,7 @@ class ClaudeCodeLauncherTest {
|
||||
return new BridgedConfig.Profile(
|
||||
profile, baseUrl, "sonnet", null, "BRIDGED_WORKER_TOKEN",
|
||||
List.of("ccs", profile), "tab", "bridged-workers", "w #{n}", null, null, null,
|
||||
null, null, null, Map.of(), null, null, true);
|
||||
null, null, null, Map.of(), null, null, true, null);
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -819,7 +819,7 @@ class ClaudeCodeLauncherTest {
|
||||
null, null, null,
|
||||
Map.of("ANTHROPIC_BASE_URL", "http://evil.example.com",
|
||||
"ANTHROPIC_AUTH_TOKEN", "sk-ant-bad", "JAVA_HOME", "/opt/jdk"),
|
||||
null, null, true);
|
||||
null, null, true, null);
|
||||
new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
|
||||
_ -> "would-be-token").spawn();
|
||||
|
||||
@@ -10,11 +10,13 @@ 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.CharterReceipt;
|
||||
import dev.ltms.bridged.peer.MemberRole;
|
||||
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.BackendQuarantine;
|
||||
import dev.ltms.bridged.placement.PlacementException;
|
||||
import dev.ltms.bridged.placement.PlacementPolicies;
|
||||
import org.junit.jupiter.api.Test;
|
||||
@@ -26,6 +28,8 @@ import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
@@ -93,6 +97,7 @@ class CompositePeerLauncherTest {
|
||||
@Override public String id() { return "pane-" + p; }
|
||||
@Override public String terminalId() { return "term-" + p; }
|
||||
@Override public String profile() { return p; }
|
||||
@Override public CharterReceipt charterReceipt() { return null; }
|
||||
};
|
||||
}
|
||||
|
||||
@@ -126,6 +131,13 @@ class CompositePeerLauncherTest {
|
||||
weight, maxLoad);
|
||||
}
|
||||
|
||||
private static BridgedConfig.Profile stubWorker(String profile, String credentialId) {
|
||||
return new BridgedConfig.Profile(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,
|
||||
null, null, credentialId);
|
||||
}
|
||||
|
||||
/**
|
||||
* An <em>order-preserving</em> profile map. Never {@code Map.of} here: its iteration order is
|
||||
* salted per JVM run, and the weighted policy breaks an exact-weight tie on candidate order —
|
||||
@@ -552,4 +564,105 @@ class CompositePeerLauncherTest {
|
||||
new SpawnRequest(null, null, null, null, null, MemberRole.DEV)));
|
||||
assertTrue(e.getMessage().contains("maxLoad"), e.getMessage());
|
||||
}
|
||||
|
||||
// ── CB-578 stage B: a BACKEND_EXHAUSTED classification quarantines the credential ──────────
|
||||
|
||||
@Test
|
||||
void explicitSpawnOntoAQuarantinedProfileIsRefusedNamingTheCredentialAndRemainingTime() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Map<String, BridgedConfig.Profile> profiles = ordered(
|
||||
"sol", stubWorker("sol", "shared-openai"),
|
||||
"terra", stubWorker("terra", "shared-openai"));
|
||||
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "sol", Set.of());
|
||||
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
|
||||
quarantine.quarantine("shared-openai");
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "sol", profiles,
|
||||
PlacementPolicies.fixed(), _ -> 0, null, quarantine);
|
||||
|
||||
PlacementException e = assertThrows(PlacementException.class,
|
||||
() -> composite.spawn(new SpawnRequest("sol", null, null)));
|
||||
assertTrue(e.getMessage().contains("sol"), "message names the profile: " + e.getMessage());
|
||||
assertTrue(e.getMessage().contains("shared-openai"), "message names the credential: " + e.getMessage());
|
||||
assertTrue(e.getMessage().contains("1800"), "message names roughly when it lifts: " + e.getMessage());
|
||||
assertEquals(0, adapter.spawnCount("sol"), "the quarantined profile is never delegated to");
|
||||
}
|
||||
|
||||
/**
|
||||
* The part CB-578 stage B calls out as easy to get wrong: sol and terra are two different
|
||||
* profiles sharing one OpenAI credential. Quarantining because of an exhaustion classified on
|
||||
* ONE of them must lock out the other too, or the fleet just walks onto the same dead account
|
||||
* under the sibling's name.
|
||||
*/
|
||||
@Test
|
||||
void twoProfilesSharingACredentialAreBothQuarantinedByOneExhaustionEvent() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Map<String, BridgedConfig.Profile> profiles = ordered(
|
||||
"sol", stubWorker("sol", "shared-openai"),
|
||||
"terra", stubWorker("terra", "shared-openai"));
|
||||
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "sol", Set.of());
|
||||
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
|
||||
// Only "sol" was classified BACKEND_EXHAUSTED — but the two profiles share one credential.
|
||||
quarantine.quarantine("shared-openai");
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "sol", profiles,
|
||||
PlacementPolicies.fixed(), _ -> 0, null, quarantine);
|
||||
|
||||
assertThrows(PlacementException.class, () -> composite.spawn(new SpawnRequest("sol", null, null)),
|
||||
"sol was the one classified exhausted");
|
||||
assertThrows(PlacementException.class, () -> composite.spawn(new SpawnRequest("terra", null, null)),
|
||||
"terra shares sol's credential, so it must be locked out too");
|
||||
assertEquals(0, adapter.spawnCount("sol"));
|
||||
assertEquals(0, adapter.spawnCount("terra"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void placementSkipsAQuarantinedProfileAndRoutesToAnUnquarantinedOne() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Map<String, BridgedConfig.Profile> profiles = ordered(
|
||||
"sol", stubWorker("sol", "shared-openai"),
|
||||
"b", stubWorker("b"));
|
||||
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "sol", Set.of());
|
||||
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
|
||||
quarantine.quarantine("shared-openai");
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "sol", profiles,
|
||||
PlacementPolicies.weighted(), _ -> 0, null, quarantine);
|
||||
|
||||
PeerHandle h = composite.spawn(new SpawnRequest(null, null, null));
|
||||
assertEquals("b", h.profile(), "sol is quarantined, so an unqualified spawn must land on b");
|
||||
assertEquals(0, adapter.spawnCount("sol"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void aQuarantineLiftsOnTheInjectedClockAndTheProfileBecomesSpawnableAgain() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Map<String, BridgedConfig.Profile> profiles = ordered(
|
||||
"sol", stubWorker("sol", "shared-openai"),
|
||||
"b", stubWorker("b"));
|
||||
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "sol", Set.of());
|
||||
AtomicLong nowNanos = new AtomicLong(0L);
|
||||
BackendQuarantine quarantine = new BackendQuarantine(nowNanos::get, TimeUnit.MINUTES.toNanos(30));
|
||||
quarantine.quarantine("shared-openai");
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "sol", profiles,
|
||||
PlacementPolicies.fixed(), _ -> 0, null, quarantine);
|
||||
|
||||
assertThrows(PlacementException.class, () -> composite.spawn(new SpawnRequest("sol", null, null)),
|
||||
"still inside the cooldown");
|
||||
|
||||
nowNanos.set(TimeUnit.MINUTES.toNanos(31));
|
||||
|
||||
PeerHandle h = composite.spawn(new SpawnRequest("sol", null, null));
|
||||
assertEquals("sol", h.profile(), "the cooldown expired on the injected clock — sol is spawnable again");
|
||||
assertEquals(1, adapter.spawnCount("sol"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void aFleetWithNoExhaustedPatternAnywhereBehavesExactlyAsBeforeQuarantineExisted() {
|
||||
// BackendQuarantine.none() is the inert stand-in every constructor already defaults to when
|
||||
// no quarantine is wired — the 2/5/6-arg constructors used throughout this file all exercise
|
||||
// it. This test pins that an explicit .none() also never refuses a spawn, for any profile.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
PeerLauncher composite = composite(herdr);
|
||||
|
||||
assertDoesNotThrow(() -> composite.spawn(new SpawnRequest("claude", null, null)));
|
||||
assertDoesNotThrow(() -> composite.spawn(new SpawnRequest("gemini", null, null)));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,13 +1,19 @@
|
||||
package dev.ltms.bridged.member;
|
||||
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.Logger;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import dev.ltms.bridged.config.BridgedConfig;
|
||||
import dev.ltms.bridged.herdr.AgentControl;
|
||||
import dev.ltms.bridged.herdr.FakeHerdr;
|
||||
import dev.ltms.bridged.herdr.WorkspaceControl;
|
||||
import dev.ltms.bridged.peer.Capability;
|
||||
import dev.ltms.bridged.peer.CharterReceipt;
|
||||
import dev.ltms.bridged.peer.MemberRole;
|
||||
import dev.ltms.bridged.peer.SpawnRequest;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
@@ -17,6 +23,8 @@ import java.util.concurrent.atomic.AtomicReference;
|
||||
import java.util.function.Supplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
class HerdrPeerLauncherCharterTest {
|
||||
|
||||
@@ -38,10 +46,72 @@ class HerdrPeerLauncherCharterTest {
|
||||
"a role charter does not depend on an MCP mount");
|
||||
}
|
||||
|
||||
@Test
|
||||
void panePlacementSpawnLogNeverContainsTheCharterText() {
|
||||
// A pane-placement spawn used to log the whole argv (CB-571), and the charter travels
|
||||
// inside argv — so the charter text leaked to the daemon log. Prove the legacy pane path
|
||||
// now redacts it to its digest.
|
||||
String secret = "TOP SECRET charter marker 99x"; // distinctive, so a leak is unambiguous
|
||||
AtomicReference<BridgedConfig.Fleet> fleet = new AtomicReference<>(fleet(Map.of("dev", secret)));
|
||||
CharterArgLauncher launcher = new CharterArgLauncher(fleet::get);
|
||||
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
|
||||
Level previous = logger.getLevel();
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
logger.setLevel(Level.INFO); // the test logback sets dev.ltms.bridged to WARN; a leak lives at INFO
|
||||
try {
|
||||
launcher.spawn(new SpawnRequest("mcp", null, null, null, null, MemberRole.DEV));
|
||||
|
||||
String all = String.join("\n", appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
|
||||
assertFalse(all.contains(secret),
|
||||
"the pane-placement spawn log must not contain the charter text; got:\n" + all);
|
||||
// The "mcp" profile composes role + reply charter; the digest must match that composed
|
||||
// string (the exact bytes the adapter receives), proving the redaction hashes and
|
||||
// removes the real, full charter — not some placeholder.
|
||||
String composed = secret + "\n\n" + HerdrPeerLauncher.REPLY_CHARTER;
|
||||
assertTrue(all.contains("<charter sha256=" + CharterReceipt.digestOf(composed) + ">"),
|
||||
"the charter argv argument should be replaced by its digest; got:\n" + all);
|
||||
} finally {
|
||||
logger.setLevel(previous);
|
||||
logger.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
|
||||
private static BridgedConfig.Fleet fleet(Map<String, String> charters) {
|
||||
return new BridgedConfig.Fleet(Map.of(), Map.of(), Map.of(), Map.of(), charters, null);
|
||||
}
|
||||
|
||||
private static BridgedConfig.Profile profile(String name, String mcpUrl) {
|
||||
return new BridgedConfig.Profile(name, "http://gx00.gw:8000", null, null,
|
||||
"BRIDGED_WORKER_TOKEN", List.of("test"), "pane", null, null, mcpUrl, null, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* A launcher whose {@code buildLaunch} hands the composed charter to herdr as one argv element
|
||||
* (what the claude-cod adapter does), so a pane-placement spawn log would print it unless the
|
||||
* base redacts it.
|
||||
*/
|
||||
private static final class CharterArgLauncher extends HerdrPeerLauncher {
|
||||
|
||||
CharterArgLauncher(Supplier<BridgedConfig.Fleet> fleet) {
|
||||
super("test", new AgentControl(new FakeHerdr()), new WorkspaceControl(new FakeHerdr()),
|
||||
Map.of("mcp", profile("mcp", "http://bridge")),
|
||||
"mcp", _ -> null, 0, () -> 0L, () -> { }, fleet);
|
||||
}
|
||||
|
||||
@Override
|
||||
protected Launch buildLaunch(BridgedConfig.Profile cfg, LaunchSpec spec) {
|
||||
return new Launch(Map.of(), List.of("test", spec.charter() == null ? "none" : spec.charter()));
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<Capability> capabilities() {
|
||||
return Set.of();
|
||||
}
|
||||
}
|
||||
|
||||
private static final class CapturingLauncher extends HerdrPeerLauncher {
|
||||
private final List<LaunchSpec> specs = new ArrayList<>();
|
||||
|
||||
@@ -62,10 +132,5 @@ class HerdrPeerLauncherCharterTest {
|
||||
public Set<Capability> capabilities() {
|
||||
return Set.of();
|
||||
}
|
||||
|
||||
private static BridgedConfig.Profile profile(String name, String mcpUrl) {
|
||||
return new BridgedConfig.Profile(name, "http://gx00.gw:8000", null, null,
|
||||
"BRIDGED_WORKER_TOKEN", List.of("test"), "pane", null, null, mcpUrl, null, null);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -7,6 +7,7 @@ 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.CharterReceipt;
|
||||
import dev.ltms.bridged.peer.PeerHandle;
|
||||
import dev.ltms.bridged.peer.PeerUnreachableException;
|
||||
import dev.ltms.bridged.peer.SpawnRequest;
|
||||
@@ -322,6 +323,26 @@ class OpenCodeLauncherTest {
|
||||
assertFalse(herdr.called("agent.get"), "no polling when the gate is disabled");
|
||||
}
|
||||
|
||||
@Test
|
||||
void handleCarriesTheRealCharterReceiptNotTheInterfaceDefault(@TempDir Path root) {
|
||||
// The base's WorkerHandle computes a real CharterReceipt (CB-571), but the opencode adapter
|
||||
// wraps it in SessionAwareHandle for lazy session discovery. Before this fix that decorator
|
||||
// did not override charterReceipt(), so it silently inherited PeerHandle's `null` default
|
||||
// and the real receipt sitting on its delegate was lost.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
BridgedConfig.Fleet fleet = new BridgedConfig.Fleet(Map.of(), Map.of(), Map.of(), Map.of(),
|
||||
Map.of("dev", "role rule"), null);
|
||||
PeerHandle handle = service(herdr, root,
|
||||
opencodeCfg("google/gemini-2.5-pro", "http://127.0.0.1:8765/mcp", null), () -> fleet)
|
||||
.spawn(new SpawnRequest(null, null, null));
|
||||
|
||||
assertNotNull(handle.charterReceipt(),
|
||||
"an opencode spawn's charterReceipt() must not silently be null");
|
||||
String composed = "role rule\n\n" + HerdrPeerLauncher.REPLY_CHARTER;
|
||||
assertEquals(CharterReceipt.digestOf(composed), handle.charterReceipt().charterSha256(),
|
||||
"the receipt on the wrapped handle must match the exact composed charter bytes");
|
||||
}
|
||||
|
||||
// --- CB-508: pinned OpenAI-compatible endpoint (e.g. a local vLLM) ---------------------------
|
||||
|
||||
/** A profile with a baseUrl but no model provider prefix cannot be resolved — fail loudly. */
|
||||
|
||||
@@ -5,6 +5,8 @@ import dev.ltms.bridged.herdr.AgentStatus;
|
||||
import dev.ltms.bridged.herdr.FakeHerdr;
|
||||
import dev.ltms.bridged.herdr.HerdrException;
|
||||
import dev.ltms.bridged.inject.CompletionResolver;
|
||||
import dev.ltms.bridged.inject.ExhaustedPatternLookup;
|
||||
import dev.ltms.bridged.inject.ExhaustionSink;
|
||||
import dev.ltms.bridged.mcp.PrimaryRegistry;
|
||||
import dev.ltms.bridged.inject.Injector;
|
||||
import org.junit.jupiter.api.BeforeEach;
|
||||
@@ -31,7 +33,8 @@ class MessageServiceTest {
|
||||
private final FakeHerdr herdr = new FakeHerdr().readText("BUILD GREEN: 391 files");
|
||||
private final AgentControl agents = new AgentControl(herdr);
|
||||
private final Rendezvous rendezvous = new Rendezvous();
|
||||
private final CompletionResolver completion = new CompletionResolver(agents, rendezvous);
|
||||
private final CompletionResolver completion =
|
||||
new CompletionResolver(agents, rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
private final Injector injector = new Injector(agents, completion);
|
||||
private final InMemoryReplyInbox inbox = new InMemoryReplyInbox();
|
||||
private final MessageService messages = new MessageService(agents, injector, rendezvous, inbox);
|
||||
|
||||
@@ -114,6 +114,17 @@ class RendezvousTest {
|
||||
"the first resolution wins; the stored value is unchanged");
|
||||
}
|
||||
|
||||
@Test
|
||||
void resolveExhaustedTwiceIsANoOpTheSecondTime() {
|
||||
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.open(W);
|
||||
assertTrue(rendezvous.resolveExhausted(waiter, "first reason"), "the first classification resolves");
|
||||
assertFalse(rendezvous.resolveExhausted(waiter, "second reason"),
|
||||
"a second exhausted resolution on an already-resolved waiter returns false");
|
||||
assertEquals(Rendezvous.Kind.BACKEND_EXHAUSTED, waiter.getNow(null).kind());
|
||||
assertEquals("first reason", waiter.getNow(null).text(),
|
||||
"the first resolution wins; the stored value is unchanged");
|
||||
}
|
||||
|
||||
@Test
|
||||
void closeAskRemovesTheTurn() {
|
||||
Rendezvous.AskTicket t = rendezvous.openAsk(W);
|
||||
|
||||
@@ -0,0 +1,57 @@
|
||||
package dev.ltms.bridged.peer;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
class CharterReceiptTest {
|
||||
|
||||
@Test
|
||||
void recordsNoCharterConfiguredDistinctFromCharterDelivered() {
|
||||
// No role charter configured — only the reply charter is composed. Source is "none", but
|
||||
// text was still delivered, so absent() is false and the digest is present.
|
||||
CharterReceipt viaReply = CharterReceipt.compose(MemberRole.DEV, "s", null, "reply charter");
|
||||
// A role charter was configured AND delivered.
|
||||
CharterReceipt delivered = CharterReceipt.compose(MemberRole.DEV, "s", "role charter",
|
||||
"role charter\n\nreply charter");
|
||||
|
||||
// The two cases must not collapse: the no-role-charter case reports "none", the delivered
|
||||
// case reports the config key, and their digests differ.
|
||||
assertEquals(CharterReceipt.NO_SOURCE, viaReply.charterSource());
|
||||
assertEquals("fleet.charters.dev", delivered.charterSource());
|
||||
assertNotEquals(viaReply.charterSource(), delivered.charterSource());
|
||||
assertNotEquals(viaReply.charterSha256(), delivered.charterSha256());
|
||||
|
||||
// Both actually delivered text — the distinction is the source and digest, not absence.
|
||||
assertFalse(viaReply.absent());
|
||||
assertFalse(delivered.absent());
|
||||
}
|
||||
|
||||
@Test
|
||||
void recordsExplicitAbsenceWhenNoCharterIsComposed() {
|
||||
CharterReceipt none = CharterReceipt.compose(MemberRole.REVIEWER, "s", null, null);
|
||||
assertTrue(none.absent());
|
||||
assertNull(none.charterSha256());
|
||||
assertEquals(0, none.charterBytes());
|
||||
assertEquals(CharterReceipt.NO_SOURCE, none.charterSource(),
|
||||
"no configured charter and nothing composed still reports a source, never a gap");
|
||||
}
|
||||
|
||||
@Test
|
||||
void digestIsStableForSameTextAndDiffersForDifferentText() {
|
||||
assertEquals(CharterReceipt.digestOf("charter-aaa"), CharterReceipt.digestOf("charter-aaa"),
|
||||
"the same text must always produce the same digest");
|
||||
assertNotEquals(CharterReceipt.digestOf("charter-aaa"), CharterReceipt.digestOf("charter-bbb"),
|
||||
"different text must produce a different digest");
|
||||
assertNull(CharterReceipt.digestOf(""), "blank text carries no digest");
|
||||
|
||||
// The record's fingerprint matches the standalone digest for the same composed string.
|
||||
CharterReceipt r = CharterReceipt.compose(MemberRole.DEV, "s", "role", "the composed text");
|
||||
assertEquals(CharterReceipt.digestOf("the composed text"), r.charterSha256());
|
||||
assertFalse(r.absent());
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,108 @@
|
||||
package dev.ltms.bridged.placement;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.Map;
|
||||
import java.util.OptionalLong;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertThrows;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* CB-578 stage B: the credential-keyed quarantine tracker itself, isolated from placement/spawn
|
||||
* wiring (that's {@code CompositePeerLauncherTest}). The clock is a plain {@link AtomicLong} of
|
||||
* nanos so expiry is exercised without a real sleep.
|
||||
*/
|
||||
class BackendQuarantineTest {
|
||||
|
||||
@Test
|
||||
void aFreshCredentialIsNotQuarantined() {
|
||||
BackendQuarantine q = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
|
||||
assertFalse(q.isQuarantined("shared-openai"));
|
||||
assertEquals(OptionalLong.empty(), q.remainingSeconds("shared-openai"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void quarantineBlocksTheCredentialForTheFullCooldown() {
|
||||
BackendQuarantine q = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
|
||||
q.quarantine("shared-openai");
|
||||
|
||||
assertTrue(q.isQuarantined("shared-openai"));
|
||||
assertEquals(OptionalLong.of(1800L), q.remainingSeconds("shared-openai"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void onlyTheQuarantinedCredentialIsAffected() {
|
||||
BackendQuarantine q = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
|
||||
q.quarantine("shared-openai");
|
||||
|
||||
assertFalse(q.isQuarantined("some-other-credential"),
|
||||
"an unrelated credential must not be swept into the quarantine");
|
||||
}
|
||||
|
||||
@Test
|
||||
void expiresOnTheInjectedClock() {
|
||||
AtomicLong now = new AtomicLong(0L);
|
||||
BackendQuarantine q = new BackendQuarantine(now::get, TimeUnit.MINUTES.toNanos(30));
|
||||
q.quarantine("shared-openai");
|
||||
assertTrue(q.isQuarantined("shared-openai"));
|
||||
|
||||
now.set(TimeUnit.MINUTES.toNanos(29));
|
||||
assertTrue(q.isQuarantined("shared-openai"), "still inside the cooldown");
|
||||
|
||||
now.set(TimeUnit.MINUTES.toNanos(31));
|
||||
assertFalse(q.isQuarantined("shared-openai"), "the cooldown has elapsed on the injected clock");
|
||||
assertEquals(OptionalLong.empty(), q.remainingSeconds("shared-openai"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void aRepeatQuarantineCallRestartsTheCooldownAtFullLength() {
|
||||
AtomicLong now = new AtomicLong(0L);
|
||||
BackendQuarantine q = new BackendQuarantine(now::get, TimeUnit.MINUTES.toNanos(30));
|
||||
q.quarantine("shared-openai");
|
||||
|
||||
now.set(TimeUnit.MINUTES.toNanos(20));
|
||||
q.quarantine("shared-openai");
|
||||
|
||||
now.set(TimeUnit.MINUTES.toNanos(45)); // 25 min after the second call, 45 after the first
|
||||
assertTrue(q.isQuarantined("shared-openai"),
|
||||
"a fresh exhaustion resets the cooldown to full length, not the earlier shorter wait");
|
||||
}
|
||||
|
||||
@Test
|
||||
void activeRemainingSecondsListsOnlyStillQuarantinedCredentials() {
|
||||
AtomicLong now = new AtomicLong(0L);
|
||||
BackendQuarantine q = new BackendQuarantine(now::get, TimeUnit.MINUTES.toNanos(30));
|
||||
q.quarantine("shared-openai");
|
||||
q.quarantine("another-credential");
|
||||
|
||||
now.set(TimeUnit.MINUTES.toNanos(31));
|
||||
q.quarantine("shared-openai"); // re-quarantined after the first one expired
|
||||
|
||||
Map<String, Long> active = q.activeRemainingSeconds();
|
||||
assertEquals(Map.of("shared-openai", 1800L), active,
|
||||
"the expired credential is dropped; the re-quarantined one is reported");
|
||||
}
|
||||
|
||||
@Test
|
||||
void noneReportsNothingQuarantinedWhenNeverToldTo() {
|
||||
// .none() is the stand-in for a caller whose code path never calls #quarantine at all (e.g.
|
||||
// the 2/5/6-arg CompositePeerLauncher constructors) — not a guarantee that a call to
|
||||
// #quarantine on it is a no-op. Left alone, as those call sites leave it, nothing is ever
|
||||
// quarantined.
|
||||
BackendQuarantine q = BackendQuarantine.none();
|
||||
|
||||
assertFalse(q.isQuarantined("anything"));
|
||||
assertTrue(q.activeRemainingSeconds().isEmpty());
|
||||
}
|
||||
|
||||
@Test
|
||||
void aNonPositiveCooldownIsRejected() {
|
||||
assertThrows(IllegalArgumentException.class, () -> new BackendQuarantine(() -> 0L, 0L));
|
||||
assertThrows(IllegalArgumentException.class, () -> new BackendQuarantine(() -> 0L, -1L));
|
||||
}
|
||||
}
|
||||
@@ -24,7 +24,13 @@ class PlacementPolicyTest {
|
||||
private static PlacementContext ctx(List<PlacementCandidate> candidates,
|
||||
Function<String, Integer> liveCount,
|
||||
Set<String> unreachable) {
|
||||
return new PlacementContext("b", candidates, liveCount, unreachable);
|
||||
return ctx(candidates, liveCount, unreachable, Set.of());
|
||||
}
|
||||
|
||||
private static PlacementContext ctx(List<PlacementCandidate> candidates,
|
||||
Function<String, Integer> liveCount,
|
||||
Set<String> unreachable, Set<String> quarantined) {
|
||||
return new PlacementContext("b", candidates, liveCount, unreachable, quarantined);
|
||||
}
|
||||
|
||||
private static PlacementContext ctx(List<PlacementCandidate> candidates,
|
||||
@@ -46,17 +52,37 @@ class PlacementPolicyTest {
|
||||
PlacementPolicy policy = PlacementPolicies.fixed();
|
||||
PlacementContext ctx = new PlacementContext(null,
|
||||
List.of(PlacementCandidate.profile("a"), PlacementCandidate.profile("b")),
|
||||
noSessions(), Set.of());
|
||||
noSessions(), Set.of(), 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());
|
||||
PlacementContext ctx = new PlacementContext(null, List.of(), noSessions(), Set.of(), Set.of());
|
||||
assertThrows(PlacementException.class, () -> policy.select(ctx));
|
||||
}
|
||||
|
||||
@Test
|
||||
void fixedSkipsQuarantinedDefault() {
|
||||
PlacementPolicy policy = PlacementPolicies.fixed();
|
||||
PlacementContext ctx = new PlacementContext("b",
|
||||
List.of(PlacementCandidate.profile("a"), PlacementCandidate.profile("b")),
|
||||
noSessions(), Set.of(), Set.of("b"));
|
||||
assertEquals("a", policy.select(ctx).profile(),
|
||||
"the default 'b' is quarantined, so fixed falls through to the first un-quarantined candidate");
|
||||
}
|
||||
|
||||
@Test
|
||||
void fixedThrowsWhenDefaultAndEveryCandidateQuarantined() {
|
||||
PlacementPolicy policy = PlacementPolicies.fixed();
|
||||
PlacementContext ctx = new PlacementContext("b",
|
||||
List.of(PlacementCandidate.profile("a"), PlacementCandidate.profile("b")),
|
||||
noSessions(), Set.of(), Set.of("a", "b"));
|
||||
PlacementException e = assertThrows(PlacementException.class, () -> policy.select(ctx));
|
||||
assertTrue(e.getMessage().contains("quarantined"), e.getMessage());
|
||||
}
|
||||
|
||||
@Test
|
||||
void roundRobinCyclesThroughAvailableProfiles() {
|
||||
PlacementPolicy policy = PlacementPolicies.roundRobin();
|
||||
@@ -82,6 +108,18 @@ class PlacementPolicyTest {
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void roundRobinSkipsQuarantinedProfiles() {
|
||||
PlacementPolicy policy = PlacementPolicies.roundRobin();
|
||||
List<PlacementCandidate> candidates = List.of(
|
||||
PlacementCandidate.profile("a"),
|
||||
PlacementCandidate.profile("b"));
|
||||
PlacementContext ctx = ctx(candidates, noSessions(), Set.of(), Set.of("a"));
|
||||
for (int i = 0; i < 5; i++) {
|
||||
assertEquals("b", policy.select(ctx).profile(), "a is quarantined, so every pick lands on b");
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void roundRobinThrowsWhenAllAtMaxLoad() {
|
||||
PlacementPolicy policy = PlacementPolicies.roundRobin();
|
||||
@@ -150,6 +188,29 @@ class PlacementPolicyTest {
|
||||
assertTrue(e.getMessage().contains("maxLoad"), e.getMessage());
|
||||
}
|
||||
|
||||
@Test
|
||||
void weightedSkipsQuarantinedProfile() {
|
||||
PlacementPolicy policy = PlacementPolicies.weighted();
|
||||
List<PlacementCandidate> candidates = List.of(
|
||||
PlacementCandidate.profile("a", 1.0f, null),
|
||||
PlacementCandidate.profile("b", 1.0f, null));
|
||||
PlacementContext ctx = ctx(candidates, noSessions(), Set.of(), Set.of("a"));
|
||||
for (int i = 0; i < 5; i++) {
|
||||
assertEquals("b", policy.select(ctx).profile(), "a is quarantined, so every pick lands on b");
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void weightedThrowsWhenAllQuarantined() {
|
||||
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(), Set.of("a", "b"))));
|
||||
assertTrue(e.getMessage().contains("quarantined"), e.getMessage());
|
||||
}
|
||||
|
||||
@Test
|
||||
void weightedThrowsWhenAllUnreachable() {
|
||||
PlacementPolicy policy = PlacementPolicies.weighted();
|
||||
|
||||
@@ -11,6 +11,8 @@ import dev.ltms.bridged.herdr.FakeHerdr;
|
||||
import dev.ltms.bridged.herdr.WorkspaceControl;
|
||||
import dev.ltms.bridged.member.ClaudeCodeLauncher;
|
||||
import dev.ltms.bridged.msg.TestTurnTokens;
|
||||
import dev.ltms.bridged.peer.CharterReceipt;
|
||||
import dev.ltms.bridged.peer.MemberRole;
|
||||
import dev.ltms.bridged.peer.PeerUnreachableException;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.LoggerFactory;
|
||||
@@ -91,6 +93,24 @@ class SessionManagerTest {
|
||||
assertEquals(2, sessions.roster().size(), "both sessions are registered");
|
||||
}
|
||||
|
||||
@Test
|
||||
void rosterViewExposesTheCharterReceiptButNeverTheCharterText() {
|
||||
// The roster (bridge_list and GET /members both render through rosterView) must let a lead
|
||||
// see which charter a member got, without ever carrying the charter prose itself (CB-571).
|
||||
MemberSession s = new MemberSession("p1", "term1", "prof", MemberRole.DEV, "/cwd", null,
|
||||
0, 0, 0, MemberSession.State.READY, null, null,
|
||||
CharterReceipt.compose(MemberRole.DEV, "prof", "role charter", "role charter\n\nreply"));
|
||||
|
||||
Map<String, Object> view = SessionManager.rosterView(s, null);
|
||||
|
||||
assertEquals("fleet.charters.dev", view.get("charterSource"),
|
||||
"the config key that supplied the role charter is reported");
|
||||
assertEquals(CharterReceipt.digestOf("role charter\n\nreply"), view.get("charterSha256"),
|
||||
"the digest of the exact composed charter bytes is reported");
|
||||
assertFalse(view.values().toString().contains("role charter"),
|
||||
"the roster row must not embed the charter text itself");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aNullTerminalFromThePrimaryIsANoOpEvenWithSessionsRegistered() {
|
||||
// The primary resolves to a Principal with no terminal, and BridgeMcp's context extractor
|
||||
|
||||
Reference in New Issue
Block a user