Compare commits

..

1 Commits

Author SHA1 Message Date
Dai Ha ea41bbf6b9 fleetd#329: fix silent async-ticket bugs in MessageService (F1/F2/F3)
CI / contract (pull_request) Successful in 50s
CI / build (pull_request) Failing after 1m28s
F2 (sendAsync executor catch): log when completeExceptionally returns
false, so an exception thrown after finishAsyncTask already completed
the ticket's future is no longer silently lost.

F1 (answer()'s stranded async ticket): reuse the Task reference answer()
already looked up before rendezvous.answerAsk(), instead of a second
asyncTasksByTurn lookup by turnId in finishAsyncTask. The second lookup
raced ask()'s unlocked timeout cleanup, which could forget turnId first
and leave the ticket stuck PENDING even though answer() itself returned
REPLIED. The #282 chained-ask guard is unaffected: it is still keyed on
result.outcome() == QUESTION, not on this lookup. Removed the now-unused
finishAsyncTask(String, Reply) overload.

F3 (reply()'s orphan recovery path): read orphan.turnId once instead of
twice, closing the same double-read shape fleetd #324 fixed in
finishAsyncTask.

Each fix has its own test plus a test-only race hook (mirroring #324's
finishAsyncTaskRaceHook) to force the exact interleaving deterministically.
Mutation-tested each fix by reverting it, confirming the real failure
(swallowed exception / PENDING ticket / NullPointerException), then
restoring it.

mvn clean install: Tests run: 1345, Failures: 0, Errors: 0, Skipped: 0,
BUILD SUCCESS.
2026-09-04 15:19:38 +07:00
6 changed files with 306 additions and 841 deletions
@@ -22,19 +22,19 @@ import java.util.function.Supplier;
* choice rather than an accident of where the field was initialised.
*
* <h2>Not every key can change under a running daemon</h2>
* Keys fall into four classes, and the difference is about what already exists when the reload
* Keys fall into three classes, and the difference is about what already exists when the reload
* happens — not about how important the key is.
*
* <ul>
* <li><strong>Hot</strong> — re-read per use, so a reload takes effect on the next spawn:
* {@code placement:}, and an existing profile's {@code weight} / {@code maxLoad}. Both are
* read through a supplier on {@code CompositePeerLauncher}, which is what makes them hot —
* not the fact that they are config. Most of {@code fleet:} — every role pool
* ({@code architects}/{@code developers}/{@code reviewers}), {@code charters}, and
* {@code tabLabel} — is read the same live way, through the same supplier
* ({@code () -> config.get().fleet()}). <strong>But {@code fleet:} as a whole is NOT in this
* class</strong>: {@code fleet.leaders} inside the same key is frozen, which is exactly what
* makes {@code fleet:} split rather than hot — see below.</li>
* {@code fleet:} (every role pool, {@code charters}, and {@code tabLabel}),
* {@code placement:}, and an existing profile's {@code weight} / {@code maxLoad}. Those
* three are read through a supplier on {@code CompositePeerLauncher}, which is what makes
* them hot — not the fact that they are config. <strong>This does NOT include
* {@code fleet.leaders}</strong>: {@code Fleetd.main} reads {@code cfg.fleet().leaders()}
* once at startup to build the {@code LeadTabScanner} and the {@code LeadLauncher}, and
* neither is reconstructed on reload — so a lead added, removed, or re-{@code tab}'d under
* {@code fleet.leaders} needs a restart, the same as any deferred key below.</li>
* <li><strong>Deferred</strong> — accepted into the new snapshot, but the wiring built at startup
* keeps the old value until a restart: {@code lifecycle:}, {@code leadHeartbeat:},
* {@code spawnReadyTimeoutMs} / {@code spawnReadyPollMs}, {@code quarantineCooldownSeconds}
@@ -71,99 +71,35 @@ import java.util.function.Supplier;
* {@code HerdrPeerLauncher} takes {@code Map.copyOf(profiles)} at construction and resolves
* each spawn out of that copy, so those never reach a launch until the daemon restarts. A
* reload logs these rather than pretending they applied.</li>
* <li><strong>Split</strong> (fleetd #330; extended to a third key by fleetd #333) — read
* <em>both</em> ways at different sites, so the key does not fit any class above as a whole:
* {@code health:}, {@code coordinator:} and {@code fleet:}. Each is read off the startup
* snapshot to build a long-lived object, and read live off {@link #get()} at a different,
* unrelated site — so half of a reload's effect already applies while the other half waits
* for a restart, and a bare "config reloaded" would under-claim by exactly that half.
* <ul>
* <li>{@code health:} — the monitor itself ({@code enabled}, {@code intervalSeconds},
* {@code workingSuspectAfterSeconds}) is built once at {@code Fleetd.java:556-563}
* and never rebuilt, so a changed value needs a restart to actually start, stop, or
* retime it. The coverage string {@code fleet_profiles} reports
* ({@code Fleetd.java:648-650}) is read live off {@link #get()} on every call, so it
* already reflects the new value.</li>
* <li>{@code coordinator:} — the {@code LeadMailbox} connection ({@code uri},
* {@code uriEnv}, {@code selfId}, {@code prefetch}) is opened once at
* {@code Fleetd.java:502} and never reopened, so a changed value needs a restart —
* {@code selfId} in particular names this daemon's own AMQP inbox queue, and a peer
* lead that learned the old name would not discover a new one on its own. The broker
* URI env-var <em>name</em> that {@code MemberEnvAllowList} keeps out of a member's
* environment is read live off {@link #get()} on every spawn
* ({@code HerdrPeerLauncher.java:1530}), so it already applies.</li>
* <li>{@code fleet:} (fleetd #333) — {@code fleet.leaders} is the frozen half:
* {@code Fleetd.java:281} reads {@code cfg.fleet().leaders()} off the startup snapshot
* for two long-lived objects built right after it and never rebuilt — the
* {@code LeadTabScanner}'s {@code tab label → lead name} map ({@code Fleetd.java:301},
* wired into {@code CallerResolver.withLeadsAndMembers} at {@code Fleetd.java:620/624},
* which is how a caller's pane is recognised as a lead at all) and, when herdr answered
* at startup, {@code LeadLauncher(...).ensureLeads()} ({@code Fleetd.java:315}), which
* auto-launches each declared lead up to its {@code instances} count. So a lead added,
* removed, or given a new {@code tab:} label under {@code fleet.leaders} needs a
* restart — until then it is invisible to identity resolution, and this is exactly the
* scenario fleetd #333 named: an operator edits a lead's {@code tab:} to match a
* renamed pane, sees "config reloaded", and the pane keeps resolving as a worker,
* because {@code CallerResolver} is still matching against the old label. The live
* half is the rest of {@code fleet:} — {@code architects}/{@code developers}/
* {@code reviewers}, {@code charters}, {@code tabLabel} — read live through the same
* {@code CompositePeerLauncher} supplier the Hot bullet above names, so a reload that
* only touches those already applies with nothing to report. Because the hot and frozen
* halves of {@code fleet:} are disjoint sub-fields rather than the same fields read two
* ways (contrast {@code coordinator.uriEnv} above), {@link #changedSplitKeys} compares
* {@code fleet.leaders} alone, not the whole {@code Fleet} record — comparing the whole
* record would report "split" for a {@code tabLabel}-only change that is actually fully
* hot, over-claiming in exactly the direction this class exists to avoid under-claiming
* in.</li>
* </ul>
* A split change is still accepted — {@link Outcome#applied()} stays {@code true}, the same
* as a deferred change — because the live half genuinely took effect; refusing the whole
* reload would leave the operator worse off than today. {@link Outcome#split()} names the
* key and says which half is which each time, rather than trying to score "how changed" a
* mixed key is or handle "both halves changed in one reload" as a special case.</li>
* <li><strong>Cold</strong> — cannot change at all under a running daemon: {@code bind:},
* {@code herdrSocket:}, {@code broker:} and {@code auth:}. The socket is bound, the broker
* connection is open, and the auth mode decides who may reach the port that is already
* listening.</li>
* </ul>
*
* <p><strong>The denominator, measured on 2026-09-04 (fleetd #330; recounted for fleetd #333).</strong>
* {@code FleetConfig} has 22 top-level record components: 5 cold, 11 deferred, 3 split, 3
* hot-excluded. Three of them are named nowhere in this file, and the reason is the same for all
* three: {@code placement}, {@code memberCredentials} and {@code memberLoginShell} are
* <strong>hot</strong> and correctly absent — all three are read live off {@code config.get()}
* (placement through the {@code CompositePeerLauncher} supplier the Hot bullet names;
* {@code memberCredentials}/{@code memberLoginShell} at spawn time, {@code Fleetd.java:198, 205, 729}
* and {@code HerdrPeerLauncher#configuredMemberLoginShell}), so a reload takes effect on the next
* spawn with no entry needed here.
* {@code health} and {@code coordinator} used to be a third kind — <strong>undecided</strong>, not
* hot — until fleetd #330 added the <strong>split</strong> class above and gave them a home. A
* reload touching either used to report a bare "config reloaded", which under-claimed; now it names
* the key and says which half is which. {@code fleet} was the same story in reverse: fleetd #330's
* own fact-find named {@code health}/{@code coordinator} as "the complete set of split-shaped keys"
* and filed {@code fleet.leaders}'s restart requirement as a documented caveat sitting in the
* <strong>hot-excluded</strong> escape hatch instead — correctly documented, but in the one bucket
* this file's own coverage test cannot check the truth of (see that test's javadoc). fleetd #333
* moved it into <strong>split</strong>, where {@link #changedSplitKeys} actually reports it.
* <p>The point of writing the count down: "not mentioned in this file" looks identical for a key
* that is correctly hot and for a key nobody triaged. Three times now — {@code worktreeGroup} (#323),
* {@code primary}/{@code configReload} (#326), and {@code fleet.leaders} sitting in the escape hatch
* (#333) — the second kind hid among the first. A top-level coverage checker in the
* {@link ConfigRefProfileCoverageTest} shape (one level up, over {@code FleetConfig} itself rather
* than {@code FleetConfig.Profile}) proves this file's four classes exhaust the record's components
* — see {@code ConfigRefTopLevelCoverageTest}. That test proves the record's <em>shape</em> is fully
* triaged; it does NOT prove a {@code SPLIT_KEYS}/{@code COLD_KEYS} member has any reporting code
* behind it at all — {@code ConfigRefTopLevelReportingCoverageTest} is what fleetd #333 added for
* that, after measuring that a {@code SPLIT_KEYS} entry with its reporting branch deleted passes
* both this file's own "kept in step" assert and {@code ConfigRefTopLevelCoverageTest} unchanged.
* <p><strong>The denominator, measured on 2026-09-04 (fleetd #326).</strong> {@code FleetConfig} has
* 22 top-level record components. Four of them are named nowhere in this file, and the reason
* differs per key, so do not read "absent" as "forgotten":
* <ul>
* <li>{@code memberCredentials} and {@code memberLoginShell} are <strong>hot</strong> and
* correctly absent — both are read live off {@code config.get()} at spawn time
* ({@code Fleetd.java:198, 205, 729} and {@code HerdrPeerLauncher#configuredMemberLoginShell}),
* so a reload takes effect on the next spawn with no entry needed here.</li>
* <li>{@code health} and {@code coordinator} are <strong>undecided</strong>, not hot. Each is read
* both ways at different sites, so neither fits the three classes above as a whole key. Until
* that is settled a reload touching them reports a bare "config reloaded", which under-claims.
* Deciding it is what a top-level coverage checker (the {@link ConfigRefProfileCoverageTest}
* shape, one level up) is waiting on — without it, such a checker cannot express the answer.</li>
* </ul>
* The point of writing the count down: "not mentioned in this file" looks identical for a key that
* is correctly hot and for a key nobody triaged. Twice now — {@code worktreeGroup} (#323) and
* {@code primary}/{@code configReload} (#326) — the second kind hid among the first.
*
* <p><strong>A cold change refuses the whole reload.</strong> Not the hot half applied and the cold
* half warned about: that would leave the running daemon in a state matching no file on disk, which
* is the worst thing a reload can do to an operator debugging one. Refusing keeps the invariant that
* the live config is always some version of the file, and the message names the keys that must
* change through a restart. A split change does <em>not</em> refuse, for a different reason than a
* deferred change does not: its live half genuinely took effect, so refusing would throw that away
* and leave the operator worse off than the partial-but-honest report {@link Outcome#split()} gives.
* change through a restart.
*
* <p>A reload that fails to parse or fails validation is also refused, and the previous config keeps
* running. A config file being edited is normally read once mid-save; degrading a working daemon
@@ -173,24 +109,10 @@ public final class ConfigRef implements Supplier<FleetConfig> {
private static final Logger log = LoggerFactory.getLogger(ConfigRef.class);
/**
* Keys that cannot change under a running daemon — see the class doc.
*
* <p>Package-private (not {@code private}) so {@code ConfigRefTopLevelCoverageTest} can fold it
* into the top-level triage it checks, the same way it reads {@link #SPLIT_KEYS}.
*/
static final Set<String> COLD_KEYS =
/** Keys that cannot change under a running daemon — see the class doc. */
private static final Set<String> COLD_KEYS =
Set.of("bind", "herdrSocket", "memberHerdrSocket", "broker", "auth");
/**
* Keys read BOTH off the startup snapshot and live off {@link #get()} at different sites, so
* neither the hot, deferred nor cold class fits them as a whole — see the class doc's Split
* bullet (fleetd #330). A changed split key is accepted ({@link Outcome#applied()} stays
* {@code true}) and reported by name, with a message naming which half is live and which needs
* a restart.
*/
static final Set<String> SPLIT_KEYS = Set.of("health", "coordinator", "fleet");
private final Path path;
private final AtomicReference<FleetConfig> current;
@@ -218,37 +140,25 @@ public final class ConfigRef implements Supplier<FleetConfig> {
/**
* What a reload attempt did.
*
* <p>{@code split} is a separate field from {@code deferred} rather than a differently-worded
* entry inside it, because the two carry different guarantees for any caller that branches on
* them rather than just printing {@link #summary()}: every {@code deferred} entry means "this
* key's whole change waits for a restart", while every {@code split} entry means "part of this
* key's change already applied, and the message says which part" — collapsing them would force
* a caller to re-parse the message to tell those apart. See the class doc's Split bullet
* (fleetd #330) for why the key needs this at all.
*
* @param applied true when the new config is now live
* @param coldKeys cold keys whose value changed, which is why an unapplied reload was refused
* @param deferred keys that changed and were accepted, but whose effect waits entirely on a
* restart
* @param split split keys that changed and were accepted, each named with which half of it
* is already live and which half waits for a restart
* @param deferred keys that changed and were accepted, but whose effect waits for a restart
* @param error the parse or validation failure that refused the reload, else {@code null}
*/
public record Outcome(boolean applied, List<String> coldKeys, List<String> deferred,
List<String> split, String error) {
String error) {
public Outcome {
coldKeys = List.copyOf(coldKeys);
deferred = List.copyOf(deferred);
split = List.copyOf(split);
}
static Outcome refusedCold(List<String> keys) {
return new Outcome(false, keys, List.of(), List.of(), null);
return new Outcome(false, keys, List.of(), null);
}
static Outcome failed(String error) {
return new Outcome(false, List.of(), List.of(), List.of(), error);
return new Outcome(false, List.of(), List.of(), error);
}
/** A one-line summary for the operator — the reason, not just the verdict. */
@@ -260,18 +170,11 @@ public final class ConfigRef implements Supplier<FleetConfig> {
return "config reload refused — these keys cannot change under a running daemon: "
+ String.join(", ", coldKeys) + ". Restart fleetd to apply them.";
}
if (deferred.isEmpty() && split.isEmpty()) {
return "config reloaded";
}
StringBuilder out = new StringBuilder("config reloaded");
if (!deferred.isEmpty()) {
out.append("; these changes need a restart to take effect: ")
.append(String.join(", ", deferred));
return "config reloaded; these changes need a restart to take effect: "
+ String.join(", ", deferred);
}
if (!split.isEmpty()) {
out.append("; partially live — ").append(String.join(" | ", split));
}
return out.toString();
return "config reloaded";
}
}
@@ -312,21 +215,14 @@ public final class ConfigRef implements Supplier<FleetConfig> {
}
List<String> deferred = changedDeferredKeys(old, fresh);
List<String> split = changedSplitKeys(old, fresh);
current.set(fresh);
Outcome out = new Outcome(true, List.of(), deferred, split, null);
Outcome out = new Outcome(true, List.of(), deferred, null);
log.info(out.summary());
return out;
}
/**
* Cold keys whose value differs between the running config and the candidate.
*
* <p>Package-private (not {@code private}) so {@code ConfigRefTopLevelReportingCoverageTest}
* can call it directly with a reflection-built {@code FleetConfig} pair, the same reason
* {@link #sameLaunchSettings} is package-private — see that test's class doc (fleetd #333).
*/
static List<String> changedColdKeys(FleetConfig old, FleetConfig fresh) {
/** Cold keys whose value differs between the running config and the candidate. */
private static List<String> changedColdKeys(FleetConfig old, FleetConfig fresh) {
List<String> changed = new ArrayList<>();
if (!Objects.equals(old.bind(), fresh.bind())) {
changed.add("bind");
@@ -429,73 +325,6 @@ public final class ConfigRef implements Supplier<FleetConfig> {
return changed;
}
/**
* Split keys whose value differs between the running config and the candidate — see the class
* doc's Split bullet (fleetd #330, extended for {@code fleet:} by fleetd #333). Unlike
* {@link #changedDeferredKeys}, this does not try to tell which sub-field moved for {@code
* health:} or {@code coordinator:}: any change to either gets the same fixed message, because
* the message already names both halves every time, so there is no "which half changed"
* question left for the caller to answer. {@code fleet:} is different on purpose — see below.
*
* <p>Package-private (not {@code private}) so {@code ConfigRefTopLevelReportingCoverageTest}
* can call it directly with a reflection-built {@code FleetConfig} pair, the same reason
* {@link #sameLaunchSettings} is package-private — see that test's class doc (fleetd #333). That
* test exists because membership in {@link #SPLIT_KEYS} proves nothing about this method on its
* own: fleetd #333 measured that dropping the {@code coordinator} branch out of this method
* while leaving {@code "coordinator"} in {@code SPLIT_KEYS} left the whole suite green except a
* hand-written {@code ConfigRefTest} case — neither {@code ConfigRefTopLevelCoverageTest} (it
* only reads the set) nor the "kept in step" assert below (it only checks the reported keys are
* a SUBSET of {@code SPLIT_KEYS}, never that every {@code SPLIT_KEYS} member has a branch here)
* would have caught it.
*/
static List<String> changedSplitKeys(FleetConfig old, FleetConfig fresh) {
List<String> changed = new ArrayList<>();
if (!Objects.equals(old.health(), fresh.health())) {
changed.add("health: the monitor itself (enabled, interval, workingSuspectAfter) is "
+ "frozen at startup and needs a restart; the coverage status fleet_profiles "
+ "reports is read live and already applied");
}
if (!Objects.equals(old.coordinator(), fresh.coordinator())) {
changed.add("coordinator: the LeadMailbox connection (uri, uriEnv, selfId, prefetch) is "
+ "opened once and needs a restart; the broker URI env-var name kept out of a "
+ "member's environment is read live on every spawn and already applied");
}
// fleetd #333: unlike health/coordinator above, most of `fleet:` (architects, developers,
// reviewers, charters, tabLabel) is genuinely hot — ConfigRefTest.aHotChangeIsAppliedAndRead-
// ThroughGet and aCharterChangeIsHotAndReachesTheLiveConfig prove it reaches the live config
// with no restart note. Only fleet.leaders is frozen (Fleetd.java:281 reads
// cfg.fleet().leaders() off the startup snapshot to build both the LeadTabScanner's
// tab-label-to-name map, wired into CallerResolver.withLeadsAndMembers at Fleetd.java:620/624,
// and — when herdr answered — LeadLauncher(...).ensureLeads() at Fleetd.java:315, which
// auto-launches each lead up to its `instances` count; neither is rebuilt on reload). So this
// compares fleet.leaders alone, not the whole Fleet record: comparing the whole record would
// report "split" for a tabLabel-only or charters-only change that is actually fully hot,
// which is the over-claim mirror of the under-claim bug this class exists to prevent.
if (!Objects.equals(leadersOf(old), leadersOf(fresh))) {
changed.add("fleet: fleet.leaders (each lead's tab, workspace, cwd, profile and "
+ "instances count) is read once at startup to build the LeadTabScanner's "
+ "identity map and to auto-launch leads, and neither is rebuilt on reload, so a "
+ "lead added, removed, or given a new tab: label needs a restart — until then it "
+ "stays unrecognised, and a caller from its new tab resolves as a worker, not a "
+ "lead; the rest of fleet: (architects, developers, reviewers, charters, "
+ "tabLabel) is read live through the supplier on CompositePeerLauncher and "
+ "already applied");
}
// Kept in step with SPLIT_KEYS the same way changedColdKeys is kept in step with COLD_KEYS —
// every message here must be traceable to one of the split keys the class doc documents.
// NOTE what this does NOT prove, per the javadoc above: it does not catch a SPLIT_KEYS
// member with no branch above at all, only a branch whose message is mis-worded relative to
// the set. ConfigRefTopLevelReportingCoverageTest is what proves the former.
assert changed.stream().allMatch(m -> SPLIT_KEYS.stream().anyMatch(k -> m.startsWith(k + ":")))
: "a split entry was reported that does not start with a SPLIT_KEYS name: " + changed;
return changed;
}
/** {@code cfg.fleet().leaders()}, defensively, in case a caller hands in a non-defaulted config. */
private static Map<String, FleetConfig.Leader> leadersOf(FleetConfig cfg) {
return cfg.fleet() == null ? Map.of() : cfg.fleet().leaders();
}
/**
* {@link FleetConfig.Profile} record components deliberately left out of
* {@link #sameLaunchSettings} because they are read <em>live</em>, not baked in at spawn — see
@@ -466,8 +466,18 @@ public final class MessageService {
if (candidates.size() == 1) {
Task orphan = candidates.get(0);
if (orphan.future.complete(new Reply(Outcome.REPLIED, content))) {
if (orphan.turnId != null) {
asyncTasksByTurn.remove(orphan.turnId, orphan);
// fleetd #329 (F3): read orphan.turnId once. It used to be read twice, under no
// lock — once for this null check, once as the removal key — the same double-read
// shape fleetd #324 fixed in finishAsyncTask. clearAsyncQuestion's unlocked
// forgetTurn=true path (ask()'s timeout cleanup) can null the field between the two
// reads; capturing it once removes the torn read here too.
String turnId = orphan.turnId;
if (turnId != null) {
if (replyOrphanTurnIdRaceHookForTest != null) {
// Test-only (fleetd #329, F3): see the field's own javadoc.
replyOrphanTurnIdRaceHookForTest.run();
}
asyncTasksByTurn.remove(turnId, orphan);
}
count(FleetMetrics.REPLIES, "path", "async-recovered");
return true; // the ticket itself took it — no inbox stranding at all
@@ -1002,13 +1012,36 @@ public final class MessageService {
// Measured when #282 was merged: this guard is DEFENCE IN DEPTH, not the thing
// that makes the chained ask work. ask() calls markAsyncQuestion (:860) before
// resolveQuestion (:861), so by the time this thread wakes, the task has already
// moved to the new turnId and finishAsyncTask(oldTurnId, ...) finds nothing. Removing
// this guard alone leaves the test green. Keep it anyway: it mirrors sendAsync's
// sibling guard, and that sibling's own comment (:1017) warns the two orderings are
// not something to rely on. Do NOT delete it as dead code without re-checking that
// ordering, and do not treat it as the sole protection either.
if (result.outcome() != Outcome.QUESTION) {
finishAsyncTask(turnId, result);
// moved to the new turnId — this outcome check, not a Task lookup, is what tells the
// two cases apart (see fleetd #329 below). Removing this guard alone leaves the test
// green. Keep it anyway: it mirrors sendAsync's sibling guard, and that sibling's own
// comment (:1017) warns the two orderings are not something to rely on. Do NOT delete
// it as dead code without re-checking that ordering, and do not treat it as the sole
// protection either.
//
// fleetd #329 (F1): complete the SAME Task object this method already looked up at
// :991, rather than re-resolving it from turnId a second time. The old
// finishAsyncTask(turnId, result) did its own asyncTasksByTurn.get(turnId) here, and
// that second lookup races ask()'s own timeout path: ask()'s ticket.answer().get(...)
// can time out (or lose that exact race) at essentially the same instant this method's
// rendezvous.answerAsk(turnId, ...) above already succeeded, running
// markAskTimedOut + clearAsyncQuestion(turnId, true) with no lock at all — which
// forgets turnId (removes it from asyncTasksByTurn, nulls Task.turnId) before this
// thread ever gets here. The worker's real reply then arrived, this method's own wait
// woke up with it, and the by-turnId lookup found nothing: the async ticket's future
// was never completed, so fleet_poll{ticket} stayed PENDING forever even though
// answer() itself correctly returned REPLIED. Reusing the reference captured at :991
// — before that race window opens — sidesteps the second lookup entirely: it is the
// identical Task whatever asyncTasksByTurn or Task.turnId say by the time we reach
// this point, so finishAsyncTask(Task, Reply) can still complete its future and detach
// it using whatever turnId it now reads. The #282 chained-ask case is unaffected
// because it is still gated purely by result.outcome() == QUESTION above, which does
// not depend on this lookup — widening what "task" means here cannot complete a ticket
// the chained ask deliberately left open. When task is null (this was never an async
// ticket — a blocking fleet_ask's answer() call has no Task at all), there is nothing
// to complete, matching the old lookup-miss behaviour.
if (result.outcome() != Outcome.QUESTION && task != null) {
finishAsyncTask(task, result);
}
return result;
} catch (TimeoutException e) {
@@ -1059,7 +1092,7 @@ public final class MessageService {
// this fires exactly once, from whichever path completes it: finishAsyncTask(task, result)
// below on any non-QUESTION outcome of send() — a worker's fleet_reply, the CB-106
// completion fallback, a CB-109 wedge, TIMED_OUT, BUSY, or BACKEND_EXHAUSTED — the same
// finishAsyncTask reached via answer()'s finishAsyncTask(turnId, result) once a QUESTION
// finishAsyncTask reached via answer()'s own finishAsyncTask(task, result) once a QUESTION
// is resolved, completeExceptionally(t) just below when send() itself throws, or a CB-516
// abandon() on teardown. Without this, MessageService.reply's rendezvous fast path (the
// one an async ticket always takes) never told the push loop anything happened — see the
@@ -1079,7 +1112,22 @@ public final class MessageService {
finishAsyncTask(task, result);
}
} catch (Throwable t) {
task.future.completeExceptionally(t);
// fleetd #329 (F2): finishAsyncTask above already completes task.future — on its very
// first line — before doing anything else, so anything that throws afterward (inside
// finishAsyncTask's own cleanup, or from a future addition to this try block) lands
// here with the future already resolved. completeExceptionally on an already-completed
// future is a silent no-op: it returns false and does nothing, so without the check
// below the exception simply vanished — no log, no metric, nothing. Measured (see the
// ticket): temporarily reintroducing the fleetd #324 NPE reproduced 19 real exceptions
// on the ordinary path, with 74/74 tests staying green and not one log line produced.
// Do not stop completing the future first — finishAsyncTask completing before it
// cleans up is what makes a late failure harmless to the ticket's own result — only
// add the missing visibility for the case where that step, or whatever ran after it,
// has already lost the race to report through the future.
if (!task.future.completeExceptionally(t)) {
log.error("async send {} -> {} threw after its ticket was already resolved",
ticket, target, t);
}
}
});
pruneTerminalTickets();
@@ -1247,6 +1295,10 @@ public final class MessageService {
*/
private void finishAsyncTask(Task task, Reply result) {
task.future.complete(result);
if (afterFinishAsyncTaskCompleteHookForTest != null) {
// Test-only (fleetd #329, F2): see the field's own javadoc.
afterFinishAsyncTaskCompleteHookForTest.run();
}
String turnId = task.turnId;
if (turnId != null) {
if (finishAsyncTaskRaceHook != null) {
@@ -1276,6 +1328,49 @@ public final class MessageService {
this.finishAsyncTaskRaceHook = hook;
}
/**
* Null in production; test seam for fleetd #329 (F2) — invoked from {@link #finishAsyncTask(Task,
* Reply)} unconditionally, immediately after {@code task.future.complete(result)} runs (before
* {@link Task#turnId} is even read, so it fires regardless of whether this task was ever asked).
* A test installs this to force an exception into exactly the shape fleetd #329 identified:
* something throws inside {@link #sendAsync}'s executor task after the async ticket's future is
* already resolved, so the surrounding {@code catch (Throwable t)} can only report failure through
* {@code completeExceptionally} — a silent no-op on an already-completed future. Engineering a
* real exception to land in that exact post-completion window is what the ticket itself had to do
* by temporarily deleting a production guard (fleetd #324's {@code turnId != null} check); this
* hook drives the identical shape deterministically instead.
*/
private volatile Runnable afterFinishAsyncTaskCompleteHookForTest;
/**
* Test-only (fleetd #329, F2): install {@link #afterFinishAsyncTaskCompleteHookForTest}.
* Package-private so the test, in the same package, can reach it without widening any production
* API.
*/
void setAfterFinishAsyncTaskCompleteHookForTest(Runnable hook) {
this.afterFinishAsyncTaskCompleteHookForTest = hook;
}
/**
* Null in production; test seam for fleetd #329 (F3) — invoked from {@link #reply} right after
* the single local read of {@code orphan.turnId} passes its null-check and before that (now-local)
* value is used as the {@code asyncTasksByTurn} removal key. Mirrors {@link
* #finishAsyncTaskRaceHook} exactly, for the structurally identical double-read fleetd #324 fixed
* in {@link #finishAsyncTask}: a test installs this to force, deterministically, {@code
* clearAsyncQuestion}'s unlocked {@code forgetTurn=true} path nulling {@link Task#turnId} in that
* exact window, and to confirm the single-read fix tolerates it (the captured local is used
* unconditionally, so a hook that nulls the field afterward cannot affect this call).
*/
private volatile Runnable replyOrphanTurnIdRaceHookForTest;
/**
* Test-only (fleetd #329, F3): install {@link #replyOrphanTurnIdRaceHookForTest}. Package-private
* so the test, in the same package, can reach it without widening any production API.
*/
void setReplyOrphanTurnIdRaceHookForTest(Runnable hook) {
this.replyOrphanTurnIdRaceHookForTest = hook;
}
/**
* Test-only (fleetd #324): run the exact production cleanup {@link #ask}'s own timeout path runs
* unlocked — {@link #clearAsyncQuestion(String, boolean)} with {@code forgetTurn=true} — so a test
@@ -1285,14 +1380,6 @@ public final class MessageService {
clearAsyncQuestion(turnId, true);
}
/** Complete the async ticket correlated to a specific answered turn. */
private void finishAsyncTask(String turnId, Reply result) {
Task task = asyncTasksByTurn.get(turnId);
if (task != null) {
finishAsyncTask(task, result);
}
}
/** A new send must not open a waiter while an async ticket owns this worker's paused turn. */
private boolean hasAsyncQuestion(String target) {
return asyncTasksByTurn.values().stream().anyMatch(task -> target.equals(task.target));
@@ -570,228 +570,6 @@ class ConfigRefTest {
assertEquals(30, ref.get().configReload().intervalSeconds());
}
/**
* fleetd #330: {@code health:} is read both ways — {@code Fleetd.java:556-563} builds the
* monitor off the startup snapshot and never rebuilds it, but {@code Fleetd.java:648-650} reads
* {@code config.get().health()} live on every {@code fleet_profiles} call. A changed value is
* neither purely hot nor purely deferred, so it gets its own {@code split} report naming both
* halves rather than a bare "config reloaded" (which would hide the frozen half) or a plain
* {@code deferred} entry (which would hide that the coverage string already applied).
*/
@Test
void changingHealthIsReportedAsSplit(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml("""
health:
enabled: true
intervalSeconds: 30
"""));
ConfigRef ref = refFor(f);
Files.writeString(f, yaml("""
health:
enabled: true
intervalSeconds: 90
"""));
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied());
assertTrue(out.deferred().isEmpty(), out.deferred().toString());
assertEquals(1, out.split().size(), out.split().toString());
assertTrue(out.split().getFirst().startsWith("health:"), out.split().toString());
assertTrue(out.split().getFirst().contains("restart"), out.split().toString());
assertTrue(out.split().getFirst().contains("live"), out.split().toString());
assertTrue(out.summary().contains("partially live"), out.summary());
// The snapshot still carries the new value — the monitor itself is what waits for a restart.
assertEquals(90, ref.get().health().intervalSeconds());
}
/**
* fleetd #330: {@code coordinator:} is the other split key — {@code Fleetd.java:502} opens the
* {@code LeadMailbox} off the startup snapshot and never reopens it, but
* {@code HerdrPeerLauncher.java:1530} reads {@code config.get().coordinator()} live on every
* spawn to keep the broker URI env-var name out of a member's environment.
*/
@Test
void changingCoordinatorIsReportedAsSplit(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml("""
coordinator:
selfId: mac-a
"""));
ConfigRef ref = refFor(f);
Files.writeString(f, yaml("""
coordinator:
selfId: mac-b
"""));
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied());
assertTrue(out.deferred().isEmpty(), out.deferred().toString());
assertEquals(1, out.split().size(), out.split().toString());
assertTrue(out.split().getFirst().startsWith("coordinator:"), out.split().toString());
assertTrue(out.split().getFirst().contains("restart"), out.split().toString());
assertTrue(out.split().getFirst().contains("live"), out.split().toString());
// The snapshot still carries the new value — the LeadMailbox connection is what waits for a
// restart; selfId names this daemon's own inbox queue and a peer cannot discover a rename.
assertEquals("mac-b", ref.get().coordinator().selfId());
}
/**
* fleetd #333: {@code fleet:} is split too — {@code fleet.leaders} is read only at
* {@code Fleetd.java:281} to build the {@code LeadTabScanner}'s identity map (fed into
* {@code CallerResolver}) and to auto-launch leads via {@code LeadLauncher.ensureLeads()}, and
* neither is rebuilt on reload, while the rest of {@code fleet:} (role pools, charters,
* tabLabel) is read live through the {@code CompositePeerLauncher} supplier. Before this fix
* {@code fleet} sat in the coverage checker's hot-excluded escape hatch, so a reload that
* changed only {@code fleet.leaders} reported a bare "config reloaded" — exactly the
* under-claim the split class exists to prevent for {@code health}/{@code coordinator}.
*/
@Test
void changingFleetLeadersIsReportedAsSplit(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml("""
fleet:
leaders:
opus:
tab: "lead: opus-a"
profile: sonnet
"""));
ConfigRef ref = refFor(f);
Files.writeString(f, yaml("""
fleet:
leaders:
opus:
tab: "lead: opus-b"
profile: sonnet
"""));
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied());
assertTrue(out.deferred().isEmpty(), out.deferred().toString());
assertEquals(1, out.split().size(), out.split().toString());
assertTrue(out.split().getFirst().startsWith("fleet:"), out.split().toString());
assertTrue(out.split().getFirst().contains("restart"), out.split().toString());
assertTrue(out.split().getFirst().contains("live"), out.split().toString());
assertTrue(out.summary().contains("partially live"), out.summary());
// The snapshot still carries the new value — the LeadTabScanner's identity map and
// LeadLauncher's auto-launch are what wait for a restart; a reload rebuilds neither.
assertEquals("lead: opus-b", ref.get().fleet().leaders().get("opus").tab());
}
/**
* fleetd #333: {@code fleet.leaders} is the ONLY frozen part of {@code fleet:}. A reload that
* changes {@code tabLabel} (or charters, or a role pool) without touching {@code fleet.leaders}
* must stay fully hot with nothing reported — proving {@link ConfigRef#changedSplitKeys}
* compares {@code fleet.leaders} specifically rather than the whole {@code Fleet} record, which
* would over-claim "needs a restart" for a change that is genuinely all live (the mirror
* mistake of the under-claim this class exists to prevent).
*/
@Test
void changingFleetTabLabelWithoutLeadersStaysFullyHot(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml("""
fleet:
tabLabel: "{role}: {profile} #{n}"
"""));
ConfigRef ref = refFor(f);
Files.writeString(f, yaml("""
fleet:
tabLabel: "[{profile}] {role}"
"""));
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied());
assertTrue(out.deferred().isEmpty(), out.deferred().toString());
assertTrue(out.split().isEmpty(), out.split().toString());
assertEquals("config reloaded", out.summary());
assertEquals("[{profile}] {role}", ref.get().fleet().tabLabel());
}
/**
* A split change must not refuse the reload (invariant 2 of fleetd #330) and
* {@code Outcome.applied()} must stay {@code true} (invariant 3) — unlike a cold change, the
* live half of a split key genuinely took effect, so refusing would throw that away.
*/
@Test
void aSplitChangeDoesNotRefuseTheReload(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml("coordinator:\n selfId: mac-a\n"));
ConfigRef ref = refFor(f);
Files.writeString(f, yaml("coordinator:\n selfId: mac-b\n"));
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied());
assertTrue(out.error() == null);
assertTrue(out.coldKeys().isEmpty());
}
/**
* Both split keys can change in one reload — the report names both, and a caller reading
* {@code split} does not have to guess which half of which key already applied.
*/
@Test
void changingBothSplitKeysReportsBoth(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml("""
health:
enabled: true
coordinator:
selfId: mac-a
"""));
ConfigRef ref = refFor(f);
Files.writeString(f, yaml("""
health:
enabled: false
coordinator:
selfId: mac-b
"""));
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied());
assertEquals(2, out.split().size(), out.split().toString());
assertTrue(out.split().stream().anyMatch(s -> s.startsWith("health:")), out.split().toString());
assertTrue(out.split().stream().anyMatch(s -> s.startsWith("coordinator:")), out.split().toString());
}
/**
* A split key and a deferred key changing in the same reload must both show up, each in its own
* list — proving the two fields do not step on each other and {@link ConfigRef.Outcome#summary()}
* reports both halves of the message.
*/
@Test
void aSplitChangeAndADeferredChangeCoexist(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml("""
coordinator:
selfId: mac-a
lifecycle:
drainTimeoutSeconds: 30
"""));
ConfigRef ref = refFor(f);
Files.writeString(f, yaml("""
coordinator:
selfId: mac-b
lifecycle:
drainTimeoutSeconds: 60
"""));
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied());
assertEquals(java.util.List.of("lifecycle"), out.deferred());
assertEquals(1, out.split().size(), out.split().toString());
assertTrue(out.split().getFirst().startsWith("coordinator:"), out.split().toString());
assertTrue(out.summary().contains("need a restart") || out.summary().contains("needs a restart"),
out.summary());
assertTrue(out.summary().contains("partially live"), out.summary());
}
@Test
void aFixedRefHasNoFileAndRefusesToReload() {
FleetConfig cfg = new FleetConfig(null, null, null, null, null, null,
@@ -1,171 +0,0 @@
package dev.ltms.fleet.config;
import org.junit.jupiter.api.Test;
import java.lang.reflect.RecordComponent;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Set;
import java.util.TreeSet;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #330: a new top-level {@link FleetConfig} record component must be triaged into a reload
* class before it ships, or it repeats fleetd #323 ({@code worktreeGroup} missing from {@code
* ConfigRef.changedDeferredKeys}) and fleetd #326 ({@code primary}/{@code configReload} missing the
* same way) — a key silently absent from {@link ConfigRef}'s reload machinery, so a reload changing
* only that key reports a bare "config reloaded" for a change the running daemon never picked up.
*
* <p>This is the top-level counterpart of {@link ConfigRefProfileCoverageTest}: instead of
* enumerating {@link FleetConfig.Profile}'s record components, it enumerates {@link FleetConfig}'s
* own — {@code bind}, {@code health}, {@code memberCredentials}, and so on — and requires each to
* fall into exactly one of four homes: {@link ConfigRef#COLD_KEYS}, {@link #DEFERRED_TOP_LEVEL_KEYS}
* (compared in {@code ConfigRef.changedDeferredKeys}), {@link ConfigRef#SPLIT_KEYS}, or
* {@link #HOT_EXCLUDED_TOP_LEVEL_KEYS} (the escape hatch: read live off the config supplier, so no
* reload bookkeeping is needed for it at all).
*
* <h2>What this checker can and cannot prove</h2>
* It proves the record's <em>shape</em> is fully triaged: every one of {@code FleetConfig}'s
* components sits in exactly one of the four sets, none sits in two, and the escape hatch
* ({@link #HOT_EXCLUDED_TOP_LEVEL_KEYS}) cannot silently grow without a visible diff to this file.
* That is what "a new component cannot be added without someone triaging it" means in practice.
*
* <p>It CANNOT prove that any of the citations are <em>true</em>. "Compared in {@code
* changedDeferredKeys}" and "read live off {@code config.get()}" are facts about {@code
* ConfigRef.java}, {@code Fleetd.java} and {@code HerdrPeerLauncher.java} that a reflection-only
* test over {@code FleetConfig}'s shape has no way to inspect — this test would pass identically
* whether or not the cited line still does what the comment next to it says. Trust the citation
* because a person read the source (the exact call sites are named next to each set below), not
* because this test is green. {@link ConfigRefTest} is what behaviourally proves the deferred and
* split keys it covers actually get reported; {@link ConfigRefProfileCoverageTest} does the same,
* behaviourally, for {@code FleetConfig.Profile}'s own fields.
*/
class ConfigRefTopLevelCoverageTest {
private static final RecordComponent[] COMPONENTS = FleetConfig.class.getRecordComponents();
/**
* Top-level components whose change {@code ConfigRef.changedDeferredKeys} reads and reports on
* — verified by reading that method as of fleetd #330, not derived from this test.
* {@code spawnReadyTimeoutMs}/{@code spawnReadyPollMs} are compared together and reported under
* one combined label ({@code "spawnReady*"}); {@code profiles} is compared twice over — once for
* added/removed profile names, once for an existing profile's launch settings — and that second
* comparison excludes {@code weight}/{@code maxLoad}/{@code credentialId} as hot sub-fields,
* which is what {@link ConfigRefProfileCoverageTest} exists to keep honest at the sub-field
* level. {@code profiles} itself still belongs here, not in the hot-exclusion set below: most of
* a profile's fields are NOT read live, so citing "read live off the config supplier" for the
* whole top-level key would be false.
*/
private static final Set<String> DEFERRED_TOP_LEVEL_KEYS = Set.of(
"guard", "worktreeRoot", "worktreeGroup", "primary", "configReload",
"leadHeartbeat", "lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs",
"quarantineCooldownSeconds", "profiles");
/**
* The escape hatch: top-level components with no reload bookkeeping at all, because every read
* of them goes live through {@link ConfigRef#get()} rather than off a startup snapshot. A
* component belongs here ONLY if that is true — never because adding it here makes this test
* pass. fleetd #323 is the cautionary tale for exactly this pattern: the identical hatch on
* {@code ConfigRef.LAUNCH_SETTINGS_EXCLUDED} let two profile fields be silently re-broken with
* the whole suite green, and it was only caught by mutating the checker itself (see this class's
* own mutation test below, and {@code ConfigRefProfileCoverageTest}'s equivalent).
*
* <ul>
* <li>{@code placement} — read live by the placement policy on every spawn (class doc, Hot
* bullet; {@code ConfigRefTest.aConsumerHoldingTheRefSeesTheNewValue} proves it
* behaviourally).</li>
* <li>{@code memberCredentials} — read live at {@code Fleetd.java:198, 205, 729}.</li>
* <li>{@code memberLoginShell} — read live at
* {@code HerdrPeerLauncher#configuredMemberLoginShell}.</li>
* </ul>
*
* <p>{@code fleet} used to sit here too, on the strength of most of it (role pools, charters,
* tabLabel) being read the same live way — but {@code fleet.leaders} inside the same key is
* read only at startup and genuinely needs a restart, which is a real gap this escape hatch
* cannot represent: it excuses a whole top-level key from reload bookkeeping, and {@code fleet}
* needed exactly half of it excused. fleetd #333 moved it into {@link ConfigRef#SPLIT_KEYS}
* instead, where {@code changedSplitKeys} reports the frozen half by name. That is the
* cautionary tale this test's own javadoc already told: this checker proves the record's shape
* is triaged, never that a bucket a key sits in is the right one — a person has to read the
* source, which is exactly how fleetd #333 found {@code fleet} sitting in the wrong bucket
* while this test stayed green throughout.</p>
*/
private static final Set<String> HOT_EXCLUDED_TOP_LEVEL_KEYS =
Set.of("placement", "memberCredentials", "memberLoginShell");
@Test
void everyTopLevelComponentIsAccountedForInExactlyOneClass() {
Set<String> allNames = new TreeSet<>();
for (RecordComponent rc : COMPONENTS) {
allNames.add(rc.getName());
}
Set<String> cold = ConfigRef.COLD_KEYS;
Set<String> split = ConfigRef.SPLIT_KEYS;
Set<String> deferred = DEFERRED_TOP_LEVEL_KEYS;
Set<String> hot = HOT_EXCLUDED_TOP_LEVEL_KEYS;
// Typo guard on each set — the same check ConfigRefProfileCoverageTest runs on
// LAUNCH_SETTINGS_EXCLUDED. A name that does not exist on FleetConfig is a silent no-op.
assertTrue(allNames.containsAll(cold),
"ConfigRef.COLD_KEYS names a component that does not exist on FleetConfig: " + cold);
assertTrue(allNames.containsAll(split),
"ConfigRef.SPLIT_KEYS names a component that does not exist on FleetConfig: " + split);
assertTrue(allNames.containsAll(deferred),
"DEFERRED_TOP_LEVEL_KEYS names a component that does not exist on FleetConfig: " + deferred);
assertTrue(allNames.containsAll(hot),
"HOT_EXCLUDED_TOP_LEVEL_KEYS names a component that does not exist on FleetConfig: " + hot);
// The escape hatch is pinned. Growing it requires editing this line — a visible, deliberate
// diff, not a quiet one. See the field javadoc above for what "belongs here" actually means.
assertEquals(Set.of("placement", "memberCredentials", "memberLoginShell"), hot,
"HOT_EXCLUDED_TOP_LEVEL_KEYS changed. A component belongs here ONLY if it is read "
+ "live off the config supplier, never because adding it makes this test "
+ "pass. If you are adding one to silence this test, that is fleetd #323 "
+ "happening again: account for it in ConfigRef.changedDeferredKeys (or "
+ "COLD_KEYS/SPLIT_KEYS) instead. If it really is read live, name where and "
+ "update this assertion and the field javadoc together.");
// No component may sit in two buckets at once — the denominator check below could not catch
// that on its own (two buckets double-booking one key still sums to the right total if
// another key is simultaneously missing), so check every pair directly and name the culprit.
record Bucket(String name, Set<String> keys) {}
List<Bucket> buckets = List.of(
new Bucket("COLD_KEYS", cold), new Bucket("SPLIT_KEYS", split),
new Bucket("DEFERRED_TOP_LEVEL_KEYS", deferred), new Bucket("HOT_EXCLUDED_TOP_LEVEL_KEYS", hot));
for (int i = 0; i < buckets.size(); i++) {
for (int j = i + 1; j < buckets.size(); j++) {
Set<String> overlap = new LinkedHashSet<>(buckets.get(i).keys());
overlap.retainAll(buckets.get(j).keys());
assertEquals(Set.of(), overlap, "a component is in both " + buckets.get(i).name()
+ " and " + buckets.get(j).name() + ": " + overlap);
}
}
Set<String> union = new TreeSet<>();
union.addAll(cold);
union.addAll(split);
union.addAll(deferred);
union.addAll(hot);
System.out.printf(
"FleetConfig top-level coverage — %d components total: %d cold %s, %d deferred %s, "
+ "%d split %s, %d hot-excluded %s%n",
allNames.size(), cold.size(), cold, deferred.size(), deferred, split.size(), split,
hot.size(), hot);
Set<String> missing = new TreeSet<>(allNames);
missing.removeAll(union);
assertEquals(Set.of(), missing,
"these FleetConfig components are in none of COLD_KEYS, DEFERRED_TOP_LEVEL_KEYS, "
+ "SPLIT_KEYS or HOT_EXCLUDED_TOP_LEVEL_KEYS — triage each one into whichever "
+ "actually describes it: " + missing);
assertEquals(allNames.size(), cold.size() + deferred.size() + split.size() + hot.size(),
"counts don't sum to the component total even though every component was found in "
+ "the union — " + allNames.size() + " components, " + cold.size()
+ " cold + " + deferred.size() + " deferred + " + split.size() + " split + "
+ hot.size() + " hot-excluded");
}
}
@@ -1,219 +0,0 @@
package dev.ltms.fleet.config;
import org.junit.jupiter.api.Test;
import java.lang.reflect.Constructor;
import java.lang.reflect.RecordComponent;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.TreeSet;
import static org.junit.jupiter.api.Assertions.assertEquals;
/**
* fleetd #333, finding F2: {@link ConfigRefTopLevelCoverageTest} proves every {@link FleetConfig}
* top-level component sits in exactly one of {@link ConfigRef#COLD_KEYS}, {@code
* DEFERRED_TOP_LEVEL_KEYS}, {@link ConfigRef#SPLIT_KEYS} or the hot-excluded set. It does
* <strong>not</strong> prove that a key's membership in {@code COLD_KEYS} or {@code SPLIT_KEYS}
* corresponds to any actual comparison in {@link ConfigRef}: a key can sit in either set with no
* branch in {@code changedColdKeys}/{@code changedSplitKeys} checking it, and both
* {@link ConfigRefTopLevelCoverageTest} and the "kept in step" {@code assert} inside each of those
* methods stay green, because neither one reads the method body — the coverage test only reads
* set membership, and the assert only checks that reported entries are a SUBSET of the set, never
* that every set member produced a reported entry.
*
* <p>Measured directly, live, while fixing fleetd #333: dropping the {@code coordinator} branch out
* of {@code ConfigRef.changedSplitKeys} while leaving {@code "coordinator"} in
* {@link ConfigRef#SPLIT_KEYS} left {@link ConfigRefTopLevelCoverageTest} and the in-method assert
* both green — only a hand-written behavioural case in {@link ConfigRefTest} caught it, because it
* happened to name that exact key. This is the {@link ConfigRefProfileCoverageTest} mechanism one
* level up, generalised over every {@code COLD_KEYS}/{@code SPLIT_KEYS} member rather than one
* hand-picked field: enumerate {@link FleetConfig}'s own record components by reflection, build "a
* config where only {@code <key>} differs" for each cold/split key, and call the real
* {@link ConfigRef#changedColdKeys}/{@link ConfigRef#changedSplitKeys} methods (made
* package-private for exactly this, the same reason {@link ConfigRef#sameLaunchSettings} already
* is) to prove each one is actually reported — not assumed from a set literal.
*
* <h2>What this deliberately does NOT cover</h2>
* {@code DEFERRED_TOP_LEVEL_KEYS} is not exercised here. That bucket carries the identical
* one-way risk in principle — a key added to it with no matching branch in
* {@code changedDeferredKeys} would pass {@link ConfigRefTopLevelCoverageTest} exactly the way
* {@code coordinator} passed it above — but this test stays narrow to {@code COLD_KEYS} and
* {@code SPLIT_KEYS} for two reasons. First, that is where fleetd #333 actually found and measured
* the gap (F1 was a live instance of it). Second, most of {@code DEFERRED_TOP_LEVEL_KEYS} already
* carries an individual behavioural test in {@link ConfigRefTest} naming it by key — {@code
* worktreeGroup}, {@code primary}, {@code configReload}, {@code profiles}' launch settings /
* weight-maxLoad / exhaustedPattern / errorPattern / ideProjectDir — which is the same protection
* this class gives {@code COLD_KEYS}/{@code SPLIT_KEYS}, just written by hand per key instead of
* generated by reflection over the whole set. {@code guard}, {@code lifecycle},
* {@code leadHeartbeat}, {@code spawnReadyTimeoutMs}/{@code spawnReadyPollMs} and
* {@code quarantineCooldownSeconds} do NOT have a dedicated behavioural test naming them, so the
* one-way gap fleetd's own memory notes ("pre-existing on COLD_KEYS and on the test's own
* DEFERRED_TOP_LEVEL_KEYS") is real and not fully closed by this class — extending this mechanism's
* {@code BASE}/{@code ALT} map to cover every top-level component and adding a third
* {@code everyDeferredKeyIsActuallyReportedByChangedDeferredKeys} test is the natural next step, left
* for whoever next finds a deferred key with the same shape as this ticket's {@code fleet.leaders}.
*/
class ConfigRefTopLevelReportingCoverageTest {
private static final RecordComponent[] COMPONENTS = FleetConfig.class.getRecordComponents();
/**
* One valid value per top-level {@link FleetConfig} component — "the a value". Components not
* exercised by either test below ({@code profiles}, {@code guard}, {@code worktreeRoot}, …) are
* left {@code null}/empty; {@link FleetConfig}'s compact constructor only normalizes
* {@code profiles}, so every other field accepts {@code null} unmutated.
*/
private static final Map<String, Object> BASE = baseValues();
/** The same shape, each value distinct from {@link #BASE} — "the b value" — for COLD_KEYS/SPLIT_KEYS only. */
private static final Map<String, Object> ALT = altValues();
private static Map<String, Object> baseValues() {
Map<String, Object> v = new LinkedHashMap<>();
v.put("bind", new FleetConfig.Bind("127.0.0.1", 8765));
v.put("herdrSocket", "~/.config/herdr/a.sock");
v.put("memberHerdrSocket", "~/.config/herdr/member-a.sock");
v.put("profiles", Map.of());
v.put("guard", null);
v.put("worktreeRoot", null);
v.put("lifecycle", null);
v.put("spawnReadyTimeoutMs", null);
v.put("spawnReadyPollMs", null);
v.put("broker", new FleetConfig.Broker("amqp://a", null, 1));
v.put("primary", null);
v.put("fleet", new FleetConfig.Fleet(
Map.of("opus", new FleetConfig.Leader("sonnet", "lead: opus-a", 1, null, 10,
"claude", null, null, null)),
Map.of(), Map.of(), Map.of(), Map.of(), "{role}: {profile} #{n}"));
v.put("leadHeartbeat", null);
v.put("health", new FleetConfig.Health(true, 30, 600, null, null));
v.put("placement", null);
v.put("auth", new FleetConfig.Auth("loopback-trust", null));
v.put("configReload", null);
v.put("quarantineCooldownSeconds", null);
v.put("memberCredentials", null);
v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-a", null, "self-a", 1));
v.put("worktreeGroup", null);
v.put("memberLoginShell", null);
assertNamesMatchComponents(v);
return v;
}
private static Map<String, Object> altValues() {
Map<String, Object> v = new LinkedHashMap<>();
v.put("bind", new FleetConfig.Bind("127.0.0.2", 8766));
v.put("herdrSocket", "~/.config/herdr/b.sock");
v.put("memberHerdrSocket", "~/.config/herdr/member-b.sock");
v.put("profiles", Map.of());
v.put("guard", null);
v.put("worktreeRoot", null);
v.put("lifecycle", null);
v.put("spawnReadyTimeoutMs", null);
v.put("spawnReadyPollMs", null);
v.put("broker", new FleetConfig.Broker("amqp://b", null, 2));
v.put("primary", null);
// Differs from BASE.fleet only in fleet.leaders (a different tab for "opus") — the frozen
// sub-field ConfigRef.changedSplitKeys actually compares. A Fleet that instead differed only
// in tabLabel would correctly NOT be reported (see
// ConfigRefTest.changingFleetTabLabelWithoutLeadersStaysFullyHot) and would wrongly fail this
// test — that is by design, not a gap: this map exists to prove fleet.leaders is covered.
v.put("fleet", new FleetConfig.Fleet(
Map.of("opus", new FleetConfig.Leader("sonnet", "lead: opus-b", 1, null, 10,
"claude", null, null, null)),
Map.of(), Map.of(), Map.of(), Map.of(), "{role}: {profile} #{n}"));
v.put("leadHeartbeat", null);
v.put("health", new FleetConfig.Health(false, 90, 900, null, null));
v.put("placement", null);
v.put("auth", new FleetConfig.Auth("token", "TOKEN_ENV"));
v.put("configReload", null);
v.put("quarantineCooldownSeconds", null);
v.put("memberCredentials", null);
v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-b", null, "self-b", 2));
v.put("worktreeGroup", null);
v.put("memberLoginShell", null);
assertNamesMatchComponents(v);
return v;
}
private static void assertNamesMatchComponents(Map<String, Object> values) {
Set<String> names = new TreeSet<>();
for (RecordComponent rc : COMPONENTS) {
names.add(rc.getName());
}
assertEquals(names, new TreeSet<>(values.keySet()),
"this test's value map has drifted from FleetConfig's actual top-level components — "
+ "update BASE/ALT alongside the record");
}
private static FleetConfig configOf(Map<String, Object> values) throws ReflectiveOperationException {
Class<?>[] types = Arrays.stream(COMPONENTS).map(RecordComponent::getType).toArray(Class<?>[]::new);
Object[] args = Arrays.stream(COMPONENTS).map(rc -> values.get(rc.getName())).toArray();
Constructor<FleetConfig> ctor = FleetConfig.class.getDeclaredConstructor(types);
return ctor.newInstance(args);
}
/** {@code BASE} with exactly one named top-level component swapped for its {@code ALT} value. */
private static FleetConfig mutate(String componentName) throws ReflectiveOperationException {
Map<String, Object> values = new LinkedHashMap<>(BASE);
values.put(componentName, ALT.get(componentName));
return configOf(values);
}
/**
* The mechanism fleetd #333 F2 asked for, applied to {@link ConfigRef#COLD_KEYS}: mutate each
* cold key in isolation and prove {@code changedColdKeys} actually names it, not just that
* {@code COLD_KEYS} claims it does.
*/
@Test
void everyColdKeyIsActuallyReportedByChangedColdKeys() throws ReflectiveOperationException {
FleetConfig base = configOf(BASE);
List<String> uncovered = new ArrayList<>();
for (String key : new TreeSet<>(ConfigRef.COLD_KEYS)) {
FleetConfig mutated = mutate(key);
if (!ConfigRef.changedColdKeys(base, mutated).contains(key)) {
uncovered.add(key);
}
}
System.out.printf(
"ConfigRef.changedColdKeys reporting coverage — %d COLD_KEYS, %d verified%n",
ConfigRef.COLD_KEYS.size(), ConfigRef.COLD_KEYS.size() - uncovered.size());
assertEquals(List.of(), uncovered,
"these keys are in ConfigRef.COLD_KEYS but mutating them alone produces no matching "
+ "entry from changedColdKeys — a set entry with no comparison behind it: "
+ uncovered);
}
/**
* The same mechanism applied to {@link ConfigRef#SPLIT_KEYS}: mutate each split key in
* isolation and prove {@code changedSplitKeys} actually names it (message starting with
* {@code "<key>:"}), not just that {@code SPLIT_KEYS} claims it does. This is the exact check
* that would have failed fleetd #333's own reproduction — dropping the {@code coordinator}
* branch from {@code changedSplitKeys} while {@code "coordinator"} stayed in {@code SPLIT_KEYS}
* — see the mutation proof in the fleetd #333 PR description; the earlier checkers here (the
* shape test and the in-method assert) do not.
*/
@Test
void everySplitKeyIsActuallyReportedByChangedSplitKeys() throws ReflectiveOperationException {
FleetConfig base = configOf(BASE);
List<String> uncovered = new ArrayList<>();
for (String key : new TreeSet<>(ConfigRef.SPLIT_KEYS)) {
FleetConfig mutated = mutate(key);
List<String> split = ConfigRef.changedSplitKeys(base, mutated);
if (split.stream().noneMatch(s -> s.startsWith(key + ":"))) {
uncovered.add(key);
}
}
System.out.printf(
"ConfigRef.changedSplitKeys reporting coverage — %d SPLIT_KEYS, %d verified%n",
ConfigRef.SPLIT_KEYS.size(), ConfigRef.SPLIT_KEYS.size() - uncovered.size());
assertEquals(List.of(), uncovered,
"these keys are in ConfigRef.SPLIT_KEYS but mutating them alone produces no matching "
+ "entry from changedSplitKeys — a set entry with no comparison behind it, "
+ "exactly the fleetd #333 F2 shape: " + uncovered);
}
}
@@ -1,5 +1,9 @@
package dev.ltms.fleet.msg;
import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.FakeHerdr;
@@ -11,6 +15,7 @@ import dev.ltms.fleet.mcp.PrimaryRegistry;
import dev.ltms.fleet.inject.Injector;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.slf4j.LoggerFactory;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
@@ -1041,6 +1046,105 @@ class MessageServiceTest {
"the reply completed its own ticket directly and never touched the inbox");
}
/**
* fleetd #329 (F1). {@code answer()} completes the async ticket by looking {@code turnId} up in
* {@code asyncTasksByTurn} a SECOND time (the first is at :991, purely to re-register the
* {@code asyncTasksByWaiter} entry for #282's chained-ask case). That second lookup races
* {@code ask()}'s own unlocked timeout cleanup ({@code clearAsyncQuestion(turnId, true)}, run from
* {@code markAskTimedOut} + forgetting): {@code ask()}'s {@code ticket.answer().get(timeoutMillis)}
* can time out at essentially the same instant {@code answer()}'s {@code rendezvous.answerAsk}
* call above already succeeded and unblocked the worker. fleetd #324's single-read fix does not
* help here — that fixed a torn read of one already-held {@link MessageService} internal
* {@code Task}; this is a second, independent map lookup by a caller that no longer holds the
* {@code Task} it already found once.
*
* <p>This test does not wait for that real race to land on its own schedule — it drives the exact
* sequence the ticket describes (worker asks, primary answers, worker's real reply arrives) and
* fires the identical production cleanup {@code ask()}'s timeout path runs
* ({@code clearAsyncQuestion(turnId, true)}, via the {@code forgetTurnForTest} seam fleetd #324
* already merged) at the point between "the worker's own {@code ask()} call has unblocked" and
* "the worker's real {@code fleet_reply} arrives" — the exact window fleetd #329 names.
*
* <p>What this proves: given that exact interleaving, the async ticket must still resolve to the
* worker's real reply, not stay stuck {@link MessageService.Phase#PENDING} forever. What it does
* not prove: that the interleaving itself is reachable in production on its own timing — that is
* established by reading the code (see the ticket's "path in"), not by this test, for the same
* reason fleetd #324's own race test says so.
*/
@Test
void aReplyRacingAsksTimeoutCleanupStillCompletesTheAsyncTicket() throws Exception {
String ticket = messages.sendAsync(T, "task that asks");
awaitWaiting();
injectDelivery();
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
String turnId = asking.turnId();
CompletableFuture<MessageService.Reply> answer =
CompletableFuture.supplyAsync(() -> messages.answer(turnId, "config.yaml", 5000));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer(),
"the worker's own ask() call must have already unblocked with the primary's answer "
+ "before we force the race below");
// ask()'s own unlocked timeout cleanup can forget this exact turnId at essentially the same
// instant answer() has already unblocked the worker and is now waiting on the resumed turn's
// real reply — reproduce that interleaving directly instead of trying to win a real race.
messages.forgetTurnForTest(turnId);
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome(),
"the primary's own answer() call must still see the worker's real reply");
MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE);
assertEquals("PR opened: https://example/pulls/42", done.reply(),
"fleet_poll{ticket} must return the worker's real reply, not stay PENDING forever "
+ "just because ask()'s timeout cleanup forgot this turnId first");
}
/**
* fleetd #329 (F3). {@link MessageService#reply} reads the same orphan {@code Task}'s {@code
* turnId} twice on this path (once to check it is non-null, once as the {@code
* asyncTasksByTurn.remove} key) — the exact double-read shape fleetd #324 fixed in {@code
* finishAsyncTask}. This test drives the same real sequence as {@code
* aReplyAfterAnswerTimesOutStillCompletesTheAsyncTicket} (worker asks, primary answers, the
* primary's own bounded wait for the resumed turn expires) to reach a task with {@code turnId}
* genuinely stamped and no live rendezvous waiter open — the state {@code askAnsweredAsyncTasks}
* matches here — then fires the identical production cleanup {@code ask()}'s own timeout path
* runs ({@code clearAsyncQuestion(turnId, true)}, via the {@code forgetTurnForTest} seam fleetd
* #324 merged) at the point between the check and the removal use, via a dedicated test hook
* mirroring {@code finishAsyncTaskRaceHook}.
*/
@Test
void replyToAnOrphanedTaskSurvivesTurnIdGoingNullBetweenItsTwoReads() throws Exception {
String ticket = messages.sendAsync(T, "task that asks");
awaitWaiting();
injectDelivery();
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
String turnId = asking.turnId();
MessageService.Reply answerReply = messages.answer(turnId, "config.yaml", 150);
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome(),
"the primary's own bounded wait must give up first, leaving turnId stamped with no "
+ "live waiter — the state reply()'s F3 code path matches");
messages.setReplyOrphanTurnIdRaceHookForTest(() -> messages.forgetTurnForTest(turnId));
try {
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE);
assertEquals("PR opened: https://example/pulls/42", done.reply(),
"the ticket must still resolve to the worker's real reply despite the forced race");
} finally {
messages.setReplyOrphanTurnIdRaceHookForTest(null);
}
}
/**
* fleetd #307's ambiguity guard: an ask timeout frees its target ({@code hasAsyncQuestion}
* becomes false the instant it lapses — proven above), so a second, independent delegation can
@@ -1151,6 +1255,63 @@ class MessageServiceTest {
assertEquals("done", done.reply());
}
// --- fleetd #329 (F2): an exception after the ticket's future completes must reach a log -----
//
// sendAsync's own executor task ends with catch (Throwable t) { task.future.completeExceptionally(t); }
// — but finishAsyncTask completes that same future on its first line, so anything that throws
// afterward hits an already-completed future: completeExceptionally returns false and does
// nothing, and (before this fix) nothing logged it either. Measured on the ticket: temporarily
// reintroducing the fleetd #324 NPE reproduced 19 real exceptions on the ordinary path with
// 74/74 tests staying green and zero log lines. There is no reachable production call site where
// finishAsyncTask throws after completing the future (fleetd #324 already closed the one that
// used to), so this test drives the shape directly via a dedicated test-only hook
// (afterFinishAsyncTaskCompleteHookForTest) rather than trying to engineer a real exception into
// that narrow window — the same technique fleetd #324's own race test and this ticket's F1/F3
// tests use for their own races.
private static ListAppender<ILoggingEvent> attachMessageServiceLog() {
Logger logger = (Logger) LoggerFactory.getLogger(MessageService.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
return appender;
}
private static void detachMessageServiceLog(ListAppender<ILoggingEvent> appender) {
((Logger) LoggerFactory.getLogger(MessageService.class)).detachAppender(appender);
}
@Test
void anExceptionAfterTheTicketFutureCompletesStillReachesTheLog() throws Exception {
ListAppender<ILoggingEvent> appender = attachMessageServiceLog();
try {
messages.setAfterFinishAsyncTaskCompleteHookForTest(() -> {
throw new RuntimeException("PROBE-329-F2");
});
String ticket = messages.sendAsync(T, "do the task");
awaitWaiting();
assertTrue(rendezvous.resolve(T, "done"));
// finishAsyncTask's own first line already completed the future with the real result
// BEFORE the hook threw — F2 is a visibility gap, not a correctness gap for the ticket
// itself, so the ticket's own outcome must be unaffected by the swallowed exception.
MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE);
assertEquals("done", done.reply());
assertTrue(appender.list.stream().anyMatch(e ->
e.getLevel() == Level.ERROR
&& e.getFormattedMessage().contains(ticket)
&& e.getThrowableProxy() != null
&& "PROBE-329-F2".equals(e.getThrowableProxy().getMessage())),
"an exception thrown after the ticket's future already completed must still "
+ "reach the log, not vanish silently");
} finally {
messages.setAfterFinishAsyncTaskCompleteHookForTest(null);
detachMessageServiceLog(appender);
}
}
// --- CB-582: fleet_status pendingAsk() ------------------------------------------------------
@Test