CB-303: session lifecycle limits — idle_ttl reaper, context_cap, graceful drain

Verified on primary: ide_diagnostics clean (incl. weak warnings), mvn clean install
BUILD SUCCESS, 173 tests. Delegated impl (worker/cb-303-80ec1a-3, 3 parts), primary-gated.
This commit is contained in:
Dai Ha
2026-07-17 10:08:55 +02:00
7 changed files with 459 additions and 19 deletions
+10
View File
@@ -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
@@ -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 {}",
@@ -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<String, Worker> 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);
}
}
@@ -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);
}
@@ -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();
}
}
@@ -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);
}
}
@@ -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();
}
}