diff --git a/fleetd/src/main/java/dev/ltms/fleet/placement/BackendOutagePolicy.java b/fleetd/src/main/java/dev/ltms/fleet/placement/BackendOutagePolicy.java
new file mode 100644
index 0000000..12cc865
--- /dev/null
+++ b/fleetd/src/main/java/dev/ltms/fleet/placement/BackendOutagePolicy.java
@@ -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.
+ *
+ *
Two classified errors on the same credential — never the profile name, never the error
+ * text, see {@code FleetConfig.Profile#effectiveCredentialId()} like {@link BackendQuarantine} —
+ * from two distinct targets 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.
+ *
+ *
Correlation and cool-off live in one class 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.
+ *
+ *
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 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()}).
+ *
+ * Returns a populated {@link Incident} only at the exact moment a second distinct target
+ * 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 record(String credentialId, String target, String reason) {
+ Objects.requireNonNull(credentialId, "credentialId");
+ Objects.requireNonNull(target, "target");
+ Objects.requireNonNull(reason, "reason");
+ long now = nowNanos.getAsLong();
+
+ AtomicReference 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.
+ *
+ * {@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 targets,
+ int evidenceCount,
+ long windowSeconds,
+ long remainingCoolOffSeconds,
+ List 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 targets;
+ final List reasons;
+ final long coolOffUntilNanos; // 0 means "not cooling off"
+
+ private CredentialState(long firstEventNanos, Set targets, List 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 targets = new LinkedHashSet<>();
+ targets.add(target);
+ List reasons = new ArrayList<>();
+ reasons.add(reason);
+ return new CredentialState(nowNanos, targets, reasons, 0L);
+ }
+
+ CredentialState withAdditionalEvidence(String target, String reason) {
+ Set newTargets = new LinkedHashSet<>(targets);
+ newTargets.add(target);
+ List 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;
+ }
+ }
+}
diff --git a/fleetd/src/test/java/dev/ltms/fleet/placement/BackendOutagePolicyTest.java b/fleetd/src/test/java/dev/ltms/fleet/placement/BackendOutagePolicyTest.java
new file mode 100644
index 0000000..1d3e39f
--- /dev/null
+++ b/fleetd/src/test/java/dev/ltms/fleet/placement/BackendOutagePolicyTest.java
@@ -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 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 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 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 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 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 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 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 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 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 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 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 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);
+ }
+}