CB-402 Increment 1: extract HerdrPeerLauncher abstract base
Behaviour-preserving refactor ahead of the second-adapter work. All herdr transport shared by any peer kind — tab/pane placement, the CB-306 spawn-readiness gate, unique naming + CB-117 orphan reap, teardown, listing, cwd resolution, and the peer-neutral git-forge grant — moves into a new abstract HerdrPeerLauncher (Template-Method base). ClaudeCodeLauncher becomes a final subclass supplying only the two Claude-specific seams: the `claude` name prefix and buildLaunch(), which encodes the subscription boundary (ANTHROPIC_BASE_URL + SubscriptionGuard assert, inline --mcp-config and --append-system-prompt reply charter). The base owns the injectable clock + sleeper for the readiness gate; the poll interval is baked into the sleeper, so the vestigial spawnReadyPollMs field is dropped from the base and from the full testability constructor (the explicit sleeper already encodes it). The 6-arg and 8-arg production constructors keep their signatures; three full-ctor test sites drop the now-unused poll argument. No behaviour change: 242 tests green, both refactored files 0 IDE problems.
This commit is contained in:
@@ -4,73 +4,39 @@ import dev.ltms.bridged.config.BridgedConfig;
|
|||||||
import dev.ltms.bridged.guard.SubscriptionGuard;
|
import dev.ltms.bridged.guard.SubscriptionGuard;
|
||||||
import dev.ltms.bridged.herdr.Agent;
|
import dev.ltms.bridged.herdr.Agent;
|
||||||
import dev.ltms.bridged.herdr.AgentControl;
|
import dev.ltms.bridged.herdr.AgentControl;
|
||||||
import dev.ltms.bridged.herdr.HerdrException;
|
|
||||||
import dev.ltms.bridged.herdr.Tab;
|
|
||||||
import dev.ltms.bridged.herdr.Workspace;
|
|
||||||
import dev.ltms.bridged.herdr.WorkspaceControl;
|
import dev.ltms.bridged.herdr.WorkspaceControl;
|
||||||
import dev.ltms.bridged.peer.Capability;
|
import dev.ltms.bridged.peer.Capability;
|
||||||
import dev.ltms.bridged.peer.PeerHandle;
|
|
||||||
import dev.ltms.bridged.peer.PeerLauncher;
|
|
||||||
import dev.ltms.bridged.peer.PeerUnreachableException;
|
|
||||||
import dev.ltms.bridged.peer.SpawnRequest;
|
|
||||||
import org.slf4j.Logger;
|
|
||||||
import org.slf4j.LoggerFactory;
|
|
||||||
|
|
||||||
import java.security.SecureRandom;
|
|
||||||
import java.util.ArrayList;
|
|
||||||
import java.util.EnumSet;
|
import java.util.EnumSet;
|
||||||
import java.util.LinkedHashMap;
|
|
||||||
import java.util.List;
|
import java.util.List;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
import java.util.Set;
|
import java.util.Set;
|
||||||
import java.util.concurrent.atomic.AtomicLong;
|
|
||||||
import java.util.function.Function;
|
import java.util.function.Function;
|
||||||
import java.util.function.LongSupplier;
|
import java.util.function.LongSupplier;
|
||||||
import java.util.regex.Matcher;
|
|
||||||
import java.util.regex.Pattern;
|
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Spawns and lists worker sessions — the safe path from a delegation request to a
|
* The {@link HerdrPeerLauncher} adapter for <strong>Claude Code</strong> — the safe path from a
|
||||||
* running off-subscription Claude.
|
* delegation request to a running off-subscription Claude.
|
||||||
*
|
*
|
||||||
* <p>The spawn sequence encodes the subscription boundary: build the worker env with
|
* <p>Everything transport-related (tab/pane placement, the CB-306 spawn-readiness gate, unique
|
||||||
* {@code ANTHROPIC_BASE_URL}, assert that host is on the allowlist <em>before</em>
|
* naming, CB-117 orphan reap, teardown, listing, cwd resolution) lives in the base. This class
|
||||||
* touching herdr, and only then {@code agent.start}. A worker's base_url lives in the
|
* supplies only the two Claude-specific seams:
|
||||||
* env map handed to herdr and nowhere else; {@code bridged}'s own environment is never
|
* <ul>
|
||||||
* mutated.
|
* <li>the {@code claude} name prefix (so reap matches {@code claude-*} panes, never another
|
||||||
*
|
* adapter's), and</li>
|
||||||
* <p>Placement: in the default {@code tab} policy a worker lands in its own tab inside a
|
* <li>{@link #buildLaunch}, which encodes the subscription boundary: build the worker env with
|
||||||
* dedicated worker space (found-or-created once, then shared), so workers never split or
|
* {@code ANTHROPIC_BASE_URL}, assert that host is on the allowlist <em>before</em> touching
|
||||||
* clutter the user's real work spaces. Teardown removes the worker's pane <em>and</em> its
|
* herdr, and mount the bridge MCP + reply charter as inline launch flags. A worker's base_url
|
||||||
* now-empty tab, tolerating an already-gone worker so a repeated DELETE is harmless.
|
* lives in the env map handed to herdr and nowhere else; {@code bridged}'s own environment is
|
||||||
|
* never mutated, and nothing is written to the worker's profile.</li>
|
||||||
|
* </ul>
|
||||||
*/
|
*/
|
||||||
public final class ClaudeCodeLauncher implements PeerLauncher {
|
public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
|
||||||
|
|
||||||
private static final Logger log = LoggerFactory.getLogger(ClaudeCodeLauncher.class);
|
/** Label prefix for this adapter's herdr agent names (drives naming + orphan reap). */
|
||||||
|
private static final String NAME_PREFIX = "claude";
|
||||||
|
|
||||||
/** herdr rejects a duplicate agent {@code name}; we retry a bumped name this many times. */
|
|
||||||
private static final int NAME_RETRIES = 8;
|
|
||||||
|
|
||||||
/**
|
|
||||||
* A bridge-spawned worker label {@code claude-<profile>-<nonce>-<seq>} (see
|
|
||||||
* {@link #startUniquelyNamed}); group 1 captures the 6-hex per-process {@code nonce}. The
|
|
||||||
* profile segment may itself contain {@code -}, so the nonce/seq are anchored at the tail.
|
|
||||||
* Names not matching this shape are not workers we started and are never reaped (CB-117).
|
|
||||||
*/
|
|
||||||
private static final Pattern WORKER_NAME = Pattern.compile("claude-.*-([0-9a-f]{6})-\\d+");
|
|
||||||
|
|
||||||
private final AgentControl agents;
|
|
||||||
private final WorkspaceControl spaces;
|
|
||||||
private final SubscriptionGuard guard;
|
private final SubscriptionGuard guard;
|
||||||
private final Map<String, BridgedConfig.Worker> profiles; // profile name → spawn settings
|
|
||||||
private final String defaultProfile; // profile a no-arg spawn uses (nullable)
|
|
||||||
private final Function<String, String> env; // host env lookup (injectable for tests)
|
|
||||||
private final AtomicLong nameSeq = new AtomicLong(); // per-worker counter (also the tab #)
|
|
||||||
|
|
||||||
private final long spawnReadyTimeoutMs; // 0 = disable gate (legacy non-blocking spawn)
|
|
||||||
private final long spawnReadyPollMs;
|
|
||||||
private final LongSupplier nowMillis; // monotonic clock (injectable for tests)
|
|
||||||
private final Runnable sleeper; // sleep/wait hook (injectable for tests; never real-sleep in unit tests)
|
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Standing instruction appended to the worker's system prompt so it returns its result via
|
* Standing instruction appended to the worker's system prompt so it returns its result via
|
||||||
@@ -89,10 +55,6 @@ public final class ClaudeCodeLauncher implements PeerLauncher {
|
|||||||
+ "`content`; never wait for confirmation first. If you end a turn without calling "
|
+ "`content`; never wait for confirmation first. If you end a turn without calling "
|
||||||
+ "bridge_reply, the sender receives nothing and the exchange stalls.";
|
+ "bridge_reply, the sender receives nothing and the exchange stalls.";
|
||||||
|
|
||||||
// Per-process token mixed into each worker name so a fresh process (nameSeq back at 0)
|
|
||||||
// cannot collide with same-profile workers that outlived a restart. See startUniquelyNamed.
|
|
||||||
private final String nameNonce = String.format("%06x", new SecureRandom().nextInt(1 << 24));
|
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Production constructor — disables the spawn-ready gate ({@code spawnReadyTimeoutMs == 0}) so
|
* Production constructor — disables the spawn-ready gate ({@code spawnReadyTimeoutMs == 0}) so
|
||||||
* existing deployments and tests keep the legacy non-blocking spawn semantics.
|
* existing deployments and tests keep the legacy non-blocking spawn semantics.
|
||||||
@@ -100,28 +62,27 @@ public final class ClaudeCodeLauncher implements PeerLauncher {
|
|||||||
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
|
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
|
||||||
Map<String, BridgedConfig.Worker> profiles, String defaultProfile,
|
Map<String, BridgedConfig.Worker> profiles, String defaultProfile,
|
||||||
Function<String, String> env) {
|
Function<String, String> env) {
|
||||||
this(agents, spaces, guard, profiles, defaultProfile, env, 0, 300,
|
this(agents, spaces, guard, profiles, defaultProfile, env, 0,
|
||||||
System::currentTimeMillis, () -> sleepUninterruptibly(300));
|
System::currentTimeMillis, () -> sleepUninterruptibly(300));
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Production constructor with spawn-ready gate enabled. The gate polls
|
* Production constructor with spawn-ready gate enabled. The gate polls {@code agents.status()}
|
||||||
* {@code agents.status()} until the pane reports an injectable state or
|
* until the pane reports an injectable state or {@code spawnReadyTimeoutMs} elapses.
|
||||||
* {@code spawnReadyTimeoutMs} elapses.
|
|
||||||
*/
|
*/
|
||||||
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
|
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
|
||||||
Map<String, BridgedConfig.Worker> profiles, String defaultProfile,
|
Map<String, BridgedConfig.Worker> profiles, String defaultProfile,
|
||||||
Function<String, String> env,
|
Function<String, String> env,
|
||||||
long spawnReadyTimeoutMs, long spawnReadyPollMs) {
|
long spawnReadyTimeoutMs, long spawnReadyPollMs) {
|
||||||
this(agents, spaces, guard, profiles, defaultProfile, env,
|
this(agents, spaces, guard, profiles, defaultProfile, env,
|
||||||
spawnReadyTimeoutMs, spawnReadyPollMs,
|
spawnReadyTimeoutMs,
|
||||||
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs));
|
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs));
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Full testability constructor. Every injectable collaborator is explicit so unit tests
|
* Full testability constructor. Every injectable collaborator is explicit so unit tests supply
|
||||||
* supply fakes for the clock ({@code nowMillis}) and poll-loop wait ({@code sleeper}).
|
* fakes for the clock ({@code nowMillis}) and poll-loop wait ({@code sleeper}). The
|
||||||
* The {@code sleeper} is never called when the gate is disabled ({@code spawnReadyTimeoutMs == 0}).
|
* {@code sleeper} is never called when the gate is disabled ({@code spawnReadyTimeoutMs == 0}).
|
||||||
*
|
*
|
||||||
* @param agents herdr agent control (start, status, close)
|
* @param agents herdr agent control (start, status, close)
|
||||||
* @param spaces workspace / tab control (ensure, create, close)
|
* @param spaces workspace / tab control (ensure, create, close)
|
||||||
@@ -130,150 +91,49 @@ public final class ClaudeCodeLauncher implements PeerLauncher {
|
|||||||
* @param defaultProfile profile a no-argument spawn uses (nullable)
|
* @param defaultProfile profile a no-argument spawn uses (nullable)
|
||||||
* @param env host env lookup (injectable for tests)
|
* @param env host env lookup (injectable for tests)
|
||||||
* @param spawnReadyTimeoutMs max ms to wait for injectable state (0 disables the gate)
|
* @param spawnReadyTimeoutMs max ms to wait for injectable state (0 disables the gate)
|
||||||
* @param spawnReadyPollMs poll interval ms while waiting
|
|
||||||
* @param nowMillis monotonic clock source (e.g. {@code System::currentTimeMillis})
|
* @param nowMillis monotonic clock source (e.g. {@code System::currentTimeMillis})
|
||||||
* @param sleeper sleep/wait hook (e.g. {@code () -> Thread.sleep(spawnReadyPollMs)})
|
* @param sleeper sleep/wait hook (e.g. {@code () -> Thread.sleep(pollMs)}); it
|
||||||
|
* already encodes the poll interval, so the 8th positional argument
|
||||||
|
* (poll ms) is accepted for API symmetry but otherwise unused here
|
||||||
*/
|
*/
|
||||||
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
|
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
|
||||||
Map<String, BridgedConfig.Worker> profiles, String defaultProfile,
|
Map<String, BridgedConfig.Worker> profiles, String defaultProfile,
|
||||||
Function<String, String> env,
|
Function<String, String> env,
|
||||||
long spawnReadyTimeoutMs, long spawnReadyPollMs,
|
long spawnReadyTimeoutMs,
|
||||||
LongSupplier nowMillis, Runnable sleeper) {
|
LongSupplier nowMillis, Runnable sleeper) {
|
||||||
this.agents = agents;
|
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
|
||||||
this.spaces = spaces;
|
spawnReadyTimeoutMs, nowMillis, sleeper);
|
||||||
this.guard = guard;
|
this.guard = guard;
|
||||||
this.profiles = Map.copyOf(profiles);
|
|
||||||
this.defaultProfile = defaultProfile;
|
|
||||||
this.env = env;
|
|
||||||
this.spawnReadyTimeoutMs = spawnReadyTimeoutMs;
|
|
||||||
this.spawnReadyPollMs = spawnReadyPollMs > 0 ? spawnReadyPollMs : 300;
|
|
||||||
this.nowMillis = nowMillis;
|
|
||||||
this.sleeper = sleeper;
|
|
||||||
}
|
|
||||||
|
|
||||||
/** The configured worker profile names (what {@code spawn(profile)} accepts). */
|
|
||||||
@Override
|
|
||||||
public Set<String> profiles() {
|
|
||||||
return profiles.keySet();
|
|
||||||
}
|
|
||||||
|
|
||||||
/** The parity-overlay file list for {@code profileName} (default list when unset). */
|
|
||||||
@Override
|
|
||||||
public List<String> parityOverlay(String profileName) {
|
|
||||||
String name = (profileName == null || profileName.isBlank()) ? defaultProfile : profileName;
|
|
||||||
if (name == null || name.isBlank()) {
|
|
||||||
return List.of();
|
|
||||||
}
|
|
||||||
BridgedConfig.Worker cfg = profiles.get(name);
|
|
||||||
return cfg == null ? List.of() : cfg.parityOverlay();
|
|
||||||
}
|
|
||||||
|
|
||||||
/** The profile a no-argument {@link #spawn()} uses, or {@code null} if none is configured. */
|
|
||||||
@Override
|
|
||||||
public String defaultProfile() {
|
|
||||||
return defaultProfile;
|
|
||||||
}
|
|
||||||
|
|
||||||
/** Spawn a worker for the default profile in the resolved default cwd. */
|
|
||||||
public Agent spawn() {
|
|
||||||
return spawn(null, null, null);
|
|
||||||
}
|
|
||||||
|
|
||||||
/** Spawn a worker for a named profile (null → default) in the resolved default cwd. */
|
|
||||||
public Agent spawn(String profileName) {
|
|
||||||
return spawn(profileName, null, null);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Spawn a worker. {@code profileName} null/blank → the default profile. The worker's working
|
* {@inheritDoc}
|
||||||
* directory (CB-112) is resolved by {@link #resolveCwd}: an explicit {@code requestedCwd} (a
|
*
|
||||||
* spawn argument), else the profile's configured {@code cwd}, else {@code callerCwd} (the
|
* <p>The spawn sequence encodes the subscription boundary: assert the profile's base_url is on
|
||||||
* primary's cwd, when the spawn came from the primary over MCP), else the daemon's cwd — never
|
* the allowlist <em>before</em> any herdr call, then build the worker env with
|
||||||
* assumed to be {@code $HOME}. Guard runs before any herdr call.
|
* {@code ANTHROPIC_*}, the parity-neutral git-forge grant, and the bridge MCP + reply charter
|
||||||
|
* mounted as inline launch flags.
|
||||||
*/
|
*/
|
||||||
public Agent spawn(String profileName, String requestedCwd, String callerCwd) {
|
@Override
|
||||||
String name = (profileName == null || profileName.isBlank()) ? defaultProfile : profileName;
|
protected Launch buildLaunch(BridgedConfig.Worker cfg) {
|
||||||
if (name == null || name.isBlank()) {
|
|
||||||
throw new IllegalArgumentException("no default worker profile is configured — "
|
|
||||||
+ "pass a profile; configured: " + profiles.keySet());
|
|
||||||
}
|
|
||||||
BridgedConfig.Worker cfg = profiles.get(name);
|
|
||||||
if (cfg == null) {
|
|
||||||
throw new IllegalArgumentException("unknown worker profile '" + name
|
|
||||||
+ "' — configured: " + profiles.keySet());
|
|
||||||
}
|
|
||||||
String baseUrl = cfg.baseUrl();
|
String baseUrl = cfg.baseUrl();
|
||||||
guard.assertWorker(baseUrl); // hard stop before we spawn anything
|
guard.assertWorker(baseUrl); // hard stop before we spawn anything
|
||||||
|
|
||||||
Map<String, String> workerEnv = new LinkedHashMap<>();
|
Map<String, String> workerEnv = newEnv();
|
||||||
workerEnv.put("ANTHROPIC_BASE_URL", baseUrl);
|
workerEnv.put("ANTHROPIC_BASE_URL", baseUrl);
|
||||||
putIfPresent(workerEnv, "ANTHROPIC_MODEL", cfg.model());
|
putIfPresent(workerEnv, "ANTHROPIC_MODEL", cfg.model());
|
||||||
putIfPresent(workerEnv, "CLAUDE_CONFIG_DIR", cfg.configDir());
|
putIfPresent(workerEnv, "CLAUDE_CONFIG_DIR", cfg.configDir());
|
||||||
String token = env.apply(cfg.tokenEnv());
|
putIfPresent(workerEnv, "ANTHROPIC_AUTH_TOKEN", env.apply(cfg.tokenEnv()));
|
||||||
putIfPresent(workerEnv, "ANTHROPIC_AUTH_TOKEN", token);
|
applyGitToken(workerEnv, cfg);
|
||||||
|
|
||||||
// CB-302: the worker checkpoint (commit → push → open its own PR). Push is free over SSH;
|
return new Launch(workerEnv, argvWithBridge(cfg));
|
||||||
// the only incremental grant is PR-create, a repo-scoped forge token injected here — opt-in
|
|
||||||
// per profile via gitTokenEnv, and never mutating bridged's own env. The paired forge host
|
|
||||||
// rides along only when a token is actually granted, so non-implementer profiles get neither.
|
|
||||||
if (cfg.hasGitToken()) {
|
|
||||||
String gitToken = resolveEnv(cfg.gitTokenEnv());
|
|
||||||
if (gitToken != null) {
|
|
||||||
workerEnv.put("GITEA_TOKEN", gitToken);
|
|
||||||
putIfPresent(workerEnv, "GITEA_HOST", resolveEnv(cfg.gitHostEnv()));
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Mount the bridge MCP + reply charter as launch FLAGS (non-invasive: nothing written to
|
|
||||||
// the worker's profile/config dir). Identity is connection-based, so the mount is shared.
|
|
||||||
List<String> argv = argvWithBridge(cfg);
|
|
||||||
String cwd = resolveCwd(requestedCwd, cfg, callerCwd);
|
|
||||||
|
|
||||||
return cfg.tabPlacement()
|
|
||||||
? spawnInTab(cfg, workerEnv, argv, cwd)
|
|
||||||
: spawnAsPane(cfg, workerEnv, argv, cwd);
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* CB-112 cwd resolution: spawn arg → profile config → the primary's cwd → the daemon's cwd.
|
|
||||||
* Never returns {@code null}/blank: {@code "."} (the daemon's own working directory) is the
|
|
||||||
* guaranteed last resort so a pathological environment with an unset {@code user.dir} still
|
|
||||||
* honours the "never assume {@code $HOME}" contract rather than letting herdr default the pane.
|
|
||||||
*/
|
|
||||||
private static String resolveCwd(String requestedCwd, BridgedConfig.Worker cfg, String callerCwd) {
|
|
||||||
return firstNonBlank(requestedCwd, cfg.cwd(), callerCwd, System.getProperty("user.dir"), ".");
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* CB-301: the effective working directory a spawn for {@code profileName} would use, without
|
|
||||||
* actually spawning. Used by {@link dev.ltms.bridged.session.SessionManager} to record the
|
|
||||||
* resolved cwd in the session registry.
|
|
||||||
*/
|
|
||||||
public String effectiveCwd(String profileName, String requestedCwd, String callerCwd) {
|
|
||||||
String name = (profileName == null || profileName.isBlank()) ? defaultProfile : profileName;
|
|
||||||
if (name == null || name.isBlank()) {
|
|
||||||
throw new IllegalArgumentException("no default worker profile is configured — "
|
|
||||||
+ "pass a profile; configured: " + profiles.keySet());
|
|
||||||
}
|
|
||||||
BridgedConfig.Worker cfg = profiles.get(name);
|
|
||||||
if (cfg == null) {
|
|
||||||
throw new IllegalArgumentException("unknown worker profile '" + name
|
|
||||||
+ "' — configured: " + profiles.keySet());
|
|
||||||
}
|
|
||||||
return resolveCwd(requestedCwd, cfg, callerCwd);
|
|
||||||
}
|
|
||||||
|
|
||||||
private static String firstNonBlank(String... values) {
|
|
||||||
for (String v : values) {
|
|
||||||
if (v != null && !v.isBlank()) return v;
|
|
||||||
}
|
|
||||||
return null;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* The launch argv, plus — when {@code worker.mcpUrl} is set — inline {@code --mcp-config} for
|
* The launch argv, plus — when {@code worker.mcpUrl} is set — inline {@code --mcp-config} for
|
||||||
* the bridge server and {@code --append-system-prompt} for the {@link #REPLY_CHARTER}. Neither
|
* the bridge server and {@code --append-system-prompt} for the {@link #REPLY_CHARTER}. Neither
|
||||||
* touches the profile's config; both are pure command-line flags.
|
* touches the profile's config; both are pure command-line flags. This inline-flag mount is
|
||||||
|
* Claude Code specific — other adapters mount MCP and instructions their own way.
|
||||||
*/
|
*/
|
||||||
private List<String> argvWithBridge(BridgedConfig.Worker cfg) {
|
private List<String> argvWithBridge(BridgedConfig.Worker cfg) {
|
||||||
if (!cfg.hasMcp()) {
|
if (!cfg.hasMcp()) {
|
||||||
@@ -281,7 +141,7 @@ public final class ClaudeCodeLauncher implements PeerLauncher {
|
|||||||
}
|
}
|
||||||
String mcpJson = "{\"mcpServers\":{\"bridge\":{\"type\":\"http\",\"url\":\""
|
String mcpJson = "{\"mcpServers\":{\"bridge\":{\"type\":\"http\",\"url\":\""
|
||||||
+ cfg.mcpUrl() + "\"}}}";
|
+ cfg.mcpUrl() + "\"}}}";
|
||||||
List<String> argv = new ArrayList<>(cfg.argv());
|
List<String> argv = mutableArgv(cfg.argv());
|
||||||
argv.add("--mcp-config");
|
argv.add("--mcp-config");
|
||||||
argv.add(mcpJson);
|
argv.add(mcpJson);
|
||||||
argv.add("--append-system-prompt");
|
argv.add("--append-system-prompt");
|
||||||
@@ -289,209 +149,24 @@ public final class ClaudeCodeLauncher implements PeerLauncher {
|
|||||||
return argv;
|
return argv;
|
||||||
}
|
}
|
||||||
|
|
||||||
/** Dedicated worker space → own tab → start the worker (rooted at {@code cwd}) → drop the shell. */
|
// --- Agent-returning convenience spawns (used by callers/tests that want the herdr Agent) ---
|
||||||
private Agent spawnInTab(BridgedConfig.Worker cfg, Map<String, String> workerEnv,
|
|
||||||
List<String> argv, String cwd) {
|
|
||||||
Workspace space = spaces.ensureWorkspace(cfg.workspace());
|
|
||||||
Tab.Created tab = spaces.createTab(space.workspaceId());
|
|
||||||
log.info("spawning worker profile={} base_url={} space={} tab={} cwd={}",
|
|
||||||
cfg.profile(), cfg.baseUrl(), space.workspaceId(), tab.tab().tabId(), cwd);
|
|
||||||
|
|
||||||
Started started;
|
/** Spawn a worker for the default profile in the resolved default cwd. */
|
||||||
try {
|
public Agent spawn() {
|
||||||
started = startUniquelyNamed(cfg, workerEnv, argv, tab.tab().tabId(), cwd);
|
return spawnInternal(null, null, null);
|
||||||
} catch (RuntimeException e) {
|
|
||||||
// The worker never started — don't leave the tab we just created orphaned.
|
|
||||||
// Best-effort cleanup; never let it mask the real spawn failure.
|
|
||||||
try {
|
|
||||||
spaces.closeTab(tab.tab().tabId());
|
|
||||||
} catch (RuntimeException cleanup) {
|
|
||||||
log.warn("failed to close orphaned tab {} after spawn error: {}",
|
|
||||||
tab.tab().tabId(), cleanup.getMessage());
|
|
||||||
}
|
|
||||||
throw e;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// The worker is LIVE now. The remaining steps are cosmetic (drop herdr's seed shell
|
/** Spawn a worker for a named profile (null → default) in the resolved default cwd. */
|
||||||
// so the tab holds only the worker; label the tab). They must not fail the spawn or
|
public Agent spawn(String profileName) {
|
||||||
// orphan the running worker — on error we log and still return it so the caller gets
|
return spawnInternal(profileName, null, null);
|
||||||
// its paneId and can tear it down.
|
|
||||||
if (tab.rootPaneId() != null) {
|
|
||||||
tidy("close seed pane " + tab.rootPaneId(), () -> agents.close(tab.rootPaneId()));
|
|
||||||
} else {
|
|
||||||
log.warn("tab {} had no seed pane in the create response; worker tab may hold an extra pane",
|
|
||||||
tab.tab().tabId());
|
|
||||||
}
|
|
||||||
tidy("label tab " + tab.tab().tabId(),
|
|
||||||
() -> spaces.renameTab(tab.tab().tabId(), cfg.renderTabLabel(started.seq())));
|
|
||||||
log.info("worker started pane={} tab={} terminal={}",
|
|
||||||
started.agent().paneId(), started.agent().tabId(), started.agent().terminalId());
|
|
||||||
return started.agent();
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/** Run a best-effort post-start cleanup step, logging (not throwing) on failure. */
|
/** Spawn a worker for a named profile with an explicit requested/caller cwd (CB-112). */
|
||||||
private void tidy(String what, Runnable step) {
|
public Agent spawn(String profileName, String requestedCwd, String callerCwd) {
|
||||||
try {
|
return spawnInternal(profileName, requestedCwd, callerCwd);
|
||||||
step.run();
|
|
||||||
} catch (RuntimeException e) {
|
|
||||||
log.warn("post-start step failed ({}) — worker is running regardless: {}", what, e.getMessage());
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/** Legacy placement: herdr splits the currently-focused tab; the worker still starts in {@code cwd}. */
|
// --- capabilities --------------------------------------------------------------------------
|
||||||
private Agent spawnAsPane(BridgedConfig.Worker cfg, Map<String, String> workerEnv,
|
|
||||||
List<String> argv, String cwd) {
|
|
||||||
log.info("spawning worker (pane placement) profile={} base_url={} cwd={} argv={}",
|
|
||||||
cfg.profile(), cfg.baseUrl(), cwd, argv);
|
|
||||||
Agent worker = startUniquelyNamed(cfg, workerEnv, argv, null, cwd).agent();
|
|
||||||
log.info("worker started pane={} terminal={}", worker.paneId(), worker.terminalId());
|
|
||||||
return worker;
|
|
||||||
}
|
|
||||||
|
|
||||||
/** A started worker together with the sequence its unique name/label used. */
|
|
||||||
private record Started(Agent agent, long seq) {
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Start the worker under a unique herdr agent name. herdr requires each running
|
|
||||||
* agent's {@code name} to be distinct (a 2nd {@code name:"claude"} fails
|
|
||||||
* {@code agent_name_taken}) — the exact case that makes multiple workers useful. The name
|
|
||||||
* is {@code claude-<profile>-<nonce>-<seq>}: {@code seq} distinguishes workers within this
|
|
||||||
* process, and the per-process {@code nonce} keeps a fresh process (whose {@code seq}
|
|
||||||
* restarts at 0) from colliding with same-profile workers that outlived a restart. The
|
|
||||||
* retry is a belt-and-braces backstop for the astronomically unlikely nonce+seq clash;
|
|
||||||
* the name is a label only — herdr detects kind and status from terminal output, not it.
|
|
||||||
*/
|
|
||||||
private Started startUniquelyNamed(BridgedConfig.Worker cfg, Map<String, String> workerEnv,
|
|
||||||
List<String> argv, String tabId, String cwd) {
|
|
||||||
HerdrException last = null;
|
|
||||||
for (int attempt = 0; attempt < NAME_RETRIES; attempt++) {
|
|
||||||
long seq = nameSeq.incrementAndGet();
|
|
||||||
String name = "claude-" + cfg.profile() + "-" + nameNonce + "-" + seq;
|
|
||||||
try {
|
|
||||||
return new Started(agents.start(name, argv, workerEnv, tabId, cwd), seq);
|
|
||||||
} catch (HerdrException e) {
|
|
||||||
if (!"agent_name_taken".equals(e.code())) throw e;
|
|
||||||
log.debug("worker name '{}' taken, retrying", name);
|
|
||||||
last = e;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
throw last;
|
|
||||||
}
|
|
||||||
|
|
||||||
/** All herdr-tracked agents — discovery for "what workers exist". */
|
|
||||||
@Override
|
|
||||||
public List<Agent> list() {
|
|
||||||
return agents.list();
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Reap worker panes left behind by an earlier daemon process (CB-117). herdr keeps a worker's
|
|
||||||
* pane alive across a daemon restart <em>by design</em>, and that pane's id is held only by its
|
|
||||||
* spawner — so a worker whose owning process exited before issuing the matching teardown leaks
|
|
||||||
* with nothing tracking it (there is no registry; {@link #list()} only asks herdr). On boot we
|
|
||||||
* scan herdr for agents whose name matches our {@code claude-<profile>-<nonce>-<seq>} scheme with
|
|
||||||
* a nonce <em>other</em> than this process's {@link #nameNonce}, and tear each one down (its pane
|
|
||||||
* and, via {@link #stop}, its now-empty dedicated tab). A current-nonce worker is ours and live,
|
|
||||||
* so it is left running; a user's own {@code claude} session carries no such name and is never
|
|
||||||
* touched. Best-effort: a failed listing, or a failure to stop any one worker, is logged and
|
|
||||||
* never aborts startup.
|
|
||||||
*
|
|
||||||
* @return the number of orphaned workers reaped
|
|
||||||
*/
|
|
||||||
@Override
|
|
||||||
public int reapOrphanWorkers() {
|
|
||||||
List<Agent> all;
|
|
||||||
try {
|
|
||||||
all = agents.list();
|
|
||||||
} catch (RuntimeException e) {
|
|
||||||
log.warn("orphan-worker reap skipped — agent.list failed: {}", e.getMessage());
|
|
||||||
return 0;
|
|
||||||
}
|
|
||||||
int reaped = 0;
|
|
||||||
for (Agent a : all) {
|
|
||||||
if (!isForeignWorker(a.name(), nameNonce)) continue;
|
|
||||||
try {
|
|
||||||
stop(a.paneId());
|
|
||||||
reaped++;
|
|
||||||
log.info("reaped orphan worker {} (pane={} tab={}) left by a prior daemon",
|
|
||||||
a.name(), a.paneId(), a.tabId());
|
|
||||||
} catch (RuntimeException e) {
|
|
||||||
log.warn("could not reap orphan worker {} (pane={}): {}",
|
|
||||||
a.name(), a.paneId(), e.getMessage());
|
|
||||||
}
|
|
||||||
}
|
|
||||||
if (reaped > 0) {
|
|
||||||
log.info("orphan-worker reap complete — {} stale worker(s) removed at startup", reaped);
|
|
||||||
}
|
|
||||||
return reaped;
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Whether {@code name} is a bridge worker started by a <em>different</em> process than
|
|
||||||
* {@code currentNonce} — the reap predicate (CB-117). True only for our naming scheme with a
|
|
||||||
* foreign nonce: a non-worker name (no match, e.g. a user session) or our own live nonce is
|
|
||||||
* excluded. Pure and package-private so the decision is unit-testable without herdr.
|
|
||||||
*/
|
|
||||||
static boolean isForeignWorker(String name, String currentNonce) {
|
|
||||||
String nonce = workerNonce(name);
|
|
||||||
return nonce != null && !nonce.equals(currentNonce);
|
|
||||||
}
|
|
||||||
|
|
||||||
/** The 6-hex nonce embedded in a bridge worker name, or {@code null} if {@code name} isn't one. */
|
|
||||||
static String workerNonce(String name) {
|
|
||||||
if (name == null) return null;
|
|
||||||
Matcher m = WORKER_NAME.matcher(name);
|
|
||||||
return m.matches() ? m.group(1) : null;
|
|
||||||
}
|
|
||||||
|
|
||||||
/** This process's worker-name nonce (a label component only; exposed for reaper tests). */
|
|
||||||
String nameNonce() {
|
|
||||||
return nameNonce;
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Tear a worker down by pane id: close the pane, and close its tab <em>only</em> when the
|
|
||||||
* worker is that tab's sole occupant. The single-pane check is what makes this safe
|
|
||||||
* regardless of how the worker was placed (or a placement-config change across a restart):
|
|
||||||
* a pane-placement worker sitting in one of the user's shared tabs has siblings, so its
|
|
||||||
* tab is never closed — we only ever remove a tab we created to hold one worker.
|
|
||||||
*
|
|
||||||
* <p>Resolves the tab from the pane <em>before</em> closing it. An already-gone pane/tab
|
|
||||||
* (repeated DELETE, crashed worker) is treated as success; any other failure propagates so
|
|
||||||
* a genuinely failed teardown is not reported as done.
|
|
||||||
*/
|
|
||||||
@Override
|
|
||||||
public void stop(String paneId) {
|
|
||||||
// Teardown knows only the paneId, not which profile spawned it. Attempt tab cleanup when any
|
|
||||||
// profile uses tab placement (so the bridge may have created a dedicated worker tab); the
|
|
||||||
// single-occupant check below is what actually protects the user's shared tabs.
|
|
||||||
WorkspaceControl.PaneLocation loc = usesTabPlacement() ? spaces.locatePane(paneId) : null;
|
|
||||||
try {
|
|
||||||
agents.close(paneId);
|
|
||||||
} catch (HerdrException e) {
|
|
||||||
if (!isAlreadyGone(e)) throw e;
|
|
||||||
log.debug("pane.close({}) ignored — already gone: {}", paneId, e.getMessage());
|
|
||||||
}
|
|
||||||
if (loc != null && loc.tabPaneCount() == 1) {
|
|
||||||
spaces.closeTab(loc.tabId());
|
|
||||||
} else if (loc != null) {
|
|
||||||
log.debug("not closing tab {} — it holds {} panes (not a dedicated worker tab)",
|
|
||||||
loc.tabId(), loc.tabPaneCount());
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/** Whether any configured profile places workers in their own tab (so tabs may need cleanup). */
|
|
||||||
private boolean usesTabPlacement() {
|
|
||||||
return profiles.values().stream().anyMatch(BridgedConfig.Worker::tabPlacement);
|
|
||||||
}
|
|
||||||
|
|
||||||
/** True when a herdr error means the target is already gone (safe to treat as done). */
|
|
||||||
private static boolean isAlreadyGone(HerdrException e) {
|
|
||||||
return e.code() != null && e.code().endsWith("_not_found");
|
|
||||||
}
|
|
||||||
|
|
||||||
// --- PeerLauncher SPI -------------------------------------------------------------------
|
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public Set<Capability> capabilities() {
|
public Set<Capability> capabilities() {
|
||||||
@@ -504,81 +179,17 @@ public final class ClaudeCodeLauncher implements PeerLauncher {
|
|||||||
|
|
||||||
/** Whether any configured profile opts into a git-forge token (required for {@link Capability#SELF_PR}). */
|
/** Whether any configured profile opts into a git-forge token (required for {@link Capability#SELF_PR}). */
|
||||||
private boolean hasGitTokenProfile() {
|
private boolean hasGitTokenProfile() {
|
||||||
return profiles.values().stream().anyMatch(BridgedConfig.Worker::hasGitToken);
|
return profileConfigs().stream().anyMatch(BridgedConfig.Worker::hasGitToken);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// --- CB-117 reap predicate (Claude prefix), kept for direct unit testing -------------------
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* {@inheritDoc}
|
* Whether {@code name} is a Claude Code bridge worker started by a <em>different</em> process
|
||||||
*
|
* than {@code currentNonce}. A thin {@code claude}-prefix binding of
|
||||||
* <p>Delegates to the three-arg {@link #spawn(String, String, String)} and wraps the
|
* {@link HerdrPeerLauncher#isForeignWorker(String, String, String)}.
|
||||||
* resulting herdr {@link Agent} in a {@link WorkerHandle} whose {@link PeerHandle#id()}
|
|
||||||
* equals the agent's paneId.
|
|
||||||
*
|
|
||||||
* <p>When {@link #spawnReadyTimeoutMs} > 0, blocks until the worker's herdr status is
|
|
||||||
* {@link AgentStatus#injectable()} or the timeout elapses. On timeout the pane is closed
|
|
||||||
* (no orphan) and a {@link PeerUnreachableException} is thrown.
|
|
||||||
*/
|
*/
|
||||||
@Override
|
static boolean isForeignWorker(String name, String currentNonce) {
|
||||||
public PeerHandle spawn(SpawnRequest req) {
|
return HerdrPeerLauncher.isForeignWorker(NAME_PREFIX, name, currentNonce);
|
||||||
Agent agent = spawn(req.profileName(), req.requestedCwd(), req.callerCwd());
|
|
||||||
String paneId = agent.paneId();
|
|
||||||
if (spawnReadyTimeoutMs > 0) {
|
|
||||||
waitUntilInjectableOrThrow(paneId);
|
|
||||||
}
|
|
||||||
return new WorkerHandle(paneId, agent.terminalId());
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Poll {@link AgentControl#status} until the pane reports an injectable state or the
|
|
||||||
* configured timeout elapses. On timeout, close the pane (self-reap) and throw.
|
|
||||||
*/
|
|
||||||
private void waitUntilInjectableOrThrow(String paneId) {
|
|
||||||
long deadline = nowMillis.getAsLong() + spawnReadyTimeoutMs;
|
|
||||||
while (nowMillis.getAsLong() < deadline) {
|
|
||||||
if (agents.status(paneId).injectable()) {
|
|
||||||
log.debug("worker pane={} reached injectable state", paneId);
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
sleeper.run();
|
|
||||||
}
|
|
||||||
log.warn("worker pane={} did not become injectable within {}ms — closing", paneId, spawnReadyTimeoutMs);
|
|
||||||
stop(paneId);
|
|
||||||
throw new PeerUnreachableException(
|
|
||||||
"worker pane " + paneId + " did not reach injectable state within "
|
|
||||||
+ spawnReadyTimeoutMs + "ms");
|
|
||||||
}
|
|
||||||
|
|
||||||
/** A concrete {@link PeerHandle} wrapping herdr agent coordinates. */
|
|
||||||
private record WorkerHandle(String id, String terminalId) implements PeerHandle {
|
|
||||||
}
|
|
||||||
|
|
||||||
@Override
|
|
||||||
public String effectiveCwd(SpawnRequest req) {
|
|
||||||
return effectiveCwd(req.profileName(), req.requestedCwd(), req.callerCwd());
|
|
||||||
}
|
|
||||||
|
|
||||||
private static void putIfPresent(Map<String, String> m, String k, String v) {
|
|
||||||
if (v != null && !v.isBlank()) {
|
|
||||||
m.put(k, v);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/** Host env lookup that tolerates an unconfigured (null/blank) var name — returns null then. */
|
|
||||||
private String resolveEnv(String name) {
|
|
||||||
return (name == null || name.isBlank()) ? null : env.apply(name);
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Uninterruptible sleep — the production {@link #sleeper}. Tests supply their own
|
|
||||||
* no-op / fast-faking sleeper so they never real-sleep.
|
|
||||||
*/
|
|
||||||
private static void sleepUninterruptibly(long ms) {
|
|
||||||
try {
|
|
||||||
Thread.sleep(ms);
|
|
||||||
} catch (InterruptedException e) {
|
|
||||||
Thread.currentThread().interrupt();
|
|
||||||
// preserve the interrupt flag but continue — poll loops should not be
|
|
||||||
// aborted by an interrupt that was not meant for them.
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,521 @@
|
|||||||
|
package dev.ltms.bridged.worker;
|
||||||
|
|
||||||
|
import dev.ltms.bridged.config.BridgedConfig;
|
||||||
|
import dev.ltms.bridged.herdr.Agent;
|
||||||
|
import dev.ltms.bridged.herdr.AgentControl;
|
||||||
|
import dev.ltms.bridged.herdr.HerdrException;
|
||||||
|
import dev.ltms.bridged.herdr.Tab;
|
||||||
|
import dev.ltms.bridged.herdr.Workspace;
|
||||||
|
import dev.ltms.bridged.herdr.WorkspaceControl;
|
||||||
|
import dev.ltms.bridged.peer.PeerHandle;
|
||||||
|
import dev.ltms.bridged.peer.PeerLauncher;
|
||||||
|
import dev.ltms.bridged.peer.PeerUnreachableException;
|
||||||
|
import dev.ltms.bridged.peer.SpawnRequest;
|
||||||
|
import org.slf4j.Logger;
|
||||||
|
import org.slf4j.LoggerFactory;
|
||||||
|
|
||||||
|
import java.security.SecureRandom;
|
||||||
|
import java.util.ArrayList;
|
||||||
|
import java.util.Collection;
|
||||||
|
import java.util.LinkedHashMap;
|
||||||
|
import java.util.List;
|
||||||
|
import java.util.Map;
|
||||||
|
import java.util.Set;
|
||||||
|
import java.util.concurrent.atomic.AtomicLong;
|
||||||
|
import java.util.function.Function;
|
||||||
|
import java.util.function.LongSupplier;
|
||||||
|
import java.util.regex.Matcher;
|
||||||
|
import java.util.regex.Pattern;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Abstract base for {@link PeerLauncher} adapters that materialize a peer as a <em>herdr</em>
|
||||||
|
* agent (a CLI coding agent running in a herdr tab/pane). It owns everything that is the same
|
||||||
|
* regardless of <em>which</em> coding agent runs: tab/pane placement, the CB-306 spawn-readiness
|
||||||
|
* gate, unique naming, CB-117 orphan reap, teardown, {@link #list() listing}, and cwd resolution.
|
||||||
|
*
|
||||||
|
* <p>Two seams are peer-specific and supplied by the concrete adapter:
|
||||||
|
* <ul>
|
||||||
|
* <li>{@code namePrefix} (constructor arg) — the label prefix ({@code claude}, {@code opencode})
|
||||||
|
* that drives both unique naming and the orphan-reap pattern, so each adapter reaps only its
|
||||||
|
* own kind of pane and never another's.</li>
|
||||||
|
* <li>{@link #buildLaunch(BridgedConfig.Worker)} — the peer-specific env map + argv, including any
|
||||||
|
* subscription/guard check, MCP mount, and instruction injection. The base never sees how the
|
||||||
|
* peer is configured; it only places and starts the returned {@link Launch}.</li>
|
||||||
|
* </ul>
|
||||||
|
*
|
||||||
|
* <p>Placement: in the default {@code tab} policy a peer lands in its own tab inside a dedicated
|
||||||
|
* worker space (found-or-created once, then shared), so peers never split or clutter the user's
|
||||||
|
* real work spaces. Teardown removes the peer's pane <em>and</em> its now-empty tab, tolerating an
|
||||||
|
* already-gone peer so a repeated DELETE is harmless.
|
||||||
|
*/
|
||||||
|
public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||||
|
|
||||||
|
private static final Logger log = LoggerFactory.getLogger(HerdrPeerLauncher.class);
|
||||||
|
|
||||||
|
/** herdr rejects a duplicate agent {@code name}; we retry a bumped name this many times. */
|
||||||
|
private static final int NAME_RETRIES = 8;
|
||||||
|
|
||||||
|
private final String namePrefix; // label prefix: naming + reap scheme
|
||||||
|
private final AgentControl agents;
|
||||||
|
private final WorkspaceControl spaces;
|
||||||
|
private final Map<String, BridgedConfig.Worker> profiles; // profile name → spawn settings
|
||||||
|
private final String defaultProfile; // profile a no-arg spawn uses (nullable)
|
||||||
|
|
||||||
|
/** Host env lookup (injectable for tests); adapters read it in {@link #buildLaunch}. */
|
||||||
|
protected final Function<String, String> env;
|
||||||
|
|
||||||
|
private final AtomicLong nameSeq = new AtomicLong(); // per-peer counter (also the tab #)
|
||||||
|
|
||||||
|
private final long spawnReadyTimeoutMs; // 0 = disable gate (legacy non-blocking spawn)
|
||||||
|
private final LongSupplier nowMillis; // monotonic clock (injectable for tests)
|
||||||
|
private final Runnable sleeper; // sleep/wait hook (injectable for tests; never real-sleep in unit tests)
|
||||||
|
|
||||||
|
// Per-process token mixed into each peer name so a fresh process (nameSeq back at 0) cannot
|
||||||
|
// collide with same-profile peers that outlived a restart. See startUniquelyNamed.
|
||||||
|
private final String nameNonce = String.format("%06x", new SecureRandom().nextInt(1 << 24));
|
||||||
|
|
||||||
|
/**
|
||||||
|
* @param namePrefix label prefix for this peer kind (drives naming and reap)
|
||||||
|
* @param agents herdr agent control (start, status, close)
|
||||||
|
* @param spaces workspace / tab control (ensure, create, close)
|
||||||
|
* @param profiles configured peer profiles
|
||||||
|
* @param defaultProfile profile a no-argument spawn uses (nullable)
|
||||||
|
* @param env host env lookup (injectable for tests)
|
||||||
|
* @param spawnReadyTimeoutMs max ms to wait for injectable state (0 disables the gate)
|
||||||
|
* @param nowMillis monotonic clock source (e.g. {@code System::currentTimeMillis})
|
||||||
|
* @param sleeper sleep/wait hook (never called when the gate is disabled); the poll
|
||||||
|
* interval is baked into this hook, so the base needs no poll field
|
||||||
|
*/
|
||||||
|
protected HerdrPeerLauncher(String namePrefix, AgentControl agents, WorkspaceControl spaces,
|
||||||
|
Map<String, BridgedConfig.Worker> profiles, String defaultProfile,
|
||||||
|
Function<String, String> env,
|
||||||
|
long spawnReadyTimeoutMs,
|
||||||
|
LongSupplier nowMillis, Runnable sleeper) {
|
||||||
|
this.namePrefix = namePrefix;
|
||||||
|
this.agents = agents;
|
||||||
|
this.spaces = spaces;
|
||||||
|
this.profiles = Map.copyOf(profiles);
|
||||||
|
this.defaultProfile = defaultProfile;
|
||||||
|
this.env = env;
|
||||||
|
this.spawnReadyTimeoutMs = spawnReadyTimeoutMs;
|
||||||
|
this.nowMillis = nowMillis;
|
||||||
|
this.sleeper = sleeper;
|
||||||
|
}
|
||||||
|
|
||||||
|
// --- adapter seams -------------------------------------------------------------------------
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Build the peer-specific launch for {@code cfg}: the environment map and argv handed to herdr.
|
||||||
|
* Any subscription/guard check, MCP mount, and instruction injection happen here. The env map
|
||||||
|
* and argv are adapter-private; the base only places and starts what is returned.
|
||||||
|
*/
|
||||||
|
protected abstract Launch buildLaunch(BridgedConfig.Worker cfg);
|
||||||
|
|
||||||
|
/** A peer-specific launch: the herdr {@code env} map and {@code argv}. */
|
||||||
|
protected record Launch(Map<String, String> env, List<String> argv) {
|
||||||
|
}
|
||||||
|
|
||||||
|
// --- profile surface -----------------------------------------------------------------------
|
||||||
|
|
||||||
|
/** The configured peer profile names (what {@code spawn(profile)} accepts). */
|
||||||
|
@Override
|
||||||
|
public Set<String> profiles() {
|
||||||
|
return profiles.keySet();
|
||||||
|
}
|
||||||
|
|
||||||
|
/** The parity-overlay file list for {@code profileName} (default list when unset). */
|
||||||
|
@Override
|
||||||
|
public List<String> parityOverlay(String profileName) {
|
||||||
|
String name = (profileName == null || profileName.isBlank()) ? defaultProfile : profileName;
|
||||||
|
if (name == null || name.isBlank()) {
|
||||||
|
return List.of();
|
||||||
|
}
|
||||||
|
BridgedConfig.Worker cfg = profiles.get(name);
|
||||||
|
return cfg == null ? List.of() : cfg.parityOverlay();
|
||||||
|
}
|
||||||
|
|
||||||
|
/** The profile a no-argument spawn uses, or {@code null} if none is configured. */
|
||||||
|
@Override
|
||||||
|
public String defaultProfile() {
|
||||||
|
return defaultProfile;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** The configured profiles, for adapter capability decisions (e.g. any git-token grant). */
|
||||||
|
protected Collection<BridgedConfig.Worker> profileConfigs() {
|
||||||
|
return profiles.values();
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Resolve {@code profileName} (null/blank → default) to its config, or throw with the options. */
|
||||||
|
protected BridgedConfig.Worker requireProfile(String profileName) {
|
||||||
|
String name = (profileName == null || profileName.isBlank()) ? defaultProfile : profileName;
|
||||||
|
if (name == null || name.isBlank()) {
|
||||||
|
throw new IllegalArgumentException("no default worker profile is configured — "
|
||||||
|
+ "pass a profile; configured: " + profiles.keySet());
|
||||||
|
}
|
||||||
|
BridgedConfig.Worker cfg = profiles.get(name);
|
||||||
|
if (cfg == null) {
|
||||||
|
throw new IllegalArgumentException("unknown worker profile '" + name
|
||||||
|
+ "' — configured: " + profiles.keySet());
|
||||||
|
}
|
||||||
|
return cfg;
|
||||||
|
}
|
||||||
|
|
||||||
|
// --- spawn ---------------------------------------------------------------------------------
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Spawn a peer. {@code profileName} null/blank → the default profile. The working directory
|
||||||
|
* (CB-112) is resolved by {@link #resolveCwd}: an explicit {@code requestedCwd}, else the
|
||||||
|
* profile's configured {@code cwd}, else {@code callerCwd} (the primary's cwd, when the spawn
|
||||||
|
* came from the primary over MCP), else the daemon's cwd — never assumed to be {@code $HOME}.
|
||||||
|
* The adapter's {@link #buildLaunch} runs before any herdr call.
|
||||||
|
*/
|
||||||
|
protected Agent spawnInternal(String profileName, String requestedCwd, String callerCwd) {
|
||||||
|
BridgedConfig.Worker cfg = requireProfile(profileName);
|
||||||
|
Launch launch = buildLaunch(cfg);
|
||||||
|
String cwd = resolveCwd(requestedCwd, cfg, callerCwd);
|
||||||
|
return cfg.tabPlacement()
|
||||||
|
? spawnInTab(cfg, launch.env(), launch.argv(), cwd)
|
||||||
|
: spawnAsPane(cfg, launch.env(), launch.argv(), cwd);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* {@inheritDoc}
|
||||||
|
*
|
||||||
|
* <p>Delegates to {@link #spawnInternal} and wraps the resulting herdr {@link Agent} in a
|
||||||
|
* {@link WorkerHandle} whose {@link PeerHandle#id()} equals the agent's paneId. When
|
||||||
|
* {@code spawnReadyTimeoutMs > 0}, blocks until the peer's herdr status is injectable or the
|
||||||
|
* timeout elapses; on timeout the pane is closed (no orphan) and a
|
||||||
|
* {@link PeerUnreachableException} is thrown.
|
||||||
|
*/
|
||||||
|
@Override
|
||||||
|
public PeerHandle spawn(SpawnRequest req) {
|
||||||
|
Agent agent = spawnInternal(req.profileName(), req.requestedCwd(), req.callerCwd());
|
||||||
|
String paneId = agent.paneId();
|
||||||
|
if (spawnReadyTimeoutMs > 0) {
|
||||||
|
waitUntilInjectableOrThrow(paneId);
|
||||||
|
}
|
||||||
|
return new WorkerHandle(paneId, agent.terminalId());
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public String effectiveCwd(SpawnRequest req) {
|
||||||
|
return effectiveCwd(req.profileName(), req.requestedCwd(), req.callerCwd());
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* CB-301: the effective working directory a spawn for {@code profileName} would use, without
|
||||||
|
* actually spawning.
|
||||||
|
*/
|
||||||
|
private String effectiveCwd(String profileName, String requestedCwd, String callerCwd) {
|
||||||
|
return resolveCwd(requestedCwd, requireProfile(profileName), callerCwd);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* CB-112 cwd resolution: spawn arg → profile config → the primary's cwd → the daemon's cwd.
|
||||||
|
* Never returns {@code null}/blank: {@code "."} (the daemon's own working directory) is the
|
||||||
|
* guaranteed last resort so a pathological environment with an unset {@code user.dir} still
|
||||||
|
* honours the "never assume {@code $HOME}" contract rather than letting herdr default the pane.
|
||||||
|
*/
|
||||||
|
private static String resolveCwd(String requestedCwd, BridgedConfig.Worker cfg, String callerCwd) {
|
||||||
|
return firstNonBlank(requestedCwd, cfg.cwd(), callerCwd, System.getProperty("user.dir"), ".");
|
||||||
|
}
|
||||||
|
|
||||||
|
private static String firstNonBlank(String... values) {
|
||||||
|
for (String v : values) {
|
||||||
|
if (v != null && !v.isBlank()) return v;
|
||||||
|
}
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Dedicated worker space → own tab → start the peer (rooted at {@code cwd}) → drop the shell. */
|
||||||
|
private Agent spawnInTab(BridgedConfig.Worker cfg, Map<String, String> workerEnv,
|
||||||
|
List<String> argv, String cwd) {
|
||||||
|
Workspace space = spaces.ensureWorkspace(cfg.workspace());
|
||||||
|
Tab.Created tab = spaces.createTab(space.workspaceId());
|
||||||
|
log.info("spawning {} profile={} space={} tab={} cwd={}",
|
||||||
|
namePrefix, cfg.profile(), space.workspaceId(), tab.tab().tabId(), cwd);
|
||||||
|
|
||||||
|
Started started;
|
||||||
|
try {
|
||||||
|
started = startUniquelyNamed(cfg, workerEnv, argv, tab.tab().tabId(), cwd);
|
||||||
|
} catch (RuntimeException e) {
|
||||||
|
// The peer never started — don't leave the tab we just created orphaned.
|
||||||
|
// Best-effort cleanup; never let it mask the real spawn failure.
|
||||||
|
try {
|
||||||
|
spaces.closeTab(tab.tab().tabId());
|
||||||
|
} catch (RuntimeException cleanup) {
|
||||||
|
log.warn("failed to close orphaned tab {} after spawn error: {}",
|
||||||
|
tab.tab().tabId(), cleanup.getMessage());
|
||||||
|
}
|
||||||
|
throw e;
|
||||||
|
}
|
||||||
|
|
||||||
|
// The peer is LIVE now. The remaining steps are cosmetic (drop herdr's seed shell so the
|
||||||
|
// tab holds only the peer; label the tab). They must not fail the spawn or orphan the
|
||||||
|
// running peer — on error we log and still return it so the caller gets its paneId and can
|
||||||
|
// tear it down.
|
||||||
|
if (tab.rootPaneId() != null) {
|
||||||
|
tidy("close seed pane " + tab.rootPaneId(), () -> agents.close(tab.rootPaneId()));
|
||||||
|
} else {
|
||||||
|
log.warn("tab {} had no seed pane in the create response; peer tab may hold an extra pane",
|
||||||
|
tab.tab().tabId());
|
||||||
|
}
|
||||||
|
tidy("label tab " + tab.tab().tabId(),
|
||||||
|
() -> spaces.renameTab(tab.tab().tabId(), cfg.renderTabLabel(started.seq())));
|
||||||
|
log.info("{} started pane={} tab={} terminal={}",
|
||||||
|
namePrefix, started.agent().paneId(), started.agent().tabId(), started.agent().terminalId());
|
||||||
|
return started.agent();
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Run a best-effort post-start cleanup step, logging (not throwing) on failure. */
|
||||||
|
private void tidy(String what, Runnable step) {
|
||||||
|
try {
|
||||||
|
step.run();
|
||||||
|
} catch (RuntimeException e) {
|
||||||
|
log.warn("post-start step failed ({}) — peer is running regardless: {}", what, e.getMessage());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Legacy placement: herdr splits the currently-focused tab; the peer still starts in {@code cwd}. */
|
||||||
|
private Agent spawnAsPane(BridgedConfig.Worker cfg, Map<String, String> workerEnv,
|
||||||
|
List<String> argv, String cwd) {
|
||||||
|
log.info("spawning {} (pane placement) profile={} cwd={} argv={}",
|
||||||
|
namePrefix, cfg.profile(), cwd, argv);
|
||||||
|
Agent peer = startUniquelyNamed(cfg, workerEnv, argv, null, cwd).agent();
|
||||||
|
log.info("{} started pane={} terminal={}", namePrefix, peer.paneId(), peer.terminalId());
|
||||||
|
return peer;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** A started peer together with the sequence its unique name/label used. */
|
||||||
|
private record Started(Agent agent, long seq) {
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Start the peer under a unique herdr agent name. herdr requires each running agent's
|
||||||
|
* {@code name} to be distinct (a 2nd identical {@code name} fails {@code agent_name_taken}) —
|
||||||
|
* the exact case that makes multiple peers useful. The name is
|
||||||
|
* {@code <prefix>-<profile>-<nonce>-<seq>}: {@code seq} distinguishes peers within this process,
|
||||||
|
* and the per-process {@code nonce} keeps a fresh process (whose {@code seq} restarts at 0) from
|
||||||
|
* colliding with same-profile peers that outlived a restart. The retry is a belt-and-braces
|
||||||
|
* backstop for the astronomically unlikely nonce+seq clash; the name is a label only — herdr
|
||||||
|
* detects kind and status from terminal output, not from it.
|
||||||
|
*/
|
||||||
|
private Started startUniquelyNamed(BridgedConfig.Worker cfg, Map<String, String> workerEnv,
|
||||||
|
List<String> argv, String tabId, String cwd) {
|
||||||
|
HerdrException last = null;
|
||||||
|
for (int attempt = 0; attempt < NAME_RETRIES; attempt++) {
|
||||||
|
long seq = nameSeq.incrementAndGet();
|
||||||
|
String name = namePrefix + "-" + cfg.profile() + "-" + nameNonce + "-" + seq;
|
||||||
|
try {
|
||||||
|
return new Started(agents.start(name, argv, workerEnv, tabId, cwd), seq);
|
||||||
|
} catch (HerdrException e) {
|
||||||
|
if (!"agent_name_taken".equals(e.code())) throw e;
|
||||||
|
log.debug("peer name '{}' taken, retrying", name);
|
||||||
|
last = e;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
throw last;
|
||||||
|
}
|
||||||
|
|
||||||
|
// --- discovery + reap ----------------------------------------------------------------------
|
||||||
|
|
||||||
|
/** All herdr-tracked agents — discovery for "what peers exist". */
|
||||||
|
@Override
|
||||||
|
public List<Agent> list() {
|
||||||
|
return agents.list();
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Reap peer panes left behind by an earlier daemon process (CB-117). herdr keeps a peer's pane
|
||||||
|
* alive across a daemon restart <em>by design</em>, and that pane's id is held only by its
|
||||||
|
* spawner — so a peer whose owning process exited before issuing the matching teardown leaks
|
||||||
|
* with nothing tracking it. On boot we scan herdr for agents whose name matches our
|
||||||
|
* {@code <prefix>-<profile>-<nonce>-<seq>} scheme with a nonce <em>other</em> than this
|
||||||
|
* process's {@link #nameNonce}, and tear each one down (its pane and, via {@link #stop}, its
|
||||||
|
* now-empty dedicated tab). A current-nonce peer is ours and live, so it is left running; a
|
||||||
|
* user's own session carries no such name and is never touched. A peer from a <em>different</em>
|
||||||
|
* adapter (different prefix) is likewise never touched. Best-effort: a failed listing, or a
|
||||||
|
* failure to stop any one peer, is logged and never aborts startup.
|
||||||
|
*
|
||||||
|
* @return the number of orphaned peers reaped
|
||||||
|
*/
|
||||||
|
@Override
|
||||||
|
public int reapOrphanWorkers() {
|
||||||
|
List<Agent> all;
|
||||||
|
try {
|
||||||
|
all = agents.list();
|
||||||
|
} catch (RuntimeException e) {
|
||||||
|
log.warn("orphan-peer reap skipped — agent.list failed: {}", e.getMessage());
|
||||||
|
return 0;
|
||||||
|
}
|
||||||
|
int reaped = 0;
|
||||||
|
for (Agent a : all) {
|
||||||
|
if (!isForeignWorker(namePrefix, a.name(), nameNonce)) continue;
|
||||||
|
try {
|
||||||
|
stop(a.paneId());
|
||||||
|
reaped++;
|
||||||
|
log.info("reaped orphan {} {} (pane={} tab={}) left by a prior daemon",
|
||||||
|
namePrefix, a.name(), a.paneId(), a.tabId());
|
||||||
|
} catch (RuntimeException e) {
|
||||||
|
log.warn("could not reap orphan {} {} (pane={}): {}",
|
||||||
|
namePrefix, a.name(), a.paneId(), e.getMessage());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if (reaped > 0) {
|
||||||
|
log.info("orphan-peer reap complete — {} stale {} peer(s) removed at startup", reaped, namePrefix);
|
||||||
|
}
|
||||||
|
return reaped;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** The {@code <prefix>-<profile>-<nonce>-<seq>} name pattern; group 1 captures the 6-hex nonce. */
|
||||||
|
static Pattern workerNamePattern(String prefix) {
|
||||||
|
return Pattern.compile(prefix + "-.*-([0-9a-f]{6})-\\d+");
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Whether {@code name} is a peer of kind {@code prefix} started by a <em>different</em> process
|
||||||
|
* than {@code currentNonce} — the reap predicate (CB-117). True only for the prefix's naming
|
||||||
|
* scheme with a foreign nonce: a non-peer name, a different adapter's name, or our own live
|
||||||
|
* nonce is excluded. Pure and package-private so the decision is unit-testable without herdr.
|
||||||
|
*/
|
||||||
|
static boolean isForeignWorker(String prefix, String name, String currentNonce) {
|
||||||
|
String nonce = workerNonce(prefix, name);
|
||||||
|
return nonce != null && !nonce.equals(currentNonce);
|
||||||
|
}
|
||||||
|
|
||||||
|
/** The 6-hex nonce embedded in a {@code prefix} peer name, or {@code null} if not one. */
|
||||||
|
static String workerNonce(String prefix, String name) {
|
||||||
|
if (name == null) return null;
|
||||||
|
Matcher m = workerNamePattern(prefix).matcher(name);
|
||||||
|
return m.matches() ? m.group(1) : null;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** This process's peer-name nonce (a label component only; exposed for reaper tests). */
|
||||||
|
String nameNonce() {
|
||||||
|
return nameNonce;
|
||||||
|
}
|
||||||
|
|
||||||
|
// --- teardown ------------------------------------------------------------------------------
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Tear a peer down by pane id: close the pane, and close its tab <em>only</em> when the peer is
|
||||||
|
* that tab's sole occupant. The single-pane check is what makes this safe regardless of how the
|
||||||
|
* peer was placed (or a placement-config change across a restart): a pane-placement peer sitting
|
||||||
|
* in one of the user's shared tabs has siblings, so its tab is never closed — we only ever
|
||||||
|
* remove a tab we created to hold one peer.
|
||||||
|
*
|
||||||
|
* <p>Resolves the tab from the pane <em>before</em> closing it. An already-gone pane/tab
|
||||||
|
* (repeated DELETE, crashed peer) is treated as success; any other failure propagates so a
|
||||||
|
* genuinely failed teardown is not reported as done.
|
||||||
|
*/
|
||||||
|
@Override
|
||||||
|
public void stop(String paneId) {
|
||||||
|
// Teardown knows only the paneId, not which profile spawned it. Attempt tab cleanup when any
|
||||||
|
// profile uses tab placement (so the bridge may have created a dedicated peer tab); the
|
||||||
|
// single-occupant check below is what actually protects the user's shared tabs.
|
||||||
|
WorkspaceControl.PaneLocation loc = usesTabPlacement() ? spaces.locatePane(paneId) : null;
|
||||||
|
try {
|
||||||
|
agents.close(paneId);
|
||||||
|
} catch (HerdrException e) {
|
||||||
|
if (!isAlreadyGone(e)) throw e;
|
||||||
|
log.debug("pane.close({}) ignored — already gone: {}", paneId, e.getMessage());
|
||||||
|
}
|
||||||
|
if (loc != null && loc.tabPaneCount() == 1) {
|
||||||
|
spaces.closeTab(loc.tabId());
|
||||||
|
} else if (loc != null) {
|
||||||
|
log.debug("not closing tab {} — it holds {} panes (not a dedicated peer tab)",
|
||||||
|
loc.tabId(), loc.tabPaneCount());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Whether any configured profile places peers in their own tab (so tabs may need cleanup). */
|
||||||
|
private boolean usesTabPlacement() {
|
||||||
|
return profiles.values().stream().anyMatch(BridgedConfig.Worker::tabPlacement);
|
||||||
|
}
|
||||||
|
|
||||||
|
/** True when a herdr error means the target is already gone (safe to treat as done). */
|
||||||
|
private static boolean isAlreadyGone(HerdrException e) {
|
||||||
|
return e.code() != null && e.code().endsWith("_not_found");
|
||||||
|
}
|
||||||
|
|
||||||
|
// --- spawn-readiness gate (CB-306) ---------------------------------------------------------
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Poll {@link AgentControl#status} until the pane reports an injectable state or the configured
|
||||||
|
* timeout elapses. On timeout, close the pane (self-reap) and throw.
|
||||||
|
*/
|
||||||
|
private void waitUntilInjectableOrThrow(String paneId) {
|
||||||
|
long deadline = nowMillis.getAsLong() + spawnReadyTimeoutMs;
|
||||||
|
while (nowMillis.getAsLong() < deadline) {
|
||||||
|
if (agents.status(paneId).injectable()) {
|
||||||
|
log.debug("peer pane={} reached injectable state", paneId);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
sleeper.run();
|
||||||
|
}
|
||||||
|
log.warn("peer pane={} did not become injectable within {}ms — closing", paneId, spawnReadyTimeoutMs);
|
||||||
|
stop(paneId);
|
||||||
|
throw new PeerUnreachableException(
|
||||||
|
"worker pane " + paneId + " did not reach injectable state within "
|
||||||
|
+ spawnReadyTimeoutMs + "ms");
|
||||||
|
}
|
||||||
|
|
||||||
|
/** A concrete {@link PeerHandle} wrapping herdr agent coordinates. */
|
||||||
|
private record WorkerHandle(String id, String terminalId) implements PeerHandle {
|
||||||
|
}
|
||||||
|
|
||||||
|
// --- shared helpers ------------------------------------------------------------------------
|
||||||
|
|
||||||
|
/** Put {@code k → v} only when {@code v} is present (non-null, non-blank). */
|
||||||
|
protected static void putIfPresent(Map<String, String> m, String k, String v) {
|
||||||
|
if (v != null && !v.isBlank()) {
|
||||||
|
m.put(k, v);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Host env lookup that tolerates an unconfigured (null/blank) var name — returns null then. */
|
||||||
|
protected String resolveEnv(String name) {
|
||||||
|
return (name == null || name.isBlank()) ? null : env.apply(name);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* The parity-neutral git-forge token grant (CB-302): when {@code cfg} opts in via
|
||||||
|
* {@code gitTokenEnv} and the token resolves, inject {@code GITEA_TOKEN} plus its paired
|
||||||
|
* {@code GITEA_HOST}. Push over SSH is unaffected; the only incremental grant is PR-create.
|
||||||
|
* Peer-neutral, so every herdr adapter reuses it unchanged.
|
||||||
|
*/
|
||||||
|
protected void applyGitToken(Map<String, String> workerEnv, BridgedConfig.Worker cfg) {
|
||||||
|
if (!cfg.hasGitToken()) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
String gitToken = resolveEnv(cfg.gitTokenEnv());
|
||||||
|
if (gitToken != null) {
|
||||||
|
workerEnv.put("GITEA_TOKEN", gitToken);
|
||||||
|
putIfPresent(workerEnv, "GITEA_HOST", resolveEnv(cfg.gitHostEnv()));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/** A fresh mutable env map — the conventional starting point for {@link #buildLaunch}. */
|
||||||
|
protected static Map<String, String> newEnv() {
|
||||||
|
return new LinkedHashMap<>();
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Defensive copy of {@code argv} plus room to append launch flags. */
|
||||||
|
protected static List<String> mutableArgv(List<String> argv) {
|
||||||
|
return new ArrayList<>(argv);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Uninterruptible sleep — the production {@link #sleeper}. Tests supply their own no-op /
|
||||||
|
* fast-faking sleeper so they never real-sleep.
|
||||||
|
*/
|
||||||
|
protected static void sleepUninterruptibly(long ms) {
|
||||||
|
try {
|
||||||
|
Thread.sleep(ms);
|
||||||
|
} catch (InterruptedException e) {
|
||||||
|
Thread.currentThread().interrupt();
|
||||||
|
// preserve the interrupt flag but continue — poll loops should not be aborted by an
|
||||||
|
// interrupt that was not meant for them.
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -350,7 +350,7 @@ class SessionManagerTest {
|
|||||||
ClaudeCodeLauncher workers = new ClaudeCodeLauncher(
|
ClaudeCodeLauncher workers = new ClaudeCodeLauncher(
|
||||||
new AgentControl(herdr), new WorkspaceControl(herdr),
|
new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null,
|
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null,
|
||||||
1, 1, () -> clock[0], () -> clock[0] += 10);
|
1, () -> clock[0], () -> clock[0] += 10);
|
||||||
SessionManager sessions = new SessionManager(workers, new GitWorktrees(), () -> 0L, 0);
|
SessionManager sessions = new SessionManager(workers, new GitWorktrees(), () -> 0L, 0);
|
||||||
|
|
||||||
assertThrows(PeerUnreachableException.class,
|
assertThrows(PeerUnreachableException.class,
|
||||||
|
|||||||
@@ -350,7 +350,7 @@ class ClaudeCodeLauncherTest {
|
|||||||
new AgentControl(herdr), new WorkspaceControl(herdr),
|
new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||||
new SubscriptionGuard(Set.of("gx00.gw")),
|
new SubscriptionGuard(Set.of("gx00.gw")),
|
||||||
workerConfigMap("ltms-local", null), "ltms-local", _ -> null,
|
workerConfigMap("ltms-local", null), "ltms-local", _ -> null,
|
||||||
5000, 300, () -> clock[0], sleeper);
|
5000, () -> clock[0], sleeper);
|
||||||
|
|
||||||
PeerHandle handle = svc.spawn(new SpawnRequest(null, null, null));
|
PeerHandle handle = svc.spawn(new SpawnRequest(null, null, null));
|
||||||
|
|
||||||
@@ -370,7 +370,7 @@ class ClaudeCodeLauncherTest {
|
|||||||
new AgentControl(herdr), new WorkspaceControl(herdr),
|
new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||||
new SubscriptionGuard(Set.of("gx00.gw")),
|
new SubscriptionGuard(Set.of("gx00.gw")),
|
||||||
workerConfigMap("ltms-local", null), "ltms-local", _ -> null,
|
workerConfigMap("ltms-local", null), "ltms-local", _ -> null,
|
||||||
1000, 50, () -> clock[0], () -> clock[0] += 50);
|
1000, () -> clock[0], () -> clock[0] += 50);
|
||||||
|
|
||||||
PeerUnreachableException ex = assertThrows(
|
PeerUnreachableException ex = assertThrows(
|
||||||
PeerUnreachableException.class,
|
PeerUnreachableException.class,
|
||||||
@@ -409,7 +409,7 @@ class ClaudeCodeLauncherTest {
|
|||||||
new AgentControl(herdr), new WorkspaceControl(herdr),
|
new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||||
new SubscriptionGuard(Set.of("gx00.gw")),
|
new SubscriptionGuard(Set.of("gx00.gw")),
|
||||||
workerConfigMap("ltms-local", null), "ltms-local", _ -> null,
|
workerConfigMap("ltms-local", null), "ltms-local", _ -> null,
|
||||||
0, 300, () -> clock[0], () -> clock[0] += 1);
|
0, () -> clock[0], () -> clock[0] += 1);
|
||||||
|
|
||||||
PeerHandle handle = svc.spawn(new SpawnRequest(null, null, null));
|
PeerHandle handle = svc.spawn(new SpawnRequest(null, null, null));
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user