Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 554395b104 |
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user