From 4877992a7075e60bd05e29797e83ba9981380824 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 3 Sep 2026 10:09:26 +0700 Subject: [PATCH 1/4] fleetd #234: key the opencode model check on the resolved session id, and make the spawn-time quarantine actually happen Defect 1: OpenCodeSessionDiscovery.actualModelForDirectory queried WHERE directory = ?, the same heuristic sessionIdForDirectory uses. Since a default fleet_spawn (no worktree:) shares the lead's cwd with every other worker and every past session ever run there, the model read-back could silently compare against a DIFFERENT session's row. Renamed to actualModelForSessionId(sessionId), keyed on the primary key id instead, and made OpenCodeLauncher's SessionAwareHandle cache the resolved id once non-null (AtomicReference) so a later sibling row in the same directory can never flip which session's evidence is read. sessionIdForDirectory (#209) is left directory-based on purpose, with a comment explaining why the heuristic is unavoidable at that layer. Defect 2: the ERROR log claimed "quarantining this profile's credential" but Fleetd's ExhaustionSink lambda resolved target -> roster -> profile -> credential, while OpenCodeLauncher's model-mismatch check fires from agentSessionId() during SessionManager.acquire(), before the session is registered in the roster -- the lookup found nothing and silently no-opped. Added a default 3-arg ExhaustionSink.onExhausted(target, reason, profile) overload (defaults to the 2-arg method, so CompletionResolver's two call sites are unchanged); OpenCodeLauncher now passes its own already-known profile name; Fleetd's sink became an anonymous class that tries the roster first, falls back to the hint, and logs loudly at ERROR naming target/reason when neither resolves, instead of silently no-oping. Both fixes proven by mutation: reverting each independently makes its new test fail with a real assertion message, restoring makes it pass again. mvn clean install: Tests run: 1127, Failures: 0, Errors: 0, Skipped: 0, BUILD SUCCESS. --- .../src/main/java/dev/ltms/fleet/Fleetd.java | 51 ++++-- .../dev/ltms/fleet/inject/ExhaustionSink.java | 26 +++ .../ltms/fleet/member/OpenCodeLauncher.java | 51 +++++- .../member/OpenCodeSessionDiscovery.java | 66 +++++--- .../fleet/member/OpenCodeLauncherTest.java | 153 ++++++++++++++++++ .../member/OpenCodeSessionDiscoveryTest.java | 45 ++++-- 6 files changed, 338 insertions(+), 54 deletions(-) diff --git a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java index fa9541f..a8eed1b 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java +++ b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java @@ -345,17 +345,46 @@ public final class Fleetd { // model-mismatch check — it needs nothing profile-specific from the caller beyond `target` // (a herdr terminal id) and `reason`, so reusing it here is exactly "the existing // ExhaustionSink path", not a new mechanism. - ExhaustionSink exhaustionSink = (target, reason) -> sessions.roster().stream() - .filter(session -> target.equals(session.terminalId())) - .findFirst() - .map(MemberSession::profile) - .map(profileName -> config.get().profiles().get(profileName)) - .ifPresent(profile -> { - String credentialId = profile.effectiveCredentialId(); - quarantine.quarantine(credentialId); - log.warn("credential '{}' quarantined for {}s (profile '{}'): {}", credentialId, - cfg.quarantineCooldownSeconds(), profile.profile(), reason); - }); + // + // fleetd #234: that check fires from SessionAwareHandle.agentSessionId(), which runs during + // SessionManager.acquire() BEFORE this session is registered in sessions.roster() — so the + // roster-only lookup below used to find nothing, .ifPresent silently no-op'd, and the + // ERROR the check had just logged ("quarantining this profile's credential") was a lie: + // nothing was quarantined, and nothing said so. Two changes: (1) OpenCodeLauncher now + // passes its OWN profile name via ExhaustionSink's 3-arg overload — it already has the + // FleetConfig.Profile in hand and does not need the roster at all — used here as a + // fallback whenever the roster lookup misses; (2) if a profile still cannot be resolved + // (neither the roster nor the hint names a configured one), this logs loudly at ERROR + // instead of silently doing nothing — a control that cannot act must say so. + ExhaustionSink exhaustionSink = new ExhaustionSink() { + @Override + public void onExhausted(String target, String reason) { + onExhausted(target, reason, null); + } + + @Override + public void onExhausted(String target, String reason, String profileHint) { + String profileName = sessions.roster().stream() + .filter(session -> target.equals(session.terminalId())) + .findFirst() + .map(MemberSession::profile) + .orElse(profileHint); + FleetConfig.Profile profile = profileName == null ? null : config.get().profiles().get(profileName); + if (profile == null) { + log.error("quarantine requested for target '{}' ({}) but no profile could be " + + "resolved — the target is not (yet) in the roster, and {} — " + + "credential NOT quarantined (fleetd #234)", + target, reason, + profileHint == null ? "no profile hint was given" + : "the hinted profile '" + profileHint + "' is not configured"); + return; + } + String credentialId = profile.effectiveCredentialId(); + quarantine.quarantine(credentialId); + log.warn("credential '{}' quarantined for {}s (profile '{}'): {}", credentialId, + cfg.quarantineCooldownSeconds(), profile.profile(), reason); + } + }; // fleetd #175: point the forwarding sink handed to OpenCodeLauncher above at the real one, // now that `sessions` exists to resolve target -> session -> profile. exhaustionSinkRef.set(exhaustionSink); diff --git a/fleetd/src/main/java/dev/ltms/fleet/inject/ExhaustionSink.java b/fleetd/src/main/java/dev/ltms/fleet/inject/ExhaustionSink.java index e1a9220..711e57e 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/inject/ExhaustionSink.java +++ b/fleetd/src/main/java/dev/ltms/fleet/inject/ExhaustionSink.java @@ -20,6 +20,32 @@ public interface ExhaustionSink { */ void onExhausted(String target, String reason); + /** + * Same notification, plus a profile name the CALLER already knows — for a caller whose + * {@code target} is not yet resolvable through whatever roster the sink's implementation + * consults (fleetd #234). {@link dev.ltms.fleet.member.OpenCodeLauncher}'s model-mismatch check + * fires from {@code SessionAwareHandle.agentSessionId()}, which runs during {@code + * SessionManager.acquire()} before that session is registered — a target -> session -> + * profile lookup finds nothing at that point. That launcher already has its own {@code + * FleetConfig.Profile} in hand and does not need the roster to know which profile to + * quarantine, so it calls this overload instead of leaving the sink to guess. + * + *

