CB-303 part 3: graceful drain on shutdown

This commit is contained in:
Dai Ha
2026-07-17 09:59:57 +02:00
parent 954351a80b
commit 09d3948acf
3 changed files with 43 additions and 10 deletions
@@ -53,7 +53,6 @@ public final class Bridged {
: UnixSocketHerdrClient.defaultSocketPath();
UnixSocketHerdrClient herdr = UnixSocketHerdrClient.connect(socket, new com.fasterxml.jackson.databind.ObjectMapper());
Runtime.getRuntime().addShutdownHook(new Thread(herdr::close));
AgentControl agents = new AgentControl(herdr);
WorkspaceControl spaces = new WorkspaceControl(herdr);
@@ -74,13 +73,14 @@ public final class Bridged {
SessionManager sessions = new SessionManager(workers, new GitWorktrees(cfg.worktreeRoot()), contextCap);
// CB-303 part 1: idle-ttl reaper — only when configured, defaults to disabled.
SessionReaper reaper = null;
final SessionReaper reaper;
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));
} else {
reaper = null;
}
// Status-gated injector (CB-103): the single writer into workers, fed by a poller.
@@ -113,17 +113,27 @@ public final class Bridged {
Injector injector = new Injector(agents, turnListener, presence::isPresent, presence::forget);
StatusPoller poller = new StatusPoller(agents, injector, INJECT_POLL_MILLIS);
poller.start();
Runtime.getRuntime().addShutdownHook(new Thread(poller::stop));
MessageService messages = new MessageService(agents, injector, rendezvous);
Runtime.getRuntime().addShutdownHook(new Thread(messages::close));
// MCP server face (CB-105): bridge_send/bridge_reply/bridge_status, mounted at /mcp.
// Caller identity is resolved from the connection (peer PID → herdr pane), not arguments.
ConnectionIdentity identity = new ConnectionIdentity(
new PaneLocator(herdr), new LsofPeerPidLookup(), new LsofProcessCwdLookup());
BridgeMcp mcp = new BridgeMcp(messages, rendezvous, workers, sessions, identity, presence);
Runtime.getRuntime().addShutdownHook(new Thread(mcp::close));
// CB-303 part 3: single ordered shutdown hook. Drain sessions first while herdr is still
// open (so releases reach the daemon), then stop poller/message/mcp/reaper, and close herdr
// last. This replaces the earlier independent hooks that could race and close herdr early.
Runtime.getRuntime().addShutdownHook(new Thread(() -> {
sessions.close(cfg.lifecycle() != null ? cfg.lifecycle().drainTimeoutSeconds() : null);
poller.stop();
messages.close();
mcp.close();
if (reaper != null) reaper.stop();
herdr.close();
}));
Javalin app = new BridgedApp(herdr, workers, sessions, messages, rendezvous, presence, mcp.servlet()).build();
app.start(cfg.bind().host(), cfg.bind().port());
log.info("bridged listening on {}:{}, herdr socket {}",
@@ -299,18 +299,17 @@ public final class SessionManager implements TurnListener {
* abort the rest.
*/
void drainAll(long timeoutNanos) {
long deadline = nowNanos.getAsLong() + timeoutNanos;
long pollNanos = TimeUnit.MILLISECONDS.toNanos(50);
long deadline = System.nanoTime() + timeoutNanos;
for (WorkerSession s : roster()) {
try {
if (s.state() == WorkerSession.State.BUSY) {
while (nowNanos.getAsLong() < deadline) {
while (System.nanoTime() < deadline) {
WorkerSession current = registry.get(s.paneId());
if (current == null || current.state() != WorkerSession.State.BUSY) {
break;
}
try {
long remaining = deadline - nowNanos.getAsLong();
long remaining = deadline - System.nanoTime();
Thread.sleep(Math.min(TimeUnit.NANOSECONDS.toMillis(remaining), 50));
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
@@ -11,6 +11,7 @@ import org.junit.jupiter.api.Test;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.TimeUnit;
import java.util.function.LongSupplier;
import static org.junit.jupiter.api.Assertions.*;
@@ -302,6 +303,29 @@ class SessionManagerTest {
"forced release tears the worker pane down exactly once");
}
@Test
void drainAllReleasesBusyAndReadySessionsAndWaitsForBusy() {
long[] clock = {0};
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr, () -> clock[0]);
WorkerSession ready = sessions.acquire("ltms-local", "/ready", "/caller", "ownerR");
WorkerSession busy = sessions.acquire("ltms-local", "/busy", "/caller", "ownerB");
sessions.asPresence().markPresent(ready.terminalId());
sessions.asPresence().markPresent(busy.terminalId());
sessions.onDelivered(busy.terminalId());
sessions.drainAll(TimeUnit.MILLISECONDS.toNanos(100));
assertTrue(sessions.roster().isEmpty(), "drain clears the roster");
assertTrue(sessions.get(ready.paneId()).isEmpty(), "ready session is released");
assertTrue(sessions.get(busy.paneId()).isEmpty(), "busy session is released after timeout");
assertEquals(1, paneCloseCallsFor(herdr, ready.paneId()),
"ready worker pane is torn down");
assertEquals(1, paneCloseCallsFor(herdr, busy.paneId()),
"busy worker pane is torn down");
}
private static long paneCloseCallsFor(FakeHerdr herdr, String paneId) {
return herdr.calls.stream()
.filter(c -> "pane.close".equals(c.method()))