From 541df872726bb0dc5f52338f79352ab445b4e26d Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 15 Aug 2026 09:55:06 +0200 Subject: [PATCH] CB-578 stage A: classify a usage-limit refusal instead of a completed reply A backend that refuses on a subscription usage limit leaves the pane healthy but the turn ends with no bridge_reply; the completion fallback used to scrape and hand that refusal back as if it were a real answer. CompletionResolver now matches the scrape against a per-profile exhaustedPattern (config, never a vendor string) and resolves the send as Rendezvous.Kind/Outcome.BACKEND_EXHAUSTED with a reason carrying the matched line, kept distinct from GONE/WORKER_FAILED. A profile with no pattern configured is unaffected. Coverage is logged at startup via CompletionResolver.coverage(...), naming which profiles have a pattern and which don't, following FleetHealthMonitor.coverage's pattern. --- bridged/bridged.example.yaml | 7 ++ .../main/java/dev/ltms/bridged/Bridged.java | 21 +++- .../ltms/bridged/config/BridgedConfig.java | 28 ++++- .../bridged/inject/CompletionResolver.java | 73 ++++++++++- .../inject/ExhaustedPatternLookup.java | 27 ++++ .../java/dev/ltms/bridged/mcp/BridgeMcp.java | 5 + .../dev/ltms/bridged/msg/MessageService.java | 18 ++- .../java/dev/ltms/bridged/msg/Rendezvous.java | 20 +++ .../dev/ltms/bridged/rest/BridgedApp.java | 5 +- .../bridged/config/BridgedConfigTest.java | 3 + .../inject/CompletionResolverTest.java | 117 +++++++++++++++--- .../ltms/bridged/lead/LeadLauncherTest.java | 2 +- .../member/ClaudeCodeLauncherTest.java | 4 +- .../ltms/bridged/msg/MessageServiceTest.java | 4 +- .../dev/ltms/bridged/msg/RendezvousTest.java | 11 ++ 15 files changed, 313 insertions(+), 32 deletions(-) create mode 100644 bridged/src/main/java/dev/ltms/bridged/inject/ExhaustedPatternLookup.java diff --git a/bridged/bridged.example.yaml b/bridged/bridged.example.yaml index 35630a2..3fcf852 100644 --- a/bridged/bridged.example.yaml +++ b/bridged/bridged.example.yaml @@ -138,6 +138,12 @@ herdrSocket: ~/.config/herdr/herdr.sock # SSH is unaffected). The token value itself is never stored in this file. # gitHostEnv → host env var holding the forge host (default GITEA_HOST). Injected as # GITEA_HOST *only* alongside a resolved gitTokenEnv. +# exhaustedPattern → regex matched against a completion-fallback scrape (CB-578 stage A) to +# classify a turn that ended with no bridge_reply as the backend having +# refused on a subscription usage limit, rather than a real answer. Opt-in — +# omit and this profile's completion fallback behaves exactly as before. +# Every backend words its refusal differently, so this is config, never a +# vendor string baked into bridged itself. # env → extra environment for this profile's workers, as a literal key/value map # (CB-511). Use it to give workers a toolchain. # @@ -171,6 +177,7 @@ profiles: maxLoad: 2 # max live workers on this profile (omit for unlimited) # gitTokenEnv: GITEA_TOKEN # opt-in: let this profile's workers open their own PR (CB-302) # gitHostEnv: GITEA_HOST # defaults to GITEA_HOST; injected only with gitTokenEnv + # exhaustedPattern: "usage limit has been reached" # opt-in: classify a usage-limit refusal (CB-578) # configDir: /Users/me/.ccs/instances/gx10 # CLAUDE_CONFIG_DIR — inherit that profile's skills/MCP # cwd: /Users/me/src/myrepo # pin the working dir; omit to inherit the primary's # parityOverlay: [".claude/settings.local.json", ".env", ".envrc"] # never add .mcp.json — see above diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java index 55c5b04..97211d2 100644 --- a/bridged/src/main/java/dev/ltms/bridged/Bridged.java +++ b/bridged/src/main/java/dev/ltms/bridged/Bridged.java @@ -13,6 +13,7 @@ import dev.ltms.bridged.herdr.PaneLocator; import dev.ltms.bridged.herdr.UnixSocketHerdrClient; import dev.ltms.bridged.herdr.WorkspaceControl; import dev.ltms.bridged.inject.CompletionResolver; +import dev.ltms.bridged.inject.ExhaustedPatternLookup; import dev.ltms.bridged.inject.Injector; import dev.ltms.bridged.inject.StatusPoller; import dev.ltms.bridged.inject.TurnListener; @@ -60,6 +61,7 @@ import java.util.concurrent.atomic.AtomicReference; import java.util.function.Function; import java.util.function.Predicate; import java.util.function.Supplier; +import java.util.regex.Pattern; import java.util.stream.Collectors; /** @@ -253,7 +255,24 @@ public final class Bridged { // The blocking message endpoint (CB-104) is the producer; the poller is inert until then. // CB-106: a confirmed turn completion resolves a blocked send whose worker never replied. Rendezvous rendezvous = new Rendezvous(); - CompletionResolver completion = new CompletionResolver(agents, rendezvous); + // CB-578 stage A: classify a completion-fallback scrape that matches a profile's configured + // usage-limit refusal as BACKEND_EXHAUSTED rather than handing it back as a real answer. + // Compiled once at startup, keyed by profile name; a profile with no exhaustedPattern is + // simply absent here, so its workers keep today's completion-fallback behaviour unchanged. + Map exhaustedPatternsByProfile = new LinkedHashMap<>(); + cfg.profiles().forEach((name, profile) -> { + if (profile.hasExhaustedPattern()) { + exhaustedPatternsByProfile.put(name, Pattern.compile(profile.exhaustedPattern())); + } + }); + ExhaustedPatternLookup exhaustedPatterns = target -> sessions.roster().stream() + .filter(session -> target.equals(session.terminalId())) + .findFirst() + .map(session -> exhaustedPatternsByProfile.get(session.profile())) + .orElse(null); + log.info("backend-exhausted classification (CB-578 stage A): {}", + CompletionResolver.coverage(cfg.profiles().keySet(), exhaustedPatternsByProfile.keySet())); + CompletionResolver completion = new CompletionResolver(agents, rendezvous, exhaustedPatterns); // CB-113: deliver only to an available worker (its MCP is connected), never its boot window. // CB-301: the manager's presence bridge records availability and drives SPAWNING → READY. MemberPresence presence = sessions.asPresence(); diff --git a/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java b/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java index 30f74b2..4d257a9 100644 --- a/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java +++ b/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java @@ -181,6 +181,12 @@ public record BridgedConfig( * opposite intents). For the same reason, an {@code env:} entry naming * {@code ANTHROPIC_BASE_URL} or {@code ANTHROPIC_AUTH_TOKEN} is refused at * config load (CB-542): on the subscription path no guard would vet it. + * @param exhaustedPattern regex matched against a completion-fallback scrape (CB-578 stage A) to + * classify a turn that ended with no {@code bridge_reply} as the backend + * having refused on a subscription usage limit, rather than a real answer. + * {@code null}/blank ⇒ the classification never fires for this profile and + * today's completion-fallback behaviour is unchanged. Every backend words + * its refusal differently, so this is config, never a vendor string in code. */ @JsonIgnoreProperties(ignoreUnknown = true) public record Profile(String profile, String baseUrl, String model, @@ -193,7 +199,8 @@ public record BridgedConfig( Map env, Float weight, Integer maxLoad, - Boolean subscription) { + Boolean subscription, + String exhaustedPattern) { /** Peer kind spawned by {@link dev.ltms.bridged.member.ClaudeCodeLauncher} (the default). */ public static final String KIND_CLAUDE_CODE = "claude-code"; @@ -236,6 +243,9 @@ public record BridgedConfig( weight = (weight == null || weight <= 0.0f) ? 1.0f : weight; maxLoad = (maxLoad == null || maxLoad <= 0) ? null : maxLoad; subscription = (subscription != null && subscription) ? Boolean.TRUE : Boolean.FALSE; + // exhaustedPattern stays null when unset/blank (opt-in) — no defaulting, no vendor + // wording: an unconfigured profile keeps today's completion-fallback behaviour exactly. + exhaustedPattern = (exhaustedPattern == null || exhaustedPattern.isBlank()) ? null : exhaustedPattern; } /** @@ -248,7 +258,7 @@ public record BridgedConfig( String placement, String workspace, String tabLabel, String mcpUrl, String cwd, List parityOverlay) { this(profile, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel, - mcpUrl, cwd, parityOverlay, null, null, null, null, null, null, null); + mcpUrl, cwd, parityOverlay, null, null, null, null, null, null, null, null); } /** @@ -260,7 +270,7 @@ public record BridgedConfig( String placement, String workspace, String tabLabel, String mcpUrl, String cwd, List parityOverlay, String gitTokenEnv, String gitHostEnv) { this(profile, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel, - mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, null, null, null, null, null); + mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, null, null, null, null, null, null); } /** @@ -273,13 +283,14 @@ public record BridgedConfig( String cwd, List parityOverlay, String gitTokenEnv, String gitHostEnv, String kind) { this(profile, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel, - mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, null, null, null, null); + mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, null, null, null, null, null); } /** A copy with {@code profile} set — used to default a profile to its {@code workers} key. */ public Profile withProfile(String p) { return new Profile(p, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel, - mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, env, weight, maxLoad, subscription); + mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, env, weight, maxLoad, subscription, + exhaustedPattern); } /** True when this profile is served by the Claude Code adapter (the default kind). */ @@ -312,7 +323,12 @@ public record BridgedConfig( String cwd, List parityOverlay, String gitTokenEnv, String gitHostEnv, String kind, Map env, Float weight, Integer maxLoad) { this(profile, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel, - mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, env, weight, maxLoad, null); + mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, env, weight, maxLoad, null, null); + } + + /** True when this profile's CB-578 stage A backend-exhausted classification is configured. */ + public boolean hasExhaustedPattern() { + return exhaustedPattern != null; } /** True when this profile's workers are granted a forge token to open their own PR (CB-302). */ diff --git a/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java b/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java index 5259ffd..82c0b6a 100644 --- a/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java +++ b/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java @@ -6,8 +6,13 @@ import dev.ltms.bridged.msg.TurnToken; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.util.List; +import java.util.Objects; +import java.util.Set; +import java.util.TreeSet; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; +import java.util.regex.Pattern; /** * The CB-106 completion fallback: bridges the {@link Injector}'s turn-completion signal to the @@ -61,6 +66,7 @@ public final class CompletionResolver implements TurnListener { private final AgentControl agents; private final Rendezvous rendezvous; + private final ExhaustedPatternLookup exhaustedPatterns; /** * Per-target record of the turn currently in flight: the exact {@link Rendezvous} waiter its @@ -78,9 +84,16 @@ public final class CompletionResolver implements TurnListener { private final ConcurrentHashMap inFlight = new ConcurrentHashMap<>(); - public CompletionResolver(AgentControl agents, Rendezvous rendezvous) { + /** + * @param exhaustedPatterns CB-578 stage A: per-target lookup for a profile's configured + * usage-limit refusal pattern. Required — there is deliberately no + * defaulting overload; a caller that does not want the classification + * must pass an explicit inert value ({@link ExhaustedPatternLookup#none()}). + */ + public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns) { this.agents = agents; this.rendezvous = rendezvous; + this.exhaustedPatterns = Objects.requireNonNull(exhaustedPatterns, "exhaustedPatterns"); } @Override @@ -157,11 +170,12 @@ public final class CompletionResolver implements TurnListener { return; } String tail; + String assistantBlock = null; int originalLength = 0; boolean clipped = false; boolean scrapeFailed = false; try { - String assistantBlock = lastAssistantBlock(agents.read(target, SCRAPE_SOURCE)); + assistantBlock = lastAssistantBlock(agents.read(target, SCRAPE_SOURCE)); originalLength = assistantBlock.strip().length(); clipped = originalLength > MAX_SCRAPE_CHARS; tail = clip(assistantBlock); @@ -184,6 +198,22 @@ public final class CompletionResolver implements TurnListener { target); return; // keep the in-flight record: a later genuine completion still needs it } + // CB-578 stage A: a turn that ended with no bridge_reply AND whose scrape matches the + // backend's configured usage-limit pattern is a refusal, not an answer. Classify it as + // BACKEND_EXHAUSTED rather than handing the caller a scrape that reads like a real reply. + if (!scrapeFailed) { + Pattern exhausted = exhaustedPatterns.patternFor(target); + String matchedLine = exhausted == null ? null : firstMatchingLine(assistantBlock, exhausted); + if (matchedLine != null) { + String reason = "backend exhausted (usage limit): " + matchedLine; + if (rendezvous.resolveExhausted(waiter, reason)) { + inFlight.remove(target, turn); + log.warn("completion for {} classified BACKEND_EXHAUSTED (no bridge_reply; scrape " + + "matched the profile's exhausted pattern): {}", target, reason); + } + return; + } + } String completion = clipped ? tail + "\n" + CLIPPED_PANE_TAIL_MARKER : tail; if (rendezvous.resolveCompletion(waiter, completion)) { inFlight.remove(target, turn); @@ -232,6 +262,45 @@ public final class CompletionResolver implements TurnListener { } } + /** + * The first line of {@code text} matching {@code pattern}, stripped — the CB-578 stage A + * evidence carried in a {@code BACKEND_EXHAUSTED} reason so the operator sees the real refusal + * text, never a generic label. {@code null} if no line matches. + */ + static String firstMatchingLine(String text, Pattern pattern) { + if (text == null || text.isEmpty()) return null; + for (String line : text.split("\n", -1)) { + if (pattern.matcher(line).find()) { + return line.strip(); + } + } + return null; + } + + /** + * Coverage summary for the CB-578 stage A exhausted-pattern classification, logged at startup + * the way {@link dev.ltms.bridged.health.FleetHealthMonitor#coverage} is — so an operator can + * see whether the classification is on, and for which profiles, without reading every + * profile's config by hand. + * + * @param allProfiles every configured profile name + * @param configuredProfiles the subset of {@code allProfiles} that carry an exhausted pattern + */ + public static String coverage(Set allProfiles, Set configuredProfiles) { + if (configuredProfiles.isEmpty()) { + return "off (no profile has an exhaustedPattern configured; profiles: " + sorted(allProfiles) + ")"; + } + Set 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 sorted(Set names) { + return names.stream().sorted().toList(); + } + private static String clip(String s) { if (s == null) return ""; String trimmed = s.strip(); diff --git a/bridged/src/main/java/dev/ltms/bridged/inject/ExhaustedPatternLookup.java b/bridged/src/main/java/dev/ltms/bridged/inject/ExhaustedPatternLookup.java new file mode 100644 index 0000000..79a4209 --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/inject/ExhaustedPatternLookup.java @@ -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. + * + *

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; + } +} diff --git a/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java index af11447..da37ea9 100644 --- a/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java +++ b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java @@ -458,6 +458,11 @@ public final class BridgeMcp { "[worker finished without a structured bridge_reply — transcript tail follows]\n" + r.text()); // The worker ran the turn then wedged (CB-109) — surface the error context. case WORKER_FAILED -> text("[worker failed — turn ended in an unrecoverable state]\n" + r.text()); + // The backend refused on a subscription usage limit (CB-578 stage A) — the worker's + // pane stayed healthy, but its account is exhausted. Distinct from WORKER_FAILED so the + // primary gets the real cause, not a generic wedge. + case BACKEND_EXHAUSTED -> text("[backend exhausted — the worker's account refused on a " + + "usage limit]\n" + r.text()); // The worker paused mid-turn to ask (CB-205) — tell the primary how to answer in-turn. case QUESTION -> text("[question] the worker paused to ask before it can finish:\n" + r.text() + "\n\nAnswer it by calling bridge_send again with turnId=\"" + r.turnId() diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java index c83c0bd..5db11cc 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java @@ -69,6 +69,14 @@ public final class MessageService { * failure context (e.g. the error screen). Terminal, but not a successful completion. */ WORKER_FAILED, + /** + * The turn finished without a {@code bridge_reply} and the scrape matched the backend's + * configured usage-limit refusal pattern (CB-578 stage A); {@code text} is the reason, + * carrying the matched line. The worker's pane is healthy — only its account is refusing — + * so this is never reported as a completed reply, and is kept distinct from + * {@link #WORKER_FAILED} (a wedged worker) and a session simply going {@code GONE}. + */ + BACKEND_EXHAUSTED, /** * The worker paused mid-turn to ask the primary a question (CB-205); {@code text} is the * question and {@code turnId} correlates the answer. Not terminal — the primary answers with @@ -280,6 +288,7 @@ public final class MessageService { case COMPLETED_UNREPLIED -> "completion_fallback"; case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, BUSY -> "timeout"; case WORKER_FAILED -> "failed"; + case BACKEND_EXHAUSTED -> "backend_exhausted"; case STALE_TURN, QUESTION -> null; // not a completed delegation }; } @@ -591,9 +600,11 @@ public final class MessageService { String source = r.outcome() == Outcome.REPLIED ? "reply" : "transcript"; return new TaskView(ticket, Phase.DONE, r.text(), source, null, null); } - // A wedged worker (CB-109) carries the error context as its reason; the timeout/busy - // outcomes carry none, so fall back to the outcome name. - String detail = r.outcome() == Outcome.WORKER_FAILED && r.text() != null + // A wedged worker (CB-109) or a backend-exhausted classification (CB-578 stage A) carries + // the real cause as its reason; the timeout/busy outcomes carry none, so fall back to the + // outcome name. + boolean carriesReason = r.outcome() == Outcome.WORKER_FAILED || r.outcome() == Outcome.BACKEND_EXHAUSTED; + String detail = carriesReason && r.text() != null ? r.text() : "no reply — " + r.outcome().name().toLowerCase(); return new TaskView(ticket, Phase.FAILED, null, null, detail, null); @@ -669,6 +680,7 @@ public final class MessageService { case REPLY -> Outcome.REPLIED; case COMPLETION -> Outcome.COMPLETED_UNREPLIED; case FAILED -> Outcome.WORKER_FAILED; + case BACKEND_EXHAUSTED -> Outcome.BACKEND_EXHAUSTED; case QUESTION -> Outcome.QUESTION; }; } diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/Rendezvous.java b/bridged/src/main/java/dev/ltms/bridged/msg/Rendezvous.java index 5136537..cd6c7f9 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/Rendezvous.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/Rendezvous.java @@ -33,6 +33,13 @@ public final class Rendezvous { COMPLETION, /** The worker ran the turn then wedged (CB-109); {@code text} is the failure context. */ FAILED, + /** + * The turn finished without a {@code bridge_reply}, and the scrape matched the backend's + * configured usage-limit refusal pattern (CB-578 stage A); {@code text} is the reason, + * carrying the matched line. The pane is healthy — only the account is refusing — so this + * is kept separate from a session simply going {@code GONE}. + */ + BACKEND_EXHAUSTED, /** * The worker paused mid-turn to ask the primary a question (CB-205 reverse rendezvous); * {@code text} is the question and {@code turnId} correlates the primary's answer back to @@ -224,6 +231,19 @@ public final class Rendezvous { return waiter != null && waiter.complete(new Resolution(Kind.FAILED, reason)); } + /** + * Resolve a specific captured {@code waiter} as {@link Kind#BACKEND_EXHAUSTED} (CB-578 stage A): + * the turn finished with no {@code bridge_reply} and the scrape matched the backend's configured + * usage-limit pattern; {@code reason} carries the matched line. Like + * {@link #resolveCompletion(CompletableFuture, String)} it targets the exact captured send + * (CB-116). A no-op if that waiter was already resolved — first resolution wins. + * + * @return {@code true} if this call resolved the waiter, {@code false} if it was null or already resolved + */ + public boolean resolveExhausted(CompletableFuture waiter, String reason) { + return waiter != null && waiter.complete(new Resolution(Kind.BACKEND_EXHAUSTED, reason)); + } + private boolean complete(String session, Resolution resolution) { CompletableFuture waiter = waiters.get(session); return waiter != null && waiter.complete(resolution); diff --git a/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java b/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java index 9aa922e..f2b2345 100644 --- a/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java +++ b/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java @@ -403,9 +403,12 @@ public final class BridgedApp { case TIMED_OUT_QUEUED -> "queued"; case BUSY -> "busy"; case WORKER_FAILED -> "failed"; + case BACKEND_EXHAUSTED -> "backend_exhausted"; default -> "done"; // unreachable (terminal outcomes handled above) }, - "detail", reply.outcome() == MessageService.Outcome.WORKER_FAILED && reply.text() != null + "detail", (reply.outcome() == MessageService.Outcome.WORKER_FAILED + || reply.outcome() == MessageService.Outcome.BACKEND_EXHAUSTED) + && reply.text() != null ? reply.text() : "no reply within " + timeout + "ms; poll status or retry")); } diff --git a/bridged/src/test/java/dev/ltms/bridged/config/BridgedConfigTest.java b/bridged/src/test/java/dev/ltms/bridged/config/BridgedConfigTest.java index 475bc37..6274396 100644 --- a/bridged/src/test/java/dev/ltms/bridged/config/BridgedConfigTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/config/BridgedConfigTest.java @@ -1140,6 +1140,7 @@ class BridgedConfigTest { gitHostEnv: GITEA_HOST weight: 0.5 maxLoad: 2 + exhaustedPattern: "usage limit has been reached" placement: weighted lifecycle: idleTtlSeconds: 300 @@ -1171,6 +1172,8 @@ class BridgedConfigTest { assertEquals("GITEA_HOST", w.gitHostEnv()); assertEquals(0.5f, w.weight(), 0.0001f, "weight binds as a float"); assertEquals(2, w.maxLoad(), "maxLoad binds as an integer"); + assertTrue(w.hasExhaustedPattern(), "exhaustedPattern binds and enables the CB-578 stage A classification"); + assertEquals("usage limit has been reached", w.exhaustedPattern()); assertEquals("weighted", cfg.placement(), "placement binds at the top level"); assertEquals(300, cfg.lifecycle().idleTtlSeconds()); diff --git a/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java b/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java index 6cf1a40..a778b00 100644 --- a/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java @@ -12,6 +12,9 @@ import dev.ltms.bridged.msg.TurnToken; import org.junit.jupiter.api.Test; import org.slf4j.LoggerFactory; +import java.util.Set; +import java.util.regex.Pattern; + import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertTrue; @@ -23,7 +26,7 @@ class CompletionResolverTest { void skipsTheScrapeWhenNoSendIsWaiting() { FakeHerdr herdr = new FakeHerdr(); Rendezvous rendezvous = new Rendezvous(); // no waiter opened - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); resolver.resolve("term_a", null); // no in-flight turn captured for this target @@ -35,7 +38,7 @@ class CompletionResolverTest { void failSkipsTheScrapeWhenNoSendIsWaiting() { FakeHerdr herdr = new FakeHerdr(); Rendezvous rendezvous = new Rendezvous(); // no waiter opened - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); resolver.fail("term_a", null); // no in-flight turn, and no registered waiter to fall back to @@ -47,7 +50,7 @@ class CompletionResolverTest { void captureBaselineSkipsTheReadWhenNoSendIsWaiting() { FakeHerdr herdr = new FakeHerdr().readText("⏺ X\n❯ "); Rendezvous rendezvous = new Rendezvous(); // no waiter opened - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); resolver.captureBaseline("term_a", TestTurnTokens.inert("term_a")); // no send to attribute a later completion to @@ -136,7 +139,7 @@ class CompletionResolverTest { // send must NOT be resolved with the stale answer. FakeHerdr herdr = new FakeHerdr().readText("⏺ 391\n❯ "); Rendezvous rendezvous = new Rendezvous(); - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); var waiter = rendezvous.open("term_a"); // a send is blocked on this turn // The turn as captured at delivery: its waiter, and the previous turn's answer still on screen. @@ -151,7 +154,7 @@ class CompletionResolverTest { void resolvesACompletionWhoseScrapeChangedSinceDelivery() { FakeHerdr herdr = new FakeHerdr().readText("⏺ No, 391 = 17 × 23.\n❯ "); // the worker's real answer Rendezvous rendezvous = new Rendezvous(); - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); var waiter = rendezvous.open("term_a"); // Delivery baseline was the previous turn's "391"; the scrape now differs → resolve. @@ -168,7 +171,7 @@ class CompletionResolverTest { String block = "⏺ " + "x".repeat(CompletionResolver.MAX_SCRAPE_CHARS + 1) + "\n❯ "; FakeHerdr herdr = new FakeHerdr().readText(block); Rendezvous rendezvous = new Rendezvous(); - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); var waiter = rendezvous.open("term_a"); resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null)); @@ -182,7 +185,7 @@ class CompletionResolverTest { void leavesAnUnclippedCompletionPaneTailUnmarked() { FakeHerdr herdr = new FakeHerdr().readText("⏺ complete report\n❯ "); Rendezvous rendezvous = new Rendezvous(); - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); var waiter = rendezvous.open("term_a"); resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null)); @@ -194,7 +197,7 @@ class CompletionResolverTest { void resolvesSynchronouslyBeforePostTurnContextClearing() { FakeHerdr herdr = new FakeHerdr().readText("⏺ previous answer\n❯ "); Rendezvous rendezvous = new Rendezvous(); - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); var waiter = rendezvous.open("term_a"); resolver.captureBaseline("term_a", new TurnToken("term_a", waiter)); herdr.readText("⏺ answer that /clear would erase\n❯ "); @@ -216,7 +219,7 @@ class CompletionResolverTest { String longBlock = "⏺ " + "x".repeat(CompletionResolver.MAX_SCRAPE_CHARS + 500) + "\n❯ "; FakeHerdr herdr = new FakeHerdr().readText(longBlock); Rendezvous rendezvous = new Rendezvous(); - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); var waiter = rendezvous.open("term_a"); // a send is blocked on this turn resolver.captureBaseline("term_a", new TurnToken("term_a", waiter)); // baseline is the clipped >cap block @@ -236,7 +239,7 @@ class CompletionResolverTest { // No delivery baseline (e.g. the pre-turn read failed) ⇒ never suppress; the completion resolves. FakeHerdr herdr = new FakeHerdr().readText("⏺ hello\n❯ "); Rendezvous rendezvous = new Rendezvous(); - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); var waiter = rendezvous.open("term_a"); resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null)); @@ -254,7 +257,7 @@ class CompletionResolverTest { // byte-identical guard would wrongly match the empty tail and suppress. FakeHerdr herdr = new FakeHerdr().healthy(false); // agent.read throws HerdrException Rendezvous rendezvous = new Rendezvous(); - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); var waiter = rendezvous.open("term_a"); var turn = new CompletionResolver.InFlight(waiter, ""); // empty pane baselined at delivery @@ -274,7 +277,7 @@ class CompletionResolverTest { // fail must not overwrite that value, and must not even scrape the worker — nobody needs it. FakeHerdr herdr = new FakeHerdr().readText("an error screen"); Rendezvous rendezvous = new Rendezvous(); - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); var waiter = rendezvous.open("term_a"); var turn = new CompletionResolver.InFlight(waiter, null); @@ -295,7 +298,7 @@ class CompletionResolverTest { // fail falls back to the waiter currently registered on the Rendezvous and fails it. FakeHerdr herdr = new FakeHerdr().readText("stuck on an error screen"); Rendezvous rendezvous = new Rendezvous(); - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); var waiter = rendezvous.open("term_a"); // send registered, but no captureBaseline ever ran resolver.fail("term_a", null); // no in-flight turn → fall back to the registered waiter @@ -321,7 +324,7 @@ class CompletionResolverTest { try { FakeHerdr herdr = new FakeHerdr().readText("stuck on an error screen"); Rendezvous rendezvous = new Rendezvous(); - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); var waiter = rendezvous.open("term_a"); resolver.fail("term_a", null); @@ -349,7 +352,7 @@ class CompletionResolverTest { // scrape to turn N+1; targeting turn N's captured waiter makes the late completion a no-op. FakeHerdr herdr = new FakeHerdr().readText("⏺ turn N answer\n❯ "); Rendezvous rendezvous = new Rendezvous(); - CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); var waiterN = rendezvous.open("term_a"); // turn N's send // The turn as the injector captured it at delivery (waiter + pre-turn baseline). @@ -370,4 +373,88 @@ class CompletionResolverTest { "turn N stays resolved by its own reply"); assertTrue(rendezvous.isWaiting("term_a"), "turn N+1 is still awaiting its own resolution"); } + + // --- CB-578 stage A: backend-exhausted classification --------------------------------- + + @Test + void classifiesAMatchingScrapeAsBackendExhaustedInsteadOfACompletedReply() { + String block = "⏺ Working on it...\nThe usage limit has been reached. Try again later.\n❯ "; + FakeHerdr herdr = new FakeHerdr().readText(block); + Rendezvous rendezvous = new Rendezvous(); + ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached"); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns); + + var waiter = rendezvous.open("term_a"); + resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null)); + + assertTrue(waiter.isDone(), "a matching scrape still resolves the blocked send"); + assertEquals(Rendezvous.Kind.BACKEND_EXHAUSTED, waiter.getNow(null).kind(), + "not reported as a completed reply — the classification is distinct"); + } + + @Test + void theExhaustedReasonCarriesTheMatchedLine() { + String block = "⏺ Working on it...\nThe usage limit has been reached. Try again later.\n❯ "; + FakeHerdr herdr = new FakeHerdr().readText(block); + Rendezvous rendezvous = new Rendezvous(); + ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached"); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns); + + var waiter = rendezvous.open("term_a"); + resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null)); + + assertEquals("backend exhausted (usage limit): The usage limit has been reached. Try again later.", + waiter.getNow(null).text(), "the reason names the real cause and carries the matched line"); + } + + @Test + void aNonMatchingScrapeResolvesAsAnOrdinaryCompletion() { + FakeHerdr herdr = new FakeHerdr().readText("⏺ complete report\n❯ "); + Rendezvous rendezvous = new Rendezvous(); + ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached"); + CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns); + + var waiter = rendezvous.open("term_a"); + resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null)); + + assertEquals(Rendezvous.Kind.COMPLETION, waiter.getNow(null).kind(), + "a scrape that does not match the pattern is an ordinary completion"); + assertEquals("complete report", waiter.getNow(null).text()); + } + + @Test + void aProfileWithNoConfiguredPatternKeepsTodaysCompletionFallbackUnchanged() { + // Even a scrape that WOULD have matched some other profile's pattern must resolve as a + // plain completion when this target's own profile has none configured (CB-578 criterion 4). + String block = "⏺ The usage limit has been reached.\n❯ "; + FakeHerdr herdr = new FakeHerdr().readText(block); + Rendezvous rendezvous = new Rendezvous(); + CompletionResolver resolver = + new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none()); + + var waiter = rendezvous.open("term_a"); + resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null)); + + assertEquals(Rendezvous.Kind.COMPLETION, waiter.getNow(null).kind(), + "no pattern configured for this target's profile ⇒ unchanged completion-fallback behaviour"); + assertEquals("The usage limit has been reached.", waiter.getNow(null).text()); + } + + @Test + void coverageIsOffWhenNoProfileHasAPatternConfigured() { + assertEquals("off (no profile has an exhaustedPattern configured; profiles: [terra])", + CompletionResolver.coverage(Set.of("terra"), Set.of())); + } + + @Test + void coverageIsFullWhenEveryProfileHasAPatternConfigured() { + assertEquals("full (all profiles configured: [gx10, terra])", + CompletionResolver.coverage(Set.of("terra", "gx10"), Set.of("terra", "gx10"))); + } + + @Test + void coverageIsPartialAndNamesWhichProfilesAreConfigured() { + assertEquals("partial (configured: [terra]; not configured: [gx10])", + CompletionResolver.coverage(Set.of("terra", "gx10"), Set.of("terra"))); + } } diff --git a/bridged/src/test/java/dev/ltms/bridged/lead/LeadLauncherTest.java b/bridged/src/test/java/dev/ltms/bridged/lead/LeadLauncherTest.java index 3b3d771..81897f1 100644 --- a/bridged/src/test/java/dev/ltms/bridged/lead/LeadLauncherTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/lead/LeadLauncherTest.java @@ -28,7 +28,7 @@ class LeadLauncherTest { List.of("ccs", "ltms"), "tab", "bridged-workers", null, "http://127.0.0.1:8765/mcp", null, null, null, null, null, - Map.of("CLAUDE_CODE_AUTO_COMPACT_WINDOW", "300000"), null, null, true); + Map.of("CLAUDE_CODE_AUTO_COMPACT_WINDOW", "300000"), null, null, true, null); } private static BridgedConfig configWith(BridgedConfig.Leader lead) { diff --git a/bridged/src/test/java/dev/ltms/bridged/member/ClaudeCodeLauncherTest.java b/bridged/src/test/java/dev/ltms/bridged/member/ClaudeCodeLauncherTest.java index b50dbba..ca5ccad 100644 --- a/bridged/src/test/java/dev/ltms/bridged/member/ClaudeCodeLauncherTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/member/ClaudeCodeLauncherTest.java @@ -733,7 +733,7 @@ class ClaudeCodeLauncherTest { return new BridgedConfig.Profile( profile, baseUrl, "sonnet", null, "BRIDGED_WORKER_TOKEN", List.of("ccs", profile), "tab", "bridged-workers", "w #{n}", null, null, null, - null, null, null, Map.of(), null, null, true); + null, null, null, Map.of(), null, null, true, null); } @Test @@ -819,7 +819,7 @@ class ClaudeCodeLauncherTest { null, null, null, Map.of("ANTHROPIC_BASE_URL", "http://evil.example.com", "ANTHROPIC_AUTH_TOKEN", "sk-ant-bad", "JAVA_HOME", "/opt/jdk"), - null, null, true); + null, null, true, null); new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr), new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> "would-be-token").spawn(); diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java index e08a3e7..eb5de30 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java @@ -5,6 +5,7 @@ import dev.ltms.bridged.herdr.AgentStatus; import dev.ltms.bridged.herdr.FakeHerdr; import dev.ltms.bridged.herdr.HerdrException; import dev.ltms.bridged.inject.CompletionResolver; +import dev.ltms.bridged.inject.ExhaustedPatternLookup; import dev.ltms.bridged.mcp.PrimaryRegistry; import dev.ltms.bridged.inject.Injector; import org.junit.jupiter.api.BeforeEach; @@ -31,7 +32,8 @@ class MessageServiceTest { private final FakeHerdr herdr = new FakeHerdr().readText("BUILD GREEN: 391 files"); private final AgentControl agents = new AgentControl(herdr); private final Rendezvous rendezvous = new Rendezvous(); - private final CompletionResolver completion = new CompletionResolver(agents, rendezvous); + private final CompletionResolver completion = + new CompletionResolver(agents, rendezvous, ExhaustedPatternLookup.none()); private final Injector injector = new Injector(agents, completion); private final InMemoryReplyInbox inbox = new InMemoryReplyInbox(); private final MessageService messages = new MessageService(agents, injector, rendezvous, inbox); diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/RendezvousTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/RendezvousTest.java index 14199ee..1265144 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/RendezvousTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/RendezvousTest.java @@ -114,6 +114,17 @@ class RendezvousTest { "the first resolution wins; the stored value is unchanged"); } + @Test + void resolveExhaustedTwiceIsANoOpTheSecondTime() { + CompletableFuture waiter = rendezvous.open(W); + assertTrue(rendezvous.resolveExhausted(waiter, "first reason"), "the first classification resolves"); + assertFalse(rendezvous.resolveExhausted(waiter, "second reason"), + "a second exhausted resolution on an already-resolved waiter returns false"); + assertEquals(Rendezvous.Kind.BACKEND_EXHAUSTED, waiter.getNow(null).kind()); + assertEquals("first reason", waiter.getNow(null).text(), + "the first resolution wins; the stored value is unchanged"); + } + @Test void closeAskRemovesTheTurn() { Rendezvous.AskTicket t = rendezvous.openAsk(W);