Compare commits
15 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| efd9cdb983 | |||
| aabecce901 | |||
| 6754b4edbc | |||
| 787ae0ed7a | |||
| e3050efe8b | |||
| 11998cd626 | |||
| 8e5394f63f | |||
| 428a12af62 | |||
| a332dfdb2c | |||
| d0f4ae057b | |||
| 7b3beaa209 | |||
| b3b2bf3da6 | |||
| 8bb2aa0be4 | |||
| 4ce3149bfd | |||
| 38544d467c |
@@ -159,6 +159,29 @@ fails.
|
||||
|
||||
`{action: "cancel", token}` drops a pending request without rolling.
|
||||
|
||||
**A change is coming: the roll will restart the process instead of sending `/clear` (fleetd #726
|
||||
unit 2, written 2026-10-04).**
|
||||
|
||||
Today a roll types `/clear` into your pane. Your `claude` process keeps running, so a newer CLI on
|
||||
disk is never loaded. Unit 2 replaces that: the daemon ends the old pane, launches a fresh one,
|
||||
waits for the new terminal to be recognised as a lead, and only then sends the bootstrap text.
|
||||
|
||||
**Everything below about `/clear` is accurate while the old jar is running.** Unit 2 was not merged
|
||||
when this note was written, and a merge is not a deployment.
|
||||
|
||||
**How to tell which one is live: read your own tool list.** If `fleet_handover`'s description says
|
||||
it will "clear your pane", the daemon is serving the old behaviour. If it names a restart, the new
|
||||
behaviour is live. The description comes from the running daemon, so it cannot disagree with the
|
||||
code that is actually loaded.
|
||||
|
||||
Two things change for you once it is live. The `status` outcomes are different: three new failures
|
||||
replace the `/clear` ones. And the "never observed as WORKING after 8 consecutive IDLE/DONE polls"
|
||||
warning described below can no longer appear, because that wait is deleted — so if you still see
|
||||
it, the old jar is running. `TURN_NEVER_SETTLED` does not change, and still means nothing was
|
||||
touched.
|
||||
|
||||
**Delete this note and rewrite the `/clear` paragraphs once the new jar is live.**
|
||||
|
||||
**Things that will surprise you:**
|
||||
|
||||
- **`accepted` does not mean your pane has been cleared.** It means every gate passed and the roll
|
||||
@@ -200,13 +223,18 @@ fails.
|
||||
|
||||
- **The roll can still refuse after `confirm` returns**, and by then there is no caller to tell.
|
||||
Those outcomes are logged only, as `lead-rollover:` lines in the daemon log.
|
||||
- **The bootstrap prompt works end to end. Measured 2026-09-22.** This used to say the fix was
|
||||
unproven (fleetd #489) and told you to expect a failure. That is no longer true. The daemon log
|
||||
now holds four `lead-rollover: rolled` lines, and three of them ran on 2026-09-22 at 10:01:43,
|
||||
10:38:28 and 11:15:47. Each one cleared the old lead and started a fresh session against the
|
||||
handover file, with the configured `bootstrapText` arriving as its first message. No context was
|
||||
lost. The old `Unknown command: /clearFresh` failure from 2026-09-12 does not appear in the log
|
||||
at all. Re-measure both numbers with:
|
||||
- **The bootstrap prompt works end to end. Measured 2026-09-22, re-measured 2026-10-04.** This used
|
||||
to say the fix was unproven (fleetd #489) and told you to expect a failure. That is no longer
|
||||
true. On 2026-09-22 the daemon log held four `lead-rollover: rolled` lines. On 2026-10-04 it holds
|
||||
**20**, against a control of 86 `lead-rollover:` lines. Each roll cleared the old lead and started
|
||||
a fresh session against the handover file, with the configured `bootstrapText` arriving as its
|
||||
first message. No context was lost. The old `Unknown command: /clearFresh` failure from 2026-09-12
|
||||
does not appear in the log at all.
|
||||
|
||||
19 of the 20 carry an `elapsedMs`: median 16507 ms, maximum 48261 ms, and two above 45000 ms. That
|
||||
figure times the **whole** roll, and the wait for your own turn to end dominates it, so do not
|
||||
read it as the cost of the clear. Expect a roll to take tens of seconds, and do not treat a slow
|
||||
one as a failed one. Re-measure all of these with:
|
||||
|
||||
```bash
|
||||
grep -c "lead-rollover: rolled" fleetd/fleetd.out # successful rolls
|
||||
|
||||
@@ -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;
|
||||
@@ -134,12 +145,32 @@ public record Principal(Role role, String terminal, long pid, String name) {
|
||||
return terminal != null && terminal.equals(sessionId);
|
||||
}
|
||||
|
||||
/**
|
||||
* Stable identity used to own tickets and open turns. The unnamed primary has no owner key so
|
||||
* it can use the message layer's primary-wide ticket access rule.
|
||||
*/
|
||||
public String ownerKey() {
|
||||
return switch (role) {
|
||||
case PRIMARY -> name == null ? null : prefixed("leader", name);
|
||||
case WORKER -> prefixed("worker", terminal);
|
||||
case ARCHITECT -> prefixed("architect", terminal);
|
||||
case COLLABORATOR -> prefixed("collaborator", name);
|
||||
case OBSERVER -> prefixed("observer", terminal);
|
||||
case ANONYMOUS -> "anonymous";
|
||||
};
|
||||
}
|
||||
|
||||
private static String prefixed(String role, String identity) {
|
||||
return role + ":" + identity;
|
||||
}
|
||||
|
||||
/** Short, non-sensitive description for audit lines and error details. */
|
||||
public String describe() {
|
||||
return switch (role) {
|
||||
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()),
|
||||
@@ -459,12 +459,14 @@ public final class FleetMcp {
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_send", req.arguments()),
|
||||
str(req.arguments(), "sessionId"));
|
||||
if (denied != null) return denied;
|
||||
String caller = callerTerminal(exchange);
|
||||
Principal caller = principal(exchange);
|
||||
String callerTerminal = caller.terminal();
|
||||
String callerOwner = caller.ownerKey();
|
||||
// CB-548: only a PRIMARY caller may claim the legacy singleton "primary" fallback.
|
||||
// An architect delegates as its own pane but must never become the fallback that
|
||||
// no-delegation inbox nudges target as if it were the primary (the per-target
|
||||
// delegation map does not cure the singleton).
|
||||
recordPrimarySingleton(primaryRegistry, caller, principal(exchange));
|
||||
recordPrimarySingleton(primaryRegistry, callerTerminal, caller);
|
||||
Map<String, Object> a = req.arguments();
|
||||
String target = str(a, "sessionId");
|
||||
String content = str(a, "content");
|
||||
@@ -480,18 +482,18 @@ public final class FleetMcp {
|
||||
// Answering a worker's fleet_ask (CB-205): resolve its blocked question and
|
||||
// block for the worker's reply as it resumes the same turn. This is the same
|
||||
// delegation, so ownership is left untouched (CB-548) — never re-recorded.
|
||||
return answer(messages, turnId, content, timeoutMs(a), caller);
|
||||
return answer(messages, turnId, content, timeoutMs(a), callerOwner);
|
||||
}
|
||||
// CB-548: delegator ownership (which lead's reply nudge this worker routes to,
|
||||
// CB-532) is recorded only once the send is ACCEPTED — MessageService has won the
|
||||
// session lock and queued delivery — via the accepted-delivery callback, never at
|
||||
// request time. A concurrent sender that times out BUSY therefore cannot steal a
|
||||
// live turn's reply routing without ever owning the turn.
|
||||
Runnable onAccepted = () -> primaryRegistry.recordDelegation(target, caller);
|
||||
Runnable onAccepted = () -> primaryRegistry.recordDelegation(target, callerTerminal);
|
||||
// wait defaults to true (block for the reply); wait:false is fire-and-poll.
|
||||
return Boolean.FALSE.equals(a.get("wait"))
|
||||
? sendAsync(messages, target, content, onAccepted, workers.profiles(), caller)
|
||||
: send(messages, target, content, timeoutMs(a), onAccepted, workers.profiles(), caller);
|
||||
: send(messages, target, content, timeoutMs(a), onAccepted, workers.profiles(), callerOwner);
|
||||
};
|
||||
// fleet_reply's identity is the CONNECTION, never an argument — so the authz check
|
||||
// is "is this caller a worker at all", and it can only ever reply as itself.
|
||||
@@ -514,7 +516,7 @@ public final class FleetMcp {
|
||||
(exchange, req) -> {
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_status", req.arguments()), null);
|
||||
if (denied != null) return denied;
|
||||
return status(messages, str(req.arguments(), "sessionId"), callerTerminal(exchange));
|
||||
return status(messages, str(req.arguments(), "sessionId"), principal(exchange).ownerKey());
|
||||
};
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> pollHandler =
|
||||
(exchange, req) -> {
|
||||
@@ -524,7 +526,8 @@ public final class FleetMcp {
|
||||
// The action depends on the ARGUMENTS, not on the tool name -- see pollAction.
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_poll", a), target);
|
||||
if (denied != null) return denied;
|
||||
return poll(messages, leadChannel, str(a, "ticket"), target, coordId, callerTerminal(exchange));
|
||||
return poll(messages, leadChannel, str(a, "ticket"), target, coordId,
|
||||
principal(exchange).ownerKey());
|
||||
};
|
||||
// CB-307 Increment 3: per-msgId ack (not needed in v1 but supported by the inbox).
|
||||
// Acking removes a reply from the inbox, so it is a drain, not a read.
|
||||
@@ -833,9 +836,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());
|
||||
}
|
||||
}
|
||||
@@ -909,7 +917,7 @@ public final class FleetMcp {
|
||||
*/
|
||||
static McpSchema.CallToolResult send(MessageService messages, String sessionId, String content,
|
||||
Long timeoutMs, Runnable onAccepted, Set<String> profiles,
|
||||
String callerTerminal) {
|
||||
String callerOwner) {
|
||||
if (isBlank(sessionId) || isBlank(content)) {
|
||||
return error("sessionId and content are required");
|
||||
}
|
||||
@@ -919,7 +927,7 @@ public final class FleetMcp {
|
||||
}
|
||||
long timeout = clamp(timeoutMs == null ? DEFAULT_TIMEOUT_MS : timeoutMs);
|
||||
try {
|
||||
return formatReply(messages.send(sessionId, content, timeout, onAccepted, callerTerminal), timeout);
|
||||
return formatReply(messages.send(sessionId, content, timeout, onAccepted, callerOwner), timeout);
|
||||
} catch (HerdrException e) {
|
||||
return error("herdr error contacting session " + sessionId + ": " + e.getMessage());
|
||||
}
|
||||
@@ -929,15 +937,15 @@ public final class FleetMcp {
|
||||
* {@code fleet_send} carrying a {@code turnId}: the primary's answer to a worker's
|
||||
* {@code fleet_ask} (CB-205). Resolves the worker's blocked question and blocks for its reply as
|
||||
* it resumes the same turn — surfaced to the primary identically to a normal send.
|
||||
* {@code callerTerminal} must match the turn's recorded owner or this is refused.
|
||||
* {@code callerOwner} must match the turn's recorded owner or this is refused.
|
||||
*/
|
||||
static McpSchema.CallToolResult answer(MessageService messages, String turnId, String content, Long timeoutMs,
|
||||
String callerTerminal) {
|
||||
String callerOwner) {
|
||||
if (isBlank(turnId) || isBlank(content)) {
|
||||
return error("turnId and content are required to answer a worker's question");
|
||||
}
|
||||
long timeout = clamp(timeoutMs == null ? DEFAULT_TIMEOUT_MS : timeoutMs);
|
||||
return formatReply(messages.answer(turnId, content, timeout, callerTerminal), timeout);
|
||||
return formatReply(messages.answer(turnId, content, timeout, callerOwner), timeout);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1013,13 +1021,12 @@ public final class FleetMcp {
|
||||
}
|
||||
|
||||
/**
|
||||
* As above, recording {@code creatorTerminal} as this ticket's owner (the caller's own
|
||||
* terminal, resolved from the connection) so a later {@code fleet_poll{ticket}} only hands the
|
||||
* result back to that same caller — see {@link MessageService#poll(String, String)}.
|
||||
* As above, recording {@code creator}'s owner key so a later {@code fleet_poll{ticket}} only
|
||||
* hands the result back to the same caller — see {@link MessageService#poll(String, String)}.
|
||||
*/
|
||||
static McpSchema.CallToolResult sendAsync(MessageService messages, String sessionId, String content,
|
||||
Runnable onAccepted, Set<String> profiles,
|
||||
String creatorTerminal) {
|
||||
Runnable onAccepted, Set<String> profiles,
|
||||
Principal creator) {
|
||||
if (isBlank(sessionId) || isBlank(content)) {
|
||||
return error("sessionId and content are required");
|
||||
}
|
||||
@@ -1027,7 +1034,7 @@ public final class FleetMcp {
|
||||
if (targetError != null) {
|
||||
return targetError;
|
||||
}
|
||||
String ticket = messages.sendAsync(sessionId, content, onAccepted, creatorTerminal);
|
||||
String ticket = messages.sendAsync(sessionId, content, onAccepted, creator);
|
||||
return text("accepted — task delegated. Poll fleet_poll with ticket=" + ticket);
|
||||
}
|
||||
|
||||
@@ -1217,13 +1224,12 @@ public final class FleetMcp {
|
||||
}
|
||||
|
||||
/**
|
||||
* As above, refusing a ticket lookup whose caller's terminal differs from the terminal that
|
||||
* created it — see {@link MessageService#poll(String, String)}. {@code callerTerminal} is the
|
||||
* CALLING session's terminal id, resolved by the MCP layer from the connection, never a
|
||||
* client-supplied value.
|
||||
* As above, refusing a ticket lookup whose caller owner key differs from the key that created it
|
||||
* — see {@link MessageService#poll(String, String)}. {@code callerOwner} comes from the calling
|
||||
* connection's resolved principal, never a client-supplied value.
|
||||
*/
|
||||
static McpSchema.CallToolResult poll(MessageService messages, LeadChannel leadChannel, String ticket,
|
||||
String target, String coordId, String callerTerminal) {
|
||||
String target, String coordId, String callerOwner) {
|
||||
if (!isBlank(coordId)) {
|
||||
return pollHeldPeerMail(leadChannel, coordId);
|
||||
}
|
||||
@@ -1237,7 +1243,7 @@ public final class FleetMcp {
|
||||
if (isBlank(ticket)) {
|
||||
return error("ticket (or target) is required");
|
||||
}
|
||||
MessageService.TaskView v = messages.poll(ticket, callerTerminal);
|
||||
MessageService.TaskView v = messages.poll(ticket, callerOwner);
|
||||
if (v == null) {
|
||||
return error("unknown ticket: " + ticket + " (never issued, or expired)");
|
||||
}
|
||||
@@ -1346,18 +1352,17 @@ public final class FleetMcp {
|
||||
* {@code fleet_status}: the live lifecycle status of a worker session, plus — when the worker
|
||||
* is paused mid-turn in an async {@code fleet_ask} — the open question and how to answer it, so
|
||||
* a lead on its normal poll cadence does not need the ticket to notice. The question, its
|
||||
* {@code turnId} and its ticket id are shown only to the caller whose terminal created that
|
||||
* delegation, or to a caller with no terminal at all (the unnamed primary); any other caller
|
||||
* still sees the base status. {@code callerTerminal} is the CALLING session's terminal id,
|
||||
* resolved by the MCP layer from the connection, never a client-supplied value.
|
||||
* {@code turnId} and its ticket id are shown only to the caller whose owner key created that
|
||||
* delegation, or to the unnamed primary; any other caller still sees the base status.
|
||||
* {@code callerOwner} comes from the calling connection's resolved principal.
|
||||
*/
|
||||
static McpSchema.CallToolResult status(MessageService messages, String sessionId, String callerTerminal) {
|
||||
static McpSchema.CallToolResult status(MessageService messages, String sessionId, String callerOwner) {
|
||||
if (isBlank(sessionId)) {
|
||||
return error("sessionId is required");
|
||||
}
|
||||
try {
|
||||
String base = messages.status(sessionId).name().toLowerCase();
|
||||
MessageService.PendingAsk ask = messages.pendingAsk(sessionId, callerTerminal);
|
||||
MessageService.PendingAsk ask = messages.pendingAsk(sessionId, callerOwner);
|
||||
if (ask == null) {
|
||||
return text(base);
|
||||
}
|
||||
@@ -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
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
package dev.ltms.fleet.msg;
|
||||
|
||||
import dev.ltms.fleet.auth.Principal;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.AgentStatus;
|
||||
import dev.ltms.fleet.herdr.HerdrRouter;
|
||||
@@ -282,18 +283,15 @@ public final class MessageService {
|
||||
*/
|
||||
private volatile boolean askTimedOut;
|
||||
/**
|
||||
* The terminal of the caller whose {@code fleet_send{wait:false}} created this ticket, or
|
||||
* {@code null} when that caller had no terminal (the unnamed primary) or the ticket was
|
||||
* created through an overload that does not record one. {@link #poll(String, String)}
|
||||
* compares a polling caller's own terminal against this field before handing back the
|
||||
* ticket's state.
|
||||
* The owner key of the caller whose {@code fleet_send{wait:false}} created this ticket, or
|
||||
* {@code null} for the unnamed primary and overloads that do not record a caller.
|
||||
*/
|
||||
private final String creatorTerminal;
|
||||
private final String creatorOwner;
|
||||
|
||||
private Task(String ticket, String target, LongSupplier nowNanos, String creatorTerminal) {
|
||||
private Task(String ticket, String target, LongSupplier nowNanos, String creatorOwner) {
|
||||
this.ticket = ticket;
|
||||
this.target = target;
|
||||
this.creatorTerminal = creatorTerminal;
|
||||
this.creatorOwner = creatorOwner;
|
||||
this.createdNanos = nowNanos.getAsLong();
|
||||
future.whenComplete((reply, ex) -> completedNanos = nowNanos.getAsLong());
|
||||
}
|
||||
@@ -357,7 +355,7 @@ public final class MessageService {
|
||||
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
|
||||
* {@link #sendAsync(String, String, Runnable, Principal)}). {@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.
|
||||
@@ -938,13 +936,13 @@ public final class MessageService {
|
||||
|
||||
/**
|
||||
* Deliver {@code content} to {@code target} (a herdr {@code terminal_id}) and block until the
|
||||
* worker replies via {@link Rendezvous} or {@code timeoutMillis} elapses. {@code callerTerminal}
|
||||
* is the terminal of the caller making this call — {@code null} for the unnamed primary — and is
|
||||
* recorded as the turn's owner, the only caller {@link #answer(String, String, long, String)} will
|
||||
* worker replies via {@link Rendezvous} or {@code timeoutMillis} elapses. {@code callerOwner}
|
||||
* identifies the caller making this call and is recorded as the turn's owner. It is the only
|
||||
* caller {@link #answer(String, String, long, String)} will
|
||||
* later accept an answer from if the worker pauses mid-turn to ask.
|
||||
*/
|
||||
public Reply send(String target, String content, long timeoutMillis, String callerTerminal) {
|
||||
return send(target, content, timeoutMillis, null, callerTerminal);
|
||||
public Reply send(String target, String content, long timeoutMillis, String callerOwner) {
|
||||
return send(target, content, timeoutMillis, null, callerOwner);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -959,13 +957,13 @@ public final class MessageService {
|
||||
* acceptance means a concurrent sender that times out {@code BUSY} can never steal ownership it
|
||||
* never earned. {@code null} disables the hook.
|
||||
*/
|
||||
public Reply send(String target, String content, long timeoutMillis, Runnable onAccepted, String callerTerminal) {
|
||||
return send(target, content, timeoutMillis, onAccepted, null, callerTerminal);
|
||||
public Reply send(String target, String content, long timeoutMillis, Runnable onAccepted, String callerOwner) {
|
||||
return send(target, content, timeoutMillis, onAccepted, null, callerOwner);
|
||||
}
|
||||
|
||||
/** Run a send, optionally stopping an async task that teardown already failed before acceptance. */
|
||||
private Reply send(String target, String content, long timeoutMillis, Runnable onAccepted, Task task,
|
||||
String callerTerminal) {
|
||||
String callerOwner) {
|
||||
long deadlineNanos = System.nanoTime() + timeoutMillis * 1_000_000L;
|
||||
ReentrantLock lock = sessionLocks.computeIfAbsent(target, _ -> new ReentrantLock());
|
||||
|
||||
@@ -985,7 +983,7 @@ public final class MessageService {
|
||||
// open race). Opening first also means a throwing onAccepted (fired before enqueue) or an
|
||||
// enqueue failure is safely closed by the finally below: nothing is left queued, and the
|
||||
// failed send leaves no stale waiter behind.
|
||||
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(target, Rendezvous.Owner.of(callerTerminal));
|
||||
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(target, Rendezvous.Owner.of(callerOwner));
|
||||
// CB-640: this send now owns target's delivery, so any earlier stranded-reply or
|
||||
// still-queued fact no longer describes the live state — clear both rather than let
|
||||
// them outlive the send that supersedes them.
|
||||
@@ -1165,20 +1163,20 @@ public final class MessageService {
|
||||
* call, not a new status-gated delivery. The forward waiter is opened <em>before</em> the worker
|
||||
* is unblocked so a reply that lands the instant it resumes is not lost.
|
||||
*
|
||||
* <p>{@code callerTerminal} is the terminal of the caller making this call — {@code null} for
|
||||
* the unnamed primary. It is checked against the turn's recorded owner (the caller whose
|
||||
* <p>{@code callerOwner} identifies the caller making this call. It is checked against the
|
||||
* turn's recorded owner (the caller whose
|
||||
* accepted delegation opened it, see {@link #send(String, String, long, String)} and
|
||||
* {@link #sendAsync(String, String, Runnable, String)}) before anything else runs: a mismatch,
|
||||
* {@link #sendAsync(String, String, Runnable, Principal)}) before anything else runs: a mismatch,
|
||||
* including a turn with no owner on record at all, returns {@link Outcome#NOT_TURN_OWNER}
|
||||
* without touching the rendezvous, the session lock, or any async task bookkeeping.
|
||||
*/
|
||||
public Reply answer(String turnId, String content, long timeoutMillis, String callerTerminal) {
|
||||
public Reply answer(String turnId, String content, long timeoutMillis, String callerOwner) {
|
||||
String workerSession = rendezvous.askSession(turnId);
|
||||
if (workerSession == null) {
|
||||
return new Reply(Outcome.STALE_TURN, null); // the ask lapsed (timed out or already answered)
|
||||
}
|
||||
Rendezvous.Owner owner = rendezvous.askOwner(turnId);
|
||||
if (!Rendezvous.Owner.permits(owner, callerTerminal)) {
|
||||
if (!Rendezvous.Owner.permits(owner, callerOwner)) {
|
||||
return new Reply(Outcome.NOT_TURN_OWNER, null);
|
||||
}
|
||||
long deadlineNanos = System.nanoTime() + timeoutMillis * 1_000_000L;
|
||||
@@ -1321,16 +1319,16 @@ public final class MessageService {
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #sendAsync(String, String, Runnable)}, recording {@code creatorTerminal} as this
|
||||
* ticket's owner. {@link #poll(String, String)} refuses a later caller whose own terminal
|
||||
* differs from this one; {@code null} records no owner (a caller with no terminal — the
|
||||
* unnamed primary — is always allowed to poll the result regardless).
|
||||
* As {@link #sendAsync(String, String, Runnable)}, recording {@code creator}'s owner key as this
|
||||
* ticket's owner. The key is derived here from the resolved principal so callers cannot pass a
|
||||
* terminal address where an owner identity is required.
|
||||
*
|
||||
* @return the ticket to poll for the eventual result
|
||||
*/
|
||||
public String sendAsync(String target, String content, Runnable onAccepted, String creatorTerminal) {
|
||||
public String sendAsync(String target, String content, Runnable onAccepted, Principal creator) {
|
||||
String ticket = "task-" + ticketBootNonce + "-" + ticketSeq.incrementAndGet();
|
||||
Task task = new Task(ticket, target, nowNanos, creatorTerminal);
|
||||
String creatorOwner = creator == null ? null : creator.ownerKey();
|
||||
Task task = new Task(ticket, target, nowNanos, creatorOwner);
|
||||
tasks.put(ticket, task);
|
||||
if (pushLoop != null) {
|
||||
// CB-588: task.future only ever completes on a terminal phase (DONE or a failure) — a
|
||||
@@ -1361,7 +1359,7 @@ public final class MessageService {
|
||||
}
|
||||
asyncExecutor.submit(() -> {
|
||||
try {
|
||||
Reply result = send(target, content, ASYNC_TIMEOUT_MS, onAccepted, task, creatorTerminal);
|
||||
Reply result = send(target, content, ASYNC_TIMEOUT_MS, onAccepted, task, creatorOwner);
|
||||
if (result.outcome() == Outcome.QUESTION) {
|
||||
// Keep the accepted owner until answer() finishes it. markAsyncQuestion may run
|
||||
// just after resolveQuestion wakes this thread.
|
||||
@@ -1393,9 +1391,9 @@ public final class MessageService {
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #poll(String, String)}, with no caller terminal — the ticket's ownership is never
|
||||
* As {@link #poll(String, String)}, with no caller owner key — the unnamed primary's ticket rule
|
||||
* checked, so this overload must only be used where the caller's identity is otherwise
|
||||
* irrelevant (a test, or a surface that does not resolve a caller terminal at all).
|
||||
* irrelevant.
|
||||
*/
|
||||
public TaskView poll(String ticket) {
|
||||
return poll(ticket, null);
|
||||
@@ -1403,18 +1401,18 @@ public final class MessageService {
|
||||
|
||||
/**
|
||||
* Snapshot the state of an async delegation. Returns {@code null} for an unknown/expired ticket.
|
||||
* Refuses a {@code callerTerminal} that differs from the terminal that created the ticket (see
|
||||
* {@link #sendAsync(String, String, Runnable, String)}) with a {@link Phase#FAILED} view that
|
||||
* carries no reply text — a caller with no terminal (the unnamed primary) is never refused.
|
||||
* Refuses a {@code callerOwner} that differs from the owner that created the ticket (see
|
||||
* {@link #sendAsync(String, String, Runnable, Principal)}) with a {@link Phase#FAILED} view that
|
||||
* carries no reply text. The unnamed primary has a {@code null} owner key and is never refused.
|
||||
* Otherwise returns a {@link Phase#PENDING} view (with the live worker status as detail), a
|
||||
* {@link Phase#DONE} view carrying the reply, or a {@link Phase#FAILED} view with the reason.
|
||||
*/
|
||||
public TaskView poll(String ticket, String callerTerminal) {
|
||||
public TaskView poll(String ticket, String callerOwner) {
|
||||
Task task = tasks.get(ticket);
|
||||
if (task == null) {
|
||||
return null;
|
||||
}
|
||||
if (!ownsTicket(task, callerTerminal)) {
|
||||
if (!ownsTicket(task, callerOwner)) {
|
||||
return new TaskView(ticket, Phase.FAILED, null, null,
|
||||
"forbidden: this ticket was created by a different session", null);
|
||||
}
|
||||
@@ -1454,14 +1452,13 @@ public final class MessageService {
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether {@code callerTerminal} may read {@code task}'s state. A caller with no terminal
|
||||
* always may — that is the unnamed primary, resolved by token or loopback trust, which never
|
||||
* carries a herdr pane and must keep reading every ticket. Otherwise the caller's terminal must
|
||||
* equal the terminal recorded on the task; a task with no recorded terminal matches no
|
||||
* terminal-bearing caller.
|
||||
* Whether {@code callerOwner} may read {@code task}'s state. A {@code null} caller key is the
|
||||
* unnamed primary and may read every ticket. Other callers must match the task's owner key. This
|
||||
* differs from {@link Rendezvous.Owner#permits}: a missing rendezvous owner is not an authenticated
|
||||
* unnamed primary, so that gate refuses every caller when no owner was recorded.
|
||||
*/
|
||||
private static boolean ownsTicket(Task task, String callerTerminal) {
|
||||
return callerTerminal == null || callerTerminal.equals(task.creatorTerminal);
|
||||
private static boolean ownsTicket(Task task, String callerOwner) {
|
||||
return callerOwner == null || callerOwner.equals(task.creatorOwner);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1787,13 +1784,13 @@ public final class MessageService {
|
||||
* {@code fleet_status} uses this to show a pending question without the caller needing the
|
||||
* ticket. {@code null} when the session has no open async question (including a session mid a
|
||||
* <em>blocking</em> {@code fleet_ask}, which has no {@link Task} to look up — see
|
||||
* {@link PendingAsk}), or when {@code callerTerminal} does not own the task the question
|
||||
* {@link PendingAsk}), or when {@code callerOwner} does not own the task the question
|
||||
* belongs to (see {@link #ownsTicket(Task, String)}).
|
||||
*/
|
||||
public PendingAsk pendingAsk(String workerSession, String callerTerminal) {
|
||||
public PendingAsk pendingAsk(String workerSession, String callerOwner) {
|
||||
for (Task task : tasks.values()) {
|
||||
Reply q = task.question;
|
||||
if (q != null && workerSession.equals(task.target) && ownsTicket(task, callerTerminal)) {
|
||||
if (q != null && workerSession.equals(task.target) && ownsTicket(task, callerOwner)) {
|
||||
return new PendingAsk(task.ticket, q.text(), q.turnId());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -73,23 +73,22 @@ public final class Rendezvous {
|
||||
|
||||
/**
|
||||
* The caller whose accepted delegation opened a turn — the only caller allowed to answer it.
|
||||
* A {@code null} terminal means the unnamed primary, an authenticated caller with no pane.
|
||||
* A {@code null} owner key means the unnamed primary.
|
||||
*/
|
||||
public record Owner(String terminal) {
|
||||
public record Owner(String ownerKey) {
|
||||
public static final Owner UNNAMED_PRIMARY = new Owner(null);
|
||||
|
||||
public static Owner of(String terminal) {
|
||||
return terminal == null ? UNNAMED_PRIMARY : new Owner(terminal);
|
||||
public static Owner of(String ownerKey) {
|
||||
return ownerKey == null ? UNNAMED_PRIMARY : new Owner(ownerKey);
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether {@code callerTerminal} matches {@code owner}. A {@code null} owner means no
|
||||
* owner was ever recorded, and that state matches no caller, not even one whose own
|
||||
* terminal is {@code null} — "no record" and "recorded as the unnamed primary" are
|
||||
* different states.
|
||||
* Whether {@code callerOwner} matches {@code owner}. A {@code null} owner means no owner was
|
||||
* recorded, so it matches no caller. {@link #UNNAMED_PRIMARY} records the unnamed primary
|
||||
* with an owner object whose key is {@code null}.
|
||||
*/
|
||||
public static boolean permits(Owner owner, String callerTerminal) {
|
||||
return owner != null && java.util.Objects.equals(owner.terminal(), callerTerminal);
|
||||
public static boolean permits(Owner owner, String callerOwner) {
|
||||
return owner != null && java.util.Objects.equals(owner.ownerKey(), callerOwner);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -664,7 +664,7 @@ public final class FleetApp {
|
||||
return;
|
||||
}
|
||||
Principal caller = ctx.attribute(CALLER);
|
||||
String callerTerminal = caller == null ? null : caller.terminal();
|
||||
String callerOwner = caller == null ? null : caller.ownerKey();
|
||||
JsonNode body;
|
||||
try {
|
||||
body = mapper.readTree(ctx.body());
|
||||
@@ -690,19 +690,19 @@ public final class FleetApp {
|
||||
|
||||
// Answering a worker's fleet_ask (CB-205): always blocks, and derives the worker from turnId.
|
||||
if (turnId != null && !turnId.isBlank()) {
|
||||
writeReply(ctx, id, messages.answer(turnId, content, timeout, callerTerminal), timeout);
|
||||
writeReply(ctx, id, messages.answer(turnId, content, timeout, callerOwner), timeout);
|
||||
return;
|
||||
}
|
||||
|
||||
if (!wait) {
|
||||
// Fire-and-poll (CB-107): return a ticket immediately; the caller polls GET /tasks/{ticket}.
|
||||
String ticket = messages.sendAsync(id, content, null, callerTerminal);
|
||||
String ticket = messages.sendAsync(id, content, null, caller);
|
||||
ctx.status(202).json(Map.of("sessionId", id, "ticket", ticket, "status", "accepted"));
|
||||
return;
|
||||
}
|
||||
|
||||
try {
|
||||
writeReply(ctx, id, messages.send(id, content, timeout, callerTerminal), timeout);
|
||||
writeReply(ctx, id, messages.send(id, content, timeout, callerOwner), timeout);
|
||||
} catch (HerdrException e) {
|
||||
herdrError(ctx, e);
|
||||
}
|
||||
@@ -883,10 +883,10 @@ public final class FleetApp {
|
||||
body.put("ready", deliverable.test(id));
|
||||
// A worker paused mid-turn in an async fleet_ask is otherwise invisible to a status
|
||||
// poll — surface the open question and how to answer it, same as fleet_poll's
|
||||
// Phase.ASKING view, but only to the caller whose terminal created that delegation, or
|
||||
// to a caller with no terminal at all (the unnamed primary).
|
||||
// Phase.ASKING view, but only to the caller whose owner key created that delegation, or
|
||||
// to the unnamed primary.
|
||||
Principal caller = ctx.attribute(CALLER);
|
||||
MessageService.PendingAsk ask = messages.pendingAsk(id, caller == null ? null : caller.terminal());
|
||||
MessageService.PendingAsk ask = messages.pendingAsk(id, caller == null ? null : caller.ownerKey());
|
||||
if (ask != null) {
|
||||
body.put("question", ask.question());
|
||||
body.put("turnId", ask.turnId());
|
||||
@@ -904,7 +904,7 @@ public final class FleetApp {
|
||||
return;
|
||||
}
|
||||
Principal caller = ctx.attribute(CALLER);
|
||||
MessageService.TaskView v = messages.poll(ctx.pathParam("ticket"), caller == null ? null : caller.terminal());
|
||||
MessageService.TaskView v = messages.poll(ctx.pathParam("ticket"), caller == null ? null : caller.ownerKey());
|
||||
if (v == null) {
|
||||
ctx.status(404).json(Map.of("error", "unknown_ticket", "detail", "no such task (or it has expired)"));
|
||||
return;
|
||||
|
||||
@@ -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());
|
||||
@@ -466,6 +470,12 @@ public final class SessionManager implements TurnListener {
|
||||
MemberSession resolved = resolveAgentSessionId(removed, removedHandle);
|
||||
notifyReleased(new ReleaseDetail(resolved.terminalId(), resolved.worktree(),
|
||||
resolved.branch(), snapshotRef, resolved.agentSessionId()));
|
||||
String terminal = removed.terminalId();
|
||||
if (terminal != null && !terminal.isBlank()) {
|
||||
// Without this, a terminal stays marked present after its pane is gone, so a
|
||||
// later send to the same id would read as deliverable instead of refused.
|
||||
presence.forget(terminal);
|
||||
}
|
||||
}
|
||||
}
|
||||
// CB-581: the pane must always stop, even if the dirty check above threw. A session removed
|
||||
@@ -809,6 +819,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 +997,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");
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,31 @@
|
||||
package dev.ltms.fleet.auth;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertNull;
|
||||
|
||||
class PrincipalTest {
|
||||
|
||||
@Test
|
||||
void ownerKeyCoversEveryRole() {
|
||||
assertEquals("leader:opus", Principal.leader("opus", "term_lead", 1).ownerKey());
|
||||
assertNull(Principal.primary(2).ownerKey());
|
||||
assertEquals("worker:term_worker", Principal.worker("term_worker", 3).ownerKey());
|
||||
assertEquals("architect:term_arch", Principal.architect("opus", "term_arch", 4).ownerKey());
|
||||
assertEquals("collaborator:ops", Principal.collaborator("ops", "term_collab", 5).ownerKey());
|
||||
assertEquals("observer:term_observer", Principal.observer("term_observer", 6).ownerKey());
|
||||
assertEquals("anonymous", Principal.anonymous().ownerKey());
|
||||
}
|
||||
|
||||
@Test
|
||||
void rolePrefixesKeepLeadAndArchitectKeysDistinct() {
|
||||
String lead = Principal.leader("opus", "term_lead", 1).ownerKey();
|
||||
String architect = Principal.architect("design", "opus", 2).ownerKey();
|
||||
|
||||
assertEquals("leader:opus", lead);
|
||||
assertEquals("architect:opus", architect);
|
||||
assertNotEquals(lead, architect);
|
||||
}
|
||||
}
|
||||
@@ -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());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -504,12 +504,12 @@ class FleetMcpAuthzTest {
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code fleet_poll{ticket}} must thread the calling connection's own terminal into
|
||||
* {@code fleet_poll{ticket}} must thread the calling connection's owner key into
|
||||
* {@link MessageService#poll(String, String)}, so a worker cannot read a ticket a different
|
||||
* session created.
|
||||
*/
|
||||
@Test
|
||||
void theFleetPollHandlerActuallyThreadsCallerTerminalIntoPoll() throws Exception {
|
||||
void theFleetPollHandlerActuallyThreadsCallerOwnerIntoPoll() throws Exception {
|
||||
String source = Files.readString(MCP_SOURCE);
|
||||
|
||||
int start = source.indexOf("pollHandler =");
|
||||
@@ -525,18 +525,18 @@ class FleetMcpAuthzTest {
|
||||
"control failed: the scraped pollHandler block contains no poll(messages, ...) call "
|
||||
+ "at all -- the anchors have drifted, this test is not testing what it claims to");
|
||||
|
||||
assertTrue(handlerBlock.contains("callerTerminal(exchange)"),
|
||||
"the fleet_poll handler must thread callerTerminal(exchange) into poll(...), not omit "
|
||||
assertTrue(handlerBlock.contains("principal(exchange).ownerKey()"),
|
||||
"the fleet_poll handler must thread principal(exchange).ownerKey() into poll(...), not omit "
|
||||
+ "it or pass a literal null -- block: " + handlerBlock);
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code fleet_status} must thread the calling connection's own terminal into
|
||||
* {@code fleet_status} must thread the calling connection's owner key into
|
||||
* {@link FleetMcp#status(MessageService, String, String)}, so a caller that did not create a
|
||||
* worker's open delegation cannot read its pending question through the status handler either.
|
||||
*/
|
||||
@Test
|
||||
void theFleetStatusHandlerActuallyThreadsCallerTerminalIntoStatus() throws Exception {
|
||||
void theFleetStatusHandlerActuallyThreadsCallerOwnerIntoStatus() throws Exception {
|
||||
String source = Files.readString(MCP_SOURCE);
|
||||
|
||||
int start = source.indexOf("statusHandler =");
|
||||
@@ -552,8 +552,8 @@ class FleetMcpAuthzTest {
|
||||
"control failed: the scraped statusHandler block contains no status(messages, ...) "
|
||||
+ "call at all -- the anchors have drifted, this test is not testing what it claims to");
|
||||
|
||||
assertTrue(handlerBlock.contains("callerTerminal(exchange)"),
|
||||
"the fleet_status handler must thread callerTerminal(exchange) into status(...), not "
|
||||
assertTrue(handlerBlock.contains("principal(exchange).ownerKey()"),
|
||||
"the fleet_status handler must thread principal(exchange).ownerKey() into status(...), not "
|
||||
+ "omit it or pass a literal null -- block: " + handlerBlock);
|
||||
}
|
||||
|
||||
|
||||
@@ -329,12 +329,13 @@ class FleetMcpTest {
|
||||
|
||||
/**
|
||||
* Same hijack and control through the fire-and-poll ({@code sendAsync}) path: the owner comes
|
||||
* from the ticket's recorded creator terminal, not from a caller threaded through a live call.
|
||||
* from the resolved principal, not from a caller argument threaded through a live call.
|
||||
*/
|
||||
@Test
|
||||
void aDifferentCallersMcpAnswerIsRefusedForAnAsyncSendButTheRealOwnerSucceeds() throws Exception {
|
||||
Principal owner = Principal.worker("term_owner", 1);
|
||||
McpSchema.CallToolResult accepted =
|
||||
FleetMcp.sendAsync(messages, T, "do it", null, Set.of(), "term_owner");
|
||||
FleetMcp.sendAsync(messages, T, "do it", null, Set.of(), owner);
|
||||
String ticket = textOf(accepted).substring(textOf(accepted).indexOf("ticket=") + "ticket=".length()).trim();
|
||||
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
@@ -354,14 +355,15 @@ class FleetMcpTest {
|
||||
assertEquals(MessageService.Phase.ASKING, asking.phase());
|
||||
String turnId = asking.turnId();
|
||||
|
||||
McpSchema.CallToolResult hijacked = FleetMcp.answer(messages, turnId, "evil.yaml", 500L, "term_attacker");
|
||||
McpSchema.CallToolResult hijacked = FleetMcp.answer(messages, turnId, "evil.yaml", 500L,
|
||||
"worker:term_attacker");
|
||||
assertTrue(hijacked.isError(), "a caller that did not create this delegation must get an error");
|
||||
assertFalse(ask.isDone(), "a refused answer must not resolve the worker's blocked fleet_ask");
|
||||
assertEquals(MessageService.Phase.ASKING, messages.poll(ticket).phase(),
|
||||
"a refused answer must not advance the async ticket's phase");
|
||||
|
||||
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
|
||||
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, "term_owner"));
|
||||
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, owner.ownerKey()));
|
||||
assertEquals("config.yaml", textOf(ask.get(5, TimeUnit.SECONDS)));
|
||||
deadline = System.currentTimeMillis() + 3000;
|
||||
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
|
||||
@@ -1920,8 +1922,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 +1934,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");
|
||||
@@ -1996,13 +2013,13 @@ class FleetMcpTest {
|
||||
|
||||
/**
|
||||
* {@code fleet_status}'s pending-ask block (the question, its {@code turnId} and its ticket)
|
||||
* is shown only to the caller whose terminal created the delegation, or to a caller with no
|
||||
* terminal at all (the unnamed primary) — a different terminal-bearing caller still sees the
|
||||
* base status line, but none of the pending-ask fields.
|
||||
* is shown only to the caller whose owner key created the delegation, or to the unnamed primary.
|
||||
* A different caller still sees the base status line, but none of the pending-ask fields.
|
||||
*/
|
||||
@Test
|
||||
void statusGatesThePendingAskFieldsByTheDelegationsCreatorTerminal() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "task that asks", null, "term_creator");
|
||||
void statusGatesThePendingAskFieldsByTheDelegationsCreatorOwner() throws Exception {
|
||||
Principal creator = Principal.worker("term_creator", 1);
|
||||
String ticket = messages.sendAsync(T, "task that asks", null, creator);
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
|
||||
Thread.sleep(5);
|
||||
@@ -2020,7 +2037,7 @@ class FleetMcpTest {
|
||||
} while (asking.phase() != MessageService.Phase.ASKING && System.currentTimeMillis() < deadline);
|
||||
assertEquals(MessageService.Phase.ASKING, asking.phase());
|
||||
|
||||
String other = textOf(FleetMcp.status(messages, T, "term_other"));
|
||||
String other = textOf(FleetMcp.status(messages, T, "worker:term_other"));
|
||||
assertTrue(other.startsWith("idle"), "the base status must still be shown: " + other);
|
||||
assertFalse(other.contains("which config file?"),
|
||||
"a non-creating caller must not see the question text: " + other);
|
||||
@@ -2029,10 +2046,12 @@ class FleetMcpTest {
|
||||
assertFalse(other.contains(ticket),
|
||||
"a non-creating caller must not see the ticket: " + other);
|
||||
|
||||
String creator = textOf(FleetMcp.status(messages, T, "term_creator"));
|
||||
assertTrue(creator.contains("which config file?"), "the creator must see the question: " + creator);
|
||||
assertTrue(creator.contains(asking.turnId()), "the creator must see the turnId: " + creator);
|
||||
assertTrue(creator.contains(ticket), "the creator must see the ticket: " + creator);
|
||||
String creatorStatus = textOf(FleetMcp.status(messages, T, creator.ownerKey()));
|
||||
assertTrue(creatorStatus.contains("which config file?"),
|
||||
"the creator must see the question: " + creatorStatus);
|
||||
assertTrue(creatorStatus.contains(asking.turnId()),
|
||||
"the creator must see the turnId: " + creatorStatus);
|
||||
assertTrue(creatorStatus.contains(ticket), "the creator must see the ticket: " + creatorStatus);
|
||||
|
||||
String unnamed = textOf(FleetMcp.status(messages, T, null));
|
||||
assertTrue(unnamed.contains("which config file?"),
|
||||
@@ -2041,7 +2060,7 @@ class FleetMcpTest {
|
||||
// Clean up the still-open ask so the background thread does not linger past the test.
|
||||
String turnId = asking.turnId();
|
||||
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
|
||||
() -> messages.answer(turnId, "config.yaml", 5000, "term_creator"));
|
||||
() -> messages.answer(turnId, "config.yaml", 5000, creator.ownerKey()));
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
deadline = System.currentTimeMillis() + 3000;
|
||||
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
|
||||
@@ -2149,6 +2168,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
|
||||
|
||||
@@ -19,7 +19,7 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
* {@link MessageService#poll(String)} overload. That overload skips the ownership check in
|
||||
* {@code MessageService}'s {@code ownsTicket} entirely, so a caller of it can read any session's
|
||||
* ticket. Every production caller must go through {@link MessageService#poll(String, String)}
|
||||
* and pass a {@code callerTerminal} explicitly, even when it is {@code null}.
|
||||
* and pass a {@code callerOwner} explicitly, even when it is {@code null}.
|
||||
*
|
||||
* <p>This reads each file's own source text rather than reflecting on compiled bytecode, because
|
||||
* the risk is a future one-word edit at a call site, not a missing overload.
|
||||
@@ -48,7 +48,7 @@ class MessageServicePollUsageTest {
|
||||
+ "below proves nothing");
|
||||
|
||||
assertTrue(violations.isEmpty(), "found a call to the fail-open MessageService.poll(String) "
|
||||
+ "overload, which skips the ownership check entirely -- pass a callerTerminal "
|
||||
+ "overload, which skips the ownership check entirely -- pass a callerOwner "
|
||||
+ "explicitly (even if null) through poll(String, String) instead: " + violations);
|
||||
|
||||
// CONTROL: the arity parser actually finds the two genuine two-argument call sites (the
|
||||
@@ -58,7 +58,7 @@ class MessageServicePollUsageTest {
|
||||
assertEquals(2, twoArgSites.size(), "control failed: expected exactly the two known "
|
||||
+ "two-argument messages.poll(...) call sites, found: " + twoArgSites);
|
||||
assertTrue(twoArgSites.stream().anyMatch(s -> s.contains("FleetMcp.java")),
|
||||
"control failed: did not find the FleetMcp.java messages.poll(ticket, callerTerminal) "
|
||||
"control failed: did not find the FleetMcp.java messages.poll(ticket, callerOwner) "
|
||||
+ "site among: " + twoArgSites);
|
||||
assertTrue(twoArgSites.stream().anyMatch(s -> s.contains("FleetApp.java")),
|
||||
"control failed: did not find the FleetApp.java messages.poll(...) site among: "
|
||||
|
||||
@@ -4,6 +4,7 @@ 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.auth.Principal;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.AgentStatus;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
@@ -921,24 +922,26 @@ class MessageServiceTest {
|
||||
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 -------------------------
|
||||
// --- ticket ownership -----------------------------------------------------------------------
|
||||
|
||||
@Test
|
||||
void pollByAnotherTerminalIsRefused() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "long task", null, "term_a");
|
||||
void leadBIsRefusedFromLeadAsTicket() throws Exception {
|
||||
Principal leadA = Principal.leader("opus", "term_a", 1);
|
||||
Principal leadB = Principal.leader("sol", "term_b", 2);
|
||||
String ticket = messages.sendAsync(T, "long task", null, leadA);
|
||||
awaitWaiting();
|
||||
injector.onStatus(T, AgentStatus.IDLE); // deliver
|
||||
injector.onStatus(T, AgentStatus.WORKING); // worker works
|
||||
assertTrue(rendezvous.resolve(T, "secret async result"), "a reply resolves the async send");
|
||||
|
||||
MessageService.TaskView owner = driveAsyncTicketToDone(ticket, "term_a");
|
||||
MessageService.TaskView owner = driveAsyncTicketToDone(ticket, leadA.ownerKey());
|
||||
assertNotNull(owner, "the creator must still be able to read its own ticket");
|
||||
assertEquals(MessageService.Phase.DONE, owner.phase());
|
||||
|
||||
MessageService.TaskView refused = messages.poll(ticket, "term_b");
|
||||
assertNotNull(refused, "a different terminal gets a refusal, not silence");
|
||||
MessageService.TaskView refused = messages.poll(ticket, leadB.ownerKey());
|
||||
assertNotNull(refused, "a different lead gets a refusal, not silence");
|
||||
assertNotEquals(MessageService.Phase.DONE, refused.phase(),
|
||||
"a different terminal must never see the ticket as DONE");
|
||||
"a different lead must never see the ticket as DONE");
|
||||
assertNull(refused.reply(), "a refusal must never carry the reply text");
|
||||
assertFalse(String.valueOf(refused).contains("secret async result"),
|
||||
"the reply text must not appear anywhere in the refused view");
|
||||
@@ -946,14 +949,13 @@ class MessageServiceTest {
|
||||
|
||||
@Test
|
||||
void unnamedPrimaryStillReadsAnyTicket() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "long task", null, "term_lead");
|
||||
String ticket = messages.sendAsync(T, "long task", null,
|
||||
Principal.leader("opus", "term_lead", 1));
|
||||
awaitWaiting();
|
||||
injector.onStatus(T, AgentStatus.IDLE); // deliver
|
||||
injector.onStatus(T, AgentStatus.WORKING); // worker works
|
||||
assertTrue(rendezvous.resolve(T, "primary-visible result"), "a reply resolves the async send");
|
||||
|
||||
// callerTerminal == null is the unnamed primary (resolved by token or loopback trust, with
|
||||
// no herdr pane) — it must read a ticket a terminal-bearing lead created.
|
||||
MessageService.TaskView view = driveAsyncTicketToDone(ticket, null);
|
||||
assertNotNull(view, "the unnamed primary must be able to read any ticket");
|
||||
assertEquals(MessageService.Phase.DONE, view.phase());
|
||||
@@ -961,19 +963,44 @@ class MessageServiceTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
void creatorReadsItsOwnTicket() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "long task", null, "term_creator");
|
||||
void namedLeadCanPollItsTicketAfterItsTerminalChanges() throws Exception {
|
||||
Principal oldLead = Principal.leader("opus", "term_OLD", 1);
|
||||
Principal newLead = Principal.leader("opus", "term_NEW", 2);
|
||||
assertNotEquals(oldLead.terminal(), newLead.terminal(), "the test requires different terminals");
|
||||
String ticket = messages.sendAsync(T, "long task", null, oldLead);
|
||||
awaitWaiting();
|
||||
injector.onStatus(T, AgentStatus.IDLE); // deliver
|
||||
injector.onStatus(T, AgentStatus.WORKING); // worker works
|
||||
assertTrue(rendezvous.resolve(T, "own result"), "a reply resolves the async send");
|
||||
|
||||
MessageService.TaskView view = driveAsyncTicketToDone(ticket, "term_creator");
|
||||
assertNotNull(view, "the ticket's own creator must be able to read it");
|
||||
MessageService.TaskView view = driveAsyncTicketToDone(ticket, newLead.ownerKey());
|
||||
assertNotNull(view, "the same named lead must read the ticket from its new terminal");
|
||||
assertEquals(MessageService.Phase.DONE, view.phase());
|
||||
assertEquals("own result", view.reply());
|
||||
}
|
||||
|
||||
@Test
|
||||
void anonymousOwnerKeyIsRefusedByTheTicketGateItself() {
|
||||
Principal lead = Principal.leader("opus", "term_lead", 1);
|
||||
String ticket = messages.sendAsync(T, "long task", null, lead);
|
||||
|
||||
MessageService.TaskView refused = messages.poll(ticket, Principal.anonymous().ownerKey());
|
||||
|
||||
assertNotNull(refused);
|
||||
assertEquals(MessageService.Phase.FAILED, refused.phase());
|
||||
assertEquals("forbidden: this ticket was created by a different session", refused.detail());
|
||||
}
|
||||
|
||||
@Test
|
||||
void architectOwnershipUsesTerminalRatherThanSlot() {
|
||||
Principal oldArchitect = Principal.architect("opus", "term_OLD", 1);
|
||||
Principal newArchitect = Principal.architect("opus", "term_NEW", 2);
|
||||
String ticket = messages.sendAsync(T, "long task", null, oldArchitect);
|
||||
|
||||
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket, oldArchitect.ownerKey()).phase());
|
||||
assertEquals(MessageService.Phase.FAILED, messages.poll(ticket, newArchitect.ownerKey()).phase());
|
||||
}
|
||||
|
||||
@Test
|
||||
void pollReportsACompletedTicket() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "long task");
|
||||
@@ -1950,13 +1977,13 @@ class MessageServiceTest {
|
||||
}
|
||||
|
||||
/**
|
||||
* A caller's own terminal must match the terminal that created the delegation to see its
|
||||
* pending question; a different terminal-bearing caller sees nothing, and a caller with no
|
||||
* terminal at all (the unnamed primary) always sees it.
|
||||
* A caller's owner key must match the key that created the delegation to see its pending
|
||||
* question. The unnamed primary always sees it.
|
||||
*/
|
||||
@Test
|
||||
void pendingAskGatesTheQuestionByTheDelegationsCreatorTerminal() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "task that asks", null, "term_creator");
|
||||
void pendingAskGatesTheQuestionByTheDelegationsCreatorOwner() throws Exception {
|
||||
Principal creator = Principal.worker("term_creator", 1);
|
||||
String ticket = messages.sendAsync(T, "task that asks", null, creator);
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
@@ -1964,10 +1991,10 @@ class MessageServiceTest {
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
|
||||
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
|
||||
|
||||
assertNull(messages.pendingAsk(T, "term_other"),
|
||||
"a caller whose terminal did not create the delegation must not see the question");
|
||||
assertNull(messages.pendingAsk(T, "worker:term_other"),
|
||||
"a caller whose key did not create the delegation must not see the question");
|
||||
|
||||
MessageService.PendingAsk own = messages.pendingAsk(T, "term_creator");
|
||||
MessageService.PendingAsk own = messages.pendingAsk(T, creator.ownerKey());
|
||||
assertNotNull(own, "the creating caller must see its own open question");
|
||||
assertEquals("which config file?", own.question());
|
||||
|
||||
@@ -1976,7 +2003,7 @@ class MessageServiceTest {
|
||||
assertEquals("which config file?", unnamed.question());
|
||||
|
||||
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
|
||||
() -> messages.answer(asking.turnId(), "config.yaml", 5000, "term_creator"));
|
||||
() -> messages.answer(asking.turnId(), "config.yaml", 5000, creator.ownerKey()));
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
awaitWaiting();
|
||||
assertTrue(rendezvous.resolve(T, "done"));
|
||||
@@ -1984,13 +2011,12 @@ class MessageServiceTest {
|
||||
}
|
||||
|
||||
/**
|
||||
* A task created with no recorded creator terminal (a short {@code sendAsync} overload) must
|
||||
* not hand its open question to any caller that does have a terminal — only a caller with no
|
||||
* terminal at all may still see it.
|
||||
* A task created with no recorded owner (a short {@code sendAsync} overload) must not hand its
|
||||
* open question to a caller with an owner key. Only the unnamed primary may still see it.
|
||||
*/
|
||||
@Test
|
||||
void pendingAskDeniesATerminalBearingCallerWhenTheTaskRecordsNoCreator() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "task that asks"); // no creatorTerminal recorded
|
||||
String ticket = messages.sendAsync(T, "task that asks");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
@@ -1998,8 +2024,8 @@ class MessageServiceTest {
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
|
||||
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
|
||||
|
||||
assertNull(messages.pendingAsk(T, "term_someone"),
|
||||
"a terminal-bearing caller must not see a question whose task records no creator");
|
||||
assertNull(messages.pendingAsk(T, "worker:term_someone"),
|
||||
"a caller with an owner key must not see a question whose task records no creator");
|
||||
assertNotNull(messages.pendingAsk(T, null),
|
||||
"the unnamed primary must still see it even with no recorded creator");
|
||||
|
||||
@@ -2055,7 +2081,8 @@ class MessageServiceTest {
|
||||
*/
|
||||
@Test
|
||||
void aDifferentCallersAnswerIsRefusedForAnAsyncSendDelegationButTheRealOwnerSucceeds() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "task that asks", null, "term_owner");
|
||||
Principal owner = Principal.worker("term_owner", 1);
|
||||
String ticket = messages.sendAsync(T, "task that asks", null, owner);
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
@@ -2063,7 +2090,8 @@ class MessageServiceTest {
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
|
||||
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
|
||||
|
||||
MessageService.Reply hijacked = messages.answer(asking.turnId(), "evil.yaml", 500, "term_attacker");
|
||||
MessageService.Reply hijacked = messages.answer(asking.turnId(), "evil.yaml", 500,
|
||||
"worker:term_attacker");
|
||||
assertEquals(MessageService.Outcome.NOT_TURN_OWNER, hijacked.outcome(),
|
||||
"a caller that did not create this delegation must be refused, not served");
|
||||
assertFalse(ask.isDone(), "a refused answer must not resolve the worker's blocked fleet_ask");
|
||||
@@ -2071,7 +2099,38 @@ class MessageServiceTest {
|
||||
"a refused answer must not advance the async ticket's phase");
|
||||
|
||||
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
|
||||
() -> messages.answer(asking.turnId(), "config.yaml", 5000, "term_owner"));
|
||||
() -> messages.answer(asking.turnId(), "config.yaml", 5000, owner.ownerKey()));
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
awaitWaiting();
|
||||
assertTrue(rendezvous.resolve(T, "done"));
|
||||
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
|
||||
assertEquals("done", awaitTicketPhase(ticket, MessageService.Phase.DONE).reply());
|
||||
}
|
||||
|
||||
@Test
|
||||
void namedLeadCanSeeAndAnswerAnAskAfterItsTerminalChangesWhileLeadBIsRefused() throws Exception {
|
||||
Principal oldLead = Principal.leader("opus", "term_OLD", 1);
|
||||
Principal newLead = Principal.leader("opus", "term_NEW", 2);
|
||||
Principal leadB = Principal.leader("sol", "term_SOL", 3);
|
||||
assertNotEquals(oldLead.terminal(), newLead.terminal(), "the test requires different terminals");
|
||||
String ticket = messages.sendAsync(T, "task that asks", null, oldLead);
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
CompletableFuture<MessageService.AskResult> ask =
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
|
||||
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
|
||||
|
||||
assertNotNull(messages.pendingAsk(T, newLead.ownerKey()),
|
||||
"the same lead at its new terminal must see the pending ask");
|
||||
assertNull(messages.pendingAsk(T, leadB.ownerKey()),
|
||||
"another lead must not see the pending ask");
|
||||
assertEquals(MessageService.Outcome.NOT_TURN_OWNER,
|
||||
messages.answer(asking.turnId(), "evil.yaml", 500, leadB.ownerKey()).outcome());
|
||||
assertFalse(ask.isDone(), "another lead must not resolve the worker's ask");
|
||||
|
||||
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
|
||||
() -> messages.answer(asking.turnId(), "config.yaml", 5000, newLead.ownerKey()));
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
awaitWaiting();
|
||||
assertTrue(rendezvous.resolve(T, "done"));
|
||||
|
||||
@@ -241,18 +241,18 @@ class RendezvousTest {
|
||||
@Test
|
||||
void ownerPermitsIsFailClosedOnARecordAndThreeStatesAreDistinct() {
|
||||
assertFalse(Rendezvous.Owner.permits(null, null),
|
||||
"no owner on record refuses even a caller with no terminal");
|
||||
assertFalse(Rendezvous.Owner.permits(null, "term_a"),
|
||||
"no owner on record refuses a terminal-bearing caller too");
|
||||
"no owner on record refuses even the unnamed primary");
|
||||
assertFalse(Rendezvous.Owner.permits(null, "worker:term_a"),
|
||||
"no owner on record refuses a caller with an owner key too");
|
||||
assertTrue(Rendezvous.Owner.permits(Rendezvous.Owner.UNNAMED_PRIMARY, null),
|
||||
"the unnamed primary owner matches a caller with no terminal");
|
||||
assertFalse(Rendezvous.Owner.permits(Rendezvous.Owner.UNNAMED_PRIMARY, "term_a"),
|
||||
"the unnamed primary owner does not match a terminal-bearing caller");
|
||||
assertTrue(Rendezvous.Owner.permits(Rendezvous.Owner.of("term_a"), "term_a"),
|
||||
"a named owner matches the same terminal");
|
||||
assertFalse(Rendezvous.Owner.permits(Rendezvous.Owner.of("term_a"), "term_b"),
|
||||
"a named owner refuses a different terminal");
|
||||
assertFalse(Rendezvous.Owner.permits(Rendezvous.Owner.of("term_a"), null),
|
||||
"the recorded unnamed primary matches a caller with a null owner key");
|
||||
assertFalse(Rendezvous.Owner.permits(Rendezvous.Owner.UNNAMED_PRIMARY, "worker:term_a"),
|
||||
"the recorded unnamed primary does not match another owner key");
|
||||
assertTrue(Rendezvous.Owner.permits(Rendezvous.Owner.of("worker:term_a"), "worker:term_a"),
|
||||
"an owner matches the same key");
|
||||
assertFalse(Rendezvous.Owner.permits(Rendezvous.Owner.of("worker:term_a"), "worker:term_b"),
|
||||
"an owner refuses a different key");
|
||||
assertFalse(Rendezvous.Owner.permits(Rendezvous.Owner.of("worker:term_a"), null),
|
||||
"a named owner refuses the unnamed primary");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
@@ -146,11 +148,11 @@ class FleetAppAuthTest {
|
||||
|
||||
/**
|
||||
* {@code GET /tasks/{ticket}} must resolve its caller the same way {@code allow(...)} does
|
||||
* and thread that terminal into {@link MessageService#poll(String, String)}, not the
|
||||
* and thread that owner key into {@link MessageService#poll(String, String)}, not the
|
||||
* no-check overload that ignores who is asking.
|
||||
*/
|
||||
@Test
|
||||
void theTaskStatusRouteActuallyThreadsTheCallersTerminalIntoPoll() throws Exception {
|
||||
void theTaskStatusRouteActuallyThreadsTheCallersOwnerKeyIntoPoll() throws Exception {
|
||||
String source = Files.readString(REST_SOURCE);
|
||||
|
||||
int start = source.indexOf("private void taskStatus(Context ctx) {");
|
||||
@@ -166,8 +168,8 @@ class FleetAppAuthTest {
|
||||
"control failed: the scraped taskStatus block contains no messages.poll( call at all "
|
||||
+ "-- the anchors have drifted, this test is not testing what it claims to");
|
||||
|
||||
assertTrue(handlerBlock.contains("caller.terminal()"),
|
||||
"the taskStatus route must thread the resolved caller's terminal into messages.poll(...), "
|
||||
assertTrue(handlerBlock.contains("caller.ownerKey()"),
|
||||
"the taskStatus route must thread the resolved caller's owner key into messages.poll(...), "
|
||||
+ "not the no-check overload -- block: " + handlerBlock);
|
||||
assertTrue(handlerBlock.contains("ctx.attribute(CALLER)"),
|
||||
"the taskStatus route must resolve its caller the same way allow(...) does, not via a "
|
||||
@@ -176,11 +178,11 @@ class FleetAppAuthTest {
|
||||
|
||||
/**
|
||||
* {@code GET /sessions/{id}/status} must resolve its caller the same way {@code allow(...)}
|
||||
* does and thread that terminal into {@link MessageService#pendingAsk(String, String)}, not
|
||||
* does and thread that owner key into {@link MessageService#pendingAsk(String, String)}, not
|
||||
* the no-check overload that ignores who is asking.
|
||||
*/
|
||||
@Test
|
||||
void theSessionStatusRouteActuallyThreadsTheCallersTerminalIntoPendingAsk() throws Exception {
|
||||
void theSessionStatusRouteActuallyThreadsTheCallersOwnerKeyIntoPendingAsk() throws Exception {
|
||||
String source = Files.readString(REST_SOURCE);
|
||||
|
||||
int start = source.indexOf("private void sessionStatus(Context ctx) {");
|
||||
@@ -196,8 +198,8 @@ class FleetAppAuthTest {
|
||||
"control failed: the scraped sessionStatus block contains no messages.pendingAsk( "
|
||||
+ "call at all -- the anchors have drifted, this test is not testing what it claims to");
|
||||
|
||||
assertTrue(handlerBlock.contains("caller.terminal()"),
|
||||
"the sessionStatus route must thread the resolved caller's terminal into "
|
||||
assertTrue(handlerBlock.contains("caller.ownerKey()"),
|
||||
"the sessionStatus route must thread the resolved caller's owner key into "
|
||||
+ "messages.pendingAsk(...), not the no-check overload -- block: " + handlerBlock);
|
||||
assertTrue(handlerBlock.contains("ctx.attribute(CALLER)"),
|
||||
"the sessionStatus route must resolve its caller the same way allow(...) does, not via "
|
||||
@@ -222,7 +224,8 @@ class FleetAppAuthTest {
|
||||
Javalin otherWorkerApp = startOnSharedService(messages, herdr, 9001L); // -> term_shell
|
||||
Javalin primaryApp = startOnSharedService(messages, herdr, 999_999L); // no pane -> primary
|
||||
try {
|
||||
String ticket = messages.sendAsync("term_a", "long task", null, "term_a");
|
||||
String ticket = messages.sendAsync("term_a", "long task", null,
|
||||
Principal.worker("term_a", FakeHerdr.WORKER_PID));
|
||||
|
||||
HttpResponse<String> refused = send(otherWorkerApp.port(), "GET", "/tasks/" + ticket, null, null);
|
||||
assertEquals(200, refused.statusCode());
|
||||
@@ -287,12 +290,12 @@ class FleetAppAuthTest {
|
||||
|
||||
/**
|
||||
* {@code GET /sessions/{id}/status} shows a worker's pending {@code fleet_ask} question, its
|
||||
* {@code turnId} and its ticket only to the caller whose terminal created that delegation, or
|
||||
* to a caller with no terminal at all (the unnamed primary) — a different terminal-bearing
|
||||
* caller still sees the base status line, but none of the pending-ask fields.
|
||||
* {@code turnId} and its ticket only to the caller whose owner key created that delegation, or
|
||||
* to the unnamed primary. A different caller still sees the base status line, but none of the
|
||||
* pending-ask fields.
|
||||
*/
|
||||
@Test
|
||||
void restStatusGatesThePendingAskFieldsByTheDelegationsCreatorTerminal() throws Exception {
|
||||
void restStatusGatesThePendingAskFieldsByTheDelegationsCreatorOwner() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
Injector injector = new Injector(agents);
|
||||
@@ -304,7 +307,8 @@ class FleetAppAuthTest {
|
||||
Javalin primaryApp = startOnSharedService(messages, herdr, 999_999L); // no pane -> primary
|
||||
try {
|
||||
ObjectMapper mapper = new ObjectMapper();
|
||||
String ticket = messages.sendAsync("term_target", "task that asks", null, "term_a");
|
||||
String ticket = messages.sendAsync("term_target", "task that asks", null,
|
||||
Principal.worker("term_a", FakeHerdr.WORKER_PID));
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
while (!rendezvous.isWaiting("term_target") && System.currentTimeMillis() < deadline) {
|
||||
Thread.sleep(5);
|
||||
@@ -343,7 +347,7 @@ class FleetAppAuthTest {
|
||||
|
||||
// Clean up the still-open ask so the background thread does not linger past the test.
|
||||
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
|
||||
() -> messages.answer(turnId, "config.yaml", 5000, "term_a"));
|
||||
() -> messages.answer(turnId, "config.yaml", 5000, "worker:term_a"));
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
deadline = System.currentTimeMillis() + 3000;
|
||||
while (!rendezvous.isWaiting("term_target") && System.currentTimeMillis() < deadline) {
|
||||
@@ -520,6 +524,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 +543,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();
|
||||
@@ -1405,6 +1452,105 @@ class SessionManagerTest {
|
||||
+ "dirty check threw");
|
||||
}
|
||||
|
||||
// --- fleetd #736: a release must forget the member's presence entry, not just its registry
|
||||
// row ---------------------------------------------------------------------------------------
|
||||
|
||||
@Test
|
||||
void releaseByPaneIdForgetsThePresenceEntry() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
SessionManager sessions = sessionManager(herdr);
|
||||
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
|
||||
String terminal = session.terminalId();
|
||||
sessions.asPresence().markPresent(terminal);
|
||||
assertTrue(sessions.asPresence().isPresent(terminal), "present before the release");
|
||||
|
||||
sessions.release(session.paneId());
|
||||
|
||||
assertFalse(sessions.asPresence().isPresent(terminal),
|
||||
"release must forget the terminal's presence, not just remove its registry row");
|
||||
}
|
||||
|
||||
@Test
|
||||
void reapIdleForgetsThePresenceEntryToo() {
|
||||
long[] clock = {0};
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
SessionManager sessions = sessionManager(herdr, () -> clock[0]);
|
||||
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
|
||||
String terminal = session.terminalId();
|
||||
sessions.asPresence().markPresent(terminal);
|
||||
assertTrue(sessions.asPresence().isPresent(terminal), "present before the reap");
|
||||
|
||||
clock[0] = 11;
|
||||
assertEquals(1, sessions.reapIdle(10), "READY session past TTL is reaped");
|
||||
|
||||
assertFalse(sessions.asPresence().isPresent(terminal),
|
||||
"the idle-reap release path (releaseIfCurrent) goes through the same teardown "
|
||||
+ "funnel as an explicit release, so it must forget presence too");
|
||||
}
|
||||
|
||||
@Test
|
||||
void shutdownDrainAlsoForgetsThePresenceEntry() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
SessionManager sessions = sessionManager(herdr);
|
||||
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
|
||||
String terminal = session.terminalId();
|
||||
sessions.asPresence().markPresent(terminal);
|
||||
assertTrue(sessions.asPresence().isPresent(terminal), "present before the drain");
|
||||
|
||||
sessions.drainAll(TimeUnit.MILLISECONDS.toNanos(100));
|
||||
|
||||
assertFalse(sessions.asPresence().isPresent(terminal),
|
||||
"a shutdown drain still ends the member's process, so presence must be cleared "
|
||||
+ "exactly as it is for any other release cause");
|
||||
}
|
||||
|
||||
@Test
|
||||
void releaseOfAnUnknownPaneIdDoesNotThrow() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
SessionManager sessions = sessionManager(herdr);
|
||||
|
||||
assertDoesNotThrow(() -> sessions.release("no-such-pane"),
|
||||
"releasing a pane id that was never registered must be a no-op, not a throw");
|
||||
}
|
||||
|
||||
@Test
|
||||
void releaseStillForgetsPresenceWhenDirtyCheckThrows() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
RecordingWorktrees worktrees = new RecordingWorktrees();
|
||||
SessionManager sessions = sessionManager(herdr, worktrees);
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("fleetd-736", null));
|
||||
String terminal = s.terminalId();
|
||||
sessions.asPresence().markPresent(terminal);
|
||||
worktrees.failHasUncommittedWith(new WorktreeException("git status exited 128"));
|
||||
|
||||
assertDoesNotThrow(() -> sessions.release(s.paneId()),
|
||||
"a throwing dirty check must not abort the release");
|
||||
|
||||
assertFalse(sessions.asPresence().isPresent(terminal),
|
||||
"presence must be forgotten even when the dirty check throws, which pins the "
|
||||
+ "forget call to the finally block that runs no matter what happened above");
|
||||
}
|
||||
|
||||
@Test
|
||||
void releaseLeavesADifferentStillLiveMembersPresenceUntouched() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
SessionManager sessions = sessionManager(herdr);
|
||||
MemberSession released = sessions.acquire("ltms-local", null, "/caller/a", "ownerA");
|
||||
MemberSession stillLive = sessions.acquire("ltms-local", null, "/caller/b", "ownerB");
|
||||
sessions.asPresence().markPresent(released.terminalId());
|
||||
sessions.asPresence().markPresent(stillLive.terminalId());
|
||||
assertTrue(sessions.asPresence().isPresent(stillLive.terminalId()),
|
||||
"present before the release of the other member");
|
||||
|
||||
sessions.release(released.paneId());
|
||||
|
||||
assertFalse(sessions.asPresence().isPresent(released.terminalId()),
|
||||
"the released terminal is forgotten");
|
||||
assertTrue(sessions.asPresence().isPresent(stillLive.terminalId()),
|
||||
"a still-live member's presence must survive an unrelated release");
|
||||
}
|
||||
|
||||
// --- fleetd #316: the dirty check must be re-taken after the worker is stopped, not trusted
|
||||
// stale from before it ------------------------------------------------------------------------
|
||||
|
||||
|
||||
@@ -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();
|
||||
|
||||
Reference in New Issue
Block a user