Compare commits
5 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 6d82ca95a4 | |||
| 8067ee4ec4 | |||
| 388aba7632 | |||
| ad587eafa3 | |||
| 08ce9aef11 |
@@ -634,6 +634,17 @@ 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.
|
||||
# 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)
|
||||
|
||||
@@ -177,14 +177,14 @@ public final class Fleetd {
|
||||
claudeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
|
||||
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
|
||||
() -> config.get().fleet(),
|
||||
() -> config.get().memberCredentials()));
|
||||
() -> config.get().memberCredentials(), null, config::get));
|
||||
}
|
||||
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().memberCredentials(), config::get));
|
||||
}
|
||||
AtomicReference<Function<String, Integer>> liveCountRef = new AtomicReference<>(_ -> 0);
|
||||
// CB-578 stage B: one quarantine tracker for the whole daemon, shared between the launcher
|
||||
@@ -224,7 +224,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()),
|
||||
SessionManager sessions = new SessionManager(workers, new GitWorktrees(cfg.worktreeRoot(), cfg.worktreeGroup()),
|
||||
System::nanoTime, contextCap, clearAfterTurn);
|
||||
liveCountRef.set(profileName -> (int) sessions.roster().stream()
|
||||
.filter(s -> profileName.equals(s.profile()))
|
||||
|
||||
@@ -76,6 +76,14 @@ 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}.
|
||||
*/
|
||||
@JsonIgnoreProperties(ignoreUnknown = true)
|
||||
public record FleetConfig(
|
||||
@@ -98,7 +106,20 @@ public record FleetConfig(
|
||||
ConfigReload configReload,
|
||||
Integer quarantineCooldownSeconds,
|
||||
MemberCredentials memberCredentials,
|
||||
Coordinator coordinator) {
|
||||
Coordinator coordinator,
|
||||
String worktreeGroup) {
|
||||
|
||||
/** 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);
|
||||
}
|
||||
|
||||
/** Back-compat form before the {@code coordinator:} block was added. */
|
||||
public FleetConfig(Bind bind, String herdrSocket, Map<String, Profile> profiles, Guard guard,
|
||||
@@ -109,7 +130,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);
|
||||
configReload, quarantineCooldownSeconds, memberCredentials, null, null);
|
||||
}
|
||||
|
||||
/** Back-compat form before the CB-596 {@code memberCredentials:} block was added. */
|
||||
@@ -1325,7 +1346,7 @@ public record FleetConfig(
|
||||
"bind", "herdrSocket", "memberHerdrSocket", "profiles", "guard", "worktreeRoot",
|
||||
"lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs", "broker", "primary", "fleet",
|
||||
"leadHeartbeat", "health", "placement", "auth", "configReload", "quarantineCooldownSeconds",
|
||||
"memberCredentials", "coordinator");
|
||||
"memberCredentials", "coordinator", "worktreeGroup");
|
||||
|
||||
/** Load and validate config from {@code path}. */
|
||||
public static FleetConfig load(Path path) {
|
||||
@@ -1943,9 +1964,11 @@ 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.
|
||||
return new FleetConfig(b, herdrSocket, memberHerdrSocket, profiles, g, worktreeRoot, l, timeout, pollMs,
|
||||
broker, primary, f, leadHeartbeat, health, placementOrDefault, a, configReload,
|
||||
quarantineCooldown, mc, coordinator);
|
||||
quarantineCooldown, mc, coordinator, worktreeGroup);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -2,8 +2,11 @@ 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
|
||||
@@ -12,23 +15,46 @@ import java.util.Map;
|
||||
* 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. This scans agent panes; a spawn-time {@code pid→terminal} cache is
|
||||
* the obvious optimization once wired into {@code ClaudeCodeLauncher}.
|
||||
* 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}.
|
||||
*
|
||||
* <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).
|
||||
* 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.
|
||||
*/
|
||||
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;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -37,7 +63,13 @@ 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;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -49,8 +81,9 @@ public final class PaneLocator {
|
||||
if (pid <= 0) {
|
||||
return null;
|
||||
}
|
||||
Set<Long> ancestry = ancestorsOf(pid);
|
||||
for (HerdrClient herdr : herdrs) {
|
||||
String terminal = terminalForPid(herdr, pid);
|
||||
String terminal = terminalForPid(herdr, ancestry);
|
||||
if (terminal != null) {
|
||||
return terminal;
|
||||
}
|
||||
@@ -58,28 +91,54 @@ public final class PaneLocator {
|
||||
return null;
|
||||
}
|
||||
|
||||
private static String terminalForPid(HerdrClient herdr, long pid) {
|
||||
/**
|
||||
* {@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) {
|
||||
for (JsonNode pane : herdr.call("pane.list", Map.of()).path("panes")) {
|
||||
String paneId = pane.path("pane_id").asText(null);
|
||||
if (paneId != null && paneOwnsPid(herdr, paneId, pid)) {
|
||||
if (paneId != null && paneOwnsAnyOf(herdr, paneId, ancestry)) {
|
||||
return pane.path("terminal_id").asText(null);
|
||||
}
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
private static boolean paneOwnsPid(HerdrClient herdr, String paneId, long pid) {
|
||||
private static boolean paneOwnsAnyOf(HerdrClient herdr, String paneId, Set<Long> ancestry) {
|
||||
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 (info.path("shell_pid").asLong(-1) == pid) {
|
||||
if (ancestry.contains(info.path("shell_pid").asLong(-1))) {
|
||||
return true;
|
||||
}
|
||||
for (JsonNode p : info.path("foreground_processes")) {
|
||||
if (p.path("pid").asLong(-1) == pid) {
|
||||
if (ancestry.contains(p.path("pid").asLong(-1))) {
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,21 @@
|
||||
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());
|
||||
}
|
||||
@@ -94,13 +94,26 @@ 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) {
|
||||
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) {
|
||||
this(agents, spaces, guard, profiles, defaultProfile, env,
|
||||
spawnReadyTimeoutMs,
|
||||
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs),
|
||||
fleet, memberCredentials);
|
||||
fleet, memberCredentials, hostEnvNames, config);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -153,10 +166,22 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
|
||||
Function<String, String> env,
|
||||
long spawnReadyTimeoutMs,
|
||||
LongSupplier nowMillis, Runnable sleeper,
|
||||
Supplier<FleetConfig.Fleet> fleet,
|
||||
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
|
||||
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) {
|
||||
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
|
||||
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials);
|
||||
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials, hostEnvNames, config);
|
||||
this.guard = guard;
|
||||
}
|
||||
|
||||
@@ -174,7 +199,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);
|
||||
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials, hostEnvNames, null);
|
||||
this.guard = guard;
|
||||
}
|
||||
|
||||
|
||||
@@ -170,6 +170,8 @@ 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)
|
||||
@@ -238,8 +240,19 @@ 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<Set<String>> hostEnvNames, Supplier<FleetConfig> config) {
|
||||
this.fleet = fleet;
|
||||
this.namePrefix = namePrefix;
|
||||
this.agents = agents;
|
||||
@@ -252,6 +265,7 @@ 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 -------------------------------------------------------------------------
|
||||
@@ -1028,9 +1042,11 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
}
|
||||
|
||||
/** Put {@link #BLOCKED_CREDENTIAL_SENTINEL} over every blocked name in the pane-creation env map. */
|
||||
private static void overlayBlockedCredentials(Map<String, String> workerEnv,
|
||||
FleetConfig.MemberCredentials creds) {
|
||||
for (String name : creds.blockedSet()) {
|
||||
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) {
|
||||
workerEnv.put(name, BLOCKED_CREDENTIAL_SENTINEL);
|
||||
}
|
||||
}
|
||||
@@ -1097,15 +1113,21 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
* 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()));
|
||||
MemberEnvAllowList.derive(profiles.values(), creds.allowSet(), brokerUriEnvNames));
|
||||
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
|
||||
|
||||
@@ -40,8 +40,9 @@ 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} 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
|
||||
* 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
|
||||
* (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
|
||||
@@ -106,6 +107,16 @@ 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) {
|
||||
@@ -124,9 +135,37 @@ 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
|
||||
@@ -146,4 +185,16 @@ 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());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -116,12 +116,23 @@ 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) {
|
||||
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) {
|
||||
this(agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs,
|
||||
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs),
|
||||
defaultConfigRoot(), defaultDiscoveryRoot(), fleet, memberCredentials);
|
||||
defaultConfigRoot(), defaultDiscoveryRoot(), fleet, memberCredentials, config);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -182,8 +193,20 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
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);
|
||||
}
|
||||
|
||||
/** Full testability 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,
|
||||
LongSupplier nowMillis, Runnable sleeper, Path configRoot, Path discoveryRoot,
|
||||
Supplier<FleetConfig.Fleet> fleet,
|
||||
Supplier<FleetConfig.MemberCredentials> memberCredentials,
|
||||
Supplier<FleetConfig> config) {
|
||||
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
|
||||
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials);
|
||||
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials, null, config);
|
||||
this.configRoot = configRoot;
|
||||
this.discovery = new OpenCodeSessionDiscovery(discoveryRoot);
|
||||
}
|
||||
|
||||
@@ -366,21 +366,40 @@ public final class MessageService {
|
||||
}
|
||||
|
||||
/**
|
||||
* 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.
|
||||
* 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.
|
||||
*
|
||||
* <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.
|
||||
*
|
||||
* @return always {@code true} — the reply either resolved a live send or was queued
|
||||
* @return always {@code true} — the reply resolved a live send, completed a parked ticket, 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.
|
||||
Task orphan = askAnsweredAsyncTask(session);
|
||||
if (orphan != null && 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
|
||||
}
|
||||
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.
|
||||
@@ -394,6 +413,24 @@ public final class MessageService {
|
||||
return true; // held, not lost
|
||||
}
|
||||
|
||||
/**
|
||||
* The 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). {@code null} 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}).
|
||||
*/
|
||||
private Task askAnsweredAsyncTask(String target) {
|
||||
for (Task task : tasks.values()) {
|
||||
if (target.equals(task.target) && task.question == null && task.turnId != null
|
||||
&& !task.future.isDone()) {
|
||||
return task;
|
||||
}
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
/** 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) {
|
||||
@@ -435,19 +472,43 @@ 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}.
|
||||
*
|
||||
* @return true if a live waiter was failed
|
||||
* <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 one is still parked
|
||||
* waiting for it (see {@link #askAnsweredAsyncTask}), 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).
|
||||
*
|
||||
* @return true if a live waiter or an async task was failed (never true for one recovered as a
|
||||
* reply — see the note above)
|
||||
*/
|
||||
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;
|
||||
Reply recovered = null; // lazily drained at most once, only if a task actually needs it
|
||||
for (Task task : tasks.values()) {
|
||||
if (target.equals(task.target) && task.question == null
|
||||
&& task.future.complete(new Reply(Outcome.WORKER_FAILED, reason))) {
|
||||
asyncFailed = true;
|
||||
if (!target.equals(task.target) || task.question != null || task.future.isDone()) {
|
||||
continue;
|
||||
}
|
||||
if (hadStrandedReply && recovered == null) {
|
||||
recovered = recoverStrandedReply(target);
|
||||
}
|
||||
Reply outcome = recovered != null ? 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);
|
||||
}
|
||||
}
|
||||
}
|
||||
if (failed) {
|
||||
@@ -456,6 +517,21 @@ 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).
|
||||
*/
|
||||
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.
|
||||
|
||||
@@ -23,6 +23,7 @@ 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;
|
||||
|
||||
/**
|
||||
@@ -88,24 +89,64 @@ 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);
|
||||
this(null, (String) null);
|
||||
}
|
||||
|
||||
/** @param configuredRoot nullable absolute or relative path; null/blank derives a sibling of the repo root. */
|
||||
/**
|
||||
* @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.
|
||||
*/
|
||||
public GitWorktrees(String configuredRoot) {
|
||||
this(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, _ -> {});
|
||||
}
|
||||
|
||||
/** 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
|
||||
@@ -607,6 +648,126 @@ 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
|
||||
|
||||
@@ -472,6 +472,10 @@ 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={}: {}",
|
||||
|
||||
@@ -98,4 +98,21 @@ 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,6 +1036,35 @@ 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");
|
||||
|
||||
@@ -0,0 +1,27 @@
|
||||
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,7 +1,11 @@
|
||||
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). */
|
||||
@@ -63,4 +67,123 @@ 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() {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+45
-2
@@ -161,6 +161,34 @@ 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
|
||||
@@ -260,15 +288,24 @@ 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) {
|
||||
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) {
|
||||
super("test", new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
Map.of("test", profile()), "test",
|
||||
name -> "SHELL".equals(name) ? shell : null,
|
||||
0, () -> 0L, () -> { }, null, creds, hostEnvNames);
|
||||
0, () -> 0L, () -> { }, null, creds, hostEnvNames, config);
|
||||
}
|
||||
|
||||
@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"));
|
||||
}
|
||||
|
||||
@@ -278,6 +315,12 @@ 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();
|
||||
}
|
||||
|
||||
/** The generated directory is a temp directory; make sure the test does not leave a pile. */
|
||||
@Test
|
||||
void theGeneratedDirectoryIsRemovedWhenThePaneIsStopped() {
|
||||
|
||||
@@ -121,6 +121,34 @@ 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() {
|
||||
|
||||
@@ -710,6 +710,68 @@ 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");
|
||||
|
||||
@@ -30,12 +30,19 @@ 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();
|
||||
@@ -136,6 +143,13 @@ 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
|
||||
@@ -203,4 +217,17 @@ 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,4 +924,179 @@ 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,6 +148,10 @@ class SessionManagerTest {
|
||||
return 0;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shareWithGroup(String repoRoot, String worktreePath) {
|
||||
}
|
||||
|
||||
List<String> removeCalls() {
|
||||
return List.copyOf(removeCalls);
|
||||
}
|
||||
|
||||
@@ -163,6 +163,37 @@ 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