Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 8d5bc3ee89 |
@@ -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).
|
||||
|
||||
@@ -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.");
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user