From 83f2aea60fa47e060844c2883de1691fc663a8b0 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Fri, 4 Sep 2026 13:21:53 +0700 Subject: [PATCH] #308: refuse a spawn once the shutdown drain has started, and sweep stragglers MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit drainAll iterated a one-shot registry snapshot with nothing to refuse a new fleet_spawn while the drain was still running (mcp.close() only runs 8 calls after sessions.close() in the shutdown hook). A session registered in that window was never visited by the drain loop: its pane kept running and its worktree was never preserved, with the in-memory registry gone at exit. Fix, both mechanisms as the issue asked for (neither alone is complete): - SessionManager.acquire now checks a `draining` flag, flipped true at the very start of drainAll before the registry snapshot is even taken, and throws the new ShuttingDownException (invariant 3: fail loudly, say why). FleetMcp.spawn and FleetApp.spawnMember surface it as a clean error/503 rather than an uncaught RuntimeException. - The flag alone cannot close the whole race: a caller already past the check can still be mid-launcher.spawn() (a real herdr round trip) when drainAll snapshots the registry. drainAll now re-reads the registry once its main pass finishes and drains whatever straggler landed there too, bounded by the SAME whole-drain deadline (invariant 1: timeoutNanos stays a budget for the whole drain, never extended for a straggler). - ReleaseCause.SHUTDOWN still preserves worktrees for both the initial pass and the sweep (invariant 2, unchanged release() path). Tests (SessionManagerTest): a guard test proving acquire() throws once drainAll has started, and a race test using a launcher double that blocks the second spawn() and the first stop() call to force, deterministically, the exact interleaving where a spawn passes the guard before drainAll flips it and only registers after the initial snapshot — proving the post-loop sweep catches it. Shape check (SessionManager.java only, not fixed): reapIdle has the same shape — a decision made from a roster() snapshot, then acted on via release(s.paneId()) with no re-check of the session's current state. --- .../java/dev/ltms/fleet/mcp/FleetMcp.java | 5 + .../java/dev/ltms/fleet/rest/FleetApp.java | 6 + .../ltms/fleet/session/SessionManager.java | 51 ++++- .../fleet/session/ShuttingDownException.java | 18 ++ .../fleet/session/SessionManagerTest.java | 185 ++++++++++++++++++ 5 files changed, 264 insertions(+), 1 deletion(-) create mode 100644 fleetd/src/main/java/dev/ltms/fleet/session/ShuttingDownException.java diff --git a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java index e042e1e..9dda812 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java @@ -22,6 +22,7 @@ import dev.ltms.fleet.placement.BackendQuarantine; import dev.ltms.fleet.placement.PlacementException; import dev.ltms.fleet.session.SessionManager; import dev.ltms.fleet.session.MemberSession; +import dev.ltms.fleet.session.ShuttingDownException; import dev.ltms.fleet.session.WorktreeRequest; import dev.ltms.fleet.peer.MemberRole; import dev.ltms.fleet.peer.PeerLauncher; @@ -962,6 +963,10 @@ public final class FleetMcp { return text(json(memberView(member))); } catch (GuardException e) { return error("subscription boundary: " + e.getMessage()); + } catch (ShuttingDownException e) { + // fleetd #308: the daemon's shutdown drain has already started — refuse loudly rather + // than register a session drainAll will never see again. + return error("shutting down: " + e.getMessage()); } catch (PlacementException e) { // CB-599: no candidate had capacity (maxLoad, quarantine, or all-exhausted) — distinct // from "profile does not exist" below. diff --git a/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java b/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java index 2ebf041..e536cd9 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java +++ b/fleetd/src/main/java/dev/ltms/fleet/rest/FleetApp.java @@ -19,6 +19,7 @@ import dev.ltms.fleet.peer.PeerUnreachableException; import dev.ltms.fleet.placement.PlacementException; import dev.ltms.fleet.msg.MessageService; import dev.ltms.fleet.session.SessionManager; +import dev.ltms.fleet.session.ShuttingDownException; import dev.ltms.fleet.peer.MemberRole; import dev.ltms.fleet.session.MemberSession; import dev.ltms.fleet.session.WorktreeRequest; @@ -501,6 +502,11 @@ public final class FleetApp { ctx.status(201).json(view(member)); } catch (GuardException e) { ctx.status(403).json(Map.of("error", "subscription_boundary", "detail", e.getMessage())); + } catch (ShuttingDownException e) { + // fleetd #308: the daemon's shutdown drain has already started — 503, not a bare 500, + // so this reads the same as PlacementException below: valid request, refused because + // of a transient daemon state rather than a bad argument. + ctx.status(503).json(Map.of("error", "shutting_down", "detail", e.getMessage())); } catch (PlacementException e) { // CB-599: no candidate had capacity (maxLoad, quarantine, or all-exhausted) — a benign, // likely-transient refusal, distinct from "profile does not exist" below. 503: the 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..229ed00 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/session/SessionManager.java +++ b/fleetd/src/main/java/dev/ltms/fleet/session/SessionManager.java @@ -21,6 +21,7 @@ import java.util.Optional; import java.util.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicLong; import java.util.function.Consumer; import java.util.function.LongSupplier; @@ -71,6 +72,16 @@ public final class SessionManager implements TurnListener { */ private volatile String fleetRepoRoot; + /** + * fleetd #308: flips true the instant {@link #drainAll} starts, before its registry snapshot + * is even taken — so a spawn already in flight sees the refusal as early as a plain flag can + * make it. This alone cannot close the race completely: a caller that read {@code false} just + * before the flip can still land in the registry after the snapshot. {@link #drainAll}'s + * post-loop sweep is what catches that straggler; the two mechanisms are deliberately paired, + * see {@link #drainAll}'s javadoc. + */ + private final AtomicBoolean draining = new AtomicBoolean(false); + /** CB-520: notified with a terminalId on every acquire; no-op until wired. */ private final List> acquireListeners = new java.util.concurrent.CopyOnWriteArrayList<>(); /** CB-516: notified with a {@link ReleaseDetail} on every release; no-op until wired. */ @@ -181,6 +192,14 @@ public final class SessionManager implements TurnListener { public MemberSession acquire(String profile, MemberRole role, String requestedCwd, String callerCwd, String ownerTerminal, WorktreeRequest wt, String sessionName, String resumeSessionId) { + // fleetd #308: refuse before anything else runs — no slot reservation, no launcher spawn — + // so a caller learns the daemon is going down instead of getting a session drainAll will + // never see again. Checked here because every other acquire(...) overload delegates to + // this one, so this is the single point every spawn path passes through. + if (draining.get()) { + throw new ShuttingDownException("fleetd is shutting down; refusing to spawn a session " + + "the shutdown drain would never see"); + } MemberRole memberRole = (role == null) ? MemberRole.DEV : role; requireResumeCapability(profile, resumeSessionId); // CB-619 / fleetd #123: an explicit profile bypasses placement (CompositePeerLauncher only @@ -884,10 +903,40 @@ public final class SessionManager implements TurnListener { * a reason to delete a worker's only copy of its uncommitted work. A session still {@code BUSY} * when the timeout expired is abandoned mid-turn and logged loudly so an operator can find its * kept worktree. + * + *

