Compare commits

..

5 Commits

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

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

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

Accepted the implementer's deviation from the brief. The brief asked for a
terminal HealthState BACKEND_EXHAUSTED. FleetHealth.decide() is a pure
classifier over HealthSnapshot, which carries only booleans, AgentStatus and
MemberSession.State — never pane text. FleetHealth's own javadoc already says
ERROR_ON_SCREEN 'is not decided yet because it needs a bounded pane detection
read'. A second declared-but-unproduced value would repeat a gap the class
already documents as a problem, so the signal was put where the evidence
actually lives: the completion scrape.
2026-08-15 10:00:34 +02:00
Dai Ha 541df87272 CB-578 stage A: classify a usage-limit refusal instead of a completed reply
CI / build (pull_request) Failing after 59s
CI / contract (pull_request) Successful in 1m10s
A backend that refuses on a subscription usage limit leaves the pane healthy but
the turn ends with no bridge_reply; the completion fallback used to scrape and
hand that refusal back as if it were a real answer. CompletionResolver now
matches the scrape against a per-profile exhaustedPattern (config, never a
vendor string) and resolves the send as Rendezvous.Kind/Outcome.BACKEND_EXHAUSTED
with a reason carrying the matched line, kept distinct from GONE/WORKER_FAILED.
A profile with no pattern configured is unaffected. Coverage is logged at
startup via CompletionResolver.coverage(...), naming which profiles have a
pattern and which don't, following FleetHealthMonitor.coverage's pattern.
2026-08-15 09:55:06 +02:00
17 changed files with 601 additions and 52 deletions
+7
View File
@@ -138,6 +138,12 @@ herdrSocket: ~/.config/herdr/herdr.sock
# SSH is unaffected). The token value itself is never stored in this file.
# 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
@@ -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<String, Pattern> exhaustedPatternsByProfile = new LinkedHashMap<>();
cfg.profiles().forEach((name, profile) -> {
if (profile.hasExhaustedPattern()) {
exhaustedPatternsByProfile.put(name, Pattern.compile(profile.exhaustedPattern()));
}
});
ExhaustedPatternLookup exhaustedPatterns = target -> sessions.roster().stream()
.filter(session -> target.equals(session.terminalId()))
.findFirst()
.map(session -> exhaustedPatternsByProfile.get(session.profile()))
.orElse(null);
log.info("backend-exhausted classification (CB-578 stage A): {}",
CompletionResolver.coverage(cfg.profiles().keySet(), exhaustedPatternsByProfile.keySet()));
CompletionResolver completion = new CompletionResolver(agents, rendezvous, exhaustedPatterns);
// CB-113: deliver only to an available worker (its MCP is connected), never its boot window.
// CB-301: the manager's presence bridge records availability and drives SPAWNING → READY.
MemberPresence presence = sessions.asPresence();
@@ -181,6 +181,12 @@ public record BridgedConfig(
* opposite intents). For the same reason, an {@code env:} entry naming
* {@code ANTHROPIC_BASE_URL} or {@code ANTHROPIC_AUTH_TOKEN} is refused at
* config load (CB-542): on the subscription path no guard would vet it.
* @param exhaustedPattern regex matched against a completion-fallback scrape (CB-578 stage A) to
* classify a turn that ended with no {@code bridge_reply} as the backend
* having refused on a subscription usage limit, rather than a real answer.
* {@code null}/blank ⇒ the classification never fires for this profile and
* today's completion-fallback behaviour is unchanged. Every backend words
* its refusal differently, so this is config, never a vendor string in code.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record Profile(String profile, String baseUrl, String model,
@@ -193,7 +199,8 @@ public record BridgedConfig(
Map<String, String> env,
Float weight,
Integer maxLoad,
Boolean subscription) {
Boolean subscription,
String exhaustedPattern) {
/** Peer kind spawned by {@link dev.ltms.bridged.member.ClaudeCodeLauncher} (the default). */
public static final String KIND_CLAUDE_CODE = "claude-code";
@@ -236,6 +243,9 @@ public record BridgedConfig(
weight = (weight == null || weight <= 0.0f) ? 1.0f : weight;
maxLoad = (maxLoad == null || maxLoad <= 0) ? null : maxLoad;
subscription = (subscription != null && subscription) ? Boolean.TRUE : Boolean.FALSE;
// exhaustedPattern stays null when unset/blank (opt-in) — no defaulting, no vendor
// wording: an unconfigured profile keeps today's completion-fallback behaviour exactly.
exhaustedPattern = (exhaustedPattern == null || exhaustedPattern.isBlank()) ? null : exhaustedPattern;
}
/**
@@ -248,7 +258,7 @@ public record BridgedConfig(
String placement, String workspace, String tabLabel, String mcpUrl,
String cwd, List<String> parityOverlay) {
this(profile, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel,
mcpUrl, cwd, parityOverlay, null, null, null, null, null, null, null);
mcpUrl, cwd, parityOverlay, null, null, null, null, null, null, null, null);
}
/**
@@ -260,7 +270,7 @@ public record BridgedConfig(
String placement, String workspace, String tabLabel, String mcpUrl,
String cwd, List<String> parityOverlay, String gitTokenEnv, String gitHostEnv) {
this(profile, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel,
mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, null, null, null, null, null);
mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, null, null, null, null, null, null);
}
/**
@@ -273,13 +283,14 @@ public record BridgedConfig(
String cwd, List<String> parityOverlay, String gitTokenEnv, String gitHostEnv,
String kind) {
this(profile, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel,
mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, null, null, null, null);
mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, null, null, null, null, null);
}
/** A copy with {@code profile} set — used to default a profile to its {@code workers} key. */
public Profile withProfile(String p) {
return new Profile(p, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel,
mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, env, weight, maxLoad, subscription);
mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, env, weight, maxLoad, subscription,
exhaustedPattern);
}
/** True when this profile is served by the Claude Code adapter (the default kind). */
@@ -312,7 +323,12 @@ public record BridgedConfig(
String cwd, List<String> parityOverlay, String gitTokenEnv, String gitHostEnv,
String kind, Map<String, String> env, Float weight, Integer maxLoad) {
this(profile, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel,
mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, env, weight, maxLoad, null);
mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, env, weight, maxLoad, null, null);
}
/** True when this profile's CB-578 stage A backend-exhausted classification is configured. */
public boolean hasExhaustedPattern() {
return exhaustedPattern != null;
}
/** True when this profile's workers are granted a forge token to open their own PR (CB-302). */
@@ -6,8 +6,13 @@ import dev.ltms.bridged.msg.TurnToken;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.List;
import java.util.Objects;
import java.util.Set;
import java.util.TreeSet;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.regex.Pattern;
/**
* The CB-106 completion fallback: bridges the {@link Injector}'s turn-completion signal to the
@@ -61,6 +66,7 @@ public final class CompletionResolver implements TurnListener {
private final AgentControl agents;
private final Rendezvous rendezvous;
private final ExhaustedPatternLookup exhaustedPatterns;
/**
* Per-target record of the turn currently in flight: the exact {@link Rendezvous} waiter its
@@ -78,9 +84,16 @@ public final class CompletionResolver implements TurnListener {
private final ConcurrentHashMap<String, InFlight> inFlight = new ConcurrentHashMap<>();
public CompletionResolver(AgentControl agents, Rendezvous rendezvous) {
/**
* @param exhaustedPatterns CB-578 stage A: per-target lookup for a profile's configured
* usage-limit refusal pattern. Required — there is deliberately no
* defaulting overload; a caller that does not want the classification
* must pass an explicit inert value ({@link ExhaustedPatternLookup#none()}).
*/
public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns) {
this.agents = agents;
this.rendezvous = rendezvous;
this.exhaustedPatterns = Objects.requireNonNull(exhaustedPatterns, "exhaustedPatterns");
}
@Override
@@ -157,11 +170,12 @@ public final class CompletionResolver implements TurnListener {
return;
}
String tail;
String assistantBlock = null;
int originalLength = 0;
boolean clipped = false;
boolean scrapeFailed = false;
try {
String assistantBlock = lastAssistantBlock(agents.read(target, SCRAPE_SOURCE));
assistantBlock = lastAssistantBlock(agents.read(target, SCRAPE_SOURCE));
originalLength = assistantBlock.strip().length();
clipped = originalLength > MAX_SCRAPE_CHARS;
tail = clip(assistantBlock);
@@ -184,6 +198,22 @@ public final class CompletionResolver implements TurnListener {
target);
return; // keep the in-flight record: a later genuine completion still needs it
}
// CB-578 stage A: a turn that ended with no bridge_reply AND whose scrape matches the
// backend's configured usage-limit pattern is a refusal, not an answer. Classify it as
// BACKEND_EXHAUSTED rather than handing the caller a scrape that reads like a real reply.
if (!scrapeFailed) {
Pattern exhausted = exhaustedPatterns.patternFor(target);
String matchedLine = exhausted == null ? null : firstMatchingLine(assistantBlock, exhausted);
if (matchedLine != null) {
String reason = "backend exhausted (usage limit): " + matchedLine;
if (rendezvous.resolveExhausted(waiter, reason)) {
inFlight.remove(target, turn);
log.warn("completion for {} classified BACKEND_EXHAUSTED (no bridge_reply; scrape "
+ "matched the profile's exhausted pattern): {}", target, reason);
}
return;
}
}
String completion = clipped ? tail + "\n" + CLIPPED_PANE_TAIL_MARKER : tail;
if (rendezvous.resolveCompletion(waiter, completion)) {
inFlight.remove(target, turn);
@@ -232,6 +262,45 @@ public final class CompletionResolver implements TurnListener {
}
}
/**
* The first line of {@code text} matching {@code pattern}, stripped — the CB-578 stage A
* evidence carried in a {@code BACKEND_EXHAUSTED} reason so the operator sees the real refusal
* text, never a generic label. {@code null} if no line matches.
*/
static String firstMatchingLine(String text, Pattern pattern) {
if (text == null || text.isEmpty()) return null;
for (String line : text.split("\n", -1)) {
if (pattern.matcher(line).find()) {
return line.strip();
}
}
return null;
}
/**
* Coverage summary for the CB-578 stage A exhausted-pattern classification, logged at startup
* the way {@link dev.ltms.bridged.health.FleetHealthMonitor#coverage} is — so an operator can
* see whether the classification is on, and for which profiles, without reading every
* profile's config by hand.
*
* @param allProfiles every configured profile name
* @param configuredProfiles the subset of {@code allProfiles} that carry an exhausted pattern
*/
public static String coverage(Set<String> allProfiles, Set<String> configuredProfiles) {
if (configuredProfiles.isEmpty()) {
return "off (no profile has an exhaustedPattern configured; profiles: " + sorted(allProfiles) + ")";
}
Set<String> unconfigured = new TreeSet<>(allProfiles);
unconfigured.removeAll(configuredProfiles);
return unconfigured.isEmpty()
? "full (all profiles configured: " + sorted(allProfiles) + ")"
: "partial (configured: " + sorted(configuredProfiles) + "; not configured: " + sorted(unconfigured) + ")";
}
private static List<String> sorted(Set<String> names) {
return names.stream().sorted().toList();
}
private static String clip(String s) {
if (s == null) return "";
String trimmed = s.strip();
@@ -0,0 +1,27 @@
package dev.ltms.bridged.inject;
import java.util.regex.Pattern;
/**
* Per-target lookup for a profile's configured usage-limit refusal pattern (CB-578 stage A): how
* {@link CompletionResolver} tells a backend that refused on a subscription usage limit — the
* worker's pane stays healthy, but the account is exhausted — apart from a genuine completion.
*
* <p>The pattern is always profile config, never a vendor string in Java source: every backend
* words its refusal differently, so a hardcoded sentence would only ever match one of them.
*/
@FunctionalInterface
public interface ExhaustedPatternLookup {
/** The compiled pattern configured for {@code target}'s profile, or {@code null} if none. */
Pattern patternFor(String target);
/**
* Inert lookup — no profile has a pattern configured, so the classification never fires and
* the completion fallback behaves exactly as before CB-578 stage A. The explicit stand-in a
* caller (or a test not exercising this feature) passes instead of a defaulting overload.
*/
static ExhaustedPatternLookup none() {
return target -> null;
}
}
@@ -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()
@@ -69,6 +69,14 @@ public final class MessageService {
* failure context (e.g. the error screen). Terminal, but not a successful completion.
*/
WORKER_FAILED,
/**
* The turn finished without a {@code bridge_reply} and the scrape matched the backend's
* configured usage-limit refusal pattern (CB-578 stage A); {@code text} is the reason,
* carrying the matched line. The worker's pane is healthy — only its account is refusing —
* so this is never reported as a completed reply, and is kept distinct from
* {@link #WORKER_FAILED} (a wedged worker) and a session simply going {@code GONE}.
*/
BACKEND_EXHAUSTED,
/**
* The worker paused mid-turn to ask the primary a question (CB-205); {@code text} is the
* question and {@code turnId} correlates the answer. Not terminal — the primary answers with
@@ -280,6 +288,7 @@ public final class MessageService {
case COMPLETED_UNREPLIED -> "completion_fallback";
case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, BUSY -> "timeout";
case WORKER_FAILED -> "failed";
case BACKEND_EXHAUSTED -> "backend_exhausted";
case STALE_TURN, QUESTION -> null; // not a completed delegation
};
}
@@ -591,9 +600,11 @@ public final class MessageService {
String source = r.outcome() == Outcome.REPLIED ? "reply" : "transcript";
return new TaskView(ticket, Phase.DONE, r.text(), source, null, null);
}
// A wedged worker (CB-109) carries the error context as its reason; the timeout/busy
// outcomes carry none, so fall back to the outcome name.
String detail = r.outcome() == Outcome.WORKER_FAILED && r.text() != null
// A wedged worker (CB-109) or a backend-exhausted classification (CB-578 stage A) carries
// the real cause as its reason; the timeout/busy outcomes carry none, so fall back to the
// outcome name.
boolean carriesReason = r.outcome() == Outcome.WORKER_FAILED || r.outcome() == Outcome.BACKEND_EXHAUSTED;
String detail = carriesReason && r.text() != null
? r.text()
: "no reply — " + r.outcome().name().toLowerCase();
return new TaskView(ticket, Phase.FAILED, null, null, detail, null);
@@ -669,6 +680,7 @@ public final class MessageService {
case REPLY -> Outcome.REPLIED;
case COMPLETION -> Outcome.COMPLETED_UNREPLIED;
case FAILED -> Outcome.WORKER_FAILED;
case BACKEND_EXHAUSTED -> Outcome.BACKEND_EXHAUSTED;
case QUESTION -> Outcome.QUESTION;
};
}
@@ -33,6 +33,13 @@ public final class Rendezvous {
COMPLETION,
/** The worker ran the turn then wedged (CB-109); {@code text} is the failure context. */
FAILED,
/**
* The turn finished without a {@code bridge_reply}, and the scrape matched the backend's
* configured usage-limit refusal pattern (CB-578 stage A); {@code text} is the reason,
* carrying the matched line. The pane is healthy — only the account is refusing — so this
* is kept separate from a session simply going {@code GONE}.
*/
BACKEND_EXHAUSTED,
/**
* The worker paused mid-turn to ask the primary a question (CB-205 reverse rendezvous);
* {@code text} is the question and {@code turnId} correlates the primary's answer back to
@@ -224,6 +231,19 @@ public final class Rendezvous {
return waiter != null && waiter.complete(new Resolution(Kind.FAILED, reason));
}
/**
* Resolve a specific captured {@code waiter} as {@link Kind#BACKEND_EXHAUSTED} (CB-578 stage A):
* the turn finished with no {@code bridge_reply} and the scrape matched the backend's configured
* usage-limit pattern; {@code reason} carries the matched line. Like
* {@link #resolveCompletion(CompletableFuture, String)} it targets the exact captured send
* (CB-116). A no-op if that waiter was already resolved — first resolution wins.
*
* @return {@code true} if this call resolved the waiter, {@code false} if it was null or already resolved
*/
public boolean resolveExhausted(CompletableFuture<Resolution> waiter, String reason) {
return waiter != null && waiter.complete(new Resolution(Kind.BACKEND_EXHAUSTED, reason));
}
private boolean complete(String session, Resolution resolution) {
CompletableFuture<Resolution> waiter = waiters.get(session);
return waiter != null && waiter.complete(resolution);
@@ -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"));
}
@@ -197,27 +197,44 @@ public final class SessionManager implements TurnListener {
MemberSession removed = registry.remove(paneId);
boolean preserveWorktree = cause == ReleaseCause.SHUTDOWN;
if (removed != null) {
memberLifecycle.released(removed.terminalId());
log.debug("releasing session pane={} terminal={} state={} cause={}",
removed.paneId(), removed.terminalId(), removed.state(), cause);
if (preserveWorktree && removed.worktree() != null) {
logPreservedForShutdown(removed);
} else if (removed.worktree() != null && worktrees.hasUncommitted(removed.worktree())) {
// CB-576: a release that would otherwise remove the worktree finds it holding
// uncommitted work the bridge cannot see. A worker that ends a turn without
// committing (normally because it stopped to ask a question or refused the turn)
// has its only copy of that work in the worktree. Remove would --force-delete it,
// so preserve the directory and tell an operator where to find it.
try {
memberLifecycle.released(removed.terminalId());
log.debug("releasing session pane={} terminal={} state={} cause={}",
removed.paneId(), removed.terminalId(), removed.state(), cause);
if (preserveWorktree && removed.worktree() != null) {
logPreservedForShutdown(removed);
} else if (removed.worktree() != null && worktrees.hasUncommitted(removed.worktree())) {
// CB-576: a release that would otherwise remove the worktree finds it holding
// uncommitted work the bridge cannot see. A worker that ends a turn without
// committing (normally because it stopped to ask a question or refused the turn)
// has its only copy of that work in the worktree. Remove would --force-delete it,
// so preserve the directory and tell an operator where to find it.
preserveWorktree = true;
log.warn("release {} preserves dirty worktree {} for pane={} terminal={}: "
+ "the worktree holds uncommitted changes that --force remove would destroy",
cause, removed.worktree(), removed.paneId(), removed.terminalId());
}
} catch (RuntimeException e) {
// CB-581: hasUncommitted shells out to `git status` and can throw on a non-zero
// exit. We can no longer tell whether the worktree holds uncommitted work, so fail
// toward the safe answer and preserve it — deleting on a guess can destroy work
// that has no other copy (CB-576), while keeping it on a false alarm only costs
// disk. The exception must not propagate: the pane still has to stop below.
preserveWorktree = true;
log.warn("release {} preserves dirty worktree {} for pane={} terminal={}: "
+ "the worktree holds uncommitted changes that --force remove would destroy",
cause, removed.worktree(), removed.paneId(), removed.terminalId());
log.warn("release {} could not tell whether worktree {} for pane={} terminal={} has "
+ "uncommitted changes; preserving it rather than risk destroying unsaved work: {}",
cause, removed.worktree(), removed.paneId(), removed.terminalId(), e.toString());
} finally {
// CB-516/CB-581: a send still waiting on this worker can never be answered now, no
// matter what happened above. Tell the listener BEFORE the pane is torn down, so a
// blocked caller fails fast with a real reason instead of sitting on a rendezvous
// nothing will ever resolve.
notifyReleased(removed.terminalId());
}
// CB-516: a send still waiting on this worker can never be answered now. Tell the
// listener BEFORE the pane is torn down, so a blocked caller fails fast with a real
// reason instead of sitting on a rendezvous nothing will ever resolve.
notifyReleased(removed.terminalId());
}
// CB-581: the pane must always stop, even if the dirty check above threw. A session removed
// from the registry with no pane stop is an orphaned pane — a live terminal burning a fleet
// slot that no longer appears in the roster and can never be reclaimed.
launcher.stop(paneId);
if (removed != null && !preserveWorktree && removed.worktree() != null) {
worktrees.remove(worktrees.repoRoot(removed.cwd()), removed.worktree());
@@ -539,8 +556,15 @@ public final class SessionManager implements TurnListener {
log.debug("reaping idle session terminal={} pane={}: idle {}s exceeds the {}s ttl",
s.terminalId(), s.paneId(), TimeUnit.NANOSECONDS.toSeconds(idleNanos),
TimeUnit.NANOSECONDS.toSeconds(idleTtlNanos));
release(s.paneId());
reaped++;
// CB-581: one session that fails to release must not abort the whole reaping pass —
// match drainAll's per-session try/catch so the rest of the roster still gets reaped.
try {
release(s.paneId());
reaped++;
} catch (RuntimeException e) {
log.warn("reap failed for pane={} terminal={} worktree={}; continuing with "
+ "remaining sessions", s.paneId(), s.terminalId(), s.worktree(), e);
}
}
}
return reaped;
@@ -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());
@@ -12,6 +12,9 @@ import dev.ltms.bridged.msg.TurnToken;
import org.junit.jupiter.api.Test;
import org.slf4j.LoggerFactory;
import java.util.Set;
import java.util.regex.Pattern;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
@@ -23,7 +26,7 @@ class CompletionResolverTest {
void skipsTheScrapeWhenNoSendIsWaiting() {
FakeHerdr herdr = new FakeHerdr();
Rendezvous rendezvous = new Rendezvous(); // no waiter opened
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
resolver.resolve("term_a", null); // no in-flight turn captured for this target
@@ -35,7 +38,7 @@ class CompletionResolverTest {
void failSkipsTheScrapeWhenNoSendIsWaiting() {
FakeHerdr herdr = new FakeHerdr();
Rendezvous rendezvous = new Rendezvous(); // no waiter opened
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
resolver.fail("term_a", null); // no in-flight turn, and no registered waiter to fall back to
@@ -47,7 +50,7 @@ class CompletionResolverTest {
void captureBaselineSkipsTheReadWhenNoSendIsWaiting() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ X\n❯ ");
Rendezvous rendezvous = new Rendezvous(); // no waiter opened
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
resolver.captureBaseline("term_a", TestTurnTokens.inert("term_a")); // no send to attribute a later completion to
@@ -136,7 +139,7 @@ class CompletionResolverTest {
// send must NOT be resolved with the stale answer.
FakeHerdr herdr = new FakeHerdr().readText("⏺ 391\n❯ ");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
var waiter = rendezvous.open("term_a"); // a send is blocked on this turn
// The turn as captured at delivery: its waiter, and the previous turn's answer still on screen.
@@ -151,7 +154,7 @@ class CompletionResolverTest {
void resolvesACompletionWhoseScrapeChangedSinceDelivery() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ No, 391 = 17 × 23.\n❯ "); // the worker's real answer
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
var waiter = rendezvous.open("term_a");
// Delivery baseline was the previous turn's "391"; the scrape now differs → resolve.
@@ -168,7 +171,7 @@ class CompletionResolverTest {
String block = "⏺ " + "x".repeat(CompletionResolver.MAX_SCRAPE_CHARS + 1) + "\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
@@ -182,7 +185,7 @@ class CompletionResolverTest {
void leavesAnUnclippedCompletionPaneTailUnmarked() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ complete report\n❯ ");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
@@ -194,7 +197,7 @@ class CompletionResolverTest {
void resolvesSynchronouslyBeforePostTurnContextClearing() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ previous answer\n❯ ");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
var waiter = rendezvous.open("term_a");
resolver.captureBaseline("term_a", new TurnToken("term_a", waiter));
herdr.readText("⏺ answer that /clear would erase\n❯ ");
@@ -216,7 +219,7 @@ class CompletionResolverTest {
String longBlock = "⏺ " + "x".repeat(CompletionResolver.MAX_SCRAPE_CHARS + 500) + "\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(longBlock);
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
var waiter = rendezvous.open("term_a"); // a send is blocked on this turn
resolver.captureBaseline("term_a", new TurnToken("term_a", waiter)); // baseline is the clipped >cap block
@@ -236,7 +239,7 @@ class CompletionResolverTest {
// No delivery baseline (e.g. the pre-turn read failed) ⇒ never suppress; the completion resolves.
FakeHerdr herdr = new FakeHerdr().readText("⏺ hello\n❯ ");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
@@ -254,7 +257,7 @@ class CompletionResolverTest {
// byte-identical guard would wrongly match the empty tail and suppress.
FakeHerdr herdr = new FakeHerdr().healthy(false); // agent.read throws HerdrException
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
var waiter = rendezvous.open("term_a");
var turn = new CompletionResolver.InFlight(waiter, ""); // empty pane baselined at delivery
@@ -274,7 +277,7 @@ class CompletionResolverTest {
// fail must not overwrite that value, and must not even scrape the worker — nobody needs it.
FakeHerdr herdr = new FakeHerdr().readText("an error screen");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
var waiter = rendezvous.open("term_a");
var turn = new CompletionResolver.InFlight(waiter, null);
@@ -295,7 +298,7 @@ class CompletionResolverTest {
// fail falls back to the waiter currently registered on the Rendezvous and fails it.
FakeHerdr herdr = new FakeHerdr().readText("stuck on an error screen");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
var waiter = rendezvous.open("term_a"); // send registered, but no captureBaseline ever ran
resolver.fail("term_a", null); // no in-flight turn → fall back to the registered waiter
@@ -321,7 +324,7 @@ class CompletionResolverTest {
try {
FakeHerdr herdr = new FakeHerdr().readText("stuck on an error screen");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
var waiter = rendezvous.open("term_a");
resolver.fail("term_a", null);
@@ -349,7 +352,7 @@ class CompletionResolverTest {
// scrape to turn N+1; targeting turn N's captured waiter makes the late completion a no-op.
FakeHerdr herdr = new FakeHerdr().readText("⏺ turn N answer\n❯ ");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
var waiterN = rendezvous.open("term_a"); // turn N's send
// The turn as the injector captured it at delivery (waiter + pre-turn baseline).
@@ -370,4 +373,88 @@ class CompletionResolverTest {
"turn N stays resolved by its own reply");
assertTrue(rendezvous.isWaiting("term_a"), "turn N+1 is still awaiting its own resolution");
}
// --- CB-578 stage A: backend-exhausted classification ---------------------------------
@Test
void classifiesAMatchingScrapeAsBackendExhaustedInsteadOfACompletedReply() {
String block = "⏺ Working on it...\nThe usage limit has been reached. Try again later.\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached");
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns);
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
assertTrue(waiter.isDone(), "a matching scrape still resolves the blocked send");
assertEquals(Rendezvous.Kind.BACKEND_EXHAUSTED, waiter.getNow(null).kind(),
"not reported as a completed reply — the classification is distinct");
}
@Test
void theExhaustedReasonCarriesTheMatchedLine() {
String block = "⏺ Working on it...\nThe usage limit has been reached. Try again later.\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached");
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns);
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
assertEquals("backend exhausted (usage limit): The usage limit has been reached. Try again later.",
waiter.getNow(null).text(), "the reason names the real cause and carries the matched line");
}
@Test
void aNonMatchingScrapeResolvesAsAnOrdinaryCompletion() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ complete report\n❯ ");
Rendezvous rendezvous = new Rendezvous();
ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached");
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns);
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
assertEquals(Rendezvous.Kind.COMPLETION, waiter.getNow(null).kind(),
"a scrape that does not match the pattern is an ordinary completion");
assertEquals("complete report", waiter.getNow(null).text());
}
@Test
void aProfileWithNoConfiguredPatternKeepsTodaysCompletionFallbackUnchanged() {
// Even a scrape that WOULD have matched some other profile's pattern must resolve as a
// plain completion when this target's own profile has none configured (CB-578 criterion 4).
String block = "⏺ The usage limit has been reached.\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver =
new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none());
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
assertEquals(Rendezvous.Kind.COMPLETION, waiter.getNow(null).kind(),
"no pattern configured for this target's profile ⇒ unchanged completion-fallback behaviour");
assertEquals("The usage limit has been reached.", waiter.getNow(null).text());
}
@Test
void coverageIsOffWhenNoProfileHasAPatternConfigured() {
assertEquals("off (no profile has an exhaustedPattern configured; profiles: [terra])",
CompletionResolver.coverage(Set.of("terra"), Set.of()));
}
@Test
void coverageIsFullWhenEveryProfileHasAPatternConfigured() {
assertEquals("full (all profiles configured: [gx10, terra])",
CompletionResolver.coverage(Set.of("terra", "gx10"), Set.of("terra", "gx10")));
}
@Test
void coverageIsPartialAndNamesWhichProfilesAreConfigured() {
assertEquals("partial (configured: [terra]; not configured: [gx10])",
CompletionResolver.coverage(Set.of("terra", "gx10"), Set.of("terra")));
}
}
@@ -28,7 +28,7 @@ class LeadLauncherTest {
List.of("ccs", "ltms"), "tab", "bridged-workers", null,
"http://127.0.0.1:8765/mcp", null, null,
null, null, null,
Map.of("CLAUDE_CODE_AUTO_COMPACT_WINDOW", "300000"), null, null, true);
Map.of("CLAUDE_CODE_AUTO_COMPACT_WINDOW", "300000"), null, null, true, null);
}
private static BridgedConfig configWith(BridgedConfig.Leader lead) {
@@ -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();
@@ -5,6 +5,7 @@ import dev.ltms.bridged.herdr.AgentStatus;
import dev.ltms.bridged.herdr.FakeHerdr;
import dev.ltms.bridged.herdr.HerdrException;
import dev.ltms.bridged.inject.CompletionResolver;
import dev.ltms.bridged.inject.ExhaustedPatternLookup;
import dev.ltms.bridged.mcp.PrimaryRegistry;
import dev.ltms.bridged.inject.Injector;
import org.junit.jupiter.api.BeforeEach;
@@ -31,7 +32,8 @@ class MessageServiceTest {
private final FakeHerdr herdr = new FakeHerdr().readText("BUILD GREEN: 391 files");
private final AgentControl agents = new AgentControl(herdr);
private final Rendezvous rendezvous = new Rendezvous();
private final CompletionResolver completion = new CompletionResolver(agents, rendezvous);
private final CompletionResolver completion =
new CompletionResolver(agents, rendezvous, ExhaustedPatternLookup.none());
private final Injector injector = new Injector(agents, completion);
private final InMemoryReplyInbox inbox = new InMemoryReplyInbox();
private final MessageService messages = new MessageService(agents, injector, rendezvous, inbox);
@@ -114,6 +114,17 @@ class RendezvousTest {
"the first resolution wins; the stored value is unchanged");
}
@Test
void resolveExhaustedTwiceIsANoOpTheSecondTime() {
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.open(W);
assertTrue(rendezvous.resolveExhausted(waiter, "first reason"), "the first classification resolves");
assertFalse(rendezvous.resolveExhausted(waiter, "second reason"),
"a second exhausted resolution on an already-resolved waiter returns false");
assertEquals(Rendezvous.Kind.BACKEND_EXHAUSTED, waiter.getNow(null).kind());
assertEquals("first reason", waiter.getNow(null).text(),
"the first resolution wins; the stored value is unchanged");
}
@Test
void closeAskRemovesTheTurn() {
Rendezvous.AskTicket t = rendezvous.openAsk(W);
@@ -46,6 +46,82 @@ class SessionManagerTest {
return sessionManager(herdr, clock, 0);
}
private SessionManager sessionManager(FakeHerdr herdr, Worktrees worktrees) {
return sessionManager(herdr, worktrees, System::nanoTime);
}
private SessionManager sessionManager(FakeHerdr herdr, Worktrees worktrees, LongSupplier clock) {
BridgedConfig.Profile cfg = new BridgedConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "BRIDGED_WORKER_TOKEN",
List.of("ccs", "ltms-local"), "tab", "bridged-workers",
"worker: {profile} #{n}", null, null, null);
ClaudeCodeLauncher workers = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
return new SessionManager(workers, worktrees, clock);
}
/**
* CB-581: a {@link Worktrees} test double whose {@code hasUncommitted} and {@code remove} can
* be told to throw, so {@link SessionManager#release} can be exercised against exactly the
* failure {@code GitWorktrees} produces when {@code git status}/{@code git worktree remove}
* exits non-zero.
*/
private static final class RecordingWorktrees implements Worktrees {
private final List<String> removeCalls = new java.util.ArrayList<>();
private final java.util.Set<String> failRemoveFor = new java.util.HashSet<>();
private volatile boolean dirty = false;
private volatile RuntimeException hasUncommittedFailure;
RecordingWorktrees dirty(boolean dirty) {
this.dirty = dirty;
return this;
}
RecordingWorktrees failHasUncommittedWith(RuntimeException e) {
this.hasUncommittedFailure = e;
return this;
}
RecordingWorktrees failRemoveFor(String worktreePath) {
failRemoveFor.add(worktreePath);
return this;
}
@Override
public String add(String repoRoot, String branch, String baseRef) {
return "/wt/" + branch.replace('/', '_');
}
@Override
public void remove(String repoRoot, String worktreePath) {
if (failRemoveFor.contains(worktreePath)) {
throw new WorktreeException("simulated remove failure for " + worktreePath);
}
removeCalls.add(worktreePath);
}
@Override
public boolean hasUncommitted(String worktreePath) {
if (hasUncommittedFailure != null) {
throw hasUncommittedFailure;
}
return dirty;
}
@Override
public void overlayParity(String repoRoot, String worktreePath, List<String> overlay) {
}
@Override
public String repoRoot(String cwd) {
return "/repo";
}
List<String> removeCalls() {
return List.copyOf(removeCalls);
}
}
private SessionManager sessionManager(FakeHerdr herdr, LongSupplier clock, int contextCap) {
return sessionManager(herdr, clock, contextCap, false);
}
@@ -525,4 +601,172 @@ class SessionManagerTest {
"a listener failure must never prevent the teardown it is reacting to");
assertTrue(sessions.get(s.paneId()).isEmpty(), "and the session is still deregistered");
}
// --- CB-581: a throw inside release() must not orphan the pane or abort reapIdle -----------
@Test
void releasePreservesWorktreeWhenDirtyCheckThrows() {
FakeHerdr herdr = new FakeHerdr();
RecordingWorktrees worktrees = new RecordingWorktrees();
SessionManager sessions = sessionManager(herdr, worktrees);
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-581a", null));
worktrees.failHasUncommittedWith(new WorktreeException("git status exited 128"));
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
ch.qos.logback.classic.Logger sessionLog =
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(SessionManager.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.setContext(ctx);
appender.start();
sessionLog.addAppender(appender);
sessionLog.setLevel(Level.WARN);
try {
assertDoesNotThrow(() -> sessions.release(s.paneId()),
"a throwing dirty check must not abort the release");
assertTrue(worktrees.removeCalls().isEmpty(),
"the worktree is preserved when its dirty state cannot be determined");
String warn = appender.list.stream()
.filter(e -> e.getLevel().equals(Level.WARN))
.map(ILoggingEvent::getFormattedMessage)
.filter(m -> m.contains(s.worktree()))
.findFirst()
.orElse("no warn logged naming the worktree");
assertTrue(warn.contains(s.paneId()), "the WARN names the pane: " + warn);
assertTrue(warn.contains(s.terminalId()), "the WARN names the terminal: " + warn);
} finally {
sessionLog.detachAppender(appender);
}
}
@Test
void releaseStillStopsThePaneWhenDirtyCheckThrows() {
FakeHerdr herdr = new FakeHerdr();
RecordingWorktrees worktrees = new RecordingWorktrees();
SessionManager sessions = sessionManager(herdr, worktrees);
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-581b", null));
worktrees.failHasUncommittedWith(new WorktreeException("git status exited 128"));
sessions.release(s.paneId());
assertEquals(1, paneCloseCallsFor(herdr, "w9:pRoot_1"),
"the pane is stopped exactly once even though the dirty check threw");
}
@Test
void releaseStillNotifiesTheListenerWhenDirtyCheckThrows() {
FakeHerdr herdr = new FakeHerdr();
RecordingWorktrees worktrees = new RecordingWorktrees();
SessionManager sessions = sessionManager(herdr, worktrees);
java.util.List<String> released = new java.util.concurrent.CopyOnWriteArrayList<>();
sessions.onRelease(released::add);
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-581c", null));
worktrees.failHasUncommittedWith(new WorktreeException("git status exited 128"));
sessions.release(s.paneId());
assertEquals(java.util.List.of(s.terminalId()), released,
"a blocked caller must still be told the terminal was released, even though the "
+ "dirty check threw");
}
@Test
void reapIdleSurvivesOneSessionThatFailsToRelease() {
long[] clock = {0};
FakeHerdr herdr = new FakeHerdr();
RecordingWorktrees worktrees = new RecordingWorktrees();
SessionManager sessions = sessionManager(herdr, worktrees, () -> clock[0]);
MemberSession a = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-581d", null));
MemberSession b = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-581e", null));
MemberSession c = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-581f", null));
sessions.asPresence().markPresent(a.terminalId());
sessions.asPresence().markPresent(b.terminalId());
sessions.asPresence().markPresent(c.terminalId());
// The middle session's worktree removal fails — release() propagates that, so this is the
// one call reapIdle's per-session guard must survive without skipping the rest of the pass.
worktrees.failRemoveFor(b.worktree());
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
ch.qos.logback.classic.Logger sessionLog =
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(SessionManager.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.setContext(ctx);
appender.start();
sessionLog.addAppender(appender);
sessionLog.setLevel(Level.WARN);
int reaped;
try {
clock[0] = 100;
reaped = sessions.reapIdle(10);
String warn = appender.list.stream()
.filter(e -> e.getLevel().equals(Level.WARN))
.map(ILoggingEvent::getFormattedMessage)
.filter(m -> m.contains(b.paneId()))
.findFirst()
.orElse("no reap-failure WARN logged");
assertTrue(warn.contains(b.terminalId()), "the WARN names the failed session's terminal: " + warn);
assertTrue(warn.contains(b.worktree()), "the WARN names the failed session's worktree: " + warn);
} finally {
sessionLog.detachAppender(appender);
}
assertEquals(2, reaped, "the middle session's failure is logged, not counted as reaped");
assertTrue(sessions.get(a.paneId()).isEmpty(), "the first session is still released");
assertTrue(sessions.get(c.paneId()).isEmpty(), "the third session is still released");
assertTrue(sessions.get(b.paneId()).isEmpty(),
"the middle session is still deregistered even though its worktree removal threw");
assertEquals(1, paneCloseCallsFor(herdr, "w9:pRoot_1"), "the first pane is stopped");
assertEquals(1, paneCloseCallsFor(herdr, "w9:pRoot_2"),
"the middle pane is still stopped even though its worktree removal failed");
assertEquals(1, paneCloseCallsFor(herdr, "w9:pRoot_3"), "the third pane is stopped");
}
@Test
void unchangedRegressionCleanCompletedReleaseStillRemovesTheWorktree() {
FakeHerdr herdr = new FakeHerdr();
RecordingWorktrees worktrees = new RecordingWorktrees();
SessionManager sessions = sessionManager(herdr, worktrees);
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-581g", null));
sessions.release(s.paneId());
assertEquals(List.of(s.worktree()), worktrees.removeCalls(),
"COMPLETED release of a clean worktree still removes it");
}
@Test
void unchangedRegressionDirtyCompletedReleaseStillPreservesTheWorktree() {
FakeHerdr herdr = new FakeHerdr();
RecordingWorktrees worktrees = new RecordingWorktrees().dirty(true);
SessionManager sessions = sessionManager(herdr, worktrees);
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-581h", null));
sessions.release(s.paneId());
assertTrue(worktrees.removeCalls().isEmpty(),
"COMPLETED release of a dirty worktree still preserves it");
}
@Test
void unchangedRegressionShutdownDrainStillPreservesTheWorktree() {
FakeHerdr herdr = new FakeHerdr();
RecordingWorktrees worktrees = new RecordingWorktrees();
SessionManager sessions = sessionManager(herdr, worktrees);
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-581i", null));
sessions.asPresence().markPresent(s.terminalId());
sessions.drainAll(TimeUnit.MILLISECONDS.toNanos(100));
assertTrue(worktrees.removeCalls().isEmpty(), "SHUTDOWN drain still preserves the worktree");
}
}