Compare commits

...

12 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 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 1ef93e57cc fleetd #727: give a lead launch the three protections every member spawn gets
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 50s
CI / build (pull_request) Failing after 1m59s
Extract checkPaneCommandFits, the agent_pane_busy retry, and the agent_name_taken
retry out of HerdrPeerLauncher into a shared dev.ltms.fleet.herdr.ResilientAgentLaunch,
and route LeadLauncher.launch through the same seam instead of a bare agents.start
call. The lead's agent name now carries a per-process nonce and a per-start sequence
number (like a member's), so a stale agent_name_taken from an earlier crashed session
no longer blocks a legitimate relaunch outright.
2026-10-04 18:14:06 +02:00
20 changed files with 1112 additions and 191 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
}
@@ -0,0 +1,146 @@
package dev.ltms.fleet.herdr;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.nio.charset.StandardCharsets;
import java.util.List;
import java.util.function.IntFunction;
/**
* The three protections every {@code agent.start} caller needs against herdr's pane-typed launch
* surface (fleetd #220, #727): a byte-limit check on the assembled command line, a bounded retry
* on {@code agent_pane_busy} (the target pane's shell has not reached its prompt yet), and a
* bounded retry on {@code agent_name_taken} with a fresh name each attempt. One implementation —
* every caller of {@code agent.start}, lead or member, goes through this seam rather than carrying
* its own copy.
*/
public final class ResilientAgentLaunch {
private static final Logger log = LoggerFactory.getLogger(ResilientAgentLaunch.class);
private ResilientAgentLaunch() {
}
/**
* The pty line buffer herdr types a launch command into: BSD/macOS {@code MAX_CANON}. Not a
* fleetd choice and not configurable — see {@link #checkFits}.
*/
public static final int PANE_COMMAND_BYTE_LIMIT = 1024;
/** Per-argument allowance for the separating space and a shell quote pair fleetd cannot see. */
private static final int QUOTING_OVERHEAD_PER_ARG = 3;
/** herdr rejects a duplicate agent {@code name}; a caller retries a bumped name this many times. */
public static final int NAME_RETRIES = 8;
/**
* Retries for {@code agent.start} against a seed pane whose shell has not reached its prompt
* yet — {@code tab.create}/{@code pane.split} return as soon as the pane exists, and herdr
* refuses to start an agent in a pane that is not "an available shell" ({@code agent_pane_busy}).
*/
public static final int SHELL_READY_RETRIES = 20;
/** Raised by {@link #checkFits} when the assembled command cannot fit the pane line. */
public static final class TooLargeException extends RuntimeException {
public TooLargeException(String message) {
super(message);
}
}
/**
* Verify the assembled launch line fits the pane herdr types it into. herdr does not exec the
* launch command — it TYPES it into the pane as one line, and a pty line buffer holds only
* {@value #PANE_COMMAND_BYTE_LIMIT} bytes. Everything past that byte is dropped with no error
* anywhere: herdr answers "agent started", the backend exits on the mangled argument it was
* handed, the pane closes, and the only symptom is a readiness timeout with no reason. That is
* how fleetd #214 broke every claude-code spawn — one 50-byte flag pushed a 978-byte command to
* 1028, and the tail that got cut was {@code --autocompact 250000}.
*
* <p>So measure it here and refuse, loudly and immediately, rather than start something that
* cannot work. The estimate is deliberately conservative: fleetd cannot see herdr's quoting, so
* every argument is charged its own bytes plus a separator and a quote pair. An over-estimate
* costs a clear error at a length that was already unsafe; an under-estimate would let the
* silent truncation back in.
*
* @param label names the launch in the refusal message (a profile name)
* @param argv the full argv, including the executable at index 0
* @throws TooLargeException naming the limit, the estimate, and the longest argument
*/
public static void checkFits(String label, List<String> argv) {
int bytes = 0;
String longest = null;
int longestBytes = 0;
for (String arg : argv) {
int argBytes = arg == null ? 0 : arg.getBytes(StandardCharsets.UTF_8).length;
bytes += argBytes + QUOTING_OVERHEAD_PER_ARG;
if (argBytes > longestBytes) {
longestBytes = argBytes;
longest = arg;
}
}
if (bytes <= PANE_COMMAND_BYTE_LIMIT) {
return;
}
String culprit = longest == null ? "<none>"
: longest.substring(0, Math.min(longest.length(), 60)) + (longest.length() > 60 ? "…" : "");
throw new TooLargeException(
"launch command for " + label + " is about " + bytes + " bytes, over the "
+ PANE_COMMAND_BYTE_LIMIT + "-byte limit of the pane line herdr types it into. "
+ "The pty would drop the tail silently and the backend would exit on a mangled "
+ "argument. Longest argument is " + longestBytes + " bytes: " + culprit
+ " — move it off the command line (a file flag) or shorten it.");
}
/**
* Start an agent into {@code paneId}, retrying {@code agent_pane_busy} up to {@code retries}
* times with {@code sleeper} run between attempts.
*
* @throws HerdrException the last {@code agent_pane_busy} failure once {@code retries} is
* spent, or immediately for any other herdr failure
*/
public static Agent startAwaitingShellPrompt(AgentControl agents, String name, String kind,
List<String> args, String paneId,
int retries, Runnable sleeper) {
HerdrException busy = null;
for (int attempt = 0; attempt < retries; attempt++) {
try {
return agents.start(name, kind, args, paneId);
} catch (HerdrException e) {
if (!"agent_pane_busy".equals(e.code())) throw e;
log.debug("pane {} not at its shell prompt yet, retrying agent.start", paneId);
busy = e;
sleeper.run();
}
}
throw busy;
}
/**
* Start an agent under a freshly generated name each attempt, retrying {@code agent_name_taken}
* up to {@code nameRetries} times — herdr refuses a duplicate {@code name} outright, so a stale
* registry entry (a crashed session, a name the registry has not yet released) must not block a
* legitimate relaunch. Each attempt also carries its own {@link #startAwaitingShellPrompt} retry.
*
* @param nameForAttempt called once per attempt (0-based) to produce that attempt's name
* @throws HerdrException the last {@code agent_name_taken} failure once {@code nameRetries} is
* spent, or immediately for any other herdr failure
*/
public static Agent startUniquelyNamed(AgentControl agents, String kind, List<String> args,
String paneId, IntFunction<String> nameForAttempt,
int nameRetries, int shellReadyRetries, Runnable sleeper) {
HerdrException last = null;
for (int attempt = 0; attempt < nameRetries; attempt++) {
String name = nameForAttempt.apply(attempt);
try {
return startAwaitingShellPrompt(agents, name, kind, args, paneId,
shellReadyRetries, sleeper);
} catch (HerdrException e) {
if (!"agent_name_taken".equals(e.code())) throw e;
log.debug("agent name '{}' taken, retrying", name);
last = e;
}
}
throw last;
}
}
@@ -5,6 +5,7 @@ import dev.ltms.fleet.herdr.Agent;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.HerdrException;
import dev.ltms.fleet.herdr.PendingCloseMarker;
import dev.ltms.fleet.herdr.ResilientAgentLaunch;
import dev.ltms.fleet.herdr.Tab;
import dev.ltms.fleet.herdr.Workspace;
import dev.ltms.fleet.herdr.WorkspaceControl;
@@ -12,6 +13,7 @@ import dev.ltms.fleet.launch.ClaudeCodeArguments;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.security.SecureRandom;
import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.LinkedHashSet;
@@ -19,6 +21,7 @@ import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
import java.util.concurrent.atomic.AtomicLong;
import dev.ltms.fleet.peer.PeerLauncher;
/**
@@ -80,9 +83,19 @@ 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;
private final Runnable sleeper;
// Per-process token mixed into each lead agent name so a fresh daemon process (seq back at 0)
// cannot collide with a same-name lead that outlived a restart — the same scheme
// HerdrPeerLauncher uses for members (fleetd #727).
private final String nameNonce = String.format("%06x", new SecureRandom().nextInt(1 << 24));
private final AtomicLong nameSeq = new AtomicLong();
/**
* @param agents herdr agent control (start, list)
@@ -90,9 +103,27 @@ public final class LeadLauncher {
* @param cfg the loaded config — {@code fleet.leaders}, {@code profiles} and each lead's tab
*/
public LeadLauncher(AgentControl agents, WorkspaceControl spaces, FleetConfig cfg) {
this(agents, spaces, cfg, () -> sleepUninterruptibly(300));
}
/**
* Test seam: as above, plus an injectable {@code sleeper} for the {@code agent_pane_busy}
* retry (fleetd #727), so a test can prove the retry budget without a real sleep.
*/
LeadLauncher(AgentControl agents, WorkspaceControl spaces, FleetConfig cfg, Runnable sleeper) {
this.agents = agents;
this.spaces = spaces;
this.cfg = cfg;
this.sleeper = sleeper;
}
/** Uninterruptible sleep — the production {@link #sleeper} between {@code agent_pane_busy} retries. */
private static void sleepUninterruptibly(long ms) {
try {
Thread.sleep(ms);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
/**
@@ -173,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++;
}
}
@@ -198,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}
@@ -306,8 +403,17 @@ public final class LeadLauncher {
return null;
}
/** Start one lead. Returns false (having logged) rather than throwing on any failure. */
private boolean launch(String name, FleetConfig.Leader lead, FleetConfig.Profile profile) {
/**
* 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
* types it into, a stale {@code agent_name_taken} (a crashed session's name the registry has
* not yet released) is retried under a fresh per-attempt name rather than refusing the whole
* 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 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();
@@ -323,8 +429,12 @@ public final class LeadLauncher {
// Same shape as the member launchers: herdr resolves the executable from `kind`, so
// argv[0] (the configured launcher, e.g. `ccs`) is dropped and only the rest is passed.
List<String> argv = leadArgv(profile);
Agent started = agents.start("lead-" + name, herdrKind(profile),
argv.isEmpty() ? argv : argv.subList(1, argv.size()), tab.rootPaneId());
ResilientAgentLaunch.checkFits(profile.profile(), argv);
List<String> args = argv.isEmpty() ? argv : argv.subList(1, argv.size());
Agent started = ResilientAgentLaunch.startUniquelyNamed(agents, herdrKind(profile), args,
tab.rootPaneId(),
attempt -> "lead-" + name + "-" + nameNonce + "-" + nameSeq.incrementAndGet(),
ResilientAgentLaunch.NAME_RETRIES, ResilientAgentLaunch.SHELL_READY_RETRIES, sleeper);
// Label AFTER the start succeeds. A label written before would survive a failed start
// and then read back as a live lead on the next boot, which is the exact staleness the
@@ -334,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());
@@ -346,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
@@ -6,6 +6,7 @@ import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.herdr.HerdrException;
import dev.ltms.fleet.herdr.ResilientAgentLaunch;
import dev.ltms.fleet.herdr.Tab;
import dev.ltms.fleet.herdr.Workspace;
import dev.ltms.fleet.herdr.WorkspaceControl;
@@ -67,16 +68,6 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
private static final Logger log = LoggerFactory.getLogger(HerdrPeerLauncher.class);
/** herdr rejects a duplicate agent {@code name}; we retry a bumped name this many times. */
private static final int NAME_RETRIES = 8;
/**
* Retries for {@code agent.start} against a seed pane whose shell has not reached its prompt
* yet — {@code tab.create}/{@code pane.split} return as soon as the pane exists, and herdr
* refuses to start an agent in a pane that is not "an available shell" ({@code agent_pane_busy}).
*/
private static final int SHELL_READY_RETRIES = 20;
private final String namePrefix; // label prefix: naming + reap scheme
private final AgentControl agents;
private final WorkspaceControl spaces;
@@ -766,90 +757,26 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
// Protocol 19 resolves the executable from the agent kind (== namePrefix here), so
// argv[0] — the configured executable — is dropped and only the extra args are passed.
List<String> args = argv.isEmpty() ? argv : argv.subList(1, argv.size());
checkPaneCommandFits(cfg, argv);
HerdrException last = null;
for (int attempt = 0; attempt < NAME_RETRIES; attempt++) {
long seq = nameSeq.incrementAndGet();
String name = namePrefix + "-" + cfg.profile() + "-" + nameNonce + "-" + seq;
try {
return new Started(startAwaitingShellPrompt(name, args, paneId), seq);
} catch (HerdrException e) {
if (!"agent_name_taken".equals(e.code())) throw e;
log.debug("peer name '{}' taken, retrying", name);
last = e;
}
try {
ResilientAgentLaunch.checkFits(cfg.profile(), argv);
} catch (ResilientAgentLaunch.TooLargeException e) {
throw new PeerUnreachableException(e.getMessage());
}
throw last;
}
/**
* fleetd #220: herdr does not exec the launch command — it TYPES it into the pane as one line,
* and a pty line buffer holds only {@value #PANE_COMMAND_BYTE_LIMIT} bytes (BSD/macOS {@code
* MAX_CANON}). Everything past that byte is dropped. Nothing reports it: herdr answers "agent
* started", the backend exits on the mangled argument it was handed, the pane closes, and the
* only symptom is {@link #waitUntilInjectableOrThrow} timing out 20 seconds later with no
* reason. That is exactly how #214 broke every claude-code spawn — one 50-byte flag pushed a
* 978-byte command to 1028, and the tail that got cut was {@code --autocompact 250000}.
*
* <p>So measure it here and refuse, loudly and immediately, rather than spawn something that
* cannot work. The estimate is deliberately conservative: fleetd cannot see herdr's quoting, so
* every argument is charged its own bytes plus a separator and a quote pair. An over-estimate
* costs a clear error at a length that was already unsafe; an under-estimate would let the
* silent truncation back in.
*
* @throws PeerUnreachableException when the command cannot fit — the same failure the spawn
* would have hit anyway, named at the point it is still
* explainable
*/
private void checkPaneCommandFits(FleetConfig.Profile cfg, List<String> argv) {
int bytes = 0;
String longest = null;
int longestBytes = 0;
for (String arg : argv) {
int argBytes = arg == null ? 0 : arg.getBytes(java.nio.charset.StandardCharsets.UTF_8).length;
bytes += argBytes + QUOTING_OVERHEAD_PER_ARG;
if (argBytes > longestBytes) {
longestBytes = argBytes;
longest = arg;
}
}
if (bytes <= PANE_COMMAND_BYTE_LIMIT) {
return;
}
String culprit = longest == null ? "<none>"
: longest.substring(0, Math.min(longest.length(), 60)) + (longest.length() > 60 ? "…" : "");
throw new PeerUnreachableException(
"launch command for profile " + cfg.profile() + " is about " + bytes + " bytes, over the "
+ PANE_COMMAND_BYTE_LIMIT + "-byte limit of the pane line herdr types it into. "
+ "The pty would drop the tail silently and the backend would exit on a mangled "
+ "argument. Longest argument is " + longestBytes + " bytes: " + culprit
+ " — move it off the command line (a file flag) or shorten it.");
long[] lastSeq = {0};
Agent agent = ResilientAgentLaunch.startUniquelyNamed(agents, namePrefix, args, paneId,
attempt -> {
lastSeq[0] = nameSeq.incrementAndGet();
return namePrefix + "-" + cfg.profile() + "-" + nameNonce + "-" + lastSeq[0];
},
ResilientAgentLaunch.NAME_RETRIES, ResilientAgentLaunch.SHELL_READY_RETRIES, sleeper);
return new Started(agent, lastSeq[0]);
}
/**
* The pty line buffer herdr types a launch command into: BSD/macOS {@code MAX_CANON}. Not a
* fleetd choice and not configurable — see {@link #checkPaneCommandFits}.
* fleetd choice and not configurable — see {@link ResilientAgentLaunch#checkFits}.
*/
static final int PANE_COMMAND_BYTE_LIMIT = 1024;
/** Per-argument allowance for the separating space and a shell quote pair fleetd cannot see. */
private static final int QUOTING_OVERHEAD_PER_ARG = 3;
/** Start the agent into {@code paneId}, waiting out the seed shell's boot with the sleeper. */
private Agent startAwaitingShellPrompt(String name, List<String> args, String paneId) {
HerdrException busy = null;
for (int attempt = 0; attempt < SHELL_READY_RETRIES; attempt++) {
try {
return agents.start(name, namePrefix, args, paneId);
} catch (HerdrException e) {
if (!"agent_pane_busy".equals(e.code())) throw e;
log.debug("pane {} not at its shell prompt yet, retrying agent.start", paneId);
busy = e;
sleeper.run();
}
}
throw busy;
}
static final int PANE_COMMAND_BYTE_LIMIT = ResilientAgentLaunch.PANE_COMMAND_BYTE_LIMIT;
// --- discovery + reap ----------------------------------------------------------------------
@@ -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);
@@ -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");
}
@@ -1,10 +1,16 @@
package dev.ltms.fleet.lead;
import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.ResilientAgentLaunch;
import dev.ltms.fleet.herdr.WorkspaceControl;
import org.junit.jupiter.api.Test;
import org.slf4j.LoggerFactory;
import java.util.LinkedHashMap;
import java.util.List;
@@ -64,6 +70,21 @@ class LeadLauncherTest {
return new LeadLauncher(new AgentControl(herdr), new WorkspaceControl(herdr), cfg);
}
/** As {@link #launcher}, plus a fast no-op sleeper so a busy-retry test never real-sleeps. */
private static LeadLauncher fastLauncher(FakeHerdr herdr, FleetConfig cfg) {
return new LeadLauncher(new AgentControl(herdr), new WorkspaceControl(herdr), cfg, () -> { });
}
/** {@link #opusProfile()} with one argv element long enough to overflow the pane line limit. */
private static FleetConfig.Profile hugeArgvProfile() {
return new FleetConfig.Profile(
"opus", null, "claude-opus-5", null, "FLEETD_WORKER_TOKEN",
List.of("ccs", "x".repeat(1500)), "tab", "fleet", null,
"http://127.0.0.1:8765/mcp", null, null,
null, null, null,
Map.of("CLAUDE_CODE_AUTO_COMPACT_WINDOW", "300000"), null, null, true, null);
}
@SuppressWarnings("unchecked")
private static List<String> startedArgs(FakeHerdr herdr) {
return (List<String>) ((Map<String, Object>) herdr.lastCall("agent.start").params()).get("args");
@@ -86,7 +107,11 @@ class LeadLauncherTest {
assertEquals(1, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads());
assertTrue(herdr.called("agent.start"), "a lead must actually be started");
assertEquals("lead-opus", startedName(herdr));
// fleetd #727: the name carries a per-process nonce and a per-start sequence number — the
// same unique-naming scheme HerdrPeerLauncher uses for members — rather than the fixed
// "lead-opus" a stale registry entry could block a legitimate relaunch under.
assertTrue(startedName(herdr).matches("lead-opus-[0-9a-f]{6}-\\d+"),
"name is lead-<name>-<nonce>-<seq>: " + startedName(herdr));
}
/** The tab is labelled with the configured `tab:` so the scanner finds the lead on the next resolve. */
@@ -436,4 +461,247 @@ class LeadLauncherTest {
assertEquals(0, launcher(herdr, cfg).ensureLeads());
assertTrue(herdr.calls.isEmpty(), "nothing declared ⇒ nothing scanned");
}
// ── fleetd #727: the same three launch protections every member spawn gets ───────────────────
/**
* A freshly created pane may not have redrawn its prompt yet, so herdr answers
* {@code agent_pane_busy}. The lead launch must wait it out rather than fail on the first miss
* — exactly the retry {@code HerdrPeerLauncher} already gives every member.
*/
@Test
void aLeadLaunchRetriesWhileTheSeedShellBoots() {
FakeHerdr herdr = new FakeHerdr().agentPaneBusyTimes(2);
assertEquals(1, fastLauncher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads(),
"the lead must still start once the shell is ready");
assertEquals(3, herdr.calls.stream().filter(c -> c.method().equals("agent.start")).count(),
"two busy rejections, then the successful start");
}
/**
* The busy retry is bounded, not an infinite poll. If the pane never becomes ready the launch
* must eventually give up and log the failure, not hang the daemon's reconcile loop forever —
* proven here by a budget that would still be busy on attempt
* {@value ResilientAgentLaunch#SHELL_READY_RETRIES} and a call count that stops exactly there.
*/
@Test
void aLeadLaunchGivesUpAfterTheBoundedBusyBudgetRatherThanLoopingForever() {
FakeHerdr herdr = new FakeHerdr().agentPaneBusyTimes(999);
assertEquals(0, fastLauncher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads(),
"a pane that never becomes ready must not be reported as a started lead");
assertEquals(ResilientAgentLaunch.SHELL_READY_RETRIES,
herdr.calls.stream().filter(c -> c.method().equals("agent.start")).count(),
"the retry budget is bounded: it stops after exactly SHELL_READY_RETRIES attempts");
}
/**
* herdr refuses a duplicate agent {@code name} outright. A stale registry entry — a crashed
* lead session, or a name the registry has not yet released — must not permanently block a
* legitimate relaunch, so each retry attempt carries a fresh per-start name.
*/
@Test
@SuppressWarnings("unchecked")
void aLeadLaunchRetriesUnderAFreshNameWhenTheOldNameIsStillTaken() {
FakeHerdr herdr = new FakeHerdr().agentNameTakenTimes(2);
assertEquals(1, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads(),
"the lead must still start once a free name is found");
List<String> names = herdr.calls.stream()
.filter(c -> c.method().equals("agent.start"))
.map(c -> ((Map<String, Object>) c.params()).get("name").toString())
.toList();
assertEquals(3, names.size(), "2 rejected + 1 success");
assertEquals(3, Set.copyOf(names).size(), "each attempt must use a distinct name");
}
/**
* The name-collision retry is bounded too. If the name is taken on every attempt, the launch
* must give up rather than keep minting new names forever — the per-start naming scheme means
* every failed attempt was refused outright by herdr (no process exists under a name herdr
* refused), so a bounded, exhausted retry never leaves a second process running: nothing is
* running at all.
*/
@Test
void aLeadLaunchNeverEndsUpWithASecondProcessWhenTheNameStaysTaken() {
FakeHerdr herdr = new FakeHerdr().agentNameTakenTimes(999);
assertEquals(0, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads(),
"a name that is never free must not be reported as a started lead");
assertEquals(ResilientAgentLaunch.NAME_RETRIES,
herdr.calls.stream().filter(c -> c.method().equals("agent.start")).count(),
"the retry budget is bounded: it stops after exactly NAME_RETRIES attempts, "
+ "never racing a duplicate into existence");
}
/**
* fleetd #220/#727: herdr types the launch command into the pane as one line, and a pty line
* buffer holds only 1024 bytes — past that the tail is dropped with no error at all, and the
* backend exits on a mangled argument. The lead launch must refuse an over-long command outright
* rather than let it be typed and silently truncated.
*/
@Test
void anOverlongLeadArgvIsRefusedRatherThanTypedAndTruncated() {
FakeHerdr herdr = new FakeHerdr();
Logger logger = (Logger) LoggerFactory.getLogger(LeadLauncher.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
int started;
try {
started = launcher(herdr, configWith(lead("opus", "lead: opus", 1), hugeArgvProfile()))
.ensureLeads();
} finally {
logger.detachAppender(appender);
}
assertEquals(0, started, "an over-long command must never be reported as a started lead");
assertFalse(herdr.called("agent.start"),
"nothing may be started — a truncated command is worse than no spawn");
String warn = appender.list.stream()
.filter(e -> e.getLevel().equals(Level.WARN))
.map(ILoggingEvent::getFormattedMessage)
.filter(m -> m.contains("failed to launch"))
.findFirst()
.orElseThrow(() -> new AssertionError("expected a WARN naming the launch failure: "
+ appender.list));
assertTrue(warn.contains("1024"), "names the limit: " + warn);
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
@@ -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,