Compare commits

..

8 Commits

Author SHA1 Message Date
Dai Ha 886ce1521d Merge remote-tracking branch 'origin/main' into worker/737-9c61d3-4
CI / shell-tests (pull_request) Failing after 10s
CI / contract (pull_request) Successful in 55s
CI / build (pull_request) Failing after 1m59s
2026-10-05 06:00:49 +02:00
Dai Ha d41aff4012 fleetd #737 unit 5 correction: outstanding() reports terminal-phase tickets too
A DONE/FAILED ticket nobody has polled yet is destroyed on a timer by
pruneTerminalTickets' completion-based TTL, while a PENDING ticket is not
going anywhere. Drop the !future.isDone() filter so outstanding() reports
every ticket the caller owns that is still in tasks, with its real phase.
2026-10-05 06:00:43 +02:00
Dai Ha 70a735b638 fleetd #737 unit 5: fleet_handover{open} reports outstanding tickets and open asks
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 51s
CI / build (pull_request) Failing after 2m10s
MessageService.outstanding(callerOwner) lists the caller's non-terminal
delegations and the subset paused in fleet_ask, filtered by the same
ownsTicket rule poll() already uses. FleetMcp threads the caller's owner
key into handover/handoverOpen and adds outstandingTickets/openAsks to
the open() JSON, alongside the unchanged token/handoverPath/requestedAtMillis.
2026-10-05 05:56:06 +02:00
Dai Ha 803c91ea6c fleetd #737 unit 3: name the resolving accessor in LeadRollover's javadoc
CI / shell-tests (push) Failing after 6s
CI / contract (push) Successful in 1m0s
CI / build (push) Failing after 2m7s
The heartbeat loop now resolves its nudge target through
PrimaryRegistry#currentPrimaryTerminal(), so the sentence describing what a
background loop with no caller uses named the raw accessor instead.
2026-10-05 05:37:32 +02:00
Dai Ha d438a74575 Merge branch 'worker/737-20d1d9-1' 2026-10-05 05:37:32 +02:00
Dai Ha 7467ffa252 fleetd #737 unit 6: drop the stale unnamed-primary claim from the REST status comment
CI / shell-tests (push) Failing after 6s
CI / contract (push) Successful in 55s
CI / build (push) Failing after 2m9s
The comment described the pre-#737 rule, where a null caller key read every
ticket. The pendingAsk gate now matches the unnamed primary's key like any
other, so the exemption it named no longer exists.
2026-10-05 05:32:59 +02:00
Dai Ha f4176ae455 fleetd #737 unit 3: probe the fallback terminal and resolve it by lead name
CI / shell-tests (pull_request) Failing after 10s
CI / contract (pull_request) Successful in 55s
CI / build (pull_request) Failing after 1m55s
ReplyPushLoop.resolveLiveLead dropped a dead per-target delegation and
retried PrimaryRegistry.nudgeTargetFor, but returned that fallback without
checking isLive. The fallback is now probed the same way the first lead is,
and the method returns empty rather than trust a dead terminal.

