Compare commits

..

1 Commits

Author SHA1 Message Date
Dai Ha 8d5bc3ee89 fleetd #416: fleet_list must enumerate the STARTUP profile set
CI / contract (pull_request) Successful in 1m29s
CI / build (pull_request) Successful in 1m52s
`profiles` is a DEFERRED key: HerdrPeerLauncher takes Map.copyOf(profiles) once at
construction, so a profile added only to the hot-reloaded map can never be spawned.
CapacitySource was built with `() -> config.get().profiles().keySet()` — the live map —
so fleet_list reported a hot-added profile as free while fleet_spawn on that same
profile failed with "unknown worker profile". fleet_list's contract for `free` says it
runs "the same check the spawn gate itself runs"; it did not.

The set now comes from the startup snapshot, the same shape as the coordinator.peers
wiring three lines below, whose comment already stated the rule. maxLoad stays live on
purpose — ConfigRef documents it as hot, like credentialId and weight — so a maxLoad
edit still takes effect without a restart.

Found by the fleet01 lead, who proved it with a live ghost profile rather than an
argument. Implemented by a sonnet member; its reply was lost to an empty scrape, so the
work was salvaged uncommitted from its worktree and both proof steps were run by the
lead instead:

  full build                          Tests run: 1505, Failures: 0, 0 compile errors
  mutation A: live keySet restored    liveOnlyProfileIsNotListed FAILS (alone)
  mutation B: permanently empty set   startupProfileIsListed FAILS (alone)

