CB-303 part 2: context_cap turn budget

This commit is contained in:
Dai Ha
2026-07-17 09:56:46 +02:00
parent 8d51066ddd
commit 954351a80b
3 changed files with 103 additions and 7 deletions
@@ -65,7 +65,13 @@ 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.
SessionReaper reaper = null;
@@ -40,23 +40,35 @@ public final class SessionManager implements TurnListener {
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(), System::nanoTime);
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);
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;
}
/**
@@ -208,16 +220,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. */
@@ -33,13 +33,17 @@ class SessionManagerTest {
}
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);
return new SessionManager(workers, new GitWorktrees(), clock, contextCap);
}
@Test
@@ -254,4 +258,54 @@ class SessionManagerTest {
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");
}
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();
}
}