Compare commits

..

1 Commits

Author SHA1 Message Date
Dai Ha e5cb51a90e #324: read task.turnId once in finishAsyncTask to stop an NPE from ask()'s unlocked forgetting
CI / contract (pull_request) Successful in 1m28s
CI / build (pull_request) Successful in 2m1s
answer() holds sessionLocks while finishAsyncTask reads the volatile Task.turnId twice — once to
check it is non-null, once as the ConcurrentHashMap.remove key. ask()'s own timeout path mutates
the same field with no lock, via clearAsyncQuestion(turnId, true). volatile makes each read fresh
but not the pair atomic, so the field can go null between the two reads and remove(null, task)
throws NullPointerException on the lead's own answer() call, even though the reply already
completed on the line above.

Capture task.turnId into a local once and use that for both the check and the removal.

Added a package-private test seam (finishAsyncTaskRaceHook + forgetTurnForTest) so a test can force
the exact interleaving deterministically, by running the identical clearAsyncQuestion(turnId, true)
cleanup ask() uses, at the point between finishAsyncTask's former two reads. Both are inert (null)
in production.
2026-09-04 14:38:52 +07:00
5 changed files with 108 additions and 418 deletions
@@ -39,28 +39,12 @@ import java.util.function.Supplier;
* keeps the old value until a restart: {@code lifecycle:}, {@code leadHeartbeat:},
* {@code spawnReadyTimeoutMs} / {@code spawnReadyPollMs}, {@code quarantineCooldownSeconds}
* (CB-578 stage B — baked once into the {@code BackendQuarantine} built at startup),
* {@code guard:}, {@code worktreeRoot:} and {@code worktreeGroup:} (both baked once into the
* {@code GitWorktrees} built at {@code Fleetd.java:251} and never rebuilt — fleetd #323
* instance 2 found {@code worktreeGroup} missing from this list and from
* {@link #changedDeferredKeys}), {@code primary:} (fleetd #326 — {@code Fleetd.java:506, 519,
* 520} read {@code cfg.primary()} only off the startup snapshot to build {@code
* PrimaryRegistry} and size {@code ReplyPushLoop}'s reminder cap/backoff, and neither is
* rebuilt on reload; a lead whose pinned terminal changed under a running daemon stays
* unresolved as primary until a restart), {@code configReload:} (fleetd #326 — {@code
* Fleetd.java:679-680} read it only at startup to decide whether to build a {@code
* ConfigWatcher} at all and with what interval; the watcher that would apply a later change is
* itself built once, so a running watcher keeps polling on its original enabled flag and
* interval regardless of what a reload changes it to, the same shape as {@code lifecycle} —
* not cold, because no already-open resource goes inconsistent with the new value, the watcher
* (if any) simply keeps its old settings), adding or removing a profile (a new backend needs its own launcher,
* {@code guard:}, {@code worktreeRoot:}, adding or removing a profile (a new backend needs its own launcher,
* which is constructed once), <em>and an existing profile's launch settings</em> —
* {@code model}, {@code baseUrl}, {@code argv}, {@code env}, {@code mcpUrl},
* {@code exhaustedPattern} (CB-578 stage A — compiled once into {@code Fleetd.main}'s
* pattern map at startup), {@code errorPattern} (fleetd #201 Unit 5 — compiled once into
* {@code Fleetd.main}'s backend-error pattern map at startup, the same way),
* {@code ideProjectDir} / {@code ideOpenCommand} / {@code autoCompactWindow} (fleetd #323
* instance 1 — all three are read at spawn off the same frozen profile map and were missing
* from {@link #sameLaunchSettings}), and the rest of {@link #sameLaunchSettings}.
* {@code Fleetd.main}'s backend-error pattern map at startup, the same way), and the rest.
* {@code credentialId} (CB-578 stage B) is NOT on
* this list — it is read live off the config supplier at every quarantine check and
* exhaustion event, exactly like {@code weight} / {@code maxLoad}, so it is hot instead.
@@ -237,27 +221,6 @@ public final class ConfigRef implements Supplier<FleetConfig> {
if (!Objects.equals(old.worktreeRoot(), fresh.worktreeRoot())) {
changed.add("worktreeRoot");
}
// Baked into the same GitWorktrees as worktreeRoot (Fleetd.java:251) and never rebuilt
// either — see the class doc. Missing this check was fleetd #323 instance 2: a reload
// that changed only worktreeGroup reported "config reloaded" with nothing deferred, and
// newly provisioned worktrees kept the old sharing behaviour.
if (!Objects.equals(old.worktreeGroup(), fresh.worktreeGroup())) {
changed.add("worktreeGroup");
}
// fleetd #326: Fleetd.java:506, 519, 520 read cfg.primary() only off the startup snapshot
// (PrimaryRegistry's pinned terminal, ReplyPushLoop's reminder cap and backoff) — neither is
// rebuilt on reload, so a changed pin needs a restart before a lead resolves as primary again.
if (!Objects.equals(old.primary(), fresh.primary())) {
changed.add("primary");
}
// fleetd #326: Fleetd.java:679-680 read cfg.configReload() only at startup to decide whether
// to build a ConfigWatcher at all and with what interval — the watcher that would apply a
// later change is itself built once, so a running watcher keeps its original enabled flag and
// interval regardless of what a reload changes it to. Not cold: no already-open resource goes
// inconsistent with the new value, a watcher (if any) simply keeps polling on the old settings.
if (!Objects.equals(old.configReload(), fresh.configReload())) {
changed.add("configReload");
}
if (!Objects.equals(old.spawnReadyTimeoutMs(), fresh.spawnReadyTimeoutMs())
|| !Objects.equals(old.spawnReadyPollMs(), fresh.spawnReadyPollMs())) {
changed.add("spawnReady*");
@@ -302,34 +265,14 @@ public final class ConfigRef implements Supplier<FleetConfig> {
}
/**
* {@link FleetConfig.Profile} record components deliberately left out of
* {@link #sameLaunchSettings} because they are read <em>live</em>, not baked in at spawn — see
* the class doc's <em>Hot</em> bullet. {@code weight} and {@code maxLoad} are read live by the
* placement policy on every spawn; {@code credentialId} is read live by
* {@code CompositePeerLauncher} and the CB-578 stage B exhaustion sink. Nothing else is
* excluded — see {@code sameLaunchSettingsComparesEveryProfileComponentOrExcludesIt} in
* {@code ConfigRefProfileCoverageTest}, which enumerates every {@code Profile} record component
* by reflection and fails the build if one is neither compared below nor named here.
* Whether two versions of a profile would launch a peer identically. Compares every component
* the launcher reads at spawn; {@code weight}, {@code maxLoad} and {@code credentialId} are
* excluded because those are read live (by the placement policy and, for credentialId, by
* {@code CompositePeerLauncher}/the CB-578 stage B exhaustion sink) and really do take effect on
* the next spawn.
*/
static final Set<String> LAUNCH_SETTINGS_EXCLUDED = Set.of("weight", "maxLoad", "credentialId");
/**
* Whether two versions of a profile would launch a peer identically.
*
* <p>This must compare every {@link FleetConfig.Profile} record component except the three in
* {@link #LAUNCH_SETTINGS_EXCLUDED}. That is not a claim this javadoc can make good on by
* itself — a javadoc saying "compares every component" is exactly what fleetd #323 found to be
* false for three fields (and a sibling method's field list, for a fourth). The actual
* guarantee comes from {@code ConfigRefProfileCoverageTest}: it enumerates every record
* component of {@code FleetConfig.Profile} by reflection, mutates each one not in
* {@code LAUNCH_SETTINGS_EXCLUDED} on a base profile, and asserts this method reports a
* difference — so a new component that is neither compared here nor added to
* {@code LAUNCH_SETTINGS_EXCLUDED} (with a reason) fails that test by name, rather than
* silently reporting "config reloaded" for a value the daemon never picked up.
*/
static boolean sameLaunchSettings(FleetConfig.Profile a, FleetConfig.Profile b) {
return Objects.equals(a.profile(), b.profile())
&& Objects.equals(a.baseUrl(), b.baseUrl())
private static boolean sameLaunchSettings(FleetConfig.Profile a, FleetConfig.Profile b) {
return Objects.equals(a.baseUrl(), b.baseUrl())
&& Objects.equals(a.model(), b.model())
&& Objects.equals(a.configDir(), b.configDir())
&& Objects.equals(a.tokenEnv(), b.tokenEnv())
@@ -355,16 +298,6 @@ public final class ConfigRef implements Supplier<FleetConfig> {
// fleetd #201 Unit 5: errorPattern is compiled once into Fleetd.main's backend-error
// pattern map at startup (see BackendErrorPatternLookup wiring), the same way
// exhaustedPattern is — a reload never re-reads it either.
&& Objects.equals(a.errorPattern(), b.errorPattern())
// fleetd #323 instance 1: ideProjectDir and ideOpenCommand are read at spawn off the
// same frozen profile map as ideMcpUrl above (ClaudeCodeLauncher.java:267/269,
// OpenCodeLauncher.java:474/480/486) and were missing from this comparison.
&& Objects.equals(a.ideProjectDir(), b.ideProjectDir())
&& Objects.equals(a.ideOpenCommand(), b.ideOpenCommand())
// fleetd #323 instance 1: autoCompactWindow is read at spawn the same way
// (ClaudeCodeLauncher.java:926, OpenCodeLauncher.java:650). Comparing it here only
// makes the reload REPORT that a restart is needed — it deliberately does not make
// autoCompactWindow take effect live, which is a separate, larger change.
&& Objects.equals(a.autoCompactWindow(), b.autoCompactWindow());
&& Objects.equals(a.errorPattern(), b.errorPattern());
}
}
@@ -1230,14 +1230,61 @@ public final class MessageService {
}
}
/** Complete and detach an async ticket after its worker's actual terminal reply. */
/**
* Complete and detach an async ticket after its worker's actual terminal reply.
*
* <p><strong>fleetd #324.</strong> {@code task.turnId} is read into {@code turnId} exactly once.
* It used to be read twice — once for the null check, once as the removal key — and {@code
* volatile} makes each of those reads individually fresh but does not make the pair atomic.
* {@link #answer} calls this while holding {@code sessionLocks} for the target; {@link #ask}'s
* own timeout path calls {@link #clearAsyncQuestion} (which nulls {@link Task#turnId}) under no
* lock at all. When that unlocked null-out landed between the two reads here, the second read saw
* {@code null} and {@code asyncTasksByTurn.remove(null, task)} threw {@code NullPointerException}
* on the lead's own {@code answer()} call — even though {@code task.future.complete(result)} on
* the line above had already run, so the answer was in fact delivered. Capturing the field once
* removes the torn read; see the ticket for why the wider asymmetry between the locked and
* unlocked sides is not fixed by this alone.
*/
private void finishAsyncTask(Task task, Reply result) {
task.future.complete(result);
if (task.turnId != null) {
asyncTasksByTurn.remove(task.turnId, task);
String turnId = task.turnId;
if (turnId != null) {
if (finishAsyncTaskRaceHook != null) {
// Test-only (fleetd #324): see the field's own javadoc.
finishAsyncTaskRaceHook.run();
}
asyncTasksByTurn.remove(turnId, task);
}
}
/**
* Null in production; test seam for fleetd #324 — invoked from {@link #finishAsyncTask(Task,
* Reply)} right after {@code task.turnId}'s null-check passes and before the (now-local) value is
* used for the removal. A test installs this to force, deterministically, the exact interleaving
* that a real race between this method and {@link #ask}'s unlocked timeout cleanup can otherwise
* only produce by chance: firing it here reproduces "the field went null between the check and the
* use" against the pre-fix code, and demonstrates the fix tolerates it (the captured local is used
* unconditionally, so a hook that nulls the field afterward cannot affect this call).
*/
private volatile Runnable finishAsyncTaskRaceHook;
/**
* Test-only (fleetd #324): install {@link #finishAsyncTaskRaceHook}. Package-private so the test,
* in the same package, can reach it without widening any production API.
*/
void setFinishAsyncTaskRaceHookForTest(Runnable hook) {
this.finishAsyncTaskRaceHook = 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
* can reproduce that specific mutation instead of hand-rolling an approximation of it.
*/
void forgetTurnForTest(String turnId) {
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);
@@ -1,202 +0,0 @@
package dev.ltms.fleet.config;
import org.junit.jupiter.api.Test;
import java.lang.reflect.Constructor;
import java.lang.reflect.RecordComponent;
import java.util.Arrays;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.TreeSet;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #323: {@code ConfigRef.sameLaunchSettings} javadoc used to claim it "compares every
* component the launcher reads at spawn". It did not — {@code ideProjectDir}, {@code
* ideOpenCommand} and {@code autoCompactWindow} were all baked in at daemon startup (see
* {@code ClaudeCodeLauncher}/{@code OpenCodeLauncher}) and missing from the comparison, so a reload
* that changed only one of them reported "config reloaded" and the running daemon kept the old
* value.
*
* <p>This class is the mechanism the issue asked for: it enumerates every record component of
* {@link FleetConfig.Profile} by reflection and proves — by actually mutating a base profile one
* field at a time and calling the real method — that each component is either compared by
* {@link ConfigRef#sameLaunchSettings} or named in {@link ConfigRef#LAUNCH_SETTINGS_EXCLUDED} with
* a reason. A new profile field that is neither fails this test by name, not a hand-maintained list
* going stale.
*/
class ConfigRefProfileCoverageTest {
private static final RecordComponent[] COMPONENTS = FleetConfig.Profile.class.getRecordComponents();
/**
* One valid, non-blank value per record component — "the a value". None of these trip any
* defaulting/normalization in {@code Profile}'s compact constructor (see {@code
* FleetConfig.java}), so what goes in is what {@code sameLaunchSettings} sees back out.
*/
private static final Map<String, Object> BASE = baseValues();
/** The same shape, each value distinct from {@link #BASE} — "the b value". */
private static final Map<String, Object> ALT = altValues();
private static Map<String, Object> baseValues() {
Map<String, Object> v = new LinkedHashMap<>();
v.put("profile", "sonnet");
v.put("baseUrl", "http://gx00.gw:8000");
v.put("model", "sonnet");
v.put("configDir", "/config/a");
v.put("tokenEnv", "TOKEN_A");
v.put("argv", List.of("claude", "--flag-a"));
v.put("placement", "tab");
v.put("workspace", "workspace-a");
v.put("tabLabel", "label-a");
v.put("mcpUrl", "http://mcp-a");
v.put("cwd", "/cwd/a");
v.put("parityOverlay", List.of(".env", ".env.a"));
v.put("gitTokenEnv", "GIT_TOKEN_A");
v.put("gitHostEnv", "GITEA_HOST_A");
v.put("kind", "claude-code");
v.put("env", Map.of("K", "A"));
v.put("weight", 1.0f);
v.put("maxLoad", 5);
v.put("subscription", Boolean.TRUE);
v.put("exhaustedPattern", "usage limit a");
v.put("credentialId", "cred-a");
v.put("ideMcpUrl", "http://ide-mcp-a");
v.put("ideProjectDir", "modules/a");
v.put("ideOpenCommand", "open-cmd-a {dir}");
v.put("autoCompactWindow", 150000);
v.put("errorPattern", "error a");
assertNamesMatchComponents(v);
return v;
}
private static Map<String, Object> altValues() {
Map<String, Object> v = new LinkedHashMap<>();
v.put("profile", "sonnet-b");
v.put("baseUrl", "http://gx01.gw:8000");
v.put("model", "haiku");
v.put("configDir", "/config/b");
v.put("tokenEnv", "TOKEN_B");
v.put("argv", List.of("claude", "--flag-b"));
v.put("placement", "weighted");
v.put("workspace", "workspace-b");
v.put("tabLabel", "label-b");
v.put("mcpUrl", "http://mcp-b");
v.put("cwd", "/cwd/b");
v.put("parityOverlay", List.of(".env", ".env.b"));
v.put("gitTokenEnv", "GIT_TOKEN_B");
v.put("gitHostEnv", "GITEA_HOST_B");
v.put("kind", "opencode");
v.put("env", Map.of("K", "B"));
v.put("weight", 2.0f);
v.put("maxLoad", 9);
v.put("subscription", Boolean.FALSE);
v.put("exhaustedPattern", "usage limit b");
v.put("credentialId", "cred-b");
v.put("ideMcpUrl", "http://ide-mcp-b");
v.put("ideProjectDir", "modules/b");
v.put("ideOpenCommand", "open-cmd-b {dir}");
v.put("autoCompactWindow", 250000);
v.put("errorPattern", "error b");
assertNamesMatchComponents(v);
return v;
}
private static void assertNamesMatchComponents(Map<String, Object> values) {
Set<String> componentNames = new TreeSet<>();
for (RecordComponent rc : COMPONENTS) {
componentNames.add(rc.getName());
}
assertEquals(componentNames, new TreeSet<>(values.keySet()),
"this test's value map has drifted from FleetConfig.Profile's actual components — "
+ "update BASE/ALT alongside the record");
}
private static FleetConfig.Profile profileOf(Map<String, Object> values) throws ReflectiveOperationException {
Class<?>[] types = Arrays.stream(COMPONENTS).map(RecordComponent::getType).toArray(Class<?>[]::new);
Object[] args = Arrays.stream(COMPONENTS)
.map(rc -> values.get(rc.getName()))
.toArray();
Constructor<FleetConfig.Profile> ctor = FleetConfig.Profile.class.getDeclaredConstructor(types);
return ctor.newInstance(args);
}
/** {@code BASE} with exactly one named component swapped for its {@code ALT} value. */
private static FleetConfig.Profile mutate(String componentName) throws ReflectiveOperationException {
Map<String, Object> values = new LinkedHashMap<>(BASE);
values.put(componentName, ALT.get(componentName));
return profileOf(values);
}
/**
* The mechanism fleetd #323 asked for: enumerate {@link FleetConfig.Profile}'s record
* components, mutate each non-excluded one, and prove {@code sameLaunchSettings} actually
* notices — not just that some hand-maintained list claims it does. Prints the denominator
* (total / compared / excluded) the issue required: a checker that cannot state its own
* denominator is the failure this repo keeps hitting.
*/
@Test
void sameLaunchSettingsComparesEveryProfileComponentOrExcludesIt() throws ReflectiveOperationException {
int total = COMPONENTS.length;
Set<String> excluded = ConfigRef.LAUNCH_SETTINGS_EXCLUDED;
Set<String> allNames = new TreeSet<>();
for (RecordComponent rc : COMPONENTS) {
allNames.add(rc.getName());
}
assertTrue(allNames.containsAll(excluded),
"ConfigRef.LAUNCH_SETTINGS_EXCLUDED names a component that does not exist on "
+ "FleetConfig.Profile — check for a typo: " + excluded);
// The exclusion set is this mechanism's own escape hatch, so it has to be pinned too.
// Found by mutation while verifying fleetd #323: moving autoCompactWindow and
// ideOpenCommand OUT of the comparison and INTO the exclusion set left the whole suite
// green — the loop below simply skips them, and the denominator assertion still balances.
// That is exactly the lazy move a failing coverage test invites, and it silently restores
// the #323 bug. Only ideProjectDir and worktreeGroup were saved by a behavioural test in
// ConfigRefTest; the other two had none. So: growing this set now requires editing this
// line as well, which is a visible, deliberate diff rather than a quiet one.
assertEquals(Set.of("weight", "maxLoad", "credentialId"), excluded,
"ConfigRef.LAUNCH_SETTINGS_EXCLUDED changed. A component belongs in it ONLY if it "
+ "is read live off the config supplier, not baked into a launcher at "
+ "startup. If you are adding one to silence this test, that is fleetd #323 "
+ "happening again: compare it in sameLaunchSettings instead. If it really "
+ "is read live, name where it is read and update this assertion.");
FleetConfig.Profile base = profileOf(BASE);
List<String> uncovered = new java.util.ArrayList<>();
int compared = 0;
for (RecordComponent rc : COMPONENTS) {
String name = rc.getName();
if (excluded.contains(name)) {
continue;
}
FleetConfig.Profile mutated = mutate(name);
if (ConfigRef.sameLaunchSettings(base, mutated)) {
uncovered.add(name);
} else {
compared++;
}
}
System.out.printf(
"ConfigRef.sameLaunchSettings coverage — %d Profile components total, %d compared, "
+ "%d excluded (%s)%n",
total, compared, excluded.size(), excluded);
assertEquals(List.of(), uncovered,
"these FleetConfig.Profile components changed but ConfigRef.sameLaunchSettings "
+ "reported no difference — add each one to the comparison (it is read at "
+ "spawn and baked in until a restart) or to ConfigRef.LAUNCH_SETTINGS_EXCLUDED "
+ "with a reason it is genuinely read live: " + uncovered);
assertEquals(total, compared + excluded.size(),
"every FleetConfig.Profile record component must be either compared or excluded — "
+ total + " components, " + compared + " compared, " + excluded.size()
+ " excluded");
}
}
@@ -434,142 +434,6 @@ class ConfigRefTest {
assertEquals("provider 5xx", ref.get().profiles().get("sonnet").errorPattern());
}
/**
* fleetd #323 instance 1: {@code ideProjectDir} is read at spawn off the frozen profile map
* (see {@code ClaudeCodeLauncher}/{@code OpenCodeLauncher}) exactly like {@code model}, but was
* missing from {@code sameLaunchSettings} — a reload changing only this field used to report a
* bare "config reloaded" and the running daemon kept launching with the old value.
*/
@Test
void changingAProfilesIdeProjectDirIsReportedAsDeferred(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
bind:
host: 127.0.0.1
port: 8765
herdrSocket: ~/.config/herdr/herdr.sock
profiles:
sonnet:
baseUrl: http://gx00.gw:8000
model: sonnet
ideProjectDir: fleetd
guard:
offSubscriptionHosts:
- gx00.gw
""");
ConfigRef ref = refFor(f);
Files.writeString(f, """
bind:
host: 127.0.0.1
port: 8765
herdrSocket: ~/.config/herdr/herdr.sock
profiles:
sonnet:
baseUrl: http://gx00.gw:8000
model: sonnet
ideProjectDir: fleetd-renamed
guard:
offSubscriptionHosts:
- gx00.gw
""");
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied());
assertEquals(1, out.deferred().size(), out.deferred().toString());
assertTrue(out.deferred().getFirst().contains("sonnet"), out.deferred().toString());
assertTrue(out.deferred().getFirst().contains("launch settings"), out.deferred().toString());
// The snapshot still carries the new value — a restart is what makes it take effect.
assertEquals("fleetd-renamed", ref.get().profiles().get("sonnet").ideProjectDir());
}
/**
* fleetd #323 instance 2: {@code worktreeGroup} is baked into the same {@code GitWorktrees}
* as {@code worktreeRoot} (Fleetd.java:251) and never rebuilt, but only {@code worktreeRoot}
* was on {@code changedDeferredKeys} — a reload changing only the group reported a bare
* "config reloaded" and newly provisioned worktrees kept the old sharing behaviour.
*/
@Test
void changingWorktreeGroupIsReportedAsDeferred(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml("worktreeGroup: devgroup\n"));
ConfigRef ref = refFor(f);
Files.writeString(f, yaml("worktreeGroup: devgroup2\n"));
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied());
assertEquals(java.util.List.of("worktreeGroup"), out.deferred());
assertTrue(out.summary().contains("needs a restart") || out.summary().contains("need a restart"),
out.summary());
// The snapshot still carries the new value — a restart is what makes it take effect.
assertEquals("devgroup2", ref.get().worktreeGroup());
}
/**
* fleetd #326: {@code primary} is read only off the startup snapshot — {@code Fleetd.java:506,
* 519, 520} feed {@code PrimaryRegistry} and {@code ReplyPushLoop} at construction and neither is
* rebuilt on reload — but it was missing from {@link ConfigRef#changedDeferredKeys}, so a reload
* that only changed the pinned primary terminal reported a bare "config reloaded" while a lead
* whose tab no longer matched stayed demoted to worker.
*/
@Test
void changingPrimaryIsReportedAsDeferred(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml("""
primary:
terminal: term-a
"""));
ConfigRef ref = refFor(f);
Files.writeString(f, yaml("""
primary:
terminal: term-b
"""));
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied());
assertEquals(java.util.List.of("primary"), out.deferred());
assertTrue(out.summary().contains("need") && out.summary().contains("restart"), out.summary());
// The snapshot still carries the new value — a restart is what makes it take effect.
assertEquals("term-b", ref.get().primary().terminal());
}
/**
* fleetd #326: {@code configReload} itself is read only at startup ({@code Fleetd.java:679-680})
* to decide whether to build a {@code ConfigWatcher} at all, and with what interval — the watcher
* that would apply a later change is itself built once, so it is deferred rather than cold (see
* {@link ConfigRef}'s class doc: cold means an already-open resource would go inconsistent with
* the new value, and there is no such resource here — a running watcher just keeps polling on its
* original enabled/interval until a restart, exactly like {@code lifecycle} or {@code guard}).
* Before this fix, turning reload off (or changing its interval) through a reload reported a bare
* "config reloaded" — the obvious joke the issue names.
*/
@Test
void changingConfigReloadIsReportedAsDeferred(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml("""
configReload:
enabled: true
intervalSeconds: 10
"""));
ConfigRef ref = refFor(f);
Files.writeString(f, yaml("""
configReload:
enabled: false
intervalSeconds: 30
"""));
ConfigRef.Outcome out = ref.reload();
assertTrue(out.applied());
assertEquals(java.util.List.of("configReload"), out.deferred());
assertTrue(out.summary().contains("need") && out.summary().contains("restart"), out.summary());
// The snapshot still carries the new value — a restart is what makes it take effect.
assertFalse(ref.get().configReload().isEnabled());
assertEquals(30, ref.get().configReload().intervalSeconds());
}
@Test
void aFixedRefHasNoFileAndRefusesToReload() {
FleetConfig cfg = new FleetConfig(null, null, null, null, null, null,
@@ -940,6 +940,54 @@ class MessageServiceTest {
assertEquals("PR opened: https://example/pulls/42", view.reply());
}
/**
* fleetd #324: {@code answer()} holds {@code sessionLocks} for the target and, once the worker's
* real terminal reply arrives, calls {@code finishAsyncTask}, which used to read the volatile
* {@code task.turnId} twice — once to check it is non-null, once as the key for
* {@code asyncTasksByTurn.remove}. {@code ask()}'s own timeout path mutates the same field with no
* lock at all. This test does not wait for a real race to land in that narrow window between the
* two reads — instead it drives the exact sequence the ticket describes (worker asks, primary
* answers, worker's real reply arrives) and, via a package-private test hook wired to fire at
* precisely that point, runs the identical production cleanup {@code ask()}'s timeout catch block
* runs ({@code clearAsyncQuestion(turnId, true)}) so the field goes {@code null} between the two
* reads deterministically rather than by chance.
*
* <p>What this proves: given that exact interleaving, {@code answer()} must not throw and the
* ticket must still resolve to the worker's real reply. What it does not prove: that the
* interleaving itself is reachable in production — that is established by reading the code (see
* the ticket), not by this test, since forcing it via a hook is not the same as two independent
* threads racing on their own schedules.
*/
@Test
void finishAsyncTaskSurvivesTurnIdGoingNullBetweenItsTwoReads() 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();
// Fire ask()'s own unlocked timeout cleanup at the moment finishAsyncTask has already checked
// task.turnId is non-null but has not yet used it — the exact torn-read window fleetd #324
// describes.
messages.setFinishAsyncTaskRaceHookForTest(() -> messages.forgetTurnForTest(turnId));
CompletableFuture<MessageService.Reply> answer =
CompletableFuture.supplyAsync(() -> messages.answer(turnId, "config.yaml", 5000));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
awaitWaiting(); // answer() opened its own forward waiter for the resumed worker turn
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome(),
"the lead's own answer() call must not throw because ask()'s timeout cleanup raced it");
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");
}
@Test
void unansweredAsyncQuestionReturnsTheTicketToPendingAndReleasesItsTarget() throws Exception {
String ticket = messages.sendAsync(T, "task that asks");