From 2be287ea03ba590e2f8acf1ac4acf7800c893825 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 1 Oct 2026 16:38:00 +0200 Subject: [PATCH] fleetd #612 step 4 (ranks 1&2): behavioural assembly tests for the exhaustion wirings MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- .../FleetdExhaustedPatternAssemblyTest.java | 211 ++++++++++++++++++ ...nCodeExhaustionForwardingAssemblyTest.java | 190 ++++++++++++++++ 2 files changed, 401 insertions(+) create mode 100644 fleetd/src/test/java/dev/ltms/fleet/FleetdExhaustedPatternAssemblyTest.java create mode 100644 fleetd/src/test/java/dev/ltms/fleet/FleetdOpenCodeExhaustionForwardingAssemblyTest.java diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdExhaustedPatternAssemblyTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdExhaustedPatternAssemblyTest.java new file mode 100644 index 0000000..ee37adb --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/FleetdExhaustedPatternAssemblyTest.java @@ -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. + * + *

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. + * + *

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 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 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(); + } + } +} diff --git a/fleetd/src/test/java/dev/ltms/fleet/FleetdOpenCodeExhaustionForwardingAssemblyTest.java b/fleetd/src/test/java/dev/ltms/fleet/FleetdOpenCodeExhaustionForwardingAssemblyTest.java new file mode 100644 index 0000000..9412fdb --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/FleetdOpenCodeExhaustionForwardingAssemblyTest.java @@ -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. + * + *

{@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. + * + *

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 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 byProfile = (Map) + 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(); + } + } +}