CB-557: adopt the member taxonomy in the API surface

Every spawned peer is now a member with a role, and the role travels with it
from the spawn call to the roster.

MCP:
  bridge_spawn gains role: architect | dev | reviewer (default dev). An
    unknown role is refused with the valid spellings in the message.
  bridge_list returns "members" instead of "workers"; each row carries both
    role (what it is for) and profile (which backend it runs on).
  The spawn result echoes the role back, so a spawn that fell back to dev is
    visible rather than silent.

REST:
  GET/POST /members and DELETE /members/{paneId} replace /workers.
  POST accepts role= as a query param or a body field; an unknown role is 400.

Code:
  dev.ltms.bridged.worker package -> dev.ltms.bridged.member
  WorkerSession   -> MemberSession, plus a MemberRole role component
  WorkerPresence  -> MemberPresence
  SessionManager.acquire gains a role parameter; the existing overloads keep
    working and default to DEV, which is exactly what "worker" used to mean.

ClaudeCodeLauncher and OpenCodeLauncher keep their names on purpose — they
are named after the backend, not the role.

Not done here: the launch charter is still one string for every role, so a
member is told its role by nobody yet. That is the next ticket.

mvn clean install: 583 tests, 0 failures, 0 errors, BUILD SUCCESS.
This commit is contained in:
Dai Ha
2026-08-14 07:08:26 +02:00
parent d14a624421
commit 4875127daa
29 changed files with 365 additions and 230 deletions
@@ -13,7 +13,7 @@ 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.inject.MemberPresence;
import dev.ltms.bridged.auth.MemberRegistry;
import dev.ltms.bridged.auth.CallerResolver;
import dev.ltms.bridged.mcp.BridgeMcp;
@@ -35,11 +35,11 @@ import dev.ltms.bridged.session.GitWorktrees;
import dev.ltms.bridged.session.SessionManager;
import dev.ltms.bridged.peer.PeerLauncher;
import dev.ltms.bridged.session.SessionReaper;
import dev.ltms.bridged.worker.ClaudeCodeLauncher;
import dev.ltms.bridged.member.ClaudeCodeLauncher;
import dev.ltms.bridged.placement.PlacementPolicies;
import dev.ltms.bridged.worker.CompositePeerLauncher;
import dev.ltms.bridged.worker.HerdrPeerLauncher;
import dev.ltms.bridged.worker.OpenCodeLauncher;
import dev.ltms.bridged.member.CompositePeerLauncher;
import dev.ltms.bridged.member.HerdrPeerLauncher;
import dev.ltms.bridged.member.OpenCodeLauncher;
import io.javalin.Javalin;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -224,7 +224,7 @@ public final class Bridged {
CompletionResolver completion = new CompletionResolver(agents, rendezvous);
// CB-113: deliver only to an available worker (its MCP is connected), never its boot window.
// CB-301: the manager's presence bridge records availability and drives SPAWNING → READY.
WorkerPresence presence = sessions.asPresence();
MemberPresence presence = sessions.asPresence();
TurnListener turnListener = new TurnListener() {
@Override
public void onTurnComplete(String target) {
@@ -394,7 +394,7 @@ public final class Bridged {
* hazard is a property of spawning. A lead is never spawned: the operator started it and named it
* (or labelled its tab) only once it was up, so there is no boot window to guard.
*
* <p>A lead is also never enrolled in {@link WorkerPresence} — {@code BridgeMcp} marks presence
* <p>A lead is also never enrolled in {@link MemberPresence} — {@code BridgeMcp} marks presence
* only for a worker, deliberately, since that map doubles as the worker roster's availability
* signal and a lead counted there would show up as an available worker. So without the second
* disjunct a lead is permanently un-deliverable: every lead→lead send sat on the gate for
@@ -403,7 +403,7 @@ public final class Bridged {
* <p>The lead set is read through the supplier on each call rather than snapshotted, so a lead
* discovered by {@code leadScan} after startup becomes deliverable without a restart.
*/
static Predicate<String> deliverableTo(WorkerPresence presence, Supplier<Map<String, String>> leads) {
static Predicate<String> deliverableTo(MemberPresence presence, Supplier<Map<String, String>> leads) {
return target -> presence.isPresent(target) || leads.get().containsKey(target);
}
@@ -156,7 +156,7 @@ public record BridgedConfig(
* non-positive ⇒ unlimited. Live means any session the registry still owns
* (acquired and not yet released), in any state.
* @param kind which peer launcher spawns this profile: {@code "claude-code"} (default —
* the {@link dev.ltms.bridged.worker.ClaudeCodeLauncher}) or {@code "opencode"}.
* the {@link dev.ltms.bridged.member.ClaudeCodeLauncher}) or {@code "opencode"}.
* The {@code CompositePeerLauncher} routes {@code spawn}/reap by this value, so
* each adapter drives only its own kind. Normalised to lower-case; blank ⇒ the
* default. It selects the adapter, not the transport — placement, tabs, cwd, and
@@ -186,7 +186,7 @@ public record BridgedConfig(
Integer maxLoad,
Boolean subscription) {
/** Peer kind spawned by {@link dev.ltms.bridged.worker.ClaudeCodeLauncher} (the default). */
/** Peer kind spawned by {@link dev.ltms.bridged.member.ClaudeCodeLauncher} (the default). */
public static final String KIND_CLAUDE_CODE = "claude-code";
/** Peer kind spawned by the opencode adapter (CB-402). */
public static final String KIND_OPENCODE = "opencode";
@@ -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 class WorkerPresence {
public class MemberPresence {
private final Set<String> present = ConcurrentHashMap.newKeySet();
@@ -9,14 +9,15 @@ import dev.ltms.bridged.guard.GuardException;
import dev.ltms.bridged.herdr.Agent;
import dev.ltms.bridged.metrics.BridgedMetrics;
import dev.ltms.bridged.metrics.Metrics;
import dev.ltms.bridged.inject.WorkerPresence;
import dev.ltms.bridged.inject.MemberPresence;
import dev.ltms.bridged.herdr.HerdrException;
import dev.ltms.bridged.msg.MessageService;
import dev.ltms.bridged.msg.Rendezvous;
import dev.ltms.bridged.peer.PeerUnreachableException;
import dev.ltms.bridged.session.SessionManager;
import dev.ltms.bridged.session.WorkerSession;
import dev.ltms.bridged.session.MemberSession;
import dev.ltms.bridged.session.WorktreeRequest;
import dev.ltms.bridged.peer.MemberRole;
import dev.ltms.bridged.peer.PeerLauncher;
import io.modelcontextprotocol.common.McpTransportContext;
import io.modelcontextprotocol.json.McpJsonMapper;
@@ -78,7 +79,7 @@ public final class BridgeMcp {
* without an auth fixture.
*/
public BridgeMcp(MessageService messages, PeerLauncher workers,
SessionManager sessions, ConnectionIdentity identity, WorkerPresence presence,
SessionManager sessions, ConnectionIdentity identity, MemberPresence presence,
PrimaryRegistry primaryRegistry) {
this(messages, workers, sessions, identity, presence, primaryRegistry, null, null);
}
@@ -91,7 +92,7 @@ public final class BridgeMcp {
* @param metrics registry for auth-failure counting; may be {@code null}
*/
public BridgeMcp(MessageService messages, PeerLauncher workers,
SessionManager sessions, ConnectionIdentity identity, WorkerPresence presence,
SessionManager sessions, ConnectionIdentity identity, MemberPresence presence,
PrimaryRegistry primaryRegistry, CallerResolver callers, Metrics metrics) {
McpJsonMapper json = new JacksonMcpJsonMapperSupplier().get();
this.transport = HttpServletStreamableServerTransportProvider.builder()
@@ -200,7 +201,7 @@ public final class BridgeMcp {
// CB-301: carry the caller's identity as the session owner (null for the primary).
// CB-301-ext: optional isolated worktree for parallel implementers.
String callerCwd = identity.cwdForPid(callerPid(exchange));
return spawn(sessions, str(a, "profile"), str(a, "cwd"), callerCwd,
return spawn(sessions, str(a, "profile"), str(a, "role"), str(a, "cwd"), callerCwd,
callerTerminal(exchange), worktreeRequest(a));
})
.toolCall(listTool(), (exchange, _) -> {
@@ -606,23 +607,33 @@ public final class BridgeMcp {
/** {@code bridge_spawn} without cwd/caller context (default resolution). */
static McpSchema.CallToolResult spawn(SessionManager sessions, String profile) {
return spawn(sessions, profile, null, null, null, null);
return spawn(sessions, profile, null, null, 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.
* {@code bridge_spawn}: launch a guard-checked member for {@code profile} (blank → the default
* profile) under {@code role} (blank → {@code dev}), and return its session id + pane id. The
* member'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.
* CB-301-ext: {@code worktreeRequest} non-null provisions an isolated git worktree.
*
* <p>{@code role} and {@code profile} are independent: the role picks the contract, the profile
* picks the backend. A reviewer on the same profile as the dev it reviews is a normal spawn.
*/
static McpSchema.CallToolResult spawn(SessionManager sessions, String profile,
static McpSchema.CallToolResult spawn(SessionManager sessions, String profile, String role,
String requestedCwd, String callerCwd,
String ownerTerminal, WorktreeRequest worktreeRequest) {
MemberRole memberRole;
try {
WorkerSession worker = sessions.acquire(isBlank(profile) ? null : profile,
memberRole = isBlank(role) ? MemberRole.DEV : MemberRole.parse(role);
} catch (IllegalArgumentException e) {
return error(e.getMessage());
}
try {
MemberSession member = sessions.acquire(isBlank(profile) ? null : profile, memberRole,
requestedCwd, callerCwd, ownerTerminal, worktreeRequest);
return text(json(workerView(worker)));
return text(json(memberView(member)));
} catch (GuardException e) {
return error("subscription boundary: " + e.getMessage());
} catch (IllegalArgumentException e) {
@@ -702,7 +713,7 @@ public final class BridgeMcp {
List<Map<String, Object>> out = sessions.roster().stream()
.map(s -> SessionManager.rosterView(s, live.get(s.terminalId())))
.toList();
return text(json(Map.of("leads", leadRows, "workers", out)));
return text(json(Map.of("leads", leadRows, "members", out)));
} catch (HerdrException e) {
return error("herdr error listing the fleet: " + e.getMessage());
}
@@ -743,10 +754,14 @@ public final class BridgeMcp {
}
/** CB-301 projection from the authoritative session registry. */
private static Map<String, Object> workerView(WorkerSession s) {
private static Map<String, Object> memberView(MemberSession s) {
Map<String, Object> m = new LinkedHashMap<>();
m.put("sessionId", s.terminalId());
m.put("paneId", s.paneId());
// Echo the role back so the caller can see what it actually got, not what it meant to ask
// for — a spawn that silently fell back to dev is otherwise invisible.
m.put("role", s.role() == null ? "dev" : s.role().wireName());
m.put("profile", s.profile());
m.put("status", s.state().name().toLowerCase());
if (s.worktree() != null) {
m.put("worktree", s.worktree());
@@ -823,14 +838,20 @@ public final class BridgeMcp {
private static McpSchema.Tool spawnTool() {
return tool("bridge_spawn",
"Spawn a new off-subscription worker session. Pass a profile (from bridge_profiles) to "
+ "pick the backend, or omit it for the default. The worker opens your current "
+ "directory by default; pass cwd to pin a different one. Pass worktree:true (with "
+ "ticket) or worktree:<ticket-slug> to provision an isolated git worktree. "
+ "Returns the worker's sessionId (use with bridge_send) and paneId (use with bridge_stop).",
"Spawn a new off-subscription member session. A member has two independent attributes: "
+ "role (what it is for) and profile (which backend it runs on). Pass role to pick "
+ "the contract — 'dev' implements a unit and opens its own PR, 'reviewer' reviews a "
+ "diff it did not write, 'architect' refines a ticket before anyone builds it; omit "
+ "it for 'dev'. Pass profile (from bridge_profiles) to pick the backend, or omit it "
+ "for the default. The two are independent: a reviewer may run on the same profile "
+ "as the dev it reviews. The member opens your current directory by default; pass "
+ "cwd to pin a different one. Pass worktree:true (with ticket) or "
+ "worktree:<ticket-slug> to provision an isolated git worktree. Returns the member's "
+ "sessionId (use with bridge_send) and paneId (use with bridge_stop).",
objectSchema(Map.of(
"profile", stringProp("Worker profile to spawn (omit for the default profile)"),
"cwd", stringProp("Working directory for the worker (omit to inherit yours)"),
"role", stringProp("What the member is for: architect, dev or reviewer (default dev)"),
"profile", stringProp("Which backend to run it on (omit for the default profile)"),
"cwd", stringProp("Working directory for the member (omit to inherit yours)"),
"worktree", Map.of("type", "string", "description", "'true' or a ticket slug — requests an isolated git worktree"),
"ticket", stringProp("Ticket slug when worktree:true")),
List.of()));
@@ -847,10 +868,11 @@ public final class BridgeMcp {
"List the whole fleet the bridge tracks, in two parts. 'leads' are your PEERS — other "
+ "orchestrators, each with its sessionId (the address to bridge_send to), "
+ "name, live status, and 'self': true on your own row; this is how you "
+ "discover a peer lead without being told its address. 'workers' are the "
+ "sessions delegated to — each with sessionId, paneId, profile, state, "
+ "optional worktree/branch/owner, and live herdr status. An empty 'workers' "
+ "means no workers are spawned; it says nothing about peers.",
+ "discover a peer lead without being told its address. 'members' are the "
+ "sessions delegated to — each with sessionId, paneId, role (architect/dev/"
+ "reviewer), profile (the backend it runs on), state, optional "
+ "worktree/branch/owner, and live herdr status. An empty 'members' "
+ "means no members are spawned; it says nothing about peers.",
objectSchema(Map.of(), List.of()));
}
@@ -1,4 +1,4 @@
package dev.ltms.bridged.worker;
package dev.ltms.bridged.member;
import dev.ltms.bridged.config.BridgedConfig;
import dev.ltms.bridged.guard.SubscriptionGuard;
@@ -1,4 +1,4 @@
package dev.ltms.bridged.worker;
package dev.ltms.bridged.member;
import dev.ltms.bridged.config.BridgedConfig;
import dev.ltms.bridged.herdr.Agent;
@@ -1,4 +1,4 @@
package dev.ltms.bridged.worker;
package dev.ltms.bridged.member;
import dev.ltms.bridged.config.BridgedConfig;
import dev.ltms.bridged.herdr.Agent;
@@ -1,4 +1,4 @@
package dev.ltms.bridged.worker;
package dev.ltms.bridged.member;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.node.ObjectNode;
@@ -1,4 +1,4 @@
package dev.ltms.bridged.worker;
package dev.ltms.bridged.member;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
@@ -2,7 +2,7 @@ package dev.ltms.bridged.metrics;
import dev.ltms.bridged.msg.ReplyInbox;
import dev.ltms.bridged.session.SessionManager;
import dev.ltms.bridged.session.WorkerSession;
import dev.ltms.bridged.session.MemberSession;
import java.util.LinkedHashMap;
import java.util.Map;
@@ -74,7 +74,7 @@ public final class BridgedMetrics {
// One gauge per state so a scrape shows the whole census even when a state is empty —
// an absent series and a zero series read very differently on a dashboard.
for (WorkerSession.State state : WorkerSession.State.values()) {
for (MemberSession.State state : MemberSession.State.values()) {
String label = state.name().toLowerCase();
m.gauge(SESSIONS, () -> countIn(sessions, state), "state", label);
}
@@ -83,7 +83,7 @@ public final class BridgedMetrics {
// port's non-destructive read — scraping metrics must never ack a reply out of the inbox.
m.collector(INBOX_DEPTH, "target", () -> {
Map<String, Number> depths = new LinkedHashMap<>();
for (WorkerSession s : sessions.roster()) {
for (MemberSession s : sessions.roster()) {
String target = s.terminalId();
if (target == null) {
continue;
@@ -95,7 +95,7 @@ public final class BridgedMetrics {
return m;
}
private static long countIn(SessionManager sessions, WorkerSession.State state) {
private static long countIn(SessionManager sessions, MemberSession.State state) {
return sessions.roster().stream().filter(s -> s.state() == state).count();
}
}
@@ -5,7 +5,7 @@ import dev.ltms.bridged.herdr.AgentStatus;
import dev.ltms.bridged.mcp.PrimaryRegistry;
import dev.ltms.bridged.metrics.BridgedMetrics;
import dev.ltms.bridged.metrics.Metrics;
import dev.ltms.bridged.session.WorkerSession;
import dev.ltms.bridged.session.MemberSession;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -55,7 +55,7 @@ public final class LeadHeartbeatLoop {
private final PrimaryRegistry primaryRegistry;
private final AgentControl agents;
private final ReplyInbox inbox;
private final Supplier<List<WorkerSession>> roster;
private final Supplier<List<MemberSession>> roster;
private final ReplyPushLoop pushLoop;
private final ScheduledExecutorService scheduler;
private final LongSupplier clock;
@@ -71,7 +71,7 @@ public final class LeadHeartbeatLoop {
/** Constructor with an injectable clock and no metric registry (unit tests, or wiring that opts out). */
public LeadHeartbeatLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
Supplier<List<WorkerSession>> roster, ReplyPushLoop pushLoop,
Supplier<List<MemberSession>> roster, ReplyPushLoop pushLoop,
ScheduledExecutorService scheduler, LongSupplier clock,
long idleAfterNanos, long backoffMs, int quietNudgeCap) {
this(primaryRegistry, agents, inbox, roster, pushLoop, scheduler, clock,
@@ -80,7 +80,7 @@ public final class LeadHeartbeatLoop {
/** As above, with a metric registry (the CB-512 pattern) so nudge outcomes are counted. */
public LeadHeartbeatLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
Supplier<List<WorkerSession>> roster, ReplyPushLoop pushLoop,
Supplier<List<MemberSession>> roster, ReplyPushLoop pushLoop,
ScheduledExecutorService scheduler, LongSupplier clock,
long idleAfterNanos, long backoffMs, int quietNudgeCap, Metrics metrics) {
this.primaryRegistry = primaryRegistry;
@@ -328,12 +328,12 @@ public final class LeadHeartbeatLoop {
* sources the metrics collector uses, so the heartbeat reports facts the operator can cross-check
* against {@code /metrics}. Package-private and pure (reads only) so it is unit-testable.
*/
static FleetState snapshot(ReplyInbox inbox, Supplier<List<WorkerSession>> roster) {
List<WorkerSession> sessions = roster.get();
static FleetState snapshot(ReplyInbox inbox, Supplier<List<MemberSession>> roster) {
List<MemberSession> sessions = roster.get();
int replies = 0;
int done = 0;
List<String> replyTargets = new ArrayList<>();
for (WorkerSession s : sessions) {
for (MemberSession s : sessions) {
if (s.terminalId() == null) {
continue;
}
@@ -342,7 +342,7 @@ public final class LeadHeartbeatLoop {
replies += depth;
replyTargets.add(s.terminalId());
}
if (s.state() == WorkerSession.State.DONE) {
if (s.state() == MemberSession.State.DONE) {
done++;
}
}
@@ -11,11 +11,12 @@ import dev.ltms.bridged.metrics.Metrics;
import dev.ltms.bridged.herdr.Agent;
import dev.ltms.bridged.herdr.HerdrClient;
import dev.ltms.bridged.herdr.HerdrException;
import dev.ltms.bridged.inject.WorkerPresence;
import dev.ltms.bridged.inject.MemberPresence;
import dev.ltms.bridged.peer.PeerUnreachableException;
import dev.ltms.bridged.msg.MessageService;
import dev.ltms.bridged.session.SessionManager;
import dev.ltms.bridged.session.WorkerSession;
import dev.ltms.bridged.peer.MemberRole;
import dev.ltms.bridged.session.MemberSession;
import dev.ltms.bridged.session.WorktreeRequest;
import dev.ltms.bridged.peer.PeerLauncher;
import io.javalin.Javalin;
@@ -55,7 +56,7 @@ public final class BridgedApp {
private final PeerLauncher workers;
private final SessionManager sessions; // CB-301: authoritative session registry
private final MessageService messages;
private final WorkerPresence presence; // CB-113: which workers are MCP-connected (available)
private final MemberPresence presence; // CB-113: which workers are MCP-connected (available)
private final HttpServlet mcpServlet; // MCP Streamable-HTTP endpoint, mounted at /mcp (nullable)
private final CallerResolver auth; // CB-501: null → authz not enforced (legacy behaviour)
private final Metrics metrics; // CB-502: null → /metrics not exposed
@@ -67,7 +68,7 @@ public final class BridgedApp {
* behaviour without each needing an auth fixture.
*/
public BridgedApp(HerdrClient herdr, PeerLauncher workers, SessionManager sessions,
MessageService messages, WorkerPresence presence,
MessageService messages, MemberPresence presence,
HttpServlet mcpServlet) {
this(herdr, workers, sessions, messages, presence, mcpServlet, null, null);
}
@@ -79,7 +80,7 @@ public final class BridgedApp {
* the endpoint
*/
public BridgedApp(HerdrClient herdr, PeerLauncher workers, SessionManager sessions,
MessageService messages, WorkerPresence presence,
MessageService messages, MemberPresence presence,
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics) {
this.herdr = herdr;
this.workers = workers;
@@ -115,10 +116,10 @@ public final class BridgedApp {
}
app.get("/sessions", this::sessions);
app.get("/agents", this::agents);
app.get("/workers", this::listWorkers); // CB-304: registry roster + live herdr status
app.get("/profiles", this::profiles); // configured worker profiles
app.post("/workers", this::spawnWorker); // optional ?profile= or {"profile":…}
app.delete("/workers/{paneId}", this::stopWorker);
app.get("/members", this::listMembers); // CB-304: registry roster + live herdr status
app.get("/profiles", this::profiles); // configured backend profiles
app.post("/members", this::spawnMember); // optional ?role=&profile= or {"role":…,"profile":…}
app.delete("/members/{paneId}", this::stopMember);
app.post("/sessions/{id}/message", this::sendMessage); // bridge_send (primary; blocking, wait:false, or answer via turnId)
app.post("/sessions/{id}/reply", this::replyMessage); // bridge_reply (worker)
app.get("/sessions/{id}/replies", this::drainReplies); // drain reply inbox (CB-307)
@@ -221,7 +222,7 @@ public final class BridgedApp {
}
/** CB-304: bridge-owned roster merged with live herdr status by paneId. */
private void listWorkers(Context ctx) {
private void listMembers(Context ctx) {
if (!allow(ctx, Authz.Action.READ, null)) {
return;
}
@@ -251,10 +252,11 @@ public final class BridgedApp {
* body) picks which configured profile; omitted → the default. 403 if the base_url would breach
* the subscription boundary, 400 for an unknown profile.
*/
private void spawnWorker(Context ctx) {
private void spawnMember(Context ctx) {
if (!allow(ctx, Authz.Action.SPAWN, null)) {
return;
}
String role = ctx.queryParam("role");
String profile = ctx.queryParam("profile");
String cwd = ctx.queryParam("cwd");
String worktree = ctx.queryParam("worktree");
@@ -265,6 +267,7 @@ public final class BridgedApp {
String body = ctx.body();
if (!body.isBlank()) {
JsonNode b = mapper.readTree(body);
if (role == null || role.isBlank()) role = b.path("role").asText(null);
if (profile == null || profile.isBlank()) profile = b.path("profile").asText(null);
if (cwd == null || cwd.isBlank()) cwd = b.path("cwd").asText(null);
if (worktree == null || worktree.isBlank()) worktree = b.path("worktree").asText(null);
@@ -275,10 +278,18 @@ public final class BridgedApp {
}
}
WorktreeRequest wt = worktreeRequest(worktree, ticket);
MemberRole memberRole;
try {
memberRole = (role == null || role.isBlank()) ? MemberRole.DEV : MemberRole.parse(role);
} catch (IllegalArgumentException e) {
ctx.status(400).json(Map.of("error", "unknown_role", "detail", e.getMessage()));
return;
}
try {
// No MCP caller over REST, so callerCwd and ownerTerminal are null.
WorkerSession worker = sessions.acquire(blankToNull(profile), blankToNull(cwd), null, null, wt);
ctx.status(201).json(view(worker));
MemberSession member = sessions.acquire(blankToNull(profile), memberRole,
blankToNull(cwd), null, null, wt);
ctx.status(201).json(view(member));
} catch (GuardException e) {
ctx.status(403).json(Map.of("error", "subscription_boundary", "detail", e.getMessage()));
} catch (IllegalArgumentException e) {
@@ -306,7 +317,7 @@ public final class BridgedApp {
}
/** Tear a worker down by pane id. */
private void stopWorker(Context ctx) {
private void stopMember(Context ctx) {
String paneId = ctx.pathParam("paneId");
if (!allow(ctx, Authz.Action.STOP, paneId)) {
return;
@@ -544,7 +555,7 @@ public final class BridgedApp {
}
/** CB-301 projection of an authoritative bridge-owned session. */
private static Map<String, Object> view(WorkerSession s) {
private static Map<String, Object> view(MemberSession s) {
Map<String, Object> m = new LinkedHashMap<>();
m.put("terminalId", s.terminalId());
m.put("paneId", s.paneId());
@@ -1,5 +1,7 @@
package dev.ltms.bridged.session;
import dev.ltms.bridged.peer.MemberRole;
/**
* 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
@@ -10,7 +12,11 @@ package dev.ltms.bridged.session;
* dev.ltms.bridged.peer.PeerHandle#id()}, a UUID, and is distinct from the
* launcher-private herdr pane coordinate.
* @param terminalId herdr terminal handle — the {@code target} for send/read/status
* @param profile the worker profile name that spawned this session
* @param profile the profile name that spawned this session — WHICH BACKEND
* @param role the contract this member runs under — WHAT IT IS FOR. Independent of
* {@code profile}: a reviewer may run on the same profile as the dev
* whose diff it reads. {@code null} only for a session recorded before
* the role was known
* @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
@@ -18,10 +24,11 @@ package dev.ltms.bridged.session;
* @param turnCount number of delegated turns that have been delivered to this session
* @param state current lifecycle state in the one-shot FSM
*/
public record WorkerSession(
public record MemberSession(
String paneId,
String terminalId,
String profile,
MemberRole role,
String cwd,
String ownerTerminal,
long spawnedAtNanos,
@@ -42,20 +49,20 @@ public record WorkerSession(
}
/** Return a copy of this session in {@code state}. */
public WorkerSession withState(State state) {
return new WorkerSession(paneId, terminalId, profile, cwd, ownerTerminal, spawnedAtNanos,
public MemberSession withState(State state) {
return new MemberSession(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
lastActivityAtNanos, turnCount, state, worktree, branch);
}
/** Return a copy with {@code lastActivityAtNanos} updated to {@code nowNanos}. */
public WorkerSession withActivity(long nowNanos) {
return new WorkerSession(paneId, terminalId, profile, cwd, ownerTerminal, spawnedAtNanos,
public MemberSession withActivity(long nowNanos) {
return new MemberSession(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
nowNanos, turnCount, state, worktree, branch);
}
/** Return a copy with the turn count incremented and activity timestamped at {@code nowNanos}. */
public WorkerSession bumpTurn(long nowNanos) {
return new WorkerSession(paneId, terminalId, profile, cwd, ownerTerminal, spawnedAtNanos,
public MemberSession bumpTurn(long nowNanos) {
return new MemberSession(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
nowNanos, turnCount + 1, state, worktree, branch);
}
}
@@ -2,7 +2,8 @@ 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.inject.MemberPresence;
import dev.ltms.bridged.peer.MemberRole;
import dev.ltms.bridged.peer.PeerHandle;
import dev.ltms.bridged.peer.PeerLauncher;
import dev.ltms.bridged.peer.SpawnRequest;
@@ -31,7 +32,7 @@ import java.util.function.LongSupplier;
* 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
* {@code READY → BUSY → DONE} (or {@code FAILED}). It exposes a {@link MemberPresence} view via
* {@link #asPresence()}: any MCP contact from a worker marks it present and simultaneously
* transitions the session {@code SPAWNING → READY}.
*/
@@ -41,8 +42,8 @@ public final class SessionManager implements TurnListener {
private final PeerLauncher launcher;
private final Worktrees worktrees;
private final ConcurrentHashMap<String /*paneId*/, WorkerSession> registry = new ConcurrentHashMap<>();
private final WorkerPresence presence;
private final ConcurrentHashMap<String /*paneId*/, MemberSession> registry = new ConcurrentHashMap<>();
private final MemberPresence presence;
private final SecureRandom nonceRandom = new SecureRandom();
private final AtomicLong nonceSeq = new AtomicLong();
private final LongSupplier nowNanos;
@@ -90,48 +91,66 @@ public final class SessionManager implements TurnListener {
}
/**
* The single {@link WorkerPresence} view of this manager: it records availability and forwards
* The single {@link MemberPresence} 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
* {@code BridgeMcp} where they previously accepted a plain {@link MemberPresence}. 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() {
public MemberPresence asPresence() {
return presence;
}
/**
* Spawn a worker and register it as {@link WorkerSession.State#SPAWNING}. The caller's
* Spawn a worker and register it as {@link MemberSession.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,
public MemberSession acquire(String profile, String requestedCwd, String callerCwd,
String ownerTerminal) {
return acquire(profile, requestedCwd, callerCwd, ownerTerminal, null);
}
/**
* Spawn a worker, optionally inside a fresh git worktree. When {@code wt} is non-null the
* worktree is provisioned, parity-overlaid, and its path becomes the worker's cwd. On any
* failure before registration the worktree is removed so no dangling checkout is left.
* Spawn a member, optionally inside a fresh git worktree, defaulting the role to
* {@link MemberRole#DEV}.
*
* <p>{@code DEV} is the right default because it is exactly what the old "worker" meant: an
* unqualified spawn is a unit of implementation work. An architect or a reviewer is always
* asked for on purpose, so neither is ever what a caller silently gets.
*/
public WorkerSession acquire(String profile, String requestedCwd, String callerCwd,
public MemberSession acquire(String profile, String requestedCwd, String callerCwd,
String ownerTerminal, WorktreeRequest wt) {
return acquire(profile, MemberRole.DEV, requestedCwd, callerCwd, ownerTerminal, wt);
}
/**
* Spawn a member, optionally inside a fresh git worktree. When {@code wt} is non-null the
* worktree is provisioned, parity-overlaid, and its path becomes the member's cwd. On any
* failure before registration the worktree is removed so no dangling checkout is left.
*
* @param profile which backend to run on — a {@code profiles:} key
* @param role which contract the member runs under; never {@code null}
*/
public MemberSession acquire(String profile, MemberRole role, String requestedCwd, String callerCwd,
String ownerTerminal, WorktreeRequest wt) {
MemberRole memberRole = (role == null) ? MemberRole.DEV : role;
if (wt == null) {
SpawnRequest req = new SpawnRequest(profile, requestedCwd, callerCwd);
PeerHandle handle = launcher.spawn(req);
String resolvedProfile = resolveProfile(handle, profile);
String cwd = launcher.effectiveCwd(new SpawnRequest(resolvedProfile, requestedCwd, callerCwd));
long now = nowNanos.getAsLong();
WorkerSession session = new WorkerSession(
MemberSession session = new MemberSession(
handle.id(),
handle.terminalId(),
resolvedProfile,
memberRole,
cwd,
ownerTerminal,
now,
now,
0,
WorkerSession.State.SPAWNING,
MemberSession.State.SPAWNING,
null,
null);
registry.put(handle.id(), session);
@@ -140,7 +159,7 @@ public final class SessionManager implements TurnListener {
notifyAcquired(session.terminalId());
return session;
}
return acquireWithWorktree(profile, requestedCwd, callerCwd, ownerTerminal, wt);
return acquireWithWorktree(profile, memberRole, requestedCwd, callerCwd, ownerTerminal, wt);
}
/** Tear a worker down by pane id and remove it from the registry. Idempotent. */
@@ -162,7 +181,7 @@ public final class SessionManager implements TurnListener {
* is a logged path an operator can reclaim, the cost of a deleted one is unrecoverable work.
*/
private void release(String paneId, ReleaseCause cause) {
WorkerSession removed = registry.remove(paneId);
MemberSession removed = registry.remove(paneId);
boolean preserveWorktree = cause == ReleaseCause.SHUTDOWN;
if (removed != null) {
log.debug("releasing session pane={} terminal={} state={} cause={}",
@@ -198,8 +217,8 @@ public final class SessionManager implements TurnListener {
* when the drain timeout expired was abandoned mid-turn — that work may be uncommitted and is
* the only copy — so the message is loud and points at the path an operator needs to reclaim.
*/
private void logPreservedForShutdown(WorkerSession session) {
if (session.state() == WorkerSession.State.BUSY) {
private void logPreservedForShutdown(MemberSession session) {
if (session.state() == MemberSession.State.BUSY) {
log.warn("shutdown drain abandoned BUSY session pane={} terminal={} mid-turn; "
+ "worktree preserved at {}", session.paneId(), session.terminalId(),
session.worktree());
@@ -263,7 +282,7 @@ public final class SessionManager implements TurnListener {
}
}
private WorkerSession acquireWithWorktree(String profile, String requestedCwd, String callerCwd,
private MemberSession acquireWithWorktree(String profile, MemberRole memberRole, String requestedCwd, String callerCwd,
String ownerTerminal, WorktreeRequest wt) {
String preResolvedProfile = (profile == null || profile.isBlank())
? launcher.defaultProfile() : profile;
@@ -294,16 +313,17 @@ public final class SessionManager implements TurnListener {
String resolvedProfile = resolveProfile(handle, profile);
String cwd = launcher.effectiveCwd(new SpawnRequest(resolvedProfile, path, callerCwd));
long now = nowNanos.getAsLong();
WorkerSession session = new WorkerSession(
MemberSession session = new MemberSession(
handle.id(),
handle.terminalId(),
resolvedProfile,
memberRole,
cwd,
ownerTerminal,
now,
now,
0,
WorkerSession.State.SPAWNING,
MemberSession.State.SPAWNING,
path,
branch);
registry.put(handle.id(), session);
@@ -341,8 +361,8 @@ public final class SessionManager implements TurnListener {
* 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);
public MemberSession recycle(String paneId) {
MemberSession old = registry.get(paneId);
if (old == null) {
throw new IllegalArgumentException("no session for paneId " + paneId);
}
@@ -351,12 +371,12 @@ public final class SessionManager implements TurnListener {
}
/** The session for {@code paneId}, if it is still registered and not released. */
public Optional<WorkerSession> get(String paneId) {
public Optional<MemberSession> get(String paneId) {
return Optional.ofNullable(registry.get(paneId));
}
/** Bridge-owned roster: all registered sessions (acquired minus released). */
public List<WorkerSession> roster() {
public List<MemberSession> roster() {
return List.copyOf(registry.values());
}
@@ -364,11 +384,15 @@ public final class SessionManager implements TurnListener {
* CB-304 merged roster+live view. The registry is authoritative for worktree, branch,
* profile, owner, and state; the optional live agent supplies the herdr-reported status.
*/
public static Map<String, Object> rosterView(WorkerSession session, Agent live) {
public static Map<String, Object> rosterView(MemberSession session, Agent live) {
Map<String, Object> m = new LinkedHashMap<>();
m.put("sessionId", session.terminalId());
m.put("paneId", session.paneId());
// Both axes, always: profile says which backend this member runs on, role says what it is
// for. A lead reading the roster needs both — two rows may share a profile and still be
// allowed to do entirely different things.
m.put("profile", session.profile());
m.put("role", session.role() == null ? "dev" : session.role().wireName());
m.put("state", session.state().name().toLowerCase());
if (session.worktree() != null) {
m.put("worktree", session.worktree());
@@ -385,7 +409,7 @@ public final class SessionManager implements TurnListener {
/** Lifecycle hook: worker became available on the bridge MCP. */
void onReady(String terminalId) {
transitionByTerminal(terminalId, WorkerSession.State.SPAWNING, WorkerSession.State.READY);
transitionByTerminal(terminalId, MemberSession.State.SPAWNING, MemberSession.State.READY);
}
/**
@@ -395,13 +419,13 @@ public final class SessionManager implements TurnListener {
*/
@Override
public void onDelivered(String target) {
WorkerSession current = findByTerminal(target);
MemberSession current = findByTerminal(target);
if (current == null) return;
if (current.state() != WorkerSession.State.READY && current.state() != WorkerSession.State.DONE) {
if (current.state() != MemberSession.State.READY && current.state() != MemberSession.State.DONE) {
return;
}
long now = nowNanos.getAsLong();
WorkerSession updated = current.withState(WorkerSession.State.BUSY).bumpTurn(now);
MemberSession updated = current.withState(MemberSession.State.BUSY).bumpTurn(now);
if (replace(current, updated)) {
log.debug("session transitioned terminal={} pane={} {} -> BUSY turn={}",
target, current.paneId(), current.state(), updated.turnCount());
@@ -417,8 +441,8 @@ public final class SessionManager implements TurnListener {
@Override
public boolean hasPostTurnAction(String target) {
if (!clearAfterTurn) return false;
WorkerSession current = findByTerminal(target);
return current != null && current.state() == WorkerSession.State.BUSY
MemberSession current = findByTerminal(target);
return current != null && current.state() == MemberSession.State.BUSY
&& (contextCap <= 0 || current.turnCount() < contextCap);
}
@@ -428,10 +452,10 @@ public final class SessionManager implements TurnListener {
}
private boolean completeTurn(String target, boolean startContextReset) {
WorkerSession current = findByTerminal(target);
if (current == null || current.state() != WorkerSession.State.BUSY) return false;
MemberSession current = findByTerminal(target);
if (current == null || current.state() != MemberSession.State.BUSY) return false;
long now = nowNanos.getAsLong();
WorkerSession updated = current.withState(WorkerSession.State.DONE).withActivity(now);
MemberSession updated = current.withState(MemberSession.State.DONE).withActivity(now);
if (replace(current, updated)) {
log.debug("session transitioned terminal={} pane={} BUSY -> DONE turn={}",
target, current.paneId(), updated.turnCount());
@@ -458,10 +482,10 @@ public final class SessionManager implements TurnListener {
/** Lifecycle hook: the worker vanished or was dropped mid-life. */
void onFailed(String target) {
WorkerSession current = findByTerminal(target);
MemberSession current = findByTerminal(target);
if (current == null) return;
if (current.state() == WorkerSession.State.RELEASED) return;
if (replace(current, current.withState(WorkerSession.State.FAILED))) {
if (current.state() == MemberSession.State.RELEASED) return;
if (replace(current, current.withState(MemberSession.State.FAILED))) {
log.debug("session marked failed terminal={} pane={}", target, current.paneId());
}
}
@@ -474,8 +498,8 @@ public final class SessionManager implements TurnListener {
int reapIdle(long idleTtlNanos) {
long now = nowNanos.getAsLong();
int reaped = 0;
for (WorkerSession s : roster()) {
if (s.state() != WorkerSession.State.READY && s.state() != WorkerSession.State.DONE) {
for (MemberSession s : roster()) {
if (s.state() != MemberSession.State.READY && s.state() != MemberSession.State.DONE) {
continue;
}
if (now - s.lastActivityAtNanos() > idleTtlNanos) {
@@ -500,12 +524,12 @@ public final class SessionManager implements TurnListener {
*/
void drainAll(long timeoutNanos) {
long deadline = System.nanoTime() + timeoutNanos;
for (WorkerSession s : roster()) {
for (MemberSession s : roster()) {
try {
if (s.state() == WorkerSession.State.BUSY) {
if (s.state() == MemberSession.State.BUSY) {
while (System.nanoTime() < deadline) {
WorkerSession current = registry.get(s.paneId());
if (current == null || current.state() != WorkerSession.State.BUSY) {
MemberSession current = registry.get(s.paneId());
if (current == null || current.state() != MemberSession.State.BUSY) {
break;
}
try {
@@ -544,22 +568,22 @@ public final class SessionManager implements TurnListener {
* <p>A null {@code terminalId} is a normal input, not a caller bug: every lifecycle hook here is
* fed from the MCP transport, where the <em>primary</em> resolves to a {@link
* dev.ltms.bridged.auth.Principal} with no terminal. {@code BridgeMcp} documents that contact as
* a no-op, and {@link dev.ltms.bridged.inject.WorkerPresence#markPresent} honours it — but
* a no-op, and {@link dev.ltms.bridged.inject.MemberPresence#markPresent} honours it — but
* {@code PresenceBridge} then forwards the same null here. Matching on a null id can never
* succeed anyway (a registered session always has a terminal), so answer "no match" rather than
* throwing: an NPE on this path takes down an unrelated tool call for the primary.
*/
private WorkerSession findByTerminal(String terminalId) {
private MemberSession findByTerminal(String terminalId) {
if (terminalId == null) return null;
for (WorkerSession s : registry.values()) {
for (MemberSession 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);
private void transitionByTerminal(String terminalId, MemberSession.State from,
MemberSession.State to) {
MemberSession current = findByTerminal(terminalId);
if (current == null || current.state() != from) return;
long now = nowNanos.getAsLong();
if (replace(current, current.withState(to).withActivity(now))) {
@@ -568,12 +592,12 @@ public final class SessionManager implements TurnListener {
}
}
private boolean replace(WorkerSession expected, WorkerSession updated) {
private boolean replace(MemberSession expected, MemberSession updated) {
return registry.replace(expected.paneId(), expected, updated);
}
/** WorkerPresence bridge that also drives the manager's READY transition. */
private static final class PresenceBridge extends WorkerPresence {
/** MemberPresence bridge that also drives the manager's READY transition. */
private static final class PresenceBridge extends MemberPresence {
private final SessionManager sessions;
PresenceBridge(SessionManager sessions) {