From 54d907c31430bfadb674a2fa19c1c0fdf851eafa Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 16 Jul 2026 19:30:56 +0200 Subject: [PATCH] =?UTF-8?q?CB-301:=20SessionManager=20=E2=80=94=20authorit?= =?UTF-8?q?ative=20worker=20session=20registry=20+=20one-shot=20FSM?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Adds dev.ltms.bridged.session with WorkerSession (immutable record) and SessionManager wrapping WorkerService: a ConcurrentHashMap registry keyed by paneId, the one-shot lifecycle FSM (SPAWNING->READY->BUSY->DONE, ->FAILED on drop/turn-failure, ->RELEASED on teardown), ownership (ownerTerminal), and recycle = release + fresh acquire (no-reuse invariant). Driven by TurnListener (BUSY/DONE/FAILED) and a WorkerPresence bridge (READY). Wiring: Bridged.main constructs it and composes it into the TurnListener alongside CompletionResolver; bridge_spawn / POST /workers route through acquire (carrying caller identity as owner); bridge_stop / DELETE /workers route through release. WorkerService gains effectiveCwd(); WorkerPresence de-finalized so the manager can present a READY-driving view. asPresence() returns a single cached bridge (a fresh one per call would fragment the shared present set). roster() is the registry snapshot; the live herdr join is left for CB-304. 6 fake-based acceptance tests; full suite green (155/155). Delegated to an off-subscription worker against docs/CB-301-Session-Manager.md; primary verified (ide diagnostics clean, mvn clean install green) + fixed the asPresence caching bug. --- .../main/java/dev/ltms/bridged/Bridged.java | 33 +++- .../ltms/bridged/inject/WorkerPresence.java | 2 +- .../java/dev/ltms/bridged/mcp/BridgeMcp.java | 38 ++-- .../dev/ltms/bridged/rest/BridgedApp.java | 27 ++- .../ltms/bridged/session/SessionManager.java | 185 ++++++++++++++++++ .../ltms/bridged/session/WorkerSession.java | 39 ++++ .../ltms/bridged/worker/WorkerService.java | 19 ++ .../dev/ltms/bridged/herdr/FakeHerdr.java | 8 +- .../dev/ltms/bridged/mcp/BridgeMcpTest.java | 25 ++- .../dev/ltms/bridged/rest/BridgedAppTest.java | 11 +- .../bridged/session/SessionManagerTest.java | 150 ++++++++++++++ 11 files changed, 501 insertions(+), 36 deletions(-) create mode 100644 bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java create mode 100644 bridged/src/main/java/dev/ltms/bridged/session/WorkerSession.java create mode 100644 bridged/src/test/java/dev/ltms/bridged/session/SessionManagerTest.java diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java index 5371411..caeea80 100644 --- a/bridged/src/main/java/dev/ltms/bridged/Bridged.java +++ b/bridged/src/main/java/dev/ltms/bridged/Bridged.java @@ -9,6 +9,7 @@ import dev.ltms.bridged.herdr.WorkspaceControl; import dev.ltms.bridged.inject.CompletionResolver; import dev.ltms.bridged.inject.Injector; import dev.ltms.bridged.inject.StatusPoller; +import dev.ltms.bridged.inject.TurnListener; import dev.ltms.bridged.inject.WorkerPresence; import dev.ltms.bridged.mcp.BridgeMcp; import dev.ltms.bridged.mcp.ConnectionIdentity; @@ -17,6 +18,7 @@ import dev.ltms.bridged.mcp.LsofProcessCwdLookup; import dev.ltms.bridged.msg.MessageService; import dev.ltms.bridged.msg.Rendezvous; import dev.ltms.bridged.rest.BridgedApp; +import dev.ltms.bridged.session.SessionManager; import dev.ltms.bridged.worker.WorkerService; import io.javalin.Javalin; import org.slf4j.Logger; @@ -59,14 +61,37 @@ public final class Bridged { // the previous process — reap those leaked orphans now, before we start serving. workers.reapOrphanWorkers(); + // CB-301: authoritative session registry + lifecycle FSM on top of WorkerService. + SessionManager sessions = new SessionManager(workers); + // Status-gated injector (CB-103): the single writer into workers, fed by a poller. // The blocking message endpoint (CB-104) is the producer; the poller is inert until then. // CB-106: a confirmed turn completion resolves a blocked send whose worker never replied. Rendezvous rendezvous = new Rendezvous(); CompletionResolver completion = new CompletionResolver(agents, rendezvous); // CB-113: deliver only to an available worker (its MCP is connected), never its boot window. - WorkerPresence presence = new WorkerPresence(); - Injector injector = new Injector(agents, completion, presence::isPresent, presence::forget); + // CB-301: the manager's presence bridge records availability and drives SPAWNING → READY. + WorkerPresence presence = sessions.asPresence(); + TurnListener turnListener = new TurnListener() { + @Override + public void onTurnComplete(String target) { + completion.onTurnComplete(target); + sessions.onTurnComplete(target); + } + + @Override + public void onDelivered(String target) { + completion.onDelivered(target); + sessions.onDelivered(target); + } + + @Override + public void onTurnFailed(String target) { + completion.onTurnFailed(target); + sessions.onTurnFailed(target); + } + }; + 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)); @@ -78,9 +103,9 @@ public final class Bridged { // 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, identity, presence); + BridgeMcp mcp = new BridgeMcp(messages, rendezvous, workers, sessions, identity, presence); Runtime.getRuntime().addShutdownHook(new Thread(mcp::close)); - Javalin app = new BridgedApp(herdr, workers, messages, rendezvous, presence, mcp.servlet()).build(); + 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 {}", cfg.bind().host(), cfg.bind().port(), socket); diff --git a/bridged/src/main/java/dev/ltms/bridged/inject/WorkerPresence.java b/bridged/src/main/java/dev/ltms/bridged/inject/WorkerPresence.java index 5b31c16..1fcbbd2 100644 --- a/bridged/src/main/java/dev/ltms/bridged/inject/WorkerPresence.java +++ b/bridged/src/main/java/dev/ltms/bridged/inject/WorkerPresence.java @@ -15,7 +15,7 @@ import java.util.Set; * mounts the bridge MCP is never marked present — its sends stay queued until they time out, which is * correct (it could not have replied anyway). */ -public final class WorkerPresence { +public class WorkerPresence { private final Set present = ConcurrentHashMap.newKeySet(); diff --git a/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java index f9ccf2a..ce7c0b0 100644 --- a/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java +++ b/bridged/src/main/java/dev/ltms/bridged/mcp/BridgeMcp.java @@ -6,6 +6,8 @@ import dev.ltms.bridged.inject.WorkerPresence; import dev.ltms.bridged.herdr.HerdrException; import dev.ltms.bridged.msg.MessageService; import dev.ltms.bridged.msg.Rendezvous; +import dev.ltms.bridged.session.SessionManager; +import dev.ltms.bridged.session.WorkerSession; import dev.ltms.bridged.worker.WorkerService; import io.modelcontextprotocol.common.McpTransportContext; import io.modelcontextprotocol.json.McpJsonMapper; @@ -55,7 +57,7 @@ public final class BridgeMcp { private final McpSyncServer server; public BridgeMcp(MessageService messages, Rendezvous rendezvous, WorkerService workers, - ConnectionIdentity identity, WorkerPresence presence) { + SessionManager sessions, ConnectionIdentity identity, WorkerPresence presence) { McpJsonMapper json = new JacksonMcpJsonMapperSupplier().get(); this.transport = HttpServletStreamableServerTransportProvider.builder() .jsonMapper(json) @@ -102,11 +104,13 @@ public final class BridgeMcp { .toolCall(spawnTool(), (exchange, req) -> { Map a = req.arguments(); // CB-112: worker inherits the primary's cwd unless the call pins one. + // CB-301: carry the caller's identity as the session owner (null for the primary). String callerCwd = identity.cwdForPid(callerPid(exchange)); - return spawn(workers, str(a, "profile"), str(a, "cwd"), callerCwd); + return spawn(sessions, str(a, "profile"), str(a, "cwd"), callerCwd, + callerTerminal(exchange)); }) .toolCall(listTool(), (_, _) -> listWorkers(workers)) - .toolCall(stopTool(), (_, req) -> stop(workers, str(req.arguments(), "paneId"))) + .toolCall(stopTool(), (_, req) -> stop(sessions, str(req.arguments(), "paneId"))) .toolCall(profilesTool(), (_, _) -> profiles(workers)) .build(); } @@ -275,22 +279,25 @@ public final class BridgeMcp { } } - // --- fleet management logic (CB-108) ------------------------------------------------------- + // --- fleet management logic (CB-108 / CB-301) -------------------------------------------- /** {@code bridge_spawn} without cwd/caller context (default resolution). */ - static McpSchema.CallToolResult spawn(WorkerService workers, String profile) { - return spawn(workers, profile, null, null); + static McpSchema.CallToolResult spawn(SessionManager sessions, String profile) { + return spawn(sessions, profile, null, null, null); } /** * {@code bridge_spawn}: launch a guard-checked worker for {@code profile} (blank → the default * profile) and return its session id + pane id. The worker's cwd is {@code requestedCwd} if given, * else the profile's config, else {@code callerCwd} (the primary's directory), else the daemon's. + * CB-301: the session is registered with {@code ownerTerminal} as its owner. */ - static McpSchema.CallToolResult spawn(WorkerService workers, String profile, - String requestedCwd, String callerCwd) { + static McpSchema.CallToolResult spawn(SessionManager sessions, String profile, + String requestedCwd, String callerCwd, + String ownerTerminal) { try { - Agent worker = workers.spawn(isBlank(profile) ? null : profile, requestedCwd, callerCwd); + WorkerSession worker = sessions.acquire(isBlank(profile) ? null : profile, + requestedCwd, callerCwd, ownerTerminal); return text(json(workerView(worker))); } catch (GuardException e) { return error("subscription boundary: " + e.getMessage()); @@ -319,12 +326,12 @@ public final class BridgeMcp { } /** {@code bridge_stop}: tear a worker down by its pane id. */ - static McpSchema.CallToolResult stop(WorkerService workers, String paneId) { + static McpSchema.CallToolResult stop(SessionManager sessions, String paneId) { if (isBlank(paneId)) { return error("paneId is required"); } try { - workers.stop(paneId); + sessions.release(paneId); return text("stopped " + paneId); } catch (HerdrException e) { return error("herdr error stopping " + paneId + ": " + e.getMessage()); @@ -340,6 +347,15 @@ public final class BridgeMcp { return m; } + /** CB-301 projection from the authoritative session registry. */ + private static Map workerView(WorkerSession s) { + Map m = new LinkedHashMap<>(); + m.put("sessionId", s.terminalId()); + m.put("paneId", s.paneId()); + m.put("status", s.state().name().toLowerCase()); + return m; + } + private static String json(Object o) { try { return MAPPER.writeValueAsString(o); diff --git a/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java b/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java index f3c64d4..a39ee9c 100644 --- a/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java +++ b/bridged/src/main/java/dev/ltms/bridged/rest/BridgedApp.java @@ -9,6 +9,8 @@ import dev.ltms.bridged.herdr.HerdrException; import dev.ltms.bridged.inject.WorkerPresence; import dev.ltms.bridged.msg.MessageService; import dev.ltms.bridged.msg.Rendezvous; +import dev.ltms.bridged.session.SessionManager; +import dev.ltms.bridged.session.WorkerSession; import dev.ltms.bridged.worker.WorkerService; import io.javalin.Javalin; import io.javalin.http.Context; @@ -40,16 +42,19 @@ public final class BridgedApp { private final HerdrClient herdr; private final WorkerService workers; + private final SessionManager sessions; // CB-301: authoritative session registry private final MessageService messages; private final Rendezvous rendezvous; private final WorkerPresence presence; // CB-113: which workers are MCP-connected (available) private final HttpServlet mcpServlet; // MCP Streamable-HTTP endpoint, mounted at /mcp (nullable) private final ObjectMapper mapper = new ObjectMapper(); - public BridgedApp(HerdrClient herdr, WorkerService workers, MessageService messages, - Rendezvous rendezvous, WorkerPresence presence, HttpServlet mcpServlet) { + public BridgedApp(HerdrClient herdr, WorkerService workers, SessionManager sessions, + MessageService messages, Rendezvous rendezvous, WorkerPresence presence, + HttpServlet mcpServlet) { this.herdr = herdr; this.workers = workers; + this.sessions = sessions; this.messages = messages; this.rendezvous = rendezvous; this.presence = presence; @@ -145,8 +150,8 @@ public final class BridgedApp { } } try { - // No MCP caller over REST, so callerCwd is null → the daemon cwd is the last resort. - Agent worker = workers.spawn(blankToNull(profile), blankToNull(cwd), null); + // No MCP caller over REST, so callerCwd and ownerTerminal are null. + WorkerSession worker = sessions.acquire(blankToNull(profile), blankToNull(cwd), null, null); ctx.status(201).json(view(worker)); } catch (GuardException e) { ctx.status(403).json(Map.of("error", "subscription_boundary", "detail", e.getMessage())); @@ -161,7 +166,7 @@ public final class BridgedApp { /** Tear a worker down by pane id. */ private void stopWorker(Context ctx) { - workers.stop(ctx.pathParam("paneId")); + sessions.release(ctx.pathParam("paneId")); ctx.status(204); } @@ -362,4 +367,16 @@ public final class BridgedApp { m.put("status", a.status().name().toLowerCase()); return m; } + + /** CB-301 projection of an authoritative bridge-owned session. */ + private static Map view(WorkerSession s) { + Map m = new LinkedHashMap<>(); + m.put("terminalId", s.terminalId()); + m.put("paneId", s.paneId()); + m.put("profile", s.profile()); + m.put("cwd", s.cwd()); + m.put("ownerTerminal", s.ownerTerminal()); + m.put("state", s.state().name().toLowerCase()); + return m; + } } diff --git a/bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java b/bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java new file mode 100644 index 0000000..7c85636 --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java @@ -0,0 +1,185 @@ +package dev.ltms.bridged.session; + +import dev.ltms.bridged.herdr.Agent; +import dev.ltms.bridged.inject.TurnListener; +import dev.ltms.bridged.inject.WorkerPresence; +import dev.ltms.bridged.worker.WorkerService; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.util.List; +import java.util.Optional; +import java.util.concurrent.ConcurrentHashMap; + +/** + * Authoritative in-daemon registry of the worker sessions this {@code bridged} process spawned. + * Wraps {@link WorkerService} (it still performs subscription-guarded spawn/teardown) and adds + * lifecycle tracking, ownership, and deterministic teardown on top. + * + *

The state machine is intentionally one-shot / no-reuse: every acquired worker is fresh, + * and a finished or released worker is torn down, never pooled. {@link #recycle} is a convenience + * for {@code release + acquire} with a new distinct pane id. + * + *

The manager implements {@link TurnListener} so the injector's turn boundaries drive + * {@code READY → BUSY → DONE} (or {@code FAILED}). It exposes a {@link WorkerPresence} view via + * {@link #asPresence()}: any MCP contact from a worker marks it present and simultaneously + * transitions the session {@code SPAWNING → READY}. + */ +public final class SessionManager implements TurnListener { + + private static final Logger log = LoggerFactory.getLogger(SessionManager.class); + + private final WorkerService workerService; + private final ConcurrentHashMap registry = new ConcurrentHashMap<>(); + private final WorkerPresence presence; + + public SessionManager(WorkerService workerService) { + this.workerService = workerService; + this.presence = new PresenceBridge(this); + } + + /** + * The single {@link WorkerPresence} view of this manager: it records availability and forwards + * the signal to the {@code SPAWNING → READY} transition. Pass this to the {@code Injector} and + * {@code BridgeMcp} where they previously accepted a plain {@link WorkerPresence}. The same + * instance is returned every call — presence is shared state, so a fresh bridge per call would + * fragment the {@code present} set and lose signals across callers. + */ + public WorkerPresence asPresence() { + return presence; + } + + /** + * Spawn a worker and register it as {@link WorkerSession.State#SPAWNING}. The caller's + * identity is recorded as {@code ownerTerminal} ({@code null} for daemon/anon callers). + */ + public WorkerSession acquire(String profile, String requestedCwd, String callerCwd, + String ownerTerminal) { + Agent worker = workerService.spawn(profile, requestedCwd, callerCwd); + String resolvedProfile = (profile == null || profile.isBlank()) + ? workerService.defaultProfile() : profile; + WorkerSession session = new WorkerSession( + worker.paneId(), + worker.terminalId(), + resolvedProfile, + resolveCwd(requestedCwd, profile, callerCwd), + ownerTerminal, + System.nanoTime(), + WorkerSession.State.SPAWNING); + registry.put(session.paneId(), session); + log.debug("acquired session pane={} terminal={} profile={} owner={}", + session.paneId(), session.terminalId(), session.profile(), session.ownerTerminal()); + return session; + } + + /** Tear a worker down by pane id and remove it from the registry. Idempotent. */ + public void release(String paneId) { + WorkerSession removed = registry.remove(paneId); + if (removed != null) { + log.debug("releasing session pane={} terminal={} state={}", + removed.paneId(), removed.terminalId(), removed.state()); + } + workerService.stop(paneId); + } + + /** + * Release the old session and acquire a fresh one with the same profile and working directory. + * The new session is guaranteed to have a pane id distinct from the old one (no-reuse invariant). + */ + public WorkerSession recycle(String paneId) { + WorkerSession old = registry.get(paneId); + if (old == null) { + throw new IllegalArgumentException("no session for paneId " + paneId); + } + release(paneId); + return acquire(old.profile(), old.cwd(), old.cwd(), old.ownerTerminal()); + } + + /** The session for {@code paneId}, if it is still registered and not released. */ + public Optional get(String paneId) { + return Optional.ofNullable(registry.get(paneId)); + } + + /** Bridge-owned roster: all registered sessions (acquired minus released). */ + public List roster() { + return List.copyOf(registry.values()); + } + + /** Lifecycle hook: worker became available on the bridge MCP. */ + void onReady(String terminalId) { + transitionByTerminal(terminalId, WorkerSession.State.SPAWNING, WorkerSession.State.READY); + } + + /** Lifecycle hook: a message was delivered into the worker — it is now busy on a turn. */ + @Override + public void onDelivered(String target) { + transitionByTerminal(target, WorkerSession.State.READY, WorkerSession.State.BUSY); + } + + /** Lifecycle hook: the worker's delegated turn completed successfully. */ + @Override + public void onTurnComplete(String target) { + transitionByTerminal(target, WorkerSession.State.BUSY, WorkerSession.State.DONE); + } + + /** Lifecycle hook: the worker's delegated turn failed. */ + @Override + public void onTurnFailed(String target) { + onFailed(target); + } + + /** Lifecycle hook: the worker vanished or was dropped mid-life. */ + void onFailed(String target) { + WorkerSession current = findByTerminal(target); + if (current == null) return; + if (current.state() == WorkerSession.State.RELEASED) return; + if (replace(current, current.withState(WorkerSession.State.FAILED))) { + log.debug("session marked failed terminal={} pane={}", target, current.paneId()); + } + } + + /** Number of sessions currently registered. */ + public int size() { + return registry.size(); + } + + private WorkerSession findByTerminal(String terminalId) { + for (WorkerSession s : registry.values()) { + if (terminalId.equals(s.terminalId())) return s; + } + return null; + } + + private void transitionByTerminal(String terminalId, WorkerSession.State from, + WorkerSession.State to) { + WorkerSession current = findByTerminal(terminalId); + if (current == null || current.state() != from) return; + if (replace(current, current.withState(to))) { + log.debug("session transitioned terminal={} pane={} {} -> {}", + terminalId, current.paneId(), from, to); + } + } + + private boolean replace(WorkerSession expected, WorkerSession updated) { + return registry.replace(expected.paneId(), expected, updated); + } + + private String resolveCwd(String requestedCwd, String profileName, String callerCwd) { + return workerService.effectiveCwd(profileName, requestedCwd, callerCwd); + } + + /** WorkerPresence bridge that also drives the manager's READY transition. */ + private static final class PresenceBridge extends WorkerPresence { + private final SessionManager sessions; + + PresenceBridge(SessionManager sessions) { + this.sessions = sessions; + } + + @Override + public void markPresent(String terminal) { + super.markPresent(terminal); + sessions.onReady(terminal); + } + } +} diff --git a/bridged/src/main/java/dev/ltms/bridged/session/WorkerSession.java b/bridged/src/main/java/dev/ltms/bridged/session/WorkerSession.java new file mode 100644 index 0000000..1679443 --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/session/WorkerSession.java @@ -0,0 +1,39 @@ +package dev.ltms.bridged.session; + +/** + * A bridge-owned worker session — the authoritative in-daemon record of a worker this + * process spawned. Immutable; state transitions are performed by replacing the record in + * {@link SessionManager}'s registry. + * + * @param paneId herdr pane handle — the registry key and the argument to teardown + * @param terminalId herdr terminal handle — the {@code target} for send/read/status + * @param profile the worker profile name that spawned this session + * @param cwd the resolved working directory the worker started in + * @param ownerTerminal the caller that requested this worker ({@code null} = daemon/anon) + * @param spawnedAtNanos {@link System#nanoTime()} when the session was registered + * @param state current lifecycle state in the one-shot FSM + */ +public record WorkerSession( + String paneId, + String terminalId, + String profile, + String cwd, + String ownerTerminal, + long spawnedAtNanos, + State state) { + + /** One-shot worker lifecycle states. */ + public enum State { + SPAWNING, + READY, + BUSY, + DONE, + FAILED, + RELEASED + } + + /** Return a copy of this session in {@code state}. */ + public WorkerSession withState(State state) { + return new WorkerSession(paneId, terminalId, profile, cwd, ownerTerminal, spawnedAtNanos, state); + } +} diff --git a/bridged/src/main/java/dev/ltms/bridged/worker/WorkerService.java b/bridged/src/main/java/dev/ltms/bridged/worker/WorkerService.java index 1a159f2..ab719ce 100644 --- a/bridged/src/main/java/dev/ltms/bridged/worker/WorkerService.java +++ b/bridged/src/main/java/dev/ltms/bridged/worker/WorkerService.java @@ -160,6 +160,25 @@ public final class WorkerService { return firstNonBlank(requestedCwd, cfg.cwd(), callerCwd, System.getProperty("user.dir"), "."); } + /** + * CB-301: the effective working directory a spawn for {@code profileName} would use, without + * actually spawning. Used by {@link dev.ltms.bridged.session.SessionManager} to record the + * resolved cwd in the session registry. + */ + public String effectiveCwd(String profileName, String requestedCwd, String callerCwd) { + String name = (profileName == null || profileName.isBlank()) ? defaultProfile : profileName; + if (name == null || name.isBlank()) { + throw new IllegalArgumentException("no default worker profile is configured — " + + "pass a profile; configured: " + profiles.keySet()); + } + BridgedConfig.Worker cfg = profiles.get(name); + if (cfg == null) { + throw new IllegalArgumentException("unknown worker profile '" + name + + "' — configured: " + profiles.keySet()); + } + return resolveCwd(requestedCwd, cfg, callerCwd); + } + private static String firstNonBlank(String... values) { for (String v : values) { if (v != null && !v.isBlank()) return v; diff --git a/bridged/src/test/java/dev/ltms/bridged/herdr/FakeHerdr.java b/bridged/src/test/java/dev/ltms/bridged/herdr/FakeHerdr.java index 675b44d..85d7fee 100644 --- a/bridged/src/test/java/dev/ltms/bridged/herdr/FakeHerdr.java +++ b/bridged/src/test/java/dev/ltms/bridged/herdr/FakeHerdr.java @@ -142,10 +142,12 @@ public final class FakeHerdr implements HerdrClient { "herdr error [agent_name_taken]: agent name already used", "agent_name_taken", null); } - yield mapper.readTree(""" + long n = starts - agentNameTakenFor; + yield mapper.readTree((""" {"type":"agent_started","agent":{ - "terminal_id":"term_new","name":"claude","agent_status":"unknown", - "workspace_id":"w9","tab_id":"w9:t2","pane_id":"w9:pW"}}"""); + "terminal_id":"term_new_%d","name":"claude","agent_status":"unknown", + "workspace_id":"w9","tab_id":"w9:t2","pane_id":"w9:pW_%d"}}""") + .formatted(n, n)); } case "workspace.create" -> mapper.readTree(""" {"type":"workspace_created", diff --git a/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java b/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java index 81cd6e8..726c856 100644 --- a/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/mcp/BridgeMcpTest.java @@ -8,6 +8,7 @@ import dev.ltms.bridged.herdr.WorkspaceControl; import dev.ltms.bridged.inject.Injector; import dev.ltms.bridged.msg.MessageService; import dev.ltms.bridged.msg.Rendezvous; +import dev.ltms.bridged.session.SessionManager; import dev.ltms.bridged.worker.WorkerService; import io.modelcontextprotocol.spec.McpSchema; import org.junit.jupiter.api.Test; @@ -43,6 +44,10 @@ class BridgeMcpTest { new SubscriptionGuard(allow), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> "tok"); } + private static SessionManager sessionManager(FakeHerdr h, String baseUrl, Set allow) { + return new SessionManager(workerService(h, baseUrl, allow)); + } + @Test void sendThenReplyRoundTrips() throws Exception { // bridge_send blocks; bridge_reply resolves it with the worker's structured answer. @@ -180,18 +185,20 @@ class BridgeMcpTest { @Test void spawnReturnsTheNewWorkersSessionAndPane() { FakeHerdr h = new FakeHerdr(); - McpSchema.CallToolResult res = BridgeMcp.spawn(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), null); + McpSchema.CallToolResult res = BridgeMcp.spawn( + sessionManager(h, "http://gx00.gw:8000", Set.of("gx00.gw")), null); assertNotEquals(Boolean.TRUE, res.isError()); String out = textOf(res); - assertTrue(out.contains("\"sessionId\":\"term_new\""), out); - assertTrue(out.contains("\"paneId\":\"w9:pW\""), out); + assertTrue(out.contains("\"sessionId\":\"term_new_1\""), out); + assertTrue(out.contains("\"paneId\":\"w9:pW_1\""), out); + assertTrue(out.contains("\"status\":\"spawning\""), out); } @Test void spawnRejectsAnOffAllowlistProfileWithoutTouchingHerdr() { FakeHerdr h = new FakeHerdr(); McpSchema.CallToolResult res = - BridgeMcp.spawn(workerService(h, "https://api.anthropic.com", Set.of("gx00.gw")), null); + BridgeMcp.spawn(sessionManager(h, "https://api.anthropic.com", Set.of("gx00.gw")), null); assertTrue(res.isError()); assertTrue(textOf(res).contains("subscription boundary")); assertFalse(h.called("agent.start"), "the guard must block before any spawn"); @@ -201,7 +208,7 @@ class BridgeMcpTest { void spawnRejectsAnUnknownProfileAsAnError() { FakeHerdr h = new FakeHerdr(); McpSchema.CallToolResult res = - BridgeMcp.spawn(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), "nope"); + BridgeMcp.spawn(sessionManager(h, "http://gx00.gw:8000", Set.of("gx00.gw")), "nope"); assertTrue(res.isError()); assertTrue(textOf(res).contains("unknown worker profile"), textOf(res)); } @@ -210,7 +217,7 @@ class BridgeMcpTest { void spawnPassesTheRequestedCwdToTheWorker() { FakeHerdr h = new FakeHerdr(); McpSchema.CallToolResult res = BridgeMcp.spawn( - workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), null, "/req/dir", null); + sessionManager(h, "http://gx00.gw:8000", Set.of("gx00.gw")), null, "/req/dir", null, null); assertNotEquals(Boolean.TRUE, res.isError()); @SuppressWarnings("unchecked") Map start = (Map) h.lastCall("agent.start").params(); @@ -238,7 +245,8 @@ class BridgeMcpTest { @Test void stopTearsDownAWorkerByPane() { FakeHerdr h = new FakeHerdr(); - McpSchema.CallToolResult res = BridgeMcp.stop(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), "w9:pW"); + McpSchema.CallToolResult res = BridgeMcp.stop( + sessionManager(h, "http://gx00.gw:8000", Set.of("gx00.gw")), "w9:pW"); assertNotEquals(Boolean.TRUE, res.isError()); assertEquals("stopped w9:pW", textOf(res)); assertTrue(h.called("pane.close")); @@ -247,7 +255,8 @@ class BridgeMcpTest { @Test void stopRequiresAPaneId() { FakeHerdr h = new FakeHerdr(); - assertTrue(BridgeMcp.stop(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), " ").isError()); + assertTrue(BridgeMcp.stop( + sessionManager(h, "http://gx00.gw:8000", Set.of("gx00.gw")), " ").isError()); } @Test diff --git a/bridged/src/test/java/dev/ltms/bridged/rest/BridgedAppTest.java b/bridged/src/test/java/dev/ltms/bridged/rest/BridgedAppTest.java index 9c3fce1..d615237 100644 --- a/bridged/src/test/java/dev/ltms/bridged/rest/BridgedAppTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/rest/BridgedAppTest.java @@ -12,6 +12,7 @@ import dev.ltms.bridged.inject.StatusPoller; import dev.ltms.bridged.inject.WorkerPresence; import dev.ltms.bridged.msg.MessageService; import dev.ltms.bridged.msg.Rendezvous; +import dev.ltms.bridged.session.SessionManager; import dev.ltms.bridged.worker.WorkerService; import io.javalin.Javalin; import org.junit.jupiter.api.AfterEach; @@ -37,7 +38,7 @@ class BridgedAppTest { private final ObjectMapper mapper = new ObjectMapper(); private final HttpClient http = HttpClient.newHttpClient(); - private final WorkerPresence presence = new WorkerPresence(); + private WorkerPresence presence; private Javalin app; private StatusPoller poller; @@ -60,12 +61,14 @@ class BridgedAppTest { agents, new WorkspaceControl(herdr), new SubscriptionGuard(allow), Map.of(wcfg.profile(), wcfg), wcfg.profile(), k -> "BRIDGED_WORKER_TOKEN".equals(k) ? "tok-abc" : null); + SessionManager sessions = new SessionManager(workers); + this.presence = sessions.asPresence(); Injector injector = new Injector(agents); poller = new StatusPoller(agents, injector, 5); // delivers when the fake reports idle poller.start(); Rendezvous rendezvous = new Rendezvous(); MessageService messages = new MessageService(agents, injector, rendezvous); - app = new BridgedApp(herdr, workers, messages, rendezvous, this.presence, null) + app = new BridgedApp(herdr, workers, sessions, messages, rendezvous, this.presence, null) .build().start("127.0.0.1", 0); return app.port(); } @@ -146,8 +149,8 @@ class BridgedAppTest { HttpResponse res = req(port, "POST", "/workers"); assertEquals(201, res.statusCode()); JsonNode body = mapper.readTree(res.body()); - assertEquals("w9:pW", body.get("paneId").asText()); - assertEquals("w9:t2", body.get("tabId").asText()); + assertEquals("w9:pW_1", body.get("paneId").asText()); + assertEquals("spawning", body.get("state").asText()); // Subscription boundary: agent.start carried base_url + token in its env map. Map start = params(herdr, "agent.start"); diff --git a/bridged/src/test/java/dev/ltms/bridged/session/SessionManagerTest.java b/bridged/src/test/java/dev/ltms/bridged/session/SessionManagerTest.java new file mode 100644 index 0000000..5aebbcf --- /dev/null +++ b/bridged/src/test/java/dev/ltms/bridged/session/SessionManagerTest.java @@ -0,0 +1,150 @@ +package dev.ltms.bridged.session; + +import dev.ltms.bridged.config.BridgedConfig; +import dev.ltms.bridged.guard.SubscriptionGuard; +import dev.ltms.bridged.herdr.AgentControl; +import dev.ltms.bridged.herdr.FakeHerdr; +import dev.ltms.bridged.herdr.WorkspaceControl; +import dev.ltms.bridged.worker.WorkerService; +import org.junit.jupiter.api.Test; + +import java.util.List; +import java.util.Map; +import java.util.Set; + +import static org.junit.jupiter.api.Assertions.*; + +/** + * CB-301 acceptance tests for the authoritative session registry and one-shot lifecycle FSM. + * No live herdr — everything runs against the same {@link FakeHerdr} the rest of the project uses. + */ +class SessionManagerTest { + + private SessionManager sessionManager(FakeHerdr herdr) { + BridgedConfig.Worker cfg = new BridgedConfig.Worker( + "ltms-local", "http://gx00.gw:8000", "coder", null, "BRIDGED_WORKER_TOKEN", + List.of("ccs", "ltms-local"), "tab", "bridged-workers", + "worker: {profile} #{n}", null, null); + WorkerService workers = new WorkerService(new AgentControl(herdr), new WorkspaceControl(herdr), + new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null); + return new SessionManager(workers); + } + + @Test + void acquireRegistersSpawningSessionWithDistinctPaneId() { + FakeHerdr herdr = new FakeHerdr(); + SessionManager sessions = sessionManager(herdr); + + WorkerSession a = sessions.acquire("ltms-local", "/work/a", "/caller/a", "term_primary"); + WorkerSession b = sessions.acquire("ltms-local", "/work/b", "/caller/b", "term_primary"); + + assertEquals(WorkerSession.State.SPAWNING, a.state(), "fresh session starts spawning"); + assertEquals("ltms-local", a.profile()); + assertEquals("/work/a", a.cwd(), "explicit requested cwd is recorded"); + assertEquals("term_primary", a.ownerTerminal()); + assertTrue(a.spawnedAtNanos() > 0); + assertNotNull(a.paneId()); + assertNotNull(a.terminalId()); + + assertNotEquals(a.paneId(), b.paneId(), "no pane reuse"); + assertNotEquals(a.terminalId(), b.terminalId(), "no terminal reuse"); + assertEquals(2, sessions.roster().size(), "both sessions are registered"); + } + + @Test + void presenceMovesSpawningToReadyAndDeliveredTurnMovesToDone() { + FakeHerdr herdr = new FakeHerdr(); + SessionManager sessions = sessionManager(herdr); + WorkerSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary"); + String terminal = session.terminalId(); + + sessions.asPresence().markPresent(terminal); + assertEquals(WorkerSession.State.READY, sessions.get(session.paneId()).orElseThrow().state(), + "MCP presence moves SPAWNING → READY"); + assertTrue(sessions.asPresence().isPresent(terminal), "presence is also recorded"); + + sessions.onDelivered(terminal); + assertEquals(WorkerSession.State.BUSY, sessions.get(session.paneId()).orElseThrow().state(), + "delivery moves READY → BUSY"); + + sessions.onTurnComplete(terminal); + assertEquals(WorkerSession.State.DONE, sessions.get(session.paneId()).orElseThrow().state(), + "turn completion moves BUSY → DONE"); + } + + @Test + void releaseTearsDownWorkerAndRemovesFromRosterAndIsIdempotent() { + FakeHerdr herdr = new FakeHerdr(); + SessionManager sessions = sessionManager(herdr); + WorkerSession session = sessions.acquire("ltms-local", null, "/caller", null); + String paneId = session.paneId(); + + sessions.release(paneId); + + assertTrue(herdr.called("pane.close"), "release tears the worker pane down"); + assertTrue(sessions.get(paneId).isEmpty(), "released session is no longer retrievable"); + assertTrue(sessions.roster().isEmpty(), "released session is no longer in the roster"); + + assertDoesNotThrow(() -> sessions.release(paneId), "a second release is harmless"); + } + + @Test + void onTurnFailedMovesSessionToFailed() { + FakeHerdr herdr = new FakeHerdr(); + SessionManager sessions = sessionManager(herdr); + WorkerSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary"); + String terminal = session.terminalId(); + sessions.asPresence().markPresent(terminal); + sessions.onDelivered(terminal); + + sessions.onTurnFailed(terminal); + + WorkerSession updated = sessions.get(session.paneId()).orElseThrow(); + assertEquals(WorkerSession.State.FAILED, updated.state(), "turn failure moves to FAILED"); + assertTrue(sessions.roster().contains(updated), "FAILED is still in acquired-minus-released roster"); + } + + @Test + void recycleProducesNewPaneIdAndOldOneIsGone() { + FakeHerdr herdr = new FakeHerdr(); + SessionManager sessions = sessionManager(herdr); + WorkerSession oldSession = sessions.acquire("ltms-local", null, "/caller", "term_primary"); + String oldPane = oldSession.paneId(); + String oldTerminal = oldSession.terminalId(); + + WorkerSession fresh = sessions.recycle(oldPane); + + assertNotEquals(oldPane, fresh.paneId(), "recycle yields a new pane id"); + assertNotEquals(oldTerminal, fresh.terminalId(), "recycle yields a new terminal id"); + assertEquals(oldSession.profile(), fresh.profile(), "profile is preserved"); + assertEquals(oldSession.cwd(), fresh.cwd(), "cwd is preserved"); + assertEquals(oldSession.ownerTerminal(), fresh.ownerTerminal(), "owner is preserved"); + + assertTrue(sessions.get(oldPane).isEmpty(), "old pane is deregistered"); + assertEquals(1, sessions.roster().size(), "only the fresh session remains"); + assertEquals(fresh.paneId(), sessions.roster().getFirst().paneId()); + + long paneCloseCount = herdr.calls.stream() + .filter(c -> "pane.close".equals(c.method())) + .filter(c -> oldPane.equals(((Map) c.params()).get("pane_id"))) + .count(); + assertEquals(1, paneCloseCount, "the old worker was torn down"); + } + + @Test + void rosterReflectsAcquiredMinusReleased() { + FakeHerdr herdr = new FakeHerdr(); + SessionManager sessions = sessionManager(herdr); + WorkerSession a = sessions.acquire("ltms-local", "/a", "/caller", "ownerA"); + WorkerSession b = sessions.acquire("ltms-local", "/b", "/caller", "ownerB"); + + assertEquals(2, sessions.roster().size()); + assertTrue(sessions.roster().stream().anyMatch(s -> s.paneId().equals(a.paneId()))); + assertTrue(sessions.roster().stream().anyMatch(s -> s.paneId().equals(b.paneId()))); + + sessions.release(a.paneId()); + + assertEquals(1, sessions.roster().size()); + assertEquals(b.paneId(), sessions.roster().getFirst().paneId()); + } +}