fleetd #308: {@code roster()} is a one-shot snapshot (see its javadoc), and nothing used + * to stop a new session from registering after it was taken — {@link #acquire} stayed open for + * as long as this drain waited on a {@code BUSY} session, up to the whole {@code timeoutNanos} + * budget. Two things close that window, deliberately paired because neither alone is complete: + * {@link #draining} is flipped true before the snapshot is even taken, so {@link #acquire} + * refuses (invariant 3: loudly, via {@link ShuttingDownException}) as much of the window as a + * plain flag can close; and the sweep below re-reads the registry once the initial snapshot has + * fully drained and drains whatever a straggler — a caller that read the flag as {@code false} + * a moment before it flipped — still managed to register. The sweep shares the same + * {@code deadline} rather than getting its own: {@code timeoutNanos} is a budget for the WHOLE + * drain (see above), and a straggler must not buy the drain more time than the flag it lost the + * race against would have. In the ordinary case the sweep finds nothing and costs one empty + * {@link #roster()} call. */ void drainAll(long timeoutNanos) { long deadline = System.nanoTime() + timeoutNanos; - for (MemberSession s : roster()) { + draining.set(true); + drainSnapshot(roster(), deadline); + List stragglers = roster(); + if (!stragglers.isEmpty()) { + log.warn("drain sweep found {} session(s) registered after the drain snapshot was " + + "taken (raced past the shutdown guard); draining them too", stragglers.size()); + drainSnapshot(stragglers, deadline); + } + } + + /** + * Drain exactly the sessions in {@code snapshot}, waiting out a {@code BUSY} one against the + * shared whole-drain {@code deadline} before releasing it. Shared by {@link #drainAll}'s main + * pass and its post-loop straggler sweep (fleetd #308) so both honor the same one budget. + */ + private void drainSnapshot(List snapshot, long deadline) { + for (MemberSession s : snapshot) { try { if (s.state() == MemberSession.State.BUSY) { while (System.nanoTime() < deadline) { diff --git a/fleetd/src/main/java/dev/ltms/fleet/session/ShuttingDownException.java b/fleetd/src/main/java/dev/ltms/fleet/session/ShuttingDownException.java new file mode 100644 index 0000000..5652bdb --- /dev/null +++ b/fleetd/src/main/java/dev/ltms/fleet/session/ShuttingDownException.java @@ -0,0 +1,18 @@ +package dev.ltms.fleet.session; + +/** + * Thrown by {@link SessionManager#acquire} when a spawn is requested after the daemon's shutdown + * drain has already begun (fleetd #308). + * + *

