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 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 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