Compare commits
6 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| cf54aed451 | |||
| 826e0aeb2a | |||
| 0f08b93659 | |||
| 432c1d92d1 | |||
| 60b7e67b42 | |||
| 3fbd43fe3f |
@@ -3,11 +3,13 @@ package dev.ltms.fleet.member;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.herdr.Agent;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.AgentStatus;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.herdr.HerdrException;
|
||||
import dev.ltms.fleet.herdr.Tab;
|
||||
import dev.ltms.fleet.herdr.Workspace;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
import dev.ltms.fleet.inject.StatusRefiner;
|
||||
import dev.ltms.fleet.peer.Capability;
|
||||
import dev.ltms.fleet.peer.CharterReceipt;
|
||||
import dev.ltms.fleet.peer.MemberRole;
|
||||
@@ -151,6 +153,14 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
private final LongSupplier nowMillis; // monotonic clock (injectable for tests)
|
||||
private final Runnable sleeper; // sleep/wait hook (injectable for tests; never real-sleep in unit tests)
|
||||
|
||||
/**
|
||||
* fleetd #176 fix 2: the same UNKNOWN-refinement {@code dev.ltms.fleet.inject.StatusPoller}
|
||||
* uses, reused here for the spawn-readiness gate. Constructed once from {@link #agents} — see
|
||||
* {@link #refinedInjectable(String, Agent)} for the corroboration that keeps it from firing on
|
||||
* a dead pane's bare shell prompt.
|
||||
*/
|
||||
private final StatusRefiner statusRefiner;
|
||||
|
||||
// Per-process token mixed into each peer name so a fresh process (nameSeq back at 0) cannot
|
||||
// collide with same-profile peers that outlived a restart. See startUniquelyNamed.
|
||||
private final String nameNonce = String.format("%06x", new SecureRandom().nextInt(1 << 24));
|
||||
@@ -272,6 +282,7 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
this.memberCredentials = memberCredentials;
|
||||
this.hostEnvNames = hostEnvNames != null ? hostEnvNames : () -> System.getenv().keySet();
|
||||
this.config = config;
|
||||
this.statusRefiner = new StatusRefiner(agents);
|
||||
}
|
||||
|
||||
// --- adapter seams -------------------------------------------------------------------------
|
||||
@@ -909,16 +920,35 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
// --- spawn-readiness gate (CB-306) ---------------------------------------------------------
|
||||
|
||||
/**
|
||||
* Poll {@link AgentControl#status} until the pane reports an injectable state or the configured
|
||||
* Poll {@link AgentControl#get} until the pane reports an injectable state or the configured
|
||||
* timeout elapses. On timeout, close the pane (self-reap) and throw.
|
||||
*
|
||||
* <p>fleetd #176 fix 1: {@code agents.get} was previously unguarded here, so a herdr
|
||||
* {@code *_not_found} answer — which is what happens when the backend process EXITED rather
|
||||
* than being slow — propagated as a raw {@link HerdrException} instead of the
|
||||
* {@link PeerUnreachableException} every other failure path of this gate produces, and skipped
|
||||
* teardown ({@link #stop}) entirely, leaking the pane/tab. {@link #failFastOnGoneBackend} closes
|
||||
* that gap: it stops waiting immediately (never burns the rest of the timeout), runs the same
|
||||
* teardown the timeout path below runs, and throws with a message that says the backend exited
|
||||
* rather than that the pane was slow. Any other {@link HerdrException} still propagates
|
||||
* unchanged — this gate does not know how to recover from it.
|
||||
*/
|
||||
private void waitUntilInjectableOrThrow(String paneId) {
|
||||
long deadline = nowMillis.getAsLong() + spawnReadyTimeoutMs;
|
||||
Object lastStatus = null;
|
||||
long start = nowMillis.getAsLong();
|
||||
long deadline = start + spawnReadyTimeoutMs;
|
||||
AgentStatus lastStatus = null;
|
||||
while (nowMillis.getAsLong() < deadline) {
|
||||
var status = agents.status(paneId);
|
||||
lastStatus = status;
|
||||
if (status.injectable()) {
|
||||
Agent sample;
|
||||
try {
|
||||
sample = agents.get(paneId);
|
||||
} catch (HerdrException e) {
|
||||
if (isAlreadyGone(e)) {
|
||||
failFastOnGoneBackend(paneId, e, nowMillis.getAsLong() - start);
|
||||
}
|
||||
throw e; // any other herdr failure is not ours to interpret — let it propagate
|
||||
}
|
||||
lastStatus = sample.status();
|
||||
if (lastStatus.injectable() || refinedInjectable(paneId, sample)) {
|
||||
log.debug("peer pane={} reached injectable state", paneId);
|
||||
return;
|
||||
}
|
||||
@@ -937,6 +967,66 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
+ spawnReadyTimeoutMs + "ms");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #176 fix 1: the backend process exited while the gate was still waiting — herdr
|
||||
* answered {@code *_not_found} instead of ever reporting an injectable status. Fails
|
||||
* immediately (never burns the rest of {@link #spawnReadyTimeoutMs}), runs the exact same
|
||||
* teardown {@link #waitUntilInjectableOrThrow}'s timeout path runs, and throws a
|
||||
* {@link PeerUnreachableException} whose message says the process exited rather than that the
|
||||
* pane was slow.
|
||||
*
|
||||
* @throws PeerUnreachableException always — this method never returns normally
|
||||
*/
|
||||
private void failFastOnGoneBackend(String paneId, HerdrException cause, long elapsedMs) {
|
||||
String tail = readPaneQuietly(paneId); // read before stop() closes the pane
|
||||
log.warn("peer pane={} backend process exited after {}ms while waiting for injectable "
|
||||
+ "state (herdr: {}) — closing. Pane tail:\n{}",
|
||||
paneId, elapsedMs, cause.getMessage(), tail);
|
||||
stop(paneId);
|
||||
throw new PeerUnreachableException(
|
||||
"worker pane " + paneId + " backend process exited after " + elapsedMs
|
||||
+ "ms while waiting to become injectable (herdr reported: " + cause.getMessage()
|
||||
+ "). Pane tail:\n" + tail);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #176 fix 2: resolve a raw {@link AgentStatus#UNKNOWN} sample into a trustworthy
|
||||
* injectable state via pane content, the same refinement {@code StatusPoller} applies during a
|
||||
* peer's working life — but corroborated, because the trap this gate is exposed to that the
|
||||
* poller is not: a pane whose backend has already exited settles at a plain shell prompt, and
|
||||
* that prompt commonly contains the same {@code ❯} glyph {@link StatusRefiner#classify} treats
|
||||
* as "idle at the Claude Code TUI prompt". Naively wiring the refiner in would turn "the backend
|
||||
* died" into "ready to inject" — strictly worse than today's timeout.
|
||||
*
|
||||
* <p>Two guards, both required:
|
||||
* <ul>
|
||||
* <li>{@link StatusRefiner#classify} is written for the Claude Code TUI only (its own javadoc
|
||||
* says so), so refinement only ever runs for the {@code claude} adapter — never for
|
||||
* {@code opencode} or any future non-Claude backend, whatever its pane looks like.</li>
|
||||
* <li>The refined status is accepted only when the <em>same</em> {@code agents.get} sample
|
||||
* ({@code sample}, taken once by the caller) still reports a non-null
|
||||
* {@link Agent#agentType()} — herdr's own "a supported backend is still detected here"
|
||||
* signal, read from the very sample the raw status came from so the two can never
|
||||
* disagree. A bare shell prompt left by an exited backend reports no {@code agentType}.</li>
|
||||
* </ul>
|
||||
*
|
||||
* <p>Called only when {@code sample.status() == UNKNOWN}, so a healthy spawn — which never sees
|
||||
* {@code UNKNOWN} — triggers zero extra herdr calls; only a persistently-{@code unknown} pane
|
||||
* pays the one extra {@code agent.read} {@link StatusRefiner#refine} performs.
|
||||
*/
|
||||
private boolean refinedInjectable(String paneId, Agent sample) {
|
||||
if (sample.status() != AgentStatus.UNKNOWN) {
|
||||
return false;
|
||||
}
|
||||
if (!"claude".equals(namePrefix)) {
|
||||
return false; // StatusRefiner.classify reads a Claude Code TUI prompt specifically
|
||||
}
|
||||
if (sample.agentType() == null) {
|
||||
return false; // no corroborating liveness signal — could be a dead pane's bare shell
|
||||
}
|
||||
return statusRefiner.refine(paneId, AgentStatus.UNKNOWN, agents).injectable();
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #220: the pane's recent output, clipped, for the readiness-gate timeout log — or a
|
||||
* short note when it cannot be read. Best-effort by construction: this runs on a path that is
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -44,10 +44,13 @@ public final class FakeHerdr implements HerdrClient {
|
||||
private String agentSendErrorCode = null;
|
||||
private boolean noPanes = false;
|
||||
private volatile String agentStatus = "idle"; // steady-state agent.get status
|
||||
private volatile String agentType = "claude"; // detected agent kind on agent.get; null = undetected
|
||||
private volatile String readText = "worker transcript tail"; // canned agent.read output
|
||||
private int pinnedStarts = 0; // how many upcoming agent.start calls report a fixed pane
|
||||
private String pinnedStartTerminal;
|
||||
private String pinnedStartPane;
|
||||
private volatile int agentGetOkCalls = Integer.MAX_VALUE; // how many agent.get calls succeed first
|
||||
private volatile String agentGetFailCode = null; // error code every agent.get call after that reports
|
||||
|
||||
public FakeHerdr healthy(boolean h) {
|
||||
this.healthy = h;
|
||||
@@ -103,6 +106,27 @@ public final class FakeHerdr implements HerdrClient {
|
||||
return this;
|
||||
}
|
||||
|
||||
/**
|
||||
* Set the detected agent kind ({@code "agent"} field) that {@code agent.get} reports —
|
||||
* {@code null} models herdr not (or no longer) detecting a supported backend in the pane, e.g.
|
||||
* a bare shell prompt (fleetd #176 fix 2 corroboration test).
|
||||
*/
|
||||
public FakeHerdr agentType(String type) {
|
||||
this.agentType = type;
|
||||
return this;
|
||||
}
|
||||
|
||||
/**
|
||||
* Make {@code agent.get} succeed normally for its first {@code okCalls} invocations, then fail
|
||||
* every call after that with {@code code} — fleetd #176 fix 1's "backend exited mid-wait"
|
||||
* fixture. {@code okCalls == 0} fails from the very first call.
|
||||
*/
|
||||
public FakeHerdr agentGetFailsWithAfter(int okCalls, String code) {
|
||||
this.agentGetOkCalls = okCalls;
|
||||
this.agentGetFailCode = code;
|
||||
return this;
|
||||
}
|
||||
|
||||
/** The text {@code agent.read} returns (the CB-106 completion scrape). */
|
||||
public FakeHerdr readText(String text) {
|
||||
this.readText = text;
|
||||
@@ -208,10 +232,21 @@ public final class FakeHerdr implements HerdrClient {
|
||||
}
|
||||
yield mapper.readTree("{\"type\":\"ok\"}");
|
||||
}
|
||||
case "agent.get" -> mapper.readTree(("""
|
||||
{"type":"agent_info","agent":{"terminal_id":"term_a","agent":"claude",
|
||||
case "agent.get" -> {
|
||||
if (agentGetFailCode != null) {
|
||||
long getCalls = calls.stream().filter(c -> c.method().equals("agent.get")).count();
|
||||
if (getCalls > agentGetOkCalls) {
|
||||
throw new HerdrException(
|
||||
"herdr error [" + agentGetFailCode + "]: agent target not found",
|
||||
agentGetFailCode, null);
|
||||
}
|
||||
}
|
||||
String agentField = agentType == null ? "null" : "\"" + agentType + "\"";
|
||||
yield mapper.readTree(("""
|
||||
{"type":"agent_info","agent":{"terminal_id":"term_a","agent":%s,
|
||||
"agent_status":"%s","workspace_id":"w2","tab_id":"w2:t7","pane_id":"w2:p7"}}""")
|
||||
.formatted(agentStatus));
|
||||
.formatted(agentField, agentStatus));
|
||||
}
|
||||
case "agent.read" -> mapper.readTree(mapper.writeValueAsString(
|
||||
java.util.Map.of("type", "agent_read", "read", java.util.Map.of("text", readText))));
|
||||
case "agent.start" -> {
|
||||
|
||||
@@ -963,6 +963,152 @@ class ClaudeCodeLauncherTest {
|
||||
assertDoesNotThrow(() -> UUID.fromString(handle.id()));
|
||||
}
|
||||
|
||||
// --- fleetd #176 fix 1: fail fast when the backend process exits mid-wait -------------------
|
||||
|
||||
@Test
|
||||
void spawnFailsFastWhenBackendProcessExitsMidWaitInsteadOfBurningTheTimeout() {
|
||||
// First agent.get sees UNKNOWN (one normal tick); the second reports the pane gone, exactly
|
||||
// what herdr answers when the backend process has already exited. The timeout is generous
|
||||
// (60s) so a test that wrongly falls through to the old unguarded call — and therefore
|
||||
// waits out the whole window — is unambiguously distinguishable from one that fails fast.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
herdr.agentStatus("unknown");
|
||||
herdr.agentGetFailsWithAfter(1, "pane_not_found");
|
||||
long[] clock = {0};
|
||||
|
||||
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(
|
||||
new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")),
|
||||
workerConfigMap("ltms-local", null), "ltms-local", _ -> null,
|
||||
60_000, () -> clock[0], () -> clock[0] += 50);
|
||||
|
||||
PeerUnreachableException ex = assertThrows(
|
||||
PeerUnreachableException.class,
|
||||
() -> svc.spawn(new SpawnRequest(null, null, null)));
|
||||
|
||||
assertTrue(ex.getMessage().contains("w9:pRoot_1"),
|
||||
"exception message references the paneId: " + ex.getMessage());
|
||||
assertTrue(ex.getMessage().toLowerCase().contains("exited"),
|
||||
"exception message says the backend exited, not that the pane was slow: "
|
||||
+ ex.getMessage());
|
||||
assertTrue(clock[0] < 60_000,
|
||||
"the gate must not burn the rest of the 60s timeout: clock only reached " + clock[0]);
|
||||
long getCalls = herdr.calls.stream().filter(c -> c.method().equals("agent.get")).count();
|
||||
assertEquals(2, getCalls,
|
||||
"exactly one normal poll then the not_found answer — no further polling after that: "
|
||||
+ getCalls);
|
||||
assertEquals(1, paneCloseCount(herdr, "w9:pRoot_1"),
|
||||
"the pane is torn down (no orphan) even on the fail-fast path");
|
||||
}
|
||||
|
||||
@Test
|
||||
void spawnLetsAnUnrelatedHerdrErrorPropagateUnchanged() {
|
||||
// Fix 1 must only special-case a "*_not_found" answer. Any other herdr failure keeps
|
||||
// propagating as-is — this gate does not know how to recover from it.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
herdr.agentStatus("unknown");
|
||||
herdr.agentGetFailsWithAfter(0, "internal_error");
|
||||
long[] clock = {0};
|
||||
|
||||
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(
|
||||
new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")),
|
||||
workerConfigMap("ltms-local", null), "ltms-local", _ -> null,
|
||||
60_000, () -> clock[0], () -> clock[0] += 50);
|
||||
|
||||
dev.ltms.fleet.herdr.HerdrException ex = assertThrows(
|
||||
dev.ltms.fleet.herdr.HerdrException.class,
|
||||
() -> svc.spawn(new SpawnRequest(null, null, null)));
|
||||
|
||||
assertEquals("internal_error", ex.code());
|
||||
assertEquals(0, paneCloseCount(herdr, "w9:pRoot_1"),
|
||||
"an error this gate does not recognize is not this gate's teardown to run");
|
||||
}
|
||||
|
||||
// --- fleetd #176 fix 2: corroborated UNKNOWN refinement --------------------------------------
|
||||
|
||||
@Test
|
||||
void refinedIdleIsNotAcceptedWhenAgentTypeIsNull() {
|
||||
// The trap fix 2 must close: a pane sitting at a bare shell prompt after its backend exited
|
||||
// still contains the same "❯" glyph StatusRefiner.classify treats as "idle at the Claude
|
||||
// Code TUI prompt". Without the agentType corroboration this would be misread as injectable
|
||||
// and the gate would hand back a peer that never started. herdr reports no agentType for
|
||||
// that bare shell.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
herdr.agentStatus("unknown"); // always UNKNOWN
|
||||
herdr.agentType(null); // no supported backend detected — could be a bare shell
|
||||
herdr.readText("some-host:~ user$ ❯ "); // looks exactly like an idle Claude Code prompt
|
||||
long[] clock = {0};
|
||||
|
||||
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(
|
||||
new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")),
|
||||
workerConfigMap("ltms-local", null), "ltms-local", _ -> null,
|
||||
500, () -> clock[0], () -> clock[0] += 50);
|
||||
|
||||
PeerUnreachableException ex = assertThrows(
|
||||
PeerUnreachableException.class,
|
||||
() -> svc.spawn(new SpawnRequest(null, null, null)));
|
||||
|
||||
assertTrue(clock[0] >= 500,
|
||||
"a null agentType must not let the '❯' prompt refine to injectable — the gate has "
|
||||
+ "to wait out the full timeout: clock only reached " + clock[0]);
|
||||
assertEquals(1, paneCloseCount(herdr, "w9:pRoot_1"),
|
||||
"the pane is torn down on timeout, same as any other never-injectable spawn");
|
||||
}
|
||||
|
||||
@Test
|
||||
void refinedIdleIsAcceptedWhenAgentTypeCorroboratesLiveness() {
|
||||
// The positive case: a genuinely live claude pane that herdr misreports as UNKNOWN (CB-115)
|
||||
// still resolves to injectable once agentType corroborates that a supported backend is
|
||||
// detected in the same sample.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
herdr.agentStatus("unknown"); // always UNKNOWN at the raw level
|
||||
herdr.agentType("claude"); // herdr still detects a live claude backend
|
||||
herdr.readText("? for shortcuts"); // StatusRefiner.classify's idle footer marker
|
||||
long[] clock = {0};
|
||||
|
||||
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(
|
||||
new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")),
|
||||
workerConfigMap("ltms-local", null), "ltms-local", _ -> null,
|
||||
60_000, () -> clock[0], () -> clock[0] += 50);
|
||||
|
||||
PeerHandle handle = svc.spawn(new SpawnRequest(null, null, null));
|
||||
|
||||
assertNotNull(handle, "a corroborated refined-IDLE sample lets the spawn succeed");
|
||||
assertTrue(clock[0] < 60_000,
|
||||
"refinement must resolve well before the timeout: clock reached " + clock[0]);
|
||||
assertEquals(0, paneCloseCount(herdr, "w9:pRoot_1"),
|
||||
"no pane close — the peer is genuinely injectable");
|
||||
}
|
||||
|
||||
@Test
|
||||
void refinementNeverFiresWhenRawStatusIsAlreadyInjectable() {
|
||||
// Acceptance criterion 4: refinement must trigger only on a raw UNKNOWN sample. Seed the
|
||||
// pane content with an active-generation marker that StatusRefiner.classify would read as
|
||||
// WORKING (never injectable) if refine() were wrongly invoked here — proving that a raw
|
||||
// IDLE status short-circuits before refine() (and its extra agent.read call) ever runs.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
herdr.agentStatus("idle"); // already injectable at the raw level
|
||||
herdr.agentType("claude");
|
||||
herdr.readText("esc to interrupt"); // would classify as WORKING if refine() ran anyway
|
||||
|
||||
long[] clock = {0};
|
||||
ClaudeCodeLauncher gated = new ClaudeCodeLauncher(
|
||||
new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")),
|
||||
workerConfigMap("ltms-local", null), "ltms-local", _ -> null,
|
||||
5000, () -> clock[0], () -> clock[0] += 300);
|
||||
|
||||
PeerHandle handle = gated.spawn(new SpawnRequest(null, null, null));
|
||||
|
||||
assertNotNull(handle, "an already-injectable raw status succeeds without ever refining");
|
||||
assertFalse(herdr.called("agent.read"),
|
||||
"refine() must never run (and so never call agent.read) when the raw status is "
|
||||
+ "already injectable");
|
||||
}
|
||||
|
||||
// --- CB-511: worker environment seeding -----------------------------------------------------
|
||||
|
||||
@Test
|
||||
|
||||
@@ -396,6 +396,46 @@ class OpenCodeLauncherTest {
|
||||
assertNotNull(ex.getMessage());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #176 fix 2's per-adapter guard: {@code StatusRefiner.classify} is written for the
|
||||
* Claude Code TUI only, and the spawn-readiness gate must never run it against an opencode
|
||||
* pane. {@code agentType("opencode")} deliberately satisfies the OTHER guard (the corroborating
|
||||
* liveness check) so it cannot be what makes this test pass — only the {@code namePrefix}
|
||||
* check can be. If that check were ever removed, this pane's {@code ❯} content would refine
|
||||
* straight to IDLE and the gate would report ready before the backend actually was.
|
||||
*/
|
||||
@Test
|
||||
void opencodePaneIsNeverRefinedEvenWhenItsContentLooksLikeAnIdleClaudePrompt(@TempDir Path root) {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
herdr.agentStatus("unknown"); // always UNKNOWN
|
||||
herdr.agentType("opencode"); // non-null — satisfies the liveness guard on its own
|
||||
herdr.readText("some-host:~ user$ ❯ "); // content StatusRefiner.classify reads as IDLE
|
||||
long[] clock = {0};
|
||||
OpenCodeLauncher svc = new OpenCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
Map.of("gemini", opencodeCfg(null, null, null)), "gemini", _ -> null,
|
||||
1000, () -> clock[0], () -> clock[0] += 50, root, root);
|
||||
|
||||
assertThrows(PeerUnreachableException.class,
|
||||
() -> svc.spawn(new SpawnRequest(null, null, null)));
|
||||
|
||||
assertTrue(clock[0] >= 1000,
|
||||
"an opencode pane must never refine to injectable, however its content reads — "
|
||||
+ "the gate has to wait out the full timeout: clock only reached " + clock[0]);
|
||||
// A blanket "agent.read is never called" does not hold here: the timeout path itself reads
|
||||
// the pane tail (source=recent) for its own log message, on every timeout, regardless of
|
||||
// adapter — see HerdrPeerLauncher.readPaneQuietly. So assert on the refiner's OWN probe
|
||||
// source (StatusRefiner.PROBE_SOURCE = "detection") instead — that call happens only inside
|
||||
// StatusRefiner.refine, so its absence proves the refiner itself was never reached for this
|
||||
// opencode pane, not merely that its answer was discarded.
|
||||
long detectionReads = herdr.calls.stream()
|
||||
.filter(c -> c.method().equals("agent.read"))
|
||||
.filter(c -> "detection".equals(((Map<?, ?>) c.params()).get("source")))
|
||||
.count();
|
||||
assertEquals(0, detectionReads,
|
||||
"the refiner's own pane probe (source=detection) must never run against an "
|
||||
+ "opencode pane — the namePrefix guard has to stop it before that call");
|
||||
}
|
||||
|
||||
@Test
|
||||
void spawnReturnsHandleWhenGateDisabled(@TempDir Path root) {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user