Compare commits
17 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| be835aa259 | |||
| 80506c79d0 | |||
| 6677ec8c63 | |||
| ca0c965932 | |||
| e8ab933cbc | |||
| 2adb950a12 | |||
| 7f0c4a8464 | |||
| 682991a846 | |||
| abe617c48c | |||
| 886ce1521d | |||
| d41aff4012 | |||
| 70a735b638 | |||
| 803c91ea6c | |||
| d438a74575 | |||
| 7467ffa252 | |||
| f4176ae455 | |||
| 6ab3a81af7 |
@@ -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
|
||||
|
||||
@@ -502,3 +502,33 @@ to replace them.
|
||||
|
||||
Prefer the unnamed lambda parameter `_` for required-but-unused params; a non-public
|
||||
`static void main(String[])` is valid (JEP 512) and boots via `java -jar`.
|
||||
|
||||
## Code quality — five rules, and what each already cost (enforced)
|
||||
|
||||
Measured at `7f0c4a8`: 124 main files, 36,278 lines, of which **17,850 are code** — 44% is comment,
|
||||
and only **two** files exceed 1000 *code* lines. Encapsulation and inheritance are already sound (0
|
||||
public mutable fields; 12 `extends`, 8 of them exceptions; 120 records). So there is **no Clean Code
|
||||
section, no SOLID list and no pattern catalogue** here: two architect reviews rejected those
|
||||
independently as text that would change no behaviour. These five rules are the whole standard.
|
||||
|
||||
1. **A comment states the current contract or a current maintainer constraint — nothing else.** No
|
||||
tickets, history, dates, measurements or review rationale; those go in the commit message or the
|
||||
MR description. Source code only — this rule never applies to Markdown.
|
||||
2. **A javadoc block stops at 30 lines.** Longer means it is a design argument, so it moves to
|
||||
`docs/<subject>.md` and is linked in one line. The longest here is 235 lines
|
||||
(`config/ConfigRef.java`) and the knowledge in it is load-bearing: **move it, never delete it.**
|
||||
This project has **no ADR** — subject pages under `docs/` are the destination.
|
||||
3. **A comment in main source never names a test class.** There are 44 such names in 76 places and
|
||||
**2 are already dead**, because a name inside `{@code}` is invisible to the compiler and rots in
|
||||
silence. Say what the code guarantees; the test is found by looking.
|
||||
4. **Never relieve a testing problem by reshaping production code.** `Fleetd` carries 54 static
|
||||
factories, `MessageService` carries 7 `volatile` race hooks, and 8 tests assert on main source as
|
||||
*text*. Make the part injectable instead. `FleetdAssembly.assembleAndStart` is 452 lines and may
|
||||
not grow; no new source-text test may be added.
|
||||
5. **No new package cycle, and no widening of a recorded one.** Five pairs are frozen as an exact
|
||||
edge baseline in `PackageCyclesTest` — four of them involve `msg`.
|
||||
|
||||
Rules 1, 2, 3 and 5 have build checks, and Gitea CI runs them on every PR, so they bind members too.
|
||||
**Rule 4's judgement half has no mechanism**: a cap stops a count growing, but no test tells a good
|
||||
decomposition from a bad one. That half is a review obligation, and saying so is deliberate — a rule
|
||||
dressed as a gate it does not have is worse than an honest review item.
|
||||
|
||||
@@ -336,22 +336,6 @@ public final class Fleetd {
|
||||
}, reasonByCredential::get);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #415 (review follow-up): package-private factory for the CB-578 stage A {@code
|
||||
* exhaustedPattern} startup coverage line, paired explicitly with {@link
|
||||
* CompletionResolver.UnsetMeaning#OFF} — {@code exhaustedPattern} has no fallback, so a
|
||||
* profile with none configured really does have the classification off.
|
||||
*
|
||||
* <p>Extracted out of {@code main} for the same reason {@link #capacitySource} and {@link
|
||||
* #worktreeBranchLookup} were: {@code coverage()}'s own tests ({@code CompletionResolverTest})
|
||||
* prove it words {@code OFF} and {@link CompletionResolver.UnsetMeaning#BUILT_IN_DEFAULT}
|
||||
* correctly when a test supplies the meaning itself — they cannot prove {@code main} pairs the
|
||||
* right meaning with the right key, which is the actual fleetd #415 defect. <b>Measured:</b>
|
||||
* swapping the {@code UnsetMeaning} arguments between this method and {@link
|
||||
* #errorPatternCoverageLine} — recreating #415's defect with the two keys exchanged — compiled
|
||||
* with 0 errors and left all 1506 existing tests green before {@code
|
||||
* FleetdPatternCoverageLineTest} was added to catch exactly that swap.
|
||||
*/
|
||||
/**
|
||||
* fleetd #446 follow-up: the criterion-2 WARNING text — "name the fix, not just the fact" —
|
||||
* for a profile whose {@code model:} is configured. Extracted out of the {@code
|
||||
@@ -386,19 +370,22 @@ public final class Fleetd {
|
||||
+ "s quarantine above is the only thing keeping new spawns off it for now";
|
||||
}
|
||||
|
||||
/**
|
||||
* Package-private factory for the {@code exhaustedPattern} startup coverage line, paired
|
||||
* explicitly with {@link CompletionResolver.UnsetMeaning#OFF} — {@code exhaustedPattern} has
|
||||
* no fallback, so a profile with none configured really does have the classification off.
|
||||
*/
|
||||
static String exhaustedPatternCoverageLine(Set<String> allProfiles, Set<String> configuredProfiles) {
|
||||
return CompletionResolver.coverage("exhaustedPattern", CompletionResolver.UnsetMeaning.OFF,
|
||||
allProfiles, configuredProfiles);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #415 (review follow-up): the {@code errorPattern} counterpart of {@link
|
||||
* #exhaustedPatternCoverageLine}, paired explicitly with {@link
|
||||
* CompletionResolver.UnsetMeaning#BUILT_IN_DEFAULT} — an unset {@code errorPattern} still runs
|
||||
* backend-error classification against {@code CompletionResolver}'s built-in {@code
|
||||
* BACKEND_ERROR} pattern, so the empty case is not "off". See {@link
|
||||
* #exhaustedPatternCoverageLine}'s javadoc for the measured swap mutation this pairing guards
|
||||
* against.
|
||||
* The {@code errorPattern} counterpart of {@link #exhaustedPatternCoverageLine}, paired
|
||||
* explicitly with {@link CompletionResolver.UnsetMeaning#BUILT_IN_DEFAULT} — an unset
|
||||
* {@code errorPattern} still runs backend-error classification against
|
||||
* {@code CompletionResolver}'s built-in {@code BACKEND_ERROR} pattern, so the empty case is
|
||||
* not "off".
|
||||
*/
|
||||
static String errorPatternCoverageLine(Set<String> allProfiles, Set<String> configuredProfiles) {
|
||||
return CompletionResolver.coverage("errorPattern", CompletionResolver.UnsetMeaning.BUILT_IN_DEFAULT,
|
||||
@@ -734,8 +721,8 @@ public final class Fleetd {
|
||||
* {@code CompletionResolver} constructor call — provably untested wiring, the whole reason
|
||||
* fleetd #248 exists: dropping that one argument (passing {@code _ -> null} instead) compiled
|
||||
* clean and left every test green. Extracted here, {@code main} now calls this factory instead
|
||||
* of building the lambda inline, and a source assertion on that call site
|
||||
* ({@code FleetdCompletionResolverWiringTest}) proves the argument is still actually passed.
|
||||
* of building the lambda inline, so the argument reaching the {@code CompletionResolver}
|
||||
* constructor is a named, directly testable call rather than an inline lambda.
|
||||
*
|
||||
* <p>Takes the roster as a plain {@link Supplier} — not a {@link SessionManager} — so this is
|
||||
* directly testable with a hand-built session list; no real {@code SessionManager} (launcher,
|
||||
@@ -904,6 +891,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) {
|
||||
|
||||
@@ -2095,13 +2095,6 @@ public record FleetConfig(
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* The top-level keys in {@code yaml} that this build does not understand, sorted. Package-private
|
||||
* so the guardrail is asserted directly rather than through a log appender.
|
||||
*
|
||||
* @return empty when everything is known, or when {@code yaml} is not a mapping at all (a
|
||||
* malformed file is {@code readValue}'s error to report, not this method's)
|
||||
*/
|
||||
/**
|
||||
* Top-level keys renamed by the member taxonomy, mapped old → new.
|
||||
*
|
||||
@@ -2636,6 +2629,13 @@ public record FleetConfig(
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* The top-level keys in {@code yaml} that this build does not understand, sorted. Package-private
|
||||
* so the guardrail is asserted directly rather than through a log appender.
|
||||
*
|
||||
* @return empty when everything is known, or when {@code yaml} is not a mapping at all (a
|
||||
* malformed file is {@code readValue}'s error to report, not this method's)
|
||||
*/
|
||||
static List<String> unknownTopLevelKeys(String yaml) {
|
||||
Map<?, ?> raw;
|
||||
try {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -349,11 +349,9 @@ public final class FleetMcp {
|
||||
* overload instead of failing to compile, and every existing test — none of which exercises
|
||||
* {@code Fleetd.main} itself — stays green while the live daemon quietly answers
|
||||
* {@code NOT_CONFIGURED} to {@code fleet_handover} forever. Collapsing every overload into one
|
||||
* required-everything constructor turns that mistake into a compile error instead: this
|
||||
* project's own antidote for a defaulted parameter surviving as an untested decision (see
|
||||
* {@code FleetdCompletionResolverWiringTest} / {@code FleetdLeadRolloverWiringTest}'s own
|
||||
* javadoc for the same lesson applied to a different seam). A caller that genuinely wants a
|
||||
* feature off must now say so explicitly at the call site — {@code null},
|
||||
* required-everything constructor turns that mistake into a compile error instead, the same
|
||||
* defaulted-parameter guard this codebase applies to other required wiring. A caller that
|
||||
* genuinely wants a feature off must now say so explicitly at the call site — {@code null},
|
||||
* {@link OutageSource#none()}, {@link LeadSeatSource#none()}, {@code List.of()} are all still
|
||||
* perfectly fine values, just never an implicit default reached by omission.
|
||||
*
|
||||
@@ -489,7 +487,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 +598,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 +735,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 +1357,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 +1481,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 +1505,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 +1524,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 +2657,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;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -445,27 +445,6 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
+ " distinct candidate(s): " + String.join(", ", unreachable));
|
||||
}
|
||||
|
||||
/**
|
||||
* Refuse an explicit-profile spawn when the profile is at its {@code maxLoad} cap.
|
||||
*
|
||||
* <p>maxLoad is a documented, unconditional capacity limit (see {@code FleetConfig.Profile#maxLoad}),
|
||||
* and the charter makes explicit-profile spawns the normal path — so enforcing it only in placement
|
||||
* ({@code PlacementPolicyUtil}, package-private, hence not linked) would leave the cap dead config
|
||||
* on every call that names a profile. Same rule as placement: {@code live >= cap} is at capacity.
|
||||
*
|
||||
* <p>Deliberately no fallback to another profile: the caller named {@code profile} for a cost/model
|
||||
* reason, and silently re-routing a paid-tier (subscription) request elsewhere is worse than
|
||||
* refusing it. A caller that wants placement should omit the profile and let the policy pick.
|
||||
*
|
||||
* <p>Known TOCTOU limitation — documented, not fixed. {@link #liveCount} is read outside any lock and
|
||||
* {@code SessionManager} registers a session only after {@code launcher.spawn} returns, so two
|
||||
* genuinely concurrent spawns can both pass this check. The race already exists on the placement
|
||||
* path. Closing it needs slot reservation in the registry; serializing spawn here would block on
|
||||
* the readiness gate and is a far worse trade.
|
||||
*
|
||||
* @param profile the profile the caller explicitly named
|
||||
* @throws PlacementException when the profile is at capacity
|
||||
*/
|
||||
/**
|
||||
* Refuse an explicit-profile spawn whose credential is quarantined (CB-578 stage B): a prior
|
||||
* {@code BACKEND_EXHAUSTED} classification on this profile, or on another profile sharing its
|
||||
@@ -527,6 +506,27 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
.collect(Collectors.toSet());
|
||||
}
|
||||
|
||||
/**
|
||||
* Refuse an explicit-profile spawn when the profile is at its {@code maxLoad} cap.
|
||||
*
|
||||
* <p>maxLoad is a documented, unconditional capacity limit (see {@code FleetConfig.Profile#maxLoad}),
|
||||
* and the charter makes explicit-profile spawns the normal path — so enforcing it only in placement
|
||||
* ({@code PlacementPolicyUtil}, package-private, hence not linked) would leave the cap dead config
|
||||
* on every call that names a profile. Same rule as placement: {@code live >= cap} is at capacity.
|
||||
*
|
||||
* <p>No fallback to another profile: the caller named {@code profile} for a cost/model
|
||||
* reason, and silently re-routing a paid-tier (subscription) request elsewhere is worse than
|
||||
* refusing it. A caller that wants placement should omit the profile and let the policy pick.
|
||||
*
|
||||
* <p>Known TOCTOU limitation — documented, not fixed. {@link #liveCount} is read outside any lock and
|
||||
* {@code SessionManager} registers a session only after {@code launcher.spawn} returns, so two
|
||||
* genuinely concurrent spawns can both pass this check. The race already exists on the placement
|
||||
* path. Closing it needs slot reservation in the registry; serializing spawn here would block on
|
||||
* the readiness gate and is a far worse trade.
|
||||
*
|
||||
* @param profile the profile the caller explicitly named
|
||||
* @throws PlacementException when the profile is at capacity
|
||||
*/
|
||||
private void enforceMaxLoad(String profile) {
|
||||
// Absent config, or a config whose maxLoad normalized to null (ABSENT ⇒ unlimited at load),
|
||||
// means no cap — never cap what wasn't configured. Note "non-positive ⇒ unlimited" was true
|
||||
|
||||
@@ -474,7 +474,6 @@ public final class EnvAllowListScrub {
|
||||
}
|
||||
}
|
||||
|
||||
/** Best-effort recursive delete; failures are swallowed — JVM-exit cleanup is the backstop. */
|
||||
/**
|
||||
* Remove generated directories left behind by an earlier daemon process.
|
||||
*
|
||||
@@ -511,6 +510,7 @@ public final class EnvAllowListScrub {
|
||||
}
|
||||
}
|
||||
|
||||
/** Best-effort recursive delete; failures are swallowed — JVM-exit cleanup is the backstop. */
|
||||
static void deleteRecursively(Path dir) {
|
||||
if (dir == null || !Files.exists(dir)) {
|
||||
return;
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -442,38 +442,6 @@ public final class GitWorktrees implements Worktrees {
|
||||
ENVIRONMENT_CREDENTIAL_HELPER);
|
||||
}
|
||||
|
||||
/**
|
||||
* {@link #configureEnvironmentCredentialHelper} only ever fires for an HTTPS origin — Git never
|
||||
* consults a {@code credential.helper} for an SSH transport. This repo's own origin is
|
||||
* {@code ssh://git@git.ltms.dev:2224/fleet/fleetd.git}, so a member sitting on that origin never
|
||||
* reaches the helper and the repo-scoped {@code WORKER_GITEA_TOKEN} is simply not used.
|
||||
*
|
||||
* <p>An earlier version of this javadoc justified the rewrite by claiming a member <em>cannot</em>
|
||||
* push once {@code memberCredentials.policy: allow-list} blocks {@code SSH_AUTH_SOCK}, because
|
||||
* "there is no private key file on this host, only an ssh-agent socket". That premise is false
|
||||
* (fleetd #184): {@code ssh -G} resolves a readable, passphrase-free {@code IdentityFile} outside
|
||||
* {@code ~/.ssh}, and a member — same OS user — pushes over SSH with the socket blanked. The
|
||||
* rewrite is still worth having, but for the reason below rather than that one: it routes the
|
||||
* member through its own scoped token instead of the operator's ssh identity, which is what makes
|
||||
* a member's pushes attributable and revocable.
|
||||
*
|
||||
* <p>The fix is a <em>worktree-scoped</em> URL rewrite: {@code url.<https-base>.insteadOf
|
||||
* <ssh-base>}, set with {@code --worktree} so it lands only in
|
||||
* {@code <worktree>/.git/worktrees/<name>/config.worktree} (enabled by
|
||||
* {@code extensions.worktreeConfig}, already turned on above) and never touches the shared
|
||||
* repo-level config the primary checkout also reads. {@code insteadOf} — not
|
||||
* {@code pushInsteadOf} — because a member may also need to fetch or rebase, and both should go
|
||||
* through the member's own token for the same reason.
|
||||
*
|
||||
* <p>The host (and, for the rewrite's SSH-side match, the port) come from parsing the origin
|
||||
* itself — never a hardcoded forge host, which is exactly what #177 removed. An origin that is
|
||||
* already {@code https://} is left alone; the credential helper already covers it. An origin
|
||||
* that is neither {@code ssh://} nor {@code https://} — including the scp-like shorthand
|
||||
* ({@code git@host:path}, no scheme) — is left untouched deliberately: that shorthand's
|
||||
* {@code host:path} split is defined by the user's ssh_config aliases, not by URI syntax, so
|
||||
* guessing at it risks rewriting to the wrong place. A repo provisioned from that form keeps
|
||||
* today's (broken, if the policy blocks the agent) SSH-only behaviour rather than a wrong rewrite.
|
||||
*/
|
||||
/**
|
||||
* Blank the user-info of a remote URL before it reaches a log. A remote URL is not obviously a
|
||||
* credential channel, which is exactly why one has leaked here three times ({@code git remote -v}
|
||||
@@ -485,6 +453,30 @@ public final class GitWorktrees implements Worktrees {
|
||||
return url == null ? null : url.replaceAll("://[^@/]*@", "://<redacted>@");
|
||||
}
|
||||
|
||||
/**
|
||||
* {@link #configureEnvironmentCredentialHelper} only fires for an HTTPS origin — Git never
|
||||
* consults a {@code credential.helper} for an SSH transport. This repo's own origin is
|
||||
* {@code ssh://git@git.ltms.dev:2224/fleet/fleetd.git}, so a member on that origin never
|
||||
* reaches the helper, and the repo-scoped token goes unused without a separate rewrite.
|
||||
*
|
||||
* <p>This method routes the member through its own scoped token instead of the operator's ssh
|
||||
* identity, which is what makes a member's pushes attributable and revocable.
|
||||
*
|
||||
* <p>The fix is a <em>worktree-scoped</em> URL rewrite: {@code url.<https-base>.insteadOf
|
||||
* <ssh-base>}, set with {@code --worktree} so it lands only in
|
||||
* {@code <worktree>/.git/worktrees/<name>/config.worktree} and never touches the shared
|
||||
* repo-level config the primary checkout also reads. {@code insteadOf} — not
|
||||
* {@code pushInsteadOf} — because a member may also need to fetch or rebase through its own
|
||||
* token.
|
||||
*
|
||||
* <p>The host (and, for the rewrite's SSH-side match, the port) come from parsing the origin
|
||||
* itself, never a hardcoded forge host. An origin already {@code https://} is left alone; the
|
||||
* credential helper already covers it. An origin that is neither {@code ssh://} nor
|
||||
* {@code https://} — including the scp-like shorthand ({@code git@host:path}, no scheme) — is
|
||||
* left untouched: that shorthand's {@code host:path} split is defined by the user's ssh_config
|
||||
* aliases, not by URI syntax, so guessing at it risks rewriting to the wrong place, and that
|
||||
* origin keeps SSH-only push behaviour instead.
|
||||
*/
|
||||
private void configureHttpsUrlRewriteForSshOrigin(String repoRoot, String worktreePath) {
|
||||
if (exitCode("git", "-C", repoRoot, "config", "--get", "remote.origin.url") != 0) {
|
||||
return;
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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());
|
||||
|
||||
Reference in New Issue
Block a user