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.
This commit is contained in:
Dai Ha
2026-08-15 10:00:34 +02:00
15 changed files with 313 additions and 32 deletions
+7
View File
@@ -138,6 +138,12 @@ herdrSocket: ~/.config/herdr/herdr.sock
# SSH is unaffected). The token value itself is never stored in this file. # 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 # gitHostEnv → host env var holding the forge host (default GITEA_HOST). Injected as
# GITEA_HOST *only* alongside a resolved gitTokenEnv. # 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 # env → extra environment for this profile's workers, as a literal key/value map
# (CB-511). Use it to give workers a toolchain. # (CB-511). Use it to give workers a toolchain.
# #
@@ -171,6 +177,7 @@ profiles:
maxLoad: 2 # max live workers on this profile (omit for unlimited) 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) # 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 # 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 # 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 # 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 # parityOverlay: [".claude/settings.local.json", ".env", ".envrc"] # never add .mcp.json — see above
@@ -13,6 +13,7 @@ import dev.ltms.bridged.herdr.PaneLocator;
import dev.ltms.bridged.herdr.UnixSocketHerdrClient; import dev.ltms.bridged.herdr.UnixSocketHerdrClient;
import dev.ltms.bridged.herdr.WorkspaceControl; import dev.ltms.bridged.herdr.WorkspaceControl;
import dev.ltms.bridged.inject.CompletionResolver; import dev.ltms.bridged.inject.CompletionResolver;
import dev.ltms.bridged.inject.ExhaustedPatternLookup;
import dev.ltms.bridged.inject.Injector; import dev.ltms.bridged.inject.Injector;
import dev.ltms.bridged.inject.StatusPoller; import dev.ltms.bridged.inject.StatusPoller;
import dev.ltms.bridged.inject.TurnListener; 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.Function;
import java.util.function.Predicate; import java.util.function.Predicate;
import java.util.function.Supplier; import java.util.function.Supplier;
import java.util.regex.Pattern;
import java.util.stream.Collectors; import java.util.stream.Collectors;
/** /**
@@ -253,7 +255,24 @@ public final class Bridged {
// The blocking message endpoint (CB-104) is the producer; the poller is inert until then. // 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. // CB-106: a confirmed turn completion resolves a blocked send whose worker never replied.
Rendezvous rendezvous = new Rendezvous(); 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-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. // CB-301: the manager's presence bridge records availability and drives SPAWNING → READY.
MemberPresence presence = sessions.asPresence(); MemberPresence presence = sessions.asPresence();
@@ -181,6 +181,12 @@ public record BridgedConfig(
* opposite intents). For the same reason, an {@code env:} entry naming * opposite intents). For the same reason, an {@code env:} entry naming
* {@code ANTHROPIC_BASE_URL} or {@code ANTHROPIC_AUTH_TOKEN} is refused at * {@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. * 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) @JsonIgnoreProperties(ignoreUnknown = true)
public record Profile(String profile, String baseUrl, String model, public record Profile(String profile, String baseUrl, String model,
@@ -193,7 +199,8 @@ public record BridgedConfig(
Map<String, String> env, Map<String, String> env,
Float weight, Float weight,
Integer maxLoad, Integer maxLoad,
Boolean subscription) { Boolean subscription,
String exhaustedPattern) {
/** Peer kind spawned by {@link dev.ltms.bridged.member.ClaudeCodeLauncher} (the default). */ /** Peer kind spawned by {@link dev.ltms.bridged.member.ClaudeCodeLauncher} (the default). */
public static final String KIND_CLAUDE_CODE = "claude-code"; 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; weight = (weight == null || weight <= 0.0f) ? 1.0f : weight;
maxLoad = (maxLoad == null || maxLoad <= 0) ? null : maxLoad; maxLoad = (maxLoad == null || maxLoad <= 0) ? null : maxLoad;
subscription = (subscription != null && subscription) ? Boolean.TRUE : Boolean.FALSE; 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 placement, String workspace, String tabLabel, String mcpUrl,
String cwd, List<String> parityOverlay) { String cwd, List<String> parityOverlay) {
this(profile, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel, 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 placement, String workspace, String tabLabel, String mcpUrl,
String cwd, List<String> parityOverlay, String gitTokenEnv, String gitHostEnv) { String cwd, List<String> parityOverlay, String gitTokenEnv, String gitHostEnv) {
this(profile, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel, 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 cwd, List<String> parityOverlay, String gitTokenEnv, String gitHostEnv,
String kind) { String kind) {
this(profile, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel, 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. */ /** A copy with {@code profile} set — used to default a profile to its {@code workers} key. */
public Profile withProfile(String p) { public Profile withProfile(String p) {
return new Profile(p, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel, 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). */ /** 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 cwd, List<String> parityOverlay, String gitTokenEnv, String gitHostEnv,
String kind, Map<String, String> env, Float weight, Integer maxLoad) { String kind, Map<String, String> env, Float weight, Integer maxLoad) {
this(profile, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel, 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). */ /** True when this profile's workers are granted a forge token to open their own PR (CB-302). */
@@ -6,8 +6,13 @@ import dev.ltms.bridged.msg.TurnToken;
import org.slf4j.Logger; import org.slf4j.Logger;
import org.slf4j.LoggerFactory; 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.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap; 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 * 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 AgentControl agents;
private final Rendezvous rendezvous; private final Rendezvous rendezvous;
private final ExhaustedPatternLookup exhaustedPatterns;
/** /**
* Per-target record of the turn currently in flight: the exact {@link Rendezvous} waiter its * 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<>(); 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.agents = agents;
this.rendezvous = rendezvous; this.rendezvous = rendezvous;
this.exhaustedPatterns = Objects.requireNonNull(exhaustedPatterns, "exhaustedPatterns");
} }
@Override @Override
@@ -157,11 +170,12 @@ public final class CompletionResolver implements TurnListener {
return; return;
} }
String tail; String tail;
String assistantBlock = null;
int originalLength = 0; int originalLength = 0;
boolean clipped = false; boolean clipped = false;
boolean scrapeFailed = false; boolean scrapeFailed = false;
try { try {
String assistantBlock = lastAssistantBlock(agents.read(target, SCRAPE_SOURCE)); assistantBlock = lastAssistantBlock(agents.read(target, SCRAPE_SOURCE));
originalLength = assistantBlock.strip().length(); originalLength = assistantBlock.strip().length();
clipped = originalLength > MAX_SCRAPE_CHARS; clipped = originalLength > MAX_SCRAPE_CHARS;
tail = clip(assistantBlock); tail = clip(assistantBlock);
@@ -184,6 +198,22 @@ public final class CompletionResolver implements TurnListener {
target); target);
return; // keep the in-flight record: a later genuine completion still needs it 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; String completion = clipped ? tail + "\n" + CLIPPED_PANE_TAIL_MARKER : tail;
if (rendezvous.resolveCompletion(waiter, completion)) { if (rendezvous.resolveCompletion(waiter, completion)) {
inFlight.remove(target, turn); 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) { private static String clip(String s) {
if (s == null) return ""; if (s == null) return "";
String trimmed = s.strip(); 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;
}
}
@@ -458,6 +458,11 @@ public final class BridgeMcp {
"[worker finished without a structured bridge_reply — transcript tail follows]\n" + r.text()); "[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. // 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()); 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. // 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() 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() + "\n\nAnswer it by calling bridge_send again with turnId=\"" + r.turnId()
@@ -69,6 +69,14 @@ public final class MessageService {
* failure context (e.g. the error screen). Terminal, but not a successful completion. * failure context (e.g. the error screen). Terminal, but not a successful completion.
*/ */
WORKER_FAILED, 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 * 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 * 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 COMPLETED_UNREPLIED -> "completion_fallback";
case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, BUSY -> "timeout"; case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, BUSY -> "timeout";
case WORKER_FAILED -> "failed"; case WORKER_FAILED -> "failed";
case BACKEND_EXHAUSTED -> "backend_exhausted";
case STALE_TURN, QUESTION -> null; // not a completed delegation 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"; String source = r.outcome() == Outcome.REPLIED ? "reply" : "transcript";
return new TaskView(ticket, Phase.DONE, r.text(), source, null, null); 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 // A wedged worker (CB-109) or a backend-exhausted classification (CB-578 stage A) carries
// outcomes carry none, so fall back to the outcome name. // the real cause as its reason; the timeout/busy outcomes carry none, so fall back to the
String detail = r.outcome() == Outcome.WORKER_FAILED && r.text() != null // outcome name.
boolean carriesReason = r.outcome() == Outcome.WORKER_FAILED || r.outcome() == Outcome.BACKEND_EXHAUSTED;
String detail = carriesReason && r.text() != null
? r.text() ? r.text()
: "no reply — " + r.outcome().name().toLowerCase(); : "no reply — " + r.outcome().name().toLowerCase();
return new TaskView(ticket, Phase.FAILED, null, null, detail, null); return new TaskView(ticket, Phase.FAILED, null, null, detail, null);
@@ -669,6 +680,7 @@ public final class MessageService {
case REPLY -> Outcome.REPLIED; case REPLY -> Outcome.REPLIED;
case COMPLETION -> Outcome.COMPLETED_UNREPLIED; case COMPLETION -> Outcome.COMPLETED_UNREPLIED;
case FAILED -> Outcome.WORKER_FAILED; case FAILED -> Outcome.WORKER_FAILED;
case BACKEND_EXHAUSTED -> Outcome.BACKEND_EXHAUSTED;
case QUESTION -> Outcome.QUESTION; case QUESTION -> Outcome.QUESTION;
}; };
} }
@@ -33,6 +33,13 @@ public final class Rendezvous {
COMPLETION, COMPLETION,
/** The worker ran the turn then wedged (CB-109); {@code text} is the failure context. */ /** The worker ran the turn then wedged (CB-109); {@code text} is the failure context. */
FAILED, 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); * 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 * {@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)); 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) { private boolean complete(String session, Resolution resolution) {
CompletableFuture<Resolution> waiter = waiters.get(session); CompletableFuture<Resolution> waiter = waiters.get(session);
return waiter != null && waiter.complete(resolution); return waiter != null && waiter.complete(resolution);
@@ -403,9 +403,12 @@ public final class BridgedApp {
case TIMED_OUT_QUEUED -> "queued"; case TIMED_OUT_QUEUED -> "queued";
case BUSY -> "busy"; case BUSY -> "busy";
case WORKER_FAILED -> "failed"; case WORKER_FAILED -> "failed";
case BACKEND_EXHAUSTED -> "backend_exhausted";
default -> "done"; // unreachable (terminal outcomes handled above) 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() ? reply.text()
: "no reply within " + timeout + "ms; poll status or retry")); : "no reply within " + timeout + "ms; poll status or retry"));
} }
@@ -1140,6 +1140,7 @@ class BridgedConfigTest {
gitHostEnv: GITEA_HOST gitHostEnv: GITEA_HOST
weight: 0.5 weight: 0.5
maxLoad: 2 maxLoad: 2
exhaustedPattern: "usage limit has been reached"
placement: weighted placement: weighted
lifecycle: lifecycle:
idleTtlSeconds: 300 idleTtlSeconds: 300
@@ -1171,6 +1172,8 @@ class BridgedConfigTest {
assertEquals("GITEA_HOST", w.gitHostEnv()); assertEquals("GITEA_HOST", w.gitHostEnv());
assertEquals(0.5f, w.weight(), 0.0001f, "weight binds as a float"); assertEquals(0.5f, w.weight(), 0.0001f, "weight binds as a float");
assertEquals(2, w.maxLoad(), "maxLoad binds as an integer"); 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("weighted", cfg.placement(), "placement binds at the top level");
assertEquals(300, cfg.lifecycle().idleTtlSeconds()); assertEquals(300, cfg.lifecycle().idleTtlSeconds());
@@ -12,6 +12,9 @@ import dev.ltms.bridged.msg.TurnToken;
import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Test;
import org.slf4j.LoggerFactory; 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.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue; import static org.junit.jupiter.api.Assertions.assertTrue;
@@ -23,7 +26,7 @@ class CompletionResolverTest {
void skipsTheScrapeWhenNoSendIsWaiting() { void skipsTheScrapeWhenNoSendIsWaiting() {
FakeHerdr herdr = new FakeHerdr(); FakeHerdr herdr = new FakeHerdr();
Rendezvous rendezvous = new Rendezvous(); // no waiter opened 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 resolver.resolve("term_a", null); // no in-flight turn captured for this target
@@ -35,7 +38,7 @@ class CompletionResolverTest {
void failSkipsTheScrapeWhenNoSendIsWaiting() { void failSkipsTheScrapeWhenNoSendIsWaiting() {
FakeHerdr herdr = new FakeHerdr(); FakeHerdr herdr = new FakeHerdr();
Rendezvous rendezvous = new Rendezvous(); // no waiter opened 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 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() { void captureBaselineSkipsTheReadWhenNoSendIsWaiting() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ X\n❯ "); FakeHerdr herdr = new FakeHerdr().readText("⏺ X\n❯ ");
Rendezvous rendezvous = new Rendezvous(); // no waiter opened 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 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. // send must NOT be resolved with the stale answer.
FakeHerdr herdr = new FakeHerdr().readText("⏺ 391\n❯ "); FakeHerdr herdr = new FakeHerdr().readText("⏺ 391\n❯ ");
Rendezvous rendezvous = new Rendezvous(); 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 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. // 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() { void resolvesACompletionWhoseScrapeChangedSinceDelivery() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ No, 391 = 17 × 23.\n❯ "); // the worker's real answer FakeHerdr herdr = new FakeHerdr().readText("⏺ No, 391 = 17 × 23.\n❯ "); // the worker's real answer
Rendezvous rendezvous = new Rendezvous(); 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 waiter = rendezvous.open("term_a");
// Delivery baseline was the previous turn's "391"; the scrape now differs → resolve. // 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❯ "; String block = "⏺ " + "x".repeat(CompletionResolver.MAX_SCRAPE_CHARS + 1) + "\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(block); FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous(); 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 waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null)); resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
@@ -182,7 +185,7 @@ class CompletionResolverTest {
void leavesAnUnclippedCompletionPaneTailUnmarked() { void leavesAnUnclippedCompletionPaneTailUnmarked() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ complete report\n❯ "); FakeHerdr herdr = new FakeHerdr().readText("⏺ complete report\n❯ ");
Rendezvous rendezvous = new Rendezvous(); 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 waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null)); resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
@@ -194,7 +197,7 @@ class CompletionResolverTest {
void resolvesSynchronouslyBeforePostTurnContextClearing() { void resolvesSynchronouslyBeforePostTurnContextClearing() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ previous answer\n❯ "); FakeHerdr herdr = new FakeHerdr().readText("⏺ previous answer\n❯ ");
Rendezvous rendezvous = new Rendezvous(); 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 waiter = rendezvous.open("term_a");
resolver.captureBaseline("term_a", new TurnToken("term_a", waiter)); resolver.captureBaseline("term_a", new TurnToken("term_a", waiter));
herdr.readText("⏺ answer that /clear would erase\n❯ "); herdr.readText("⏺ answer that /clear would erase\n❯ ");
@@ -216,7 +219,7 @@ class CompletionResolverTest {
String longBlock = "⏺ " + "x".repeat(CompletionResolver.MAX_SCRAPE_CHARS + 500) + "\n❯ "; String longBlock = "⏺ " + "x".repeat(CompletionResolver.MAX_SCRAPE_CHARS + 500) + "\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(longBlock); FakeHerdr herdr = new FakeHerdr().readText(longBlock);
Rendezvous rendezvous = new Rendezvous(); 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 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 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. // No delivery baseline (e.g. the pre-turn read failed) ⇒ never suppress; the completion resolves.
FakeHerdr herdr = new FakeHerdr().readText("⏺ hello\n❯ "); FakeHerdr herdr = new FakeHerdr().readText("⏺ hello\n❯ ");
Rendezvous rendezvous = new Rendezvous(); 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 waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null)); 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. // byte-identical guard would wrongly match the empty tail and suppress.
FakeHerdr herdr = new FakeHerdr().healthy(false); // agent.read throws HerdrException FakeHerdr herdr = new FakeHerdr().healthy(false); // agent.read throws HerdrException
Rendezvous rendezvous = new Rendezvous(); 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 waiter = rendezvous.open("term_a");
var turn = new CompletionResolver.InFlight(waiter, ""); // empty pane baselined at delivery 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. // fail must not overwrite that value, and must not even scrape the worker — nobody needs it.
FakeHerdr herdr = new FakeHerdr().readText("an error screen"); FakeHerdr herdr = new FakeHerdr().readText("an error screen");
Rendezvous rendezvous = new Rendezvous(); 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 waiter = rendezvous.open("term_a");
var turn = new CompletionResolver.InFlight(waiter, null); 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. // fail falls back to the waiter currently registered on the Rendezvous and fails it.
FakeHerdr herdr = new FakeHerdr().readText("stuck on an error screen"); FakeHerdr herdr = new FakeHerdr().readText("stuck on an error screen");
Rendezvous rendezvous = new Rendezvous(); 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 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 resolver.fail("term_a", null); // no in-flight turn → fall back to the registered waiter
@@ -321,7 +324,7 @@ class CompletionResolverTest {
try { try {
FakeHerdr herdr = new FakeHerdr().readText("stuck on an error screen"); FakeHerdr herdr = new FakeHerdr().readText("stuck on an error screen");
Rendezvous rendezvous = new Rendezvous(); 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 waiter = rendezvous.open("term_a");
resolver.fail("term_a", null); 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. // 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❯ "); FakeHerdr herdr = new FakeHerdr().readText("⏺ turn N answer\n❯ ");
Rendezvous rendezvous = new Rendezvous(); 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 var waiterN = rendezvous.open("term_a"); // turn N's send
// The turn as the injector captured it at delivery (waiter + pre-turn baseline). // 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"); "turn N stays resolved by its own reply");
assertTrue(rendezvous.isWaiting("term_a"), "turn N+1 is still awaiting its own resolution"); 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, List.of("ccs", "ltms"), "tab", "bridged-workers", null,
"http://127.0.0.1:8765/mcp", null, null, "http://127.0.0.1:8765/mcp", null, null,
null, 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) { private static BridgedConfig configWith(BridgedConfig.Leader lead) {
@@ -733,7 +733,7 @@ class ClaudeCodeLauncherTest {
return new BridgedConfig.Profile( return new BridgedConfig.Profile(
profile, baseUrl, "sonnet", null, "BRIDGED_WORKER_TOKEN", profile, baseUrl, "sonnet", null, "BRIDGED_WORKER_TOKEN",
List.of("ccs", profile), "tab", "bridged-workers", "w #{n}", null, null, null, 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 @Test
@@ -819,7 +819,7 @@ class ClaudeCodeLauncherTest {
null, null, null, null, null, null,
Map.of("ANTHROPIC_BASE_URL", "http://evil.example.com", Map.of("ANTHROPIC_BASE_URL", "http://evil.example.com",
"ANTHROPIC_AUTH_TOKEN", "sk-ant-bad", "JAVA_HOME", "/opt/jdk"), "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 ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
_ -> "would-be-token").spawn(); _ -> "would-be-token").spawn();
@@ -5,6 +5,7 @@ import dev.ltms.bridged.herdr.AgentStatus;
import dev.ltms.bridged.herdr.FakeHerdr; import dev.ltms.bridged.herdr.FakeHerdr;
import dev.ltms.bridged.herdr.HerdrException; import dev.ltms.bridged.herdr.HerdrException;
import dev.ltms.bridged.inject.CompletionResolver; import dev.ltms.bridged.inject.CompletionResolver;
import dev.ltms.bridged.inject.ExhaustedPatternLookup;
import dev.ltms.bridged.mcp.PrimaryRegistry; import dev.ltms.bridged.mcp.PrimaryRegistry;
import dev.ltms.bridged.inject.Injector; import dev.ltms.bridged.inject.Injector;
import org.junit.jupiter.api.BeforeEach; 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 FakeHerdr herdr = new FakeHerdr().readText("BUILD GREEN: 391 files");
private final AgentControl agents = new AgentControl(herdr); private final AgentControl agents = new AgentControl(herdr);
private final Rendezvous rendezvous = new Rendezvous(); 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 Injector injector = new Injector(agents, completion);
private final InMemoryReplyInbox inbox = new InMemoryReplyInbox(); private final InMemoryReplyInbox inbox = new InMemoryReplyInbox();
private final MessageService messages = new MessageService(agents, injector, rendezvous, inbox); 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"); "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 @Test
void closeAskRemovesTheTurn() { void closeAskRemovesTheTurn() {
Rendezvous.AskTicket t = rendezvous.openAsk(W); Rendezvous.AskTicket t = rendezvous.openAsk(W);