diff --git a/fleetd/src/main/java/dev/ltms/fleet/inject/BackendErrorPatternLookup.java b/fleetd/src/main/java/dev/ltms/fleet/inject/BackendErrorPatternLookup.java new file mode 100644 index 0000000..ae5a363 --- /dev/null +++ b/fleetd/src/main/java/dev/ltms/fleet/inject/BackendErrorPatternLookup.java @@ -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. + * + *
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; + } +} diff --git a/fleetd/src/main/java/dev/ltms/fleet/inject/BackendErrorSink.java b/fleetd/src/main/java/dev/ltms/fleet/inject/BackendErrorSink.java new file mode 100644 index 0000000..3370afd --- /dev/null +++ b/fleetd/src/main/java/dev/ltms/fleet/inject/BackendErrorSink.java @@ -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. + * + *
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) -> { }; + } +} diff --git a/fleetd/src/main/java/dev/ltms/fleet/inject/CompletionResolver.java b/fleetd/src/main/java/dev/ltms/fleet/inject/CompletionResolver.java index 20ca1b2..f289783 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/inject/CompletionResolver.java +++ b/fleetd/src/main/java/dev/ltms/fleet/inject/CompletionResolver.java @@ -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. + * + *
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