Compare commits

...

8 Commits

Author SHA1 Message Date
Dai Ha 8e5394f63f fleetd #705 option 1: narrow the unconfigured-pane floor to OBSERVER
CI / shell-tests (pull_request) Failing after 6s
CI / contract (pull_request) Successful in 58s
CI / build (pull_request) Failing after 1m57s
Adds Role.OBSERVER as the bottom rung CallerResolver falls to when a
herdr pane matches no live roster entry, lead, architect slot, or
collaborator tab. An observer may only READ/METRICS and REPLY/ASK on
its own pane. Widens the presence gate so an observer's MCP contact
still marks it deliverable, matching what already happens for a
worker or architect, so a pane that outlives a daemon restart is not
left permanently undeliverable.

Ships as defence in depth alongside the already-merged ticket-owner
check (#712/#716), which closed the reachable exploit this ticket
reported.
2026-10-04 19:46:54 +02:00
Dai Ha a332dfdb2c Merge remote-tracking branch 'origin/worker/726-ea34a0-2'
CI / shell-tests (push) Failing after 7s
CI / contract (push) Successful in 54s
CI / build (push) Failing after 1m53s
2026-10-04 19:09:03 +02:00
Dai Ha d0f4ae057b fleetd #726 unit 3 fix: release the single-flight claim when continuationRunner rejects the hand-off
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 1m1s
CI / build (pull_request) Failing after 1m51s
confirm() takes the per-lead-terminal claim before handing the roll to
continuationRunner, and the only release path was runRollover's own
finally. If continuationRunner.accept itself throws, runRollover never
starts, so that finally never runs, and nothing else ever writes
rollingByTerminal — the claim is held forever and the terminal can never
be rolled again. This differs from fleetd #615, which covers a throw
INSIDE the continuation (runRollover already catches that and still
releases the claim) — this is a throw from the hand-off itself, which
is not reachable with today's virtual-thread runner but would be with
a bounded executor's RejectedExecutionException.

confirm() now catches that throw, releases the claim, and overwrites
the IN_PROGRESS outcome with a terminal FAILED one, matching how a
throw inside the continuation is already surfaced.
2026-10-04 19:06:57 +02:00
Dai Ha 7b3beaa209 Merge remote-tracking branch 'origin/worker/726-10cbf0-1'
CI / shell-tests (push) Failing after 10s
CI / contract (push) Successful in 49s
CI / build (push) Failing after 1m53s
2026-10-04 19:06:47 +02:00
Dai Ha b3b2bf3da6 fleetd #726 unit 1 review fixes: correct relaunch's javadoc and dedupe its resolve logic
CI / shell-tests (pull_request) Failing after 10s
CI / contract (pull_request) Successful in 57s
CI / build (pull_request) Failing after 1m49s
relaunch's javadoc said it reads the live config; it actually reads the
FleetConfig snapshot this launcher was constructed with (fleet.leaders
is the frozen half), so say that and note a profile/tab edit needs a
daemon restart.

relaunch's recognise-only refusal reused ensureLeads()'s log wording,
which claims the lead 'is not live' — true in ensureLeads()'s context
(reached only after a short live count), false in relaunch's (which
never counts liveness, by design). Dropped that clause.

Pulled the declared/creatable/profile-configured resolution shared by
ensureLeads() and relaunch() into one private resolveLaunchable(name)
helper (returns a new ResolvedLead(lead, profile) record, or null
having logged), so the three refusals and their wording live in one
place instead of two copies that can drift. Behaviour-preserving:
ensureLeads() keeps its own liveness-count logic around the shared
resolve, and the existing 37 LeadLauncherTest cases are unchanged and
still pass.
2026-10-04 18:59:56 +02:00
Dai Ha 8bb2aa0be4 Merge remote-tracking branch 'origin/worker/729-5961c6-3'
CI / shell-tests (push) Failing after 8s
CI / contract (push) Successful in 57s
CI / build (push) Failing after 2m0s
2026-10-04 18:59:54 +02:00
Dai Ha 4ce3149bfd fleetd #726 unit 3: make LeadRollover.confirm() single-flight per lead terminal
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 55s
CI / build (pull_request) Failing after 1m48s
Two open() calls for the same lead terminal minted two tokens that both
passed confirm()'s ownership check, so both could reach the deferred
continuation and roll the same lead twice. confirm() now claims a
per-lead-terminal slot (an atomic put-if-absent) once every other gate has
passed, refusing a concurrent confirm with the new ROLL_ALREADY_RUNNING
reason; runRollover releases the claim in a finally, on both the success
and the thrown-exception path.
2026-10-04 18:53:53 +02:00
Dai Ha 38544d467c fleetd #726 unit 1: give LeadLauncher a public single-lead relaunch seam
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 45s
CI / build (pull_request) Failing after 1m48s
Adds LeadLauncher.relaunch(name), which starts exactly the named lead
from the live config, outside of ensureLeads()'s instances bookkeeping.
It retries the whole launch attempt (not just the agent_name_taken/
agent_pane_busy cases ResilientAgentLaunch already retries inside one
agents.start call) up to RELAUNCH_ATTEMPTS times.

