CB-301: SessionManager — authoritative worker session registry + one-shot FSM

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.
This commit is contained in:
Dai Ha
2026-07-16 19:30:56 +02:00
parent 82c7d6553a
commit 54d907c314
11 changed files with 501 additions and 36 deletions
@@ -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);
@@ -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<String> present = ConcurrentHashMap.newKeySet();
@@ -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<String, Object> 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<String, Object> workerView(WorkerSession s) {
Map<String, Object> 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);
@@ -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<String, Object> view(WorkerSession s) {
Map<String, Object> 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;
}
}
@@ -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.
*
* <p>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.
*
* <p>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<String /*paneId*/, WorkerSession> 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<WorkerSession> get(String paneId) {
return Optional.ofNullable(registry.get(paneId));
}
/** Bridge-owned roster: all registered sessions (acquired minus released). */
public List<WorkerSession> 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);
}
}
}
@@ -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);
}
}
@@ -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;