{@link SessionManager#drainAll} snapshots the registry once and tears down exactly what is + * in that snapshot. A session registered after the snapshot is invisible to the drain loop: its + * pane is left running and its worktree is never preserved, and nothing else ever reclaims + * either — the daemon's in-memory registry dies with the process. Refusing the spawn here, + * loudly, is what stops that session from ever being created in the first place, rather than + * silently handing the caller a session the daemon can no longer manage. + */ +public final class ShuttingDownException extends RuntimeException { + public ShuttingDownException(String message) { + super(message); + } +} 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..3aab11d 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/session/SessionManagerTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/session/SessionManagerTest.java @@ -26,7 +26,12 @@ import org.slf4j.LoggerFactory; import java.util.List; import java.util.Map; import java.util.Set; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; import java.util.function.LongSupplier; import static org.junit.jupiter.api.Assertions.*; @@ -824,6 +829,186 @@ class SessionManagerTest { .count(); } + // --- fleetd #308: a spawn accepted while the shutdown drain is running must not orphan --- + + @Test + void acquireRefusesANewSpawnOnceDrainAllHasStarted() { + FakeHerdr herdr = new FakeHerdr(); + SessionManager sessions = sessionManager(herdr); + + sessions.drainAll(TimeUnit.MILLISECONDS.toNanos(50)); // empty roster — returns immediately, + // but the shutdown guard it flips must stay tripped for the life of the process. + + ShuttingDownException e = assertThrows(ShuttingDownException.class, + () -> sessions.acquire("ltms-local", null, "/caller", "term_primary"), + "a spawn requested after the drain has begun must be refused loudly (invariant 3), " + + "not silently registered into a registry the drain will never revisit"); + assertNotNull(e.getMessage()); + assertFalse(e.getMessage().isBlank(), "the refusal must say why, not just that it failed"); + assertTrue(sessions.roster().isEmpty(), "the refused spawn must never reach the registry"); + } + + /** + * fleetd #308: the guard above closes most of the shutdown-race window, but it cannot close + * all of it — a caller that already passed the {@code draining} check before {@code drainAll} + * flips it can still be mid-{@code launcher.spawn()} (a real herdr round trip, not + * instantaneous) when {@code drainAll} takes its registry snapshot. This test forces exactly + * that interleaving with a launcher double that blocks the second {@code spawn()} call and the + * first {@code stop()} call until released, then proves the post-loop sweep in {@code + * drainAll} still finds and tears down the straggler that lands in the registry afterward. + */ + @Test + void drainAllSweepsAStragglerThatRegisteredAfterTheInitialSnapshot() throws Exception { + FakeHerdr herdr = new FakeHerdr(); + 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 delegate = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr), + new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null); + RaceLauncher race = new RaceLauncher(delegate); + SessionManager sessions = new SessionManager(race); + + // Registered normally, before the drain starts — the first spawn call, never blocked. + MemberSession ready = sessions.acquire("ltms-local", "/ready", "/caller", "ownerR"); + + ExecutorService exec = Executors.newFixedThreadPool(2); + try { + // The straggler's acquire() reads `draining == false` (checked before this call ever + // touches the launcher) and then blocks inside its own spawn() — the second spawn call. + Future straggler = exec.submit(() -> + sessions.acquire("ltms-local", "/late", "/caller", "ownerLate")); + + assertTrue(race.enteredSecondSpawn.await(5, TimeUnit.SECONDS), + "the straggler must have passed the shutdown guard and reached spawn() before " + + "drainAll ever runs"); + assertEquals(1, sessions.roster().size(), + "the straggler is still inside spawn() — not registered yet"); + + // drainAll flips `draining`, snapshots the registry (only `ready` is in it), and starts + // releasing that snapshot — its first release() call stops `ready`'s pane, which this + // launcher double blocks on so the interleaving below is deterministic, not a timing bet. + Future drain = exec.submit(() -> sessions.drainAll(TimeUnit.SECONDS.toNanos(5))); + + assertTrue(race.enteredFirstStop.await(5, TimeUnit.SECONDS), + "drainAll must be stopping the ready session's pane — proof its initial " + + "registry snapshot has already been taken"); + + // Only now does the straggler's spawn complete and register — strictly after the + // snapshot drainAll's main pass is working from. + race.releaseSecondSpawn.countDown(); + MemberSession registered = straggler.get(5, TimeUnit.SECONDS); + + // Let drainAll finish releasing `ready`; it then re-checks the registry and must find + // (and drain) the straggler that just landed in it. + race.releaseFirstStop.countDown(); + drain.get(5, TimeUnit.SECONDS); + + assertTrue(sessions.roster().isEmpty(), + "the post-loop sweep must drain the straggler too, not just the initial snapshot"); + assertNotNull(registered.paneId()); + long paneCloseCalls = herdr.calls.stream().filter(c -> "pane.close".equals(c.method())).count(); + assertEquals(2, paneCloseCalls, + "both ready's pane AND the straggler's pane must actually be stopped — a pane " + + "left running is exactly the orphan this ticket is about"); + } finally { + exec.shutdownNow(); + } + } + + /** + * Delegates every call while blocking the SECOND {@code spawn()} call and the FIRST + * {@code stop()} call until the test releases them — used to force the fleetd #308 race + * deterministically instead of betting on real thread-scheduling timing. + */ + private static final class RaceLauncher implements PeerLauncher { + private final PeerLauncher delegate; + private final AtomicInteger spawnCalls = new AtomicInteger(); + private final AtomicInteger stopCalls = new AtomicInteger(); + final CountDownLatch enteredSecondSpawn = new CountDownLatch(1); + final CountDownLatch releaseSecondSpawn = new CountDownLatch(1); + final CountDownLatch enteredFirstStop = new CountDownLatch(1); + final CountDownLatch releaseFirstStop = new CountDownLatch(1); + + RaceLauncher(PeerLauncher delegate) { + this.delegate = delegate; + } + + private static void awaitOrFail(CountDownLatch latch) { + try { + if (!latch.await(5, TimeUnit.SECONDS)) { + throw new AssertionError("RaceLauncher latch timed out"); + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new AssertionError("RaceLauncher latch interrupted", e); + } + } + + @Override + public Set capabilities() { + return delegate.capabilities(); + } + + @Override + public Set capabilitiesFor(String profileName) { + return delegate.capabilitiesFor(profileName); + } + + @Override + public PeerHandle spawn(SpawnRequest req) { + if (spawnCalls.incrementAndGet() == 2) { + enteredSecondSpawn.countDown(); + awaitOrFail(releaseSecondSpawn); + } + return delegate.spawn(req); + } + + @Override + public Set profiles() { + return delegate.profiles(); + } + + @Override + public String defaultProfile() { + return delegate.defaultProfile(); + } + + @Override + public String effectiveCwd(SpawnRequest req) { + return delegate.effectiveCwd(req); + } + + @Override + public List parityOverlay(String profileName) { + return delegate.parityOverlay(profileName); + } + + @Override + public List list() { + return delegate.list(); + } + + @Override + public int reapOrphanWorkers() { + return delegate.reapOrphanWorkers(); + } + + @Override + public void stop(String id) { + if (stopCalls.incrementAndGet() == 1) { + enteredFirstStop.countDown(); + awaitOrFail(releaseFirstStop); + } + delegate.stop(id); + } + + @Override + public boolean clearContext(String id) { + return delegate.clearContext(id); + } + } + private static List promptTexts(FakeHerdr herdr) { return herdr.calls.stream() .filter(c -> "agent.prompt".equals(c.method()))