Compare commits

...

16 Commits

Author SHA1 Message Date
Dai Ha e3050efe8b fleetd #705: correct the stale reason on the TASK_READ gate
CI / shell-tests (push) Failing after 8s
CI / contract (push) Successful in 50s
CI / build (push) Failing after 1m48s
The comment said ticket ids are a sequential counter with no owner check, so a
holder could walk every ticket and read another session's reply. PRs #712 and
#716 added that owner check: MessageService.ownsTicket compares a ticket's
creatorTerminal to the caller on every read.

The rule is still right, so only the reason changes. This matters now because
fleetd #737 is deciding ticket ownership across a lead handover, and a reader
who believed the old text could delete the TASK_READ restriction on the grounds
that its stated reason no longer applies.

Comment-only. mvn -o clean install: Tests run: 2083, Failures: 0, Errors: 0,
BUILD SUCCESS. Flagged by the #705 option-1 worker as out of its scope, which
was the right call.
2026-10-04 19:59:45 +02:00
Dai Ha 11998cd626 Merge remote-tracking branch 'origin/worker/705-observer-14c258-6'
CI / shell-tests (push) Failing after 8s
CI / contract (push) Successful in 56s
CI / build (push) Failing after 1m54s
2026-10-04 19:52:34 +02:00
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 428a12af62 fleetd #722: reconcile presence that arrives before a session's registry entry
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 50s
CI / build (pull_request) Failing after 1m48s
CI / shell-tests (push) Failing after 9s
CI / contract (push) Successful in 47s
CI / build (push) Failing after 1m46s
A member whose MCP contact lands between launcher.spawn() and registry.put()
had its presence marked, but the SPAWNING -> READY transition that markPresent
triggers found no registry entry yet and silently did nothing. The mark then
persisted while registration left the session in SPAWNING, with nothing to
retry the transition. That left the session undeliverable to reclaim/seat
accounting even though it was present and deliverable.

Add SessionManager.reconcilePresence, called right after registry.put in both
the plain-spawn and worktree-spawn paths, to retry the transition for a
terminal already marked present. One private helper serves both call sites.