launch() now returns the started Agent (null on failure) instead of a
boolean, so relaunch() and ensureLeads() share the same primitive.
2026-10-04 18:52:38 +02:00
16 changed files with 715 additions and 99 deletions
+19 -14
View File
@@ -27,26 +27,31 @@ through its `fleet_*` tools. No session addresses a peer, a broker, or the netwo
**Every role reads this file.** A member runs in a git worktree of this same repo, so it inherits
this `CLAUDE.md` verbatim, and every rule below is role-conditional.
**Call `fleet_whoami`.** It returns `primary`, `worker`, `architect`, or `collaborator`, resolved by
the daemon from your connection — unforgeable, and the same resolution its authorization gate uses.
A worker also carries its `sessionId`, `profile`, `worktree` and `branch`; an architect carries the
slot name it was bound to; a collaborator carries its registry name and its own `sessionId`, and
**no `leader` key** — a collaborator is a named peer, not a primary. Don't infer what you can ask.
**Call `fleet_whoami`.** It returns `primary`, `worker`, `architect`, `collaborator`, or `observer`,
resolved by the daemon from your connection — unforgeable, and the same resolution its authorization
gate uses. A worker also carries its `sessionId`, `profile`, `worktree` and `branch`; an architect
carries the slot name it was bound to; a collaborator carries its registry name and its own
`sessionId`, and **no `leader` key** — a collaborator is a named peer, not a primary. An **observer**
carries only its own `sessionId`: a pane the daemon could not place as any of the above, authorized
to `READ`/`METRICS` and to `REPLY`/`ASK` on its own pane and nothing more — never `SEND`, never a
ticket. Don't infer what you can ask.
Only if that call is unavailable, fall back to these — each is one-way, so keep reading until one
fires: the reply charter in your system prompt (*"You are a spawned member in the
claude-bridge fleet"*) ⇒ **spawned member**; fleet tools prefixed `mcp__fleet__*` ⇒ **spawned
member** (the launcher fixes that mount name; a primary's mount is named by whoever wrote its
`.mcp.json`, so it varies — and a member spawned before CB-632 still says `mcp__bridge__*`); `ANTHROPIC_BASE_URL` set ⇒ **spawned member** (Claude-model members run
on a clean env, so its *absence* proves nothing). None of these separate a worker from an architect —
only `fleet_whoami` does. **And none of them fires for a collaborator at all**: every signal in the
ladder detects a *spawned* member, while a collaborator is a tab a person opened by hand, so it has
no charter, no fixed mount name and a normal environment. A collaborator that cannot call
`fleet_whoami` therefore falls to the line below and acts as a worker. That is the safe direction —
it under-privileges, and the refusals are loud — but it means a collaborator has no way to learn
what it is except by asking. **Still unsure ⇒ act as a worker**, the most restricted member role. The
two mistakes are not symmetric: a primary acting as a worker is refused by the authorization gate —
loud and self-correcting — while a member acting as the primary ends its turn with no `fleet_reply`,
on a clean env, so its *absence* proves nothing). None of these separate a worker from an architect,
or a worker from an **observer** — an observer is just as unspawned as a collaborator and carries
none of these signals either, so only `fleet_whoami` tells the two apart. **And none of them fires
for a collaborator at all**: every signal in the ladder detects a *spawned* member, while a
collaborator is a tab a person opened by hand, so it has no charter, no fixed mount name and a
normal environment. A collaborator — or an observer — that cannot call `fleet_whoami` therefore falls
to the line below and acts as a worker. That is the safe direction — it under-privileges, and the
refusals are loud — but it means a collaborator or an observer has no way to learn what it is except
by asking. **Still unsure ⇒ act as a worker**, the most restricted member role this ladder can name.
The two mistakes are not symmetric: a primary acting as a worker is refused by the authorization gate
— loud and self-correcting — while a member acting as the primary ends its turn with no `fleet_reply`,
and the sender silently receives nothing. Fail toward the recoverable error.
### Invariants — every role, no exceptions
@@ -226,9 +226,8 @@ public final class Fleetd {
* is the same way — a person's own tab, matched to a configured name, never spawned.
*
* <p>Neither a lead nor a collaborator is ever enrolled in {@link MemberPresence} — {@code
* FleetMcp} marks presence for every spawned member (worker and architect), deliberately, since
* that map doubles as the member roster's availability signal and a lead or collaborator counted
* there would show up as an available member. So without the second and third disjuncts a lead or
* FleetMcp} marks presence only for a worker, an architect, or the unconfigured-pane floor,
* never for a lead or a collaborator. So without the second and third disjuncts a lead or
* collaborator is permanently un-deliverable: every send to one sat on the gate for
* {@code READINESS_GRACE_POLLS} (~60s) and then failed having never been typed into the pane.
*
@@ -138,13 +138,14 @@ public final class Authz {
// fleet_whoami — and carries no secrets: no ticket reply, no pending question, and no
// other session's turn state. Those live under TASK_READ. METRICS is the separate
// Prometheus scrape. Both are open to every authenticated role, including a
// collaborator.
// collaborator and the unconfigured-pane floor.
case READ, METRICS -> caller.isPrimary() || caller.isWorker() || caller.isArchitect()
|| caller.isCollaborator();
|| caller.isCollaborator() || caller.isObserver();
// Ticket polling and session status, open to every role READ is open to except a
// collaborator: ticket ids are a sequential counter with no owner check, so a holder
// could walk every ticket and read another session's delegation reply.
// collaborator or an observer: ticket ids are a sequential counter with no owner
// check, so a holder could walk every ticket and read another session's delegation
// reply.
case TASK_READ -> caller.isPrimary() || caller.isWorker() || caller.isArchitect();
// fleetd #421: reading held lead-to-lead mail is the primary's alone. An architect
@@ -40,7 +40,7 @@ import java.util.function.Supplier;
* the case the previous step does not catch: a binding with no live spawned-member session.</li>
* <li>A loopback peer PID that maps to an operator-labelled collaborator tab ⇒
* {@link Role#COLLABORATOR}, carrying that collaborator's name.</li>
* <li>A loopback peer PID that maps to any other herdr pane ⇒ {@link Role#WORKER}. This is
* <li>A loopback peer PID that maps to any other herdr pane ⇒ {@link Role#OBSERVER}. This is
* unforgeable (the OS reports the PID, herdr owns the PID→pane map) and is honoured
* regardless of auth mode, so enabling auth never breaks the fleet.</li>
* <li>Otherwise, under {@code token} mode, a valid bearer token ⇒ {@link Role#PRIMARY}.</li>
@@ -323,7 +323,7 @@ public final class CallerResolver {
// of the above keeps that stronger role.
return Principal.collaborator(collaborator, c.terminal(), c.pid());
}
return Principal.worker(c.terminal(), c.pid()); // unforgeable; never token-gated
return Principal.observer(c.terminal(), c.pid()); // unforgeable; never token-gated
}
if (tokenMode) {
@@ -88,6 +88,13 @@ public record Principal(Role role, String terminal, long pid, String name) {
return new Principal(Role.COLLABORATOR, terminal, pid, name);
}
/**
* The unconfigured-pane floor: a loopback caller whose pane matched no other role.
*/
public static Principal observer(String terminal, long pid) {
return new Principal(Role.OBSERVER, terminal, pid);
}
public boolean isPrimary() {
return role == Role.PRIMARY;
}
@@ -104,11 +111,15 @@ public record Principal(Role role, String terminal, long pid, String name) {
return role == Role.WORKER;
}
public boolean isObserver() {
return role == Role.OBSERVER;
}
/**
* 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.
* <p>Both workers and architects are spawned members. A lead is not: it is a peer the
* operator started and named, never a pane this daemon spawned.
*/
public boolean isSpawnedMember() {
return role == Role.WORKER || role == Role.ARCHITECT;
@@ -140,6 +151,7 @@ public record Principal(Role role, String terminal, long pid, String name) {
case WORKER -> "worker:" + terminal;
case ARCHITECT -> "architect:" + name;
case COLLABORATOR -> "collaborator:" + name;
case OBSERVER -> "observer:" + terminal;
case PRIMARY -> name == null ? "primary" : "leader:" + name;
case ANONYMOUS -> "anonymous";
};
@@ -45,6 +45,17 @@ public enum Role {
*/
COLLABORATOR,
/**
* A loopback pane that resolved to none of the roles above: not a live spawned member, not a
* configured lead, not a bound architect slot, not a configured collaborator tab. Unforgeable
* like a worker's — derived from the connection's pane, never from a request argument, and
* honoured regardless of auth mode. May {@code READ} and {@code METRICS}, and {@code REPLY}/
* {@code ASK} only as its own pane; may not {@code SPAWN}/{@code STOP}/{@code DRAIN}/
* {@code HANDOVER}, {@code SEND}, poll a ticket ({@code TASK_READ}), or reach the coordination
* broker ({@code COORD_SEND}/{@code COORD_READ}).
*/
OBSERVER,
/** Authenticated as nothing. Authorized for nothing but {@code /healthz}. */
ANONYMOUS
}
@@ -83,6 +83,9 @@ public final class LeadLauncher {
private static final Logger log = LoggerFactory.getLogger(LeadLauncher.class);
/** Attempts {@link #relaunch(String)} makes before giving up and returning {@code null}. */
static final int RELAUNCH_ATTEMPTS = 3;
private final AgentControl agents;
private final WorkspaceControl spaces;
private final FleetConfig cfg;
@@ -201,24 +204,14 @@ public final class LeadLauncher {
log.info("lead '{}': {} live, {} wanted — nothing to start", name, running, wanted);
continue;
}
if (!lead.isCreatable()) {
// A lead with a `tab:` but no `profile:` is recognise-only by design: the operator
// opens it by hand. Say so once rather than looking like a silent failure.
log.info("lead '{}' is not live, and names no profile — it can be recognised but not "
+ "launched. Add `profile:` under fleet.leaders.{} to have fleetd start it.",
name, name);
continue;
}
FleetConfig.Profile profile = cfg.profiles().get(lead.profile());
if (profile == null) {
log.warn("lead '{}' names profile '{}', which is not configured — not launching",
name, lead.profile());
ResolvedLead resolved = resolveLaunchable(name);
if (resolved == null) {
continue;
}
for (int i = running; i < wanted; i++) {
if (launch(name, lead, profile)) {
if (launch(name, resolved.lead(), resolved.profile()) != null) {
started++;
}
}
@@ -226,6 +219,82 @@ public final class LeadLauncher {
return started;
}
/** A declared lead paired with the profile it launches on — {@link #resolveLaunchable}'s result. */
private record ResolvedLead(FleetConfig.Leader lead, FleetConfig.Profile profile) {
}
/**
* The declared {@code Leader} and its {@code Profile} for {@code name}, read from the config
* snapshot this launcher was constructed with.
*
* @return the resolved pair, or {@code null} (having logged) if {@code name} is not declared
* under {@code fleet.leaders}, that lead names no {@code profile:} (a {@code tab:}-only,
* recognise-only lead), or its {@code profile:} is not configured. Shared by
* {@link #ensureLeads()} and {@link #relaunch(String)} so the three refusals and their
* wording live in one place.
*/
private ResolvedLead resolveLaunchable(String name) {
FleetConfig.Leader lead = cfg.fleet().leaders().get(name);
if (lead == null) {
log.warn("lead '{}' is not declared under fleet.leaders — not launching", name);
return null;
}
if (!lead.isCreatable()) {
// A lead with a `tab:` but no `profile:` is recognise-only by design: the operator
// opens it by hand. Say so once rather than looking like a silent failure.
log.info("lead '{}' names no profile — it can be recognised but not launched. Add "
+ "`profile:` under fleet.leaders.{} to have fleetd start it.", name, name);
return null;
}
FleetConfig.Profile profile = cfg.profiles().get(lead.profile());
if (profile == null) {
log.warn("lead '{}' names profile '{}', which is not configured — not launching",
name, lead.profile());
return null;
}
return new ResolvedLead(lead, profile);
}
/**
* Start the named lead from the config snapshot this launcher was constructed with — not a
* live read, so a lead's {@code profile:} or {@code tab:} edited in config needs a daemon
* restart to take effect here — outside of {@link #ensureLeads()}'s {@code instances}
* bookkeeping.
*
* @return the started {@link Agent}, or {@code null} if {@code name} is not declared under
* {@code fleet.leaders}, that lead names no {@code profile:} (a {@code tab:}-only,
* recognise-only lead), its {@code profile:} is not configured, or every attempt up to
* {@link #RELAUNCH_ATTEMPTS} failed to start it. Never throws.
*
* <p>Does not count how many instances of this lead are already live. {@link #ensureLeads()}'s
* count exists to avoid starting a second orchestrator; the caller of this method has already
* decided to replace the lead and owns that decision.
*
* <p>Retries the whole launch attempt — not only the {@code agent_name_taken}/
* {@code agent_pane_busy} cases {@link ResilientAgentLaunch} already retries inside one
* {@code agents.start} call — up to {@link #RELAUNCH_ATTEMPTS} times, sleeping via the
* injected sleeper between attempts, and returns the agent from the first attempt that
* succeeds.
*/
public Agent relaunch(String name) {
ResolvedLead resolved = resolveLaunchable(name);
if (resolved == null) {
return null;
}
for (int attempt = 1; attempt <= RELAUNCH_ATTEMPTS; attempt++) {
Agent started = launch(name, resolved.lead(), resolved.profile());
if (started != null) {
return started;
}
if (attempt < RELAUNCH_ATTEMPTS) {
sleeper.run();
}
}
return null;
}
/**
* How many live leads exist per configured name, and which of that name's labelled tabs are
* <em>not</em> live: a running agent in a tab labelled with that lead's exact {@code tab}
@@ -335,7 +404,7 @@ public final class LeadLauncher {
}
/**
* Start one lead. Returns false (having logged) rather than throwing on any failure.
* Start one lead. Returns null (having logged) rather than throwing on any failure.
*
* <p>Goes through the same {@link ResilientAgentLaunch} seam every member spawn uses
* (fleetd #727): the assembled argv is refused outright if it cannot fit the pane line herdr
@@ -344,7 +413,7 @@ public final class LeadLauncher {
* relaunch, and a seed pane whose shell has not reached its prompt yet ({@code
* agent_pane_busy}) is retried rather than failing on the first miss.
*/
private boolean launch(String name, FleetConfig.Leader lead, FleetConfig.Profile profile) {
private Agent launch(String name, FleetConfig.Leader lead, FleetConfig.Profile profile) {
String label = lead.tabLabel();
String cwd = (lead.cwd() == null || lead.cwd().isBlank())
? System.getProperty("user.dir") : lead.cwd();
@@ -375,7 +444,7 @@ public final class LeadLauncher {
log.info("lead '{}' launched: profile={} tab={} pane={} terminal={} label='{}' cwd={}",
name, profile.profile(), tab.tab().tabId(), started.paneId(),
started.terminalId(), label, cwd);
return true;
return started;
} catch (RuntimeException e) {
log.warn("lead '{}' failed to launch on profile '{}': {}",
name, profile.profile(), e.getMessage());
@@ -387,7 +456,7 @@ public final class LeadLauncher {
tab.tab().tabId(), cleanup.getMessage());
}
}
return false;
return null;
}
}
@@ -164,7 +164,13 @@ public final class LeadRollover {
* The handover file's modified time is not after {@link #open}'s request timestamp, or is
* older than {@code maxDocAgeSeconds}.
*/
HANDOVER_STALE
HANDOVER_STALE,
/**
* This lead terminal already has a roll running: an earlier {@link #confirm} call claimed
* it and that roll's continuation has not released it yet. {@code detail} names the lead
* terminal and the token that holds the claim.
*/
ROLL_ALREADY_RUNNING
}
/**
@@ -305,6 +311,15 @@ public final class LeadRollover {
*/
private final Consumer<Runnable> continuationRunner;
private final Map<String, PendingRollover> pending = new ConcurrentHashMap<>();
/**
* Lead terminal → the token of the roll currently holding that terminal exclusive, for
* {@link #confirm}'s single-flight claim. {@link #confirm} claims an entry here with an
* atomic put-if-absent once every other gate has passed, refusing with {@link
* RefusalReason#ROLL_ALREADY_RUNNING} when a claim is already held; {@link #runRollover}
* releases it in a {@code finally}, on both the success and the thrown-exception path. A
* terminal absent from this map has no roll currently in flight for it.
*/
private final Map<String, String> rollingByTerminal = new ConcurrentHashMap<>();
/**
* Finished tokens → what actually happened, for {@link #status}. Bounded by {@link
* #OUTCOME_HISTORY_CAP}, oldest evicted first ({@code removeEldestEntry} on an insertion-order
@@ -487,6 +502,16 @@ public final class LeadRollover {
return docCheck;
}
// Single-flight claim: atomic put-if-absent, taken only after every other gate has
// passed, so a refused confirm() never takes it. A non-null previous value means a
// different, still-running roll already holds this lead terminal.
String holder = rollingByTerminal.putIfAbsent(p.leadTerminal(), token);
if (holder != null) {
return RollDecision.refused(RefusalReason.ROLL_ALREADY_RUNNING,
"lead terminal " + p.leadTerminal() + " already has a roll running under token "
+ holder);
}
// Record IN_PROGRESS BEFORE removing from `pending` — see RollState#IN_PROGRESS and
// OUTCOME_HISTORY_CAP's javadoc. This ordering means `token` is written into `outcomes`
// while it is STILL present in `pending`; status() checks `outcomes` first (see that
@@ -499,7 +524,23 @@ public final class LeadRollover {
pending.remove(token);
log.info("lead-rollover: confirmed token={} lead={} — roll scheduled once the calling turn ends",
token, callerTerminal);
continuationRunner.accept(() -> runRollover(p, cfg));
try {
continuationRunner.accept(() -> runRollover(p, cfg));
} catch (RuntimeException e) {
// continuationRunner can reject the hand-off itself (e.g. a bounded executor's
// RejectedExecutionException) before runRollover ever starts, so runRollover's own
// finally — the only other place that releases rollingByTerminal — never runs either.
// Release the claim here and overwrite the IN_PROGRESS entry with a terminal outcome,
// or this lead terminal could never be rolled again and status() would report
// IN_PROGRESS forever for a roll that in fact never started.
log.warn("lead-rollover: continuationRunner rejected token={} lead={}: {} — the roll "
+ "never started; releasing its claim and reporting it as FAILED",
token, callerTerminal, e.toString(), e);
rollingByTerminal.remove(p.leadTerminal(), token);
outcomes.put(token, new RollStatus(RollState.FAILED,
"continuationRunner rejected this roll before it ever started: " + e.toString()
+ " — the roll never ran; open() a fresh rollover request"));
}
return RollDecision.approved();
}
@@ -539,6 +580,12 @@ public final class LeadRollover {
"the roll's continuation threw " + e.toString() + " — the roll is dead and will "
+ "not retry itself; check the daemon log for the stack trace, then open() "
+ "a fresh rollover request"));
} finally {
// Release the single-flight claim on both the normal return and the thrown-exception
// path above — a release only on success would leave this lead terminal unrollable
// forever after one failure. The conditional two-argument remove only clears the
// entry this roll itself holds, never a different roll's claim on the same terminal.
rollingByTerminal.remove(p.leadTerminal(), p.token());
}
}
@@ -334,7 +334,7 @@ public final class FleetMcp {
* caller explicitly saying so — never by omitting a {@link CallerResolver} the way the old
* {@code callers == null} idiom allowed. {@code callers} itself is required either way: even
* under {@link #UNENFORCED}, the one real {@link CallerResolver} still resolves every caller's
* {@link Principal} (so {@code markSpawnedMemberPresent}/{@code recordPrimarySingleton} see a
* {@link Principal} (so {@code markTrackedCallerPresent}/{@code recordPrimarySingleton} see a
* real identity), and {@link #denyFor} is the only thing that changes.
*/
public enum AuthorizationMode { ENFORCED, UNENFORCED }
@@ -442,10 +442,10 @@ public final class FleetMcp {
// fall back to here — AuthorizationMode governs enforcement, not identity.
Principal p = callers.resolve(req.getRemoteAddr(), req.getRemotePort(),
req.getHeader("Authorization"));
// 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);
// Guards on the ROLE, not on the terminal being null — this marks presence for
// a worker, an architect, or the unconfigured-pane floor, and excludes a lead
// or a collaborator even though each carries its own pane too.
markTrackedCallerPresent(p, presence);
return McpTransportContext.create(Map.of(
CALLER_TERMINAL, orEmpty(p.terminal()),
CALLER_PID, Long.toString(p.pid()),
@@ -833,9 +833,14 @@ public final class FleetMcp {
return identity;
}
/** Mark a connected spawned member available for the injector readiness gate. */
static void markSpawnedMemberPresent(Principal caller, MemberPresence presence) {
if (caller.isSpawnedMember()) {
/**
* Mark a caller present for the injector readiness gate, when its deliverability depends on
* proving a live MCP contact: a worker, an architect, or the unconfigured-pane floor. A lead
* or a collaborator is excluded — each is already deliverable through its own named-registry
* entry.
*/
static void markTrackedCallerPresent(Principal caller, MemberPresence presence) {
if (caller.isSpawnedMember() || caller.isObserver()) {
presence.markPresent(caller.terminal());
}
}
@@ -1412,6 +1417,14 @@ public final class FleetMcp {
}
return text(json(m));
}
if (caller.isObserver()) {
// No architect slot, collaborator name, or lead name to report — only the pane itself,
// so a peer that already knows this terminal can still address it.
if (caller.terminal() != null) {
m.put("sessionId", caller.terminal());
}
return text(json(m));
}
if (!caller.isWorker()) {
// CB-530: which lead, once more than one pane is configured as one. `role` deliberately
// still reads "primary" — the fallback ladder in CLAUDE.md keys on it, and a lead IS a
@@ -15,6 +15,7 @@ class AuthzTest {
private static final Principal ARCH_DESIGN = Principal.architect("lead-designer", "term_design", 400);
private static final Principal ARCH_OTHER = Principal.architect("reviewer", "term_review", 500);
private static final Principal COLLABORATOR = Principal.collaborator("ops", "term_collab", 600);
private static final Principal OBSERVER = Principal.observer("term_observer", 700);
@Test
void anonymousIsAuthorizedForNothing() {
@@ -266,4 +267,49 @@ class AuthzTest {
assertTrue(WORKER_A.isSpawnedMember());
assertTrue(ARCH_DESIGN.isSpawnedMember());
}
// ── the observer matrix ─────────────────────────────────────────────────────────────────────
@Test
void anObserverMayReadAndScrapeMetrics() {
assertTrue(Authz.permits(OBSERVER, READ, null));
assertTrue(Authz.permits(OBSERVER, METRICS, null));
}
@Test
void anObserverMayReplyAndAskOnlyAsItsOwnPane() {
assertTrue(Authz.permits(OBSERVER, REPLY, "term_observer"), "its own pane is its own");
assertTrue(Authz.permits(OBSERVER, ASK, "term_observer"));
assertFalse(Authz.permits(OBSERVER, REPLY, "term_design"),
"an observer must not reply on another pane");
assertFalse(Authz.permits(OBSERVER, REPLY, null),
"an absent target must not pass the own-session rule");
}
/**
* Every action beyond READ/METRICS/REPLY/ASK, asserted denied for an observer — including
* {@code TASK_READ}, which is the entire point of this role: an unconfigured pane must not be
* able to poll a ticket or read another session's status.
*/
@Test
void anObserverIsDeniedEverythingBeyondReadMetricsReplyAndAsk() {
for (Authz.Action a : Authz.Action.values()) {
if (a == READ || a == METRICS || a == REPLY || a == ASK) {
continue;
}
assertFalse(Authz.permits(OBSERVER, a, "term_observer", target -> true),
"an observer must not " + a + " even when the classifier accepts every target");
}
}
@Test
void anObserverIsNotCountedAsAnyOtherRole() {
assertFalse(OBSERVER.isPrimary());
assertFalse(OBSERVER.isWorker());
assertFalse(OBSERVER.isArchitect());
assertFalse(OBSERVER.isCollaborator());
assertFalse(OBSERVER.isSpawnedMember());
assertTrue(OBSERVER.isObserver());
}
}
@@ -57,17 +57,23 @@ class CallerResolverTest {
return members;
}
/**
* With no roster wired up at all (the simple constructor), a loopback pane that owns a herdr
* pane but is not recognised as a live spawned member lands on the {@link Role#OBSERVER} floor
* — unforgeable and never token-gated, exactly like a worker's own identity, because it comes
* from the same connection-derived pane mapping.
*/
@Test
void aLoopbackWorkerPaneResolvesToWorkerRegardlessOfAuthMode() {
void aLoopbackPaneWithNoLiveRosterResolvesToObserverRegardlessOfAuthMode() {
Principal underTrust = new CallerResolver(workerIdentity()).resolve("127.0.0.1", 42, null);
Principal underToken = new CallerResolver(workerIdentity(), true, "s3cret")
.resolve("127.0.0.1", 42, null);
assertEquals(Role.WORKER, underTrust.role());
assertEquals(Role.OBSERVER, underTrust.role());
assertEquals("term_a", underTrust.terminal());
assertEquals(Role.WORKER, underToken.role(),
"worker identity is unforgeable and must never be token-gated — otherwise enabling "
+ "auth would lock the whole fleet out of fleet_reply");
assertEquals(Role.OBSERVER, underToken.role(),
"the floor is unforgeable and must never be token-gated — otherwise enabling auth "
+ "would lock every unconfigured pane out of even READ");
assertEquals("term_a", underToken.terminal());
}
@@ -91,20 +97,20 @@ class CallerResolverTest {
}
@Test
void otherPanesRemainWorkersWhenAPinIsSet() {
void otherPanesRemainAtTheFloorWhenAPinIsSet() {
Principal p = CallerResolver.pinnedTo(workerIdentity(), false, null, "term_someone_else")
.resolve("127.0.0.1", 42, null);
assertEquals(Role.WORKER, p.role());
assertEquals(Role.OBSERVER, p.role());
assertEquals("term_a", p.terminal());
}
/** The pin is optional config, so an absent or whitespace one must change nothing at all. */
@Test
void aBlankPinLeavesWorkerResolutionUntouched() {
assertEquals(Role.WORKER,
void aBlankPinLeavesFloorResolutionUntouched() {
assertEquals(Role.OBSERVER,
CallerResolver.pinnedTo(workerIdentity(), false, null, " ").resolve("127.0.0.1", 42, null).role());
assertEquals(Role.WORKER,
assertEquals(Role.OBSERVER,
CallerResolver.pinnedTo(workerIdentity(), false, null, null).resolve("127.0.0.1", 42, null).role());
}
@@ -225,13 +231,13 @@ class CallerResolverTest {
* mid-scan teardown into a refusal — the real match is still found and resolves as a worker.
*/
@Test
void aHerdrErrorOnANonOwningPaneStillResolvesTheRealWorker() {
void aHerdrErrorOnANonOwningPaneStillResolvesTheRealPane() {
FakeHerdr vanishedElsewhere = new FakeHerdr().processInfoFailsForPane("w2:p9", "pane_not_found");
ConnectionIdentity id = new ConnectionIdentity(new PaneLocator(vanishedElsewhere), _ -> FakeHerdr.WORKER_PID);
Principal p = new CallerResolver(id).resolve("127.0.0.1", 55555, null);
assertEquals(Role.WORKER, p.role());
assertEquals(Role.OBSERVER, p.role());
assertEquals("term_a", p.terminal());
}
@@ -275,12 +281,12 @@ class CallerResolverTest {
}
@Test
void aPaneAbsentFromTheRegistryIsStillAWorker() {
void aPaneAbsentFromTheRegistryFallsToTheObserverFloor() {
Principal p = new CallerResolver(workerIdentity(), false, null,
Map.of("term_elsewhere", "gpt-sol-5.6"))
.resolve("127.0.0.1", 42, null);
assertEquals(Role.WORKER, p.role());
assertEquals(Role.OBSERVER, p.role());
assertEquals("term_a", p.terminal());
assertNull(p.name());
}
@@ -306,16 +312,16 @@ class CallerResolverTest {
}
@Test
void anEmptyRegistryLeavesEveryPaneAWorker() {
void anEmptyRegistryLeavesEveryPaneAtTheObserverFloor() {
Map<String, String> noLeads = null;
assertEquals(Role.WORKER,
assertEquals(Role.OBSERVER,
new CallerResolver(workerIdentity(), false, null, Map.of())
.resolve("127.0.0.1", 42, null).role());
assertEquals(Role.WORKER,
assertEquals(Role.OBSERVER,
new CallerResolver(workerIdentity(), false, null, noLeads)
.resolve("127.0.0.1", 42, null).role());
// CB-531: and the same for the live-registry form, whose supplier may also be absent.
assertEquals(Role.WORKER,
// And the same for the live-registry form, whose supplier may also be absent.
assertEquals(Role.OBSERVER,
CallerResolver.withLeads(workerIdentity(), false, null, null)
.resolve("127.0.0.1", 42, null).role());
}
@@ -388,7 +394,7 @@ class CallerResolverTest {
Map<String, String> live = new java.util.HashMap<>();
CallerResolver r = CallerResolver.withLeads(workerIdentity(), false, null, () -> live);
assertEquals(Role.WORKER, r.resolve("127.0.0.1", 42, null).role());
assertEquals(Role.OBSERVER, r.resolve("127.0.0.1", 42, null).role());
live.put("term_a", "gpt-sol-5.6"); // the scanner sees a newly-labelled tab
@@ -405,7 +411,7 @@ class CallerResolverTest {
mutable.put("term_a", "sneaky");
assertEquals(Role.WORKER, r.resolve("127.0.0.1", 42, null).role());
assertEquals(Role.OBSERVER, r.resolve("127.0.0.1", 42, null).role());
}
// ── CB-548: architect slots ─────────────────────────────────────────────────────────────────
@@ -435,13 +441,13 @@ class CallerResolverTest {
}
@Test
void anUnboundPaneStillResolvesAsAWorker() {
void anUnboundPaneResolvesToTheObserverFloor() {
MemberRegistry members = new MemberRegistry(new FleetConfig.Fleet(Map.of(),
Map.of("lead-designer", new FleetConfig.Slot("sonnet")), Map.of(), Map.of(), null));
Principal p = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
Map::of, members).resolve("127.0.0.1", 42, null);
assertEquals(Role.WORKER, p.role());
assertEquals(Role.OBSERVER, p.role());
assertNull(p.name());
}
@@ -467,7 +473,7 @@ class CallerResolverTest {
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
Map::of, members);
assertEquals(Role.WORKER, r.resolve("127.0.0.1", 42, null).role());
assertEquals(Role.OBSERVER, r.resolve("127.0.0.1", 42, null).role());
assertTrue(members.bind("architect:lead-designer", "term_a")); // the later lifecycle binds the slot
@@ -494,12 +500,13 @@ class CallerResolverTest {
}
@Test
void aBoundNonArchitectSlotStillResolvesAsAWorker() {
void aBoundNonArchitectSlotResolvesToTheObserverFloorNotArchitect() {
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");
assertEquals(Role.OBSERVER, p.role(), "a dev binding must never grant architect rights, "
+ "and this construction path wires no roster to recognise it as the live dev it is");
}
@Test
@@ -510,17 +517,17 @@ class CallerResolverTest {
assertThrows(IllegalArgumentException.class, () -> new CallerResolver(id, true, " "));
}
@Test
void aWorkerOnAnyLoopbackSourceAddressIsStillAWorkerNotThePrimary() {
// fleetd #305: the escalation. ConnectionIdentity used to accept only 127.0.0.1, so a
// worker connecting from 127.0.0.2 resolved to no terminal, and this resolver's own
// (wider) loopback check then made it the PRIMARY — granting spawn, stop, send and drain.
// Measured on the Linux fleet host: binding a source of 127.0.0.2 succeeds there, so the
// path is real and not theoretical.
void aPaneOnAnyLoopbackSourceAddressIsStillAtTheFloorNotThePrimary() {
// fleetd #305: the escalation this guards against. ConnectionIdentity used to accept only
// 127.0.0.1, so a pane connecting from 127.0.0.2 resolved to no terminal, and this
// resolver's own (wider) loopback check then made it the PRIMARY — granting spawn, stop,
// send and drain. Measured on the Linux fleet host: binding a source of 127.0.0.2 succeeds
// there, so the path is real and not theoretical.
CallerResolver r = new CallerResolver(workerIdentity(), false, null);
for (String src : new String[]{"127.0.0.1", "127.0.0.2", "127.1.2.3", "::ffff:127.0.0.2"}) {
Principal p = r.resolve(src, 55555, null);
assertEquals(Role.WORKER, p.role(), "a worker must stay a worker from source " + src);
assertEquals("term_a", p.terminal(), "worker terminal from source " + src);
assertEquals(Role.OBSERVER, p.role(), "the pane must stay off PRIMARY from source " + src);
assertEquals("term_a", p.terminal(), "pane terminal from source " + src);
}
}
@@ -826,14 +833,14 @@ class CallerResolverTest {
assertEquals("term_a", p.terminal());
}
/** Regression: an empty collaborator registry leaves every pane exactly as before. */
/** Regression: an empty collaborator registry leaves every pane at the unconfigured-pane floor. */
@Test
void anEmptyCollaboratorRegistryLeavesEveryPaneAsBefore() {
void anEmptyCollaboratorRegistryLeavesEveryPaneAtTheObserverFloor() {
Principal p = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null, Map::of,
new MemberRegistry(null), t -> null, Map::of)
.resolve("127.0.0.1", 42, null);
assertEquals(Role.WORKER, p.role());
assertEquals(Role.OBSERVER, p.role());
assertNull(p.name());
}
@@ -862,6 +869,39 @@ class CallerResolverTest {
assertFalse(r.knownLeadOrCollaborator().test("term_other"));
}
// ── fleetd #705: narrowing the unconfigured-pane floor to OBSERVER ──────────────────────────
/**
* The case this ticket exists for: a pane the resolver cannot place as a live spawned member,
* a lead, a bound architect slot, or a configured collaborator must land on the narrow
* {@link Role#OBSERVER} floor, never the {@link Role#WORKER} the old fallback granted.
*
* <p>The second assertion is the control the ticket requires: a terminal the roster DOES
* recognise as a live spawned member must still resolve its own role. Without it, this test
* would also pass if the fix accidentally turned every caller into an observer.
*/
@Test
void anUnconfiguredPaneResolvesObserverButARegisteredMemberStillResolvesItsOwnRole() {
Principal unconfigured = new CallerResolver(workerIdentity()).resolve("127.0.0.1", 42, null);
assertEquals(Role.OBSERVER, unconfigured.role(),
"a pane matching none of the configured or live-roster roles must fall to the "
+ "floor, not WORKER");
assertEquals("term_a", unconfigured.terminal());
Principal registered = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
Map::of, new MemberRegistry(null),
t -> "term_a".equals(t) ? MemberRole.DEV : null, Map::of)
.resolve("127.0.0.1", 42, null);
assertEquals(Role.WORKER, registered.role(),
"control: a live spawned member must keep resolving its own role, never the "
+ "unconfigured-pane floor");
}
@Test
void describeNamesTheObserverByItsPane() {
assertEquals("observer:term_a", Principal.observer("term_a", 1).describe());
}
@Test
void knownLeadOrCollaboratorIsFalseForASpawnedMembersTerminal() {
// The exact scenario a collaborator's SEND must never reach: a live spawned member's own
@@ -295,9 +295,10 @@ class MemberRegistryLiveTest {
assertTrue(out.applied(), "the reload must actually take effect: " + out.summary());
Principal after = resolver.resolve("127.0.0.1", 42, null);
assertEquals(Role.WORKER, after.role(),
"removing the slot from config must demote the bound session to worker on its "
+ "NEXT request — this is the ticket's whole point");
assertEquals(Role.OBSERVER, after.role(),
"removing the slot from config must demote the bound session on its NEXT request — "
+ "this harness wires no live roster for term_a, so the demotion lands on "
+ "the unconfigured-pane floor");
assertEquals("term_a", after.terminal(), "same pane, same terminal — only the role changed");
}
@@ -571,4 +571,137 @@ class LeadLauncherTest {
assertTrue(warn.contains("opus"), "names the profile: " + warn);
assertTrue(warn.contains("x".repeat(60)), "names the culprit argument: " + warn);
}
// ── fleetd #726 unit 1: the single-lead relaunch seam ─────────────────────────────────────
/**
* The returned agent's {@code terminalId()}/{@code paneId()} are the ones the fake
* {@code AgentControl} actually started — not a coincidental field left over from the caller.
* {@code paneId()} echoes the exact {@code pane_id} the launch's own {@code agent.start} call
* carried (protocol 19: the agent starts into the pane it is asked to), and {@code
* terminalId()} is herdr's own generated id, which the fake always shapes as {@code
* term_new_<n>}.
*/
@Test
void relaunchReturnsTheStartedAgent() {
FakeHerdr herdr = new FakeHerdr();
dev.ltms.fleet.herdr.Agent started =
launcher(herdr, configWith(lead("opus", "lead: opus", 1))).relaunch("opus");
assertNotNull(started, "a launchable, configured lead must start");
Object startedPaneIdParam = ((Map<?, ?>) herdr.lastCall("agent.start").params()).get("pane_id");
assertEquals(startedPaneIdParam, started.paneId(),
"paneId() must be the pane the agent.start call actually targeted");
assertTrue(started.terminalId() != null && started.terminalId().startsWith("term_new_"),
"terminalId() must be herdr's own generated id: " + started.terminalId());
}
/** The new tab is labelled with the lead's configured {@code tab:}, and AFTER the start. */
@Test
void relaunchLabelsTheNewTabAfterStarting() {
FakeHerdr herdr = new FakeHerdr();
dev.ltms.fleet.herdr.Agent started =
launcher(herdr, configWith(lead("opus", "lead: opus", 1))).relaunch("opus");
assertNotNull(started);
assertEquals("lead: opus", ((Map<?, ?>) herdr.lastCall("tab.rename").params()).get("label"));
int startIndex = indexOfLastCall(herdr, "agent.start");
int renameIndex = indexOfLastCall(herdr, "tab.rename");
assertTrue(renameIndex > startIndex,
"the tab must be renamed AFTER the start succeeds, not before: start=" + startIndex
+ " rename=" + renameIndex);
}
private static int indexOfLastCall(FakeHerdr herdr, String method) {
int idx = -1;
List<FakeHerdr.Call> calls = herdr.calls;
for (int i = 0; i < calls.size(); i++) {
if (calls.get(i).method().equals(method)) {
idx = i;
}
}
return idx;
}
@Test
void relaunchOfAnUnknownLeadNameReturnsNullAndStartsNothing() {
FakeHerdr herdr = new FakeHerdr();
dev.ltms.fleet.herdr.Agent started =
launcher(herdr, configWith(lead("opus", "lead: opus", 1))).relaunch("not-declared");
assertNull(started);
assertFalse(herdr.called("agent.start"));
assertFalse(herdr.called("workspace.create"));
assertFalse(herdr.called("tab.create"));
}
@Test
void relaunchOfARecogniseOnlyLeadReturnsNullAndStartsNothing() {
FakeHerdr herdr = new FakeHerdr();
dev.ltms.fleet.herdr.Agent started =
launcher(herdr, configWith(lead(null, "lead: dead", 1))).relaunch("opus");
assertNull(started);
assertFalse(herdr.called("agent.start"));
}
@Test
void relaunchWithAnUnconfiguredProfileReturnsNullAndStartsNothing() {
FakeHerdr herdr = new FakeHerdr();
dev.ltms.fleet.herdr.Agent started =
launcher(herdr, configWith(lead("nope", "lead: opus", 1))).relaunch("opus");
assertNull(started);
assertFalse(herdr.called("agent.start"));
}
/**
* The outer retry {@link LeadLauncher#relaunch(String)} owns, separate from {@code
* ResilientAgentLaunch}'s internal {@code agent_name_taken} retry: a failed attempt must not
* be the end of the whole relaunch. Each of the first two attempts exhausts {@code
* ResilientAgentLaunch.NAME_RETRIES} name attempts (every one of them rejected), so each
* attempt's own tab is created and then closed; the third attempt's first name is free.
*/
@Test
void relaunchRetriesTheWholeAttemptAndSucceedsOnTheThird() {
FakeHerdr herdr = new FakeHerdr()
.agentNameTakenTimes(2 * ResilientAgentLaunch.NAME_RETRIES);
dev.ltms.fleet.herdr.Agent started =
fastLauncher(herdr, configWith(lead("opus", "lead: opus", 1))).relaunch("opus");
assertNotNull(started, "the third attempt's first name is free — it must succeed");
assertEquals(3, herdr.calls.stream().filter(c -> c.method().equals("tab.create")).count(),
"one tab per attempt: three attempts");
assertEquals(2, herdr.calls.stream().filter(c -> c.method().equals("tab.close")).count(),
"the two failed attempts' tabs must be closed");
}
/**
* Every attempt fails outright (a herdr error {@code ResilientAgentLaunch} does not retry at
* all) — {@link LeadLauncher#relaunch(String)} must give up after exactly {@code
* RELAUNCH_ATTEMPTS} and must not leak any of the tabs it created along the way.
*/
@Test
void relaunchGivesUpAfterExactlyRelaunchAttemptsAndLeaksNoTab() {
FakeHerdr herdr = new FakeHerdr().agentStartFailsWith("some_other_error");
dev.ltms.fleet.herdr.Agent started =
fastLauncher(herdr, configWith(lead("opus", "lead: opus", 1))).relaunch("opus");
assertNull(started, "every attempt failed — relaunch must give up, not hang or guess");
assertEquals(LeadLauncher.RELAUNCH_ATTEMPTS,
herdr.calls.stream().filter(c -> c.method().equals("agent.start")).count(),
"exactly RELAUNCH_ATTEMPTS attempts, no more, no fewer");
long tabsCreated = herdr.calls.stream().filter(c -> c.method().equals("tab.create")).count();
long tabsClosed = herdr.calls.stream().filter(c -> c.method().equals("tab.close")).count();
assertEquals(LeadLauncher.RELAUNCH_ATTEMPTS, tabsCreated);
assertEquals(tabsCreated, tabsClosed, "every tab this method created must be closed — no leaks");
}
}
@@ -1305,8 +1305,12 @@ class LeadRolloverTest {
int rolls = LeadRollover.OUTCOME_HISTORY_CAP + 50;
String[] tokens = new String[rolls];
for (int i = 0; i < rolls; i++) {
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full #" + i);
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
// a distinct terminal per iteration: confirm() single-flights per lead terminal, so
// reusing LEAD here would refuse every roll after the first instead of building up
// OUTCOME_HISTORY_CAP+50 simultaneously in-flight entries.
String terminal = "term_cap_" + i;
LeadRollover.PendingRollover pending = rollover.open(terminal, "context is full #" + i);
LeadRollover.RollDecision decision = rollover.confirm(terminal, pending.token(), true);
assertTrue(decision.accepted(), "roll #" + i + " should have been approved: " + decision.detail());
tokens[i] = pending.token();
}
@@ -1447,4 +1451,195 @@ class LeadRolloverTest {
assertTrue(status.detail().contains("HerdrException"), "the detail must name the exception "
+ "so an operator reading status() has something to act on: " + status.detail());
}
// ---- fleetd #726 unit 3: confirm() single-flights per lead terminal -----------------------
@Test
@DisplayName("[SINGLE-FLIGHT 1] a second confirm() for the SAME lead terminal is refused with "
+ "ROLL_ALREADY_RUNNING while the first roll is still in flight, continuationRunner ran "
+ "exactly once, and the claim is released once that first roll finishes")
void secondConfirmForSameTerminalIsRefusedWhileFirstRollIsStillInFlight() throws IOException {
FakeHerdr herdr = new FakeHerdr(); // default idle — the held roll WOULD complete once run
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
HoldingRunner runner = new HoldingRunner();
LeadRollover rollover = newRolloverWithHoldingRunner(herdr, cfg(handover.toString()),
fixedClock(clock), runner);
LeadRollover.PendingRollover first = rollover.open(LEAD, "context is full");
LeadRollover.RollDecision firstDecision = rollover.confirm(LEAD, first.token(), true);
assertTrue(firstDecision.accepted(), "expected approval; got: " + firstDecision.reason()
+ " / " + firstDecision.detail());
assertEquals(1, runner.heldCount(), "sanity: the first roll is held, not run yet");
LeadRollover.PendingRollover second = rollover.open(LEAD, "a second request for the same terminal");
LeadRollover.RollDecision secondDecision = rollover.confirm(LEAD, second.token(), true);
assertFalse(secondDecision.accepted());
assertEquals(LeadRollover.RefusalReason.ROLL_ALREADY_RUNNING, secondDecision.reason());
assertTrue(secondDecision.detail().contains(LEAD), "the refusal detail must name the lead "
+ "terminal: " + secondDecision.detail());
assertTrue(secondDecision.detail().contains(first.token()), "the refusal detail must name "
+ "the token holding the claim: " + secondDecision.detail());
assertEquals(1, runner.heldCount(), "the refused confirm() must never have reached "
+ "continuationRunner — it ran exactly once, for the first roll only");
runner.runNext(); // let the first (and only) held roll finish
assertEquals(2, promptCallCount(herdr), "continuationRunner ran exactly once: /clear then "
+ "bootstrapText, for the first roll only");
// The claim must have been released once the first roll's continuation finished — a
// fresh request for the SAME terminal is now approved.
LeadRollover.PendingRollover third = rollover.open(LEAD, "retry after the first roll finished");
LeadRollover.RollDecision thirdDecision = rollover.confirm(LEAD, third.token(), true);
assertTrue(thirdDecision.accepted(), "the claim must have been released once the first "
+ "roll finished: " + thirdDecision.reason() + " / " + thirdDecision.detail());
}
@Test
@DisplayName("[SINGLE-FLIGHT 2] two DIFFERENT lead terminals confirm without either being "
+ "refused, and the continuation runs once per terminal")
void differentLeadTerminalsConfirmIndependently() throws IOException {
FakeHerdr herdr = new FakeHerdr(); // default idle — both rolls complete
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
LeadRollover rollover = newRollover(herdr, cfg(handover.toString()), fixedClock(clock));
LeadRollover.PendingRollover a = rollover.open(LEAD, "context is full");
LeadRollover.PendingRollover b = rollover.open(OTHER_LEAD, "context is full too");
LeadRollover.RollDecision aDecision = rollover.confirm(LEAD, a.token(), true);
LeadRollover.RollDecision bDecision = rollover.confirm(OTHER_LEAD, b.token(), true);
assertTrue(aDecision.accepted(), "expected approval; got: " + aDecision.reason() + " / " + aDecision.detail());
assertTrue(bDecision.accepted(), "expected approval; got: " + bDecision.reason() + " / " + bDecision.detail());
assertEquals(4, promptCallCount(herdr), "both rolls ran to completion — /clear and "
+ "bootstrapText, once per terminal");
}
@Test
@DisplayName("[SINGLE-FLIGHT 3] a confirm() refused for OPERATOR_NOT_CONFIRMED does not take "
+ "the single-flight claim — a later valid confirm on the same terminal is still approved")
void refusedConfirmDoesNotTakeTheClaim() throws IOException {
FakeHerdr herdr = new FakeHerdr();
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
LeadRollover rollover = newRollover(herdr, cfg(handover.toString(), true), fixedClock(clock));
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
LeadRollover.RollDecision refused = rollover.confirm(LEAD, pending.token(), false);
assertFalse(refused.accepted());
assertEquals(LeadRollover.RefusalReason.OPERATOR_NOT_CONFIRMED, refused.reason());
// the same token stays pending after a gate refusal — retrying it with operator
// confirmation must succeed, proving the earlier refusal never took the claim.
LeadRollover.RollDecision approved = rollover.confirm(LEAD, pending.token(), true);
assertTrue(approved.accepted(), "a confirm() refused for an existing gate must never have "
+ "taken the single-flight claim: " + approved.reason() + " / " + approved.detail());
}
@Test
@DisplayName("[SINGLE-FLIGHT 4] a roll whose continuation throws a RuntimeException still "
+ "releases the claim — a fresh confirm() on that terminal is approved afterwards")
void throwingContinuationStillReleasesTheClaim() throws IOException {
FakeHerdr fake = new FakeHerdr(); // default idle — the turn-settle wait passes immediately
HerdrClient throwsOnClear = new HerdrClient() {
@Override
public JsonNode call(String method, Object params) throws HerdrException {
if ("agent.prompt".equals(method) && String.valueOf(params).contains("/clear")) {
throw new HerdrException("simulated herdr transport failure sending /clear");
}
return fake.call(method, params);
}
@Override
public void close() {
fake.close();
}
};
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
LeadRollover rollover = newRollover(throwsOnClear, cfg(handover.toString()), fixedClock(clock));
LeadRollover.PendingRollover first = rollover.open(LEAD, "context is full");
LeadRollover.RollDecision decision = rollover.confirm(LEAD, first.token(), true);
assertTrue(decision.accepted(), "every synchronous gate passes; the throw happens only "
+ "inside the deferred continuation, which this test's synchronous runner has "
+ "already run to completion by the time confirm() returns");
LeadRollover.RollStatus status = rollover.status(first.token());
assertEquals(LeadRollover.RollState.FAILED, status.state(), "sanity: the continuation threw "
+ "and left a terminal FAILED outcome: " + status.detail());
LeadRollover.PendingRollover second = rollover.open(LEAD, "retry after the failure");
LeadRollover.RollDecision retryDecision = rollover.confirm(LEAD, second.token(), true);
assertTrue(retryDecision.accepted(), "the claim must have been released even though the "
+ "continuation threw — a release only on the success path would leave this lead "
+ "terminal unrollable forever after one failure: " + retryDecision.reason() + " / "
+ retryDecision.detail());
}
@Test
@DisplayName("[SINGLE-FLIGHT 5] cancel(token) after a successful confirm(token) returns false "
+ "and does not release the claim held by that roll")
void cancelAfterConfirmDoesNotReleaseTheClaim() throws IOException {
FakeHerdr herdr = new FakeHerdr();
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
HoldingRunner runner = new HoldingRunner();
LeadRollover rollover = newRolloverWithHoldingRunner(herdr, cfg(handover.toString()),
fixedClock(clock), runner);
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertTrue(decision.accepted(), "expected approval; got: " + decision.reason() + " / " + decision.detail());
assertFalse(rollover.cancel(pending.token()), "confirm() already removed the token from "
+ "pending, so cancel() must find nothing for it");
LeadRollover.PendingRollover another = rollover.open(LEAD,
"a second request while the first roll is still held");
LeadRollover.RollDecision anotherDecision = rollover.confirm(LEAD, another.token(), true);
assertFalse(anotherDecision.accepted(), "cancel() on the already-confirmed token must not "
+ "have released the claim held by that roll");
assertEquals(LeadRollover.RefusalReason.ROLL_ALREADY_RUNNING, anotherDecision.reason());
}
@Test
@DisplayName("[SINGLE-FLIGHT 6] a continuationRunner that rejects the hand-off still releases "
+ "the claim and leaves a terminal FAILED outcome, not a stuck IN_PROGRESS")
void continuationRunnerThatRejectsTheHandOffStillReleasesTheClaim() throws IOException {
FakeHerdr herdr = new FakeHerdr();
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
AgentControl agents = new AgentControl(herdr);
// A continuationRunner that rejects every hand-off, standing in for e.g. a bounded
// executor's RejectedExecutionException — the continuation never runs, so runRollover's
// own finally never gets a chance to release the claim either.
LeadRollover rollover = new LeadRollover(agents, () -> cfg(handover.toString()), _ -> null,
fixedClock(clock), () -> { },
r -> { throw new RuntimeException("simulated continuationRunner rejection"); });
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertTrue(decision.accepted(), "every gate passed before the hand-off itself rejected — "
+ "confirm() treats that rejection the same way it treats a throw INSIDE the "
+ "continuation: logged and surfaced only through status(), never as a refusal here: "
+ decision.reason() + " / " + decision.detail());
LeadRollover.RollStatus status = rollover.status(pending.token());
assertEquals(LeadRollover.RollState.FAILED, status.state(), "a rejected hand-off must leave "
+ "a TERMINAL outcome, never a stuck IN_PROGRESS — status() would otherwise have no "
+ "way to tell a roll that never started from one still genuinely running: got "
+ status.state() + " / " + status.detail());
assertNotEquals(LeadRollover.RollState.IN_PROGRESS, status.state());
LeadRollover.PendingRollover retry = rollover.open(LEAD, "retry after the rejected hand-off");
LeadRollover.RollDecision retryDecision = rollover.confirm(LEAD, retry.token(), true);
assertTrue(retryDecision.accepted(), "the claim must have been released even though "
+ "continuationRunner itself threw before the continuation ever ran — otherwise "
+ "this lead terminal could never be rolled again: " + retryDecision.reason() + " / "
+ retryDecision.detail());
}
}
@@ -1920,8 +1920,8 @@ class FleetMcpTest {
Principal architect = Principal.architect("lead-designer", "term_design", 400);
MemberPresence presence = new MemberPresence();
FleetMcp.markSpawnedMemberPresent(worker, presence);
FleetMcp.markSpawnedMemberPresent(architect, presence);
FleetMcp.markTrackedCallerPresent(worker, presence);
FleetMcp.markTrackedCallerPresent(architect, presence);
assertTrue(presence.isPresent("term_worker"));
assertTrue(presence.isPresent("term_design"));
@@ -1932,12 +1932,27 @@ class FleetMcpTest {
Principal lead = Principal.leader("opus", "term_lead", 100);
MemberPresence presence = new MemberPresence();
FleetMcp.markSpawnedMemberPresent(lead, presence);
FleetMcp.markSpawnedMemberPresent(Principal.anonymous(), presence);
FleetMcp.markTrackedCallerPresent(lead, presence);
FleetMcp.markTrackedCallerPresent(Principal.anonymous(), presence);
assertFalse(presence.isPresent("term_lead"));
}
/**
* The item whose absence would be silent: an observer's own MCP contact must still mark
* presence, or a pane resolving to the unconfigured-pane floor would sit on the injector
* readiness gate forever once something addresses it.
*/
@Test
void anObserverContactMarksPresence() {
Principal observer = Principal.observer("term_observer", 800);
MemberPresence presence = new MemberPresence();
FleetMcp.markTrackedCallerPresent(observer, presence);
assertTrue(presence.isPresent("term_observer"));
}
@Test
void statusReportsLiveAgentStatus() {
FakeHerdr blocked = new FakeHerdr().agentStatus("blocked");
@@ -2149,6 +2164,27 @@ class FleetMcpTest {
assertTrue(leadOut.contains("\"leader\":\"opus\""), leadOut);
}
/**
* An observer reports its own role and pane, never a {@code leader} key. Without an explicit
* branch it would reach the lead branch by elimination and look right only because the
* {@code leader} key is guarded on a non-null name — this pins the branch rather than the
* accident.
*/
@Test
void whoamiReportsAnObserverNotALead() {
FakeHerdr h = new FakeHerdr();
SessionManager sessions = sessionManager(h, "http://gx00.gw:8000", Set.of("gx00.gw"));
McpSchema.CallToolResult res = FleetMcp.whoami(
Principal.observer("term_observer", 900), sessions);
assertNotEquals(Boolean.TRUE, res.isError());
String out = textOf(res);
assertTrue(out.contains("\"role\":\"observer\""), out);
assertTrue(out.contains("\"sessionId\":\"term_observer\""), out);
assertFalse(out.contains("leader"), out);
}
/**
* CB-548: an architect SEND delegates as its own pane (recording the per-target delegation) but
* must NEVER become the legacy singleton "primary" fallback — the per-target map does not cure
@@ -8,6 +8,7 @@ import dev.ltms.fleet.auth.CallerResolver;
import dev.ltms.fleet.auth.Authz;
import dev.ltms.fleet.auth.MemberRegistry;
import dev.ltms.fleet.auth.Principal;
import dev.ltms.fleet.auth.Role;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.AgentControl;
@@ -20,6 +21,7 @@ import dev.ltms.fleet.metrics.FleetMetrics;
import dev.ltms.fleet.metrics.Metrics;
import dev.ltms.fleet.msg.MessageService;
import dev.ltms.fleet.msg.Rendezvous;
import dev.ltms.fleet.peer.MemberRole;
import dev.ltms.fleet.session.FakeWorktrees;
import dev.ltms.fleet.session.SessionManager;
import dev.ltms.fleet.member.ClaudeCodeLauncher;
@@ -520,6 +522,11 @@ class FleetAppAuthTest {
* As {@link #startOnSharedService(MessageService, FakeHerdr, long)}, but {@code leadTerminals}
* resolves the given pid's terminal to a named lead (a caller with SEND permission) instead of
* a plain worker, for a test that needs a terminal-bearing caller able to create a ticket.
*
* <p>Every connecting pane not already claimed by {@code leadTerminals} is wired into the live
* roster as a spawned worker, so a caller's resolved role matches what its own test expects:
* a {@link Role#WORKER}, never the unconfigured-pane {@link Role#OBSERVER} floor a roster-less
* resolver would otherwise fall to.
*/
private Javalin startOnSharedService(MessageService messages, FakeHerdr herdr, long pid,
Map<String, String> leadTerminals) {
@@ -534,7 +541,8 @@ class FleetAppAuthTest {
SessionManager sessions = new SessionManager(workers, new FakeWorktrees());
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> pid);
CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, false, null,
() -> leadTerminals, new MemberRegistry(null));
() -> leadTerminals, new MemberRegistry(null),
t -> leadTerminals.containsKey(t) ? null : MemberRole.DEV, Map::of);
Metrics appMetrics = FleetMetrics.create(sessions, new dev.ltms.fleet.msg.InMemoryReplyInbox());
return new FleetApp(herdr, workers, sessions, messages, sessions.asPresence(), null,