diff --git a/bridged/bridged.example.yaml b/bridged/bridged.example.yaml index f619a79..ad284c8 100644 --- a/bridged/bridged.example.yaml +++ b/bridged/bridged.example.yaml @@ -53,3 +53,13 @@ guard: - gx00.gw - gx01.gw - ollama.ltms.dev + +# Session lifecycle limits (CB-303). All knobs are opt-in; omit or set to null to keep +# the feature disabled. By default the daemon never reaps, caps, or drains sessions. +# idleTtlSeconds → reap READY/DONE sessions idle longer than this (never BUSY/SPAWNING) +# contextCap → force-release a session after this many delegated turns +# drainTimeoutSeconds → seconds to wait for BUSY sessions on shutdown before forced teardown +# lifecycle: +# idleTtlSeconds: 300 +# contextCap: 10 +# drainTimeoutSeconds: 5 diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java index 77180ae..d0ec531 100644 --- a/bridged/src/main/java/dev/ltms/bridged/Bridged.java +++ b/bridged/src/main/java/dev/ltms/bridged/Bridged.java @@ -20,6 +20,7 @@ import dev.ltms.bridged.msg.Rendezvous; import dev.ltms.bridged.rest.BridgedApp; import dev.ltms.bridged.session.GitWorktrees; import dev.ltms.bridged.session.SessionManager; +import dev.ltms.bridged.session.SessionReaper; import dev.ltms.bridged.worker.WorkerService; import io.javalin.Javalin; import org.slf4j.Logger; @@ -52,7 +53,6 @@ public final class Bridged { : UnixSocketHerdrClient.defaultSocketPath(); UnixSocketHerdrClient herdr = UnixSocketHerdrClient.connect(socket, new com.fasterxml.jackson.databind.ObjectMapper()); - Runtime.getRuntime().addShutdownHook(new Thread(herdr::close)); AgentControl agents = new AgentControl(herdr); WorkspaceControl spaces = new WorkspaceControl(herdr); @@ -64,7 +64,24 @@ public final class Bridged { // CB-301: authoritative session registry + lifecycle FSM on top of WorkerService. // CB-301-ext: worktree provisioning seam, optionally rooted at a configured directory. - SessionManager sessions = new SessionManager(workers, new GitWorktrees(cfg.worktreeRoot())); + // CB-303 part 2: context cap is opt-in and disabled (0) when absent/null. + int contextCap = 0; + if (cfg.lifecycle() != null && cfg.lifecycle().contextCap() != null + && cfg.lifecycle().contextCap() > 0) { + contextCap = cfg.lifecycle().contextCap(); + } + SessionManager sessions = new SessionManager(workers, new GitWorktrees(cfg.worktreeRoot()), contextCap); + + // CB-303 part 1: idle-ttl reaper — only when configured, defaults to disabled. + final SessionReaper reaper; + if (cfg.lifecycle() != null + && cfg.lifecycle().idleTtlSeconds() != null + && cfg.lifecycle().idleTtlSeconds() > 0) { + reaper = new SessionReaper(sessions, cfg.lifecycle().idleTtlSeconds()); + reaper.start(); + } else { + reaper = null; + } // Status-gated injector (CB-103): the single writer into workers, fed by a poller. // The blocking message endpoint (CB-104) is the producer; the poller is inert until then. @@ -96,17 +113,27 @@ public final class Bridged { Injector injector = new Injector(agents, turnListener, presence::isPresent, presence::forget); StatusPoller poller = new StatusPoller(agents, injector, INJECT_POLL_MILLIS); poller.start(); - Runtime.getRuntime().addShutdownHook(new Thread(poller::stop)); MessageService messages = new MessageService(agents, injector, rendezvous); - Runtime.getRuntime().addShutdownHook(new Thread(messages::close)); // MCP server face (CB-105): bridge_send/bridge_reply/bridge_status, mounted at /mcp. // Caller identity is resolved from the connection (peer PID → herdr pane), not arguments. ConnectionIdentity identity = new ConnectionIdentity( new PaneLocator(herdr), new LsofPeerPidLookup(), new LsofProcessCwdLookup()); BridgeMcp mcp = new BridgeMcp(messages, rendezvous, workers, sessions, identity, presence); - Runtime.getRuntime().addShutdownHook(new Thread(mcp::close)); + + // CB-303 part 3: single ordered shutdown hook. Drain sessions first while herdr is still + // open (so releases reach the daemon), then stop poller/message/mcp/reaper, and close herdr + // last. This replaces the earlier independent hooks that could race and close herdr early. + Runtime.getRuntime().addShutdownHook(new Thread(() -> { + sessions.close(cfg.lifecycle() != null ? cfg.lifecycle().drainTimeoutSeconds() : null); + poller.stop(); + messages.close(); + mcp.close(); + if (reaper != null) reaper.stop(); + herdr.close(); + })); + Javalin app = new BridgedApp(herdr, workers, sessions, messages, rendezvous, presence, mcp.servlet()).build(); app.start(cfg.bind().host(), cfg.bind().port()); log.info("bridged listening on {}:{}, herdr socket {}", diff --git a/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java b/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java index 1394d5f..0252f1e 100644 --- a/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java +++ b/bridged/src/main/java/dev/ltms/bridged/config/BridgedConfig.java @@ -27,6 +27,7 @@ import java.util.Set; * @param guard subscription-boundary allowlist * @param worktreeRoot nullable root directory for provisioned worktrees; defaults to a sibling * of the repo root + * @param lifecycle session lifecycle limits ({@code null} = all disabled) */ @JsonIgnoreProperties(ignoreUnknown = true) public record BridgedConfig( @@ -36,7 +37,8 @@ public record BridgedConfig( Map workers, String defaultWorker, Guard guard, - String worktreeRoot) { + String worktreeRoot, + Lifecycle lifecycle) { @JsonIgnoreProperties(ignoreUnknown = true) public record Bind(String host, int port) { @@ -144,6 +146,21 @@ public record BridgedConfig( } } + /** + * Session lifecycle limits. All knobs are opt-in: {@code null} or {@code 0} disables the + * feature so existing configs keep the previous behaviour. + * + * @param idleTtlSeconds max seconds a {@code READY}/{@code DONE} session may sit idle + * before it is reaped ({@code null} → disabled) + * @param contextCap max delegated turns a session serves before force-release + * ({@code null} → disabled) + * @param drainTimeoutSeconds seconds to wait for {@code BUSY} sessions to finish before + * forced teardown on shutdown (default 5 when unset) + */ + @JsonIgnoreProperties(ignoreUnknown = true) + public record Lifecycle(Integer idleTtlSeconds, Integer contextCap, Integer drainTimeoutSeconds) { + } + /** * Subscription boundary. Only these hosts may back a worker's * {@code ANTHROPIC_BASE_URL}; the primary must carry none. @@ -211,6 +228,7 @@ public record BridgedConfig( public BridgedConfig withDefaults() { Bind b = bind != null ? bind : new Bind(null, 0); Guard g = guard != null ? guard : new Guard(List.of()); - return new BridgedConfig(b, herdrSocket, worker, workers, defaultWorker, g, worktreeRoot); + Lifecycle l = lifecycle != null ? lifecycle : new Lifecycle(null, null, null); + return new BridgedConfig(b, herdrSocket, worker, workers, defaultWorker, g, worktreeRoot, l); } } diff --git a/bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java b/bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java index 7165c64..f1883a2 100644 --- a/bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java +++ b/bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java @@ -13,7 +13,9 @@ import java.util.List; import java.util.Map; import java.util.Optional; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicLong; +import java.util.function.LongSupplier; /** * Authoritative in-daemon registry of the worker sessions this {@code bridged} process spawned. @@ -39,16 +41,36 @@ public final class SessionManager implements TurnListener { private final WorkerPresence presence; private final SecureRandom nonceRandom = new SecureRandom(); private final AtomicLong nonceSeq = new AtomicLong(); + private final LongSupplier nowNanos; + private final int contextCap; /** Backward-compatible constructor: shared-tree sessions, production git seam. */ public SessionManager(WorkerService workerService) { - this(workerService, new GitWorktrees()); + this(workerService, new GitWorktrees(), System::nanoTime, 0); } + /** Backward-compatible constructor with an injectable worktree seam. */ public SessionManager(WorkerService workerService, Worktrees worktrees) { + this(workerService, worktrees, System::nanoTime, 0); + } + + /** Test constructor with an injectable clock. */ + public SessionManager(WorkerService workerService, Worktrees worktrees, LongSupplier nowNanos) { + this(workerService, worktrees, nowNanos, 0); + } + + /** Production constructor with a configured context turn cap. */ + public SessionManager(WorkerService workerService, Worktrees worktrees, int contextCap) { + this(workerService, worktrees, System::nanoTime, contextCap); + } + + public SessionManager(WorkerService workerService, Worktrees worktrees, LongSupplier nowNanos, + int contextCap) { this.workerService = workerService; this.worktrees = worktrees; this.presence = new PresenceBridge(this); + this.nowNanos = nowNanos; + this.contextCap = contextCap; } /** @@ -82,13 +104,16 @@ public final class SessionManager implements TurnListener { Agent worker = workerService.spawn(profile, requestedCwd, callerCwd); String resolvedProfile = (profile == null || profile.isBlank()) ? workerService.defaultProfile() : profile; + long now = nowNanos.getAsLong(); WorkerSession session = new WorkerSession( worker.paneId(), worker.terminalId(), resolvedProfile, resolveCwd(requestedCwd, profile, callerCwd), ownerTerminal, - System.nanoTime(), + now, + now, + 0, WorkerSession.State.SPAWNING, null, null); @@ -135,13 +160,16 @@ public final class SessionManager implements TurnListener { } throw e; } + long now = nowNanos.getAsLong(); WorkerSession session = new WorkerSession( worker.paneId(), worker.terminalId(), resolvedProfile, resolveCwd(path, profile, callerCwd), ownerTerminal, - System.nanoTime(), + now, + now, + 0, WorkerSession.State.SPAWNING, path, branch); @@ -217,16 +245,40 @@ public final class SessionManager implements TurnListener { transitionByTerminal(terminalId, WorkerSession.State.SPAWNING, WorkerSession.State.READY); } - /** Lifecycle hook: a message was delivered into the worker — it is now busy on a turn. */ + /** + * Lifecycle hook: a message was delivered into the worker — it is now busy on a turn. + * The turn count is bumped and the activity timestamp is refreshed. A {@code DONE} session + * can be re-delivered for multi-turn reuse until it is released. + */ @Override public void onDelivered(String target) { - transitionByTerminal(target, WorkerSession.State.READY, WorkerSession.State.BUSY); + WorkerSession current = findByTerminal(target); + if (current == null) return; + if (current.state() != WorkerSession.State.READY && current.state() != WorkerSession.State.DONE) { + return; + } + long now = nowNanos.getAsLong(); + WorkerSession updated = current.withState(WorkerSession.State.BUSY).bumpTurn(now); + if (replace(current, updated)) { + log.debug("session transitioned terminal={} pane={} {} -> BUSY turn={}", + target, current.paneId(), current.state(), updated.turnCount()); + } } /** Lifecycle hook: the worker's delegated turn completed successfully. */ @Override public void onTurnComplete(String target) { - transitionByTerminal(target, WorkerSession.State.BUSY, WorkerSession.State.DONE); + WorkerSession current = findByTerminal(target); + if (current == null || current.state() != WorkerSession.State.BUSY) return; + long now = nowNanos.getAsLong(); + WorkerSession updated = current.withState(WorkerSession.State.DONE).withActivity(now); + if (replace(current, updated)) { + log.debug("session transitioned terminal={} pane={} BUSY -> DONE turn={}", + target, current.paneId(), updated.turnCount()); + } + if (contextCap > 0 && updated.turnCount() >= contextCap) { + release(current.paneId()); + } } /** Lifecycle hook: the worker's delegated turn failed. */ @@ -245,6 +297,67 @@ public final class SessionManager implements TurnListener { } } + /** + * Best-effort reap of sessions that have been idle longer than {@code idleTtlNanos}. Only + * {@code READY} and {@code DONE} sessions are eligible — never a {@code SPAWNING} or + * {@code BUSY} worker. Returns the number of sessions released. + */ + int reapIdle(long idleTtlNanos) { + long now = nowNanos.getAsLong(); + int reaped = 0; + for (WorkerSession s : roster()) { + if (s.state() != WorkerSession.State.READY && s.state() != WorkerSession.State.DONE) { + continue; + } + if (now - s.lastActivityAtNanos() > idleTtlNanos) { + release(s.paneId()); + reaped++; + } + } + return reaped; + } + + /** + * Gracefully drain all registered sessions. For each session that is {@code BUSY}, poll up to + * {@code timeoutNanos} for it to leave {@code BUSY}, then release it regardless. Non-busy + * sessions are released immediately. A failure releasing one session is logged and does not + * abort the rest. + */ + void drainAll(long timeoutNanos) { + long deadline = System.nanoTime() + timeoutNanos; + for (WorkerSession s : roster()) { + try { + if (s.state() == WorkerSession.State.BUSY) { + while (System.nanoTime() < deadline) { + WorkerSession current = registry.get(s.paneId()); + if (current == null || current.state() != WorkerSession.State.BUSY) { + break; + } + try { + long remaining = deadline - System.nanoTime(); + Thread.sleep(Math.min(TimeUnit.NANOSECONDS.toMillis(remaining), 50)); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + break; + } + } + } + release(s.paneId()); + } catch (RuntimeException e) { + log.warn("drain failed for pane={}; continuing with remaining sessions", s.paneId(), e); + } + } + } + + /** + * Close this manager by draining all sessions. The timeout comes from configuration when set, + * otherwise a sensible default. + */ + public void close(Integer drainTimeoutSeconds) { + int seconds = (drainTimeoutSeconds != null && drainTimeoutSeconds > 0) ? drainTimeoutSeconds : 5; + drainAll(TimeUnit.SECONDS.toNanos(seconds)); + } + /** Number of sessions currently registered. */ public int size() { return registry.size(); @@ -261,7 +374,8 @@ public final class SessionManager implements TurnListener { WorkerSession.State to) { WorkerSession current = findByTerminal(terminalId); if (current == null || current.state() != from) return; - if (replace(current, current.withState(to))) { + long now = nowNanos.getAsLong(); + if (replace(current, current.withState(to).withActivity(now))) { log.debug("session transitioned terminal={} pane={} {} -> {}", terminalId, current.paneId(), from, to); } diff --git a/bridged/src/main/java/dev/ltms/bridged/session/SessionReaper.java b/bridged/src/main/java/dev/ltms/bridged/session/SessionReaper.java new file mode 100644 index 0000000..43bb764 --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/session/SessionReaper.java @@ -0,0 +1,70 @@ +package dev.ltms.bridged.session; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.util.concurrent.TimeUnit; + +/** + * Periodic virtual-thread reaper that tears down {@code READY}/{@code DONE} sessions which have + * exceeded their idle TTL. Modeled on {@link dev.ltms.bridged.inject.StatusPoller}: a single + * virtual-thread loop, idempotent start/stop, and no {@code ScheduledExecutorService}. + */ +public final class SessionReaper { + + private static final Logger log = LoggerFactory.getLogger(SessionReaper.class); + private static final long DEFAULT_INTERVAL_MILLIS = 5000; + + private final SessionManager sessions; + private final long idleTtlNanos; + private final long intervalMillis; + private volatile boolean running; + private Thread thread; + + /** Construct a reaper with the default 5-second polling interval. */ + public SessionReaper(SessionManager sessions, long idleTtlSeconds) { + this(sessions, idleTtlSeconds, DEFAULT_INTERVAL_MILLIS); + } + + /** Construct a reaper with an explicit polling interval (useful for tests). */ + public SessionReaper(SessionManager sessions, long idleTtlSeconds, long intervalMillis) { + this.sessions = sessions; + this.idleTtlNanos = TimeUnit.SECONDS.toNanos(idleTtlSeconds); + this.intervalMillis = intervalMillis; + } + + /** Start the reaper loop on a virtual thread. Idempotent. */ + public synchronized void start() { + if (running) return; + running = true; + thread = Thread.ofVirtual().name("session-reaper").start(this::loop); + log.info("session reaper started (idle ttl {}s, interval {}ms)", + TimeUnit.NANOSECONDS.toSeconds(idleTtlNanos), intervalMillis); + } + + private void loop() { + while (running) { + try { + sessions.reapIdle(idleTtlNanos); + } catch (RuntimeException e) { + log.warn("session reaper iteration failed; continuing", e); + } + sleep(); + } + } + + private void sleep() { + try { + Thread.sleep(intervalMillis); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + running = false; + } + } + + /** Stop the reaper loop. Idempotent. */ + public synchronized void stop() { + running = false; + if (thread != null) thread.interrupt(); + } +} diff --git a/bridged/src/main/java/dev/ltms/bridged/session/WorkerSession.java b/bridged/src/main/java/dev/ltms/bridged/session/WorkerSession.java index cd29a1e..3c96113 100644 --- a/bridged/src/main/java/dev/ltms/bridged/session/WorkerSession.java +++ b/bridged/src/main/java/dev/ltms/bridged/session/WorkerSession.java @@ -10,8 +10,10 @@ package dev.ltms.bridged.session; * @param profile the worker profile name that spawned this session * @param cwd the resolved working directory the worker started in * @param ownerTerminal the caller that requested this worker ({@code null} = daemon/anon) - * @param spawnedAtNanos {@link System#nanoTime()} when the session was registered - * @param state current lifecycle state in the one-shot FSM + * @param spawnedAtNanos {@link System#nanoTime()} when the session was registered + * @param lastActivityAtNanos {@link System#nanoTime()} of the most recent lifecycle event + * @param turnCount number of delegated turns that have been delivered to this session + * @param state current lifecycle state in the one-shot FSM */ public record WorkerSession( String paneId, @@ -20,6 +22,8 @@ public record WorkerSession( String cwd, String ownerTerminal, long spawnedAtNanos, + long lastActivityAtNanos, + int turnCount, State state, String worktree, String branch) { @@ -36,7 +40,19 @@ public record WorkerSession( /** Return a copy of this session in {@code state}. */ public WorkerSession withState(State state) { - return new WorkerSession(paneId, terminalId, profile, cwd, ownerTerminal, spawnedAtNanos, state, - worktree, branch); + return new WorkerSession(paneId, terminalId, profile, cwd, ownerTerminal, spawnedAtNanos, + lastActivityAtNanos, turnCount, state, worktree, branch); + } + + /** Return a copy with {@code lastActivityAtNanos} updated to {@code nowNanos}. */ + public WorkerSession withActivity(long nowNanos) { + return new WorkerSession(paneId, terminalId, profile, cwd, ownerTerminal, spawnedAtNanos, + nowNanos, turnCount, state, worktree, branch); + } + + /** Return a copy with the turn count incremented and activity timestamped at {@code nowNanos}. */ + public WorkerSession bumpTurn(long nowNanos) { + return new WorkerSession(paneId, terminalId, profile, cwd, ownerTerminal, spawnedAtNanos, + nowNanos, turnCount + 1, state, worktree, branch); } } 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 40eeb0d..b2fff3f 100644 --- a/bridged/src/test/java/dev/ltms/bridged/session/SessionManagerTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/session/SessionManagerTest.java @@ -11,11 +11,14 @@ import org.junit.jupiter.api.Test; import java.util.List; import java.util.Map; import java.util.Set; +import java.util.concurrent.TimeUnit; +import java.util.function.LongSupplier; import static org.junit.jupiter.api.Assertions.*; /** - * CB-301 acceptance tests for the authoritative session registry and one-shot lifecycle FSM. + * CB-301 / CB-303 acceptance tests for the authoritative session registry, one-shot lifecycle FSM, + * and configurable lifecycle limits (idle TTL, context cap, drain). * No live herdr — everything runs against the same {@link FakeHerdr} the rest of the project uses. */ class SessionManagerTest { @@ -30,6 +33,20 @@ class SessionManagerTest { return new SessionManager(workers); } + private SessionManager sessionManager(FakeHerdr herdr, LongSupplier clock) { + return sessionManager(herdr, clock, 0); + } + + private SessionManager sessionManager(FakeHerdr herdr, LongSupplier clock, int contextCap) { + BridgedConfig.Worker cfg = new BridgedConfig.Worker( + "ltms-local", "http://gx00.gw:8000", "coder", null, "BRIDGED_WORKER_TOKEN", + List.of("ccs", "ltms-local"), "tab", "bridged-workers", + "worker: {profile} #{n}", null, null, null); + WorkerService workers = new WorkerService(new AgentControl(herdr), new WorkspaceControl(herdr), + new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null); + return new SessionManager(workers, new GitWorktrees(), clock, contextCap); + } + @Test void acquireRegistersSpawningSessionWithDistinctPaneId() { FakeHerdr herdr = new FakeHerdr(); @@ -147,4 +164,172 @@ class SessionManagerTest { assertEquals(1, sessions.roster().size()); assertEquals(b.paneId(), sessions.roster().getFirst().paneId()); } + + // --- CB-303 lifecycle limits ---------------------------------------------------- + + @Test + void reapIdleDoesNothingWhenNoSessions() { + FakeHerdr herdr = new FakeHerdr(); + SessionManager sessions = sessionManager(herdr, () -> 0L); + + assertEquals(0, sessions.reapIdle(10)); + assertTrue(sessions.roster().isEmpty()); + } + + @Test + void readySessionPastIdleTtlIsReaped() { + long[] clock = {0}; + FakeHerdr herdr = new FakeHerdr(); + SessionManager sessions = sessionManager(herdr, () -> clock[0]); + WorkerSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary"); + String terminal = session.terminalId(); + sessions.asPresence().markPresent(terminal); + + clock[0] = 11; + assertEquals(1, sessions.reapIdle(10), "READY session past TTL is reaped"); + assertTrue(sessions.get(session.paneId()).isEmpty(), "reaped session is removed from registry"); + assertTrue(herdr.called("pane.close"), "reaped session tears the pane down"); + } + + @Test + void readySessionWithinIdleTtlSurvives() { + long[] clock = {0}; + FakeHerdr herdr = new FakeHerdr(); + SessionManager sessions = sessionManager(herdr, () -> clock[0]); + WorkerSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary"); + String terminal = session.terminalId(); + sessions.asPresence().markPresent(terminal); + + clock[0] = 5; + assertEquals(0, sessions.reapIdle(10), "READY session within TTL is not reaped"); + assertEquals(WorkerSession.State.READY, + sessions.get(session.paneId()).orElseThrow().state(), + "READY session survives"); + } + + @Test + void busySessionPastIdleTtlIsNotReaped() { + long[] clock = {0}; + FakeHerdr herdr = new FakeHerdr(); + SessionManager sessions = sessionManager(herdr, () -> clock[0]); + WorkerSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary"); + String terminal = session.terminalId(); + sessions.asPresence().markPresent(terminal); + sessions.onDelivered(terminal); + + clock[0] = 100; + assertEquals(0, sessions.reapIdle(10), "BUSY session past TTL is never reaped"); + assertEquals(WorkerSession.State.BUSY, + sessions.get(session.paneId()).orElseThrow().state(), + "BUSY session remains"); + } + + @Test + void doneSessionPastIdleTtlIsReaped() { + long[] clock = {0}; + FakeHerdr herdr = new FakeHerdr(); + SessionManager sessions = sessionManager(herdr, () -> clock[0]); + WorkerSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary"); + String terminal = session.terminalId(); + sessions.asPresence().markPresent(terminal); + sessions.onDelivered(terminal); + sessions.onTurnComplete(terminal); + + clock[0] = 21; + assertEquals(1, sessions.reapIdle(20), "DONE session past TTL is reaped"); + assertTrue(sessions.get(session.paneId()).isEmpty(), "DONE session is removed"); + } + + @Test + void reapIdleReturnsCorrectCountAndSkipsBusy() { + long[] clock = {0}; + FakeHerdr herdr = new FakeHerdr(); + SessionManager sessions = sessionManager(herdr, () -> clock[0]); + + WorkerSession ready = sessions.acquire("ltms-local", "/ready", "/caller", "owner1"); + WorkerSession busy = sessions.acquire("ltms-local", "/busy", "/caller", "owner2"); + sessions.asPresence().markPresent(ready.terminalId()); + sessions.asPresence().markPresent(busy.terminalId()); + sessions.onDelivered(busy.terminalId()); + + clock[0] = 50; + assertEquals(1, sessions.reapIdle(30), "only READY past TTL is reaped"); + assertTrue(sessions.get(ready.paneId()).isEmpty(), "READY session is gone"); + assertEquals(WorkerSession.State.BUSY, + sessions.get(busy.paneId()).orElseThrow().state(), + "BUSY session is still registered"); + } + + @Test + void contextCapDisabledSessionSurvivesMultipleTurns() { + FakeHerdr herdr = new FakeHerdr(); + SessionManager sessions = sessionManager(herdr, () -> 0L, 0); + WorkerSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary"); + String terminal = session.terminalId(); + sessions.asPresence().markPresent(terminal); + + sessions.onDelivered(terminal); + sessions.onTurnComplete(terminal); + sessions.onDelivered(terminal); + sessions.onTurnComplete(terminal); + + WorkerSession updated = sessions.get(session.paneId()).orElseThrow(); + assertEquals(WorkerSession.State.DONE, updated.state(), "session finishes second turn"); + assertEquals(2, updated.turnCount(), "turn count tracks both deliveries"); + long releaseCloseCount = paneCloseCallsFor(herdr, session.paneId()); + assertEquals(0, releaseCloseCount, "cap disabled — no forced release of the worker pane"); + } + + @Test + void contextCapTwoReleasesAfterSecondComplete() { + FakeHerdr herdr = new FakeHerdr(); + SessionManager sessions = sessionManager(herdr, () -> 0L, 2); + WorkerSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary"); + String terminal = session.terminalId(); + sessions.asPresence().markPresent(terminal); + + sessions.onDelivered(terminal); + sessions.onTurnComplete(terminal); + assertEquals(WorkerSession.State.DONE, + sessions.get(session.paneId()).orElseThrow().state(), + "first turn completes without release"); + + sessions.onDelivered(terminal); + sessions.onTurnComplete(terminal); + + assertTrue(sessions.get(session.paneId()).isEmpty(), "session released after cap reached"); + assertTrue(sessions.roster().isEmpty(), "released session leaves roster"); + assertEquals(1, paneCloseCallsFor(herdr, session.paneId()), + "forced release tears the worker pane down exactly once"); + } + + @Test + void drainAllReleasesBusyAndReadySessionsAndWaitsForBusy() { + long[] clock = {0}; + FakeHerdr herdr = new FakeHerdr(); + SessionManager sessions = sessionManager(herdr, () -> clock[0]); + + WorkerSession ready = sessions.acquire("ltms-local", "/ready", "/caller", "ownerR"); + WorkerSession busy = sessions.acquire("ltms-local", "/busy", "/caller", "ownerB"); + sessions.asPresence().markPresent(ready.terminalId()); + sessions.asPresence().markPresent(busy.terminalId()); + sessions.onDelivered(busy.terminalId()); + + sessions.drainAll(TimeUnit.MILLISECONDS.toNanos(100)); + + assertTrue(sessions.roster().isEmpty(), "drain clears the roster"); + assertTrue(sessions.get(ready.paneId()).isEmpty(), "ready session is released"); + assertTrue(sessions.get(busy.paneId()).isEmpty(), "busy session is released after timeout"); + assertEquals(1, paneCloseCallsFor(herdr, ready.paneId()), + "ready worker pane is torn down"); + assertEquals(1, paneCloseCallsFor(herdr, busy.paneId()), + "busy worker pane is torn down"); + } + + private static long paneCloseCallsFor(FakeHerdr herdr, String paneId) { + return herdr.calls.stream() + .filter(c -> "pane.close".equals(c.method())) + .filter(c -> paneId.equals(((Map) c.params()).get("pane_id"))) + .count(); + } }