Defaults to the two-arg overload, discarding {@code profile} — the correct behaviour for + * every caller that has not been updated to supply one: {@link + * dev.ltms.fleet.inject.CompletionResolver}'s two call sites always call {@code target} that + * IS live in the roster at the time of the call, so they need no hint and keep working exactly + * as before. A sink that wants to use the hint (see {@code Fleetd.main}'s wiring) overrides this + * method directly rather than relying on the default. + * + * @param target as {@link #onExhausted(String, String)} + * @param reason as {@link #onExhausted(String, String)} + * @param profile the profile the caller already knows should be quarantined, or {@code null} + * when the caller has no better answer than {@code target} alone + */ + default void onExhausted(String target, String reason, String profile) { + onExhausted(target, reason); + } + /** * Inert sink — nothing happens on exhaustion. The explicit stand-in a caller (or a test not * exercising this feature) passes instead of a defaulting overload, exactly like diff --git a/fleetd/src/main/java/dev/ltms/fleet/member/OpenCodeLauncher.java b/fleetd/src/main/java/dev/ltms/fleet/member/OpenCodeLauncher.java index 006da0d..e5aaefc 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/member/OpenCodeLauncher.java +++ b/fleetd/src/main/java/dev/ltms/fleet/member/OpenCodeLauncher.java @@ -22,6 +22,7 @@ import java.util.List; import java.util.Map; import java.util.Set; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicReference; import java.util.function.BooleanSupplier; import java.util.function.Function; import java.util.function.LongSupplier; @@ -706,6 +707,17 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher { private final ExhaustionSink exhaustionSink; /** CAS'd true the first (and only) time a model mismatch is reported for this handle. */ private final AtomicBoolean modelMismatchReported = new AtomicBoolean(); + /** + * The session id, once {@link OpenCodeSessionDiscovery#sessionIdForDirectory} first + * resolves a non-null answer for this handle (fleetd #234). Sticky on purpose: {@code + * directory} is a shared-cwd heuristic (see {@link OpenCodeSessionDiscovery}'s class + * javadoc) that can start returning a DIFFERENT row once another session shares the same + * directory and writes a newer one — re-deriving it on every call would let this handle's + * identity silently drift to a sibling's session. Once resolved, this IS the answer, and + * {@link #checkModelMatch} reads only the row this id names, never "whatever is newest in + * the directory right now." + */ + private final AtomicReference resolvedSessionId = new AtomicReference<>(); SessionAwareHandle(PeerHandle delegate, OpenCodeSessionDiscovery discovery, String cwd, FleetConfig.Profile cfg, @@ -759,17 +771,27 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher { } return null; } + // fleetd #234: once resolved, stay resolved. Re-deriving from `directory` on every call + // would let this handle's identity drift to a sibling session that later shares the + // same cwd and writes a newer row — see resolvedSessionId's javadoc. + String cached = resolvedSessionId.get(); + if (cached != null) { + return cached; + } // Lazy + retried, never a spawn-time blocker: opencode writes the session record only // when the session is first persisted, so null here is the correct interim answer and // the caller re-calls later (each call re-scans, picking up a record that has since // appeared). String id = discovery.sessionIdForDirectory(cwd); + if (id != null) { + resolvedSessionId.compareAndSet(null, id); + } // fleetd #175: check on the SAME tick — while the caller (SessionManager's late-resolve // step) is still re-polling because the id is unknown, the row this id came from (once // it exists) is exactly the row that also carries the actual model. Once id resolves, // the caller stops calling agentSessionId() for this session, so this is naturally a // once-only check that happens right when the row first appears. - checkModelMatch(); + checkModelMatch(id); return id; } @@ -777,16 +799,25 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher { * Verify the live opencode session is running the model {@link #cfg} requested (fleetd * #175) and, on a real mismatch, log an ERROR and quarantine through {@link * #exhaustionSink}. A no-op when there is nothing to compare against — no model configured, - * already reported once for this handle, or the actual model is still UNKNOWN (no row yet, - * unreadable database, or unparseable evidence). UNKNOWN must never be treated as a - * mismatch: that is the single most important safety rule here — a false positive would + * already reported once for this handle, {@code sessionId} itself is not resolved yet + * (fleetd #234: absent evidence, not a mismatch), or the actual model is still UNKNOWN (no + * row yet, unreadable database, or unparseable evidence). UNKNOWN must never be treated as + * a mismatch: that is the single most important safety rule here — a false positive would * quarantine a perfectly working profile's credential. + * + * @param sessionId the id {@link #agentSessionId()} just resolved (or had cached) for THIS + * handle — the model is read back for this exact session (fleetd #234's + * {@link OpenCodeSessionDiscovery#actualModelForSessionId}), never + * re-derived from {@code directory} */ - private void checkModelMatch() { + private void checkModelMatch(String sessionId) { if (modelMismatchReported.get() || cfg.model() == null || cfg.model().isBlank()) { return; } - OpenCodeSessionDiscovery.ActualModel actual = discovery.actualModelForDirectory(cwd); + if (sessionId == null || sessionId.isBlank()) { + return; // id not resolved yet — UNKNOWN, never a mismatch (fleetd #175's rule) + } + OpenCodeSessionDiscovery.ActualModel actual = discovery.actualModelForSessionId(sessionId); if (actual == null) { return; // UNKNOWN evidence — never a mismatch } @@ -817,10 +848,16 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher { + "falls back to a default model, which may be a PAID credential " + "(fleetd #175); quarantining this profile's credential", cfg.profile(), cfg.model(), actualDisplay); + // fleetd #234: pass our OWN profile name too. This check fires from agentSessionId(), + // called during SessionManager.acquire() BEFORE this session is registered in + // sessions.roster() — a roster-only sink (Fleetd's target -> session -> profile lookup) + // finds nothing at this point and silently no-ops (defect 2). We already know exactly + // which profile to quarantine without the roster; the sink is passed it explicitly. exhaustionSink.onExhausted(delegate.terminalId(), "opencode model mismatch: profile '" + cfg.profile() + "' requested '" + cfg.model() + "' but the live session is running '" + actualDisplay - + "' (fleetd #175)"); + + "' (fleetd #175)", + cfg.profile()); } @Override diff --git a/fleetd/src/main/java/dev/ltms/fleet/member/OpenCodeSessionDiscovery.java b/fleetd/src/main/java/dev/ltms/fleet/member/OpenCodeSessionDiscovery.java index 7989cb8..a64aec7 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/member/OpenCodeSessionDiscovery.java +++ b/fleetd/src/main/java/dev/ltms/fleet/member/OpenCodeSessionDiscovery.java @@ -32,11 +32,19 @@ import java.util.concurrent.atomic.AtomicBoolean; * here, so a layout change, or a switch to the HTTP server, changes exactly one class and nothing * in {@link OpenCodeLauncher}. * - *

The determinism that makes this useful is structural, not a guess: every fleetd worker runs - * in its own unique git worktree, so the row's {@code directory} (its project root) equals the - * worker's cwd identifies its session unambiguously. We match on {@code directory} rather - * than diffing {@code opencode session list} before/after — that races under concurrent spawns, and - * the CLI listing does not even show the directory. + *

The {@code directory} match is a heuristic, not an identity — fleetd #234. A + * worktree is opt-in: {@code fleet_spawn} only provisions one when the caller passes {@code + * worktree:}; the default spawn inherits the lead's own cwd, which every other worker spawned the + * same way (and every past session ever run there) shares. {@code directory} therefore does + * not identify a session unambiguously in general — only in the special case of a fresh, + * unique worktree does "most recently updated row for this directory" reliably mean "this worker's + * own row." {@link #sessionIdForDirectory} still has to use this heuristic (the id has to come from + * somewhere, and nothing else is available at this layer — see that method's javadoc), but a caller + * that already holds a resolved id must never re-derive evidence about that same session via + * {@code directory} again; see {@link #actualModelForSessionId}, which looks up by {@code id} + * instead for exactly this reason. We match on {@code directory} rather than diffing + * {@code opencode session list} before/after — that races under concurrent spawns, and the CLI + * listing does not even show the directory. * *

All reads are best-effort and never throw: a missing or unreadable database, a query that * fails, or a directory with no row yet all yield {@code null}, and the caller (the session @@ -89,8 +97,16 @@ final class OpenCodeSessionDiscovery { /** * The opencode session id whose row references {@code directory} (the worker's cwd), or * {@code null} when no row matches yet. When several rows share the directory — e.g. repeated - * spawns into the same worktree — the row with the highest {@code time_updated} wins: it is - * the session the pane most likely corresponds to. + * spawns into the same worktree, OR several workers sharing one cwd because none of them was + * given a worktree (fleetd #234) — the row with the highest {@code time_updated} wins: it is + * the session the pane most likely corresponds to. That "most likely" is a real caveat, not a + * formality: when the directory is shared, this can and does pick another session's row (see + * the class javadoc). This is the one place fleetd resolves an opencode session id at all — + * nothing else is available at this layer to disambiguate further (no {@code opencode session + * list} entry names the directory, and diffing before/after races under concurrent spawns) — so + * the heuristic stays here unchanged. What must never happen is a SECOND, independent piece of + * evidence about the same session being re-derived via {@code directory} once an id has already + * come out of this method; see {@link #actualModelForSessionId}. * *

Never throws: a missing {@code opencode.db}, a locked/unreadable database, a query * failure, or a directory that has not been persisted yet all resolve to {@code null} rather @@ -134,24 +150,36 @@ final class OpenCodeSessionDiscovery { } /** - * The model opencode actually ran the {@code directory}'s most-recent session on (fleetd - * #175), read from the same row {@link #sessionIdForDirectory} matches — but via its own - * query and its own connection, deliberately kept independent so a database whose schema - * predates the {@code model} column (or any other read failure on this column alone) can - * never take {@link #sessionIdForDirectory}'s id resolution down with it. That would be a - * regression of the id-resolution feature #209 shipped; this method degrades on its own. + * The model opencode actually ran the session {@code sessionId} on (fleetd #175/#234), read by + * primary-key lookup — the ONE row that id names, and no other. Deliberately keyed on + * {@code id} rather than {@code directory}: two independent {@code WHERE directory = ? ORDER BY + * time_updated DESC LIMIT 1} queries (one for the id, one for the model) can each pick a + * DIFFERENT row once more than one session shares a directory (fleetd #234 — the default + * no-worktree spawn shares the lead's cwd with every other worker and every past session ever + * run there), silently comparing a profile's requested model against a session that is not even + * the one whose id was returned. Keying on {@code id} instead makes that impossible: the model + * read back is always the SAME session {@link #sessionIdForDirectory} (or a cached copy of its + * answer) already resolved. + * + *

Kept as its own query and its own connection, independent from {@link + * #sessionIdForDirectory}: a database whose schema predates the {@code model} column (or any + * other read failure on this column alone) can never take id resolution down with it. That + * would be a regression of the id-resolution feature #209 shipped; this method degrades on its + * own. * *

Never throws, and every failure mode — no matching row, a missing/unreadable database, a * missing {@code model} column, a null/blank {@code model} value, or JSON that does not parse * into {@code {"id": "...", "providerID": "..."}} with a non-blank {@code id} — resolves to * {@code null}. That is UNKNOWN evidence, not a mismatch signal: the caller must never - * quarantine a profile on the strength of a {@code null} here. + * quarantine a profile on the strength of a {@code null} here. A blank/null {@code sessionId} + * (the id is not resolved yet) is UNKNOWN too, for the same reason — never call this with one. * - * @param directory the worker's cwd, as resolved for this spawn + * @param sessionId the session id already resolved by {@link #sessionIdForDirectory} for this + * spawn — never re-derived from {@code directory} here * @return the actual model, or {@code null} when unknown */ - ActualModel actualModelForDirectory(String directory) { - if (directory == null || directory.isBlank()) { + ActualModel actualModelForSessionId(String sessionId) { + if (sessionId == null || sessionId.isBlank()) { return null; } if (!Files.isRegularFile(databasePath)) { @@ -159,10 +187,10 @@ final class OpenCodeSessionDiscovery { // exact condition — do not double-log it here. return null; } - String sql = "SELECT model FROM session WHERE directory = ? ORDER BY time_updated DESC LIMIT 1"; + String sql = "SELECT model FROM session WHERE id = ?"; try (Connection connection = openReadOnly(); PreparedStatement statement = connection.prepareStatement(sql)) { - statement.setString(1, directory); + statement.setString(1, sessionId); try (ResultSet rows = statement.executeQuery()) { if (rows.next()) { return parseModel(rows.getString("model")); diff --git a/fleetd/src/test/java/dev/ltms/fleet/member/OpenCodeLauncherTest.java b/fleetd/src/test/java/dev/ltms/fleet/member/OpenCodeLauncherTest.java index 3cb9745..bdde1e5 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/member/OpenCodeLauncherTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/member/OpenCodeLauncherTest.java @@ -17,6 +17,7 @@ import dev.ltms.fleet.peer.MemberRole; import dev.ltms.fleet.peer.PeerHandle; import dev.ltms.fleet.peer.PeerUnreachableException; import dev.ltms.fleet.peer.SpawnRequest; +import dev.ltms.fleet.placement.BackendQuarantine; import dev.ltms.fleet.session.MemberSession; import dev.ltms.fleet.session.SessionManager; import org.junit.jupiter.api.Test; @@ -34,6 +35,7 @@ import java.util.Optional; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; import java.util.function.Supplier; import static org.junit.jupiter.api.Assertions.*; @@ -52,6 +54,19 @@ class OpenCodeLauncherTest { null, null, gitTokenEnv, null, FleetConfig.Profile.KIND_OPENCODE); } + /** + * A profile carrying an explicit {@code credentialId} (fleetd #234 defect 2) — distinct from + * the profile's own name, so a test can assert on the credential precisely rather than relying + * on {@code effectiveCredentialId()}'s profile-name fallback. + */ + private static FleetConfig.Profile opencodeCfgWithCredential(String profileName, String model, + String credentialId) { + return new FleetConfig.Profile(profileName, null, model, null, "FLEETD_WORKER_TOKEN", + List.of("opencode"), "tab", "fleetd-workers", "opencode: {model} #{n}", null, + null, List.of(), null, null, FleetConfig.Profile.KIND_OPENCODE, Map.of(), 1.0f, + null, false, null, credentialId, null); + } + /** Gate-disabled launcher whose per-spawn config dirs land under an inspectable temp root. */ private static OpenCodeLauncher service(FakeHerdr herdr, Path configRoot, FleetConfig.Profile cfg) { return new OpenCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr), @@ -1133,4 +1148,142 @@ class OpenCodeLauncherTest { assertEquals("ses_x", handle.agentSessionId()); assertTrue(exhausted.isEmpty(), "no model configured → nothing to compare: " + exhausted); } + + // --- fleetd #234 defect 1: key the model check on the RESOLVED id, not the shared directory -- + + /** + * The exact shape fleetd #234 reported: {@code fleet_spawn} with no {@code worktree:} shares + * the lead's cwd across every worker, so more than one session row can exist for the SAME + * {@code directory}. Once THIS handle's own session id is resolved, a sibling member spawned + * later into the same shared directory — writing a NEWER, unrelated row — must never make the + * already-resolved session look mismatched. Today's code re-derives "the newest row in this + * directory" on every call (both for the id AND, independently, for the model), so it would + * pick up the sibling's row on the second call and flag a false mismatch AND flip the returned + * id. The fix (fleetd #234) makes the id sticky once resolved and reads the model back for + * exactly that id (see {@link OpenCodeSessionDiscovery#actualModelForSessionId}) — never + * "whatever is newest in the directory right now." + */ + @Test + void modelCheckReadsTheResolvedSessionsOwnRowNotWhateverIsNewestInTheSharedDirectory( + @TempDir Path configRoot, @TempDir Path discRoot) throws Exception { + List exhausted = new ArrayList<>(); + ExhaustionSink sink = (target, reason) -> exhausted.add(reason); + FleetConfig.Profile cfg = opencodeCfg("openai/gpt-5.6-terra", null, null); + PeerHandle handle = serviceWithSink(new FakeHerdr(), configRoot, discRoot, cfg, sink) + .spawn(new SpawnRequest(null, "/work/dir", null)); + + // Our own session's row, correctly matching the profile's requested model. + OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "ses_ours", "/work/dir", 1000L, + "{\"id\":\"gpt-5.6-terra\",\"providerID\":\"openai\"}"); + assertEquals("ses_ours", handle.agentSessionId(), "resolves to our own session"); + assertTrue(exhausted.isEmpty(), "matching model → no mismatch on first resolve: " + exhausted); + + // A sibling member, spawned later into the SAME shared directory (no worktree, fleetd + // #234's default), writes a newer row running a totally different model. + OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "ses_sibling", "/work/dir", 9000L, + "{\"id\":\"deepseek-v4-flash\",\"providerID\":\"gx\"}"); + + assertEquals("ses_ours", handle.agentSessionId(), + "the session id, once resolved, must not flip to a sibling sharing the directory"); + assertTrue(exhausted.isEmpty(), + "a sibling's later, unrelated row in the same shared directory must never be read " + + "as OUR session's model: " + exhausted); + } + + // --- fleetd #234 defect 2: the quarantine the ERROR announces must actually happen ----------- + + /** + * fleetd #234: {@link OpenCodeLauncher.SessionAwareHandle#checkModelMatch} fires from {@code + * agentSessionId()}, which {@code SessionManager.acquire()} calls to build the very first + * {@code MemberSession} record — BEFORE that session is put into the registry {@code + * sessions.roster()} reads. A sink that resolves {@code target -> profile} ONLY through the + * roster (today's {@code Fleetd.java} code, before this fix) therefore finds nothing at this + * exact moment and silently does not quarantine, even though it just logged an ERROR saying it + * would. This test drives the REAL path — {@code SessionManager.acquire()} — not the sink + * directly, because the bug is entirely about this ordering; a direct-sink test cannot see it + * (and is exactly why #175's own test suite, which only ever called the sink directly or after + * registration, never caught this). + * + *

The sink under test mirrors {@code Fleetd.main()}'s real wiring after the fix: resolve + * via the roster first (unchanged for {@code CompletionResolver}'s two call sites), falling + * back to the profile hint {@link OpenCodeLauncher} now supplies via {@link + * ExhaustionSink#onExhausted(String, String, String)} when the roster lookup misses. + */ + @Test + void aSpawnTimeModelMismatchActuallyQuarantinesTheCredentialThroughTheRealAcquirePath( + @TempDir Path configRoot, @TempDir Path discRoot) throws Exception { + FleetConfig.Profile cfg = opencodeCfgWithCredential( + "terra", "opencode/nemotron-3-ultra-free", "openai-shared"); + Map profiles = Map.of(cfg.profile(), cfg); + BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.SECONDS.toNanos(1800)); + + ExhaustionSink sink = new ExhaustionSink() { + @Override + public void onExhausted(String target, String reason) { + onExhausted(target, reason, null); // no roster resolution modelled here — see below + } + + @Override + public void onExhausted(String target, String reason, String profileHint) { + FleetConfig.Profile profile = profileHint == null ? null : profiles.get(profileHint); + if (profile != null) { + quarantine.quarantine(profile.effectiveCredentialId()); + } + } + }; + + FakeHerdr herdr = new FakeHerdr(); + OpenCodeLauncher launcher = serviceWithSink(herdr, configRoot, discRoot, cfg, sink); + SessionManager sessions = new SessionManager(launcher); + + // The mismatching row exists BEFORE the spawn — reproducing fleetd #234's exact timing: + // opencode's session table already carries evidence by the moment acquire() first asks. + OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "ses_x", "/work/dir", 1000L, + "{\"id\":\"gpt-5.6-sol\",\"providerID\":\"openai\"}"); + + assertFalse(quarantine.isQuarantined("openai-shared"), "nothing quarantined before the spawn"); + + // The real production entrypoint: acquire() builds the MemberSession by calling + // handle.agentSessionId() BEFORE registry.put() runs. + MemberSession acquired = sessions.acquire(cfg.profile(), "/work/dir", null, null); + + assertEquals("ses_x", acquired.agentSessionId(), "the id itself still resolves correctly"); + assertTrue(quarantine.isQuarantined("openai-shared"), + "the mismatch fires DURING acquire(), before roster registration, and must still " + + "reach the quarantine via the profile hint — not silently no-op"); + } + + /** + * The other half of the same proof: a sink that resolves {@code target -> profile} ONLY + * through the roster (i.e. ignores the profile hint entirely, ~today's pre-fix {@code + * Fleetd.java}) drops the SAME spawn-time mismatch silently — the credential is never + * quarantined even though {@link OpenCodeLauncher} logged the mismatch ERROR. This is the + * failure fleetd #234 reported, reproduced through the real {@code SessionManager.acquire()} + * path rather than asserted by inspecting the fix. + */ + @Test + void aRosterOnlySinkSilentlyDropsTheSpawnTimeQuarantine(@TempDir Path configRoot, + @TempDir Path discRoot) throws Exception { + FleetConfig.Profile cfg = opencodeCfgWithCredential( + "terra", "opencode/nemotron-3-ultra-free", "openai-shared"); + BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.SECONDS.toNanos(1800)); + + // Deliberately ignores the profile hint — the pre-fix shape: only a roster lookup (modelled + // here as always empty, since acquire() has not registered the session yet either way). + ExhaustionSink rosterOnlySink = (target, reason) -> { }; + + FakeHerdr herdr = new FakeHerdr(); + OpenCodeLauncher launcher = serviceWithSink(herdr, configRoot, discRoot, cfg, rosterOnlySink); + SessionManager sessions = new SessionManager(launcher); + + OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "ses_x", "/work/dir", 1000L, + "{\"id\":\"gpt-5.6-sol\",\"providerID\":\"openai\"}"); + + MemberSession acquired = sessions.acquire(cfg.profile(), "/work/dir", null, null); + + assertEquals("ses_x", acquired.agentSessionId(), "the id itself still resolves correctly"); + assertFalse(quarantine.isQuarantined("openai-shared"), + "a roster-only sink cannot see this target yet — the quarantine silently never " + + "happens, which is exactly fleetd #234 defect 2"); + } } diff --git a/fleetd/src/test/java/dev/ltms/fleet/member/OpenCodeSessionDiscoveryTest.java b/fleetd/src/test/java/dev/ltms/fleet/member/OpenCodeSessionDiscoveryTest.java index c3add98..426369c 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/member/OpenCodeSessionDiscoveryTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/member/OpenCodeSessionDiscoveryTest.java @@ -146,14 +146,14 @@ class OpenCodeSessionDiscoveryTest { "an unreadable database resolves to null, not an exception"); } - // --- fleetd #175: actualModelForDirectory / the model JSON column --------------------------- + // --- fleetd #175/#234: actualModelForSessionId / the model JSON column ----------------------- @Test void parsesTheModelJsonIntoProviderAndId(@TempDir Path root) throws Exception { writeRecord(root, "ses_aaa", "/w/a", 1000L, "{\"id\":\"gpt-5.6-terra\",\"providerID\":\"openai\"}"); OpenCodeSessionDiscovery.ActualModel actual = - new OpenCodeSessionDiscovery(root).actualModelForDirectory("/w/a"); + new OpenCodeSessionDiscovery(root).actualModelForSessionId("ses_aaa"); assertNotNull(actual, "a well-formed model JSON parses"); assertEquals("gpt-5.6-terra", actual.id()); @@ -161,28 +161,39 @@ class OpenCodeSessionDiscoveryTest { } @Test - void prefersTheModelOfTheMostRecentlyUpdatedRow(@TempDir Path root) throws Exception { - writeRecord(root, "ses_old", "/w/a", 1000L, "{\"id\":\"old-model\",\"providerID\":\"openai\"}"); - writeRecord(root, "ses_new", "/w/a", 5000L, "{\"id\":\"new-model\",\"providerID\":\"openai\"}"); + void looksUpByIdEvenWhenAnotherRowInTheSameDirectoryIsNewer(@TempDir Path root) throws Exception { + writeRecord(root, "ses_ours", "/work/dir", 1000L, "{\"id\":\"gpt-5.6-terra\",\"providerID\":\"openai\"}"); + writeRecord(root, "ses_sibling", "/work/dir", 9000L, "{\"id\":\"deepseek-v4-flash\",\"providerID\":\"gx\"}"); - assertEquals("new-model", - new OpenCodeSessionDiscovery(root).actualModelForDirectory("/w/a").id(), - "the model of the row with the highest time_updated wins, same as the id"); + OpenCodeSessionDiscovery discovery = new OpenCodeSessionDiscovery(root); + assertEquals("gpt-5.6-terra", discovery.actualModelForSessionId("ses_ours").id(), + "querying by id reads OUR row, not the directory's newest row"); + assertEquals("deepseek-v4-flash", discovery.actualModelForSessionId("ses_sibling").id(), + "each id resolves to its own row independently of time_updated ordering"); } @Test - void aNonMatchingDirectoryYieldsUnknownModelRatherThanAMismatch(@TempDir Path root) throws Exception { + void anUnknownSessionIdYieldsUnknownModelRatherThanAMismatch(@TempDir Path root) throws Exception { writeRecord(root, "ses_aaa", "/w/a", 1000L, "{\"id\":\"x\",\"providerID\":\"y\"}"); - assertNull(new OpenCodeSessionDiscovery(root).actualModelForDirectory("/w/other"), - "no row for this cwd yet → unknown, not a wrong model"); + assertNull(new OpenCodeSessionDiscovery(root).actualModelForSessionId("ses_no_such_row"), + "no row for this id yet → unknown, not a wrong model"); + } + + @Test + void aBlankOrNullSessionIdYieldsUnknownModel(@TempDir Path root) throws Exception { + writeRecord(root, "ses_aaa", "/w/a", 1000L, "{\"id\":\"x\",\"providerID\":\"y\"}"); + + OpenCodeSessionDiscovery discovery = new OpenCodeSessionDiscovery(root); + assertNull(discovery.actualModelForSessionId(null)); + assertNull(discovery.actualModelForSessionId(" ")); } @Test void aNullModelColumnYieldsUnknownWithoutThrowing(@TempDir Path root) throws Exception { writeRecord(root, "ses_aaa", "/w/a", 1000L, null); - assertNull(new OpenCodeSessionDiscovery(root).actualModelForDirectory("/w/a"), + assertNull(new OpenCodeSessionDiscovery(root).actualModelForSessionId("ses_aaa"), "a row with no model value yet is unknown, not a mismatch"); } @@ -190,7 +201,7 @@ class OpenCodeSessionDiscoveryTest { void unparseableModelJsonYieldsUnknownWithoutThrowing(@TempDir Path root) throws Exception { writeRecord(root, "ses_aaa", "/w/a", 1000L, "this is not json"); - assertNull(new OpenCodeSessionDiscovery(root).actualModelForDirectory("/w/a"), + assertNull(new OpenCodeSessionDiscovery(root).actualModelForSessionId("ses_aaa"), "JSON that fails to parse resolves to unknown, never an exception"); } @@ -198,13 +209,13 @@ class OpenCodeSessionDiscoveryTest { void modelJsonMissingIdYieldsUnknown(@TempDir Path root) throws Exception { writeRecord(root, "ses_aaa", "/w/a", 1000L, "{\"providerID\":\"openai\"}"); - assertNull(new OpenCodeSessionDiscovery(root).actualModelForDirectory("/w/a"), + assertNull(new OpenCodeSessionDiscovery(root).actualModelForSessionId("ses_aaa"), "no id in the JSON → unknown, since id is what a caller actually compares"); } @Test void aMissingDatabaseYieldsUnknownModelWithoutThrowing(@TempDir Path root) { - assertNull(new OpenCodeSessionDiscovery(root).actualModelForDirectory("/w/a")); + assertNull(new OpenCodeSessionDiscovery(root).actualModelForSessionId("ses_aaa")); } /** @@ -212,7 +223,7 @@ class OpenCodeSessionDiscoveryTest { * shape a real opencode upgrade/downgrade could produce. This must degrade to UNKNOWN for the * model, and — the property that actually matters — must NOT take id resolution down with it. * A combined single query for both columns would fail this test; that is why - * {@link OpenCodeSessionDiscovery#actualModelForDirectory} runs its own independent query. + * {@link OpenCodeSessionDiscovery#actualModelForSessionId} runs its own independent query. */ @Test void aMissingModelColumnYieldsUnknownButIdResolutionStillWorks(@TempDir Path root) throws Exception { @@ -234,7 +245,7 @@ class OpenCodeSessionDiscoveryTest { OpenCodeSessionDiscovery discovery = new OpenCodeSessionDiscovery(root); assertEquals("ses_aaa", discovery.sessionIdForDirectory("/w/a"), "id resolution must survive a database with no model column at all"); - assertNull(discovery.actualModelForDirectory("/w/a"), + assertNull(discovery.actualModelForSessionId("ses_aaa"), "no model column → unknown, not a throw and not a mismatch"); } } From c325054242201e3d7d93b5c4de91d543c13c9dec Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 3 Sep 2026 10:22:31 +0700 Subject: [PATCH 2/4] fleetd #234 round 2: fix the ExhaustionSink forwarding hop Fleetd.java actually uses The round-1 fix was dead on the real production path. Fleetd.java:177 builds a forwarding sink (needed because the adapters are constructed before `sessions` exists, breaking a genuine cycle) as a LAMBDA: ExhaustionSink forwardingExhaustionSink = (target, reason) -> exhaustionSinkRef.get().onExhausted(target, reason); A lambda can only implement the interface's one abstract method (the 2-arg overload), so it silently inherited the 3-arg overload's default body, which drops the profile hint and calls back into the 2-arg method. OpenCodeLauncher is constructed with this forwarder, so the hint it supplies (its own already-known profile name) was thrown away before it ever reached the real sink built later in Fleetd.main -- reproducing the exact silent no-op round 1 was sent to fix. The 1127 tests from round 1 all injected a sink directly into OpenCodeLauncher and never went through this forwarding hop, so none of them could see it. Fix: forwardingExhaustionSink is now an anonymous class overriding both overloads, each delegating to whatever exhaustionSinkRef currently holds. Audited every other ExhaustionSink value in main/: the only other one is ExhaustionSink.none() (a lambda), which is safe regardless of arity since both its 2-arg body and the inherited 3-arg default are true no-ops. New tests: - ExhaustionSinkForwardingHazardTest: isolates the hazard at the interface level (a lambda forwarder drops the hint; an anonymous-class forwarder does not), independent of Fleetd.java's specific wiring. - OpenCodeLauncherTest#theSpawnTimeQuarantineSurvivesTheFleetdStyleForwardingHop: replicates Fleetd.java's actual construction order (forwarder built and handed to the launcher first, real sink built and pointed at via the AtomicReference afterward) and drives the quarantine through it via the real SessionManager.acquire() path. Both proven by mutation: temporarily rewriting each fixed forwarder back into the pre-fix lambda makes its test fail with a real assertion message (both matched exactly: "expected: but was: " for the interface proof, "expected: but was: " for the composed-wiring test); restoring makes it pass again. No reverts were committed. mvn clean install: Tests run: 1130, Failures: 0, Errors: 0, Skipped: 0, BUILD SUCCESS. --- .../src/main/java/dev/ltms/fleet/Fleetd.java | 20 +++- .../ExhaustionSinkForwardingHazardTest.java | 96 +++++++++++++++++++ .../fleet/member/OpenCodeLauncherTest.java | 77 +++++++++++++++ 3 files changed, 192 insertions(+), 1 deletion(-) create mode 100644 fleetd/src/test/java/dev/ltms/fleet/inject/ExhaustionSinkForwardingHazardTest.java diff --git a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java index a8eed1b..5580446 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java +++ b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java @@ -174,7 +174,25 @@ public final class Fleetd { // genuine cycle. Break it exactly like liveCountRef below: a forwarding sink built now, // pointed at the real one once it exists. AtomicReference exhaustionSinkRef = new AtomicReference<>(ExhaustionSink.none()); - ExhaustionSink forwardingExhaustionSink = (target, reason) -> exhaustionSinkRef.get().onExhausted(target, reason); + // fleetd #234: this MUST be an anonymous class, not a lambda. A lambda can only implement + // the interface's single abstract method (the 2-arg overload) — it then inherits the 3-arg + // method's DEFAULT body, which drops the profile hint and calls back into the 2-arg + // overload here, so OpenCodeLauncher's hint (its own already-known profile name, passed + // for exactly the reason explained below at the real sink's construction) never reaches + // the real sink at all. That silently reproduced the very bug this hint exists to fix: a + // lambda here means the roster-miss path always fires, even after the hint was supplied. + // Both overloads must forward explicitly to whatever exhaustionSinkRef currently holds. + ExhaustionSink forwardingExhaustionSink = new ExhaustionSink() { + @Override + public void onExhausted(String target, String reason) { + exhaustionSinkRef.get().onExhausted(target, reason); + } + + @Override + public void onExhausted(String target, String reason, String profile) { + exhaustionSinkRef.get().onExhausted(target, reason, profile); + } + }; // The claude-code adapter is the always-present default; keep it even with no profiles (so a // bridge configured with no workers, or opencode-only, still has a well-defined base adapter) // unless opencode is the only kind configured. diff --git a/fleetd/src/test/java/dev/ltms/fleet/inject/ExhaustionSinkForwardingHazardTest.java b/fleetd/src/test/java/dev/ltms/fleet/inject/ExhaustionSinkForwardingHazardTest.java new file mode 100644 index 0000000..c788234 --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/inject/ExhaustionSinkForwardingHazardTest.java @@ -0,0 +1,96 @@ +package dev.ltms.fleet.inject; + +import org.junit.jupiter.api.Test; + +import java.util.concurrent.atomic.AtomicReference; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNull; + +/** + * fleetd #234 follow-up: a lambda implementing {@link ExhaustionSink} can only ever implement the + * interface's single abstract method — the 2-arg {@link ExhaustionSink#onExhausted(String, String)} + * — so it silently inherits the 3-arg overload's {@code default} body, which drops whatever profile + * hint a caller supplied and calls back into the 2-arg method instead. {@code Fleetd.main} built + * exactly this shape at {@code Fleetd.java:177} — a forwarding sink standing in for the real one + * until {@code sessions} exists (a genuine construction-order cycle) — as a lambda, so every hint + * {@link dev.ltms.fleet.member.OpenCodeLauncher} passed through it was thrown away before it ever + * reached the real sink. The fix (defect 2, round 2) reported ERROR and correctly declined to + * quarantine in that situation — which is exactly why the earlier {@code OpenCodeLauncherTest} + * mutation tests, which inject a sink directly into the launcher and never go through this + * forwarding hop, could not see the bug: the hop itself was the defect. + * + *

This test pins the hazard at the interface level, independent of {@code Fleetd.java}'s + * specific wiring: ANY forwarding sink standing in front of another {@link ExhaustionSink} must be + * an anonymous class (or otherwise override both overloads) — a lambda there is a silent regression + * of this exact bug. See {@code OpenCodeLauncherTest#theSpawnTimeQuarantineSurvivesTheFleetdStyleForwardingHop} + * for the composed, production-shaped reproduction (launcher mismatch → forwarder → real sink → + * quarantine). + */ +class ExhaustionSinkForwardingHazardTest { + + @Test + void aLambdaForwarderDropsTheProfileHintBeforeItReachesTheRealSink() { + AtomicReference hintSeenByRealSink = new AtomicReference<>("NEVER CALLED"); + ExhaustionSink real = new ExhaustionSink() { + @Override + public void onExhausted(String target, String reason) { + onExhausted(target, reason, null); + } + + @Override + public void onExhausted(String target, String reason, String profileHint) { + hintSeenByRealSink.set(profileHint); + } + }; + + AtomicReference ref = new AtomicReference<>(ExhaustionSink.none()); + // The broken shape: a lambda can only implement the 2-arg method. + ExhaustionSink brokenForwarder = (target, reason) -> ref.get().onExhausted(target, reason); + ref.set(real); + + brokenForwarder.onExhausted("term_x", "model mismatch", "gx"); + + assertNull(hintSeenByRealSink.get(), + "documents the hazard: a lambda forwarder can only implement the 2-arg overload, so " + + "it forwards through that overload alone — the real sink's 3-arg method " + + "still runs (via its own default), but with the hint already discarded"); + } + + @Test + void profileHintSurvivesTheForwardingHopUsedInProduction() { + AtomicReference hintSeenByRealSink = new AtomicReference<>("NEVER CALLED"); + ExhaustionSink real = new ExhaustionSink() { + @Override + public void onExhausted(String target, String reason) { + onExhausted(target, reason, null); + } + + @Override + public void onExhausted(String target, String reason, String profileHint) { + hintSeenByRealSink.set(profileHint); + } + }; + + AtomicReference ref = new AtomicReference<>(ExhaustionSink.none()); + // The fixed shape (fleetd #234, matching Fleetd.java:177): an anonymous class overriding + // BOTH overloads, each forwarding to the currently-held sink. + ExhaustionSink forwarding = new ExhaustionSink() { + @Override + public void onExhausted(String target, String reason) { + ref.get().onExhausted(target, reason); + } + + @Override + public void onExhausted(String target, String reason, String profile) { + ref.get().onExhausted(target, reason, profile); + } + }; + ref.set(real); + + forwarding.onExhausted("term_x", "model mismatch", "gx"); + + assertEquals("gx", hintSeenByRealSink.get(), + "the profile hint must reach the real sink through the forwarding hop"); + } +} diff --git a/fleetd/src/test/java/dev/ltms/fleet/member/OpenCodeLauncherTest.java b/fleetd/src/test/java/dev/ltms/fleet/member/OpenCodeLauncherTest.java index bdde1e5..bb6a563 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/member/OpenCodeLauncherTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/member/OpenCodeLauncherTest.java @@ -1286,4 +1286,81 @@ class OpenCodeLauncherTest { "a roster-only sink cannot see this target yet — the quarantine silently never " + "happens, which is exactly fleetd #234 defect 2"); } + + /** + * fleetd #234, round 2: the two tests above inject a sink DIRECTLY into {@link OpenCodeLauncher}, + * which is not what {@code Fleetd.main} actually does. Production has one more hop: {@code + * Fleetd.java} builds the adapters (including {@code OpenCodeLauncher}) before {@code sessions} + * exists — a genuine construction-order cycle — so it hands the launcher a forwarding + * sink pointed at an {@code AtomicReference}, and only later {@code .set(...)}s + * that reference to the real sink once {@code sessions} is built. That forwarding sink was + * written as a lambda ({@code Fleetd.java:177}), which can only implement the 2-arg overload — + * it silently inherited the 3-arg method's default body, dropping {@link OpenCodeLauncher}'s + * profile hint on every call, so the roster-miss branch fired even after the hint fix above + * shipped. This test reproduces that exact shape (construct with the forwarder, {@code .set()} + * the real sink afterward, same order as {@code Fleetd.main}) and would fail if the forwarder + * were ever written as a lambda again — see the class javadoc on {@code + * ExhaustionSinkForwardingHazardTest} for the isolated proof of the underlying mechanism. + */ + @Test + void theSpawnTimeQuarantineSurvivesTheFleetdStyleForwardingHop(@TempDir Path configRoot, + @TempDir Path discRoot) throws Exception { + FleetConfig.Profile cfg = opencodeCfgWithCredential( + "terra", "opencode/nemotron-3-ultra-free", "openai-shared"); + Map profiles = Map.of(cfg.profile(), cfg); + BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.SECONDS.toNanos(1800)); + + // Fleetd.java:176 — the forwarding sink is built BEFORE the real one can exist, and the + // launcher below is constructed against this forwarder, exactly like Fleetd.main. + java.util.concurrent.atomic.AtomicReference exhaustionSinkRef = + new java.util.concurrent.atomic.AtomicReference<>(ExhaustionSink.none()); + // Fleetd.java:177 fixed shape — an anonymous class forwarding BOTH overloads. Rewriting + // this as `(target, reason) -> exhaustionSinkRef.get().onExhausted(target, reason)` is + // exactly the regression this test exists to catch. + ExhaustionSink forwardingExhaustionSink = new ExhaustionSink() { + @Override + public void onExhausted(String target, String reason) { + exhaustionSinkRef.get().onExhausted(target, reason); + } + + @Override + public void onExhausted(String target, String reason, String profile) { + exhaustionSinkRef.get().onExhausted(target, reason, profile); + } + }; + + FakeHerdr herdr = new FakeHerdr(); + OpenCodeLauncher launcher = serviceWithSink(herdr, configRoot, discRoot, cfg, forwardingExhaustionSink); + SessionManager sessions = new SessionManager(launcher); + + // Fleetd.java:359-390 — the real sink is only built and pointed to AFTER `sessions` exists, + // same order as production. + ExhaustionSink realSink = new ExhaustionSink() { + @Override + public void onExhausted(String target, String reason) { + onExhausted(target, reason, null); // no roster resolution modelled — mirrors a miss + } + + @Override + public void onExhausted(String target, String reason, String profileHint) { + FleetConfig.Profile profile = profileHint == null ? null : profiles.get(profileHint); + if (profile != null) { + quarantine.quarantine(profile.effectiveCredentialId()); + } + } + }; + exhaustionSinkRef.set(realSink); + + OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "ses_x", "/work/dir", 1000L, + "{\"id\":\"gpt-5.6-sol\",\"providerID\":\"openai\"}"); + + assertFalse(quarantine.isQuarantined("openai-shared"), "nothing quarantined before the spawn"); + + MemberSession acquired = sessions.acquire(cfg.profile(), "/work/dir", null, null); + + assertEquals("ses_x", acquired.agentSessionId(), "the id itself still resolves correctly"); + assertTrue(quarantine.isQuarantined("openai-shared"), + "the profile hint must survive the Fleetd-style forwarding hop and reach the real " + + "sink — a lambda forwarder drops it and this must go red"); + } } From c935b181dd69a6fa662b53f8852c9f17b4b4ad0f Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 3 Sep 2026 10:30:37 +0700 Subject: [PATCH 3/4] fleetd #234 round 3: make ExhaustionSink's forwarder a shared factory, not a rebuilt-per-caller shape Round 2's tests never reached Fleetd.java at all: both new tests declared their OWN local copy of the forwarding shape instead of calling production's. Mutating Fleetd.java's real forwarder back into the broken lambda left those copies untouched, so the whole suite stayed green while production had regressed to exactly the bug being fixed -- proven live by the reviewer. Fix: extracted the forwarding shape into one named factory, ExhaustionSink.forwardingTo(Supplier target), with the "why a lambda here is wrong" explanation moved onto it (the one place the shape is now written). Fleetd.java's forwarder collapses to one line: ExhaustionSink forwardingExhaustionSink = ExhaustionSink.forwardingTo(exhaustionSinkRef::get); Both new tests now call this same factory instead of rebuilding an anonymous class inline, so they exercise the identical object production builds: - ExhaustionSinkForwardingHazardTest: calls ExhaustionSink.forwardingTo directly and asserts the hint reaches the real sink through it. - OpenCodeLauncherTest#theSpawnTimeQuarantineSurvivesTheFleetdStyleForwardingHop: same factory call, inside the full Fleetd-shaped construction order (forwarder built first, real sink pointed at via the AtomicReference afterward), driven through the real SessionManager.acquire() path. Mutation proof, this time on production code only: deleted the factory's 3-arg override (falls back to the interface default, dropping the hint) -- both new tests go red with no test file touched: ExhaustionSinkForwardingHazardTest...: expected: but was: OpenCodeLauncherTest...ForwardingHop: expected: but was: Tests run: 68, Failures: 2 Restored, re-ran: green (Tests run: 68, Failures: 0). Confirmed Fleetd.java carries no lambda ExhaustionSink anywhere (grep). ExhaustionSink.none() stays a lambda on purpose -- both its overloads are true no-ops regardless of arity, so there is no hint to drop. mvn clean install: Tests run: 1129, Failures: 0, Errors: 0, Skipped: 0, BUILD SUCCESS. --- .../src/main/java/dev/ltms/fleet/Fleetd.java | 25 ++----- .../dev/ltms/fleet/inject/ExhaustionSink.java | 47 +++++++++++- .../ExhaustionSinkForwardingHazardTest.java | 75 ++++--------------- .../fleet/member/OpenCodeLauncherTest.java | 41 +++++----- 4 files changed, 83 insertions(+), 105 deletions(-) diff --git a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java index 5580446..a3fbe44 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java +++ b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java @@ -174,25 +174,12 @@ public final class Fleetd { // genuine cycle. Break it exactly like liveCountRef below: a forwarding sink built now, // pointed at the real one once it exists. AtomicReference exhaustionSinkRef = new AtomicReference<>(ExhaustionSink.none()); - // fleetd #234: this MUST be an anonymous class, not a lambda. A lambda can only implement - // the interface's single abstract method (the 2-arg overload) — it then inherits the 3-arg - // method's DEFAULT body, which drops the profile hint and calls back into the 2-arg - // overload here, so OpenCodeLauncher's hint (its own already-known profile name, passed - // for exactly the reason explained below at the real sink's construction) never reaches - // the real sink at all. That silently reproduced the very bug this hint exists to fix: a - // lambda here means the roster-miss path always fires, even after the hint was supplied. - // Both overloads must forward explicitly to whatever exhaustionSinkRef currently holds. - ExhaustionSink forwardingExhaustionSink = new ExhaustionSink() { - @Override - public void onExhausted(String target, String reason) { - exhaustionSinkRef.get().onExhausted(target, reason); - } - - @Override - public void onExhausted(String target, String reason, String profile) { - exhaustionSinkRef.get().onExhausted(target, reason, profile); - } - }; + // fleetd #234, round 3: this MUST go through ExhaustionSink.forwardingTo, never a lambda or + // a hand-written anonymous class here — see that factory's javadoc for why a lambda at this + // call site silently drops OpenCodeLauncher's profile hint. Routing through the shared + // factory also lets a test call the exact same object this line builds, instead of + // asserting a copy of its shape. + ExhaustionSink forwardingExhaustionSink = ExhaustionSink.forwardingTo(exhaustionSinkRef::get); // The claude-code adapter is the always-present default; keep it even with no profiles (so a // bridge configured with no workers, or opencode-only, still has a well-defined base adapter) // unless opencode is the only kind configured. diff --git a/fleetd/src/main/java/dev/ltms/fleet/inject/ExhaustionSink.java b/fleetd/src/main/java/dev/ltms/fleet/inject/ExhaustionSink.java index 711e57e..5a61d11 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/inject/ExhaustionSink.java +++ b/fleetd/src/main/java/dev/ltms/fleet/inject/ExhaustionSink.java @@ -1,5 +1,7 @@ package dev.ltms.fleet.inject; +import java.util.function.Supplier; + /** * Notified when {@link CompletionResolver} actually delivers a {@code BACKEND_EXHAUSTED} * classification to a waiting send (CB-578 stage B) — never on a race that lost (see @@ -49,9 +51,52 @@ public interface ExhaustionSink { /** * Inert sink — nothing happens on exhaustion. The explicit stand-in a caller (or a test not * exercising this feature) passes instead of a defaulting overload, exactly like - * {@link ExhaustedPatternLookup#none()}. + * {@link ExhaustedPatternLookup#none()}. Safe as a lambda regardless of arity: its 2-arg body + * is empty, and the inherited 3-arg {@code default} just calls that same empty body — there is + * no hint to drop because this sink does nothing either way. */ static ExhaustionSink none() { return (target, reason) -> { }; } + + /** + * A sink that forwards BOTH overloads to whatever {@code target} currently supplies (fleetd + * #234, round 3). Exists to break a genuine construction-order cycle: {@code Fleetd.main} + * builds its adapters (including {@link dev.ltms.fleet.member.OpenCodeLauncher}) before + * {@code sessions} exists, so it cannot hand them the real sink yet — it hands them a + * forwarder pointed at an {@code AtomicReference} that starts at {@link #none()} + * and gets {@code .set()} to the real sink once {@code sessions} is built. {@code target} is + * evaluated on every call, never cached, so the forwarder keeps working after the reference is + * repointed. + * + *

This is the one place that forwarding shape is written, on purpose. A + * forwarder written directly at a call site as a lambda — + * {@code (t, r) -> target.get().onExhausted(t, r)} — implements only the interface's single + * abstract method (the 2-arg overload) and silently inherits the 3-arg overload's {@code + * default} body, which discards whatever profile hint the real caller supplied and calls back + * into the 2-arg method instead. {@link dev.ltms.fleet.member.OpenCodeLauncher}'s model-mismatch + * check relies on that hint reaching the real sink (it fires before its session is registered + * in the roster, so the roster alone cannot resolve which profile to quarantine) — a lambda + * forwarder here silently reproduces that exact bug. Routing every forwarder through this one + * factory, instead of re-writing the shape at each call site, means a test can pin the shape + * once and have both {@code Fleetd.java} and the test call the identical object — see {@code + * ExhaustionSinkForwardingHazardTest} and {@code OpenCodeLauncherTest} for the coverage this + * makes possible. + * + * @param target supplies the sink to forward to, evaluated fresh on every call + * @return a sink whose every overload delegates to {@code target.get()}'s matching overload + */ + static ExhaustionSink forwardingTo(Supplier target) { + return new ExhaustionSink() { + @Override + public void onExhausted(String t, String r) { + target.get().onExhausted(t, r); + } + + @Override + public void onExhausted(String t, String r, String p) { + target.get().onExhausted(t, r, p); + } + }; + } } diff --git a/fleetd/src/test/java/dev/ltms/fleet/inject/ExhaustionSinkForwardingHazardTest.java b/fleetd/src/test/java/dev/ltms/fleet/inject/ExhaustionSinkForwardingHazardTest.java index c788234..6c42381 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/inject/ExhaustionSinkForwardingHazardTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/inject/ExhaustionSinkForwardingHazardTest.java @@ -5,32 +5,21 @@ import org.junit.jupiter.api.Test; import java.util.concurrent.atomic.AtomicReference; import static org.junit.jupiter.api.Assertions.assertEquals; -import static org.junit.jupiter.api.Assertions.assertNull; /** - * fleetd #234 follow-up: a lambda implementing {@link ExhaustionSink} can only ever implement the - * interface's single abstract method — the 2-arg {@link ExhaustionSink#onExhausted(String, String)} - * — so it silently inherits the 3-arg overload's {@code default} body, which drops whatever profile - * hint a caller supplied and calls back into the 2-arg method instead. {@code Fleetd.main} built - * exactly this shape at {@code Fleetd.java:177} — a forwarding sink standing in for the real one - * until {@code sessions} exists (a genuine construction-order cycle) — as a lambda, so every hint - * {@link dev.ltms.fleet.member.OpenCodeLauncher} passed through it was thrown away before it ever - * reached the real sink. The fix (defect 2, round 2) reported ERROR and correctly declined to - * quarantine in that situation — which is exactly why the earlier {@code OpenCodeLauncherTest} - * mutation tests, which inject a sink directly into the launcher and never go through this - * forwarding hop, could not see the bug: the hop itself was the defect. - * - *

This test pins the hazard at the interface level, independent of {@code Fleetd.java}'s - * specific wiring: ANY forwarding sink standing in front of another {@link ExhaustionSink} must be - * an anonymous class (or otherwise override both overloads) — a lambda there is a silent regression - * of this exact bug. See {@code OpenCodeLauncherTest#theSpawnTimeQuarantineSurvivesTheFleetdStyleForwardingHop} - * for the composed, production-shaped reproduction (launcher mismatch → forwarder → real sink → - * quarantine). + * fleetd #234, round 3: calls the REAL production factory, {@link ExhaustionSink#forwardingTo}, + * rather than rebuilding its shape locally. A round-2 version of this test built its own copy of + * the forwarder inline — mutating {@code Fleetd.java}'s actual forwarder back into a lambda left + * that copy untouched, so the test kept passing while production had regressed to exactly the bug + * it was meant to catch. Calling the shared factory here means the same object under test IS the + * object {@code Fleetd.java} builds via {@code ExhaustionSink.forwardingTo(exhaustionSinkRef::get)} + * — a defect in the factory body, or a call site reverting to a hand-written lambda instead of the + * factory, has exactly one place left to hide, and this test reaches into it. */ class ExhaustionSinkForwardingHazardTest { @Test - void aLambdaForwarderDropsTheProfileHintBeforeItReachesTheRealSink() { + void forwardingToDeliversTheProfileHintToWhateverSinkTheSupplierCurrentlyReturns() { AtomicReference hintSeenByRealSink = new AtomicReference<>("NEVER CALLED"); ExhaustionSink real = new ExhaustionSink() { @Override @@ -44,53 +33,15 @@ class ExhaustionSinkForwardingHazardTest { } }; + // Same shape Fleetd.java uses: a reference that starts at none() and is repointed later. AtomicReference ref = new AtomicReference<>(ExhaustionSink.none()); - // The broken shape: a lambda can only implement the 2-arg method. - ExhaustionSink brokenForwarder = (target, reason) -> ref.get().onExhausted(target, reason); - ref.set(real); - - brokenForwarder.onExhausted("term_x", "model mismatch", "gx"); - - assertNull(hintSeenByRealSink.get(), - "documents the hazard: a lambda forwarder can only implement the 2-arg overload, so " - + "it forwards through that overload alone — the real sink's 3-arg method " - + "still runs (via its own default), but with the hint already discarded"); - } - - @Test - void profileHintSurvivesTheForwardingHopUsedInProduction() { - AtomicReference hintSeenByRealSink = new AtomicReference<>("NEVER CALLED"); - ExhaustionSink real = new ExhaustionSink() { - @Override - public void onExhausted(String target, String reason) { - onExhausted(target, reason, null); - } - - @Override - public void onExhausted(String target, String reason, String profileHint) { - hintSeenByRealSink.set(profileHint); - } - }; - - AtomicReference ref = new AtomicReference<>(ExhaustionSink.none()); - // The fixed shape (fleetd #234, matching Fleetd.java:177): an anonymous class overriding - // BOTH overloads, each forwarding to the currently-held sink. - ExhaustionSink forwarding = new ExhaustionSink() { - @Override - public void onExhausted(String target, String reason) { - ref.get().onExhausted(target, reason); - } - - @Override - public void onExhausted(String target, String reason, String profile) { - ref.get().onExhausted(target, reason, profile); - } - }; + ExhaustionSink forwarding = ExhaustionSink.forwardingTo(ref::get); ref.set(real); forwarding.onExhausted("term_x", "model mismatch", "gx"); assertEquals("gx", hintSeenByRealSink.get(), - "the profile hint must reach the real sink through the forwarding hop"); + "ExhaustionSink.forwardingTo must deliver the profile hint to the sink the supplier " + + "currently returns — the exact object Fleetd.java's forwarder is built from"); } } diff --git a/fleetd/src/test/java/dev/ltms/fleet/member/OpenCodeLauncherTest.java b/fleetd/src/test/java/dev/ltms/fleet/member/OpenCodeLauncherTest.java index bb6a563..ee67781 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/member/OpenCodeLauncherTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/member/OpenCodeLauncherTest.java @@ -1293,14 +1293,16 @@ class OpenCodeLauncherTest { * Fleetd.java} builds the adapters (including {@code OpenCodeLauncher}) before {@code sessions} * exists — a genuine construction-order cycle — so it hands the launcher a forwarding * sink pointed at an {@code AtomicReference}, and only later {@code .set(...)}s - * that reference to the real sink once {@code sessions} is built. That forwarding sink was - * written as a lambda ({@code Fleetd.java:177}), which can only implement the 2-arg overload — - * it silently inherited the 3-arg method's default body, dropping {@link OpenCodeLauncher}'s - * profile hint on every call, so the roster-miss branch fired even after the hint fix above - * shipped. This test reproduces that exact shape (construct with the forwarder, {@code .set()} - * the real sink afterward, same order as {@code Fleetd.main}) and would fail if the forwarder - * were ever written as a lambda again — see the class javadoc on {@code - * ExhaustionSinkForwardingHazardTest} for the isolated proof of the underlying mechanism. + * that reference to the real sink once {@code sessions} is built. + * + *

Round 3: a round-2 version of this test built its own copy of that + * forwarder's shape as an anonymous class. Mutating {@code Fleetd.java}'s ACTUAL forwarder back + * into a broken lambda left this test's own copy untouched, so it kept passing while production + * had regressed to the exact bug being fixed. This version instead calls {@link + * ExhaustionSink#forwardingTo}, the same factory {@code Fleetd.java} calls — the identical + * object, not a rebuilt copy of its shape — so a regression at either the {@code Fleetd.java} + * call site or inside {@code forwardingTo} itself has nowhere left to hide. See {@code + * ExhaustionSinkForwardingHazardTest} for the same factory exercised in isolation. */ @Test void theSpawnTimeQuarantineSurvivesTheFleetdStyleForwardingHop(@TempDir Path configRoot, @@ -1311,23 +1313,16 @@ class OpenCodeLauncherTest { BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.SECONDS.toNanos(1800)); // Fleetd.java:176 — the forwarding sink is built BEFORE the real one can exist, and the - // launcher below is constructed against this forwarder, exactly like Fleetd.main. + // launcher below is constructed against this forwarder, exactly like Fleetd.main. Calling + // the SAME factory Fleetd.java calls — ExhaustionSink.forwardingTo — rather than rebuilding + // the forwarder's shape here is the whole point (fleetd #234, round 3): a round-2 version of + // this test built its own copy, so mutating Fleetd.java's real forwarder back into a lambda + // left this test untouched. Routed through the shared factory, a regression at either the + // Fleetd.java call site (reverting to a hand-written lambda) or inside the factory body + // itself has nowhere left to hide from this test. java.util.concurrent.atomic.AtomicReference exhaustionSinkRef = new java.util.concurrent.atomic.AtomicReference<>(ExhaustionSink.none()); - // Fleetd.java:177 fixed shape — an anonymous class forwarding BOTH overloads. Rewriting - // this as `(target, reason) -> exhaustionSinkRef.get().onExhausted(target, reason)` is - // exactly the regression this test exists to catch. - ExhaustionSink forwardingExhaustionSink = new ExhaustionSink() { - @Override - public void onExhausted(String target, String reason) { - exhaustionSinkRef.get().onExhausted(target, reason); - } - - @Override - public void onExhausted(String target, String reason, String profile) { - exhaustionSinkRef.get().onExhausted(target, reason, profile); - } - }; + ExhaustionSink forwardingExhaustionSink = ExhaustionSink.forwardingTo(exhaustionSinkRef::get); FakeHerdr herdr = new FakeHerdr(); OpenCodeLauncher launcher = serviceWithSink(herdr, configRoot, discRoot, cfg, forwardingExhaustionSink); From 31b028e860882defaab05b82592242ca980a02ca Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 3 Sep 2026 10:46:27 +0700 Subject: [PATCH 4/4] fleetd #234 round 4: invert ExhaustionSink's abstract method so the bug class is unrepresentable Round 3's factory fixed the two known call sites but the underlying shape was still there: a lambda written against ExhaustionSink binds to whichever overload is abstract, and the 2-arg form held that position, so ANY lambda -- a call-site forwarder, a hand-built test double, a future caller who has never heard of fleetd #234 -- could still silently take the hint-dropping default. Two rounds shipped exactly that mistake in two different places. Fix: made the 3-arg onExhausted(target, reason, profile) the interface's single abstract method; the 2-arg form is now a default that delegates with a null profile. A lambda declared against ExhaustionSink today is forced by the compiler to take three parameters -- there is no overload left for it to bind to that can drop the hint. This is enforced by the type system, not by a test that has to remember to check for it. Knock-on changes: - ExhaustionSink.none() -- a 3-arg lambda, still a genuine no-op, now safe by construction rather than by care. - ExhaustionSink.forwardingTo(...) -- collapses to a one-line 3-arg lambda; kept as a named factory (round 3's lesson: a test must call the real object, not rebuild its shape). - Fleetd.java's real sink and the two OpenCodeLauncherTest sinks that used to be anonymous classes overriding both overloads are now plain lambdas too -- the 2-arg override each carried was pure boilerplate once the interface provides it as a default. - CompletionResolver.java itself: UNCHANGED, zero diff (confirmed via `git diff --stat` before staging) -- its two call sites still call the 2-arg onExhausted(target, reason), which is now the default and behaves identically. CompletionResolverTest (41 tests, 0 failures) proves this; its five ExhaustionSink lambdas needed a mechanical third parameter added to keep compiling against the new abstract method, no assertion changed. Mutation proof, re-run against the new shape: forwardingTo's body edited to call the 2-arg default instead of passing the hint through (the equivalent of round 3's "delete the 3-arg override" now that there is only one method to break) -- both new tests go red with the same assertions as round 3: ExhaustionSinkForwardingHazardTest...: expected: but was: OpenCodeLauncherTest...ForwardingHop: expected: but was: Tests run: 68, Failures: 2 Restored, re-ran: green (Tests run: 109, Failures: 0, including CompletionResolverTest). Compiler proof (not committed -- a scratch file outside the worktree, compiled with the real ExhaustionSink.java on the classpath, then deleted): ExhaustionSink forwarder = (target, reason) -> System.out.println(target + reason); error: incompatible types: incompatible parameter types in lambda expression A 2-arg lambda against this interface no longer compiles at all. mvn clean install: Tests run: 1129, Failures: 0, Errors: 0, Skipped: 0, BUILD SUCCESS. --- .../src/main/java/dev/ltms/fleet/Fleetd.java | 60 +++++----- .../dev/ltms/fleet/inject/ExhaustionSink.java | 112 ++++++++---------- .../fleet/inject/CompletionResolverTest.java | 10 +- .../ExhaustionSinkForwardingHazardTest.java | 32 ++--- .../fleet/member/OpenCodeLauncherTest.java | 58 ++++----- 5 files changed, 119 insertions(+), 153 deletions(-) diff --git a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java index a3fbe44..2b578fb 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java +++ b/fleetd/src/main/java/dev/ltms/fleet/Fleetd.java @@ -174,11 +174,12 @@ public final class Fleetd { // genuine cycle. Break it exactly like liveCountRef below: a forwarding sink built now, // pointed at the real one once it exists. AtomicReference exhaustionSinkRef = new AtomicReference<>(ExhaustionSink.none()); - // fleetd #234, round 3: this MUST go through ExhaustionSink.forwardingTo, never a lambda or - // a hand-written anonymous class here — see that factory's javadoc for why a lambda at this - // call site silently drops OpenCodeLauncher's profile hint. Routing through the shared - // factory also lets a test call the exact same object this line builds, instead of - // asserting a copy of its shape. + // fleetd #234, round 4: routed through the shared ExhaustionSink.forwardingTo factory + // rather than written inline here — not because a lambda at this call site is unsafe + // anymore (it is not: the 3-arg overload is now the interface's single abstract method, so + // there is no 2-arg overload left for any lambda to silently bind to instead), but so a + // test can call the exact same object this line builds, instead of asserting a copy of its + // shape (round 3's lesson). ExhaustionSink forwardingExhaustionSink = ExhaustionSink.forwardingTo(exhaustionSinkRef::get); // The claude-code adapter is the always-present default; keep it even with no profiles (so a // bridge configured with no workers, or opencode-only, still has a well-defined base adapter) @@ -361,34 +362,29 @@ public final class Fleetd { // fallback whenever the roster lookup misses; (2) if a profile still cannot be resolved // (neither the roster nor the hint names a configured one), this logs loudly at ERROR // instead of silently doing nothing — a control that cannot act must say so. - ExhaustionSink exhaustionSink = new ExhaustionSink() { - @Override - public void onExhausted(String target, String reason) { - onExhausted(target, reason, null); - } - - @Override - public void onExhausted(String target, String reason, String profileHint) { - String profileName = sessions.roster().stream() - .filter(session -> target.equals(session.terminalId())) - .findFirst() - .map(MemberSession::profile) - .orElse(profileHint); - FleetConfig.Profile profile = profileName == null ? null : config.get().profiles().get(profileName); - if (profile == null) { - log.error("quarantine requested for target '{}' ({}) but no profile could be " - + "resolved — the target is not (yet) in the roster, and {} — " - + "credential NOT quarantined (fleetd #234)", - target, reason, - profileHint == null ? "no profile hint was given" - : "the hinted profile '" + profileHint + "' is not configured"); - return; - } - String credentialId = profile.effectiveCredentialId(); - quarantine.quarantine(credentialId); - log.warn("credential '{}' quarantined for {}s (profile '{}'): {}", credentialId, - cfg.quarantineCooldownSeconds(), profile.profile(), reason); + // fleetd #234, round 4: the 3-arg overload is now ExhaustionSink's single abstract method, + // so this is safely a lambda — there is no separate 2-arg overload left for it to bind to + // instead and silently drop profileHint (that was rounds 1-3's whole hazard). + ExhaustionSink exhaustionSink = (target, reason, profileHint) -> { + String profileName = sessions.roster().stream() + .filter(session -> target.equals(session.terminalId())) + .findFirst() + .map(MemberSession::profile) + .orElse(profileHint); + FleetConfig.Profile profile = profileName == null ? null : config.get().profiles().get(profileName); + if (profile == null) { + log.error("quarantine requested for target '{}' ({}) but no profile could be " + + "resolved — the target is not (yet) in the roster, and {} — " + + "credential NOT quarantined (fleetd #234)", + target, reason, + profileHint == null ? "no profile hint was given" + : "the hinted profile '" + profileHint + "' is not configured"); + return; } + String credentialId = profile.effectiveCredentialId(); + quarantine.quarantine(credentialId); + log.warn("credential '{}' quarantined for {}s (profile '{}'): {}", credentialId, + cfg.quarantineCooldownSeconds(), profile.profile(), reason); }; // fleetd #175: point the forwarding sink handed to OpenCodeLauncher above at the real one, // now that `sessions` exists to resolve target -> session -> profile. diff --git a/fleetd/src/main/java/dev/ltms/fleet/inject/ExhaustionSink.java b/fleetd/src/main/java/dev/ltms/fleet/inject/ExhaustionSink.java index 5a61d11..dc39cc6 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/inject/ExhaustionSink.java +++ b/fleetd/src/main/java/dev/ltms/fleet/inject/ExhaustionSink.java @@ -12,91 +12,81 @@ import java.util.function.Supplier; * profiles or credentials, so mapping {@code target} to whatever should be quarantined is entirely * the sink's job — see {@code Fleetd.main}'s wiring, which resolves target → session → profile → * {@code effectiveCredentialId()} and calls {@code BackendQuarantine.quarantine} on it. + * + *

The 3-arg overload is the single abstract method — fleetd #234, round 4. Two + * earlier rounds each shipped a caller that silently dropped the profile hint (see {@link + * #onExhausted(String, String, String)}): a lambda written against this interface can only ever + * implement whichever overload is abstract, and while the 2-arg form held that position, EVERY + * lambda site — a call-site forwarder in {@code Fleetd.java}, a hand-built test double — bound to + * it and silently inherited the profile-dropping default, whether or not its author remembered the + * hazard. Making the 3-arg form abstract instead removes the shape entirely: a lambda declared + * against this interface today is forced by the compiler to take {@code (target, reason, + * profile)}, so there is no overload left for it to bind to that can drop the hint. This is a type + * change, not a test — it holds even for a caller that has never heard of fleetd #234. */ @FunctionalInterface public interface ExhaustionSink { /** - * @param target the herdr terminal id whose turn was classified {@code BACKEND_EXHAUSTED} - * @param reason the matched-line reason carried by the classification + * @param target the herdr terminal id whose turn was classified {@code BACKEND_EXHAUSTED} + * @param reason the matched-line reason carried by the classification + * @param profile the profile the caller already knows should be quarantined, or {@code null} + * when the caller has no better answer than {@code target} alone (fleetd #234): + * a caller whose {@code target} is not yet resolvable through whatever roster + * the sink's implementation consults — {@link + * dev.ltms.fleet.member.OpenCodeLauncher}'s model-mismatch check fires from + * {@code SessionAwareHandle.agentSessionId()}, which runs during {@code + * SessionManager.acquire()} before that session is registered, so a + * target -> session -> profile lookup finds nothing at that point. That launcher + * already has its own {@code FleetConfig.Profile} in hand and does not need the + * roster to know which profile to quarantine, so it supplies this directly + * instead of leaving the sink to guess. */ - void onExhausted(String target, String reason); + void onExhausted(String target, String reason, String profile); /** - * Same notification, plus a profile name the CALLER already knows — for a caller whose - * {@code target} is not yet resolvable through whatever roster the sink's implementation - * consults (fleetd #234). {@link dev.ltms.fleet.member.OpenCodeLauncher}'s model-mismatch check - * fires from {@code SessionAwareHandle.agentSessionId()}, which runs during {@code - * SessionManager.acquire()} before that session is registered — a target -> session -> - * profile lookup finds nothing at that point. That launcher already has its own {@code - * FleetConfig.Profile} in hand and does not need the roster to know which profile to - * quarantine, so it calls this overload instead of leaving the sink to guess. + * Convenience for a caller with no profile to offer — every existing call site that predates + * the hint (fleetd #234): {@link dev.ltms.fleet.inject.CompletionResolver}'s two call sites + * always call with a {@code target} that IS live in the roster at the time of the call, so they + * need no hint and keep working exactly as before, unchanged by this default. * - *

Defaults to the two-arg overload, discarding {@code profile} — the correct behaviour for - * every caller that has not been updated to supply one: {@link - * dev.ltms.fleet.inject.CompletionResolver}'s two call sites always call {@code target} that - * IS live in the roster at the time of the call, so they need no hint and keep working exactly - * as before. A sink that wants to use the hint (see {@code Fleetd.main}'s wiring) overrides this - * method directly rather than relying on the default. - * - * @param target as {@link #onExhausted(String, String)} - * @param reason as {@link #onExhausted(String, String)} - * @param profile the profile the caller already knows should be quarantined, or {@code null} - * when the caller has no better answer than {@code target} alone + * @param target as {@link #onExhausted(String, String, String)} + * @param reason as {@link #onExhausted(String, String, String)} */ - default void onExhausted(String target, String reason, String profile) { - onExhausted(target, reason); + default void onExhausted(String target, String reason) { + onExhausted(target, reason, null); } /** * Inert sink — nothing happens on exhaustion. The explicit stand-in a caller (or a test not * exercising this feature) passes instead of a defaulting overload, exactly like - * {@link ExhaustedPatternLookup#none()}. Safe as a lambda regardless of arity: its 2-arg body - * is empty, and the inherited 3-arg {@code default} just calls that same empty body — there is - * no hint to drop because this sink does nothing either way. + * {@link ExhaustedPatternLookup#none()}. Safe as a lambda: the 3-arg form is now the interface's + * single abstract method, so a lambda here has no other overload to silently bind to instead — + * it simply does nothing with all three arguments. */ static ExhaustionSink none() { - return (target, reason) -> { }; + return (target, reason, profile) -> { }; } /** - * A sink that forwards BOTH overloads to whatever {@code target} currently supplies (fleetd - * #234, round 3). Exists to break a genuine construction-order cycle: {@code Fleetd.main} - * builds its adapters (including {@link dev.ltms.fleet.member.OpenCodeLauncher}) before - * {@code sessions} exists, so it cannot hand them the real sink yet — it hands them a - * forwarder pointed at an {@code AtomicReference} that starts at {@link #none()} - * and gets {@code .set()} to the real sink once {@code sessions} is built. {@code target} is - * evaluated on every call, never cached, so the forwarder keeps working after the reference is - * repointed. + * A sink that forwards to whatever {@code target} currently supplies (fleetd #234). Exists to + * break a genuine construction-order cycle: {@code Fleetd.main} builds its adapters (including + * {@link dev.ltms.fleet.member.OpenCodeLauncher}) before {@code sessions} exists, so it cannot + * hand them the real sink yet — it hands them a forwarder pointed at an {@code + * AtomicReference} that starts at {@link #none()} and gets {@code .set()} to the + * real sink once {@code sessions} is built. {@code target} is evaluated on every call, never + * cached, so the forwarder keeps working after the reference is repointed. * - *

This is the one place that forwarding shape is written, on purpose. A - * forwarder written directly at a call site as a lambda — - * {@code (t, r) -> target.get().onExhausted(t, r)} — implements only the interface's single - * abstract method (the 2-arg overload) and silently inherits the 3-arg overload's {@code - * default} body, which discards whatever profile hint the real caller supplied and calls back - * into the 2-arg method instead. {@link dev.ltms.fleet.member.OpenCodeLauncher}'s model-mismatch - * check relies on that hint reaching the real sink (it fires before its session is registered - * in the roster, so the roster alone cannot resolve which profile to quarantine) — a lambda - * forwarder here silently reproduces that exact bug. Routing every forwarder through this one - * factory, instead of re-writing the shape at each call site, means a test can pin the shape - * once and have both {@code Fleetd.java} and the test call the identical object — see {@code - * ExhaustionSinkForwardingHazardTest} and {@code OpenCodeLauncherTest} for the coverage this - * makes possible. + *

Now safe as a one-line lambda (round 4): forwarding the single 3-arg abstract method + * forwards everything a caller can supply — there is no separate 2-arg overload left for a + * forwarder to bind to instead and silently lose the hint. Kept as a named factory rather than + * written inline at each call site anyway, so a test can call the exact object {@code + * Fleetd.java} builds instead of asserting a rebuilt copy of its shape (round 3's lesson). * * @param target supplies the sink to forward to, evaluated fresh on every call - * @return a sink whose every overload delegates to {@code target.get()}'s matching overload + * @return a sink whose call delegates to {@code target.get()} */ static ExhaustionSink forwardingTo(Supplier target) { - return new ExhaustionSink() { - @Override - public void onExhausted(String t, String r) { - target.get().onExhausted(t, r); - } - - @Override - public void onExhausted(String t, String r, String p) { - target.get().onExhausted(t, r, p); - } - }; + return (t, r, p) -> target.get().onExhausted(t, r, p); } } diff --git a/fleetd/src/test/java/dev/ltms/fleet/inject/CompletionResolverTest.java b/fleetd/src/test/java/dev/ltms/fleet/inject/CompletionResolverTest.java index d837f1f..01e8939 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/inject/CompletionResolverTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/inject/CompletionResolverTest.java @@ -499,7 +499,7 @@ class CompletionResolverTest { Rendezvous rendezvous = new Rendezvous(); ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached"); java.util.List notified = new java.util.ArrayList<>(); - ExhaustionSink sink = (target, reason) -> notified.add(target + ": " + reason); + ExhaustionSink sink = (target, reason, profile) -> notified.add(target + ": " + reason); CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, sink); var waiter = rendezvous.open("term_a"); @@ -520,7 +520,7 @@ class CompletionResolverTest { Rendezvous rendezvous = new Rendezvous(); ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached"); java.util.List notified = new java.util.ArrayList<>(); - ExhaustionSink sink = (target, reason) -> notified.add(target + ": " + reason); + ExhaustionSink sink = (target, reason, profile) -> notified.add(target + ": " + reason); CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, sink); var waiter = rendezvous.open("term_a"); @@ -648,7 +648,7 @@ class CompletionResolverTest { Rendezvous rendezvous = new Rendezvous(); ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached"); java.util.List notified = new java.util.ArrayList<>(); - ExhaustionSink sink = (target, reason) -> notified.add(target + ": " + reason); + ExhaustionSink sink = (target, reason, profile) -> notified.add(target + ": " + reason); CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, sink); var waiter = rendezvous.open("term_a"); @@ -701,7 +701,7 @@ class CompletionResolverTest { Rendezvous rendezvous = new Rendezvous(); ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached"); java.util.List notified = new java.util.ArrayList<>(); - ExhaustionSink sink = (target, reason) -> notified.add(target + ": " + reason); + ExhaustionSink sink = (target, reason, profile) -> notified.add(target + ": " + reason); CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, sink); var waiter = rendezvous.open("term_a"); @@ -728,7 +728,7 @@ class CompletionResolverTest { Rendezvous rendezvous = new Rendezvous(); ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached"); java.util.List notified = new java.util.ArrayList<>(); - ExhaustionSink sink = (target, reason) -> notified.add(target + ": " + reason); + ExhaustionSink sink = (target, reason, profile) -> notified.add(target + ": " + reason); CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, patterns, sink); var waiter = rendezvous.open("term_a"); diff --git a/fleetd/src/test/java/dev/ltms/fleet/inject/ExhaustionSinkForwardingHazardTest.java b/fleetd/src/test/java/dev/ltms/fleet/inject/ExhaustionSinkForwardingHazardTest.java index 6c42381..2a2e538 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/inject/ExhaustionSinkForwardingHazardTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/inject/ExhaustionSinkForwardingHazardTest.java @@ -7,31 +7,25 @@ import java.util.concurrent.atomic.AtomicReference; import static org.junit.jupiter.api.Assertions.assertEquals; /** - * fleetd #234, round 3: calls the REAL production factory, {@link ExhaustionSink#forwardingTo}, - * rather than rebuilding its shape locally. A round-2 version of this test built its own copy of - * the forwarder inline — mutating {@code Fleetd.java}'s actual forwarder back into a lambda left - * that copy untouched, so the test kept passing while production had regressed to exactly the bug - * it was meant to catch. Calling the shared factory here means the same object under test IS the - * object {@code Fleetd.java} builds via {@code ExhaustionSink.forwardingTo(exhaustionSinkRef::get)} - * — a defect in the factory body, or a call site reverting to a hand-written lambda instead of the - * factory, has exactly one place left to hide, and this test reaches into it. + * fleetd #234. Round 3: calls the REAL production factory, {@link ExhaustionSink#forwardingTo}, + * rather than rebuilding its shape locally — a round-2 version of this test built its own copy of + * the forwarder inline, so mutating {@code Fleetd.java}'s actual forwarder back into a lambda left + * that copy untouched and the test kept passing while production had regressed. Calling the shared + * factory here means the same object under test IS the object {@code Fleetd.java} builds via + * {@code ExhaustionSink.forwardingTo(exhaustionSinkRef::get)}. + * + *

Round 4: the 3-arg overload is now {@link ExhaustionSink}'s single abstract method, so {@code + * forwardingTo} itself is a one-line lambda and every sink below can safely be one too — there is + * no 2-arg overload left for any of them to silently bind to instead. This test still catches a + * regression inside {@code forwardingTo}'s body (e.g. one that calls the 2-arg default and drops + * the hint that way) because it still goes through the shared factory rather than a rebuilt copy. */ class ExhaustionSinkForwardingHazardTest { @Test void forwardingToDeliversTheProfileHintToWhateverSinkTheSupplierCurrentlyReturns() { AtomicReference hintSeenByRealSink = new AtomicReference<>("NEVER CALLED"); - ExhaustionSink real = new ExhaustionSink() { - @Override - public void onExhausted(String target, String reason) { - onExhausted(target, reason, null); - } - - @Override - public void onExhausted(String target, String reason, String profileHint) { - hintSeenByRealSink.set(profileHint); - } - }; + ExhaustionSink real = (target, reason, profileHint) -> hintSeenByRealSink.set(profileHint); // Same shape Fleetd.java uses: a reference that starts at none() and is repointed later. AtomicReference ref = new AtomicReference<>(ExhaustionSink.none()); diff --git a/fleetd/src/test/java/dev/ltms/fleet/member/OpenCodeLauncherTest.java b/fleetd/src/test/java/dev/ltms/fleet/member/OpenCodeLauncherTest.java index ee67781..6fd791b 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/member/OpenCodeLauncherTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/member/OpenCodeLauncherTest.java @@ -925,7 +925,7 @@ class OpenCodeLauncherTest { // credentialId — the profile that actually escaped the fleet's accounting. FleetConfig.Profile cfg = opencodeCfg("opencode/nemotron-3-ultra-free", null, null); List exhausted = new ArrayList<>(); - ExhaustionSink sink = (target, reason) -> exhausted.add(target + "|" + reason); + ExhaustionSink sink = (target, reason, profile) -> exhausted.add(target + "|" + reason); OpenCodeLauncher launcher = serviceWithSink(herdr, configRoot, discRoot, cfg, sink); SessionManager sessions = new SessionManager(launcher); @@ -960,7 +960,7 @@ class OpenCodeLauncherTest { void aProviderPrefixedModelMatchingBothIdAndProviderIsNotAMismatch(@TempDir Path configRoot, @TempDir Path discRoot) throws Exception { List exhausted = new ArrayList<>(); - ExhaustionSink sink = (target, reason) -> exhausted.add(reason); + ExhaustionSink sink = (target, reason, profile) -> exhausted.add(reason); FleetConfig.Profile cfg = opencodeCfg("openai/gpt-5.6-terra", null, null); PeerHandle handle = serviceWithSink(new FakeHerdr(), configRoot, discRoot, cfg, sink) .spawn(new SpawnRequest(null, "/work/dir", null)); @@ -976,7 +976,7 @@ class OpenCodeLauncherTest { void aGxProviderPrefixedModelMatchingBothIdAndProviderIsNotAMismatch(@TempDir Path configRoot, @TempDir Path discRoot) throws Exception { List exhausted = new ArrayList<>(); - ExhaustionSink sink = (target, reason) -> exhausted.add(reason); + ExhaustionSink sink = (target, reason, profile) -> exhausted.add(reason); FleetConfig.Profile cfg = opencodeCfg("gx/deepseek-v4-flash", null, null); PeerHandle handle = serviceWithSink(new FakeHerdr(), configRoot, discRoot, cfg, sink) .spawn(new SpawnRequest(null, "/work/dir", null)); @@ -998,7 +998,7 @@ class OpenCodeLauncherTest { void aMissingProviderIdInTheEvidenceIsUnknownNotAMismatchWhenTheIdMatches( @TempDir Path configRoot, @TempDir Path discRoot) throws Exception { List exhausted = new ArrayList<>(); - ExhaustionSink sink = (target, reason) -> exhausted.add(reason); + ExhaustionSink sink = (target, reason, profile) -> exhausted.add(reason); FleetConfig.Profile cfg = opencodeCfg("openai/gpt-5.6-terra", null, null); PeerHandle handle = serviceWithSink(new FakeHerdr(), configRoot, discRoot, cfg, sink) .spawn(new SpawnRequest(null, "/work/dir", null)); @@ -1019,7 +1019,7 @@ class OpenCodeLauncherTest { void aMissingProviderIdInTheEvidenceStillCatchesARealIdMismatch( @TempDir Path configRoot, @TempDir Path discRoot) throws Exception { List exhausted = new ArrayList<>(); - ExhaustionSink sink = (target, reason) -> exhausted.add(reason); + ExhaustionSink sink = (target, reason, profile) -> exhausted.add(reason); FleetConfig.Profile cfg = opencodeCfg("openai/gpt-5.6-terra", null, null); PeerHandle handle = serviceWithSink(new FakeHerdr(), configRoot, discRoot, cfg, sink) .spawn(new SpawnRequest(null, "/work/dir", null)); @@ -1043,7 +1043,7 @@ class OpenCodeLauncherTest { void aBareModelWithNoProviderPrefixMatchesOnIdAloneAndIsNotAMismatch(@TempDir Path configRoot, @TempDir Path discRoot) throws Exception { List exhausted = new ArrayList<>(); - ExhaustionSink sink = (target, reason) -> exhausted.add(reason); + ExhaustionSink sink = (target, reason, profile) -> exhausted.add(reason); FleetConfig.Profile cfg = opencodeCfg("deepseek-v4-flash", null, null); PeerHandle handle = serviceWithSink(new FakeHerdr(), configRoot, discRoot, cfg, sink) .spawn(new SpawnRequest(null, "/work/dir", null)); @@ -1065,7 +1065,7 @@ class OpenCodeLauncherTest { void aRealIdMismatchLogsAnErrorNamingBothModelsAndQuarantinesThroughTheSink( @TempDir Path configRoot, @TempDir Path discRoot) throws Exception { List exhausted = new ArrayList<>(); - ExhaustionSink sink = (target, reason) -> exhausted.add(target + "|" + reason); + ExhaustionSink sink = (target, reason, profile) -> exhausted.add(target + "|" + reason); FleetConfig.Profile cfg = opencodeCfg("opencode/nemotron-3-ultra-free", null, null); Logger logger = (Logger) LoggerFactory.getLogger(OpenCodeLauncher.class); @@ -1112,7 +1112,7 @@ class OpenCodeLauncherTest { void unknownOrUnparseableModelEvidenceNeverQuarantines(@TempDir Path configRoot, @TempDir Path discRoot) throws Exception { List exhausted = new ArrayList<>(); - ExhaustionSink sink = (target, reason) -> exhausted.add(reason); + ExhaustionSink sink = (target, reason, profile) -> exhausted.add(reason); FleetConfig.Profile cfg = opencodeCfg("openai/gpt-5.6-terra", null, null); PeerHandle handle = serviceWithSink(new FakeHerdr(), configRoot, discRoot, cfg, sink) .spawn(new SpawnRequest(null, "/work/dir", null)); @@ -1138,7 +1138,7 @@ class OpenCodeLauncherTest { void aProfileWithNoConfiguredModelIsNeverCheckedForAMismatch(@TempDir Path configRoot, @TempDir Path discRoot) throws Exception { List exhausted = new ArrayList<>(); - ExhaustionSink sink = (target, reason) -> exhausted.add(reason); + ExhaustionSink sink = (target, reason, profile) -> exhausted.add(reason); FleetConfig.Profile cfg = opencodeCfg(null, null, null); PeerHandle handle = serviceWithSink(new FakeHerdr(), configRoot, discRoot, cfg, sink) .spawn(new SpawnRequest(null, "/work/dir", null)); @@ -1167,7 +1167,7 @@ class OpenCodeLauncherTest { void modelCheckReadsTheResolvedSessionsOwnRowNotWhateverIsNewestInTheSharedDirectory( @TempDir Path configRoot, @TempDir Path discRoot) throws Exception { List exhausted = new ArrayList<>(); - ExhaustionSink sink = (target, reason) -> exhausted.add(reason); + ExhaustionSink sink = (target, reason, profile) -> exhausted.add(reason); FleetConfig.Profile cfg = opencodeCfg("openai/gpt-5.6-terra", null, null); PeerHandle handle = serviceWithSink(new FakeHerdr(), configRoot, discRoot, cfg, sink) .spawn(new SpawnRequest(null, "/work/dir", null)); @@ -1217,18 +1217,12 @@ class OpenCodeLauncherTest { Map profiles = Map.of(cfg.profile(), cfg); BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.SECONDS.toNanos(1800)); - ExhaustionSink sink = new ExhaustionSink() { - @Override - public void onExhausted(String target, String reason) { - onExhausted(target, reason, null); // no roster resolution modelled here — see below - } - - @Override - public void onExhausted(String target, String reason, String profileHint) { - FleetConfig.Profile profile = profileHint == null ? null : profiles.get(profileHint); - if (profile != null) { - quarantine.quarantine(profile.effectiveCredentialId()); - } + // fleetd #234, round 4: safely a lambda now — the 3-arg overload is the interface's single + // abstract method, so there is no 2-arg overload left to bind to instead. + ExhaustionSink sink = (target, reason, profileHint) -> { + FleetConfig.Profile profile = profileHint == null ? null : profiles.get(profileHint); + if (profile != null) { + quarantine.quarantine(profile.effectiveCredentialId()); } }; @@ -1270,7 +1264,7 @@ class OpenCodeLauncherTest { // Deliberately ignores the profile hint — the pre-fix shape: only a roster lookup (modelled // here as always empty, since acquire() has not registered the session yet either way). - ExhaustionSink rosterOnlySink = (target, reason) -> { }; + ExhaustionSink rosterOnlySink = (target, reason, profile) -> { }; FakeHerdr herdr = new FakeHerdr(); OpenCodeLauncher launcher = serviceWithSink(herdr, configRoot, discRoot, cfg, rosterOnlySink); @@ -1329,19 +1323,11 @@ class OpenCodeLauncherTest { SessionManager sessions = new SessionManager(launcher); // Fleetd.java:359-390 — the real sink is only built and pointed to AFTER `sessions` exists, - // same order as production. - ExhaustionSink realSink = new ExhaustionSink() { - @Override - public void onExhausted(String target, String reason) { - onExhausted(target, reason, null); // no roster resolution modelled — mirrors a miss - } - - @Override - public void onExhausted(String target, String reason, String profileHint) { - FleetConfig.Profile profile = profileHint == null ? null : profiles.get(profileHint); - if (profile != null) { - quarantine.quarantine(profile.effectiveCredentialId()); - } + // same order as production. Safely a lambda (round 4): see the note on the forwarder above. + ExhaustionSink realSink = (target, reason, profileHint) -> { + FleetConfig.Profile profile = profileHint == null ? null : profiles.get(profileHint); + if (profile != null) { + quarantine.quarantine(profile.effectiveCredentialId()); } }; exhaustionSinkRef.set(realSink);