diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java index 30e92b2..d0ec531 100644 --- a/bridged/src/main/java/dev/ltms/bridged/Bridged.java +++ b/bridged/src/main/java/dev/ltms/bridged/Bridged.java @@ -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 {}", diff --git a/bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java b/bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java index 9ba082d..1e74fb9 100644 --- a/bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java +++ b/bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java @@ -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(); diff --git a/bridged/src/test/java/dev/ltms/bridged/session/SessionManagerTest.java b/bridged/src/test/java/dev/ltms/bridged/session/SessionManagerTest.java index e062db6..b2fff3f 100644 --- a/bridged/src/test/java/dev/ltms/bridged/session/SessionManagerTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/session/SessionManagerTest.java @@ -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()))