From a49671ceb959a996fabf58abb0a16833e51c7bea Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Fri, 4 Sep 2026 13:17:51 +0700 Subject: [PATCH] #310: prevent idle reap from stopping delivered workers --- .../ltms/fleet/session/SessionManager.java | 47 ++++++++++++++++--- .../fleet/session/SessionManagerTest.java | 27 ++++++++++- 2 files changed, 65 insertions(+), 9 deletions(-) diff --git a/fleetd/src/main/java/dev/ltms/fleet/session/SessionManager.java b/fleetd/src/main/java/dev/ltms/fleet/session/SessionManager.java index 75ab15c..cf8d8d5 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/session/SessionManager.java +++ b/fleetd/src/main/java/dev/ltms/fleet/session/SessionManager.java @@ -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 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 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); diff --git a/fleetd/src/test/java/dev/ltms/fleet/session/SessionManagerTest.java b/fleetd/src/test/java/dev/ltms/fleet/session/SessionManagerTest.java index e353519..16eba34 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/session/SessionManagerTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/session/SessionManagerTest.java @@ -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 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};