Compare commits
8 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 32ebf065ac | |||
| 1178b3f684 | |||
| 2a434ced2f | |||
| cc919aa2b6 | |||
| 39c7ce76f3 | |||
| b40f477210 | |||
| 5c56cb347f | |||
| e694deace3 |
@@ -169,6 +169,12 @@ public final class Fleetd {
|
||||
}
|
||||
});
|
||||
List<HerdrPeerLauncher> adapters = new ArrayList<>();
|
||||
// fleetd #175: the daemon's real ExhaustionSink can only be built once `sessions` exists
|
||||
// (below), but `sessions` needs `workers`, which needs the adapters built right here — a
|
||||
// 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);
|
||||
// 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.
|
||||
@@ -184,7 +190,7 @@ public final class Fleetd {
|
||||
opencodeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
|
||||
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
|
||||
() -> config.get().fleet(),
|
||||
() -> config.get().memberCredentials(), config::get));
|
||||
() -> config.get().memberCredentials(), config::get, forwardingExhaustionSink));
|
||||
}
|
||||
AtomicReference<Function<String, Integer>> liveCountRef = new AtomicReference<>(_ -> 0);
|
||||
// CB-578 stage B: one quarantine tracker for the whole daemon, shared between the launcher
|
||||
@@ -334,6 +340,11 @@ public final class Fleetd {
|
||||
// CREDENTIAL — not the profile name — so a profile sharing that credential (e.g. two models
|
||||
// on one OpenAI account) is refused too, not just the one that happened to report it. Reads
|
||||
// the profile config live off `config`, so a credentialId edit is hot: no restart needed.
|
||||
//
|
||||
// fleetd #175: this sink is now also the quarantine target for OpenCodeLauncher's
|
||||
// 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()
|
||||
@@ -342,10 +353,12 @@ public final class Fleetd {
|
||||
.ifPresent(profile -> {
|
||||
String credentialId = profile.effectiveCredentialId();
|
||||
quarantine.quarantine(credentialId);
|
||||
log.warn("credential '{}' quarantined for {}s (profile '{}' classified "
|
||||
+ "BACKEND_EXHAUSTED): {}", 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);
|
||||
AgentControl agents = router.memberAgents();
|
||||
CompletionResolver completion = new CompletionResolver(agents, rendezvous, exhaustedPatterns, exhaustionSink);
|
||||
// CB-113: deliver only to an available worker (its MCP is connected), never its boot window.
|
||||
|
||||
@@ -5,6 +5,9 @@ import dev.ltms.fleet.peer.MemberRole;
|
||||
/** Optional session lifecycle hook for live member-slot bindings. */
|
||||
public interface MemberLifecycle {
|
||||
|
||||
/** A slot held before a member process starts. */
|
||||
record SlotReservation(String slot, String profile) { }
|
||||
|
||||
MemberLifecycle NONE = new MemberLifecycle() {
|
||||
@Override
|
||||
public MemberRole acquired(MemberRole role, String profile, String terminal) {
|
||||
@@ -19,6 +22,20 @@ public interface MemberLifecycle {
|
||||
public void requireSlotFor(MemberRole role, String profile) {
|
||||
// no registry configured — nothing to validate against, so nothing is refused
|
||||
}
|
||||
|
||||
@Override
|
||||
public SlotReservation reserve(MemberRole role, String profile) {
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean bind(SlotReservation reservation, String terminal) {
|
||||
return false;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void release(SlotReservation reservation) {
|
||||
}
|
||||
};
|
||||
|
||||
/**
|
||||
@@ -29,10 +46,9 @@ public interface MemberLifecycle {
|
||||
* role — never {@code role} — when a slot-bound role (architect) could not be bound.
|
||||
* Callers must record THIS value on the session, never the requested {@code role}, so
|
||||
* a later roster read never reports a role the session does not hold (CB-619). In
|
||||
* normal operation this fallback should not happen once {@link #requireSlotFor} has
|
||||
* refused every unbindable spawn upfront — but a slot can still be lost between that
|
||||
* check and this call to a concurrent spawn racing for the same slot, so the honest
|
||||
* answer is still needed here too.
|
||||
* normal operation this fallback should not happen once a reservation has been bound.
|
||||
* It remains the honest answer if a caller has no reservation, or if binding a
|
||||
* reservation unexpectedly fails.
|
||||
*/
|
||||
MemberRole acquired(MemberRole role, String profile, String terminal);
|
||||
|
||||
@@ -46,4 +62,13 @@ public interface MemberLifecycle {
|
||||
* @throws IllegalArgumentException naming the role, the profile, and the pools that do carry it
|
||||
*/
|
||||
void requireSlotFor(MemberRole role, String profile);
|
||||
|
||||
/** Reserve a matching slot before launch, or refuse before a charter can be delivered. */
|
||||
SlotReservation reserve(MemberRole role, String profile);
|
||||
|
||||
/** Convert a reservation into a live terminal binding. */
|
||||
boolean bind(SlotReservation reservation, String terminal);
|
||||
|
||||
/** Return an unbound reservation after a failed launch. */
|
||||
void release(SlotReservation reservation);
|
||||
}
|
||||
|
||||
@@ -56,8 +56,10 @@ public final class MemberRegistry implements MemberLifecycle {
|
||||
}
|
||||
|
||||
private final Map<String, Entry> slots;
|
||||
/** Live {@code terminal_id → qualified slot key}; guarded by {@code this}. */
|
||||
/** Live {@code terminal_id → qualified slot key}; guarded by {@code terminalToSlot}. */
|
||||
private final Map<String, String> terminalToSlot = new HashMap<>();
|
||||
/** Slot keys held between reservation and the terminal binding. Guarded by terminalToSlot. */
|
||||
private final java.util.Set<String> reservedSlots = new java.util.HashSet<>();
|
||||
|
||||
/** Flatten every role pool in {@code fleet} into one registry. Leaders are not members. */
|
||||
public MemberRegistry(FleetConfig.Fleet fleet) {
|
||||
@@ -167,7 +169,7 @@ public final class MemberRegistry implements MemberLifecycle {
|
||||
if (existingSlot != null) {
|
||||
return slot.equals(existingSlot); // already this slot (idempotent) or a different one
|
||||
}
|
||||
if (terminalToSlot.containsValue(slot)) {
|
||||
if (terminalToSlot.containsValue(slot) || reservedSlots.contains(slot)) {
|
||||
return false; // slot already hosts a terminal — no second one
|
||||
}
|
||||
terminalToSlot.put(terminal, slot);
|
||||
@@ -275,6 +277,50 @@ public final class MemberRegistry implements MemberLifecycle {
|
||||
+ " on one of those profiles instead");
|
||||
}
|
||||
|
||||
@Override
|
||||
public SlotReservation reserve(MemberRole role, String profile) {
|
||||
if (role != MemberRole.ARCHITECT) {
|
||||
return null;
|
||||
}
|
||||
synchronized (terminalToSlot) {
|
||||
for (Entry entry : slotsFor(MemberRole.ARCHITECT).values()) {
|
||||
if ((profile == null || profile.isBlank() || Objects.equals(profile, entry.profile()))
|
||||
&& !terminalToSlot.containsValue(entry.key()) && reservedSlots.add(entry.key())) {
|
||||
return new SlotReservation(entry.key(), entry.profile());
|
||||
}
|
||||
}
|
||||
}
|
||||
throw new IllegalArgumentException("no free architect slot for profile '" + profile
|
||||
+ "' — every matching slot is already bound or reserved");
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean bind(SlotReservation reservation, String terminal) {
|
||||
if (reservation == null || terminal == null || terminal.isBlank()) {
|
||||
return false;
|
||||
}
|
||||
synchronized (terminalToSlot) {
|
||||
if (!reservedSlots.remove(reservation.slot())) {
|
||||
return false;
|
||||
}
|
||||
if (!isSlot(reservation.slot()) || terminalToSlot.containsKey(terminal)
|
||||
|| terminalToSlot.containsValue(reservation.slot())) {
|
||||
return false;
|
||||
}
|
||||
terminalToSlot.put(terminal, reservation.slot());
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void release(SlotReservation reservation) {
|
||||
if (reservation != null) {
|
||||
synchronized (terminalToSlot) {
|
||||
reservedSlots.remove(reservation.slot());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** Unbind a released terminal using the compare-safe registry operation. */
|
||||
@Override
|
||||
public void released(String terminal) {
|
||||
|
||||
@@ -6,6 +6,7 @@ import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.herdr.Agent;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
import dev.ltms.fleet.inject.ExhaustionSink;
|
||||
import dev.ltms.fleet.peer.Capability;
|
||||
import dev.ltms.fleet.peer.CharterReceipt;
|
||||
import dev.ltms.fleet.peer.PeerHandle;
|
||||
@@ -73,6 +74,19 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
*/
|
||||
private final OpenCodeSessionDiscovery discovery;
|
||||
|
||||
/**
|
||||
* fleetd #175: notified when {@link SessionAwareHandle} detects, on the same late-resolve
|
||||
* read that discovers the session id, that the live opencode session is running a DIFFERENT
|
||||
* model than the profile requested — opencode does not fail on an unknown {@code -m}, it
|
||||
* silently falls back to a default (potentially paid) model. Reused exactly as
|
||||
* {@code CompletionResolver}'s {@code BACKEND_EXHAUSTED} path uses it: this launcher supplies
|
||||
* only the herdr terminal id and a reason string; mapping target → session → profile →
|
||||
* credential stays entirely the sink's job (see {@code Fleetd.main}'s wiring). Defaults to
|
||||
* {@link ExhaustionSink#none()} for a caller (an older constructor, or a test not exercising
|
||||
* this) that opts out.
|
||||
*/
|
||||
private final ExhaustionSink exhaustionSink;
|
||||
|
||||
/**
|
||||
* Production constructor — disables the spawn-ready gate ({@code spawnReadyTimeoutMs == 0}) so it
|
||||
* matches the legacy non-blocking spawn semantics. Config dirs are created under the JVM temp dir.
|
||||
@@ -132,9 +146,25 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
Supplier<FleetConfig.Fleet> fleet,
|
||||
Supplier<FleetConfig.MemberCredentials> memberCredentials,
|
||||
Supplier<FleetConfig> config) {
|
||||
this(agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs, spawnReadyPollMs,
|
||||
fleet, memberCredentials, config, ExhaustionSink.none());
|
||||
}
|
||||
|
||||
/**
|
||||
* Production constructor, plus the live config for URI environment exclusions and the fleetd
|
||||
* #175 model-mismatch quarantine sink.
|
||||
*/
|
||||
public OpenCodeLauncher(AgentControl agents, WorkspaceControl spaces,
|
||||
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
|
||||
Function<String, String> env, long spawnReadyTimeoutMs, long spawnReadyPollMs,
|
||||
Supplier<FleetConfig.Fleet> fleet,
|
||||
Supplier<FleetConfig.MemberCredentials> memberCredentials,
|
||||
Supplier<FleetConfig> config,
|
||||
ExhaustionSink exhaustionSink) {
|
||||
this(agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs,
|
||||
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs),
|
||||
defaultConfigRoot(), defaultDiscoveryRoot(), fleet, memberCredentials, config);
|
||||
defaultConfigRoot(), defaultDiscoveryRoot(), fleet, memberCredentials, config,
|
||||
exhaustionSink);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -182,6 +212,7 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
spawnReadyTimeoutMs, nowMillis, sleeper, fleet);
|
||||
this.configRoot = configRoot;
|
||||
this.discovery = new OpenCodeSessionDiscovery(discoveryRoot);
|
||||
this.exhaustionSink = ExhaustionSink.none();
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -207,10 +238,29 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
Supplier<FleetConfig.Fleet> fleet,
|
||||
Supplier<FleetConfig.MemberCredentials> memberCredentials,
|
||||
Supplier<FleetConfig> config) {
|
||||
this(agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs, nowMillis, sleeper,
|
||||
configRoot, discoveryRoot, fleet, memberCredentials, config, ExhaustionSink.none());
|
||||
}
|
||||
|
||||
/**
|
||||
* Full testability constructor, plus the live config for URI environment exclusions and the
|
||||
* fleetd #175 model-mismatch quarantine sink — the constructor a test drives directly to
|
||||
* observe {@link ExhaustionSink#onExhausted} without going through {@code Fleetd.main}'s
|
||||
* wiring.
|
||||
*/
|
||||
public OpenCodeLauncher(AgentControl agents, WorkspaceControl spaces,
|
||||
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
|
||||
Function<String, String> env, long spawnReadyTimeoutMs,
|
||||
LongSupplier nowMillis, Runnable sleeper, Path configRoot, Path discoveryRoot,
|
||||
Supplier<FleetConfig.Fleet> fleet,
|
||||
Supplier<FleetConfig.MemberCredentials> memberCredentials,
|
||||
Supplier<FleetConfig> config,
|
||||
ExhaustionSink exhaustionSink) {
|
||||
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
|
||||
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials, null, config);
|
||||
this.configRoot = configRoot;
|
||||
this.discovery = new OpenCodeSessionDiscovery(discoveryRoot);
|
||||
this.exhaustionSink = exhaustionSink == null ? ExhaustionSink.none() : exhaustionSink;
|
||||
}
|
||||
|
||||
private static Path defaultConfigRoot() {
|
||||
@@ -622,8 +672,13 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
@Override
|
||||
public PeerHandle spawn(SpawnRequest req) {
|
||||
PeerHandle inner = super.spawn(req);
|
||||
return new SessionAwareHandle(inner, discovery, effectiveCwd(req),
|
||||
this::memberHerdrSocketConfigured, discoveryUnavailableWarned);
|
||||
// fleetd #175: the same profile config buildLaunch resolved for this spawn (requireProfile
|
||||
// is deterministic on req.profileName(), so re-resolving here costs a map lookup, not a
|
||||
// second decision) — SessionAwareHandle needs cfg.model() to know what THIS session should
|
||||
// be running.
|
||||
FleetConfig.Profile cfg = requireProfile(req.profileName());
|
||||
return new SessionAwareHandle(inner, discovery, effectiveCwd(req), cfg,
|
||||
this::memberHerdrSocketConfigured, discoveryUnavailableWarned, exhaustionSink);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -633,22 +688,37 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
* are untouched — only the opencode-specific identity answer is added. {@code sessionName()}
|
||||
* stays null: opencode has no display-name seam, so the logical name lives only in the bridge's
|
||||
* roster (see the SESSION_NAME capability).
|
||||
*
|
||||
* <p>fleetd #175: {@link #agentSessionId()} is also where the model-mismatch check lives (see
|
||||
* {@link #checkModelMatch()}) — it is the one method the real late-resolve path
|
||||
* ({@code SessionManager.resolveAgentSessionId}, via {@code get}/{@code rosterResolved}/
|
||||
* {@code release}) actually calls, and only while the session id is still unknown. Putting the
|
||||
* check anywhere else risks repeating PR #203's mistake: a check that runs before opencode has
|
||||
* written the row it needs, and so never fires.
|
||||
*/
|
||||
private static final class SessionAwareHandle implements PeerHandle {
|
||||
private final PeerHandle delegate;
|
||||
private final OpenCodeSessionDiscovery discovery;
|
||||
private final String cwd;
|
||||
private final FleetConfig.Profile cfg;
|
||||
private final BooleanSupplier discoveryUnavailable;
|
||||
private final AtomicBoolean discoveryUnavailableWarned;
|
||||
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();
|
||||
|
||||
SessionAwareHandle(PeerHandle delegate, OpenCodeSessionDiscovery discovery, String cwd,
|
||||
FleetConfig.Profile cfg,
|
||||
BooleanSupplier discoveryUnavailable,
|
||||
AtomicBoolean discoveryUnavailableWarned) {
|
||||
AtomicBoolean discoveryUnavailableWarned,
|
||||
ExhaustionSink exhaustionSink) {
|
||||
this.delegate = delegate;
|
||||
this.discovery = discovery;
|
||||
this.cwd = cwd;
|
||||
this.cfg = cfg;
|
||||
this.discoveryUnavailable = discoveryUnavailable;
|
||||
this.discoveryUnavailableWarned = discoveryUnavailableWarned;
|
||||
this.exhaustionSink = exhaustionSink;
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -693,7 +763,64 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
// 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).
|
||||
return discovery.sessionIdForDirectory(cwd);
|
||||
String id = discovery.sessionIdForDirectory(cwd);
|
||||
// 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();
|
||||
return id;
|
||||
}
|
||||
|
||||
/**
|
||||
* 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
|
||||
* quarantine a perfectly working profile's credential.
|
||||
*/
|
||||
private void checkModelMatch() {
|
||||
if (modelMismatchReported.get() || cfg.model() == null || cfg.model().isBlank()) {
|
||||
return;
|
||||
}
|
||||
OpenCodeSessionDiscovery.ActualModel actual = discovery.actualModelForDirectory(cwd);
|
||||
if (actual == null) {
|
||||
return; // UNKNOWN evidence — never a mismatch
|
||||
}
|
||||
String[] requestedParts = splitProviderModel(cfg.model());
|
||||
String requestedProvider = requestedParts == null ? null : requestedParts[0];
|
||||
String requestedId = requestedParts == null ? cfg.model() : requestedParts[1];
|
||||
boolean idMatches = requestedId.equals(actual.id());
|
||||
// Compare the provider ONLY when BOTH sides have one. requestedProvider == null covers
|
||||
// a bare profile model with no "/" — the profile never asked for a specific provider.
|
||||
// actual.provider() == null covers opencode's model JSON having an id but no providerID
|
||||
// (a real shape parseModel accepts) — that is missing evidence, not a contradiction, and
|
||||
// acceptance rule 4 says missing evidence is UNKNOWN, never a mismatch. Narrowing the
|
||||
// provider comparison this way keeps the id comparison (the part that actually caught the
|
||||
// xf bug) fully intact — a genuine id mismatch is still caught either way (fleetd #175
|
||||
// review round 2).
|
||||
boolean providerMatches = requestedProvider == null || actual.provider() == null
|
||||
|| requestedProvider.equals(actual.provider());
|
||||
if (idMatches && providerMatches) {
|
||||
return;
|
||||
}
|
||||
if (!modelMismatchReported.compareAndSet(false, true)) {
|
||||
return; // another thread already reported this exact mismatch
|
||||
}
|
||||
String actualDisplay = actual.provider() == null
|
||||
? actual.id() : actual.provider() + "/" + actual.id();
|
||||
log.error("opencode profile '{}' requested model '{}' but the live session is actually "
|
||||
+ "running '{}' — opencode does not fail on an unknown -m, it silently "
|
||||
+ "falls back to a default model, which may be a PAID credential "
|
||||
+ "(fleetd #175); quarantining this profile's credential",
|
||||
cfg.profile(), cfg.model(), actualDisplay);
|
||||
exhaustionSink.onExhausted(delegate.terminalId(),
|
||||
"opencode model mismatch: profile '" + cfg.profile() + "' requested '"
|
||||
+ cfg.model() + "' but the live session is running '" + actualDisplay
|
||||
+ "' (fleetd #175)");
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
@@ -1,9 +1,12 @@
|
||||
package dev.ltms.fleet.member;
|
||||
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.sqlite.SQLiteConfig;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.sql.Connection;
|
||||
@@ -44,6 +47,18 @@ import java.util.concurrent.atomic.AtomicBoolean;
|
||||
final class OpenCodeSessionDiscovery {
|
||||
|
||||
private static final Logger log = LoggerFactory.getLogger(OpenCodeSessionDiscovery.class);
|
||||
private static final ObjectMapper MAPPER = new ObjectMapper();
|
||||
|
||||
/**
|
||||
* The model opencode actually ran a session on, parsed from the {@code session.model} JSON
|
||||
* column (fleetd #175). {@code id} is never null/blank on a non-null {@code ActualModel} —
|
||||
* {@link #parseModel} returns {@code null} instead when {@code id} cannot be determined, so a
|
||||
* caller only ever sees a fully-known record or {@code null} (UNKNOWN). {@code provider} may
|
||||
* still be {@code null} on its own when the profile that requested the session named no
|
||||
* provider prefix, or opencode's JSON omitted {@code providerID}.
|
||||
*/
|
||||
record ActualModel(String provider, String id) {
|
||||
}
|
||||
|
||||
private final Path storageRoot; // e.g. ~/.local/share/opencode (injectable for tests)
|
||||
private final Path databasePath;
|
||||
@@ -117,4 +132,73 @@ final class OpenCodeSessionDiscovery {
|
||||
storageRoot, directory);
|
||||
return null;
|
||||
}
|
||||
|
||||
/**
|
||||
* 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.
|
||||
*
|
||||
* <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.
|
||||
*
|
||||
* @param directory the worker's cwd, as resolved for this spawn
|
||||
* @return the actual model, or {@code null} when unknown
|
||||
*/
|
||||
ActualModel actualModelForDirectory(String directory) {
|
||||
if (directory == null || directory.isBlank()) {
|
||||
return null;
|
||||
}
|
||||
if (!Files.isRegularFile(databasePath)) {
|
||||
// sessionIdForDirectory already WARNs once (shared warnedMissingDatabase) for this
|
||||
// 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";
|
||||
try (Connection connection = openReadOnly();
|
||||
PreparedStatement statement = connection.prepareStatement(sql)) {
|
||||
statement.setString(1, directory);
|
||||
try (ResultSet rows = statement.executeQuery()) {
|
||||
if (rows.next()) {
|
||||
return parseModel(rows.getString("model"));
|
||||
}
|
||||
}
|
||||
} catch (SQLException e) {
|
||||
// Locked/corrupt database, or a `model` column this schema version does not have —
|
||||
// never fatal, and never a mismatch signal. See the class doc above.
|
||||
log.debug("opencode session model unreadable at {}: {}", databasePath, e.toString());
|
||||
return null;
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
/**
|
||||
* Parse opencode's {@code model} column — {@code {"id":"...","providerID":"..."}} — into an
|
||||
* {@link ActualModel}, or {@code null} when {@code json} is null/blank, is not valid JSON, or
|
||||
* parses without a non-blank {@code id}. {@code providerID} may be absent; that alone does not
|
||||
* make the record unknown, since a caller comparing against a profile with no provider prefix
|
||||
* never looks at it.
|
||||
*/
|
||||
private static ActualModel parseModel(String json) {
|
||||
if (json == null || json.isBlank()) {
|
||||
return null;
|
||||
}
|
||||
try {
|
||||
JsonNode node = MAPPER.readTree(json);
|
||||
String id = node.path("id").asText(null);
|
||||
if (id == null || id.isBlank()) {
|
||||
return null;
|
||||
}
|
||||
String provider = node.path("providerID").asText(null);
|
||||
return new ActualModel(provider, id);
|
||||
} catch (IOException e) {
|
||||
log.debug("opencode session model JSON unparseable: {}", e.toString());
|
||||
return null;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -13,6 +13,7 @@ import java.nio.charset.StandardCharsets;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.nio.file.StandardCopyOption;
|
||||
import java.nio.file.attribute.PosixFilePermissions;
|
||||
import java.security.SecureRandom;
|
||||
import java.util.ArrayList;
|
||||
import java.util.HashSet;
|
||||
@@ -161,6 +162,7 @@ public final class GitWorktrees implements Worktrees {
|
||||
} catch (IOException e) {
|
||||
throw new WorktreeException("cannot create worktree root " + root + ": " + e.getMessage(), e);
|
||||
}
|
||||
shareRootWithGroup(root);
|
||||
String wt = path.toAbsolutePath().toString();
|
||||
log.info("adding worktree branch={} path={} base={}", branch, wt, base);
|
||||
removeUserInfoFromHttpsOrigin(repoRoot);
|
||||
@@ -712,6 +714,54 @@ public final class GitWorktrees implements Worktrees {
|
||||
group, repoRoot, worktreePath, touched);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #224: make {@code worktreeRoot} ITSELF group-traversable — established once, here,
|
||||
* where the root is created, never at a use site. {@link #shareWithGroup} shares each worktree
|
||||
* (and the repo's common git dir) with {@link #group}, but never the PARENT directory that
|
||||
* contains every worktree — and under {@code memberHerdrSocket:} the member pane runs as a
|
||||
* different OS user, which needs the execute bit on every ancestor directory to reach anything
|
||||
* underneath, no matter how carefully each child is shared. Without this, a member cannot read
|
||||
* the ephemeral {@code opencode.json} #219 places under this root, cannot reach its own
|
||||
* worktree, and cannot read #213's ZDOTDIR scrub when placed here either.
|
||||
*
|
||||
* <p>No-op — no process spawned — when {@link #group} is null/blank, so behaviour with
|
||||
* {@code worktreeGroup:} unset (today's only live mode) is unchanged. Only {@code root} itself
|
||||
* is touched (non-recursive): each child underneath is shared individually, either by
|
||||
* {@link #shareWithGroup} for a worktree or by the launcher that generates it (fleetd #213/#219)
|
||||
* for a scrub/config directory — sharing this level again would just duplicate that policy in
|
||||
* the wrong layer.
|
||||
*
|
||||
* <p>Fails loudly, naming {@code root}, its mode at the time of the attempt, and {@link #group}:
|
||||
* a member that starts and then cannot see its own checkout is worse than a refused spawn, since
|
||||
* nothing about that failure mode points at a directory's permission bits.
|
||||
*/
|
||||
private void shareRootWithGroup(Path root) {
|
||||
if (group == null) {
|
||||
return;
|
||||
}
|
||||
String mode = currentPosixMode(root);
|
||||
try {
|
||||
shareGroupRunner.apply(new String[]{"chgrp", group, root.toString()});
|
||||
shareGroupRunner.apply(new String[]{"chmod", "g+x", root.toString()});
|
||||
} catch (WorktreeException e) {
|
||||
throw new WorktreeException("cannot make worktree root " + root + " (mode " + mode
|
||||
+ ") group-traversable for group '" + group + "': " + e.getMessage()
|
||||
+ " — the group must exist, and the fleetd operator (" + System.getProperty("user.name")
|
||||
+ ") must be a member of it", e);
|
||||
}
|
||||
log.info("worktreeGroup={} made worktree root {} group-traversable (was mode {})", group, root, mode);
|
||||
}
|
||||
|
||||
/** {@code root}'s current POSIX permission string, or {@code "unknown"} on a filesystem that does
|
||||
* not support POSIX permissions — used only to name the mode in a refusal message. */
|
||||
private static String currentPosixMode(Path root) {
|
||||
try {
|
||||
return PosixFilePermissions.toString(Files.getPosixFilePermissions(root));
|
||||
} catch (IOException | UnsupportedOperationException e) {
|
||||
return "unknown";
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* The repo's <em>common</em> git directory as an absolute path — where {@code objects},
|
||||
* {@code refs} and {@code worktrees} actually live. {@code git rev-parse --git-common-dir}
|
||||
|
||||
@@ -190,49 +190,44 @@ public final class SessionManager implements TurnListener {
|
||||
if (profile != null && !profile.isBlank()) {
|
||||
memberLifecycle.requireSlotFor(memberRole, profile);
|
||||
}
|
||||
MemberLifecycle.SlotReservation reservation = memberLifecycle.reserve(memberRole, profile);
|
||||
String launchProfile = reservation == null ? profile : reservation.profile();
|
||||
if (wt == null) {
|
||||
// CB-557: the role must ride on the SpawnRequest, not stay a local. The launcher needs it
|
||||
// to pick the profile out of that role's pool and to label the tab; a role kept only on
|
||||
// the MemberSession is recorded after the spawn it was supposed to steer.
|
||||
SpawnRequest req = new SpawnRequest(profile, requestedCwd, callerCwd, sessionName, resumeSessionId, memberRole);
|
||||
SpawnRequest req = new SpawnRequest(launchProfile, requestedCwd, callerCwd, sessionName, resumeSessionId, memberRole);
|
||||
PeerHandle handle;
|
||||
boolean bound = false;
|
||||
try {
|
||||
handle = launcher.spawn(req);
|
||||
String resolvedProfile = resolveProfile(handle, launchProfile);
|
||||
String cwd = launcher.effectiveCwd(new SpawnRequest(resolvedProfile, requestedCwd, callerCwd));
|
||||
long now = nowNanos.getAsLong();
|
||||
MemberRole actualRole = acquired(memberRole, resolvedProfile, handle.terminalId(), reservation);
|
||||
bound = reservation == null || actualRole == MemberRole.ARCHITECT;
|
||||
MemberSession session = new MemberSession(
|
||||
handle.id(), handle.terminalId(), resolvedProfile, actualRole, cwd, ownerTerminal, now, now, 0,
|
||||
MemberSession.State.SPAWNING, null, null, handle.charterReceipt(), handle.agentSessionId());
|
||||
registry.put(handle.id(), session);
|
||||
handles.put(handle.id(), handle);
|
||||
log.debug("acquired session id={} terminal={} profile={} owner={}",
|
||||
handle.id(), handle.terminalId(), session.profile(), session.ownerTerminal());
|
||||
notifyAcquired(session.terminalId());
|
||||
return session;
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("spawn failed for profile={} role={}: {}", profile, memberRole, e.getMessage());
|
||||
if (!bound) memberLifecycle.release(reservation);
|
||||
log.warn("spawn failed for profile={} role={}: {}", launchProfile, memberRole, e.getMessage());
|
||||
throw e;
|
||||
}
|
||||
String resolvedProfile = resolveProfile(handle, profile);
|
||||
String cwd = launcher.effectiveCwd(new SpawnRequest(resolvedProfile, requestedCwd, callerCwd));
|
||||
long now = nowNanos.getAsLong();
|
||||
// CB-619: bind (or fail to bind) BEFORE the session is recorded, and store whatever role
|
||||
// this call actually returns — never the requested memberRole — so the session's role,
|
||||
// what GET /members and fleet_list report, is never a lie about what this terminal holds.
|
||||
MemberRole actualRole = memberLifecycle.acquired(memberRole, resolvedProfile, handle.terminalId());
|
||||
MemberSession session = new MemberSession(
|
||||
handle.id(),
|
||||
handle.terminalId(),
|
||||
resolvedProfile,
|
||||
actualRole,
|
||||
cwd,
|
||||
ownerTerminal,
|
||||
now,
|
||||
now,
|
||||
0,
|
||||
MemberSession.State.SPAWNING,
|
||||
null,
|
||||
null,
|
||||
handle.charterReceipt(),
|
||||
handle.agentSessionId());
|
||||
registry.put(handle.id(), session);
|
||||
handles.put(handle.id(), handle);
|
||||
log.debug("acquired session id={} terminal={} profile={} owner={}",
|
||||
handle.id(), handle.terminalId(), session.profile(), session.ownerTerminal());
|
||||
notifyAcquired(session.terminalId());
|
||||
return session;
|
||||
}
|
||||
return acquireWithWorktree(profile, memberRole, requestedCwd, callerCwd, ownerTerminal, wt,
|
||||
sessionName, resumeSessionId);
|
||||
try {
|
||||
return acquireWithWorktree(launchProfile, memberRole, requestedCwd, callerCwd, ownerTerminal, wt,
|
||||
sessionName, resumeSessionId, reservation);
|
||||
} catch (RuntimeException e) {
|
||||
memberLifecycle.release(reservation);
|
||||
throw e;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -478,7 +473,8 @@ public final class SessionManager implements TurnListener {
|
||||
|
||||
private MemberSession acquireWithWorktree(String profile, MemberRole memberRole, String requestedCwd, String callerCwd,
|
||||
String ownerTerminal, WorktreeRequest wt,
|
||||
String sessionName, String resumeSessionId) {
|
||||
String sessionName, String resumeSessionId,
|
||||
MemberLifecycle.SlotReservation reservation) {
|
||||
String preResolvedProfile = (profile == null || profile.isBlank())
|
||||
? launcher.defaultProfile() : profile;
|
||||
// CB-507: resolve through the launcher's CB-112 chain (requested → profile cwd → caller →
|
||||
@@ -521,7 +517,7 @@ public final class SessionManager implements TurnListener {
|
||||
long now = nowNanos.getAsLong();
|
||||
// CB-619: see the no-worktree path above — bind before recording, and store the returned
|
||||
// actual role, so this session's role is never a lie about what it actually holds.
|
||||
MemberRole actualRole = memberLifecycle.acquired(memberRole, resolvedProfile, handle.terminalId());
|
||||
MemberRole actualRole = acquired(memberRole, resolvedProfile, handle.terminalId(), reservation);
|
||||
MemberSession session = new MemberSession(
|
||||
handle.id(),
|
||||
handle.terminalId(),
|
||||
@@ -545,6 +541,18 @@ public final class SessionManager implements TurnListener {
|
||||
return session;
|
||||
}
|
||||
|
||||
private MemberRole acquired(MemberRole role, String profile, String terminal,
|
||||
MemberLifecycle.SlotReservation reservation) {
|
||||
if (reservation == null) {
|
||||
return memberLifecycle.acquired(role, profile, terminal);
|
||||
}
|
||||
if (memberLifecycle.bind(reservation, terminal)) {
|
||||
return role;
|
||||
}
|
||||
memberLifecycle.release(reservation);
|
||||
return memberLifecycle.acquired(role, profile, terminal);
|
||||
}
|
||||
|
||||
private String slug(String raw) {
|
||||
return raw == null ? "ticket" : raw.toLowerCase().replaceAll("[^a-z0-9]+", "-").replaceAll("^-+|-+$", "");
|
||||
}
|
||||
|
||||
@@ -114,6 +114,68 @@ class EnvAllowListScrubTest {
|
||||
assertNull(EnvAllowListScrub.readReport(dir));
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #224 criterion 5 (a gap the #221 reviewer flagged): {@link EnvAllowListScrub#shareWithGroup}
|
||||
* lists one FLAT level of {@code dir} and shares every file it finds there — which does cover the
|
||||
* OPTIONAL files a caller may or may not have written before calling it ({@code member-charter.md},
|
||||
* {@code ide-rules.md} — both written by {@code OpenCodeLauncher#writeConfig}), but nothing
|
||||
* pinned that down.
|
||||
* Without this test, a future change that writes a file AFTER the sharing call, or into a
|
||||
* subdirectory, would pass every existing test while quietly leaving that file unreadable to a
|
||||
* different-uid member.
|
||||
*
|
||||
* <p>This drives {@code shareWithGroup} directly against a directory holding several flat files —
|
||||
* not only {@code opencode.json}, but also the two optional ones named above plus a third,
|
||||
* unrelated file, so the assertion is "every flat file", not "the two files someone thought of".
|
||||
*/
|
||||
@Test
|
||||
void shareWithGroupCoversEveryFlatFileIncludingTheOptionalOnes(@TempDir Path dir) throws Exception {
|
||||
String group = currentUserGroup();
|
||||
Files.writeString(dir.resolve("opencode.json"), "{}\n");
|
||||
Files.writeString(dir.resolve("member-charter.md"), "# charter\n");
|
||||
Files.writeString(dir.resolve("ide-rules.md"), "# ide rules\n");
|
||||
Files.writeString(dir.resolve("another-flat-file.txt"), "unrelated\n");
|
||||
|
||||
EnvAllowListScrub.shareWithGroup(dir, group);
|
||||
|
||||
assertEquals("rwxr-x---", java.nio.file.attribute.PosixFilePermissions.toString(
|
||||
Files.getPosixFilePermissions(dir)),
|
||||
"the directory itself must be group-traversable+readable, owner-only writable");
|
||||
for (String name : List.of("opencode.json", "member-charter.md", "ide-rules.md", "another-flat-file.txt")) {
|
||||
Path file = dir.resolve(name);
|
||||
assertEquals("rw-r-----", java.nio.file.attribute.PosixFilePermissions.toString(
|
||||
Files.getPosixFilePermissions(file)),
|
||||
name + " must be group-readable, never group-writable");
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* The current process's REAL primary group — resolved via {@code id -gn}, never by reading a
|
||||
* directory's owning group (fleetd #225: that reads wherever Maven happened to be started from,
|
||||
* not the process's own group, and the two diverge outside a home checkout). Skips (never fails)
|
||||
* when {@code id} is unavailable or its primary group cannot be resolved on this host.
|
||||
*/
|
||||
private static String currentUserGroup() {
|
||||
String out;
|
||||
boolean ok;
|
||||
try {
|
||||
Process p = new ProcessBuilder("id", "-gn").redirectErrorStream(true).start();
|
||||
String raw = new String(p.getInputStream().readAllBytes(), StandardCharsets.UTF_8).trim();
|
||||
out = raw;
|
||||
ok = p.waitFor(5, java.util.concurrent.TimeUnit.SECONDS) && p.exitValue() == 0 && !raw.isBlank();
|
||||
} catch (IOException e) {
|
||||
out = null;
|
||||
ok = false;
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
out = null;
|
||||
ok = false;
|
||||
}
|
||||
assumeTrue(ok, "cannot resolve this process's real primary group via `id -gn` on this host "
|
||||
+ "— skipping a POSIX-group-dependent test rather than failing it");
|
||||
return out;
|
||||
}
|
||||
|
||||
/**
|
||||
* The same equality, for a shell that is INTERACTIVE but NOT a login shell — the shape herdr
|
||||
* opens on Linux.
|
||||
|
||||
+31
-6
@@ -18,7 +18,6 @@ import org.slf4j.LoggerFactory;
|
||||
import java.io.IOException;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.nio.file.attribute.PosixFileAttributeView;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
@@ -513,11 +512,37 @@ class HerdrPeerLauncherAllowListWiringTest {
|
||||
+ "java.io.tmpdir, unchanged from before this fix: " + dir);
|
||||
}
|
||||
|
||||
/** The current process's own primary group — resolvable on whatever host runs this test. */
|
||||
private static String currentUserGroup() throws IOException {
|
||||
PosixFileAttributeView view = Files.getFileAttributeView(Path.of("."), PosixFileAttributeView.class);
|
||||
assumeTrue(view != null, "this host's filesystem does not support POSIX group ownership");
|
||||
return view.readAttributes().group().getName();
|
||||
/**
|
||||
* The current process's REAL primary group — resolved via {@code id -gn}, never by reading a
|
||||
* directory's owning group (fleetd #225). Those two coincide only by accident: a directory's
|
||||
* group is whichever group happened to own the path Maven was started from — {@code staff} in
|
||||
* a home checkout, {@code wheel} under {@code /private/tmp} on macOS — and the fix-up this test
|
||||
* exercises then fails for real when the operator is not a member of that borrowed group,
|
||||
* exactly the case {@code assumeTrue(view != null, ...)} never covered (it only detects a
|
||||
* filesystem with no POSIX groups at all, not a resolvable-but-wrong one). Skips (never fails)
|
||||
* when {@code id} is unavailable or its primary group cannot be resolved on this host.
|
||||
*/
|
||||
private static String currentUserGroup() {
|
||||
String out;
|
||||
boolean ok;
|
||||
try {
|
||||
Process p = new ProcessBuilder("id", "-gn").redirectErrorStream(true).start();
|
||||
try (java.io.BufferedReader r = new java.io.BufferedReader(
|
||||
new java.io.InputStreamReader(p.getInputStream(), java.nio.charset.StandardCharsets.UTF_8))) {
|
||||
out = r.lines().collect(java.util.stream.Collectors.joining("\n")).trim();
|
||||
}
|
||||
ok = p.waitFor(5, java.util.concurrent.TimeUnit.SECONDS) && p.exitValue() == 0 && !out.isBlank();
|
||||
} catch (IOException e) {
|
||||
out = null;
|
||||
ok = false;
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
out = null;
|
||||
ok = false;
|
||||
}
|
||||
assumeTrue(ok, "cannot resolve this process's real primary group via `id -gn` on this host "
|
||||
+ "— skipping a POSIX-group-dependent test rather than failing it");
|
||||
return out;
|
||||
}
|
||||
|
||||
/** Spawn once through the real launcher path, capturing every INFO+ line this class logs. */
|
||||
|
||||
@@ -10,12 +10,15 @@ import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
import dev.ltms.fleet.inject.ExhaustionSink;
|
||||
import dev.ltms.fleet.peer.Capability;
|
||||
import dev.ltms.fleet.peer.CharterReceipt;
|
||||
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.session.MemberSession;
|
||||
import dev.ltms.fleet.session.SessionManager;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
import org.slf4j.LoggerFactory;
|
||||
@@ -23,10 +26,11 @@ import org.slf4j.LoggerFactory;
|
||||
import java.io.IOException;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.nio.file.attribute.PosixFileAttributeView;
|
||||
import java.nio.file.attribute.PosixFilePermissions;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Optional;
|
||||
import java.util.concurrent.ExecutorService;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.Future;
|
||||
@@ -657,11 +661,37 @@ class OpenCodeLauncherTest {
|
||||
0, System::currentTimeMillis, () -> { }, configRoot, discoveryRoot, null, null, config);
|
||||
}
|
||||
|
||||
/** The current process's own primary group — resolvable on whatever host runs this test. */
|
||||
private static String currentUserGroup() throws IOException {
|
||||
PosixFileAttributeView view = Files.getFileAttributeView(Path.of("."), PosixFileAttributeView.class);
|
||||
assumeTrue(view != null, "this host's filesystem does not support POSIX group ownership");
|
||||
return view.readAttributes().group().getName();
|
||||
/**
|
||||
* The current process's REAL primary group — resolved via {@code id -gn}, never by reading a
|
||||
* directory's owning group (fleetd #225). Those two coincide only by accident: a directory's
|
||||
* group is whichever group happened to own the path Maven was started from — {@code staff} in
|
||||
* a home checkout, {@code wheel} under {@code /private/tmp} on macOS — and the fix-up this test
|
||||
* exercises then fails for real when the operator is not a member of that borrowed group,
|
||||
* exactly the case {@code assumeTrue(view != null, ...)} never covered (it only detects a
|
||||
* filesystem with no POSIX groups at all, not a resolvable-but-wrong one). Skips (never fails)
|
||||
* when {@code id} is unavailable or its primary group cannot be resolved on this host.
|
||||
*/
|
||||
private static String currentUserGroup() {
|
||||
String out;
|
||||
boolean ok;
|
||||
try {
|
||||
Process p = new ProcessBuilder("id", "-gn").redirectErrorStream(true).start();
|
||||
try (java.io.BufferedReader r = new java.io.BufferedReader(
|
||||
new java.io.InputStreamReader(p.getInputStream(), java.nio.charset.StandardCharsets.UTF_8))) {
|
||||
out = r.lines().collect(java.util.stream.Collectors.joining("\n")).trim();
|
||||
}
|
||||
ok = p.waitFor(5, java.util.concurrent.TimeUnit.SECONDS) && p.exitValue() == 0 && !out.isBlank();
|
||||
} catch (IOException e) {
|
||||
out = null;
|
||||
ok = false;
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
out = null;
|
||||
ok = false;
|
||||
}
|
||||
assumeTrue(ok, "cannot resolve this process's real primary group via `id -gn` on this host "
|
||||
+ "— skipping a POSIX-group-dependent test rather than failing it");
|
||||
return out;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -811,4 +841,256 @@ class OpenCodeLauncherTest {
|
||||
assertTrue(warnings.get(0).contains("memberHerdrSocket"),
|
||||
"the WARN must name memberHerdrSocket as the reason — got: " + warnings.get(0));
|
||||
}
|
||||
|
||||
// --- fleetd #175: opencode silently substitutes a model on an unknown -m flag ----------------
|
||||
|
||||
private static OpenCodeLauncher serviceWithSink(FakeHerdr herdr, Path configRoot, Path discoveryRoot,
|
||||
FleetConfig.Profile cfg, ExhaustionSink sink) {
|
||||
return new OpenCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null,
|
||||
0, System::currentTimeMillis, () -> { }, configRoot, discoveryRoot,
|
||||
null, null, null, sink);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #175 acceptance criterion 1: the check must fire on the REAL late-resolve path —
|
||||
* {@code SessionManager}'s {@code get}/{@code roster} re-polling a retained {@link PeerHandle}
|
||||
* (fleetd #209) — not on a handle built and queried directly. A handle whose row does not exist
|
||||
* yet, then does, proves the check runs exactly where production runs it: PR #203 shipped a
|
||||
* check that ran inside {@code spawn()}, before opencode had written the row, and every test
|
||||
* passed anyway because none of them drove it through this path. This test would have caught
|
||||
* that: {@code sessions.get(...)} before the row exists must show no mismatch, and the SAME
|
||||
* call, re-driven after the row appears, must be what fires the sink.
|
||||
*/
|
||||
@Test
|
||||
void theRealSessionManagerLateResolvePathCatchesAModelMismatch(@TempDir Path configRoot,
|
||||
@TempDir Path discRoot) throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
// xf's real shape (fleetd #175): weight:80, model "opencode/nemotron-3-ultra-free", no
|
||||
// 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);
|
||||
OpenCodeLauncher launcher = serviceWithSink(herdr, configRoot, discRoot, cfg, sink);
|
||||
|
||||
SessionManager sessions = new SessionManager(launcher);
|
||||
MemberSession acquired = sessions.acquire(cfg.profile(), "/work/dir", null, null);
|
||||
|
||||
// Real late-resolve path, driven BEFORE opencode has written its session row — same shape
|
||||
// as production the instant a pane goes ready.
|
||||
Optional<MemberSession> beforeRow = sessions.get(acquired.paneId());
|
||||
assertTrue(beforeRow.isPresent());
|
||||
assertNull(beforeRow.get().agentSessionId(), "no opencode row yet");
|
||||
assertTrue(exhausted.isEmpty(), "no row yet → nothing to compare, the sink must stay silent");
|
||||
|
||||
// opencode writes its row late, running gpt-5.6-sol (a PAID credential) instead of the
|
||||
// withdrawn free model the profile actually asked for — the exact fleetd #175 scenario.
|
||||
OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "ses_x", "/work/dir", 1000L,
|
||||
"{\"id\":\"gpt-5.6-sol\",\"providerID\":\"openai\"}");
|
||||
|
||||
// Drive the SAME real late-resolve path again: sessions.get() -> resolveAgentSessionId ->
|
||||
// the retained PeerHandle's agentSessionId() -> discovery -> checkModelMatch, all in one
|
||||
// call, unmodified SessionManager code from fleetd #209.
|
||||
Optional<MemberSession> afterRow = sessions.get(acquired.paneId());
|
||||
assertEquals("ses_x", afterRow.get().agentSessionId(),
|
||||
"the id itself still resolves correctly alongside the model check");
|
||||
|
||||
assertEquals(1, exhausted.size(), "the mismatch must fire exactly once through the real path");
|
||||
assertTrue(exhausted.get(0).contains(cfg.model()), "reports the requested model: " + exhausted.get(0));
|
||||
assertTrue(exhausted.get(0).contains("gpt-5.6-sol"), "reports the actual model: " + exhausted.get(0));
|
||||
}
|
||||
|
||||
/** THE TRAP, row 1: a provider-prefixed request matching the DB's id AND provider is a match. */
|
||||
@Test
|
||||
void aProviderPrefixedModelMatchingBothIdAndProviderIsNotAMismatch(@TempDir Path configRoot,
|
||||
@TempDir Path discRoot) throws Exception {
|
||||
List<String> 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));
|
||||
OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "ses_x", "/work/dir", 1000L,
|
||||
"{\"id\":\"gpt-5.6-terra\",\"providerID\":\"openai\"}");
|
||||
|
||||
assertEquals("ses_x", handle.agentSessionId());
|
||||
assertTrue(exhausted.isEmpty(), "id and provider both match → never a mismatch: " + exhausted);
|
||||
}
|
||||
|
||||
/** THE TRAP, row 2: same shape as row 1 with a different provider/id pair (the gx profile). */
|
||||
@Test
|
||||
void aGxProviderPrefixedModelMatchingBothIdAndProviderIsNotAMismatch(@TempDir Path configRoot,
|
||||
@TempDir Path discRoot) throws Exception {
|
||||
List<String> exhausted = new ArrayList<>();
|
||||
ExhaustionSink sink = (target, reason) -> 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));
|
||||
OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "ses_x", "/work/dir", 1000L,
|
||||
"{\"id\":\"deepseek-v4-flash\",\"providerID\":\"gx\"}");
|
||||
|
||||
assertEquals("ses_x", handle.agentSessionId());
|
||||
assertTrue(exhausted.isEmpty(), "id and provider both match → never a mismatch: " + exhausted);
|
||||
}
|
||||
|
||||
/**
|
||||
* Incomplete evidence, not THE TRAP's provider mismatch: opencode's {@code model} JSON had an
|
||||
* {@code id} but no {@code providerID} at all (a real shape {@code parseModel} accepts — see
|
||||
* {@code OpenCodeSessionDiscoveryTest}). The id matches; the provider dimension is simply
|
||||
* unknown, not contradicted. A provider-prefixed profile must NOT be quarantined on this —
|
||||
* that would quarantine on incomplete evidence, which acceptance rule 4 forbids.
|
||||
*/
|
||||
@Test
|
||||
void aMissingProviderIdInTheEvidenceIsUnknownNotAMismatchWhenTheIdMatches(
|
||||
@TempDir Path configRoot, @TempDir Path discRoot) throws Exception {
|
||||
List<String> 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));
|
||||
OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "ses_x", "/work/dir", 1000L,
|
||||
"{\"id\":\"gpt-5.6-terra\"}");
|
||||
|
||||
assertEquals("ses_x", handle.agentSessionId());
|
||||
assertTrue(exhausted.isEmpty(),
|
||||
"id matches and provider is simply unknown (absent), never a mismatch: " + exhausted);
|
||||
}
|
||||
|
||||
/**
|
||||
* The other half of the same fix: a missing {@code providerID} must NOT blind the check to a
|
||||
* genuine id mismatch. This is what proves the fix narrows the comparison rather than switching
|
||||
* the whole check off whenever {@code providerID} happens to be absent.
|
||||
*/
|
||||
@Test
|
||||
void aMissingProviderIdInTheEvidenceStillCatchesARealIdMismatch(
|
||||
@TempDir Path configRoot, @TempDir Path discRoot) throws Exception {
|
||||
List<String> 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));
|
||||
OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "ses_x", "/work/dir", 1000L,
|
||||
"{\"id\":\"gpt-5.6-sol\"}");
|
||||
|
||||
assertEquals("ses_x", handle.agentSessionId());
|
||||
assertEquals(1, exhausted.size(),
|
||||
"the id genuinely differs, so this must still quarantine even with providerID absent: "
|
||||
+ exhausted);
|
||||
assertTrue(exhausted.get(0).contains("gpt-5.6-sol") && exhausted.get(0).contains(cfg.model()),
|
||||
"reports both the requested and actual model: " + exhausted.get(0));
|
||||
}
|
||||
|
||||
/**
|
||||
* THE TRAP, row 3: a profile that names no provider prefix (bare {@code "deepseek-v4-flash"})
|
||||
* must match on id alone — the profile never asked for a specific provider, so opencode
|
||||
* resolving it to {@code gx} is not evidence of anything wrong.
|
||||
*/
|
||||
@Test
|
||||
void aBareModelWithNoProviderPrefixMatchesOnIdAloneAndIsNotAMismatch(@TempDir Path configRoot,
|
||||
@TempDir Path discRoot) throws Exception {
|
||||
List<String> exhausted = new ArrayList<>();
|
||||
ExhaustionSink sink = (target, reason) -> 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));
|
||||
OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "ses_x", "/work/dir", 1000L,
|
||||
"{\"id\":\"deepseek-v4-flash\",\"providerID\":\"gx\"}");
|
||||
|
||||
assertEquals("ses_x", handle.agentSessionId());
|
||||
assertTrue(exhausted.isEmpty(),
|
||||
"no provider was requested, so opencode's own provider resolution is not a mismatch: "
|
||||
+ exhausted);
|
||||
}
|
||||
|
||||
/**
|
||||
* THE TRAP, row 4 (not a row in the table, the reason the table exists): a mismatched id, with
|
||||
* the profile's requested and the actual model both named in the ERROR and the sink's reason —
|
||||
* the exact fleetd #175 scenario (a withdrawn model silently falls back to a paid one).
|
||||
*/
|
||||
@Test
|
||||
void aRealIdMismatchLogsAnErrorNamingBothModelsAndQuarantinesThroughTheSink(
|
||||
@TempDir Path configRoot, @TempDir Path discRoot) throws Exception {
|
||||
List<String> exhausted = new ArrayList<>();
|
||||
ExhaustionSink sink = (target, reason) -> exhausted.add(target + "|" + reason);
|
||||
FleetConfig.Profile cfg = opencodeCfg("opencode/nemotron-3-ultra-free", null, null);
|
||||
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(OpenCodeLauncher.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
PeerHandle handle;
|
||||
try {
|
||||
handle = serviceWithSink(new FakeHerdr(), configRoot, discRoot, cfg, sink)
|
||||
.spawn(new SpawnRequest(null, "/work/dir", null));
|
||||
OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "ses_x", "/work/dir", 1000L,
|
||||
"{\"id\":\"gpt-5.6-sol\",\"providerID\":\"openai\"}");
|
||||
assertEquals("ses_x", handle.agentSessionId());
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
}
|
||||
|
||||
assertEquals(1, exhausted.size(), "a real id mismatch must reach the sink exactly once");
|
||||
assertTrue(exhausted.get(0).contains("opencode/nemotron-3-ultra-free"),
|
||||
"the sink reason must name the requested model: " + exhausted.get(0));
|
||||
assertTrue(exhausted.get(0).contains("gpt-5.6-sol"),
|
||||
"the sink reason must name the actual model: " + exhausted.get(0));
|
||||
|
||||
List<ILoggingEvent> errors = appender.list.stream()
|
||||
.filter(e -> e.getLevel() == Level.ERROR)
|
||||
.toList();
|
||||
assertEquals(1, errors.size(), "exactly one ERROR for the mismatch: " + appender.list);
|
||||
String message = errors.get(0).getFormattedMessage();
|
||||
assertTrue(message.contains("opencode/nemotron-3-ultra-free") && message.contains("gpt-5.6-sol")
|
||||
&& message.contains(cfg.profile()),
|
||||
"the ERROR must name the requested model, the actual model, AND the profile: " + message);
|
||||
|
||||
// The check runs at most once per handle even if agentSessionId() is polled again.
|
||||
handle.agentSessionId();
|
||||
assertEquals(1, exhausted.size(), "no duplicate quarantine on a repeated call");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #175's central safety rule: absent, empty, or unparseable model evidence is UNKNOWN,
|
||||
* never a mismatch — it must never quarantine a working profile. Covers every "no real
|
||||
* evidence yet" shape: no row at all, a row with a null model, and a row with unparseable JSON.
|
||||
*/
|
||||
@Test
|
||||
void unknownOrUnparseableModelEvidenceNeverQuarantines(@TempDir Path configRoot,
|
||||
@TempDir Path discRoot) throws Exception {
|
||||
List<String> 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));
|
||||
|
||||
// No row yet at all.
|
||||
assertNull(handle.agentSessionId());
|
||||
assertTrue(exhausted.isEmpty(), "no row yet is UNKNOWN, not a mismatch: " + exhausted);
|
||||
|
||||
// A row exists (for a DIFFERENT directory) so a poll on ours still finds nothing.
|
||||
OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "ses_other", "/work/other", 500L,
|
||||
"{\"id\":\"gpt-5.6-sol\",\"providerID\":\"openai\"}");
|
||||
assertNull(handle.agentSessionId());
|
||||
assertTrue(exhausted.isEmpty(), "a row for a different directory is UNKNOWN, not a mismatch: "
|
||||
+ exhausted);
|
||||
}
|
||||
|
||||
/**
|
||||
* A profile with no {@code model:} configured has nothing to compare against — the check must
|
||||
* stay silent no matter what opencode actually ran, since there is no requested value to be
|
||||
* wrong about.
|
||||
*/
|
||||
@Test
|
||||
void aProfileWithNoConfiguredModelIsNeverCheckedForAMismatch(@TempDir Path configRoot,
|
||||
@TempDir Path discRoot) throws Exception {
|
||||
List<String> exhausted = new ArrayList<>();
|
||||
ExhaustionSink sink = (target, reason) -> 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));
|
||||
OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "ses_x", "/work/dir", 1000L,
|
||||
"{\"id\":\"anything-at-all\",\"providerID\":\"anyone\"}");
|
||||
|
||||
assertEquals("ses_x", handle.agentSessionId());
|
||||
assertTrue(exhausted.isEmpty(), "no model configured → nothing to compare: " + exhausted);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -27,19 +27,30 @@ class OpenCodeSessionDiscoveryTest {
|
||||
* and insert one row. Static so {@link OpenCodeLauncherTest} can reuse it.
|
||||
*/
|
||||
static void writeRecord(Path root, String id, String directory, long timeUpdated) throws Exception {
|
||||
writeRecord(root, id, directory, timeUpdated, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* Same as {@link #writeRecord(Path, String, String, long)}, plus the {@code model} column
|
||||
* fleetd #175 reads — the raw JSON opencode writes, e.g.
|
||||
* {@code {"id":"gpt-5.6-terra","providerID":"openai"}}. {@code modelJson} may be {@code null}.
|
||||
*/
|
||||
static void writeRecord(Path root, String id, String directory, long timeUpdated, String modelJson)
|
||||
throws Exception {
|
||||
Path db = root.resolve("opencode.db");
|
||||
try (Connection connection = DriverManager.getConnection("jdbc:sqlite:" + db)) {
|
||||
try (Statement statement = connection.createStatement()) {
|
||||
statement.execute("CREATE TABLE IF NOT EXISTS session ("
|
||||
+ "id TEXT PRIMARY KEY, directory TEXT, time_updated INTEGER)");
|
||||
+ "id TEXT PRIMARY KEY, directory TEXT, time_updated INTEGER, model TEXT)");
|
||||
}
|
||||
// Bound parameters, not string interpolation: the class under test uses a
|
||||
// PreparedStatement, and a hand-escaped INSERT here is a pattern someone copies out.
|
||||
try (PreparedStatement insert = connection.prepareStatement(
|
||||
"INSERT INTO session (id, directory, time_updated) VALUES (?, ?, ?)")) {
|
||||
"INSERT INTO session (id, directory, time_updated, model) VALUES (?, ?, ?, ?)")) {
|
||||
insert.setString(1, id);
|
||||
insert.setString(2, directory);
|
||||
insert.setLong(3, timeUpdated);
|
||||
insert.setString(4, modelJson);
|
||||
insert.executeUpdate();
|
||||
}
|
||||
}
|
||||
@@ -134,4 +145,96 @@ class OpenCodeSessionDiscoveryTest {
|
||||
assertNull(new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"),
|
||||
"an unreadable database resolves to null, not an exception");
|
||||
}
|
||||
|
||||
// --- fleetd #175: actualModelForDirectory / 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");
|
||||
|
||||
assertNotNull(actual, "a well-formed model JSON parses");
|
||||
assertEquals("gpt-5.6-terra", actual.id());
|
||||
assertEquals("openai", actual.provider());
|
||||
}
|
||||
|
||||
@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\"}");
|
||||
|
||||
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");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aNonMatchingDirectoryYieldsUnknownModelRatherThanAMismatch(@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");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aNullModelColumnYieldsUnknownWithoutThrowing(@TempDir Path root) throws Exception {
|
||||
writeRecord(root, "ses_aaa", "/w/a", 1000L, null);
|
||||
|
||||
assertNull(new OpenCodeSessionDiscovery(root).actualModelForDirectory("/w/a"),
|
||||
"a row with no model value yet is unknown, not a mismatch");
|
||||
}
|
||||
|
||||
@Test
|
||||
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"),
|
||||
"JSON that fails to parse resolves to unknown, never an exception");
|
||||
}
|
||||
|
||||
@Test
|
||||
void modelJsonMissingIdYieldsUnknown(@TempDir Path root) throws Exception {
|
||||
writeRecord(root, "ses_aaa", "/w/a", 1000L, "{\"providerID\":\"openai\"}");
|
||||
|
||||
assertNull(new OpenCodeSessionDiscovery(root).actualModelForDirectory("/w/a"),
|
||||
"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"));
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #175 site 2: a schema that predates the {@code model} column entirely — the exact
|
||||
* 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.
|
||||
*/
|
||||
@Test
|
||||
void aMissingModelColumnYieldsUnknownButIdResolutionStillWorks(@TempDir Path root) throws Exception {
|
||||
Path db = root.resolve("opencode.db");
|
||||
try (Connection connection = DriverManager.getConnection("jdbc:sqlite:" + db)) {
|
||||
try (Statement statement = connection.createStatement()) {
|
||||
statement.execute("CREATE TABLE session (id TEXT PRIMARY KEY, directory TEXT, "
|
||||
+ "time_updated INTEGER)");
|
||||
}
|
||||
try (PreparedStatement insert = connection.prepareStatement(
|
||||
"INSERT INTO session (id, directory, time_updated) VALUES (?, ?, ?)")) {
|
||||
insert.setString(1, "ses_aaa");
|
||||
insert.setString(2, "/w/a");
|
||||
insert.setLong(3, 1000L);
|
||||
insert.executeUpdate();
|
||||
}
|
||||
}
|
||||
|
||||
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"),
|
||||
"no model column → unknown, not a throw and not a mismatch");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1099,4 +1099,82 @@ class GitWorktreesTest {
|
||||
assertTrue(e.getMessage().contains("cb185-nonexistent-group-zz"),
|
||||
"exception must name the missing/refused group: " + e.getMessage());
|
||||
}
|
||||
|
||||
// --- fleetd #224: worktreeRoot itself must be group-traversable, established in add() ---------
|
||||
|
||||
/**
|
||||
* {@code add} must make the worktree ROOT itself group-traversable when a group is configured —
|
||||
* established once here, where the root is created, never at a use site (never inside a
|
||||
* launcher). A recording runner stands in for chgrp/chmod, the same seam
|
||||
* {@link #shareWithGroupRunsConfigThenChgrpChmodSetgidPerPath} uses for the per-worktree share.
|
||||
*/
|
||||
@Test
|
||||
void addSharesWorktreeRootWithGroupWhenConfigured(@TempDir Path tmp) throws Exception {
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
Path root = tmp.resolve("wts");
|
||||
List<List<String>> recorded = new java.util.ArrayList<>();
|
||||
java.util.function.Function<String[], String> recordingRunner = cmd -> {
|
||||
recorded.add(joined(cmd));
|
||||
return "";
|
||||
};
|
||||
GitWorktrees gitWorktrees = new GitWorktrees(root.toString(), "devteam", _ -> {}, recordingRunner);
|
||||
|
||||
gitWorktrees.add(repo.toString(), "cb-224-branch", "HEAD");
|
||||
|
||||
assertTrue(recorded.contains(List.of("chgrp", "devteam", root.toString())),
|
||||
"the worktree root itself must be chgrp'd to the configured group: " + recorded);
|
||||
assertTrue(recorded.contains(List.of("chmod", "g+x", root.toString())),
|
||||
"the worktree root itself must gain group-execute so a different-uid member can "
|
||||
+ "traverse into it: " + recorded);
|
||||
}
|
||||
|
||||
/** {@code worktreeGroup} unset (today's only live mode) ⇒ {@code add} spawns no share process
|
||||
* for the root at all — behaviour must be byte-identical to before fleetd #224. */
|
||||
@Test
|
||||
void addSharesNothingForTheRootWhenNoGroupConfigured(@TempDir Path tmp) throws Exception {
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
Path root = tmp.resolve("wts");
|
||||
List<List<String>> recorded = new java.util.ArrayList<>();
|
||||
java.util.function.Function<String[], String> recordingRunner = cmd -> {
|
||||
recorded.add(joined(cmd));
|
||||
return "";
|
||||
};
|
||||
GitWorktrees gitWorktrees = new GitWorktrees(root.toString(), null, _ -> {}, recordingRunner);
|
||||
|
||||
gitWorktrees.add(repo.toString(), "cb-224-nogroup", "HEAD");
|
||||
|
||||
assertTrue(recorded.isEmpty(), "no group configured must spawn no share process for the "
|
||||
+ "root at all: " + recorded);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #224, acceptance criterion 2/3: when the root cannot be made group-traversable — here
|
||||
* because the configured group does not exist, the same real-failure shape
|
||||
* {@link #shareWithGroupThrowsNamingTheGroupWhenChgrpFails} drives for the per-worktree share —
|
||||
* the spawn is refused with a message naming the root, its current mode, and the group. This
|
||||
* drives the REAL {@code chgrp} (no recording runner), and the refusal happens before {@code git
|
||||
* worktree add} ever runs, so no partial worktree is left behind either.
|
||||
*/
|
||||
@Test
|
||||
void addRefusesWhenWorktreeRootCannotBeMadeGroupTraversable(@TempDir Path tmp) throws Exception {
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
Path root = tmp.resolve("wts");
|
||||
GitWorktrees gitWorktrees = new GitWorktrees(root.toString(), "cb224-nonexistent-group-zz");
|
||||
|
||||
WorktreeException e = assertThrows(WorktreeException.class,
|
||||
() -> gitWorktrees.add(repo.toString(), "cb-224-refuse", "HEAD"));
|
||||
|
||||
assertTrue(e.getMessage().contains(root.toString()),
|
||||
"refusal must name the worktree root: " + e.getMessage());
|
||||
assertTrue(e.getMessage().contains("cb224-nonexistent-group-zz"),
|
||||
"refusal must name the missing/refused group: " + e.getMessage());
|
||||
assertTrue(Files.isDirectory(root), "the root is created before the group check runs");
|
||||
String mode = java.nio.file.attribute.PosixFilePermissions.toString(Files.getPosixFilePermissions(root));
|
||||
assertTrue(e.getMessage().contains(mode),
|
||||
"refusal must name the root's current mode (" + mode + "): " + e.getMessage());
|
||||
try (java.util.stream.Stream<Path> children = Files.list(root)) {
|
||||
assertTrue(children.findAny().isEmpty(), "no worktree must be left behind under the root: "
|
||||
+ "the refusal must happen before `git worktree add` ever runs");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -5,6 +5,7 @@ import ch.qos.logback.classic.LoggerContext;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import dev.ltms.fleet.auth.MemberRegistry;
|
||||
import dev.ltms.fleet.auth.MemberLifecycle;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
@@ -335,18 +336,31 @@ class SessionManagerTest {
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-619 / fleetd #123: {@code requireSlotFor} closes the config-gap case (no slot at all
|
||||
* carries the profile) before anything spawns, but a profile that DOES carry a slot can still
|
||||
* lose the bind to a concurrent spawn racing for the same slot. This drives that residual case
|
||||
* through the REAL path — {@link SessionManager#acquire} against the real {@link
|
||||
* dev.ltms.fleet.member.ClaudeCodeLauncher} and {@link FakeHerdr} — never {@link
|
||||
* dev.ltms.fleet.auth.MemberRegistry#bind} directly for the session under test (only the
|
||||
* precondition uses it, to occupy the slot before the real spawn happens). The session that
|
||||
* loses the race must be held as a plain {@code dev}, never left claiming {@code architect} in
|
||||
* the roster, and the daemon log must say so at WARN.
|
||||
* fleetd #226: a contended slot is refused through the real {@link SessionManager#acquire}
|
||||
* path before the real launcher can hand an architect charter to a process.
|
||||
*/
|
||||
@Test
|
||||
void aContendedArchitectSlotRefusesBeforeTheLauncherStartsAMember() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
SessionManager sessions = sessionManager(herdr);
|
||||
MemberRegistry members = architectRegistry();
|
||||
assertTrue(members.bind("architect:opus", "term_already_bound"),
|
||||
"precondition: occupy the sole architect slot before the real spawn under test");
|
||||
sessions.setMemberLifecycle(members);
|
||||
|
||||
assertThrows(IllegalArgumentException.class, () -> sessions.acquire("ltms-local", MemberRole.ARCHITECT,
|
||||
null, "/caller", "term_primary", null));
|
||||
|
||||
assertFalse(herdr.called("agent.start"), "a refused slot must never reach the launcher");
|
||||
assertTrue(sessions.roster().isEmpty(), "no member exists to receive the wrong charter");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aSecondArchitectOnAnAlreadyBoundProfileIsHeldAsDevNotArchitectAndWarnsLoudly() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
SessionManager sessions = sessionManager(herdr);
|
||||
MemberRegistry members = architectRegistry();
|
||||
sessions.setMemberLifecycle(bindFailureAfterReservation(members));
|
||||
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
|
||||
ch.qos.logback.classic.Logger registryLog = (ch.qos.logback.classic.Logger)
|
||||
LoggerFactory.getLogger(MemberRegistry.class);
|
||||
@@ -356,35 +370,75 @@ class SessionManagerTest {
|
||||
registryLog.addAppender(appender);
|
||||
registryLog.setLevel(Level.WARN);
|
||||
try {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
SessionManager sessions = sessionManager(herdr);
|
||||
MemberRegistry members = new MemberRegistry(
|
||||
new FleetConfig.Fleet(Map.of(), Map.of("opus", new FleetConfig.Slot("ltms-local")),
|
||||
Map.of(), Map.of(), null));
|
||||
assertTrue(members.bind("architect:opus", "term_already_bound"),
|
||||
"precondition: occupy the sole architect slot before the real spawn under test");
|
||||
sessions.setMemberLifecycle(members);
|
||||
|
||||
MemberSession session = sessions.acquire("ltms-local", MemberRole.ARCHITECT, null,
|
||||
"/caller", "term_primary", null);
|
||||
|
||||
assertEquals(MemberRole.DEV, session.role(),
|
||||
"the slot is taken, so this session must be held as a plain member, never a lie");
|
||||
assertEquals("dev", SessionManager.rosterView(session, null).get("role"),
|
||||
"the roster must report what this session actually holds, not what it asked for");
|
||||
|
||||
assertEquals(MemberRole.DEV, session.role(), "a failed reservation bind must use the fallback");
|
||||
String warn = appender.list.stream()
|
||||
.filter(e -> e.getLevel().equals(Level.WARN))
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.findFirst()
|
||||
.orElse("no slot-exhaustion WARN logged");
|
||||
assertTrue(warn.contains("ltms-local"), "the log names the profile: " + warn);
|
||||
assertTrue(warn.contains(session.terminalId()), "the log names the terminal: " + warn);
|
||||
assertTrue(warn.contains("ltms-local"), "the WARN names the profile: " + warn);
|
||||
assertTrue(warn.contains(session.terminalId()), "the WARN names the terminal: " + warn);
|
||||
} finally {
|
||||
registryLog.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void reservedArchitectSlotBindsTheLaunchedMember() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
SessionManager sessions = sessionManager(herdr);
|
||||
sessions.setMemberLifecycle(architectRegistry());
|
||||
|
||||
MemberSession session = sessions.acquire("ltms-local", MemberRole.ARCHITECT, null,
|
||||
"/caller", "term_primary", null);
|
||||
|
||||
assertEquals(MemberRole.ARCHITECT, session.role());
|
||||
}
|
||||
|
||||
private static MemberRegistry architectRegistry() {
|
||||
return new MemberRegistry(new FleetConfig.Fleet(Map.of(), Map.of("opus", new FleetConfig.Slot("ltms-local")),
|
||||
Map.of(), Map.of(), null));
|
||||
}
|
||||
|
||||
private static MemberLifecycle bindFailureAfterReservation(MemberRegistry members) {
|
||||
return new MemberLifecycle() {
|
||||
@Override
|
||||
public MemberRole acquired(MemberRole role, String profile, String terminal) {
|
||||
return members.acquired(role, profile, terminal);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void released(String terminal) {
|
||||
members.released(terminal);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void requireSlotFor(MemberRole role, String profile) {
|
||||
members.requireSlotFor(role, profile);
|
||||
}
|
||||
|
||||
@Override
|
||||
public SlotReservation reserve(MemberRole role, String profile) {
|
||||
return members.reserve(role, profile);
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean bind(SlotReservation reservation, String terminal) {
|
||||
members.release(reservation);
|
||||
assertTrue(members.bind(reservation.slot(), "term_racer"), "the racer takes the released slot");
|
||||
return false;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void release(SlotReservation reservation) {
|
||||
members.release(reservation);
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
@Test
|
||||
void rosterReflectsAcquiredMinusReleased() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
@@ -651,6 +705,32 @@ class SessionManagerTest {
|
||||
"no session is registered when spawn times out (roster empty)");
|
||||
}
|
||||
|
||||
@Test
|
||||
void failedArchitectLaunchReleasesItsReservationForTheNextLaunch() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
herdr.agentStatus("unknown");
|
||||
long[] clock = {0};
|
||||
FleetConfig.Profile cfg = new FleetConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
|
||||
List.of("ccs", "ltms-local"), "tab", "fleetd-workers",
|
||||
"worker: {profile} #{n}", null, null, null);
|
||||
ClaudeCodeLauncher workers = new ClaudeCodeLauncher(
|
||||
new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null,
|
||||
1, () -> clock[0], () -> clock[0] += 10);
|
||||
SessionManager sessions = new SessionManager(workers, new GitWorktrees(), () -> 0L, 0);
|
||||
sessions.setMemberLifecycle(architectRegistry());
|
||||
|
||||
assertThrows(PeerUnreachableException.class, () -> sessions.acquire("ltms-local", MemberRole.ARCHITECT,
|
||||
null, "/caller", "term_primary", null));
|
||||
|
||||
herdr.agentStatus("idle");
|
||||
MemberSession retry = sessions.acquire("ltms-local", MemberRole.ARCHITECT, null,
|
||||
"/caller", "term_primary", null);
|
||||
assertEquals(MemberRole.ARCHITECT, retry.role(),
|
||||
"the failed launch returned its reservation instead of silently shrinking the slot pool");
|
||||
}
|
||||
|
||||
// --- CB-516: release must notify, so a blocked send can be failed --------------------------
|
||||
|
||||
@Test
|
||||
|
||||
@@ -95,10 +95,9 @@ class WorktreeSessionManagerTest {
|
||||
null, "/caller/proj", null, null);
|
||||
assertEquals("architect:architect", members.slotForTerminal(replacement.terminalId()));
|
||||
|
||||
MemberSession overflow = sessions.acquire("ltms-local", MemberRole.ARCHITECT,
|
||||
null, "/caller/proj", null, null);
|
||||
assertNull(members.slotForTerminal(overflow.terminalId()),
|
||||
"a full slot pool must not stop the architect spawn");
|
||||
assertThrows(IllegalArgumentException.class, () -> sessions.acquire("ltms-local", MemberRole.ARCHITECT,
|
||||
null, "/caller/proj", null, null),
|
||||
"a full slot pool must stop the spawn before it can receive an architect charter");
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
Reference in New Issue
Block a user