Tests cover both orderings (contact-then-register and register-then-contact)
for both spawn paths, plus a terminal never marked present staying in
SPAWNING. The contact-then-register tests use a new PresenceRacingLauncher
test double that marks presence from inside spawn(), before acquire()'s own
registry.put runs.
2026-10-04 19:23:02 +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 1fc9e85bf1 fleetd #729: fold a per-boot nonce into every turnId
CI / shell-tests (pull_request) Failing after 6s
CI / contract (pull_request) Successful in 54s
CI / build (pull_request) Failing after 1m57s
askSeq restarts at 0 on every daemon boot, so a turnId (session#n)
minted by one Rendezvous instance could be minted again by a later
instance and resolve to an unrelated ask. Fold a per-instance nonce
into the mint, the same way #719 fixed MessageService's ticket ids.
2026-10-04 18:53:05 +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
Dai Ha 809b7d9b20 Merge remote-tracking branch 'origin/worker/727-ee14ed-3'
CI / shell-tests (push) Failing after 10s
CI / contract (push) Successful in 47s
CI / build (push) Failing after 1m56s
2026-10-04 18:29:28 +02:00
Dai Ha 8cf7215d56 Merge remote-tracking branch 'origin/worker/719-bdd95e-4'
CI / shell-tests (push) Failing after 6s
CI / contract (push) Successful in 51s
CI / build (push) Failing after 1m58s
2026-10-04 18:19:55 +02:00
Dai Ha cf0c9b9316 fleetd #719: make the foreign-id test reach a colliding sequence number
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 52s
CI / build (pull_request) Failing after 1m49s
other.sendAsync had never been called, so other's tasks map was empty and
poll(ticket) returned null regardless of whether the nonce existed — the
test passed against an empty map, not against a colliding id. Mint once on
other so it reaches the same sequence number as the first instance, making
the test exercise the actual collision the nonce guards against.
2026-10-04 18:14:01 +02:00
Dai Ha 337dbd491e fleetd #719: fold a per-boot nonce into every ticket id
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 1m1s
CI / build (pull_request) Failing after 2m0s
ticketSeq restarted at zero on every daemon boot with no persistence, so a
ticket id minted in one boot could be reused by a later boot and resolve to
an unrelated Task instead of failing to resolve at all. Mint each
MessageService instance's own short nonce once and fold it into every ticket
(task-<nonce>-<n>), so an id from one instance can never match another's id
space.

Adds a disjoint-id-space test and a foreign-instance-ticket test (with the
positive control) in MessageServiceTest.
2026-10-04 18:05:06 +02:00
24 changed files with 1010 additions and 101 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,15 @@ 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. MessageService compares a ticket's creator to the
// caller on every read as well, so dropping this gate would not expose another
// session's reply — it would move the refusal later and widen what a caller that
// never orchestrates can probe.
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
@@ -355,6 +355,14 @@ public final class MessageService {
*/
private final ConcurrentHashMap<String, Boolean> queuedDeliveries = new ConcurrentHashMap<>();
private final AtomicLong ticketSeq = new AtomicLong();
/**
* Minted once per {@code MessageService} instance and folded into every ticket id (see
* {@link #sendAsync(String, String, Runnable, String)}). {@link #ticketSeq} alone restarts at
* zero for every instance, so without this a ticket id can be reused across instances and
* resolve to an unrelated {@link Task} with no error; this nonce makes that impossible, because
* an id minted by one instance can never match the id space of another.
*/
private final String ticketBootNonce = UUID.randomUUID().toString().substring(0, 6);
private final ExecutorService asyncExecutor = Executors.newThreadPerTaskExecutor(
Thread.ofVirtual().name("bridge-async-", 0).factory());
@@ -1321,7 +1329,7 @@ public final class MessageService {
* @return the ticket to poll for the eventual result
*/
public String sendAsync(String target, String content, Runnable onAccepted, String creatorTerminal) {
String ticket = "task-" + ticketSeq.incrementAndGet();
String ticket = "task-" + ticketBootNonce + "-" + ticketSeq.incrementAndGet();
Task task = new Task(ticket, target, nowNanos, creatorTerminal);
tasks.put(ticket, task);
if (pushLoop != null) {
@@ -1,5 +1,6 @@
package dev.ltms.fleet.msg;
import java.util.UUID;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicLong;
@@ -101,6 +102,14 @@ public final class Rendezvous {
/** Reverse rendezvous (CB-205): worker questions awaiting the primary's answer, keyed by {@code turnId}. */
private final ConcurrentHashMap<String, AskWaiter> asks = new ConcurrentHashMap<>();
private final AtomicLong askSeq = new AtomicLong();
/**
* Minted once per {@code Rendezvous} instance and folded into every {@code turnId} (see
* {@link #openAsk(String)}). {@link #askSeq} alone restarts at zero for every instance, so
* without this a {@code turnId} minted by one instance could be minted again by another and
* resolve to an unrelated ask with no error; this nonce makes that impossible, because an id
* minted by one instance can never match the id space of another.
*/
private final String askBootNonce = UUID.randomUUID().toString().substring(0, 6);
/** Per-session index of the currently-open ask, so duplicate fleet_ask calls coalesce onto one turn. */
private final ConcurrentHashMap<String, String> openAsksBySession = new ConcurrentHashMap<>();
@@ -187,7 +196,7 @@ public final class Rendezvous {
while (true) {
AskWaiter[] minted = { null };
String turnId = openAsksBySession.computeIfAbsent(session, _ -> {
String newTurnId = session + "#" + askSeq.incrementAndGet();
String newTurnId = session + "#" + askBootNonce + "-" + askSeq.incrementAndGet();
CompletableFuture<String> answer = new CompletableFuture<>();
AskWaiter waiter = new AskWaiter(session, answer, ownerOf(session));
asks.put(newTurnId, waiter);
@@ -255,6 +255,10 @@ public final class SessionManager implements TurnListener {
handle.id(), handle.terminalId(), resolvedProfile, actualRole, cwd, ownerTerminal, now, now, 0,
MemberSession.State.SPAWNING, null, null, handle.charterReceipt(), handle.agentSessionId());
registry.put(handle.id(), session);
// A presence contact that already arrived for this terminal found no registry
// entry to transition and gave up silently. Retry it now that one exists; remove
// this call and such a session stays in SPAWNING even though it is present.
reconcilePresence(handle.terminalId());
handles.put(handle.id(), handle);
log.debug("acquired session id={} terminal={} profile={} owner={}",
handle.id(), handle.terminalId(), session.profile(), session.ownerTerminal());
@@ -809,6 +813,10 @@ public final class SessionManager implements TurnListener {
handle.charterReceipt(),
handle.agentSessionId());
registry.put(handle.id(), session);
// A presence contact that already arrived for this terminal found no registry entry to
// transition and gave up silently. Retry it now that one exists; remove this call and
// such a session stays in SPAWNING even though it is present.
reconcilePresence(handle.terminalId());
handles.put(handle.id(), handle);
log.debug("acquired worktree session id={} terminal={} profile={} branch={} path={}",
handle.id(), handle.terminalId(), session.profile(), session.branch(), session.worktree());
@@ -983,6 +991,20 @@ public final class SessionManager implements TurnListener {
transitionByTerminal(terminalId, MemberSession.State.SPAWNING, MemberSession.State.READY);
}
/**
* Completes a newly registered session's {@code SPAWNING -> READY} transition when {@code
* terminalId} was already marked present before this ran. A terminal never marked present is
* left in {@code SPAWNING}; it reaches {@code READY} normally through {@link #onReady} once
* its own contact arrives. Callers must run this only once the session's registry entry is
* already visible — {@link #onReady}'s transition matches against that entry, and reconciling
* before the entry exists finds nothing to transition.
*/
private void reconcilePresence(String terminalId) {
if (terminalId != null && !terminalId.isBlank() && presence.isPresent(terminalId)) {
onReady(terminalId);
}
}
/**
* Lifecycle hook: a message was delivered into the worker — it is now busy on a turn.
* The turn count is bumped and the activity timestamp is refreshed. A {@code DONE} session
@@ -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
@@ -888,6 +888,39 @@ class MessageServiceTest {
assertNull(messages.poll("task-999999"), "a ticket that was never minted is unknown");
}
// --- fleetd #719: a per-boot nonce keeps one instance's ticket ids out of another's space ---
/** A second, fully independent instance — its own agents/injector/rendezvous/inbox, not shared. */
private MessageService newIndependentInstance() {
FakeHerdr otherHerdr = new FakeHerdr().readText("BUILD GREEN: 391 files");
AgentControl otherAgents = new AgentControl(otherHerdr);
Injector otherInjector = new Injector(otherAgents);
return new MessageService(otherAgents, otherInjector, new Rendezvous(), new InMemoryReplyInbox());
}
@Test
void twoInstancesMintDisjointTicketIds() {
MessageService other = newIndependentInstance();
String ticketFromThis = messages.sendAsync(T, "task on first instance", null, null);
String ticketFromOther = other.sendAsync(T, "task on second instance", null, null);
assertNotEquals(ticketFromThis, ticketFromOther,
"each instance mints its own id space, so even a first ticket from each must differ");
}
@Test
void foreignInstanceTicketDoesNotResolve() {
MessageService other = newIndependentInstance();
String ticket = messages.sendAsync(T, "task on first instance", null, null);
// `other` must reach the same sequence number, or this test passes against an empty map
// instead of against a colliding id.
other.sendAsync(T, "task on second instance", null, null);
// control: the id resolves in the instance that minted it, so a null below cannot be
// explained by broken plumbing — only by the ticket being foreign to `other`.
assertNotNull(messages.poll(ticket), "the minting instance must still resolve its own ticket");
assertNull(other.poll(ticket), "a ticket minted by a different instance must not resolve here");
}
// --- fleetd #705: a ticket's creator terminal gates who may poll it -------------------------
@Test
@@ -191,6 +191,53 @@ class RendezvousTest {
assertNull(rendezvous.askOwner(t.turnId()), "a closed ask no longer reports an owner");
}
// ── fleetd #729: per-boot nonce guards turnId against cross-instance reuse ────────────────
@Test
void twoInstancesMintDisjointTurnIds() {
Rendezvous other = new Rendezvous();
Rendezvous.AskTicket fromThis = rendezvous.openAsk(W);
Rendezvous.AskTicket fromOther = other.openAsk(W);
assertNotEquals(fromThis.turnId(), fromOther.turnId(),
"each instance mints its own id space, so even a first ask from each must differ");
}
@Test
void foreignInstanceTurnIdDoesNotResolve() {
Rendezvous other = new Rendezvous();
// `other` must reach the same sequence number as `rendezvous` (two asks each, the first
// closed so the second mints fresh), or this test passes against an empty map instead of
// against a colliding id.
Rendezvous.AskTicket firstFromThis = rendezvous.openAsk(W);
rendezvous.closeAsk(firstFromThis.turnId());
Rendezvous.AskTicket secondFromThis = rendezvous.openAsk(W);
Rendezvous.AskTicket firstFromOther = other.openAsk(W);
other.closeAsk(firstFromOther.turnId());
other.openAsk(W);
// control: the id resolves in the instance that minted it, so a false below cannot be
// explained by broken plumbing — only by the turnId being foreign to `other`.
assertTrue(rendezvous.answerAsk(secondFromThis.turnId(), "answer from this instance"),
"the minting instance must still resolve its own turnId");
assertFalse(other.answerAsk(secondFromThis.turnId(), "answer from other instance"),
"a turnId minted by a different instance must not resolve here");
}
@Test
void openAskStillCoalescesDuplicatesAndStillMintsDistinctIdsPerAsk() {
Rendezvous.AskTicket t1 = rendezvous.openAsk(W);
Rendezvous.AskTicket t2 = rendezvous.openAsk(W);
assertEquals(t1.turnId(), t2.turnId(),
"a second openAsk while one is open still coalesces onto the same turn");
assertFalse(t2.fresh(), "the coalesced ask is still reported as not fresh");
rendezvous.closeAsk(t1.turnId());
Rendezvous.AskTicket t3 = rendezvous.openAsk(W);
assertNotEquals(t1.turnId(), t3.turnId(), "two asks from the same session still get different turnIds");
}
@Test
void ownerPermitsIsFailClosedOnARecordAndThreeStatesAreDistinct() {
assertFalse(Rendezvous.Owner.permits(null, null),
@@ -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,
@@ -0,0 +1,92 @@
package dev.ltms.fleet.session;
import dev.ltms.fleet.inject.MemberPresence;
import dev.ltms.fleet.peer.Capability;
import dev.ltms.fleet.peer.PeerHandle;
import dev.ltms.fleet.peer.PeerLauncher;
import dev.ltms.fleet.peer.SpawnRequest;
import dev.ltms.fleet.placement.PlacementDecision;
import java.util.List;
import java.util.Set;
/**
* {@link PeerLauncher} decorator that marks presence for a spawned terminal before returning its
* handle to the caller — the contact-then-register ordering fleetd #722 covers, where the
* terminal's MCP contact lands before {@link SessionManager#acquire} runs its own
* {@code registry.put}. The presence view is set after construction, via {@link #presence},
* because it is owned by the {@link SessionManager} this launcher is passed into.
*/
final class PresenceRacingLauncher implements PeerLauncher {
private final PeerLauncher delegate;
volatile MemberPresence presence;
PresenceRacingLauncher(PeerLauncher delegate) {
this.delegate = delegate;
}
@Override
public PeerHandle spawn(SpawnRequest req) {
PeerHandle handle = delegate.spawn(req);
presence.markPresent(handle.terminalId());
return handle;
}
@Override
public PeerHandle spawn(SpawnRequest req, PlacementDecision decision) {
PeerHandle handle = delegate.spawn(req, decision);
presence.markPresent(handle.terminalId());
return handle;
}
@Override
public Set<Capability> capabilities() {
return delegate.capabilities();
}
@Override
public Set<Capability> capabilitiesFor(String profileName) {
return delegate.capabilitiesFor(profileName);
}
@Override
public Set<String> profiles() {
return delegate.profiles();
}
@Override
public String defaultProfile() {
return delegate.defaultProfile();
}
@Override
public String effectiveCwd(SpawnRequest req) {
return delegate.effectiveCwd(req);
}
@Override
public List<String> parityOverlay(String profileName) {
return delegate.parityOverlay(profileName);
}
@Override
public List<?> list() {
return delegate.list();
}
@Override
public int reapOrphanWorkers() {
return delegate.reapOrphanWorkers();
}
@Override
public void stop(String id) {
delegate.stop(id);
}
@Override
public boolean clearContext(String id) {
return delegate.clearContext(id);
}
}
@@ -422,6 +422,53 @@ class SessionManagerTest {
"turn completion moves BUSY → DONE");
}
// --- fleetd #722: registration and presence must reach READY whichever lands first --------
@Test
void registerThenContactReachesReadyForPlainSpawn() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
sessions.asPresence().markPresent(session.terminalId());
assertEquals(MemberSession.State.READY, sessions.get(session.paneId()).orElseThrow().state(),
"a presence contact that arrives after registration reaches READY");
}
@Test
void contactThenRegisterStillReachesReadyForPlainSpawn() {
// The racing launcher marks presence for the spawned terminal from inside spawn() —
// before SessionManager.acquire's own registry.put runs — modeling an MCP contact that
// lands in that window.
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile cfg = new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
List.of("ccs", "ltms-local"), "tab", "fleetd-workers",
"worker: {profile} #{n}", null, null, null);
ClaudeCodeLauncher workers = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
PresenceRacingLauncher race = new PresenceRacingLauncher(workers);
SessionManager sessions = new SessionManager(race);
race.presence = sessions.asPresence();
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
assertEquals(MemberSession.State.READY, sessions.get(session.paneId()).orElseThrow().state(),
"a presence contact that lands before registry.put must still reach READY");
}
@Test
void aTerminalNeverMarkedPresentStaysSpawningAfterRegistration() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
assertEquals(MemberSession.State.SPAWNING, sessions.get(session.paneId()).orElseThrow().state(),
"registration alone must not advance a terminal that was never marked present");
}
@Test
void releaseTearsDownWorkerAndRemovesFromRosterAndIsIdempotent() {
FakeHerdr herdr = new FakeHerdr();
@@ -159,6 +159,40 @@ class WorktreeSessionManagerTest {
assertEquals(expectedPath, s.cwd(), "session cwd is the worktree path");
}
// --- fleetd #722: registration and presence must reach READY whichever lands first --------
@Test
void registerThenContactReachesReadyForWorktreeSpawn() {
FakeHerdr herdr = new FakeHerdr();
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt");
SessionManager sessions = new SessionManager(workerService(herdr), worktrees);
MemberSession session = sessions.acquire("ltms-local", null, "/caller/proj", "term_primary",
new WorktreeRequest("cb-722", null));
sessions.asPresence().markPresent(session.terminalId());
assertEquals(MemberSession.State.READY, sessions.get(session.paneId()).orElseThrow().state(),
"a presence contact that arrives after worktree registration reaches READY");
}
@Test
void contactThenRegisterStillReachesReadyForWorktreeSpawn() {
// The racing launcher marks presence for the spawned terminal from inside spawn() —
// before SessionManager.acquireWithWorktree's own registry.put runs — modeling an MCP
// contact that lands in that window.
FakeHerdr herdr = new FakeHerdr();
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt");
PresenceRacingLauncher race = new PresenceRacingLauncher(workerService(herdr));
SessionManager sessions = new SessionManager(race, worktrees);
race.presence = sessions.asPresence();
MemberSession session = sessions.acquire("ltms-local", null, "/caller/proj", "term_primary",
new WorktreeRequest("cb-722", null));
assertEquals(MemberSession.State.READY, sessions.get(session.paneId()).orElseThrow().state(),
"a presence contact that lands before worktree registration must still reach READY");
}
@Test
void worktreeArchitectAcquireAlsoBindsItsSlot() {
FakeHerdr herdr = new FakeHerdr();