fleetd #201/#227 unit 2: credential outage policy

Adds BackendOutagePolicy — a credential-keyed state machine on an injected monotonic
clock. Two classified backend errors from two DISTINCT targets on one credential inside
60 seconds mint one incident and start a 60-second cool-off. Errors during cool-off
neither extend it nor mint another; expiry clears evidence, so two fresh errors rearm.

Correlated on credentialId, never on profile name or error text. Deliberately not
BackendQuarantine: that restarts a 1800-second cooldown per exhaustion, and its name
would make every refusal say 'backend exhausted', which is a different condition.

Threshold counts distinct targets rather than raw events (lead decision): the classifier
is a heuristic and a valid member report can quote an 'API Error:' line, so one member
repeating that line must not remove a healthy credential's capacity. A real outage hits
every member on the credential, so true detection is unaffected.

Verified by the lead: 1135 tests green; reverting evidenceCount() to reasons.size()
turns two BackendOutagePolicyTest cases red with 0 compile errors.
This commit is contained in:
Dai Ha
2026-09-03 10:49:10 +07:00
2 changed files with 466 additions and 0 deletions
@@ -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;
}
}
}
@@ -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);
}
}