Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| fc2e26c0c5 |
@@ -0,0 +1,52 @@
|
||||
# CB-175 report
|
||||
|
||||
## Change
|
||||
|
||||
`OpenCodeLauncher` reads the newest opencode session record for the worker cwd after spawn readiness.
|
||||
It compares the requested `provider/model` selector with `model.providerID/model.id` from the record.
|
||||
`variant` is not compared because a profile selector has no variant part.
|
||||
|
||||
An absent record, malformed record, or incomplete model object is unknown evidence. It does not log
|
||||
an error or quarantine the profile.
|
||||
|
||||
On a real mismatch, fleetd logs an ERROR with the requested and resolved selectors. The mismatch goes
|
||||
through `ExhaustionSink` into the existing `BackendQuarantine` and uses the profile's
|
||||
`effectiveCredentialId()`.
|
||||
|
||||
I chose a permanent, process-lifetime quarantine. A withdrawn selector cannot become correct after a
|
||||
cooldown. A timed retry could silently use the paid fallback again. `fleet_list` will show the usual
|
||||
quarantine state, with a very large remaining time, until fleetd restarts after an operator fixes the
|
||||
profile.
|
||||
|
||||
## Tests
|
||||
|
||||
Added tests for an exact match, a mismatch, and missing or unreadable session storage.
|
||||
|
||||
I proved the mismatch test fails without the quarantine call. I commented out the call and ran:
|
||||
|
||||
```text
|
||||
mvn -Dtest=OpenCodeLauncherTest#differentResolvedModelPermanentlyQuarantinesTheProfile test
|
||||
```
|
||||
|
||||
The result was:
|
||||
|
||||
```text
|
||||
[ERROR] Tests run: 1, Failures: 1, Errors: 0, Skipped: 0
|
||||
org.opentest4j.AssertionFailedError: a fallback model must block later spawns ==> expected: <true> but was: <false>
|
||||
[INFO] BUILD FAILURE
|
||||
```
|
||||
|
||||
I restored the call. I then ran `mvn clean install` in `fleetd/` without a pipe. Its result was:
|
||||
|
||||
```text
|
||||
[INFO] Tests run: 1040, Failures: 0, Errors: 0, Skipped: 0
|
||||
[INFO] BUILD SUCCESS
|
||||
```
|
||||
|
||||
## Limits and scope
|
||||
|
||||
I could not spawn a real opencode member or restart fleetd. I did not test this end to end against a
|
||||
live opencode session database.
|
||||
|
||||
I confirmed `ClaudeCodeLauncher` passes `--model` but does not read back the resolved model. I did
|
||||
not change it because it is outside this ticket's scope.
|
||||
@@ -117,15 +117,6 @@ herdrSocket: ~/.config/herdr/herdr.sock
|
||||
# Optional socket for member panes. Omit this to use herdrSocket for both leads and members.
|
||||
# memberHerdrSocket: /Users/member/.config/herdr/herdr.sock
|
||||
|
||||
# fleetd #213: the login shell the member OS user (memberHerdrSocket above) actually runs. ONLY
|
||||
# read when memberHerdrSocket is set — fleetd's own $SHELL says nothing about a pane running
|
||||
# under a different OS user, and there is no channel to ask herdr for that user's shell, so this
|
||||
# must be told rather than guessed. Absent, blank, or anything not ending in "zsh" is treated the
|
||||
# same as "not zsh": the memberCredentials.policy: allow-list ZDOTDIR scrub (see worktreeGroup
|
||||
# below) is skipped in favour of the weaker CB-596 sentinel overlay — a degraded control, never a
|
||||
# refusal to spawn. When memberHerdrSocket is absent this key is never consulted at all.
|
||||
# memberLoginShell: /bin/zsh
|
||||
|
||||
# How member sessions are spawned. Define one or more named profiles (backends) under
|
||||
# `profiles`; each key is the profile name (also the ccs profile). A profile says only WHICH
|
||||
# BACKEND — model, CLI adapter, credentials, cost. It says nothing about what a member spawned on
|
||||
@@ -643,24 +634,6 @@ guard:
|
||||
# to a sibling directory of the repo root.
|
||||
# worktreeRoot: /Users/me/src/.bridged-worktrees
|
||||
|
||||
# Worktree group sharing (fleetd #185 stage 3). OPTIONAL, off by default. Names an OS group
|
||||
# that a provisioned worktree's repo is made group-writable for (git config
|
||||
# core.sharedRepository group, plus a one-time chgrp/chmod/setgid fix-up), so a member spawned
|
||||
# under a DIFFERENT OS user (see memberHerdrSocket) can write its own worktree, its
|
||||
# per-worktree git metadata, and its own commit objects — without it, every file GitWorktrees
|
||||
# creates is owned by fleetd's own uid and unwritable by another user.
|
||||
# CAUTION: this isolates credentials, not the repository — a member in the group can still
|
||||
# write the operator's git objects and refs in the shared repo. The operator running fleetd
|
||||
# must already be a member of the named group, or every provisioning spawn fails loudly.
|
||||
#
|
||||
# fleetd #213: this is also the ONE group the memberCredentials.policy: allow-list ZDOTDIR scrub
|
||||
# reuses when memberHerdrSocket is set — deliberately not a second config key. Under
|
||||
# memberHerdrSocket, the scrub directory is generated under worktreeRoot (never java.io.tmpdir,
|
||||
# which the member OS user cannot reach) and shared read-only with this group. If worktreeGroup
|
||||
# is unset while memberHerdrSocket is set, the scrub cannot be guaranteed reachable by the member,
|
||||
# so fleetd falls back to the weaker CB-596 sentinel overlay instead (a WARN names the gap).
|
||||
# worktreeGroup: fleet-workers
|
||||
|
||||
# Session lifecycle limits (CB-303). All knobs are opt-in; omit or set to null to keep
|
||||
# the feature disabled. By default the daemon never reaps, caps, or drains sessions.
|
||||
# idleTtlSeconds → reap READY/DONE sessions idle longer than this (never BUSY/SPAWNING)
|
||||
|
||||
@@ -28,7 +28,6 @@
|
||||
<testcontainers.version>1.20.4</testcontainers.version>
|
||||
<commons-compress.version>1.27.1</commons-compress.version>
|
||||
<commons-lang3.version>3.18.0</commons-lang3.version>
|
||||
<sqlite-jdbc.version>3.53.4.0</sqlite-jdbc.version>
|
||||
</properties>
|
||||
|
||||
<!--
|
||||
@@ -45,12 +44,6 @@
|
||||
3.0-rc5; bumping Jackson 3 to the patched 3.2.x breaks the SDK (annotation mismatch).
|
||||
Only the loopback /mcp endpoint parses this JSON, from trusted local Claude clients.
|
||||
The 11.0.23 -> 11.0.25 bump did clear jetty CVE-2024-8184 (5.9) and CVE-2024-6763.
|
||||
|
||||
fleetd #206: org.xerial:sqlite-jdbc 3.53.4.0 (added for OpenCodeSessionDiscovery) — the
|
||||
only known advisory against this artifact is CVE-2023-32697 (RCE via an attacker-controlled
|
||||
JDBC URL), fixed in 3.41.2.2; 3.53.4.0 is well past that fix and OSV.dev reports no open
|
||||
advisory against it. Checked via the OSV.dev API (no Mend.io/JetBrains IDE MCP mount
|
||||
available from this worktree) on 2026-08-31.
|
||||
-->
|
||||
|
||||
<!-- Force the latest patched Jetty 11.x across all Javalin-pulled Jetty modules (no version
|
||||
@@ -130,17 +123,6 @@
|
||||
<version>${amqp.version}</version>
|
||||
</dependency>
|
||||
|
||||
<!-- fleetd #206: opencode moved its session store from a JSON tree to SQLite
|
||||
(opencode.db). This is the JDBC driver OpenCodeSessionDiscovery uses to read it
|
||||
read-only. Ships bundled native libraries (linux/mac/windows, several archs), so it
|
||||
is a heavier jar than most deps here — see the pom's dependency-security note below
|
||||
for the size/CVE tradeoff actually measured. -->
|
||||
<dependency>
|
||||
<groupId>org.xerial</groupId>
|
||||
<artifactId>sqlite-jdbc</artifactId>
|
||||
<version>${sqlite-jdbc.version}</version>
|
||||
</dependency>
|
||||
|
||||
<!-- Logging -->
|
||||
<dependency>
|
||||
<groupId>org.slf4j</groupId>
|
||||
|
||||
@@ -168,6 +168,17 @@ public final class Fleetd {
|
||||
claudeProfiles.put(name, w);
|
||||
}
|
||||
});
|
||||
// One tracker covers timed backend exhaustion and permanent model-selector mismatches. The
|
||||
// latter cannot heal on a retry, so OpenCodeLauncher uses quarantinePermanently through the
|
||||
// sink below rather than letting a cooldown reopen a paid fallback.
|
||||
BackendQuarantine quarantine = new BackendQuarantine(System::nanoTime,
|
||||
TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds()));
|
||||
ExhaustionSink modelMismatchSink = (profileName, reason) -> {
|
||||
FleetConfig.Profile profile = config.get().profiles().get(profileName);
|
||||
if (profile != null) {
|
||||
quarantine.quarantinePermanently(profile.effectiveCredentialId());
|
||||
}
|
||||
};
|
||||
List<HerdrPeerLauncher> adapters = new ArrayList<>();
|
||||
// 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)
|
||||
@@ -177,22 +188,20 @@ public final class Fleetd {
|
||||
claudeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
|
||||
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
|
||||
() -> config.get().fleet(),
|
||||
() -> config.get().memberCredentials(), null, config::get));
|
||||
() -> config.get().memberCredentials()));
|
||||
}
|
||||
if (!opencodeProfiles.isEmpty()) {
|
||||
adapters.add(new OpenCodeLauncher(router.memberAgents(), router.memberSpaces(),
|
||||
opencodeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
|
||||
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
|
||||
() -> config.get().fleet(),
|
||||
() -> config.get().memberCredentials(), config::get));
|
||||
() -> config.get().memberCredentials(), modelMismatchSink));
|
||||
}
|
||||
AtomicReference<Function<String, Integer>> liveCountRef = new AtomicReference<>(_ -> 0);
|
||||
// CB-578 stage B: one quarantine tracker for the whole daemon, shared between the launcher
|
||||
// (checked at spawn) and the exhaustion sink wired in below (written on BACKEND_EXHAUSTED).
|
||||
// The cooldown is deferred (see FleetConfig#quarantineCooldownSeconds): it is read once
|
||||
// here, at startup, and a config reload only changes it for a daemon restart.
|
||||
BackendQuarantine quarantine = new BackendQuarantine(System::nanoTime,
|
||||
TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds()));
|
||||
PeerLauncher workers = new CompositePeerLauncher(
|
||||
adapters,
|
||||
cfg.effectiveDefaultProfile(),
|
||||
@@ -224,7 +233,7 @@ public final class Fleetd {
|
||||
contextCap = cfg.lifecycle().contextCap();
|
||||
}
|
||||
boolean clearAfterTurn = cfg.lifecycle() != null && cfg.lifecycle().clearAfterTurn();
|
||||
SessionManager sessions = new SessionManager(workers, new GitWorktrees(cfg.worktreeRoot(), cfg.worktreeGroup()),
|
||||
SessionManager sessions = new SessionManager(workers, new GitWorktrees(cfg.worktreeRoot()),
|
||||
System::nanoTime, contextCap, clearAfterTurn);
|
||||
liveCountRef.set(profileName -> (int) sessions.roster().stream()
|
||||
.filter(s -> profileName.equals(s.profile()))
|
||||
|
||||
@@ -76,26 +76,6 @@ import java.util.Set;
|
||||
* stay on {@code broker}'s vhost). {@code null} → no lead mailbox is opened.
|
||||
* Config parsing + accessors only — nothing here wires it into a live
|
||||
* {@code LeadMailbox}; that is a separate ticket. See {@link Coordinator}.
|
||||
* @param worktreeGroup optional OS group name (fleetd #185 stage 3) that makes a provisioned
|
||||
* worktree's repo group-shared, so a member running as a different OS user
|
||||
* (see {@code memberHerdrSocket}) can write its own worktree, its per-worktree
|
||||
* git metadata, and its own commit objects. {@code null}/blank/empty ⇒ off,
|
||||
* today's behaviour unchanged (every file stays owned by fleetd's own uid).
|
||||
* <strong>This isolates credentials, not the repository</strong>: a member in
|
||||
* the group can still write the operator's git objects and refs in the shared
|
||||
* repo. See {@link dev.ltms.fleet.session.Worktrees#shareWithGroup}.
|
||||
* @param memberLoginShell fleetd #213: the login shell the member's OS user actually runs, ONLY
|
||||
* meaningful (and only ever read) when {@code memberHerdrSocket} is
|
||||
* configured — that mode spawns member panes under a different OS user than
|
||||
* fleetd's own process, so fleetd's own {@code $SHELL} says nothing about what
|
||||
* that pane runs. There is no channel to ask herdr for another user's shell, so
|
||||
* this must be told, never guessed. {@code null}/blank (or a value not ending
|
||||
* in {@code zsh}) is treated the same as "not zsh": the {@code
|
||||
* memberCredentials.policy: allow-list} ZDOTDIR scrub is skipped in favour of
|
||||
* the CB-596 sentinel overlay — a degraded control, never a refusal to spawn.
|
||||
* When {@code memberHerdrSocket} is NOT configured this field is never
|
||||
* consulted at all; fleetd keeps reading its own {@code $SHELL}, exactly as
|
||||
* before this field existed.
|
||||
*/
|
||||
@JsonIgnoreProperties(ignoreUnknown = true)
|
||||
public record FleetConfig(
|
||||
@@ -118,33 +98,7 @@ public record FleetConfig(
|
||||
ConfigReload configReload,
|
||||
Integer quarantineCooldownSeconds,
|
||||
MemberCredentials memberCredentials,
|
||||
Coordinator coordinator,
|
||||
String worktreeGroup,
|
||||
String memberLoginShell) {
|
||||
|
||||
/** Back-compat form before the {@code memberLoginShell} key was added. */
|
||||
public FleetConfig(Bind bind, String herdrSocket, String memberHerdrSocket, Map<String, Profile> profiles,
|
||||
Guard guard, String worktreeRoot, Lifecycle lifecycle, Integer spawnReadyTimeoutMs,
|
||||
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
|
||||
LeadHeartbeat leadHeartbeat, Health health, String placement, Auth auth,
|
||||
ConfigReload configReload, Integer quarantineCooldownSeconds,
|
||||
MemberCredentials memberCredentials, Coordinator coordinator, String worktreeGroup) {
|
||||
this(bind, herdrSocket, memberHerdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
|
||||
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth,
|
||||
configReload, quarantineCooldownSeconds, memberCredentials, coordinator, worktreeGroup, null);
|
||||
}
|
||||
|
||||
/** Back-compat form before the {@code worktreeGroup} key was added. */
|
||||
public FleetConfig(Bind bind, String herdrSocket, String memberHerdrSocket, Map<String, Profile> profiles,
|
||||
Guard guard, String worktreeRoot, Lifecycle lifecycle, Integer spawnReadyTimeoutMs,
|
||||
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
|
||||
LeadHeartbeat leadHeartbeat, Health health, String placement, Auth auth,
|
||||
ConfigReload configReload, Integer quarantineCooldownSeconds,
|
||||
MemberCredentials memberCredentials, Coordinator coordinator) {
|
||||
this(bind, herdrSocket, memberHerdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
|
||||
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth,
|
||||
configReload, quarantineCooldownSeconds, memberCredentials, coordinator, null, null);
|
||||
}
|
||||
Coordinator coordinator) {
|
||||
|
||||
/** Back-compat form before the {@code coordinator:} block was added. */
|
||||
public FleetConfig(Bind bind, String herdrSocket, Map<String, Profile> profiles, Guard guard,
|
||||
@@ -155,7 +109,7 @@ public record FleetConfig(
|
||||
MemberCredentials memberCredentials) {
|
||||
this(bind, herdrSocket, null, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
|
||||
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth,
|
||||
configReload, quarantineCooldownSeconds, memberCredentials, null, null);
|
||||
configReload, quarantineCooldownSeconds, memberCredentials, null);
|
||||
}
|
||||
|
||||
/** Back-compat form before the CB-596 {@code memberCredentials:} block was added. */
|
||||
@@ -1371,7 +1325,7 @@ public record FleetConfig(
|
||||
"bind", "herdrSocket", "memberHerdrSocket", "profiles", "guard", "worktreeRoot",
|
||||
"lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs", "broker", "primary", "fleet",
|
||||
"leadHeartbeat", "health", "placement", "auth", "configReload", "quarantineCooldownSeconds",
|
||||
"memberCredentials", "coordinator", "worktreeGroup", "memberLoginShell");
|
||||
"memberCredentials", "coordinator");
|
||||
|
||||
/** Load and validate config from {@code path}. */
|
||||
public static FleetConfig load(Path path) {
|
||||
@@ -1989,14 +1943,9 @@ public record FleetConfig(
|
||||
: new MemberCredentials(null, List.of(), List.of());
|
||||
// coordinator is left as-is, like broker/primary above: null keeps no LeadMailbox opened,
|
||||
// and this ticket's Coordinator is config-only anyway (nothing yet reads it at startup).
|
||||
// worktreeGroup is left as-is (fleetd #185 stage 3): null/blank is "off", and there is no
|
||||
// sane non-null default — an OS group name is operator-specific.
|
||||
// memberLoginShell is left as-is (fleetd #213), like worktreeGroup: null/blank is "not
|
||||
// configured", and there is no sane non-null default — a member's login shell is
|
||||
// operator-specific and only meaningful when memberHerdrSocket is also set.
|
||||
return new FleetConfig(b, herdrSocket, memberHerdrSocket, profiles, g, worktreeRoot, l, timeout, pollMs,
|
||||
broker, primary, f, leadHeartbeat, health, placementOrDefault, a, configReload,
|
||||
quarantineCooldown, mc, coordinator, worktreeGroup, memberLoginShell);
|
||||
quarantineCooldown, mc, coordinator);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -2,11 +2,8 @@ package dev.ltms.fleet.herdr;
|
||||
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
|
||||
import java.util.LinkedHashSet;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.OptionalLong;
|
||||
import java.util.Set;
|
||||
|
||||
/**
|
||||
* Resolves which herdr pane a process belongs to — the herdr half of connection-based MCP
|
||||
@@ -15,46 +12,23 @@ import java.util.Set;
|
||||
* is calling without the worker sending anything spoofable.
|
||||
*
|
||||
* <p>herdr owns the PID→pane truth: {@code pane.process_info} reports each pane's {@code shell_pid}
|
||||
* and foreground process PIDs. A pid that is neither of those directly — e.g. a grandchild a
|
||||
* worker spawned, such as a {@code python3} or {@code curl} helper that opens its own MCP
|
||||
* connection — is resolved by walking its ancestry (via {@link ParentResolver}) up to the root and
|
||||
* matching any ancestor against a pane's {@code shell_pid} or foreground pids (CB-161). Without
|
||||
* this walk such a pid matches no pane, and the caller falls through to loopback-trust and is
|
||||
* resolved as the primary — a worker→primary privilege escalation.
|
||||
*
|
||||
* <p>This scans agent panes; a spawn-time {@code pid→terminal} cache is the obvious optimization
|
||||
* once wired into {@code ClaudeCodeLauncher}.
|
||||
* and foreground process PIDs. This scans agent panes; a spawn-time {@code pid→terminal} cache is
|
||||
* the obvious optimization once wired into {@code ClaudeCodeLauncher}.
|
||||
*
|
||||
* <p>CB-185 split the fleet across two herdr daemons — lead operations on one, members on the
|
||||
* other ({@code memberHerdrSocket}). A caller's pane can live on <em>either</em> daemon (a lead's
|
||||
* MCP connection resolves against the lead daemon; a member's against the member daemon), so this
|
||||
* must be able to search more than one client. {@link #PaneLocator(HerdrClient, HerdrClient)}
|
||||
* searches the lead client first, then the member client, and collapses to a single scan when the
|
||||
* two are the same object (the historical single-daemon deployment). The caller's ancestor set is
|
||||
* computed once per {@link #terminalForPid} call and reused across every client searched — it
|
||||
* does not depend on which daemon a pane happens to live on.
|
||||
* two are the same object (the historical single-daemon deployment).
|
||||
*/
|
||||
public final class PaneLocator {
|
||||
|
||||
/**
|
||||
* Bound on how many ancestor generations {@link #ancestorsOf} walks. This runs on every MCP
|
||||
* call, so a cycle or a pathologically deep process tree must not hang identity resolution;
|
||||
* 32 generations is far more than any real worker→helper process tree needs.
|
||||
*/
|
||||
private static final int MAX_ANCESTRY_DEPTH = 32;
|
||||
|
||||
private final List<HerdrClient> herdrs;
|
||||
private final ParentResolver parentResolver;
|
||||
|
||||
/** Search only this client — the single-daemon deployment. */
|
||||
public PaneLocator(HerdrClient herdr) {
|
||||
this(herdr, ParentResolver.PROCESS_HANDLE);
|
||||
}
|
||||
|
||||
/** Search only this client, resolving ancestry through {@code parentResolver} — for tests. */
|
||||
public PaneLocator(HerdrClient herdr, ParentResolver parentResolver) {
|
||||
this.herdrs = List.of(herdr);
|
||||
this.parentResolver = parentResolver;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -63,13 +37,7 @@ public final class PaneLocator {
|
||||
* collapses to one client and one scan, exactly {@link #PaneLocator(HerdrClient)}'s behaviour.
|
||||
*/
|
||||
public PaneLocator(HerdrClient lead, HerdrClient member) {
|
||||
this(lead, member, ParentResolver.PROCESS_HANDLE);
|
||||
}
|
||||
|
||||
/** Two-daemon deployment, resolving ancestry through {@code parentResolver} — for tests. */
|
||||
public PaneLocator(HerdrClient lead, HerdrClient member, ParentResolver parentResolver) {
|
||||
this.herdrs = lead == member ? List.of(lead) : List.of(lead, member);
|
||||
this.parentResolver = parentResolver;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -81,9 +49,8 @@ public final class PaneLocator {
|
||||
if (pid <= 0) {
|
||||
return null;
|
||||
}
|
||||
Set<Long> ancestry = ancestorsOf(pid);
|
||||
for (HerdrClient herdr : herdrs) {
|
||||
String terminal = terminalForPid(herdr, ancestry);
|
||||
String terminal = terminalForPid(herdr, pid);
|
||||
if (terminal != null) {
|
||||
return terminal;
|
||||
}
|
||||
@@ -91,54 +58,28 @@ public final class PaneLocator {
|
||||
return null;
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code pid} itself plus its ancestor chain, walked through {@link #parentResolver} up to
|
||||
* {@link #MAX_ANCESTRY_DEPTH} generations or pid 1, whichever comes first. A vanished ancestor
|
||||
* ({@link ParentResolver#parentOf} returning empty) ends the walk without error — it just means
|
||||
* the chain is shorter than the bound. A cycle in a fake resolver is caught by the "already
|
||||
* seen" check and also ends the walk, so this can never loop.
|
||||
*/
|
||||
private Set<Long> ancestorsOf(long pid) {
|
||||
Set<Long> ancestry = new LinkedHashSet<>();
|
||||
long current = pid;
|
||||
for (int depth = 0; depth < MAX_ANCESTRY_DEPTH; depth++) {
|
||||
if (current <= 0 || !ancestry.add(current)) {
|
||||
break; // vanished/invalid pid, or a cycle back to a pid already recorded
|
||||
}
|
||||
if (current == 1) {
|
||||
break; // reached the root of the process tree
|
||||
}
|
||||
OptionalLong parent = parentResolver.parentOf(current);
|
||||
if (parent.isEmpty()) {
|
||||
break; // vanished ancestor — not an error, just the end of the chain
|
||||
}
|
||||
current = parent.getAsLong();
|
||||
}
|
||||
return ancestry;
|
||||
}
|
||||
|
||||
private static String terminalForPid(HerdrClient herdr, Set<Long> ancestry) {
|
||||
private static String terminalForPid(HerdrClient herdr, long pid) {
|
||||
for (JsonNode pane : herdr.call("pane.list", Map.of()).path("panes")) {
|
||||
String paneId = pane.path("pane_id").asText(null);
|
||||
if (paneId != null && paneOwnsAnyOf(herdr, paneId, ancestry)) {
|
||||
if (paneId != null && paneOwnsPid(herdr, paneId, pid)) {
|
||||
return pane.path("terminal_id").asText(null);
|
||||
}
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
private static boolean paneOwnsAnyOf(HerdrClient herdr, String paneId, Set<Long> ancestry) {
|
||||
private static boolean paneOwnsPid(HerdrClient herdr, String paneId, long pid) {
|
||||
JsonNode info;
|
||||
try {
|
||||
info = herdr.call("pane.process_info", Map.of("pane_id", paneId)).path("process_info");
|
||||
} catch (HerdrException e) {
|
||||
return false; // pane vanished mid-scan — just skip it
|
||||
}
|
||||
if (ancestry.contains(info.path("shell_pid").asLong(-1))) {
|
||||
if (info.path("shell_pid").asLong(-1) == pid) {
|
||||
return true;
|
||||
}
|
||||
for (JsonNode p : info.path("foreground_processes")) {
|
||||
if (ancestry.contains(p.path("pid").asLong(-1))) {
|
||||
if (p.path("pid").asLong(-1) == pid) {
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,21 +0,0 @@
|
||||
package dev.ltms.fleet.herdr;
|
||||
|
||||
import java.util.OptionalLong;
|
||||
|
||||
/**
|
||||
* Resolves a pid's parent pid — the seam {@link PaneLocator} walks a process's ancestry through,
|
||||
* so its tests can drive the walk from a fake pid→parent map instead of spawning real processes.
|
||||
*
|
||||
* <p>{@link #PROCESS_HANDLE} is the production implementation, backed by {@link ProcessHandle}.
|
||||
*/
|
||||
public interface ParentResolver {
|
||||
|
||||
/** The parent pid of {@code pid}, or empty if {@code pid} is gone or has no known parent. */
|
||||
OptionalLong parentOf(long pid);
|
||||
|
||||
/** Production resolver: asks the JVM's {@link ProcessHandle} view of the OS process tree. */
|
||||
ParentResolver PROCESS_HANDLE = pid -> ProcessHandle.of(pid)
|
||||
.flatMap(ProcessHandle::parent)
|
||||
.map(parent -> OptionalLong.of(parent.pid()))
|
||||
.orElse(OptionalLong.empty());
|
||||
}
|
||||
@@ -951,10 +951,7 @@ public final class FleetMcp {
|
||||
.sorted(Map.Entry.comparingByValue())
|
||||
.map(e -> leadView(e.getKey(), e.getValue(), live.get(e.getKey()), selfTerm))
|
||||
.toList();
|
||||
// fleetd #209: this is the caller-driven fleet_list read that actually reports
|
||||
// agentSessionId (via memberCapacityView -> SessionManager.rosterView), so it uses the
|
||||
// resolving roster; the heartbeat/health/metrics timers stay on the plain sessions.roster().
|
||||
List<MemberSession> roster = sessions.rosterResolved();
|
||||
List<MemberSession> roster = sessions.roster();
|
||||
List<Map<String, Object>> out = roster.stream()
|
||||
.map(s -> memberCapacityView(s, live.get(s.terminalId()), messages, capacity.clock().getAsLong()))
|
||||
.toList();
|
||||
|
||||
@@ -94,26 +94,13 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
|
||||
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
|
||||
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
|
||||
Function<String, String> env,
|
||||
long spawnReadyTimeoutMs, long spawnReadyPollMs,
|
||||
Supplier<FleetConfig.Fleet> fleet,
|
||||
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
|
||||
this(agents, spaces, guard, profiles, defaultProfile, env, spawnReadyTimeoutMs, spawnReadyPollMs,
|
||||
fleet, memberCredentials, null, null);
|
||||
}
|
||||
|
||||
/** Production constructor, plus the live config for URI environment exclusions. */
|
||||
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
|
||||
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
|
||||
Function<String, String> env,
|
||||
long spawnReadyTimeoutMs, long spawnReadyPollMs,
|
||||
Supplier<FleetConfig.Fleet> fleet,
|
||||
Supplier<FleetConfig.MemberCredentials> memberCredentials,
|
||||
Supplier<Set<String>> hostEnvNames,
|
||||
Supplier<FleetConfig> config) {
|
||||
long spawnReadyTimeoutMs, long spawnReadyPollMs,
|
||||
Supplier<FleetConfig.Fleet> fleet,
|
||||
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
|
||||
this(agents, spaces, guard, profiles, defaultProfile, env,
|
||||
spawnReadyTimeoutMs,
|
||||
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs),
|
||||
fleet, memberCredentials, hostEnvNames, config);
|
||||
fleet, memberCredentials);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -166,22 +153,10 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
|
||||
Function<String, String> env,
|
||||
long spawnReadyTimeoutMs,
|
||||
LongSupplier nowMillis, Runnable sleeper,
|
||||
Supplier<FleetConfig.Fleet> fleet,
|
||||
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
|
||||
this(agents, spaces, guard, profiles, defaultProfile, env, spawnReadyTimeoutMs, nowMillis, sleeper,
|
||||
fleet, memberCredentials, null, null);
|
||||
}
|
||||
|
||||
/** Full testability constructor, plus the live config for URI environment exclusions. */
|
||||
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
|
||||
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
|
||||
Function<String, String> env, long spawnReadyTimeoutMs,
|
||||
LongSupplier nowMillis, Runnable sleeper, Supplier<FleetConfig.Fleet> fleet,
|
||||
Supplier<FleetConfig.MemberCredentials> memberCredentials,
|
||||
Supplier<Set<String>> hostEnvNames,
|
||||
Supplier<FleetConfig> config) {
|
||||
Supplier<FleetConfig.Fleet> fleet,
|
||||
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
|
||||
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
|
||||
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials, hostEnvNames, config);
|
||||
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials);
|
||||
this.guard = guard;
|
||||
}
|
||||
|
||||
@@ -199,7 +174,7 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
|
||||
Supplier<FleetConfig.MemberCredentials> memberCredentials,
|
||||
Supplier<Set<String>> hostEnvNames) {
|
||||
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
|
||||
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials, hostEnvNames, null);
|
||||
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials, hostEnvNames);
|
||||
this.guard = guard;
|
||||
}
|
||||
|
||||
|
||||
@@ -7,9 +7,6 @@ import java.io.IOException;
|
||||
import java.io.UncheckedIOException;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.nio.file.attribute.GroupPrincipal;
|
||||
import java.nio.file.attribute.PosixFileAttributeView;
|
||||
import java.nio.file.attribute.PosixFilePermissions;
|
||||
import java.time.Duration;
|
||||
import java.time.Instant;
|
||||
import java.util.ArrayList;
|
||||
@@ -119,70 +116,6 @@ public final class EnvAllowListScrub {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #213: as {@link #generate(Path, Set)}, plus share the generated directory with
|
||||
* {@code group} — the member OS user's group (the operator's existing {@code worktreeGroup:}
|
||||
* name, reused rather than inventing a second one) — so a member running under a different OS
|
||||
* user than fleetd's own process can still read what it needs from a directory placed outside
|
||||
* {@code java.io.tmpdir}. {@code group} null/blank ⇒ identical to {@link #generate(Path, Set)};
|
||||
* this is the single-daemon (no {@code memberHerdrSocket}) shape, where the pane is fleetd's own
|
||||
* uid and no group sharing is needed.
|
||||
*
|
||||
* @throws UncheckedIOException also when {@code group} does not resolve on this host, or a
|
||||
* group-ownership/permission call is refused — the same "fail
|
||||
* loudly rather than start unprotected" contract as above: a scrub
|
||||
* the configured member user cannot even read is not a working
|
||||
* control.
|
||||
*/
|
||||
public static Path generate(Path parentDir, Set<String> allowedNames, String group) {
|
||||
Path dir = generate(parentDir, allowedNames);
|
||||
if (group != null && !group.isBlank()) {
|
||||
shareWithGroup(dir, group);
|
||||
}
|
||||
return dir;
|
||||
}
|
||||
|
||||
/**
|
||||
* chgrp/chmod-equivalent over the freshly generated directory and the startup files already
|
||||
* written into it: owner keeps full access, {@code group} gets traverse+read on the directory
|
||||
* ({@code rwxr-x---}, so a login shell under that group can find and source the files) and
|
||||
* read-only on each file ({@code rw-r-----}) — deliberately no group WRITE anywhere, since a
|
||||
* member never needs to add or change fleetd's own generated scrub. (The scrub script's own
|
||||
* report write inside the pane consequently fails closed rather than open — see {@code
|
||||
* scrub.zsh}'s trailing {@code 2>/dev/null} — which {@link
|
||||
* dev.ltms.fleet.member.HerdrPeerLauncher#releaseZdotdir} already treats as "cannot be
|
||||
* confirmed to have run" rather than success.)
|
||||
*/
|
||||
private static void shareWithGroup(Path dir, String group) {
|
||||
try {
|
||||
GroupPrincipal principal = dir.getFileSystem().getUserPrincipalLookupService()
|
||||
.lookupPrincipalByGroupName(group);
|
||||
setGroupAndPermissions(dir, principal, "rwxr-x---");
|
||||
try (Stream<Path> entries = Files.list(dir)) {
|
||||
for (Path file : entries.toList()) {
|
||||
setGroupAndPermissions(file, principal, "rw-r-----");
|
||||
}
|
||||
}
|
||||
} catch (IOException e) {
|
||||
throw new UncheckedIOException("cannot share generated ZDOTDIR " + dir + " with group '"
|
||||
+ group + "' — the group must exist, and the fleetd operator ("
|
||||
+ System.getProperty("user.name") + ") must be a member of it", e);
|
||||
} catch (UnsupportedOperationException e) {
|
||||
throw new UncheckedIOException("cannot share generated ZDOTDIR " + dir + " with group '"
|
||||
+ group + "' — this filesystem does not support POSIX group ownership",
|
||||
new IOException(e));
|
||||
}
|
||||
}
|
||||
|
||||
private static void setGroupAndPermissions(Path path, GroupPrincipal group, String perms) throws IOException {
|
||||
PosixFileAttributeView view = Files.getFileAttributeView(path, PosixFileAttributeView.class);
|
||||
if (view == null) {
|
||||
throw new IOException("POSIX file attributes are not supported for " + path);
|
||||
}
|
||||
view.setGroup(group);
|
||||
Files.setPosixFilePermissions(path, PosixFilePermissions.fromString(perms));
|
||||
}
|
||||
|
||||
/** One operator-sourcing startup file: source the {@code $HOME} counterpart, change nothing else. */
|
||||
private static String homeSourcingFile(String name) {
|
||||
return """
|
||||
|
||||
@@ -112,12 +112,6 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
* daemon's own process is started the same way (a login shell sourcing the same secret store —
|
||||
* see CB-592's investigation of {@code secrets.sh}), so on a single-host deployment its env
|
||||
* mirrors what the pane's login shell is about to export.
|
||||
*
|
||||
* <p>fleetd #185 stage 2: that mirroring assumption holds only while the member pane runs under
|
||||
* the SAME OS user as the daemon. When {@code memberHerdrSocket:} is configured, member panes
|
||||
* run on a second herdr owned by a different user — different {@code $HOME}, different {@code
|
||||
* secrets.sh}, different environment entirely — so this field's data no longer describes what a
|
||||
* member pane inherits. See {@link #logCredentialGap} for how that mode is handled.
|
||||
*/
|
||||
private final Supplier<Set<String>> hostEnvNames;
|
||||
|
||||
@@ -176,8 +170,6 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
|
||||
/** Guards {@link #warnNonZsh} to one WARN per launcher instance, not one per spawn. */
|
||||
private final AtomicBoolean nonZshShellWarned = new AtomicBoolean();
|
||||
/** Live config provides URI environment names that must never enter member panes. */
|
||||
private final Supplier<FleetConfig> config;
|
||||
|
||||
/**
|
||||
* @param namePrefix label prefix for this peer kind (drives naming and reap)
|
||||
@@ -246,19 +238,8 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
long spawnReadyTimeoutMs,
|
||||
LongSupplier nowMillis, Runnable sleeper,
|
||||
Supplier<FleetConfig.Fleet> fleet,
|
||||
Supplier<FleetConfig.MemberCredentials> memberCredentials,
|
||||
Supplier<Set<String>> hostEnvNames) {
|
||||
this(namePrefix, agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs, nowMillis,
|
||||
sleeper, fleet, memberCredentials, hostEnvNames, null);
|
||||
}
|
||||
|
||||
/** As above, plus the live full config for secret-bearing URI environment names. */
|
||||
protected HerdrPeerLauncher(String namePrefix, AgentControl agents, WorkspaceControl spaces,
|
||||
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
|
||||
Function<String, String> env, long spawnReadyTimeoutMs,
|
||||
LongSupplier nowMillis, Runnable sleeper, Supplier<FleetConfig.Fleet> fleet,
|
||||
Supplier<FleetConfig.MemberCredentials> memberCredentials,
|
||||
Supplier<Set<String>> hostEnvNames, Supplier<FleetConfig> config) {
|
||||
Supplier<Set<String>> hostEnvNames) {
|
||||
this.fleet = fleet;
|
||||
this.namePrefix = namePrefix;
|
||||
this.agents = agents;
|
||||
@@ -271,7 +252,6 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
this.sleeper = sleeper;
|
||||
this.memberCredentials = memberCredentials;
|
||||
this.hostEnvNames = hostEnvNames != null ? hostEnvNames : () -> System.getenv().keySet();
|
||||
this.config = config;
|
||||
}
|
||||
|
||||
// --- adapter seams -------------------------------------------------------------------------
|
||||
@@ -1048,11 +1028,9 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
}
|
||||
|
||||
/** Put {@link #BLOCKED_CREDENTIAL_SENTINEL} over every blocked name in the pane-creation env map. */
|
||||
private void overlayBlockedCredentials(Map<String, String> workerEnv,
|
||||
FleetConfig.MemberCredentials creds) {
|
||||
Set<String> blocked = new java.util.TreeSet<>(creds.blockedSet());
|
||||
blocked.addAll(brokerUriEnvNames());
|
||||
for (String name : blocked) {
|
||||
private static void overlayBlockedCredentials(Map<String, String> workerEnv,
|
||||
FleetConfig.MemberCredentials creds) {
|
||||
for (String name : creds.blockedSet()) {
|
||||
workerEnv.put(name, BLOCKED_CREDENTIAL_SENTINEL);
|
||||
}
|
||||
}
|
||||
@@ -1077,29 +1055,6 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
* name. It stays a one-off decision because it is a live handle to the operator's ssh-agent, not
|
||||
* a value — a member holding it can sign with every key the agent holds, so letting it ride in
|
||||
* on the generic {@code allow:} list would hand that out for an unrelated reason.
|
||||
*
|
||||
* <p>fleetd #213: which shell decides the zsh gate, and where the generated directory lives,
|
||||
* both depend on whether {@code memberHerdrSocket:} is configured — see {@link
|
||||
* #memberHerdrSocketConfigured()}'s javadoc for why fleetd's own {@code $SHELL} and {@code
|
||||
* java.io.tmpdir} describe the wrong process once member panes run under a different OS user.
|
||||
* <ul>
|
||||
* <li>{@code memberHerdrSocket} ABSENT (today's only mode): byte-identical to before this
|
||||
* fix — fleetd's own {@code $SHELL} decides zsh, and the directory is generated under
|
||||
* {@code java.io.tmpdir}.</li>
|
||||
* <li>{@code memberHerdrSocket} PRESENT: the configured {@code memberLoginShell:} decides
|
||||
* zsh instead — fleetd's own {@code $SHELL} is never consulted, since it names a
|
||||
* different user's shell, not the member's. Absent/non-zsh falls back exactly like the
|
||||
* non-zsh case below. When it IS zsh, the directory still cannot go under {@code
|
||||
* java.io.tmpdir} (mode 0700, unreadable by another uid — the exact gap fleetd #213
|
||||
* exists to close), so it is generated under {@code worktreeRoot} instead and shared
|
||||
* read-only with {@code worktreeGroup} — the same group {@link
|
||||
* dev.ltms.fleet.session.Worktrees#shareWithGroup} already uses, reused rather than
|
||||
* inventing a second group key. Either one missing means the scrub cannot be guaranteed
|
||||
* reachable by the member, which is the same "cannot guarantee the scrub runs" case as a
|
||||
* non-zsh shell, so it gets the identical fallback.</li>
|
||||
* </ul>
|
||||
* In every branch: never refuse to spawn. A degraded credential control must not become an
|
||||
* outage for an opt-in feature.
|
||||
*/
|
||||
private Path applyEnvironmentAllowListPolicy(FleetConfig.Profile cfg, Launch launch) {
|
||||
FleetConfig.MemberCredentials creds = memberCredentials == null ? null : memberCredentials.get();
|
||||
@@ -1107,14 +1062,8 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
return null;
|
||||
}
|
||||
Set<String> allowed = derivedAllowedNames(creds, launch);
|
||||
boolean memberHerdrSocket = memberHerdrSocketConfigured();
|
||||
// fleetd #213 defect 1: under memberHerdrSocket the member pane runs as a DIFFERENT OS
|
||||
// user, so fleetd's own $SHELL says nothing about what that pane runs — resolveEnv("SHELL")
|
||||
// must not even be called on this path, only the explicit memberLoginShell: config can
|
||||
// answer it. With memberHerdrSocket absent, nothing here changes: fleetd's own $SHELL is
|
||||
// still the input, exactly as before this fix.
|
||||
String loginShell = memberHerdrSocket ? configuredMemberLoginShell() : resolveEnv("SHELL");
|
||||
boolean zsh = isZshShell(loginShell);
|
||||
String loginShell = resolveEnv("SHELL");
|
||||
boolean zsh = loginShell != null && (loginShell.endsWith("/zsh") || loginShell.equals("zsh"));
|
||||
if (!zsh) {
|
||||
// A non-zsh login shell ignores ZDOTDIR entirely: NO scrub would run, so pretending
|
||||
// otherwise would be worse than saying so. Warn loudly and fall back to the CB-596
|
||||
@@ -1128,35 +1077,13 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
logCredentialGap(creds, null);
|
||||
return null;
|
||||
}
|
||||
Path parentDir;
|
||||
String group = null;
|
||||
if (memberHerdrSocket) {
|
||||
// fleetd #213 defect 2: java.io.tmpdir is fleetd's own per-user temp dir (mode 0700 on
|
||||
// macOS) — a member running as a different uid cannot even traverse it, let alone read
|
||||
// the generated files. worktreeRoot is the only configured location a different-uid
|
||||
// member can be given access to, and only WITH worktreeGroup to grant that access —
|
||||
// absent either, the scrub cannot be guaranteed reachable, so this falls back exactly
|
||||
// like the non-zsh case above rather than generating a directory nothing can read.
|
||||
parentDir = memberScrubParentDir();
|
||||
group = memberGroup();
|
||||
if (parentDir == null || group == null) {
|
||||
warnCannotShareScrubDirectory();
|
||||
overlayBlockedCredentials(launch.env(), creds);
|
||||
logCredentialGap(creds, null);
|
||||
return null;
|
||||
}
|
||||
} else {
|
||||
parentDir = Path.of(System.getProperty("java.io.tmpdir"));
|
||||
}
|
||||
// Only reached when the scrub is actually about to run — the count below describes that
|
||||
// scrub, so it must not be logged before this gate (see the non-zsh branch above). Same
|
||||
// reasoning gates logCredentialGap's wording: passing the derived `allowed` set (non-null)
|
||||
// here, and ONLY here, is what tells it the scrub will really blank an unkept name — #192.
|
||||
logAllowListCoverage(allowed);
|
||||
logCredentialGap(creds, allowed);
|
||||
Path dir = memberHerdrSocket
|
||||
? EnvAllowListScrub.generate(parentDir, allowed, group)
|
||||
: EnvAllowListScrub.generate(parentDir, allowed);
|
||||
Path dir = EnvAllowListScrub.generate(Path.of(System.getProperty("java.io.tmpdir")), allowed);
|
||||
launch.env().put("ZDOTDIR", dir.toAbsolutePath().toString());
|
||||
log.info("memberCredentials policy=allow-list: profile={} generated ZDOTDIR {} — derived "
|
||||
+ "allow-list holds {} name(s); the pane reports allowed N of M at release",
|
||||
@@ -1164,76 +1091,21 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
return dir;
|
||||
}
|
||||
|
||||
/** True when {@code shell} is a zsh login shell path or bare name — the ZDOTDIR gate. */
|
||||
private static boolean isZshShell(String shell) {
|
||||
return shell != null && (shell.endsWith("/zsh") || shell.equals("zsh"));
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #213: the configured {@code memberLoginShell:}, or {@code null} when unconfigured.
|
||||
* Called ONLY from the {@code memberHerdrSocket}-configured branch of {@link
|
||||
* #applyEnvironmentAllowListPolicy} — fleetd's own {@code $SHELL} is never read on that path.
|
||||
* {@link #config} being {@code null} (an older test call site, or a launcher that never
|
||||
* threaded the full config through) is treated the same as "not configured".
|
||||
*/
|
||||
private String configuredMemberLoginShell() {
|
||||
FleetConfig cfg = config == null ? null : config.get();
|
||||
return cfg == null ? null : cfg.memberLoginShell();
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #213: {@code worktreeRoot}, as the ZDOTDIR scrub's parent directory under {@code
|
||||
* memberHerdrSocket}, or {@code null} when unconfigured — the same "cannot guarantee the scrub
|
||||
* runs" gap as {@link #memberGroup()} being unset (see {@link
|
||||
* #applyEnvironmentAllowListPolicy}). Deliberately no sibling-of-repo-root default here, unlike
|
||||
* {@code GitWorktrees}' own {@code worktreeRoot} resolution: that default is a convenience for
|
||||
* provisioning a worktree that will exist regardless, whereas an unconfigured value here means
|
||||
* fleetd has no operator-endorsed location to put a credential-bearing directory a different OS
|
||||
* user must reach, so falling back to the overlay is the honest answer, not a guess.
|
||||
*/
|
||||
private Path memberScrubParentDir() {
|
||||
FleetConfig cfg = config == null ? null : config.get();
|
||||
if (cfg == null || cfg.worktreeRoot() == null || cfg.worktreeRoot().isBlank()) {
|
||||
return null;
|
||||
}
|
||||
return Path.of(cfg.worktreeRoot());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #213: the configured {@code worktreeGroup:}, or {@code null} when unset/blank. Reuses
|
||||
* the group {@link dev.ltms.fleet.session.Worktrees#shareWithGroup} already establishes for
|
||||
* provisioned worktrees, rather than a second group key — see {@link
|
||||
* #applyEnvironmentAllowListPolicy}.
|
||||
*/
|
||||
private String memberGroup() {
|
||||
FleetConfig cfg = config == null ? null : config.get();
|
||||
if (cfg == null || cfg.worktreeGroup() == null || cfg.worktreeGroup().isBlank()) {
|
||||
return null;
|
||||
}
|
||||
return cfg.worktreeGroup();
|
||||
}
|
||||
|
||||
/**
|
||||
* The full kept-name set for this spawn: the profile-derived names, unioned with {@code
|
||||
* memberCredentials.allow:} (CB-633 follow-up — previously ignored by this whole policy), the
|
||||
* ssh-agent handle when explicitly allowed, and the exact keys of THIS launch's own env map.
|
||||
*/
|
||||
private Set<String> derivedAllowedNames(FleetConfig.MemberCredentials creds, Launch launch) {
|
||||
Set<String> brokerUriEnvNames = brokerUriEnvNames();
|
||||
Set<String> allowed = new java.util.TreeSet<>(
|
||||
MemberEnvAllowList.derive(profiles.values(), creds.allowSet(), brokerUriEnvNames));
|
||||
MemberEnvAllowList.derive(profiles.values(), creds.allowSet()));
|
||||
if (creds.sshAuthSockAllowed()) {
|
||||
allowed.add(SSH_AUTH_SOCK);
|
||||
} // blocked by default: absent from the set ⇒ blanked by the scrub like any other name
|
||||
allowed.addAll(launch.env().keySet());
|
||||
allowed.removeAll(brokerUriEnvNames);
|
||||
return allowed;
|
||||
}
|
||||
|
||||
private Set<String> brokerUriEnvNames() {
|
||||
return MemberEnvAllowList.brokerUriEnvNames(config == null ? null : config.get());
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-633 follow-up: one INFO line per allow-list spawn WHOSE SCRUB ACTUALLY RUNS, so an operator
|
||||
* can read a single log line and know the scrub ran and how much of the visible environment it
|
||||
@@ -1270,30 +1142,6 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
}
|
||||
}
|
||||
|
||||
/** Guards {@link #warnCannotShareScrubDirectory} to one WARN per launcher instance. */
|
||||
private final AtomicBoolean cannotShareScrubDirWarned = new AtomicBoolean();
|
||||
|
||||
/**
|
||||
* fleetd #213: {@code memberHerdrSocket} is configured and the member login shell IS zsh, but
|
||||
* {@code worktreeRoot} and/or {@code worktreeGroup} is missing, so the generated ZDOTDIR cannot
|
||||
* be placed anywhere the member's OS user can reach — {@code java.io.tmpdir} is fleetd's own
|
||||
* 0700 temp dir, unreadable by another uid, which is the exact gap this ticket exists to close.
|
||||
* Say so once per launcher instance, instead of either generating a directory nothing can read
|
||||
* (protection theatre) or refusing to spawn (turning a degraded credential control into an
|
||||
* outage for an opt-in feature).
|
||||
*/
|
||||
private void warnCannotShareScrubDirectory() {
|
||||
if (cannotShareScrubDirWarned.compareAndSet(false, true)) {
|
||||
log.warn("memberCredentials policy=allow-list: memberHerdrSocket is configured and the "
|
||||
+ "member login shell is zsh, but worktreeRoot and/or worktreeGroup is not "
|
||||
+ "configured — the generated ZDOTDIR cannot be placed where the member's OS "
|
||||
+ "user can read it (java.io.tmpdir is fleetd's own, unreadable by another uid), "
|
||||
+ "so the scrub cannot be guaranteed to run. Falling back to the CB-596 sentinel "
|
||||
+ "overlay. Configure both worktreeRoot and worktreeGroup to enable the "
|
||||
+ "allow-list scrub under memberHerdrSocket.");
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-633 teardown half: read the pane's scrub report (the denominator report the generated
|
||||
* scrub wrote) and delete the directory. Called from {@link #stop}, which is the one funnel
|
||||
@@ -1360,67 +1208,6 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
*/
|
||||
private final AtomicBoolean allowListGapLogged = new AtomicBoolean();
|
||||
|
||||
/**
|
||||
* fleetd #185 stage 2: guards {@link #warnUnknownMemberEnvironment} to one WARN per launcher
|
||||
* instance, not one per spawn — the same one-per-instance shape as {@link #unprotectedGapLogged}
|
||||
* and {@link #allowListGapLogged}, kept as its own flag for the same reason those two are split:
|
||||
* this mode is orthogonal to which of the other two branches would otherwise have fired.
|
||||
*/
|
||||
private final AtomicBoolean unknownMemberEnvironmentWarned = new AtomicBoolean();
|
||||
|
||||
/**
|
||||
* fleetd #185 stage 2: whether {@code memberHerdrSocket:} is configured, i.e. member panes run
|
||||
* on a second herdr owned by a different OS user than the daemon's own process. Re-read from the
|
||||
* live config on every call (same hot-reload shape as {@link #memberCredentials}), never cached,
|
||||
* so a config reload takes effect on the next spawn without a restart.
|
||||
*
|
||||
* <p>{@link #config} is {@code null} on any call site that never threaded the full config
|
||||
* through (every production {@code HerdrPeerLauncher} does; a handful of older tests do not) —
|
||||
* treated the same as "not configured", which is the correct, permissive default: it is exactly
|
||||
* today's single-daemon behaviour.
|
||||
*/
|
||||
private boolean memberHerdrSocketConfigured() {
|
||||
if (config == null) {
|
||||
return false;
|
||||
}
|
||||
FleetConfig cfg = config.get();
|
||||
return cfg != null && cfg.memberHerdrSocket() != null && !cfg.memberHerdrSocket().isBlank();
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #185 stage 2: the single replacement WARN for {@link #logCredentialGap}'s usual
|
||||
* conclusions when {@code memberHerdrSocket:} is configured. {@link #hostEnvNames} (and
|
||||
* everything derived from it — {@code known}/{@code allow} coverage, the allow-list scrub's
|
||||
* derived set) describes the DAEMON's own environment; under this config key member panes run as
|
||||
* a different OS user with a different environment entirely, so neither "every member pane
|
||||
* inherits them UNBLOCKED" nor "the scrub blanks them" is evidence-backed here — both would be
|
||||
* reporting on the wrong process. Logged once, names the config key, and states the honest
|
||||
* conclusion: the gap for member panes is UNKNOWN, not clean, so {@code memberCredentials} cannot
|
||||
* be verified from this daemon. The one count it does report is scoped explicitly to fleetd's own
|
||||
* environment, never presented as if it said anything about the member's — see {@link
|
||||
* #logCredentialGap}'s javadoc for why this branch exists.
|
||||
*/
|
||||
private void warnUnknownMemberEnvironment(FleetConfig.MemberCredentials creds) {
|
||||
if (!unknownMemberEnvironmentWarned.compareAndSet(false, true)) {
|
||||
return;
|
||||
}
|
||||
Set<String> covered = new HashSet<>(creds.known());
|
||||
covered.addAll(creds.allow());
|
||||
Set<String> hostNames = hostEnvNames.get();
|
||||
long gapInFleetdsOwnEnv = hostNames.stream()
|
||||
.filter(name -> CREDENTIAL_SHAPED_NAME.matcher(name).matches())
|
||||
.filter(name -> !covered.contains(name))
|
||||
.count();
|
||||
log.warn("memberCredentials gap: memberHerdrSocket is configured, so member panes run under "
|
||||
+ "a different OS user than fleetd's own process, with a different environment "
|
||||
+ "entirely — fleetd has no channel to read that user's environment. {} of the "
|
||||
+ "{} names in fleetd's OWN environment are credential-shaped and not on "
|
||||
+ "known:/allow:, but that count describes fleetd's process, not the member "
|
||||
+ "herdr's. The credential gap for member panes is UNKNOWN, not clean, and "
|
||||
+ "memberCredentials cannot be verified from here.",
|
||||
gapInFleetdsOwnEnv, hostNames.size());
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-596 criterion 4: a credential-shaped host env var name on neither {@code known} nor
|
||||
* {@code allow} is not silently allowed — it is reported. {@link #hostEnvNames} enumerates the
|
||||
@@ -1449,22 +1236,8 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
* (same severity, and same guard, as the deny-by-default case — a name genuinely reaching a
|
||||
* member unprotected is equally serious whichever path put it there), and the names it says are
|
||||
* blanked keep the INFO.
|
||||
*
|
||||
* <p>fleetd #185 stage 2: everything above assumes the member pane runs under the same OS user
|
||||
* as the daemon, so {@link #hostEnvNames} mirrors what the pane inherits — see that field's
|
||||
* javadoc. When {@code memberHerdrSocket:} is configured that assumption is false: the member
|
||||
* pane runs on a second herdr owned by a <em>different</em> user, and neither conclusion below
|
||||
* ("inherits them UNBLOCKED" / "the scrub blanks them") is backed by evidence about that user's
|
||||
* environment. So this method checks that first and, when configured, reports the honest
|
||||
* "unknown, not clean" conclusion instead — see {@link #warnUnknownMemberEnvironment}. When
|
||||
* {@code memberHerdrSocket:} is absent (the default, and the only mode this host runs) this
|
||||
* branch is never taken and every line below is unchanged.
|
||||
*/
|
||||
private void logCredentialGap(FleetConfig.MemberCredentials creds, Set<String> effectiveAllowed) {
|
||||
if (memberHerdrSocketConfigured()) {
|
||||
warnUnknownMemberEnvironment(creds);
|
||||
return;
|
||||
}
|
||||
Set<String> covered = new HashSet<>(creds.known());
|
||||
covered.addAll(creds.allow());
|
||||
List<String> gap = hostEnvNames.get().stream()
|
||||
|
||||
@@ -40,9 +40,8 @@ import java.util.TreeSet;
|
||||
* operator's own explicit list. Before this, {@code policy: allow-list} silently ignored every name
|
||||
* an operator wrote under {@code allow:} unless a profile happened to carry it too, which meant
|
||||
* turning the policy on could blank credentials working members already depended on. {@code
|
||||
* SSH_AUTH_SOCK} and configured broker URI environment names are exceptions: even when the operator
|
||||
* lists them under {@code allow:}, they are excluded here. {@code SSH_AUTH_SOCK} is added back ONLY
|
||||
* by the caller when {@code sshAuthSock: allow} is explicitly set
|
||||
* SSH_AUTH_SOCK} is the one exception: even when the operator lists it under {@code allow:}, it is
|
||||
* excluded here and added back ONLY by the caller when {@code sshAuthSock: allow} is explicitly set
|
||||
* (see {@link #SSH_AUTH_SOCK}'s javadoc) — it is a live handle to the operator's own ssh-agent, not
|
||||
* a value, so treating it like any other allow-listed name would hand a member every key the
|
||||
* operator's agent holds the moment they typed the name under {@code allow:} for an unrelated
|
||||
@@ -107,16 +106,6 @@ public final class MemberEnvAllowList {
|
||||
* run-to-run.
|
||||
*/
|
||||
public static Set<String> derive(Collection<FleetConfig.Profile> profiles, Set<String> configuredAllow) {
|
||||
return derive(profiles, configuredAllow, Set.of());
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #derive(Collection, Set)}, while excluding names that fleetd knows carry credentials.
|
||||
* A configured broker URI contains its AMQP password inline, so it must never reach a member,
|
||||
* even when an operator put its variable name in {@code memberCredentials.allow:}.
|
||||
*/
|
||||
public static Set<String> derive(Collection<FleetConfig.Profile> profiles, Set<String> configuredAllow,
|
||||
Set<String> excludedNames) {
|
||||
Set<String> derived = new TreeSet<>(INFRASTRUCTURE_PASSTHROUGH);
|
||||
if (profiles != null) {
|
||||
for (FleetConfig.Profile p : profiles) {
|
||||
@@ -135,37 +124,9 @@ public final class MemberEnvAllowList {
|
||||
}
|
||||
}
|
||||
}
|
||||
if (excludedNames != null) {
|
||||
derived.removeAll(excludedNames);
|
||||
}
|
||||
return Set.copyOf(derived);
|
||||
}
|
||||
|
||||
/**
|
||||
* The host environment names whose values are AMQP URIs with inline passwords. Both broker
|
||||
* connections belong to fleetd, never to a member pane. Blank and absent configuration changes
|
||||
* nothing.
|
||||
*
|
||||
* <p><b>How strong this exclusion is depends on the policy, and the difference matters.</b>
|
||||
* Under {@code policy: allow-list} it is enforced by the generated ZDOTDIR scrub, which runs
|
||||
* AFTER the pane's shell has sourced the operator's chain — so a login shell that re-exports the
|
||||
* name is still blanked. Under the deny-list policy there is no scrub: the name is only removed
|
||||
* from the pre-shell env map, and a login shell that sources the operator's secret store
|
||||
* re-exports it. That is the long-standing weakness of deny-list (a sourced file can undo it),
|
||||
* not something this exclusion introduces, but it means deny-list deployments do NOT get this
|
||||
* guarantee. The same caveat applies to the non-zsh path, which has no scrub at all — see
|
||||
* {@code HerdrPeerLauncher#applyEnvironmentAllowListPolicy}.
|
||||
*/
|
||||
public static Set<String> brokerUriEnvNames(FleetConfig config) {
|
||||
if (config == null) {
|
||||
return Set.of();
|
||||
}
|
||||
Set<String> names = new TreeSet<>();
|
||||
addUriEnvIfPresent(names, config.broker());
|
||||
addUriEnvIfPresent(names, config.coordinator());
|
||||
return Set.copyOf(names);
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether {@code name} survives the scrub when {@code allowedNames} is the derived set: an exact
|
||||
* match, or an infrastructure-prefixed name ({@code LC_*}). Prefix rules live ONLY here and in
|
||||
@@ -185,16 +146,4 @@ public final class MemberEnvAllowList {
|
||||
into.add(name);
|
||||
}
|
||||
}
|
||||
|
||||
private static void addUriEnvIfPresent(Set<String> into, FleetConfig.Broker broker) {
|
||||
if (broker != null && broker.hasUriEnv()) {
|
||||
into.add(broker.uriEnv());
|
||||
}
|
||||
}
|
||||
|
||||
private static void addUriEnvIfPresent(Set<String> into, FleetConfig.Coordinator coordinator) {
|
||||
if (coordinator != null && coordinator.hasUriEnv()) {
|
||||
into.add(coordinator.uriEnv());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
@@ -71,6 +72,9 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
*/
|
||||
private final OpenCodeSessionDiscovery discovery;
|
||||
|
||||
/** Receives the profile when a resolved model differs from its requested selector. */
|
||||
private final ExhaustionSink modelMismatchSink;
|
||||
|
||||
/**
|
||||
* 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.
|
||||
@@ -116,23 +120,25 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
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) {
|
||||
this(agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs, spawnReadyPollMs,
|
||||
fleet, memberCredentials, null);
|
||||
}
|
||||
|
||||
/** Production constructor, plus the live config for URI environment exclusions. */
|
||||
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) {
|
||||
long spawnReadyTimeoutMs, long spawnReadyPollMs,
|
||||
Supplier<FleetConfig.Fleet> fleet,
|
||||
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
|
||||
this(agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs,
|
||||
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs),
|
||||
defaultConfigRoot(), defaultDiscoveryRoot(), fleet, memberCredentials, config);
|
||||
defaultConfigRoot(), defaultDiscoveryRoot(), fleet, memberCredentials, ExhaustionSink.none());
|
||||
}
|
||||
|
||||
/** Production constructor with permanent-quarantine wiring for a model mismatch. */
|
||||
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,
|
||||
ExhaustionSink modelMismatchSink) {
|
||||
this(agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs,
|
||||
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs),
|
||||
defaultConfigRoot(), defaultDiscoveryRoot(), fleet, memberCredentials, modelMismatchSink);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -174,12 +180,10 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
Function<String, String> env,
|
||||
long spawnReadyTimeoutMs,
|
||||
LongSupplier nowMillis, Runnable sleeper,
|
||||
Path configRoot, Path discoveryRoot,
|
||||
Supplier<FleetConfig.Fleet> fleet) {
|
||||
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
|
||||
spawnReadyTimeoutMs, nowMillis, sleeper, fleet);
|
||||
this.configRoot = configRoot;
|
||||
this.discovery = new OpenCodeSessionDiscovery(discoveryRoot);
|
||||
Path configRoot, Path discoveryRoot,
|
||||
Supplier<FleetConfig.Fleet> fleet) {
|
||||
this(agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs, nowMillis, sleeper,
|
||||
configRoot, discoveryRoot, fleet, null, ExhaustionSink.none());
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -190,25 +194,32 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
Function<String, String> env,
|
||||
long spawnReadyTimeoutMs,
|
||||
LongSupplier nowMillis, Runnable sleeper,
|
||||
Path configRoot, Path discoveryRoot,
|
||||
Supplier<FleetConfig.Fleet> fleet,
|
||||
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
|
||||
Path configRoot, Path discoveryRoot,
|
||||
Supplier<FleetConfig.Fleet> fleet,
|
||||
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
|
||||
this(agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs, nowMillis, sleeper,
|
||||
configRoot, discoveryRoot, fleet, memberCredentials, null);
|
||||
configRoot, discoveryRoot, fleet, memberCredentials, ExhaustionSink.none());
|
||||
}
|
||||
|
||||
/** Full testability constructor, plus the live config for URI environment exclusions. */
|
||||
/**
|
||||
* Full constructor with the model-mismatch quarantine callback. The callback is an
|
||||
* {@link ExhaustionSink} so model mismatches use the existing quarantine path rather than a
|
||||
* second state tracker.
|
||||
*/
|
||||
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,
|
||||
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 modelMismatchSink) {
|
||||
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
|
||||
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials, null, config);
|
||||
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials);
|
||||
this.configRoot = configRoot;
|
||||
this.discovery = new OpenCodeSessionDiscovery(discoveryRoot);
|
||||
this.modelMismatchSink = modelMismatchSink;
|
||||
}
|
||||
|
||||
private static Path defaultConfigRoot() {
|
||||
@@ -521,9 +532,34 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
@Override
|
||||
public PeerHandle spawn(SpawnRequest req) {
|
||||
PeerHandle inner = super.spawn(req);
|
||||
verifyResolvedModel(requireProfile(req.profileName()), effectiveCwd(req));
|
||||
return new SessionAwareHandle(inner, discovery, effectiveCwd(req));
|
||||
}
|
||||
|
||||
/**
|
||||
* Read the session record once after the spawn-readiness gate. opencode writes the model it
|
||||
* actually selected there. No record, incomplete model object, or a selector without one slash
|
||||
* is unknown evidence, so it must not quarantine a working profile.
|
||||
*/
|
||||
private void verifyResolvedModel(FleetConfig.Profile cfg, String cwd) {
|
||||
String[] requested = splitProviderModel(cfg.model());
|
||||
if (requested == null) {
|
||||
return;
|
||||
}
|
||||
OpenCodeSessionDiscovery.SessionRecord record = discovery.sessionForDirectory(cwd);
|
||||
if (record == null || record.providerId() == null || record.modelId() == null) {
|
||||
return;
|
||||
}
|
||||
String actual = record.providerId() + "/" + record.modelId();
|
||||
if (cfg.model().equals(actual)) {
|
||||
return;
|
||||
}
|
||||
log.error("opencode model mismatch for profile '{}': requested '{}' but resolved '{}'",
|
||||
cfg.profile(), cfg.model(), actual);
|
||||
modelMismatchSink.onExhausted(cfg.profile(), "opencode model mismatch: requested "
|
||||
+ cfg.model() + ", resolved " + actual);
|
||||
}
|
||||
|
||||
/**
|
||||
* A {@link PeerHandle} that delegates everything to the base's worker handle but resolves
|
||||
* {@link #agentSessionId()} lazily through opencode session discovery. Delegate-only, so the
|
||||
|
||||
@@ -1,120 +1,143 @@
|
||||
package dev.ltms.fleet.member;
|
||||
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.sqlite.SQLiteConfig;
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.sql.Connection;
|
||||
import java.sql.PreparedStatement;
|
||||
import java.sql.ResultSet;
|
||||
import java.sql.SQLException;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.stream.Stream;
|
||||
|
||||
/**
|
||||
* Resolves the opencode session id for a fleetd worker from opencode's on-disk storage — the
|
||||
* only place this adapter touches opencode's private layout, and deliberately the <em>only</em>
|
||||
* class that does.
|
||||
*
|
||||
* <p><strong>Why this is isolated behind one seam.</strong> The layout is version-coupled and not
|
||||
* a stable contract: opencode persists its session state in a SQLite database at
|
||||
* {@code <storageRoot>/opencode.db} (a {@code session} table, one row per session, keyed by id and
|
||||
* carrying a {@code directory} column). That schema can move between opencode releases exactly
|
||||
* like the JSON-file layout it replaced did (opencode migrated off a one-JSON-file-per-session
|
||||
* tree under {@code <storageRoot>/storage/session/<projectID>/ses_*.json} in January 2026 — that
|
||||
* tree is now a frozen migration artefact nothing writes, which is why this class no longer reads
|
||||
* it). opencode also ships a headless HTTP server that may supersede both of these entirely.
|
||||
* Everything this adapter knows about that private storage — its shape and column names — lives
|
||||
* here, so a layout change, or a switch to the HTTP server, changes exactly one class and nothing
|
||||
* in {@link OpenCodeLauncher}.
|
||||
* <p><strong>Why this is isolated behind one seam.</strong> The layout is version-coupled and not a
|
||||
* stable contract: opencode writes one JSON file per session under
|
||||
* {@code <storageRoot>/session/<projectID>/<ses_*.json>}, and each record carries a
|
||||
* {@code "version"} field (e.g. {@code "1.1.31"}), so the exact directory shape, file naming, and
|
||||
* field names can move between opencode releases. opencode also ships a headless HTTP server that
|
||||
* may supersede file scanning entirely. Everything this adapter knows about that private storage —
|
||||
* its shape, naming, and field names — lives here, so a layout change, or a switch to the HTTP
|
||||
* server, changes exactly one class and nothing in {@link OpenCodeLauncher}.
|
||||
*
|
||||
* <p>The determinism that makes this useful is structural, not a guess: every fleetd worker runs
|
||||
* in its own unique git worktree, so the row's {@code directory} (its project root) equals the
|
||||
* in its own unique git worktree, so the record's {@code directory} (its project root) equals the
|
||||
* worker's cwd identifies <em>its</em> session unambiguously. We match on {@code directory} rather
|
||||
* than diffing {@code opencode session list} before/after — that races under concurrent spawns, and
|
||||
* the CLI listing does not even show the directory.
|
||||
*
|
||||
* <p>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
|
||||
* handle) treats that as "identity not resolved yet" and retries later. The database is opened
|
||||
* read-only and never written to: opencode itself may be running and writing it concurrently (WAL
|
||||
* mode), and this class must never disturb that.
|
||||
* <p>All reads are best-effort and never throw: a missing or unreadable storage root, a record that
|
||||
* fails to parse, or a directory with no record yet all yield {@code null}, and the caller (the
|
||||
* session handle) treats that as "identity not resolved yet" and retries later.
|
||||
*/
|
||||
final class OpenCodeSessionDiscovery {
|
||||
|
||||
private static final Logger log = LoggerFactory.getLogger(OpenCodeSessionDiscovery.class);
|
||||
|
||||
private final Path storageRoot; // e.g. ~/.local/share/opencode (injectable for tests)
|
||||
private final Path databasePath;
|
||||
private final AtomicBoolean warnedMissingDatabase = new AtomicBoolean(false);
|
||||
private final ObjectMapper json;
|
||||
|
||||
OpenCodeSessionDiscovery(Path storageRoot) {
|
||||
this.storageRoot = storageRoot;
|
||||
this.databasePath = storageRoot.resolve("opencode.db");
|
||||
this.json = new ObjectMapper();
|
||||
}
|
||||
|
||||
/**
|
||||
* A connection to {@link #databasePath} opened with SQLite's {@code SQLITE_OPEN_READONLY}
|
||||
* flag: it never creates the file, never writes, and never touches WAL or journal mode.
|
||||
* opencode may be running and writing this database concurrently, and this class must never
|
||||
* disturb it.
|
||||
*
|
||||
* <p>Package-private so a test can hold the connection and prove it refuses a write. That is
|
||||
* the only way to pin this property: making the file unwritable does <em>not</em> work,
|
||||
* because SQLite silently downgrades a read-write open of an unwritable file to read-only, so
|
||||
* such a test passes whether or not the flag is set.
|
||||
*/
|
||||
Connection openReadOnly() throws SQLException {
|
||||
SQLiteConfig config = new SQLiteConfig();
|
||||
config.setReadOnly(true);
|
||||
return config.createConnection("jdbc:sqlite:" + databasePath);
|
||||
}
|
||||
|
||||
/**
|
||||
* 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 opencode session id whose record references {@code directory} (the worker's cwd), or
|
||||
* {@code null} when no record matches yet. When several records share the directory — e.g.
|
||||
* repeated spawns into the same worktree — the <em>most recently modified</em> one wins: it is
|
||||
* the session the pane most likely corresponds to.
|
||||
*
|
||||
* <p>Never throws: a missing {@code opencode.db}, a locked/unreadable database, a query
|
||||
* failure, or a directory that has not been persisted yet all resolve to {@code null} rather
|
||||
* than failing a spawn. A fleetd worker's session row is written lazily (when the session is
|
||||
* first persisted), so {@code null} here is the normal answer right after the pane is ready,
|
||||
* and the caller retries later.
|
||||
* <p>Never throws: a missing {@code storageRoot}, an unreadable/malformed record, or a
|
||||
* directory that has not been persisted yet all resolve to {@code null} rather than failing a
|
||||
* spawn. A fleetd worker's session record is written lazily (when the session is first
|
||||
* persisted), so {@code null} here is the normal answer right after the pane is ready, and the
|
||||
* caller retries later.
|
||||
*
|
||||
* @param directory the worker's cwd, as resolved for this spawn
|
||||
* @return the matching session id, or {@code null} if none is known yet
|
||||
*/
|
||||
String sessionIdForDirectory(String directory) {
|
||||
SessionRecord record = sessionForDirectory(directory);
|
||||
return record == null ? null : record.id();
|
||||
}
|
||||
|
||||
/**
|
||||
* The newest session record for {@code directory}, or {@code null} when opencode has not written
|
||||
* one yet. This is the single storage seam for both session identity and resolved-model checks.
|
||||
*/
|
||||
SessionRecord sessionForDirectory(String directory) {
|
||||
if (directory == null || directory.isBlank()) {
|
||||
return null;
|
||||
}
|
||||
if (!Files.isRegularFile(databasePath)) {
|
||||
if (warnedMissingDatabase.compareAndSet(false, true)) {
|
||||
log.warn("opencode session database not found at {} — opencode's on-disk layout "
|
||||
+ "may have moved again; session discovery will keep returning null",
|
||||
databasePath);
|
||||
}
|
||||
Path sessionRoot = storageRoot.resolve("session");
|
||||
if (!Files.isDirectory(sessionRoot)) {
|
||||
return null;
|
||||
}
|
||||
String sql = "SELECT id 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 rows.getString("id");
|
||||
SessionRecord best = null;
|
||||
long bestMtime = Long.MIN_VALUE;
|
||||
try (Stream<Path> projectDirs = Files.list(sessionRoot)) {
|
||||
for (Path projectDir : projectDirs.filter(Files::isDirectory).toList()) {
|
||||
try (Stream<Path> records = Files.list(projectDir)) {
|
||||
for (Path record : records.toList()) {
|
||||
SessionRecord matched = matchRecord(record, directory);
|
||||
if (matched == null) {
|
||||
continue;
|
||||
}
|
||||
long mtime = lastModifiedEpochMillis(record);
|
||||
if (mtime > bestMtime) {
|
||||
bestMtime = mtime;
|
||||
best = matched;
|
||||
}
|
||||
}
|
||||
} catch (IOException ignored) {
|
||||
// one project dir unreadable — skip it; another may still match
|
||||
}
|
||||
}
|
||||
} catch (SQLException e) {
|
||||
// Locked, corrupt, or otherwise unreadable — never fatal to a spawn. Not the
|
||||
// "database moved" signal (the file exists), so this stays below WARN.
|
||||
log.debug("opencode session database unreadable at {}: {}", databasePath, e.toString());
|
||||
} catch (IOException ignored) {
|
||||
// storage root vanished or became unreadable — "no session known yet"
|
||||
return null;
|
||||
}
|
||||
log.debug("no opencode session row for directory (root={}, directory={})",
|
||||
storageRoot, directory);
|
||||
return null;
|
||||
return best;
|
||||
}
|
||||
|
||||
/**
|
||||
* The record's session id when it references {@code directory}, else {@code null}. A record
|
||||
* that is not JSON, lacks {@code id}/{@code directory}, or points at a different directory is
|
||||
* simply not our session; a malformed one is skipped, never fatal.
|
||||
*/
|
||||
private SessionRecord matchRecord(Path record, String directory) {
|
||||
try {
|
||||
JsonNode node = json.readTree(record.toFile());
|
||||
JsonNode id = node == null ? null : node.get("id");
|
||||
JsonNode dir = node == null ? null : node.get("directory");
|
||||
if (id == null || dir == null || !directory.equals(dir.asText())) {
|
||||
return null;
|
||||
}
|
||||
JsonNode model = node.path("model");
|
||||
String providerId = text(model, "providerID");
|
||||
String modelId = text(model, "id");
|
||||
return new SessionRecord(id.asText(), providerId, modelId);
|
||||
} catch (IOException e) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
private static String text(JsonNode node, String name) {
|
||||
JsonNode value = node.get(name);
|
||||
return value == null || value.isNull() || value.asText().isBlank() ? null : value.asText();
|
||||
}
|
||||
|
||||
/** The opencode fields fleetd reads from one session record. Null model fields mean unknown. */
|
||||
record SessionRecord(String id, String providerId, String modelId) {
|
||||
}
|
||||
|
||||
/** The record's last-modified epoch ms, or {@code Long.MIN_VALUE} if unreadable (never wins). */
|
||||
private static long lastModifiedEpochMillis(Path record) {
|
||||
try {
|
||||
return Files.getLastModifiedTime(record).toMillis();
|
||||
} catch (IOException e) {
|
||||
return Long.MIN_VALUE;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -9,7 +9,6 @@ import dev.ltms.fleet.metrics.Metrics;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.UUID;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
@@ -168,12 +167,6 @@ public final class MessageService {
|
||||
private static final class Task {
|
||||
private final String ticket;
|
||||
private final String target;
|
||||
/**
|
||||
* When this task was created (#137 fix): the tiebreaker for which of several open tasks on
|
||||
* one target gets a recovered reply in {@link #abandon} — the oldest, since it is the one
|
||||
* that has been waiting longest.
|
||||
*/
|
||||
private final long createdNanos;
|
||||
private final CompletableFuture<Reply> future = new CompletableFuture<>();
|
||||
/**
|
||||
* When {@link #future} resolved, or {@code null} while it is still pending — the clock
|
||||
@@ -191,7 +184,6 @@ public final class MessageService {
|
||||
private Task(String ticket, String target, LongSupplier nowNanos) {
|
||||
this.ticket = ticket;
|
||||
this.target = target;
|
||||
this.createdNanos = nowNanos.getAsLong();
|
||||
future.whenComplete((reply, ex) -> completedNanos = nowNanos.getAsLong());
|
||||
}
|
||||
}
|
||||
@@ -374,61 +366,21 @@ public final class MessageService {
|
||||
}
|
||||
|
||||
/**
|
||||
* Route a worker's explicit {@code fleet_reply}: resolve an open send, complete an async ticket
|
||||
* still parked waiting on this exact turn's answer, or — only once neither applies — queue it in
|
||||
* the inbox. Unlike the bare {@link Rendezvous#resolve}, a no-waiter result is <em>not</em> a
|
||||
* failure — the reply is held for later drain.
|
||||
* Route a worker's explicit {@code fleet_reply}: resolve an open send, or queue it in the
|
||||
* inbox if no send is currently open. Unlike the bare {@link Rendezvous#resolve}, a no-waiter
|
||||
* result is <em>not</em> a failure — the reply is held for later drain.
|
||||
*
|
||||
* <p><strong>Do NOT use this for mid-turn questions.</strong> {@code fleet_ask} /
|
||||
* {@link Rendezvous#resolveQuestion} must keep today's {@code NO_WAITER} behaviour — questions
|
||||
* are interactive and must never be queued.
|
||||
*
|
||||
* <p><strong>Ambiguous match also falls to the inbox.</strong> {@link #askAnsweredAsyncTasks}
|
||||
* cannot actually return more than one entry today (see its own javadoc for why — in short,
|
||||
* {@link #hasAsyncQuestion} keeps a target BUSY, so no second task can reach this state, for as
|
||||
* long as an earlier one's {@code turnId} is still stamped). That is an emergent guarantee from
|
||||
* two other facts, not one this method enforces, so this branch stays in as defence in depth
|
||||
* rather than being removed as dead code: if it ever weakens, returning whichever candidate a
|
||||
* {@code ConcurrentHashMap} iteration reaches first would let a genuine reply complete the
|
||||
* <em>wrong</em> ticket — silently handing the lead something that reads like a correct answer to
|
||||
* a delegation the worker never touched, which is worse than a failure because the lead acts on
|
||||
* it. When more than one candidate exists, guessing is not safe: fall back to the inbox exactly
|
||||
* as the zero-candidate case does, and let {@link #abandon} apply the eventual recovery
|
||||
* deterministically instead.
|
||||
*
|
||||
* @return always {@code true} — the reply resolved a live send, completed a parked ticket, or
|
||||
* was queued
|
||||
* @return always {@code true} — the reply either resolved a live send or was queued
|
||||
*/
|
||||
public boolean reply(String session, String content) {
|
||||
if (rendezvous.resolve(session, content)) {
|
||||
count(FleetMetrics.REPLIES, "path", "rendezvous");
|
||||
return true; // a live send took it — unchanged fast path
|
||||
}
|
||||
// #137: no live rendezvous waiter, but this may be the worker's real fleet_reply resuming a
|
||||
// turn that {@link #answer} already gave up waiting on. answer()'s own bounded wait (the
|
||||
// primary's fleet_send{turnId} call, capped well under a minute) can time out and close its
|
||||
// waiter long before the worker — now actually resuming real work — finishes and replies. That
|
||||
// reply used to have nowhere to land but the session inbox, leaving the async ticket's future
|
||||
// unresolved forever: fleet_poll{ticket} stayed PENDING until fleet_stop's abandon() forced it
|
||||
// FAILED with a misleading "session released before it replied" reason, even though the reply
|
||||
// had, in fact, arrived. Completing the matching ticket directly here means fleet_poll{ticket}
|
||||
// sees the real reply instead.
|
||||
List<Task> candidates = askAnsweredAsyncTasks(session);
|
||||
if (candidates.size() == 1) {
|
||||
Task orphan = candidates.get(0);
|
||||
if (orphan.future.complete(new Reply(Outcome.REPLIED, content))) {
|
||||
if (orphan.turnId != null) {
|
||||
asyncTasksByTurn.remove(orphan.turnId, orphan);
|
||||
}
|
||||
count(FleetMetrics.REPLIES, "path", "async-recovered");
|
||||
return true; // the ticket itself took it — no inbox stranding at all
|
||||
}
|
||||
} else if (candidates.size() > 1) {
|
||||
List<String> tickets = candidates.stream().map(t -> t.ticket).toList();
|
||||
log.warn("reply from {} matches {} open async tickets {} — cannot tell which one it "
|
||||
+ "answers, queuing to the inbox instead of guessing", session, candidates.size(),
|
||||
tickets);
|
||||
}
|
||||
inbox.publish(session, UUID.randomUUID().toString(), content);
|
||||
// CB-640: record the stranding itself (not just the reply text) so fleet health can see a
|
||||
// worker whose replies keep missing their waiter, not only the queue depth this leaves behind.
|
||||
@@ -442,41 +394,6 @@ public final class MessageService {
|
||||
return true; // held, not lost
|
||||
}
|
||||
|
||||
/**
|
||||
* Every still-open async task on {@code target} whose {@code fleet_ask} was already answered —
|
||||
* its {@link Task#turnId} is stamped but its {@link Task#question} was cleared by {@link #answer}
|
||||
* — yet whose future is not resolved yet (#137). Empty if no such task exists, including the
|
||||
* common case where {@code target}'s worker never used {@code fleet_ask} at all (a task that was
|
||||
* never asked has {@code turnId == null}, so it can never match here and only ever completes
|
||||
* through the ordinary rendezvous fast path in {@link #reply}).
|
||||
*
|
||||
* <p><strong>Returns at most one entry today — verified, not assumed.</strong> {@link #send}
|
||||
* refuses to open a waiter on {@code target} while {@link #hasAsyncQuestion} is true, and that
|
||||
* check matches ANY task whose {@code turnId} is still stamped in {@code asyncTasksByTurn} —
|
||||
* not only while its question is still open. {@link #answer} deliberately leaves that stamp in
|
||||
* place ({@code clearAsyncQuestion(turnId, false)}) until the resumed turn's own future actually
|
||||
* resolves, at which point {@link #finishAsyncTask} both removes the stamp AND completes that
|
||||
* task's future in the same call. So a second task can never reach "{@code turnId} stamped, future
|
||||
* still open" — the exact pair this method matches on — while a first one already holds it: by
|
||||
* the time the stamp is gone, so is the eligibility. This is an emergent property of those two
|
||||
* facts holding together, not something this method (or its callers) enforces on its own — flip
|
||||
* {@code forgetTurn} to {@code true} in that one {@link #answer} call and it silently stops being
|
||||
* true, with nothing left to fail loudly. The callers below still handle "more than one" as
|
||||
* defence in depth against exactly that, not because they exercise it today: {@link #reply}
|
||||
* treats it as unresolvable and falls back to the inbox; {@link #abandon} would pick the oldest
|
||||
* deterministically (its own {@code matching} list has no such guarantee — see its javadoc).
|
||||
*/
|
||||
private List<Task> askAnsweredAsyncTasks(String target) {
|
||||
List<Task> candidates = new ArrayList<>();
|
||||
for (Task task : tasks.values()) {
|
||||
if (target.equals(task.target) && task.question == null && task.turnId != null
|
||||
&& !task.future.isDone()) {
|
||||
candidates.add(task);
|
||||
}
|
||||
}
|
||||
return candidates;
|
||||
}
|
||||
|
||||
/** Record a counter sample when a registry is wired; a no-op in unit tests. */
|
||||
private void count(String name, String... labels) {
|
||||
if (metrics != null) {
|
||||
@@ -518,108 +435,19 @@ public final class MessageService {
|
||||
* <p>Resolving the waiter as a failure — rather than letting it time out — also means the
|
||||
* outcome is counted, so a torn-down delegation stops being invisible to {@code /metrics}.
|
||||
*
|
||||
* <p><strong>#137 defence in depth.</strong> {@link #reply} already hands a worker's real
|
||||
* {@code fleet_reply} straight to the async ticket it belongs to whenever exactly one is still
|
||||
* parked waiting for it (see {@link #askAnsweredAsyncTasks}), so by the time a session is
|
||||
* released its tasks are normally already resolved — this loop's {@code complete} calls are then
|
||||
* harmless no-ops (a {@link CompletableFuture} can only resolve once). But should some other path
|
||||
* someday strand a reply in the inbox without completing its ticket, checking
|
||||
* {@link #hasStrandedReply(String)} here — before ever writing a failure — means a torn-down
|
||||
* session whose worker in fact replied is still reported {@code REPLIED} with that reply's own
|
||||
* text, never the misleading "the worker session was released before it replied" (which also
|
||||
* means the snapshot/worktree recovery hint that follows it never prints once a reply exists).
|
||||
*
|
||||
* <p><strong>At most one task gets the recovered reply — and here, unlike {@link #reply}'s
|
||||
* {@link #askAnsweredAsyncTasks}, {@code matching.size() >= 2} alone is reachable today.</strong>
|
||||
* This method's {@code matching} filter has no {@code turnId != null} requirement, so it matches
|
||||
* any plain (never-asked) open task too — and {@link #sendAsync} does not limit a target to one
|
||||
* of those: a second {@code fleet_send{wait:false}} at a target that is still busy returns its own
|
||||
* ticket immediately and simply parks its {@link #send} behind the target's session lock for up
|
||||
* to {@link #ASYNC_TIMEOUT_MS}, exactly as {@code abandonFailsEveryPendingAsyncTicketForTheReleasedTarget}
|
||||
* already proves. Before this fix, the loop below drained the strand once and then reused that
|
||||
* same {@code Reply} for <em>every</em> task it walked past — so two open tasks really did both
|
||||
* complete {@code REPLIED} with the same text (see the pre-fix loop in commit 97f6c33's parent).
|
||||
* A stranded reply is one worker answer, so it can settle at most one open task on this target —
|
||||
* never every open task, and never a guess. When more than one task is still open here, the
|
||||
* recovered reply goes to the <em>oldest</em> (lowest {@link Task#createdNanos}) — it has been
|
||||
* waiting longest, so it is the one most likely to be what the reply actually answers. Every
|
||||
* other open task keeps the ordinary {@code WORKER_FAILED} path it would take without a stranded
|
||||
* reply at all.
|
||||
*
|
||||
* <p><strong>{@code matching.size() >= 2} together with {@code hadStrandedReply} is a different
|
||||
* question, and today it is defence in depth rather than a path this codebase's public API can
|
||||
* drive.</strong> This class has exactly two sites that ever acquire a target's entry in
|
||||
* {@code sessionLocks} — {@link #send} and {@link #answer} — and both open a {@link Rendezvous}
|
||||
* waiter for that same target as the very first thing they do after acquiring the lock, then hold
|
||||
* lock and waiter together for the rest of their critical section ({@link #send} also clears
|
||||
* {@link #strandedReplies} right there, the instant it opens its waiter — before it ever enqueues
|
||||
* delivery). So "the session lock is held" and "a live waiter is open for it" are the same fact
|
||||
* throughout this class, and {@link #reply}'s fast path always resolves a currently-open waiter
|
||||
* directly rather than stranding. The two facts this method wants therefore cannot be produced
|
||||
* side by side: while the lock is held, a real reply resolves the open waiter directly and never
|
||||
* reaches {@link #strandedReplies}; the instant the lock is free, any parked matching task's own
|
||||
* {@link #send} that is scheduled next wins it and, by opening its waiter, clears the strand again
|
||||
* before this method ever runs. There is no way to hold that lock open-but-unaccepted from outside
|
||||
* {@link #send}/{@link #answer} to freeze a window in between. Constructing both facts at once
|
||||
* through {@code sendAsync}/{@code reply}/{@code ask}/{@code answer} would need a race against
|
||||
* virtual-thread scheduling, not a deterministic sequence — so the oldest-wins code below stays as
|
||||
* defence in depth against a regression to that mechanism (e.g. clearing {@link #strandedReplies}
|
||||
* on a narrower condition than "any acceptance"), not because today's test suite exercises the
|
||||
* conjunction. {@code matching.size() >= 2} alone, without a strand, is exactly what
|
||||
* {@code abandonFailsEveryPendingAsyncTicketForTheReleasedTarget} already covers.
|
||||
*
|
||||
* <p><strong>Do not "fix" that gap with a test that reaches past this class.</strong> A test can
|
||||
* build both facts by calling {@link Rendezvous#close} itself on the waiter an accepted
|
||||
* {@link #send} is still blocked on: the send keeps the lock, no waiter is registered any more, a
|
||||
* second parked task stays open, and the next {@link #reply} then strands. That was checked, and
|
||||
* such a test does go red against the pre-fix loop. But it only goes red because it broke the
|
||||
* lock-and-waiter invariant above from outside — no caller of this class ever does that — so it
|
||||
* pins a state production cannot reach, and would read to the next person as if it could.
|
||||
*
|
||||
* @return true if a live waiter or an async task was failed (never true for one recovered as a
|
||||
* reply — see the note above)
|
||||
* @return true if a live waiter was failed
|
||||
*/
|
||||
public boolean abandon(String target, String reason) {
|
||||
boolean hadStrandedReply = hasStrandedReply(target);
|
||||
// CB-640: the session is gone — nothing will ever accept or deliver into it now.
|
||||
strandedReplies.remove(target);
|
||||
queuedDeliveries.remove(target);
|
||||
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(target);
|
||||
boolean failed = waiter != null && !waiter.isDone() && rendezvous.resolveFailure(waiter, reason);
|
||||
boolean asyncFailed = false;
|
||||
|
||||
List<Task> matching = new ArrayList<>();
|
||||
for (Task task : tasks.values()) {
|
||||
if (target.equals(task.target) && task.question == null && !task.future.isDone()) {
|
||||
matching.add(task);
|
||||
}
|
||||
}
|
||||
Task recoveryTask = null;
|
||||
if (hadStrandedReply && !matching.isEmpty()) {
|
||||
recoveryTask = matching.get(0);
|
||||
for (Task candidate : matching) {
|
||||
if (candidate.createdNanos < recoveryTask.createdNanos) {
|
||||
recoveryTask = candidate;
|
||||
}
|
||||
}
|
||||
}
|
||||
Reply recovered = recoveryTask != null ? recoverStrandedReply(target) : null;
|
||||
for (Task task : matching) {
|
||||
boolean isRecovery = task == recoveryTask && recovered != null;
|
||||
Reply outcome = isRecovery ? recovered : new Reply(Outcome.WORKER_FAILED, reason);
|
||||
if (task.future.complete(outcome)) {
|
||||
if (outcome.outcome() == Outcome.WORKER_FAILED) {
|
||||
asyncFailed = true;
|
||||
} else if (task.turnId != null) {
|
||||
asyncTasksByTurn.remove(task.turnId, task);
|
||||
}
|
||||
} else if (isRecovery) {
|
||||
// The recovered reply was already drained out of the inbox, but this task resolved
|
||||
// through another path (e.g. a concurrent reply() or a second abandon() racing this
|
||||
// one) between us choosing it and completing it here. Put the reply back rather than
|
||||
// lose it silently — it may still belong to some other still-open task, or the next
|
||||
// caller that drains this target's inbox.
|
||||
inbox.publish(target, UUID.randomUUID().toString(), recovered.text());
|
||||
if (target.equals(task.target) && task.question == null
|
||||
&& task.future.complete(new Reply(Outcome.WORKER_FAILED, reason))) {
|
||||
asyncFailed = true;
|
||||
}
|
||||
}
|
||||
if (failed) {
|
||||
@@ -628,25 +456,6 @@ public final class MessageService {
|
||||
return failed || asyncFailed;
|
||||
}
|
||||
|
||||
/**
|
||||
* Drain {@code target}'s inbox and hand its content back as a {@link Outcome#REPLIED} result
|
||||
* (#137 defence in depth for {@link #abandon}) — {@code null} if it turned out empty (the
|
||||
* stranding fact raced away, e.g. a lead's own {@code fleet_poll} on the raw session already
|
||||
* drained it first). When more than one message is queued, only the newest is the worker's actual
|
||||
* final answer ({@link #drainReplies} returns them oldest-first).
|
||||
*
|
||||
* <p>This does drain (removes the messages from the inbox) before the caller knows whether the
|
||||
* task it is recovering for will actually accept them — {@link #abandon} is the one that puts a
|
||||
* reply back if its {@code complete} call turns out to lose the race.
|
||||
*/
|
||||
private Reply recoverStrandedReply(String target) {
|
||||
var messages = drainReplies(target);
|
||||
if (messages.isEmpty()) {
|
||||
return null;
|
||||
}
|
||||
return new Reply(Outcome.REPLIED, messages.get(messages.size() - 1).content());
|
||||
}
|
||||
|
||||
/**
|
||||
* Acknowledge a specific reply by {@code msgId} for {@code target}. Removes it from the inbox
|
||||
* so that a subsequent drain or peek no longer returns it.
|
||||
|
||||
@@ -76,6 +76,20 @@ public final class BackendQuarantine {
|
||||
quarantinedUntilNanos.put(credentialId, nowNanos.getAsLong() + cooldownNanos);
|
||||
}
|
||||
|
||||
/**
|
||||
* Quarantine {@code credentialId} until the daemon restarts. This is for a configuration error
|
||||
* that cannot heal with time, unlike an exhausted backend. A model selector that opencode silently
|
||||
* resolves to another model stays wrong until an operator changes the profile, so a cooldown would
|
||||
* start the same unsafe work again.
|
||||
*/
|
||||
public void quarantinePermanently(String credentialId) {
|
||||
Objects.requireNonNull(credentialId, "credentialId");
|
||||
if (inert) {
|
||||
return;
|
||||
}
|
||||
quarantinedUntilNanos.put(credentialId, Long.MAX_VALUE);
|
||||
}
|
||||
|
||||
/** Whether {@code credentialId} is quarantined right now. */
|
||||
public boolean isQuarantined(String credentialId) {
|
||||
return remainingNanos(credentialId) > 0;
|
||||
@@ -110,6 +124,7 @@ public final class BackendQuarantine {
|
||||
}
|
||||
|
||||
private static long toSecondsRoundedUp(long nanos) {
|
||||
return (nanos + 999_999_999L) / 1_000_000_000L;
|
||||
long seconds = nanos / 1_000_000_000L;
|
||||
return seconds + (nanos % 1_000_000_000L == 0 ? 0 : 1);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -315,18 +315,10 @@ public final class FleetApp {
|
||||
.map(Agent.class::cast)
|
||||
.filter(a -> a.terminalId() != null)
|
||||
.collect(Collectors.toMap(Agent::terminalId, Function.identity(), (_, b) -> b));
|
||||
// fleetd #209: this REST roster reports agentSessionId via SessionManager.rosterView, so it
|
||||
// uses the resolving roster read (caller-driven, not a timer) rather than the plain one.
|
||||
List<Map<String, Object>> out = sessions.rosterResolved().stream()
|
||||
List<Map<String, Object>> out = sessions.roster().stream()
|
||||
.map(s -> SessionManager.rosterView(s, live.get(s.terminalId())))
|
||||
.toList();
|
||||
Map<String, Object> body = new LinkedHashMap<>();
|
||||
// fleetd #199: the endpoint became /members in the CB-634 rename but the body key stayed
|
||||
// "workers", so a caller that read "members" saw an empty fleet and reported no members at
|
||||
// all. "members" is the canonical key; "workers" stays as a deprecated alias so an existing
|
||||
// REST consumer keeps working — the out-of-band path a lead falls back to when its MCP mount
|
||||
// drops reads this endpoint. Drop the alias once nothing reads it.
|
||||
body.put("members", out);
|
||||
body.put("workers", out);
|
||||
// CB-586: operator visibility for the refs/wip snapshot store without shelling into the
|
||||
// repo — how many snapshot refs exist and roughly what they cost. Present only once a
|
||||
|
||||
@@ -23,7 +23,6 @@ import java.util.Set;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import java.util.function.Consumer;
|
||||
import java.util.function.Function;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
/**
|
||||
@@ -89,64 +88,24 @@ public final class GitWorktrees implements Worktrees {
|
||||
);
|
||||
|
||||
private final String configuredRoot;
|
||||
/** OS group name for {@link #shareWithGroup} (fleetd #185 stage 3); {@code null} ⇒ feature off. */
|
||||
private final String group;
|
||||
private final Consumer<String> afterWorktreeAdded;
|
||||
/** How {@link #shareWithGroup}'s processes (git config / chgrp / chmod / find) actually run.
|
||||
* Defaults to the real {@link #exec(String...)}. Package-private test seam so a unit test can
|
||||
* prove "no group configured ⇒ zero processes spawned" and inspect exactly what a configured
|
||||
* group runs, without a real second OS user or OS group on this host. */
|
||||
private final Function<String[], String> shareGroupRunner;
|
||||
private final SecureRandom random = new SecureRandom();
|
||||
private final AtomicLong seq = new AtomicLong();
|
||||
|
||||
/** Default constructor: worktree root is derived per-repo as {@code <repoRoot>/../.bridged-worktrees}. */
|
||||
public GitWorktrees() {
|
||||
this(null, (String) null);
|
||||
this(null);
|
||||
}
|
||||
|
||||
/**
|
||||
* @param configuredRoot nullable absolute or relative path; null/blank derives a sibling of
|
||||
* the repo root. No {@code worktreeGroup} configured — {@link #shareWithGroup}
|
||||
* is a no-op.
|
||||
*/
|
||||
/** @param configuredRoot nullable absolute or relative path; null/blank derives a sibling of the repo root. */
|
||||
public GitWorktrees(String configuredRoot) {
|
||||
this(configuredRoot, (String) null);
|
||||
}
|
||||
|
||||
/**
|
||||
* @param configuredRoot nullable absolute or relative path; null/blank derives a sibling of
|
||||
* the repo root.
|
||||
* @param group optional OS group name (fleetd #185 stage 3, {@code worktreeGroup:} in
|
||||
* config); null/blank ⇒ {@link #shareWithGroup} is a no-op.
|
||||
*/
|
||||
public GitWorktrees(String configuredRoot, String group) {
|
||||
this(configuredRoot, group, _ -> {});
|
||||
this(configuredRoot, _ -> {});
|
||||
}
|
||||
|
||||
/** Test seam for changing a real worktree between its creation and its security check. */
|
||||
GitWorktrees(String configuredRoot, Consumer<String> afterWorktreeAdded) {
|
||||
this(configuredRoot, null, afterWorktreeAdded);
|
||||
}
|
||||
|
||||
/** Test seam combining a configurable {@code group} with {@link #afterWorktreeAdded}. */
|
||||
GitWorktrees(String configuredRoot, String group, Consumer<String> afterWorktreeAdded) {
|
||||
this(configuredRoot, group, afterWorktreeAdded, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* Full test seam: also overrides how {@link #shareWithGroup}'s processes run (fleetd #185
|
||||
* stage 3), so a unit test can prove "no group configured ⇒ no process spawned" and inspect
|
||||
* exactly what commands a configured group runs, without a real second OS user/group.
|
||||
*
|
||||
* @param shareGroupRunner {@code null} ⇒ the real {@link #exec(String...)}.
|
||||
*/
|
||||
GitWorktrees(String configuredRoot, String group, Consumer<String> afterWorktreeAdded,
|
||||
Function<String[], String> shareGroupRunner) {
|
||||
this.configuredRoot = configuredRoot;
|
||||
this.group = (group == null || group.isBlank()) ? null : group;
|
||||
this.afterWorktreeAdded = afterWorktreeAdded == null ? _ -> {} : afterWorktreeAdded;
|
||||
this.shareGroupRunner = shareGroupRunner != null ? shareGroupRunner : this::exec;
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -648,126 +607,6 @@ public final class GitWorktrees implements Worktrees {
|
||||
return deleted;
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritDoc}
|
||||
*
|
||||
* <p>fleetd #185 stage 3. No-op — no process spawned, nothing logged — when {@link #group} is
|
||||
* null/blank. Otherwise:
|
||||
* <ol>
|
||||
* <li>{@code git -C repoRoot config core.sharedRepository group} so every future write by
|
||||
* either uid stays group-writable;</li>
|
||||
* <li>a one-time {@code chgrp}/{@code chmod g+rwX} fix-up over the worktree directory and,
|
||||
* under the repo's <em>common</em> git directory, {@code objects}, {@code refs},
|
||||
* {@code logs}, {@code worktrees} and {@code packed-refs} — with setgid
|
||||
* ({@code chmod g+s}) applied only to the directories among them, so files created later
|
||||
* inherit the group;</li>
|
||||
* <li>one INFO line naming the group and the paths touched.</li>
|
||||
* </ol>
|
||||
*
|
||||
* <p><b>Every path is skipped when it does not exist.</b> {@code .git/logs} is absent in a repo
|
||||
* with {@code core.logAllRefUpdates=false} or one that has had no ref update yet, and
|
||||
* {@code packed-refs} is absent until refs are packed. Passing a missing path to {@code chgrp}
|
||||
* exits non-zero, which would fail <em>every</em> provisioning spawn with a message blaming a
|
||||
* group that is in fact fine.
|
||||
*
|
||||
* <p><b>The git directory is resolved, not assumed.</b> {@code <repoRoot>/.git} is a
|
||||
* <em>file</em>, not a directory, when the checkout is itself a linked worktree — the very
|
||||
* thing this class creates for every member. {@code git rev-parse --git-common-dir} gives the
|
||||
* real shared store, and it may answer relatively, so it is resolved against {@code repoRoot}.
|
||||
*
|
||||
* <p><b>The fix-up re-runs on every spawn, by design.</b> {@code core.sharedRepository=group}
|
||||
* governs only what git writes <em>after</em> it is set; the walk is what covers everything
|
||||
* already on disk. It is not redundant work to optimise away — dropping it silently leaves
|
||||
* pre-existing objects unreadable to the member. It costs three walks of the object store per
|
||||
* spawn (about 3000 files in this repo, well under a second, but it grows with the repo).
|
||||
*
|
||||
* <p>This only fixes up file ownership/permissions on the operator's shared repo so a
|
||||
* different-uid member can write to it — it isolates credentials, not the repository. A member
|
||||
* in the group can still write the operator's git objects and refs.
|
||||
*
|
||||
* <p>Fails loudly: a missing group, or a {@code chgrp}/{@code chmod} refused because the
|
||||
* operator is not a member of it, becomes a {@link WorktreeException} naming the group — never
|
||||
* a silent skip that leaves a member unable to work with nothing in the log to explain why.
|
||||
*/
|
||||
@Override
|
||||
public void shareWithGroup(String repoRoot, String worktreePath) {
|
||||
if (group == null) {
|
||||
return;
|
||||
}
|
||||
List<String> touched = new ArrayList<>();
|
||||
try {
|
||||
shareGroupRunner.apply(new String[]{"git", "-C", repoRoot, "config", "core.sharedRepository", "group"});
|
||||
String commonDir = gitCommonDir(repoRoot);
|
||||
shareGroupPathIfPresent(worktreePath, true, touched);
|
||||
for (String name : List.of("objects", "refs", "logs", "worktrees")) {
|
||||
shareGroupPathIfPresent(commonDir + "/" + name, true, touched);
|
||||
}
|
||||
shareGroupPathIfPresent(commonDir + "/packed-refs", false, touched);
|
||||
} catch (WorktreeException e) {
|
||||
throw new WorktreeException("cannot share worktree with 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={} shared repoRoot={} worktreePath={} paths={}",
|
||||
group, repoRoot, worktreePath, touched);
|
||||
}
|
||||
|
||||
/**
|
||||
* 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}
|
||||
* answers relative to {@code repoRoot} in the ordinary case ({@code .git}) and absolutely for a
|
||||
* linked worktree, so the answer is resolved against {@code repoRoot} either way. Never
|
||||
* hardcode {@code repoRoot + "/.git"}: that is a FILE when the checkout is itself a linked
|
||||
* worktree.
|
||||
*/
|
||||
private String gitCommonDir(String repoRoot) {
|
||||
String answer = shareGroupRunner.apply(
|
||||
new String[]{"git", "-C", repoRoot, "rev-parse", "--git-common-dir"});
|
||||
String trimmed = answer == null ? "" : answer.trim();
|
||||
if (trimmed.isEmpty()) {
|
||||
trimmed = ".git";
|
||||
}
|
||||
return Path.of(repoRoot).resolve(trimmed).normalize().toString();
|
||||
}
|
||||
|
||||
/**
|
||||
* {@link #shareGroupPath} when {@code path} exists, recording it in {@code touched}; otherwise
|
||||
* nothing at all. A missing path is normal, not an error — see {@link #shareWithGroup}'s
|
||||
* javadoc for which ones are routinely absent and why passing them to {@code chgrp} would fail
|
||||
* every spawn.
|
||||
*/
|
||||
private void shareGroupPathIfPresent(String path, boolean recursive, List<String> touched) {
|
||||
if (!Files.exists(Path.of(path))) {
|
||||
return;
|
||||
}
|
||||
shareGroupPath(path, recursive);
|
||||
touched.add(path);
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code chgrp}/{@code chmod g+rwX} {@code path} to {@link #group}. When {@code recursive},
|
||||
* also walks the directories under {@code path} (including {@code path} itself, when it is a
|
||||
* directory) and sets setgid on each — directories only, per the javadoc on
|
||||
* {@link #shareWithGroup}.
|
||||
*/
|
||||
private void shareGroupPath(String path, boolean recursive) {
|
||||
List<String> chgrp = new ArrayList<>(List.of("chgrp"));
|
||||
if (recursive) chgrp.add("-R");
|
||||
chgrp.add(group);
|
||||
chgrp.add(path);
|
||||
shareGroupRunner.apply(chgrp.toArray(new String[0]));
|
||||
|
||||
List<String> chmod = new ArrayList<>(List.of("chmod"));
|
||||
if (recursive) chmod.add("-R");
|
||||
chmod.add("g+rwX");
|
||||
chmod.add(path);
|
||||
shareGroupRunner.apply(chmod.toArray(new String[0]));
|
||||
|
||||
if (recursive) {
|
||||
shareGroupRunner.apply(new String[]{"find", path, "-type", "d", "-exec", "chmod", "g+s", "{}", "+"});
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Every {@code refs/wip/*} ref (see {@link WipRef}). The committer date is read as a unix
|
||||
* count of seconds and converted to millis. {@code %00} (NUL) separates the fields because a
|
||||
|
||||
@@ -89,15 +89,4 @@ public record MemberSession(
|
||||
return new MemberSession(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
|
||||
nowNanos, turnCount + 1, state, worktree, branch, charterReceipt, agentSessionId);
|
||||
}
|
||||
|
||||
/**
|
||||
* Return a copy with {@code agentSessionId} resolved to a non-null value (fleetd #209). Some
|
||||
* adapters (opencode) cannot answer {@link dev.ltms.fleet.peer.PeerHandle#agentSessionId()} at
|
||||
* spawn time — the peer has not persisted its session record yet — so the id is discovered on
|
||||
* a later poll and swapped into the otherwise-immutable session via this wither.
|
||||
*/
|
||||
public MemberSession withAgentSessionId(String agentSessionId) {
|
||||
return new MemberSession(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
|
||||
lastActivityAtNanos, turnCount, state, worktree, branch, charterReceipt, agentSessionId);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -46,15 +46,6 @@ public final class SessionManager implements TurnListener {
|
||||
private final PeerLauncher launcher;
|
||||
private final Worktrees worktrees;
|
||||
private final ConcurrentHashMap<String /*paneId*/, MemberSession> registry = new ConcurrentHashMap<>();
|
||||
/**
|
||||
* fleetd #209: the live {@link PeerHandle} for every registered pane, retained solely so
|
||||
* {@link #resolveAgentSessionId} can re-poll {@link PeerHandle#agentSessionId()} after spawn.
|
||||
* The handle used to go out of scope at the end of the spawn method, so a launcher that answers
|
||||
* the id lazily (opencode — the on-disk session row is written after the pane is created) could
|
||||
* never be re-asked, and {@code fleet_list}/{@code fleet_spawn resumeSessionId} never saw it.
|
||||
* Populated on every spawn path, removed on {@link #release}.
|
||||
*/
|
||||
private final ConcurrentHashMap<String /*paneId*/, PeerHandle> handles = new ConcurrentHashMap<>();
|
||||
private final MemberPresence presence;
|
||||
private final SecureRandom nonceRandom = new SecureRandom();
|
||||
private final AtomicLong nonceSeq = new AtomicLong();
|
||||
@@ -214,7 +205,6 @@ public final class SessionManager implements TurnListener {
|
||||
handle.charterReceipt(),
|
||||
handle.agentSessionId());
|
||||
registry.put(handle.id(), session);
|
||||
handles.put(handle.id(), handle);
|
||||
memberLifecycle.acquired(session.role(), session.profile(), session.terminalId());
|
||||
log.debug("acquired session id={} terminal={} profile={} owner={}",
|
||||
handle.id(), handle.terminalId(), session.profile(), session.ownerTerminal());
|
||||
@@ -267,10 +257,6 @@ public final class SessionManager implements TurnListener {
|
||||
*/
|
||||
private void release(String paneId, ReleaseCause cause) {
|
||||
MemberSession removed = registry.remove(paneId);
|
||||
// fleetd #209: remove right alongside the registry entry so a released session's handle is
|
||||
// never leaked — but keep the local reference below, so the id can still be resolved for
|
||||
// the ReleaseDetail this teardown notifies with.
|
||||
PeerHandle removedHandle = handles.remove(paneId);
|
||||
boolean preserveWorktree = cause == ReleaseCause.SHUTDOWN;
|
||||
String snapshotRef = null;
|
||||
if (removed != null) {
|
||||
@@ -317,11 +303,8 @@ public final class SessionManager implements TurnListener {
|
||||
// too, so a failed ticket's detail can point a lead at the same tree to re-dispatch.
|
||||
// CB-584 (issue #65 criterion 5): carry agentSessionId alongside them, so a lead can
|
||||
// also resume the member's conversation, not just re-dispatch onto its files.
|
||||
// fleetd #209: a late-resolving adapter (opencode) may only now have an id — resolve
|
||||
// one last time so a released member's detail carries the id it now has.
|
||||
MemberSession resolved = resolveAgentSessionId(removed, removedHandle);
|
||||
notifyReleased(new ReleaseDetail(resolved.terminalId(), resolved.worktree(),
|
||||
resolved.branch(), snapshotRef, resolved.agentSessionId()));
|
||||
notifyReleased(new ReleaseDetail(removed.terminalId(), removed.worktree(),
|
||||
removed.branch(), snapshotRef, removed.agentSessionId()));
|
||||
}
|
||||
}
|
||||
// CB-581: the pane must always stop, even if the dirty check above threw. A session removed
|
||||
@@ -489,10 +472,6 @@ public final class SessionManager implements TurnListener {
|
||||
try {
|
||||
path = worktrees.add(repoRoot, branch, wt.baseRef());
|
||||
worktrees.overlayParity(repoRoot, path, launcher.parityOverlay(preResolvedProfile));
|
||||
// fleetd #185 stage 3: MUST run after overlayParity, not folded into add() — overlayParity
|
||||
// copies more files into the worktree after add() returns, so sharing the group any earlier
|
||||
// leaves those overlay files operator-owned and read-only for a different-uid member.
|
||||
worktrees.shareWithGroup(repoRoot, path);
|
||||
handle = launcher.spawn(new SpawnRequest(profile, path, callerCwd, sessionName, resumeSessionId, memberRole));
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("spawn failed for profile={} role={} branch={} path={}: {}",
|
||||
@@ -525,7 +504,6 @@ public final class SessionManager implements TurnListener {
|
||||
handle.charterReceipt(),
|
||||
handle.agentSessionId());
|
||||
registry.put(handle.id(), session);
|
||||
handles.put(handle.id(), handle);
|
||||
memberLifecycle.acquired(session.role(), session.profile(), session.terminalId());
|
||||
log.debug("acquired worktree session id={} terminal={} profile={} branch={} path={}",
|
||||
handle.id(), handle.terminalId(), session.profile(), session.branch(), session.worktree());
|
||||
@@ -557,73 +535,16 @@ public final class SessionManager implements TurnListener {
|
||||
return launcher.defaultProfile();
|
||||
}
|
||||
|
||||
/**
|
||||
* The session for {@code paneId}, if it is still registered and not released. fleetd #209:
|
||||
* resolves a still-unknown {@code agentSessionId} against the retained handle before returning,
|
||||
* so {@code fleet_status} sees an id a lazy-resolving adapter has since written.
|
||||
*/
|
||||
/** The session for {@code paneId}, if it is still registered and not released. */
|
||||
public Optional<MemberSession> get(String paneId) {
|
||||
return Optional.ofNullable(registry.get(paneId)).map(this::resolveAgentSessionId);
|
||||
return Optional.ofNullable(registry.get(paneId));
|
||||
}
|
||||
|
||||
/**
|
||||
* Fleet-owned roster: all registered sessions (acquired minus released). Deliberately does
|
||||
* <strong>not</strong> resolve {@code agentSessionId} (fleetd #209 follow-up) — this is the
|
||||
* roster supplier on the heartbeat and health-tick timers ({@code LeadHeartbeatLoop},
|
||||
* {@code FleetHealthMonitor} in {@code Fleetd}), on the placement/exhaustion paths, and on the
|
||||
* metrics scrape ({@code FleetMetrics}), all called far more often than any caller actually
|
||||
* reads {@code agentSessionId}. Resolving here would mean every tick opens a lazy-resolving
|
||||
* adapter's (opencode's) on-disk session store once per member whose id is still unknown — and
|
||||
* for a member whose id never appears, that cost never stops, for the life of the process. Use
|
||||
* {@link #rosterResolved()} instead wherever the id must be current.
|
||||
*/
|
||||
/** Fleet-owned roster: all registered sessions (acquired minus released). */
|
||||
public List<MemberSession> roster() {
|
||||
return List.copyOf(registry.values());
|
||||
}
|
||||
|
||||
/**
|
||||
* {@link #roster()}, with each session's still-unknown {@code agentSessionId} re-resolved
|
||||
* against its retained handle (fleetd #209) — so a caller that actually reports the id (
|
||||
* {@code fleet_list}, the REST roster) sees one a lazy-resolving adapter (opencode) has since
|
||||
* written, rather than the null frozen in at spawn time. Reserved for caller-driven reads, not
|
||||
* timers: see {@link #roster()}'s javadoc for why the plain roster must stay non-resolving.
|
||||
*/
|
||||
public List<MemberSession> rosterResolved() {
|
||||
return registry.values().stream().map(this::resolveAgentSessionId).toList();
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve {@code session}'s {@code agentSessionId} if still unknown, re-polling the retained
|
||||
* {@link PeerHandle} for this pane (fleetd #209). A no-op — returning {@code session} unchanged
|
||||
* — once the id is already known, once no handle is retained for this pane (never spawned, or
|
||||
* already released), or if the handle throws while answering. A resolved id is best-effort
|
||||
* CAS-swapped into the registry via {@link #replace}; a lost race just means another caller
|
||||
* already applied the same update, so the resolved value is returned either way.
|
||||
*/
|
||||
private MemberSession resolveAgentSessionId(MemberSession session) {
|
||||
return resolveAgentSessionId(session, handles.get(session.paneId()));
|
||||
}
|
||||
|
||||
private MemberSession resolveAgentSessionId(MemberSession session, PeerHandle handle) {
|
||||
if (session.agentSessionId() != null || handle == null) {
|
||||
return session;
|
||||
}
|
||||
String resolved;
|
||||
try {
|
||||
resolved = handle.agentSessionId();
|
||||
} catch (RuntimeException e) {
|
||||
log.debug("agentSessionId lookup failed for pane={} terminal={}: {}",
|
||||
session.paneId(), session.terminalId(), e.toString());
|
||||
return session;
|
||||
}
|
||||
if (resolved == null) {
|
||||
return session;
|
||||
}
|
||||
MemberSession updated = session.withAgentSessionId(resolved);
|
||||
replace(session, updated); // best-effort; a lost CAS just means the resolved value stands anyway
|
||||
return updated;
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-304 merged roster+live view. The registry is authoritative for worktree, branch,
|
||||
* profile, owner, and state; the optional live agent supplies the herdr-reported status.
|
||||
|
||||
@@ -98,21 +98,4 @@ public interface Worktrees {
|
||||
/** CB-586: the operator-visible census of {@code refs/wip/*} in one repository. */
|
||||
record WipRefStats(int count, long costBytes) {
|
||||
}
|
||||
|
||||
/**
|
||||
* Make {@code repoRoot}'s git store and {@code worktreePath} writable by the configured group
|
||||
* (fleetd #185 stage 3), so a member spawned as a different OS user (see
|
||||
* {@code memberHerdrSocket}) can write its own worktree, its per-worktree git metadata, and
|
||||
* its own commit objects. No-op when no group is configured.
|
||||
*
|
||||
* <p><strong>This isolates credentials, not the repository.</strong> A member in the group can
|
||||
* still write the operator's git objects and refs in the shared repo — this only fixes file
|
||||
* ownership/permissions so a different-uid member can work at all, it grants no narrower access
|
||||
* than that.
|
||||
*
|
||||
* @param repoRoot the repository whose git store ({@code .git/objects}, {@code refs},
|
||||
* {@code logs}, {@code worktrees}, {@code packed-refs}) needs sharing
|
||||
* @param worktreePath the linked worktree's own directory
|
||||
*/
|
||||
void shareWithGroup(String repoRoot, String worktreePath);
|
||||
}
|
||||
|
||||
@@ -1036,35 +1036,6 @@ class FleetConfigTest {
|
||||
assertEquals(LeadMailbox.DEFAULT_PREFETCH, noEnv.prefetchOrDefault());
|
||||
}
|
||||
|
||||
@Test
|
||||
void absentWorktreeGroupLeavesItNull(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("no-worktree-group.yaml");
|
||||
Files.writeString(f, "bind:\n port: 8080\n");
|
||||
|
||||
FleetConfig cfg = FleetConfig.load(f);
|
||||
assertNull(cfg.worktreeGroup(), "no worktreeGroup: key → null → GitWorktrees.shareWithGroup is a no-op");
|
||||
}
|
||||
|
||||
@Test
|
||||
void worktreeGroupKeyParses(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("worktree-group.yaml");
|
||||
Files.writeString(f, "bind:\n port: 8080\nworktreeGroup: fleet-workers\n");
|
||||
|
||||
FleetConfig cfg = FleetConfig.load(f);
|
||||
assertEquals("fleet-workers", cfg.worktreeGroup());
|
||||
}
|
||||
|
||||
@Test
|
||||
void absentWorktreeGroupSurvivesTheBackCompatConstructorChain() {
|
||||
// fleetd #185 stage 3: withDefaults() (and every pre-existing call site) must not silently
|
||||
// drop a live worktreeGroup by routing through a back-compat constructor that defaults it
|
||||
// to null.
|
||||
FleetConfig cfg = new FleetConfig(null, null, null, Map.of(), null, null, null, null, null,
|
||||
null, null, null, null, null, null, null, null, null, null, null, "fleet-workers");
|
||||
assertEquals("fleet-workers", cfg.withDefaults().worktreeGroup(),
|
||||
"withDefaults() must carry a configured worktreeGroup through unchanged");
|
||||
}
|
||||
|
||||
@Test
|
||||
void absentPrimaryBlockLeavesPrimaryNull(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("no-primary.yaml");
|
||||
|
||||
@@ -1,27 +0,0 @@
|
||||
package dev.ltms.fleet.herdr;
|
||||
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
import java.util.OptionalLong;
|
||||
|
||||
/**
|
||||
* Fake {@link ParentResolver} backed by an explicit pid→parent map — lets {@link PaneLocatorTest}
|
||||
* drive {@link PaneLocator}'s ancestry walk (grandchild pids, cycles) without spawning real OS
|
||||
* processes.
|
||||
*/
|
||||
final class FakeParentResolver implements ParentResolver {
|
||||
|
||||
private final Map<Long, Long> parents = new HashMap<>();
|
||||
|
||||
/** {@code pid}'s parent is {@code parentPid}. A pid with no entry here has no known parent. */
|
||||
FakeParentResolver parent(long pid, long parentPid) {
|
||||
parents.put(pid, parentPid);
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public OptionalLong parentOf(long pid) {
|
||||
Long parent = parents.get(pid);
|
||||
return parent == null ? OptionalLong.empty() : OptionalLong.of(parent);
|
||||
}
|
||||
}
|
||||
@@ -1,11 +1,7 @@
|
||||
package dev.ltms.fleet.herdr;
|
||||
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
/** Unit tests for PID → pane resolution (the herdr half of connection-based MCP identity). */
|
||||
@@ -67,123 +63,4 @@ class PaneLocatorTest {
|
||||
long paneListCalls = shared.calls.stream().filter(c -> c.method().equals("pane.list")).count();
|
||||
assertEquals(1, paneListCalls, "same-object lead/member must scan exactly once, not twice");
|
||||
}
|
||||
|
||||
// --- ancestry walk (CB-161: grandchild pids matched no pane, resolving as primary) --------
|
||||
|
||||
@Test
|
||||
void stillResolvesAPidThatIsExactlyThePaneShellPid() {
|
||||
// Regression: a pid with no parent chain at all — no ancestry walk is needed to match it.
|
||||
OnePaneHerdr pane = new OnePaneHerdr("term_x", "pX", 5000, 6000);
|
||||
PaneLocator loc = new PaneLocator(pane, new FakeParentResolver());
|
||||
assertEquals("term_x", loc.terminalForPid(5000));
|
||||
}
|
||||
|
||||
@Test
|
||||
void stillResolvesAPidThatIsExactlyAForegroundPid() {
|
||||
// Regression: same as above, but matching via the foreground-processes list.
|
||||
OnePaneHerdr pane = new OnePaneHerdr("term_x", "pX", 5000, 6000);
|
||||
PaneLocator loc = new PaneLocator(pane, new FakeParentResolver());
|
||||
assertEquals("term_x", loc.terminalForPid(6000));
|
||||
}
|
||||
|
||||
@Test
|
||||
void resolvesAGrandchildPidTwoLevelsBelowTheShellPid() {
|
||||
// The bug: a helper process a worker spawns (python3, curl, ...) is a grandchild of the
|
||||
// pane's shell — not the shell_pid and not a foreground pid directly. Before the fix,
|
||||
// paneOwnsPid only checked direct pid equality, so this pid matched no pane and the
|
||||
// caller fell through to loopback-trust as the primary.
|
||||
OnePaneHerdr pane = new OnePaneHerdr("term_x", "pX", 5000, 6000);
|
||||
FakeParentResolver parents = new FakeParentResolver()
|
||||
.parent(7002, 7001) // grandchild -> child
|
||||
.parent(7001, 5000); // child -> shell (the pane's shell_pid)
|
||||
PaneLocator loc = new PaneLocator(pane, parents);
|
||||
assertEquals("term_x", loc.terminalForPid(7002));
|
||||
}
|
||||
|
||||
@Test
|
||||
void nullForAPidWhoseAncestryMatchesNoPane() {
|
||||
// Must not break the other direction: a pid that truly belongs to nothing here (e.g. the
|
||||
// real primary) must still resolve to null. Resolving everything to a worker would demote
|
||||
// the actual lead and refuse every orchestration call.
|
||||
OnePaneHerdr pane = new OnePaneHerdr("term_x", "pX", 5000, 6000);
|
||||
FakeParentResolver parents = new FakeParentResolver()
|
||||
.parent(9002, 9001)
|
||||
.parent(9001, 9000); // chain never reaches 5000 or 6000
|
||||
PaneLocator loc = new PaneLocator(pane, parents);
|
||||
assertNull(loc.terminalForPid(9002));
|
||||
}
|
||||
|
||||
@Test
|
||||
void ancestryWalkTerminatesOnACycleInsteadOfHanging() {
|
||||
// A fake (or corrupted) parent map that cycles must not hang identity resolution, which
|
||||
// runs on every MCP call. The walk must still terminate and correctly resolve to null.
|
||||
OnePaneHerdr pane = new OnePaneHerdr("term_x", "pX", 5000, 6000);
|
||||
FakeParentResolver parents = new FakeParentResolver()
|
||||
.parent(100, 101)
|
||||
.parent(101, 100); // cycle, never reaches the pane's pids
|
||||
PaneLocator loc = new PaneLocator(pane, parents);
|
||||
assertNull(loc.terminalForPid(100));
|
||||
}
|
||||
|
||||
@Test
|
||||
void ancestrySetIsComputedOnceAcrossBothClientsInTheTwoDaemonConstructor() {
|
||||
// CB-185: the two-daemon constructor searches lead then member. The ancestor set is
|
||||
// per-caller, not per-client — it must be walked once and reused, not recomputed for
|
||||
// each client searched.
|
||||
OnePaneHerdr pane = new OnePaneHerdr("term_x", "pX", 5000, 6000);
|
||||
FakeParentResolver parents = new FakeParentResolver()
|
||||
.parent(7002, 7001)
|
||||
.parent(7001, 5000);
|
||||
AtomicInteger calls = new AtomicInteger();
|
||||
ParentResolver counting = pid -> {
|
||||
calls.incrementAndGet();
|
||||
return parents.parentOf(pid);
|
||||
};
|
||||
HerdrClient noPanes = new FakeHerdr().withNoPanes();
|
||||
PaneLocator two = new PaneLocator(noPanes, pane, counting);
|
||||
assertEquals("term_x", two.terminalForPid(7002));
|
||||
assertEquals(3, calls.get(), "ancestry must be walked once (3 lookups: 7002, 7001, 5000), "
|
||||
+ "not re-walked per herdr client");
|
||||
}
|
||||
|
||||
/** Minimal single-pane {@link HerdrClient} fake, purpose-built for the ancestry tests above. */
|
||||
private static final class OnePaneHerdr implements HerdrClient {
|
||||
private final ObjectMapper mapper = new ObjectMapper();
|
||||
private final String terminalId;
|
||||
private final String paneId;
|
||||
private final long shellPid;
|
||||
private final long foregroundPid;
|
||||
|
||||
OnePaneHerdr(String terminalId, String paneId, long shellPid, long foregroundPid) {
|
||||
this.terminalId = terminalId;
|
||||
this.paneId = paneId;
|
||||
this.shellPid = shellPid;
|
||||
this.foregroundPid = foregroundPid;
|
||||
}
|
||||
|
||||
@Override
|
||||
public JsonNode call(String method, Object params) {
|
||||
try {
|
||||
return switch (method) {
|
||||
case "pane.list" -> mapper.readTree(("""
|
||||
{"type":"pane_list","panes":[
|
||||
{"pane_id":"%s","terminal_id":"%s","workspace_id":"w1","tab_id":"w1:t1","agent":"claude"}]}""")
|
||||
.formatted(paneId, terminalId));
|
||||
case "pane.process_info" -> mapper.readTree(("""
|
||||
{"type":"pane_process_info","process_info":{"pane_id":"%s","shell_pid":%d,
|
||||
"foreground_processes":[{"pid":%d,"name":"node","argv0":"claude"}]}}""")
|
||||
.formatted(paneId, shellPid, foregroundPid));
|
||||
default -> throw new HerdrException("OnePaneHerdr has no canned response for " + method);
|
||||
};
|
||||
} catch (HerdrException e) {
|
||||
throw e;
|
||||
} catch (Exception e) {
|
||||
throw new HerdrException("OnePaneHerdr decode failed for " + method, e);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+3
-430
@@ -12,26 +12,20 @@ import dev.ltms.fleet.peer.Capability;
|
||||
import dev.ltms.fleet.peer.MemberRole;
|
||||
import dev.ltms.fleet.peer.SpawnRequest;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
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;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.Supplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
import static org.junit.jupiter.api.Assumptions.assumeTrue;
|
||||
|
||||
/**
|
||||
* CB-633: proves the allow-list scrub is actually WIRED INTO the spawn path — not merely that its
|
||||
@@ -125,18 +119,6 @@ class HerdrPeerLauncherAllowListWiringTest {
|
||||
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST, allow, List.of(), null);
|
||||
}
|
||||
|
||||
/**
|
||||
* Same as {@link #allowList()} but with an operator-configured {@code known:} list, so the
|
||||
* fallback overlay (the non-zsh / no-memberLoginShell path) has something visible to shadow —
|
||||
* {@link FleetConfig.MemberCredentials#blockedSet()} is {@code known - allow}, so an empty
|
||||
* {@code known} (what {@link #allowList()} uses) blocks nothing and a fallback test would have
|
||||
* no sentinel entry to assert on.
|
||||
*/
|
||||
private static Supplier<FleetConfig.MemberCredentials> allowListWithKnown(List<String> known) {
|
||||
return () -> new FleetConfig.MemberCredentials(
|
||||
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST, List.of(), known, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-633 follow-up: a name that lives ONLY in {@code memberCredentials.allow:} — no profile
|
||||
* mentions it — must survive the scrub the real spawn path generates. Calling {@code
|
||||
@@ -179,34 +161,6 @@ class HerdrPeerLauncherAllowListWiringTest {
|
||||
+ "it under allow: — sshAuthSock is unset here, so it defaults to block");
|
||||
}
|
||||
|
||||
@Test
|
||||
void brokerUriEnvStaysBlockedWhenListedInMemberCredentialsAllow() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
WiringLauncher launcher = new WiringLauncher(herdr,
|
||||
allowListWithAllow(List.of("BROKER_CONNECTION_URI")), "/bin/zsh", null,
|
||||
() -> config("BROKER_CONNECTION_URI"));
|
||||
|
||||
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
|
||||
|
||||
Path dir = Path.of(launcher.env.get("ZDOTDIR"));
|
||||
assertFalse(readAll(dir.resolve(EnvAllowListScrub.SCRUB_FILE)).contains("'BROKER_CONNECTION_URI'"),
|
||||
"broker.uriEnv must not reach a member even when listed in memberCredentials.allow:");
|
||||
}
|
||||
|
||||
@Test
|
||||
void brokerUriEnvIsDeniedUnderTheDenyListPolicyEvenWhenAllowed() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
WiringLauncher launcher = new WiringLauncher(herdr,
|
||||
() -> new FleetConfig.MemberCredentials(null, List.of("BROKER_CONNECTION_URI"), List.of(), null),
|
||||
"/bin/bash", null,
|
||||
() -> config("BROKER_CONNECTION_URI"));
|
||||
|
||||
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
|
||||
|
||||
assertEquals("blocked-by-fleetd-cb596-see-gitea-issue-82", launcher.env.get("BROKER_CONNECTION_URI"),
|
||||
"the deny-list overlay must deny broker.uriEnv even when allow: names it");
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-633 follow-up criterion 3: on every allow-list spawn the daemon logs one INFO line, shaped
|
||||
* "member credentials: allowed N of M", with real counts — not constants. Real path: the count
|
||||
@@ -275,268 +229,6 @@ class HerdrPeerLauncherAllowListWiringTest {
|
||||
+ appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #185 stage 2 pin: with {@code memberHerdrSocket:} absent (today's only mode, and the
|
||||
* default — this host runs no other), the gap detector's WARN/INFO conclusions read exactly as
|
||||
* they did before this fix. Real path: {@code FLEETD_WORKER_TOKEN} (the test profile's own
|
||||
* {@code tokenEnv}) is a name the derived allow-list keeps, so it gets the "UNBLOCKED" WARN;
|
||||
* {@code SOME_UNKNOWN_SECRET_TOKEN} is not derived from anywhere, so it gets the "scrub blanks
|
||||
* them" INFO. This is the exact shape #185 stage 2 must not touch on this path.
|
||||
*/
|
||||
@Test
|
||||
void gapConclusionsAreByteIdenticalWhenMemberHerdrSocketIsAbsent() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Set<String> hostEnvNames = Set.of("FLEETD_WORKER_TOKEN", "SOME_UNKNOWN_SECRET_TOKEN");
|
||||
WiringLauncher launcher = new WiringLauncher(herdr, allowList(), "/bin/zsh", () -> hostEnvNames,
|
||||
() -> config(null));
|
||||
|
||||
List<String> messages = spawnAndCaptureLogs(launcher);
|
||||
|
||||
assertTrue(messages.contains("memberCredentials gap: 1 credential-shaped env var name(s) are on "
|
||||
+ "neither known: nor allow: — the derived allow-list keeps them anyway (a "
|
||||
+ "profile's gitTokenEnv/gitHostEnv/tokenEnv/env: names one, or this spawn "
|
||||
+ "injects it), so every member pane inherits them UNBLOCKED — [FLEETD_WORKER_TOKEN]. "
|
||||
+ "Add each to memberCredentials.known (or .allow if a member legitimately needs "
|
||||
+ "it), or remove it from whatever profile setting derives it in."),
|
||||
"expected the pre-existing UNBLOCKED WARN unchanged, got: " + messages);
|
||||
assertTrue(messages.contains("memberCredentials gap: 1 credential-shaped env var name(s) are on "
|
||||
+ "neither known: nor allow: — [SOME_UNKNOWN_SECRET_TOKEN]. The allow-list scrub "
|
||||
+ "blanks them anyway (they are not on the derived allow-list), so no member pane "
|
||||
+ "keeps them; add each to memberCredentials.known or .allow to make that explicit."),
|
||||
"expected the pre-existing 'scrub blanks them' INFO unchanged, got: " + messages);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #185 stage 2: with {@code memberHerdrSocket:} configured, member panes run under a
|
||||
* different OS user — {@link HerdrPeerLauncher#hostEnvNames} describes fleetd's own process, not
|
||||
* that user's. Neither "inherits them UNBLOCKED" nor "scrub blanks them" is evidence-backed
|
||||
* there, so neither may print; the single unknown-environment WARN must, naming the config key.
|
||||
*/
|
||||
@Test
|
||||
void gapDetectorReportsUnknownInsteadOfAConclusionWhenMemberHerdrSocketIsConfigured() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Set<String> hostEnvNames = Set.of("FLEETD_WORKER_TOKEN", "SOME_UNKNOWN_SECRET_TOKEN");
|
||||
WiringLauncher launcher = new WiringLauncher(herdr, allowList(), "/bin/zsh", () -> hostEnvNames,
|
||||
() -> configWithMemberHerdrSocket("/tmp/other-user-herdr.sock"));
|
||||
|
||||
List<String> messages = spawnAndCaptureLogs(launcher);
|
||||
|
||||
assertTrue(messages.stream().anyMatch(m -> m.contains("memberHerdrSocket")
|
||||
&& m.contains("UNKNOWN") && m.contains("cannot be verified")),
|
||||
"expected the unknown-member-environment WARN naming memberHerdrSocket, got: " + messages);
|
||||
assertFalse(messages.stream().anyMatch(m -> m.contains("UNBLOCKED")),
|
||||
"the 'inherits them UNBLOCKED' conclusion must not print once the evidence is about "
|
||||
+ "the wrong (daemon's own) environment — got: " + messages);
|
||||
assertFalse(messages.stream().anyMatch(m -> m.contains("scrub blanks them")),
|
||||
"the 'scrub blanks them' conclusion must not print once the evidence is about the "
|
||||
+ "wrong (daemon's own) environment — got: " + messages);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #185 stage 2: the unknown-environment WARN is a standing fact about this launcher's
|
||||
* configuration, not per-spawn news — it must fire once per launcher instance, the same shape as
|
||||
* every other one-time WARN in this class (e.g. {@code warnNonZsh}).
|
||||
*/
|
||||
@Test
|
||||
void theUnknownEnvironmentWarnFiresOnceNotOncePerSpawn() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Set<String> hostEnvNames = Set.of("FLEETD_WORKER_TOKEN", "SOME_UNKNOWN_SECRET_TOKEN");
|
||||
WiringLauncher launcher = new WiringLauncher(herdr, allowList(), "/bin/zsh", () -> hostEnvNames,
|
||||
() -> configWithMemberHerdrSocket("/tmp/other-user-herdr.sock"));
|
||||
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
|
||||
Level original = logger.getLevel();
|
||||
logger.setLevel(Level.WARN);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
|
||||
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
logger.setLevel(original);
|
||||
}
|
||||
|
||||
long count = appender.list.stream()
|
||||
.filter(e -> e.getFormattedMessage().contains("cannot be verified from here"))
|
||||
.count();
|
||||
assertEquals(1, count, "the unknown-member-environment WARN must fire once per launcher "
|
||||
+ "instance, not once per spawn — got " + count + " occurrence(s) among: "
|
||||
+ appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
|
||||
}
|
||||
|
||||
/**
|
||||
* Hard constraint: the gap detector must never log an env var VALUE, only its NAME. {@code
|
||||
* SOME_UNKNOWN_SECRET_TOKEN} resolves to a distinctive canary value through the same {@code env}
|
||||
* lookup the launcher uses elsewhere (SHELL, PATH, token resolution) — proving the value IS
|
||||
* resolvable does not mean the detector reads it, since {@link HerdrPeerLauncher#hostEnvNames}
|
||||
* (names only) is its data source, never {@code env.apply(name)} for those names.
|
||||
*/
|
||||
@Test
|
||||
void theGapDetectorNeverLogsAnEnvVarValueOnlyItsName() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
String canary = "sekrit-value-CANARY-9f3a1b7c";
|
||||
Set<String> hostEnvNames = Set.of("FLEETD_WORKER_TOKEN", "SOME_UNKNOWN_SECRET_TOKEN");
|
||||
WiringLauncher launcher = new WiringLauncher(herdr, allowList(), "/bin/zsh", () -> hostEnvNames,
|
||||
() -> config(null), Map.of("SOME_UNKNOWN_SECRET_TOKEN", canary));
|
||||
|
||||
List<String> messages = spawnAndCaptureLogs(launcher);
|
||||
|
||||
assertTrue(messages.stream().anyMatch(m -> m.contains("SOME_UNKNOWN_SECRET_TOKEN")),
|
||||
"expected the credential-shaped NAME to appear in the log, got: " + messages);
|
||||
assertFalse(messages.stream().anyMatch(m -> m.contains(canary)),
|
||||
"the log must never contain an env var VALUE, only its NAME — got: " + messages);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #213 defect 1, acceptance criterion 1: {@code memberHerdrSocket} configured and {@code
|
||||
* memberLoginShell} configured as non-zsh must fall back to the sentinel overlay exactly like a
|
||||
* non-zsh {@code $SHELL} does today — and the "generated ZDOTDIR" INFO must not appear, since no
|
||||
* scrub actually runs. The WiringLauncher's own {@code env("SHELL")} is deliberately set to
|
||||
* {@code /bin/zsh} — the OPPOSITE of what {@code memberLoginShell} says — so a launcher that
|
||||
* (incorrectly) fell back to fleetd's own {@code $SHELL} here would wrongly pass the gate and
|
||||
* fail this test.
|
||||
*/
|
||||
@Test
|
||||
void memberHerdrSocketWithNonZshMemberLoginShellFallsBackToTheOverlay() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
WiringLauncher launcher = new WiringLauncher(herdr, allowListWithKnown(List.of("SOME_TOKEN")),
|
||||
"/bin/zsh", null,
|
||||
() -> configWithMemberHerdrSocket("/tmp/other-user-herdr.sock", "/bin/bash"));
|
||||
|
||||
List<String> messages = spawnAndCaptureLogs(launcher);
|
||||
|
||||
assertEquals("blocked-by-fleetd-cb596-see-gitea-issue-82", launcher.env.get("SOME_TOKEN"),
|
||||
"a non-zsh memberLoginShell must fall back to the CB-596 sentinel overlay, exactly "
|
||||
+ "like a non-zsh $SHELL does when memberHerdrSocket is absent");
|
||||
assertFalse(launcher.env.containsKey("ZDOTDIR"),
|
||||
"no scrub directory may be generated when the configured member login shell is not zsh");
|
||||
assertFalse(messages.stream().anyMatch(m -> m.contains("generated ZDOTDIR")),
|
||||
"the 'generated ZDOTDIR' INFO must not appear when the scrub never runs — got: " + messages);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #213 defect 1, acceptance criterion 2: {@code memberHerdrSocket} configured and NO
|
||||
* {@code memberLoginShell} configured must fall back exactly like criterion 1 above — AND
|
||||
* fleetd's own {@code $SHELL} must never even be consulted (not merely "not decisive"). The
|
||||
* fixture's {@code env} function reports {@code /bin/zsh} for {@code SHELL} — a value that would
|
||||
* WRONGLY pass the zsh gate if the fix regressed to reading it — while flagging whether it was
|
||||
* ever asked for at all, so this test fails loudly on either kind of regression.
|
||||
*/
|
||||
@Test
|
||||
void memberHerdrSocketWithNoMemberLoginShellFallsBackAndNeverConsultsFleetdsOwnShell() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
AtomicBoolean shellQueried = new AtomicBoolean(false);
|
||||
Function<String, String> env = name -> {
|
||||
if ("SHELL".equals(name)) {
|
||||
shellQueried.set(true);
|
||||
return "/bin/zsh"; // would wrongly pass the zsh gate if this ever leaked through
|
||||
}
|
||||
return null;
|
||||
};
|
||||
WiringLauncher launcher = new WiringLauncher(herdr, allowListWithKnown(List.of("SOME_TOKEN")), env,
|
||||
() -> configWithMemberHerdrSocket("/tmp/other-user-herdr.sock", null));
|
||||
|
||||
List<String> messages = spawnAndCaptureLogs(launcher);
|
||||
|
||||
assertFalse(shellQueried.get(), "fleetd's own $SHELL must never be consulted once "
|
||||
+ "memberHerdrSocket is configured — only memberLoginShell: may decide the gate");
|
||||
assertEquals("blocked-by-fleetd-cb596-see-gitea-issue-82", launcher.env.get("SOME_TOKEN"),
|
||||
"no memberLoginShell configured must fall back to the sentinel overlay, same as a "
|
||||
+ "configured non-zsh shell");
|
||||
assertFalse(messages.stream().anyMatch(m -> m.contains("generated ZDOTDIR")),
|
||||
"no scrub may run without a configured memberLoginShell — got: " + messages);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #213 defect 2, acceptance criterion 3: with {@code memberHerdrSocket} configured, a
|
||||
* zsh {@code memberLoginShell}, and {@code worktreeRoot}/{@code worktreeGroup} both configured,
|
||||
* the generated scrub directory must live under {@code worktreeRoot} — NEVER under {@code
|
||||
* java.io.tmpdir}, which is fleetd's own 0700 temp dir and unreadable by the member's different
|
||||
* OS user. {@code worktreeGroup} is set to the CURRENT process's own primary group so {@code
|
||||
* EnvAllowListScrub}'s group-sharing step resolves on whatever host runs this test, rather than
|
||||
* hardcoding a group name that may not exist here.
|
||||
*/
|
||||
@Test
|
||||
void memberHerdrSocketWithZshMemberLoginShellPutsTheScrubOutsideJavaIoTmpdir(@TempDir Path worktreeRoot)
|
||||
throws IOException {
|
||||
String group = currentUserGroup();
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
WiringLauncher launcher = new WiringLauncher(herdr, allowList(), "/bin/bash-should-be-ignored", null,
|
||||
() -> configWithMemberHerdrSocketRootAndGroup("/tmp/other-user-herdr.sock", "/bin/zsh",
|
||||
worktreeRoot.toString(), group));
|
||||
|
||||
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
|
||||
|
||||
String zdotdir = launcher.env.get("ZDOTDIR");
|
||||
assertNotNull(zdotdir, "a zsh memberLoginShell with worktreeRoot+worktreeGroup configured "
|
||||
+ "must still generate a ZDOTDIR");
|
||||
Path dir = Path.of(zdotdir);
|
||||
// The immediate PARENT is asserted (not merely startsWith(java.io.tmpdir)), because a
|
||||
// JUnit @TempDir is itself carved out of the JVM's java.io.tmpdir — startsWith alone would
|
||||
// pass by coincidence of the test fixture, not because the launcher used worktreeRoot.
|
||||
assertEquals(worktreeRoot.toAbsolutePath().normalize(), dir.getParent(),
|
||||
"the generated scrub directory's parent must be the configured worktreeRoot, not "
|
||||
+ "System.getProperty(\"java.io.tmpdir\") — got parent " + dir.getParent());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #213, acceptance criterion 4: with {@code memberHerdrSocket} absent (today's only
|
||||
* mode), behaviour must be byte-identical to before this fix — fleetd's own {@code $SHELL}
|
||||
* still decides the gate (proven here, not merely assumed, by flagging the lookup), and the
|
||||
* scrub still lands under {@code java.io.tmpdir}.
|
||||
*/
|
||||
@Test
|
||||
void memberHerdrSocketAbsentStillConsultsFleetdsOwnShellAndBehavesAsBefore() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
AtomicBoolean shellQueried = new AtomicBoolean(false);
|
||||
Function<String, String> env = name -> {
|
||||
if ("SHELL".equals(name)) {
|
||||
shellQueried.set(true);
|
||||
return "/bin/zsh";
|
||||
}
|
||||
return null;
|
||||
};
|
||||
WiringLauncher launcher = new WiringLauncher(herdr, allowList(), env, null); // memberHerdrSocket absent
|
||||
|
||||
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
|
||||
|
||||
assertTrue(shellQueried.get(), "with memberHerdrSocket absent, fleetd's own $SHELL must "
|
||||
+ "still decide the zsh gate, unchanged from before this fix");
|
||||
String zdotdir = launcher.env.get("ZDOTDIR");
|
||||
assertNotNull(zdotdir, "SHELL=/bin/zsh with memberHerdrSocket absent must still generate a "
|
||||
+ "ZDOTDIR, as before this fix");
|
||||
Path dir = Path.of(zdotdir);
|
||||
assertTrue(dir.startsWith(Path.of(System.getProperty("java.io.tmpdir"))),
|
||||
"with memberHerdrSocket absent the scrub directory must still be generated under "
|
||||
+ "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();
|
||||
}
|
||||
|
||||
/** Spawn once through the real launcher path, capturing every INFO+ line this class logs. */
|
||||
private static List<String> spawnAndCaptureLogs(HerdrPeerLauncher launcher) {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
|
||||
Level original = logger.getLevel();
|
||||
logger.setLevel(Level.INFO);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
logger.setLevel(original);
|
||||
}
|
||||
return appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList();
|
||||
}
|
||||
|
||||
private static String readAll(Path p) {
|
||||
try {
|
||||
return Files.readString(p);
|
||||
@@ -568,52 +260,15 @@ class HerdrPeerLauncherAllowListWiringTest {
|
||||
|
||||
/** Plus an injectable {@code hostEnvNames} source, for the "allowed N of M" log line test. */
|
||||
WiringLauncher(FakeHerdr herdr, Supplier<FleetConfig.MemberCredentials> creds, String shell,
|
||||
Supplier<Set<String>> hostEnvNames) {
|
||||
this(herdr, creds, shell, hostEnvNames, null);
|
||||
}
|
||||
|
||||
WiringLauncher(FakeHerdr herdr, Supplier<FleetConfig.MemberCredentials> creds, String shell,
|
||||
Supplier<Set<String>> hostEnvNames, Supplier<FleetConfig> config) {
|
||||
this(herdr, creds, shell, hostEnvNames, config, Map.of());
|
||||
}
|
||||
|
||||
/**
|
||||
* Plus a host-env value map (name → value), resolved through the same {@code env} lookup
|
||||
* every adapter uses for {@code SHELL}/{@code PATH}/token resolution — fleetd #185 stage 2's
|
||||
* "the gap detector never logs a value" tests use this to prove a value that IS resolvable
|
||||
* for a credential-shaped name never reaches the log, since the detector only ever reads
|
||||
* {@code hostEnvNames} (names), never {@code env.apply(name)} (values), for those names.
|
||||
*/
|
||||
WiringLauncher(FakeHerdr herdr, Supplier<FleetConfig.MemberCredentials> creds, String shell,
|
||||
Supplier<Set<String>> hostEnvNames, Supplier<FleetConfig> config,
|
||||
Map<String, String> extraEnvValues) {
|
||||
Supplier<Set<String>> hostEnvNames) {
|
||||
super("test", new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
Map.of("test", profile()), "test",
|
||||
name -> "SHELL".equals(name) ? shell
|
||||
: (extraEnvValues != null && extraEnvValues.containsKey(name))
|
||||
? extraEnvValues.get(name) : null,
|
||||
0, () -> 0L, () -> { }, null, creds, hostEnvNames, config);
|
||||
}
|
||||
|
||||
/**
|
||||
* Full control over the {@code env} lookup, bypassing the {@code shell}/{@code
|
||||
* extraEnvValues} convenience above entirely — fleetd #213's "fleetd's own $SHELL must
|
||||
* never be consulted" tests need to OBSERVE whether {@code SHELL} was ever looked up, which
|
||||
* a plain value substitution cannot do.
|
||||
*/
|
||||
WiringLauncher(FakeHerdr herdr, Supplier<FleetConfig.MemberCredentials> creds,
|
||||
Function<String, String> env, Supplier<FleetConfig> config) {
|
||||
super("test", new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
Map.of("test", profile()), "test",
|
||||
env, 0, () -> 0L, () -> { }, null, creds, null, config);
|
||||
name -> "SHELL".equals(name) ? shell : null,
|
||||
0, () -> 0L, () -> { }, null, creds, hostEnvNames);
|
||||
}
|
||||
|
||||
@Override
|
||||
protected Launch buildLaunch(FleetConfig.Profile cfg, LaunchSpec spec) {
|
||||
Map<String, String> launchEnv = baseEnv(cfg);
|
||||
launchEnv.putAll(env);
|
||||
env.clear();
|
||||
env.putAll(launchEnv);
|
||||
return new Launch(env, List.of("test"));
|
||||
}
|
||||
|
||||
@@ -623,88 +278,6 @@ class HerdrPeerLauncherAllowListWiringTest {
|
||||
}
|
||||
}
|
||||
|
||||
private static FleetConfig config(String brokerUriEnv) {
|
||||
return new FleetConfig(null, null, null, Map.of(), null, null, null, null, null,
|
||||
new FleetConfig.Broker(null, brokerUriEnv, null), null, null, null, null, null,
|
||||
null, null, null, null, null).withDefaults();
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #185 stage 2: a config with {@code memberHerdrSocket:} set — member panes run on a
|
||||
* second herdr owned by a different OS user, so {@link HerdrPeerLauncher#hostEnvNames} no
|
||||
* longer describes what a member pane inherits.
|
||||
*/
|
||||
private static FleetConfig configWithMemberHerdrSocket(String memberHerdrSocket) {
|
||||
return new FleetConfig(null, null, memberHerdrSocket, Map.of(), null, null, null, null, null,
|
||||
null, null, null, null, null, null, null, null, null, null, null).withDefaults();
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #213: as {@link #configWithMemberHerdrSocket(String)}, plus the {@code
|
||||
* memberLoginShell:} the member's OS user actually runs — the config key {@link
|
||||
* HerdrPeerLauncher#applyEnvironmentAllowListPolicy} must consult instead of fleetd's own
|
||||
* {@code $SHELL} once {@code memberHerdrSocket} is configured.
|
||||
*/
|
||||
private static FleetConfig configWithMemberHerdrSocket(String memberHerdrSocket, String memberLoginShell) {
|
||||
return new FleetConfig(
|
||||
null, // bind
|
||||
null, // herdrSocket
|
||||
memberHerdrSocket, // memberHerdrSocket
|
||||
Map.of(), // profiles
|
||||
null, // guard
|
||||
null, // worktreeRoot
|
||||
null, // lifecycle
|
||||
null, // spawnReadyTimeoutMs
|
||||
null, // spawnReadyPollMs
|
||||
null, // broker
|
||||
null, // primary
|
||||
null, // fleet
|
||||
null, // leadHeartbeat
|
||||
null, // health
|
||||
null, // placement
|
||||
null, // auth
|
||||
null, // configReload
|
||||
null, // quarantineCooldownSeconds
|
||||
null, // memberCredentials
|
||||
null, // coordinator
|
||||
null, // worktreeGroup
|
||||
memberLoginShell // memberLoginShell
|
||||
).withDefaults();
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #213: as above, plus {@code worktreeRoot:}/{@code worktreeGroup:} — both required for
|
||||
* the ZDOTDIR scrub to run at all once {@code memberHerdrSocket} is configured; either missing
|
||||
* falls back to the sentinel overlay, same as a non-zsh {@code memberLoginShell}.
|
||||
*/
|
||||
private static FleetConfig configWithMemberHerdrSocketRootAndGroup(String memberHerdrSocket,
|
||||
String memberLoginShell, String worktreeRoot, String worktreeGroup) {
|
||||
return new FleetConfig(
|
||||
null, // bind
|
||||
null, // herdrSocket
|
||||
memberHerdrSocket, // memberHerdrSocket
|
||||
Map.of(), // profiles
|
||||
null, // guard
|
||||
worktreeRoot, // worktreeRoot
|
||||
null, // lifecycle
|
||||
null, // spawnReadyTimeoutMs
|
||||
null, // spawnReadyPollMs
|
||||
null, // broker
|
||||
null, // primary
|
||||
null, // fleet
|
||||
null, // leadHeartbeat
|
||||
null, // health
|
||||
null, // placement
|
||||
null, // auth
|
||||
null, // configReload
|
||||
null, // quarantineCooldownSeconds
|
||||
null, // memberCredentials
|
||||
null, // coordinator
|
||||
worktreeGroup, // worktreeGroup
|
||||
memberLoginShell // memberLoginShell
|
||||
).withDefaults();
|
||||
}
|
||||
|
||||
/** The generated directory is a temp directory; make sure the test does not leave a pile. */
|
||||
@Test
|
||||
void theGeneratedDirectoryIsRemovedWhenThePaneIsStopped() {
|
||||
|
||||
@@ -121,34 +121,6 @@ class MemberEnvAllowListTest {
|
||||
assertTrue(derived.contains("OTHER_NAME"), "other allow: names are unaffected");
|
||||
}
|
||||
|
||||
@Test
|
||||
void configuredBrokerAndCoordinatorUriEnvNamesAreExcludedEvenWhenAllowed() {
|
||||
FleetConfig config = config("BROKER_CONNECTION_URI", "COORDINATOR_CONNECTION_URI");
|
||||
Set<String> excluded = MemberEnvAllowList.brokerUriEnvNames(config);
|
||||
|
||||
Set<String> derived = MemberEnvAllowList.derive(List.of(),
|
||||
Set.of("BROKER_CONNECTION_URI", "COORDINATOR_CONNECTION_URI", "OTHER_NAME"), excluded);
|
||||
|
||||
assertFalse(derived.contains("BROKER_CONNECTION_URI"),
|
||||
"broker.uriEnv is secret-bearing and must not ride in on allow:");
|
||||
assertFalse(derived.contains("COORDINATOR_CONNECTION_URI"),
|
||||
"coordinator.uriEnv has the same inline-password shape");
|
||||
assertTrue(derived.contains("OTHER_NAME"), "unrelated allow: entries are unaffected");
|
||||
}
|
||||
|
||||
@Test
|
||||
void absentOrBlankBrokerUriEnvAddsNoExclusions() {
|
||||
assertTrue(MemberEnvAllowList.brokerUriEnvNames(config(null, null)).isEmpty());
|
||||
assertTrue(MemberEnvAllowList.brokerUriEnvNames(config(" ", "")).isEmpty());
|
||||
}
|
||||
|
||||
private static FleetConfig config(String brokerUriEnv, String coordinatorUriEnv) {
|
||||
return new FleetConfig(null, null, null, Map.of(), null, null, null, null, null,
|
||||
new FleetConfig.Broker(null, brokerUriEnv, null), null, null, null, null, null,
|
||||
null, null, null, null,
|
||||
new FleetConfig.Coordinator(null, coordinatorUriEnv, null, null)).withDefaults();
|
||||
}
|
||||
|
||||
/** {@code LC_*} categories are infrastructure by prefix; everything else needs an exact match. */
|
||||
@Test
|
||||
void keepsMatchesExactlyPlusTheLocalePrefixRule() {
|
||||
|
||||
@@ -12,6 +12,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 org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
@@ -22,6 +23,7 @@ import java.util.Map;
|
||||
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.*;
|
||||
@@ -61,6 +63,23 @@ class OpenCodeLauncherTest {
|
||||
0, System::currentTimeMillis, () -> { }, configRoot, configRoot, null, () -> creds);
|
||||
}
|
||||
|
||||
private static OpenCodeLauncher serviceWithModelMismatchSink(FakeHerdr herdr, Path root,
|
||||
FleetConfig.Profile cfg,
|
||||
BackendQuarantine quarantine) {
|
||||
return new OpenCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null,
|
||||
0, System::currentTimeMillis, () -> { }, root, root, null, null,
|
||||
(profile, _) -> quarantine.quarantinePermanently(profile));
|
||||
}
|
||||
|
||||
private static void writeResolvedModelRecord(Path root, String directory, String providerId,
|
||||
String modelId) throws Exception {
|
||||
Path record = Files.createDirectories(root.resolve("session").resolve("p1")).resolve("ses_a.json");
|
||||
Files.writeString(record, "{\"id\":\"ses_a\",\"directory\":\"" + directory
|
||||
+ "\",\"model\":{\"id\":\"" + modelId + "\",\"providerID\":\""
|
||||
+ providerId + "\",\"variant\":\"high\"}}");
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
private static Map<String, Object> lastStart(FakeHerdr herdr) {
|
||||
return (Map<String, Object>) herdr.lastCall("agent.start").params();
|
||||
@@ -211,6 +230,50 @@ class OpenCodeLauncherTest {
|
||||
"--auto is unconditional: a model-less worker still must never block on approval");
|
||||
}
|
||||
|
||||
// --- CB-175: verify opencode's recorded resolved model after spawn readiness ----------------
|
||||
|
||||
@Test
|
||||
void matchingResolvedModelDoesNotQuarantineTheProfile(@TempDir Path root) throws Exception {
|
||||
FleetConfig.Profile cfg = opencodeCfg("opencode/x-preview-f-free", null, null);
|
||||
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(1));
|
||||
writeResolvedModelRecord(root, root.toString(), "opencode", "x-preview-f-free");
|
||||
|
||||
serviceWithModelMismatchSink(new FakeHerdr(), root, cfg, quarantine)
|
||||
.spawn(new SpawnRequest(null, root.toString(), null, null, null, MemberRole.DEV));
|
||||
|
||||
assertFalse(quarantine.isQuarantined("gemini"), "the exact provider/model match is safe");
|
||||
}
|
||||
|
||||
@Test
|
||||
void differentResolvedModelPermanentlyQuarantinesTheProfile(@TempDir Path root) throws Exception {
|
||||
FleetConfig.Profile cfg = opencodeCfg("opencode/x-preview-f-free", null, null);
|
||||
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(1));
|
||||
writeResolvedModelRecord(root, root.toString(), "openai", "gpt-5.6-sol");
|
||||
|
||||
serviceWithModelMismatchSink(new FakeHerdr(), root, cfg, quarantine)
|
||||
.spawn(new SpawnRequest(null, root.toString(), null, null, null, MemberRole.DEV));
|
||||
|
||||
assertTrue(quarantine.isQuarantined("gemini"), "a fallback model must block later spawns");
|
||||
assertEquals(Long.MAX_VALUE / 1_000_000_000L + 1, quarantine.remainingSeconds("gemini").orElseThrow(),
|
||||
"a withdrawn selector cannot become safe after the normal cooldown");
|
||||
}
|
||||
|
||||
@Test
|
||||
void missingOrUnreadableSessionDatabaseDoesNotQuarantine(@TempDir Path root) throws Exception {
|
||||
FleetConfig.Profile cfg = opencodeCfg("opencode/x-preview-f-free", null, null);
|
||||
BackendQuarantine missing = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(1));
|
||||
serviceWithModelMismatchSink(new FakeHerdr(), root, cfg, missing)
|
||||
.spawn(new SpawnRequest(null, root.toString(), null, null, null, MemberRole.DEV));
|
||||
assertFalse(missing.isQuarantined("gemini"), "a missing database is unknown evidence");
|
||||
|
||||
Path broken = Files.createDirectories(root.resolve("session").resolve("p1")).resolve("ses_a.json");
|
||||
Files.writeString(broken, "not JSON");
|
||||
BackendQuarantine unreadable = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(1));
|
||||
serviceWithModelMismatchSink(new FakeHerdr(), root, cfg, unreadable)
|
||||
.spawn(new SpawnRequest(null, root.toString(), null, null, null, MemberRole.DEV));
|
||||
assertFalse(unreadable.isQuarantined("gemini"), "an unreadable database is unknown evidence");
|
||||
}
|
||||
|
||||
// --- CB-617: --agent <role> when the role has an agent-definition file --------------------
|
||||
|
||||
@Test
|
||||
@@ -334,7 +397,8 @@ class OpenCodeLauncherTest {
|
||||
assertNull(handle.agentSessionId(), "no record yet → null, not a spawn-time block");
|
||||
// Once the record appears (here: same cwd), lazy discovery resolves it — the handle's
|
||||
// session id matches its own worktree, not another's.
|
||||
OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "ses_resolved", "/work/dir", 1000L);
|
||||
OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "p1", "ses_a.json",
|
||||
"ses_resolved", "/work/dir", 1000L);
|
||||
assertEquals("ses_resolved", handle.agentSessionId(),
|
||||
"agentSessionId() re-scans and picks up a record that has since been written");
|
||||
}
|
||||
|
||||
@@ -5,89 +5,68 @@ import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.sql.Connection;
|
||||
import java.sql.DriverManager;
|
||||
import java.sql.PreparedStatement;
|
||||
import java.sql.SQLException;
|
||||
import java.sql.Statement;
|
||||
import java.nio.file.attribute.FileTime;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
/**
|
||||
* {@link OpenCodeSessionDiscovery} matches an opencode session row by the worker's cwd (its
|
||||
* {@code directory}) against opencode's {@code opencode.db} SQLite database. These tests build a
|
||||
* SYNTHETIC database themselves, in a JUnit temp directory — never the operator's real
|
||||
* {@code ~/.local/share/opencode/opencode.db}, which a live opencode process may be writing.
|
||||
* {@link OpenCodeSessionDiscovery} matches an opencode session record by the worker's cwd (its
|
||||
* {@code directory}) against opencode's on-disk storage. These tests populate a TEMP storage root
|
||||
* themselves — never the operator's real {@code ~/.local/share/opencode}.
|
||||
*/
|
||||
class OpenCodeSessionDiscoveryTest {
|
||||
|
||||
/**
|
||||
* Create {@code <root>/opencode.db} with a minimal {@code session} table (just the columns
|
||||
* {@link OpenCodeSessionDiscovery} reads: {@code id}, {@code directory}, {@code time_updated})
|
||||
* and insert one row. Static so {@link OpenCodeLauncherTest} can reuse it.
|
||||
* Write a session record {@code {"id":..., "directory":...}} under
|
||||
* {@code <root>/session/<projectID>/<fileName>} and stamp it with a known last-modified time,
|
||||
* so "most recently modified wins" is deterministic. Static so the launcher test can reuse it.
|
||||
*/
|
||||
static void writeRecord(Path root, String id, String directory, long timeUpdated) 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)");
|
||||
}
|
||||
// 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.setString(1, id);
|
||||
insert.setString(2, directory);
|
||||
insert.setLong(3, timeUpdated);
|
||||
insert.executeUpdate();
|
||||
}
|
||||
}
|
||||
static void writeRecord(Path root, String projectId, String fileName, String id,
|
||||
String directory, long lastModifiedEpochMillis) throws Exception {
|
||||
Path dir = root.resolve("session").resolve(projectId);
|
||||
Files.createDirectories(dir);
|
||||
Path file = dir.resolve(fileName);
|
||||
Files.writeString(file, "{\"id\":\"" + id + "\",\"directory\":\"" + directory
|
||||
+ "\",\"projectID\":\"" + projectId + "\",\"version\":\"1.1.31\"}");
|
||||
Files.setLastModifiedTime(file, FileTime.fromMillis(lastModifiedEpochMillis));
|
||||
}
|
||||
|
||||
@Test
|
||||
void findsTheRowWhoseDirectoryEqualsTheCwd(@TempDir Path root) throws Exception {
|
||||
writeRecord(root, "ses_aaa", "/w/a", 1000L);
|
||||
writeRecord(root, "ses_bbb", "/w/b", 2000L);
|
||||
void findsTheRecordWhoseDirectoryEqualsTheCwd(@TempDir Path root) throws Exception {
|
||||
writeRecord(root, "p1", "ses_a.json", "ses_aaa", "/w/a", 1000L);
|
||||
writeRecord(root, "p2", "ses_b.json", "ses_bbb", "/w/b", 2000L);
|
||||
|
||||
assertEquals("ses_bbb", new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/b"),
|
||||
"the row whose directory equals the cwd is the one found");
|
||||
"the record whose directory equals the cwd is the one found");
|
||||
assertEquals("ses_aaa", new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void aNonMatchingDirectoryYieldsNullRatherThanAMismatch(@TempDir Path root) throws Exception {
|
||||
writeRecord(root, "ses_aaa", "/w/a", 1000L);
|
||||
writeRecord(root, "p1", "ses_a.json", "ses_aaa", "/w/a", 1000L);
|
||||
|
||||
assertNull(new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/other"),
|
||||
"no row for this cwd yet → null, not a wrong session");
|
||||
"no record for this cwd yet → null, not a wrong session");
|
||||
}
|
||||
|
||||
@Test
|
||||
void prefersTheMostRecentlyUpdatedRowWhenSeveralMatch(@TempDir Path root) throws Exception {
|
||||
writeRecord(root, "ses_old", "/w/a", 1000L);
|
||||
writeRecord(root, "ses_new", "/w/a", 5000L);
|
||||
void prefersTheMostRecentlyModifiedRecordWhenSeveralMatch(@TempDir Path root) throws Exception {
|
||||
writeRecord(root, "p1", "old.json", "ses_old", "/w/a", 1000L);
|
||||
writeRecord(root, "p2", "new.json", "ses_new", "/w/a", 5000L);
|
||||
|
||||
assertEquals("ses_new", new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"),
|
||||
"the row with the highest time_updated for the cwd wins");
|
||||
"the freshest record for the cwd wins");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aMissingDatabaseYieldsNullWithoutThrowing(@TempDir Path root) {
|
||||
// No opencode.db at all under the root.
|
||||
void aMissingOrEmptyStorageRootYieldsNullWithoutThrowing(@TempDir Path root) throws Exception {
|
||||
// Missing: no session dir at all under the root.
|
||||
assertNull(new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void anEmptyDatabaseYieldsNullWithoutThrowing(@TempDir Path root) throws Exception {
|
||||
Path db = root.resolve("opencode.db");
|
||||
try (Connection connection = DriverManager.getConnection("jdbc:sqlite:" + db);
|
||||
Statement statement = connection.createStatement()) {
|
||||
statement.execute("CREATE TABLE session (id TEXT PRIMARY KEY, directory TEXT, "
|
||||
+ "time_updated INTEGER)");
|
||||
}
|
||||
|
||||
assertNull(new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"));
|
||||
// Present but empty: a session dir with nothing in it produces no match, not a throw.
|
||||
Path emptyRoot = root.resolve("empty");
|
||||
Files.createDirectories(emptyRoot.resolve("session"));
|
||||
assertNull(new OpenCodeSessionDiscovery(emptyRoot).sessionIdForDirectory("/w/a"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -97,41 +76,15 @@ class OpenCodeSessionDiscoveryTest {
|
||||
assertNull(discovery.sessionIdForDirectory(" "));
|
||||
}
|
||||
|
||||
/**
|
||||
* The one line standing between fleetd and writing the operator's live {@code opencode.db} —
|
||||
* 841MB, with a running opencode writing it — is {@code config.setReadOnly(true)} in
|
||||
* {@link OpenCodeSessionDiscovery#openReadOnly()}. Delete it and every other test in this class
|
||||
* still passes, so this is the test that guards it.
|
||||
*
|
||||
* <p>It asks the connection to write, and requires a refusal. The obvious alternative — make
|
||||
* the database file unwritable and check the read still works — proves nothing: SQLite silently
|
||||
* downgrades a read-write open of an unwritable file to read-only, so that test passes either
|
||||
* way. It was tried and watched pass with the flag removed.
|
||||
*/
|
||||
@Test
|
||||
void theDatabaseIsOpenedReadOnly(@TempDir Path root) throws Exception {
|
||||
writeRecord(root, "ses_aaa", "/w/a", 1000L);
|
||||
void aMalformedRecordIsSkippedRatherThanFatal(@TempDir Path root) throws Exception {
|
||||
// A record that fails to parse must not abort the scan of its siblings.
|
||||
Path dir = root.resolve("session").resolve("p1");
|
||||
Files.createDirectories(dir);
|
||||
Files.writeString(dir.resolve("broken.json"), "{not valid json");
|
||||
writeRecord(root, "p1", "good.json", "ses_good", "/w/a", 1000L);
|
||||
|
||||
try (Connection connection = new OpenCodeSessionDiscovery(root).openReadOnly();
|
||||
Statement statement = connection.createStatement()) {
|
||||
SQLException refused = assertThrows(SQLException.class,
|
||||
() -> statement.executeUpdate("INSERT INTO session (id, directory, time_updated) "
|
||||
+ "VALUES ('ses_zzz', '/w/z', 1)"),
|
||||
"the connection must REFUSE a write — opencode is writing this database live");
|
||||
assertTrue(refused.getMessage().toLowerCase().contains("readonly")
|
||||
|| refused.getMessage().toLowerCase().contains("read-only"),
|
||||
"the refusal must be about read-only, not some other error: " + refused.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void aCorruptDatabaseFileYieldsNullWithoutThrowing(@TempDir Path root) throws Exception {
|
||||
// A file at opencode.db that is not a SQLite database at all — the open/query must fail
|
||||
// safe, never fatal to a spawn.
|
||||
Path db = root.resolve("opencode.db");
|
||||
Files.writeString(db, "this is not a sqlite database");
|
||||
|
||||
assertNull(new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"),
|
||||
"an unreadable database resolves to null, not an exception");
|
||||
assertEquals("ses_good", new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"),
|
||||
"an unreadable record is skipped; a later valid one still matches");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -688,32 +688,6 @@ class MessageServiceTest {
|
||||
assertFailedTicket(third, "agent target term_a not found");
|
||||
}
|
||||
|
||||
// --- #137 follow-up: abandon() must not guess when more than one task is open ---------------
|
||||
//
|
||||
// A test combining a genuine stranded reply (hasStrandedReply(T)==true) with two simultaneously
|
||||
// open matching tasks was attempted here and removed after investigation showed the combination
|
||||
// is not reachable through the public API today, not merely hard to time right:
|
||||
//
|
||||
// This class has exactly two call sites that ever hold a target's entry in the session-lock map
|
||||
// (send() and answer()), and both open a Rendezvous waiter for that same target as the first thing
|
||||
// they do after acquiring the lock, holding lock and waiter together for their whole critical
|
||||
// section. So "the lock is held" and "a live waiter is open" are the same fact throughout this
|
||||
// class. reply()'s fast path always resolves a currently-open waiter directly instead of
|
||||
// stranding — so a strand can only be created while NO task is accepted (lock free), and the
|
||||
// instant the lock is next taken (by any parked matching task's own send(), the moment it is
|
||||
// scheduled), that acceptance clears strandedReplies again (see send()'s CB-640 comment) before
|
||||
// abandon() can ever observe both facts together. Confirmed empirically too: an earlier version of
|
||||
// this test stranded a reply, then created an "accepted" task (awaitWaiting()) followed by a
|
||||
// "parked" one — and the accepted task's own acceptance silently cleared the strand it was
|
||||
// supposed to be racing against, so the parked task came back WORKER_FAILED instead of DONE, not
|
||||
// because the fix was missing but because the test's premise could not be constructed.
|
||||
//
|
||||
// The reachable half — matching.size() >= 2 alone, no strand — is exactly what
|
||||
// abandonFailsEveryPendingAsyncTicketForTheReleasedTarget already covers (all fail, none guess).
|
||||
// The oldest-wins code in abandon() stays as defence in depth (see its own javadoc) against a
|
||||
// regression that would make the conjunction reachable, e.g. clearing strandedReplies on a
|
||||
// narrower condition than "any acceptance" — not because this suite exercises it today.
|
||||
|
||||
@Test
|
||||
void abandonDoesNotFailAnAsyncTicketWaitingForAnAnswer() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "task that asks");
|
||||
@@ -736,68 +710,6 @@ class MessageServiceTest {
|
||||
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
|
||||
}
|
||||
|
||||
// --- #137: a fleet_ask round-trip must not orphan the ticket's own reply -------------------
|
||||
//
|
||||
// The primary's fleet_send{turnId} answer call is itself bounded (a real MCP call, capped well
|
||||
// under a minute) — far shorter than a resumed turn can genuinely take to finish real work. These
|
||||
// drive the exact real delegation path (async send -> worker asks -> primary answers -> primary's
|
||||
// own wait gives up -> worker's real fleet_reply arrives afterwards) rather than calling a reply
|
||||
// sink directly, since the bug is specifically about which sink the resumed turn's reply reaches.
|
||||
|
||||
@Test
|
||||
void aReplyAfterAnswerTimesOutStillCompletesTheAsyncTicket() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "task that asks");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
CompletableFuture<MessageService.AskResult> ask =
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
|
||||
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
|
||||
|
||||
// The primary answers, but its own bounded wait for the worker's resumed turn is short and
|
||||
// expires before the worker (still genuinely working) gets back to it.
|
||||
MessageService.Reply answerReply = messages.answer(asking.turnId(), "config.yaml", 150);
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome(),
|
||||
"the primary's own bounded wait gives up before the worker finishes resuming");
|
||||
|
||||
// The worker keeps working past that window and only now calls fleet_reply.
|
||||
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
|
||||
|
||||
MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE);
|
||||
assertEquals("PR opened: https://example/pulls/42", done.reply(),
|
||||
"fleet_poll{ticket} must return the worker's real reply, not stay pending forever");
|
||||
assertEquals("reply", done.replySource());
|
||||
assertFalse(messages.hasStrandedReply(T),
|
||||
"the reply completed its own ticket directly and never touched the inbox");
|
||||
}
|
||||
|
||||
@Test
|
||||
void fleetStopAfterAnOrphanedReplyDoesNotFailTheTicket() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "task that asks");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
CompletableFuture<MessageService.AskResult> ask =
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
|
||||
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
|
||||
|
||||
MessageService.Reply answerReply = messages.answer(asking.turnId(), "config.yaml", 150);
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome());
|
||||
|
||||
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
|
||||
|
||||
// fleet_stop tears the worker's session down right after the reply landed — this must never
|
||||
// report the misleading "the worker session was released before it replied": a reply is
|
||||
// exactly what happened.
|
||||
assertFalse(messages.abandon(T, "the worker session was released before it replied"),
|
||||
"a reply already arrived, so nothing here is a genuine failure");
|
||||
|
||||
MessageService.TaskView view = awaitTicketPhase(ticket, MessageService.Phase.DONE);
|
||||
assertEquals("PR opened: https://example/pulls/42", view.reply());
|
||||
}
|
||||
|
||||
@Test
|
||||
void unansweredAsyncQuestionReturnsTheTicketToPendingAndReleasesItsTarget() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "task that asks");
|
||||
|
||||
@@ -218,17 +218,9 @@ class FleetAppTest {
|
||||
|
||||
HttpResponse<String> res = req(port, "GET", "/members");
|
||||
assertEquals(200, res.statusCode());
|
||||
JsonNode body = mapper.readTree(res.body());
|
||||
// fleetd #199: "members" is the canonical key. The endpoint is /members, so a caller that
|
||||
// reads "members" must not see an empty fleet. "workers" is kept only as a deprecated alias
|
||||
// and must carry the same rows — assert both, or the alias can silently drift.
|
||||
JsonNode members = body.get("members");
|
||||
assertNotNull(members, "GET /members must return its rows under \"members\"");
|
||||
assertEquals(1, members.size());
|
||||
JsonNode workers = body.get("workers");
|
||||
assertNotNull(workers, "the deprecated \"workers\" alias is still emitted");
|
||||
assertEquals(members, workers, "the alias must carry the same rows as \"members\"");
|
||||
JsonNode w = members.get(0);
|
||||
JsonNode workers = mapper.readTree(res.body()).get("workers");
|
||||
assertEquals(1, workers.size());
|
||||
JsonNode w = workers.get(0);
|
||||
assertEquals(spawned.get("terminalId").asText(), w.get("sessionId").asText());
|
||||
assertEquals(paneId, w.get("paneId").asText());
|
||||
assertEquals("ltms-local", w.get("profile").asText());
|
||||
|
||||
@@ -30,19 +30,12 @@ public final class FakeWorktrees implements Worktrees {
|
||||
public record PruneCall(String repoRoot, long minAgeMillis) {
|
||||
}
|
||||
|
||||
public record ShareCall(String repoRoot, String worktreePath) {
|
||||
}
|
||||
|
||||
private final List<AddCall> addCalls = new CopyOnWriteArrayList<>();
|
||||
private final List<RemoveCall> removeCalls = new CopyOnWriteArrayList<>();
|
||||
private final List<OverlayCall> overlayCalls = new CopyOnWriteArrayList<>();
|
||||
private final List<RepoRootCall> repoRootCalls = new CopyOnWriteArrayList<>();
|
||||
private final List<SnapshotCall> snapshotCalls = new CopyOnWriteArrayList<>();
|
||||
private final List<PruneCall> pruneCalls = new CopyOnWriteArrayList<>();
|
||||
private final List<ShareCall> shareCalls = new CopyOnWriteArrayList<>();
|
||||
/** Tags every {@code overlayParity}/{@code shareWithGroup} call in call order, so a test can
|
||||
* pin that sharing runs after the overlay copy (fleetd #185 stage 3). */
|
||||
private final List<String> overlayShareOrder = new CopyOnWriteArrayList<>();
|
||||
private final Set<String> existingPaths = ConcurrentHashMap.newKeySet();
|
||||
private final Set<String> trackedPaths = ConcurrentHashMap.newKeySet();
|
||||
private final AtomicLong snapshotSeq = new AtomicLong();
|
||||
@@ -143,13 +136,6 @@ public final class FakeWorktrees implements Worktrees {
|
||||
}
|
||||
overlayCalls.add(new OverlayCall(repoRoot, worktreePath, List.copyOf(overlay),
|
||||
List.copyOf(copied), List.copyOf(skipped)));
|
||||
overlayShareOrder.add("overlay:" + worktreePath);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shareWithGroup(String repoRoot, String worktreePath) {
|
||||
shareCalls.add(new ShareCall(repoRoot, worktreePath));
|
||||
overlayShareOrder.add("share:" + worktreePath);
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -217,17 +203,4 @@ public final class FakeWorktrees implements Worktrees {
|
||||
public SnapshotCall lastSnapshot() {
|
||||
return snapshotCalls.isEmpty() ? null : snapshotCalls.getLast();
|
||||
}
|
||||
|
||||
public List<ShareCall> shareCalls() {
|
||||
return List.copyOf(shareCalls);
|
||||
}
|
||||
|
||||
public ShareCall lastShare() {
|
||||
return shareCalls.isEmpty() ? null : shareCalls.getLast();
|
||||
}
|
||||
|
||||
/** Call-order tags ({@code "overlay:<path>"}/{@code "share:<path>"}) — see field javadoc. */
|
||||
public List<String> overlayShareOrder() {
|
||||
return List.copyOf(overlayShareOrder);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -924,179 +924,4 @@ class GitWorktreesTest {
|
||||
assertEquals(2, stats.count(), "two snapshot refs are reported");
|
||||
assertTrue(stats.costBytes() > 0, "the cost of the snapshots is a positive byte count");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #185 stage 3: a recording {@link java.util.function.Function} test seam stands in for
|
||||
* every process {@link GitWorktrees#shareWithGroup} would run — no real second OS user/group
|
||||
* exists on this host, so these are unit tests against that seam, not a live-group integration
|
||||
* test (out of scope per the ticket).
|
||||
*/
|
||||
private static List<String> joined(String[] command) {
|
||||
return List.of(command);
|
||||
}
|
||||
|
||||
/** Remove {@code path} and anything under it. Tolerates an already-absent path. */
|
||||
private static void deleteRecursively(Path path) throws Exception {
|
||||
if (!Files.exists(path)) {
|
||||
return;
|
||||
}
|
||||
if (Files.isDirectory(path)) {
|
||||
try (java.util.stream.Stream<Path> children = Files.list(path)) {
|
||||
for (Path child : children.toList()) {
|
||||
deleteRecursively(child);
|
||||
}
|
||||
}
|
||||
}
|
||||
Files.delete(path);
|
||||
}
|
||||
|
||||
/** {@code worktreeGroup} absent ⇒ zero processes spawned and no git config written. */
|
||||
@Test
|
||||
void shareWithGroupIsNoopWhenNoGroupConfigured(@TempDir Path tmp) throws Exception {
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
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(tmp.resolve("wts").toString(), null, _ -> {}, recordingRunner);
|
||||
|
||||
gitWorktrees.shareWithGroup(repo.toString(), repo.resolve("some-worktree").toString());
|
||||
|
||||
assertTrue(recorded.isEmpty(), "no group configured must spawn no process at all: " + recorded);
|
||||
}
|
||||
|
||||
/** A configured group runs {@code git config core.sharedRepository group} first, then
|
||||
* chgrp/chmod/setgid over every path {@link GitWorktrees#shareWithGroup} documents. */
|
||||
@Test
|
||||
void shareWithGroupRunsConfigThenChgrpChmodSetgidPerPath(@TempDir Path tmp) throws Exception {
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
String repoRoot = repo.toString();
|
||||
Path worktree = Files.createDirectories(repo.resolve("some-worktree"));
|
||||
String worktreePath = worktree.toString();
|
||||
Files.createDirectories(repo.resolve(".git/worktrees"));
|
||||
List<List<String>> recorded = new java.util.ArrayList<>();
|
||||
java.util.function.Function<String[], String> recordingRunner = cmd -> {
|
||||
recorded.add(joined(cmd));
|
||||
// What real git answers for an ordinary (non-linked) checkout: relative to repoRoot.
|
||||
return List.of(cmd).contains("--git-common-dir") ? ".git\n" : "";
|
||||
};
|
||||
GitWorktrees gitWorktrees =
|
||||
new GitWorktrees(tmp.resolve("wts").toString(), "devteam", _ -> {}, recordingRunner);
|
||||
|
||||
gitWorktrees.shareWithGroup(repoRoot, worktreePath);
|
||||
|
||||
assertEquals(List.of("git", "-C", repoRoot, "config", "core.sharedRepository", "group"), recorded.get(0),
|
||||
"core.sharedRepository must be set first, so it keeps working after the one-time fix-up");
|
||||
assertTrue(recorded.contains(List.of("git", "-C", repoRoot, "rev-parse", "--git-common-dir")),
|
||||
"the git dir must be asked for, never hardcoded as <repoRoot>/.git — that is a FILE "
|
||||
+ "when the checkout is itself a linked worktree: " + recorded);
|
||||
|
||||
for (String dir : List.of(worktreePath, repoRoot + "/.git/objects", repoRoot + "/.git/refs",
|
||||
repoRoot + "/.git/logs", repoRoot + "/.git/worktrees")) {
|
||||
assertTrue(recorded.contains(List.of("chgrp", "-R", "devteam", dir)), "missing chgrp -R for " + dir);
|
||||
assertTrue(recorded.contains(List.of("chmod", "-R", "g+rwX", dir)), "missing chmod -R for " + dir);
|
||||
assertTrue(recorded.contains(List.of("find", dir, "-type", "d", "-exec", "chmod", "g+s", "{}", "+")),
|
||||
"missing setgid find pass for " + dir);
|
||||
}
|
||||
// packed-refs does not exist in a freshly-init'd repo (only git gc / pack-refs creates it) —
|
||||
// tolerated absence, so it must not appear at all: no recursive/-R treatment for a plain file.
|
||||
String packedRefs = repoRoot + "/.git/packed-refs";
|
||||
assertTrue(recorded.stream().noneMatch(c -> c.contains(packedRefs)),
|
||||
"packed-refs is absent here and must be skipped, not chgrp'd: " + recorded);
|
||||
}
|
||||
|
||||
/**
|
||||
* A path that does not exist is skipped, never handed to {@code chgrp}. {@code .git/logs} is
|
||||
* absent whenever {@code core.logAllRefUpdates} is false or no ref has been updated yet, and
|
||||
* {@code chgrp} on a missing path exits non-zero — which would fail EVERY provisioning spawn
|
||||
* with a message blaming a group that is in fact fine.
|
||||
*/
|
||||
@Test
|
||||
void shareWithGroupSkipsPathsThatDoNotExist(@TempDir Path tmp) throws Exception {
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
String repoRoot = repo.toString();
|
||||
deleteRecursively(repo.resolve(".git/logs"));
|
||||
assertFalse(Files.exists(repo.resolve(".git/logs")), "fixture: .git/logs must be gone");
|
||||
List<List<String>> recorded = new java.util.ArrayList<>();
|
||||
java.util.function.Function<String[], String> recordingRunner = cmd -> {
|
||||
recorded.add(joined(cmd));
|
||||
return List.of(cmd).contains("--git-common-dir") ? ".git\n" : "";
|
||||
};
|
||||
GitWorktrees gitWorktrees =
|
||||
new GitWorktrees(tmp.resolve("wts").toString(), "devteam", _ -> {}, recordingRunner);
|
||||
|
||||
gitWorktrees.shareWithGroup(repoRoot, repo.resolve("no-such-worktree").toString());
|
||||
|
||||
String logs = repoRoot + "/.git/logs";
|
||||
assertTrue(recorded.stream().noneMatch(c -> c.contains(logs)),
|
||||
"a missing .git/logs must be skipped, not chgrp'd: " + recorded);
|
||||
assertTrue(recorded.stream().noneMatch(c -> c.contains(repo.resolve("no-such-worktree").toString())),
|
||||
"a missing worktree path must be skipped too: " + recorded);
|
||||
assertTrue(recorded.contains(List.of("chgrp", "-R", "devteam", repoRoot + "/.git/objects")),
|
||||
"paths that DO exist are still shared: " + recorded);
|
||||
}
|
||||
|
||||
/**
|
||||
* The git store is located by {@code rev-parse --git-common-dir}, not by appending
|
||||
* {@code /.git}. When git answers with an absolute path — what it does for a linked worktree,
|
||||
* where {@code <repoRoot>/.git} is a file — every shared path must follow that answer.
|
||||
*/
|
||||
@Test
|
||||
void shareWithGroupFollowsAnAbsoluteGitCommonDir(@TempDir Path tmp) throws Exception {
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
Path realGitDir = repo.resolve(".git");
|
||||
List<List<String>> recorded = new java.util.ArrayList<>();
|
||||
java.util.function.Function<String[], String> recordingRunner = cmd -> {
|
||||
recorded.add(joined(cmd));
|
||||
return List.of(cmd).contains("--git-common-dir") ? realGitDir + "\n" : "";
|
||||
};
|
||||
GitWorktrees gitWorktrees =
|
||||
new GitWorktrees(tmp.resolve("wts").toString(), "devteam", _ -> {}, recordingRunner);
|
||||
|
||||
gitWorktrees.shareWithGroup(tmp.resolve("some/linked/worktree").toString(),
|
||||
repo.resolve("wt").toString());
|
||||
|
||||
assertTrue(recorded.contains(List.of("chgrp", "-R", "devteam", realGitDir + "/objects")),
|
||||
"objects must be taken from the reported common dir, not <repoRoot>/.git: " + recorded);
|
||||
}
|
||||
|
||||
/** {@code packed-refs}, when present, is chgrp/chmod'd but never setgid'd (it is a file, not a dir). */
|
||||
@Test
|
||||
void shareWithGroupIncludesPackedRefsWhenPresent(@TempDir Path tmp) throws Exception {
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
String repoRoot = repo.toString();
|
||||
Path packedRefsPath = repo.resolve(".git/packed-refs");
|
||||
Files.writeString(packedRefsPath, "");
|
||||
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(tmp.resolve("wts").toString(), "devteam", _ -> {}, recordingRunner);
|
||||
|
||||
gitWorktrees.shareWithGroup(repoRoot, repo.resolve("some-worktree").toString());
|
||||
|
||||
String packedRefs = packedRefsPath.toString();
|
||||
assertTrue(recorded.contains(List.of("chgrp", "devteam", packedRefs)),
|
||||
"packed-refs must be chgrp'd non-recursively when present: " + recorded);
|
||||
assertTrue(recorded.contains(List.of("chmod", "g+rwX", packedRefs)),
|
||||
"packed-refs must be chmod'd non-recursively when present: " + recorded);
|
||||
assertTrue(recorded.stream().noneMatch(c -> c.contains("find") && c.contains(packedRefs)),
|
||||
"packed-refs (a file) must never get the recursive setgid pass: " + recorded);
|
||||
}
|
||||
|
||||
/** A group that does not exist (or that the operator is not a member of) fails loudly, naming it. */
|
||||
@Test
|
||||
void shareWithGroupThrowsNamingTheGroupWhenChgrpFails(@TempDir Path tmp) throws Exception {
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
GitWorktrees gitWorktrees = new GitWorktrees(tmp.resolve("wts").toString(), "cb185-nonexistent-group-zz");
|
||||
String wt = new GitWorktrees(tmp.resolve("wts").toString()).add(repo.toString(), "cb-185-share", "HEAD");
|
||||
|
||||
WorktreeException e = assertThrows(WorktreeException.class,
|
||||
() -> gitWorktrees.shareWithGroup(repo.toString(), wt));
|
||||
assertTrue(e.getMessage().contains("cb185-nonexistent-group-zz"),
|
||||
"exception must name the missing/refused group: " + e.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -148,10 +148,6 @@ class SessionManagerTest {
|
||||
return 0;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shareWithGroup(String repoRoot, String worktreePath) {
|
||||
}
|
||||
|
||||
List<String> removeCalls() {
|
||||
return List.copyOf(removeCalls);
|
||||
}
|
||||
@@ -931,235 +927,4 @@ class SessionManagerTest {
|
||||
assertTrue(e.getMessage().contains("SESSION_RESUME"), e.getMessage());
|
||||
assertTrue(e.getMessage().contains("stub-profile"), e.getMessage());
|
||||
}
|
||||
|
||||
// ── fleetd #209: agentSessionId resolved lazily against the retained handle ─────────────────
|
||||
//
|
||||
// The opencode adapter cannot answer PeerHandle.agentSessionId() at spawn time — the on-disk
|
||||
// session row is written only after the pane is live — so the id must be re-polled on a LATER
|
||||
// call, against the SAME handle instance the launcher returned at spawn. SessionManager used
|
||||
// to let that handle go out of scope at the end of the spawn method, so no caller ever re-asked
|
||||
// it and fleet_list/fleet_status never saw the id. LazyIdHandle below reproduces exactly that
|
||||
// shape: null on the first N calls (the spawn-time call included), a real id after.
|
||||
|
||||
/**
|
||||
* A {@link PeerHandle} whose {@link #agentSessionId()} answers {@code null} for its first
|
||||
* {@code nullCalls} invocations, then either a fixed id or a configured throw on every call
|
||||
* after that — the shape of the opencode bug (fleetd #209): the session row is not written
|
||||
* until after the pane is live, so early polls come back empty and a later one finds it.
|
||||
*/
|
||||
private static final class LazyIdHandle implements PeerHandle {
|
||||
private final String id;
|
||||
private final String terminalId;
|
||||
private final int nullCalls;
|
||||
private final String resolvedId;
|
||||
private final java.util.concurrent.atomic.AtomicInteger calls =
|
||||
new java.util.concurrent.atomic.AtomicInteger();
|
||||
private volatile RuntimeException throwAfter;
|
||||
|
||||
LazyIdHandle(String id, String terminalId, int nullCalls, String resolvedId) {
|
||||
this.id = id;
|
||||
this.terminalId = terminalId;
|
||||
this.nullCalls = nullCalls;
|
||||
this.resolvedId = resolvedId;
|
||||
}
|
||||
|
||||
/** After the null calls are exhausted, throw instead of answering the resolved id. */
|
||||
LazyIdHandle throwing(RuntimeException e) {
|
||||
this.throwAfter = e;
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String id() {
|
||||
return id;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String terminalId() {
|
||||
return terminalId;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String agentSessionId() {
|
||||
int n = calls.incrementAndGet();
|
||||
if (n <= nullCalls) {
|
||||
return null;
|
||||
}
|
||||
if (throwAfter != null) {
|
||||
throw throwAfter;
|
||||
}
|
||||
return resolvedId;
|
||||
}
|
||||
|
||||
@Override
|
||||
public CharterReceipt charterReceipt() {
|
||||
return null;
|
||||
}
|
||||
|
||||
int callCount() {
|
||||
return calls.get();
|
||||
}
|
||||
}
|
||||
|
||||
/** A minimal {@link PeerLauncher} that hands out pre-built {@link LazyIdHandle}s, one per spawn. */
|
||||
private static final class LazyIdLauncher implements PeerLauncher {
|
||||
private final java.util.Deque<LazyIdHandle> queued = new java.util.ArrayDeque<>();
|
||||
|
||||
LazyIdLauncher queue(LazyIdHandle handle) {
|
||||
queued.add(handle);
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<Capability> capabilities() {
|
||||
return Set.of(Capability.WORKTREE, Capability.SESSION_RESUME);
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<Capability> capabilitiesFor(String profileName) {
|
||||
return capabilities();
|
||||
}
|
||||
|
||||
@Override
|
||||
public PeerHandle spawn(SpawnRequest req) {
|
||||
LazyIdHandle handle = queued.poll();
|
||||
if (handle == null) {
|
||||
throw new IllegalStateException("no queued LazyIdHandle for this spawn");
|
||||
}
|
||||
return handle;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<String> profiles() {
|
||||
return Set.of("lazy");
|
||||
}
|
||||
|
||||
@Override
|
||||
public String defaultProfile() {
|
||||
return "lazy";
|
||||
}
|
||||
|
||||
@Override
|
||||
public String effectiveCwd(SpawnRequest req) {
|
||||
return "/cwd";
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<String> parityOverlay(String profileName) {
|
||||
return List.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<?> list() {
|
||||
return List.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public int reapOrphanWorkers() {
|
||||
return 0;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void stop(String id) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean clearContext(String id) {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void plainRosterDoesNotResolveAgentSessionId() {
|
||||
// fleetd #209 follow-up: roster() sits on the heartbeat/health-tick timers (and the metrics
|
||||
// scrape), so it must never trigger the resolve lookup — for opencode that lookup opens an
|
||||
// on-disk session database, and a member whose id never appears would pay that cost forever.
|
||||
// rosterResolved() is the one to use when a caller actually reports the id.
|
||||
LazyIdHandle handle = new LazyIdHandle("p0", "t0", 1, "oc-session-0");
|
||||
SessionManager sessions = new SessionManager(new LazyIdLauncher().queue(handle));
|
||||
MemberSession acquired = sessions.acquire("lazy", "/cwd", "/caller", null);
|
||||
assertEquals(1, handle.callCount(), "sanity: only the spawn-time call happened so far");
|
||||
|
||||
List<MemberSession> roster = sessions.roster();
|
||||
|
||||
assertEquals(1, roster.size());
|
||||
assertEquals(acquired.paneId(), roster.getFirst().paneId());
|
||||
assertNull(roster.getFirst().agentSessionId(), "the plain roster must not resolve the id");
|
||||
assertEquals(1, handle.callCount(),
|
||||
"roster() must never call agentSessionId() again — it sits on the heartbeat/health timers");
|
||||
}
|
||||
|
||||
@Test
|
||||
void rosterResolvedResolvesALateAgentSessionIdFromTheRetainedHandle() {
|
||||
LazyIdHandle handle = new LazyIdHandle("p1", "t1", 1, "oc-session-1");
|
||||
SessionManager sessions = new SessionManager(new LazyIdLauncher().queue(handle));
|
||||
|
||||
MemberSession acquired = sessions.acquire("lazy", "/cwd", "/caller", null);
|
||||
assertNull(acquired.agentSessionId(),
|
||||
"opencode has not written its session row yet at spawn time");
|
||||
|
||||
List<MemberSession> roster = sessions.rosterResolved();
|
||||
assertEquals(1, roster.size());
|
||||
assertEquals("oc-session-1", roster.getFirst().agentSessionId(),
|
||||
"fleet_list must see the id once the adapter can answer it");
|
||||
|
||||
Map<String, Object> view = SessionManager.rosterView(roster.getFirst(), null);
|
||||
assertEquals("oc-session-1", view.get("agentSessionId"),
|
||||
"rosterView renders whatever rosterResolved() resolved");
|
||||
}
|
||||
|
||||
@Test
|
||||
void getResolvesALateAgentSessionIdFromTheRetainedHandle() {
|
||||
LazyIdHandle handle = new LazyIdHandle("p2", "t2", 1, "oc-session-2");
|
||||
SessionManager sessions = new SessionManager(new LazyIdLauncher().queue(handle));
|
||||
MemberSession acquired = sessions.acquire("lazy", "/cwd", "/caller", null);
|
||||
|
||||
MemberSession resolved = sessions.get(acquired.paneId()).orElseThrow();
|
||||
|
||||
assertEquals("oc-session-2", resolved.agentSessionId(),
|
||||
"fleet_status (single-session lookup) must also see the late-resolved id");
|
||||
}
|
||||
|
||||
@Test
|
||||
void releaseCarriesALateResolvedAgentSessionIdIntoTheReleaseDetail() {
|
||||
LazyIdHandle handle = new LazyIdHandle("p3", "t3", 1, "oc-session-3");
|
||||
SessionManager sessions = new SessionManager(new LazyIdLauncher().queue(handle));
|
||||
MemberSession acquired = sessions.acquire("lazy", "/cwd", "/caller", null);
|
||||
java.util.List<String> released = new java.util.concurrent.CopyOnWriteArrayList<>();
|
||||
sessions.onRelease(detail -> released.add(detail.agentSessionId()));
|
||||
|
||||
sessions.release(acquired.paneId());
|
||||
|
||||
assertEquals(java.util.List.of("oc-session-3"), released,
|
||||
"a released member's detail carries the id it has since resolved, not the null "
|
||||
+ "frozen in at spawn time");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aThrowingHandleDoesNotBreakRosterResolved() {
|
||||
LazyIdHandle handle = new LazyIdHandle("p4", "t4", 1, "oc-session-4")
|
||||
.throwing(new RuntimeException("sqlite locked"));
|
||||
SessionManager sessions = new SessionManager(new LazyIdLauncher().queue(handle));
|
||||
sessions.acquire("lazy", "/cwd", "/caller", null);
|
||||
|
||||
List<MemberSession> roster = assertDoesNotThrow(sessions::rosterResolved,
|
||||
"a handle that throws resolving its id must not break the roster read");
|
||||
|
||||
assertEquals(1, roster.size());
|
||||
assertNull(roster.getFirst().agentSessionId(), "the id stays unresolved when the lookup throws");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aResolvedAgentSessionIdIsNotLookedUpAgain() {
|
||||
LazyIdHandle handle = new LazyIdHandle("p5", "t5", 1, "oc-session-5");
|
||||
SessionManager sessions = new SessionManager(new LazyIdLauncher().queue(handle));
|
||||
sessions.acquire("lazy", "/cwd", "/caller", null);
|
||||
assertEquals(1, handle.callCount(), "sanity: only the spawn-time call happened so far");
|
||||
|
||||
sessions.rosterResolved();
|
||||
assertEquals(2, handle.callCount(), "the first rosterResolved() read resolves the id");
|
||||
|
||||
sessions.rosterResolved();
|
||||
assertEquals(2, handle.callCount(), "once resolved, the id must not be looked up again");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -163,37 +163,6 @@ class WorktreeSessionManagerTest {
|
||||
"tracked copied paths are --skip-worktree'd");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #185 stage 3, THE TRAP: {@code overlayParity} copies more files into the worktree
|
||||
* AFTER {@code add} returns, so {@code shareWithGroup} must run after it, not folded into
|
||||
* {@code add()} — otherwise every overlay file lands operator-owned and unwritable for a
|
||||
* different-uid member, with a green test suite hiding it.
|
||||
*/
|
||||
@Test
|
||||
void shareWithGroupRunsAfterOverlayParityNotBeforeIt() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt")
|
||||
.track(".envrc");
|
||||
SessionManager sessions = new SessionManager(workerService(herdr), worktrees);
|
||||
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("cb-185", null));
|
||||
|
||||
assertEquals(1, worktrees.overlayCalls().size(), "overlayParity ran exactly once");
|
||||
assertEquals(1, worktrees.shareCalls().size(), "shareWithGroup ran exactly once");
|
||||
FakeWorktrees.OverlayCall overlay = worktrees.lastOverlay();
|
||||
FakeWorktrees.ShareCall share = worktrees.lastShare();
|
||||
assertEquals(s.worktree(), overlay.worktreePath());
|
||||
assertEquals(s.worktree(), share.worktreePath());
|
||||
|
||||
List<String> order = worktrees.overlayShareOrder();
|
||||
int overlayIndex = order.indexOf("overlay:" + s.worktree());
|
||||
int shareIndex = order.indexOf("share:" + s.worktree());
|
||||
assertTrue(overlayIndex >= 0 && shareIndex >= 0, "both calls must be recorded: " + order);
|
||||
assertTrue(overlayIndex < shareIndex,
|
||||
"shareWithGroup MUST run after overlayParity, not before/inside add(): " + order);
|
||||
}
|
||||
|
||||
@Test
|
||||
void releaseRemovesWorktreeButDoesNotDeleteBranch() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
|
||||
Reference in New Issue
Block a user