PrimaryRegistry records a delegating lead's name alongside its learned
terminal and resolves the name back to its current terminal at nudge time,
through a name-to-terminal lookup backed by the live lead-tab scan. A name
with no current match falls back to the terminal that was actually learned,
so an unnamed primary, an off-host lead, or a non-herdr lead keeps working
exactly as before. LeadHeartbeatLoop now reads the resolved current terminal
instead of the raw learned one.
2026-10-05 05:29:06 +02:00
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
18 changed files with 855 additions and 286 deletions
@@ -904,6 +904,32 @@ public final class Fleetd {
throw (T) t;
}
/**
* The {@link PrimaryRegistry} lookup for "which terminal currently hosts the lead named
* {@code name}" — the inverse of {@code liveLeadTerminals} (terminal id → lead name), read live
* on every call so a lead discovered, rolled, or lost since the last call is reflected without
* a restart. Returns {@code null} when no currently recognised lead carries that name — a name
* that is not a lead at all (an architect slot, a collaborator), or a lead whose tab the scan
* cannot currently place (just rolled, off-host, non-herdr).
*
* @param liveLeadTerminals terminal id → lead name for every CURRENTLY recognised lead, normally
* the same {@code leads} supplier {@code main} already builds for
* {@code HerdrRouter}/{@link #leadSeatLookup}
*/
static Function<String, String> currentTerminalForName(Supplier<Map<String, String>> liveLeadTerminals) {
return name -> {
if (name == null) {
return null;
}
for (var entry : liveLeadTerminals.get().entrySet()) {
if (name.equals(entry.getValue())) {
return entry.getKey();
}
}
return null;
};
}
/**
* fleetd #480: construct the {@link LeadRollover} executor only when {@code leadRollover:} is
* present at startup — the same presence gate {@code leadHeartbeat:} uses just above this
@@ -393,7 +393,8 @@ final class FleetdAssembly {
ports.leadMailboxOpener());
// CB-307: learn the primary's terminal from orchestration tool calls (or pin from config).
String pinnedPrimaryTerminal = cfg.primary() != null ? cfg.primary().terminal() : null;
PrimaryRegistry primaryRegistry = new PrimaryRegistry(pinnedPrimaryTerminal);
PrimaryRegistry primaryRegistry = new PrimaryRegistry(pinnedPrimaryTerminal,
Fleetd.currentTerminalForName(leads));
// CB-532: `primary.terminal` is superseded and no longer needed for either of its jobs.
if (pinnedPrimaryTerminal != null && !pinnedPrimaryTerminal.isBlank()) {
log.warn("primary.terminal is DEPRECATED (CB-532) and can be deleted: identity now comes "
@@ -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) {
@@ -93,7 +93,7 @@ import java.util.function.Supplier;
*
* <p><strong>Identity is resolved by the caller, never looked up here — a second fleetd #480
* correction.</strong> The first version resolved the pane to clear via {@code
* PrimaryRegistry#primaryTerminal()}. That is correct for a background loop with no caller (see
* PrimaryRegistry#currentPrimaryTerminal()}. That is correct for a background loop with no caller (see
* {@code dev.ltms.fleet.msg.LeadHeartbeatLoop}), but wrong here and a violation of this project's
* own charter invariant 3 — "identity comes from the connection, never an argument." This daemon
* can hold more than one labelled lead tab (see {@code LeadLauncher}'s fleetd #359 two-reading
@@ -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());
}
}
@@ -489,7 +489,13 @@ public final class FleetMcp {
// 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, callerTerminal);
// Only a lead's name is ever resolvable back to a current terminal (PrimaryRegistry
// only looks it up among currently recognised leads) — an architect or collaborator
// name would never match there anyway, but passing null for them keeps the intent
// explicit rather than relying on that lookup to filter it out.
String delegatorName = caller.isPrimary() ? caller.name() : null;
Runnable onAccepted = () ->
primaryRegistry.recordDelegation(target, callerTerminal, delegatorName);
// 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)
@@ -594,7 +600,8 @@ public final class FleetMcp {
(exchange, req) -> {
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_handover", req.arguments()), null);
if (denied != null) return denied;
return handover(leadRollover, callerTerminal(exchange), req.arguments());
Principal caller = principal(exchange);
return handover(leadRollover, messages, caller.terminal(), caller.ownerKey(), req.arguments());
};
McpSchema.Tool fleetSend = sendTool();
@@ -730,7 +737,7 @@ public final class FleetMcp {
*/
static void recordPrimarySingleton(PrimaryRegistry registry, String callerTerminal, Principal caller) {
if (caller != null && caller.isPrimary()) {
registry.record(callerTerminal);
registry.record(callerTerminal, caller.name());
}
}
@@ -1352,9 +1359,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)) {
@@ -1475,14 +1483,15 @@ public final class FleetMcp {
* #474 charter tool-surface gate), so every action here must degrade to a clean, structured
* refusal naming {@code NOT_CONFIGURED} rather than ever throwing.
*/
static McpSchema.CallToolResult handover(LeadRollover leadRollover, String callerTerminal,
static McpSchema.CallToolResult handover(LeadRollover leadRollover, MessageService messages,
String callerTerminal, String callerOwner,
Map<String, Object> args) {
String action = str(args, "action");
if (isBlank(action)) {
return error("action is required: \"open\", \"confirm\", \"cancel\" or \"status\"");
}
return switch (action) {
case "open" -> handoverOpen(leadRollover, callerTerminal, str(args, "reason"));
case "open" -> handoverOpen(leadRollover, messages, callerTerminal, callerOwner, str(args, "reason"));
case "confirm" -> handoverConfirm(leadRollover, callerTerminal, str(args, "token"),
truthy(args, "operatorConfirmed"));
case "cancel" -> handoverCancel(leadRollover, str(args, "token"));
@@ -1498,7 +1507,8 @@ public final class FleetMcp {
* IllegalStateException} for that), both degrade to the same clean {@code NOT_CONFIGURED}
* refusal — never an escaping exception.
*/
private static McpSchema.CallToolResult handoverOpen(LeadRollover leadRollover, String callerTerminal,
private static McpSchema.CallToolResult handoverOpen(LeadRollover leadRollover, MessageService messages,
String callerTerminal, String callerOwner,
String reason) {
if (leadRollover == null) {
return notConfigured();
@@ -1516,6 +1526,9 @@ public final class FleetMcp {
m.put("token", p.token());
m.put("handoverPath", p.handoverPath());
m.put("requestedAtMillis", p.requestedAtMillis());
MessageService.Outstanding outstanding = messages.outstanding(callerOwner);
m.put("outstandingTickets", outstanding.tickets());
m.put("openAsks", outstanding.asks());
return text(json(m));
} catch (IllegalStateException e) {
// leadRollover: was removed from config by a hot reload since this FleetMcp was
@@ -2646,7 +2659,9 @@ public final class FleetMcp {
+ "then use this to have fleetd end your pane's process and relaunch a fresh "
+ "lead session bootstrapped against it. Four actions: 'open' (requests a "
+ "token and the handoverPath you must write the handover file to before "
+ "confirming), 'confirm' (validates every gate and — only if every one "
+ "confirming — the response also lists outstandingTickets and openAsks, "
+ "your own async delegations and fleet_ask turns, so their ids can go into "
+ "the handover file too), 'confirm' (validates every gate and — only if every one "
+ "passes — schedules the roll; it does NOT itself end your pane, the roll "
+ "runs once this call's own turn ends), 'cancel' (drops a pending request "
+ "without rolling), and 'status' (read-only: what happened to a token after "
@@ -6,6 +6,7 @@ import org.slf4j.LoggerFactory;
import java.util.Optional;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.Function;
/**
* Single-slot, thread-safe registry for the primary's herdr {@code terminal_id}.
@@ -18,13 +19,25 @@ import java.util.concurrent.atomic.AtomicReference;
* <p>The push loop ({@code ReplyPushLoop}) uses {@link #isKnown()} to decide
* whether active nudging is possible; an empty registry means the primary is
* off-host or non-herdr and delivery falls back to pull.
*
* <p><strong>A learned terminal can go stale; a configured lead's name cannot.</strong> A lead that
* is rolled (a fresh pane replacing the old one) keeps its name but gets a new {@code terminal_id}.
* So every terminal this class learns — the singleton and each per-target delegation — is recorded
* together with the delegating lead's name, when the caller carries one. {@link
* #currentPrimaryTerminal()} and {@link #nudgeTargetFor(String)} resolve that name back to a
* terminal through the live {@code currentTerminalForName} lookup before falling back to the
* terminal that was actually recorded. A caller with no name (an unnamed primary, an architect, a
* collaborator — none of those are leads a lookup keyed on lead names can resolve) is tracked by
* terminal alone, exactly as before this indirection existed.
*/
public final class PrimaryRegistry {
private static final Logger log = LoggerFactory.getLogger(PrimaryRegistry.class);
private final AtomicReference<String> terminal = new AtomicReference<>();
private final AtomicReference<String> primaryName = new AtomicReference<>();
private final boolean pinned;
private final Function<String, String> currentTerminalForName;
/**
* CB-532: worker terminal → the lead that delegated to it. The single slot above answers "who is
@@ -33,12 +46,32 @@ public final class PrimaryRegistry {
* other lead's delegations. This map answers the question that actually matters — "who is
* waiting on THIS worker" — and is what lets {@code primary.terminal} be retired.
*/
private final ConcurrentHashMap<String, String> leadByTarget = new ConcurrentHashMap<>();
private final ConcurrentHashMap<String, Delegation> leadByTarget = new ConcurrentHashMap<>();
/** A recorded delegator: the terminal learned from call traffic, and its name, if it has one. */
private record Delegation(String terminal, String name) {
}
/**
* @param pinnedTerminal an optional pinned terminal from config ({@code null}/blank = unpinned)
*/
public PrimaryRegistry(String pinnedTerminal) {
this(pinnedTerminal, name -> null);
}
/**
* As above, with a live {@code lead name → current terminal} lookup — normally the inverse of
* the same {@code terminal_id → lead name} supplier {@code CallerResolver} and the lead-tab
* scan already read. A lookup that cannot place a name (it is not a currently recognised lead,
* or no lookup is wired) returns {@code null}, and every resolution here falls back to the
* terminal that was actually recorded.
*
* @param pinnedTerminal an optional pinned terminal from config ({@code null}/blank =
* unpinned)
* @param currentTerminalForName lead name → its current terminal, or {@code null} if that name
* is not a currently recognised lead
*/
public PrimaryRegistry(String pinnedTerminal, Function<String, String> currentTerminalForName) {
if (pinnedTerminal != null && !pinnedTerminal.isBlank()) {
this.terminal.set(pinnedTerminal);
this.pinned = true;
@@ -46,19 +79,31 @@ public final class PrimaryRegistry {
} else {
this.pinned = false;
}
this.currentTerminalForName = currentTerminalForName != null ? currentTerminalForName : name -> null;
}
/**
* Record a terminal_id. No-op when:
* Record a terminal_id, with no lead name. No-op when:
* <ul>
* <li>the registry is pinned (config override),
* <li>{@code terminalId} is {@code null} or blank (non-herdr caller).
* </ul>
*/
public void record(String terminalId) {
record(terminalId, null);
}
/**
* As {@link #record(String)}, additionally recording the caller's name — present for a
* configured lead, {@code null} for an unnamed primary. The name is what lets {@link
* #currentPrimaryTerminal()} keep nudging the same lead across a roll even though its terminal
* changed.
*/
public void record(String terminalId, String name) {
if (pinned) return;
if (terminalId == null || terminalId.isBlank()) return;
String prev = terminal.getAndSet(terminalId);
primaryName.set(blankToNull(name));
if (prev == null) {
log.debug("primary terminal learned: {}", terminalId);
} else if (!prev.equals(terminalId)) {
@@ -67,7 +112,8 @@ public final class PrimaryRegistry {
}
/**
* Record that {@code leadTerminal} owns the accepted delegation of worker {@code target} (CB-532).
* Record that {@code leadTerminal} owns the accepted delegation of worker {@code target}
* (CB-532), with no lead name.
*
* <p>Called from the {@code MessageService} accepted-delivery hook — only after a send has won
* the session's send lock and queued delivery — where both halves are known (CB-548). It is
@@ -77,10 +123,19 @@ public final class PrimaryRegistry {
* lead that most recently delegated to it, which is the one waiting.
*/
public void recordDelegation(String target, String leadTerminal) {
recordDelegation(target, leadTerminal, null);
}
/**
* As {@link #recordDelegation(String, String)}, additionally recording the delegating lead's
* name when the caller carries one. See {@link #record(String, String)} for why the name
* matters.
*/
public void recordDelegation(String target, String leadTerminal, String leadName) {
if (target == null || target.isBlank() || leadTerminal == null || leadTerminal.isBlank()) {
return;
}
leadByTarget.put(target, leadTerminal);
leadByTarget.put(target, new Delegation(leadTerminal, blankToNull(leadName)));
}
/** Forget a worker's delegating lead — call on release, so a torn-down session leaks nothing. */
@@ -99,19 +154,59 @@ public final class PrimaryRegistry {
* recorded delegation there is no right answer, so this returns empty rather than guessing —
* delivery degrades to pull, which is exactly what the durable inbox is for, instead of
* interrupting the wrong lead with someone else's result.
*
* <p>A delegation recorded with a name is resolved to that lead's <em>current</em> terminal
* first — see {@link #currentTerminalForName} — so a lead that has since been rolled is still
* reachable here, not just the pane that delegated the work originally.
*/
public Optional<String> nudgeTargetFor(String target) {
String lead = target == null ? null : leadByTarget.get(target);
return lead != null ? Optional.of(lead) : Optional.ofNullable(terminal.get());
Delegation delegation = target == null ? null : leadByTarget.get(target);
if (delegation != null) {
return Optional.of(resolveCurrent(delegation.terminal(), delegation.name()));
}
return currentPrimaryTerminal();
}
/** The known primary terminal, or empty if not yet learned (and not pinned). */
/**
* The known primary terminal, or empty if not yet learned (and not pinned) — the raw value as
* it was recorded, with no attempt to resolve a named lead's current pane. Callers that need a
* nudge destination which survives a lead roll want {@link #currentPrimaryTerminal()} instead.
*/
public Optional<String> primaryTerminal() {
return Optional.ofNullable(terminal.get());
}
/**
* The terminal to nudge for the singleton primary right now: the recorded name resolved to its
* current terminal when one was recorded and is still a recognised lead, otherwise the terminal
* that was actually recorded — empty only when nothing has been learned or pinned at all.
*/
public Optional<String> currentPrimaryTerminal() {
String learned = terminal.get();
if (learned == null) {
return Optional.empty();
}
return Optional.of(resolveCurrent(learned, primaryName.get()));
}
/** {@code true} once a terminal has been recorded (or was pinned at construction). */
public boolean isKnown() {
return terminal.get() != null;
}
/**
* {@code learnedTerminal}, unless {@code name} is non-null and {@code currentTerminalForName}
* currently places that name at a different, live terminal — in which case the live one wins.
*/
private String resolveCurrent(String learnedTerminal, String name) {
if (name == null) {
return learnedTerminal;
}
String current = currentTerminalForName.apply(name);
return current != null && !current.isBlank() ? current : learnedTerminal;
}
private static String blankToNull(String s) {
return s == null || s.isBlank() ? null : s;
}
}
@@ -307,12 +307,12 @@ public final class LeadHeartbeatLoop {
* (mirroring {@link ReplyPushLoop#tick(String)}) so tests can drive it directly with a fake clock and a
* fake {@link AgentControl} instead of racing the scheduler thread. */
void tick() {
boolean leadKnown = primaryRegistry.primaryTerminal().isPresent();
boolean leadKnown = primaryRegistry.currentPrimaryTerminal().isPresent();
FleetState fleet = snapshot(inbox, roster);
AgentStatus status = AgentStatus.UNKNOWN;
LeadContextGauge.Reading reading = LeadContextGauge.Reading.unknown();
if (leadKnown) {
String leadTerminal = primaryRegistry.primaryTerminal().orElseThrow();
String leadTerminal = primaryRegistry.currentPrimaryTerminal().orElseThrow();
try {
status = agents.status(leadTerminal);
} catch (RuntimeException e) {
@@ -367,7 +367,7 @@ public final class LeadHeartbeatLoop {
// itself should be built from. Otherwise a HIGH stretch that is still latched would never see the
// notice at all, defeating the very check this fixes.
String notice = contextNotice(contextHighNudge, reading, contextNotified, requireOperatorConfirm);
var lead = primaryRegistry.primaryTerminal();
var lead = primaryRegistry.currentPrimaryTerminal();
boolean sent = lead.isPresent() && trySend(lead.get(), fleet.nudgeText() + notice, notice);
// The latch becomes true only when all three hold: decide() chose to notify, a notice was
// actually included in the text, and the send reached the pane without throwing. Whenever no
@@ -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);
}
/**
@@ -1797,6 +1814,67 @@ public final class MessageService {
return null;
}
/**
* One ticket {@code callerOwner} created, still present in {@link #tasks}, surfaced by
* {@link #outstanding} so a lead can carry its id into a handover file. {@link #phase} is the
* same value {@link #poll} would report right now, terminal phases included: a {@code DONE} or
* {@code FAILED} ticket stays in {@link #tasks} — and so stays reported here — until
* {@link #pruneTerminalTickets} evicts it.
*/
public record OutstandingTicket(String ticket, Phase phase, String target) {
}
/**
* One worker session paused in {@code fleet_ask}, with the {@code turnId} that answers it,
* surfaced by {@link #outstanding} alongside {@link OutstandingTicket}.
*/
public record OutstandingAsk(String ticket, String turnId, String workerSession) {
}
/** The outstanding tickets and open asks a single call to {@link #outstanding} reports. */
public record Outstanding(List<OutstandingTicket> tickets, List<OutstandingAsk> asks) {
}
/**
* Every ticket {@code callerOwner} created that is still in {@link #tasks} — including a
* finished one nobody has polled yet, since {@link #pruneTerminalTickets} discards its reply
* on a timer and a lead that does not carry its id forward can no longer read it after losing
* its session's context — plus the subset of those whose worker is paused in
* {@code fleet_ask}. Filtered by the same ownership rule as {@link #poll}:
* {@link #ownsTicket(Task, String)}.
*/
public Outstanding outstanding(String callerOwner) {
List<OutstandingTicket> tickets = new ArrayList<>();
List<OutstandingAsk> asks = new ArrayList<>();
for (Task task : tasks.values()) {
if (!ownsTicket(task, callerOwner)) {
continue;
}
Reply question = task.question;
Phase phase;
if (task.future.isDone()) {
phase = terminalPhase(task.future);
} else if (question != null) {
phase = Phase.ASKING;
asks.add(new OutstandingAsk(task.ticket, question.turnId(), task.target));
} else {
phase = Phase.PENDING;
}
tickets.add(new OutstandingTicket(task.ticket, phase, task.target));
}
return new Outstanding(tickets, asks);
}
/** As {@link #poll}'s own terminal-result handling, reduced to just the {@link Phase}. */
private static Phase terminalPhase(CompletableFuture<Reply> future) {
try {
Reply r = future.getNow(null);
return r != null && r.completed() ? Phase.DONE : Phase.FAILED;
} catch (CompletionException | java.util.concurrent.CancellationException e) {
return Phase.FAILED;
}
}
/** Release the async executor. */
public void close() {
asyncExecutor.shutdown();
@@ -439,6 +439,15 @@ public final class ReplyPushLoop {
* — timeout, transport error, a decode error — is treated as still live and the binding is left
* alone, because guessing wrong here is unrecoverable while guessing "live" merely costs one more
* retry on the next tick, which {@link #decide} already tolerates.
*
* <p><strong>The fallback is probed too.</strong> {@code PrimaryRegistry.nudgeTargetFor} already
* resolves a named delegator to its current terminal before this method ever sees it, which
* keeps a rolled lead's per-target binding live. What that resolution cannot fix is a caller
* that was never recorded with a name at all — an unnamed primary, or a lead whose tab the
* scanner cannot currently see — where the fallback it returns is still the raw terminal last
* learned from call traffic. This method returns that fallback only after the same liveness
* check, and gives up for this tick (an empty result, exactly like "no lead known at all") rather
* than hand a caller a second stale address un-probed.
*/
private Optional<String> resolveLiveLead(String target) {
Optional<String> lead = primaryRegistry.nudgeTargetFor(target);
@@ -448,7 +457,12 @@ public final class ReplyPushLoop {
log.debug("push: lead {} delegated to for {} is no longer live, forgetting the stale binding "
+ "and falling back", lead.get(), target);
primaryRegistry.forgetDelegation(target);
return primaryRegistry.nudgeTargetFor(target);
Optional<String> fallback = primaryRegistry.nudgeTargetFor(target);
if (fallback.isEmpty() || isLive(fallback.get())) {
return fallback;
}
log.debug("push: fallback lead {} for {} is also not live, skipping this tick", fallback.get(), target);
return Optional.empty();
}
/**
@@ -883,8 +883,7 @@ 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 owner key created that delegation, or
// to the unnamed primary.
// Phase.ASKING view, but only to the caller whose owner key created that delegation.
Principal caller = ctx.attribute(CALLER);
MessageService.PendingAsk ask = messages.pendingAsk(id, caller == null ? null : caller.ownerKey());
if (ask != null) {
@@ -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());
}
}
@@ -56,9 +56,14 @@ class FleetMcpHandoverTest {
private static final String LEAD = "term_lead";
private static final String OTHER_LEAD = "term_other_lead";
private static final String LEAD_OWNER = "leader:lead";
private final FakeHerdr herdr = new FakeHerdr();
private final AgentControl agents = new AgentControl(herdr);
/** Fed to every direct {@code FleetMcp.handover} call below — none of this class's own tests
* exercise ticket/ask ownership, so a single instance with no delegations is enough. */
private final MessageService messages = new MessageService(agents, new Injector(agents),
new Rendezvous(), new InMemoryReplyInbox());
private FleetMcp mcp;
@AfterEach
@@ -151,16 +156,16 @@ class FleetMcpHandoverTest {
@DisplayName("with leadRollover: absent, every action returns a clean NOT_CONFIGURED refusal and never throws")
void nullLeadRolloverRefusesCleanlyForEveryAction() {
McpSchema.CallToolResult open = assertDoesNotThrow(
() -> FleetMcp.handover(null, LEAD, Map.of("action", "open")));
() -> FleetMcp.handover(null, messages, LEAD, LEAD_OWNER, Map.of("action", "open")));
assertFalse(open.isError(), "a refusal is not a protocol error: " + textOf(open));
assertTrue(textOf(open).contains("NOT_CONFIGURED"), textOf(open));
McpSchema.CallToolResult confirm = assertDoesNotThrow(() -> FleetMcp.handover(null, LEAD,
McpSchema.CallToolResult confirm = assertDoesNotThrow(() -> FleetMcp.handover(null, messages, LEAD, LEAD_OWNER,
Map.of("action", "confirm", "token", "whatever")));
assertFalse(confirm.isError());
assertTrue(textOf(confirm).contains("NOT_CONFIGURED"), textOf(confirm));
McpSchema.CallToolResult cancel = assertDoesNotThrow(() -> FleetMcp.handover(null, LEAD,
McpSchema.CallToolResult cancel = assertDoesNotThrow(() -> FleetMcp.handover(null, messages, LEAD, LEAD_OWNER,
Map.of("action", "cancel", "token", "whatever")));
assertFalse(cancel.isError());
assertTrue(textOf(cancel).contains("NOT_CONFIGURED"), textOf(cancel));
@@ -169,11 +174,11 @@ class FleetMcpHandoverTest {
@Test
@DisplayName("a blank/unknown action is a clean tool error, never an exception")
void unknownActionIsACleanError() {
McpSchema.CallToolResult missing = assertDoesNotThrow(() -> FleetMcp.handover(null, LEAD, Map.of()));
McpSchema.CallToolResult missing = assertDoesNotThrow(() -> FleetMcp.handover(null, messages, LEAD, LEAD_OWNER, Map.of()));
assertTrue(missing.isError());
McpSchema.CallToolResult bogus = assertDoesNotThrow(
() -> FleetMcp.handover(null, LEAD, Map.of("action", "bogus")));
() -> FleetMcp.handover(null, messages, LEAD, LEAD_OWNER, Map.of("action", "bogus")));
assertTrue(bogus.isError());
}
@@ -229,7 +234,7 @@ class FleetMcpHandoverTest {
Files.writeString(handover, "not written yet");
LeadRollover rollover = newRollover(handover.toString());
McpSchema.CallToolResult openResult = FleetMcp.handover(rollover, LEAD, Map.of("action", "open"));
McpSchema.CallToolResult openResult = FleetMcp.handover(rollover, messages, LEAD, LEAD_OWNER, Map.of("action", "open"));
assertFalse(openResult.isError(), textOf(openResult));
String token = extractToken(textOf(openResult));
@@ -238,13 +243,13 @@ class FleetMcpHandoverTest {
Thread.sleep(50);
Files.writeString(handover, "the real handover content");
McpSchema.CallToolResult wrongCaller = FleetMcp.handover(rollover, OTHER_LEAD,
McpSchema.CallToolResult wrongCaller = FleetMcp.handover(rollover, messages, OTHER_LEAD, "leader:other-lead",
Map.of("action", "confirm", "token", token));
assertFalse(wrongCaller.isError(), "a refusal is a legitimate outcome, not a protocol error");
assertTrue(textOf(wrongCaller).contains("NOT_YOUR_ROLLOVER"),
"a different lead terminal confirming must surface NOT_YOUR_ROLLOVER: " + textOf(wrongCaller));
McpSchema.CallToolResult confirmed = FleetMcp.handover(rollover, LEAD,
McpSchema.CallToolResult confirmed = FleetMcp.handover(rollover, messages, LEAD, LEAD_OWNER,
Map.of("action", "confirm", "token", token));
assertFalse(confirmed.isError(), textOf(confirmed));
assertTrue(textOf(confirmed).contains("\"accepted\":true"),
@@ -258,7 +263,7 @@ class FleetMcpHandoverTest {
void cancelUnknownTokenIsCleanNotAFailure() {
LeadRollover rollover = newRollover(tmp.resolve("h.md").toString());
McpSchema.CallToolResult r = FleetMcp.handover(rollover, LEAD,
McpSchema.CallToolResult r = FleetMcp.handover(rollover, messages, LEAD, LEAD_OWNER,
Map.of("action", "cancel", "token", "does-not-exist"));
assertFalse(r.isError());
assertTrue(textOf(r).contains("\"cancelled\":false"), textOf(r));
@@ -269,9 +274,9 @@ class FleetMcpHandoverTest {
@DisplayName("cancel on a token actually opened reports cancelled:true")
void cancelKnownTokenSucceeds() {
LeadRollover rollover = newRollover(tmp.resolve("h.md").toString());
String token = extractToken(textOf(FleetMcp.handover(rollover, LEAD, Map.of("action", "open"))));
String token = extractToken(textOf(FleetMcp.handover(rollover, messages, LEAD, LEAD_OWNER, Map.of("action", "open"))));
McpSchema.CallToolResult r = FleetMcp.handover(rollover, LEAD,
McpSchema.CallToolResult r = FleetMcp.handover(rollover, messages, LEAD, LEAD_OWNER,
Map.of("action", "cancel", "token", token));
assertFalse(r.isError());
assertTrue(textOf(r).contains("\"cancelled\":true"), textOf(r));
@@ -282,7 +287,7 @@ class FleetMcpHandoverTest {
@Test
@DisplayName("status on a null LeadRollover is a clean NOT_CONFIGURED refusal, never a throw")
void statusWithNullLeadRolloverRefusesCleanly() {
McpSchema.CallToolResult r = assertDoesNotThrow(() -> FleetMcp.handover(null, LEAD,
McpSchema.CallToolResult r = assertDoesNotThrow(() -> FleetMcp.handover(null, messages, LEAD, LEAD_OWNER,
Map.of("action", "status", "token", "whatever")));
assertFalse(r.isError());
assertTrue(textOf(r).contains("NOT_CONFIGURED"), textOf(r));
@@ -293,7 +298,7 @@ class FleetMcpHandoverTest {
void statusOnUnknownTokenReportsUnknown() {
LeadRollover rollover = newRollover(tmp.resolve("h.md").toString());
McpSchema.CallToolResult r = FleetMcp.handover(rollover, LEAD,
McpSchema.CallToolResult r = FleetMcp.handover(rollover, messages, LEAD, LEAD_OWNER,
Map.of("action", "status", "token", "does-not-exist"));
assertFalse(r.isError());
assertTrue(textOf(r).contains("\"state\":\"UNKNOWN\""), textOf(r));
@@ -303,9 +308,9 @@ class FleetMcpHandoverTest {
@DisplayName("status on a token that is still pending (opened, not confirmed) reports PENDING")
void statusOnPendingTokenReportsPending() {
LeadRollover rollover = newRollover(tmp.resolve("h.md").toString());
String token = extractToken(textOf(FleetMcp.handover(rollover, LEAD, Map.of("action", "open"))));
String token = extractToken(textOf(FleetMcp.handover(rollover, messages, LEAD, LEAD_OWNER, Map.of("action", "open"))));
McpSchema.CallToolResult r = FleetMcp.handover(rollover, LEAD,
McpSchema.CallToolResult r = FleetMcp.handover(rollover, messages, LEAD, LEAD_OWNER,
Map.of("action", "status", "token", token));
assertFalse(r.isError());
assertTrue(textOf(r).contains("\"state\":\"PENDING\""), textOf(r));
@@ -323,5 +328,40 @@ class FleetMcpHandoverTest {
"the tool's own description must advertise the 'status' action: " + tool.description());
}
// --- unit 5: open() reports outstanding tickets and open asks ------------------------------
@Test
@DisplayName("open() keeps token/handoverPath/requestedAtMillis and reports empty outstanding collections when the caller has nothing")
void openReportsEmptyOutstandingCollectionsWhenCallerHasNothing() {
LeadRollover rollover = newRollover(tmp.resolve("h.md").toString());
McpSchema.CallToolResult r = FleetMcp.handover(rollover, messages, LEAD, LEAD_OWNER,
Map.of("action", "open"));
assertFalse(r.isError(), textOf(r));
String json = textOf(r);
assertTrue(json.contains("\"token\":"), json);
assertTrue(json.contains("\"handoverPath\":"), json);
assertTrue(json.contains("\"requestedAtMillis\":"), json);
assertTrue(json.contains("\"outstandingTickets\":[]"),
"a caller with nothing gets an empty array, not an absent key: " + json);
assertTrue(json.contains("\"openAsks\":[]"),
"a caller with nothing gets an empty array, not an absent key: " + json);
}
@Test
@DisplayName("open() reports an owned pending ticket with its phase and target")
void openReportsAnOwnedPendingTicketWithItsPhase() {
LeadRollover rollover = newRollover(tmp.resolve("h.md").toString());
String ticket = messages.sendAsync("term_worker", "a task", null, Principal.leader("lead", LEAD, 1));
McpSchema.CallToolResult r = FleetMcp.handover(rollover, messages, LEAD, LEAD_OWNER,
Map.of("action", "open"));
assertFalse(r.isError(), textOf(r));
String json = textOf(r);
assertTrue(json.contains("\"ticket\":\"" + ticket + "\""), json);
assertTrue(json.contains("\"phase\":\"PENDING\""), json);
assertTrue(json.contains("\"target\":\"term_worker\""), json);
}
// --- acceptance 7 (wiring) is covered by FleetdLeadRolloverWiringTest, unchanged -----------
}
@@ -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();
@@ -169,4 +169,101 @@ class PrimaryRegistryTest {
assertTrue(reg.nudgeTargetFor("term_worker").isEmpty());
assertTrue(reg.nudgeTargetFor(null).isEmpty());
}
// ── fleetd #737 unit 3: a named lead's terminal is resolved live, not just recorded ─────────
/**
* The whole point of carrying a name: a lead that has been rolled keeps its name but gets a
* fresh terminal. {@code currentTerminalForName} stands in for the live lead-tab scan here —
* it reports the lead now sits on a different terminal than the one that was recorded — and
* {@code nudgeTargetFor} must follow the name to that current terminal, not the stale one.
*/
@Test
void nudgeTargetForFollowsARolledLeadsNameToItsCurrentTerminal() {
var reg = new PrimaryRegistry(null, name -> "opus".equals(name) ? "term_opus_after_roll" : null);
reg.recordDelegation("term_worker", "term_opus_before_roll", "opus");
assertEquals("term_opus_after_roll", reg.nudgeTargetFor("term_worker").orElseThrow(),
"the name must be resolved to the lead's CURRENT terminal, not the one recorded "
+ "at delegation time");
}
/**
* The lookup cannot place every name — an architect/collaborator name (never a lead), or a lead
* whose tab the scan cannot currently see (just rolled, off-host, non-herdr). Either way the
* terminal actually recorded is still the right thing to try, exactly as before this unit.
*/
@Test
void nudgeTargetForFallsBackToTheRecordedTerminalWhenTheNameCannotBePlaced() {
var reg = new PrimaryRegistry(null, name -> null); // nothing is ever currently recognised
reg.recordDelegation("term_worker", "term_lead_recorded", "opus");
assertEquals("term_lead_recorded", reg.nudgeTargetFor("term_worker").orElseThrow());
}
/** The 2-arg {@code recordDelegation} overload records no name, so resolution never applies. */
@Test
void recordDelegationWithNoNameIsNeverResolvedByLookup() {
var reg = new PrimaryRegistry(null, name -> {
throw new AssertionError("a delegation recorded with no name must never consult the lookup");
});
reg.recordDelegation("term_worker", "term_lead");
assertEquals("term_lead", reg.nudgeTargetFor("term_worker").orElseThrow());
}
/** As {@link #nudgeTargetForFollowsARolledLeadsNameToItsCurrentTerminal}, for the singleton. */
@Test
void currentPrimaryTerminalFollowsARolledLeadsNameToItsCurrentTerminal() {
var reg = new PrimaryRegistry(null, name -> "sol".equals(name) ? "term_sol_after_roll" : null);
reg.record("term_sol_before_roll", "sol");
assertEquals("term_sol_after_roll", reg.currentPrimaryTerminal().orElseThrow());
assertEquals("term_sol_before_roll", reg.primaryTerminal().orElseThrow(),
"primaryTerminal() stays the raw recorded value — currentPrimaryTerminal() is the "
+ "one that resolves live");
}
@Test
void currentPrimaryTerminalFallsBackWhenTheNameCannotBePlaced() {
var reg = new PrimaryRegistry(null, name -> null);
reg.record("term_sol", "sol");
assertEquals("term_sol", reg.currentPrimaryTerminal().orElseThrow());
}
@Test
void currentPrimaryTerminalWithNoNameRecordedIsTheRawTerminal() {
var reg = new PrimaryRegistry(null, name -> {
throw new AssertionError("no name was ever recorded, the lookup must not be consulted");
});
reg.record("term_x");
assertEquals("term_x", reg.currentPrimaryTerminal().orElseThrow());
}
@Test
void currentPrimaryTerminalIsEmptyWhenNothingWasEverLearned() {
var reg = new PrimaryRegistry(null, name -> "anything");
assertTrue(reg.currentPrimaryTerminal().isEmpty());
}
/** A pin never carries a name, so a pinned registry's singleton resolution is always a no-op. */
@Test
void currentPrimaryTerminalForAPinIsNeverResolvedByLookup() {
var reg = new PrimaryRegistry("term_pinned", name -> {
throw new AssertionError("a pin carries no name, the lookup must not be consulted");
});
assertEquals("term_pinned", reg.currentPrimaryTerminal().orElseThrow());
}
/** {@code nudgeTargetFor}'s fallback to the singleton is the resolved one, not the raw one. */
@Test
void nudgeTargetForWithNoDelegationFallsBackToTheResolvedSingleton() {
var reg = new PrimaryRegistry(null, name -> "opus".equals(name) ? "term_opus_after_roll" : null);
reg.record("term_opus_before_roll", "opus");
assertEquals("term_opus_after_roll", reg.nudgeTargetFor("term_never_seen").orElseThrow());
}
}
@@ -648,6 +648,42 @@ class LeadHeartbeatLoopTest {
"the notice text must appear exactly once across all three sends: " + herdr.sentTexts());
}
// ── fleetd #737 unit 3: tick() nudges the lead's CURRENT terminal, not the learned one ───────
/**
* {@code tick()} reads {@code primaryRegistry.currentPrimaryTerminal()} both to check the lead's
* status and to send the nudge. Here the registry learned the lead's terminal under its name
* before a roll; {@code currentTerminalForName} stands in for the live lead-tab scan and reports
* the lead now sits on a different terminal. A correct tick must follow the name and nudge the
* new terminal — nudging the old one would mean the heartbeat lost the lead across its own roll.
*/
@Test
void tickNudgesTheLeadsCurrentTerminalAfterARoll() {
var herdr = new FailableHerdrClient("term_lead_after_roll");
var now = new AtomicLong(NOW);
InMemoryReplyInbox inbox = new InMemoryReplyInbox();
inbox.own(WORKER);
inbox.publish(WORKER, "m1", "hello");
@SuppressWarnings("unchecked")
List<MemberSession>[] rosterBox = new List[]{List.of(new MemberSession("p1", WORKER, "prof",
MemberRole.DEV, "/cwd", null, 0, 0, 0, MemberSession.State.READY, null, null))};
AgentControl agents = new AgentControl(herdr);
PrimaryRegistry registry = new PrimaryRegistry(null,
name -> "opus".equals(name) ? "term_lead_after_roll" : null);
registry.record("term_lead_before_roll", "opus");
ReplyPushLoop pushLoop = new ReplyPushLoop(registry, agents, inbox, scheduler, 5, 100_000);
LeadHeartbeatLoop loop = new LeadHeartbeatLoop(registry, agents, inbox, () -> rosterBox[0], pushLoop,
scheduler, now::get, IDLE_AFTER_NANOS, 100_000L, 0);
loop.tick(); // opens the idle window
now.addAndGet(TimeUnit.SECONDS.toNanos(400));
loop.tick(); // past the quiet period, pending reply -> INJECT
assertEquals(List.of("term_lead_after_roll"), herdr.promptTargets(),
"the heartbeat must read and nudge the lead's CURRENT terminal, not the one learned "
+ "before the roll");
}
/**
* Fake herdr client for the four tests above: always reports {@code lead} as IDLE, records the
* {@code text} of every {@code agent.prompt} call, and can be told to throw on the very next
@@ -657,6 +693,7 @@ class LeadHeartbeatLoopTest {
private static final ObjectMapper MAPPER = new ObjectMapper();
private final String lead;
private final List<String> sentTexts = new ArrayList<>();
private final List<String> promptTargets = new ArrayList<>();
private boolean throwOnNextSend = false;
FailableHerdrClient(String lead) {
@@ -671,6 +708,11 @@ class LeadHeartbeatLoopTest {
return List.copyOf(sentTexts);
}
/** Every terminal an {@code agent.prompt} call named, in call order. */
List<String> promptTargets() {
return List.copyOf(promptTargets);
}
@Override
@SuppressWarnings("unchecked")
public JsonNode call(String method, Object params) {
@@ -681,12 +723,13 @@ class LeadHeartbeatLoopTest {
.put("agent_status", "idle"));
}
if ("agent.prompt".equals(method)) {
Map<String, Object> p = params instanceof Map ? (Map<String, Object>) params : Map.of();
if (throwOnNextSend) {
throwOnNextSend = false;
throw new RuntimeException("simulated transient herdr send failure");
}
Map<String, Object> p = params instanceof Map ? (Map<String, Object>) params : Map.of();
sentTexts.add(String.valueOf(p.get("text")));
promptTargets.add(String.valueOf(p.get("target")));
}
return MAPPER.createObjectNode();
}
@@ -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());
@@ -2037,6 +2137,118 @@ class MessageServiceTest {
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
}
// --- outstanding(): fleet_handover{open}'s list of open tickets and asks -------------------
@Test
void outstandingReportsOwnedPendingTicketsWithPhases() {
Principal lead = Principal.leader("opus", "term_lead", 1);
String ticket1 = messages.sendAsync(T, "task one", null, lead);
String ticket2 = messages.sendAsync(T, "task two", null, lead);
MessageService.Outstanding outstanding = messages.outstanding(lead.ownerKey());
assertEquals(2, outstanding.tickets().size());
assertTrue(outstanding.tickets().stream().anyMatch(t ->
ticket1.equals(t.ticket()) && t.phase() == MessageService.Phase.PENDING && T.equals(t.target())),
"ticket1 must be reported PENDING: " + outstanding.tickets());
assertTrue(outstanding.tickets().stream().anyMatch(t ->
ticket2.equals(t.ticket()) && t.phase() == MessageService.Phase.PENDING && T.equals(t.target())),
"ticket2 must be reported PENDING: " + outstanding.tickets());
assertTrue(outstanding.asks().isEmpty(), "neither ticket has an open question");
}
@Test
void outstandingReportsAnOpenAsksTurnId() throws Exception {
Principal lead = Principal.leader("opus", "term_lead", 1);
String ticket = messages.sendAsync(T, "task that asks", null, lead);
awaitWaiting();
injectDelivery();
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
MessageService.Outstanding outstanding = messages.outstanding(lead.ownerKey());
assertEquals(1, outstanding.tickets().size());
assertEquals(MessageService.Phase.ASKING, outstanding.tickets().get(0).phase());
assertEquals(1, outstanding.asks().size());
MessageService.OutstandingAsk openAsk = outstanding.asks().get(0);
assertEquals(ticket, openAsk.ticket());
assertEquals(asking.turnId(), openAsk.turnId());
assertEquals(T, openAsk.workerSession());
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(asking.turnId(), "config.yaml", 5000, lead.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: {@code leadA} must still see its own ticket and ask, so {@code leadB}
* seeing neither is the owner filter at work and not {@code outstanding} refusing everyone.
*/
@Test
void outstandingDoesNotLeakAcrossNamedLeads() throws Exception {
Principal leadA = Principal.leader("opus", "term_a", 1);
Principal leadB = Principal.leader("sol", "term_b", 2);
String ticket = messages.sendAsync(T, "task that asks", null, leadA);
awaitWaiting();
injectDelivery();
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
MessageService.Outstanding seenByA = messages.outstanding(leadA.ownerKey());
assertEquals(1, seenByA.tickets().size(), "lead A must see its own ticket");
assertEquals(1, seenByA.asks().size(), "lead A must see its own open ask");
MessageService.Outstanding seenByB = messages.outstanding(leadB.ownerKey());
assertTrue(seenByB.tickets().isEmpty(), "lead B must not see lead A's ticket");
assertTrue(seenByB.asks().isEmpty(), "lead B must not see lead A's open ask");
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(asking.turnId(), "config.yaml", 5000, leadA.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());
}
@Test
void outstandingIsEmptyCollectionsNotNullForACallerWithNothing() {
MessageService.Outstanding outstanding =
messages.outstanding(Principal.leader("opus", "term_lead", 1).ownerKey());
assertNotNull(outstanding.tickets(), "a caller with no delegations still gets a list, not null");
assertNotNull(outstanding.asks(), "a caller with no delegations still gets a list, not null");
assertTrue(outstanding.tickets().isEmpty());
assertTrue(outstanding.asks().isEmpty());
}
/**
* A ticket is destroyed on a timer once it goes terminal ({@link #pruneTerminalTickets}'s TTL
* runs from completion), so it is the one case where carrying the id forward actually matters
* — a still-PENDING ticket is in no such danger, its worker is still running.
*/
@Test
void outstandingReportsACompletedUncollectedTicketWithATerminalPhase() throws Exception {
Principal lead = Principal.leader("opus", "term_lead", 1);
String ticket = messages.sendAsync(T, "long task", null, lead);
awaitWaiting();
injectDelivery();
assertTrue(rendezvous.resolve(T, "async result"), "a reply resolves the async send");
awaitTicketPhase(ticket, MessageService.Phase.DONE);
MessageService.Outstanding outstanding = messages.outstanding(lead.ownerKey());
assertEquals(1, outstanding.tickets().size());
MessageService.OutstandingTicket done = outstanding.tickets().get(0);
assertEquals(ticket, done.ticket());
assertEquals(MessageService.Phase.DONE, done.phase());
assertEquals(T, done.target());
}
// --- fleetd #715: answer() is gated on the caller that owns the turn -----------------------
/**
@@ -304,6 +304,40 @@ class ReplyPushLoopTest {
+ "still live");
}
// --- fleetd #737 unit 3: the fallback nudgeTargetFor returns must be probed too --------------
/**
* {@code resolveLiveLead} forgets a dead per-target binding and asks {@code PrimaryRegistry}
* again for a fallback. That fallback can be dead too — here PRIMARY, the pinned singleton, is
* affirmatively gone alongside DEAD_LEAD. A correct {@code resolveLiveLead} probes it exactly
* like the first lead and gives up for this tick rather than trust it unchecked.
*
* <p>This is checked through {@code onReplyQueued}/{@code isActive} rather than a send count:
* {@link ReplyPushLoop#onReplyQueued} only registers pending work and starts a schedule once
* {@code resolveLiveLead} returns a present value — a dead fallback that was trusted unprobed
* would already make this true, synchronously, with no tick or send needed to observe it. Pairs
* with {@link #aStaleLeadBindingFallsBackToTheLiveLeadInsteadOfNudgingADeadTerminal} as the
* positive control: same stale-DEAD_LEAD setup, but there the fallback (PRIMARY) is live and the
* nudge does fire — proving this test's "nothing happens" result comes from the fallback being
* dead, not from the assertion being unable to observe a nudge at all.
*/
@Test
void aDoublyDeadFallbackIsNeverTrustedAndStartsNoSchedule() {
registry.recordDelegation(WORKER, DEAD_LEAD);
var rec = new AllDeadHerdrClient(Set.of(DEAD_LEAD, PRIMARY));
agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello");
var loop = loop(1, 50);
loop.onReplyQueued(WORKER);
assertFalse(loop.isActive(),
"both the per-target binding and the fallback are dead, so resolveLiveLead must "
+ "return empty and onReplyQueued must never register pending work or start "
+ "a schedule — an unprobed fallback would start one here");
assertEquals(0, rec.promptTargets().size(), "nobody live was found, so nothing was ever sent");
}
// --- nudge format --------------------------------------------------------------------------
@Test
@@ -1248,6 +1282,52 @@ class ReplyPushLoopTest {
}
}
/**
* Fake herdr client for fleetd #737 unit 3: every terminal named in {@code deadTargets} reports
* {@code agent_not_found} from {@code agent.get} — unlike {@link DeadLeadHerdrClient}, which can
* only make one terminal dead, this can make a per-target binding AND its fallback dead in the
* same test. {@code agent.prompt} is recorded unconditionally (no liveness check of its own),
* so a test can tell "resolveLiveLead probed and correctly found nobody live" (no prompt call)
* apart from "resolveLiveLead trusted a dead fallback and sent into it anyway" (a prompt call to
* a terminal this fake has already declared gone).
*/
private static final class AllDeadHerdrClient implements HerdrClient {
private final Set<String> deadTargets;
private final List<String> promptTargets = Collections.synchronizedList(new ArrayList<>());
AllDeadHerdrClient(Set<String> deadTargets) {
this.deadTargets = deadTargets;
}
@Override
@SuppressWarnings("unchecked")
public JsonNode call(String method, Object params) {
Map<String, Object> p = params instanceof Map ? (Map<String, Object>) params : Map.of();
if ("agent.get".equals(method)) {
String target = String.valueOf(p.get("target"));
if (deadTargets.contains(target)) {
throw new HerdrException("no such agent: " + target, "agent_not_found", null);
}
return MAPPER.createObjectNode()
.set("agent", MAPPER.createObjectNode()
.put("terminal_id", target)
.put("agent_status", "idle"));
}
if ("agent.prompt".equals(method)) {
promptTargets.add(String.valueOf(p.get("target")));
}
return MAPPER.createObjectNode();
}
List<String> promptTargets() {
return List.copyOf(promptTargets);
}
@Override
public void close() {
}
}
/**
* Fake herdr client for fleetd #368 review: {@code flakyTarget}'s FIRST {@code agent.get} call
* fails with a transient, non-{@code agent_not_found} {@code HerdrException} — a transport-level
@@ -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());