Compare commits
18 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| ed2027b202 | |||
| 6b7caba248 | |||
| f8b0d42a5c | |||
| d1e7d71eee | |||
| 7fd914df1a | |||
| a196d34455 | |||
| 766772763f | |||
| eab8185d7b | |||
| 7f672f0fb8 | |||
| e2801b9bbc | |||
| ea02c7b248 | |||
| ce74e164c6 | |||
| e60f892efd | |||
| ce05886831 | |||
| be123d0ac7 | |||
| 2d09c8b027 | |||
| ed54f0224e | |||
| 8d5bc3ee89 |
@@ -69,6 +69,7 @@ import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Optional;
|
||||
import java.util.Set;
|
||||
import java.util.TreeSet;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
@@ -231,6 +232,13 @@ public final class Fleetd {
|
||||
profileName -> liveCountRef.get().apply(profileName),
|
||||
quarantine,
|
||||
outagePolicy);
|
||||
// fleetd #422 follow-up: say which of the three model-gate states the daemon booted into —
|
||||
// no models: block at all, a block armed with nothing off, or a block with N off — the same
|
||||
// way exhaustedPatternCoverageLine/errorPatternCoverageLine report CB-578 stage A/fleetd
|
||||
// #201 Unit 5 coverage just below. Read from workers.modelGateState() (never a separate
|
||||
// config.get().models() here) so this line and fleet_profiles' modelGateArmed can never
|
||||
// disagree about what CompositePeerLauncher's spawn gate actually enforces.
|
||||
log.info("model gate (fleetd #422): {}", modelGateCoverageLine(workers.modelGateState()));
|
||||
// CB-504: under supervision (launchd/systemd) fleetd can start before herdr's socket
|
||||
// exists. The client itself is lazy — it connects per call — but the orphan reap below is
|
||||
// the first thing that actually talks to herdr, so without this wait a boot-order race
|
||||
@@ -352,7 +360,11 @@ public final class Fleetd {
|
||||
// pane resolves to an architect until the later spawn lifecycle binds one. The registry is
|
||||
// what CallerResolver resolves against and what that lifecycle will read profiles from;
|
||||
// nothing here spawns a slot.
|
||||
MemberRegistry members = new MemberRegistry(cfg.fleet());
|
||||
// fleetd #424: MemberRegistry.live re-reads fleet.architects through `config` on every
|
||||
// reserve/requireSlotFor call, so a reload that removes or adds an architect slot governs
|
||||
// the next spawn with no restart — the frozen `new MemberRegistry(cfg.fleet())` this used
|
||||
// to be let a "revoked" slot keep granting new architect spawns forever.
|
||||
MemberRegistry members = MemberRegistry.live(() -> config.get().fleet());
|
||||
sessions.setMemberLifecycle(members);
|
||||
if (!members.slots().isEmpty()) {
|
||||
log.info("member slots: {} configured {} — none bound yet (a slot is idle until the "
|
||||
@@ -380,8 +392,7 @@ public final class Fleetd {
|
||||
.map(session -> exhaustedPatternsByProfile.get(session.profile()))
|
||||
.orElse(null);
|
||||
log.info("backend-exhausted classification (CB-578 stage A): {}",
|
||||
CompletionResolver.coverage("exhaustedPattern", cfg.profiles().keySet(),
|
||||
exhaustedPatternsByProfile.keySet()));
|
||||
exhaustedPatternCoverageLine(cfg.profiles().keySet(), exhaustedPatternsByProfile.keySet()));
|
||||
// fleetd #201 Unit 5: classify a completion-fallback scrape that matches a profile's
|
||||
// configured backend-error refusal (a credential outage, a provider 5xx) as a backend error
|
||||
// rather than handing it back as a real answer. Compiled once at startup, keyed by profile
|
||||
@@ -401,8 +412,7 @@ public final class Fleetd {
|
||||
BackendErrorPatternLookup backendErrorPatterns = backendErrorPatternLookup(sessions::roster,
|
||||
errorPatternsByProfile);
|
||||
log.info("backend-error classification (fleetd #201 Unit 5): {}",
|
||||
CompletionResolver.coverage("errorPattern", cfg.profiles().keySet(),
|
||||
errorPatternsByProfile.keySet()));
|
||||
errorPatternCoverageLine(cfg.profiles().keySet(), errorPatternsByProfile.keySet()));
|
||||
// CB-578 stage B: on a classification that actually wins, quarantine the exhausted profile's
|
||||
// CREDENTIAL — not the profile name — so a profile sharing that credential (e.g. two models
|
||||
// on one OpenAI account) is refused too, not just the one that happened to report it. Reads
|
||||
@@ -670,11 +680,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 +817,93 @@ public final class Fleetd {
|
||||
}, quarantine, profile -> startupExhaustedPatterns.containsKey(profile));
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #415 (review follow-up): package-private factory for the CB-578 stage A {@code
|
||||
* exhaustedPattern} startup coverage line, paired explicitly with {@link
|
||||
* CompletionResolver.UnsetMeaning#OFF} — {@code exhaustedPattern} has no fallback, so a
|
||||
* profile with none configured really does have the classification off.
|
||||
*
|
||||
* <p>Extracted out of {@code main} for the same reason {@link #capacitySource} and {@link
|
||||
* #worktreeBranchLookup} were: {@code coverage()}'s own tests ({@code CompletionResolverTest})
|
||||
* prove it words {@code OFF} and {@link CompletionResolver.UnsetMeaning#BUILT_IN_DEFAULT}
|
||||
* correctly when a test supplies the meaning itself — they cannot prove {@code main} pairs the
|
||||
* right meaning with the right key, which is the actual fleetd #415 defect. <b>Measured:</b>
|
||||
* swapping the {@code UnsetMeaning} arguments between this method and {@link
|
||||
* #errorPatternCoverageLine} — recreating #415's defect with the two keys exchanged — compiled
|
||||
* with 0 errors and left all 1506 existing tests green before {@code
|
||||
* FleetdPatternCoverageLineTest} was added to catch exactly that swap.
|
||||
*/
|
||||
static String exhaustedPatternCoverageLine(Set<String> allProfiles, Set<String> configuredProfiles) {
|
||||
return CompletionResolver.coverage("exhaustedPattern", CompletionResolver.UnsetMeaning.OFF,
|
||||
allProfiles, configuredProfiles);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #415 (review follow-up): the {@code errorPattern} counterpart of {@link
|
||||
* #exhaustedPatternCoverageLine}, paired explicitly with {@link
|
||||
* CompletionResolver.UnsetMeaning#BUILT_IN_DEFAULT} — an unset {@code errorPattern} still runs
|
||||
* backend-error classification against {@code CompletionResolver}'s built-in {@code
|
||||
* BACKEND_ERROR} pattern, so the empty case is not "off". See {@link
|
||||
* #exhaustedPatternCoverageLine}'s javadoc for the measured swap mutation this pairing guards
|
||||
* against.
|
||||
*/
|
||||
static String errorPatternCoverageLine(Set<String> allProfiles, Set<String> configuredProfiles) {
|
||||
return CompletionResolver.coverage("errorPattern", CompletionResolver.UnsetMeaning.BUILT_IN_DEFAULT,
|
||||
allProfiles, configuredProfiles);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #422 follow-up: package-private factory for the startup line reporting which of the
|
||||
* three central {@code models.allow:} gate states the daemon booted into. Extracted the same
|
||||
* way {@link #exhaustedPatternCoverageLine}/{@link #errorPatternCoverageLine} are, so a
|
||||
* dedicated test can call it directly rather than parsing log output, and so {@code main}'s
|
||||
* only source for this line is {@link PeerLauncher#modelGateState()} — never a second,
|
||||
* independently-derived read of {@code cfg.models()} that could disagree with what {@code
|
||||
* CompositePeerLauncher}'s spawn gate actually enforces (the fleetd #404 lesson).
|
||||
*
|
||||
* <p>Unlike the two pattern-key lines above, there is no {@code UnsetMeaning} choice to make
|
||||
* here: {@link PeerLauncher.ModelGateState#configured()} already states, unambiguously, whether
|
||||
* an empty {@link PeerLauncher.ModelGateState#off()} means "no {@code models:} block to gate
|
||||
* with" or "a block armed and currently reporting zero off" — the exact two states a bare
|
||||
* {@code disabledModels()} read could not tell apart before this ticket.
|
||||
*/
|
||||
static String modelGateCoverageLine(PeerLauncher.ModelGateState state) {
|
||||
if (!state.configured()) {
|
||||
return "not configured (no models: block — nothing is gated, and nothing can be)";
|
||||
}
|
||||
Set<String> off = state.off();
|
||||
return off.isEmpty()
|
||||
? "armed (models: block present; 0 models currently turned off)"
|
||||
: "armed (" + off.size() + " model(s) turned off: " + new TreeSet<>(off) + ")";
|
||||
}
|
||||
|
||||
/**
|
||||
* 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).
|
||||
|
||||
@@ -11,6 +11,7 @@ import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
import java.util.function.Supplier;
|
||||
|
||||
/**
|
||||
* The architect-slot registry (CB-548): every gateway-local architect name and the strong-model
|
||||
@@ -19,16 +20,43 @@ import java.util.Objects;
|
||||
*
|
||||
* <p>Two halves, split by who owns each:
|
||||
* <ul>
|
||||
* <li><b>slots</b> — configured once, keyed by the gateway-local unique name; each carries the
|
||||
* {@code profile} reference the spawn lifecycle reads when it stands the slot up. A read-only
|
||||
* snapshot taken at construction.</li>
|
||||
* <li><b>terminal bindings</b> — owned by this registry and initially <em>empty</em>. Config
|
||||
* declares no architect terminal, so at startup every slot is idle and nothing resolves to an
|
||||
* architect; a session only becomes one when the spawn lifecycle {@linkplain #bind(String,
|
||||
* String) binds} its terminal to a slot. {@link CallerResolver} reads this through
|
||||
* {@link #snapshot()} to turn a pane into an {@link Role#ARCHITECT}.</li>
|
||||
* <li><b>slots</b> — read from {@code fleet.architects}/{@code developers}/{@code reviewers}
|
||||
* (see {@link #slots()}), each carrying the {@code profile} reference the spawn lifecycle
|
||||
* reads when it stands the slot up. <strong>Live, since fleetd #424</strong>: {@link #live}
|
||||
* re-reads {@code fleet:} on every call, through a supplier the same shape as
|
||||
* {@code CompositePeerLauncher}'s (see {@code ConfigRef}'s class doc) — so a config reload
|
||||
* that removes or adds an architect slot governs the <em>next</em> spawn with no restart.
|
||||
* Only {@link #MemberRegistry(FleetConfig.Fleet)} freezes the pool at construction, and that
|
||||
* constructor exists for tests and for the (rare) case of wiring a fixed, code-built config.</li>
|
||||
* <li><b>terminal bindings</b> — owned by this registry, initially <em>empty</em>, and
|
||||
* <strong>never</strong> touched by a reload. Config declares no architect terminal, so at
|
||||
* startup every slot is idle and nothing resolves to an architect; a session only becomes one
|
||||
* when the spawn lifecycle {@linkplain #bind(String, String) binds} its terminal to a slot.
|
||||
* {@link CallerResolver} reads this through {@link #snapshot()} to turn a pane into an
|
||||
* {@link Role#ARCHITECT}.</li>
|
||||
* </ul>
|
||||
*
|
||||
* <p><strong>The binding rule (fleetd #424): config governs what a bound slot still grants, as
|
||||
* well as what may be bound next.</strong> Removing a slot from config revokes it — that is the
|
||||
* ticket's entire point ("Revoking an architect slot does not revoke it"). Revoking it means an
|
||||
* architect already bound to that slot loses the ARCHITECT privilege on its very next request:
|
||||
* {@link #roleForSlot} and {@link #nameForSlot} read {@link #slots()} directly, with no cache, so
|
||||
* the moment a slot drops out of config, {@link CallerResolver#resolve} (which calls both on every
|
||||
* request from a bound pane, {@code CallerResolver.java:220}) can no longer confirm the pane's slot
|
||||
* is an architect slot, and the pane falls through to {@code Principal.worker(...)}. What does
|
||||
* <em>not</em> change is the {@code terminalToSlot} <em>occupancy</em> — the binding created by
|
||||
* {@link #bind} is untouched by a reload, on purpose: unbinding it here would double-book the slot
|
||||
* key (a second terminal could then bind to the "freed" key while the first is still the terminal
|
||||
* the operator actually meant to demote) and would silently break {@link #unbind}'s compare-safe
|
||||
* contract, which needs the original {@code terminal → slot} pair intact to remove it cleanly. So
|
||||
* the demoted session keeps occupying its slot — {@link #slotForTerminal} and {@link #snapshot()}
|
||||
* still name it — it just no longer resolves as an architect through that occupancy, and a fresh
|
||||
* spawn still cannot bind to the same key while it is occupied ({@link #reserve}/
|
||||
* {@link #requireSlotFor} refuse it anyway, since it is gone from {@link #slots()}). The demoted
|
||||
* session's own turn is unaffected: {@code fleet_reply}'s authorization
|
||||
* ({@code Authz.Action.REPLY}) is {@code caller.ownsSession(targetSession)} — identity by terminal,
|
||||
* not by role — so a demoted architect can still end its own turn normally.
|
||||
*
|
||||
* <p>Spawning/lifecycle is deliberately a separate unit: this class only owns the bindings and
|
||||
* exposes the map the resolver resolves against plus the profile lookup lifecycle will call.
|
||||
* Nothing here creates or manages an architect session.
|
||||
@@ -55,14 +83,42 @@ public final class MemberRegistry implements MemberLifecycle {
|
||||
}
|
||||
}
|
||||
|
||||
private final Map<String, Entry> slots;
|
||||
private final Supplier<FleetConfig.Fleet> fleet;
|
||||
/** Live {@code terminal_id → qualified slot key}; guarded by {@code terminalToSlot}. */
|
||||
private final Map<String, String> terminalToSlot = new HashMap<>();
|
||||
/** Slot keys held between reservation and the terminal binding. Guarded by terminalToSlot. */
|
||||
private final java.util.Set<String> reservedSlots = new java.util.HashSet<>();
|
||||
|
||||
/** Flatten every role pool in {@code fleet} into one registry. Leaders are not members. */
|
||||
/**
|
||||
* Freeze the pool at construction — for tests, and for the rare case of wiring a fixed,
|
||||
* code-built config. Production wiring should prefer {@link #live}, which re-reads
|
||||
* {@code fleet:} on every call.
|
||||
*/
|
||||
public MemberRegistry(FleetConfig.Fleet fleet) {
|
||||
this(() -> fleet);
|
||||
}
|
||||
|
||||
private MemberRegistry(Supplier<FleetConfig.Fleet> fleet) {
|
||||
this.fleet = fleet;
|
||||
}
|
||||
|
||||
/**
|
||||
* Live variant (fleetd #424): {@code fleet} is read fresh on every {@link #slots()} call — pass
|
||||
* {@code () -> config.get().fleet()}, the same supplier shape {@code CompositePeerLauncher}
|
||||
* already uses for placement — so a reload that adds or removes an architect slot governs the
|
||||
* next spawn's {@link #reserve}/{@link #requireSlotFor} check with no restart. A separate,
|
||||
* private constructor rather than a same-arity public overload of
|
||||
* {@link #MemberRegistry(FleetConfig.Fleet)}: a {@code FleetConfig.Fleet} and a
|
||||
* {@code Supplier<FleetConfig.Fleet>} overload are ambiguous for a literal {@code null} — the
|
||||
* same reason {@code CallerResolver.withLeads} is a static factory rather than a fourth
|
||||
* constructor overload.
|
||||
*/
|
||||
public static MemberRegistry live(Supplier<FleetConfig.Fleet> fleet) {
|
||||
return new MemberRegistry(Objects.requireNonNull(fleet, "fleet"));
|
||||
}
|
||||
|
||||
/** Flatten every role pool in {@code fleet} into one map. Leaders are not members. */
|
||||
private static Map<String, Entry> flatten(FleetConfig.Fleet fleet) {
|
||||
Map<String, Entry> flat = new LinkedHashMap<>();
|
||||
if (fleet != null) {
|
||||
for (MemberRole role : MemberRole.values()) {
|
||||
@@ -74,18 +130,22 @@ public final class MemberRegistry implements MemberLifecycle {
|
||||
});
|
||||
}
|
||||
}
|
||||
this.slots = Collections.unmodifiableMap(flat);
|
||||
return Collections.unmodifiableMap(flat);
|
||||
}
|
||||
|
||||
/** The configured slots, keyed by qualified {@link Entry#key()}. Unmodifiable snapshot. */
|
||||
/**
|
||||
* The configured slots, keyed by qualified {@link Entry#key()}. Unmodifiable snapshot of
|
||||
* {@code fleet:} <em>as of this call</em> — see the class doc for which constructor makes that
|
||||
* live versus frozen.
|
||||
*/
|
||||
public Map<String, Entry> slots() {
|
||||
return slots;
|
||||
return flatten(fleet.get());
|
||||
}
|
||||
|
||||
/** The slots belonging to {@code role}, in definition order. */
|
||||
/** The slots belonging to {@code role}, in definition order, as of this call. */
|
||||
public Map<String, Entry> slotsFor(MemberRole role) {
|
||||
Map<String, Entry> out = new LinkedHashMap<>();
|
||||
slots.forEach((key, e) -> {
|
||||
slots().forEach((key, e) -> {
|
||||
if (e.role() == role) {
|
||||
out.put(key, e);
|
||||
}
|
||||
@@ -117,31 +177,49 @@ public final class MemberRegistry implements MemberLifecycle {
|
||||
}
|
||||
|
||||
/**
|
||||
* The strong-model profile a slot runs under — what the spawn lifecycle reads.
|
||||
* The strong-model profile a slot runs under, as of this call.
|
||||
*
|
||||
* <p>Nothing in {@code src/main} calls this (fleetd #431 — grepped both the {@code
|
||||
* .profileForSlot(} and the {@code ::profileForSlot} form). This javadoc used to say "what the
|
||||
* spawn lifecycle reads", and that seam does not exist: the spawn lifecycle takes its profile
|
||||
* from the {@link MemberLifecycle.SlotReservation} that {@code reserve} returns, never from
|
||||
* here. Kept and pinned rather than deleted because it is the natural accessor for that seam
|
||||
* if one is added; live for the same reason as {@link #roleForSlot}, so a reload cannot leave
|
||||
* it answering for the old config.
|
||||
*
|
||||
* @return the slot's configured {@code profile}, or {@code null} if the slot is unknown or
|
||||
* declares none
|
||||
*/
|
||||
public String profileForSlot(String slotName) {
|
||||
Entry e = slots.get(slotName);
|
||||
Entry e = slots().get(slotName);
|
||||
return (e == null || e.profile() == null) ? null : e.profile();
|
||||
}
|
||||
|
||||
/** The role a qualified slot key belongs to, or {@code null} when the key is unknown. */
|
||||
/**
|
||||
* The role a qualified slot key belongs to, or {@code null} when the key is not currently
|
||||
* configured. Deliberately live, with no cache (fleetd #424, see the class doc's binding rule):
|
||||
* removing a slot from config must make {@link CallerResolver#resolve} stop granting the
|
||||
* ARCHITECT role for it on the very next request from a terminal that was bound to it, which is
|
||||
* the ticket's whole point — revoking a slot must actually revoke it, not just refuse the next
|
||||
* spawn.
|
||||
*/
|
||||
public MemberRole roleForSlot(String slotName) {
|
||||
Entry e = slots.get(slotName);
|
||||
Entry e = slots().get(slotName);
|
||||
return e == null ? null : e.role();
|
||||
}
|
||||
|
||||
/** The unqualified configured name for a slot, or {@code null} if it is unknown. */
|
||||
/**
|
||||
* The unqualified configured name for a slot, or {@code null} if it is not currently configured.
|
||||
* Live for the same reason as {@link #roleForSlot} — see the class doc's binding rule.
|
||||
*/
|
||||
public String nameForSlot(String slotName) {
|
||||
Entry e = slots.get(slotName);
|
||||
Entry e = slots().get(slotName);
|
||||
return e == null ? null : e.name();
|
||||
}
|
||||
|
||||
/** True when {@code slotName} is a configured architect slot. */
|
||||
public boolean isSlot(String slotName) {
|
||||
return slots.containsKey(slotName);
|
||||
return slots().containsKey(slotName);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -32,9 +32,26 @@ import java.util.function.Supplier;
|
||||
* not the fact that they are config. Most of {@code fleet:} — every role pool
|
||||
* ({@code architects}/{@code developers}/{@code reviewers}), {@code charters}, and
|
||||
* {@code tabLabel} — is read the same live way, through the same supplier
|
||||
* ({@code () -> config.get().fleet()}). <strong>But {@code fleet:} as a whole is NOT in this
|
||||
* class</strong>: {@code fleet.leaders} inside the same key is frozen, which is exactly what
|
||||
* makes {@code fleet:} split rather than hot — see below.</li>
|
||||
* ({@code () -> config.get().fleet()}). {@code architects} in particular is hot for
|
||||
* <strong>two independent consumers</strong> (fleetd #424): {@code CompositePeerLauncher}
|
||||
* reads it live for placement (which profile an unqualified architect spawn may land on), and
|
||||
* {@code MemberRegistry} separately reads it live, through its own instance of the same
|
||||
* supplier shape, for identity — both which slot a spawn may bind to <em>and</em> what a slot
|
||||
* already bound still grants. Removing an architect slot from config therefore revokes the
|
||||
* {@link dev.ltms.fleet.auth.Role#ARCHITECT} role on the bound pane's very next request; only
|
||||
* the slot <em>occupancy</em> survives, so the demoted session still holds its slot key until
|
||||
* it unbinds. See {@code MemberRegistry}'s class doc for that binding rule.
|
||||
* <strong>But {@code fleet:} as a whole is NOT in this class</strong>: {@code fleet.leaders}
|
||||
* inside the same key is frozen, which is exactly what makes {@code fleet:} split rather than
|
||||
* hot — see below. {@code models:} (fleetd #422) joined this class whole: {@link
|
||||
* FleetConfig#validateModels()} re-runs fully against the fresh config on every {@link
|
||||
* #reload()} (via {@link FleetConfig#validateAll()}), refusing a bad edit outright rather than
|
||||
* caching a stale copy anywhere, and the on/off half added by fleetd #422 is read live both by
|
||||
* {@code CompositePeerLauncher}'s spawn gate ({@code enforceModelEnabled} and its candidate
|
||||
* filter) and by {@code fleet_profiles}/{@code GET /profiles} (via
|
||||
* {@code PeerLauncher.disabledModels()}). Nothing about {@code models:} is baked into an
|
||||
* object built at startup, so — unlike the deferred keys below — there is no frozen half left
|
||||
* to report; it moved here from deferred rather than joining split.</li>
|
||||
* <li><strong>Deferred</strong> — accepted into the new snapshot, but the wiring built at startup
|
||||
* keeps the old value until a restart: {@code lifecycle:}, {@code leadHeartbeat:},
|
||||
* {@code idleSleepGuard:} ({@code Fleetd.java} reads it once, at startup, to decide whether
|
||||
@@ -43,12 +60,6 @@ import java.util.function.Supplier;
|
||||
* running daemon keeps whatever this was at startup regardless of a later edit),
|
||||
* {@code spawnReadyTimeoutMs} / {@code spawnReadyPollMs}, {@code quarantineCooldownSeconds}
|
||||
* (CB-578 stage B — baked once into the {@code BackendQuarantine} built at startup),
|
||||
* {@code models:} (fleetd ticket "central allow-list of usable models" — {@link
|
||||
* FleetConfig#validateModels()} re-runs against the fresh config in {@link #reload()}
|
||||
* (via {@link FleetConfig#validateAll()}), so a
|
||||
* models.allow: edit that would refuse to boot still refuses the reload; a change that
|
||||
* passes has nothing built at startup to rebuild, so it is reported deferred rather than
|
||||
* silently accepted with no report at all),
|
||||
* {@code guard:}, {@code worktreeRoot:}, {@code worktreeGroup:} and {@code memberSkills:}
|
||||
* (all three of the latter baked once into the {@code GitWorktrees} built at
|
||||
* {@code Fleetd.java:251} and never rebuilt — fleetd #323 instance 2 found
|
||||
@@ -141,16 +152,18 @@ import java.util.function.Supplier;
|
||||
* </ul>
|
||||
*
|
||||
* <p><strong>The denominator, measured on 2026-09-04 (fleetd #330; recounted for fleetd #333);
|
||||
* recounted again for fleetd #362, again after {@code idleSleepGuard:} was added, and again after
|
||||
* {@code models:} was added.</strong>
|
||||
* {@code FleetConfig} has 25 top-level record components: 5 cold, 14 deferred, 3 split, 3
|
||||
* hot-excluded. Three of them are named nowhere in this file, and the reason is the same for all
|
||||
* three: {@code placement}, {@code memberCredentials} and {@code memberLoginShell} are
|
||||
* <strong>hot</strong> and correctly absent — all three are read live off {@code config.get()}
|
||||
* recounted again for fleetd #362, again after {@code idleSleepGuard:} was added, again after
|
||||
* {@code models:} was added as deferred, and again for fleetd #422, which moved {@code models:}
|
||||
* from deferred to hot-excluded once its on/off half was read live everywhere.</strong>
|
||||
* {@code FleetConfig} has 25 top-level record components: 5 cold, 13 deferred, 3 split, 4
|
||||
* hot-excluded. Four of them are named nowhere in this file, and the reason is the same for all
|
||||
* four: {@code placement}, {@code memberCredentials}, {@code memberLoginShell} and {@code models}
|
||||
* are <strong>hot</strong> and correctly absent — all four are read live off {@code config.get()}
|
||||
* (placement through the {@code CompositePeerLauncher} supplier the Hot bullet names;
|
||||
* {@code memberCredentials}/{@code memberLoginShell} 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.
|
||||
* and {@code HerdrPeerLauncher#configuredMemberLoginShell}; {@code models} the same way, through the
|
||||
* Hot bullet's {@code models:} paragraph), so a reload takes effect on the next spawn (or, for
|
||||
* {@code models}, the next reported status) 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
|
||||
@@ -226,7 +239,7 @@ public final class ConfigRef implements Supplier<FleetConfig> {
|
||||
static final Set<String> DEFERRED_KEYS = Set.of(
|
||||
"guard", "worktreeRoot", "worktreeGroup", "memberSkills", "primary", "configReload",
|
||||
"leadHeartbeat", "lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs",
|
||||
"quarantineCooldownSeconds", "profiles", "idleSleepGuard", "models");
|
||||
"quarantineCooldownSeconds", "profiles", "idleSleepGuard");
|
||||
|
||||
private final Path path;
|
||||
private final AtomicReference<FleetConfig> current;
|
||||
@@ -446,15 +459,6 @@ public final class ConfigRef implements Supplier<FleetConfig> {
|
||||
if (!Objects.equals(old.idleSleepGuard(), fresh.idleSleepGuard())) {
|
||||
changed.add("idleSleepGuard");
|
||||
}
|
||||
// fleetd ticket "central allow-list of usable models": validateModels() runs again in
|
||||
// reload() above (via validateAll()), so a bad edit is already refused as cold-adjacent
|
||||
// (the whole reload is refused via the catch block, never partially applied). A GOOD edit
|
||||
// to the allow-list
|
||||
// itself has nothing built at startup to rebuild — it only ever mattered to the validation
|
||||
// call that already ran — so report it deferred rather than silently swallowing the change.
|
||||
if (!Objects.equals(old.models(), fresh.models())) {
|
||||
changed.add("models");
|
||||
}
|
||||
if (!Objects.equals(old.spawnReadyTimeoutMs(), fresh.spawnReadyTimeoutMs())
|
||||
|| !Objects.equals(old.spawnReadyPollMs(), fresh.spawnReadyPollMs())) {
|
||||
changed.add("spawnReady*");
|
||||
@@ -529,26 +533,35 @@ public final class ConfigRef implements Supplier<FleetConfig> {
|
||||
+ "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");
|
||||
}
|
||||
// fleetd #333: unlike health/coordinator above, most of `fleet:` (architects, developers,
|
||||
// reviewers, charters, tabLabel) is genuinely hot — ConfigRefTest.aHotChangeIsAppliedAndRead-
|
||||
// fleetd #333: unlike health/coordinator above, most of `fleet:` (developers, reviewers,
|
||||
// charters, tabLabel) is genuinely hot — ConfigRefTest.aHotChangeIsAppliedAndRead-
|
||||
// ThroughGet and aCharterChangeIsHotAndReachesTheLiveConfig prove it reaches the live config
|
||||
// with no restart note. Only fleet.leaders is frozen (Fleetd.java:281 reads
|
||||
// cfg.fleet().leaders() off the startup snapshot to build both the LeadTabScanner's
|
||||
// tab-label-to-name map, wired into CallerResolver.withLeadsAndMembers at Fleetd.java:620/624,
|
||||
// and — when herdr answered — LeadLauncher(...).ensureLeads() at Fleetd.java:315, which
|
||||
// auto-launches each lead up to its `instances` count; neither is rebuilt on reload). So this
|
||||
// compares fleet.leaders alone, not the whole Fleet record: comparing the whole record would
|
||||
// report "split" for a tabLabel-only or charters-only change that is actually fully hot,
|
||||
// which is the over-claim mirror of the under-claim bug this class exists to prevent.
|
||||
// with no restart note. `architects` is hot too, and — since fleetd #424 — hot for BOTH of
|
||||
// its consumers, not just the one this comment used to name: CompositePeerLauncher reads it
|
||||
// live for PLACEMENT through the () -> config.get().fleet() supplier named in the class doc's
|
||||
// Hot bullet, and MemberRegistry separately reads it live for IDENTITY (which slot a spawn
|
||||
// may bind to, AND what a slot already bound still grants) through its own instance of that
|
||||
// same supplier shape — see MemberRegistry.live and its class doc for the binding rule:
|
||||
// removing a slot revokes ARCHITECT on the bound pane's very next request, and only the slot
|
||||
// OCCUPANCY survives, so the demoted session keeps its slot key until it unbinds. Only
|
||||
// fleet.leaders is frozen (Fleetd.java:281 reads cfg.fleet().leaders() off the startup
|
||||
// snapshot to build both the LeadTabScanner's tab-label-to-name map, wired into
|
||||
// CallerResolver.withLeadsAndMembers at Fleetd.java:620/624, and — when herdr answered —
|
||||
// LeadLauncher(...).ensureLeads() at Fleetd.java:315, which auto-launches each lead up to its
|
||||
// `instances` count; neither is rebuilt on reload). So this compares fleet.leaders alone, not
|
||||
// the whole Fleet record: comparing the whole record would report "split" for a tabLabel-only
|
||||
// or architects-only change that is actually fully hot, which is the over-claim mirror of the
|
||||
// under-claim bug this class exists to prevent.
|
||||
if (!Objects.equals(leadersOf(old), leadersOf(fresh))) {
|
||||
changed.add("fleet: fleet.leaders (each lead's tab, workspace, cwd, profile and "
|
||||
+ "instances count) is read once at startup to build the LeadTabScanner's "
|
||||
+ "identity map and to auto-launch leads, and neither is rebuilt on reload, so a "
|
||||
+ "lead added, removed, or given a new tab: label needs a restart — until then it "
|
||||
+ "stays unrecognised, and a caller from its new tab resolves as a worker, not a "
|
||||
+ "lead; the rest of fleet: (architects, developers, reviewers, charters, "
|
||||
+ "tabLabel) is read live through the supplier on CompositePeerLauncher and "
|
||||
+ "already applied");
|
||||
+ "lead; the rest of fleet: (developers, reviewers, charters, tabLabel) is read "
|
||||
+ "live through the supplier on CompositePeerLauncher, and architects is read "
|
||||
+ "live through that same supplier for placement AND through a separate supplier "
|
||||
+ "on MemberRegistry for spawn-time identity — both 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 split keys the class doc documents.
|
||||
|
||||
@@ -133,8 +133,10 @@ import java.util.regex.PatternSyntaxException;
|
||||
* nothing and every existing config keeps working exactly as it does today.
|
||||
* When non-empty, a profile whose {@code model:} is not one of {@link
|
||||
* Models#ids()} fails config load, naming both the model and the profile —
|
||||
* see {@link #validateModels()}. This block only decides what may be
|
||||
* CONFIGURED; nothing here enforces it at spawn time. See {@link Models}.
|
||||
* see {@link #validateModels()}. This block decides what may be CONFIGURED;
|
||||
* fleetd #422 added the separate on/off question — whether a configured model
|
||||
* may be spawned onto RIGHT NOW ({@link Models.ModelEntry#enabled}) — enforced
|
||||
* live at spawn by {@code CompositePeerLauncher}, not here. See {@link Models}.
|
||||
*/
|
||||
@JsonIgnoreProperties(ignoreUnknown = true)
|
||||
public record FleetConfig(
|
||||
@@ -1362,9 +1364,14 @@ public record FleetConfig(
|
||||
* the set of permitted models — only editing {@code models.allow:} itself can. This is the
|
||||
* invariant the ticket asked for: the two blocks are validated in one direction only.
|
||||
*
|
||||
* <p><b>Out of scope here, deliberately:</b> nothing in this block is read at spawn time —
|
||||
* enforcing it against a live spawn, an on/off runtime switch, and any interaction with {@code
|
||||
* BackendQuarantine} are separate units. This block is config-load validation only.
|
||||
* <p><b>Spawn-time enforcement (fleetd #422, units 2+3) lives outside this record</b> —
|
||||
* {@code CompositePeerLauncher.enforceModelEnabled} and its candidate-set filter read this
|
||||
* block LIVE (through the same kind of supplier {@code weight}/{@code maxLoad} already use),
|
||||
* so the on/off state below is hot: no restart needed. This block itself still only decides
|
||||
* what may be CONFIGURED (membership in {@link #allow}); {@link ModelEntry#enabled} decides
|
||||
* whether a member of that list is currently spawnable. The two questions are deliberately
|
||||
* separate — see {@link ModelEntry}'s javadoc for why turning a model off must never mean
|
||||
* removing it from {@link #allow}.
|
||||
*
|
||||
* @param allow the permitted models, each its own {@link ModelEntry} rather than a bare
|
||||
* string — see that record's javadoc for why. {@code null}/empty ⇒ the block is
|
||||
@@ -1378,24 +1385,56 @@ public record FleetConfig(
|
||||
}
|
||||
|
||||
/**
|
||||
* One permitted model, named as a record rather than a bare string on purpose: a later unit
|
||||
* needs to hang an on/off state and a load-limit state off each entry, and a bare {@code
|
||||
* List<String>} cannot grow those fields without changing the YAML shape underneath every
|
||||
* operator who already wrote one. {@link #model()} is intentionally a single flat,
|
||||
* opaque-string namespace — a bare Claude id ({@code claude-sonnet-5}) and an opencode
|
||||
* provider-prefixed id ({@code openai/gpt-5.6-terra}) both fit it unchanged, because
|
||||
* {@link FleetConfig#validateModels()} only ever compares a profile's {@code model:} value
|
||||
* against this string for exact equality; it never parses a provider prefix or branches on
|
||||
* a profile's {@code kind:}.
|
||||
* One permitted model, named as a record rather than a bare string on purpose: this ticket
|
||||
* (fleetd #422) is the "later unit" the original comment here predicted — it hangs an on/off
|
||||
* state ({@link #enabled}) off each entry, and a bare {@code List<String>} could not have
|
||||
* grown that field without changing the YAML shape underneath every operator who already
|
||||
* wrote one. {@link #model()} is intentionally a single flat, opaque-string namespace — a
|
||||
* bare Claude id ({@code claude-sonnet-5}) and an opencode provider-prefixed id
|
||||
* ({@code openai/gpt-5.6-terra}) both fit it unchanged, because {@link
|
||||
* FleetConfig#validateModels()} only ever compares a profile's {@code model:} value against
|
||||
* this string for exact equality; it never parses a provider prefix or branches on a
|
||||
* profile's {@code kind:}.
|
||||
*
|
||||
* @param model the model id exactly as a {@code profiles:} entry's {@code model:} field
|
||||
* would name it
|
||||
* @param model the model id exactly as a {@code profiles:} entry's {@code model:} field
|
||||
* would name it
|
||||
* @param enabled {@code false} turns spawning onto this model off; {@code null} (the field
|
||||
* omitted — every config written before fleetd #422 is this shape) or
|
||||
* {@code true} leaves it on. Turning a model off must NEVER remove it from
|
||||
* {@link Models#allow} — {@link FleetConfig#validateModels()} checks
|
||||
* <em>membership</em> only, never the on/off state, so an off entry stays a
|
||||
* valid thing for a {@code profiles:} entry to name; only
|
||||
* {@code CompositePeerLauncher}'s spawn-time gate reads {@link #enabled}.
|
||||
* Collapsing the two — turning a model off by deleting its {@code allow:}
|
||||
* entry — would make {@link FleetConfig#validateModels()} refuse the whole
|
||||
* config reload the moment a still-configured profile names it, which is
|
||||
* exactly the restart-to-flip-a-switch problem this field exists to avoid.
|
||||
* <p>A model id can be named by more than one {@code profiles:} entry (e.g.
|
||||
* {@code deepseek-v4-flash} backs both {@code local} and {@code
|
||||
* local-direct} in the live config) — turning it off disables every profile
|
||||
* that names it, on purpose: the model is what a subscription's rate limit
|
||||
* actually constrains, not any one profile alias for it.
|
||||
*/
|
||||
@JsonIgnoreProperties(ignoreUnknown = true)
|
||||
public record ModelEntry(String model) {
|
||||
public record ModelEntry(String model, Boolean enabled) {
|
||||
public ModelEntry {
|
||||
model = (model == null || model.isBlank()) ? null : model.trim();
|
||||
}
|
||||
|
||||
/**
|
||||
* Back-compat form before {@link #enabled} was added (fleetd #422) — the model is
|
||||
* unconditionally on, exactly as every {@code ModelEntry} behaved before this field
|
||||
* existed. Keeps pre-#422 call sites (and any YAML that omits {@code enabled:})
|
||||
* compiling and behaving identically.
|
||||
*/
|
||||
public ModelEntry(String model) {
|
||||
this(model, null);
|
||||
}
|
||||
|
||||
/** {@code true} unless {@link #enabled} is explicitly {@code false} — absent means on. */
|
||||
public boolean isEnabled() {
|
||||
return !Boolean.FALSE.equals(enabled);
|
||||
}
|
||||
}
|
||||
|
||||
/** {@link #allow}'s model ids, as a set for membership checks. Blank/null entries are dropped. */
|
||||
@@ -1408,6 +1447,23 @@ public record FleetConfig(
|
||||
}
|
||||
return Collections.unmodifiableSet(ids);
|
||||
}
|
||||
|
||||
/**
|
||||
* Model ids currently turned off (fleetd #422: {@link ModelEntry#isEnabled()} {@code
|
||||
* false}). Read live by {@code CompositePeerLauncher}'s spawn gate and by {@code
|
||||
* fleet_profiles}/{@code GET /profiles} — both must read this same accessor off the same
|
||||
* live config so the two surfaces cannot disagree about which model is off (the fleetd
|
||||
* #404 lesson: a status field must read the source the behaviour reads).
|
||||
*/
|
||||
public Set<String> offIds() {
|
||||
Set<String> off = new java.util.LinkedHashSet<>();
|
||||
for (ModelEntry e : allow) {
|
||||
if (e != null && e.model() != null && !e.isEnabled()) {
|
||||
off.add(e.model());
|
||||
}
|
||||
}
|
||||
return Collections.unmodifiableSet(off);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -698,17 +698,57 @@ public final class CompletionResolver implements TurnListener {
|
||||
}
|
||||
|
||||
/**
|
||||
* Coverage summary for the CB-578 stage A exhausted-pattern classification, logged at startup
|
||||
* What an unset pattern key means for the classification it configures (fleetd#415).
|
||||
* {@code coverage()} cannot infer this from the key's name — the two keys it currently
|
||||
* describes disagree on it, and a string comparison on the name would just move the same bug
|
||||
* to a new spot — so every caller must state it explicitly.
|
||||
*
|
||||
* <p><strong>This alone does not prove a caller passes the right one for its key.</strong> A
|
||||
* test that calls {@code coverage()} directly and supplies the meaning itself only proves this
|
||||
* enum is worded correctly, never that {@code Fleetd}'s two call sites pair each key with its
|
||||
* true meaning — that pairing is #415's actual defect. Measured on review: swapping the two
|
||||
* {@code UnsetMeaning} arguments at those call sites (giving {@code exhaustedPattern} the
|
||||
* built-in-default wording and {@code errorPattern} the off wording — #415's exact defect with
|
||||
* the keys exchanged) compiled with 0 errors and left all 1506 existing tests green. See
|
||||
* {@code dev.ltms.fleet.Fleetd#exhaustedPatternCoverageLine}/{@code #errorPatternCoverageLine}
|
||||
* and {@code FleetdPatternCoverageLineTest}, which exists specifically to catch that swap.
|
||||
*/
|
||||
public enum UnsetMeaning {
|
||||
/** No fallback exists: a profile with no configured pattern truly has this classification off. */
|
||||
OFF,
|
||||
/** A built-in pattern applies when unset: the classification still runs for that profile. */
|
||||
BUILT_IN_DEFAULT
|
||||
}
|
||||
|
||||
/**
|
||||
* Coverage summary for a fleetd#201/CB-578-style pattern-key classification, logged at startup
|
||||
* the way {@link dev.ltms.fleet.health.FleetHealthMonitor#coverage} is — so an operator can
|
||||
* see whether the classification is on, and for which profiles, without reading every
|
||||
* profile's config by hand.
|
||||
*
|
||||
* <p>fleetd#415: this method measures <em>pattern coverage</em> — how many profiles set the
|
||||
* key — which is not the same thing as <em>feature state</em> for a key with a fallback. For
|
||||
* {@code errorPattern}, an empty {@code configuredProfiles} still runs the classification
|
||||
* against {@code CompletionResolver}'s built-in compatibility pattern ({@link #BACKEND_ERROR}
|
||||
* at line ~84); for {@code exhaustedPattern} there is no fallback, so empty really does mean
|
||||
* off. {@code unsetMeaning} is the single, required source of that fact — see
|
||||
* {@link dev.ltms.fleet.config.FleetConfig#rejectMalformedProfilePatterns} lines ~2029-2032 for
|
||||
* where it is documented for config authors. It is a required parameter, not a defaulted
|
||||
* overload: a third pattern key added later must supply one to compile at all, rather than
|
||||
* silently inheriting whichever wording this method happened to default to.
|
||||
*
|
||||
* @param allProfiles every configured profile name
|
||||
* @param configuredProfiles the subset of {@code allProfiles} that carry an exhausted pattern
|
||||
* @param configuredProfiles the subset of {@code allProfiles} that carry the pattern
|
||||
*/
|
||||
public static String coverage(String patternKey, Set<String> allProfiles, Set<String> configuredProfiles) {
|
||||
public static String coverage(String patternKey, UnsetMeaning unsetMeaning, Set<String> allProfiles,
|
||||
Set<String> configuredProfiles) {
|
||||
if (configuredProfiles.isEmpty()) {
|
||||
return "off (no profile has an " + patternKey + " configured; profiles: " + sorted(allProfiles) + ")";
|
||||
return switch (unsetMeaning) {
|
||||
case OFF -> "off (no profile has an " + patternKey + " configured; profiles: "
|
||||
+ sorted(allProfiles) + ")";
|
||||
case BUILT_IN_DEFAULT -> "built-in default for all profiles (no profile customises "
|
||||
+ patternKey + "; profiles: " + sorted(allProfiles) + ")";
|
||||
};
|
||||
}
|
||||
Set<String> unconfigured = new TreeSet<>(allProfiles);
|
||||
unconfigured.removeAll(configuredProfiles);
|
||||
|
||||
@@ -37,6 +37,7 @@ import io.modelcontextprotocol.spec.McpSchema;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import jakarta.servlet.http.HttpServlet;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
@@ -1179,6 +1180,23 @@ public final class FleetMcp {
|
||||
if (!coolingOff.isEmpty()) {
|
||||
result.put("coolingOff", coolingOff);
|
||||
}
|
||||
// fleetd #422: read the exact same accessor CompositePeerLauncher's spawn gate reads
|
||||
// (PeerLauncher.modelGateState(), which for the composite is models0() read live) — never a
|
||||
// separately-derived answer, so this status can never overstate or understate what the gate
|
||||
// actually enforces (the fleetd #404 lesson).
|
||||
//
|
||||
// fleetd #422 follow-up: "armed" and "off" come from the ONE modelGateState() call below,
|
||||
// never two independent reads of the gate — a reload landing between two separate reads
|
||||
// could otherwise make them disagree. modelGateArmed is reported unconditionally (never
|
||||
// omitted like quarantined/coolingOff above) precisely so a lead can tell "no models: block
|
||||
// at all" (false) apart from "a models: block with nothing currently off" (true, with
|
||||
// modelsOff simply absent below) — the two states PeerLauncher.disabledModels() alone
|
||||
// cannot distinguish, both reporting an empty set.
|
||||
PeerLauncher.ModelGateState modelGate = workers.modelGateState();
|
||||
result.put("modelGateArmed", modelGate.configured());
|
||||
if (!modelGate.off().isEmpty()) {
|
||||
result.put("modelsOff", new ArrayList<>(modelGate.off()));
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
|
||||
@@ -90,6 +90,19 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
private final Supplier<Map<String, FleetConfig.Profile>> profileConfigs;
|
||||
private final Supplier<PlacementPolicy> placementPolicy;
|
||||
|
||||
/**
|
||||
* fleetd #422: the central model allow-list's on/off state, read live per spawn — same reason
|
||||
* {@link #profileConfigs} is a supplier rather than a captured map (see the class doc above and
|
||||
* {@link #enforceModelEnabled}). A caller with no {@code models:} block to read from (the
|
||||
* simpler, map-based constructors used throughout this class's own tests) wires this to a
|
||||
* constant {@code null}, which {@link #models0} treats as "nothing configured, gate never
|
||||
* fires" — the pre-#422 behaviour.
|
||||
*/
|
||||
private final Supplier<FleetConfig.Models> models;
|
||||
|
||||
/** The value {@link #models0} normalizes a {@code null} supplier result to. */
|
||||
private static final FleetConfig.Models NO_MODELS_CONFIGURED = new FleetConfig.Models(List.of());
|
||||
|
||||
/** CB-578 stage B: credential cooldown, checked before an explicit spawn and filtered into placement. */
|
||||
private final BackendQuarantine quarantine;
|
||||
|
||||
@@ -205,12 +218,44 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
FleetConfig.Fleet fleet,
|
||||
BackendQuarantine quarantine,
|
||||
BackendOutagePolicy outagePolicy) {
|
||||
// No models: block to read from a plain profiles Map — fleetd #422's gate is wired to a
|
||||
// constant null (see enforceModelEnabled/models0), the pre-#422 behaviour for every caller
|
||||
// of this overload. Use the 9-arg overload below to test the gate against a LIVE supplier.
|
||||
this(delegates, defaultProfile, profileConfigs, placementPolicy, liveCount, fleet, quarantine,
|
||||
outagePolicy, constant(null));
|
||||
}
|
||||
|
||||
/**
|
||||
* As above, plus a LIVE model-gate source (fleetd #422). The 8-arg overload above wires
|
||||
* {@code models} to a constant {@code null} because it has only a static {@code Map<String,
|
||||
* Profile>}, never a full {@code FleetConfig}, to read one from; this overload exists so a test
|
||||
* can prove {@link #enforceModelEnabled} and its candidate filter re-read {@link
|
||||
* FleetConfig.Models} on every call rather than a value captured once at construction — the
|
||||
* exact distinction fleetd #422 exists to get right (see the class doc's CB-559 note on {@link
|
||||
* #profileConfigs}, which this follows). Production wiring uses the
|
||||
* {@code Supplier<FleetConfig>} constructor below instead, which already threads a live
|
||||
* {@code config.get().models()} through.
|
||||
*
|
||||
* @param models required — pass a supplier returning {@code null} for a caller that has no
|
||||
* {@code models:} block to gate against, never a defaulting overload (the same
|
||||
* "explicit opt-out, never a silent default" rule {@code quarantine}/
|
||||
* {@code outagePolicy} already follow).
|
||||
*/
|
||||
public CompositePeerLauncher(List<HerdrPeerLauncher> delegates,
|
||||
String defaultProfile,
|
||||
Map<String, FleetConfig.Profile> profileConfigs,
|
||||
PlacementPolicy placementPolicy,
|
||||
Function<String, Integer> liveCount,
|
||||
FleetConfig.Fleet fleet,
|
||||
BackendQuarantine quarantine,
|
||||
BackendOutagePolicy outagePolicy,
|
||||
Supplier<FleetConfig.Models> models) {
|
||||
// LinkedHashMap, not Map.copyOf: candidates() promises definition order and the weighted
|
||||
// policy breaks exact-weight ties on it, so a salted iteration order would make placement
|
||||
// differ from one JVM run to the next.
|
||||
this(delegates, defaultProfile,
|
||||
constant(Collections.unmodifiableMap(new LinkedHashMap<>(profileConfigs))),
|
||||
constant(placementPolicy), liveCount, constant(fleet), quarantine, outagePolicy);
|
||||
constant(placementPolicy), liveCount, constant(fleet), quarantine, outagePolicy, models);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -250,7 +295,10 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
liveCount,
|
||||
() -> config.get().fleet(),
|
||||
quarantine,
|
||||
outagePolicy);
|
||||
outagePolicy,
|
||||
// fleetd #422: read live, same as profiles/placement/fleet above — a models.allow
|
||||
// edit (on/off or otherwise) is visible to the very next spawn, no restart needed.
|
||||
() -> config.get().models());
|
||||
}
|
||||
|
||||
/** The all-suppliers form every other constructor funnels into. */
|
||||
@@ -261,7 +309,8 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
Function<String, Integer> liveCount,
|
||||
Supplier<FleetConfig.Fleet> fleet,
|
||||
BackendQuarantine quarantine,
|
||||
BackendOutagePolicy outagePolicy) {
|
||||
BackendOutagePolicy outagePolicy,
|
||||
Supplier<FleetConfig.Models> models) {
|
||||
this.fleet = fleet;
|
||||
this.quarantine = Objects.requireNonNull(quarantine, "quarantine");
|
||||
this.outagePolicy = Objects.requireNonNull(outagePolicy, "outagePolicy");
|
||||
@@ -273,6 +322,7 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
this.profileConfigs = profileConfigs;
|
||||
this.placementPolicy = placementPolicy;
|
||||
this.liveCount = liveCount;
|
||||
this.models = Objects.requireNonNull(models, "models");
|
||||
Map<String, HerdrPeerLauncher> index = new LinkedHashMap<>();
|
||||
for (HerdrPeerLauncher d : this.delegates) {
|
||||
for (String profile : d.profiles()) {
|
||||
@@ -303,6 +353,16 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
return m == null ? Map.of() : m;
|
||||
}
|
||||
|
||||
/**
|
||||
* The currently-configured {@code models:} block, never null (fleetd #422). Read fresh on
|
||||
* every call, the same reason {@link #profiles0} is — a config reload's on/off edit must reach
|
||||
* the very next spawn.
|
||||
*/
|
||||
private FleetConfig.Models models0() {
|
||||
FleetConfig.Models m = models.get();
|
||||
return m == null ? NO_MODELS_CONFIGURED : m;
|
||||
}
|
||||
|
||||
/** The adapter owning {@code profileName} (null/blank → the default). Throws on an unknown profile. */
|
||||
private HerdrPeerLauncher route(String profileName) {
|
||||
String resolved = (profileName == null || profileName.isBlank()) ? defaultProfile : profileName;
|
||||
@@ -332,6 +392,7 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
enforceNotQuarantined(requestedProfile);
|
||||
enforceNotCoolingOff(requestedProfile);
|
||||
enforceMaxLoad(requestedProfile);
|
||||
enforceModelEnabled(requestedProfile);
|
||||
PeerHandle handle = d.spawn(req);
|
||||
spawnedBy.put(handle.id(), d);
|
||||
return handle;
|
||||
@@ -349,8 +410,11 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
Set<String> quarantined = quarantinedProfiles(candidates);
|
||||
// fleetd #201 Unit 5: a distinct set from quarantined — see PlacementContext.coolingOff.
|
||||
Set<String> coolingOff = coolingOffProfiles(candidates);
|
||||
// fleetd #422: read live per spawn, same as quarantined/coolingOff above — a config reload
|
||||
// that flips a model's enabled state is visible to the very next unqualified spawn.
|
||||
Set<String> modelOff = modelOffProfiles(candidates);
|
||||
PlacementContext ctx = new PlacementContext(roleDefault, candidates, liveCount, unreachable,
|
||||
quarantined, coolingOff);
|
||||
quarantined, coolingOff, modelOff);
|
||||
|
||||
int maxAttempts = candidates.isEmpty() ? 1 : candidates.size();
|
||||
for (int attempt = 0; attempt < maxAttempts; attempt++) {
|
||||
@@ -380,7 +444,7 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
unreachable.add(chosen.profile());
|
||||
// Update the context for the next selection so the policy excludes this profile.
|
||||
ctx = new PlacementContext(roleDefault, candidates, liveCount, unreachable,
|
||||
quarantined, coolingOff);
|
||||
quarantined, coolingOff, modelOff);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -492,6 +556,49 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Refuse an explicit-profile spawn whose {@code model:} the operator has turned off in the
|
||||
* central {@code models.allow:} list (fleetd #422). Deliberately worded apart from {@link
|
||||
* #enforceNotQuarantined} and {@link #enforceNotCoolingOff}: those two report a BACKEND-reported
|
||||
* outage (exhaustion, repeated errors); this one reports an OPERATOR decision, so the message
|
||||
* says "turned off" and names the model, never "quarantined" or "cooling off". A FOURTH,
|
||||
* independent reason to refuse a spawn — never layered onto {@code BackendQuarantine} or {@code
|
||||
* BackendOutagePolicy}, which would misattribute an operator's own choice to the backend.
|
||||
*
|
||||
* <p>Reads {@link #models0()} fresh on every call — the same liveness {@link #profiles0()}
|
||||
* already has — so flipping {@code enabled: false} and reloading takes effect on the very next
|
||||
* spawn, no restart (criterion 4). A profile naming no model, or a model absent from {@code
|
||||
* models.allow:} entirely (nothing to gate against), is never refused here.
|
||||
*
|
||||
* @throws PlacementException naming the model and the profile, distinct from quarantine/cool-off
|
||||
*/
|
||||
private void enforceModelEnabled(String profile) {
|
||||
FleetConfig.Profile cfg = profiles0().get(profile);
|
||||
String model = (cfg == null) ? null : cfg.model();
|
||||
if (model == null) {
|
||||
return;
|
||||
}
|
||||
if (models0().offIds().contains(model)) {
|
||||
throw new PlacementException("worker profile '" + profile + "' names model '" + model
|
||||
+ "', which the operator has turned off in models.allow — refusing spawn");
|
||||
}
|
||||
}
|
||||
|
||||
/** The subset of {@code candidates} whose {@code model:} is currently turned off (fleetd #422). */
|
||||
private Set<String> modelOffProfiles(List<PlacementCandidate> candidates) {
|
||||
Set<String> off = models0().offIds();
|
||||
if (off.isEmpty()) {
|
||||
return Set.of();
|
||||
}
|
||||
return candidates.stream()
|
||||
.map(PlacementCandidate::profile)
|
||||
.filter(p -> {
|
||||
FleetConfig.Profile cfg = profiles0().get(p);
|
||||
return cfg != null && cfg.model() != null && off.contains(cfg.model());
|
||||
})
|
||||
.collect(Collectors.toSet());
|
||||
}
|
||||
|
||||
/**
|
||||
* The profile names {@code role} may be placed on, in definition order.
|
||||
*
|
||||
@@ -536,6 +643,32 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
return route(profileName).parityOverlay(profileName);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #422 follow-up: the single live read that answers both "is the models.allow: gate
|
||||
* armed" and "which models are off", off the exact same accessor ({@link #models0()}) {@link
|
||||
* #enforceModelEnabled} and {@link #modelOffProfiles} read — so {@code fleet_profiles}/{@code
|
||||
* GET /profiles} (via {@link PeerLauncher#disabledModels()}, which now delegates here) can
|
||||
* never report a different answer than the gate enforces (the fleetd #404 lesson), and the
|
||||
* startup log line built from this can never disagree with either.
|
||||
*
|
||||
* <p>{@link #models0()} itself normalizes a {@code null} {@link #models} read to the shared
|
||||
* {@link #NO_MODELS_CONFIGURED} sentinel — deliberately the one object no config-supplied
|
||||
* {@code Models} instance can ever be identical to, since it is private to this class — so
|
||||
* comparing by reference here recovers exactly the fact {@code models0()}'s normalization
|
||||
* would otherwise erase: whether the live source was {@code null} (no {@code models:} block,
|
||||
* armed = false) or a real, config-supplied block (armed = true, even one whose {@code allow:}
|
||||
* is itself empty or absent — {@link FleetConfig.Models}'s "absent or empty allow: is off"
|
||||
* wording governs config-load validation, a distinct question from whether this gate is armed
|
||||
* for reporting).
|
||||
*/
|
||||
@Override
|
||||
public PeerLauncher.ModelGateState modelGateState() {
|
||||
FleetConfig.Models m = models0();
|
||||
return m == NO_MODELS_CONFIGURED
|
||||
? PeerLauncher.ModelGateState.notConfigured()
|
||||
: PeerLauncher.ModelGateState.armed(m.offIds());
|
||||
}
|
||||
|
||||
@Override
|
||||
public void stop(String id) {
|
||||
HerdrPeerLauncher d = spawnedBy.get(id);
|
||||
|
||||
@@ -192,4 +192,66 @@ public interface PeerLauncher {
|
||||
* @return {@code true} when a reset was sent and its status transition must settle before reuse
|
||||
*/
|
||||
boolean clearContext(String id);
|
||||
|
||||
/**
|
||||
* Model ids the operator has currently turned off in the central {@code models.allow:} list
|
||||
* (fleetd #422) — empty for a launcher with nothing to gate against. {@code fleet_profiles}/
|
||||
* {@code GET /profiles} (via {@code FleetMcp.profilesView}) call this to report which models
|
||||
* are off, and MUST read this exact accessor rather than deriving their own answer: the fleetd
|
||||
* #404 lesson is that a status field reading a different source than the behaviour it describes
|
||||
* can drift from what the gate ({@code CompositePeerLauncher.enforceModelEnabled} and its
|
||||
* candidate filter) actually enforces. A default of {@code Set.of()} keeps every other {@link
|
||||
* PeerLauncher} implementer (the herdr adapters, and the two test-fake implementers) unchanged.
|
||||
*
|
||||
* <p>fleetd #422 follow-up: this alone cannot tell "no {@code models:} block at all" from "a
|
||||
* {@code models:} block where nothing is currently off" — both report an empty set here. Delegates
|
||||
* to {@link #modelGateState()} so the two facts always come from the one read {@link
|
||||
* #modelGateState()}'s implementer makes; do not override this method separately from that one.
|
||||
*/
|
||||
default Set<String> disabledModels() {
|
||||
return modelGateState().off();
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether the central {@code models.allow:} gate (fleetd #422) is armed at all, together with
|
||||
* which model ids are currently off — fleetd #422 follow-up. {@link #disabledModels()} alone
|
||||
* cannot distinguish two states that both report an empty set: a host with no {@code models:}
|
||||
* block (nothing is gated, and nothing can be) and a host WITH a {@code models:} block where
|
||||
* nothing is currently turned off (the gate is armed and reporting zero). This method exists so
|
||||
* a caller — the startup log, {@code fleet_profiles}/{@code GET /profiles} — can tell the two
|
||||
* apart, the same reason {@code CompletionResolver.UnsetMeaning} exists: an accessor that can
|
||||
* legitimately report "empty" must never let a caller guess why.
|
||||
*
|
||||
* <p>Default {@link ModelGateState#notConfigured()} — every launcher without a {@code models:}
|
||||
* block to read from (the herdr adapters, and the two test-fake implementers), matching {@link
|
||||
* #disabledModels()}'s own default of an empty set.
|
||||
*/
|
||||
default ModelGateState modelGateState() {
|
||||
return ModelGateState.notConfigured();
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #422 follow-up: the result of {@link #modelGateState()} — see that method's javadoc
|
||||
* for why "armed" and "off" must be reported together from one read rather than as two
|
||||
* separately-derived facts that a reload landing between them could make disagree.
|
||||
*
|
||||
* @param configured {@code true} when a {@code models:} block exists at all (armed), regardless
|
||||
* of whether anything in it is currently turned off; {@code false} when there
|
||||
* is no block to gate against
|
||||
* @param off the model ids currently turned off; always empty when {@code configured} is
|
||||
* {@code false}
|
||||
*/
|
||||
record ModelGateState(boolean configured, Set<String> off) {
|
||||
public ModelGateState {
|
||||
off = Set.copyOf(off);
|
||||
}
|
||||
|
||||
public static ModelGateState notConfigured() {
|
||||
return new ModelGateState(false, Set.of());
|
||||
}
|
||||
|
||||
public static ModelGateState armed(Set<String> off) {
|
||||
return new ModelGateState(true, off);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -5,14 +5,12 @@ import java.util.List;
|
||||
|
||||
/**
|
||||
* Backward-compatible placement: an unqualified spawn always resolves to the configured default
|
||||
* profile, exactly as {@code CompositePeerLauncher} did before CB-518. This ignores caps
|
||||
* ({@code maxLoad}) so that a pre-existing config behaves identically after upgrade — capacity
|
||||
* gating for automatic placement is deliberately out of scope for {@code fixed}, exactly as it
|
||||
* always has been. Reachability is a narrower exception (fleetd #315, below): a profile is never
|
||||
* checked for reachability up front, only skipped once it has already failed in <em>this same</em>
|
||||
* spawn call's retry loop — see the unreachable case below.
|
||||
* profile, exactly as {@code CompositePeerLauncher} did before CB-518. Reachability is a narrower
|
||||
* exception (fleetd #315, below): a profile is never checked for reachability up front, only
|
||||
* skipped once it has already failed in <em>this same</em> spawn call's retry loop — see the
|
||||
* unreachable case below.
|
||||
*
|
||||
* <p>Four exceptions walk past the default instead of returning it unconditionally:
|
||||
* <p>Six exceptions walk past the default instead of returning it unconditionally:
|
||||
* <ul>
|
||||
* <li>Quarantine (CB-578 stage B): a quarantined default is a credential that just refused on
|
||||
* a usage limit, not a transient capacity or reachability concern.
|
||||
@@ -20,6 +18,24 @@ import java.util.List;
|
||||
* ({@code BackendOutagePolicy}) — a separate, shorter-lived source from quarantine. When a
|
||||
* profile is both quarantined and cooling off, only the quarantine reason is reported
|
||||
* (exhaustion takes priority), matching {@code CompositePeerLauncher}'s explicit-spawn order.
|
||||
* <li>At cap (fleetd #435): a profile whose live count has reached its {@code maxLoad}
|
||||
* ({@link PlacementPolicyUtil#atCap}) — a documented, unconditional capacity limit (see
|
||||
* {@code FleetConfig.Profile#maxLoad}), so {@code fixed} must gate on it exactly as {@code
|
||||
* weighted}/{@code round-robin} already do via {@link PlacementPolicyUtil#available}. Before
|
||||
* this fix {@code fixed} built its own {@link PlacementCandidate} for the default with {@code
|
||||
* maxLoad} forced to {@code null}, so a capped default was chosen anyway on every unqualified
|
||||
* spawn — the cap was advisory, not enforced, for the one placement policy every config uses
|
||||
* by default. Reported only when quarantine and cooling off are both absent, matching {@code
|
||||
* CompositePeerLauncher}'s explicit-spawn check order (quarantine, then cooling off, then max
|
||||
* load, then model-off).
|
||||
* <li>Model off (fleetd #422): a profile whose {@code model:} the operator has turned off in
|
||||
* {@code models.allow:} — an operator decision, never a backend-reported outage, so it is a
|
||||
* fifth, independent source from quarantine, cooling off, and at-cap (never merged with any
|
||||
* of them), exactly as {@code CompositePeerLauncher.enforceModelEnabled} and {@link
|
||||
* PlacementPolicyUtil#available} treat it. When a profile is model-off <em>and</em> quarantined,
|
||||
* cooling off, or at cap, only the higher-priority reason is reported, matching {@code
|
||||
* CompositePeerLauncher}'s explicit-spawn check order (quarantine, then cooling off, then max
|
||||
* load, then model-off).
|
||||
* <li>Unreachable (fleetd #315): {@code CompositePeerLauncher.spawn} retries a failed candidate
|
||||
* on the next one and rebuilds the {@link PlacementContext} so {@code ctx.unreachable()}
|
||||
* names every profile that already failed with {@code PeerUnreachableException} in this same
|
||||
@@ -32,9 +48,10 @@ import java.util.List;
|
||||
* {@code weighted}/{@code round-robin} skip it — an explicit {@code fleet_spawn} naming
|
||||
* the profile is unaffected, only this automatic fallback walk.
|
||||
* </ul>
|
||||
* A fleet where nothing is ever quarantined, cooling off, unreachable, or weight-0 never exercises
|
||||
* any of these paths, so today's behaviour is unchanged — in particular, the very first selection
|
||||
* of a spawn call always sees an empty {@code unreachable} set, so the first choice is untouched.
|
||||
* A fleet where nothing is ever quarantined, cooling off, at cap, model-off, unreachable, or
|
||||
* weight-0 never exercises any of these paths, so today's behaviour is unchanged — in particular,
|
||||
* the very first selection of a spawn call always sees an empty {@code unreachable} set, so the
|
||||
* first choice is untouched.
|
||||
*/
|
||||
final class FixedPlacementPolicy implements PlacementPolicy {
|
||||
|
||||
@@ -42,12 +59,14 @@ final class FixedPlacementPolicy implements PlacementPolicy {
|
||||
public PlacementCandidate select(PlacementContext ctx) {
|
||||
String d = ctx.defaultProfile();
|
||||
if (d != null && !d.isBlank() && !ctx.quarantined().contains(d) && !ctx.coolingOff().contains(d)
|
||||
&& !ctx.unreachable().contains(d) && !weightExcluded(ctx, d)) {
|
||||
&& !ctx.modelOff().contains(d) && !ctx.unreachable().contains(d) && !weightExcluded(ctx, d)
|
||||
&& !capExcluded(ctx, d)) {
|
||||
return new PlacementCandidate(d, null, 1.0f, null);
|
||||
}
|
||||
for (PlacementCandidate c : ctx.candidates()) {
|
||||
if (!ctx.quarantined().contains(c.profile()) && !ctx.coolingOff().contains(c.profile())
|
||||
&& !ctx.unreachable().contains(c.profile()) && !c.excluded()) {
|
||||
&& !ctx.modelOff().contains(c.profile()) && !ctx.unreachable().contains(c.profile())
|
||||
&& !c.excluded() && !PlacementPolicyUtil.atCap(ctx, c)) {
|
||||
return new PlacementCandidate(c.profile(), null, c.weight(), c.maxLoad());
|
||||
}
|
||||
}
|
||||
@@ -56,9 +75,18 @@ final class FixedPlacementPolicy implements PlacementPolicy {
|
||||
// Exhaustion quarantine takes priority: reported only when quarantine is absent, so the
|
||||
// message never claims "cooling off" for a profile that is really backend-exhausted.
|
||||
boolean dCoolingOff = !dQuarantined && ctx.coolingOff().contains(d);
|
||||
// fleetd #435: at-cap sits between cooling off and model-off, matching
|
||||
// CompositePeerLauncher's explicit-spawn check order (quarantine, cooling off, max load,
|
||||
// then model-off) — reported only when quarantine/cooling-off are both absent.
|
||||
boolean dAtCap = !dQuarantined && !dCoolingOff && capExcluded(ctx, d);
|
||||
// fleetd #422: model-off is a fifth, independent source (an operator decision) — but
|
||||
// quarantine/cooling-off/at-cap still take priority when more than one applies, matching
|
||||
// CompositePeerLauncher's explicit-spawn check order (quarantine, cooling off, max load,
|
||||
// then model-off).
|
||||
boolean dModelOff = !dQuarantined && !dCoolingOff && !dAtCap && ctx.modelOff().contains(d);
|
||||
boolean dUnreachable = ctx.unreachable().contains(d);
|
||||
boolean dWeightExcluded = weightExcluded(ctx, d);
|
||||
if (dQuarantined || dCoolingOff || dUnreachable || dWeightExcluded) {
|
||||
if (dQuarantined || dCoolingOff || dAtCap || dModelOff || dUnreachable || dWeightExcluded) {
|
||||
List<String> reasons = new ArrayList<>();
|
||||
if (dQuarantined) {
|
||||
reasons.add("is quarantined (backend exhausted)");
|
||||
@@ -66,6 +94,14 @@ final class FixedPlacementPolicy implements PlacementPolicy {
|
||||
if (dCoolingOff) {
|
||||
reasons.add("is cooling off after repeated backend errors");
|
||||
}
|
||||
if (dAtCap) {
|
||||
PlacementCandidate c = candidateFor(ctx, d);
|
||||
int live = ctx.liveCount().apply(d);
|
||||
reasons.add("is at maxLoad (" + live + " live >= " + c.maxLoad() + " cap)");
|
||||
}
|
||||
if (dModelOff) {
|
||||
reasons.add("names a model the operator has turned off in models.allow");
|
||||
}
|
||||
if (dUnreachable) {
|
||||
reasons.add("is unreachable");
|
||||
}
|
||||
@@ -78,18 +114,38 @@ final class FixedPlacementPolicy implements PlacementPolicy {
|
||||
}
|
||||
if (!ctx.candidates().isEmpty()) {
|
||||
throw new PlacementException("all worker profiles are excluded from automatic "
|
||||
+ "selection (quarantined, cooling off, unreachable, or weight-0)");
|
||||
+ "selection (quarantined, cooling off, at cap, model-off, unreachable, or weight-0)");
|
||||
}
|
||||
throw new PlacementException("no worker profiles configured");
|
||||
}
|
||||
|
||||
/** Whether {@code profile} carries {@code weight <= 0} (CB-554) among {@code ctx}'s candidates. */
|
||||
private static boolean weightExcluded(PlacementContext ctx, String profile) {
|
||||
PlacementCandidate c = candidateFor(ctx, profile);
|
||||
return c != null && c.excluded();
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether {@code profile} has reached its {@code maxLoad} cap (fleetd #435), using the shared
|
||||
* {@link PlacementPolicyUtil#atCap} definition — the same one {@code weighted}/{@code
|
||||
* round-robin} already consult via {@link PlacementPolicyUtil#available}. Looked up by name,
|
||||
* the same way {@link #weightExcluded} is: the default fast path above builds its own {@link
|
||||
* PlacementCandidate} with {@code maxLoad} forced to {@code null} (it carries no cap of its
|
||||
* own), so the candidate actually configured for {@code profile} has to be found in {@code
|
||||
* ctx.candidates()} first.
|
||||
*/
|
||||
private static boolean capExcluded(PlacementContext ctx, String profile) {
|
||||
PlacementCandidate c = candidateFor(ctx, profile);
|
||||
return c != null && PlacementPolicyUtil.atCap(ctx, c);
|
||||
}
|
||||
|
||||
/** The configured candidate named {@code profile} in {@code ctx}, or {@code null} if none. */
|
||||
private static PlacementCandidate candidateFor(PlacementContext ctx, String profile) {
|
||||
for (PlacementCandidate c : ctx.candidates()) {
|
||||
if (c.profile().equals(profile)) {
|
||||
return c.excluded();
|
||||
return c;
|
||||
}
|
||||
}
|
||||
return false;
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -22,11 +22,29 @@ import java.util.function.Function;
|
||||
* the two apart so its refusal message says "cooling off", not "exhausted",
|
||||
* when only this one is active. A profile can be in both sets at once; when it
|
||||
* is, exhaustion quarantine is reported (it takes priority).
|
||||
* @param modelOff profiles whose {@code model:} is currently turned off in {@code
|
||||
* models.allow:} (fleetd #422) — an operator decision, not a backend-reported
|
||||
* outage, so a SEPARATE, independent source from both {@code quarantined} and
|
||||
* {@code coolingOff}. A profile can be in this set together with either (or
|
||||
* both) of the others; {@link PlacementPolicyUtil} counts it into its own
|
||||
* bucket rather than merging it into theirs, the same reason
|
||||
* {@code coolingOff} is kept apart from {@code quarantined}.
|
||||
*/
|
||||
public record PlacementContext(String defaultProfile,
|
||||
List<PlacementCandidate> candidates,
|
||||
Function<String, Integer> liveCount,
|
||||
Set<String> unreachable,
|
||||
Set<String> quarantined,
|
||||
Set<String> coolingOff) {
|
||||
Set<String> coolingOff,
|
||||
Set<String> modelOff) {
|
||||
|
||||
/**
|
||||
* Back-compat form before the fleetd #422 model on/off gate was added — no candidate's model
|
||||
* is off. Keeps pre-#422 call sites (tests included) compiling and behaving identically.
|
||||
*/
|
||||
public PlacementContext(String defaultProfile, List<PlacementCandidate> candidates,
|
||||
Function<String, Integer> liveCount, Set<String> unreachable,
|
||||
Set<String> quarantined, Set<String> coolingOff) {
|
||||
this(defaultProfile, candidates, liveCount, unreachable, quarantined, coolingOff, Set.of());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -11,28 +11,40 @@ final class PlacementPolicyUtil {
|
||||
private PlacementPolicyUtil() {
|
||||
}
|
||||
|
||||
/**
|
||||
* True when {@code c} has reached its {@code maxLoad} cap: {@code liveCount(c.profile()) >=
|
||||
* c.maxLoad()}. A {@code null} maxLoad means unlimited, so it is never at cap.
|
||||
*
|
||||
* <p>Extracted as the single shared definition of "at cap" (fleetd #435): before this fix it
|
||||
* was computed inline in both {@link #available} and {@link #emptyException}, and {@code
|
||||
* FixedPlacementPolicy} — not a caller of either — quietly kept its own {@code select} free of
|
||||
* any cap check at all, so a capped default profile was chosen anyway under the default
|
||||
* placement policy. Every automatic policy must call this, not re-derive it.
|
||||
*/
|
||||
static boolean atCap(PlacementContext ctx, PlacementCandidate c) {
|
||||
Integer cap = c.maxLoad();
|
||||
return cap != null && ctx.liveCount().apply(c.profile()) >= cap;
|
||||
}
|
||||
|
||||
/**
|
||||
* Candidates that are not weight-excluded (CB-554: explicit {@code weight <= 0}, checked
|
||||
* first because it is a static config choice rather than transient state), not
|
||||
* known-unreachable, not quarantined (CB-578 stage B), not cooling off after repeated backend
|
||||
* errors (fleetd #201 Unit 5 — a separate, shorter-lived source from quarantine), and have not
|
||||
* reached their maxLoad. A {@code null} maxLoad means unlimited.
|
||||
* errors (fleetd #201 Unit 5 — a separate, shorter-lived source from quarantine), not naming a
|
||||
* model the operator has turned off (fleetd #422 — a third, independent source: an operator
|
||||
* decision, never a backend-reported outage), and have not reached their maxLoad (see {@link
|
||||
* #atCap}). A {@code null} maxLoad means unlimited.
|
||||
*/
|
||||
static List<PlacementCandidate> available(PlacementContext ctx) {
|
||||
List<PlacementCandidate> out = new ArrayList<>();
|
||||
for (PlacementCandidate c : ctx.candidates()) {
|
||||
if (c.excluded() || ctx.unreachable().contains(c.profile())
|
||||
|| ctx.quarantined().contains(c.profile())
|
||||
|| ctx.coolingOff().contains(c.profile())) {
|
||||
|| ctx.coolingOff().contains(c.profile())
|
||||
|| ctx.modelOff().contains(c.profile())
|
||||
|| atCap(ctx, c)) {
|
||||
continue;
|
||||
}
|
||||
Integer cap = c.maxLoad();
|
||||
if (cap != null) {
|
||||
int live = ctx.liveCount().apply(c.profile());
|
||||
if (live >= cap) {
|
||||
continue;
|
||||
}
|
||||
}
|
||||
out.add(c);
|
||||
}
|
||||
return out;
|
||||
@@ -40,12 +52,13 @@ final class PlacementPolicyUtil {
|
||||
|
||||
/**
|
||||
* Build a clear exception describing why every candidate was dropped: all weight-0, all
|
||||
* quarantined, all cooling off, all at capacity, all unreachable, or a mix. Each candidate is
|
||||
* counted into exactly one bucket (weight-excluded first, then quarantined, then cooling off)
|
||||
* so a candidate excluded for more than one reason is never double-counted — a candidate that is
|
||||
* both quarantined (CB-578 stage B, backend exhausted) and cooling off (fleetd #201 Unit 5,
|
||||
* repeated backend errors) counts only as quarantined, matching {@code CompositePeerLauncher}'s
|
||||
* explicit-spawn ordering: exhaustion quarantine takes priority when both are active.
|
||||
* quarantined, all cooling off, all model-off, all at capacity, all unreachable, or a mix. Each
|
||||
* candidate is counted into exactly one bucket (weight-excluded first, then quarantined, then
|
||||
* cooling off, then model-off) so a candidate excluded for more than one reason is never
|
||||
* double-counted — a candidate that is both quarantined (CB-578 stage B, backend exhausted) and
|
||||
* cooling off (fleetd #201 Unit 5, repeated backend errors) counts only as quarantined, matching
|
||||
* {@code CompositePeerLauncher}'s explicit-spawn ordering: exhaustion quarantine takes priority
|
||||
* when more than one applies.
|
||||
*/
|
||||
static PlacementException emptyException(PlacementContext ctx) {
|
||||
int weightExcluded = 0;
|
||||
@@ -53,17 +66,19 @@ final class PlacementPolicyUtil {
|
||||
int unreachable = 0;
|
||||
int quarantined = 0;
|
||||
int coolingOff = 0;
|
||||
int modelOff = 0;
|
||||
for (PlacementCandidate c : ctx.candidates()) {
|
||||
Integer cap = c.maxLoad();
|
||||
if (c.excluded()) {
|
||||
weightExcluded++;
|
||||
} else if (ctx.quarantined().contains(c.profile())) {
|
||||
quarantined++;
|
||||
} else if (ctx.coolingOff().contains(c.profile())) {
|
||||
coolingOff++;
|
||||
} else if (ctx.modelOff().contains(c.profile())) {
|
||||
modelOff++;
|
||||
} else if (ctx.unreachable().contains(c.profile())) {
|
||||
unreachable++;
|
||||
} else if (cap != null && ctx.liveCount().apply(c.profile()) >= cap) {
|
||||
} else if (atCap(ctx, c)) {
|
||||
atCap++;
|
||||
}
|
||||
}
|
||||
@@ -83,6 +98,10 @@ final class PlacementPolicyUtil {
|
||||
return new PlacementException(
|
||||
"all worker profiles are cooling off after repeated backend errors");
|
||||
}
|
||||
if (modelOff == total) {
|
||||
return new PlacementException(
|
||||
"all worker profiles name a model the operator has turned off");
|
||||
}
|
||||
if (atCap == total) {
|
||||
return new PlacementException("all worker profiles are at maxLoad");
|
||||
}
|
||||
@@ -92,8 +111,9 @@ final class PlacementPolicyUtil {
|
||||
return new PlacementException("no worker profile available: " + atCap + " at maxLoad, "
|
||||
+ unreachable + " unreachable, " + quarantined + " quarantined, "
|
||||
+ coolingOff + " cooling off, "
|
||||
+ modelOff + " model-off, "
|
||||
+ weightExcluded + " weight-0, "
|
||||
+ (total - atCap - unreachable - quarantined - coolingOff - weightExcluded)
|
||||
+ (total - atCap - unreachable - quarantined - coolingOff - modelOff - weightExcluded)
|
||||
+ " remaining");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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.");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,62 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.peer.PeerLauncher;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.Set;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotEquals;
|
||||
|
||||
/**
|
||||
* fleetd #422 follow-up: {@code Fleetd.modelGateCoverageLine} is the startup-log counterpart of
|
||||
* {@code exhaustedPatternCoverageLine}/{@code errorPatternCoverageLine} — see {@code
|
||||
* FleetdPatternCoverageLineTest} for the identical shape this follows — except here there is no
|
||||
* {@code UnsetMeaning} choice for a caller to get backwards: {@link
|
||||
* PeerLauncher.ModelGateState#configured()} already states, unambiguously, whether an empty
|
||||
* {@link PeerLauncher.ModelGateState#off()} means "no {@code models:} block to gate with at all"
|
||||
* or "a block armed and currently reporting zero off". This class proves {@code
|
||||
* modelGateCoverageLine} words those two states — plus the third, N off — distinctly, so a
|
||||
* mutation that made it ignore {@code configured()} either way is caught here.
|
||||
*/
|
||||
class FleetdModelGateCoverageLineTest {
|
||||
|
||||
@Test
|
||||
@DisplayName("no models: block reports not configured, distinct from armed-with-zero")
|
||||
void noModelsBlockReportsNotConfigured() {
|
||||
String line = Fleetd.modelGateCoverageLine(PeerLauncher.ModelGateState.notConfigured());
|
||||
assertEquals("not configured (no models: block — nothing is gated, and nothing can be)", line);
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a models: block armed with nothing off reports armed, distinct from not configured")
|
||||
void armedWithNothingOffReportsArmed() {
|
||||
String line = Fleetd.modelGateCoverageLine(PeerLauncher.ModelGateState.armed(Set.of()));
|
||||
assertEquals("armed (models: block present; 0 models currently turned off)", line);
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a models: block with N off names the off models")
|
||||
void armedWithModelsOffNamesThem() {
|
||||
String line = Fleetd.modelGateCoverageLine(
|
||||
PeerLauncher.ModelGateState.armed(Set.of("deepseek-v4-flash", "claude-opus-9000")));
|
||||
assertEquals("armed (2 model(s) turned off: [claude-opus-9000, deepseek-v4-flash])", line);
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("the three states produce pairwise-distinct wording for the same empty-looking input")
|
||||
void theThreeStatesProduceDistinctWording() {
|
||||
String notConfigured = Fleetd.modelGateCoverageLine(PeerLauncher.ModelGateState.notConfigured());
|
||||
String armedZero = Fleetd.modelGateCoverageLine(PeerLauncher.ModelGateState.armed(Set.of()));
|
||||
String armedOne = Fleetd.modelGateCoverageLine(PeerLauncher.ModelGateState.armed(Set.of("x")));
|
||||
|
||||
// Pinned individually above; restated here so this test alone still catches a regression
|
||||
// even if one of the three tests above were ever deleted — the exact FleetdPatternCoverageLineTest
|
||||
// pattern, adapted from "two keys" to "three states of one gate".
|
||||
assertNotEquals(notConfigured, armedZero,
|
||||
"collapsing 'no models: block' into 'armed, zero off' is fleetd #422 follow-up's exact defect");
|
||||
assertNotEquals(armedZero, armedOne);
|
||||
assertNotEquals(notConfigured, armedOne);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,67 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.Set;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotEquals;
|
||||
|
||||
/**
|
||||
* fleetd #415 (review follow-up): {@code CompletionResolverTest} proves {@code coverage()} words
|
||||
* {@code UnsetMeaning.OFF} and {@code UnsetMeaning.BUILT_IN_DEFAULT} correctly — but every one of
|
||||
* those tests supplies the meaning itself. That proves the enum's wording, never that {@code
|
||||
* Fleetd} pairs the right meaning with the right pattern key. That pairing is #415's actual
|
||||
* defect: {@code coverage()} had no way to know what unset meant for its key, so the fix moved
|
||||
* the fact to the caller — and nothing yet proved the caller states it correctly.
|
||||
*
|
||||
* <p><b>Measured or it didn't happen:</b> swapping the two {@code UnsetMeaning} arguments at
|
||||
* {@code Fleetd}'s two coverage call sites — giving {@code exhaustedPattern} the built-in-default
|
||||
* wording and {@code errorPattern} the off wording, #415's exact defect with the keys exchanged —
|
||||
* compiled with 0 errors and left all 1506 existing tests green. This class exists to turn that
|
||||
* swap red.
|
||||
*
|
||||
* <p>It calls {@link Fleetd#exhaustedPatternCoverageLine} and {@link Fleetd#errorPatternCoverageLine}
|
||||
* directly rather than reading {@code Fleetd.java} as source text (the shape {@code
|
||||
* FleetdCompletionResolverWiringTest} uses for a different wiring gap): those two methods are the
|
||||
* extracted call sites {@code main} actually invokes, following the same {@code static} factory +
|
||||
* dedicated-test pattern as {@link Fleetd#capacitySource} and {@link Fleetd#worktreeBranchLookup}.
|
||||
*/
|
||||
class FleetdPatternCoverageLineTest {
|
||||
|
||||
private static final Set<String> PROFILES = Set.of("terra", "gx10");
|
||||
|
||||
@Test
|
||||
@DisplayName("exhaustedPatternCoverageLine says off when no profile configures exhaustedPattern")
|
||||
void exhaustedPatternCoverageLineSaysOffWhenNoProfileConfiguresIt() {
|
||||
assertEquals("off (no profile has an exhaustedPattern configured; profiles: [gx10, terra])",
|
||||
Fleetd.exhaustedPatternCoverageLine(PROFILES, Set.of()));
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("errorPatternCoverageLine says built-in default when no profile configures errorPattern")
|
||||
void errorPatternCoverageLineSaysBuiltInDefaultWhenNoProfileConfiguresIt() {
|
||||
assertEquals("built-in default for all profiles (no profile customises errorPattern; "
|
||||
+ "profiles: [gx10, terra])",
|
||||
Fleetd.errorPatternCoverageLine(PROFILES, Set.of()));
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("the two keys produce different wording for the identical empty-coverage input")
|
||||
void theTwoKeysProduceDifferentWordingForTheSameEmptyInput() {
|
||||
String exhaustedLine = Fleetd.exhaustedPatternCoverageLine(PROFILES, Set.of());
|
||||
String errorLine = Fleetd.errorPatternCoverageLine(PROFILES, Set.of());
|
||||
|
||||
// Pinned individually above; restated here so this test alone still catches a swap even
|
||||
// if one of the two tests above were ever deleted.
|
||||
assertEquals("off (no profile has an exhaustedPattern configured; profiles: [gx10, terra])",
|
||||
exhaustedLine);
|
||||
assertEquals("built-in default for all profiles (no profile customises errorPattern; "
|
||||
+ "profiles: [gx10, terra])", errorLine);
|
||||
assertNotEquals(exhaustedLine, errorLine,
|
||||
"swapping which UnsetMeaning pairs with which pattern key at Fleetd's call sites "
|
||||
+ "must be caught here — that pairing, not coverage()'s own wording in isolation, "
|
||||
+ "is fleetd #415's actual defect");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,360 @@
|
||||
package dev.ltms.fleet.auth;
|
||||
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.PaneLocator;
|
||||
import dev.ltms.fleet.mcp.ConnectionIdentity;
|
||||
import dev.ltms.fleet.peer.MemberRole;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.Map;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
/**
|
||||
* fleetd #424 — revoking (or granting) an architect slot must take effect on the next spawn with
|
||||
* no restart. A session already bound to a slot keeps its <em>binding</em> (the {@code
|
||||
* terminalToSlot} occupancy) even after that slot drops out of config, but NOT the ARCHITECT
|
||||
* <em>privilege</em> the slot used to grant — that is revoked on the bound session's very next
|
||||
* request. See {@link MemberRegistry}'s class doc for the exact rule: "config governs what a
|
||||
* bound slot still grants, as well as what may be bound next."
|
||||
*
|
||||
* <p>Every test here drives a REAL {@link ConfigRef#reload()} against a {@code @TempDir} file and
|
||||
* asserts {@link ConfigRef.Outcome#applied()}, rather than comparing two frozen
|
||||
* {@code MemberRegistry} instances in memory — the defect this ticket fixes is specifically that
|
||||
* {@link MemberRegistry} used to ignore a live reload, so a test that never reloads cannot tell the
|
||||
* fixed registry from the broken one. {@code requireSlotFor} and {@code reserve} are pinned in
|
||||
* separate tests, in both directions (removed and added), so a registry that simply refuses (or
|
||||
* simply allows) everything cannot pass by accident — see {@link MemberRegistryTest} for the
|
||||
* registry's other invariants (bind/unbind cardinality, thread-safety), which are unaffected by
|
||||
* this ticket and still exercised against the frozen constructor.
|
||||
*/
|
||||
class MemberRegistryLiveTest {
|
||||
|
||||
private static String yaml(String fleetBlock) {
|
||||
return """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
herdrSocket: ~/.config/herdr/herdr.sock
|
||||
profiles:
|
||||
sonnet:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
model: sonnet
|
||||
opus:
|
||||
baseUrl: http://gx00.gw:8001
|
||||
model: opus
|
||||
guard:
|
||||
offSubscriptionHosts:
|
||||
- gx00.gw
|
||||
""" + fleetBlock;
|
||||
}
|
||||
|
||||
private static final String WITH_SONNET_SLOT = """
|
||||
fleet:
|
||||
architects:
|
||||
designer:
|
||||
profile: sonnet
|
||||
""";
|
||||
|
||||
/** No architect pool at all — developers is unrelated dead data for this registry (#424 out of scope). */
|
||||
private static final String WITHOUT_ARCHITECT_SLOTS = """
|
||||
fleet:
|
||||
developers:
|
||||
dev1:
|
||||
profile: sonnet
|
||||
""";
|
||||
|
||||
/** Same slot name ({@code designer}) as {@link #WITH_SONNET_SLOT}, repointed to a different profile. */
|
||||
private static final String WITH_OPUS_SLOT = """
|
||||
fleet:
|
||||
architects:
|
||||
designer:
|
||||
profile: opus
|
||||
""";
|
||||
|
||||
/** Same profile ({@code sonnet}) as {@link #WITH_SONNET_SLOT}, but the pool key is renamed. */
|
||||
private static final String WITH_RENAMED_SLOT = """
|
||||
fleet:
|
||||
architects:
|
||||
architect-lead:
|
||||
profile: sonnet
|
||||
""";
|
||||
|
||||
private static ConfigRef refFor(Path f) {
|
||||
return new ConfigRef(f, FleetConfig.load(f));
|
||||
}
|
||||
|
||||
// ── requireSlotFor is live (criteria 1, 2, 4) ──────────────────────────────────────────────
|
||||
|
||||
@Test
|
||||
void requireSlotForRefusesAProfileWhoseSlotWasRemovedByReload(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, yaml(WITH_SONNET_SLOT));
|
||||
ConfigRef ref = refFor(f);
|
||||
MemberRegistry registry = MemberRegistry.live(() -> ref.get().fleet());
|
||||
|
||||
assertDoesNotThrow(() -> registry.requireSlotFor(MemberRole.ARCHITECT, "sonnet"),
|
||||
"the slot is configured before the reload");
|
||||
|
||||
Files.writeString(f, yaml(WITHOUT_ARCHITECT_SLOTS));
|
||||
ConfigRef.Outcome out = ref.reload();
|
||||
assertTrue(out.applied(), "the reload must actually take effect: " + out.summary());
|
||||
|
||||
assertThrows(IllegalArgumentException.class,
|
||||
() -> registry.requireSlotFor(MemberRole.ARCHITECT, "sonnet"),
|
||||
"revoking the slot must refuse the NEXT spawn that names it");
|
||||
}
|
||||
|
||||
@Test
|
||||
void requireSlotForAllowsAProfileWhoseSlotWasAddedByReload(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, yaml(WITHOUT_ARCHITECT_SLOTS));
|
||||
ConfigRef ref = refFor(f);
|
||||
MemberRegistry registry = MemberRegistry.live(() -> ref.get().fleet());
|
||||
|
||||
assertThrows(IllegalArgumentException.class,
|
||||
() -> registry.requireSlotFor(MemberRole.ARCHITECT, "sonnet"),
|
||||
"no architect slot is configured yet");
|
||||
|
||||
Files.writeString(f, yaml(WITH_SONNET_SLOT));
|
||||
ConfigRef.Outcome out = ref.reload();
|
||||
assertTrue(out.applied(), "the reload must actually take effect: " + out.summary());
|
||||
|
||||
assertDoesNotThrow(() -> registry.requireSlotFor(MemberRole.ARCHITECT, "sonnet"),
|
||||
"a slot added by reload must be usable with no restart");
|
||||
}
|
||||
|
||||
// ── reserve is live too — tested separately from requireSlotFor (criteria 1, 2, 4) ────────
|
||||
|
||||
@Test
|
||||
void reserveRefusesAProfileWhoseSlotWasRemovedByReload(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, yaml(WITH_SONNET_SLOT));
|
||||
ConfigRef ref = refFor(f);
|
||||
MemberRegistry registry = MemberRegistry.live(() -> ref.get().fleet());
|
||||
|
||||
MemberLifecycle.SlotReservation before = registry.reserve(MemberRole.ARCHITECT, "sonnet");
|
||||
assertEquals("architect:designer", before.slot());
|
||||
registry.release(before); // free it back up so the reload-side reserve below starts clean
|
||||
|
||||
Files.writeString(f, yaml(WITHOUT_ARCHITECT_SLOTS));
|
||||
ConfigRef.Outcome out = ref.reload();
|
||||
assertTrue(out.applied(), "the reload must actually take effect: " + out.summary());
|
||||
|
||||
assertThrows(IllegalArgumentException.class,
|
||||
() -> registry.reserve(MemberRole.ARCHITECT, "sonnet"),
|
||||
"revoking the slot must refuse the NEXT reservation for it");
|
||||
}
|
||||
|
||||
@Test
|
||||
void reserveAllowsAProfileWhoseSlotWasAddedByReload(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, yaml(WITHOUT_ARCHITECT_SLOTS));
|
||||
ConfigRef ref = refFor(f);
|
||||
MemberRegistry registry = MemberRegistry.live(() -> ref.get().fleet());
|
||||
|
||||
assertThrows(IllegalArgumentException.class,
|
||||
() -> registry.reserve(MemberRole.ARCHITECT, "sonnet"),
|
||||
"no architect slot is configured yet");
|
||||
|
||||
Files.writeString(f, yaml(WITH_SONNET_SLOT));
|
||||
ConfigRef.Outcome out = ref.reload();
|
||||
assertTrue(out.applied(), "the reload must actually take effect: " + out.summary());
|
||||
|
||||
MemberLifecycle.SlotReservation after = registry.reserve(MemberRole.ARCHITECT, "sonnet");
|
||||
assertEquals("architect:designer", after.slot(),
|
||||
"a slot added by reload must be reservable with no restart");
|
||||
}
|
||||
|
||||
// ── profileForSlot, isSlot and nameForSlot are live too (fleetd #431) ─────────────────────────
|
||||
// #424 pinned roleForSlot against a live reload but left these three untested — proved by
|
||||
// mutating each to read a snapshot flattened once at construction: the full suite stayed green
|
||||
// for all three.
|
||||
//
|
||||
// The three differ in how much production behaviour depends on them, and the ticket first got
|
||||
// this ranking wrong. nameForSlot is wired: CallerResolver passes members::nameForSlot, next to
|
||||
// members::roleForSlot. isSlot is reached through bind, which calls it to refuse an unknown
|
||||
// slot. profileForSlot has NO caller in src/main at all — grepped both the ".profileForSlot("
|
||||
// and the "::profileForSlot" form — so there is no seam to drive its test through beyond the
|
||||
// accessor itself, and its own javadoc ("what the spawn lifecycle reads") describes a caller
|
||||
// that does not exist. These tests pin the accessors as they are; whether profileForSlot should
|
||||
// be wired or deleted is a separate question.
|
||||
|
||||
@Test
|
||||
void profileForSlotReflectsAProfileChangedByReload(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, yaml(WITH_SONNET_SLOT));
|
||||
ConfigRef ref = refFor(f);
|
||||
MemberRegistry registry = MemberRegistry.live(() -> ref.get().fleet());
|
||||
|
||||
assertEquals("sonnet", registry.profileForSlot("architect:designer"),
|
||||
"the slot's profile before the reload");
|
||||
|
||||
Files.writeString(f, yaml(WITH_OPUS_SLOT));
|
||||
ConfigRef.Outcome out = ref.reload();
|
||||
assertTrue(out.applied(), "the reload must actually take effect: " + out.summary());
|
||||
|
||||
assertEquals("opus", registry.profileForSlot("architect:designer"),
|
||||
"repointing the slot to a different profile must take effect with no restart");
|
||||
}
|
||||
|
||||
@Test
|
||||
void isSlotStopsReportingASlotRemovedByReload(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, yaml(WITH_SONNET_SLOT));
|
||||
ConfigRef ref = refFor(f);
|
||||
MemberRegistry registry = MemberRegistry.live(() -> ref.get().fleet());
|
||||
|
||||
assertTrue(registry.isSlot("architect:designer"), "the slot is configured before the reload");
|
||||
|
||||
Files.writeString(f, yaml(WITHOUT_ARCHITECT_SLOTS));
|
||||
ConfigRef.Outcome out = ref.reload();
|
||||
assertTrue(out.applied(), "the reload must actually take effect: " + out.summary());
|
||||
|
||||
assertFalse(registry.isSlot("architect:designer"),
|
||||
"removing the slot from config must make isSlot say so on the very next call, "
|
||||
+ "with no restart");
|
||||
}
|
||||
|
||||
@Test
|
||||
void isSlotStartsReportingASlotAddedByReload(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, yaml(WITHOUT_ARCHITECT_SLOTS));
|
||||
ConfigRef ref = refFor(f);
|
||||
MemberRegistry registry = MemberRegistry.live(() -> ref.get().fleet());
|
||||
|
||||
assertFalse(registry.isSlot("architect:designer"), "no architect slot is configured yet");
|
||||
|
||||
Files.writeString(f, yaml(WITH_SONNET_SLOT));
|
||||
ConfigRef.Outcome out = ref.reload();
|
||||
assertTrue(out.applied(), "the reload must actually take effect: " + out.summary());
|
||||
|
||||
assertTrue(registry.isSlot("architect:designer"),
|
||||
"a slot added by reload must be visible to isSlot with no restart");
|
||||
}
|
||||
|
||||
@Test
|
||||
void nameForSlotReflectsANameChangedByReload(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, yaml(WITH_SONNET_SLOT));
|
||||
ConfigRef ref = refFor(f);
|
||||
MemberRegistry registry = MemberRegistry.live(() -> ref.get().fleet());
|
||||
|
||||
assertEquals("designer", registry.nameForSlot("architect:designer"),
|
||||
"the configured name before the reload");
|
||||
|
||||
Files.writeString(f, yaml(WITH_RENAMED_SLOT));
|
||||
ConfigRef.Outcome out = ref.reload();
|
||||
assertTrue(out.applied(), "the reload must actually take effect: " + out.summary());
|
||||
|
||||
assertNull(registry.nameForSlot("architect:designer"),
|
||||
"the old key no longer names a configured slot — it was renamed away by the reload");
|
||||
assertEquals("architect-lead", registry.nameForSlot("architect:architect-lead"),
|
||||
"the new name must be visible under its new qualified key with no restart");
|
||||
}
|
||||
|
||||
// ── a bound architect is demoted, but the binding itself is not touched (fleetd #424) ───────
|
||||
// The lead's corrected ruling: the PRIVILEGE a slot grants is revoked on the bound session's
|
||||
// very next request, but the terminalToSlot BINDING itself is untouched by a reload — dropping
|
||||
// it would double-book the slot key and break unbind's compare-safe contract. See the class
|
||||
// doc's binding rule.
|
||||
|
||||
/** A caller identity resolving the one canned pane (terminal {@code term_a}) in {@link FakeHerdr}. */
|
||||
private static ConnectionIdentity boundPaneIdentity() {
|
||||
return new ConnectionIdentity(new PaneLocator(new FakeHerdr()), _ -> FakeHerdr.WORKER_PID);
|
||||
}
|
||||
|
||||
@Test
|
||||
void anArchitectAlreadyBoundToASlotIsDemotedByReload(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, yaml(WITH_SONNET_SLOT));
|
||||
ConfigRef ref = refFor(f);
|
||||
MemberRegistry registry = MemberRegistry.live(() -> ref.get().fleet());
|
||||
|
||||
MemberLifecycle.SlotReservation reservation = registry.reserve(MemberRole.ARCHITECT, "sonnet");
|
||||
assertTrue(registry.bind(reservation, "term_a"));
|
||||
assertEquals("architect:designer", registry.slotForTerminal("term_a"));
|
||||
|
||||
// Drive the real caller path, not the roleForSlot seam directly: CallerResolver.resolve is
|
||||
// what a live request actually goes through (CallerResolver.java:220), and a resolver that
|
||||
// ignored roleForSlot entirely would still pass a test that only checked the seam.
|
||||
CallerResolver resolver = CallerResolver.withLeadsAndMembers(
|
||||
boundPaneIdentity(), false, null, Map::of, registry);
|
||||
|
||||
Principal before = resolver.resolve("127.0.0.1", 42, null);
|
||||
assertEquals(Role.ARCHITECT, before.role(), "sanity check: the harness binds term_a as an architect");
|
||||
assertEquals("designer", before.name());
|
||||
|
||||
Files.writeString(f, yaml(WITHOUT_ARCHITECT_SLOTS));
|
||||
ConfigRef.Outcome out = ref.reload();
|
||||
assertTrue(out.applied(), "the reload must actually take effect: " + out.summary());
|
||||
|
||||
Principal after = resolver.resolve("127.0.0.1", 42, null);
|
||||
assertEquals(Role.WORKER, after.role(),
|
||||
"removing the slot from config must demote the bound session to worker on its "
|
||||
+ "NEXT request — this is the ticket's whole point");
|
||||
assertEquals("term_a", after.terminal(), "same pane, same terminal — only the role changed");
|
||||
}
|
||||
|
||||
@Test
|
||||
void theOriginalBindingStillOccupiesTheRemovedSlotSoASecondTerminalCannotClaimIt(@TempDir Path dir)
|
||||
throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, yaml(WITH_SONNET_SLOT));
|
||||
ConfigRef ref = refFor(f);
|
||||
MemberRegistry registry = MemberRegistry.live(() -> ref.get().fleet());
|
||||
|
||||
MemberLifecycle.SlotReservation reservation = registry.reserve(MemberRole.ARCHITECT, "sonnet");
|
||||
assertTrue(registry.bind(reservation, "term_a"));
|
||||
|
||||
Files.writeString(f, yaml(WITHOUT_ARCHITECT_SLOTS));
|
||||
ConfigRef.Outcome removed = ref.reload();
|
||||
assertTrue(removed.applied(), "the reload must actually take effect: " + removed.summary());
|
||||
|
||||
// The binding survives the removal untouched.
|
||||
assertEquals("architect:designer", registry.slotForTerminal("term_a"),
|
||||
"a live binding must never be retroactively unbound by a config edit");
|
||||
assertEquals(Map.of("term_a", "architect:designer"), registry.snapshot());
|
||||
|
||||
// Bring the slot back into config. If the binding had been silently dropped by the removal
|
||||
// (rather than merely losing the privilege it grants), a second terminal could now claim
|
||||
// the "freed" key — the exact double-booking the class doc's binding rule rules out.
|
||||
Files.writeString(f, yaml(WITH_SONNET_SLOT));
|
||||
ConfigRef.Outcome restored = ref.reload();
|
||||
assertTrue(restored.applied(), "the reload must actually take effect: " + restored.summary());
|
||||
|
||||
assertFalse(registry.bind("architect:designer", "term_b"),
|
||||
"the slot is still occupied by term_a — a second terminal must not bind to it");
|
||||
assertThrows(IllegalArgumentException.class,
|
||||
() -> registry.reserve(MemberRole.ARCHITECT, "sonnet"),
|
||||
"the slot is still occupied by term_a — a fresh reservation must not find it free");
|
||||
assertEquals("architect:designer", registry.slotForTerminal("term_a"),
|
||||
"the original binding is unchanged throughout");
|
||||
}
|
||||
|
||||
@Test
|
||||
void unbindStillSucceedsForTheOriginalTerminalAfterItsSlotIsRemoved(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, yaml(WITH_SONNET_SLOT));
|
||||
ConfigRef ref = refFor(f);
|
||||
MemberRegistry registry = MemberRegistry.live(() -> ref.get().fleet());
|
||||
|
||||
MemberLifecycle.SlotReservation reservation = registry.reserve(MemberRole.ARCHITECT, "sonnet");
|
||||
assertTrue(registry.bind(reservation, "term_a"));
|
||||
|
||||
Files.writeString(f, yaml(WITHOUT_ARCHITECT_SLOTS));
|
||||
ConfigRef.Outcome out = ref.reload();
|
||||
assertTrue(out.applied(), "the reload must actually take effect: " + out.summary());
|
||||
|
||||
assertTrue(registry.unbind("architect:designer", "term_a"),
|
||||
"unbind must still work for a slot that config has since removed, or a session "
|
||||
+ "that outlives its slot's removal could never release it");
|
||||
assertNull(registry.slotForTerminal("term_a"));
|
||||
assertEquals(Map.of(), registry.snapshot());
|
||||
}
|
||||
}
|
||||
@@ -79,6 +79,14 @@ class ConfigRefTopLevelCoverageTest {
|
||||
* <li>{@code memberCredentials} — read live at {@code Fleetd.java:198, 205, 729}.</li>
|
||||
* <li>{@code memberLoginShell} — read live at
|
||||
* {@code HerdrPeerLauncher#configuredMemberLoginShell}.</li>
|
||||
* <li>{@code models} (fleetd #422) — membership is re-validated in full against the fresh
|
||||
* config on every {@code ConfigRef.reload()} (via {@code FleetConfig#validateAll()}), so
|
||||
* a bad edit is refused, never cached stale; the on/off half is read live by {@code
|
||||
* CompositePeerLauncher.enforceModelEnabled} and its candidate filter, and by {@code
|
||||
* fleet_profiles}/{@code GET /profiles} through {@code PeerLauncher.disabledModels()}.
|
||||
* Unlike {@code fleet} below, nothing about {@code models} is baked into a startup-built
|
||||
* object anywhere — there is no frozen half, so it belongs here whole rather than in
|
||||
* {@code SPLIT_KEYS}.</li>
|
||||
* </ul>
|
||||
*
|
||||
* <p>{@code fleet} used to sit here too, on the strength of most of it (role pools, charters,
|
||||
@@ -93,7 +101,7 @@ class ConfigRefTopLevelCoverageTest {
|
||||
* while this test stayed green throughout.</p>
|
||||
*/
|
||||
private static final Set<String> HOT_EXCLUDED_TOP_LEVEL_KEYS =
|
||||
Set.of("placement", "memberCredentials", "memberLoginShell");
|
||||
Set.of("placement", "memberCredentials", "memberLoginShell", "models");
|
||||
|
||||
@Test
|
||||
void everyTopLevelComponentIsAccountedForInExactlyOneClass() {
|
||||
@@ -120,7 +128,7 @@ class ConfigRefTopLevelCoverageTest {
|
||||
|
||||
// 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", "memberCredentials", "memberLoginShell"), hot,
|
||||
assertEquals(Set.of("placement", "memberCredentials", "memberLoginShell", "models"), 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 "
|
||||
|
||||
@@ -2901,4 +2901,119 @@ class FleetConfigTest {
|
||||
void modelsIsAKnownTopLevelKey() {
|
||||
assertTrue(FleetConfig.KNOWN_TOP_LEVEL_KEYS.contains("models"));
|
||||
}
|
||||
|
||||
// --- fleetd #422: model on/off (spawn-time gate + hot reload) -----------------------------
|
||||
|
||||
/**
|
||||
* Criterion 3: an {@code allow:} entry written before {@code enabled:} existed — no such key in
|
||||
* the YAML at all — must behave exactly as it always did: on, and absent from {@link
|
||||
* FleetConfig.Models#offIds()}. This is the old-style fixture the ticket asks for, proven
|
||||
* through a real YAML load rather than only through the {@code ModelEntry(String)} back-compat
|
||||
* constructor.
|
||||
*/
|
||||
@Test
|
||||
void anOldStyleAllowEntryWithNoEnabledFieldStaysOn(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
profiles:
|
||||
sonnet:
|
||||
baseUrl: http://gx10.gw:8000
|
||||
model: claude-sonnet-5
|
||||
models:
|
||||
allow:
|
||||
- model: claude-sonnet-5
|
||||
""");
|
||||
|
||||
FleetConfig cfg = FleetConfig.load(f);
|
||||
assertDoesNotThrow(cfg::validateModels);
|
||||
assertTrue(cfg.models().ids().contains("claude-sonnet-5"), "membership is unaffected");
|
||||
assertTrue(cfg.models().offIds().isEmpty(), "no enabled: field ⇒ nothing is off");
|
||||
}
|
||||
|
||||
/**
|
||||
* Criterion 7 (the whole point of the ticket, mirroring criterion 3's shape): a profile naming a
|
||||
* model that is turned off ({@code enabled: false}) must still be VALID config —
|
||||
* {@code validateModels()} checks membership only, never the on/off state, so it must not throw.
|
||||
* Collapsing "off" into "removed from allow:" would make this test fail, which is exactly the
|
||||
* design trap the ticket calls out.
|
||||
*/
|
||||
@Test
|
||||
void anOffModelIsStillValidConfigForValidateModels(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
profiles:
|
||||
local:
|
||||
baseUrl: http://local.gw:8000
|
||||
model: deepseek-v4-flash
|
||||
models:
|
||||
allow:
|
||||
- model: deepseek-v4-flash
|
||||
enabled: false
|
||||
""");
|
||||
|
||||
FleetConfig cfg = FleetConfig.load(f);
|
||||
assertDoesNotThrow(cfg::validateModels,
|
||||
"an off model must stay a valid allow-list member — only the spawn gate reads enabled");
|
||||
assertTrue(cfg.models().ids().contains("deepseek-v4-flash"));
|
||||
assertTrue(cfg.models().offIds().contains("deepseek-v4-flash"));
|
||||
}
|
||||
|
||||
/**
|
||||
* Criterion 8: one {@code enabled: false} entry names a model, not a profile — every profile
|
||||
* naming that model is off, on purpose (the live example: {@code deepseek-v4-flash} backs both
|
||||
* {@code local} and {@code local-direct}). Proven here at the {@code Models}/config-load level;
|
||||
* {@code CompositePeerLauncherTest} proves the spawn-time consequence for both profiles.
|
||||
*/
|
||||
@Test
|
||||
void turningOffOneModelIsIndependentOfHowManyProfilesNameIt(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
profiles:
|
||||
local:
|
||||
baseUrl: http://local.gw:8000
|
||||
model: deepseek-v4-flash
|
||||
local-direct:
|
||||
baseUrl: http://local.gw:8001
|
||||
model: deepseek-v4-flash
|
||||
models:
|
||||
allow:
|
||||
- model: deepseek-v4-flash
|
||||
enabled: false
|
||||
""");
|
||||
|
||||
FleetConfig cfg = FleetConfig.load(f);
|
||||
assertDoesNotThrow(cfg::validateModels);
|
||||
// One allow-list entry, one off id — the fan-out to both profiles happens at the reader
|
||||
// (CompositePeerLauncher), not by duplicating the entry per profile.
|
||||
assertEquals(Set.of("deepseek-v4-flash"), cfg.models().offIds());
|
||||
assertEquals("deepseek-v4-flash", cfg.profiles().get("local").model());
|
||||
assertEquals("deepseek-v4-flash", cfg.profiles().get("local-direct").model());
|
||||
}
|
||||
|
||||
/** An {@code enabled: true} entry (explicit, not just absent) also stays on — not just null. */
|
||||
@Test
|
||||
void explicitlyEnabledTrueStaysOn(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
profiles:
|
||||
sonnet:
|
||||
baseUrl: http://gx10.gw:8000
|
||||
model: claude-sonnet-5
|
||||
models:
|
||||
allow:
|
||||
- model: claude-sonnet-5
|
||||
enabled: true
|
||||
""");
|
||||
|
||||
FleetConfig cfg = FleetConfig.load(f);
|
||||
assertTrue(cfg.models().offIds().isEmpty());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -20,6 +20,7 @@ import java.util.regex.Pattern;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/** Unit behaviour of the CB-106 completion resolver in isolation from the injector. */
|
||||
@@ -884,19 +885,22 @@ class CompletionResolverTest {
|
||||
@Test
|
||||
void coverageIsOffWhenNoProfileHasAPatternConfigured() {
|
||||
assertEquals("off (no profile has an exhaustedPattern configured; profiles: [terra])",
|
||||
CompletionResolver.coverage("exhaustedPattern", Set.of("terra"), Set.of()));
|
||||
CompletionResolver.coverage("exhaustedPattern", CompletionResolver.UnsetMeaning.OFF,
|
||||
Set.of("terra"), Set.of()));
|
||||
}
|
||||
|
||||
@Test
|
||||
void coverageIsFullWhenEveryProfileHasAPatternConfigured() {
|
||||
assertEquals("full (all profiles configured: [gx10, terra])",
|
||||
CompletionResolver.coverage("exhaustedPattern", Set.of("terra", "gx10"), Set.of("terra", "gx10")));
|
||||
CompletionResolver.coverage("exhaustedPattern", CompletionResolver.UnsetMeaning.OFF,
|
||||
Set.of("terra", "gx10"), Set.of("terra", "gx10")));
|
||||
}
|
||||
|
||||
@Test
|
||||
void coverageIsPartialAndNamesWhichProfilesAreConfigured() {
|
||||
assertEquals("partial (configured: [terra]; not configured: [gx10])",
|
||||
CompletionResolver.coverage("exhaustedPattern", Set.of("terra", "gx10"), Set.of("terra")));
|
||||
CompletionResolver.coverage("exhaustedPattern", CompletionResolver.UnsetMeaning.OFF,
|
||||
Set.of("terra", "gx10"), Set.of("terra")));
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -910,11 +914,42 @@ class CompletionResolverTest {
|
||||
*
|
||||
* <p>Every earlier test here passed the exhaustion case only, so none of them could see it. This
|
||||
* one pins that the message names the key the caller actually meant.
|
||||
*
|
||||
* <p>fleetd#415: the expected wording changed here too. {@code errorPattern} has a built-in
|
||||
* fallback ({@link CompletionResolver#BACKEND_ERROR}), so an empty {@code configuredProfiles}
|
||||
* for it is not "off" — see {@link #coverageDistinguishesOffFromBuiltInDefaultForTheSameEmptyInput}
|
||||
* for the test built specifically to pin that distinction.
|
||||
*/
|
||||
@Test
|
||||
void coverageNamesTheConfigKeyItsCallerMeansRatherThanAlwaysSayingExhaustedPattern() {
|
||||
assertEquals("off (no profile has an errorPattern configured; profiles: [gx10, terra])",
|
||||
CompletionResolver.coverage("errorPattern", Set.of("terra", "gx10"), Set.of()));
|
||||
assertEquals("built-in default for all profiles (no profile customises errorPattern; "
|
||||
+ "profiles: [gx10, terra])",
|
||||
CompletionResolver.coverage("errorPattern", CompletionResolver.UnsetMeaning.BUILT_IN_DEFAULT,
|
||||
Set.of("terra", "gx10"), Set.of()));
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd#415: {@code coverage()} measures pattern coverage (how many profiles set the key), but
|
||||
* for {@code errorPattern} the empty case is not the feature-off state — a profile with no
|
||||
* configured {@code errorPattern} still runs the classification against
|
||||
* {@link CompletionResolver#BACKEND_ERROR}. For {@code exhaustedPattern} there is no fallback,
|
||||
* so empty really is off. Same shape of input (empty {@code configuredProfiles}, one profile),
|
||||
* different {@link CompletionResolver.UnsetMeaning} — the wording must differ, or this method is
|
||||
* back to conflating pattern coverage with feature state for the one key where they disagree.
|
||||
*/
|
||||
@Test
|
||||
void coverageDistinguishesOffFromBuiltInDefaultForTheSameEmptyInput() {
|
||||
String exhaustedLine = CompletionResolver.coverage("exhaustedPattern",
|
||||
CompletionResolver.UnsetMeaning.OFF, Set.of("gx10", "terra"), Set.of());
|
||||
String errorLine = CompletionResolver.coverage("errorPattern",
|
||||
CompletionResolver.UnsetMeaning.BUILT_IN_DEFAULT, Set.of("gx10", "terra"), Set.of());
|
||||
|
||||
assertEquals("off (no profile has an exhaustedPattern configured; profiles: [gx10, terra])",
|
||||
exhaustedLine);
|
||||
assertEquals("built-in default for all profiles (no profile customises errorPattern; "
|
||||
+ "profiles: [gx10, terra])", errorLine);
|
||||
assertNotEquals(exhaustedLine, errorLine,
|
||||
"the same empty-coverage input must not read as the same feature state for both keys");
|
||||
}
|
||||
|
||||
// --- fleetd#201 Unit 1: target-keyed backend-error pattern + typed sink ----------------------
|
||||
|
||||
@@ -0,0 +1,109 @@
|
||||
package dev.ltms.fleet.mcp;
|
||||
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
import dev.ltms.fleet.member.ClaudeCodeLauncher;
|
||||
import dev.ltms.fleet.member.CompositePeerLauncher;
|
||||
import dev.ltms.fleet.member.HerdrPeerLauncher;
|
||||
import dev.ltms.fleet.peer.PeerLauncher;
|
||||
import dev.ltms.fleet.placement.BackendOutagePolicy;
|
||||
import dev.ltms.fleet.placement.BackendQuarantine;
|
||||
import dev.ltms.fleet.placement.PlacementPolicies;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
|
||||
/**
|
||||
* fleetd #422 follow-up: {@code fleet_profiles}/{@code GET /profiles} — the reporting surface a
|
||||
* lead actually reads — must let it tell apart the three states {@link
|
||||
* dev.ltms.fleet.peer.PeerLauncher#disabledModels()} alone collapses into one empty set: no
|
||||
* {@code models:} block at all, a block armed with nothing currently off, and a block with N
|
||||
* models off. See {@link dev.ltms.fleet.peer.PeerLauncher.ModelGateState}'s javadoc for why a
|
||||
* bare {@code disabledModels()} read cannot make this distinction, and {@link
|
||||
* FleetMcp#profilesView} for where {@code modelGateArmed} is added alongside the existing {@code
|
||||
* modelsOff} key.
|
||||
*
|
||||
* <p>Every assertion here goes through {@link FleetMcp#profilesView}, never {@code
|
||||
* PeerLauncher.modelGateState()} directly — {@code CompositePeerLauncherTest} already proves the
|
||||
* accessor itself; this class proves the surface a lead reads (fleet_profiles / GET /profiles)
|
||||
* renders what that accessor reports.
|
||||
*/
|
||||
class FleetProfilesModelGateStateTest {
|
||||
|
||||
private static FleetConfig.Profile profile(String name, String model) {
|
||||
return new FleetConfig.Profile(name, "http://gx00.gw:8000", model, null, "FLEETD_WORKER_TOKEN",
|
||||
null, "tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
|
||||
}
|
||||
|
||||
private static FleetMcp.QuarantineSource noQuarantine() {
|
||||
return new FleetMcp.QuarantineSource(_ -> null, BackendQuarantine.none(), _ -> false);
|
||||
}
|
||||
|
||||
private static HerdrPeerLauncher claudeAdapter(FakeHerdr h, Map<String, FleetConfig.Profile> profiles) {
|
||||
return new ClaudeCodeLauncher(new AgentControl(h), new WorkspaceControl(h),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), profiles, "local", _ -> "tok");
|
||||
}
|
||||
|
||||
/**
|
||||
* State 1: no {@code models:} block at all — a plain {@code ClaudeCodeLauncher} (no {@code
|
||||
* models:} supplier exists for it to read) has nothing to gate against, matching the fleet01
|
||||
* host measured for this ticket: {@code grep -c '^models:' fleetd.yaml} returns 0 there.
|
||||
*/
|
||||
@Test
|
||||
void noModelsBlockReportsGateNotArmed() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
Map<String, FleetConfig.Profile> profiles = Map.of("local", profile("local", "deepseek-v4-flash"));
|
||||
PeerLauncher workers = claudeAdapter(h, profiles);
|
||||
|
||||
Map<String, Object> view = FleetMcp.profilesView(workers, noQuarantine(), FleetMcp.OutageSource.none());
|
||||
|
||||
assertEquals(Boolean.FALSE, view.get("modelGateArmed"),
|
||||
"no models: block to read from — nothing is gated, and nothing can be");
|
||||
assertFalse(view.containsKey("modelsOff"), "nothing configured, so no off set to report either");
|
||||
}
|
||||
|
||||
/** State 2: a {@code models:} block is present, but nothing in it is currently turned off. */
|
||||
@Test
|
||||
void modelsBlockWithNothingOffReportsGateArmedAndZeroOff() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
Map<String, FleetConfig.Profile> profiles = Map.of("local", profile("local", "deepseek-v4-flash"));
|
||||
FleetConfig.Models models = new FleetConfig.Models(
|
||||
List.of(new FleetConfig.Models.ModelEntry("deepseek-v4-flash", true)));
|
||||
PeerLauncher workers = new CompositePeerLauncher(List.of(claudeAdapter(h, profiles)), "local", profiles,
|
||||
PlacementPolicies.fixed(), _ -> 0, null, BackendQuarantine.none(),
|
||||
new BackendOutagePolicy(System::nanoTime), () -> models);
|
||||
|
||||
Map<String, Object> view = FleetMcp.profilesView(workers, noQuarantine(), FleetMcp.OutageSource.none());
|
||||
|
||||
assertEquals(Boolean.TRUE, view.get("modelGateArmed"),
|
||||
"a models: block is present, so the gate is armed even though nothing is off yet");
|
||||
assertFalse(view.containsKey("modelsOff"),
|
||||
"nothing is off, so the key stays absent — an empty list here would be indistinguishable "
|
||||
+ "from today's modelsOff omission, exactly the ambiguity modelGateArmed exists to remove");
|
||||
}
|
||||
|
||||
/** State 3: a {@code models:} block is present with one model currently turned off. */
|
||||
@Test
|
||||
void modelsBlockWithModelsOffReportsGateArmedAndTheOffSet() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
Map<String, FleetConfig.Profile> profiles = Map.of("local", profile("local", "deepseek-v4-flash"));
|
||||
FleetConfig.Models models = new FleetConfig.Models(
|
||||
List.of(new FleetConfig.Models.ModelEntry("deepseek-v4-flash", false)));
|
||||
PeerLauncher workers = new CompositePeerLauncher(List.of(claudeAdapter(h, profiles)), "local", profiles,
|
||||
PlacementPolicies.fixed(), _ -> 0, null, BackendQuarantine.none(),
|
||||
new BackendOutagePolicy(System::nanoTime), () -> models);
|
||||
|
||||
Map<String, Object> view = FleetMcp.profilesView(workers, noQuarantine(), FleetMcp.OutageSource.none());
|
||||
|
||||
assertEquals(Boolean.TRUE, view.get("modelGateArmed"));
|
||||
assertEquals(List.of("deepseek-v4-flash"), view.get("modelsOff"));
|
||||
}
|
||||
}
|
||||
@@ -3,6 +3,7 @@ package dev.ltms.fleet.member;
|
||||
import ch.qos.logback.classic.Logger;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.Agent;
|
||||
@@ -22,8 +23,11 @@ import dev.ltms.fleet.placement.BackendQuarantine;
|
||||
import dev.ltms.fleet.placement.PlacementException;
|
||||
import dev.ltms.fleet.placement.PlacementPolicies;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.EnumSet;
|
||||
import java.util.HashMap;
|
||||
import java.util.LinkedHashMap;
|
||||
@@ -141,6 +145,13 @@ class CompositePeerLauncherTest {
|
||||
null, null, credentialId, null);
|
||||
}
|
||||
|
||||
/** Like {@link #stubWorker(String)}, but with an explicit {@code model:} for fleetd #422 tests. */
|
||||
private static FleetConfig.Profile stubWorkerModel(String profile, String model) {
|
||||
return new FleetConfig.Profile(profile, "http://gx00.gw:8000", model,
|
||||
null, "FLEETD_WORKER_TOKEN", List.of("claude"), "tab", "fleetd-workers",
|
||||
"w #{n}", null, null, null, null, null, null, null, null, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* An <em>order-preserving</em> profile map. Never {@code Map.of} here: its iteration order is
|
||||
* salted per JVM run, and the weighted policy breaks an exact-weight tie on candidate order —
|
||||
@@ -525,6 +536,31 @@ class CompositePeerLauncherTest {
|
||||
assertEquals("claude", h.profile(), "the returned handle carries the resolved default profile");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #435: {@code FixedPlacementPolicy} — the default placement policy every config uses
|
||||
* unless {@code placement:} is set — never consulted {@code maxLoad}, so an unqualified spawn
|
||||
* (a blank profile, the normal delegation path) landed on a capped default anyway. Measured on
|
||||
* 7667727: a single dev profile at {@code maxLoad: 1} with 1 live, under {@code fixed()},
|
||||
* returned "SPAWNED on profile=a". This test goes through {@code CompositePeerLauncher.spawn}
|
||||
* with a blank profile, not the policy in isolation, so it proves the caller actually reaches
|
||||
* the fixed default's new cap check rather than only the {@code select} method.
|
||||
*/
|
||||
@Test
|
||||
void fixedPolicyGatesDefaultProfileAtMaxLoadOnUnqualifiedSpawn() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Map<String, FleetConfig.Profile> profiles = ordered(
|
||||
"a", stubWorker("a", 1.0f, 1),
|
||||
"b", stubWorker("b", 1.0f, null));
|
||||
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "a", Set.of());
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(
|
||||
List.of(adapter), "a", profiles, PlacementPolicies.fixed(), name -> "a".equals(name) ? 1 : 0);
|
||||
|
||||
PeerHandle h = composite.spawn(new SpawnRequest(null, null, null));
|
||||
assertEquals("b", h.profile(),
|
||||
"the default profile a is at maxLoad, so fixed placement must fall through to b");
|
||||
assertEquals(0, adapter.spawnCount("a"), "a is never spawned — it is already at cap");
|
||||
}
|
||||
|
||||
@Test
|
||||
void weightedPolicyGatesProfileAtMaxLoad() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
@@ -1189,4 +1225,384 @@ class CompositePeerLauncherTest {
|
||||
assertDoesNotThrow(() -> composite.spawn(new SpawnRequest("claude", null, null)));
|
||||
assertDoesNotThrow(() -> composite.spawn(new SpawnRequest("gemini", null, null)));
|
||||
}
|
||||
|
||||
// ── fleetd #422: the model on/off gate — a FOURTH, independent reason to refuse a spawn ──────
|
||||
// ── (an operator decision, never a backend-reported outage) — never merged with quarantine ────
|
||||
// ── or cool-off above ─────────────────────────────────────────────────────────────────────────
|
||||
|
||||
private static final BackendOutagePolicy NO_OUTAGE = new BackendOutagePolicy(() -> 0L);
|
||||
|
||||
/** Criterion 1: an explicit spawn onto a profile whose model is off is refused. */
|
||||
@Test
|
||||
void explicitSpawnOntoAnOffModelProfileIsRefused() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Map<String, FleetConfig.Profile> profiles = ordered(
|
||||
"local", stubWorkerModel("local", "deepseek-v4-flash"),
|
||||
"sonnet", stubWorkerModel("sonnet", "claude-sonnet-5"));
|
||||
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "local", Set.of());
|
||||
FleetConfig.Models models = new FleetConfig.Models(List.of(
|
||||
new FleetConfig.Models.ModelEntry("deepseek-v4-flash", false)));
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "local", profiles,
|
||||
PlacementPolicies.fixed(), _ -> 0, null, BackendQuarantine.none(), NO_OUTAGE,
|
||||
() -> models);
|
||||
|
||||
PlacementException e = assertThrows(PlacementException.class,
|
||||
() -> composite.spawn(new SpawnRequest("local", null, null)));
|
||||
assertTrue(e.getMessage().contains("local"), "message names the profile: " + e.getMessage());
|
||||
assertTrue(e.getMessage().contains("deepseek-v4-flash"),
|
||||
"message names the model: " + e.getMessage());
|
||||
assertTrue(e.getMessage().contains("operator") && e.getMessage().contains("turned off"),
|
||||
"wording says the OPERATOR turned it off: " + e.getMessage());
|
||||
// Distinct from quarantine/cool-off wording (criterion 1's explicit requirement).
|
||||
assertFalse(e.getMessage().contains("quarantined"), "must not read like quarantine: " + e.getMessage());
|
||||
assertFalse(e.getMessage().contains("cooling off"), "must not read like cool-off: " + e.getMessage());
|
||||
assertEquals(0, adapter.spawnCount("local"), "the off-model profile is never delegated to");
|
||||
}
|
||||
|
||||
/** A profile whose model is NOT off spawns normally even while another model is off. */
|
||||
@Test
|
||||
void explicitSpawnOntoAnEnabledModelProfileSucceedsWhileAnotherModelIsOff() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Map<String, FleetConfig.Profile> profiles = ordered(
|
||||
"local", stubWorkerModel("local", "deepseek-v4-flash"),
|
||||
"sonnet", stubWorkerModel("sonnet", "claude-sonnet-5"));
|
||||
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "local", Set.of());
|
||||
FleetConfig.Models models = new FleetConfig.Models(List.of(
|
||||
new FleetConfig.Models.ModelEntry("deepseek-v4-flash", false)));
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "local", profiles,
|
||||
PlacementPolicies.fixed(), _ -> 0, null, BackendQuarantine.none(), NO_OUTAGE,
|
||||
() -> models);
|
||||
|
||||
assertDoesNotThrow(() -> composite.spawn(new SpawnRequest("sonnet", null, null)));
|
||||
assertEquals(1, adapter.spawnCount("sonnet"));
|
||||
}
|
||||
|
||||
/**
|
||||
* Criterion 8: one {@code enabled: false} entry disables EVERY profile naming that model — the
|
||||
* live example, {@code deepseek-v4-flash} backing both {@code local} and {@code local-direct}.
|
||||
*/
|
||||
@Test
|
||||
void turningOffOneModelRefusesEveryProfileThatNamesIt() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Map<String, FleetConfig.Profile> profiles = ordered(
|
||||
"local", stubWorkerModel("local", "deepseek-v4-flash"),
|
||||
"local-direct", stubWorkerModel("local-direct", "deepseek-v4-flash"));
|
||||
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "local", Set.of());
|
||||
FleetConfig.Models models = new FleetConfig.Models(List.of(
|
||||
new FleetConfig.Models.ModelEntry("deepseek-v4-flash", false)));
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "local", profiles,
|
||||
PlacementPolicies.fixed(), _ -> 0, null, BackendQuarantine.none(), NO_OUTAGE,
|
||||
() -> models);
|
||||
|
||||
assertThrows(PlacementException.class, () -> composite.spawn(new SpawnRequest("local", null, null)));
|
||||
assertThrows(PlacementException.class,
|
||||
() -> composite.spawn(new SpawnRequest("local-direct", null, null)),
|
||||
"local-direct shares local's model, so it must be refused too");
|
||||
assertEquals(0, adapter.spawnCount("local"));
|
||||
assertEquals(0, adapter.spawnCount("local-direct"));
|
||||
}
|
||||
|
||||
/** Criterion 3: an old-style entry with no {@code enabled} field never refuses a spawn. */
|
||||
@Test
|
||||
void anEntryWithNoEnabledFieldNeverRefusesASpawn() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Map<String, FleetConfig.Profile> profiles = Map.of("local", stubWorkerModel("local", "deepseek-v4-flash"));
|
||||
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "local", Set.of());
|
||||
// The back-compat single-arg ModelEntry constructor — no enabled field in the shape at all.
|
||||
FleetConfig.Models models = new FleetConfig.Models(
|
||||
List.of(new FleetConfig.Models.ModelEntry("deepseek-v4-flash")));
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "local", profiles,
|
||||
PlacementPolicies.fixed(), _ -> 0, null, BackendQuarantine.none(), NO_OUTAGE,
|
||||
() -> models);
|
||||
|
||||
assertDoesNotThrow(() -> composite.spawn(new SpawnRequest("local", null, null)));
|
||||
assertEquals(1, adapter.spawnCount("local"));
|
||||
}
|
||||
|
||||
/** A model absent from {@code models.allow:} entirely (nothing to gate against) is never refused. */
|
||||
@Test
|
||||
void aModelNotConfiguredInTheAllowListIsNeverGated() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Map<String, FleetConfig.Profile> profiles = Map.of("local", stubWorkerModel("local", "unlisted-model"));
|
||||
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "local", Set.of());
|
||||
FleetConfig.Models models = new FleetConfig.Models(List.of(
|
||||
new FleetConfig.Models.ModelEntry("deepseek-v4-flash", false)));
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "local", profiles,
|
||||
PlacementPolicies.fixed(), _ -> 0, null, BackendQuarantine.none(), NO_OUTAGE,
|
||||
() -> models);
|
||||
|
||||
assertDoesNotThrow(() -> composite.spawn(new SpawnRequest("local", null, null)));
|
||||
}
|
||||
|
||||
/** Criterion 2: an unqualified spawn skips an off-model candidate and lands on another one. */
|
||||
@Test
|
||||
void placementSkipsAnOffModelProfileAndRoutesToAnotherOne() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Map<String, FleetConfig.Profile> profiles = ordered(
|
||||
"local", stubWorkerModel("local", "deepseek-v4-flash"),
|
||||
"sonnet", stubWorkerModel("sonnet", "claude-sonnet-5"));
|
||||
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "local", Set.of());
|
||||
FleetConfig.Models models = new FleetConfig.Models(List.of(
|
||||
new FleetConfig.Models.ModelEntry("deepseek-v4-flash", false)));
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "local", profiles,
|
||||
PlacementPolicies.weighted(), _ -> 0, null, BackendQuarantine.none(), NO_OUTAGE,
|
||||
() -> models);
|
||||
|
||||
PeerHandle h = composite.spawn(new SpawnRequest(null, null, null));
|
||||
assertEquals("sonnet", h.profile(), "local's model is off, so an unqualified spawn must land on sonnet");
|
||||
assertEquals(0, adapter.spawnCount("local"));
|
||||
}
|
||||
|
||||
/**
|
||||
* Criterion 2's second half: when EVERY candidate's model is off, the placement exception must
|
||||
* name that as the cause — not a generic "no candidates" message.
|
||||
*/
|
||||
@Test
|
||||
void automaticPlacementNamesModelOffWhenEveryCandidateIsOffModel() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Map<String, FleetConfig.Profile> profiles = ordered(
|
||||
"local", stubWorkerModel("local", "deepseek-v4-flash"),
|
||||
"local-direct", stubWorkerModel("local-direct", "deepseek-v4-flash"));
|
||||
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "local", Set.of());
|
||||
FleetConfig.Models models = new FleetConfig.Models(List.of(
|
||||
new FleetConfig.Models.ModelEntry("deepseek-v4-flash", false)));
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "local", profiles,
|
||||
PlacementPolicies.weighted(), _ -> 0, null, BackendQuarantine.none(), NO_OUTAGE,
|
||||
() -> models);
|
||||
|
||||
PlacementException e = assertThrows(PlacementException.class,
|
||||
() -> composite.spawn(new SpawnRequest(null, null, null)),
|
||||
"both candidates share the off model — nothing is available");
|
||||
assertTrue(e.getMessage().contains("model") && e.getMessage().contains("turned off"),
|
||||
"message names model-off as the cause, not a generic no-candidates message: "
|
||||
+ e.getMessage());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #422 follow-up: {@code fixed} is the DEFAULT placement policy ({@code
|
||||
* PlacementPolicies.fromName} returns it for an absent/blank name), and it built its own inline
|
||||
* candidate filter instead of calling {@code PlacementPolicyUtil.available()} — so it never
|
||||
* checked {@code modelOff()}. Mirrors {@code placementSkipsAnOffModelProfileAndRoutesToAnotherOne}
|
||||
* above with only the policy swapped, to prove the gate now fires on the path most fleets use.
|
||||
*/
|
||||
@Test
|
||||
void fixedPlacementSkipsAnOffModelProfileToo() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Map<String, FleetConfig.Profile> profiles = ordered(
|
||||
"local", stubWorkerModel("local", "deepseek-v4-flash"),
|
||||
"sonnet", stubWorkerModel("sonnet", "claude-sonnet-5"));
|
||||
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "local", Set.of());
|
||||
FleetConfig.Models models = new FleetConfig.Models(List.of(
|
||||
new FleetConfig.Models.ModelEntry("deepseek-v4-flash", false)));
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "local", profiles,
|
||||
PlacementPolicies.fixed(), _ -> 0, null, BackendQuarantine.none(), NO_OUTAGE,
|
||||
() -> models);
|
||||
|
||||
PeerHandle h = composite.spawn(new SpawnRequest(null, null, null));
|
||||
assertEquals("sonnet", h.profile(), "local's model is off, so an unqualified spawn must land on sonnet");
|
||||
assertEquals(0, adapter.spawnCount("local"));
|
||||
}
|
||||
|
||||
/**
|
||||
* Criterion 4: turning a model off/on is HOT — no restart — proven through a REAL
|
||||
* {@code ConfigRef.reload()}, not a hand-rolled supplier swap. Also proves {@code models} is
|
||||
* correctly reclassified: the reload's {@link ConfigRef.Outcome#applied()} is {@code true} and
|
||||
* {@code "models"} never appears in {@link ConfigRef.Outcome#deferred()}.
|
||||
*/
|
||||
@Test
|
||||
void modelOnOffIsHotReloadedThroughARealConfigRef(@TempDir Path dir) throws Exception {
|
||||
Path yaml = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(yaml, """
|
||||
bind:
|
||||
port: 8080
|
||||
profiles:
|
||||
local:
|
||||
baseUrl: http://local.gw:8000
|
||||
model: deepseek-v4-flash
|
||||
models:
|
||||
allow:
|
||||
- model: deepseek-v4-flash
|
||||
""");
|
||||
FleetConfig initial = FleetConfig.load(yaml);
|
||||
ConfigRef configRef = new ConfigRef(yaml, initial);
|
||||
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
StubLauncher adapter = new StubLauncher("claude", herdr,
|
||||
Map.of("local", stubWorker("local")), "local", Set.of());
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(
|
||||
List.of(adapter), "local", configRef, _ -> 0, BackendQuarantine.none(), NO_OUTAGE);
|
||||
|
||||
assertDoesNotThrow(() -> composite.spawn(new SpawnRequest("local", null, null)),
|
||||
"the model starts enabled");
|
||||
assertEquals(1, adapter.spawnCount("local"));
|
||||
|
||||
Files.writeString(yaml, """
|
||||
bind:
|
||||
port: 8080
|
||||
profiles:
|
||||
local:
|
||||
baseUrl: http://local.gw:8000
|
||||
model: deepseek-v4-flash
|
||||
models:
|
||||
allow:
|
||||
- model: deepseek-v4-flash
|
||||
enabled: false
|
||||
""");
|
||||
ConfigRef.Outcome outcome = configRef.reload();
|
||||
assertTrue(outcome.applied(), "a models.allow on/off edit must apply live, never be refused");
|
||||
assertFalse(outcome.deferred().contains("models"),
|
||||
"models is hot-excluded now — it must never be reported as a deferred key");
|
||||
|
||||
PlacementException e = assertThrows(PlacementException.class,
|
||||
() -> composite.spawn(new SpawnRequest("local", null, null)),
|
||||
"the very next spawn must see the reload, with no restart");
|
||||
assertTrue(e.getMessage().contains("deepseek-v4-flash"));
|
||||
assertEquals(1, adapter.spawnCount("local"), "still just the one successful spawn from before");
|
||||
|
||||
// And back on, still hot, still no restart.
|
||||
Files.writeString(yaml, """
|
||||
bind:
|
||||
port: 8080
|
||||
profiles:
|
||||
local:
|
||||
baseUrl: http://local.gw:8000
|
||||
model: deepseek-v4-flash
|
||||
models:
|
||||
allow:
|
||||
- model: deepseek-v4-flash
|
||||
enabled: true
|
||||
""");
|
||||
ConfigRef.Outcome reEnabled = configRef.reload();
|
||||
assertTrue(reEnabled.applied());
|
||||
assertDoesNotThrow(() -> composite.spawn(new SpawnRequest("local", null, null)));
|
||||
assertEquals(2, adapter.spawnCount("local"));
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #422 follow-up, acceptance criterion 2: {@link CompositePeerLauncher#modelGateState()}
|
||||
* is LIVE — no restart — proven through a REAL {@link ConfigRef#reload()}, exactly like {@link
|
||||
* #modelOnOffIsHotReloadedThroughARealConfigRef} above proves for the on/off gate itself. This
|
||||
* single reload sequence walks through all three states the ticket asks for: no {@code models:}
|
||||
* block, a block armed with nothing off, and a block with one model off — so a reload that flips
|
||||
* between any of the three is proven live, not just the on/off edit within an already-armed block.
|
||||
*/
|
||||
@Test
|
||||
void modelGateStateIsHotReloadedThroughARealConfigRef(@TempDir Path dir) throws Exception {
|
||||
Path yaml = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(yaml, """
|
||||
bind:
|
||||
port: 8080
|
||||
profiles:
|
||||
local:
|
||||
baseUrl: http://local.gw:8000
|
||||
model: deepseek-v4-flash
|
||||
""");
|
||||
FleetConfig initial = FleetConfig.load(yaml);
|
||||
ConfigRef configRef = new ConfigRef(yaml, initial);
|
||||
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
StubLauncher adapter = new StubLauncher("claude", herdr,
|
||||
Map.of("local", stubWorker("local")), "local", Set.of());
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(
|
||||
List.of(adapter), "local", configRef, _ -> 0, BackendQuarantine.none(), NO_OUTAGE);
|
||||
|
||||
PeerLauncher.ModelGateState notConfigured = composite.modelGateState();
|
||||
assertFalse(notConfigured.configured(), "no models: block in the config at all");
|
||||
assertEquals(Set.of(), notConfigured.off());
|
||||
|
||||
Files.writeString(yaml, """
|
||||
bind:
|
||||
port: 8080
|
||||
profiles:
|
||||
local:
|
||||
baseUrl: http://local.gw:8000
|
||||
model: deepseek-v4-flash
|
||||
models:
|
||||
allow:
|
||||
- model: deepseek-v4-flash
|
||||
enabled: false
|
||||
""");
|
||||
assertTrue(configRef.reload().applied(), "adding a models: block must apply live, no restart");
|
||||
PeerLauncher.ModelGateState armedWithOneOff = composite.modelGateState();
|
||||
assertTrue(armedWithOneOff.configured(), "a models: block now exists — the gate is armed");
|
||||
assertEquals(Set.of("deepseek-v4-flash"), armedWithOneOff.off());
|
||||
|
||||
Files.writeString(yaml, """
|
||||
bind:
|
||||
port: 8080
|
||||
profiles:
|
||||
local:
|
||||
baseUrl: http://local.gw:8000
|
||||
model: deepseek-v4-flash
|
||||
models:
|
||||
allow:
|
||||
- model: deepseek-v4-flash
|
||||
enabled: true
|
||||
""");
|
||||
assertTrue(configRef.reload().applied(), "flipping the entry back on must apply live too");
|
||||
PeerLauncher.ModelGateState armedWithZeroOff = composite.modelGateState();
|
||||
assertTrue(armedWithZeroOff.configured(),
|
||||
"the block is still present — armed and reporting zero, not the same as no block at all");
|
||||
assertEquals(Set.of(), armedWithZeroOff.off());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #422 follow-up, acceptance criterion 3: the invariant is that an absent {@code
|
||||
* models:} block stays permitted and must never be fatal. Proved, not assumed — a config
|
||||
* without one loads, validates, reports the gate as not configured, AND still spawns normally
|
||||
* (no {@link PlacementException} from a gate that has nothing to check against), using the same
|
||||
* production-shaped {@code Supplier<FleetConfig>} wiring {@code Fleetd.main} actually uses.
|
||||
*/
|
||||
@Test
|
||||
void noModelsBlockConfigStillLoadsAndSpawnsNormally(@TempDir Path dir) throws Exception {
|
||||
Path yaml = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(yaml, """
|
||||
bind:
|
||||
port: 8080
|
||||
profiles:
|
||||
local:
|
||||
baseUrl: http://local.gw:8000
|
||||
model: deepseek-v4-flash
|
||||
""");
|
||||
FleetConfig cfg = FleetConfig.load(yaml);
|
||||
assertDoesNotThrow(cfg::validateAll, "a config with no models: block must load and validate cleanly");
|
||||
ConfigRef configRef = new ConfigRef(yaml, cfg);
|
||||
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
StubLauncher adapter = new StubLauncher("claude", herdr,
|
||||
Map.of("local", stubWorker("local")), "local", Set.of());
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(
|
||||
List.of(adapter), "local", configRef, _ -> 0, BackendQuarantine.none(), NO_OUTAGE);
|
||||
|
||||
assertFalse(composite.modelGateState().configured());
|
||||
assertDoesNotThrow(() -> composite.spawn(new SpawnRequest("local", null, null)),
|
||||
"no models: block means nothing to gate against — the spawn must go through");
|
||||
assertEquals(1, adapter.spawnCount("local"));
|
||||
}
|
||||
|
||||
/** {@code fleet_profiles}/{@code GET /profiles} must read the exact same live source the gate reads. */
|
||||
@Test
|
||||
void disabledModelsReportsWhatTheGateActuallyEnforces() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Map<String, FleetConfig.Profile> profiles = ordered(
|
||||
"local", stubWorkerModel("local", "deepseek-v4-flash"),
|
||||
"sonnet", stubWorkerModel("sonnet", "claude-sonnet-5"));
|
||||
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "local", Set.of());
|
||||
FleetConfig.Models models = new FleetConfig.Models(List.of(
|
||||
new FleetConfig.Models.ModelEntry("deepseek-v4-flash", false),
|
||||
new FleetConfig.Models.ModelEntry("claude-sonnet-5", true)));
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "local", profiles,
|
||||
PlacementPolicies.fixed(), _ -> 0, null, BackendQuarantine.none(), NO_OUTAGE,
|
||||
() -> models);
|
||||
|
||||
assertEquals(Set.of("deepseek-v4-flash"), composite.disabledModels());
|
||||
}
|
||||
|
||||
/** A caller with no {@code models:} block (the map-based constructors) reports nothing off. */
|
||||
@Test
|
||||
void disabledModelsIsEmptyWithNoModelsConfigured() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
PeerLauncher composite = composite(herdr);
|
||||
assertEquals(Set.of(), composite.disabledModels());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -90,6 +90,63 @@ class PlacementPolicyTest {
|
||||
assertTrue(e.getMessage().contains("quarantined"), e.getMessage());
|
||||
}
|
||||
|
||||
// --- fleetd #422: FixedPlacementPolicy must consult modelOff too, at BOTH filter sites ------
|
||||
|
||||
/** The default-profile fast path (:44) must skip a default whose model is off. */
|
||||
@Test
|
||||
void fixedSkipsModelOffDefault() {
|
||||
PlacementPolicy policy = PlacementPolicies.fixed();
|
||||
PlacementContext ctx = new PlacementContext("b",
|
||||
List.of(PlacementCandidate.profile("a"), PlacementCandidate.profile("b")),
|
||||
noSessions(), Set.of(), Set.of(), Set.of(), Set.of("b"));
|
||||
assertEquals("a", policy.select(ctx).profile(),
|
||||
"the default 'b' names an off model, so fixed falls through to the first available candidate");
|
||||
}
|
||||
|
||||
/**
|
||||
* The fallback walk (:48-53) must skip an off-model candidate too — exercised independently of
|
||||
* the default-profile fast path by using no default at all, so this is the only filter that runs.
|
||||
*/
|
||||
@Test
|
||||
void fixedFallbackWalkSkipsModelOffCandidate() {
|
||||
PlacementPolicy policy = PlacementPolicies.fixed();
|
||||
PlacementContext ctx = new PlacementContext(null,
|
||||
List.of(PlacementCandidate.profile("a"), PlacementCandidate.profile("b")),
|
||||
noSessions(), Set.of(), Set.of(), Set.of(), Set.of("a"));
|
||||
assertEquals("b", policy.select(ctx).profile(),
|
||||
"candidate 'a' names an off model, so the fallback walk skips it and picks 'b'");
|
||||
}
|
||||
|
||||
@Test
|
||||
void fixedThrowsWhenDefaultAndEveryCandidateModelOff() {
|
||||
PlacementPolicy policy = PlacementPolicies.fixed();
|
||||
PlacementContext ctx = new PlacementContext("b",
|
||||
List.of(PlacementCandidate.profile("a"), PlacementCandidate.profile("b")),
|
||||
noSessions(), Set.of(), Set.of(), Set.of(), Set.of("a", "b"));
|
||||
PlacementException e = assertThrows(PlacementException.class, () -> policy.select(ctx));
|
||||
assertTrue(e.getMessage().contains("turned off"),
|
||||
"message names model-off as the cause: " + e.getMessage());
|
||||
assertFalse(e.getMessage().contains("quarantined"), "must not read like quarantine: " + e.getMessage());
|
||||
assertFalse(e.getMessage().contains("cooling off"), "must not read like cool-off: " + e.getMessage());
|
||||
assertFalse(e.getMessage().contains("weight 0"), "must not read like weight-0: " + e.getMessage());
|
||||
}
|
||||
|
||||
/**
|
||||
* Quarantine still wins when a profile is both quarantined and model-off (mirrors {@code
|
||||
* fixedThrowsWhenDefaultAndEveryCandidateQuarantined}'s priority over cooling off).
|
||||
*/
|
||||
@Test
|
||||
void fixedReportsQuarantineNotModelOffWhenBothApply() {
|
||||
PlacementPolicy policy = PlacementPolicies.fixed();
|
||||
PlacementContext ctx = new PlacementContext("b",
|
||||
List.of(PlacementCandidate.profile("a"), PlacementCandidate.profile("b")),
|
||||
noSessions(), Set.of(), Set.of("a", "b"), Set.of(), Set.of("a", "b"));
|
||||
PlacementException e = assertThrows(PlacementException.class, () -> policy.select(ctx));
|
||||
assertTrue(e.getMessage().contains("quarantined"), e.getMessage());
|
||||
assertFalse(e.getMessage().contains("turned off"),
|
||||
"quarantine takes priority over model-off in the message: " + e.getMessage());
|
||||
}
|
||||
|
||||
@Test
|
||||
void roundRobinCyclesThroughAvailableProfiles() {
|
||||
PlacementPolicy policy = PlacementPolicies.roundRobin();
|
||||
@@ -404,6 +461,93 @@ class PlacementPolicyTest {
|
||||
assertTrue(e.getMessage().contains("weight 0"), e.getMessage());
|
||||
}
|
||||
|
||||
// --- fleetd #435: FixedPlacementPolicy must consult maxLoad too, at BOTH filter sites -------
|
||||
|
||||
/**
|
||||
* The default-profile fast path must skip a capped default. Measured on 7667727 before this
|
||||
* fix: a single dev profile at {@code maxLoad: 1} with 1 live, under {@code fixed()}, still
|
||||
* returned "SPAWNED on profile=a" — the cap was advisory for every unqualified spawn.
|
||||
*/
|
||||
@Test
|
||||
void fixedSkipsCappedDefault() {
|
||||
PlacementPolicy policy = PlacementPolicies.fixed();
|
||||
PlacementContext ctx = new PlacementContext("b",
|
||||
List.of(PlacementCandidate.profile("a", 1.0f, null),
|
||||
PlacementCandidate.profile("b", 1.0f, 1)),
|
||||
name -> "b".equals(name) ? 1 : 0, Set.of(), Set.of(), Set.of());
|
||||
assertEquals("a", policy.select(ctx).profile(),
|
||||
"the default 'b' is at its maxLoad cap, so fixed falls through to the free candidate 'a'");
|
||||
}
|
||||
|
||||
/**
|
||||
* The fallback walk must skip a capped candidate too — exercised independently of the
|
||||
* default-profile fast path by using no default at all, so this is the only filter that runs.
|
||||
*/
|
||||
@Test
|
||||
void fixedFallbackWalkSkipsCappedCandidate() {
|
||||
PlacementPolicy policy = PlacementPolicies.fixed();
|
||||
PlacementContext ctx = new PlacementContext(null,
|
||||
List.of(PlacementCandidate.profile("a", 1.0f, 1),
|
||||
PlacementCandidate.profile("b", 1.0f, null)),
|
||||
name -> "a".equals(name) ? 1 : 0, Set.of(), Set.of(), Set.of());
|
||||
assertEquals("b", policy.select(ctx).profile(),
|
||||
"candidate 'a' is at its maxLoad cap, so the fallback walk skips it and picks 'b'");
|
||||
}
|
||||
|
||||
@Test
|
||||
void fixedThrowsWhenDefaultAndEveryCandidateAtCap() {
|
||||
PlacementPolicy policy = PlacementPolicies.fixed();
|
||||
PlacementContext ctx = new PlacementContext("b",
|
||||
List.of(PlacementCandidate.profile("a", 1.0f, 1),
|
||||
PlacementCandidate.profile("b", 1.0f, 1)),
|
||||
_ -> 1, Set.of(), Set.of(), Set.of());
|
||||
PlacementException e = assertThrows(PlacementException.class, () -> policy.select(ctx));
|
||||
assertTrue(e.getMessage().contains("maxLoad"), "message names the cap: " + e.getMessage());
|
||||
assertTrue(e.getMessage().contains("1 cap"), "message names the cap value: " + e.getMessage());
|
||||
assertFalse(e.getMessage().contains("quarantined"), "must not read like quarantine: " + e.getMessage());
|
||||
assertFalse(e.getMessage().contains("turned off"), "must not read like model-off: " + e.getMessage());
|
||||
}
|
||||
|
||||
/** The mirror: an uncapped default is still chosen, so the new term cannot exclude everything. */
|
||||
@Test
|
||||
void fixedStillReturnsUncappedDefault() {
|
||||
PlacementPolicy policy = PlacementPolicies.fixed();
|
||||
PlacementContext ctx = new PlacementContext("b",
|
||||
List.of(PlacementCandidate.profile("a", 1.0f, null),
|
||||
PlacementCandidate.profile("b", 1.0f, null)),
|
||||
noSessions(), Set.of(), Set.of(), Set.of());
|
||||
assertEquals("b", policy.select(ctx).profile(), "an uncapped default is returned unconditionally");
|
||||
}
|
||||
|
||||
/** CB-585: an explicit {@code maxLoad: 0} on the default caps it at zero live members. */
|
||||
@Test
|
||||
void fixedSkipsMaxLoadZeroDefaultEvenWithZeroLiveWorkers() {
|
||||
PlacementPolicy policy = PlacementPolicies.fixed();
|
||||
PlacementContext ctx = new PlacementContext("b",
|
||||
List.of(PlacementCandidate.profile("a", 1.0f, null),
|
||||
PlacementCandidate.profile("b", 1.0f, 0)),
|
||||
_ -> 0, Set.of(), Set.of(), Set.of());
|
||||
assertEquals("a", policy.select(ctx).profile(),
|
||||
"the default 'b' has maxLoad 0, so it is already at its cap with nobody live");
|
||||
}
|
||||
|
||||
/**
|
||||
* Quarantine still wins when a profile is both quarantined and at cap (mirrors {@code
|
||||
* fixedReportsQuarantineNotModelOffWhenBothApply}'s priority over model-off).
|
||||
*/
|
||||
@Test
|
||||
void fixedReportsQuarantineNotAtCapWhenBothApply() {
|
||||
PlacementPolicy policy = PlacementPolicies.fixed();
|
||||
PlacementContext ctx = new PlacementContext("b",
|
||||
List.of(PlacementCandidate.profile("a", 1.0f, 1),
|
||||
PlacementCandidate.profile("b", 1.0f, 1)),
|
||||
_ -> 1, Set.of(), Set.of("a", "b"), Set.of());
|
||||
PlacementException e = assertThrows(PlacementException.class, () -> policy.select(ctx));
|
||||
assertTrue(e.getMessage().contains("quarantined"), e.getMessage());
|
||||
assertFalse(e.getMessage().contains("maxLoad"),
|
||||
"quarantine takes priority over at-cap in the message: " + e.getMessage());
|
||||
}
|
||||
|
||||
@Test
|
||||
void unknownPolicyNameThrows() {
|
||||
assertThrows(IllegalArgumentException.class, () -> PlacementPolicies.fromName("random"));
|
||||
|
||||
Reference in New Issue
Block a user