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));