Compare commits
14 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| d04b075996 | |||
| 644927636d | |||
| 95e45007aa | |||
| 7a583c4045 | |||
| 65bce058bb | |||
| 48b437083b | |||
| fa0612859b | |||
| 4ffbcd0b7d | |||
| a0cd053fd9 | |||
| f9a5e066b5 | |||
| e9bc192160 | |||
| e0a57988ad | |||
| b414a74c26 | |||
| 2c3796d598 |
@@ -70,9 +70,6 @@ public final class Bridged {
|
||||
|
||||
private static final Logger log = LoggerFactory.getLogger(Bridged.class);
|
||||
|
||||
/** How often the injector samples a busy worker's status while it has queued work. */
|
||||
private static final long INJECT_POLL_MILLIS = 250;
|
||||
|
||||
/** CB-504: how long to wait at startup for herdr's socket before serving degraded. */
|
||||
private static final long HERDR_WAIT_SECONDS = 30;
|
||||
private static final long HERDR_WAIT_POLL_MILLIS = 500;
|
||||
@@ -243,6 +240,7 @@ public final class Bridged {
|
||||
// what CallerResolver resolves against and what that lifecycle will read profiles from;
|
||||
// nothing here spawns a slot.
|
||||
MemberRegistry members = new MemberRegistry(cfg.fleet());
|
||||
sessions.setMemberLifecycle(members);
|
||||
if (!members.slots().isEmpty()) {
|
||||
log.info("member slots: {} configured {} — none bound yet (a slot is idle until the "
|
||||
+ "spawn lifecycle binds a live terminal to it)",
|
||||
@@ -289,7 +287,7 @@ public final class Bridged {
|
||||
};
|
||||
Injector injector = new Injector(agents, turnListener, deliverableTo(presence, leads),
|
||||
presence::forget);
|
||||
StatusPoller poller = new StatusPoller(agents, injector, INJECT_POLL_MILLIS);
|
||||
StatusPoller poller = new StatusPoller(agents, injector, Injector.POLL_INTERVAL_MILLIS);
|
||||
poller.start();
|
||||
|
||||
// CB-307: reply inbox. A broker: block (with a uri) selects the AMQP-backed durable adapter;
|
||||
@@ -375,13 +373,11 @@ public final class Bridged {
|
||||
throw new IllegalStateException("auth.mode=token but env var " + cfg.auth().tokenEnv()
|
||||
+ " is unset or empty — export it before starting bridged");
|
||||
}
|
||||
callers = CallerResolver.withLeadsAndMembers(identity, true, token, leads,
|
||||
members::snapshot);
|
||||
callers = CallerResolver.withLeadsAndMembers(identity, true, token, leads, members);
|
||||
log.info("auth: token mode (bearer required for non-worker callers, env {})",
|
||||
cfg.auth().tokenEnv());
|
||||
} else {
|
||||
callers = CallerResolver.withLeadsAndMembers(identity, false, null, leads,
|
||||
members::snapshot);
|
||||
callers = CallerResolver.withLeadsAndMembers(identity, false, null, leads, members);
|
||||
log.info("auth: loopback-trust (any loopback non-worker caller is the primary)");
|
||||
}
|
||||
|
||||
@@ -429,19 +425,20 @@ public final class Bridged {
|
||||
}
|
||||
|
||||
/**
|
||||
* The {@link Injector}'s readiness gate (CB-534): a target is deliverable if it is a worker whose
|
||||
* agent has connected the bridge MCP, <em>or</em> a lead.
|
||||
* The {@link Injector}'s readiness gate (CB-534): a target is deliverable if it is a spawned
|
||||
* member whose agent has connected the bridge MCP, <em>or</em> a lead.
|
||||
*
|
||||
* <p>The gate exists for one reason — to hold a delivery out of a <em>spawned</em> worker's boot
|
||||
* <p>The gate exists for one reason — to hold a delivery out of a <em>spawned</em> member's boot
|
||||
* window, where herdr already reports {@code idle} but the TUI would drop an injected paste. That
|
||||
* 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 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
|
||||
* {@code READINESS_GRACE_POLLS} (~60s) and then failed having never been typed into the pane.
|
||||
* for every spawned member (worker and architect), deliberately, since that map doubles as the
|
||||
* member roster's availability signal and a lead counted there would show up as an available
|
||||
* member. So without the second disjunct a lead is permanently un-deliverable: every
|
||||
* lead→lead send sat on the gate for {@code READINESS_GRACE_POLLS} (~60s) and then failed
|
||||
* having never been typed into the pane.
|
||||
*
|
||||
* <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.
|
||||
|
||||
@@ -1,10 +1,12 @@
|
||||
package dev.ltms.bridged.auth;
|
||||
|
||||
import dev.ltms.bridged.mcp.ConnectionIdentity;
|
||||
import dev.ltms.bridged.peer.MemberRole;
|
||||
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.security.MessageDigest;
|
||||
import java.util.Map;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.Supplier;
|
||||
|
||||
/**
|
||||
@@ -60,14 +62,16 @@ public final class CallerResolver {
|
||||
* in {@code Bridged} reads a constant from config, which is the degenerate live case.
|
||||
*/
|
||||
private final Supplier<Map<String, String>> architectTerminals;
|
||||
private final Function<String, MemberRole> memberSlotRoles;
|
||||
private final Function<String, String> memberSlotNames;
|
||||
|
||||
/** Loopback-trust resolver: no token required, historical behaviour. */
|
||||
public CallerResolver(ConnectionIdentity identity) {
|
||||
/** Loopback-trust resolver: no token required, historical behaviour. Test-only. */
|
||||
CallerResolver(ConnectionIdentity identity) {
|
||||
this(identity, false, null, Map.of());
|
||||
}
|
||||
|
||||
/** As {@link #CallerResolver(ConnectionIdentity, boolean, String, Map)} with no leads pinned. */
|
||||
public CallerResolver(ConnectionIdentity identity, boolean tokenMode, String token) {
|
||||
/** As {@link #CallerResolver(ConnectionIdentity, boolean, String, Map)} with no leads pinned. Test-only. */
|
||||
CallerResolver(ConnectionIdentity identity, boolean tokenMode, String token) {
|
||||
this(identity, tokenMode, token, Map.of());
|
||||
}
|
||||
|
||||
@@ -82,8 +86,8 @@ public final class CallerResolver {
|
||||
* @param pinnedPrimaryTerminal the primary's own herdr {@code terminal_id}
|
||||
* ({@code null}/blank = unpinned)
|
||||
*/
|
||||
public static CallerResolver pinnedTo(ConnectionIdentity identity, boolean tokenMode,
|
||||
String token, String pinnedPrimaryTerminal) {
|
||||
static CallerResolver pinnedTo(ConnectionIdentity identity, boolean tokenMode,
|
||||
String token, String pinnedPrimaryTerminal) {
|
||||
return new CallerResolver(identity, tokenMode, token,
|
||||
pinnedPrimaryTerminal == null || pinnedPrimaryTerminal.isBlank()
|
||||
? Map.of() : Map.of(pinnedPrimaryTerminal, "primary"));
|
||||
@@ -98,21 +102,11 @@ public final class CallerResolver {
|
||||
* {@link Role#PRIMARY} — rather than a worker. Empty = nothing pinned,
|
||||
* so every pane resolves as a worker.
|
||||
*/
|
||||
public CallerResolver(ConnectionIdentity identity, boolean tokenMode, String token,
|
||||
Map<String, String> leadTerminals) {
|
||||
CallerResolver(ConnectionIdentity identity, boolean tokenMode, String token,
|
||||
Map<String, String> leadTerminals) {
|
||||
this(identity, tokenMode, token, fixed(leadTerminals), null);
|
||||
}
|
||||
|
||||
/**
|
||||
* Map-form of both registries (CB-548): lead terminals and the initial architect terminal
|
||||
* bindings, each snapshotted at construction (a handed-over map is not offered as live state).
|
||||
*/
|
||||
public CallerResolver(ConnectionIdentity identity, boolean tokenMode, String token,
|
||||
Map<String, String> leadTerminals,
|
||||
Map<String, String> architectTerminals) {
|
||||
this(identity, tokenMode, token, fixed(leadTerminals), fixed(architectTerminals));
|
||||
}
|
||||
|
||||
/**
|
||||
* Live-registry form: {@code leadTerminals} is consulted on every resolve, so leads discovered
|
||||
* after startup (CB-531's tab scan) take effect without a restart.
|
||||
@@ -121,26 +115,26 @@ public final class CallerResolver {
|
||||
* {@link #pinnedTo}: {@code Map} and {@code Supplier} overloads are ambiguous for a literal
|
||||
* {@code null}.
|
||||
*/
|
||||
public static CallerResolver withLeads(ConnectionIdentity identity, boolean tokenMode,
|
||||
String token,
|
||||
Supplier<Map<String, String>> leadTerminals) {
|
||||
static CallerResolver withLeads(ConnectionIdentity identity, boolean tokenMode,
|
||||
String token,
|
||||
Supplier<Map<String, String>> leadTerminals) {
|
||||
return new CallerResolver(identity, tokenMode, token, leadTerminals, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* Live-registry form for both {@code leadTerminals} and the CB-548 architect registry: both
|
||||
* are consulted on every resolve, so a slot binding injected after startup takes effect
|
||||
* without a restart.
|
||||
* Live registry form that can confirm a bound slot is an architect slot.
|
||||
*
|
||||
* <p>A static factory rather than a constructor overload, for the same reason as
|
||||
* {@link #pinnedTo}: too many {@code Map}/{@code Supplier} combinations to make {@code null}
|
||||
* unambiguous.
|
||||
* <p>This is the only public construction path. It keeps terminal bindings and slot roles in
|
||||
* the same {@link MemberRegistry}, so a configured architect can resolve as an architect.
|
||||
*/
|
||||
public static CallerResolver withLeadsAndMembers(ConnectionIdentity identity,
|
||||
boolean tokenMode, String token,
|
||||
Supplier<Map<String, String>> leadTerminals,
|
||||
Supplier<Map<String, String>> architectTerminals) {
|
||||
return new CallerResolver(identity, tokenMode, token, leadTerminals, architectTerminals);
|
||||
boolean tokenMode, String token,
|
||||
Supplier<Map<String, String>> leadTerminals,
|
||||
MemberRegistry members) {
|
||||
return new CallerResolver(identity, tokenMode, token, leadTerminals,
|
||||
members == null ? null : members::snapshot,
|
||||
members == null ? null : members::roleForSlot,
|
||||
members == null ? null : members::nameForSlot);
|
||||
}
|
||||
|
||||
private static Supplier<Map<String, String>> fixed(Map<String, String> leadTerminals) {
|
||||
@@ -151,6 +145,21 @@ public final class CallerResolver {
|
||||
private CallerResolver(ConnectionIdentity identity, boolean tokenMode, String token,
|
||||
Supplier<Map<String, String>> leadTerminals,
|
||||
Supplier<Map<String, String>> architectTerminals) {
|
||||
this(identity, tokenMode, token, leadTerminals, architectTerminals, null);
|
||||
}
|
||||
|
||||
private CallerResolver(ConnectionIdentity identity, boolean tokenMode, String token,
|
||||
Supplier<Map<String, String>> leadTerminals,
|
||||
Supplier<Map<String, String>> architectTerminals,
|
||||
Function<String, MemberRole> memberSlotRoles) {
|
||||
this(identity, tokenMode, token, leadTerminals, architectTerminals, memberSlotRoles, null);
|
||||
}
|
||||
|
||||
private CallerResolver(ConnectionIdentity identity, boolean tokenMode, String token,
|
||||
Supplier<Map<String, String>> leadTerminals,
|
||||
Supplier<Map<String, String>> architectTerminals,
|
||||
Function<String, MemberRole> memberSlotRoles,
|
||||
Function<String, String> memberSlotNames) {
|
||||
if (tokenMode && (token == null || token.isBlank())) {
|
||||
throw new IllegalArgumentException(
|
||||
"auth.mode=token requires a non-empty token; check that the env var named by "
|
||||
@@ -161,6 +170,8 @@ public final class CallerResolver {
|
||||
this.expectedToken = tokenMode ? token.getBytes(StandardCharsets.UTF_8) : null;
|
||||
this.leadTerminals = leadTerminals == null ? Map::of : leadTerminals;
|
||||
this.architectTerminals = architectTerminals == null ? Map::of : architectTerminals;
|
||||
this.memberSlotRoles = memberSlotRoles == null ? _ -> null : memberSlotRoles;
|
||||
this.memberSlotNames = memberSlotNames == null ? Function.identity() : memberSlotNames;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -206,11 +217,12 @@ public final class CallerResolver {
|
||||
return Principal.leader(lead, c.terminal(), c.pid());
|
||||
}
|
||||
String slot = architectTerminals.get().get(c.terminal());
|
||||
if (slot != null) {
|
||||
if (slot != null && memberSlotRoles.apply(slot) == MemberRole.ARCHITECT) {
|
||||
// The config/live binding names this pane as an architect slot's own. Same
|
||||
// unforgeable pane mapping; the live binding, never a request argument, decides.
|
||||
// Checked before the generic worker fallback, per the CB-548 precedence order.
|
||||
return Principal.architect(slot, c.terminal(), c.pid());
|
||||
// Check the slot role too: this defence in depth prevents a bad lifecycle bind from
|
||||
// escalating a dev or reviewer into an architect. Checked before the worker fallback.
|
||||
return Principal.architect(memberSlotNames.apply(slot), c.terminal(), c.pid());
|
||||
}
|
||||
return Principal.worker(c.terminal(), c.pid()); // unforgeable; never token-gated
|
||||
}
|
||||
|
||||
@@ -0,0 +1,21 @@
|
||||
package dev.ltms.bridged.auth;
|
||||
|
||||
import dev.ltms.bridged.peer.MemberRole;
|
||||
|
||||
/** Optional session lifecycle hook for live member-slot bindings. */
|
||||
public interface MemberLifecycle {
|
||||
|
||||
MemberLifecycle NONE = new MemberLifecycle() {
|
||||
@Override
|
||||
public void acquired(MemberRole role, String profile, String terminal) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void released(String terminal) {
|
||||
}
|
||||
};
|
||||
|
||||
void acquired(MemberRole role, String profile, String terminal);
|
||||
|
||||
void released(String terminal);
|
||||
}
|
||||
@@ -2,11 +2,14 @@ package dev.ltms.bridged.auth;
|
||||
|
||||
import dev.ltms.bridged.config.BridgedConfig;
|
||||
import dev.ltms.bridged.peer.MemberRole;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.Collections;
|
||||
import java.util.HashMap;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
|
||||
/**
|
||||
* The architect-slot registry (CB-548): every gateway-local architect name and the strong-model
|
||||
@@ -29,7 +32,9 @@ import java.util.Map;
|
||||
* exposes the map the resolver resolves against plus the profile lookup lifecycle will call.
|
||||
* Nothing here creates or manages an architect session.
|
||||
*/
|
||||
public final class MemberRegistry {
|
||||
public final class MemberRegistry implements MemberLifecycle {
|
||||
|
||||
private static final Logger log = LoggerFactory.getLogger(MemberRegistry.class);
|
||||
|
||||
/**
|
||||
* One flattened {@code fleet:} entry.
|
||||
@@ -125,6 +130,12 @@ public final class MemberRegistry {
|
||||
return e == null ? null : e.role();
|
||||
}
|
||||
|
||||
/** The unqualified configured name for a slot, or {@code null} if it is unknown. */
|
||||
public String nameForSlot(String slotName) {
|
||||
Entry e = slots.get(slotName);
|
||||
return e == null ? null : e.name();
|
||||
}
|
||||
|
||||
/** True when {@code slotName} is a configured architect slot. */
|
||||
public boolean isSlot(String slotName) {
|
||||
return slots.containsKey(slotName);
|
||||
@@ -189,4 +200,33 @@ public final class MemberRegistry {
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Bind only architect sessions to a free slot with the resolved profile.
|
||||
*
|
||||
* <p>The role check is lifecycle policy. {@link CallerResolver} repeats it when resolving a
|
||||
* binding, so a later lifecycle regression cannot turn a worker into an architect.
|
||||
*/
|
||||
@Override
|
||||
public void acquired(MemberRole role, String profile, String terminal) {
|
||||
if (role != MemberRole.ARCHITECT || terminal == null || terminal.isBlank()) {
|
||||
return;
|
||||
}
|
||||
// slotsFor preserves definition order, so duplicate-profile slots use the first free one.
|
||||
for (Entry entry : slotsFor(MemberRole.ARCHITECT).values()) {
|
||||
if (Objects.equals(profile, entry.profile()) && bind(entry.key(), terminal)) {
|
||||
return;
|
||||
}
|
||||
}
|
||||
log.info("member slot: no free architect slot for profile={}; session remains a worker", profile);
|
||||
}
|
||||
|
||||
/** Unbind a released terminal using the compare-safe registry operation. */
|
||||
@Override
|
||||
public void released(String terminal) {
|
||||
String slot = slotForTerminal(terminal);
|
||||
if (slot != null) {
|
||||
unbind(slot, terminal);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -48,7 +48,7 @@ public record Principal(Role role, String terminal, long pid, String name) {
|
||||
* made a lead unaddressable: {@link #ownsSession} could never be true for it, so
|
||||
* {@code bridge_reply} was refused and one lead could send to another but never be answered.
|
||||
* The terminal now means "which pane is this caller", the presence map keys on
|
||||
* {@link #isWorker()} instead, and a lead is a peer that can both send and receive.
|
||||
* {@link #isSpawnedMember()} instead, and a lead is a peer that can both send and receive.
|
||||
*/
|
||||
public static Principal leader(String name, String terminal, long pid) {
|
||||
return new Principal(Role.PRIMARY, terminal, pid, name);
|
||||
@@ -85,6 +85,16 @@ public record Principal(Role role, String terminal, long pid, String name) {
|
||||
return role == Role.WORKER;
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether this caller is a spawned member with its own pane.
|
||||
*
|
||||
* <p>Both workers and architects are spawned members. A lead is excluded because recording it
|
||||
* as present would count it as an available member in the roster.
|
||||
*/
|
||||
public boolean isSpawnedMember() {
|
||||
return role == Role.WORKER || role == Role.ARCHITECT;
|
||||
}
|
||||
|
||||
public boolean isAnonymous() {
|
||||
return role == Role.ANONYMOUS;
|
||||
}
|
||||
|
||||
@@ -210,8 +210,14 @@ public record BridgedConfig(
|
||||
// worker the primary's IDE servers, which are bound to the primary's checkout — so its
|
||||
// navigation returned paths outside its own worktree. GitWorktrees now neutralizes that
|
||||
// file instead; a worker's tools are whatever its launcher mounts.
|
||||
// .claude/settings.local.json is the sibling that was left behind: it pre-approves tools
|
||||
// (mcp__context7__*, mcp__jetbrains, mcp__intellij-index, Workflow(code-review)) a worker
|
||||
// must never hold, and enables MCP servers by name. Its grants are currently INERT because
|
||||
// GitWorktrees.isolateToolSurface strips every worktree's server map to empty — the named
|
||||
// servers do not exist there to be enabled. This is defence in depth, not a live fix: the
|
||||
// worker stays isolated only because this separate mechanism already removes the servers.
|
||||
parityOverlay = (parityOverlay == null || parityOverlay.isEmpty())
|
||||
? List.of(".claude/settings.local.json", ".env", ".envrc")
|
||||
? List.of(".env", ".envrc")
|
||||
: List.copyOf(parityOverlay);
|
||||
// gitTokenEnv stays null when unset (opt-in). gitHostEnv defaults so operators enabling
|
||||
// checkpoints need only set gitTokenEnv; it is injected only alongside a resolved token.
|
||||
|
||||
@@ -55,6 +55,9 @@ public final class CompletionResolver implements TurnListener {
|
||||
/** Cap the scraped tail so a long transcript can't return an unbounded blob. */
|
||||
static final int MAX_SCRAPE_CHARS = 4000;
|
||||
|
||||
private static final String CLIPPED_PANE_TAIL_MARKER =
|
||||
"[Pane tail clipped: member did not call bridge_reply.]";
|
||||
|
||||
private final AgentControl agents;
|
||||
private final Rendezvous rendezvous;
|
||||
|
||||
@@ -147,9 +150,14 @@ public final class CompletionResolver implements TurnListener {
|
||||
return;
|
||||
}
|
||||
String tail;
|
||||
int originalLength = 0;
|
||||
boolean clipped = false;
|
||||
boolean scrapeFailed = false;
|
||||
try {
|
||||
tail = clip(lastAssistantBlock(agents.read(target, SCRAPE_SOURCE)));
|
||||
String assistantBlock = lastAssistantBlock(agents.read(target, SCRAPE_SOURCE));
|
||||
originalLength = assistantBlock.strip().length();
|
||||
clipped = originalLength > MAX_SCRAPE_CHARS;
|
||||
tail = clip(assistantBlock);
|
||||
} catch (RuntimeException e) {
|
||||
// The worker finished but we couldn't read its screen — still resolve the send so the
|
||||
// caller unblocks; an empty tail beats hanging until the caller's timeout.
|
||||
@@ -169,8 +177,14 @@ public final class CompletionResolver implements TurnListener {
|
||||
target);
|
||||
return; // keep the in-flight record: a later genuine completion still needs it
|
||||
}
|
||||
if (rendezvous.resolveCompletion(waiter, tail)) {
|
||||
String completion = clipped ? tail + "\n" + CLIPPED_PANE_TAIL_MARKER : tail;
|
||||
if (rendezvous.resolveCompletion(waiter, completion)) {
|
||||
inFlight.remove(target, turn);
|
||||
if (clipped) {
|
||||
log.warn("completion scrape for {} clipped from {} chars to the {} char cap; "
|
||||
+ "member did not call bridge_reply, so the pane tail is partial",
|
||||
target, originalLength, MAX_SCRAPE_CHARS);
|
||||
}
|
||||
log.debug("resolved send to {} via turn-completion fallback ({} chars scraped)",
|
||||
target, tail.length());
|
||||
}
|
||||
|
||||
@@ -78,6 +78,15 @@ public final class Injector {
|
||||
*/
|
||||
private static final int READINESS_GRACE_POLLS = 240;
|
||||
|
||||
/**
|
||||
* The single source for the injector poll cadence — how often the {@link StatusPoller} drives
|
||||
* {@link #onStatus} at. {@code Bridged} passes this to every {@link StatusPoller} it constructs,
|
||||
* and this class reads it to state the readiness grace in seconds on the CB-562 expiry log
|
||||
* instead of hardcoding "60s". One constant, so a cadence change cannot silently desync a log
|
||||
* that claims a grace duration.
|
||||
*/
|
||||
public static final long POLL_INTERVAL_MILLIS = 250;
|
||||
|
||||
private final AgentControl agents;
|
||||
private final TurnListener turnListener;
|
||||
private final Predicate<String> ready; // CB-113: a target is deliverable only when available
|
||||
@@ -264,6 +273,11 @@ public final class Injector {
|
||||
// fail every queued message and release the target (CB-114) instead of
|
||||
// polling it indefinitely with the caller's future never completing.
|
||||
notReady = new ArrayList<>(t.queue);
|
||||
log.warn("readiness grace for {} expired after {} polls ({}s): target never "
|
||||
+ "became deliverable, so failing {} queued message(s) that never "
|
||||
+ "reached its pane",
|
||||
target, READINESS_GRACE_POLLS,
|
||||
READINESS_GRACE_POLLS * POLL_INTERVAL_MILLIS / 1000, notReady.size());
|
||||
t.queue.clear();
|
||||
t.notReadySincePoll = 0;
|
||||
}
|
||||
|
||||
@@ -109,10 +109,10 @@ public final class BridgeMcp {
|
||||
? callers.resolve(req.getRemoteAddr(), req.getRemotePort(),
|
||||
req.getHeader("Authorization"))
|
||||
: legacyPrincipal(identity, req.getRemoteAddr(), req.getRemotePort());
|
||||
// CB-532: guard on the ROLE, not on the terminal being null. A named lead now
|
||||
// carries its pane too, and enrolling a lead in the worker presence map would
|
||||
// have it counted as an available worker.
|
||||
if (p.isWorker()) presence.markPresent(p.terminal());
|
||||
// CB-532: guard on the ROLE, not on the terminal being null. This excludes a
|
||||
// lead, which carries its pane too, while including every spawned member role.
|
||||
// Enrolling a lead would count it as an available member in the roster.
|
||||
markSpawnedMemberPresent(p, presence);
|
||||
return McpTransportContext.create(Map.of(
|
||||
CALLER_TERMINAL, orEmpty(p.terminal()),
|
||||
CALLER_PID, Long.toString(p.pid()),
|
||||
@@ -365,6 +365,13 @@ public final class BridgeMcp {
|
||||
return transport;
|
||||
}
|
||||
|
||||
/** Mark a connected spawned member available for the injector readiness gate. */
|
||||
static void markSpawnedMemberPresent(Principal caller, MemberPresence presence) {
|
||||
if (caller.isSpawnedMember()) {
|
||||
presence.markPresent(caller.terminal());
|
||||
}
|
||||
}
|
||||
|
||||
/** Graceful shutdown of the MCP server. */
|
||||
public void close() {
|
||||
server.closeGracefully();
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
package dev.ltms.bridged.session;
|
||||
|
||||
import dev.ltms.bridged.auth.MemberLifecycle;
|
||||
import dev.ltms.bridged.herdr.Agent;
|
||||
import dev.ltms.bridged.inject.TurnListener;
|
||||
import dev.ltms.bridged.inject.MemberPresence;
|
||||
@@ -49,6 +50,7 @@ public final class SessionManager implements TurnListener {
|
||||
private final LongSupplier nowNanos;
|
||||
private final int contextCap;
|
||||
private final boolean clearAfterTurn;
|
||||
private volatile MemberLifecycle memberLifecycle = MemberLifecycle.NONE;
|
||||
|
||||
/** CB-520: notified with a terminalId on every acquire; no-op until wired. */
|
||||
private final List<Consumer<String>> acquireListeners = new java.util.concurrent.CopyOnWriteArrayList<>();
|
||||
@@ -157,6 +159,7 @@ public final class SessionManager implements TurnListener {
|
||||
null,
|
||||
null);
|
||||
registry.put(handle.id(), session);
|
||||
memberLifecycle.acquired(session.role(), session.profile(), session.terminalId());
|
||||
log.debug("acquired session id={} terminal={} profile={} owner={}",
|
||||
handle.id(), handle.terminalId(), session.profile(), session.ownerTerminal());
|
||||
notifyAcquired(session.terminalId());
|
||||
@@ -187,6 +190,7 @@ public final class SessionManager implements TurnListener {
|
||||
MemberSession removed = registry.remove(paneId);
|
||||
boolean preserveWorktree = cause == ReleaseCause.SHUTDOWN;
|
||||
if (removed != null) {
|
||||
memberLifecycle.released(removed.terminalId());
|
||||
log.debug("releasing session pane={} terminal={} state={} cause={}",
|
||||
removed.paneId(), removed.terminalId(), removed.state(), cause);
|
||||
if (preserveWorktree && removed.worktree() != null) {
|
||||
@@ -257,6 +261,11 @@ public final class SessionManager implements TurnListener {
|
||||
}
|
||||
}
|
||||
|
||||
/** Inject the optional member-slot lifecycle after construction without changing constructors. */
|
||||
public void setMemberLifecycle(MemberLifecycle memberLifecycle) {
|
||||
this.memberLifecycle = memberLifecycle == null ? MemberLifecycle.NONE : memberLifecycle;
|
||||
}
|
||||
|
||||
/** A listener failure must never prevent the acquisition it is reacting to. */
|
||||
private void notifyAcquired(String terminalId) {
|
||||
if (terminalId == null) {
|
||||
@@ -330,6 +339,7 @@ public final class SessionManager implements TurnListener {
|
||||
path,
|
||||
branch);
|
||||
registry.put(handle.id(), session);
|
||||
memberLifecycle.acquired(session.role(), session.profile(), session.terminalId());
|
||||
log.debug("acquired worktree session id={} terminal={} profile={} branch={} path={}",
|
||||
handle.id(), handle.terminalId(), session.profile(), session.branch(), session.worktree());
|
||||
notifyAcquired(session.terminalId());
|
||||
|
||||
@@ -3,6 +3,8 @@ package dev.ltms.bridged.auth;
|
||||
import dev.ltms.bridged.herdr.FakeHerdr;
|
||||
import dev.ltms.bridged.herdr.PaneLocator;
|
||||
import dev.ltms.bridged.mcp.ConnectionIdentity;
|
||||
import dev.ltms.bridged.config.BridgedConfig;
|
||||
import dev.ltms.bridged.peer.MemberRole;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.Map;
|
||||
@@ -32,6 +34,17 @@ class CallerResolverTest {
|
||||
return identity(999_999);
|
||||
}
|
||||
|
||||
private static MemberRegistry boundMembers(String slot, MemberRole role) {
|
||||
Map<String, BridgedConfig.Slot> architects = role == MemberRole.ARCHITECT
|
||||
? Map.of("lead-designer", new BridgedConfig.Slot("sonnet")) : Map.of();
|
||||
Map<String, BridgedConfig.Slot> devs = role == MemberRole.DEV
|
||||
? Map.of("builder", new BridgedConfig.Slot("sonnet")) : Map.of();
|
||||
MemberRegistry members = new MemberRegistry(
|
||||
new BridgedConfig.Fleet(Map.of(), architects, devs, Map.of(), null));
|
||||
assertTrue(members.bind(slot, "term_a"));
|
||||
return members;
|
||||
}
|
||||
|
||||
@Test
|
||||
void aLoopbackWorkerPaneResolvesToWorkerRegardlessOfAuthMode() {
|
||||
Principal underTrust = new CallerResolver(workerIdentity()).resolve("127.0.0.1", 42, null);
|
||||
@@ -302,8 +315,9 @@ class CallerResolverTest {
|
||||
|
||||
@Test
|
||||
void aBoundArchitectPaneResolvesToArchitectBeforeTheWorkerFallback() {
|
||||
// This is the production construction path used by Bridged.
|
||||
Principal p = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
|
||||
Map::of, () -> Map.of("term_a", "lead-designer"))
|
||||
Map::of, boundMembers("architect:lead-designer", MemberRole.ARCHITECT))
|
||||
.resolve("127.0.0.1", 42, null);
|
||||
|
||||
assertEquals(Role.ARCHITECT, p.role(),
|
||||
@@ -316,7 +330,7 @@ class CallerResolverTest {
|
||||
@Test
|
||||
void anArchitectNeedsNoTokenEvenInTokenMode() {
|
||||
Principal p = CallerResolver.withLeadsAndMembers(workerIdentity(), true, "s3cret",
|
||||
Map::of, () -> Map.of("term_a", "lead-designer"))
|
||||
Map::of, boundMembers("architect:lead-designer", MemberRole.ARCHITECT))
|
||||
.resolve("127.0.0.1", 42, null);
|
||||
|
||||
assertEquals(Role.ARCHITECT, p.role(),
|
||||
@@ -325,9 +339,10 @@ class CallerResolverTest {
|
||||
|
||||
@Test
|
||||
void anUnboundPaneStillResolvesAsAWorker() {
|
||||
Map<String, String> arch = Map.of("term_elsewhere", "reviewer");
|
||||
MemberRegistry members = new MemberRegistry(new BridgedConfig.Fleet(Map.of(),
|
||||
Map.of("lead-designer", new BridgedConfig.Slot("sonnet")), Map.of(), Map.of(), null));
|
||||
Principal p = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
|
||||
Map::of, () -> arch).resolve("127.0.0.1", 42, null);
|
||||
Map::of, members).resolve("127.0.0.1", 42, null);
|
||||
|
||||
assertEquals(Role.WORKER, p.role());
|
||||
assertNull(p.name());
|
||||
@@ -337,7 +352,8 @@ class CallerResolverTest {
|
||||
@Test
|
||||
void aLeadWinsOverAnArchitectBindingForTheSamePane() {
|
||||
Principal p = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
|
||||
() -> Map.of("term_a", "opus-5.0"), () -> Map.of("term_a", "lead-designer"))
|
||||
() -> Map.of("term_a", "opus-5.0"),
|
||||
boundMembers("architect:lead-designer", MemberRole.ARCHITECT))
|
||||
.resolve("127.0.0.1", 42, null);
|
||||
|
||||
assertEquals(Role.PRIMARY, p.role(),
|
||||
@@ -349,26 +365,17 @@ class CallerResolverTest {
|
||||
/** The registry is live, like leads: a binding injected after construction is honoured. */
|
||||
@Test
|
||||
void anArchitectBoundAfterConstructionIsHonouredWithoutRebuildingTheResolver() {
|
||||
Map<String, String> live = new java.util.HashMap<>();
|
||||
MemberRegistry members = new MemberRegistry(new BridgedConfig.Fleet(Map.of(),
|
||||
Map.of("lead-designer", new BridgedConfig.Slot("sonnet")), Map.of(), Map.of(), null));
|
||||
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
|
||||
Map::of, () -> live);
|
||||
Map::of, members);
|
||||
|
||||
assertEquals(Role.WORKER, r.resolve("127.0.0.1", 42, null).role());
|
||||
|
||||
live.put("term_a", "lead-designer"); // the later lifecycle binds the slot
|
||||
assertTrue(members.bind("architect:lead-designer", "term_a")); // the later lifecycle binds the slot
|
||||
|
||||
assertEquals(Role.ARCHITECT, r.resolve("127.0.0.1", 42, null).role());
|
||||
assertEquals("lead-designer", r.members().get("term_a"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void theArchitectMapFormIsCopiedSoLaterMutationCannotGrantArchitect() {
|
||||
Map<String, String> mutable = new java.util.LinkedHashMap<>();
|
||||
CallerResolver r = new CallerResolver(workerIdentity(), false, null, Map.of(), mutable);
|
||||
|
||||
mutable.put("term_a", "sneaky");
|
||||
|
||||
assertEquals(Role.WORKER, r.resolve("127.0.0.1", 42, null).role());
|
||||
assertEquals("architect:lead-designer", r.members().get("term_a"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -381,7 +388,8 @@ class CallerResolverTest {
|
||||
@Test
|
||||
void anArchitectOwnsItsOwnPaneAndNoOther() {
|
||||
Principal arch = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
|
||||
Map::of, () -> Map.of("term_a", "lead-designer")).resolve("127.0.0.1", 42, null);
|
||||
Map::of, boundMembers("architect:lead-designer", MemberRole.ARCHITECT))
|
||||
.resolve("127.0.0.1", 42, null);
|
||||
|
||||
assertTrue(arch.ownsSession("term_a"));
|
||||
assertTrue(Authz.permits(arch, Authz.Action.REPLY, "term_a"));
|
||||
@@ -389,6 +397,15 @@ class CallerResolverTest {
|
||||
assertFalse(Authz.permits(arch, Authz.Action.REPLY, "term_b"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void aBoundNonArchitectSlotStillResolvesAsAWorker() {
|
||||
Principal p = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
|
||||
Map::of, boundMembers("dev:builder", MemberRole.DEV))
|
||||
.resolve("127.0.0.1", 42, null);
|
||||
|
||||
assertEquals(Role.WORKER, p.role(), "a dev binding must never grant architect rights");
|
||||
}
|
||||
|
||||
@Test
|
||||
void tokenModeRequiresANonEmptyConfiguredToken() {
|
||||
ConnectionIdentity id = nonWorkerIdentity();
|
||||
|
||||
@@ -1167,6 +1167,41 @@ class BridgedConfigTest {
|
||||
"a profile without the key stays off-subscription (the default)");
|
||||
}
|
||||
|
||||
@Test
|
||||
void parityOverlayDefaultsToEnvFilesOnly(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("no-overlay.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
profiles:
|
||||
gx10:
|
||||
baseUrl: http://gx10.gw:8000
|
||||
""");
|
||||
|
||||
BridgedConfig cfg = BridgedConfig.load(f);
|
||||
assertEquals(List.of(".env", ".envrc"),
|
||||
cfg.profiles().get("gx10").parityOverlay(),
|
||||
"the default parity overlay is the env files; settings.local.json is no longer copied by default");
|
||||
}
|
||||
|
||||
@Test
|
||||
void parityOverlayExplicitListIsPreservedVerbatim(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("explicit-overlay.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
profiles:
|
||||
gx10:
|
||||
baseUrl: http://gx10.gw:8000
|
||||
parityOverlay: [".claude/settings.local.json", ".env"]
|
||||
""");
|
||||
|
||||
BridgedConfig cfg = BridgedConfig.load(f);
|
||||
assertEquals(List.of(".claude/settings.local.json", ".env"),
|
||||
cfg.profiles().get("gx10").parityOverlay(),
|
||||
"an operator's explicit list survives verbatim — the default only changes when unset");
|
||||
}
|
||||
|
||||
// ── CB-542: subscription:true must not smuggle an unguarded endpoint via env: ───────────────
|
||||
|
||||
@Test
|
||||
|
||||
@@ -156,6 +156,33 @@ class CompletionResolverTest {
|
||||
assertEquals("No, 391 = 17 × 23.", waiter.getNow(null).text());
|
||||
}
|
||||
|
||||
@Test
|
||||
void marksAClippedCompletionPaneTail() {
|
||||
String block = "⏺ " + "x".repeat(CompletionResolver.MAX_SCRAPE_CHARS + 1) + "\n❯ ";
|
||||
FakeHerdr herdr = new FakeHerdr().readText(block);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
|
||||
|
||||
assertEquals("x".repeat(CompletionResolver.MAX_SCRAPE_CHARS)
|
||||
+ "\n[Pane tail clipped: member did not call bridge_reply.]",
|
||||
waiter.getNow(null).text());
|
||||
}
|
||||
|
||||
@Test
|
||||
void leavesAnUnclippedCompletionPaneTailUnmarked() {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("⏺ complete report\n❯ ");
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
|
||||
|
||||
assertEquals("complete report", waiter.getNow(null).text());
|
||||
}
|
||||
|
||||
@Test
|
||||
void resolvesSynchronouslyBeforePostTurnContextClearing() {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("⏺ previous answer\n❯ ");
|
||||
@@ -177,7 +204,8 @@ class CompletionResolverTest {
|
||||
// block while resolve compares against a clip()'d tail. For a block longer than MAX_SCRAPE_CHARS
|
||||
// the two capped representations differ even when the pane never changed, so the CB-115
|
||||
// byte-identical guard failed to fire and a stale completion could resolve the send. Both sides
|
||||
// must clip identically; here an unchanged >cap block on rapid back-to-back turns stays suppressed.
|
||||
// must clip identically. The returned-text marker is added only after this comparison, so an
|
||||
// unchanged >cap block on rapid back-to-back turns still stays suppressed.
|
||||
String longBlock = "⏺ " + "x".repeat(CompletionResolver.MAX_SCRAPE_CHARS + 500) + "\n❯ ";
|
||||
FakeHerdr herdr = new FakeHerdr().readText(longBlock);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
|
||||
@@ -1,10 +1,15 @@
|
||||
package dev.ltms.bridged.inject;
|
||||
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.LoggerContext;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import dev.ltms.bridged.herdr.AgentControl;
|
||||
import dev.ltms.bridged.herdr.AgentStatus;
|
||||
import dev.ltms.bridged.herdr.FakeHerdr;
|
||||
import dev.ltms.bridged.herdr.HerdrException;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
@@ -369,6 +374,40 @@ class InjectorTest {
|
||||
assertTrue(inj.activeTargets().isEmpty(), "the target is reclaimed, not polled forever");
|
||||
}
|
||||
|
||||
@Test
|
||||
void readinessGraceExpiryIsLogged() {
|
||||
// CB-562: the grace-expiry path used to clear the queue silently, so a message that never
|
||||
// reached the worker's pane surfaced elsewhere as an unrelated turn-stall failure. Assert the
|
||||
// expiry now names the real cause. (ListAppender capture pattern mirrors AuditLogTest.)
|
||||
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
|
||||
ch.qos.logback.classic.Logger injectorLog =
|
||||
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(Injector.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.setContext(ctx);
|
||||
appender.start();
|
||||
injectorLog.addAppender(appender);
|
||||
injectorLog.setLevel(Level.WARN);
|
||||
try {
|
||||
Injector inj = new Injector(new AgentControl(herdr), TurnListener.NOOP, _ -> false, _ -> {
|
||||
});
|
||||
inj.enqueue(T, "task");
|
||||
|
||||
for (int i = 0; i < READINESS_SAMPLES; i++) inj.onStatus(T, AgentStatus.IDLE);
|
||||
|
||||
String warn = appender.list.stream()
|
||||
.filter(e -> e.getLevel().equals(Level.WARN))
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.findFirst()
|
||||
.orElse("no grace-expiry WARN logged");
|
||||
assertTrue(warn.contains(T), "the log names the target terminal: " + warn);
|
||||
assertTrue(warn.contains("never reached"), "the log names the real cause: " + warn);
|
||||
assertTrue(warn.contains("1 queued message"),
|
||||
"the log carries the failed message count: " + warn);
|
||||
} finally {
|
||||
injectorLog.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void aWorkerThatBecomesReadyWithinTheGraceIsDeliveredNormally() {
|
||||
// The readiness grace must not fail a worker that is merely slow to boot: once it becomes
|
||||
|
||||
@@ -2,6 +2,7 @@ package dev.ltms.bridged.mcp;
|
||||
|
||||
import dev.ltms.bridged.auth.Authz;
|
||||
import dev.ltms.bridged.auth.CallerResolver;
|
||||
import dev.ltms.bridged.auth.MemberRegistry;
|
||||
import dev.ltms.bridged.auth.Principal;
|
||||
import dev.ltms.bridged.auth.Role;
|
||||
import dev.ltms.bridged.config.BridgedConfig;
|
||||
@@ -69,7 +70,8 @@ class BridgeMcpAuthzTest {
|
||||
|
||||
mcp = new BridgeMcp(messages, workers, sessions, identity, sessions.asPresence(),
|
||||
new PrimaryRegistry(null),
|
||||
enforce ? new CallerResolver(identity) : null,
|
||||
enforce ? CallerResolver.withLeadsAndMembers(identity, false, null,
|
||||
Map::of, new MemberRegistry(null)) : null,
|
||||
metrics);
|
||||
return mcp;
|
||||
}
|
||||
|
||||
@@ -12,6 +12,7 @@ import dev.ltms.bridged.msg.Rendezvous;
|
||||
import dev.ltms.bridged.session.FakeWorktrees;
|
||||
import dev.ltms.bridged.session.SessionManager;
|
||||
import dev.ltms.bridged.peer.MemberRole;
|
||||
import dev.ltms.bridged.inject.MemberPresence;
|
||||
import dev.ltms.bridged.session.MemberSession;
|
||||
import dev.ltms.bridged.session.WorktreeRequest;
|
||||
import dev.ltms.bridged.member.ClaudeCodeLauncher;
|
||||
@@ -404,6 +405,30 @@ class BridgeMcpTest {
|
||||
assertDoesNotThrow(() -> messages.ackReply("term_a", msgId));
|
||||
}
|
||||
|
||||
@Test
|
||||
void spawnedMembersAreMarkedPresent() {
|
||||
Principal worker = Principal.worker("term_worker", 200);
|
||||
Principal architect = Principal.architect("lead-designer", "term_design", 400);
|
||||
MemberPresence presence = new MemberPresence();
|
||||
|
||||
BridgeMcp.markSpawnedMemberPresent(worker, presence);
|
||||
BridgeMcp.markSpawnedMemberPresent(architect, presence);
|
||||
|
||||
assertTrue(presence.isPresent("term_worker"));
|
||||
assertTrue(presence.isPresent("term_design"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void nonMembersAreNotMarkedPresent() {
|
||||
Principal lead = Principal.leader("opus", "term_lead", 100);
|
||||
MemberPresence presence = new MemberPresence();
|
||||
|
||||
BridgeMcp.markSpawnedMemberPresent(lead, presence);
|
||||
BridgeMcp.markSpawnedMemberPresent(Principal.anonymous(), presence);
|
||||
|
||||
assertFalse(presence.isPresent("term_lead"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void statusReportsLiveAgentStatus() {
|
||||
FakeHerdr blocked = new FakeHerdr().agentStatus("blocked");
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package dev.ltms.bridged.rest;
|
||||
|
||||
import dev.ltms.bridged.auth.CallerResolver;
|
||||
import dev.ltms.bridged.auth.MemberRegistry;
|
||||
import dev.ltms.bridged.config.BridgedConfig;
|
||||
import dev.ltms.bridged.guard.SubscriptionGuard;
|
||||
import dev.ltms.bridged.herdr.AgentControl;
|
||||
@@ -65,9 +66,8 @@ class BridgedAppAuthTest {
|
||||
MessageService messages = new MessageService(agents, injector, new Rendezvous());
|
||||
|
||||
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> pid);
|
||||
CallerResolver callers = tokenMode
|
||||
? new CallerResolver(identity, true, token)
|
||||
: new CallerResolver(identity);
|
||||
CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, tokenMode, token,
|
||||
Map::of, new MemberRegistry(null));
|
||||
metrics = BridgedMetrics.create(sessions, new dev.ltms.bridged.msg.InMemoryReplyInbox());
|
||||
|
||||
app = new BridgedApp(herdr, workers, sessions, messages, sessions.asPresence(), null,
|
||||
|
||||
@@ -1,11 +1,13 @@
|
||||
package dev.ltms.bridged.session;
|
||||
|
||||
import dev.ltms.bridged.auth.MemberRegistry;
|
||||
import dev.ltms.bridged.config.BridgedConfig;
|
||||
import dev.ltms.bridged.guard.SubscriptionGuard;
|
||||
import dev.ltms.bridged.herdr.AgentControl;
|
||||
import dev.ltms.bridged.herdr.FakeHerdr;
|
||||
import dev.ltms.bridged.herdr.WorkspaceControl;
|
||||
import dev.ltms.bridged.member.ClaudeCodeLauncher;
|
||||
import dev.ltms.bridged.peer.MemberRole;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.List;
|
||||
@@ -22,6 +24,13 @@ import static org.junit.jupiter.api.Assertions.*;
|
||||
*/
|
||||
class WorktreeSessionManagerTest {
|
||||
|
||||
private static MemberRegistry members() {
|
||||
return new MemberRegistry(new BridgedConfig.Fleet(Map.of(),
|
||||
Map.of("architect", new BridgedConfig.Slot("ltms-local")),
|
||||
Map.of("dev", new BridgedConfig.Slot("ltms-local")),
|
||||
Map.of("reviewer", new BridgedConfig.Slot("ltms-local")), null));
|
||||
}
|
||||
|
||||
private static ClaudeCodeLauncher workerService(FakeHerdr herdr) {
|
||||
BridgedConfig.Profile cfg = new BridgedConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "BRIDGED_WORKER_TOKEN",
|
||||
@@ -56,6 +65,36 @@ class WorktreeSessionManagerTest {
|
||||
assertEquals("/caller/proj", startCwd(herdr), "spawn receives the caller's cwd");
|
||||
}
|
||||
|
||||
@Test
|
||||
void onlyArchitectsBindAndReleaseMakesTheirSlotReusable() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
MemberRegistry members = members();
|
||||
SessionManager sessions = new SessionManager(workerService(herdr), new FakeWorktrees());
|
||||
sessions.setMemberLifecycle(members);
|
||||
|
||||
MemberSession architect = sessions.acquire("ltms-local", MemberRole.ARCHITECT,
|
||||
null, "/caller/proj", null, null);
|
||||
MemberSession dev = sessions.acquire("ltms-local", MemberRole.DEV,
|
||||
null, "/caller/proj", null, null);
|
||||
MemberSession reviewer = sessions.acquire("ltms-local", MemberRole.REVIEWER,
|
||||
null, "/caller/proj", null, null);
|
||||
|
||||
assertEquals("architect:architect", members.slotForTerminal(architect.terminalId()));
|
||||
assertNull(members.slotForTerminal(dev.terminalId()), "a dev must never receive architect rights");
|
||||
assertNull(members.slotForTerminal(reviewer.terminalId()),
|
||||
"a reviewer must never receive architect rights");
|
||||
|
||||
sessions.release(architect.paneId());
|
||||
MemberSession replacement = sessions.acquire("ltms-local", MemberRole.ARCHITECT,
|
||||
null, "/caller/proj", null, null);
|
||||
assertEquals("architect:architect", members.slotForTerminal(replacement.terminalId()));
|
||||
|
||||
MemberSession overflow = sessions.acquire("ltms-local", MemberRole.ARCHITECT,
|
||||
null, "/caller/proj", null, null);
|
||||
assertNull(members.slotForTerminal(overflow.terminalId()),
|
||||
"a full slot pool must not stop the architect spawn");
|
||||
}
|
||||
|
||||
@Test
|
||||
void worktreeAcquireProvisionsAndRecordsPathAndBranch() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
@@ -79,6 +118,19 @@ class WorktreeSessionManagerTest {
|
||||
assertEquals(expectedPath, s.cwd(), "session cwd is the worktree path");
|
||||
}
|
||||
|
||||
@Test
|
||||
void worktreeArchitectAcquireAlsoBindsItsSlot() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
MemberRegistry members = members();
|
||||
SessionManager sessions = new SessionManager(workerService(herdr), new FakeWorktrees());
|
||||
sessions.setMemberLifecycle(members);
|
||||
|
||||
MemberSession architect = sessions.acquire("ltms-local", MemberRole.ARCHITECT,
|
||||
null, "/caller/proj", null, new WorktreeRequest("cb-548", null));
|
||||
|
||||
assertEquals("architect:architect", members.slotForTerminal(architect.terminalId()));
|
||||
}
|
||||
|
||||
@Test
|
||||
void worktreeAcquireRunsParityOverlayWithProfileDefaults() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
@@ -94,12 +146,12 @@ class WorktreeSessionManagerTest {
|
||||
FakeWorktrees.OverlayCall overlay = worktrees.lastOverlay();
|
||||
assertNotNull(overlay);
|
||||
assertEquals("/repo", overlay.repoRoot());
|
||||
assertEquals(List.of(".claude/settings.local.json", ".env", ".envrc"),
|
||||
assertEquals(List.of(".env", ".envrc"),
|
||||
overlay.requested(), "default parity overlay is used when unset");
|
||||
assertFalse(overlay.requested().contains(".mcp.json"),
|
||||
"CB-525: replicating the primary's MCP config gives a worker the primary's IDE "
|
||||
+ "servers, which navigate its edits out of its own worktree");
|
||||
assertEquals(List.of(".claude/settings.local.json", ".envrc"), overlay.copied(),
|
||||
assertEquals(List.of(".envrc"), overlay.copied(),
|
||||
"existing paths are copied; missing paths are skipped");
|
||||
assertEquals(List.of(".envrc"), overlay.skipWorktree(),
|
||||
"tracked copied paths are --skip-worktree'd");
|
||||
|
||||
Reference in New Issue
Block a user