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
5 changed files with 296 additions and 452 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 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>
@@ -71,63 +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) — 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 #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>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
@@ -137,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");
private final Path path;
private final AtomicReference<FleetConfig> current;
@@ -182,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. */
@@ -224,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";
}
}
@@ -276,9 +215,8 @@ 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;
}
@@ -387,32 +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). 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,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,155 +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());
}
/**
* 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,167 +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 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,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