From 826e0aeb2a14555fcc70146c020e123648b6e544 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 3 Sep 2026 10:30:40 +0700 Subject: [PATCH] CB-201 unit 2: credential-keyed backend outage policy MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- .../fleet/placement/BackendOutagePolicy.java | 186 +++++++++++++++ .../placement/BackendOutagePolicyTest.java | 225 ++++++++++++++++++ 2 files changed, 411 insertions(+) create mode 100644 fleetd/src/main/java/dev/ltms/fleet/placement/BackendOutagePolicy.java create mode 100644 fleetd/src/test/java/dev/ltms/fleet/placement/BackendOutagePolicyTest.java 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..5d202ab --- /dev/null +++ b/fleetd/src/main/java/dev/ltms/fleet/placement/BackendOutagePolicy.java @@ -0,0 +1,186 @@ +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} — + * 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. A single error is common (the + * classifier is a heuristic and a valid member report can quote an {@code API Error:} line), so one + * false match must never remove fleet capacity; two is the smallest threshold that protects an + * honest one-turn failure. 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 { + + /** Classified errors on one credential needed 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) with {@code reason} (the classifier's free-text reason, kept per event + * for whoever renders the eventual notice). + * + *

Returns a populated {@link Incident} only at the exact moment the threshold is crossed — + * never before, and never again while the resulting cool-off is active. 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 errors 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. + */ + 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); + } + + int evidenceCount() { + return reasons.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..88b2721 --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/placement/BackendOutagePolicyTest.java @@ -0,0 +1,225 @@ +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 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); + } +}