Compare commits

..

1 Commits

Author SHA1 Message Date
Dai Ha 554395b104 fleetd#330: split reload class for health/coordinator + top-level coverage
CI / contract (pull_request) Successful in 53s
CI / build (pull_request) Successful in 2m17s
Unit 1: ConfigRef gets a fourth reload class, `split`, for keys read both
off the startup snapshot and live off config.get() at different sites
(health:, coordinator:). A split change is accepted (Outcome.applied()
stays true) and reported by name, naming which half is live and which
needs a restart, via a new Outcome.split() field kept separate from
deferred() since the two carry different guarantees for any caller that
branches on them, not just prose in summary(). Class doc updated: four
classes now, denominator note no longer calls health/coordinator
undecided.

Unit 2: ConfigRefTopLevelCoverageTest enumerates FleetConfig's 22
top-level record components and requires each to sit in exactly one of
COLD_KEYS, a pinned "compared in changedDeferredKeys" set, SPLIT_KEYS, or
a pinned hot-exclusion escape hatch — printing its own denominator and
pinning the escape hatch's exact contents the way #323 asked for.

Deviates from the issue's starting values by one key: `profiles` moves
from the suggested Hot bucket into the deferred bucket, because
changedDeferredKeys demonstrably compares it (add/remove and launch
settings), and citing "read live off the config supplier" for the whole
key would be false — most Profile fields are not read live, only
weight/maxLoad/credentialId are (and those are already covered by
ConfigRefProfileCoverageTest). Cold=5, split=2, deferred=11, hot=4,
total=22 — verified against the record and against ConfigRef's code, not
copied from the issue.
2026-09-04 15:16:45 +07:00
5 changed files with 453 additions and 297 deletions
@@ -22,7 +22,7 @@ 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 three classes, and the difference is about what already exists when the reload
* Keys fall into four classes, and the difference is about what already exists when the reload
* happens — not about how important the key is.
*
* <ul>
@@ -71,35 +71,63 @@ 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) — read <em>both</em> ways at different sites, so the
* key does not fit any class above as a whole: {@code health:} and {@code coordinator:}.
* 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>
* </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 #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>The denominator, measured on 2026-09-04 (fleetd #330).</strong> {@code FleetConfig} has
* 22 top-level record components. Two of them are named nowhere in this file, and the reason is the
* same for both: {@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.
* {@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.
* <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. Twice now — {@code worktreeGroup} (#323) and
* {@code primary}/{@code configReload} (#326) — 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}.
*
* <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.
* 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.
*
* <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
@@ -109,10 +137,24 @@ 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. */
private static final Set<String> COLD_KEYS =
/**
* 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 =
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");
private final Path path;
private final AtomicReference<FleetConfig> current;
@@ -140,25 +182,37 @@ 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 for a restart
* @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 error the parse or validation failure that refused the reload, else {@code null}
*/
public record Outcome(boolean applied, List<String> coldKeys, List<String> deferred,
String error) {
List<String> split, 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(), null);
return new Outcome(false, keys, List.of(), List.of(), null);
}
static Outcome failed(String error) {
return new Outcome(false, List.of(), List.of(), error);
return new Outcome(false, List.of(), List.of(), List.of(), error);
}
/** A one-line summary for the operator — the reason, not just the verdict. */
@@ -170,11 +224,18 @@ 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()) {
return "config reloaded; these changes need a restart to take effect: "
+ String.join(", ", deferred);
if (deferred.isEmpty() && split.isEmpty()) {
return "config reloaded";
}
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));
}
if (!split.isEmpty()) {
out.append("; partially live — ").append(String.join(" | ", split));
}
return out.toString();
}
}
@@ -215,8 +276,9 @@ 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, null);
Outcome out = new Outcome(true, List.of(), deferred, split, null);
log.info(out.summary());
return out;
}
@@ -325,6 +387,32 @@ 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). Unlike {@link #changedDeferredKeys}, this does not try to
* tell which sub-field moved: any change to {@code health:} or {@code coordinator:} 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.
*/
private 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");
}
// 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 two split keys the class doc documents.
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;
}
/**
* {@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,18 +466,8 @@ public final class MessageService {
if (candidates.size() == 1) {
Task orphan = candidates.get(0);
if (orphan.future.complete(new Reply(Outcome.REPLIED, content))) {
// 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);
if (orphan.turnId != null) {
asyncTasksByTurn.remove(orphan.turnId, orphan);
}
count(FleetMetrics.REPLIES, "path", "async-recovered");
return true; // the ticket itself took it — no inbox stranding at all
@@ -1012,36 +1002,13 @@ 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 — 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);
// 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);
}
return result;
} catch (TimeoutException e) {
@@ -1092,7 +1059,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 own finishAsyncTask(task, result) once a QUESTION
// finishAsyncTask reached via answer()'s finishAsyncTask(turnId, 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
@@ -1112,22 +1079,7 @@ public final class MessageService {
finishAsyncTask(task, result);
}
} catch (Throwable 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);
}
task.future.completeExceptionally(t);
}
});
pruneTerminalTickets();
@@ -1295,10 +1247,6 @@ 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) {
@@ -1328,49 +1276,6 @@ 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
@@ -1380,6 +1285,14 @@ 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,6 +570,155 @@ 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());
}
/**
* 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,
@@ -0,0 +1,167 @@
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 fleet} — role pools, {@code charters} and {@code tabLabel} are read live through
* the supplier on {@code CompositePeerLauncher} (class doc, Hot bullet). The one
* documented exception, {@code fleet.leaders}, is read only at startup and genuinely needs
* a restart — a real gap in {@code changedDeferredKeys}, but one the class doc already
* carries and that fleetd #330 explicitly did not re-open (its "seven readers" fact-find
* named {@code health}/{@code coordinator} as the complete set of split-shaped keys, not
* {@code fleet}). Reported as a caveat, not fixed here.</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>
*/
private static final Set<String> HOT_EXCLUDED_TOP_LEVEL_KEYS =
Set.of("placement", "fleet", "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", "fleet", "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,9 +1,5 @@
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;
@@ -15,7 +11,6 @@ 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;
@@ -1046,105 +1041,6 @@ 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
@@ -1255,63 +1151,6 @@ 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