Compare commits

...

14 Commits

Author SHA1 Message Date
Dai Ha 8fd2d7e5e7 CB-581: fail-safe release() so a throw never orphans the pane or aborts reapIdle
CI / build (pull_request) Successful in 52s
CI / contract (pull_request) Successful in 1m1s
hasUncommitted shells out to git and can throw; release() now catches that
inside a try/finally so notifyReleased and launcher.stop always run, and
defaults to preserving the worktree on a throw (can't tell dirty vs clean,
so don't risk deleting unrecoverable work). reapIdle wraps each per-session
release in try/catch, matching drainAll, so one bad session no longer
skips the rest of the reaping pass.
2026-08-15 10:09:32 +02:00
Dai Ha c00a86b32c Merge CB-571: charter receipt on the roster, with no silent null for OpenCode
CI / contract (push) Successful in 42s
CI / build (push) Successful in 1m19s
Verified by the lead: own build of this branch merged onto main — 710 tests,
BUILD SUCCESS, exit 0.

Supersedes PR #54. A reviewer found that OpenCodeLauncher.SessionAwareHandle
wrapped the base WorkerHandle but never overrode charterReceipt(), so it
inherited the interface default of null while the real receipt sat on its
delegate — sol and terra would have shown no charterSource/charterSha256 on
the roster while Claude Code members showed both.

The fix is the root one, not the one-line override: PeerHandle.charterReceipt()
is no longer a default, so the compiler forces every implementation to answer.
This repo had shipped that same class of defect — a defaulted dependency that
compiles, passes tests, and quietly turns a feature off — eight times before
this one.
2026-08-15 10:02:36 +02:00
Dai Ha d3ae0350a2 Merge remote-tracking branch 'origin/main' into fix-charter 2026-08-15 10:01:10 +02:00
Dai Ha 2c2196a1f1 Merge CB-578 stage A: classify a usage-limit refusal instead of a completed reply
CI / build (push) Successful in 55s
CI / contract (push) Successful in 1m24s
Verified by the lead: own build of the branch merged onto main — 704 tests,
BUILD SUCCESS, exit 0. No vendor wording in any Java source (grep clean); a
profile with no exhaustedPattern keeps today's completion-fallback path exactly.

Accepted the implementer's deviation from the brief. The brief asked for a
terminal HealthState BACKEND_EXHAUSTED. FleetHealth.decide() is a pure
classifier over HealthSnapshot, which carries only booleans, AgentStatus and
MemberSession.State — never pane text. FleetHealth's own javadoc already says
ERROR_ON_SCREEN 'is not decided yet because it needs a bounded pane detection
read'. A second declared-but-unproduced value would repeat a gap the class
already documents as a problem, so the signal was put where the evidence
actually lives: the completion scrape.
2026-08-15 10:00:34 +02:00
Dai Ha 35ade14630 CB-571: make PeerHandle.charterReceipt() abstract, fix OpenCode adapter's silent null
CI / build (pull_request) Successful in 50s
CI / contract (pull_request) Successful in 1m21s
SessionAwareHandle wrapped the base's WorkerHandle but never overrode
charterReceipt(), so it silently inherited the interface default (null)
while the real receipt sat on its delegate. sol/terra never got a
charterSource/charterSha256 roster row.

Deletes the default so every PeerHandle must answer explicitly; the
compiler now catches this class of gap instead of a roster field
quietly going missing.
2026-08-15 09:55:52 +02:00
Dai Ha 541df87272 CB-578 stage A: classify a usage-limit refusal instead of a completed reply
CI / build (pull_request) Failing after 59s
CI / contract (pull_request) Successful in 1m10s
A backend that refuses on a subscription usage limit leaves the pane healthy but
the turn ends with no bridge_reply; the completion fallback used to scrape and
hand that refusal back as if it were a real answer. CompletionResolver now
matches the scrape against a per-profile exhaustedPattern (config, never a
vendor string) and resolves the send as Rendezvous.Kind/Outcome.BACKEND_EXHAUSTED
with a reason carrying the matched line, kept distinct from GONE/WORKER_FAILED.
A profile with no pattern configured is unaffected. Coverage is logged at
startup via CompletionResolver.coverage(...), naming which profiles have a
pattern and which don't, following FleetHealthMonitor.coverage's pattern.
2026-08-15 09:55:06 +02:00
Dai Ha 337b6ccd6e Merge remote-tracking branch 'origin/main' into fix-charter
# Conflicts:
#	bridged/src/test/java/dev/ltms/bridged/session/SessionManagerTest.java
2026-08-15 09:52:39 +02:00
ltms 42f46dfe9a Merge CB-579: resolve a lead by its tab name, drop the terminal-id pin
CI / build (push) Successful in 58s
CI / contract (push) Successful in 1m4s
Verified by the lead: merged onto main (9088d2b) in a scratch worktree, mvn -f bridged/pom.xml
clean install unpiped — MVN_EXIT=0, Tests run: 696, Failures: 0, BUILD SUCCESS. main alone measures
692, so this adds 4 net tests. Merges cleanly; Bridged.java auto-merged against CB-580.

Reviewed by the lead reading the full production diff and the three test files. The member's own
report was lost to the idle reaper before collection, so there was no author write-up.

Closes the live bug: LeadTabScanner.scan() no longer merges the config pin over the scan result, and
the cache no longer seeds from it, so a lead disappears once its tab is gone. The ghost this fixes
had begun throwing agent_not_found from ReplyPushLoop.decide on every tick, against a dead terminal
that still owned two live members.

Beyond the brief, and correct: `tab` is required for every leader, not only non-creatable ones,
because LeadLauncher also uses it to label a tab it creates. terminal: is rejected by a raw-YAML
check rather than by record shape — the only way to beat @JsonIgnoreProperties(ignoreUnknown = true).
The silent-default trap is avoided: the back-compat constructor still takes `tab` positionally.

Behaviour change worth knowing: the scanner's initial cache is now empty instead of the config pins,
so a herdr failure on the very first scan yields no leads until a scan succeeds. That is unavoidable
once the pins are gone, it fails loudly rather than silently, and it is covered by
aFailedFirstScanReturnsEmptyRatherThanThrowing.

Operator action required: bridged.yaml must replace fleet.leaders.<name>.terminal with tab. Already
done for this deployment.
2026-08-15 09:29:00 +02:00
ltms 9088d2b2c5 Merge CB-580: fail a ticket when its member reaches a terminal health state
CI / contract (push) Successful in 41s
CI / build (push) Successful in 1m32s
Verified by the lead on 0af902e: mvn -f bridged/pom.xml clean install, unpiped, in a scratch
worktree — MVN_EXIT=0, Tests run: 689, Failures: 0, Errors: 0, BUILD SUCCESS.

Reviewed by the lead reading the diff. The member's own report was lost to the idle reaper before
it was collected, so there was no author write-up to review against.

All three defects that sank 3b2f395 are absent:
  * failTarget is required — one constructor, Objects.requireNonNull, no defaulting overload
    anywhere in the repo, and Bridged.java:365 updated to pass messages::abandon.
  * The failure fires inside reportTransition after its `if (previous == next) return;` guard, so an
    unchanged tick cannot reach it.
  * No AtomicReference; MessageService is passed directly, so there is no empty window.

abandon(String target, String reason) confirmed as CB-568's target-wide operation that resolves
waiters as a failure rather than letting them time out.

Known limitation, accepted and covered by its own test: states.put records the new state before the
bounded retries run, so if all three attempts throw, the tickets stay pending and no later tick
retries. It is logged at WARN, and the retries carry no backoff.
2026-08-15 08:57:03 +02:00
ltms 500bfa2c33 Merge CB-576: release preserves a dirty worktree instead of deleting it
CI / build (push) Successful in 51s
CI / contract (push) Successful in 1m1s
Verified by the lead on b525b0f: mvn -f bridged/pom.xml clean install, unpiped, in a scratch
worktree — MVN_EXIT=0, Tests run: 686, Failures: 0, Errors: 0, BUILD SUCCESS.

Review accepted the required-interface-method shape (no defaulting overload) and the decision to
count untracked files as dirty — the work lost in the incident was a file that was never added.

One blocking defect was found and fixed in b525b0f: hasUncommitted called git with no existence
check, so a missing worktree threw WorktreeException from inside release() after registry.remove()
but before notifyReleased() and launcher.stop(), orphaning the pane and stranding a blocked
bridge_send caller. It now mirrors remove()'s already-gone tolerance.

Waived on merge, tracked as follow-up: a SessionManager-level test that teardown completes when the
worktree is gone, and the stronger fix behind it — reapIdle calls release() with no try/catch while
drainAll wraps it, so any exception in that window aborts the whole reaping pass.
2026-08-15 08:54:33 +02:00
Dai Ha 9ca9c43dfa CB-571: retag charter-receipt references from the taken CB-575
CI / contract (pull_request) Successful in 43s
CI / build (pull_request) Successful in 54s
CB-575 already names the merged MCP-cancellation-filter change, so the
charter-receipt comments used the wrong number. Retag to CB-571, the number
this work was authored against.
2026-08-15 08:48:57 +02:00
Dai Ha 1966c69994 CB-575: charter receipt on spawn, in the roster and in the logs
CI / contract (pull_request) Successful in 46s
CI / build (pull_request) Successful in 1m34s
Record a CharterReceipt (role, source, sha-256 digest, byte count) for every
launch, store it on the MemberSession, expose it in bridge_list and GET
/members, and log it at spawn as digest+role only. Redact the charter argv
argument in the legacy pane-placement spawn log so the charter text never
reaches the daemon log. The charter prose itself is never recorded.
2026-08-15 08:10:56 +02:00
Dai Ha 976eff8ad1 CB-579: resolve a lead by its tab name, drop the terminal-id pin
CI / build (pull_request) Successful in 52s
CI / contract (pull_request) Successful in 1m18s
Leader.terminal -> Leader.tab (exact tab label, case-insensitive match).
LeadTabScanner matches an exact tab->name map instead of stripping a
shared tabPrefix, and no longer merges configured leads into every
scan result -- a stale pin can no longer outlive its tab.
LeadLauncher.tabLabel() returns the configured tab directly; the
terminalId pinned-terminal fallback in liveLeads() is gone.
Config load now rejects a leftover fleet.leaders.*.terminal key
instead of silently ignoring it. primary.terminal is untouched.
2026-08-15 08:05:21 +02:00
Dai Ha 0af902ec43 CB-580: fail a ticket when its member reaches a terminal health state
CI / build (pull_request) Successful in 54s
CI / contract (pull_request) Successful in 1m5s
FleetHealthMonitor now requires a failTarget BiConsumer<String,String>
collaborator (no defaulting overload) and calls it exactly once when a
member transitions into GONE or NEVER_READY, via CB-568's idempotent
target-wide abandon() operation. The reason string names the real
terminal state. failTarget invocation retries up to
MAX_FAIL_TARGET_ATTEMPTS (3) within the same transition if it throws,
and never refires on a later tick where the state is unchanged.

Bridged.java wires messages::abandon as the production failTarget.
2026-08-15 07:54:28 +02:00
32 changed files with 1527 additions and 319 deletions
+26 -14
View File
@@ -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,12 @@ 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.
# 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 +177,7 @@ 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)
# 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
@@ -309,15 +319,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 +339,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 +349,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,7 @@ 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.Injector;
import dev.ltms.bridged.inject.StatusPoller;
import dev.ltms.bridged.inject.TurnListener;
@@ -60,6 +61,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;
/**
@@ -193,34 +195,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 +255,24 @@ 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()));
CompletionResolver completion = new CompletionResolver(agents, rendezvous, exhaustedPatterns);
// 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 +379,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)) {
@@ -181,6 +181,12 @@ 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.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record Profile(String profile, String baseUrl, String model,
@@ -193,7 +199,8 @@ public record BridgedConfig(
Map<String, String> env,
Float weight,
Integer maxLoad,
Boolean subscription) {
Boolean subscription,
String exhaustedPattern) {
/** Peer kind spawned by {@link dev.ltms.bridged.member.ClaudeCodeLauncher} (the default). */
public static final String KIND_CLAUDE_CODE = "claude-code";
@@ -236,6 +243,9 @@ 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;
}
/**
@@ -248,7 +258,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 +270,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 +283,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);
}
/** True when this profile is served by the Claude Code adapter (the default kind). */
@@ -312,7 +323,12 @@ 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);
}
/** True when this profile's CB-578 stage A backend-exhausted classification is configured. */
public boolean hasExhaustedPattern() {
return exhaustedPattern != null;
}
/** True when this profile's workers are granted a forge token to open their own PR (CB-302). */
@@ -458,24 +474,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 +517,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 +531,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 +732,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);
}
@@ -840,6 +857,7 @@ public record BridgedConfig(
try {
String yaml = Files.readString(path);
rejectRenamedTopLevelKeys(yaml);
rejectLeaderTerminalKey(yaml);
warnUnknownTopLevelKeys(yaml, path);
rejectDuplicateMemberSlots(yaml);
BridgedConfig cfg = YAML.readValue(yaml, BridgedConfig.class);
@@ -1064,6 +1082,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 {
@@ -1313,10 +1368,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,7 +30,11 @@ 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:},
@@ -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,7 @@ public final class CompletionResolver implements TurnListener {
private final AgentControl agents;
private final Rendezvous rendezvous;
private final ExhaustedPatternLookup exhaustedPatterns;
/**
* Per-target record of the turn currently in flight: the exact {@link Rendezvous} waiter its
@@ -78,9 +84,16 @@ 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()}).
*/
public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns) {
this.agents = agents;
this.rendezvous = rendezvous;
this.exhaustedPatterns = Objects.requireNonNull(exhaustedPatterns, "exhaustedPatterns");
}
@Override
@@ -157,11 +170,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 +198,22 @@ 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);
}
return;
}
}
String completion = clipped ? tail + "\n" + CLIPPED_PANE_TAIL_MARKER : tail;
if (rendezvous.resolveCompletion(waiter, completion)) {
inFlight.remove(target, turn);
@@ -232,6 +262,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;
}
}
@@ -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();
@@ -458,6 +458,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()
@@ -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();
}
@@ -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={}",
@@ -196,27 +197,44 @@ public final class SessionManager implements TurnListener {
MemberSession removed = registry.remove(paneId);
boolean preserveWorktree = cause == ReleaseCause.SHUTDOWN;
if (removed != null) {
memberLifecycle.released(removed.terminalId());
log.debug("releasing session pane={} terminal={} state={} cause={}",
removed.paneId(), removed.terminalId(), removed.state(), cause);
if (preserveWorktree && removed.worktree() != null) {
logPreservedForShutdown(removed);
} else if (removed.worktree() != null && worktrees.hasUncommitted(removed.worktree())) {
// CB-576: a release that would otherwise remove the worktree finds it holding
// uncommitted work the bridge cannot see. A worker that ends a turn without
// committing (normally because it stopped to ask a question or refused the turn)
// has its only copy of that work in the worktree. Remove would --force-delete it,
// so preserve the directory and tell an operator where to find it.
try {
memberLifecycle.released(removed.terminalId());
log.debug("releasing session pane={} terminal={} state={} cause={}",
removed.paneId(), removed.terminalId(), removed.state(), cause);
if (preserveWorktree && removed.worktree() != null) {
logPreservedForShutdown(removed);
} else if (removed.worktree() != null && worktrees.hasUncommitted(removed.worktree())) {
// CB-576: a release that would otherwise remove the worktree finds it holding
// uncommitted work the bridge cannot see. A worker that ends a turn without
// committing (normally because it stopped to ask a question or refused the turn)
// has its only copy of that work in the worktree. Remove would --force-delete it,
// so preserve the directory and tell an operator where to find it.
preserveWorktree = true;
log.warn("release {} preserves dirty worktree {} for pane={} terminal={}: "
+ "the worktree holds uncommitted changes that --force remove would destroy",
cause, removed.worktree(), removed.paneId(), removed.terminalId());
}
} catch (RuntimeException e) {
// CB-581: hasUncommitted shells out to `git status` and can throw on a non-zero
// exit. We can no longer tell whether the worktree holds uncommitted work, so fail
// toward the safe answer and preserve it — deleting on a guess can destroy work
// that has no other copy (CB-576), while keeping it on a false alarm only costs
// disk. The exception must not propagate: the pane still has to stop below.
preserveWorktree = true;
log.warn("release {} preserves dirty worktree {} for pane={} terminal={}: "
+ "the worktree holds uncommitted changes that --force remove would destroy",
cause, removed.worktree(), removed.paneId(), removed.terminalId());
log.warn("release {} could not tell whether worktree {} for pane={} terminal={} has "
+ "uncommitted changes; preserving it rather than risk destroying unsaved work: {}",
cause, removed.worktree(), removed.paneId(), removed.terminalId(), e.toString());
} finally {
// CB-516/CB-581: a send still waiting on this worker can never be answered now, no
// matter what happened above. Tell the listener BEFORE the pane is torn down, so a
// blocked caller fails fast with a real reason instead of sitting on a rendezvous
// nothing will ever resolve.
notifyReleased(removed.terminalId());
}
// CB-516: a send still waiting on this worker can never be answered now. Tell the
// listener BEFORE the pane is torn down, so a blocked caller fails fast with a real
// reason instead of sitting on a rendezvous nothing will ever resolve.
notifyReleased(removed.terminalId());
}
// CB-581: the pane must always stop, even if the dirty check above threw. A session removed
// from the registry with no pane stop is an orphaned pane — a live terminal burning a fleet
// slot that no longer appears in the roster and can never be reclaimed.
launcher.stop(paneId);
if (removed != null && !preserveWorktree && removed.worktree() != null) {
worktrees.remove(worktrees.repoRoot(removed.cwd()), removed.worktree());
@@ -355,7 +373,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 +440,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;
}
@@ -527,8 +556,15 @@ public final class SessionManager implements TurnListener {
log.debug("reaping idle session terminal={} pane={}: idle {}s exceeds the {}s ttl",
s.terminalId(), s.paneId(), TimeUnit.NANOSECONDS.toSeconds(idleNanos),
TimeUnit.NANOSECONDS.toSeconds(idleTtlNanos));
release(s.paneId());
reaped++;
// CB-581: one session that fails to release must not abort the whole reaping pass —
// match drainAll's per-session try/catch so the rest of the roster still gets reaped.
try {
release(s.paneId());
reaped++;
} catch (RuntimeException e) {
log.warn("reap failed for pane={} terminal={} worktree={}; continuing with "
+ "remaining sessions", s.paneId(), s.terminalId(), s.worktree(), e);
}
}
}
return reaped;
@@ -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());
@@ -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());
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());
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());
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());
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());
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());
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());
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());
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());
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());
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());
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());
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());
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());
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());
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,88 @@ 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);
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);
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 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);
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());
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"));
}
@@ -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,6 +10,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.MemberRole;
import dev.ltms.bridged.peer.PeerHandle;
import dev.ltms.bridged.peer.PeerLauncher;
@@ -93,6 +94,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; }
};
}
@@ -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,7 @@ 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.mcp.PrimaryRegistry;
import dev.ltms.bridged.inject.Injector;
import org.junit.jupiter.api.BeforeEach;
@@ -31,7 +32,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());
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());
}
}
@@ -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;
@@ -44,6 +46,82 @@ class SessionManagerTest {
return sessionManager(herdr, clock, 0);
}
private SessionManager sessionManager(FakeHerdr herdr, Worktrees worktrees) {
return sessionManager(herdr, worktrees, System::nanoTime);
}
private SessionManager sessionManager(FakeHerdr herdr, Worktrees worktrees, LongSupplier clock) {
BridgedConfig.Profile cfg = new BridgedConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "BRIDGED_WORKER_TOKEN",
List.of("ccs", "ltms-local"), "tab", "bridged-workers",
"worker: {profile} #{n}", null, null, null);
ClaudeCodeLauncher workers = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
return new SessionManager(workers, worktrees, clock);
}
/**
* CB-581: a {@link Worktrees} test double whose {@code hasUncommitted} and {@code remove} can
* be told to throw, so {@link SessionManager#release} can be exercised against exactly the
* failure {@code GitWorktrees} produces when {@code git status}/{@code git worktree remove}
* exits non-zero.
*/
private static final class RecordingWorktrees implements Worktrees {
private final List<String> removeCalls = new java.util.ArrayList<>();
private final java.util.Set<String> failRemoveFor = new java.util.HashSet<>();
private volatile boolean dirty = false;
private volatile RuntimeException hasUncommittedFailure;
RecordingWorktrees dirty(boolean dirty) {
this.dirty = dirty;
return this;
}
RecordingWorktrees failHasUncommittedWith(RuntimeException e) {
this.hasUncommittedFailure = e;
return this;
}
RecordingWorktrees failRemoveFor(String worktreePath) {
failRemoveFor.add(worktreePath);
return this;
}
@Override
public String add(String repoRoot, String branch, String baseRef) {
return "/wt/" + branch.replace('/', '_');
}
@Override
public void remove(String repoRoot, String worktreePath) {
if (failRemoveFor.contains(worktreePath)) {
throw new WorktreeException("simulated remove failure for " + worktreePath);
}
removeCalls.add(worktreePath);
}
@Override
public boolean hasUncommitted(String worktreePath) {
if (hasUncommittedFailure != null) {
throw hasUncommittedFailure;
}
return dirty;
}
@Override
public void overlayParity(String repoRoot, String worktreePath, List<String> overlay) {
}
@Override
public String repoRoot(String cwd) {
return "/repo";
}
List<String> removeCalls() {
return List.copyOf(removeCalls);
}
}
private SessionManager sessionManager(FakeHerdr herdr, LongSupplier clock, int contextCap) {
return sessionManager(herdr, clock, contextCap, false);
}
@@ -91,6 +169,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
@@ -505,4 +601,172 @@ class SessionManagerTest {
"a listener failure must never prevent the teardown it is reacting to");
assertTrue(sessions.get(s.paneId()).isEmpty(), "and the session is still deregistered");
}
// --- CB-581: a throw inside release() must not orphan the pane or abort reapIdle -----------
@Test
void releasePreservesWorktreeWhenDirtyCheckThrows() {
FakeHerdr herdr = new FakeHerdr();
RecordingWorktrees worktrees = new RecordingWorktrees();
SessionManager sessions = sessionManager(herdr, worktrees);
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-581a", null));
worktrees.failHasUncommittedWith(new WorktreeException("git status exited 128"));
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
ch.qos.logback.classic.Logger sessionLog =
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(SessionManager.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.setContext(ctx);
appender.start();
sessionLog.addAppender(appender);
sessionLog.setLevel(Level.WARN);
try {
assertDoesNotThrow(() -> sessions.release(s.paneId()),
"a throwing dirty check must not abort the release");
assertTrue(worktrees.removeCalls().isEmpty(),
"the worktree is preserved when its dirty state cannot be determined");
String warn = appender.list.stream()
.filter(e -> e.getLevel().equals(Level.WARN))
.map(ILoggingEvent::getFormattedMessage)
.filter(m -> m.contains(s.worktree()))
.findFirst()
.orElse("no warn logged naming the worktree");
assertTrue(warn.contains(s.paneId()), "the WARN names the pane: " + warn);
assertTrue(warn.contains(s.terminalId()), "the WARN names the terminal: " + warn);
} finally {
sessionLog.detachAppender(appender);
}
}
@Test
void releaseStillStopsThePaneWhenDirtyCheckThrows() {
FakeHerdr herdr = new FakeHerdr();
RecordingWorktrees worktrees = new RecordingWorktrees();
SessionManager sessions = sessionManager(herdr, worktrees);
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-581b", null));
worktrees.failHasUncommittedWith(new WorktreeException("git status exited 128"));
sessions.release(s.paneId());
assertEquals(1, paneCloseCallsFor(herdr, "w9:pRoot_1"),
"the pane is stopped exactly once even though the dirty check threw");
}
@Test
void releaseStillNotifiesTheListenerWhenDirtyCheckThrows() {
FakeHerdr herdr = new FakeHerdr();
RecordingWorktrees worktrees = new RecordingWorktrees();
SessionManager sessions = sessionManager(herdr, worktrees);
java.util.List<String> released = new java.util.concurrent.CopyOnWriteArrayList<>();
sessions.onRelease(released::add);
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-581c", null));
worktrees.failHasUncommittedWith(new WorktreeException("git status exited 128"));
sessions.release(s.paneId());
assertEquals(java.util.List.of(s.terminalId()), released,
"a blocked caller must still be told the terminal was released, even though the "
+ "dirty check threw");
}
@Test
void reapIdleSurvivesOneSessionThatFailsToRelease() {
long[] clock = {0};
FakeHerdr herdr = new FakeHerdr();
RecordingWorktrees worktrees = new RecordingWorktrees();
SessionManager sessions = sessionManager(herdr, worktrees, () -> clock[0]);
MemberSession a = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-581d", null));
MemberSession b = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-581e", null));
MemberSession c = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-581f", null));
sessions.asPresence().markPresent(a.terminalId());
sessions.asPresence().markPresent(b.terminalId());
sessions.asPresence().markPresent(c.terminalId());
// The middle session's worktree removal fails — release() propagates that, so this is the
// one call reapIdle's per-session guard must survive without skipping the rest of the pass.
worktrees.failRemoveFor(b.worktree());
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
ch.qos.logback.classic.Logger sessionLog =
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(SessionManager.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.setContext(ctx);
appender.start();
sessionLog.addAppender(appender);
sessionLog.setLevel(Level.WARN);
int reaped;
try {
clock[0] = 100;
reaped = sessions.reapIdle(10);
String warn = appender.list.stream()
.filter(e -> e.getLevel().equals(Level.WARN))
.map(ILoggingEvent::getFormattedMessage)
.filter(m -> m.contains(b.paneId()))
.findFirst()
.orElse("no reap-failure WARN logged");
assertTrue(warn.contains(b.terminalId()), "the WARN names the failed session's terminal: " + warn);
assertTrue(warn.contains(b.worktree()), "the WARN names the failed session's worktree: " + warn);
} finally {
sessionLog.detachAppender(appender);
}
assertEquals(2, reaped, "the middle session's failure is logged, not counted as reaped");
assertTrue(sessions.get(a.paneId()).isEmpty(), "the first session is still released");
assertTrue(sessions.get(c.paneId()).isEmpty(), "the third session is still released");
assertTrue(sessions.get(b.paneId()).isEmpty(),
"the middle session is still deregistered even though its worktree removal threw");
assertEquals(1, paneCloseCallsFor(herdr, "w9:pRoot_1"), "the first pane is stopped");
assertEquals(1, paneCloseCallsFor(herdr, "w9:pRoot_2"),
"the middle pane is still stopped even though its worktree removal failed");
assertEquals(1, paneCloseCallsFor(herdr, "w9:pRoot_3"), "the third pane is stopped");
}
@Test
void unchangedRegressionCleanCompletedReleaseStillRemovesTheWorktree() {
FakeHerdr herdr = new FakeHerdr();
RecordingWorktrees worktrees = new RecordingWorktrees();
SessionManager sessions = sessionManager(herdr, worktrees);
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-581g", null));
sessions.release(s.paneId());
assertEquals(List.of(s.worktree()), worktrees.removeCalls(),
"COMPLETED release of a clean worktree still removes it");
}
@Test
void unchangedRegressionDirtyCompletedReleaseStillPreservesTheWorktree() {
FakeHerdr herdr = new FakeHerdr();
RecordingWorktrees worktrees = new RecordingWorktrees().dirty(true);
SessionManager sessions = sessionManager(herdr, worktrees);
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-581h", null));
sessions.release(s.paneId());
assertTrue(worktrees.removeCalls().isEmpty(),
"COMPLETED release of a dirty worktree still preserves it");
}
@Test
void unchangedRegressionShutdownDrainStillPreservesTheWorktree() {
FakeHerdr herdr = new FakeHerdr();
RecordingWorktrees worktrees = new RecordingWorktrees();
SessionManager sessions = sessionManager(herdr, worktrees);
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-581i", null));
sessions.asPresence().markPresent(s.terminalId());
sessions.drainAll(TimeUnit.MILLISECONDS.toNanos(100));
assertTrue(worktrees.removeCalls().isEmpty(), "SHUTDOWN drain still preserves the worktree");
}
}