fleetd#201 Unit 1: typed backend-error classification in CompletionResolver
Replace the hardcoded API-Error check with a target-keyed BackendErrorPatternLookup plus a BackendErrorSink, mirroring the existing ExhaustedPatternLookup/ExhaustionSink pair. Classifies in all three paths (normal block, #211 raw-scrape fallback, and the fleetd#164 MIN_TURN_NANOS floor). The sink fires only after Rendezvous.resolveFailure wins for the exact waiter. A target with no configured pattern still falls back to the narrow (?i)\bAPI Error\s*: compatibility pattern. Existing constructors keep compiling via BackendErrorPatternLookup.legacy() / BackendErrorSink.none() defaults. Public send result is unchanged (still a failed send) — the typed sink event is the internal seam Unit 5 will consume.
This commit is contained in:
@@ -0,0 +1,35 @@
|
||||
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;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,35 @@
|
||||
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,6 +89,8 @@ 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;
|
||||
|
||||
/**
|
||||
@@ -129,10 +131,48 @@ 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, System::nanoTime);
|
||||
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);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -145,11 +185,14 @@ public final class CompletionResolver implements TurnListener {
|
||||
* package.
|
||||
*/
|
||||
public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns,
|
||||
ExhaustionSink exhaustionSink, LongSupplier nowNanos) {
|
||||
ExhaustionSink exhaustionSink, BackendErrorPatternLookup backendErrorPatterns,
|
||||
BackendErrorSink backendErrorSink, 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");
|
||||
}
|
||||
|
||||
@@ -232,7 +275,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) {
|
||||
fail(target, turn, tooFastReason(target, elapsedNanos));
|
||||
failTooFast(target, turn, waiter, elapsedNanos);
|
||||
return;
|
||||
}
|
||||
String tail;
|
||||
@@ -302,18 +345,27 @@ public final class CompletionResolver implements TurnListener {
|
||||
}
|
||||
return;
|
||||
}
|
||||
// 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);
|
||||
// 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));
|
||||
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.
|
||||
fail(target, turn, "member " + target + " ended on a backend error: " + backendError
|
||||
+ "\n--- pane tail ---\n" + tail);
|
||||
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);
|
||||
}
|
||||
return;
|
||||
}
|
||||
String completion = clipped ? tail + "\n" + CLIPPED_PANE_TAIL_MARKER : tail;
|
||||
@@ -353,14 +405,20 @@ public final class CompletionResolver implements TurnListener {
|
||||
}
|
||||
return true;
|
||||
}
|
||||
String backendError = firstMatchingLine(raw, BACKEND_ERROR);
|
||||
String backendError = firstMatchingLine(raw, backendErrorPatternOrFallback(target));
|
||||
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.
|
||||
fail(target, turn, "member " + target + " ended on a backend error: " + backendError
|
||||
+ "\n--- pane tail ---\n" + clip(raw));
|
||||
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);
|
||||
}
|
||||
return true;
|
||||
}
|
||||
return false;
|
||||
@@ -402,22 +460,53 @@ public final class CompletionResolver implements TurnListener {
|
||||
}
|
||||
|
||||
/**
|
||||
* 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".
|
||||
* 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".
|
||||
*/
|
||||
private String tooFastReason(String target, long elapsedNanos) {
|
||||
private void failTooFast(String target, InFlight turn, CompletableFuture<Rendezvous.Resolution> waiter,
|
||||
long elapsedNanos) {
|
||||
String scrape;
|
||||
try {
|
||||
scrape = clip(agents.read(target, SCRAPE_SOURCE));
|
||||
scrape = agents.read(target, SCRAPE_SOURCE);
|
||||
} catch (RuntimeException e) {
|
||||
scrape = "";
|
||||
}
|
||||
String reason = String.format(
|
||||
String clippedScrape = clip(scrape);
|
||||
String baseReason = 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);
|
||||
return scrape.isBlank() ? reason : reason + ": " + scrape;
|
||||
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;
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -4,8 +4,11 @@ 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;
|
||||
@@ -759,4 +762,202 @@ 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");
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user