Compare commits

..

1 Commits

Author SHA1 Message Date
Dai Ha 6ab3a81af7 fleetd #737 unit 6: stop the unnamed primary sharing the internal bypass
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 55s
CI / build (pull_request) Failing after 1m59s
ownsTicket treated a null callerOwner as "read everything", conflating
the internal no-check test seam with a real unnamed primary's owner
key. Give them two different values: a named INTERNAL_NO_OWNER_CHECK
marker for the test-only bypass, and null as just another owner key
that must equal the ticket's recorded creatorOwner (including a
null-to-null match, so an unnamed primary still owns its own tickets).
Applies to both ownsTicket call sites, poll and pendingAsk.
2026-10-05 05:24:37 +02:00
8 changed files with 235 additions and 248 deletions
@@ -146,8 +146,9 @@ public record Principal(Role role, String terminal, long pid, String name) {
}
/**
* 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.
* Stable identity used to own tickets and open turns. The unnamed primary has no owner key —
* {@code null} — and that is matched against a ticket's recorded owner the same way any other
* key is: it owns a ticket another unnamed primary created, and nothing else.
*/
public String ownerKey() {
return switch (role) {
@@ -149,17 +149,9 @@ public final class LeadRollover {
* calling lead's workspace, before storing it here; see this class's
* javadoc. This is the value the MCP layer hands back to the lead as
* "write your file here", so callers may rely on it always being absolute.
* @param rolloverKey the key {@link #confirm}'s single-flight claim is taken under. {@link
* #open} resolves this once, here, from {@code leadTerminal} while the
* calling lead is certainly still live: the lead's configured name, or
* {@code leadTerminal} itself when no name resolves. Carried rather than
* recomputed at release time, because by the time a roll's continuation
* releases its claim the OLD terminal may no longer resolve to any name at
* all — recomputing there would release a different key than the one the
* claim was taken under.
*/
public record PendingRollover(String token, String leadTerminal, String handoverPath,
long requestedAtMillis, String rolloverKey) {}
long requestedAtMillis) {}
/** Which check refused a {@link #confirm} call, named so a caller can act on it. */
public enum RefusalReason {
@@ -184,9 +176,9 @@ public final class LeadRollover {
*/
HANDOVER_STALE,
/**
* This lead already has a roll running: an earlier {@link #confirm} call claimed its
* single-flight key (see {@link PendingRollover#rolloverKey}) and that roll's continuation
* has not released it yet. {@code detail} names the key and the token that holds the claim.
* 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
}
@@ -348,12 +340,9 @@ public final class LeadRollover {
private final Function<String, String> leadWorkspace;
/**
* Terminal id → that lead's configured name under {@code fleet.leaders}, or {@code null} when
* the terminal names no currently-recognised lead. {@link #open} calls this on the calling
* lead's own terminal, while it is certainly still live, to resolve {@link
* PendingRollover#rolloverKey}. The deferred continuation also calls this, on the OLD terminal,
* before tearing it down, so it knows which lead to pass to {@link LeadLauncher#relaunch} — by
* that point the live roster may no longer contain the old terminal, so this lookup can return
* {@code null} here even though {@link #open}'s earlier call against the same terminal did not.
* the terminal names no currently-recognised lead. The deferred continuation calls this, on the
* OLD terminal, before tearing it down, so it knows which lead to pass to {@link
* LeadLauncher#relaunch}.
*/
private final Function<String, String> leadNameForTerminal;
/**
@@ -373,14 +362,14 @@ public final class LeadRollover {
private final Consumer<Runnable> continuationRunner;
private final Map<String, PendingRollover> pending = new ConcurrentHashMap<>();
/**
* {@link PendingRollover#rolloverKey} → the token of the roll currently holding that key
* 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
* 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 key
* absent from this map has no roll currently in flight for it.
* 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> rollingByLead = new ConcurrentHashMap<>();
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
@@ -449,9 +438,8 @@ public final class LeadRollover {
/**
* The lead says it is ready to be replaced. Generates a token and records the resolved
* handover path, this moment's wall-clock timestamp (the baseline {@link #confirm} checks the
* handover file's modified time against), {@code leadTerminal} — only that exact terminal may
* later {@link #confirm} this token — and {@link PendingRollover#rolloverKey}, resolved here
* from {@code leadTerminal} while the calling lead is certainly still live.
* handover file's modified time against), and {@code leadTerminal} — only that exact terminal
* may later {@link #confirm} this token.
*
* @param leadTerminal the calling lead's terminal id, resolved by the MCP layer from the
* connection (see this class's javadoc) — never a client-supplied value
@@ -471,9 +459,7 @@ public final class LeadRollover {
String token = UUID.randomUUID().toString();
long requestedAt = nowMillis.getAsLong();
String resolvedPath = resolveHandoverPath(cfg.handoverPath(), leadTerminal);
String leadName = leadNameForTerminal.apply(leadTerminal);
String rolloverKey = (leadName == null || leadName.isBlank()) ? leadTerminal : leadName;
PendingRollover p = new PendingRollover(token, leadTerminal, resolvedPath, requestedAt, rolloverKey);
PendingRollover p = new PendingRollover(token, leadTerminal, resolvedPath, requestedAt);
pending.put(token, p);
if (resolvedPath.equals(cfg.handoverPath())) {
log.info("lead-rollover: open token={} lead={} handoverPath={} reason={}",
@@ -579,11 +565,12 @@ public final class LeadRollover {
// 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's claim.
String holder = rollingByLead.putIfAbsent(p.rolloverKey(), token);
// 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 '" + p.rolloverKey() + "' already has a roll running under token " + holder);
"lead terminal " + p.leadTerminal() + " already has a roll running under token "
+ holder);
}
// Record IN_PROGRESS BEFORE removing from `pending` — see RollState#IN_PROGRESS and
@@ -604,14 +591,14 @@ public final class LeadRollover {
} 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 rollingByLead — never runs either.
// 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 could never be rolled again and status() would report IN_PROGRESS
// forever for a roll that in fact never started.
// 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);
rollingByLead.remove(p.rolloverKey(), token);
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"));
@@ -654,10 +641,10 @@ public final class LeadRollover {
+ "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 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 key.
rollingByLead.remove(p.rolloverKey(), p.token());
// 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());
}
}
@@ -1352,9 +1352,10 @@ 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 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.
* {@code turnId} and its ticket id are shown only to the caller whose owner key matches the
* delegation's creator — the unnamed primary matches only a delegation another unnamed
* primary created; 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 callerOwner) {
if (isBlank(sessionId)) {
@@ -12,6 +12,7 @@ import org.slf4j.LoggerFactory;
import java.util.ArrayList;
import java.util.List;
import java.util.Objects;
import java.util.UUID;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CompletionException;
@@ -1391,21 +1392,23 @@ public final class MessageService {
}
/**
* 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.
* As {@link #poll(String, String)}, but bypasses the ownership check entirely via
* {@link #INTERNAL_NO_OWNER_CHECK}. No production code calls this overload — it exists for
* tests that only need the ticket's state and have no caller identity to pass.
*/
public TaskView poll(String ticket) {
return poll(ticket, null);
return poll(ticket, INTERNAL_NO_OWNER_CHECK);
}
/**
* Snapshot the state of an async delegation. Returns {@code null} for an unknown/expired ticket.
* 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.
* carries no reply text. The unnamed primary's owner key is {@code null}, matched the same way
* as any other key — it reads a ticket another unnamed primary created, and is refused on a
* ticket a named caller created. 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 callerOwner) {
Task task = tasks.get(ticket);
@@ -1452,13 +1455,27 @@ public final class MessageService {
}
/**
* 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
* Marker passed as {@code callerOwner} to bypass the ownership check entirely. No
* {@link Principal#ownerKey()} ever produces this value — every real key is either
* {@code null} (the unnamed primary) or prefixed with its role, such as {@code "worker:"} or
* {@code "leader:"}. {@link #poll(String)} passes it; {@link #pendingAsk} has no matching
* no-check overload, so this stays package-private for the test that drives the bypass
* directly.
*/
static final String INTERNAL_NO_OWNER_CHECK = "internal:no-owner-check";
/**
* Whether {@code callerOwner} may read {@code task}'s state. {@code callerOwner} is matched
* against the task's recorded owner key by equality, including a {@code null} match — the
* unnamed primary's owner key is {@code null}, so it owns a ticket another unnamed primary
* created and nothing else, the same rule every other role follows. The only caller that
* reads any ticket is {@link #INTERNAL_NO_OWNER_CHECK}. 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 callerOwner) {
return callerOwner == null || callerOwner.equals(task.creatorOwner);
return INTERNAL_NO_OWNER_CHECK.equals(callerOwner)
|| Objects.equals(callerOwner, task.creatorOwner);
}
/**
@@ -16,7 +16,6 @@ import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
@@ -1128,13 +1127,8 @@ class LeadRolloverTest {
AgentControl agents = new AgentControl(herdr);
WorkspaceControl spaces = new WorkspaceControl(herdr);
LeadLauncher launcher = fakeLauncher(herdr, fleetConfigWithRelaunchableLead());
// Only LEAD resolves to a configured name — every other terminal (e.g. the distinct
// term_cap_N terminals evictionCountsInProgressEntriesTowardTheCap opens) falls back to
// keying its own single-flight claim on the terminal itself, exactly like a terminal
// the live roster does not recognise.
return new LeadRollover(agents, spaces, launcher, () -> config, _ -> null,
t -> LEAD.equals(t) ? LEAD_NAME : null, () -> liveLeadTerminals, nowMillis, () -> { },
runner);
_ -> LEAD_NAME, () -> liveLeadTerminals, nowMillis, () -> { }, runner);
}
@Test
@@ -1323,7 +1317,7 @@ class LeadRolloverTest {
+ "so an operator reading status() has something to act on: " + status.detail());
}
// ---- fleetd #726 unit 3: confirm() single-flights one roll at a time per lead ---------------
// ---- 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 "
@@ -1349,8 +1343,8 @@ class LeadRolloverTest {
assertFalse(secondDecision.accepted());
assertEquals(LeadRollover.RefusalReason.ROLL_ALREADY_RUNNING, secondDecision.reason());
assertTrue(secondDecision.detail().contains(LEAD_NAME), "the refusal detail must name "
+ "the single-flight key (the lead's name): " + secondDecision.detail());
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 "
@@ -1508,160 +1502,4 @@ class LeadRolloverTest {
+ "this lead terminal could never be rolled again: " + retryDecision.reason() + " / "
+ retryDecision.detail());
}
// ---- fleetd #737 unit 4: confirm() single-flights per LEAD NAME, not per lead terminal ------
@Test
@DisplayName("[NAME-KEYED 1] a second confirm() for the SAME lead is refused with "
+ "ROLL_ALREADY_RUNNING even when it is opened from a DIFFERENT terminal, while the "
+ "first roll's continuation is still in flight")
void secondConfirmForTheSameLeadIsRefusedEvenFromADifferentTerminal() throws IOException {
FakeHerdr herdr = herdrReadyForAFullRoll(); // the held roll WOULD complete once run
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
HoldingRunner runner = new HoldingRunner();
Map<String, String> liveLeadTerminals = Map.of(NEW_TERMINAL, LEAD_NAME);
AgentControl agents = new AgentControl(herdr);
WorkspaceControl spaces = new WorkspaceControl(herdr);
LeadLauncher launcher = fakeLauncher(herdr, fleetConfigWithRelaunchableLead());
// Both LEAD and OTHER_LEAD resolve to the SAME configured lead name — modelling a roll
// that has already replaced the lead's pane: the fresh pane (here, OTHER_LEAD) is still
// the SAME lead, just a different terminal id.
Function<String, String> leadNameForTerminal = _ -> LEAD_NAME;
LeadRollover rollover = new LeadRollover(agents, spaces, launcher, () -> cfg(handover.toString()),
_ -> null, leadNameForTerminal, () -> liveLeadTerminals, 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(OTHER_LEAD,
"a second request for the SAME lead, opened from a DIFFERENT terminal");
LeadRollover.RollDecision secondDecision = rollover.confirm(OTHER_LEAD, second.token(), true);
assertFalse(secondDecision.accepted(), "a different terminal resolving to the SAME lead "
+ "name must still be refused — the single-flight claim is keyed on the lead's "
+ "name, not its terminal");
assertEquals(LeadRollover.RefusalReason.ROLL_ALREADY_RUNNING, secondDecision.reason());
assertTrue(secondDecision.detail().contains(LEAD_NAME), "the refusal detail must name "
+ "the single-flight key (the lead's name): " + 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
// The claim must have been released once the first roll's continuation finished — a
// fresh request for the SAME lead, even opened from yet another terminal, is now
// approved.
LeadRollover.PendingRollover third = rollover.open(OTHER_LEAD, "retry after the first roll finished");
LeadRollover.RollDecision thirdDecision = rollover.confirm(OTHER_LEAD, third.token(), true);
assertTrue(thirdDecision.accepted(), "the claim must have been released once the first "
+ "roll finished: " + thirdDecision.reason() + " / " + thirdDecision.detail());
}
@Test
@DisplayName("[NAME-KEYED 2] the claim is released under the key it was taken under, even "
+ "when the old terminal has already dropped out of the live-lead roster by release "
+ "time — success (non-throwing) path")
void claimReleasedUnderCarriedKeyEvenWhenOldTerminalIsNoLongerInTheRosterSuccessPath() throws IOException {
// herdrReadyForAFullRoll() lets the RETRY below run its continuation to a genuine ROLLED
// completion (pane gone on the first check, a pinned relaunch) — needed because, unlike
// the other NAME-KEYED tests, this one's retry resolves a REAL name ("opus") and must not
// spin forever against this test's fixed, non-advancing clock (see this class's javadoc).
FakeHerdr herdr = herdrReadyForAFullRoll();
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
AgentControl agents = new AgentControl(herdr);
WorkspaceControl spaces = new WorkspaceControl(herdr);
// A config that CAN relaunch 'opus' — unlike emptyFleetConfig(), its fleet() is non-null,
// so resolveLaunchable(null) below returns null cleanly instead of throwing a NullPointer
// out of cfg.fleet() itself; this test needs the NON-throwing relaunch-refused exit, not
// an incidental NPE.
LeadLauncher launcher = fakeLauncher(herdr, fleetConfigWithRelaunchableLead());
Map<String, String> roster = new HashMap<>();
roster.put(LEAD, LEAD_NAME);
Map<String, String> liveLeadTerminals = Map.of(NEW_TERMINAL, LEAD_NAME);
LeadRollover rollover = new LeadRollover(agents, spaces, launcher, () -> cfg(handover.toString()),
_ -> null, roster::get, () -> liveLeadTerminals, fixedClock(clock), () -> { }, Runnable::run);
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
// The live roster drops the OLD terminal before this roll's continuation runs — the
// hazard fleetd #737 names: leadNameForTerminal reads the LIVE roster, and (with the
// synchronous runner this test injects) the continuation runs INSIDE this confirm()
// call, strictly after open() already resolved and carried the key.
roster.remove(LEAD);
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertTrue(decision.accepted(), "expected approval; got: " + decision.reason() + " / " + decision.detail());
LeadRollover.RollStatus status = rollover.status(pending.token());
assertEquals(LeadRollover.RollState.RELAUNCH_FAILED, status.state(), "sanity: leadName "
+ "resolved to null (the roster had already dropped the old terminal), so "
+ "relaunch(null) fails cleanly, through the non-throwing exit this test means "
+ "to cover: " + status.detail());
// A fresh terminal for the SAME lead name, confirmed live again — proves the claim this
// roll took under the lead's name was actually released, not left stuck under whatever
// a release that RECOMPUTES the key via leadNameForTerminal (now null for the old
// terminal) would have tried to remove instead.
roster.put(OTHER_LEAD, LEAD_NAME);
LeadRollover.PendingRollover retry = rollover.open(OTHER_LEAD, "a fresh pane for the same lead");
LeadRollover.RollDecision retryDecision = rollover.confirm(OTHER_LEAD, retry.token(), true);
assertTrue(retryDecision.accepted(), "the claim taken under the lead's name must have "
+ "been released even though the OLD terminal no longer resolved to any name at "
+ "release time: " + retryDecision.reason() + " / " + retryDecision.detail());
LeadRollover.RollStatus retryStatus = rollover.status(retry.token());
assertEquals(LeadRollover.RollState.ROLLED, retryStatus.state(), "sanity: this retry's "
+ "own claim (also 'opus') must not have been blocked by a leftover claim from "
+ "the first roll: " + retryStatus.detail());
}
@Test
@DisplayName("[NAME-KEYED 3] the claim is released under the key it was taken under on the "
+ "thrown-exception path too, even when the old terminal has already dropped out of "
+ "the live-lead roster")
void claimReleasedUnderCarriedKeyEvenWhenOldTerminalIsNoLongerInTheRosterThrowingPath() throws IOException {
FakeHerdr fake = new FakeHerdr(); // default idle — the turn-settle wait passes immediately
fake.paneCloseFailsWith("permission_denied"); // a real failure, not an already-gone code
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
AgentControl agents = new AgentControl(fake);
WorkspaceControl spaces = new WorkspaceControl(fake);
LeadLauncher launcher = fakeLauncher(fake, emptyFleetConfig());
Map<String, String> roster = new HashMap<>();
roster.put(LEAD, LEAD_NAME);
LeadRollover rollover = new LeadRollover(agents, spaces, launcher, () -> cfg(handover.toString()),
_ -> null, roster::get, Map::of, fixedClock(clock), () -> { }, Runnable::run);
LeadRollover.PendingRollover first = rollover.open(LEAD, "context is full");
// Same hazard as [NAME-KEYED 2] above, now exercised on the thrown-exception exit.
roster.remove(LEAD);
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());
// A fresh terminal for the SAME lead name, confirmed live again — proves the claim this
// roll took under the lead's name was actually released, not left stuck under whatever
// a release that RECOMPUTES the key via leadNameForTerminal (now null for the old
// terminal) would have tried to remove instead.
roster.put(OTHER_LEAD, LEAD_NAME);
LeadRollover.PendingRollover retry = rollover.open(OTHER_LEAD, "retry after the throw");
LeadRollover.RollDecision retryDecision = rollover.confirm(OTHER_LEAD, retry.token(), true);
assertTrue(retryDecision.accepted(), "the claim must have been released even though the "
+ "continuation threw AND the old terminal no longer resolved to any name at "
+ "release time: " + retryDecision.reason() + " / " + retryDecision.detail());
}
}
@@ -2013,8 +2013,10 @@ class FleetMcpTest {
/**
* {@code fleet_status}'s pending-ask block (the question, its {@code turnId} and its ticket)
* 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.
* is shown only to the caller whose owner key created the delegation. An unnamed primary is
* held to the same rule: its owner key is {@code null}, which here does not match the named
* worker that created the delegation, so it sees none of the pending-ask fields either — the
* same as any other non-creating caller.
*/
@Test
void statusGatesThePendingAskFieldsByTheDelegationsCreatorOwner() throws Exception {
@@ -2054,8 +2056,13 @@ class FleetMcpTest {
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?"),
"a caller with no terminal (the unnamed primary) must see the question: " + unnamed);
assertTrue(unnamed.startsWith("idle"), "the base status must still be shown: " + unnamed);
assertFalse(unnamed.contains("which config file?"),
"an unnamed primary must not see a question on a delegation a named worker created: " + unnamed);
assertFalse(unnamed.contains(asking.turnId()),
"a non-creating unnamed primary must not see the turnId: " + unnamed);
assertFalse(unnamed.contains(ticket),
"a non-creating unnamed primary must not see the ticket: " + unnamed);
// Clean up the still-open ask so the background thread does not linger past the test.
String turnId = asking.turnId();
@@ -948,7 +948,7 @@ class MessageServiceTest {
}
@Test
void unnamedPrimaryStillReadsAnyTicket() throws Exception {
void unnamedPrimaryIsRefusedFromANamedLeadsTicket() throws Exception {
String ticket = messages.sendAsync(T, "long task", null,
Principal.leader("opus", "term_lead", 1));
awaitWaiting();
@@ -956,8 +956,51 @@ class MessageServiceTest {
injector.onStatus(T, AgentStatus.WORKING); // worker works
assertTrue(rendezvous.resolve(T, "primary-visible result"), "a reply resolves the async send");
MessageService.TaskView view = driveAsyncTicketToDone(ticket, null);
assertNotNull(view, "the unnamed primary must be able to read any ticket");
MessageService.TaskView refused = messages.poll(ticket, Principal.primary(1).ownerKey());
assertNotNull(refused, "a different owner gets a refusal, not silence");
assertEquals(MessageService.Phase.FAILED, refused.phase());
assertEquals("forbidden: this ticket was created by a different session", refused.detail());
assertNull(refused.reply(), "a refusal must never carry the reply text");
assertFalse(String.valueOf(refused).contains("primary-visible result"),
"the reply text must not appear anywhere in the refused view");
}
/**
* Positive control for {@link #unnamedPrimaryIsRefusedFromANamedLeadsTicket}: without this,
* that test would pass just as well if {@code poll} refused every caller.
*/
@Test
void unnamedPrimaryReadsItsOwnTicket() throws Exception {
String ticket = messages.sendAsync(T, "long task", null, Principal.primary(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");
MessageService.TaskView view = driveAsyncTicketToDone(ticket, Principal.primary(2).ownerKey());
assertNotNull(view, "an unnamed primary must be able to read a ticket another unnamed primary created");
assertEquals(MessageService.Phase.DONE, view.phase());
assertEquals("primary-visible result", view.reply());
}
@Test
void theOneArgPollOverloadBypassesOwnershipEntirely() throws Exception {
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");
MessageService.TaskView view = null;
long deadline = System.currentTimeMillis() + 2000;
while (view == null || view.phase() != MessageService.Phase.DONE) {
if (System.currentTimeMillis() >= deadline) break;
view = messages.poll(ticket); // the one-arg, no-check overload -- no caller owner key at all
//noinspection BusyWait
Thread.sleep(5);
}
assertNotNull(view, "the internal bypass must read a ticket owned by a named lead");
assertEquals(MessageService.Phase.DONE, view.phase());
assertEquals("primary-visible result", view.reply());
}
@@ -1978,7 +2021,9 @@ class MessageServiceTest {
/**
* A caller's owner key must match the key that created the delegation to see its pending
* question. The unnamed primary always sees it.
* question. The unnamed primary is held to the same rule as everyone else: its key is
* {@code null}, which here does not match the named worker that created this delegation, so
* it is refused too.
*/
@Test
void pendingAskGatesTheQuestionByTheDelegationsCreatorOwner() throws Exception {
@@ -1998,10 +2043,65 @@ class MessageServiceTest {
assertNotNull(own, "the creating caller must see its own open question");
assertEquals("which config file?", own.question());
MessageService.PendingAsk unnamed = messages.pendingAsk(T, null);
assertNotNull(unnamed, "a caller with no terminal (the unnamed primary) must always see the question");
assertNull(messages.pendingAsk(T, Principal.primary(1).ownerKey()),
"an unnamed primary must not see a question on a delegation a named worker created");
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(asking.turnId(), "config.yaml", 5000, creator.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());
}
/**
* Positive control for {@link #pendingAskGatesTheQuestionByTheDelegationsCreatorOwner}:
* without this, that test's refusal would pass just as well if {@code pendingAsk} refused
* every caller. Here the delegation's creator is itself an unnamed primary (owner key
* {@code null}), so another unnamed primary's {@code null} key must still match it.
*/
@Test
void unnamedPrimarySeesItsOwnDelegationsPendingQuestion() throws Exception {
String ticket = messages.sendAsync(T, "task that asks", null, Principal.primary(1));
awaitWaiting();
injectDelivery();
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
awaitTicketPhase(ticket, MessageService.Phase.ASKING);
MessageService.PendingAsk unnamed = messages.pendingAsk(T, Principal.primary(2).ownerKey());
assertNotNull(unnamed, "an unnamed primary must see the question on a delegation another unnamed primary created");
assertEquals("which config file?", unnamed.question());
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(unnamed.turnId(), "config.yaml", 5000, null));
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());
}
/**
* {@code pendingAsk} has no public no-check overload the way {@link MessageService#poll}
* does, so this drives {@link MessageService#INTERNAL_NO_OWNER_CHECK} directly — the only way
* to exercise the bypass for this method.
*/
@Test
void pendingAskInternalBypassSeesAnyDelegationsPendingQuestion() throws Exception {
Principal creator = Principal.worker("term_creator", 1);
String ticket = messages.sendAsync(T, "task that asks", null, creator);
awaitWaiting();
injectDelivery();
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
MessageService.PendingAsk bypassed = messages.pendingAsk(T, MessageService.INTERNAL_NO_OWNER_CHECK);
assertNotNull(bypassed, "the internal bypass must see a question on a delegation a named worker created");
assertEquals("which config file?", bypassed.question());
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(asking.turnId(), "config.yaml", 5000, creator.ownerKey()));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
@@ -207,14 +207,15 @@ class FleetAppAuthTest {
}
/**
* {@code GET /tasks/{ticket}} refuses a worker whose terminal did not create the ticket,
* while the creating worker and the unnamed primary both still read it. The ticket is minted
* directly on the shared {@link MessageService}, the same way {@code MessageServiceTest}
* drives {@link MessageService#poll(String, String)}, so this exercises only the REST poll
* route's own handling of the ownership already recorded on the ticket.
* {@code GET /tasks/{ticket}} refuses a worker whose terminal did not create the ticket, and
* refuses an unnamed primary just the same: a named worker's ticket is not anyone else's to
* read, caller rank included. The ticket is minted directly on the shared
* {@link MessageService}, the same way {@code MessageServiceTest} drives
* {@link MessageService#poll(String, String)}, so this exercises only the REST poll route's
* own handling of the ownership already recorded on the ticket.
*/
@Test
void restPollRefusesADifferentWorkerButAllowsTheCreatorAndTheUnnamedPrimary() throws Exception {
void restPollRefusesADifferentWorkerAndAnUnnamedPrimaryOnANamedWorkersTicket() throws Exception {
FakeHerdr herdr = new FakeHerdr();
AgentControl agents = new AgentControl(herdr);
Injector injector = new Injector(agents);
@@ -241,8 +242,10 @@ class FleetAppAuthTest {
HttpResponse<String> primary = send(primaryApp.port(), "GET", "/tasks/" + ticket, null, null);
assertEquals(200, primary.statusCode());
assertFalse(primary.body().contains("forbidden"),
"the unnamed primary must read any ticket: " + primary.body());
assertTrue(primary.body().contains("forbidden"),
"an unnamed primary must not read a ticket a named worker created: " + primary.body());
assertFalse(primary.body().contains("\"reply\""),
"a refusal must never carry reply text: " + primary.body());
} finally {
creatorApp.stop();
otherWorkerApp.stop();
@@ -250,6 +253,33 @@ class FleetAppAuthTest {
}
}
/**
* Positive control for
* {@link #restPollRefusesADifferentWorkerAndAnUnnamedPrimaryOnANamedWorkersTicket}: without
* this, that test's refusal would pass just as well if the route refused every caller. Here
* the ticket's creator is itself an unnamed primary, so another unnamed primary reading it
* over REST must still succeed.
*/
@Test
void restPollAllowsAnUnnamedPrimaryItsOwnTicket() throws Exception {
FakeHerdr herdr = new FakeHerdr();
AgentControl agents = new AgentControl(herdr);
Injector injector = new Injector(agents);
MessageService messages = new MessageService(agents, injector, new Rendezvous());
Javalin primaryApp = startOnSharedService(messages, herdr, 999_999L); // no pane -> primary
try {
String ticket = messages.sendAsync("term_a", "long task", null, Principal.primary(FakeHerdr.WORKER_PID));
HttpResponse<String> own = send(primaryApp.port(), "GET", "/tasks/" + ticket, null, null);
assertEquals(200, own.statusCode());
assertFalse(own.body().contains("forbidden"),
"an unnamed primary must read a ticket another unnamed primary created: " + own.body());
} finally {
primaryApp.stop();
}
}
/**
* {@code POST /sessions/{id}/message} with {@code wait:false} must record the creating
* caller's own terminal on the ticket it returns, so that caller can still poll its own
@@ -290,9 +320,10 @@ 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 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.
* {@code turnId} and its ticket only to the caller whose owner key created that delegation. An
* unnamed primary is held to the same rule: its owner key is {@code null}, which here does not
* match the named worker that created the delegation, so it sees none of the pending-ask
* fields either — the same as any other non-creating caller.
*/
@Test
void restStatusGatesThePendingAskFieldsByTheDelegationsCreatorOwner() throws Exception {
@@ -321,7 +352,8 @@ class FleetAppAuthTest {
MessageService.TaskView asking;
deadline = System.currentTimeMillis() + 3000;
do {
asking = messages.poll(ticket, null);
// the no-check overload: this is test plumbing waiting for ASKING, not the gate under test
asking = messages.poll(ticket);
Thread.sleep(5);
} while (asking.phase() != MessageService.Phase.ASKING && System.currentTimeMillis() < deadline);
assertEquals(MessageService.Phase.ASKING, asking.phase());
@@ -342,8 +374,11 @@ class FleetAppAuthTest {
JsonNode primary = mapper.readTree(
send(primaryApp.port(), "GET", "/sessions/term_target/status", null, null).body());
assertEquals("which config file?", primary.get("question").asText(),
"a caller with no terminal (the unnamed primary) must see the question");
assertEquals("idle", primary.get("status").asText(), "the base status must still be shown");
assertFalse(primary.has("question"),
"an unnamed primary must not see a question on a delegation a named worker created: " + primary);
assertFalse(primary.has("turnId"), "a non-creating unnamed primary must not see the turnId: " + primary);
assertFalse(primary.has("ticket"), "a non-creating unnamed primary must not see the ticket: " + primary);
// Clean up the still-open ask so the background thread does not linger past the test.
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
@@ -398,7 +433,8 @@ class FleetAppAuthTest {
MessageService.TaskView asking;
deadline = System.currentTimeMillis() + 3000;
do {
asking = messages.poll(ticket, null);
// the no-check overload: this is test plumbing waiting for ASKING, not the gate under test
asking = messages.poll(ticket);
Thread.sleep(5);
} while (asking.phase() != MessageService.Phase.ASKING && System.currentTimeMillis() < deadline);
assertEquals(MessageService.Phase.ASKING, asking.phase());