Compare commits

..

2 Commits

Author SHA1 Message Date
Dai Ha cf54aed451 CB-201 unit 2 review fix: threshold counts distinct targets, not raw events
CI / contract (pull_request) Successful in 58s
CI / build (pull_request) Successful in 1m46s
Two errors from the same target inside the window must never trip the
outage threshold on their own (a valid member report can legitimately
quote an "API Error:" line twice) — only two DIFFERENT targets on the
same credential do. Change evidenceCount() to targets.size() instead
of reasons.size(); reasons() still keeps every event, including
same-target repeats, so it can be longer than evidenceCount(). A real
outage still hits every target on the credential, so this loses no
true-positive coverage while cutting a real false-positive path.
2026-09-03 10:42:37 +07:00
Dai Ha 826e0aeb2a CB-201 unit 2: credential-keyed backend outage policy
CI / build (pull_request) Successful in 1m11s
CI / contract (pull_request) Successful in 1m24s
Add BackendOutagePolicy: two classified backend errors on the same
credentialId within a 60s window mint one Incident and start a 60s
cool-off for that credential, one atomic ConcurrentHashMap.compute()
per credentialId so a concurrent second and third event can never
both cross the threshold. Errors during cool-off are ignored outright
(no extension, no incident); once cool-off elapses the next error
clears old evidence, requiring two fresh errors to rearm. This is a
new class, deliberately not BackendQuarantine (wrong store, wrong
1800s duration, misleading "exhausted" semantics for a 60s transient
fault). Knows nothing about panes, profiles, sessions, launchers, or
leads — takes events in, returns incidents out.
2026-09-03 10:30:40 +07:00
6 changed files with 486 additions and 380 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;
}
/**
@@ -0,0 +1,200 @@
package dev.ltms.fleet.placement;
import java.util.ArrayList;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Objects;
import java.util.Optional;
import java.util.OptionalLong;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.LongSupplier;
/**
* Fleetd #201 / #227: correlates classified backend errors (credential outage or provider 5xx,
* produced elsewhere by the classifier that turns raw pane text into a typed event — this class
* knows nothing about pane text, profiles, sessions, launchers, or leads) and decides, purely from
* counts and timing, when a credential's backend is out.
*
* <p>Two classified errors on the <b>same credential</b> — never the profile name, never the error
* text, see {@code FleetConfig.Profile#effectiveCredentialId()} like {@link BackendQuarantine} —
* <b>from two distinct targets</b> inside a 60-second window is treated as an outage: {@link
* #record} then returns an {@link Incident} and starts a 60-second cool-off for that credential. The
* threshold counts distinct targets, not raw events, on purpose: the classifier is a heuristic and a
* valid member report can quote an {@code API Error:} line, so one target repeating that line twice
* must never remove fleet capacity by itself. Two different targets independently producing a
* classified error is far less likely to be a coincidence, and a real backend outage hits every
* target on that credential anyway, so this loses nothing against the case being protected against.
* This is deliberately a different store from {@link BackendQuarantine}: that one holds a
* 1800-second exhaustion cooldown for a spent credential, and reusing it here would both use the
* wrong duration and report "backend exhausted" for what is a short transient fault.
*
* <p>Correlation and cool-off live in <b>one class</b> so that crossing the threshold and setting
* the cool-off deadline happen as a single atomic update: every {@link #record} call goes through
* {@link ConcurrentHashMap#compute}, which serializes on the credential's map bucket, so a
* concurrent second and third event on the same credential can never both observe "about to cross"
* and both mint an incident.
*
* <p>The clock is injected ({@link LongSupplier}, conventionally {@code System::nanoTime} like
* {@link BackendQuarantine}), never read inline, so the window and cool-off are testable without a
* real sleep.
*/
public final class BackendOutagePolicy {
/** Distinct targets a credential needs a classified error from to declare an outage. */
public static final int THRESHOLD = 2;
/** How long a credential's evidence stays fresh, inclusive of both endpoints. */
public static final long WINDOW_NANOS = 60_000_000_000L;
/** How long a credential sits out once the threshold is crossed. */
public static final long COOLOFF_NANOS = 60_000_000_000L;
private final ConcurrentHashMap<String, CredentialState> states = new ConcurrentHashMap<>();
private final LongSupplier nowNanos;
private final AtomicLong incidentSequence = new AtomicLong();
public BackendOutagePolicy(LongSupplier nowNanos) {
this.nowNanos = Objects.requireNonNull(nowNanos, "nowNanos");
}
/**
* Records one classified backend error for {@code credentialId} against {@code target} (e.g. a
* session or worker id — this class never interprets it, only collects it for the incident's
* affected-targets list, and counts distinct targets toward the threshold) with {@code reason}
* (the classifier's free-text reason, kept per event for whoever renders the eventual notice —
* every event's reason is kept even when the same target repeats, so {@link Incident#reasons()}
* can be longer than {@link Incident#evidenceCount()}).
*
* <p>Returns a populated {@link Incident} only at the exact moment a <b>second distinct target</b>
* is seen for this credential inside the window — never before, and never again while the
* resulting cool-off is active. A repeat error from a target already counted does not advance the
* threshold. While a credential is cooling off, a fresh error is ignored outright: it neither
* extends the deadline nor produces another incident. Once the cool-off has elapsed, the next
* error clears the old evidence and starts a brand-new window — two fresh, distinct targets are
* required to rearm.
*/
public Optional<Incident> record(String credentialId, String target, String reason) {
Objects.requireNonNull(credentialId, "credentialId");
Objects.requireNonNull(target, "target");
Objects.requireNonNull(reason, "reason");
long now = nowNanos.getAsLong();
AtomicReference<Incident> minted = new AtomicReference<>();
states.compute(credentialId, (_, existing) -> {
CredentialState state = existing;
if (state != null && state.inCoolOff()) {
if (now < state.coolOffUntilNanos) {
return state; // still cooling off: no extension, no incident
}
state = null; // cool-off elapsed: evidence is cleared, rearm from scratch
}
if (state == null || now - state.firstEventNanos > WINDOW_NANOS) {
return CredentialState.first(now, target, reason);
}
CredentialState advanced = state.withAdditionalEvidence(target, reason);
if (advanced.evidenceCount() < THRESHOLD) {
return advanced;
}
long coolOffUntilNanos = now + COOLOFF_NANOS;
minted.set(new Incident(
"outage-" + credentialId + "-" + incidentSequence.incrementAndGet(),
credentialId,
Set.copyOf(advanced.targets),
advanced.evidenceCount(),
toSecondsRoundedUp(WINDOW_NANOS),
toSecondsRoundedUp(COOLOFF_NANOS),
List.copyOf(advanced.reasons)));
return advanced.enteringCoolOff(coolOffUntilNanos);
});
return Optional.ofNullable(minted.get());
}
/** Seconds left on {@code credentialId}'s cool-off, or empty when it is not cooling off. */
public OptionalLong remainingCoolOffSeconds(String credentialId) {
Objects.requireNonNull(credentialId, "credentialId");
CredentialState state = states.get(credentialId);
if (state == null || !state.inCoolOff()) {
return OptionalLong.empty();
}
long remaining = state.coolOffUntilNanos - nowNanos.getAsLong();
return remaining > 0 ? OptionalLong.of(toSecondsRoundedUp(remaining)) : OptionalLong.empty();
}
private static long toSecondsRoundedUp(long nanos) {
return (nanos + 999_999_999L) / 1_000_000_000L;
}
/**
* One credential crossing the outage threshold. {@code remainingCoolOffSeconds} is the cool-off
* length as observed at the moment of minting — this incident is only ever produced right as the
* cool-off starts, so it is always the full {@link #COOLOFF_NANOS} rounded up.
*
* <p>{@code evidenceCount} is {@code targets.size()} — the threshold is on distinct targets, not
* raw events — while {@code reasons} keeps every event's reason, including repeats from a target
* already counted. The two are deliberately different lengths: a single target hammering the same
* classified error never grows {@code evidenceCount} past 1, but each occurrence still lands in
* {@code reasons} for whoever renders the notice.
*/
public record Incident(
String id,
String credentialId,
Set<String> targets,
int evidenceCount,
long windowSeconds,
long remainingCoolOffSeconds,
List<String> reasons) {
}
/** Evidence accumulated for one credential since its window opened, or its active cool-off. */
private static final class CredentialState {
final long firstEventNanos;
final Set<String> targets;
final List<String> reasons;
final long coolOffUntilNanos; // 0 means "not cooling off"
private CredentialState(long firstEventNanos, Set<String> targets, List<String> reasons,
long coolOffUntilNanos) {
this.firstEventNanos = firstEventNanos;
this.targets = targets;
this.reasons = reasons;
this.coolOffUntilNanos = coolOffUntilNanos;
}
static CredentialState first(long nowNanos, String target, String reason) {
Set<String> targets = new LinkedHashSet<>();
targets.add(target);
List<String> reasons = new ArrayList<>();
reasons.add(reason);
return new CredentialState(nowNanos, targets, reasons, 0L);
}
CredentialState withAdditionalEvidence(String target, String reason) {
Set<String> newTargets = new LinkedHashSet<>(targets);
newTargets.add(target);
List<String> newReasons = new ArrayList<>(reasons);
newReasons.add(reason);
return new CredentialState(firstEventNanos, newTargets, newReasons, 0L);
}
CredentialState enteringCoolOff(long coolOffUntilNanos) {
return new CredentialState(firstEventNanos, targets, reasons, coolOffUntilNanos);
}
/** Distinct targets seen so far — the threshold counts this, never {@code reasons.size()}. */
int evidenceCount() {
return targets.size();
}
boolean inCoolOff() {
return coolOffUntilNanos > 0;
}
}
}
@@ -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");
}
}
@@ -0,0 +1,266 @@
package dev.ltms.fleet.placement;
import org.junit.jupiter.api.Test;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #201 / #227: the credential-keyed outage decision itself, isolated from the classifier
* that produces events (Unit 1) and from the lead notice / spawn wiring (Units 3 and 5). The clock
* is a plain {@link AtomicLong} of nanos so the window and cool-off are exercised without a real
* sleep, exactly like {@code BackendQuarantineTest}.
*/
class BackendOutagePolicyTest {
private static final long SECOND = 1_000_000_000L;
@Test
void oneErrorCreatesNoIncidentAndNoCoolOff() {
BackendOutagePolicy policy = new BackendOutagePolicy(() -> 0L);
Optional<BackendOutagePolicy.Incident> incident =
policy.record("shared-openai", "session-1", "API Error: 429");
assertTrue(incident.isEmpty());
assertTrue(policy.remainingCoolOffSeconds("shared-openai").isEmpty());
}
@Test
void twoErrorsWithinTheWindowCreateExactlyOneIncidentAndA60SecondCoolOff() {
AtomicLong now = new AtomicLong(0L);
BackendOutagePolicy policy = new BackendOutagePolicy(now::get);
assertTrue(policy.record("shared-openai", "session-1", "API Error: 429").isEmpty());
now.set(30 * SECOND);
Optional<BackendOutagePolicy.Incident> incident =
policy.record("shared-openai", "session-2", "API Error: 500");
assertTrue(incident.isPresent());
BackendOutagePolicy.Incident value = incident.get();
assertEquals("shared-openai", value.credentialId());
assertEquals(Set.of("session-1", "session-2"), value.targets());
assertEquals(2, value.evidenceCount());
assertEquals(60L, value.windowSeconds());
assertEquals(60L, value.remainingCoolOffSeconds());
assertEquals(OptionalLongOf60(), policy.remainingCoolOffSeconds("shared-openai"));
}
@Test
void repeatedErrorsFromTheSameTargetNeverCreateAnIncidentHoweverManyTimesTheyRepeat() {
// The classifier is a heuristic and a valid member report can quote an "API Error:" line, so
// one target repeating that line inside the window must never, by itself, cost the credential
// its capacity — only a SECOND, DISTINCT target crossing the threshold does.
AtomicLong now = new AtomicLong(0L);
BackendOutagePolicy policy = new BackendOutagePolicy(now::get);
for (int i = 0; i < 5; i++) {
now.set(i * 10L * SECOND);
Optional<BackendOutagePolicy.Incident> incident =
policy.record("shared-openai", "session-1", "API Error: 429");
assertTrue(incident.isEmpty(),
"the same target repeating must never create an incident on its own (repeat #" + i + ")");
}
assertTrue(policy.remainingCoolOffSeconds("shared-openai").isEmpty());
}
@Test
void evidenceCountIsDistinctTargetsWhileReasonsKeepsEveryEvent() {
AtomicLong now = new AtomicLong(0L);
BackendOutagePolicy policy = new BackendOutagePolicy(now::get);
assertTrue(policy.record("shared-openai", "session-1", "API Error: 429").isEmpty());
now.set(10 * SECOND);
// a repeat from the already-counted target: still no incident, still only 1 distinct target
assertTrue(policy.record("shared-openai", "session-1", "API Error: 429 (again)").isEmpty());
now.set(20 * SECOND);
Optional<BackendOutagePolicy.Incident> incident =
policy.record("shared-openai", "session-2", "API Error: 500");
assertTrue(incident.isPresent());
BackendOutagePolicy.Incident value = incident.get();
assertEquals(Set.of("session-1", "session-2"), value.targets());
assertEquals(2, value.evidenceCount(), "evidenceCount is the distinct-target count");
assertEquals(3, value.reasons().size(),
"reasons keeps every event, including the same-target repeat that did not advance the count");
assertEquals(java.util.List.of("API Error: 429", "API Error: 429 (again)", "API Error: 500"),
value.reasons());
}
@Test
void twoErrorsMoreThanAMinuteApartCreateNoIncident() {
AtomicLong now = new AtomicLong(0L);
BackendOutagePolicy policy = new BackendOutagePolicy(now::get);
assertTrue(policy.record("shared-openai", "session-1", "API Error: 429").isEmpty());
now.set(60 * SECOND + 1);
Optional<BackendOutagePolicy.Incident> incident =
policy.record("shared-openai", "session-2", "API Error: 500");
assertTrue(incident.isEmpty(), "the first error's evidence must not survive past the window");
assertTrue(policy.remainingCoolOffSeconds("shared-openai").isEmpty());
}
@Test
void exactly60SecondsApartIsInsideTheWindow() {
AtomicLong now = new AtomicLong(0L);
BackendOutagePolicy policy = new BackendOutagePolicy(now::get);
assertTrue(policy.record("shared-openai", "session-1", "API Error: 429").isEmpty());
now.set(60 * SECOND); // exactly the boundary
Optional<BackendOutagePolicy.Incident> incident =
policy.record("shared-openai", "session-2", "API Error: 500");
assertTrue(incident.isPresent(), "60 seconds apart is inclusive of the window boundary");
}
@Test
void differentCredentialsNeverShareEvidence() {
AtomicLong now = new AtomicLong(0L);
BackendOutagePolicy policy = new BackendOutagePolicy(now::get);
assertTrue(policy.record("credential-a", "session-1", "API Error: 429").isEmpty());
now.set(10 * SECOND);
Optional<BackendOutagePolicy.Incident> incident =
policy.record("credential-b", "session-2", "API Error: 429");
assertTrue(incident.isEmpty(), "an unrelated credential must not be swept into the count");
}
@Test
void twoDifferentProfilesResolvingToTheSameCredentialShareEvidence() {
// The caller is responsible for resolving a profile to its credential id (see
// FleetConfig.Profile#effectiveCredentialId(), same as BackendQuarantine); this class only
// ever sees the resolved credentialId, so two different profile-originated events landing on
// the same credentialId must correlate.
AtomicLong now = new AtomicLong(0L);
BackendOutagePolicy policy = new BackendOutagePolicy(now::get);
assertTrue(policy.record("shared-openai", "session-from-profile-a", "API Error: 429").isEmpty());
now.set(5 * SECOND);
Optional<BackendOutagePolicy.Incident> incident =
policy.record("shared-openai", "session-from-profile-b", "API Error: 500");
assertTrue(incident.isPresent());
}
@Test
void errorsDuringCoolOffNeitherExtendTheDeadlineNorReturnAnotherIncident() {
AtomicLong now = new AtomicLong(0L);
BackendOutagePolicy policy = new BackendOutagePolicy(now::get);
policy.record("shared-openai", "session-1", "API Error: 429");
now.set(1 * SECOND);
policy.record("shared-openai", "session-2", "API Error: 500"); // crosses threshold, cool-off starts at t=1s, until t=61s
now.set(50 * SECOND);
Optional<BackendOutagePolicy.Incident> duringCoolOff =
policy.record("shared-openai", "session-3", "API Error: 500");
assertTrue(duringCoolOff.isEmpty(), "no second incident while cooling off");
assertEquals(OptionalLongOf(11), policy.remainingCoolOffSeconds("shared-openai"),
"the deadline must not have been pushed out by the error during cool-off");
}
@Test
void expiryClearsOldEvidenceSoTwoFreshErrorsAreRequiredToRearm() {
// Both errors land at the SAME instant (t=0) so the window boundary (60s from t=0) and the
// cool-off boundary (60s from the threshold crossing, also t=0) coincide exactly at t=60s.
// That is deliberate: with WINDOW_NANOS == COOLOFF_NANOS, checking "one fresh error" at any
// later instant would also be caught by the window naturally going stale, which would pass
// even if the explicit evidence-clear on cool-off expiry were missing. Checking exactly AT
// t=60s is the one instant where "now - firstEventNanos > WINDOW_NANOS" is false (60 is not
// > 60) but the cool-off has already elapsed ("now < coolOffUntilNanos" is false at 60 >= 60)
// — so only an explicit clear, not a side effect of the window check, keeps this from
// instantly re-mounting the old evidence and minting a bogus incident off one fresh error.
AtomicLong now = new AtomicLong(0L);
BackendOutagePolicy policy = new BackendOutagePolicy(now::get);
policy.record("shared-openai", "session-1", "API Error: 429");
policy.record("shared-openai", "session-2", "API Error: 500"); // cool-off until t=60s
now.set(60 * SECOND); // cool-off just elapsed, exactly at the window's own boundary too
Optional<BackendOutagePolicy.Incident> oneFreshError =
policy.record("shared-openai", "session-4", "API Error: 500");
assertTrue(oneFreshError.isEmpty(), "a single fresh error must not immediately re-trip");
now.set(61 * SECOND);
Optional<BackendOutagePolicy.Incident> secondFreshError =
policy.record("shared-openai", "session-5", "API Error: 500");
assertTrue(secondFreshError.isPresent(), "two fresh errors after expiry re-arm the policy");
}
@Test
void remainingSecondsRoundUp() {
AtomicLong now = new AtomicLong(0L);
BackendOutagePolicy policy = new BackendOutagePolicy(now::get);
policy.record("shared-openai", "session-1", "API Error: 429");
now.set(1); // 1 nanosecond later, still crosses the threshold
policy.record("shared-openai", "session-2", "API Error: 500");
// cool-off deadline is now(=1ns) + 60s exactly; a fraction-of-a-second-remaining check:
now.set(1 + 59 * SECOND + 1); // 59.000000001s of the cool-off has passed -> < 1s remains
assertEquals(OptionalLongOf(1), policy.remainingCoolOffSeconds("shared-openai"),
"a sub-second remainder must round up to a full second, matching BackendQuarantine");
}
@Test
void aConcurrentSecondAndThirdCallReturnOnlyOneIncidentBetweenThem() throws InterruptedException {
for (int attempt = 0; attempt < 25; attempt++) {
AtomicLong now = new AtomicLong(0L);
BackendOutagePolicy policy = new BackendOutagePolicy(now::get);
int threads = 8;
ExecutorService pool = Executors.newFixedThreadPool(threads);
CountDownLatch ready = new CountDownLatch(threads);
CountDownLatch go = new CountDownLatch(1);
AtomicInteger incidents = new AtomicInteger();
try {
for (int i = 0; i < threads; i++) {
int idx = i;
pool.submit(() -> {
ready.countDown();
try {
go.await();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
return;
}
Optional<BackendOutagePolicy.Incident> incident =
policy.record("shared-openai", "session-" + idx, "API Error: 500");
if (incident.isPresent()) {
incidents.incrementAndGet();
}
});
}
assertTrue(ready.await(5, TimeUnit.SECONDS));
go.countDown();
pool.shutdown();
assertTrue(pool.awaitTermination(5, TimeUnit.SECONDS));
assertEquals(1, incidents.get(),
"exactly one incident must be minted across every concurrent caller on attempt " + attempt);
} finally {
pool.shutdownNow();
}
}
}
private static java.util.OptionalLong OptionalLongOf60() {
return OptionalLongOf(60);
}
private static java.util.OptionalLong OptionalLongOf(long value) {
return java.util.OptionalLong.of(value);
}
}