Compare commits
25 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| ab0cc71aa4 | |||
| 4cd9046353 | |||
| 4e98a74047 | |||
| a502ba53e0 | |||
| 446cc11d1d | |||
| ddd81fe174 | |||
| ae74cc081f | |||
| cf8da1d5fa | |||
| b8aedeafcb | |||
| 3d61af6f6f | |||
| ed99c209ac | |||
| af4c88d54b | |||
| 44c735f6f5 | |||
| b540a1744b | |||
| 0f51d53098 | |||
| c670792ffe | |||
| 3982ace544 | |||
| 7180b1aad0 | |||
| c26f695402 | |||
| 48877315ca | |||
| 5c08054533 | |||
| b6db9c31f5 | |||
| e7b33fe3a0 | |||
| e3e403e5c8 | |||
| 799014e99d |
@@ -137,7 +137,7 @@ the merge — and merging on a reviewer's word is delegating it by proxy.
|
||||
| Answer a member's `fleet_ask` | `fleet_send{turnId, content}` — **not** `sessionId` |
|
||||
| Message a **peer lead** on this host | `fleet_send{sessionId: <their terminal>, content}` — `fleet_list` → `leads` reports it. Coordination only, **never** a task |
|
||||
| Message a **peer lead** on another daemon or host | `fleet_send{coordId: <their coord-id>, content}` — needs a `coordinator:` block; your own coord-id is in `fleet_list`. Coordination only, **never** a task |
|
||||
| Answer a peer lead that messaged you | `fleet_reply{content}` — the one case a lead replies |
|
||||
| Answer a peer lead that messaged you | `fleet_send{coordId}` — or `{sessionId}` if they are on this host. **Not** `fleet_reply`: it has no peer route and the publish is refused |
|
||||
| Collect a held reply | `fleet_poll{target}` · then `fleet_ack{target, msgId}` |
|
||||
| Tear down a member | `fleet_stop{paneId}` |
|
||||
|
||||
@@ -163,15 +163,21 @@ The traffic between leads is coordination and nothing else:
|
||||
3. **Verify a peer exactly as you verify yourself.** Peer status buys nothing: check the claim
|
||||
against the code, and re-run the build. A peer's correction gets the same treatment — right or
|
||||
wrong on the evidence, not on who said it. Neither of you merges the other's work unreviewed.
|
||||
**N observations are N data points only if they differ in the axis you are trusting.** This cuts
|
||||
both ways. N *failures* blamed on one cause are one data point when the cases share what you are
|
||||
not varying. N *agreeing measurements* are also one data point when they share an instrument —
|
||||
two hosts, two operators and the same formula is one formula, not two confirmations.
|
||||
4. **Ask a peer to read your project addendum.** Your addendum is instruction surface: every future
|
||||
session on your host obeys it, and a wrong one is obeyed just as faithfully as a right one. The
|
||||
author is the worst reader of their own qualifier placement — measured here, one addendum carried
|
||||
two defects and a non-author found both. If you have no peer, at least re-read it asking "which
|
||||
sentence goes false first, and would a reader reach the caveat before acting?"
|
||||
|
||||
Being messaged by a peer does not make you its worker: answer with `fleet_reply`, and push back on
|
||||
the substance if it is wrong. A peer that simply complies has thrown away the reason there are two of
|
||||
you.
|
||||
Being messaged by a peer does not make you its worker: answer the way you would open —
|
||||
`fleet_send{coordId}` for another daemon, `fleet_send{sessionId}` on this host — and push back on
|
||||
the substance if it is wrong. `fleet_reply` resolves a member's blocked `fleet_send`; a peer's
|
||||
coord-id message is durable and non-blocking, so there is nothing for it to resolve. A peer that
|
||||
simply complies has thrown away the reason there are two of you.
|
||||
|
||||
### Member (worker or architect) — the turn contract
|
||||
|
||||
|
||||
@@ -877,3 +877,32 @@ guard:
|
||||
# terminal: term_65619bd6174568
|
||||
# pushReminders: 5
|
||||
# pushBackoffMs: 15000
|
||||
|
||||
# Central allow-list of models any profiles: entry may name. Nothing checked a profile's model:
|
||||
# value before this block existed — it was a free-form string handed straight to the backend
|
||||
# adapter, and a withdrawn or misspelled name failed silently instead of at config load (opencode
|
||||
# falls back to a default model rather than erroring on an unknown -m).
|
||||
#
|
||||
# Absent, or present with an empty allow:, is OFF: no profile's model: is checked, exactly like
|
||||
# before this block existed. fleetd.yaml is gitignored on every host, so an upgrade must not force
|
||||
# every operator to enumerate their models before the daemon will start.
|
||||
#
|
||||
# The list is the authority; profiles: is checked against it, never the reverse — adding or
|
||||
# editing a profiles: entry cannot, by itself, widen what is permitted here.
|
||||
#
|
||||
# Enforcement is at CONFIG LOAD only (a bad model: fails the daemon at startup, naming both the
|
||||
# model and the profile). There is no spawn-time enforcement, no runtime on/off switch, and no
|
||||
# interaction with BackendQuarantine — those are separate, later units.
|
||||
#
|
||||
# allow → the permitted models. Each entry is its own block (not a bare string) so a later unit
|
||||
# can add an on/off state or a load limit per model without changing this shape.
|
||||
# model → the model id exactly as a profiles: entry's model: field would write it. One flat,
|
||||
# opaque-string namespace: a bare Claude id (claude-sonnet-5) and an opencode
|
||||
# provider-prefixed id (openai/gpt-5.6-terra) both fit here unchanged — the check is a
|
||||
# plain string match, never a parse of the provider prefix or a branch on kind:.
|
||||
# models:
|
||||
# allow:
|
||||
# - model: claude-sonnet-5
|
||||
# - model: claude-opus-5
|
||||
# - model: openai/gpt-5.6-terra
|
||||
# - model: amazon.nova-pro-v1:0
|
||||
|
||||
@@ -134,6 +134,10 @@ public final class Fleetd {
|
||||
// secret is reported above, so upgrading past this commit never silently drops CB-592's
|
||||
// protection.
|
||||
reportMemberCredentialsGap(cfg);
|
||||
// fleetd #395: an unset exhaustedPattern is a silent opt-out of usage-limit detection for
|
||||
// that profile — say so loudly, the same way the two reports above do, rather than let an
|
||||
// operator discover it only when a limit goes undetected.
|
||||
reportExhaustedPatternGap(cfg);
|
||||
// CB-559: `cfg` stays the startup snapshot — every validation and every piece of one-time
|
||||
// wiring below reads it, and must, because those decisions cannot be unmade. `config` is the
|
||||
// live reference the hot paths read per use. Which keys can actually move is ConfigRef's
|
||||
@@ -144,19 +148,16 @@ public final class Fleetd {
|
||||
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
||||
guard.assertPrimaryClean(System.getenv());
|
||||
|
||||
// CB-501: refuse to start if the bind is wider than the auth mode can defend. Under
|
||||
// loopback-trust, "not a known worker" means "the primary" — sound only because the OS
|
||||
// refuses remote connections to a loopback socket. This throws rather than warns so the
|
||||
// dangerous configuration cannot be reached by ignoring a log line.
|
||||
cfg.validateAuthExposure();
|
||||
cfg.validateLeadTabPrefixes();
|
||||
// CB-542: a subscription:true profile whose env: reseats ANTHROPIC_BASE_URL/AUTH_TOKEN would
|
||||
// reach an unguarded endpoint (the launcher skips SubscriptionGuard for it). Refuse at load.
|
||||
cfg.validateSubscriptionProfiles();
|
||||
cfg.validateCharters();
|
||||
// CB-548: every architect slot must name a configured workers: profile — the strong-model
|
||||
// backend the future spawn lifecycle would read. A stale reference dies here, not later.
|
||||
cfg.validateMembers();
|
||||
// Every FleetConfig.validateXxx() the operator's config can fail — CB-501's auth-exposure
|
||||
// check, CB-531's lead-tab-prefix check, CB-542's subscription-profile check, the charter
|
||||
// and member-slot checks, and the "central allow-list of usable models" check — must run
|
||||
// here, at load, before anything below opens a socket or spawns a member. fleetd ticket
|
||||
// "central allow-list of usable models" follow-up: mutation testing found six individual
|
||||
// calls here with nothing proving any of them still ran (deleting one left the full suite
|
||||
// green). validateAll() replaces them with the one call that FleetConfigValidateAllTest
|
||||
// and the Fleetd-startup tests actually pin — see FleetConfig#validateAll's javadoc for
|
||||
// why a name-by-name list here would have the same defect it replaces.
|
||||
cfg.validateAll();
|
||||
|
||||
Path socket = cfg.herdrSocket() != null && !cfg.herdrSocket().isBlank()
|
||||
? Path.of(cfg.herdrSocket())
|
||||
@@ -661,10 +662,8 @@ public final class Fleetd {
|
||||
// built a second time — two independently-constructed sources reading the SAME BackendQuarantine
|
||||
// / BackendOutagePolicy would still be able to drift (e.g. a future edit to the credentialIdFor
|
||||
// closure in only one of the two places), exactly the shape #284 was.
|
||||
FleetMcp.QuarantineSource quarantineSource = new FleetMcp.QuarantineSource(profile -> {
|
||||
var configured = config.get().profiles().get(profile);
|
||||
return configured == null ? null : configured.effectiveCredentialId();
|
||||
}, quarantine);
|
||||
FleetMcp.QuarantineSource quarantineSource = quarantineSource(config, quarantine,
|
||||
exhaustedPatternsByProfile);
|
||||
FleetMcp.OutageSource outageSource = new FleetMcp.OutageSource(profile -> {
|
||||
var configured = config.get().profiles().get(profile);
|
||||
return configured == null ? null : configured.effectiveCredentialId();
|
||||
@@ -798,6 +797,19 @@ public final class Fleetd {
|
||||
return target -> presence.isPresent(target) || leads.get().containsKey(target);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #404: production source for quarantine reporting. Credential IDs are hot, but
|
||||
* exhausted patterns are compiled once at startup for {@link CompletionResolver}, so the armed
|
||||
* field must use that same compiled map until restart.
|
||||
*/
|
||||
static FleetMcp.QuarantineSource quarantineSource(ConfigRef config, BackendQuarantine quarantine,
|
||||
Map<String, Pattern> startupExhaustedPatterns) {
|
||||
return new FleetMcp.QuarantineSource(profile -> {
|
||||
var configured = config.get().profiles().get(profile);
|
||||
return configured == null ? null : configured.effectiveCredentialId();
|
||||
}, quarantine, profile -> startupExhaustedPatterns.containsKey(profile));
|
||||
}
|
||||
|
||||
/**
|
||||
* 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).
|
||||
@@ -1323,6 +1335,52 @@ public final class Fleetd {
|
||||
+ "fleetd.yaml — see fleetd.example.yaml — and restart.");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #395: {@code exhaustedPattern} (see {@link FleetConfig.Profile#exhaustedPattern}) is
|
||||
* deliberately opt-in — {@code null}/blank means a backend refusal on that profile is never
|
||||
* classified as {@code BACKEND_EXHAUSTED}, so its credential is never quarantined. That is a
|
||||
* legitimate choice (guessing the vendor's wording would be worse), but an operator who never
|
||||
* opted a profile in should not discover the gap only when a usage limit silently goes
|
||||
* undetected. Warn once at startup, naming every unarmed profile, exactly like {@link
|
||||
* #reportMemberCredentialsGap} — never refuse to start over it.
|
||||
*
|
||||
* <p>A {@code subscription: true} profile that is unarmed gets a SECOND, louder WARN of its
|
||||
* own: it bills the operator's metered Claude plan, the case where an undetected usage limit
|
||||
* costs the most.
|
||||
*
|
||||
* <p>Package-private so a test can capture the real log via a {@link
|
||||
* ch.qos.logback.core.read.ListAppender}, the same pattern {@link
|
||||
* #reportMemberCredentialsGap}'s own test uses.
|
||||
*/
|
||||
static void reportExhaustedPatternGap(FleetConfig cfg) {
|
||||
List<String> unarmedSubscription = new ArrayList<>();
|
||||
List<String> unarmedOther = new ArrayList<>();
|
||||
cfg.profiles().forEach((name, profile) -> {
|
||||
if (!profile.hasExhaustedPattern()) {
|
||||
(profile.isSubscription() ? unarmedSubscription : unarmedOther).add(name);
|
||||
}
|
||||
});
|
||||
if (unarmedSubscription.isEmpty() && unarmedOther.isEmpty()) {
|
||||
log.info("exhaustedPattern: every configured profile has usage-limit detection armed");
|
||||
return;
|
||||
}
|
||||
List<String> allUnarmed = new ArrayList<>(unarmedSubscription);
|
||||
allUnarmed.addAll(unarmedOther);
|
||||
allUnarmed = allUnarmed.stream().sorted().toList();
|
||||
log.warn("exhaustedPattern: profile(s) {} have no exhaustedPattern configured — a "
|
||||
+ "usage-limit refusal on any of them is never detected and never "
|
||||
+ "quarantines its credential. Set exhaustedPattern (see "
|
||||
+ "fleetd.example.yaml) to arm detection for a profile.",
|
||||
allUnarmed);
|
||||
if (!unarmedSubscription.isEmpty()) {
|
||||
List<String> sortedSubscription = unarmedSubscription.stream().sorted().toList();
|
||||
log.warn("exhaustedPattern: subscription profile(s) {} run on the operator's metered "
|
||||
+ "Claude plan and have NO usage-limit detection armed — this is the "
|
||||
+ "case where a missed usage limit costs the most.",
|
||||
sortedSubscription);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Poll herdr's {@code ping} until it answers or {@link #HERDR_WAIT_SECONDS} elapses (CB-504).
|
||||
*
|
||||
|
||||
@@ -43,6 +43,12 @@ 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
|
||||
@@ -135,8 +141,9 @@ 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, and again after {@code idleSleepGuard:} was added.</strong>
|
||||
* {@code FleetConfig} has 24 top-level record components: 5 cold, 13 deferred, 3 split, 3
|
||||
* 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()}
|
||||
@@ -219,7 +226,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");
|
||||
"quarantineCooldownSeconds", "profiles", "idleSleepGuard", "models");
|
||||
|
||||
private final Path path;
|
||||
private final AtomicReference<FleetConfig> current;
|
||||
@@ -322,12 +329,12 @@ public final class ConfigRef implements Supplier<FleetConfig> {
|
||||
fresh = FleetConfig.load(path);
|
||||
// The same gate startup runs. A config that would have refused to boot must not be able
|
||||
// to slip in through a reload — that is how a daemon ends up in a state it could never
|
||||
// have started in, which is the hardest kind to debug.
|
||||
fresh.validateAuthExposure();
|
||||
fresh.validateLeadTabPrefixes();
|
||||
fresh.validateSubscriptionProfiles();
|
||||
fresh.validateCharters();
|
||||
fresh.validateMembers();
|
||||
// have started in, which is the hardest kind to debug. fleetd ticket "central allow-list
|
||||
// of usable models" follow-up: this used to be six individual validateXxx() calls, and
|
||||
// mutation testing found two of the six unpinned here even though startup pinned nothing
|
||||
// at all — see FleetConfig#validateAll's javadoc for why the fix is one reflective call,
|
||||
// not a longer hand-maintained list.
|
||||
fresh.validateAll();
|
||||
} catch (RuntimeException e) {
|
||||
String msg = e.getMessage() == null ? e.toString() : e.getMessage();
|
||||
log.warn("config reload from {} refused, keeping the running config: {}", path, msg);
|
||||
@@ -439,6 +446,15 @@ 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*");
|
||||
|
||||
@@ -16,11 +16,15 @@ import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.io.UncheckedIOException;
|
||||
import java.lang.reflect.InvocationTargetException;
|
||||
import java.lang.reflect.Method;
|
||||
import java.lang.reflect.Modifier;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Arrays;
|
||||
import java.util.Collections;
|
||||
import java.util.Comparator;
|
||||
import java.util.HashSet;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
@@ -124,6 +128,13 @@ import java.util.regex.PatternSyntaxException;
|
||||
* under a member's long turn. {@code null} (the block omitted) behaves the
|
||||
* same as an explicit {@code enabled: true}; set {@code enabled: false} to
|
||||
* turn it off. See {@link dev.ltms.fleet.power.IdleSleepGuard}.
|
||||
* @param models central allow-list of models any {@code profiles:} entry may name. {@code
|
||||
* null} or an empty {@code allow:} ⇒ off: {@link #validateModels()} checks
|
||||
* 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}.
|
||||
*/
|
||||
@JsonIgnoreProperties(ignoreUnknown = true)
|
||||
public record FleetConfig(
|
||||
@@ -150,7 +161,22 @@ public record FleetConfig(
|
||||
String worktreeGroup,
|
||||
String memberLoginShell,
|
||||
String memberSkills,
|
||||
IdleSleepGuard idleSleepGuard) {
|
||||
IdleSleepGuard idleSleepGuard,
|
||||
Models models) {
|
||||
|
||||
/** Back-compat form before the {@code models:} block was added. */
|
||||
public FleetConfig(Bind bind, String herdrSocket, String memberHerdrSocket, Map<String, Profile> profiles,
|
||||
Guard guard, String worktreeRoot, Lifecycle lifecycle, Integer spawnReadyTimeoutMs,
|
||||
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
|
||||
LeadHeartbeat leadHeartbeat, Health health, String placement, Auth auth,
|
||||
ConfigReload configReload, Integer quarantineCooldownSeconds,
|
||||
MemberCredentials memberCredentials, Coordinator coordinator, String worktreeGroup,
|
||||
String memberLoginShell, String memberSkills, IdleSleepGuard idleSleepGuard) {
|
||||
this(bind, herdrSocket, memberHerdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
|
||||
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth,
|
||||
configReload, quarantineCooldownSeconds, memberCredentials, coordinator, worktreeGroup,
|
||||
memberLoginShell, memberSkills, idleSleepGuard, null);
|
||||
}
|
||||
|
||||
/** Back-compat form before the {@code idleSleepGuard:} block was added. */
|
||||
public FleetConfig(Bind bind, String herdrSocket, String memberHerdrSocket, Map<String, Profile> profiles,
|
||||
@@ -1318,6 +1344,72 @@ public record FleetConfig(
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Central allow-list of models any {@code profiles:} entry may name (fleetd ticket: "a central
|
||||
* allow-list of usable models"). Nothing before this block checked a profile's {@code model:}
|
||||
* against anything — it was a free-form string handed straight to the backend adapter, and a
|
||||
* withdrawn or misspelled name failed silently (opencode falls back to a default model rather
|
||||
* than erroring on an unknown {@code -m}) rather than at config load, where a mistake is cheap.
|
||||
*
|
||||
* <p><b>Absent or empty {@code allow:} is "off"</b>, on purpose: this is a large deployment
|
||||
* with gitignored {@code fleetd.yaml} on more than one host, and a change that forced every
|
||||
* operator to enumerate their models before the daemon would start would break every one of
|
||||
* them on upgrade. See {@link FleetConfig#validateModels()}, which is where the allow-list is
|
||||
* actually enforced, at config load.
|
||||
*
|
||||
* <p><b>The list is the authority; a {@code profiles:} entry is checked against it, never the
|
||||
* other way around.</b> Adding or editing a {@code profiles:} entry cannot, by itself, widen
|
||||
* 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.
|
||||
*
|
||||
* @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
|
||||
* treated as absent: {@link FleetConfig#validateModels()} checks nothing.
|
||||
*/
|
||||
@JsonIgnoreProperties(ignoreUnknown = true)
|
||||
public record Models(List<ModelEntry> allow) {
|
||||
|
||||
public Models {
|
||||
allow = (allow == null) ? List.of() : List.copyOf(allow);
|
||||
}
|
||||
|
||||
/**
|
||||
* 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:}.
|
||||
*
|
||||
* @param model the model id exactly as a {@code profiles:} entry's {@code model:} field
|
||||
* would name it
|
||||
*/
|
||||
@JsonIgnoreProperties(ignoreUnknown = true)
|
||||
public record ModelEntry(String model) {
|
||||
public ModelEntry {
|
||||
model = (model == null || model.isBlank()) ? null : model.trim();
|
||||
}
|
||||
}
|
||||
|
||||
/** {@link #allow}'s model ids, as a set for membership checks. Blank/null entries are dropped. */
|
||||
public Set<String> ids() {
|
||||
Set<String> ids = new java.util.LinkedHashSet<>();
|
||||
for (ModelEntry e : allow) {
|
||||
if (e != null && e.model() != null) {
|
||||
ids.add(e.model());
|
||||
}
|
||||
}
|
||||
return Collections.unmodifiableSet(ids);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* The terminal → lead-name map seeded from the legacy singular {@code primary:} pin (CB-530).
|
||||
*
|
||||
@@ -1569,7 +1661,7 @@ public record FleetConfig(
|
||||
"lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs", "broker", "primary", "fleet",
|
||||
"leadHeartbeat", "health", "placement", "auth", "configReload", "quarantineCooldownSeconds",
|
||||
"memberCredentials", "coordinator", "worktreeGroup", "memberLoginShell", "memberSkills",
|
||||
"idleSleepGuard");
|
||||
"idleSleepGuard", "models");
|
||||
|
||||
/** Load and validate config from {@code path}. */
|
||||
public static FleetConfig load(Path path) {
|
||||
@@ -2254,10 +2346,15 @@ public record FleetConfig(
|
||||
// enabled: true — see its javadoc), so defaulting the block here would change nothing a
|
||||
// reader observes and would only obscure that "block omitted" and "block present and
|
||||
// enabled" are deliberately the same outcome.
|
||||
// models is left as-is, like broker/primary/coordinator above: null/empty is "off", and an
|
||||
// absent block must validate nothing (see Models's javadoc) — defaulting it here to an
|
||||
// empty Models would be a no-op for validateModels() either way, since an empty allow-list
|
||||
// already means "check nothing", so there is nothing to gain and one more null check to
|
||||
// avoid by leaving it exactly as configured.
|
||||
return new FleetConfig(b, herdrSocket, memberHerdrSocket, profiles, g, worktreeRoot, l, timeout, pollMs,
|
||||
broker, primary, f, leadHeartbeat, health, placementOrDefault, a, configReload,
|
||||
quarantineCooldown, mc, coordinator, worktreeGroup, memberLoginShell, memberSkills,
|
||||
idleSleepGuard);
|
||||
idleSleepGuard, models);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -2477,6 +2574,118 @@ public record FleetConfig(
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Reject a {@code profiles:} entry whose {@code model:} is not on the configured {@link
|
||||
* #models} allow-list.
|
||||
*
|
||||
* <p>Absent or empty {@code models.allow:} validates nothing — see {@link Models}'s javadoc:
|
||||
* an existing config with no such block must keep working exactly as it does today. Once the
|
||||
* operator declares at least one entry, every profile's {@code model:} (when set — a profile
|
||||
* may legitimately leave it {@code null}, e.g. a {@code subscription: true} profile relying on
|
||||
* the account's own default) must equal one of {@link Models#ids()} exactly. The comparison is
|
||||
* a flat string match: a bare Claude id and an opencode {@code provider/model} id are both
|
||||
* just opaque strings here, so nothing here needs to know which {@code kind:} a profile runs.
|
||||
*
|
||||
* <p>The check runs one direction only, by construction: it reads {@link #profiles} and
|
||||
* {@link #models}, and only ever adds to {@code bad} when a profile's model is missing from
|
||||
* the allow-list. Nothing here can be satisfied by widening a {@code profiles:} entry — only
|
||||
* editing {@code models.allow:} itself changes what passes. That is the invariant the ticket
|
||||
* asked for: the list is the authority, profiles are checked against it.
|
||||
*
|
||||
* @throws IllegalStateException when any profile names a model outside the configured
|
||||
* allow-list, naming both the model and the profile that wanted it
|
||||
*/
|
||||
public void validateModels() {
|
||||
if (models == null || models.allow().isEmpty()) {
|
||||
return;
|
||||
}
|
||||
Set<String> allowed = models.ids();
|
||||
List<String> bad = new ArrayList<>();
|
||||
profiles.forEach((name, p) -> {
|
||||
String model = p.model();
|
||||
if (model != null && !model.isBlank() && !allowed.contains(model)) {
|
||||
bad.add("profile '" + name + "' names model '" + model + "', which is not in "
|
||||
+ "models.allow: (have: " + allowed + ").");
|
||||
}
|
||||
});
|
||||
if (!bad.isEmpty()) {
|
||||
throw new IllegalStateException("refusing to start: " + String.join(" ", bad));
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Runs every validator this class declares — found by reflection, not by name.
|
||||
*
|
||||
* <p>fleetd ticket "central allow-list of usable models", follow-up: mutation testing found
|
||||
* that although each of the six validators above was well pinned on its own, nothing proved
|
||||
* either real caller ({@code Fleetd.main} and {@link ConfigRef#reload()}) still
|
||||
* invoked it — deleting a call site left the full suite green. The fix is not a seventh test
|
||||
* per caller; a hand-maintained list of six names here would have the exact same defect its
|
||||
* own javadoc would warn against: the seventh validator someone adds next month has no reason
|
||||
* to be added to it. So this method does not name any validator. It sweeps {@link
|
||||
* #getClass()}'s own public, no-argument, {@code void} methods whose name starts with {@code
|
||||
* "validate"} (excluding itself) and invokes every one it finds, via {@link
|
||||
* #invokeAllValidators}. A new {@code validateXxx()} method is therefore wired into both
|
||||
* callers the moment it is written — there is no second step to forget, and so no state in
|
||||
* which it silently never runs.
|
||||
*
|
||||
* <p>{@code Fleetd.main} and {@link ConfigRef#reload()} each call this one method instead of
|
||||
* the six individually — see the comments at those two call sites for why
|
||||
* each must run it.
|
||||
*
|
||||
* <p>Methods run in a fixed (alphabetical) order, so a config with more than one violation
|
||||
* always names the same one first, on every run.
|
||||
*
|
||||
* @throws IllegalStateException (or whatever unchecked exception a validator itself throws),
|
||||
* propagated unchanged from the first validator, in that order,
|
||||
* that finds a problem
|
||||
*/
|
||||
public void validateAll() {
|
||||
invokeAllValidators(this);
|
||||
}
|
||||
|
||||
/**
|
||||
* The reflective sweep behind {@link #validateAll()}, kept as its own method — taking any
|
||||
* {@code target}, not just {@code this} — so a test can prove the MECHANISM is generic (it
|
||||
* would sweep a seventh {@code validateXxx()} method added to any class, not just something
|
||||
* special-cased to today's six on {@link FleetConfig}) without needing to add a real, unwanted
|
||||
* seventh validator to this class just to exercise that claim. See {@code
|
||||
* FleetConfigValidateAllTest} for that proof.
|
||||
*
|
||||
* @param target an object whose public, no-argument, {@code void} methods named {@code
|
||||
* validateXxx} (any name starting with {@code "validate"}, excluding {@code
|
||||
* validateAll} itself) should all run, in alphabetical-by-name order
|
||||
*/
|
||||
static void invokeAllValidators(Object target) {
|
||||
List<Method> methods = new ArrayList<>();
|
||||
for (Method m : target.getClass().getMethods()) {
|
||||
if (Modifier.isPublic(m.getModifiers())
|
||||
&& m.getParameterCount() == 0
|
||||
&& m.getReturnType() == void.class
|
||||
&& m.getName().startsWith("validate")
|
||||
&& !m.getName().equals("validateAll")) {
|
||||
methods.add(m);
|
||||
}
|
||||
}
|
||||
methods.sort(Comparator.comparing(Method::getName));
|
||||
for (Method m : methods) {
|
||||
try {
|
||||
m.invoke(target);
|
||||
} catch (InvocationTargetException e) {
|
||||
Throwable cause = e.getCause();
|
||||
if (cause instanceof RuntimeException re) {
|
||||
throw re;
|
||||
}
|
||||
if (cause instanceof Error err) {
|
||||
throw err;
|
||||
}
|
||||
throw new IllegalStateException("validator " + m.getName() + " failed", cause);
|
||||
} catch (IllegalAccessException e) {
|
||||
throw new IllegalStateException("cannot invoke validator " + m.getName(), e);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** True for the loopback addresses and the unspecified-but-local forms we treat as same-host. */
|
||||
private static boolean isLoopbackBind(String host) {
|
||||
if (host == null || host.isBlank()) {
|
||||
|
||||
@@ -119,9 +119,28 @@ public final class FleetMcp {
|
||||
/**
|
||||
* CB-578 stage B quarantine facts used by {@code fleet_profiles}: a profile → credential id
|
||||
* lookup, plus the shared {@link BackendQuarantine} to read remaining cooldowns off.
|
||||
*
|
||||
* @param exhaustedPatternArmed fleetd #395: profile → whether that profile's {@code
|
||||
* exhaustedPattern} is configured (see {@code
|
||||
* FleetConfig.Profile#hasExhaustedPattern}), i.e. whether a backend refusal
|
||||
* on it can EVER be classified {@code BACKEND_EXHAUSTED} and quarantine its
|
||||
* credential. Bundled here, not a separate Source, because it answers the
|
||||
* exact question {@code fleet_profiles}'s quarantine facts already answer
|
||||
* for a QUARANTINED profile — "can this profile's usage limit ever be
|
||||
* caught?" — just for every profile, not only one currently caught.
|
||||
*/
|
||||
public record QuarantineSource(Function<String, String> credentialIdFor, BackendQuarantine quarantine) {
|
||||
/** Inert source — no profile is ever reported quarantined. Explicit stand-in, not a default. */
|
||||
public record QuarantineSource(Function<String, String> credentialIdFor, BackendQuarantine quarantine,
|
||||
Function<String, Boolean> exhaustedPatternArmed) {
|
||||
/**
|
||||
* Backward-compatible 2-arg form, before fleetd #395 added {@code exhaustedPatternArmed} —
|
||||
* reports every profile unarmed. Keeps every pre-existing call site (production and test)
|
||||
* compiling and behaving identically for the quarantine facts they actually asked for.
|
||||
*/
|
||||
public QuarantineSource(Function<String, String> credentialIdFor, BackendQuarantine quarantine) {
|
||||
this(credentialIdFor, quarantine, _ -> false);
|
||||
}
|
||||
|
||||
/** Inert source — no profile is ever reported quarantined or armed. Explicit stand-in, not a default. */
|
||||
public static QuarantineSource none() { return new QuarantineSource(_ -> null, BackendQuarantine.none()); }
|
||||
}
|
||||
|
||||
@@ -346,7 +365,7 @@ public final class FleetMcp {
|
||||
String self = callerTerminal(exchange);
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_reply", req.arguments()), self);
|
||||
if (denied != null) return denied;
|
||||
return reply(messages, self, str(req.arguments(), "content"));
|
||||
return reply(messages, self, principal(exchange).role(), str(req.arguments(), "content"));
|
||||
};
|
||||
// fleet_ask (CB-205): a worker's mid-turn question — identity from the CONNECTION.
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> askHandler =
|
||||
@@ -868,19 +887,26 @@ public final class FleetMcp {
|
||||
/**
|
||||
* {@code fleet_reply}: the worker returns its structured answer, resolving the awaiting send
|
||||
* or — when no send is open — queueing the reply in the inbox for later drain (CB-307).
|
||||
* {@code callerTerminal} is resolved from the connection (never an argument); a {@code null}
|
||||
* means the caller is not a known worker (e.g. the primary called it by mistake).
|
||||
* {@code callerTerminal} and {@code callerRole} are resolved from the connection (never an
|
||||
* argument). A {@code null} terminal means the caller is not a known worker. A PRIMARY with a
|
||||
* terminal is a lead and must use {@code fleet_send}, because reply has no peer-lead route.
|
||||
*
|
||||
* <p>fleetd #365: the result text names which of those actually happened
|
||||
* ({@link MessageService.ReplyOutcome#description()}) instead of the single word "delivered"
|
||||
* for both — a queued reply is a real success, but it is not the same fact as one that resolved
|
||||
* a live waiter, and the caller could not previously tell them apart.
|
||||
*/
|
||||
static McpSchema.CallToolResult reply(MessageService messages, String callerTerminal, String content) {
|
||||
static McpSchema.CallToolResult reply(MessageService messages, String callerTerminal, Role callerRole, String content) {
|
||||
if (callerTerminal == null) {
|
||||
return error("fleet_reply is for workers only — could not identify the calling worker "
|
||||
+ "from the connection");
|
||||
}
|
||||
if (callerRole == Role.PRIMARY) {
|
||||
return error("fleet_reply has no route to a peer lead. Use fleet_send{coordId: ...} for a peer on another "
|
||||
+ "daemon or fleet_send{sessionId: ...} for a peer on this host. fleet_reply resolves a member's "
|
||||
+ "blocked fleet_send, and a peer's coord-id message is durable and non-blocking, so there is "
|
||||
+ "nothing for it to resolve.");
|
||||
}
|
||||
// fleetd #302: isBlank, not == null, to match fleet_send's own guard above. MessageService
|
||||
// .reply now REJECTS blank content, and this handler is a bare BiFunction with no try/catch
|
||||
// around it — so a whitespace-only fleet_reply would leave here as an uncaught
|
||||
@@ -1109,6 +1135,14 @@ public final class FleetMcp {
|
||||
* not the other, and the two then disagree about a live outage. That is exactly what fleetd
|
||||
* #284 was, where one rule computed in two places was widened in only one and a single response
|
||||
* contradicted itself. Shared inputs do not make duplicated computation safe.
|
||||
*
|
||||
* <p>fleetd #395: also reports {@code exhaustionDetectionArmed}, one boolean per configured
|
||||
* profile — {@code true} when that profile's {@code exhaustedPattern} is set, {@code false}
|
||||
* when it is not, so an operator can tell "this profile is healthy" from "nothing can ever
|
||||
* quarantine this profile" without reading {@code fleetd.yaml}. Unlike {@code quarantined}/
|
||||
* {@code coolingOff}, this map always names every profile: an unarmed profile never enters a
|
||||
* transient state to be absent from, so silence here would read as "healthy" rather than "not
|
||||
* being watched at all".
|
||||
*/
|
||||
public static Map<String, Object> profilesView(PeerLauncher workers, QuarantineSource quarantine, OutageSource outage) {
|
||||
Map<String, Object> result = new LinkedHashMap<>();
|
||||
@@ -1116,7 +1150,9 @@ public final class FleetMcp {
|
||||
result.put("default", workers.defaultProfile() == null ? "" : workers.defaultProfile());
|
||||
Map<String, Object> quarantined = new LinkedHashMap<>();
|
||||
Map<String, Object> coolingOff = new LinkedHashMap<>();
|
||||
Map<String, Object> exhaustionDetectionArmed = new LinkedHashMap<>();
|
||||
for (String profile : workers.profiles()) {
|
||||
exhaustionDetectionArmed.put(profile, quarantine.exhaustedPatternArmed().apply(profile));
|
||||
String credentialId = quarantine.credentialIdFor().apply(profile);
|
||||
if (credentialId != null) {
|
||||
quarantine.quarantine().remainingSeconds(credentialId).ifPresent(remaining -> {
|
||||
@@ -1136,6 +1172,7 @@ public final class FleetMcp {
|
||||
});
|
||||
}
|
||||
}
|
||||
result.put("exhaustionDetectionArmed", exhaustionDetectionArmed);
|
||||
if (!quarantined.isEmpty()) {
|
||||
result.put("quarantined", quarantined);
|
||||
}
|
||||
@@ -1663,7 +1700,11 @@ public final class FleetMcp {
|
||||
"List the configured worker profiles (backends) and which one fleet_spawn uses by "
|
||||
+ "default. A 'quarantined' map is present when a backend-exhausted refusal put "
|
||||
+ "a profile's credential on cooldown — fleet_spawn onto it is refused until "
|
||||
+ "quarantinedForSeconds elapses; a profile sharing that credential is listed too.",
|
||||
+ "quarantinedForSeconds elapses; a profile sharing that credential is listed too. "
|
||||
+ "'exhaustionDetectionArmed' reports, per profile, whether a usage-limit refusal "
|
||||
+ "on it can EVER be classified and quarantined (its exhaustedPattern is "
|
||||
+ "configured) — false means that profile's credential can never be quarantined "
|
||||
+ "by this mechanism, however many usage-limit refusals it sees.",
|
||||
objectSchema(Map.of(), List.of()));
|
||||
}
|
||||
|
||||
|
||||
@@ -75,9 +75,39 @@ import java.util.stream.Stream;
|
||||
* never to the control.
|
||||
*
|
||||
* <p>The scrub also writes {@code scrub-report.txt} into its own directory: one {@code allowed N of
|
||||
* M} line (N = exports left untouched, M = exports present when the scrub ran), then the blanked
|
||||
* NAMES — never values. The launcher reads this back at teardown and logs it, because a blocked
|
||||
* count next to an unknown denominator is not a finding.
|
||||
* M failed F} line (N = exports left untouched, M = exports present when the scrub ran, F = names
|
||||
* the scrub attempted to blank but could not), then the NAMES — blanked ones bare, unblankable ones
|
||||
* {@code !}-prefixed — never values. The launcher reads this back at teardown and logs it, because a
|
||||
* blocked count next to an unknown denominator is not a finding.
|
||||
*
|
||||
* <p><b>fleetd #394:</b> plain {@code export "$n="} is a FATAL error for a zsh read-only or special
|
||||
* parameter (for example {@code UID}) — it aborts the whole sourced file, so every name still to
|
||||
* come is never blanked and the report above is never written at all. The blanking loop instead
|
||||
* routes each attempt through {@code eval}, which contains that error to the single iteration: the
|
||||
* loop always finishes, and a name that could not be blanked is counted as {@code failed} and
|
||||
* listed {@code !}-prefixed rather than silently disappearing. This is deliberately not a skip-list
|
||||
* of known-bad names — every enumerated name is still attempted, so a name nobody has thought of
|
||||
* yet still gets tried and, if it fails, still gets counted.
|
||||
*
|
||||
* <p>The blanking loop also re-asserts, on its own, the same {@code [A-Za-z_][A-Za-z0-9_]*} shape
|
||||
* check the enumeration loop already applied. Before {@code eval} was introduced a non-conforming
|
||||
* name reaching {@code export "$n="} was harmless either way — the quoting made it inert. With
|
||||
* {@code eval}, the name is spliced into a string and interpreted as shell syntax, so the enumeration
|
||||
* loop's check is no longer sufficient on its own to keep that call site safe — it is a guard on a
|
||||
* different loop, and the two must not silently drift apart. Re-checking right before the
|
||||
* {@code eval} keeps that call site safe by its own reading, independent of whatever the enumeration
|
||||
* loop does or stops doing in a later change.
|
||||
*
|
||||
* <p><b>fleetd #400:</b> {@code eval}'s exit status is not proof that the blank actually happened.
|
||||
* zsh coerces a bare {@code NAME=} assignment on an integer special parameter (measured on macOS zsh
|
||||
* 5.9: {@code SECONDS}, {@code RANDOM}, {@code SHLVL}, {@code HISTSIZE}, {@code COLUMNS},
|
||||
* {@code LINES}, {@code USERNAME}) to a number instead of failing — {@code eval} returns success,
|
||||
* the value is untouched, and a status-based classification reports it as blanked when it was not.
|
||||
* The fix classifies on the observed effect instead: after the attempt, the name's value is read
|
||||
* back with the {@code (P)} indirection flag and the decision is made from whether that is now
|
||||
* empty. This one check covers all three shapes a name can take at this point — a genuine blank, a
|
||||
* fatal read-only error {@code eval} merely contained, and this silent no-op — and the exit status
|
||||
* plays no part in the decision at all.
|
||||
*/
|
||||
public final class EnvAllowListScrub {
|
||||
|
||||
@@ -132,10 +162,16 @@ public final class EnvAllowListScrub {
|
||||
}
|
||||
|
||||
/**
|
||||
* A parsed {@code scrub-report.txt}: how many exported variables existed when the scrub ran,
|
||||
* how many were left untouched (allowed), and the NAMES that were blanked. Values never appear.
|
||||
* A parsed {@code scrub-report.txt}: how many exported variables existed when the scrub ran
|
||||
* ({@code total}), how many were left untouched ({@code allowed}), how many the scrub attempted
|
||||
* to blank but could not ({@code failed} — fleetd #394: a zsh read-only/special parameter such
|
||||
* as {@code UID} fatally errors on plain {@code export NAME=}, so those attempts go through
|
||||
* {@code eval} instead so the loop keeps going and the failure is counted rather than left
|
||||
* invisible), and the NAMES in each of the latter two categories. {@code allowed +
|
||||
* blanked.size() + unblankable.size() == total}, and {@code unblankable.size() == failed}.
|
||||
* Values never appear.
|
||||
*/
|
||||
record ScrubReport(int allowed, int total, List<String> blanked) {
|
||||
record ScrubReport(int allowed, int total, int failed, List<String> blanked, List<String> unblankable) {
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -334,15 +370,57 @@ public final class EnvAllowListScrub {
|
||||
_cb633_blank+=("$_cb633_n")
|
||||
done
|
||||
|
||||
{ for _cb633_n in "${_cb633_blank[@]}"; do export "$_cb633_n="; done; } 2>/dev/null
|
||||
# fleetd #394: plain `export "$n="` is FATAL for a zsh read-only/special parameter
|
||||
# (e.g. UID) and aborts this whole sourced file — every name still to come is never
|
||||
# blanked, and the report below is never written, silently. `eval` contains that
|
||||
# error to the single iteration instead: it still fails for that one name, but the
|
||||
# loop continues and we can tell allowed / blanked / unblankable apart afterwards.
|
||||
# This is not a skip-list of known-bad names (that would miss the next one nobody
|
||||
# thought of) — every name in _cb633_blank is still attempted, unconditionally.
|
||||
# Every name reaching this loop already passed the identical identifier check in the
|
||||
# enumeration loop above — but that guard is 20 lines away in a different loop, and
|
||||
# this line is about to splice the name into a string handed to `eval`. Before this
|
||||
# fix the name only ever reached `export` quoted ("$n="), which is inert on a
|
||||
# non-identifier string either way; `eval` makes THIS line the only thing standing
|
||||
# between such a string and code execution in the member's pane, so it re-asserts the
|
||||
# same check on its own rather than trusting a guard it does not own. Under normal
|
||||
# operation this can never fire (the enumeration guard already filtered everything
|
||||
# reaching _cb633_blank), so a name caught here is counted as unblankable rather than
|
||||
# silently dropped — it is real evidence that the upstream guard was bypassed.
|
||||
#
|
||||
# fleetd #400: the attempt's own exit status is NOT proof of its effect. zsh coerces
|
||||
# a bare `NAME=` assignment on an integer special parameter (SECONDS, RANDOM, SHLVL,
|
||||
# HISTSIZE, COLUMNS, LINES, USERNAME on this host) to a number instead of failing —
|
||||
# `eval` returns 0, the value is untouched, and the old exit-status check reported it
|
||||
# as blanked when it was not. Classify on the observed effect instead: attempt the
|
||||
# export, then read the name's value back with the `(P)` indirection flag and decide
|
||||
# from whether it is now empty. One check then covers all three shapes a name can
|
||||
# take here — a genuine blank, a fatal read-only error `eval` merely contained, and
|
||||
# this silent no-op — without the exit status entering the decision at all.
|
||||
typeset -a _cb633_ok _cb633_unblankable
|
||||
_cb633_ok=()
|
||||
_cb633_unblankable=()
|
||||
for _cb633_n in "${_cb633_blank[@]}"; do
|
||||
if [[ ! "$_cb633_n" =~ ^[A-Za-z_][A-Za-z0-9_]*$ ]]; then
|
||||
_cb633_unblankable+=("$_cb633_n")
|
||||
continue
|
||||
fi
|
||||
eval "export ${_cb633_n}=" 2>/dev/null
|
||||
if [[ -z "${(P)_cb633_n}" ]]; then
|
||||
_cb633_ok+=("$_cb633_n")
|
||||
else
|
||||
_cb633_unblankable+=("$_cb633_n")
|
||||
fi
|
||||
done
|
||||
|
||||
integer _cb633_kept=$(( _cb633_total - ${#_cb633_blank} ))
|
||||
{
|
||||
print -r -- "allowed $_cb633_kept of $_cb633_total"
|
||||
for _cb633_n in "${_cb633_blank[@]}"; do print -r -- "$_cb633_n"; done
|
||||
print -r -- "allowed $_cb633_kept of $_cb633_total failed ${#_cb633_unblankable}"
|
||||
for _cb633_n in "${_cb633_ok[@]}"; do print -r -- "$_cb633_n"; done
|
||||
for _cb633_n in "${_cb633_unblankable[@]}"; do print -r -- "!$_cb633_n"; done
|
||||
} > "$ZDOTDIR/%s" 2>/dev/null
|
||||
|
||||
unset _cb633_allowed _cb633_names _cb633_blank _cb633_n _cb633_total _cb633_kept
|
||||
unset _cb633_allowed _cb633_names _cb633_blank _cb633_ok _cb633_unblankable _cb633_n _cb633_total _cb633_kept
|
||||
""".formatted(names, MemberEnvAllowList.zshCasePattern(), REPORT_FILE);
|
||||
}
|
||||
|
||||
@@ -356,6 +434,11 @@ public final class EnvAllowListScrub {
|
||||
* Read and parse {@link #REPORT_FILE} out of a generated ZDOTDIR directory. Returns {@code null}
|
||||
* when absent or unreadable (the pane may have been torn down before its login shell ever got to
|
||||
* the scrub) — callers treat that as "no measurement available", never as success.
|
||||
*
|
||||
* <p>First line is {@code "allowed <N> of <M> failed <F>"} (fleetd #394 added the trailing
|
||||
* {@code failed <F>} — a count of names the scrub attempted to blank but could not, e.g. a zsh
|
||||
* read-only/special parameter). Every following non-blank line is a name: a bare name was
|
||||
* blanked, a {@code !}-prefixed name was attempted and failed. Values never appear on either.
|
||||
*/
|
||||
static ScrubReport readReport(Path zdotdir) {
|
||||
Path report = zdotdir.resolve(REPORT_FILE);
|
||||
@@ -368,17 +451,24 @@ public final class EnvAllowListScrub {
|
||||
return null;
|
||||
}
|
||||
String[] parts = lines.getFirst().substring("allowed ".length()).trim().split("\\s+");
|
||||
if (parts.length != 3 || !"of".equals(parts[1])) {
|
||||
if (parts.length != 5 || !"of".equals(parts[1]) || !"failed".equals(parts[3])) {
|
||||
return null;
|
||||
}
|
||||
List<String> blanked = new ArrayList<>();
|
||||
List<String> unblankable = new ArrayList<>();
|
||||
for (int i = 1; i < lines.size(); i++) {
|
||||
if (!lines.get(i).isBlank()) {
|
||||
blanked.add(lines.get(i));
|
||||
String line = lines.get(i);
|
||||
if (line.isBlank()) {
|
||||
continue;
|
||||
}
|
||||
if (line.startsWith("!")) {
|
||||
unblankable.add(line.substring(1));
|
||||
} else {
|
||||
blanked.add(line);
|
||||
}
|
||||
}
|
||||
return new ScrubReport(Integer.parseInt(parts[0]), Integer.parseInt(parts[2]),
|
||||
List.copyOf(blanked));
|
||||
Integer.parseInt(parts[4]), List.copyOf(blanked), List.copyOf(unblankable));
|
||||
} catch (IOException | NumberFormatException e) {
|
||||
return null;
|
||||
}
|
||||
|
||||
@@ -1643,11 +1643,19 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
log.warn("memberCredentials allow-list: pane {} left no scrub report in {} — the "
|
||||
+ "environment scrub cannot be confirmed to have run. Either the pane ended "
|
||||
+ "before its shell finished starting, or its shell never read our generated "
|
||||
+ "startup files, in which case that member saw the full host environment.",
|
||||
+ "startup files. Either way, we cannot tell from here whether the scrub ran, "
|
||||
+ "so we do not know what that member's environment contained.",
|
||||
paneId, dir);
|
||||
} else {
|
||||
log.info("memberCredentials allow-list: pane {} allowed {} of {} environment variables",
|
||||
paneId, report.allowed(), report.total());
|
||||
if (report.failed() > 0) {
|
||||
log.warn("memberCredentials allow-list: pane {} could not blank {} environment "
|
||||
+ "variable(s) — {} (likely a zsh read-only/special parameter) — those "
|
||||
+ "names were left in the member's environment. Confirm none of them is a "
|
||||
+ "credential.",
|
||||
paneId, report.failed(), report.unblankable());
|
||||
}
|
||||
List<String> shaped = report.blanked().stream()
|
||||
.filter(name -> CREDENTIAL_SHAPED_NAME.matcher(name).matches())
|
||||
.toList();
|
||||
|
||||
@@ -1317,6 +1317,27 @@ public final class MessageService {
|
||||
return new TaskView(ticket, Phase.FAILED, null, null, detail, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* Test seam only — carries no production behaviour, and nothing in this class calls it;
|
||||
* {@link #pruneTerminalTickets} still reads {@link Task#completedNanos} directly.
|
||||
*
|
||||
* <p>Reports whether {@code ticket}'s completion hook (the {@code whenComplete} registered in
|
||||
* {@link Task}'s constructor) has actually run yet. Exists because {@link #poll} can report
|
||||
* {@link Phase#DONE} for a ticket before that hook fires: {@code CompletableFuture.complete()}
|
||||
* publishes its result and only afterwards runs dependent actions such as {@code whenComplete}
|
||||
* (fleetd #399), so a caller that observes the future done via {@link #poll} is not thereby
|
||||
* guaranteed to also observe {@link Task#completedNanos} stamped. A test that must order both
|
||||
* events — e.g. before advancing an injected clock past the TTL, to avoid stamping the
|
||||
* *advanced* time and masking a real eviction bug — waits on this instead of on
|
||||
* {@link Phase#DONE}.
|
||||
*
|
||||
* @return {@code false} for an unknown ticket or one whose completion hook has not run yet
|
||||
*/
|
||||
boolean isCompletionStampedForTest(String ticket) {
|
||||
Task task = tasks.get(ticket);
|
||||
return task != null && task.completedNanos != null;
|
||||
}
|
||||
|
||||
/** Best-effort live worker status for a pending poll; never throws (a lookup error is just noise). */
|
||||
private String liveStatus(String target) {
|
||||
try {
|
||||
|
||||
@@ -230,6 +230,26 @@ public final class ReplyPushLoop {
|
||||
.collect(Collectors.toUnmodifiableSet());
|
||||
}
|
||||
|
||||
/**
|
||||
* Test seam only (fleetd #418): carries no production behaviour, and nothing in this class
|
||||
* calls it. Exposes {@link #pendingQuestionTurnIdsFor} — the exact state {@link #decide} reads
|
||||
* to decide whether a question keeps a lead's schedule alive.
|
||||
*
|
||||
* <p>{@code MessageService.ask()} does three things in order before a question is fully open to
|
||||
* this loop: it flips the ticket's {@code poll()} phase to {@code Phase.ASKING}, then resolves
|
||||
* the reverse-rendezvous waiter, then calls {@link #onQuestionOpened}, which is what actually
|
||||
* populates {@link #pendingQuestions}. A test that barriers on {@code Phase.ASKING} observes only
|
||||
* the first of those three steps — under load the asker thread can be descheduled between steps
|
||||
* one and three, so the barrier releases before this method's underlying map is populated, and
|
||||
* {@link #decide} correctly reports nothing pending yet. A test that must order itself after the
|
||||
* state {@link #decide} actually reads waits on this instead of on the phase.
|
||||
*
|
||||
* @return an unmodifiable snapshot; empty for a lead with no open questions
|
||||
*/
|
||||
Set<String> pendingQuestionTurnIdsForTest(String lead) {
|
||||
return pendingQuestionTurnIdsFor(lead);
|
||||
}
|
||||
|
||||
private List<PendingIncident> pendingIncidentsFor(String lead) {
|
||||
return pendingIncidents.values().stream().filter(i -> lead.equals(i.key().lead())).toList();
|
||||
}
|
||||
|
||||
@@ -0,0 +1,150 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.Logger;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
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.List;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #395: {@code exhaustedPattern} (see {@link FleetConfig.Profile#exhaustedPattern}) is
|
||||
* deliberately opt-in — an unset one leaves usage-limit detection silently OFF for that profile,
|
||||
* and nothing quarantines its credential. {@link Fleetd#reportExhaustedPatternGap} must say so at
|
||||
* startup, naming every unarmed profile, and must never fire when every profile is armed. Mirrors
|
||||
* {@link MemberCredentialsGapReportTest}'s pattern, capturing the real log via a
|
||||
* {@link ListAppender}.
|
||||
*
|
||||
* <p>The 8-profile shape in {@link #theLiveEightProfileShapeWarnsExactlyTheSixUnarmedProfiles} is
|
||||
* the live {@code fleetd.yaml} shape measured 2026-09-10 (fleetd #395's own ticket): 6 of 8
|
||||
* profiles unarmed, 2 of those 6 ({@code opus}, {@code sonnet}) running on the operator's Claude
|
||||
* subscription. {@code fleetd.yaml} itself is gitignored and unavailable to this test, so the
|
||||
* shape is reproduced as a throwaway config in a {@code @TempDir} rather than read off disk.
|
||||
*/
|
||||
class ExhaustedPatternGapReportTest {
|
||||
|
||||
private static FleetConfig load(Path dir, String yaml) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, yaml);
|
||||
return FleetConfig.load(f);
|
||||
}
|
||||
|
||||
private static ListAppender<ILoggingEvent> attach() {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
return appender;
|
||||
}
|
||||
|
||||
private static void detach(ListAppender<ILoggingEvent> appender) {
|
||||
((Logger) LoggerFactory.getLogger(Fleetd.class)).detachAppender(appender);
|
||||
}
|
||||
|
||||
/** Every profile name mentioned by a WARN-level log line, across every WARN this call produced. */
|
||||
private static List<String> warnMessages(ListAppender<ILoggingEvent> appender) {
|
||||
return appender.list.stream()
|
||||
.filter(e -> e.getLevel() == Level.WARN)
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.toList();
|
||||
}
|
||||
|
||||
@Test
|
||||
void theLiveEightProfileShapeWarnsExactlyTheSixUnarmedProfiles(@TempDir Path dir) throws Exception {
|
||||
// Reproduces the live shape measured 2026-09-10: 8 profiles, 2 armed (sol, terra), 6
|
||||
// unarmed (local, local-direct, gx, opus, sonnet, xf) — 2 of the unarmed 6 (opus, sonnet)
|
||||
// are subscription: true.
|
||||
FleetConfig cfg = load(dir, """
|
||||
profiles:
|
||||
local:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
local-direct:
|
||||
baseUrl: http://gx01.gw:8000
|
||||
gx:
|
||||
kind: opencode
|
||||
baseUrl: https://llm.ltms.dev/v1
|
||||
opus:
|
||||
subscription: true
|
||||
model: claude-opus-5
|
||||
sonnet:
|
||||
subscription: true
|
||||
model: claude-sonnet-5
|
||||
sol:
|
||||
baseUrl: https://llm.ltms.dev/v1
|
||||
exhaustedPattern: "The usage limit has been reached"
|
||||
terra:
|
||||
baseUrl: https://llm.ltms.dev/v1
|
||||
exhaustedPattern: "The usage limit has been reached"
|
||||
xf:
|
||||
baseUrl: https://llm.ltms.dev/v1
|
||||
""");
|
||||
|
||||
ListAppender<ILoggingEvent> appender = attach();
|
||||
try {
|
||||
Fleetd.reportExhaustedPatternGap(cfg);
|
||||
} finally {
|
||||
detach(appender);
|
||||
}
|
||||
|
||||
List<String> warns = warnMessages(appender);
|
||||
assertFalse(warns.isEmpty(), "6 of 8 profiles are unarmed — at least one WARN must fire");
|
||||
|
||||
// Exactly one WARN aggregates every unarmed profile, naming all 6 and none of the 2 armed.
|
||||
String aggregate = warns.stream()
|
||||
.filter(m -> m.contains("no exhaustedPattern configured"))
|
||||
.findFirst()
|
||||
.orElseThrow(() -> new AssertionError("expected an aggregate unarmed-profiles WARN: " + warns));
|
||||
for (String unarmed : List.of("local", "local-direct", "gx", "opus", "sonnet", "xf")) {
|
||||
assertTrue(aggregate.contains(unarmed), "aggregate WARN must name '" + unarmed + "': " + aggregate);
|
||||
}
|
||||
for (String armed : List.of("sol", "terra")) {
|
||||
assertFalse(aggregate.contains(armed), "aggregate WARN must NOT name armed profile '" + armed + "': " + aggregate);
|
||||
}
|
||||
|
||||
// A second, louder WARN calls out the subscription profiles specifically.
|
||||
String subscriptionWarn = warns.stream()
|
||||
.filter(m -> m.contains("metered Claude plan"))
|
||||
.findFirst()
|
||||
.orElseThrow(() -> new AssertionError("expected a subscription-specific WARN: " + warns));
|
||||
assertTrue(subscriptionWarn.contains("opus"), subscriptionWarn);
|
||||
assertTrue(subscriptionWarn.contains("sonnet"), subscriptionWarn);
|
||||
assertFalse(subscriptionWarn.contains("local-direct"),
|
||||
"the subscription WARN must not name a non-subscription profile: " + subscriptionWarn);
|
||||
}
|
||||
|
||||
@Test
|
||||
void everyProfileArmedProducesNoWarningAtAll(@TempDir Path dir) throws Exception {
|
||||
FleetConfig cfg = load(dir, """
|
||||
profiles:
|
||||
sol:
|
||||
baseUrl: https://llm.ltms.dev/v1
|
||||
exhaustedPattern: "The usage limit has been reached"
|
||||
terra:
|
||||
baseUrl: https://llm.ltms.dev/v1
|
||||
exhaustedPattern: "The usage limit has been reached"
|
||||
opus:
|
||||
subscription: true
|
||||
model: claude-opus-5
|
||||
exhaustedPattern: "5-hour limit reached"
|
||||
""");
|
||||
|
||||
ListAppender<ILoggingEvent> appender = attach();
|
||||
try {
|
||||
Fleetd.reportExhaustedPatternGap(cfg);
|
||||
} finally {
|
||||
detach(appender);
|
||||
}
|
||||
|
||||
assertTrue(warnMessages(appender).isEmpty(),
|
||||
"every profile is armed — a checker that warns anyway always fires: " + warnMessages(appender));
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,94 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.mcp.FleetMcp;
|
||||
import dev.ltms.fleet.placement.BackendQuarantine;
|
||||
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 java.util.Map;
|
||||
import java.util.regex.Pattern;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #404: {@code exhaustionDetectionArmed} must describe the startup pattern map, not the
|
||||
* reloaded config snapshot.
|
||||
*
|
||||
* <p>This test needs a reload. At startup the two snapshots agree, so a test of only a newly
|
||||
* started daemon would not detect a live {@code config.get()} lookup in the report field.
|
||||
*/
|
||||
class FleetdExhaustionDetectionArmedWiringTest {
|
||||
|
||||
private static final String NO_PATTERN = """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
herdrSocket: ~/.config/herdr/herdr.sock
|
||||
profiles:
|
||||
terra:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
model: terra
|
||||
guard:
|
||||
offSubscriptionHosts:
|
||||
- gx00.gw
|
||||
""";
|
||||
|
||||
private static final String WITH_PATTERN = """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
herdrSocket: ~/.config/herdr/herdr.sock
|
||||
profiles:
|
||||
terra:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
model: terra
|
||||
exhaustedPattern: "usage limit"
|
||||
guard:
|
||||
offSubscriptionHosts:
|
||||
- gx00.gw
|
||||
""";
|
||||
|
||||
@Test
|
||||
@DisplayName("reloading an exhaustedPattern does not arm the startup detection source")
|
||||
void reloadedPatternDoesNotChangeTheArmedFieldUntilRestart(@TempDir Path dir) throws Exception {
|
||||
Path file = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(file, NO_PATTERN);
|
||||
ConfigRef config = new ConfigRef(file, FleetConfig.load(file));
|
||||
|
||||
Files.writeString(file, WITH_PATTERN);
|
||||
assertTrue(config.reload().applied());
|
||||
assertTrue(config.get().profiles().get("terra").hasExhaustedPattern());
|
||||
|
||||
FleetMcp.QuarantineSource source = Fleetd.quarantineSource(config, BackendQuarantine.none(),
|
||||
Map.of());
|
||||
assertFalse(source.exhaustedPatternArmed().apply("terra"),
|
||||
"exhaustionDetectionArmed must use the startup pattern map, not config.get()");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a profile in the startup pattern map is reported as armed")
|
||||
void aProfileInTheStartupMapIsArmed(@TempDir Path dir) throws Exception {
|
||||
// fleetd #404, second direction. The test above only ever passes an EMPTY startup map, so
|
||||
// it cannot tell a correct lookup from one that is permanently off. Measured: replacing the
|
||||
// armed lambda with `profile -> false` left the whole suite green at 1475 tests. That
|
||||
// mutation would make #395's visibility feature dead — an operator fixing a detection gap
|
||||
// would be told the gap is still open after fixing it, forever. Both directions are needed:
|
||||
// this test is the only thing that fails when the field stops reporting armed at all.
|
||||
Path file = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(file, WITH_PATTERN);
|
||||
ConfigRef config = new ConfigRef(file, FleetConfig.load(file));
|
||||
|
||||
FleetMcp.QuarantineSource source = Fleetd.quarantineSource(config, BackendQuarantine.none(),
|
||||
Map.of("terra", Pattern.compile("usage limit")));
|
||||
|
||||
assertTrue(source.exhaustedPatternArmed().apply("terra"),
|
||||
"a profile whose pattern was compiled at startup must report armed");
|
||||
assertFalse(source.exhaustedPatternArmed().apply("sonnet"),
|
||||
"a profile absent from the startup map must not report armed");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,132 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
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.assertThrows;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* Proves the one thing no test proved before this ticket's follow-up: that {@code Fleetd.main}
|
||||
* ITSELF — not a copy of its logic, not the validator called directly — still refuses to start on
|
||||
* a bad config. Mutation testing found that deleting {@code cfg.validateAll();} (née six
|
||||
* individual {@code cfg.validateXxx();} calls) from {@code Fleetd.main} left the full 1478-test
|
||||
* suite green; every existing test called a validator directly and none exercised {@code
|
||||
* Fleetd.main} as the caller. See {@code FleetConfigValidateAllTest} for why the fix collapses
|
||||
* those six calls into one reflective {@link FleetConfig#validateAll()}, and {@code
|
||||
* ConfigRefTest} for the equivalent proof on the {@link
|
||||
* dev.ltms.fleet.config.ConfigRef#reload()} path.
|
||||
*
|
||||
* <p>This is deliberately the real {@code static void main(String[] args)} — package-private, so
|
||||
* only a test in this package can call it, which is exactly what makes this proof strong: it is
|
||||
* not a helper extracted for testability, it is the literal method {@code java -jar fleetd.jar}
|
||||
* invokes. Every fixture below is otherwise-valid and fails exactly one validator, and — because
|
||||
* {@link FleetConfig#validateAll()} runs immediately after {@code SubscriptionGuard.
|
||||
* assertPrimaryClean}, before {@code main} opens the herdr socket, binds Javalin, or touches
|
||||
* anything else with a real side effect — calling {@code Fleetd.main} with one of these configs is
|
||||
* safe: it is guaranteed to throw before reaching any of that, precisely because the config is
|
||||
* deliberately invalid.
|
||||
*/
|
||||
class FleetdStartupValidationTest {
|
||||
|
||||
private static void assertMainRefuses(Path dir, String fileName, String yaml,
|
||||
String mustContain) throws Exception {
|
||||
Path f = dir.resolve(fileName);
|
||||
Files.writeString(f, yaml);
|
||||
IllegalStateException e = assertThrows(IllegalStateException.class,
|
||||
() -> Fleetd.main(new String[]{f.toString()}),
|
||||
fileName + ": Fleetd.main must refuse this config before doing anything else");
|
||||
assertTrue(e.getMessage().contains(mustContain),
|
||||
fileName + ": expected message to contain \"" + mustContain + "\" but was: "
|
||||
+ e.getMessage());
|
||||
}
|
||||
|
||||
@Test
|
||||
void mainRefusesANonLoopbackBindWithoutTokenMode(@TempDir Path dir) throws Exception {
|
||||
assertMainRefuses(dir, "auth-exposure.yaml", """
|
||||
bind:
|
||||
host: 0.0.0.0
|
||||
port: 8765
|
||||
""", "auth.mode: token");
|
||||
}
|
||||
|
||||
@Test
|
||||
void mainRefusesALeadTabPrefixCollision(@TempDir Path dir) throws Exception {
|
||||
assertMainRefuses(dir, "lead-tab-prefixes.yaml", """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
fleet:
|
||||
tabLabel: "lead: {role} {profile}"
|
||||
leaders:
|
||||
opus:
|
||||
tab: "lead: opus"
|
||||
""", "fleet.tabLabel");
|
||||
}
|
||||
|
||||
@Test
|
||||
void mainRefusesASubscriptionProfileThatReseatsAnthropicBaseUrl(@TempDir Path dir)
|
||||
throws Exception {
|
||||
assertMainRefuses(dir, "subscription-profiles.yaml", """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
profiles:
|
||||
sonnet:
|
||||
subscription: true
|
||||
argv: ["ccs", "sonnet"]
|
||||
env:
|
||||
ANTHROPIC_BASE_URL: http://anything-not-on-the-allowlist
|
||||
""", "ANTHROPIC_BASE_URL");
|
||||
}
|
||||
|
||||
@Test
|
||||
void mainRefusesAnUnknownCharterKey(@TempDir Path dir) throws Exception {
|
||||
assertMainRefuses(dir, "charters.yaml", """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
fleet:
|
||||
charters:
|
||||
architetc: text
|
||||
""", "architetc");
|
||||
}
|
||||
|
||||
@Test
|
||||
void mainRefusesAnArchitectSlotNamingAnUnconfiguredProfile(@TempDir Path dir) throws Exception {
|
||||
assertMainRefuses(dir, "members.yaml", """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
profiles:
|
||||
gx10:
|
||||
baseUrl: http://gx10.gw:8000
|
||||
fleet:
|
||||
architects:
|
||||
lead-designer:
|
||||
profile: sonnet
|
||||
""", "lead-designer");
|
||||
}
|
||||
|
||||
@Test
|
||||
void mainRefusesAProfileNamingAModelOutsideTheAllowList(@TempDir Path dir) throws Exception {
|
||||
assertMainRefuses(dir, "models.yaml", """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
profiles:
|
||||
sonnet:
|
||||
baseUrl: http://gx10.gw:8000
|
||||
model: claude-sonnet-5
|
||||
rogue:
|
||||
baseUrl: http://gx11.gw:8000
|
||||
model: claude-opus-9000
|
||||
models:
|
||||
allow:
|
||||
- model: claude-sonnet-5
|
||||
""", "rogue");
|
||||
}
|
||||
}
|
||||
@@ -109,6 +109,8 @@ class ConfigRefTopLevelReportingCoverageTest {
|
||||
v.put("memberLoginShell", null);
|
||||
v.put("memberSkills", "/skills/a");
|
||||
v.put("idleSleepGuard", new FleetConfig.IdleSleepGuard(true));
|
||||
v.put("models", new FleetConfig.Models(
|
||||
List.of(new FleetConfig.Models.ModelEntry("model-a"))));
|
||||
assertNamesMatchComponents(v);
|
||||
return v;
|
||||
}
|
||||
@@ -151,6 +153,8 @@ class ConfigRefTopLevelReportingCoverageTest {
|
||||
v.put("memberLoginShell", null);
|
||||
v.put("memberSkills", "/skills/b");
|
||||
v.put("idleSleepGuard", new FleetConfig.IdleSleepGuard(false));
|
||||
v.put("models", new FleetConfig.Models(
|
||||
List.of(new FleetConfig.Models.ModelEntry("model-b"))));
|
||||
assertNamesMatchComponents(v);
|
||||
return v;
|
||||
}
|
||||
|
||||
@@ -2683,4 +2683,222 @@ class FleetConfigTest {
|
||||
assertTrue(cfg.idleSleepGuard().isEnabled(),
|
||||
"unlike ConfigReload/Health, this block defaults to ON even when present but empty");
|
||||
}
|
||||
|
||||
// --- models: central allow-list of usable models -----------------------------------------
|
||||
|
||||
/**
|
||||
* An absent {@code models:} block is today's behaviour exactly: no profile's {@code model:} is
|
||||
* checked against anything, whatever it says. This is the "existing config keeps working"
|
||||
* invariant — an operator on a gitignored {@code fleetd.yaml} that predates this feature must
|
||||
* not be broken by upgrading the daemon.
|
||||
*/
|
||||
@Test
|
||||
void absentModelsBlockValidatesNothing(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
profiles:
|
||||
gx10:
|
||||
baseUrl: http://gx10.gw:8000
|
||||
model: totally-unheard-of-model-xyz
|
||||
""");
|
||||
|
||||
FleetConfig cfg = FleetConfig.load(f);
|
||||
assertDoesNotThrow(cfg::validateModels);
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code models:} present but with an empty (or absent) {@code allow:} must behave exactly like
|
||||
* the block being absent — an operator adding the block for the first time with nothing in it
|
||||
* yet must not be surprised by every profile suddenly refusing to start.
|
||||
*/
|
||||
@Test
|
||||
void modelsBlockPresentButEmptyValidatesNothing(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
profiles:
|
||||
gx10:
|
||||
baseUrl: http://gx10.gw:8000
|
||||
model: whatever-the-operator-typed
|
||||
models:
|
||||
allow: []
|
||||
""");
|
||||
|
||||
FleetConfig cfg = FleetConfig.load(f);
|
||||
assertDoesNotThrow(cfg::validateModels);
|
||||
}
|
||||
|
||||
/**
|
||||
* The core of the ticket: a profile naming a model outside the configured allow-list fails
|
||||
* config load, naming both the model and the profile that wanted it.
|
||||
*/
|
||||
@Test
|
||||
void aProfileNamingAModelOutsideTheAllowListRefusesToStart(@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
|
||||
rogue:
|
||||
baseUrl: http://gx11.gw:8000
|
||||
model: claude-opus-9000
|
||||
models:
|
||||
allow:
|
||||
- model: claude-sonnet-5
|
||||
- model: claude-haiku-5
|
||||
""");
|
||||
FleetConfig cfg = FleetConfig.load(f);
|
||||
|
||||
IllegalStateException e = assertThrows(IllegalStateException.class, cfg::validateModels);
|
||||
assertTrue(e.getMessage().contains("rogue"), "the refusal names the offending profile");
|
||||
assertTrue(e.getMessage().contains("claude-opus-9000"), "the refusal names the offending model");
|
||||
assertFalse(e.getMessage().contains("'sonnet'"),
|
||||
"the profile whose model IS allowed must not be reported");
|
||||
}
|
||||
|
||||
/** A profile whose {@code model:} is on the allow-list loads fine. */
|
||||
@Test
|
||||
void aProfileNamingAModelOnTheAllowListLoadsFine(@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);
|
||||
}
|
||||
|
||||
/**
|
||||
* A profile that never sets {@code model:} (a {@code subscription: true} profile relying on the
|
||||
* account's own default is the live example) must not be refused just because an allow-list is
|
||||
* active — there is nothing to check it against.
|
||||
*/
|
||||
@Test
|
||||
void aProfileWithNoModelPassesEvenWithAnActiveAllowList(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
profiles:
|
||||
opus:
|
||||
subscription: true
|
||||
models:
|
||||
allow:
|
||||
- model: claude-sonnet-5
|
||||
""");
|
||||
|
||||
FleetConfig cfg = FleetConfig.load(f);
|
||||
assertDoesNotThrow(cfg::validateModels);
|
||||
}
|
||||
|
||||
/**
|
||||
* Reproduces the live config shape this ticket measured: a mix of {@code amazon-bedrock},
|
||||
* {@code opencode} and {@code openai}-backed profiles, some naming a bare Claude id and some a
|
||||
* provider-prefixed opencode id, in ONE allow-list. Both forms are just opaque strings compared
|
||||
* for exact equality — proves the shape decision holds against the shape actually seen live,
|
||||
* not just against a synthetic single-provider example.
|
||||
*/
|
||||
@Test
|
||||
void aBareClaudeIdAndAnOpencodeProviderPrefixedIdBothFitOneAllowList(@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
|
||||
terra:
|
||||
kind: opencode
|
||||
model: openai/gpt-5.6-terra
|
||||
nova:
|
||||
kind: opencode
|
||||
model: amazon-bedrock/amazon.nova-pro-v1:0
|
||||
models:
|
||||
allow:
|
||||
- model: claude-sonnet-5
|
||||
- model: openai/gpt-5.6-terra
|
||||
- model: amazon-bedrock/amazon.nova-pro-v1:0
|
||||
""");
|
||||
|
||||
FleetConfig cfg = FleetConfig.load(f);
|
||||
assertDoesNotThrow(cfg::validateModels);
|
||||
}
|
||||
|
||||
/** The same live shape, but one opencode profile's model is missing from the allow-list. */
|
||||
@Test
|
||||
void anUnlistedOpencodeProviderPrefixedModelRefusesToStart(@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
|
||||
terra:
|
||||
kind: opencode
|
||||
model: openai/gpt-5.6-terra-withdrawn
|
||||
models:
|
||||
allow:
|
||||
- model: claude-sonnet-5
|
||||
- model: openai/gpt-5.6-terra
|
||||
""");
|
||||
FleetConfig cfg = FleetConfig.load(f);
|
||||
|
||||
IllegalStateException e = assertThrows(IllegalStateException.class, cfg::validateModels);
|
||||
assertTrue(e.getMessage().contains("terra"), "the refusal names the offending profile");
|
||||
assertTrue(e.getMessage().contains("openai/gpt-5.6-terra-withdrawn"),
|
||||
"the refusal names the offending model, with its provider prefix intact");
|
||||
}
|
||||
|
||||
/**
|
||||
* Editing a {@code profiles:} entry alone must never be able to widen what is permitted — only
|
||||
* editing {@code models.allow:} itself can. This is the invariant the ticket states explicitly;
|
||||
* this test pins it by giving a profile a plausible-looking model that was never added to the
|
||||
* allow-list and confirming it is still refused.
|
||||
*/
|
||||
@Test
|
||||
void addingAProfileCannotWidenTheAllowListByItself(@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
|
||||
brand-new:
|
||||
baseUrl: http://gx12.gw:8000
|
||||
model: claude-sonnet-6-preview
|
||||
models:
|
||||
allow:
|
||||
- model: claude-sonnet-5
|
||||
""");
|
||||
FleetConfig cfg = FleetConfig.load(f);
|
||||
|
||||
IllegalStateException e = assertThrows(IllegalStateException.class, cfg::validateModels);
|
||||
assertTrue(e.getMessage().contains("brand-new"));
|
||||
assertTrue(e.getMessage().contains("claude-sonnet-6-preview"));
|
||||
}
|
||||
|
||||
/** {@code models} is a brand-new top-level key and must be recognized, not WARN-ed as unknown. */
|
||||
@Test
|
||||
void modelsIsAKnownTopLevelKey() {
|
||||
assertTrue(FleetConfig.KNOWN_TOP_LEVEL_KEYS.contains("models"));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,358 @@
|
||||
package dev.ltms.fleet.config;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.lang.reflect.Method;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.Set;
|
||||
import java.util.TreeSet;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertThrows;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* The gap this class exists to close: mutation testing on the fleetd ticket "central allow-list
|
||||
* of usable models" found that although {@link FleetConfig#validateModels()}'s own logic was well
|
||||
* pinned, nothing proved either real caller ({@code Fleetd.main} and {@link ConfigRef#reload()})
|
||||
* still invoked it — deleting the call site left the full suite green (1478/0/0/0). A follow-up
|
||||
* measurement (same technique — remove one call site, run the suite, not read the code) found the
|
||||
* SAME gap for all five of {@link FleetConfig}'s other validators at startup, and for four of the
|
||||
* six inside {@link ConfigRef#reload()}. This is a class of gap, not one line's mistake: every one
|
||||
* of those thirteen tests called the validator itself directly, never the real caller that was
|
||||
* supposed to.
|
||||
*
|
||||
* <p>The fix replaces the six individual {@code cfg.validateXxx()} calls at each of the two real
|
||||
* call sites with one {@link FleetConfig#validateAll()}, which reaches every validator by
|
||||
* reflection rather than by a hand-maintained list of names. A hand-maintained list of six names
|
||||
* would have exactly the defect it replaces: the seventh validator someone adds next month has no
|
||||
* reason to be added to it, and nothing would say so. This class proves TWO separate claims, and
|
||||
* keeps them separate on purpose:
|
||||
*
|
||||
* <ol>
|
||||
* <li>{@link #theSweepMechanismIsGenericNotHardcodedToFleetConfigsSixNames()} and its neighbours
|
||||
* prove the reflective sweep itself ({@link FleetConfig#invokeAllValidators}) is a general
|
||||
* mechanism — it runs whatever public, no-arg, void {@code validateXxx()} methods a class
|
||||
* happens to declare today, including a class with more of them than {@link FleetConfig}
|
||||
* has right now. This is the proof that a future, real seventh validator on {@link
|
||||
* FleetConfig} would be swept automatically, without needing to add a real (unwanted)
|
||||
* seventh validator just to exercise the claim.</li>
|
||||
* <li>{@link #validateAllReachesEveryOneOfTodaysSixValidators()} proves {@link
|
||||
* FleetConfig#validateAll()} itself is wired to that same generic mechanism and genuinely
|
||||
* reaches each of today's six real validators — reusing the exact minimal failing
|
||||
* configurations {@code FleetConfigTest} already established for each one directly, so a
|
||||
* single call to {@code validateAll()} is shown to reproduce every one of those six
|
||||
* failures.</li>
|
||||
* </ol>
|
||||
*
|
||||
* <p>Together with the direct-{@code Fleetd.main}-invocation tests in {@code
|
||||
* FleetdStartupValidationTest} (which prove the real startup call site still calls {@code
|
||||
* validateAll()}) and the {@code ConfigRefTest} reload tests (which prove the same for {@link
|
||||
* ConfigRef#reload()}), removing {@code cfg.validateAll();} from either real call site now fails
|
||||
* a test in this module.
|
||||
*
|
||||
* <p><b>What is NOT pinned, measured rather than assumed.</b> Reverting {@link
|
||||
* FleetConfig#validateAll()} to a hardcoded list of today's six method calls leaves the whole
|
||||
* suite green (measured at review: 1491 tests, 0 failures). Nothing ties {@code validateAll()} to
|
||||
* the generic sweep — claim 1 proves {@link FleetConfig#invokeAllValidators} is generic, and claim
|
||||
* 2 proves {@code validateAll()} reaches today's six, and a hardcoded list satisfies both. So the
|
||||
* reflective sweep is a convenience, not the guarantee. The guarantee is {@link
|
||||
* #fleetConfigDeclaresExactlyTheseSixValidatorsToday()}: it fails the moment a seventh validator
|
||||
* is declared, which forces whoever adds it to look at this file.
|
||||
*/
|
||||
class FleetConfigValidateAllTest {
|
||||
|
||||
// ── Claim 1: the reflective sweep is a general mechanism, not six names in disguise ──────────
|
||||
|
||||
/**
|
||||
* A throwaway fixture class, unrelated to {@link FleetConfig} in every way except shape: three
|
||||
* public, no-arg, void methods named {@code validateXxx}. Proves the sweep works on ANY class
|
||||
* with this shape, not on something special-cased to {@link FleetConfig}.
|
||||
*/
|
||||
static class ThreeValidators {
|
||||
final List<String> ran = new ArrayList<>();
|
||||
|
||||
public void validateAlpha() {
|
||||
ran.add("validateAlpha");
|
||||
}
|
||||
|
||||
public void validateBeta() {
|
||||
ran.add("validateBeta");
|
||||
}
|
||||
|
||||
public void validateGamma() {
|
||||
ran.add("validateGamma");
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void theSweepMechanismIsGenericNotHardcodedToFleetConfigsSixNames() {
|
||||
ThreeValidators target = new ThreeValidators();
|
||||
FleetConfig.invokeAllValidators(target);
|
||||
assertEquals(List.of("validateAlpha", "validateBeta", "validateGamma"), target.ran,
|
||||
"every validateXxx() method on this unrelated class must run, in alphabetical "
|
||||
+ "order — the sweep reads the class's own shape, not a name FleetConfig "
|
||||
+ "happens to know about");
|
||||
}
|
||||
|
||||
/**
|
||||
* The core of the "self-maintaining" requirement: the exact same class shape as {@link
|
||||
* ThreeValidators}, plus one more method — standing in for "a developer adds a validator next
|
||||
* month". Nothing about the sweep changes to pick it up; the new method is invoked purely
|
||||
* because it exists and matches the shape. This is what makes adding a seventh real validator
|
||||
* to {@link FleetConfig} safe without touching {@link FleetConfig#validateAll()} or either
|
||||
* call site — there is no "wire it in" step left to forget.
|
||||
*/
|
||||
static class FourValidators {
|
||||
final List<String> ran = new ArrayList<>();
|
||||
|
||||
public void validateAlpha() {
|
||||
ran.add("validateAlpha");
|
||||
}
|
||||
|
||||
public void validateBeta() {
|
||||
ran.add("validateBeta");
|
||||
}
|
||||
|
||||
public void validateGamma() {
|
||||
ran.add("validateGamma");
|
||||
}
|
||||
|
||||
public void validateDelta() {
|
||||
ran.add("validateDelta");
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void addingAFourthValidatorMethodGetsSweptWithNoOtherChange() {
|
||||
FourValidators target = new FourValidators();
|
||||
FleetConfig.invokeAllValidators(target);
|
||||
assertEquals(List.of("validateAlpha", "validateBeta", "validateDelta", "validateGamma"),
|
||||
sorted(target.ran),
|
||||
"the fourth method must be reached automatically — proving a class can grow the "
|
||||
+ "set of things it validates with no change to the sweep itself");
|
||||
}
|
||||
|
||||
private static List<String> sorted(List<String> in) {
|
||||
List<String> copy = new ArrayList<>(in);
|
||||
copy.sort(String::compareTo);
|
||||
return copy;
|
||||
}
|
||||
|
||||
@Test
|
||||
void aFailingValidatorStopsTheSweepAndPropagatesUnchanged() {
|
||||
class OneFails {
|
||||
@SuppressWarnings("unused")
|
||||
public void validateOk() {
|
||||
// passes
|
||||
}
|
||||
|
||||
public void validateBoom() {
|
||||
throw new IllegalStateException("refusing to start: boom");
|
||||
}
|
||||
}
|
||||
IllegalStateException e = assertThrows(IllegalStateException.class,
|
||||
() -> FleetConfig.invokeAllValidators(new OneFails()));
|
||||
assertEquals("refusing to start: boom", e.getMessage(),
|
||||
"the real exception must propagate unchanged, not be wrapped or swallowed");
|
||||
}
|
||||
|
||||
/**
|
||||
* Every rule the sweep's filter applies, proven independently: only public, no-arg, void
|
||||
* methods whose name starts with {@code "validate"} run, {@code validateAll} itself is
|
||||
* excluded (so a class that declares one of its own — as {@link FleetConfig} does — cannot
|
||||
* recurse into itself), and a same-shaped-but-wrongly-named or wrongly-shaped method never
|
||||
* runs. A reader who "simplifies" the filter in {@link FleetConfig#invokeAllValidators} in a
|
||||
* way that widens or narrows it breaks one of these.
|
||||
*/
|
||||
static class FilterEdgeCases {
|
||||
final List<String> ran = new ArrayList<>();
|
||||
|
||||
public void validateReal() {
|
||||
ran.add("validateReal");
|
||||
}
|
||||
|
||||
/** Wrong name — must not run. */
|
||||
public void checkSomething() {
|
||||
ran.add("checkSomething");
|
||||
}
|
||||
|
||||
/** Wrong shape — takes an argument. */
|
||||
public void validateWithArg(String ignored) {
|
||||
ran.add("validateWithArg");
|
||||
}
|
||||
|
||||
/** Wrong shape — returns something. */
|
||||
public boolean validateReturnsBoolean() {
|
||||
ran.add("validateReturnsBoolean");
|
||||
return true;
|
||||
}
|
||||
|
||||
/** Excluded by name on purpose, so the sweep cannot call itself. */
|
||||
public void validateAll() {
|
||||
ran.add("validateAll");
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void onlyPublicNoArgVoidValidateNamedMethodsRun() {
|
||||
FilterEdgeCases target = new FilterEdgeCases();
|
||||
FleetConfig.invokeAllValidators(target);
|
||||
assertEquals(List.of("validateReal"), target.ran,
|
||||
"checkSomething (wrong name), validateWithArg (wrong shape), "
|
||||
+ "validateReturnsBoolean (wrong shape), and validateAll (excluded by "
|
||||
+ "name) must all be skipped");
|
||||
}
|
||||
|
||||
// ── Claim 2: FleetConfig.validateAll() is wired to that mechanism and reaches all six today ──
|
||||
|
||||
/**
|
||||
* Reflectively enumerates {@link FleetConfig}'s own public, no-arg, void {@code validateXxx()}
|
||||
* methods (excluding {@code validateAll} itself) — the exact same filter {@link
|
||||
* FleetConfig#invokeAllValidators} applies. This is not the mechanism proof (that is claim 1,
|
||||
* above, on an unrelated class) — it is a visible denominator: today there are six, named
|
||||
* here, so a reader adding a seventh sees this assertion name the new count rather than a
|
||||
* silent pass at the old one.
|
||||
*/
|
||||
@Test
|
||||
void fleetConfigDeclaresExactlyTheseSixValidatorsToday() {
|
||||
Set<String> names = new TreeSet<>();
|
||||
for (Method m : FleetConfig.class.getMethods()) {
|
||||
if (java.lang.reflect.Modifier.isPublic(m.getModifiers())
|
||||
&& m.getParameterCount() == 0
|
||||
&& m.getReturnType() == void.class
|
||||
&& m.getName().startsWith("validate")
|
||||
&& !m.getName().equals("validateAll")) {
|
||||
names.add(m.getName());
|
||||
}
|
||||
}
|
||||
assertEquals(new TreeSet<>(Set.of("validateAuthExposure", "validateLeadTabPrefixes",
|
||||
"validateSubscriptionProfiles", "validateCharters", "validateMembers",
|
||||
"validateModels")), names,
|
||||
"FleetConfig's public validate*() methods changed. Do TWO things, in this "
|
||||
+ "order. First confirm validateAll() still delegates to "
|
||||
+ "invokeAllValidators(this) — a hardcoded list there passes every other "
|
||||
+ "test in this class, so this assertion is the only place that will ever "
|
||||
+ "make you check. Only then update the expected set to match.");
|
||||
}
|
||||
|
||||
/** A minimal, otherwise-valid file — same shape FleetConfigTest and ConfigRefTest use. */
|
||||
private static String minimalValidYaml() {
|
||||
return """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
profiles:
|
||||
sonnet:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
model: sonnet
|
||||
""";
|
||||
}
|
||||
|
||||
@Test
|
||||
void aFullyValidConfigPassesValidateAll(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, minimalValidYaml());
|
||||
assertDoesNotThrow(() -> FleetConfig.load(f).validateAll());
|
||||
}
|
||||
|
||||
/**
|
||||
* The heart of claim 2: for each of today's six real validators, a minimal file that fails
|
||||
* ONLY that one — the exact fixtures {@code FleetConfigTest} uses to test each validator
|
||||
* directly — must also fail through {@link FleetConfig#validateAll()}. If a future edit to
|
||||
* {@code validateAll()} silently dropped one validator from the sweep (e.g. a typo'd name
|
||||
* filter), exactly one of these six would start passing when it must not.
|
||||
*/
|
||||
@Test
|
||||
void validateAllReachesEveryOneOfTodaysSixValidators(@TempDir Path dir) throws Exception {
|
||||
// validateAuthExposure: a non-loopback bind without token mode.
|
||||
assertValidateAllRefuses(dir, "auth-exposure.yaml", """
|
||||
bind:
|
||||
host: 0.0.0.0
|
||||
port: 8765
|
||||
""", "auth.mode: token");
|
||||
|
||||
// validateLeadTabPrefixes: a fleet-wide tabLabel that starts with a lead's own tabPrefix.
|
||||
assertValidateAllRefuses(dir, "lead-tab-prefixes.yaml", """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
fleet:
|
||||
tabLabel: "lead: {role} {profile}"
|
||||
leaders:
|
||||
opus:
|
||||
tab: "lead: opus"
|
||||
""", "fleet.tabLabel");
|
||||
|
||||
// validateSubscriptionProfiles: subscription: true with env: reseating ANTHROPIC_BASE_URL.
|
||||
assertValidateAllRefuses(dir, "subscription-profiles.yaml", """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
profiles:
|
||||
sonnet:
|
||||
subscription: true
|
||||
argv: ["ccs", "sonnet"]
|
||||
env:
|
||||
ANTHROPIC_BASE_URL: http://anything-not-on-the-allowlist
|
||||
""", "ANTHROPIC_BASE_URL");
|
||||
|
||||
// validateCharters: a charter key that is not a role wire name.
|
||||
assertValidateAllRefuses(dir, "charters.yaml", """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
fleet:
|
||||
charters:
|
||||
architetc: text
|
||||
""", "architetc");
|
||||
|
||||
// validateMembers: an architect slot naming an unconfigured profile.
|
||||
assertValidateAllRefuses(dir, "members.yaml", """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
profiles:
|
||||
gx10:
|
||||
baseUrl: http://gx10.gw:8000
|
||||
fleet:
|
||||
architects:
|
||||
lead-designer:
|
||||
profile: sonnet
|
||||
""", "lead-designer");
|
||||
|
||||
// validateModels: a profile naming a model outside the configured allow-list.
|
||||
assertValidateAllRefuses(dir, "models.yaml", """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
profiles:
|
||||
sonnet:
|
||||
baseUrl: http://gx10.gw:8000
|
||||
model: claude-sonnet-5
|
||||
rogue:
|
||||
baseUrl: http://gx11.gw:8000
|
||||
model: claude-opus-9000
|
||||
models:
|
||||
allow:
|
||||
- model: claude-sonnet-5
|
||||
""", "rogue");
|
||||
}
|
||||
|
||||
private static void assertValidateAllRefuses(Path dir, String fileName, String yaml,
|
||||
String mustContain) throws Exception {
|
||||
Path f = dir.resolve(fileName);
|
||||
Files.writeString(f, yaml);
|
||||
FleetConfig cfg = FleetConfig.load(f);
|
||||
IllegalStateException e = assertThrows(IllegalStateException.class, cfg::validateAll,
|
||||
fileName + ": validateAll() must refuse this config");
|
||||
assertTrue(e.getMessage().contains(mustContain),
|
||||
fileName + ": expected message to contain \"" + mustContain + "\" but was: "
|
||||
+ e.getMessage());
|
||||
}
|
||||
}
|
||||
+2
@@ -97,6 +97,8 @@ class FleetConfigWithDefaultsPreservesEveryComponentTest {
|
||||
v.put("memberLoginShell", "/bin/zsh");
|
||||
v.put("memberSkills", "/skills/guard");
|
||||
v.put("idleSleepGuard", new FleetConfig.IdleSleepGuard(true));
|
||||
v.put("models", new FleetConfig.Models(
|
||||
List.of(new FleetConfig.Models.ModelEntry("model-guard"))));
|
||||
assertNamesMatchComponents(v);
|
||||
return v;
|
||||
}
|
||||
|
||||
@@ -3,6 +3,7 @@ package dev.ltms.fleet.mcp;
|
||||
import dev.ltms.fleet.auth.CallerResolver;
|
||||
import dev.ltms.fleet.auth.MemberRegistry;
|
||||
import dev.ltms.fleet.auth.Principal;
|
||||
import dev.ltms.fleet.auth.Role;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
@@ -28,6 +29,7 @@ import dev.ltms.fleet.placement.BackendQuarantine;
|
||||
import dev.ltms.fleet.placement.PlacementPolicies;
|
||||
import io.modelcontextprotocol.spec.McpSchema;
|
||||
import dev.ltms.fleet.msg.InMemoryReplyInbox;
|
||||
import dev.ltms.fleet.msg.ReplyInbox;
|
||||
import org.junit.jupiter.api.BeforeEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
@@ -77,7 +79,7 @@ class FleetMcpTest {
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertTrue(rendezvous.isWaiting(target), "send should be accepted for " + target);
|
||||
FleetMcp.reply(messages, target, "received");
|
||||
FleetMcp.reply(messages, target, Role.WORKER, "received");
|
||||
assertEquals("received", textOf(send.get(6, TimeUnit.SECONDS)));
|
||||
}
|
||||
|
||||
@@ -110,7 +112,7 @@ class FleetMcpTest {
|
||||
|
||||
// fleetd #365: a resolved live send must read distinctly from a merely-queued reply —
|
||||
// see replyWithNoPendingSendIsQueuedNotError below for the other case.
|
||||
McpSchema.CallToolResult reply = FleetMcp.reply(messages, "term_a", "LGTM");
|
||||
McpSchema.CallToolResult reply = FleetMcp.reply(messages, "term_a", Role.WORKER, "LGTM");
|
||||
assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND.description(), textOf(reply));
|
||||
|
||||
McpSchema.CallToolResult res = send.get(6, TimeUnit.SECONDS);
|
||||
@@ -136,7 +138,7 @@ class FleetMcpTest {
|
||||
}
|
||||
assertTrue(rendezvous.isWaiting("term_a"), "send should have opened its waiter");
|
||||
|
||||
McpSchema.CallToolResult reply = FleetMcp.reply(messages, "term_a", "async LGTM");
|
||||
McpSchema.CallToolResult reply = FleetMcp.reply(messages, "term_a", Role.WORKER, "async LGTM");
|
||||
assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND.description(), textOf(reply));
|
||||
|
||||
// Poll until the async send completes and reports the reply.
|
||||
@@ -183,7 +185,7 @@ class FleetMcpTest {
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertTrue(rendezvous.isWaiting("term_a"));
|
||||
FleetMcp.reply(messages, "term_a", "done");
|
||||
FleetMcp.reply(messages, "term_a", Role.WORKER, "done");
|
||||
assertEquals("done", textOf(answer.get(6, TimeUnit.SECONDS)));
|
||||
|
||||
McpSchema.CallToolResult done = FleetMcp.poll(messages, ticket, null);
|
||||
@@ -215,7 +217,7 @@ class FleetMcpTest {
|
||||
// fleet_poll{ticket} stuck PENDING forever and later force-failed with a false "session
|
||||
// released before it replied" reason. This used to land in the inbox instead (see the old
|
||||
// assertion this replaced: messages.drainReplies("term_a").getFirst()...) — that was the bug.
|
||||
FleetMcp.reply(messages, "term_a", "finished after timeout");
|
||||
FleetMcp.reply(messages, "term_a", Role.WORKER, "finished after timeout");
|
||||
assertEquals("finished after timeout", textOf(FleetMcp.poll(messages, ticket, null)));
|
||||
assertTrue(messages.drainReplies("term_a").isEmpty(),
|
||||
"the reply completed its own ticket directly and never touched the inbox");
|
||||
@@ -269,7 +271,7 @@ class FleetMcpTest {
|
||||
}
|
||||
assertEquals(MessageService.Phase.FAILED, second.phase());
|
||||
|
||||
FleetMcp.reply(messages, "term_a", "late reply");
|
||||
FleetMcp.reply(messages, "term_a", Role.WORKER, "late reply");
|
||||
assertEquals("late reply", messages.drainReplies("term_a").getFirst().content());
|
||||
|
||||
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
|
||||
@@ -278,7 +280,7 @@ class FleetMcpTest {
|
||||
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
|
||||
Thread.sleep(5);
|
||||
}
|
||||
FleetMcp.reply(messages, "term_a", "done");
|
||||
FleetMcp.reply(messages, "term_a", Role.WORKER, "done");
|
||||
assertEquals("done", textOf(answer.get(6, TimeUnit.SECONDS)));
|
||||
}
|
||||
|
||||
@@ -331,7 +333,7 @@ class FleetMcpTest {
|
||||
void replyWithNoPendingSendIsQueuedNotError() {
|
||||
// CB-307: a reply with no open send is now queued in the inbox, not an error.
|
||||
// fleetd #365: it must also no longer claim "delivered" — nothing was waiting for it.
|
||||
McpSchema.CallToolResult res = FleetMcp.reply(messages, "term_a", "orphan");
|
||||
McpSchema.CallToolResult res = FleetMcp.reply(messages, "term_a", Role.WORKER, "orphan");
|
||||
assertNotEquals(Boolean.TRUE, res.isError(), "a queued reply is not an error");
|
||||
assertEquals(MessageService.ReplyOutcome.QUEUED.description(), textOf(res));
|
||||
|
||||
@@ -341,6 +343,40 @@ class FleetMcpTest {
|
||||
assertEquals("orphan", drained.getFirst().content());
|
||||
}
|
||||
|
||||
@Test
|
||||
void replyFromLeadIsRefusedBeforeItCanPublishToTheWorkerInbox() {
|
||||
ReplyInbox inboxThatRejectsPublishes = new ReplyInbox() {
|
||||
@Override public void own(String target) { }
|
||||
@Override public void release(String target) { }
|
||||
@Override public void publish(String target, String msgId, String content) {
|
||||
fail("a lead fleet_reply must not publish to the worker inbox");
|
||||
}
|
||||
@Override public List<InboxMessage> peek(String target) { return List.of(); }
|
||||
@Override public void ack(String target, String msgId) { }
|
||||
};
|
||||
MessageService leadMessages = new MessageService(agents, new Injector(agents), new Rendezvous(),
|
||||
inboxThatRejectsPublishes);
|
||||
|
||||
McpSchema.CallToolResult res = assertDoesNotThrow(
|
||||
() -> FleetMcp.reply(leadMessages, "term_lead", Role.PRIMARY, "peer reply"));
|
||||
|
||||
assertTrue(res.isError());
|
||||
assertEquals("fleet_reply has no route to a peer lead. Use fleet_send{coordId: ...} for a peer on another "
|
||||
+ "daemon or fleet_send{sessionId: ...} for a peer on this host. fleet_reply resolves a member's "
|
||||
+ "blocked fleet_send, and a peer's coord-id message is durable and non-blocking, so there is "
|
||||
+ "nothing for it to resolve.",
|
||||
textOf(res));
|
||||
}
|
||||
|
||||
@Test
|
||||
void replyFromUnidentifiedCallerKeepsItsOwnError() {
|
||||
McpSchema.CallToolResult res = FleetMcp.reply(messages, null, Role.PRIMARY, "reply");
|
||||
|
||||
assertTrue(res.isError());
|
||||
assertEquals("fleet_reply is for workers only — could not identify the calling worker from the connection",
|
||||
textOf(res));
|
||||
}
|
||||
|
||||
@Test
|
||||
void replyWithBlankContentIsACleanToolErrorNotAnUncaughtException() {
|
||||
// fleetd #302: MessageService.reply now REJECTS blank content by throwing. fleet_reply's
|
||||
@@ -351,7 +387,7 @@ class FleetMcpTest {
|
||||
// has always used isBlank for exactly this reason.
|
||||
for (String blank : new String[] {null, "", " ", "\n\t"}) {
|
||||
McpSchema.CallToolResult res = assertDoesNotThrow(
|
||||
() -> FleetMcp.reply(messages, "term_a", blank),
|
||||
() -> FleetMcp.reply(messages, "term_a", Role.WORKER, blank),
|
||||
"blank content must be refused as a tool error, never thrown out of the handler");
|
||||
assertEquals(Boolean.TRUE, res.isError(), "blank content is an error result");
|
||||
assertTrue(textOf(res).contains("content is required"),
|
||||
@@ -364,7 +400,7 @@ class FleetMcpTest {
|
||||
@Test
|
||||
void bridgePollWithTargetDrainsReplies() {
|
||||
// A reply with no open send queues it in the inbox.
|
||||
FleetMcp.reply(messages, "term_a", "queued-msg");
|
||||
FleetMcp.reply(messages, "term_a", Role.WORKER, "queued-msg");
|
||||
|
||||
// fleet_poll with target drains the inbox.
|
||||
McpSchema.CallToolResult res = FleetMcp.poll(messages, null, "term_a");
|
||||
@@ -415,7 +451,7 @@ class FleetMcpTest {
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertTrue(rendezvous.isWaiting("term_a"), "the answer should have reopened a waiter");
|
||||
McpSchema.CallToolResult reply = FleetMcp.reply(messages, "term_a", "done");
|
||||
McpSchema.CallToolResult reply = FleetMcp.reply(messages, "term_a", Role.WORKER, "done");
|
||||
assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND.description(), textOf(reply));
|
||||
assertEquals("done", textOf(answer.get(6, TimeUnit.SECONDS)));
|
||||
}
|
||||
@@ -1254,13 +1290,13 @@ class FleetMcpTest {
|
||||
@Test
|
||||
void bridgeAckRemovesSpecificReply() {
|
||||
// Queue a reply and capture its msgId.
|
||||
FleetMcp.reply(messages, "term_a", "orphan");
|
||||
FleetMcp.reply(messages, "term_a", Role.WORKER, "orphan");
|
||||
var before = messages.drainReplies("term_a");
|
||||
assertEquals(1, before.size(), "one reply in the inbox");
|
||||
String msgId = before.getFirst().msgId();
|
||||
|
||||
// Publish the same reply again and ack it via fleet_ack surface.
|
||||
FleetMcp.reply(messages, "term_a", "orphan-again");
|
||||
FleetMcp.reply(messages, "term_a", Role.WORKER, "orphan-again");
|
||||
var peeked = messages.drainReplies("term_a");
|
||||
assertEquals(1, peeked.size(), "one fresh reply in the inbox");
|
||||
|
||||
|
||||
@@ -0,0 +1,62 @@
|
||||
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.peer.PeerLauncher;
|
||||
import io.modelcontextprotocol.spec.McpSchema;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertNotEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #395: {@code fleet_profiles} must let an operator tell "this profile's usage-limit
|
||||
* detection is armed" from "nothing can ever quarantine this profile" — see {@link
|
||||
* FleetMcp.QuarantineSource#exhaustedPatternArmed()} and {@link
|
||||
* FleetMcp#profilesView(PeerLauncher, FleetMcp.QuarantineSource, FleetMcp.OutageSource)}.
|
||||
*/
|
||||
class FleetProfilesArmedFieldTest {
|
||||
|
||||
private static PeerLauncher twoProfileLauncher(FakeHerdr h) {
|
||||
FleetConfig.Profile armed = new FleetConfig.Profile(
|
||||
"armed-profile", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN", null,
|
||||
"tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
|
||||
FleetConfig.Profile unarmed = new FleetConfig.Profile(
|
||||
"unarmed-profile", "http://gx01.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN", null,
|
||||
"tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
|
||||
Map<String, FleetConfig.Profile> profiles = new LinkedHashMap<>();
|
||||
profiles.put(armed.profile(), armed);
|
||||
profiles.put(unarmed.profile(), unarmed);
|
||||
return new ClaudeCodeLauncher(new AgentControl(h), new WorkspaceControl(h),
|
||||
new SubscriptionGuard(Set.of("gx00.gw", "gx01.gw")), profiles, armed.profile(), _ -> "tok");
|
||||
}
|
||||
|
||||
private static String textOf(McpSchema.CallToolResult r) {
|
||||
return ((McpSchema.TextContent) r.content().getFirst()).text();
|
||||
}
|
||||
|
||||
@Test
|
||||
void armedProfileReportsArmedAndUnarmedReportsUnarmed() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
PeerLauncher workers = twoProfileLauncher(h);
|
||||
FleetMcp.QuarantineSource source = new FleetMcp.QuarantineSource(
|
||||
_ -> null, dev.ltms.fleet.placement.BackendQuarantine.none(),
|
||||
profile -> "armed-profile".equals(profile));
|
||||
|
||||
McpSchema.CallToolResult res = FleetMcp.profiles(workers, source);
|
||||
assertNotEquals(Boolean.TRUE, res.isError());
|
||||
String out = textOf(res);
|
||||
|
||||
assertTrue(out.contains("\"exhaustionDetectionArmed\""), out);
|
||||
assertTrue(out.contains("\"armed-profile\":true"), out);
|
||||
assertTrue(out.contains("\"unarmed-profile\":false"), out);
|
||||
}
|
||||
}
|
||||
@@ -18,6 +18,7 @@ import java.util.regex.Matcher;
|
||||
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.assertNotNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
@@ -118,6 +119,205 @@ class EnvAllowListScrubTest {
|
||||
"allowed N of M with N <= M — the denominator is always reported");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #394: the actual defect. Plain {@code export "$n="} is FATAL for a zsh read-only or
|
||||
* special parameter (e.g. {@code UID}) and aborts the whole sourced file — every name still to
|
||||
* come is never blanked, and the {@code scrub-report.txt} below is never written at all,
|
||||
* silently ({@code 2>/dev/null} swallows the error). This plants an unblankable, exported,
|
||||
* read-only variable in the MIDDLE of the names the scrub attempts to blank, with two more
|
||||
* names after it, and asserts that both of those later names are STILL blanked and the report
|
||||
* is STILL written with the failure counted — a test that only checked names BEFORE the failure
|
||||
* point would pass today and prove nothing.
|
||||
*
|
||||
* <p>The planted name is a made-up one ({@code FLEETD_TEST_UNBLANKABLE}), not {@code UID} or
|
||||
* any other name a skip-list might already know about — invariant 1 is that the loop survives
|
||||
* ANY unblankable name, not a known one, so the test must not lean on one either.
|
||||
*
|
||||
* <p>Exercises the real artefact: {@link EnvAllowListScrub#scrubScript} is run verbatim under a
|
||||
* real {@code /bin/zsh}, not just asserted on as a Java string. The four planted names are
|
||||
* exported one at a time via {@code typeset -x}/{@code typeset -rx} immediately before the
|
||||
* script runs, in a fixed order — zsh's {@code export}/{@code typeset -x} appends to the
|
||||
* process's environment table in call order (verified empirically: a freshly-exported name
|
||||
* always sorts after every inherited one and after every earlier freshly-exported name in
|
||||
* {@code command env}'s own output), which is what makes the "middle" position deterministic
|
||||
* here, unlike relying on the OS's own inherited-environment order.
|
||||
*/
|
||||
@Test
|
||||
void unblankableNameInTheMiddleDoesNotAbortNamesAfterIt(@TempDir Path tmp) throws Exception {
|
||||
assumeTrue(Files.isExecutable(ZSH), "/bin/zsh not present — nothing to prove here");
|
||||
|
||||
// Only ZDOTDIR is allowed — it must survive the scrub itself, since the report is written
|
||||
// to "$ZDOTDIR/..." AFTER the blanking loop runs; if ZDOTDIR were blanked as a side effect,
|
||||
// the report write would silently go to the wrong place instead of testing anything.
|
||||
String script = EnvAllowListScrub.scrubScript(Set.of("ZDOTDIR"));
|
||||
String setup = """
|
||||
typeset -x FLEETD_TEST_BEFORE=1
|
||||
typeset -rx FLEETD_TEST_UNBLANKABLE=1
|
||||
typeset -x FLEETD_TEST_AFTER_A=1
|
||||
typeset -x FLEETD_TEST_AFTER_B=1
|
||||
""";
|
||||
|
||||
ProcessBuilder pb = new ProcessBuilder("/bin/zsh");
|
||||
pb.environment().clear();
|
||||
pb.environment().put("PATH", "/usr/bin:/bin");
|
||||
pb.environment().put("ZDOTDIR", tmp.toAbsolutePath().toString());
|
||||
pb.redirectError(ProcessBuilder.Redirect.DISCARD);
|
||||
Process zsh = pb.start();
|
||||
zsh.getOutputStream().write((setup + script).getBytes(StandardCharsets.UTF_8));
|
||||
zsh.getOutputStream().flush();
|
||||
zsh.getOutputStream().close();
|
||||
assertTrue(zsh.waitFor(60, java.util.concurrent.TimeUnit.SECONDS),
|
||||
"the scrub script did not exit within 60s");
|
||||
assertEquals(0, zsh.exitValue(),
|
||||
"the scrub script itself must never abort — an unblankable name must not kill the "
|
||||
+ "sourced file");
|
||||
|
||||
EnvAllowListScrub.ScrubReport report = EnvAllowListScrub.readReport(tmp);
|
||||
assertNotNull(report, "the report must still be written even though one name could not be "
|
||||
+ "blanked — a report that silently never appears is the #394 bug");
|
||||
assertTrue(report.blanked().contains("FLEETD_TEST_BEFORE"),
|
||||
"sanity: the name before the unblankable one must be blanked");
|
||||
assertTrue(report.blanked().contains("FLEETD_TEST_AFTER_A"),
|
||||
"the FIRST name AFTER the unblankable one must still be blanked — before the fix, "
|
||||
+ "the whole loop aborted at the unblankable name and every later name was "
|
||||
+ "silently left untouched");
|
||||
assertTrue(report.blanked().contains("FLEETD_TEST_AFTER_B"),
|
||||
"the SECOND name after the unblankable one must also still be blanked");
|
||||
assertTrue(report.unblankable().contains("FLEETD_TEST_UNBLANKABLE"),
|
||||
"the unblankable name is reported by name, not silently dropped");
|
||||
// fleetd #400 note: zsh itself auto-exports SHLVL on every shell start (measured: it appears
|
||||
// in `command env` even from a fully cleared parent), and a bare assignment to it is coerced
|
||||
// rather than failing — exactly the shape #400 fixes. Before that fix, eval's exit status
|
||||
// alone silently misclassified SHLVL as blanked, so this test's old "exactly one" assertion
|
||||
// passed by accident: it never actually proved SHLVL was absent from the candidates, only
|
||||
// that the old bug hid it. Now that classification reads the value back, SHLVL and
|
||||
// FLEETD_TEST_UNBLANKABLE both correctly land in unblankable() — real, unplanted evidence
|
||||
// the #400 fix works, not just the synthetic case in the dedicated #400 test above.
|
||||
assertTrue(report.unblankable().contains("SHLVL"),
|
||||
"fleetd #400: zsh's own auto-exported SHLVL must also be reported unblankable, not "
|
||||
+ "silently miscounted as blanked");
|
||||
assertEquals(report.unblankable().size(), report.failed(),
|
||||
"the failed count must equal the number of names actually reported unblankable");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #400: {@code eval}'s exit status is not proof that a name was actually blanked. zsh
|
||||
* coerces a bare {@code NAME=} assignment on an integer special parameter to a number instead of
|
||||
* failing, so {@code eval} reports success while the value stays non-empty — a status-based
|
||||
* classification calls that "blanked" when it was not. This drives all three shapes a name can
|
||||
* take through the real {@code scrubScript} in ONE run: a normal, genuinely blankable name; a
|
||||
* fatal one ({@code LINENO} — deliberately not {@code UID}, so this test does not depend on the
|
||||
* harness's uid); and the silent-no-op one the ticket is about ({@code SECONDS}, rc 0 but
|
||||
* unchanged). Under the pre-#400 exit-status check, {@code SECONDS} would land in
|
||||
* {@code blanked()} — that is the exact false receipt this fix removes.
|
||||
*
|
||||
* <p>Criterion 3: a cleared {@code ProcessBuilder} parent does not, by itself, give the child
|
||||
* zsh any of these names — {@code SECONDS}/{@code LINENO} are zsh's own built-in parameters and
|
||||
* only become CANDIDATES the enumeration loop can see (i.e. show up in {@code command env}) when
|
||||
* they arrive via the process's own environment table, not merely by existing as zsh parameters
|
||||
* inside the shell. So each is put into {@code pb.environment()} explicitly, after
|
||||
* {@code clear()} — confirmed empirically first (a throwaway probe piping
|
||||
* {@code env -i PATH=... SECONDS=999 LINENO=999 FLEETD_TEST_NORMAL=1 zsh -c 'command env | cut
|
||||
* -d= -f1'}) that all three names really appear in {@code command env}'s output under exactly
|
||||
* this construction, not relying on whatever the test-runner's own ambient environment happens
|
||||
* to contain.
|
||||
*/
|
||||
@Test
|
||||
void classifiesByObservedValueNotExitStatusAcrossAllThreeShapes(@TempDir Path tmp) throws Exception {
|
||||
assumeTrue(Files.isExecutable(ZSH), "/bin/zsh not present — nothing to prove here");
|
||||
|
||||
String script = EnvAllowListScrub.scrubScript(Set.of("ZDOTDIR"));
|
||||
|
||||
ProcessBuilder pb = new ProcessBuilder("/bin/zsh");
|
||||
pb.environment().clear();
|
||||
pb.environment().put("PATH", "/usr/bin:/bin");
|
||||
pb.environment().put("ZDOTDIR", tmp.toAbsolutePath().toString());
|
||||
// Explicitly placed in the child's environment table — see the javadoc above on why a
|
||||
// cleared parent alone does not put these on the enumeration loop's candidate list.
|
||||
pb.environment().put("SECONDS", "999"); // rc 0, value coerced/unchanged — the #400 bug
|
||||
pb.environment().put("LINENO", "999"); // fatal on assignment, eval rc != 0, contained
|
||||
pb.environment().put("FLEETD_TEST_NORMAL", "1"); // genuinely blankable, the control case
|
||||
pb.redirectError(ProcessBuilder.Redirect.DISCARD);
|
||||
Process zsh = pb.start();
|
||||
zsh.getOutputStream().write(script.getBytes(StandardCharsets.UTF_8));
|
||||
zsh.getOutputStream().flush();
|
||||
zsh.getOutputStream().close();
|
||||
assertTrue(zsh.waitFor(60, java.util.concurrent.TimeUnit.SECONDS),
|
||||
"the scrub script did not exit within 60s");
|
||||
assertEquals(0, zsh.exitValue(),
|
||||
"the scrub script must still reach its end with both a fatal name and a silent "
|
||||
+ "no-op name among the candidates");
|
||||
|
||||
EnvAllowListScrub.ScrubReport report = EnvAllowListScrub.readReport(tmp);
|
||||
assertNotNull(report, "the report must still be written");
|
||||
assertTrue(report.blanked().contains("FLEETD_TEST_NORMAL"),
|
||||
"the control case: an ordinary name is genuinely blankable and must be reported so");
|
||||
assertTrue(report.unblankable().contains("LINENO"),
|
||||
"a name fatal to assign to must be reported unblankable — sanity check that "
|
||||
+ "containment still works under the new classification");
|
||||
assertTrue(report.unblankable().contains("SECONDS"),
|
||||
"the #400 defect: eval returns rc 0 for SECONDS (zsh coerces the assignment instead "
|
||||
+ "of failing) but the value is left non-empty — classifying on the observed "
|
||||
+ "value catches this; classifying on eval's exit status would have called "
|
||||
+ "this \"blanked\" and produced a false receipt");
|
||||
assertFalse(report.blanked().contains("SECONDS"),
|
||||
"SECONDS must never appear as blanked — it was never actually emptied");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #394 follow-up: the blanking loop's {@code eval "export ${n}="} splices {@code n} into
|
||||
* a string that zsh then interprets as shell syntax. That is only safe because every name
|
||||
* reaching {@code _cb633_blank} already passed an identifier check in the ENUMERATION loop
|
||||
* (20 lines away, in a different loop) — so the fix re-asserts the identical check immediately
|
||||
* before the {@code eval} call, rather than trusting that distant guard to keep holding.
|
||||
*
|
||||
* <p>This test plants a value with an embedded newline, exploiting the exact "junk from
|
||||
* multi-line values" gap the enumeration loop's own comment already documents: {@code command
|
||||
* env}'s text output is read line-by-line, so a value's second line becomes a spurious extra
|
||||
* "name" that was never a real exported variable. The fragment used here ({@code
|
||||
* junk.fragment}) is merely non-conforming (it contains a dot) — never command-shaped; this
|
||||
* test must never demonstrate command execution and plants no command-shaped payload.
|
||||
*
|
||||
* <p>Exercises the real artefact end-to-end: {@link EnvAllowListScrub#scrubScript} runs
|
||||
* verbatim under a real {@code /bin/zsh}, exactly as {@code generate()} would produce it — this
|
||||
* is not a synthetic call into just the blanking loop.
|
||||
*/
|
||||
@Test
|
||||
void nonIdentifierJunkFromAMultilineValueIsSkippedNotBlankedOrUnblankable(@TempDir Path tmp)
|
||||
throws Exception {
|
||||
assumeTrue(Files.isExecutable(ZSH), "/bin/zsh not present — nothing to prove here");
|
||||
|
||||
String script = EnvAllowListScrub.scrubScript(Set.of("ZDOTDIR"));
|
||||
|
||||
ProcessBuilder pb = new ProcessBuilder("/bin/zsh");
|
||||
pb.environment().clear();
|
||||
pb.environment().put("PATH", "/usr/bin:/bin");
|
||||
pb.environment().put("ZDOTDIR", tmp.toAbsolutePath().toString());
|
||||
// Embedded newline: `command env`'s own text output splits this into two lines, and the
|
||||
// second ("junk.fragment") has no "=" at all, so `cut -d= -f1` returns it unchanged as a
|
||||
// spurious candidate "name" — it was never an actual exported variable by that name.
|
||||
pb.environment().put("FLEETD_TEST_MULTILINE", "keep\njunk.fragment");
|
||||
pb.redirectError(ProcessBuilder.Redirect.DISCARD);
|
||||
Process zsh = pb.start();
|
||||
zsh.getOutputStream().write(script.getBytes(StandardCharsets.UTF_8));
|
||||
zsh.getOutputStream().flush();
|
||||
zsh.getOutputStream().close();
|
||||
assertTrue(zsh.waitFor(60, java.util.concurrent.TimeUnit.SECONDS),
|
||||
"the scrub script did not exit within 60s");
|
||||
assertEquals(0, zsh.exitValue(), "the scrub script must reach its end");
|
||||
|
||||
EnvAllowListScrub.ScrubReport report = EnvAllowListScrub.readReport(tmp);
|
||||
assertNotNull(report, "the report must still be written");
|
||||
assertTrue(report.blanked().contains("FLEETD_TEST_MULTILINE"),
|
||||
"sanity: the real, identifier-shaped variable must still be blanked normally");
|
||||
assertFalse(report.blanked().contains("junk.fragment"),
|
||||
"a non-identifier fragment is not a real variable and must never be blanked");
|
||||
assertFalse(report.unblankable().contains("junk.fragment"),
|
||||
"a non-identifier fragment must never even become a candidate the blanking loop "
|
||||
+ "attempts — it must be filtered before either guard has to catch it, so "
|
||||
+ "it is neither blanked nor counted as a failed attempt");
|
||||
}
|
||||
|
||||
/** A group-shared ZDOTDIR still lets the member truncate and write its pre-created receipt. */
|
||||
@Test
|
||||
void groupSharedScrubWritesAndReadsItsReport(@TempDir Path tmp) throws Exception {
|
||||
|
||||
@@ -1802,7 +1802,11 @@ class MessageServiceTest {
|
||||
Thread asker = new Thread(() -> assertThrows(IllegalStateException.class,
|
||||
() -> service.ask(T, "which config file?", 30_000)));
|
||||
asker.start();
|
||||
awaitTicketPhaseOn(service, ticket, MessageService.Phase.ASKING);
|
||||
// Phase.ASKING (markAsyncQuestion) is only the FIRST of ask()'s three steps; the assertion
|
||||
// below depends on the THIRD (pushLoop.onQuestionOpened). Barrier on the push loop's own
|
||||
// pending-question state instead of the phase — see ReplyPushLoop#pendingQuestionTurnIdsForTest.
|
||||
MessageService.TaskView asking = awaitTicketPhaseOn(service, ticket, MessageService.Phase.ASKING);
|
||||
awaitQuestionPendingOn(pushLoop, LEAD, asking.turnId());
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, pushLoop.decide(LEAD, 0, 0, 0),
|
||||
"the open question should be the one thing keeping this lead's schedule alive");
|
||||
|
||||
@@ -1917,6 +1921,12 @@ class MessageServiceTest {
|
||||
// Only now does it reply. Under the old clock this reply was born already expired.
|
||||
assertTrue(rendezvous.resolve(T, "the long report"));
|
||||
awaitTicketPhaseOn(wiring.service(), slow, MessageService.Phase.DONE);
|
||||
// #399: same barrier fix as the sibling eviction test, for consistency — this test's
|
||||
// own assertion happens to survive a late stamp today (the clock is not advanced any
|
||||
// further between here and the sweep below, so cutoff cannot move past whichever
|
||||
// value completedNanos ends up stamped with), but DONE is still the wrong thing to
|
||||
// order on before a sweep that matters for the TTL.
|
||||
awaitCompletionStamped(wiring.service(), slow);
|
||||
|
||||
// A second delegation runs pruneTerminalTickets before it returns.
|
||||
wiring.service().sendAsync(T, "an unrelated second task");
|
||||
@@ -1946,6 +1956,11 @@ class MessageServiceTest {
|
||||
injectDelivery();
|
||||
assertTrue(rendezvous.resolve(T, "quick result"));
|
||||
awaitTicketPhaseOn(wiring.service(), done, MessageService.Phase.DONE);
|
||||
// #399: DONE can be observed before the completion hook stamps completedNanos. Wait
|
||||
// for the real stamp before advancing the clock, or the hook can run late and stamp
|
||||
// the ADVANCED time — making cutoff = advanced - TTL unreachable and hiding the very
|
||||
// eviction this test exists to pin.
|
||||
awaitCompletionStamped(wiring.service(), done);
|
||||
|
||||
// Nobody collected it, and the TTL has now passed since it FINISHED.
|
||||
clock.addAndGet(MessageService.TICKET_TTL_NANOS + TimeUnit.SECONDS.toNanos(1));
|
||||
@@ -1956,6 +1971,80 @@ class MessageServiceTest {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #409: pins the #399 ordering invariant deterministically, on the first run, without
|
||||
* relying on host load.
|
||||
*
|
||||
* <p>fleetd #399 was itself only reproducible probabilistically: the real race window between
|
||||
* {@code CompletableFuture.complete()} making {@link MessageService.Phase#DONE} observable and
|
||||
* the constructor's {@code whenComplete} hook actually stamping {@code completedNanos} (see
|
||||
* {@link MessageService#isCompletionStampedForTest}) is normally a handful of instructions wide,
|
||||
* and needed heavy background load on the host to show up in a run at all. This test does not
|
||||
* try to hit that narrow window by chance — it widens it on purpose: {@link #STAMP_DELAY_MILLIS}
|
||||
* is injected into the clock itself, so the completion hook's one read of {@code nowNanos} for
|
||||
* this ticket sleeps before returning, and everything the test does in the meantime (advance the
|
||||
* clock, run a sweep, assert) happens for certain inside that window, on any host.
|
||||
*
|
||||
* <p>Removing the {@link #awaitCompletionStamped} call below reproduces the pre-#399 ordering:
|
||||
* the test then advances the clock and runs its sweep while the hook is still asleep, so the
|
||||
* hook wakes up and stamps {@code completedNanos} with the clock's ALREADY-ADVANCED value
|
||||
* instead of the real completion time. {@code cutoff = advanced - TICKET_TTL_NANOS} can then
|
||||
* never exceed that stamp (the gap between them is fixed at exactly {@code TICKET_TTL_NANOS}),
|
||||
* so the sweep that already ran never evicts the ticket and the very next assertion — expecting
|
||||
* eviction — fails immediately. That is a deterministic, first-run RED failure, not a flaky one
|
||||
* and not a false pass: I ran the test with the barrier call removed and confirmed
|
||||
* {@code assertNull} fails because {@code poll} still returns the ticket's DONE view, which is
|
||||
* exactly this masked-eviction mechanism and not some unrelated defect in the test's own wiring.
|
||||
*/
|
||||
@Test
|
||||
void aTicketOrderedOnDoneInsteadOfTheCompletionStampSurvivesAnEvictionItMustNotSurvive() throws Exception {
|
||||
java.util.concurrent.atomic.AtomicLong clock = new java.util.concurrent.atomic.AtomicLong(1_000_000_000L);
|
||||
// Armed for exactly one call: MessageService's only three nowNanos() call sites are the Task
|
||||
// constructor's createdNanos, this constructor's whenComplete hook's completedNanos, and
|
||||
// pruneTerminalTickets' cutoff — none of which run between arming this (right before
|
||||
// triggering the reply below) and the completion hook firing, so the delayed call is
|
||||
// unambiguously that hook's stamp for `ticket`, never a createdNanos or cutoff read.
|
||||
java.util.concurrent.atomic.AtomicBoolean delayArmed = new java.util.concurrent.atomic.AtomicBoolean(false);
|
||||
java.util.function.LongSupplier gatedClock = () -> {
|
||||
if (delayArmed.compareAndSet(true, false)) {
|
||||
try {
|
||||
Thread.sleep(STAMP_DELAY_MILLIS);
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
}
|
||||
}
|
||||
return clock.get();
|
||||
};
|
||||
MessageService service = new MessageService(agents, injector, rendezvous, inbox, null, null, gatedClock);
|
||||
try {
|
||||
String ticket = service.sendAsync(T, "a quick task");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
delayArmed.set(true);
|
||||
assertTrue(rendezvous.resolve(T, "quick result"));
|
||||
|
||||
awaitTicketPhaseOn(service, ticket, MessageService.Phase.DONE);
|
||||
// The barrier under test (fleetd #409, same fix as fleetd #399): remove this one call to
|
||||
// reproduce the pre-#399 ordering — see the class-level note above for what happens then.
|
||||
awaitCompletionStamped(service, ticket);
|
||||
|
||||
// A real clock would need the whole TTL to pass; the injected one does it instantly, and
|
||||
// by now completedNanos already holds the REAL (small, unadvanced) completion time.
|
||||
clock.addAndGet(MessageService.TICKET_TTL_NANOS + TimeUnit.SECONDS.toNanos(1));
|
||||
service.sendAsync(T, "an unrelated second task"); // runs pruneTerminalTickets before returning
|
||||
|
||||
assertNull(service.poll(ticket),
|
||||
"a ticket whose real completion time is long past the advanced cutoff must be "
|
||||
+ "evicted, whatever the completion hook's clock read was delayed by");
|
||||
} finally {
|
||||
service.close();
|
||||
}
|
||||
}
|
||||
|
||||
/** Bounded delay the injected clock sleeps for in {@link #aTicketOrderedOnDoneInsteadOfTheCompletionStampSurvivesAnEvictionItMustNotSurvive}. */
|
||||
private static final long STAMP_DELAY_MILLIS = 300;
|
||||
|
||||
private MessageService.TaskView awaitTicketPhaseOn(MessageService svc, String ticket,
|
||||
MessageService.Phase phase) throws Exception {
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
@@ -1990,6 +2079,45 @@ class MessageServiceTest {
|
||||
return view;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #399: waits until {@code ticket}'s completion hook has actually stamped
|
||||
* {@code completedNanos}, not just until {@link MessageService#poll} reports
|
||||
* {@link MessageService.Phase#DONE} for it. {@code poll} can observe {@code DONE} the instant
|
||||
* the task's future resolves, before the {@code whenComplete} hook that stamps the completion
|
||||
* time has run — {@code CompletableFuture.complete()} publishes its result and only then runs
|
||||
* dependents. A test that is about to advance an injected clock past the TTL must order itself
|
||||
* after the stamp, not after {@code DONE}: winning the race the other way stamps the
|
||||
* *advanced* clock value and can hide a real eviction bug behind a false pass.
|
||||
*/
|
||||
private void awaitCompletionStamped(MessageService svc, String ticket) throws Exception {
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
while (!svc.isCompletionStampedForTest(ticket)) {
|
||||
assertTrue(System.currentTimeMillis() < deadline,
|
||||
"completedNanos for " + ticket + " was never stamped");
|
||||
Thread.sleep(5);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #418: waits until {@code turnId} actually appears in {@code pushLoop}'s own open-question
|
||||
* state for {@code lead} — {@link ReplyPushLoop#pendingQuestionTurnIdsForTest} — rather than until
|
||||
* {@link MessageService#poll} reports {@link MessageService.Phase#ASKING}. {@code Phase.ASKING} is
|
||||
* set by {@code markAsyncQuestion}, the FIRST of three steps {@code MessageService.ask()} performs;
|
||||
* {@code pushLoop.onQuestionOpened} (the one that actually publishes to {@code pendingQuestions})
|
||||
* is the THIRD. Under load the asker thread can be descheduled between those two steps, so a
|
||||
* barrier on the phase alone can release before the push loop has anything pending — a test that
|
||||
* then asserts on {@link ReplyPushLoop#decide} is asserting on state that has not been published
|
||||
* yet, not on the throwing path it is named for.
|
||||
*/
|
||||
private void awaitQuestionPendingOn(ReplyPushLoop pushLoop, String lead, String turnId) throws Exception {
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
while (!pushLoop.pendingQuestionTurnIdsForTest(lead).contains(turnId)) {
|
||||
assertTrue(System.currentTimeMillis() < deadline,
|
||||
"turnId " + turnId + " never appeared in the push loop's pending questions for " + lead);
|
||||
Thread.sleep(5);
|
||||
}
|
||||
}
|
||||
|
||||
// --- CB-640: fleet health evidence accessors --------------------------------------------
|
||||
|
||||
@Test
|
||||
|
||||
Reference in New Issue
Block a user