From ffce30afa241cb39b2204a7a7bfb9cba49338ee9 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Wed, 22 Jul 2026 05:15:21 +0200 Subject: [PATCH] CB-402 Increment 1: extract HerdrPeerLauncher abstract base MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- .../bridged/worker/ClaudeCodeLauncher.java | 517 +++-------------- .../bridged/worker/HerdrPeerLauncher.java | 521 ++++++++++++++++++ .../bridged/session/SessionManagerTest.java | 2 +- .../worker/ClaudeCodeLauncherTest.java | 6 +- 4 files changed, 589 insertions(+), 457 deletions(-) create mode 100644 bridged/src/main/java/dev/ltms/bridged/worker/HerdrPeerLauncher.java diff --git a/bridged/src/main/java/dev/ltms/bridged/worker/ClaudeCodeLauncher.java b/bridged/src/main/java/dev/ltms/bridged/worker/ClaudeCodeLauncher.java index 0428823..eecb8a3 100644 --- a/bridged/src/main/java/dev/ltms/bridged/worker/ClaudeCodeLauncher.java +++ b/bridged/src/main/java/dev/ltms/bridged/worker/ClaudeCodeLauncher.java @@ -4,73 +4,39 @@ import dev.ltms.bridged.config.BridgedConfig; import dev.ltms.bridged.guard.SubscriptionGuard; 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.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.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; /** - * Spawns and lists worker sessions — the safe path from a delegation request to a - * running off-subscription Claude. + * The {@link HerdrPeerLauncher} adapter for Claude Code — the safe path from a + * delegation request to a running off-subscription Claude. * - *

The spawn sequence encodes the subscription boundary: build the worker env with - * {@code ANTHROPIC_BASE_URL}, assert that host is on the allowlist before - * touching herdr, and only then {@code agent.start}. A worker's base_url lives in the - * env map handed to herdr and nowhere else; {@code bridged}'s own environment is never - * mutated. - * - *

Placement: in the default {@code tab} policy a worker lands in its own tab inside a - * dedicated worker space (found-or-created once, then shared), so workers never split or - * clutter the user's real work spaces. Teardown removes the worker's pane and its - * now-empty tab, tolerating an already-gone worker so a repeated DELETE is harmless. + *

Everything transport-related (tab/pane placement, the CB-306 spawn-readiness gate, unique + * naming, CB-117 orphan reap, teardown, listing, cwd resolution) lives in the base. This class + * supplies only the two Claude-specific seams: + *

*/ -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---} (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 Map profiles; // profile name → spawn settings - private final String defaultProfile; // profile a no-arg spawn uses (nullable) - private final Function 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 @@ -89,10 +55,6 @@ public final class ClaudeCodeLauncher implements PeerLauncher { + "`content`; never wait for confirmation first. If you end a turn without calling " + "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 * 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, Map profiles, String defaultProfile, Function env) { - this(agents, spaces, guard, profiles, defaultProfile, env, 0, 300, + this(agents, spaces, guard, profiles, defaultProfile, env, 0, System::currentTimeMillis, () -> sleepUninterruptibly(300)); } /** - * Production constructor with spawn-ready gate enabled. The gate polls - * {@code agents.status()} until the pane reports an injectable state or - * {@code spawnReadyTimeoutMs} elapses. + * Production constructor with spawn-ready gate enabled. The gate polls {@code agents.status()} + * until the pane reports an injectable state or {@code spawnReadyTimeoutMs} elapses. */ public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard, Map profiles, String defaultProfile, Function env, long spawnReadyTimeoutMs, long spawnReadyPollMs) { this(agents, spaces, guard, profiles, defaultProfile, env, - spawnReadyTimeoutMs, spawnReadyPollMs, + spawnReadyTimeoutMs, System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs)); } /** - * Full testability constructor. Every injectable collaborator is explicit so unit tests - * supply fakes for the clock ({@code nowMillis}) and poll-loop wait ({@code sleeper}). - * The {@code sleeper} is never called when the gate is disabled ({@code spawnReadyTimeoutMs == 0}). + * Full testability constructor. Every injectable collaborator is explicit so unit tests supply + * fakes for the clock ({@code nowMillis}) and poll-loop wait ({@code sleeper}). The + * {@code sleeper} is never called when the gate is disabled ({@code spawnReadyTimeoutMs == 0}). * * @param agents herdr agent control (start, status, 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 env host env lookup (injectable for tests) * @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 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, Map profiles, String defaultProfile, Function env, - long spawnReadyTimeoutMs, long spawnReadyPollMs, + long spawnReadyTimeoutMs, LongSupplier nowMillis, Runnable sleeper) { - this.agents = agents; - this.spaces = spaces; + super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env, + spawnReadyTimeoutMs, nowMillis, sleeper); 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 profiles() { - return profiles.keySet(); - } - - /** The parity-overlay file list for {@code profileName} (default list when unset). */ - @Override - public List 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 - * 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 - * primary's cwd, when the spawn came from the primary over MCP), else the daemon's cwd — never - * assumed to be {@code $HOME}. Guard runs before any herdr call. + * {@inheritDoc} + * + *

The spawn sequence encodes the subscription boundary: assert the profile's base_url is on + * the allowlist before any herdr call, then build the worker env with + * {@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) { - 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()); - } + @Override + protected Launch buildLaunch(BridgedConfig.Worker cfg) { String baseUrl = cfg.baseUrl(); guard.assertWorker(baseUrl); // hard stop before we spawn anything - Map workerEnv = new LinkedHashMap<>(); + Map workerEnv = newEnv(); workerEnv.put("ANTHROPIC_BASE_URL", baseUrl); putIfPresent(workerEnv, "ANTHROPIC_MODEL", cfg.model()); putIfPresent(workerEnv, "CLAUDE_CONFIG_DIR", cfg.configDir()); - String token = env.apply(cfg.tokenEnv()); - putIfPresent(workerEnv, "ANTHROPIC_AUTH_TOKEN", token); + putIfPresent(workerEnv, "ANTHROPIC_AUTH_TOKEN", env.apply(cfg.tokenEnv())); + applyGitToken(workerEnv, cfg); - // CB-302: the worker checkpoint (commit → push → open its own PR). Push is free over SSH; - // 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 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; + return new Launch(workerEnv, argvWithBridge(cfg)); } /** * 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 - * 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 argvWithBridge(BridgedConfig.Worker cfg) { if (!cfg.hasMcp()) { @@ -281,7 +141,7 @@ public final class ClaudeCodeLauncher implements PeerLauncher { } String mcpJson = "{\"mcpServers\":{\"bridge\":{\"type\":\"http\",\"url\":\"" + cfg.mcpUrl() + "\"}}}"; - List argv = new ArrayList<>(cfg.argv()); + List argv = mutableArgv(cfg.argv()); argv.add("--mcp-config"); argv.add(mcpJson); argv.add("--append-system-prompt"); @@ -289,209 +149,24 @@ public final class ClaudeCodeLauncher implements PeerLauncher { return argv; } - /** Dedicated worker space → own tab → start the worker (rooted at {@code cwd}) → drop the shell. */ - private Agent spawnInTab(BridgedConfig.Worker cfg, Map workerEnv, - List 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); + // --- Agent-returning convenience spawns (used by callers/tests that want the herdr Agent) --- - Started started; - try { - started = startUniquelyNamed(cfg, workerEnv, argv, tab.tab().tabId(), cwd); - } 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 - // so the tab holds only the worker; label the tab). They must not fail the spawn or - // orphan the running worker — 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; 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(); + /** Spawn a worker for the default profile in the resolved default cwd. */ + public Agent spawn() { + return spawnInternal(null, null, null); } - /** 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 ({}) — worker is running regardless: {}", what, e.getMessage()); - } + /** Spawn a worker for a named profile (null → default) in the resolved default cwd. */ + public Agent spawn(String profileName) { + return spawnInternal(profileName, null, null); } - /** Legacy placement: herdr splits the currently-focused tab; the worker still starts in {@code cwd}. */ - private Agent spawnAsPane(BridgedConfig.Worker cfg, Map workerEnv, - List 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; + /** Spawn a worker for a named profile with an explicit requested/caller cwd (CB-112). */ + public Agent spawn(String profileName, String requestedCwd, String callerCwd) { + return spawnInternal(profileName, requestedCwd, callerCwd); } - /** 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---}: {@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 workerEnv, - List 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 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 by design, 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---} scheme with - * a nonce other 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 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 different 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 only 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. - * - *

Resolves the tab from the pane before 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 ------------------------------------------------------------------- + // --- capabilities -------------------------------------------------------------------------- @Override public Set 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}). */ 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} - * - *

Delegates to the three-arg {@link #spawn(String, String, String)} and wraps the - * resulting herdr {@link Agent} in a {@link WorkerHandle} whose {@link PeerHandle#id()} - * equals the agent's paneId. - * - *

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. + * Whether {@code name} is a Claude Code bridge worker started by a different process + * than {@code currentNonce}. A thin {@code claude}-prefix binding of + * {@link HerdrPeerLauncher#isForeignWorker(String, String, String)}. */ - @Override - public PeerHandle spawn(SpawnRequest req) { - 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 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. - } + static boolean isForeignWorker(String name, String currentNonce) { + return HerdrPeerLauncher.isForeignWorker(NAME_PREFIX, name, currentNonce); } } diff --git a/bridged/src/main/java/dev/ltms/bridged/worker/HerdrPeerLauncher.java b/bridged/src/main/java/dev/ltms/bridged/worker/HerdrPeerLauncher.java new file mode 100644 index 0000000..cea92d1 --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/worker/HerdrPeerLauncher.java @@ -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 herdr + * agent (a CLI coding agent running in a herdr tab/pane). It owns everything that is the same + * regardless of which 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. + * + *

Two seams are peer-specific and supplied by the concrete adapter: + *

    + *
  • {@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.
  • + *
  • {@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}.
  • + *
+ * + *

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 and 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 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 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 profiles, String defaultProfile, + Function 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 env, List argv) { + } + + // --- profile surface ----------------------------------------------------------------------- + + /** The configured peer profile names (what {@code spawn(profile)} accepts). */ + @Override + public Set profiles() { + return profiles.keySet(); + } + + /** The parity-overlay file list for {@code profileName} (default list when unset). */ + @Override + public List 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 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} + * + *

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 workerEnv, + List 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 workerEnv, + List 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 ---}: {@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 workerEnv, + List 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 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 by design, 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 ---} scheme with a nonce other 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 different + * 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 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 ---} 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 different 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 only 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. + * + *

Resolves the tab from the pane before 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 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 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 newEnv() { + return new LinkedHashMap<>(); + } + + /** Defensive copy of {@code argv} plus room to append launch flags. */ + protected static List mutableArgv(List 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. + } + } +} diff --git a/bridged/src/test/java/dev/ltms/bridged/session/SessionManagerTest.java b/bridged/src/test/java/dev/ltms/bridged/session/SessionManagerTest.java index fcfcdfb..fcece67 100644 --- a/bridged/src/test/java/dev/ltms/bridged/session/SessionManagerTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/session/SessionManagerTest.java @@ -350,7 +350,7 @@ class SessionManagerTest { ClaudeCodeLauncher workers = new ClaudeCodeLauncher( new AgentControl(herdr), new WorkspaceControl(herdr), new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null, - 1, 1, () -> clock[0], () -> clock[0] += 10); + 1, () -> clock[0], () -> clock[0] += 10); SessionManager sessions = new SessionManager(workers, new GitWorktrees(), () -> 0L, 0); assertThrows(PeerUnreachableException.class, diff --git a/bridged/src/test/java/dev/ltms/bridged/worker/ClaudeCodeLauncherTest.java b/bridged/src/test/java/dev/ltms/bridged/worker/ClaudeCodeLauncherTest.java index 6fb1d09..a87936d 100644 --- a/bridged/src/test/java/dev/ltms/bridged/worker/ClaudeCodeLauncherTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/worker/ClaudeCodeLauncherTest.java @@ -350,7 +350,7 @@ class ClaudeCodeLauncherTest { new AgentControl(herdr), new WorkspaceControl(herdr), new SubscriptionGuard(Set.of("gx00.gw")), 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)); @@ -370,7 +370,7 @@ class ClaudeCodeLauncherTest { new AgentControl(herdr), new WorkspaceControl(herdr), new SubscriptionGuard(Set.of("gx00.gw")), workerConfigMap("ltms-local", null), "ltms-local", _ -> null, - 1000, 50, () -> clock[0], () -> clock[0] += 50); + 1000, () -> clock[0], () -> clock[0] += 50); PeerUnreachableException ex = assertThrows( PeerUnreachableException.class, @@ -409,7 +409,7 @@ class ClaudeCodeLauncherTest { new AgentControl(herdr), new WorkspaceControl(herdr), new SubscriptionGuard(Set.of("gx00.gw")), 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));