Compare commits

..

2 Commits

Author SHA1 Message Date
Dai Ha bbf68f3e3c CB-201: cover losing completion CAS
CI / contract (pull_request) Successful in 44s
CI / build (pull_request) Successful in 1m37s
2026-09-03 10:35:36 +07:00
Dai Ha fe2e5ede34 CB-201: retain backend failure outcome
CI / build (pull_request) Successful in 1m11s
CI / contract (pull_request) Successful in 1m23s
2026-09-03 10:26:54 +07:00
7 changed files with 295 additions and 390 deletions
@@ -1,35 +0,0 @@
package dev.ltms.fleet.inject;
import java.util.regex.Pattern;
/**
* Per-target lookup for a profile's configured backend-error pattern (fleetd#201 / #227): how
* {@link CompletionResolver} recognizes a backend that failed a turn outright (a credential
* outage, a provider 5xx) from whatever ends up on the member's pane, apart from a genuine
* completion.
*
* <p>The pattern is always profile config, never a vendor string in Java source: every backend
* words its failure differently, so a hardcoded sentence would only ever match one of them. A
* target with no configured pattern is not "off" the way {@link ExhaustedPatternLookup#none()}
* is — {@link CompletionResolver} falls back to its narrow {@code (?i)\bAPI Error\s*:}
* compatibility pattern instead, so classification still happens, just without a profile-specific
* match. {@link #legacy()} is the explicit stand-in existing (pre-fleetd#201) callers pass to get
* exactly that fallback-only behavior until a caller wires a real, profile-driven lookup
* (fleetd#201 Unit 5).
*/
@FunctionalInterface
public interface BackendErrorPatternLookup {
/** The compiled pattern configured for {@code target}'s profile, or {@code null} if none. */
Pattern patternFor(String target);
/**
* Legacy lookup — no profile has a configured pattern, so {@link CompletionResolver} classifies
* every target using only its built-in {@code (?i)\bAPI Error\s*:} compatibility pattern. The
* explicit stand-in existing constructors pass so their behavior is unchanged until a caller
* wires a real, profile-driven lookup.
*/
static BackendErrorPatternLookup legacy() {
return target -> null;
}
}
@@ -1,35 +0,0 @@
package dev.ltms.fleet.inject;
/**
* Notified when {@link CompletionResolver} actually delivers a typed backend-error classification
* to a waiting send (fleetd#201 / #227) — never on a race that lost. {@link CompletionResolver}
* calls this only after {@code Rendezvous.resolveFailure} returns {@code true} for that exact
* waiter, mirroring the win-only race rule {@link ExhaustionSink} already uses.
*
* <p>The public send result is unchanged by this classification — it is still a failed send
* ({@code Rendezvous.Kind#FAILED}); this sink is the internal seam a later stage (fleetd#201 Unit
* 5, the policy that cools a credential off after two such failures in 60 seconds and tells the
* lead) consumes. {@link CompletionResolver} knows only {@code target} (a herdr terminal id); it
* has no notion of profiles or credentials, so mapping {@code target} to whatever should be
* quarantined is entirely the sink's job — exactly like {@link ExhaustionSink}.
*/
@FunctionalInterface
public interface BackendErrorSink {
/**
* @param target the herdr terminal id whose turn was classified as a backend error
* @param matchedLine the single pane line that matched the backend-error pattern
* @param reason the full failure reason carried by the classification (the matched line
* plus whatever pane context the classifying path attaches)
*/
void onBackendError(String target, String matchedLine, String reason);
/**
* Inert sink — nothing happens on a backend-error classification. The explicit stand-in a
* caller (or a test not exercising this feature) passes instead of a defaulting overload,
* exactly like {@link ExhaustionSink#none()}.
*/
static BackendErrorSink none() {
return (target, matchedLine, reason) -> { };
}
}
@@ -89,8 +89,6 @@ public final class CompletionResolver implements TurnListener {
private final Rendezvous rendezvous;
private final ExhaustedPatternLookup exhaustedPatterns;
private final ExhaustionSink exhaustionSink;
private final BackendErrorPatternLookup backendErrorPatterns;
private final BackendErrorSink backendErrorSink;
private final LongSupplier nowNanos;
/**
@@ -131,48 +129,10 @@ public final class CompletionResolver implements TurnListener {
* classification actually resolves a waiter. Required for the same
* reason as {@code exhaustedPatterns} — pass {@link ExhaustionSink#none()}
* to opt out.
*
* <p>Transition constructor (fleetd#201 Unit 1): delegates to the full constructor below with
* {@link BackendErrorPatternLookup#legacy()} and {@link BackendErrorSink#none()}, so every
* existing caller keeps today's behavior — the narrow built-in {@code API Error:} match, no
* sink notified — until a caller wires a real, profile-driven backend-error lookup and sink
* (fleetd#201 Unit 5).
*/
public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns,
ExhaustionSink exhaustionSink) {
this(agents, rendezvous, exhaustedPatterns, exhaustionSink,
BackendErrorPatternLookup.legacy(), BackendErrorSink.none(), System::nanoTime);
}
/**
* Transition constructor (fleetd#201 Unit 1): same legacy backend-error defaults as the 4-arg
* constructor above, but with the injectable clock. Kept so existing fleetd#164 timing tests
* (and any other caller that wants a controllable clock but not the backend-error feature)
* keep compiling unchanged.
*/
public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns,
ExhaustionSink exhaustionSink, LongSupplier nowNanos) {
this(agents, rendezvous, exhaustedPatterns, exhaustionSink,
BackendErrorPatternLookup.legacy(), BackendErrorSink.none(), nowNanos);
}
/**
* Production constructor (fleetd#201 Unit 1): adds the target-keyed backend-error pattern
* lookup and sink alongside the existing exhaustion pair. Required, like {@code exhaustedPatterns}
* and {@code exhaustionSink} — pass {@link BackendErrorPatternLookup#legacy()} and
* {@link BackendErrorSink#none()} to opt out.
*
* @param backendErrorPatterns fleetd#201: per-target lookup for a profile's configured
* backend-error pattern; a target with none configured is matched
* against the built-in compatibility pattern instead (never "off").
* @param backendErrorSink fleetd#201: notified when a backend-error classification actually
* resolves a waiter — never on a race that lost.
*/
public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns,
ExhaustionSink exhaustionSink, BackendErrorPatternLookup backendErrorPatterns,
BackendErrorSink backendErrorSink) {
this(agents, rendezvous, exhaustedPatterns, exhaustionSink, backendErrorPatterns, backendErrorSink,
System::nanoTime);
this(agents, rendezvous, exhaustedPatterns, exhaustionSink, System::nanoTime);
}
/**
@@ -185,14 +145,11 @@ public final class CompletionResolver implements TurnListener {
* package.
*/
public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns,
ExhaustionSink exhaustionSink, BackendErrorPatternLookup backendErrorPatterns,
BackendErrorSink backendErrorSink, LongSupplier nowNanos) {
ExhaustionSink exhaustionSink, LongSupplier nowNanos) {
this.agents = agents;
this.rendezvous = rendezvous;
this.exhaustedPatterns = Objects.requireNonNull(exhaustedPatterns, "exhaustedPatterns");
this.exhaustionSink = Objects.requireNonNull(exhaustionSink, "exhaustionSink");
this.backendErrorPatterns = Objects.requireNonNull(backendErrorPatterns, "backendErrorPatterns");
this.backendErrorSink = Objects.requireNonNull(backendErrorSink, "backendErrorSink");
this.nowNanos = Objects.requireNonNull(nowNanos, "nowNanos");
}
@@ -275,7 +232,7 @@ public final class CompletionResolver implements TurnListener {
// screen, since that is usually the backend's own error.
long elapsedNanos = nowNanos.getAsLong() - turn.deliveredAtNanos();
if (elapsedNanos < MIN_TURN_NANOS) {
failTooFast(target, turn, waiter, elapsedNanos);
fail(target, turn, tooFastReason(target, elapsedNanos));
return;
}
String tail;
@@ -345,27 +302,18 @@ public final class CompletionResolver implements TurnListener {
}
return;
}
// fleetd#164 (part 2) / fleetd#201: a scrape that read cleanly and produced content still
// isn't a real reply when that content is the backend's own rejection (e.g. an HTTP 400
// before the worker did any work). Classify it as a failure naming the member, rather than
// handing the caller a scrape that reads like a completed answer, and — only on the
// resolution that actually wins the race, mirroring the exhaustion sink above — notify the
// typed backend-error sink so a later stage can act on repeated failures.
String backendError = firstMatchingLine(assistantBlock, backendErrorPatternOrFallback(target));
// fleetd#164 (part 2): a scrape that read cleanly and produced content still isn't a real
// reply when that content is the backend's own rejection (e.g. an HTTP 400 before the worker
// did any work). Classify it as a failure naming the member, rather than handing the caller a
// scrape that reads like a completed answer.
String backendError = firstMatchingLine(assistantBlock, BACKEND_ERROR);
if (backendError != null) {
// Carry the whole scrape, not just the matched line. The pattern is a heuristic: a member
// that forgot fleet_reply while reporting *about* a backend error matches it too. Failing
// is still right — the caller must not read a scrape as an answer — but dropping the rest
// of the pane would destroy the report, which is the same defect fleetd#164 is about.
String reason = "member " + target + " ended on a backend error: " + backendError
+ "\n--- pane tail ---\n" + tail;
if (rendezvous.resolveFailure(waiter, reason)) {
inFlight.remove(target, turn);
log.warn("failing send to {} via turn-stall fallback: {}", target, reason);
// fleetd#201 Unit 1: only on the resolution that actually won the race — a late
// duplicate must never double-count one backend failure.
backendErrorSink.onBackendError(target, backendError, reason);
}
fail(target, turn, "member " + target + " ended on a backend error: " + backendError
+ "\n--- pane tail ---\n" + tail);
return;
}
String completion = clipped ? tail + "\n" + CLIPPED_PANE_TAIL_MARKER : tail;
@@ -405,20 +353,14 @@ public final class CompletionResolver implements TurnListener {
}
return true;
}
String backendError = firstMatchingLine(raw, backendErrorPatternOrFallback(target));
String backendError = firstMatchingLine(raw, BACKEND_ERROR);
if (backendError != null) {
// Carry the pane, not just the matched line — the same fleetd#164 rule the normal path
// above applies. Here it matters more, not less: the trimmed block was empty, so the raw
// scrape is the ONLY copy of whatever the member managed to say. Clipped to the same cap
// the normal path uses, since a raw screen has no boundary trimming to bound it.
String reason = "member " + target + " ended on a backend error: " + backendError
+ "\n--- pane tail ---\n" + clip(raw);
if (rendezvous.resolveFailure(waiter, reason)) {
inFlight.remove(target, turn);
log.warn("failing send to {} via turn-stall fallback from the raw scrape: {}", target, reason);
// fleetd#201 Unit 1: only on the resolution that actually won the race.
backendErrorSink.onBackendError(target, backendError, reason);
}
fail(target, turn, "member " + target + " ended on a backend error: " + backendError
+ "\n--- pane tail ---\n" + clip(raw));
return true;
}
return false;
@@ -460,53 +402,22 @@ public final class CompletionResolver implements TurnListener {
}
/**
* fleetd#164 (floor) / fleetd#201 (classification): fail a turn that settled inside
* {@link #MIN_TURN_NANOS} — a crash signature (e.g. a backend HTTP 400 before the worker did
* anything) that a bare {@code BUSY -> DONE} transition cannot be told apart from a genuinely
* fast completion. Runs the same backend-error classification the normal and raw-scrape paths
* apply, against whatever is on screen right now: a match is a typed failure that notifies
* {@link #backendErrorSink} (only on the resolution that wins the race); a non-match stays the
* original generic too-fast failure, naming the member and both timings, with whatever the pane
* shows appended so the caller sees the cause, not just "it failed".
* fleetd#164: the failure reason for a turn that settled inside {@link #MIN_TURN_NANOS} — names
* the member and both timings, and appends whatever the pane shows (usually the backend's own
* error) so the caller sees the cause, not just "it failed".
*/
private void failTooFast(String target, InFlight turn, CompletableFuture<Rendezvous.Resolution> waiter,
long elapsedNanos) {
private String tooFastReason(String target, long elapsedNanos) {
String scrape;
try {
scrape = agents.read(target, SCRAPE_SOURCE);
scrape = clip(agents.read(target, SCRAPE_SOURCE));
} catch (RuntimeException e) {
scrape = "";
}
String clippedScrape = clip(scrape);
String baseReason = String.format(
String reason = String.format(
"member %s went BUSY -> DONE in %dms (floor %dms) — too fast to be real work, most "
+ "likely a backend error before any work started",
target, elapsedNanos / 1_000_000, MIN_TURN_NANOS / 1_000_000);
String backendError = firstMatchingLine(scrape, backendErrorPatternOrFallback(target));
if (backendError != null) {
String reason = baseReason + ": " + clippedScrape;
if (rendezvous.resolveFailure(waiter, reason)) {
inFlight.remove(target, turn);
log.warn("failing send to {} via turn-stall fallback: {}", target, reason);
// fleetd#201 Unit 1: only on the resolution that actually won the race.
backendErrorSink.onBackendError(target, backendError, reason);
}
return;
}
fail(target, turn, clippedScrape.isBlank() ? baseReason : baseReason + ": " + clippedScrape);
}
/**
* fleetd#201: the pattern to classify a backend error against for {@code target} — its
* profile's configured {@link BackendErrorPatternLookup} entry when there is one, else the
* built-in {@link #BACKEND_ERROR} compatibility pattern. A target with no configured pattern is
* never "off": it always falls back to this narrow default, exactly like the classification
* behaved before fleetd#201 (Unit 5 reports a target relying on this fallback separately, as
* legacy-default coverage rather than full per-profile coverage).
*/
private Pattern backendErrorPatternOrFallback(String target) {
Pattern configured = backendErrorPatterns.patternFor(target);
return configured != null ? configured : BACKEND_ERROR;
return scrape.isBlank() ? reason : reason + ": " + scrape;
}
/**
@@ -46,7 +46,8 @@ public record MemberSession(
String worktree,
String branch,
CharterReceipt charterReceipt,
String agentSessionId) {
String agentSessionId,
String failureReason) {
/** One-shot worker lifecycle states. */
public enum State {
@@ -54,6 +55,7 @@ public record MemberSession(
READY,
BUSY,
DONE,
BACKEND_ERROR,
FAILED,
RELEASED
}
@@ -69,25 +71,35 @@ public record MemberSession(
long lastActivityAtNanos, int turnCount, State state,
String worktree, String branch) {
this(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
lastActivityAtNanos, turnCount, state, worktree, branch, null, null);
lastActivityAtNanos, turnCount, state, worktree, branch, null, null, null);
}
/** Backward-compatible shape without a backend failure reason. */
public MemberSession(String paneId, String terminalId, String profile, MemberRole role,
String cwd, String ownerTerminal, long spawnedAtNanos,
long lastActivityAtNanos, int turnCount, State state,
String worktree, String branch, CharterReceipt charterReceipt,
String agentSessionId) {
this(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
lastActivityAtNanos, turnCount, state, worktree, branch, charterReceipt, agentSessionId, null);
}
/** Return a copy of this session in {@code state}. */
public MemberSession withState(State state) {
return new MemberSession(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
lastActivityAtNanos, turnCount, state, worktree, branch, charterReceipt, agentSessionId);
lastActivityAtNanos, turnCount, state, worktree, branch, charterReceipt, agentSessionId, failureReason);
}
/** Return a copy with {@code lastActivityAtNanos} updated to {@code nowNanos}. */
public MemberSession withActivity(long nowNanos) {
return new MemberSession(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
nowNanos, turnCount, state, worktree, branch, charterReceipt, agentSessionId);
nowNanos, turnCount, state, worktree, branch, charterReceipt, agentSessionId, failureReason);
}
/** Return a copy with the turn count incremented and activity timestamped at {@code nowNanos}. */
public MemberSession bumpTurn(long nowNanos) {
return new MemberSession(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
nowNanos, turnCount + 1, state, worktree, branch, charterReceipt, agentSessionId);
nowNanos, turnCount + 1, state, worktree, branch, charterReceipt, agentSessionId, failureReason);
}
/**
@@ -98,6 +110,12 @@ public record MemberSession(
*/
public MemberSession withAgentSessionId(String agentSessionId) {
return new MemberSession(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
lastActivityAtNanos, turnCount, state, worktree, branch, charterReceipt, agentSessionId);
lastActivityAtNanos, turnCount, state, worktree, branch, charterReceipt, agentSessionId, failureReason);
}
/** Return a copy with the durable backend failure detail. */
public MemberSession withFailureReason(String failureReason) {
return new MemberSession(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
lastActivityAtNanos, turnCount, state, worktree, branch, charterReceipt, agentSessionId, failureReason);
}
}
@@ -658,6 +658,9 @@ public final class SessionManager implements TurnListener {
m.put("profile", session.profile());
m.put("role", session.role() == null ? "dev" : session.role().wireName());
m.put("state", session.state().name().toLowerCase());
if (session.failureReason() != null) {
m.put("failureReason", session.failureReason());
}
if (session.worktree() != null) {
m.put("worktree", session.worktree());
}
@@ -718,6 +721,36 @@ public final class SessionManager implements TurnListener {
completeTurn(target, false);
}
/**
* Record that a member backend failed after a delegated ticket expired. This CAS loop accepts
* either side of the completion race: {@code BUSY -> BACKEND_ERROR} or
* {@code DONE -> BACKEND_ERROR}. A released or unknown session is never recreated.
*
* @return {@code true} when the error is recorded on a live session
*/
public boolean onBackendError(String target, String reason) {
while (true) {
MemberSession current = findByTerminal(target);
if (current == null) {
log.warn("backend error for unknown member terminal={}", target);
return false;
}
if (current.state() == MemberSession.State.RELEASED) {
return false;
}
if (current.state() == MemberSession.State.BACKEND_ERROR) {
return true;
}
MemberSession updated = current.withState(MemberSession.State.BACKEND_ERROR)
.withFailureReason(reason).withActivity(nowNanos.getAsLong());
if (replace(current, updated)) {
log.warn("member terminal={} pane={} transitioned {} -> BACKEND_ERROR: {}",
target, current.paneId(), current.state(), reason);
return true;
}
}
}
@Override
public boolean hasPostTurnAction(String target) {
if (!clearAfterTurn) return false;
@@ -736,10 +769,9 @@ public final class SessionManager implements TurnListener {
if (current == null || current.state() != MemberSession.State.BUSY) return false;
long now = nowNanos.getAsLong();
MemberSession updated = current.withState(MemberSession.State.DONE).withActivity(now);
if (replace(current, updated)) {
log.debug("session transitioned terminal={} pane={} BUSY -> DONE turn={}",
target, current.paneId(), updated.turnCount());
}
if (!replace(current, updated)) return false;
log.debug("session transitioned terminal={} pane={} BUSY -> DONE turn={}",
target, current.paneId(), updated.turnCount());
if (contextCap > 0 && updated.turnCount() >= contextCap) {
release(current.paneId());
return false;
@@ -4,11 +4,8 @@ import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.LoggerContext;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.msg.Rendezvous;
import dev.ltms.fleet.msg.TestTurnTokens;
import dev.ltms.fleet.msg.TurnToken;
@@ -762,202 +759,4 @@ class CompletionResolverTest {
assertEquals("partial (configured: [terra]; not configured: [gx10])",
CompletionResolver.coverage(Set.of("terra", "gx10"), Set.of("terra")));
}
// --- fleetd#201 Unit 1: target-keyed backend-error pattern + typed sink ----------------------
@Test
void aConfiguredBackendErrorPatternClassifiesAMatchAsAFailureAndNotifiesTheSinkOnce() {
String block = "⏺ 503 Service Unavailable: upstream credential rejected\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
BackendErrorPatternLookup patterns = target -> Pattern.compile("(?i)503 Service Unavailable");
java.util.List<String> notified = new java.util.ArrayList<>();
BackendErrorSink sink = (target, matchedLine, reason) -> notified.add(target + ": " + matchedLine);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none(), patterns, sink);
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
assertTrue(waiter.isDone(), "a configured-pattern match still resolves the blocked send");
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind());
assertEquals(1, notified.size(), "the sink fires exactly once for the winning classification");
assertEquals("term_a: 503 Service Unavailable: upstream credential rejected", notified.get(0));
}
@Test
void aLosingBackendErrorClassificationNeverNotifiesTheSink() {
// The waiter was already resolved (e.g. by the worker's own fleet_reply) before this scrape
// landed — resolveFailure loses the race and must return false, so the sink must not fire.
String block = "⏺ 503 Service Unavailable: upstream credential rejected\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
BackendErrorPatternLookup patterns = target -> Pattern.compile("(?i)503 Service Unavailable");
java.util.List<String> notified = new java.util.ArrayList<>();
BackendErrorSink sink = (target, matchedLine, reason) -> notified.add(target + ": " + matchedLine);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none(), patterns, sink);
var waiter = rendezvous.open("term_a");
var turn = new CompletionResolver.InFlight(waiter, null);
assertTrue(rendezvous.resolveCompletion(waiter, "already replied")); // fleet_reply won first
resolver.resolve("term_a", turn);
assertTrue(notified.isEmpty(), "a classification that loses the race must not fire the sink");
assertEquals("already replied", waiter.getNow(null).text(), "the earlier resolution stands untouched");
}
@Test
void aBackendErrorClassificationThatLosesARealRaceDuringTheScrapeNeverNotifiesTheSink() {
// The test above pre-resolves the waiter BEFORE calling resolve(), so it is caught by
// resolve()'s own top-of-method isDone() guard before ever reaching the win-gated
// resolveFailure call this unit adds — it proves the outcome, not the specific gate. This
// test forces the actual race window fleetd#201 must be safe in: a genuine fleet_reply lands
// WHILE the herdr scrape for this turn is in flight, i.e. AFTER resolve() has already passed
// the top guard (waiter not done yet) and committed to reading the pane, but BEFORE it
// reaches the "if (rendezvous.resolveFailure(...))" line below. The side effect is attached
// to the one real I/O step resolve() performs between those two points: the herdr
// "agent.read" call itself.
Rendezvous rendezvous = new Rendezvous();
var waiter = rendezvous.open("term_a");
ObjectMapper mapper = new ObjectMapper();
HerdrClient racingDuringScrape = new HerdrClient() {
@Override
public JsonNode call(String method, Object params) {
if ("agent.read".equals(method)) {
// The real fleet_reply that wins the race, landing mid-scrape.
assertTrue(rendezvous.resolve("term_a", "a genuine fleet_reply landed first"),
"the racing reply must land while the waiter is still open");
}
try {
return mapper.readTree(mapper.writeValueAsString(java.util.Map.of(
"type", "agent_read",
"read", java.util.Map.of("text",
"⏺ 503 Service Unavailable: upstream credential rejected\n❯ "))));
} catch (Exception e) {
throw new RuntimeException(e);
}
}
@Override
public void close() {
}
};
BackendErrorPatternLookup patterns = target -> Pattern.compile("(?i)503 Service Unavailable");
java.util.List<String> notified = new java.util.ArrayList<>();
BackendErrorSink sink = (target, matchedLine, reason) -> notified.add(target + ": " + matchedLine);
CompletionResolver resolver = new CompletionResolver(new AgentControl(racingDuringScrape), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none(), patterns, sink);
var turn = new CompletionResolver.InFlight(waiter, null); // not done yet — passes the top guard
resolver.resolve("term_a", turn);
assertTrue(notified.isEmpty(),
"a classification that loses a real race during the scrape must not fire the sink");
assertEquals(Rendezvous.Kind.REPLY, waiter.getNow(null).kind(), "the genuine fleet_reply stands");
assertEquals("a genuine fleet_reply landed first", waiter.getNow(null).text());
}
@Test
void exhaustionKeepsWinningOverBackendErrorEvenWhenBothPatternsMatchTheSameLine() {
// A line matching both a configured exhausted pattern AND a configured backend-error pattern
// must classify BACKEND_EXHAUSTED and call only the ExhaustionSink — the ordering fleetd#578
// already relies on must not change.
String block = "⏺ The usage limit has been reached. Try again later.\n❯ ";
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
ExhaustedPatternLookup exhausted = target -> Pattern.compile("usage limit has been reached");
BackendErrorPatternLookup backendErrors = target -> Pattern.compile("(?i)usage limit");
java.util.List<String> exhaustedNotified = new java.util.ArrayList<>();
java.util.List<String> backendErrorNotified = new java.util.ArrayList<>();
ExhaustionSink exhaustionSink = (target, reason) -> exhaustedNotified.add(target);
BackendErrorSink backendErrorSink = (target, matchedLine, reason) -> backendErrorNotified.add(target);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
exhausted, exhaustionSink, backendErrors, backendErrorSink);
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
assertEquals(Rendezvous.Kind.BACKEND_EXHAUSTED, waiter.getNow(null).kind(),
"exhaustion must keep winning over backend error");
assertEquals(1, exhaustedNotified.size(), "the exhaustion sink fires");
assertTrue(backendErrorNotified.isEmpty(), "the backend-error sink must never fire for this turn");
}
@Test
void aConfiguredPatternAlsoClassifiesTheRawScrapeFallbackAndNotifiesTheSink() {
// No ⏺ marker and leading TUI chrome ⇒ lastAssistantBlock() yields "", so classification must
// fall back to the raw scrape (fleetd#211) — and it must use the configured pattern too.
String block = """
╭──────────────────────────────────────╮
503 Service Unavailable: upstream credential rejected
""";
FakeHerdr herdr = new FakeHerdr().readText(block);
Rendezvous rendezvous = new Rendezvous();
BackendErrorPatternLookup patterns = target -> Pattern.compile("(?i)503 Service Unavailable");
java.util.List<String> notified = new java.util.ArrayList<>();
BackendErrorSink sink = (target, matchedLine, reason) -> notified.add(target + ": " + matchedLine);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none(), patterns, sink);
var waiter = rendezvous.open("term_a");
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
assertTrue(waiter.isDone(), "a raw-scrape configured-pattern match still resolves the blocked send");
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(),
"classified from the raw scrape even though the trimmed block was empty");
assertEquals(1, notified.size(), "the sink fires exactly once from the raw-scrape path");
assertTrue(notified.get(0).contains("503 Service Unavailable"), notified.get(0));
}
// --- fleetd#201 Unit 1: classification inside the fleetd#164 MIN_TURN_NANOS floor -------------
@Test
void aMatchingErrorInsideTheFloorIsTypedAndNotifiesTheSink() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ 503 Service Unavailable: upstream credential rejected\n❯ ");
Rendezvous rendezvous = new Rendezvous();
long[] clock = {10_000_000_000L};
BackendErrorPatternLookup patterns = target -> Pattern.compile("(?i)503 Service Unavailable");
java.util.List<String> notified = new java.util.ArrayList<>();
BackendErrorSink sink = (target, matchedLine, reason) -> notified.add(target + ": " + matchedLine);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none(), patterns, sink, () -> clock[0]);
var waiter = rendezvous.open("term_a");
var turn = new CompletionResolver.InFlight(waiter, null, clock[0]); // delivered "now"
clock[0] += CompletionResolver.MIN_TURN_NANOS - 1; // inside the floor
resolver.resolve("term_a", turn);
assertTrue(waiter.isDone());
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind());
assertEquals(1, notified.size(), "a matching error inside the floor notifies the sink");
assertTrue(notified.get(0).contains("503 Service Unavailable"), notified.get(0));
}
@Test
void aNonMatchInsideTheFloorStaysGenericAndNeverNotifiesTheSink() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ still starting up\n❯ ");
Rendezvous rendezvous = new Rendezvous();
long[] clock = {10_000_000_000L};
BackendErrorPatternLookup patterns = target -> Pattern.compile("(?i)503 Service Unavailable");
java.util.List<String> notified = new java.util.ArrayList<>();
BackendErrorSink sink = (target, matchedLine, reason) -> notified.add(target + ": " + matchedLine);
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none(), patterns, sink, () -> clock[0]);
var waiter = rendezvous.open("term_a");
var turn = new CompletionResolver.InFlight(waiter, null, clock[0]); // delivered "now"
clock[0] += CompletionResolver.MIN_TURN_NANOS - 1; // inside the floor
resolver.resolve("term_a", turn);
assertTrue(waiter.isDone());
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(),
"still a failure — the floor itself, not the pattern, is why");
assertTrue(waiter.getNow(null).text().contains("too fast to be real work"),
"a non-match inside the floor stays the generic too-fast reason: " + waiter.getNow(null).text());
assertTrue(notified.isEmpty(), "a non-match must never notify the typed sink");
}
}
@@ -300,6 +300,151 @@ class SessionManagerTest {
assertTrue(sessions.roster().contains(updated), "FAILED is still in acquired-minus-released roster");
}
@Test
void backendErrorWinsBeforeNormalCompletionFromBusy() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
sessions.asPresence().markPresent(session.terminalId());
sessions.onDelivered(session.terminalId(), TestTurnTokens.inert(session.terminalId()));
assertTrue(sessions.onBackendError(session.terminalId(), "backend exited"));
sessions.onTurnComplete(session.terminalId());
MemberSession updated = sessions.get(session.paneId()).orElseThrow();
assertEquals(MemberSession.State.BACKEND_ERROR, updated.state(),
"normal completion must not restore DONE after a backend error from BUSY");
assertEquals("backend exited", updated.failureReason());
}
@Test
void backendErrorWinsAfterNormalCompletionFromDone() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
sessions.asPresence().markPresent(session.terminalId());
sessions.onDelivered(session.terminalId(), TestTurnTokens.inert(session.terminalId()));
sessions.onTurnComplete(session.terminalId());
assertTrue(sessions.onBackendError(session.terminalId(), "backend exited"));
MemberSession updated = sessions.get(session.paneId()).orElseThrow();
assertEquals(MemberSession.State.BACKEND_ERROR, updated.state(),
"backend error must replace DONE when normal completion won first");
assertEquals("backend exited", updated.failureReason());
}
@Test
void backendErrorMemberCannotReceiveAnotherDelivery() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
sessions.asPresence().markPresent(session.terminalId());
sessions.onDelivered(session.terminalId(), TestTurnTokens.inert(session.terminalId()));
assertTrue(sessions.onBackendError(session.terminalId(), "backend exited"));
sessions.onDelivered(session.terminalId(), TestTurnTokens.inert(session.terminalId()));
assertEquals(MemberSession.State.BACKEND_ERROR, sessions.get(session.paneId()).orElseThrow().state(),
"onDelivered must refuse a backend-error member");
}
@Test
void losingCompletionDoesNotReleaseOrClearABackendErrorMember() {
assertLosingCompletionLeavesBackendErrorIntact(1, false,
"a stale completion must not release a backend-error member at the context cap");
assertLosingCompletionLeavesBackendErrorIntact(0, true,
"a stale completion must not clear the context of a backend-error member");
}
private void assertLosingCompletionLeavesBackendErrorIntact(int contextCap, boolean clearAfterTurn,
String protectedAction) {
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile cfg = new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
List.of("ccs", "ltms-local"), "tab", "fleetd-workers",
"worker: {profile} #{n}", null, null, null);
ClaudeCodeLauncher delegate = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
ClearContextSpyLauncher launcher = new ClearContextSpyLauncher(delegate);
java.util.concurrent.atomic.AtomicReference<SessionManager> manager = new java.util.concurrent.atomic.AtomicReference<>();
java.util.concurrent.atomic.AtomicReference<String> terminal = new java.util.concurrent.atomic.AtomicReference<>();
java.util.concurrent.atomic.AtomicBoolean armed = new java.util.concurrent.atomic.AtomicBoolean();
SessionManager sessions = new SessionManager(launcher, new RecordingWorktrees(), () -> {
if (armed.compareAndSet(true, false)) {
manager.get().onBackendError(terminal.get(), "backend exited during completion");
}
return 1;
}, contextCap, clearAfterTurn);
manager.set(sessions);
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
terminal.set(session.terminalId());
sessions.asPresence().markPresent(session.terminalId());
sessions.onDelivered(session.terminalId(), TestTurnTokens.inert(session.terminalId()));
armed.set(true);
assertFalse(sessions.onTurnCompleteWithPostAction(session.terminalId()),
"the completion CAS loses after the injected backend error");
assertTrue(sessions.get(session.paneId()).isPresent(), protectedAction);
MemberSession updated = sessions.get(session.paneId()).orElseThrow();
assertEquals(MemberSession.State.BACKEND_ERROR, updated.state());
assertEquals("backend exited during completion", updated.failureReason());
assertFalse(herdr.called("pane.close"), protectedAction);
assertEquals(0, launcher.clearContextCalls(), protectedAction);
}
@Test
void backendErrorForUnknownTargetIsWarnedAndDoesNotCreateASession() {
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);
try {
SessionManager sessions = sessionManager(new FakeHerdr());
assertFalse(sessions.onBackendError("term_missing", "backend exited"));
assertTrue(sessions.roster().isEmpty(), "unknown target must not create a session");
assertTrue(appender.list.stream().anyMatch(e -> e.getLevel().equals(Level.WARN)
&& e.getFormattedMessage().contains("term_missing")),
"unknown target is logged at WARN");
} finally {
sessionLog.detachAppender(appender);
}
}
@Test
void backendErrorForReleasedTargetDoesNotRecreateTheSession() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
sessions.release(session.paneId());
assertFalse(sessions.onBackendError(session.terminalId(), "backend exited"));
assertTrue(sessions.get(session.paneId()).isEmpty(), "released session must stay absent");
}
@Test
void rosterViewIncludesFailureReasonOnlyForBackendError() {
MemberSession failed = new MemberSession("pane-error", "term-error", "prof", MemberRole.DEV,
"/cwd", null, 0, 0, 1, MemberSession.State.BACKEND_ERROR, null, null,
null, null, "backend exited");
MemberSession ordinary = new MemberSession("pane-ready", "term-ready", "prof", MemberRole.DEV,
"/cwd", null, 0, 0, 0, MemberSession.State.READY, null, null);
Map<String, Object> failedView = SessionManager.rosterView(failed, null);
Map<String, Object> ordinaryView = SessionManager.rosterView(ordinary, null);
assertEquals("backend_error", failedView.get("state"));
assertEquals("backend exited", failedView.get("failureReason"));
assertFalse(ordinaryView.containsKey("failureReason"),
"ordinary rows must not gain a blank failureReason field");
}
@Test
void onTurnFailedIsLoggedAtWarnWithThePriorState() {
// CB-564: this transition used to be a bare DEBUG "session marked failed" — a symptom with no
@@ -1004,6 +1149,76 @@ class SessionManagerTest {
}
}
/** Delegates all launcher work while recording context clears. */
private static final class ClearContextSpyLauncher implements PeerLauncher {
private final PeerLauncher delegate;
private int clearContextCalls;
ClearContextSpyLauncher(PeerLauncher delegate) {
this.delegate = delegate;
}
@Override
public Set<Capability> capabilities() {
return delegate.capabilities();
}
@Override
public Set<Capability> capabilitiesFor(String profileName) {
return delegate.capabilitiesFor(profileName);
}
@Override
public PeerHandle spawn(SpawnRequest req) {
return delegate.spawn(req);
}
@Override
public Set<String> profiles() {
return delegate.profiles();
}
@Override
public String defaultProfile() {
return delegate.defaultProfile();
}
@Override
public String effectiveCwd(SpawnRequest req) {
return delegate.effectiveCwd(req);
}
@Override
public List<String> parityOverlay(String profileName) {
return delegate.parityOverlay(profileName);
}
@Override
public List<?> list() {
return delegate.list();
}
@Override
public int reapOrphanWorkers() {
return delegate.reapOrphanWorkers();
}
@Override
public void stop(String id) {
delegate.stop(id);
}
@Override
public boolean clearContext(String id) {
clearContextCalls++;
return delegate.clearContext(id);
}
int clearContextCalls() {
return clearContextCalls;
}
}
@Test
void acquireRecordsTheAgentSessionIdFromTheHandle() {
FakeHerdr herdr = new FakeHerdr();