#310: prevent idle reap from stopping delivered workers
This commit is contained in:
@@ -61,6 +61,8 @@ public final class SessionManager implements TurnListener {
|
||||
private final LongSupplier nowNanos;
|
||||
private final int contextCap;
|
||||
private final boolean clearAfterTurn;
|
||||
/** Null in production; test seam for the interval before an idle session's conditional release. */
|
||||
private final Consumer<MemberSession> beforeIdleRelease;
|
||||
private volatile MemberLifecycle memberLifecycle = MemberLifecycle.NONE;
|
||||
/**
|
||||
* CB-586: the repo root the fleet actually works in, remembered the first time a worktree
|
||||
@@ -102,13 +104,23 @@ public final class SessionManager implements TurnListener {
|
||||
}
|
||||
|
||||
public SessionManager(PeerLauncher launcher, Worktrees worktrees, LongSupplier nowNanos,
|
||||
int contextCap, boolean clearAfterTurn) {
|
||||
int contextCap, boolean clearAfterTurn) {
|
||||
this(launcher, worktrees, nowNanos, contextCap, clearAfterTurn, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* Package-private constructor for a deterministic reap/delivery race test. Production callers
|
||||
* use the constructor above, whose null hook adds no callback or lock to an ordinary reap.
|
||||
*/
|
||||
SessionManager(PeerLauncher launcher, Worktrees worktrees, LongSupplier nowNanos,
|
||||
int contextCap, boolean clearAfterTurn, Consumer<MemberSession> beforeIdleRelease) {
|
||||
this.launcher = launcher;
|
||||
this.worktrees = worktrees;
|
||||
this.presence = new PresenceFleet(this);
|
||||
this.nowNanos = nowNanos;
|
||||
this.contextCap = contextCap;
|
||||
this.clearAfterTurn = clearAfterTurn;
|
||||
this.beforeIdleRelease = beforeIdleRelease;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -272,10 +284,27 @@ public final class SessionManager implements TurnListener {
|
||||
*/
|
||||
private void release(String paneId, ReleaseCause cause) {
|
||||
MemberSession removed = registry.remove(paneId);
|
||||
releaseRemoved(paneId, removed, handles.remove(paneId), cause);
|
||||
}
|
||||
|
||||
/**
|
||||
* Tear a session down only while {@code expected} is still its registry value. A lifecycle
|
||||
* transition replaces the immutable record, so this prevents a reap based on an old READY or
|
||||
* DONE record from stopping a worker that delivery has made BUSY.
|
||||
*/
|
||||
private boolean releaseIfCurrent(MemberSession expected, ReleaseCause cause) {
|
||||
if (!registry.remove(expected.paneId(), expected)) {
|
||||
return false;
|
||||
}
|
||||
releaseRemoved(expected.paneId(), expected, handles.remove(expected.paneId()), cause);
|
||||
return true;
|
||||
}
|
||||
|
||||
private void releaseRemoved(String paneId, MemberSession removed, PeerHandle removedHandle,
|
||||
ReleaseCause cause) {
|
||||
// fleetd #209: remove right alongside the registry entry so a released session's handle is
|
||||
// never leaked — but keep the local reference below, so the id can still be resolved for
|
||||
// the ReleaseDetail this teardown notifies with.
|
||||
PeerHandle removedHandle = handles.remove(paneId);
|
||||
boolean preserveWorktree = cause == ReleaseCause.SHUTDOWN;
|
||||
String snapshotRef = null;
|
||||
if (removed != null) {
|
||||
@@ -846,14 +875,18 @@ public final class SessionManager implements TurnListener {
|
||||
}
|
||||
long idleNanos = now - s.lastActivityAtNanos();
|
||||
if (idleNanos > idleTtlNanos) {
|
||||
log.debug("reaping idle session terminal={} pane={}: idle {}s exceeds the {}s ttl",
|
||||
s.terminalId(), s.paneId(), TimeUnit.NANOSECONDS.toSeconds(idleNanos),
|
||||
TimeUnit.NANOSECONDS.toSeconds(idleTtlNanos));
|
||||
// CB-581: one session that fails to release must not abort the whole reaping pass —
|
||||
// match drainAll's per-session try/catch so the rest of the roster still gets reaped.
|
||||
try {
|
||||
release(s.paneId());
|
||||
reaped++;
|
||||
if (beforeIdleRelease != null) {
|
||||
beforeIdleRelease.accept(s);
|
||||
}
|
||||
if (releaseIfCurrent(s, ReleaseCause.COMPLETED)) {
|
||||
log.debug("reaping idle session terminal={} pane={}: idle {}s exceeds the {}s ttl",
|
||||
s.terminalId(), s.paneId(), TimeUnit.NANOSECONDS.toSeconds(idleNanos),
|
||||
TimeUnit.NANOSECONDS.toSeconds(idleTtlNanos));
|
||||
reaped++;
|
||||
}
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("reap failed for pane={} terminal={} worktree={}; continuing with "
|
||||
+ "remaining sessions", s.paneId(), s.terminalId(), s.worktree(), e);
|
||||
|
||||
@@ -178,14 +178,19 @@ class SessionManagerTest {
|
||||
}
|
||||
|
||||
private SessionManager sessionManager(FakeHerdr herdr, LongSupplier clock, int contextCap,
|
||||
boolean clearAfterTurn) {
|
||||
boolean clearAfterTurn) {
|
||||
return sessionManager(herdr, clock, contextCap, clearAfterTurn, null);
|
||||
}
|
||||
|
||||
private SessionManager sessionManager(FakeHerdr herdr, LongSupplier clock, int contextCap,
|
||||
boolean clearAfterTurn, java.util.function.Consumer<MemberSession> hook) {
|
||||
FleetConfig.Profile cfg = new FleetConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
|
||||
List.of("ccs", "ltms-local"), "tab", "fleetd-workers",
|
||||
"worker: {profile} #{n}", null, null, null);
|
||||
ClaudeCodeLauncher workers = new ClaudeCodeLauncher(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, clearAfterTurn);
|
||||
return new SessionManager(workers, new GitWorktrees(), clock, contextCap, clearAfterTurn, hook);
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -670,6 +675,24 @@ class SessionManagerTest {
|
||||
"BUSY session remains");
|
||||
}
|
||||
|
||||
@Test
|
||||
void reapIdleDoesNotReleaseSessionDeliveredAfterItsEligibilityCheck() {
|
||||
long[] clock = {0};
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
SessionManager[] manager = new SessionManager[1];
|
||||
SessionManager sessions = sessionManager(herdr, () -> clock[0], 0, false,
|
||||
session -> manager[0].onDelivered(session.terminalId(), TestTurnTokens.inert(session.terminalId())));
|
||||
manager[0] = sessions;
|
||||
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
|
||||
sessions.asPresence().markPresent(session.terminalId());
|
||||
|
||||
clock[0] = 11;
|
||||
assertEquals(0, sessions.reapIdle(10), "delivery replaces the idle snapshot before release");
|
||||
assertEquals(MemberSession.State.BUSY, sessions.get(session.paneId()).orElseThrow().state(),
|
||||
"a just-delivered session stays registered and busy");
|
||||
assertFalse(herdr.called("pane.close"), "the busy session pane is not stopped");
|
||||
}
|
||||
|
||||
@Test
|
||||
void doneSessionPastIdleTtlIsReaped() {
|
||||
long[] clock = {0};
|
||||
|
||||
Reference in New Issue
Block a user