Mutation B is the point of the second direction: per #404, a test that only ever checks
the absent case cannot tell a correct lookup from one that returns nothing at all.
2026-09-10 11:25:56 +07:00
4 changed files with 185 additions and 50 deletions
@@ -670,11 +670,8 @@ public final class Fleetd {
}, outagePolicy);
FleetMcp mcp = new FleetMcp(messages, workers, sessions, identity, presence,
primaryRegistry, callers, metrics, new FleetMcp.CapacitySource(profile -> liveCountRef.get().apply(profile),
profile -> {
var configured = config.get().profiles().get(profile);
return configured == null ? null : configured.maxLoad();
}, () -> config.get().profiles().keySet(), System::nanoTime),
primaryRegistry, callers, metrics,
capacitySource(config, cfg, profile -> liveCountRef.get().apply(profile)),
new FleetMcp.HealthCoverageSource(() -> {
var health = config.get().health();
return FleetHealthMonitor.coverage(health != null && health.isEnabled(),
@@ -810,6 +807,33 @@ public final class Fleetd {
}, quarantine, profile -> startupExhaustedPatterns.containsKey(profile));
}
/**
* fleetd #416: production source for {@code fleet_list}'s per-profile capacity facts.
*
* <p>The profile <em>set</em> ({@code configuredProfiles}) must come from {@code cfg} — the
* startup snapshot — not the live {@code config.get()}. {@code profiles} as a whole is a
* {@code DEFERRED} key ({@link ConfigRef#DEFERRED_KEYS}): {@code HerdrPeerLauncher} takes
* {@code Map.copyOf(profiles)} once at construction and a profile only added to the
* hot-reloaded map can never actually be spawned, so enumerating it live made {@code fleet_list}
* report a profile as available when {@code fleet_spawn} on that same profile fails with
* {@code unknown worker profile}. {@code fleet_list}'s own contract for {@code free} is "the
* same check the spawn gate itself runs" — the set the spawn gate can see is the startup one,
* so this must enumerate that one too, the same shape as {@code coordinator.peers} above.
*
* <p>{@code maxLoad} stays live on purpose: it is read off {@code config.get()} exactly like
* {@code credentialId} ({@link ConfigRef} documents both as hot), so an existing profile's
* {@code maxLoad} edit must still change what {@code fleet_list} reports without a restart.
*/
static FleetMcp.CapacitySource capacitySource(ConfigRef config, FleetConfig cfg,
Function<String, Integer> liveCount) {
return new FleetMcp.CapacitySource(liveCount,
profile -> {
var configured = config.get().profiles().get(profile);
return configured == null ? null : configured.maxLoad();
},
cfg.profiles()::keySet, System::nanoTime);
}
/**
* fleetd #248: package-private factory for the member worktree/branch lookup {@link
* CompletionResolver} uses to name a fallback report's worktree and branch (fleetd#241).
@@ -230,26 +230,6 @@ public final class ReplyPushLoop {
.collect(Collectors.toUnmodifiableSet());
}
/**
* Test seam only (fleetd #418): carries no production behaviour, and nothing in this class
* calls it. Exposes {@link #pendingQuestionTurnIdsFor} — the exact state {@link #decide} reads
* to decide whether a question keeps a lead's schedule alive.
*
* <p>{@code MessageService.ask()} does three things in order before a question is fully open to
* this loop: it flips the ticket's {@code poll()} phase to {@code Phase.ASKING}, then resolves
* the reverse-rendezvous waiter, then calls {@link #onQuestionOpened}, which is what actually
* populates {@link #pendingQuestions}. A test that barriers on {@code Phase.ASKING} observes only
* the first of those three steps — under load the asker thread can be descheduled between steps
* one and three, so the barrier releases before this method's underlying map is populated, and
* {@link #decide} correctly reports nothing pending yet. A test that must order itself after the
* state {@link #decide} actually reads waits on this instead of on the phase.
*
* @return an unmodifiable snapshot; empty for a lead with no open questions
*/
Set<String> pendingQuestionTurnIdsForTest(String lead) {
return pendingQuestionTurnIdsFor(lead);
}
private List<PendingIncident> pendingIncidentsFor(String lead) {
return pendingIncidents.values().stream().filter(i -> lead.equals(i.key().lead())).toList();
}
@@ -0,0 +1,155 @@
package dev.ltms.fleet;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.mcp.FleetMcp;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import java.nio.file.Files;
import java.nio.file.Path;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.junit.jupiter.api.Assertions.assertEquals;
/**
* fleetd #416: {@code fleet_list}'s {@code CapacitySource.configuredProfiles} must enumerate the
* <em>startup</em> profile set, not the live, hot-reloaded one.
*
* <p>{@code profiles} as a whole is a {@code DEFERRED} key ({@link ConfigRef#DEFERRED_KEYS}):
* {@code HerdrPeerLauncher} takes {@code Map.copyOf(profiles)} once at construction, so a profile
* only added to the hot-reloaded map can never actually be spawned. Before this fix, {@code
* Fleetd.main} built {@code CapacitySource} with {@code () -> config.get().profiles().keySet()} —
* the live map — so {@code fleet_list} would report a freshly hot-reloaded profile as available
* ({@code free > 0}) while {@code fleet_spawn} on that same profile failed with
* {@code unknown worker profile}. Measured on another host: adding a throwaway profile and letting
* it hot-reload gave {@code fleet_list} -> {@code free: 3} and {@code fleet_spawn} ->
* {@code error: unknown worker profile}.
*
* <p>This test needs a reload, the same reason {@link FleetdExhaustionDetectionArmedWiringTest}
* does: at startup the two snapshots agree, so a test of only a newly started daemon would not
* detect a live {@code config.get()} lookup for the set.
*
* <p><b>maxLoad must stay hot.</b> It is read off {@code config.get()} exactly like
* {@code credentialId} ({@link ConfigRef} documents both as hot, "read live off the config
* supplier ... exactly like weight/maxLoad"), so a reload that only changes an existing profile's
* {@code maxLoad} — no add/remove — must still change what {@code fleet_list} reports without a
* restart. A fix that freezes the whole {@code CapacitySource} against {@code cfg} (rather than
* only its {@code configuredProfiles} set) would trade this bug for its mirror image and is pinned
* wrong by {@link #reloadedMaxLoadStillChangesWhatFleetListReports}.
*/
class FleetdCapacitySourceWiringTest {
private static final String STARTUP = """
bind:
host: 127.0.0.1
port: 8765
herdrSocket: ~/.config/herdr/herdr.sock
profiles:
terra:
baseUrl: http://gx00.gw:8000
model: terra
maxLoad: 3
guard:
offSubscriptionHosts:
- gx00.gw
""";
private static final String WITH_NEW_PROFILE = """
bind:
host: 127.0.0.1
port: 8765
herdrSocket: ~/.config/herdr/herdr.sock
profiles:
terra:
baseUrl: http://gx00.gw:8000
model: terra
maxLoad: 3
ghost404:
baseUrl: http://gx00.gw:8001
model: ghost404
maxLoad: 3
guard:
offSubscriptionHosts:
- gx00.gw
""";
private static final String WITH_CHANGED_MAX_LOAD = """
bind:
host: 127.0.0.1
port: 8765
herdrSocket: ~/.config/herdr/herdr.sock
profiles:
terra:
baseUrl: http://gx00.gw:8000
model: terra
maxLoad: 9
guard:
offSubscriptionHosts:
- gx00.gw
""";
@Test
@DisplayName("a profile present only in the live (hot-reloaded) config is NOT listed")
void liveOnlyProfileIsNotListed(@TempDir Path dir) throws Exception {
Path file = dir.resolve("fleetd.yaml");
Files.writeString(file, STARTUP);
FleetConfig cfg = FleetConfig.load(file);
ConfigRef config = new ConfigRef(file, cfg);
Files.writeString(file, WITH_NEW_PROFILE);
assertTrue(config.reload().applied());
// The live snapshot now has the new profile — proves the reload really happened and this
// test is not accidentally passing because nothing changed.
assertTrue(config.get().profiles().containsKey("ghost404"));
FleetMcp.CapacitySource source = Fleetd.capacitySource(config, cfg, _ -> 0);
assertFalse(source.configuredProfiles().get().contains("ghost404"),
"a profile added only to the hot-reloaded config must not be listed by fleet_list — "
+ "HerdrPeerLauncher never learns about it until a restart, so fleet_spawn on it "
+ "would fail with 'unknown worker profile' while fleet_list claimed it free");
}
@Test
@DisplayName("a profile present in the startup set IS listed")
void startupProfileIsListed(@TempDir Path dir) throws Exception {
Path file = dir.resolve("fleetd.yaml");
Files.writeString(file, STARTUP);
FleetConfig cfg = FleetConfig.load(file);
ConfigRef config = new ConfigRef(file, cfg);
FleetMcp.CapacitySource source = Fleetd.capacitySource(config, cfg, _ -> 0);
// fleetd #416, both-directions requirement: a test that only ever passes an empty/absent
// startup set (the case above) cannot tell a correct lookup from one that is permanently
// empty (e.g. a mutation replacing the supplier with Set::of). This is the direction that
// fails if the fix regresses to reporting nothing at all.
assertTrue(source.configuredProfiles().get().contains("terra"),
"a profile present in the startup snapshot must still be listed by fleet_list");
}
@Test
@DisplayName("a hot maxLoad edit still changes what fleet_list reports")
void reloadedMaxLoadStillChangesWhatFleetListReports(@TempDir Path dir) throws Exception {
Path file = dir.resolve("fleetd.yaml");
Files.writeString(file, STARTUP);
FleetConfig cfg = FleetConfig.load(file);
ConfigRef config = new ConfigRef(file, cfg);
FleetMcp.CapacitySource source = Fleetd.capacitySource(config, cfg, _ -> 0);
assertEquals(3, source.maxLoad().apply("terra"),
"sanity: maxLoad reads 3 from the startup config before any reload");
Files.writeString(file, WITH_CHANGED_MAX_LOAD);
assertTrue(config.reload().applied());
assertEquals(9, source.maxLoad().apply("terra"),
"maxLoad must stay hot — the SAME CapacitySource instance must reflect a reloaded "
+ "maxLoad without a restart, exactly like credentialId. Freezing the whole "
+ "CapacitySource against the startup snapshot (rather than only its "
+ "configuredProfiles set) would trade fleetd #416 for its mirror image.");
}
}
@@ -1802,11 +1802,7 @@ class MessageServiceTest {
Thread asker = new Thread(() -> assertThrows(IllegalStateException.class,
() -> service.ask(T, "which config file?", 30_000)));
asker.start();
// Phase.ASKING (markAsyncQuestion) is only the FIRST of ask()'s three steps; the assertion
// below depends on the THIRD (pushLoop.onQuestionOpened). Barrier on the push loop's own
// pending-question state instead of the phase — see ReplyPushLoop#pendingQuestionTurnIdsForTest.
MessageService.TaskView asking = awaitTicketPhaseOn(service, ticket, MessageService.Phase.ASKING);
awaitQuestionPendingOn(pushLoop, LEAD, asking.turnId());
awaitTicketPhaseOn(service, ticket, MessageService.Phase.ASKING);
assertEquals(ReplyPushLoop.Action.INJECT, pushLoop.decide(LEAD, 0, 0, 0),
"the open question should be the one thing keeping this lead's schedule alive");
@@ -2098,26 +2094,6 @@ class MessageServiceTest {
}
}
/**
* fleetd #418: waits until {@code turnId} actually appears in {@code pushLoop}'s own open-question
* state for {@code lead} — {@link ReplyPushLoop#pendingQuestionTurnIdsForTest} — rather than until
* {@link MessageService#poll} reports {@link MessageService.Phase#ASKING}. {@code Phase.ASKING} is
* set by {@code markAsyncQuestion}, the FIRST of three steps {@code MessageService.ask()} performs;
* {@code pushLoop.onQuestionOpened} (the one that actually publishes to {@code pendingQuestions})
* is the THIRD. Under load the asker thread can be descheduled between those two steps, so a
* barrier on the phase alone can release before the push loop has anything pending — a test that
* then asserts on {@link ReplyPushLoop#decide} is asserting on state that has not been published
* yet, not on the throwing path it is named for.
*/
private void awaitQuestionPendingOn(ReplyPushLoop pushLoop, String lead, String turnId) throws Exception {
long deadline = System.currentTimeMillis() + 3000;
while (!pushLoop.pendingQuestionTurnIdsForTest(lead).contains(turnId)) {
assertTrue(System.currentTimeMillis() < deadline,
"turnId " + turnId + " never appeared in the push loop's pending questions for " + lead);
Thread.sleep(5);
}
}
// --- CB-640: fleet health evidence accessors --------------------------------------------
@Test