Merge #308: refuse spawns once the shutdown drain has started, and sweep stragglers
This commit is contained in:
@@ -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.
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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;
|
||||
@@ -73,6 +74,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<Consumer<String>> acquireListeners = new java.util.concurrent.CopyOnWriteArrayList<>();
|
||||
/** CB-516: notified with a {@link ReleaseDetail} on every release; no-op until wired. */
|
||||
@@ -193,6 +204,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
|
||||
@@ -922,10 +941,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.
|
||||
*
|
||||
* <p>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<MemberSession> 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<MemberSession> snapshot, long deadline) {
|
||||
for (MemberSession s : snapshot) {
|
||||
try {
|
||||
if (s.state() == MemberSession.State.BUSY) {
|
||||
while (System.nanoTime() < deadline) {
|
||||
|
||||
@@ -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).
|
||||
*
|
||||
* <p>{@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);
|
||||
}
|
||||
}
|
||||
@@ -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.*;
|
||||
@@ -847,6 +852,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<MemberSession> 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<Capability> capabilities() {
|
||||
return delegate.capabilities();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<Capability> 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<String> profiles() {
|
||||
return delegate.profiles();
|
||||
}
|
||||
|
||||
@Override
|
||||
public String defaultProfile() {
|
||||
return delegate.defaultProfile();
|
||||
}
|
||||
|
||||
@Override
|
||||
public String effectiveCwd(SpawnRequest req) {
|
||||
return delegate.effectiveCwd(req);
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<String> 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<String> promptTexts(FakeHerdr herdr) {
|
||||
return herdr.calls.stream()
|
||||
.filter(c -> "agent.prompt".equals(c.method()))
|
||||
|
||||
Reference in New Issue
Block a user