Merge CB-578 stage A: classify a usage-limit refusal instead of a completed reply
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:
@@ -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);
|
||||||
|
|||||||
Reference in New Issue
Block a user