Compare commits

...

12 Commits

Author SHA1 Message Date
Dai Ha 6677ec8c63 fleetd #749: pin PackageCyclesTest's exceptions to exact edges, not whole packages
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 1m8s
CI / build (pull_request) Failing after 2m25s
ignoreCycle() used to exempt every dependency between two packages, in both
directions, for the whole package. That meant a brand new dependency added
later between an already-excepted pair (auth/mcp, mcp/msg, inject/msg,
metrics/msg, msg/session) was silently exempted too, exactly where the
msg package makes the gate matter most.

Replace the package-wide ignore with a frozen baseline of the 45 exact
origin-class -> target-class edges that exist today between those five
pairs, and ignore only those via SliceRule.ignoreDependency(String, String).
A new dependency between a baselined pair is not in that set, so it is no
longer ignored and the existing beFreeOfCycles() check (or, when the new
edge alone would not form a cycle, a dedicated set-equality check) fails
and names the exact origin class, target class and package pair.

The set-equality check also fails on a baseline entry whose dependency no
longer exists in the code, so a removed edge cannot rot in the baseline
and mask the pair's eligibility for the ticket #131 removal steps. Rewrote
the javadoc to describe only the current contract.
2026-10-05 07:45:51 +02:00
Dai Ha 7f0c4a8464 fleetd #737: the handover skill carries the outstanding tickets forward
CI / shell-tests (push) Failing after 12s
CI / contract (push) Successful in 55s
CI / build (push) Failing after 2m7s
fleet_handover{open} now returns outstandingTickets and openAsks. The
successor keeps the authority to poll and answer them but not the ids, and
an uncollected terminal ticket loses its reply at the ticket TTL.
2026-10-05 06:09:26 +02:00
Dai Ha 682991a846 Merge branch 'worker/737-9c61d3-4'
CI / shell-tests (push) Failing after 10s
CI / contract (push) Successful in 50s
CI / build (push) Failing after 1m52s
2026-10-05 06:07:33 +02:00
Dai Ha abe617c48c Merge branch 'worker/737-a263f3-3'
CI / shell-tests (push) Failing after 6s
CI / contract (push) Successful in 1m3s
CI / build (push) Failing after 2m20s
2026-10-05 06:01:49 +02:00
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 8a1d73b39e fleetd #737 unit 4: key the rollover single-flight claim on the lead's name
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 52s
CI / build (pull_request) Failing after 1m59s
rollingByTerminal keyed the single-flight lock on p.leadTerminal(), the pane
address. A roll replaces the pane, so a second roll of the same lead opened
from the new terminal landed on a different map key and could run concurrent
with the first roll's still-in-flight continuation.

