fleetd #612 step 4 (ranks 1&2): behavioural assembly tests for the exhaustion wirings
Replaces nothing (no existing test covered these call sites through the real assembly); adds two new FleetdAssembly-driven tests, matching Unit A's pattern of inspecting FleetdRuntime's real, assembled objects rather than a copy. - FleetdExhaustedPatternAssemblyTest pins the liveExhaustedPatterns/ exhaustedPatterns call sites (rank 1 — "the worst consequence in the whole sweep": a genuine usage-limit refusal handed back as real completed work) together with publishExhaustionSink (rank 2, non-OpenCode half): drives a real CompletionResolver through a pane scrape matching a configured exhaustedPattern and asserts BACKEND_EXHAUSTED classification plus a real BackendQuarantine credential quarantine. - FleetdOpenCodeExhaustionForwardingAssemblyTest pins forwardingExhaustionSink (rank 2, OpenCode half — independent of publishExhaustionSink per the ticket) by reflectively reaching the real, assembled OpenCodeLauncher's exhaustionSink field (SessionManager.launcher -> CompositePeerLauncher. byProfile -> OpenCodeLauncher.exhaustionSink) and proving it forwards into the same production BackendQuarantine. All four call sites (FleetdAssembly.java:164,299,300,327) were each put through grep-anchor -> line-anchored sed mutation -> mvn -o compile -> full mvn -o test (named test RED) -> restore -> full mvn -o test (1885/0 GREEN). Mutating line 327 alone also fails the OpenCode test, confirming the documented construction-order dependency (forwardingExhaustionSink reads exhaustionSinkRef, which publishExhaustionSink sets) without weakening either site's independent pin.
This commit is contained in:
@@ -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();
|
||||
}
|
||||
}
|
||||
}
|
||||
+190
@@ -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();
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user