Compare commits

...

4 Commits

Author SHA1 Message Date
Dai Ha 31b028e860 fleetd #234 round 4: invert ExhaustionSink's abstract method so the bug class is unrepresentable
CI / contract (pull_request) Successful in 46s
CI / build (pull_request) Successful in 1m23s
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: <gx> but was: <null>
  OpenCodeLauncherTest...ForwardingHop:  expected: <true> but was: <false>
  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.
2026-09-03 10:46:27 +07:00
Dai Ha c935b181dd fleetd #234 round 3: make ExhaustionSink's forwarder a shared factory, not a rebuilt-per-caller shape
CI / contract (pull_request) Successful in 42s
CI / build (pull_request) Successful in 1m36s
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<ExhaustionSink> 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: <gx> but was: <null>
  OpenCodeLauncherTest...ForwardingHop:  expected: <true> but was: <false>
  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.
2026-09-03 10:30:37 +07:00
Dai Ha c325054242 fleetd #234 round 2: fix the ExhaustionSink forwarding hop Fleetd.java actually uses
CI / contract (pull_request) Successful in 1m10s
CI / build (pull_request) Successful in 1m18s
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: <gx> but was: <null>" for the interface
proof, "expected: <true> but was: <false>" 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.
2026-09-03 10:22:31 +07:00
Dai Ha 4877992a70 fleetd #234: key the opencode model check on the resolved session id, and make the spawn-time quarantine actually happen
CI / contract (pull_request) Successful in 58s
CI / build (pull_request) Successful in 1m29s
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.
2026-09-03 10:09:26 +07:00
8 changed files with 493 additions and 74 deletions
+42 -12
View File
@@ -174,7 +174,13 @@ 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<ExhaustionSink> exhaustionSinkRef = new AtomicReference<>(ExhaustionSink.none());
ExhaustionSink forwardingExhaustionSink = (target, reason) -> exhaustionSinkRef.get().onExhausted(target, reason);
// 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)
// unless opencode is the only kind configured.
@@ -345,17 +351,41 @@ 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.
// 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.
exhaustionSinkRef.set(exhaustionSink);
@@ -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
@@ -10,22 +12,81 @@ package dev.ltms.fleet.inject;
* 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.
*
* <p><strong>The 3-arg overload is the single abstract method — fleetd #234, round 4.</strong> 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 <em>forced</em> 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()} <em>before</em> 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);
/**
* 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.
*
* @param target as {@link #onExhausted(String, String, String)}
* @param reason as {@link #onExhausted(String, String, String)}
*/
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()}.
* {@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 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<ExhaustionSink>} 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.
*
* <p>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 call delegates to {@code target.get()}
*/
static ExhaustionSink forwardingTo(Supplier<ExhaustionSink> target) {
return (t, r, p) -> target.get().onExhausted(t, r, p);
}
}
@@ -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<String> 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
@@ -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}.
*
* <p>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 <em>its</em> 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.
* <p><strong>The {@code directory} match is a heuristic, not an identity — fleetd #234.</strong> 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
* <em>not</em> 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.
*
* <p>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}.
*
* <p>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.
*
* <p>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.
*
* <p>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"));
@@ -499,7 +499,7 @@ class CompletionResolverTest {
Rendezvous rendezvous = new Rendezvous();
ExhaustedPatternLookup patterns = target -> Pattern.compile("usage limit has been reached");
java.util.List<String> 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<String> 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<String> 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<String> 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<String> 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");
@@ -0,0 +1,41 @@
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;
/**
* 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)}.
*
* <p>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<String> hintSeenByRealSink = new AtomicReference<>("NEVER CALLED");
ExhaustionSink real = (target, reason, profileHint) -> hintSeenByRealSink.set(profileHint);
// Same shape Fleetd.java uses: a reference that starts at none() and is repointed later.
AtomicReference<ExhaustionSink> ref = new AtomicReference<>(ExhaustionSink.none());
ExhaustionSink forwarding = ExhaustionSink.forwardingTo(ref::get);
ref.set(real);
forwarding.onExhausted("term_x", "model mismatch", "gx");
assertEquals("gx", hintSeenByRealSink.get(),
"ExhaustionSink.forwardingTo must deliver the profile hint to the sink the supplier "
+ "currently returns — the exact object Fleetd.java's forwarder is built from");
}
}
@@ -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),
@@ -910,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<String> 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);
@@ -945,7 +960,7 @@ class OpenCodeLauncherTest {
void aProviderPrefixedModelMatchingBothIdAndProviderIsNotAMismatch(@TempDir Path configRoot,
@TempDir Path discRoot) throws Exception {
List<String> 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));
@@ -961,7 +976,7 @@ class OpenCodeLauncherTest {
void aGxProviderPrefixedModelMatchingBothIdAndProviderIsNotAMismatch(@TempDir Path configRoot,
@TempDir Path discRoot) throws Exception {
List<String> 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));
@@ -983,7 +998,7 @@ class OpenCodeLauncherTest {
void aMissingProviderIdInTheEvidenceIsUnknownNotAMismatchWhenTheIdMatches(
@TempDir Path configRoot, @TempDir Path discRoot) throws Exception {
List<String> 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));
@@ -1004,7 +1019,7 @@ class OpenCodeLauncherTest {
void aMissingProviderIdInTheEvidenceStillCatchesARealIdMismatch(
@TempDir Path configRoot, @TempDir Path discRoot) throws Exception {
List<String> 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));
@@ -1028,7 +1043,7 @@ class OpenCodeLauncherTest {
void aBareModelWithNoProviderPrefixMatchesOnIdAloneAndIsNotAMismatch(@TempDir Path configRoot,
@TempDir Path discRoot) throws Exception {
List<String> 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));
@@ -1050,7 +1065,7 @@ class OpenCodeLauncherTest {
void aRealIdMismatchLogsAnErrorNamingBothModelsAndQuarantinesThroughTheSink(
@TempDir Path configRoot, @TempDir Path discRoot) throws Exception {
List<String> 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);
@@ -1097,7 +1112,7 @@ class OpenCodeLauncherTest {
void unknownOrUnparseableModelEvidenceNeverQuarantines(@TempDir Path configRoot,
@TempDir Path discRoot) throws Exception {
List<String> 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));
@@ -1123,7 +1138,7 @@ class OpenCodeLauncherTest {
void aProfileWithNoConfiguredModelIsNeverCheckedForAMismatch(@TempDir Path configRoot,
@TempDir Path discRoot) throws Exception {
List<String> 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));
@@ -1133,4 +1148,200 @@ 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<String> exhausted = new ArrayList<>();
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));
// 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).
*
* <p>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<String, FleetConfig.Profile> profiles = Map.of(cfg.profile(), cfg);
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.SECONDS.toNanos(1800));
// 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());
}
};
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, profile) -> { };
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");
}
/**
* 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 <em>forwarding</em>
* sink pointed at an {@code AtomicReference<ExhaustionSink>}, and only later {@code .set(...)}s
* that reference to the real sink once {@code sessions} is built.
*
* <p><strong>Round 3:</strong> 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,
@TempDir Path discRoot) throws Exception {
FleetConfig.Profile cfg = opencodeCfgWithCredential(
"terra", "opencode/nemotron-3-ultra-free", "openai-shared");
Map<String, FleetConfig.Profile> 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. 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<ExhaustionSink> exhaustionSinkRef =
new java.util.concurrent.atomic.AtomicReference<>(ExhaustionSink.none());
ExhaustionSink forwardingExhaustionSink = ExhaustionSink.forwardingTo(exhaustionSinkRef::get);
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. 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);
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");
}
}
@@ -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");
}
}