CB-303 part 1: idle_ttl session reaper (injectable clock + SessionReaper)
This commit is contained in:
@@ -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;
|
||||
@@ -66,6 +67,16 @@ public final class Bridged {
|
||||
// CB-301-ext: worktree provisioning seam, optionally rooted at a configured directory.
|
||||
SessionManager sessions = new SessionManager(workers, new GitWorktrees(cfg.worktreeRoot()));
|
||||
|
||||
// CB-303 part 1: idle-ttl reaper — only when configured, defaults to disabled.
|
||||
SessionReaper reaper = null;
|
||||
if (cfg.lifecycle() != null
|
||||
&& cfg.lifecycle().idleTtlSeconds() != null
|
||||
&& cfg.lifecycle().idleTtlSeconds() > 0) {
|
||||
reaper = new SessionReaper(sessions, cfg.lifecycle().idleTtlSeconds());
|
||||
reaper.start();
|
||||
Runtime.getRuntime().addShutdownHook(new Thread(reaper::stop));
|
||||
}
|
||||
|
||||
// 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.
|
||||
// CB-106: a confirmed turn completion resolves a blocked send whose worker never replied.
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -11,7 +11,9 @@ import java.security.SecureRandom;
|
||||
import java.util.List;
|
||||
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.
|
||||
@@ -37,16 +39,24 @@ 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;
|
||||
|
||||
/** Backward-compatible constructor: shared-tree sessions, production git seam. */
|
||||
public SessionManager(WorkerService workerService) {
|
||||
this(workerService, new GitWorktrees());
|
||||
this(workerService, new GitWorktrees(), System::nanoTime);
|
||||
}
|
||||
|
||||
/** Backward-compatible constructor with an injectable worktree seam. */
|
||||
public SessionManager(WorkerService workerService, Worktrees worktrees) {
|
||||
this(workerService, worktrees, System::nanoTime);
|
||||
}
|
||||
|
||||
/** Test constructor with an injectable clock. */
|
||||
public SessionManager(WorkerService workerService, Worktrees worktrees, LongSupplier nowNanos) {
|
||||
this.workerService = workerService;
|
||||
this.worktrees = worktrees;
|
||||
this.presence = new PresenceBridge(this);
|
||||
this.nowNanos = nowNanos;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -80,13 +90,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);
|
||||
@@ -133,13 +146,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);
|
||||
@@ -220,6 +236,68 @@ 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 = nowNanos.getAsLong() + timeoutNanos;
|
||||
long pollNanos = TimeUnit.MILLISECONDS.toNanos(50);
|
||||
for (WorkerSession s : roster()) {
|
||||
try {
|
||||
if (s.state() == WorkerSession.State.BUSY) {
|
||||
while (nowNanos.getAsLong() < deadline) {
|
||||
WorkerSession current = registry.get(s.paneId());
|
||||
if (current == null || current.state() != WorkerSession.State.BUSY) {
|
||||
break;
|
||||
}
|
||||
try {
|
||||
long remaining = deadline - nowNanos.getAsLong();
|
||||
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();
|
||||
@@ -236,7 +314,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);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user