Merge CB-557: member taxonomy in the API surface
CI / contract (push) Successful in 1m0s
CI / build (push) Successful in 1m1s

This commit is contained in:
Dai Ha
2026-08-14 07:08:26 +02:00
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) {
@@ -1,6 +1,6 @@
package dev.ltms.bridged;
import dev.ltms.bridged.inject.WorkerPresence;
import dev.ltms.bridged.inject.MemberPresence;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
@@ -28,7 +28,7 @@ class BridgedDeliverabilityTest {
@Test
@DisplayName("a worker that has connected its MCP is deliverable")
void presentWorkerIsDeliverable() {
WorkerPresence presence = new WorkerPresence();
MemberPresence presence = new MemberPresence();
presence.markPresent("term_worker");
assertTrue(Bridged.deliverableTo(presence, leads(Map.of())).test("term_worker"));
@@ -37,13 +37,13 @@ class BridgedDeliverabilityTest {
@Test
@DisplayName("a worker still in its boot window is held back")
void absentWorkerIsNotDeliverable() {
assertFalse(Bridged.deliverableTo(new WorkerPresence(), leads(Map.of())).test("term_booting"));
assertFalse(Bridged.deliverableTo(new MemberPresence(), leads(Map.of())).test("term_booting"));
}
@Test
@DisplayName("a lead is deliverable without ever being marked present")
void leadIsDeliverableWithoutPresence() {
WorkerPresence presence = new WorkerPresence();
MemberPresence presence = new MemberPresence();
Predicate<String> deliverable =
Bridged.deliverableTo(presence, leads(Map.of("term_lead", "opus-5.0")));
@@ -54,7 +54,7 @@ class BridgedDeliverabilityTest {
@Test
@DisplayName("an unknown terminal is deliverable to neither")
void strangerIsNotDeliverable() {
WorkerPresence presence = new WorkerPresence();
MemberPresence presence = new MemberPresence();
presence.markPresent("term_worker");
assertFalse(Bridged.deliverableTo(presence, leads(Map.of("term_lead", "opus-5.0")))
@@ -65,7 +65,7 @@ class BridgedDeliverabilityTest {
@DisplayName("a lead discovered after startup becomes deliverable with no restart")
void leadSetIsReadThroughOnEveryCall() {
Map<String, String> discovered = new HashMap<>();
Predicate<String> deliverable = Bridged.deliverableTo(new WorkerPresence(), leads(discovered));
Predicate<String> deliverable = Bridged.deliverableTo(new MemberPresence(), leads(discovered));
assertFalse(deliverable.test("term_late"));
discovered.put("term_late", "gpt-sol-5.6"); // leadScan picks up a newly labelled tab
@@ -75,7 +75,7 @@ class BridgedDeliverabilityTest {
@Test
@DisplayName("forgetting a torn-down worker does not strip a lead of its deliverability")
void forgetDoesNotDisarmALead() {
WorkerPresence presence = new WorkerPresence();
MemberPresence presence = new MemberPresence();
Predicate<String> deliverable =
Bridged.deliverableTo(presence, leads(Map.of("term_lead", "opus-5.0")));
@@ -389,7 +389,7 @@ class InjectorTest {
@Test
void dropClearsWorkerPresence() {
// CB-114 (finding #1): a vanished worker's readiness must be forgotten so a stale entry cannot
// linger past the worker's life (WorkerPresence.forget had no caller before this).
// linger past the worker's life (MemberPresence.forget had no caller before this).
List<String> forgotten = new ArrayList<>();
Injector inj = new Injector(new AgentControl(herdr), TurnListener.NOOP, _ -> true, forgotten::add);
inj.enqueue(T, "orphan");
@@ -11,7 +11,7 @@ class WorkerPresenceTest {
@Test
void tracksPresenceAndForgets() {
WorkerPresence p = new WorkerPresence();
MemberPresence p = new MemberPresence();
assertFalse(p.isPresent("term_a"), "unseen worker is not available");
p.markPresent("term_a");
@@ -23,7 +23,7 @@ class WorkerPresenceTest {
@Test
void nullOrBlankMarkIsANoOp() {
WorkerPresence p = new WorkerPresence();
MemberPresence p = new MemberPresence();
assertDoesNotThrow(() -> {
p.markPresent(null);
p.markPresent(" ");
@@ -18,7 +18,7 @@ import dev.ltms.bridged.msg.MessageService;
import dev.ltms.bridged.msg.Rendezvous;
import dev.ltms.bridged.session.FakeWorktrees;
import dev.ltms.bridged.session.SessionManager;
import dev.ltms.bridged.worker.ClaudeCodeLauncher;
import dev.ltms.bridged.member.ClaudeCodeLauncher;
import io.modelcontextprotocol.spec.McpSchema;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Test;
@@ -11,9 +11,10 @@ import dev.ltms.bridged.msg.MessageService;
import dev.ltms.bridged.msg.Rendezvous;
import dev.ltms.bridged.session.FakeWorktrees;
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.worker.ClaudeCodeLauncher;
import dev.ltms.bridged.member.ClaudeCodeLauncher;
import io.modelcontextprotocol.spec.McpSchema;
import dev.ltms.bridged.msg.InMemoryReplyInbox;
import org.junit.jupiter.api.BeforeEach;
@@ -233,7 +234,7 @@ class BridgeMcpTest {
String out = textOf(res);
assertTrue(out.contains("\"sessionId\":\"term_new_1\""), out);
// CB-519: the "paneId" wire field now carries the host-unique opaque id, not the herdr pane.
WorkerSession s = sm.roster().getFirst();
MemberSession s = sm.roster().getFirst();
assertTrue(out.contains("\"paneId\":\"" + s.paneId() + "\""), out);
assertNotEquals("w9:pRoot_1", s.paneId(), "the id is decoupled from the herdr pane coordinate");
assertTrue(out.contains("\"status\":\"spawning\""), out);
@@ -262,7 +263,7 @@ class BridgeMcpTest {
void spawnPassesTheRequestedCwdToTheWorker() {
FakeHerdr h = new FakeHerdr();
McpSchema.CallToolResult res = BridgeMcp.spawn(
sessionManager(h, "http://gx00.gw:8000", Set.of("gx00.gw")), null, "/req/dir", null, null, null);
sessionManager(h, "http://gx00.gw:8000", Set.of("gx00.gw")), null, null, "/req/dir", null, null, null);
assertNotEquals(Boolean.TRUE, res.isError());
// Protocol 19: the requested cwd roots the worker's pane at creation (tab.create).
@SuppressWarnings("unchecked")
@@ -285,7 +286,7 @@ class BridgeMcpTest {
FakeHerdr h = new FakeHerdr();
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt");
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), worktrees);
WorkerSession s = sessions.acquire("ltms-local", null, "/caller/proj", "term_primary",
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", "term_primary",
new WorktreeRequest("cb-304", null));
McpSchema.CallToolResult res = BridgeMcp.listFleet(
@@ -333,10 +334,10 @@ class BridgeMcpTest {
McpSchema.CallToolResult res = BridgeMcp.listFleet(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, Map.of(), "");
// An absent "leads" key is what made an empty worker roster read as "no peers" (CB-535).
// An absent "leads" key is what made an empty member roster read as "no peers" (CB-535).
String out = textOf(res);
assertTrue(out.contains("\"leads\":[]"), out);
assertTrue(out.contains("\"workers\":[]"), out);
assertTrue(out.contains("\"members\":[]"), out);
}
@Test
@@ -434,7 +435,7 @@ class BridgeMcpTest {
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt");
SessionManager sessions = new SessionManager(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), worktrees);
WorkerSession s = sessions.acquire("ltms-local", null, "/caller/proj", "term_primary",
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", "term_primary",
new WorktreeRequest("cb-517", null));
McpSchema.CallToolResult res = BridgeMcp.whoami(Principal.worker(s.terminalId(), 200), sessions);
@@ -517,4 +518,73 @@ class BridgeMcpTest {
BridgeMcp.recordPrimarySingleton(legacy, "term_x", null);
assertTrue(legacy.primaryTerminal().isEmpty(), "no caller means nothing is recorded");
}
// --- member taxonomy (CB-557) ---------------------------------------------------------
@Test
void bridgeListReportsMembersNotWorkersAndCarriesEachRole() {
FakeHerdr h = new FakeHerdr();
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
sessions.acquire("ltms-local", MemberRole.REVIEWER, null, "/caller/proj", "term_primary", null);
McpSchema.CallToolResult res = BridgeMcp.listFleet(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, Map.of(), "");
String out = textOf(res);
assertTrue(out.contains("\"members\":"), "the roster half is named members: " + out);
assertFalse(out.contains("\"workers\":"), "the old key must be gone: " + out);
assertTrue(out.contains("\"role\":\"reviewer\""), out);
}
/**
* Role and profile are separate axes, so the roster has to report both. Two members on one
* backend may still be allowed to do entirely different things.
*/
@Test
void aRosterRowCarriesBothItsRoleAndItsProfile() {
FakeHerdr h = new FakeHerdr();
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
sessions.acquire("ltms-local", MemberRole.DEV, null, "/caller/proj", "term_primary", null);
sessions.acquire("ltms-local", MemberRole.REVIEWER, null, "/caller/proj", "term_primary", null);
String out = textOf(BridgeMcp.listFleet(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, Map.of(), ""));
assertTrue(out.contains("\"role\":\"dev\""), out);
assertTrue(out.contains("\"role\":\"reviewer\""), out);
assertEquals(2, out.split("\"profile\":\"ltms-local\"", -1).length - 1,
"both members share one profile — that is the point: " + out);
}
@Test
void spawnDefaultsTheRoleToDev() {
FakeHerdr h = new FakeHerdr();
McpSchema.CallToolResult res = BridgeMcp.spawn(
sessionManager(h, "http://gx00.gw:8000", Set.of("gx00.gw")), null, null, null, null, null, null);
assertNotEquals(Boolean.TRUE, res.isError());
assertTrue(textOf(res).contains("\"role\":\"dev\""), textOf(res));
}
@Test
void spawnAcceptsAnExplicitRole() {
FakeHerdr h = new FakeHerdr();
McpSchema.CallToolResult res = BridgeMcp.spawn(
sessionManager(h, "http://gx00.gw:8000", Set.of("gx00.gw")), null, "architect",
null, null, null, null);
assertNotEquals(Boolean.TRUE, res.isError());
assertTrue(textOf(res).contains("\"role\":\"architect\""), textOf(res));
}
@Test
void spawnRejectsAnUnknownRoleAndNamesTheValidOnes() {
FakeHerdr h = new FakeHerdr();
McpSchema.CallToolResult res = BridgeMcp.spawn(
sessionManager(h, "http://gx00.gw:8000", Set.of("gx00.gw")), null, "worker",
null, null, null, null);
assertEquals(Boolean.TRUE, res.isError());
assertTrue(textOf(res).contains("architect, dev, reviewer"), textOf(res));
}
}
@@ -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.GuardException;
@@ -1,4 +1,4 @@
package dev.ltms.bridged.worker;
package dev.ltms.bridged.member;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.spi.ILoggingEvent;
@@ -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;
@@ -1,4 +1,4 @@
package dev.ltms.bridged.worker;
package dev.ltms.bridged.member;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
@@ -2,7 +2,8 @@ package dev.ltms.bridged.msg;
import dev.ltms.bridged.herdr.AgentStatus;
import dev.ltms.bridged.mcp.PrimaryRegistry;
import dev.ltms.bridged.session.WorkerSession;
import dev.ltms.bridged.peer.MemberRole;
import dev.ltms.bridged.session.MemberSession;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
@@ -191,10 +192,10 @@ class LeadHeartbeatLoopTest {
InMemoryReplyInbox inbox = new InMemoryReplyInbox();
inbox.own("term_w1");
inbox.publish("term_w1", "m1", "hello");
WorkerSession done = new WorkerSession("p1", "term_w1", "prof", "/cwd", null,
0, 0, 0, WorkerSession.State.DONE, null, null);
WorkerSession ready = new WorkerSession("p2", "term_w2", "prof", "/cwd", null,
0, 0, 0, WorkerSession.State.READY, null, null);
MemberSession done = new MemberSession("p1", "term_w1", "prof", MemberRole.DEV, "/cwd", null,
0, 0, 0, MemberSession.State.DONE, null, null);
MemberSession ready = new MemberSession("p2", "term_w2", "prof", MemberRole.DEV, "/cwd", null,
0, 0, 0, MemberSession.State.READY, null, null);
LeadHeartbeatLoop.FleetState fs = LeadHeartbeatLoop.snapshot(inbox, () -> List.of(done, ready));
@@ -15,7 +15,7 @@ import dev.ltms.bridged.msg.MessageService;
import dev.ltms.bridged.msg.Rendezvous;
import dev.ltms.bridged.session.FakeWorktrees;
import dev.ltms.bridged.session.SessionManager;
import dev.ltms.bridged.worker.ClaudeCodeLauncher;
import dev.ltms.bridged.member.ClaudeCodeLauncher;
import io.javalin.Javalin;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Test;
@@ -128,9 +128,9 @@ class BridgedAppAuthTest {
void aWorkerMayNotOrchestrate() throws Exception {
int port = start(FakeHerdr.WORKER_PID, false, null);
assertEquals(403, send(port, "POST", "/workers", null, null).statusCode(),
assertEquals(403, send(port, "POST", "/members", null, null).statusCode(),
"a worker spawning workers would be escalating into the orchestrator role");
assertEquals(403, send(port, "DELETE", "/workers/w2:p7", null, null).statusCode());
assertEquals(403, send(port, "DELETE", "/members/w2:p7", null, null).statusCode());
assertEquals(403, send(port, "POST", "/sessions/term_b/message",
"{\"content\":\"hi\"}", null).statusCode());
assertEquals(403, send(port, "GET", "/sessions/term_a/replies", null, null).statusCode(),
@@ -9,7 +9,7 @@ import dev.ltms.bridged.herdr.FakeHerdr;
import dev.ltms.bridged.herdr.WorkspaceControl;
import dev.ltms.bridged.inject.Injector;
import dev.ltms.bridged.inject.StatusPoller;
import dev.ltms.bridged.inject.WorkerPresence;
import dev.ltms.bridged.inject.MemberPresence;
import dev.ltms.bridged.msg.InMemoryReplyInbox;
import dev.ltms.bridged.msg.MessageService;
import dev.ltms.bridged.msg.Rendezvous;
@@ -17,7 +17,7 @@ import dev.ltms.bridged.session.FakeWorktrees;
import dev.ltms.bridged.session.GitWorktrees;
import dev.ltms.bridged.session.SessionManager;
import dev.ltms.bridged.session.Worktrees;
import dev.ltms.bridged.worker.ClaudeCodeLauncher;
import dev.ltms.bridged.member.ClaudeCodeLauncher;
import io.javalin.Javalin;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Test;
@@ -43,7 +43,7 @@ class BridgedAppTest {
private final ObjectMapper mapper = new ObjectMapper();
private final HttpClient http = HttpClient.newHttpClient();
private WorkerPresence presence;
private MemberPresence presence;
private Javalin app;
private StatusPoller poller;
@@ -160,7 +160,7 @@ class BridgedAppTest {
FakeHerdr herdr = new FakeHerdr();
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
HttpResponse<String> res = req(port, "POST", "/workers");
HttpResponse<String> res = req(port, "POST", "/members");
assertEquals(201, res.statusCode());
JsonNode body = mapper.readTree(res.body());
// CB-519: the responded paneId is a host-unique opaque UUID, not the herdr pane coordinate.
@@ -202,12 +202,12 @@ class BridgedAppTest {
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"), "tab",
new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt"));
HttpResponse<String> spawn = req(port, "POST", "/workers?worktree=true&ticket=cb-304");
HttpResponse<String> spawn = req(port, "POST", "/members?worktree=true&ticket=cb-304");
assertEquals(201, spawn.statusCode());
JsonNode spawned = mapper.readTree(spawn.body());
String paneId = spawned.get("paneId").asText();
HttpResponse<String> res = req(port, "GET", "/workers");
HttpResponse<String> res = req(port, "GET", "/members");
assertEquals(200, res.statusCode());
JsonNode workers = mapper.readTree(res.body()).get("workers");
assertEquals(1, workers.size());
@@ -226,7 +226,7 @@ class BridgedAppTest {
void spawnWithACwdParamRootsTheWorkerThere() throws Exception {
FakeHerdr herdr = new FakeHerdr();
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
assertEquals(201, req(port, "POST", "/workers?cwd=/tmp/proj").statusCode());
assertEquals(201, req(port, "POST", "/members?cwd=/tmp/proj").statusCode());
assertEquals("/tmp/proj", params(herdr, "tab.create").get("cwd"), "the worker starts in cwd");
}
@@ -234,7 +234,7 @@ class BridgedAppTest {
void spawnWithAnUnknownProfileIs400() throws Exception {
FakeHerdr herdr = new FakeHerdr();
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
HttpResponse<String> res = req(port, "POST", "/workers?profile=nope");
HttpResponse<String> res = req(port, "POST", "/members?profile=nope");
assertEquals(400, res.statusCode());
assertEquals("unknown_profile", mapper.readTree(res.body()).get("error").asText());
assertFalse(herdr.called("agent.start"), "an unknown profile must not spawn anything");
@@ -246,7 +246,7 @@ class BridgedAppTest {
FakeHerdr herdr = new FakeHerdr().withWorkspace("w9", "bridged-workers");
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
assertEquals(201, req(port, "POST", "/workers").statusCode());
assertEquals(201, req(port, "POST", "/members").statusCode());
assertFalse(herdr.called("workspace.create"), "existing worker space must be reused, not recreated");
assertTrue(herdr.called("tab.create"), "a fresh tab is still created for the worker");
}
@@ -258,7 +258,7 @@ class BridgedAppTest {
FakeHerdr herdr = new FakeHerdr().agentNameTakenTimes(2);
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
assertEquals(201, req(port, "POST", "/workers").statusCode());
assertEquals(201, req(port, "POST", "/members").statusCode());
List<String> names = herdr.calls.stream()
.filter(c -> c.method().equals("agent.start"))
.map(c -> ((Map<String, Object>) c.params()).get("name").toString())
@@ -273,7 +273,7 @@ class BridgedAppTest {
FakeHerdr herdr = new FakeHerdr().agentNameTakenTimes(99);
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
assertEquals(500, req(port, "POST", "/workers").statusCode());
assertEquals(500, req(port, "POST", "/members").statusCode());
assertTrue(herdr.called("tab.create"), "a tab was created before the failed start");
assertEquals("w9:t2", params(herdr, "tab.close").get("tab_id"), "orphaned tab must be closed");
}
@@ -284,7 +284,7 @@ class BridgedAppTest {
// base_url points at the subscription — guard must block before any herdr call.
int port = start(herdr, "https://api.anthropic.com", Set.of("gx00.gw"));
HttpResponse<String> res = req(port, "POST", "/workers");
HttpResponse<String> res = req(port, "POST", "/members");
assertEquals(403, res.statusCode());
assertEquals("subscription_boundary", mapper.readTree(res.body()).get("error").asText());
assertFalse(herdr.called("agent.start"), "guard must stop the spawn before herdr");
@@ -296,7 +296,7 @@ class BridgedAppTest {
FakeHerdr herdr = new FakeHerdr();
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"), "pane");
assertEquals(201, req(port, "POST", "/workers").statusCode());
assertEquals(201, req(port, "POST", "/members").statusCode());
Map<String, Object> start = params(herdr, "agent.start");
assertFalse(start.containsKey("tab_id"), "pane placement must not target a tab");
assertFalse(herdr.called("workspace.create"), "pane placement uses no dedicated space");
@@ -308,7 +308,7 @@ class BridgedAppTest {
FakeHerdr herdr = new FakeHerdr();
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
assertEquals(204, req(port, "DELETE", "/workers/w9:pW").statusCode());
assertEquals(204, req(port, "DELETE", "/members/w9:pW").statusCode());
assertTrue(herdr.called("pane.close"));
// Tab resolved from the pane (pane.get), then closed.
assertEquals("w9:t2", params(herdr, "tab.close").get("tab_id"));
@@ -450,7 +450,7 @@ class BridgedAppTest {
FakeHerdr herdr = new FakeHerdr();
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"), "pane");
assertEquals(204, req(port, "DELETE", "/workers/w9:pW").statusCode());
assertEquals(204, req(port, "DELETE", "/members/w9:pW").statusCode());
assertTrue(herdr.called("pane.close"));
assertFalse(herdr.called("tab.close"), "pane placement owns no tab to close");
assertFalse(herdr.called("pane.get"), "no tab resolution in pane placement");
@@ -462,7 +462,7 @@ class BridgedAppTest {
FakeHerdr herdr = new FakeHerdr().withWorkerTabPaneCount(2);
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
assertEquals(204, req(port, "DELETE", "/workers/w9:pW").statusCode());
assertEquals(204, req(port, "DELETE", "/members/w9:pW").statusCode());
assertTrue(herdr.called("pane.close"), "the worker's own pane is still closed");
assertFalse(herdr.called("tab.close"), "must not close a tab that holds the user's other panes");
}
@@ -473,7 +473,7 @@ class BridgedAppTest {
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
// A genuine teardown failure must surface, not be reported as a successful 204.
assertEquals(500, req(port, "DELETE", "/workers/w9:pW").statusCode());
assertEquals(500, req(port, "DELETE", "/members/w9:pW").statusCode());
assertFalse(herdr.called("tab.close"), "tab is not removed when the pane close failed");
}
@@ -483,7 +483,7 @@ class BridgedAppTest {
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
// Already-gone is success; the (now-empty) tab is still cleaned up.
assertEquals(204, req(port, "DELETE", "/workers/w9:pW").statusCode());
assertEquals(204, req(port, "DELETE", "/members/w9:pW").statusCode());
assertTrue(herdr.called("tab.close"));
}
}
@@ -5,7 +5,7 @@ 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.ClaudeCodeLauncher;
import dev.ltms.bridged.member.ClaudeCodeLauncher;
import dev.ltms.bridged.peer.PeerUnreachableException;
import org.junit.jupiter.api.Test;
@@ -69,10 +69,10 @@ class SessionManagerTest {
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");
MemberSession a = sessions.acquire("ltms-local", "/work/a", "/caller/a", "term_primary");
MemberSession b = sessions.acquire("ltms-local", "/work/b", "/caller/b", "term_primary");
assertEquals(WorkerSession.State.SPAWNING, a.state(), "fresh session starts spawning");
assertEquals(MemberSession.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());
@@ -92,7 +92,7 @@ class SessionManagerTest {
// once a session existed, so this NPE'd the primary's second spawn while the first passed.
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
WorkerSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
assertDoesNotThrow(() -> sessions.asPresence().markPresent(null),
"the primary's null terminal must not blow up an unrelated tool call");
@@ -100,7 +100,7 @@ class SessionManagerTest {
assertDoesNotThrow(() -> sessions.onTurnComplete(null));
assertDoesNotThrow(() -> sessions.onTurnFailed(null));
assertEquals(WorkerSession.State.SPAWNING, sessions.get(session.paneId()).orElseThrow().state(),
assertEquals(MemberSession.State.SPAWNING, sessions.get(session.paneId()).orElseThrow().state(),
"and must not transition any registered session");
}
@@ -108,20 +108,20 @@ class SessionManagerTest {
void presenceMovesSpawningToReadyAndDeliveredTurnMovesToDone() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
WorkerSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
MemberSession 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(),
assertEquals(MemberSession.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(),
assertEquals(MemberSession.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(),
assertEquals(MemberSession.State.DONE, sessions.get(session.paneId()).orElseThrow().state(),
"turn completion moves BUSY → DONE");
}
@@ -129,7 +129,7 @@ class SessionManagerTest {
void releaseTearsDownWorkerAndRemovesFromRosterAndIsIdempotent() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
WorkerSession session = sessions.acquire("ltms-local", null, "/caller", null);
MemberSession session = sessions.acquire("ltms-local", null, "/caller", null);
String paneId = session.paneId();
sessions.release(paneId);
@@ -145,15 +145,15 @@ class SessionManagerTest {
void onTurnFailedMovesSessionToFailed() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
WorkerSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
MemberSession 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");
MemberSession updated = sessions.get(session.paneId()).orElseThrow();
assertEquals(MemberSession.State.FAILED, updated.state(), "turn failure moves to FAILED");
assertTrue(sessions.roster().contains(updated), "FAILED is still in acquired-minus-released roster");
}
@@ -161,11 +161,11 @@ class SessionManagerTest {
void recycleProducesNewPaneIdAndOldOneIsGone() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
WorkerSession oldSession = sessions.acquire("ltms-local", null, "/caller", "term_primary");
MemberSession oldSession = sessions.acquire("ltms-local", null, "/caller", "term_primary");
String oldPane = oldSession.paneId();
String oldTerminal = oldSession.terminalId();
WorkerSession fresh = sessions.recycle(oldPane);
MemberSession fresh = sessions.recycle(oldPane);
assertNotEquals(oldPane, fresh.paneId(), "recycle yields a new pane id");
assertNotEquals(oldTerminal, fresh.terminalId(), "recycle yields a new terminal id");
@@ -190,8 +190,8 @@ class SessionManagerTest {
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");
MemberSession a = sessions.acquire("ltms-local", "/a", "/caller", "ownerA");
MemberSession b = sessions.acquire("ltms-local", "/b", "/caller", "ownerB");
assertEquals(2, sessions.roster().size());
assertTrue(sessions.roster().stream().anyMatch(s -> s.paneId().equals(a.paneId())));
@@ -219,7 +219,7 @@ class SessionManagerTest {
long[] clock = {0};
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr, () -> clock[0]);
WorkerSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
String terminal = session.terminalId();
sessions.asPresence().markPresent(terminal);
@@ -234,13 +234,13 @@ class SessionManagerTest {
long[] clock = {0};
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr, () -> clock[0]);
WorkerSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
String terminal = session.terminalId();
sessions.asPresence().markPresent(terminal);
clock[0] = 5;
assertEquals(0, sessions.reapIdle(10), "READY session within TTL is not reaped");
assertEquals(WorkerSession.State.READY,
assertEquals(MemberSession.State.READY,
sessions.get(session.paneId()).orElseThrow().state(),
"READY session survives");
}
@@ -250,14 +250,14 @@ class SessionManagerTest {
long[] clock = {0};
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr, () -> clock[0]);
WorkerSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
String terminal = session.terminalId();
sessions.asPresence().markPresent(terminal);
sessions.onDelivered(terminal);
clock[0] = 100;
assertEquals(0, sessions.reapIdle(10), "BUSY session past TTL is never reaped");
assertEquals(WorkerSession.State.BUSY,
assertEquals(MemberSession.State.BUSY,
sessions.get(session.paneId()).orElseThrow().state(),
"BUSY session remains");
}
@@ -267,7 +267,7 @@ class SessionManagerTest {
long[] clock = {0};
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr, () -> clock[0]);
WorkerSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
String terminal = session.terminalId();
sessions.asPresence().markPresent(terminal);
sessions.onDelivered(terminal);
@@ -284,8 +284,8 @@ class SessionManagerTest {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr, () -> clock[0]);
WorkerSession ready = sessions.acquire("ltms-local", "/ready", "/caller", "owner1");
WorkerSession busy = sessions.acquire("ltms-local", "/busy", "/caller", "owner2");
MemberSession ready = sessions.acquire("ltms-local", "/ready", "/caller", "owner1");
MemberSession busy = sessions.acquire("ltms-local", "/busy", "/caller", "owner2");
sessions.asPresence().markPresent(ready.terminalId());
sessions.asPresence().markPresent(busy.terminalId());
sessions.onDelivered(busy.terminalId());
@@ -293,7 +293,7 @@ class SessionManagerTest {
clock[0] = 50;
assertEquals(1, sessions.reapIdle(30), "only READY past TTL is reaped");
assertTrue(sessions.get(ready.paneId()).isEmpty(), "READY session is gone");
assertEquals(WorkerSession.State.BUSY,
assertEquals(MemberSession.State.BUSY,
sessions.get(busy.paneId()).orElseThrow().state(),
"BUSY session is still registered");
}
@@ -302,7 +302,7 @@ class SessionManagerTest {
void contextCapDisabledSessionSurvivesMultipleTurns() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr, () -> 0L, 0);
WorkerSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
String terminal = session.terminalId();
sessions.asPresence().markPresent(terminal);
@@ -311,8 +311,8 @@ class SessionManagerTest {
sessions.onDelivered(terminal);
sessions.onTurnComplete(terminal);
WorkerSession updated = sessions.get(session.paneId()).orElseThrow();
assertEquals(WorkerSession.State.DONE, updated.state(), "session finishes second turn");
MemberSession updated = sessions.get(session.paneId()).orElseThrow();
assertEquals(MemberSession.State.DONE, updated.state(), "session finishes second turn");
assertEquals(2, updated.turnCount(), "turn count tracks both deliveries");
long releaseCloseCount = paneCloseCallsFor(herdr, "w9:pRoot_1"); // the real pane coordinate
assertEquals(0, releaseCloseCount, "cap disabled — no forced release of the worker pane");
@@ -322,13 +322,13 @@ class SessionManagerTest {
void contextCapTwoReleasesAfterSecondComplete() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr, () -> 0L, 2);
WorkerSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
String terminal = session.terminalId();
sessions.asPresence().markPresent(terminal);
sessions.onDelivered(terminal);
sessions.onTurnComplete(terminal);
assertEquals(WorkerSession.State.DONE,
assertEquals(MemberSession.State.DONE,
sessions.get(session.paneId()).orElseThrow().state(),
"first turn completes without release");
@@ -345,13 +345,13 @@ class SessionManagerTest {
void clearAfterTurnResetsContextWithoutDoubleCountingTheTurn() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr, () -> 0L, 0, true);
WorkerSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
sessions.asPresence().markPresent(session.terminalId());
sessions.onDelivered(session.terminalId());
assertTrue(sessions.onTurnCompleteWithPostAction(session.terminalId()));
WorkerSession updated = sessions.get(session.paneId()).orElseThrow();
MemberSession updated = sessions.get(session.paneId()).orElseThrow();
assertEquals(1, updated.turnCount(), "the reset is housekeeping, not a second delegation");
assertEquals(List.of("/clear"), promptTexts(herdr));
}
@@ -360,7 +360,7 @@ class SessionManagerTest {
void contextCapReleaseWinsOverClearAfterTurn() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr, () -> 0L, 1, true);
WorkerSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
sessions.asPresence().markPresent(session.terminalId());
sessions.onDelivered(session.terminalId());
@@ -375,13 +375,13 @@ class SessionManagerTest {
void clearAfterTurnFalsePreservesCompletionWithoutAControlPrompt() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr, () -> 0L, 0, false);
WorkerSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
sessions.asPresence().markPresent(session.terminalId());
sessions.onDelivered(session.terminalId());
sessions.onTurnComplete(session.terminalId());
assertEquals(WorkerSession.State.DONE, sessions.get(session.paneId()).orElseThrow().state());
assertEquals(MemberSession.State.DONE, sessions.get(session.paneId()).orElseThrow().state());
assertTrue(promptTexts(herdr).isEmpty());
}
@@ -391,8 +391,8 @@ class SessionManagerTest {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr, () -> clock[0]);
WorkerSession ready = sessions.acquire("ltms-local", "/ready", "/caller", "ownerR");
WorkerSession busy = sessions.acquire("ltms-local", "/busy", "/caller", "ownerB");
MemberSession ready = sessions.acquire("ltms-local", "/ready", "/caller", "ownerR");
MemberSession busy = sessions.acquire("ltms-local", "/busy", "/caller", "ownerB");
sessions.asPresence().markPresent(ready.terminalId());
sessions.asPresence().markPresent(busy.terminalId());
sessions.onDelivered(busy.terminalId());
@@ -461,7 +461,7 @@ class SessionManagerTest {
java.util.List<String> released = new java.util.concurrent.CopyOnWriteArrayList<>();
sessions.onRelease(released::add);
WorkerSession s = sessions.acquire("ltms-local", null, "/caller", null);
MemberSession s = sessions.acquire("ltms-local", null, "/caller", null);
sessions.release(s.paneId());
assertEquals(java.util.List.of(s.terminalId()), released,
@@ -488,7 +488,7 @@ class SessionManagerTest {
throw new IllegalStateException("listener blew up");
});
WorkerSession s = sessions.acquire("ltms-local", null, "/caller", null);
MemberSession s = sessions.acquire("ltms-local", null, "/caller", null);
assertDoesNotThrow(() -> sessions.release(s.paneId()),
"a listener failure must never prevent the teardown it is reacting to");
assertTrue(sessions.get(s.paneId()).isEmpty(), "and the session is still deregistered");
@@ -5,7 +5,7 @@ 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.ClaudeCodeLauncher;
import dev.ltms.bridged.member.ClaudeCodeLauncher;
import org.junit.jupiter.api.Test;
import java.util.List;
@@ -5,7 +5,7 @@ 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.ClaudeCodeLauncher;
import dev.ltms.bridged.member.ClaudeCodeLauncher;
import org.junit.jupiter.api.Test;
import java.util.List;
@@ -45,7 +45,7 @@ class WorktreeSessionManagerTest {
FakeWorktrees worktrees = new FakeWorktrees();
SessionManager sessions = new SessionManager(workerService(herdr), worktrees);
WorkerSession s = sessions.acquire("ltms-local", null, "/caller/proj", "term_primary");
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", "term_primary");
assertTrue(worktrees.addCalls().isEmpty(), "shared-tree acquire never adds a worktree");
assertTrue(worktrees.repoRootCalls().isEmpty(), "shared-tree acquire never resolves a repo root");
@@ -62,7 +62,7 @@ class WorktreeSessionManagerTest {
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt");
SessionManager sessions = new SessionManager(workerService(herdr), worktrees);
WorkerSession s = sessions.acquire("ltms-local", null, "/caller/proj", "term_primary",
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", "term_primary",
new WorktreeRequest("cb-999", null));
assertEquals(1, worktrees.addCalls().size(), "one worktree was added");
@@ -110,7 +110,7 @@ class WorktreeSessionManagerTest {
FakeHerdr herdr = new FakeHerdr();
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt");
SessionManager sessions = new SessionManager(workerService(herdr), worktrees);
WorkerSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-666", null));
String paneId = s.paneId();
@@ -131,7 +131,7 @@ class WorktreeSessionManagerTest {
FakeHerdr herdr = new FakeHerdr();
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt");
SessionManager sessions = new SessionManager(workerService(herdr), worktrees);
WorkerSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-544", null));
sessions.asPresence().markPresent(s.terminalId()); // READY (idle)
@@ -148,7 +148,7 @@ class WorktreeSessionManagerTest {
FakeHerdr herdr = new FakeHerdr();
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt");
SessionManager sessions = new SessionManager(workerService(herdr), worktrees);
WorkerSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-544", null));
String terminal = s.terminalId();
sessions.asPresence().markPresent(terminal);
@@ -167,7 +167,7 @@ class WorktreeSessionManagerTest {
FakeHerdr herdr = new FakeHerdr();
FakeWorktrees worktrees = new FakeWorktrees();
SessionManager sessions = new SessionManager(workerService(herdr), worktrees);
WorkerSession s = sessions.acquire("ltms-local", null, "/caller/proj", null);
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null);
sessions.release(s.paneId());
@@ -196,9 +196,9 @@ class WorktreeSessionManagerTest {
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt");
SessionManager sessions = new SessionManager(workerService(herdr), worktrees);
WorkerSession a = sessions.acquire("ltms-local", null, "/caller/proj", null,
MemberSession a = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-444", null));
WorkerSession b = sessions.acquire("ltms-local", null, "/caller/proj", null,
MemberSession b = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-444", null));
assertNotEquals(a.branch(), b.branch(), "branches are distinct");