PendingRollover now carries rolloverKey, resolved once in open() from
leadNameForTerminal while the lead is certainly still live, falling back to
the terminal itself when the name resolves null or blank. confirm()'s claim
and both release sites (the continuationRunner-rejection catch and
runRollover's finally) use the carried key instead of recomputing it, since
leadNameForTerminal no longer resolves the old terminal by release time.
NOT_YOUR_ROLLOVER and the pane-teardown calls stay keyed on p.leadTerminal(),
unchanged.

rollingByTerminal is renamed rollingByLead to match.
2026-10-05 05:57:59 +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
17 changed files with 979 additions and 135 deletions
+11 -2
View File
@@ -133,8 +133,17 @@ A handover file is a record of state and decisions. It is not a diary.
**Run the three steps in this order. The order is not a style choice — the wrong order is
refused.**
1. **`fleet_handover{action: "open", reason: "<why now>"}`.** It returns a `token` and the
`handoverPath` you must write to. Nothing has happened to your pane yet.
1. **`fleet_handover{action: "open", reason: "<why now>"}`.** It returns a `token`, the
`handoverPath` you must write to, and `outstandingTickets` plus `openAsks`. Nothing has happened
to your pane yet.
**Copy `outstandingTickets` and `openAsks` into the handover file.** Your successor keeps the
authority to poll those tickets and answer those asks, because both are gated on the lead's
name, which does not change when your pane does. What it does not keep is the ids — they exist
only in your context and in this response. A ticket already in a terminal phase is the urgent
one: its reply lives only in memory and is deleted once the ticket TTL passes, so an uncollected
report is lost for good. Poll those before you confirm, or name them in the file so your
successor polls them first.
**Write to exactly that path, and do not resolve it yourself.** It is always absolute, even when
the operator configured a relative `handoverPath`: fleetd resolves a relative one against your
@@ -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 "
@@ -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,9 +149,17 @@ 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) {}
long requestedAtMillis, String rolloverKey) {}
/** Which check refused a {@link #confirm} call, named so a caller can act on it. */
public enum RefusalReason {
@@ -176,9 +184,9 @@ public final class LeadRollover {
*/
HANDOVER_STALE,
/**
* This lead terminal already has a roll running: an earlier {@link #confirm} call claimed
* it and that roll's continuation has not released it yet. {@code detail} names the lead
* terminal and the token that holds the claim.
* 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.
*/
ROLL_ALREADY_RUNNING
}
@@ -340,9 +348,12 @@ 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. The deferred continuation calls this, on the
* OLD terminal, before tearing it down, so it knows which lead to pass to {@link
* LeadLauncher#relaunch}.
* 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.
*/
private final Function<String, String> leadNameForTerminal;
/**
@@ -362,14 +373,14 @@ public final class LeadRollover {
private final Consumer<Runnable> continuationRunner;
private final Map<String, PendingRollover> pending = new ConcurrentHashMap<>();
/**
* Lead terminal → the token of the roll currently holding that terminal exclusive, for
* {@link #confirm}'s single-flight claim. {@link #confirm} claims an entry here with an
* atomic put-if-absent once every other gate has passed, refusing with {@link
* {@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
* RefusalReason#ROLL_ALREADY_RUNNING} when a claim is already held; {@link #runRollover}
* releases it in a {@code finally}, on both the success and the thrown-exception path. A
* terminal absent from this map has no roll currently in flight for it.
* 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.
*/
private final Map<String, String> rollingByTerminal = new ConcurrentHashMap<>();
private final Map<String, String> rollingByLead = 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
@@ -438,8 +449,9 @@ 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), and {@code leadTerminal} — only that exact terminal
* may later {@link #confirm} this token.
* 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.
*
* @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
@@ -459,7 +471,9 @@ public final class LeadRollover {
String token = UUID.randomUUID().toString();
long requestedAt = nowMillis.getAsLong();
String resolvedPath = resolveHandoverPath(cfg.handoverPath(), leadTerminal);
PendingRollover p = new PendingRollover(token, leadTerminal, resolvedPath, requestedAt);
String leadName = leadNameForTerminal.apply(leadTerminal);
String rolloverKey = (leadName == null || leadName.isBlank()) ? leadTerminal : leadName;
PendingRollover p = new PendingRollover(token, leadTerminal, resolvedPath, requestedAt, rolloverKey);
pending.put(token, p);
if (resolvedPath.equals(cfg.handoverPath())) {
log.info("lead-rollover: open token={} lead={} handoverPath={} reason={}",
@@ -565,12 +579,11 @@ 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 terminal.
String holder = rollingByTerminal.putIfAbsent(p.leadTerminal(), token);
// different, still-running roll already holds this lead's claim.
String holder = rollingByLead.putIfAbsent(p.rolloverKey(), token);
if (holder != null) {
return RollDecision.refused(RefusalReason.ROLL_ALREADY_RUNNING,
"lead terminal " + p.leadTerminal() + " already has a roll running under token "
+ holder);
"lead '" + p.rolloverKey() + "' already has a roll running under token " + holder);
}
// Record IN_PROGRESS BEFORE removing from `pending` — see RollState#IN_PROGRESS and
@@ -591,14 +604,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 rollingByTerminal — never runs either.
// finally — the only other place that releases rollingByLead — never runs either.
// Release the claim here and overwrite the IN_PROGRESS entry with a terminal outcome,
// or this lead terminal could never be rolled again and status() would report
// IN_PROGRESS forever for a roll that in fact never started.
// or this lead could never be rolled again and status() would report IN_PROGRESS
// forever for a roll that in fact never started.
log.warn("lead-rollover: continuationRunner rejected token={} lead={}: {} — the roll "
+ "never started; releasing its claim and reporting it as FAILED",
token, callerTerminal, e.toString(), e);
rollingByTerminal.remove(p.leadTerminal(), token);
rollingByLead.remove(p.rolloverKey(), 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"));
@@ -641,10 +654,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 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());
// 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());
}
}
@@ -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());
}
}
@@ -1476,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"));
@@ -1499,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();
@@ -1517,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
@@ -2647,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
@@ -1814,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) {
@@ -1,96 +1,174 @@
package dev.ltms.fleet;
import com.tngtech.archunit.base.DescribedPredicate;
import com.tngtech.archunit.core.domain.Dependency;
import com.tngtech.archunit.core.domain.JavaClass;
import com.tngtech.archunit.core.domain.JavaClass.Predicates;
import com.tngtech.archunit.core.domain.JavaClasses;
import com.tngtech.archunit.core.importer.ClassFileImporter;
import com.tngtech.archunit.core.importer.ImportOption;
import com.tngtech.archunit.library.dependencies.SliceRule;
import com.tngtech.archunit.library.dependencies.SlicesRuleDefinition;
import org.junit.jupiter.api.Test;
import java.util.ArrayList;
import java.util.List;
import java.util.Set;
import java.util.TreeSet;
import static org.junit.jupiter.api.Assertions.fail;
/**
* fleetd #131 (CB-627): enforce package boundaries with an ArchUnit test instead of a
* Maven module split.
* Enforces package boundaries between the top-level {@code dev.ltms.fleet.*} packages.
*
* <p>This test fails the build the moment a NEW cycle appears between the top-level
* {@code dev.ltms.fleet.*} packages. Today's cycles are recorded below as explicit,
* narrow exceptions: each one ignores dependencies between exactly the two named
* packages, in both directions, and nothing else. A cycle through any other pair of
* packages -- or a brand new pair -- still fails this test.
* <p>{@link #BASELINE_EDGES} names the exact {@code origin class -> target class}
* dependencies allowed to cross a top-level package boundary. Any dependency between two
* top-level packages that is not in that set fails this test, including a brand new
* dependency between a pair of packages that already has other baselined edges. A baseline
* entry whose dependency no longer exists in the code also fails this test, so the baseline
* always names exactly today's exceptions and nothing more.
*
* <p><b>Main code only.</b> The import excludes test classes
* ({@link ImportOption.Predefined#DO_NOT_INCLUDE_TESTS}). Test code legitimately wires
* across many packages for setup and mocking; that is not part of the shipped
* architecture this rule protects. Verified: importing test classes too pulls in a much
* larger, noisier cycle set -- {@code herdr}, {@code member}, {@code peer}, {@code
* config}, {@code guard} and {@code placement} all show up in cycles that disappear the
* moment test classes are excluded. Scanning off the classpath via {@code
* importPackages(...)} (not a hardcoded {@code target/classes} path) also keeps this
* test correct regardless of the working directory the build is invoked from.
*
* <p><b>No package moves here</b> -- ticket #131 is explicit that removing a cycle is
* its own, later PR. See the comment on each exception below for which ticket step
* removes it.
* ({@link ImportOption.Predefined#DO_NOT_INCLUDE_TESTS}). Scanning off the classpath via
* {@code importPackages(...)} keeps this test correct regardless of the working directory
* the build is invoked from.
*/
class PackageCyclesTest {
/**
* Exact {@code "origin -> target"} class dependencies allowed to cross a top-level
* package boundary. Each entry is one directed edge between two specific classes; a
* two-way relationship between a pair of packages is listed as two separate entries,
* one per direction.
*/
private static final Set<String> BASELINE_EDGES = Set.of(
"dev.ltms.fleet.auth.CallerResolver -> dev.ltms.fleet.mcp.ConnectionIdentity",
"dev.ltms.fleet.auth.CallerResolver -> dev.ltms.fleet.mcp.ConnectionIdentity$Caller",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.auth.AuditLog",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.auth.Authz",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.auth.Authz$Action",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.auth.CallerResolver",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.auth.Principal",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.auth.Role",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.LeadChannel",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.LeadChannel$MailboxState",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.LeadMessage",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.MessageService",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.MessageService$AskOutcome",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.MessageService$AskResult",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.MessageService$Outcome",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.MessageService$Outstanding",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.MessageService$PendingAsk",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.MessageService$Phase",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.MessageService$Reply",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.MessageService$ReplyOutcome",
"dev.ltms.fleet.mcp.FleetMcp -> dev.ltms.fleet.msg.MessageService$TaskView",
"dev.ltms.fleet.mcp.FleetMcp$1 -> dev.ltms.fleet.msg.MessageService$AskOutcome",
"dev.ltms.fleet.mcp.FleetMcp$1 -> dev.ltms.fleet.msg.MessageService$Outcome",
"dev.ltms.fleet.mcp.FleetMcp$1 -> dev.ltms.fleet.msg.MessageService$Phase",
"dev.ltms.fleet.mcp.FleetMcp$CoordinationSource -> dev.ltms.fleet.msg.LeadChannel",
"dev.ltms.fleet.msg.LeadHeartbeatLoop -> dev.ltms.fleet.mcp.PrimaryRegistry",
"dev.ltms.fleet.msg.ReplyPushLoop -> dev.ltms.fleet.mcp.PrimaryRegistry",
"dev.ltms.fleet.inject.CompletionResolver -> dev.ltms.fleet.msg.Rendezvous",
"dev.ltms.fleet.inject.CompletionResolver -> dev.ltms.fleet.msg.Rendezvous$Resolution",
"dev.ltms.fleet.inject.CompletionResolver -> dev.ltms.fleet.msg.TurnToken",
"dev.ltms.fleet.inject.CompletionResolver$InFlight -> dev.ltms.fleet.msg.Rendezvous$Resolution",
"dev.ltms.fleet.inject.Injector -> dev.ltms.fleet.msg.TurnToken",
"dev.ltms.fleet.inject.Injector$Pending -> dev.ltms.fleet.msg.TurnToken",
"dev.ltms.fleet.inject.TurnListener -> dev.ltms.fleet.msg.TurnToken",
"dev.ltms.fleet.inject.TurnRegistrar -> dev.ltms.fleet.msg.TurnToken",
"dev.ltms.fleet.msg.MessageService -> dev.ltms.fleet.inject.Injector",
"dev.ltms.fleet.msg.MessageService -> dev.ltms.fleet.inject.Injector$Cancellation",
"dev.ltms.fleet.msg.MessageService -> dev.ltms.fleet.inject.Injector$Delivery",
"dev.ltms.fleet.metrics.FleetMetrics -> dev.ltms.fleet.msg.ReplyInbox",
"dev.ltms.fleet.msg.LeadHeartbeatLoop -> dev.ltms.fleet.metrics.Metrics",
"dev.ltms.fleet.msg.MessageService -> dev.ltms.fleet.metrics.Metrics",
"dev.ltms.fleet.msg.ReplyPushLoop -> dev.ltms.fleet.metrics.Metrics",
"dev.ltms.fleet.msg.LeadHeartbeatLoop -> dev.ltms.fleet.session.MemberSession",
"dev.ltms.fleet.msg.LeadHeartbeatLoop -> dev.ltms.fleet.session.MemberSession$State",
"dev.ltms.fleet.session.SessionManager -> dev.ltms.fleet.msg.TurnToken"
);
private static final String ROOT_PACKAGE = "dev.ltms.fleet.";
@Test
void packagesAreFreeOfCycles() {
var classes = new ClassFileImporter()
JavaClasses classes = new ClassFileImporter()
.withImportOption(ImportOption.Predefined.DO_NOT_INCLUDE_TESTS)
.importPackages("dev.ltms.fleet");
checkBaselineMatchesTodaysEdges(classes);
SliceRule rule = SlicesRuleDefinition.slices()
.matching("dev.ltms.fleet.(*)..")
.should().beFreeOfCycles();
// fleetd #131 step 1: move ConnectionIdentity so authz stops depending on the
// MCP layer. Evidence: auth/CallerResolver.java:3 imports mcp.ConnectionIdentity;
// mcp/FleetMcp.java:3-7 imports auth.AuditLog, Authz, CallerResolver, Principal,
// Role.
rule = ignoreCycle(rule, "auth", "mcp");
// fleetd #131 step 2: PrimaryRegistry is used by loops in msg; move it, or put
// an interface between msg and mcp. Evidence: msg/ReplyPushLoop.java:5 and
// msg/LeadHeartbeatLoop.java:5 import mcp.PrimaryRegistry; mcp/FleetMcp.java:15-18
// imports msg.LeadChannel, LeadMessage, MessageService, Rendezvous.
rule = ignoreCycle(rule, "mcp", "msg");
// fleetd #131 -- found while implementing this test, NOT one of the ticket's
// original three; it names its own follow-up step before removal. Evidence:
// inject/CompletionResolver.java:4-5, inject/Injector.java:6 and
// inject/TurnListener.java:3 import msg.Rendezvous / msg.TurnToken;
// msg/MessageService.java:6 imports inject.Injector.
rule = ignoreCycle(rule, "inject", "msg");
// fleetd #131 -- same as above, its own follow-up. Evidence:
// metrics/FleetMetrics.java:3 imports msg.ReplyInbox; msg/MessageService.java:7-8,
// msg/LeadHeartbeatLoop.java:6-7 and msg/ReplyPushLoop.java:6-7 import
// metrics.FleetMetrics / metrics.Metrics.
rule = ignoreCycle(rule, "metrics", "msg");
// fleetd #131 -- same as above, its own follow-up. Evidence:
// session/SessionManager.java:7 imports msg.TurnToken;
// msg/LeadHeartbeatLoop.java:8 imports session.MemberSession.
rule = ignoreCycle(rule, "msg", "session");
for (String edge : BASELINE_EDGES) {
String[] originAndTarget = edge.split(" -> ");
rule = rule.ignoreDependency(originAndTarget[0], originAndTarget[1]);
}
rule.check(classes);
}
/**
* Accepts today's known cycle between two top-level packages, and nothing else.
* Ignoring both directions removes exactly this pair from cycle detection; every
* other dependency -- including any new one added later, between these same two
* packages or any other pair -- is still checked.
* Fails with the exact offending edge when the live code and {@link #BASELINE_EDGES}
* disagree: a dependency crossing a baselined package pair that is not in the baseline,
* or a baseline entry whose dependency no longer exists.
*/
private static SliceRule ignoreCycle(SliceRule rule, String packageA, String packageB) {
return rule
.ignoreDependency(residesIn(packageA), residesIn(packageB))
.ignoreDependency(residesIn(packageB), residesIn(packageA));
private static void checkBaselineMatchesTodaysEdges(JavaClasses classes) {
Set<String> baselinedPackagePairs = new TreeSet<>();
for (String edge : BASELINE_EDGES) {
String[] originAndTarget = edge.split(" -> ");
baselinedPackagePairs.add(unorderedPair(
topLevelPackageOf(originAndTarget[0]), topLevelPackageOf(originAndTarget[1])));
}
Set<String> liveEdgesInBaselinedPairs = new TreeSet<>();
for (JavaClass javaClass : classes) {
for (Dependency dependency : javaClass.getDirectDependenciesFromSelf()) {
JavaClass origin = dependency.getOriginClass();
JavaClass target = dependency.getTargetClass();
String originPackage = topLevelPackageOf(origin.getFullName());
String targetPackage = topLevelPackageOf(target.getFullName());
if (originPackage.isEmpty() || targetPackage.isEmpty() || originPackage.equals(targetPackage)) {
continue;
}
if (baselinedPackagePairs.contains(unorderedPair(originPackage, targetPackage))) {
liveEdgesInBaselinedPairs.add(origin.getFullName() + " -> " + target.getFullName());
}
}
}
List<String> problems = new ArrayList<>();
for (String liveEdge : liveEdgesInBaselinedPairs) {
if (!BASELINE_EDGES.contains(liveEdge)) {
String[] originAndTarget = liveEdge.split(" -> ");
problems.add("new dependency not in the baseline: " + liveEdge
+ " (packages " + topLevelPackageOf(originAndTarget[0])
+ " -> " + topLevelPackageOf(originAndTarget[1]) + ")");
}
}
for (String baselineEdge : BASELINE_EDGES) {
if (!liveEdgesInBaselinedPairs.contains(baselineEdge)) {
String[] originAndTarget = baselineEdge.split(" -> ");
problems.add("stale baseline entry, no such dependency exists: " + baselineEdge
+ " (packages " + topLevelPackageOf(originAndTarget[0])
+ " -> " + topLevelPackageOf(originAndTarget[1]) + ")");
}
}
if (!problems.isEmpty()) {
fail("PackageCyclesTest baseline is out of date:\n " + String.join("\n ", problems));
}
}
private static DescribedPredicate<JavaClass> residesIn(String topLevelPackage) {
return Predicates.resideInAPackage("dev.ltms.fleet." + topLevelPackage + "..");
private static String unorderedPair(String packageA, String packageB) {
return packageA.compareTo(packageB) <= 0 ? packageA + "|" + packageB : packageB + "|" + packageA;
}
private static String topLevelPackageOf(String fullyQualifiedClassName) {
if (!fullyQualifiedClassName.startsWith(ROOT_PACKAGE)) {
return "";
}
String rest = fullyQualifiedClassName.substring(ROOT_PACKAGE.length());
int dot = rest.indexOf('.');
return dot < 0 ? "" : rest.substring(0, dot);
}
}
@@ -16,6 +16,7 @@ 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;
@@ -1127,8 +1128,13 @@ 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,
_ -> LEAD_NAME, () -> liveLeadTerminals, nowMillis, () -> { }, runner);
t -> LEAD.equals(t) ? LEAD_NAME : null, () -> liveLeadTerminals, nowMillis, () -> { },
runner);
}
@Test
@@ -1317,7 +1323,7 @@ class LeadRolloverTest {
+ "so an operator reading status() has something to act on: " + status.detail());
}
// ---- fleetd #726 unit 3: confirm() single-flights per lead terminal -----------------------
// ---- fleetd #726 unit 3: confirm() single-flights one roll at a time per lead ---------------
@Test
@DisplayName("[SINGLE-FLIGHT 1] a second confirm() for the SAME lead terminal is refused with "
@@ -1343,8 +1349,8 @@ class LeadRolloverTest {
assertFalse(secondDecision.accepted());
assertEquals(LeadRollover.RefusalReason.ROLL_ALREADY_RUNNING, secondDecision.reason());
assertTrue(secondDecision.detail().contains(LEAD), "the refusal detail must name the lead "
+ "terminal: " + secondDecision.detail());
assertTrue(secondDecision.detail().contains(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 "
@@ -1502,4 +1508,160 @@ 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 -----------
}
@@ -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();
}
@@ -2137,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