Merge PR #634: fleetd #612 ranks 1+2 — behavioural pins for the exhaustion wirings
CI / shell-tests (push) Failing after 8s
CI / build (push) Failing after 2m20s
CI / contract (push) Successful in 2m22s

Pins all four exhaustion call sites in FleetdAssembly: liveExhaustedPatterns (:299),
exhaustedPatternLookup (:300), publishExhaustionSink (:327) and the independent
OpenCode forwardingExhaustionSink (:164). Rank 1 is the worst consequence in the
#612 sweep — an inert lookup hands a genuine usage-limit refusal back to a waiting
caller as real completed work instead of BACKEND_EXHAUSTED.

Two separate tests, so rank 2's OpenCode half is pinned independently: mutating
:164 fails only the forwarding test, which is the independence the ticket asserts.

Verified by the lead beyond the worker's own proof: wiring publishExhaustionSink to
a throwaway BackendQuarantine — every symbol kept at the call site, only the
collaborator identity changed — is caught by both tests. So these pins survive
mis-wiring, not just deletion.

Test-only; no production change. Tears down each assembly via the captured
shutdown hook in a finally.
This commit is contained in:
Dai Ha
2026-10-01 16:45:20 +02:00
2 changed files with 401 additions and 0 deletions
@@ -0,0 +1,211 @@
package dev.ltms.fleet;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.inject.CompletionResolver;
import dev.ltms.fleet.msg.Rendezvous;
import dev.ltms.fleet.msg.TurnToken;
import dev.ltms.fleet.placement.BackendQuarantine;
import dev.ltms.fleet.session.MemberSession;
import io.javalin.Javalin;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.Map;
import java.util.OptionalLong;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.LongSupplier;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #612 step 4, ranks 1 and 2 (publish side) — {@link FleetdAssembly} lines
* {@code liveExhaustedPatterns}/{@code exhaustedPatterns} (CB-578 stage A, the ticket's own "worst
* consequence in the whole sweep": a genuine usage-limit refusal handed back to a waiting caller
* AS REAL COMPLETED WORK) and {@code Fleetd.publishExhaustionSink(...)} (CB-578 stage B: the
* credential that hit the limit is never quarantined). None of these three lines is driven by an
* existing test through the real assembly: {@code FleetdExhaustedPatternLookupWiringTest} and
* {@code FleetdLiveExhaustedPatternsWiringTest} (fleetd #589) call {@code Fleetd.liveExhaustedPatterns}
* / {@code Fleetd.exhaustedPatternLookup} directly as factories, never through {@link
* FleetdAssembly#assembleAndStart} — they prove the FACTORY classifies correctly, never that THIS
* call site is the one that actually got wired into the running {@link CompletionResolver}. {@link
* FleetdBackendQuarantineAssemblyTest} drives {@code BackendQuarantine.withEscalation(...)}
* directly, a different call site from {@code publishExhaustionSink} here.
*
* <p>This test drives the REAL assembled {@link CompletionResolver} ({@link
* FleetdRuntime#completion()}) with a profile carrying a configured {@code exhaustedPattern},
* through a pane scrape that matches it, and asserts both halves of the production consequence:
* (1) the resolution is {@link Rendezvous.Kind#BACKEND_EXHAUSTED}, never a plain completion handed
* back as real work, and (2) the profile's credential is actually quarantined afterward, through
* the REAL {@link BackendQuarantine} the same assembly built ({@link
* FleetdRuntime#mcp()}{@code .quarantineSource().quarantine()}) — never a copy.
*
* <p>Same {@code ControllableResourcePorts} shape as {@code FleetdCompletionResolverAssemblyTest}:
* a fake, advanceable {@code nanoClock} so {@code CompletionResolver.MIN_TURN_NANOS} clears without
* a real sleep, and {@link FakeHerdr#readText} to drive the pane scrape.
*/
class FleetdExhaustedPatternAssemblyTest {
private static final class ControllableResourcePorts implements ResourcePorts {
final FakeHerdr herdr;
final AtomicLong nowNanos = new AtomicLong(1_000_000_000L); // arbitrary non-zero start
Runnable shutdownHook;
ControllableResourcePorts(FakeHerdr herdr) {
this.herdr = herdr;
}
void advanceSeconds(long seconds) {
nowNanos.addAndGet(TimeUnit.SECONDS.toNanos(seconds));
}
@Override
public Map<String, String> environment() {
return Map.of();
}
@Override
public HerdrClient connectHerdr(Path socketPath) {
return herdr;
}
@Override
public Fleetd.AmqpOpener replyInboxOpener() {
return (uri, prefetch) -> {
throw new UnsupportedOperationException("replyInboxOpener must not be called — no broker: block");
};
}
@Override
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
return (uri, selfCoordId, prefetch) -> {
throw new UnsupportedOperationException("leadMailboxOpener must not be called — no coordinator: block");
};
}
@Override
public LongSupplier nanoClock() {
return nowNanos::get;
}
@Override
public LongSupplier wallClockNanos() {
return nowNanos::get;
}
@Override
public ScheduledExecutorService newScheduler(String purpose) {
return Executors.newSingleThreadScheduledExecutor();
}
@Override
public void addShutdownHook(Runnable hook) {
this.shutdownHook = hook;
}
@Override
public void startHttp(Javalin app, String host, int port) {
// Deliberately never bind a real port.
}
}
private static FleetConfig writeConfig(Path dir, int cooldownSeconds) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
bind:
host: 127.0.0.1
port: 8765
idleSleepGuard:
enabled: false
quarantineCooldownSeconds: %d
profiles:
exhaustprofile:
baseUrl: http://exhausthost.local:8000
model: sonnet
exhaustedPattern: "usage limit reached"
guard:
offSubscriptionHosts:
- exhausthost.local
""".formatted(cooldownSeconds));
return FleetConfig.load(f);
}
/**
* fleetd #589's own description of this gap ({@code Fleetd#exhaustedPatternLookup}'s javadoc):
* "the worst consequence in the whole #589 sweep" — a genuine usage-limit refusal stops being
* classified as {@code BACKEND_EXHAUSTED} and is handed back to a waiting {@code fleet_send} as
* if it were real completed work. Pins {@code FleetdAssembly}'s {@code liveExhaustedPatterns}
* AND {@code exhaustedPatterns} lines (rank 1) together with {@code publishExhaustionSink}
* (rank 2, the non-OpenCode half) in one flow: classify, then quarantine.
*/
@Test
@DisplayName("[BEHAVIOURAL] a scrape matching the profile's exhaustedPattern resolves "
+ "BACKEND_EXHAUSTED (never a plain completion) and quarantines the credential")
void assembledResolverClassifiesExhaustionAndQuarantinesTheCredential(@TempDir Path dir) throws Exception {
int cooldownSeconds = 120;
FleetConfig cfg = writeConfig(dir, cooldownSeconds);
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
ControllableResourcePorts ports = new ControllableResourcePorts(new FakeHerdr());
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
try {
MemberSession session = runtime.sessions().acquire("exhaustprofile", null, dir.toString(), null);
String target = session.terminalId();
CompletionResolver completion = runtime.completion();
CompletableFuture<Rendezvous.Resolution> waiter = new CompletableFuture<>();
ports.herdr.readText("idle, nothing yet");
completion.onDelivered(target, new TurnToken(target, waiter, null));
// The matched text must START the pane line (CompletionResolver.startsWithExhaustion) —
// no preceding sentence — for the quarantine side-effect to fire, same as production.
ports.herdr.readText("usage limit reached: try again in a few hours");
ports.advanceSeconds(3); // clear CompletionResolver.MIN_TURN_NANOS (2s), no real sleep
completion.resolveBeforePostAction(target);
// CONTROL: the waiter must have resolved synchronously at all — if the assembled
// CompletionResolver were never actually driven (e.g. a wiring break upstream silently
// left the resolver unreachable), this fails loudly before the real assertions below
// ever run, rather than passing on an untouched waiter.
Rendezvous.Resolution resolution = waiter.getNow(null);
assertTrue(resolution != null, "CONTROL: the waiter must have resolved synchronously — "
+ "if this is null, the assembled resolver was never actually exercised");
assertEquals(Rendezvous.Kind.BACKEND_EXHAUSTED, resolution.kind(),
"a scrape matching the profile's configured exhaustedPattern must classify as "
+ "BACKEND_EXHAUSTED, not a plain completion handed back as real work — "
+ "replacing FleetdAssembly's liveExhaustedPatterns/exhaustedPatterns "
+ "lines with their inert forms (Map.of() / target -> null) must fail "
+ "this assertion; got: " + resolution);
assertTrue(resolution.text().contains("usage limit reached"), resolution.text());
BackendQuarantine quarantine = runtime.mcp().quarantineSource().quarantine();
assertTrue(quarantine.isQuarantined("exhaustprofile"),
"the real publishExhaustionSink-built sink must have quarantined the profile's "
+ "credential (effectiveCredentialId() == the profile name here, no "
+ "credentialId configured) — replacing FleetdAssembly's "
+ "publishExhaustionSink call site with a hardcoded ExhaustionSink.none() "
+ "must fail this assertion, since nothing would ever call "
+ "quarantine.quarantine(...)");
OptionalLong remaining = quarantine.remainingSeconds("exhaustprofile");
assertTrue(remaining.isPresent() && remaining.getAsLong() > 0
&& remaining.getAsLong() <= cooldownSeconds,
"a fresh quarantine must block for at most the configured base cooldown: " + remaining);
} finally {
if (ports.shutdownHook != null) ports.shutdownHook.run();
}
}
}
@@ -0,0 +1,190 @@
package dev.ltms.fleet;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.inject.ExhaustionSink;
import dev.ltms.fleet.member.CompositePeerLauncher;
import dev.ltms.fleet.member.HerdrPeerLauncher;
import dev.ltms.fleet.peer.PeerLauncher;
import dev.ltms.fleet.placement.BackendQuarantine;
import io.javalin.Javalin;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import java.lang.reflect.Field;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.Map;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.LongSupplier;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #612 step 4, rank 2 (OpenCode half) — {@link FleetdAssembly}'s {@code
* forwardingExhaustionSink} line ({@code Fleetd.forwardingExhaustionSink(exhaustionSinkRef)}),
* handed to {@link dev.ltms.fleet.member.OpenCodeLauncher} so its fleetd #175 model-mismatch check
* can quarantine a credential before {@code sessions} exists to build the real sink (the
* construction-order cycle documented at that call site). The ticket calls this independent from
* {@code publishExhaustionSink} (pinned by {@link FleetdExhaustedPatternAssemblyTest}): a credential
* that hits a usage limit through THIS path is never quarantined if {@code forwardingExhaustionSink}
* is swapped for a hardcoded {@link ExhaustionSink#none()} at that call site — the OpenCode
* launcher's own quarantine check keeps compiling and keeps "running", but it permanently talks to
* a sink that does nothing, independent of whatever {@code publishExhaustionSink} does later.
*
* <p>{@code FleetdExhaustionSinkForwardingWiringTest} (fleetd #589) already proves {@code
* Fleetd.forwardingExhaustionSink(ref)} forwards to whatever {@code ref} holds — as a bare factory
* call, never through {@link FleetdAssembly#assembleAndStart}. It proves nothing about whether
* THIS call site is the one FleetdAssembly actually wires into the real {@code OpenCodeLauncher}
* it builds, which is exactly the #602/#606-shaped gap this ticket exists to close.
*
* <p>No accessor on {@link FleetdRuntime} reaches the adapter instances (by design — see that
* class's own javadoc: only the final collaborators it owns directly are exposed), so this test
* reaches the REAL, assembled {@code OpenCodeLauncher}'s {@code exhaustionSink} field the same way
* {@code SessionManager}/{@code CompositePeerLauncher} wire it internally: a short, targeted
* reflective walk ({@code SessionManager.launcher} → {@code CompositePeerLauncher.byProfile} →
* {@code OpenCodeLauncher.exhaustionSink}) onto the exact object the assembly built — never a copy,
* and never a read of the source text. Reflection is used the same way elsewhere in this suite
* (e.g. {@code StatusPollerWatchdogTest}) to reach a private collaborator a production constructor
* intentionally does not expose a public accessor for.
*/
class FleetdOpenCodeExhaustionForwardingAssemblyTest {
private static final class ControllableResourcePorts implements ResourcePorts {
final FakeHerdr herdr = new FakeHerdr();
final AtomicLong nowNanos = new AtomicLong(1_000_000_000L);
Runnable shutdownHook;
@Override
public Map<String, String> environment() {
return Map.of();
}
@Override
public HerdrClient connectHerdr(Path socketPath) {
return herdr;
}
@Override
public Fleetd.AmqpOpener replyInboxOpener() {
return (uri, prefetch) -> {
throw new UnsupportedOperationException("replyInboxOpener must not be called — no broker: block");
};
}
@Override
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
return (uri, selfCoordId, prefetch) -> {
throw new UnsupportedOperationException("leadMailboxOpener must not be called — no coordinator: block");
};
}
@Override
public LongSupplier nanoClock() {
return nowNanos::get;
}
@Override
public LongSupplier wallClockNanos() {
return nowNanos::get;
}
@Override
public ScheduledExecutorService newScheduler(String purpose) {
return Executors.newSingleThreadScheduledExecutor();
}
@Override
public void addShutdownHook(Runnable hook) {
this.shutdownHook = hook;
}
@Override
public void startHttp(Javalin app, String host, int port) {
// Deliberately never bind a real port.
}
}
private static FleetConfig writeConfig(Path dir, int cooldownSeconds) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
bind:
host: 127.0.0.1
port: 8765
idleSleepGuard:
enabled: false
quarantineCooldownSeconds: %d
profiles:
gemini:
kind: opencode
model: google/gemini-2.5-pro
""".formatted(cooldownSeconds));
return FleetConfig.load(f);
}
/** Reach a declared field by name on {@code target}'s runtime class, bypassing the access check. */
private static Object readField(Object target, Class<?> declaringClass, String fieldName) throws Exception {
Field field = declaringClass.getDeclaredField(fieldName);
field.setAccessible(true);
return field.get(target);
}
@Test
@DisplayName("[BEHAVIOURAL] the real assembled OpenCodeLauncher's exhaustionSink field forwards "
+ "an onExhausted call into the real daemon's BackendQuarantine")
void assembledOpenCodeLauncherExhaustionSinkQuarantinesTheCredential(@TempDir Path dir) throws Exception {
int cooldownSeconds = 90;
FleetConfig cfg = writeConfig(dir, cooldownSeconds);
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
ControllableResourcePorts ports = new ControllableResourcePorts();
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
try {
PeerLauncher launcherField = (PeerLauncher) readField(runtime.sessions(),
runtime.sessions().getClass(), "launcher");
// CONTROL: the composite launcher must actually be the real production type with a
// "gemini" -> OpenCodeLauncher entry — if this fails, nothing below exercised the real
// assembly at all, rather than silently passing on an empty/wrong object.
assertTrue(launcherField instanceof CompositePeerLauncher,
"CONTROL: SessionManager.launcher must be the real CompositePeerLauncher the "
+ "assembly built, got: " + launcherField);
@SuppressWarnings("unchecked")
Map<String, HerdrPeerLauncher> byProfile = (Map<String, HerdrPeerLauncher>)
readField(launcherField, CompositePeerLauncher.class, "byProfile");
HerdrPeerLauncher adapter = byProfile.get("gemini");
assertTrue(adapter != null && adapter.getClass().getSimpleName().equals("OpenCodeLauncher"),
"CONTROL: the 'gemini' profile must resolve to a real OpenCodeLauncher adapter, "
+ "got: " + adapter);
ExhaustionSink sink = (ExhaustionSink) readField(adapter, adapter.getClass(), "exhaustionSink");
assertTrue(sink != null, "CONTROL: OpenCodeLauncher.exhaustionSink must never be null");
// The exact call OpenCodeLauncher.SessionAwareHandle#checkModelMatch makes on a real
// model mismatch (fleetd #175): target, reason, and its own already-known profile name.
sink.onExhausted("term_gemini_1", "opencode model mismatch (test)", "gemini");
BackendQuarantine quarantine = runtime.mcp().quarantineSource().quarantine();
assertTrue(quarantine.isQuarantined("gemini"),
"the real forwardingExhaustionSink-wired field must have delegated into the "
+ "published production sink, which quarantines the profile's credential "
+ "('gemini' here — no credentialId configured) — replacing "
+ "FleetdAssembly's forwardingExhaustionSink call site with a hardcoded "
+ "ExhaustionSink.none() must fail this assertion, since the field read "
+ "above would then BE the inert no-op and nothing would ever reach "
+ "quarantine.quarantine(...)");
assertEquals(cooldownSeconds, quarantine.remainingSeconds("gemini").orElseThrow(
() -> new AssertionError("credential must report a remaining cooldown")));
} finally {
if (ports.shutdownHook != null) ports.shutdownHook.run();
}
}
}