Compare commits
53 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 | |||
| 69e09b10fa | |||
| b9d09e044e | |||
| fd8650cda4 | |||
| 11050e24ed | |||
| 769f282408 | |||
| 6f71f40047 | |||
| fb36c5238f | |||
| cc9cdc938b | |||
| 5fede82468 | |||
| 1515025804 | |||
| 7754f53662 | |||
| a507f7b31b | |||
| 71c322f104 | |||
| fde2c15627 | |||
| 2830735644 | |||
| 4e3ac91a22 | |||
| 29d3f0b41f | |||
| cbc732444f | |||
| 6f828b8c38 | |||
| 145a8c8862 | |||
| 7057291739 | |||
| 2302b3bc11 | |||
| 3a004dc1b3 | |||
| a9a3c12232 | |||
| 94ec77a1bc | |||
| 49f285cfda | |||
| 5a12ae7930 | |||
| 127e6832a9 |
@@ -7,6 +7,14 @@
|
||||
> wiki ([Use Cases](https://git.ltms.dev/fleet/fleetd/wiki/7-Use-Cases) → *The portable
|
||||
> CLAUDE.md block*); improvements go to the template first, then out to each project. Anything
|
||||
> specific to *this* repo lives under §Project addendum below, never inline above it.
|
||||
>
|
||||
> **Anything you measure in an addendum is perishable.** Date it, give the command that
|
||||
> re-measures it and what each outcome means, and tell the reader to delete the section once
|
||||
> it stops reproducing. The four parts work together: deciding what would falsify a claim is
|
||||
> the expensive step, and a reader in the middle of another task will not pay it, so a bare
|
||||
> "verify before relying on this" costs the same space and does nothing. The case this is for
|
||||
> is a note that goes stale as a live restriction — it will tell a future session it cannot do
|
||||
> the thing at the moment doing it becomes the job.
|
||||
|
||||
If no `fleet_*` MCP tools are mounted in this session, this section does not apply — skip it.
|
||||
|
||||
@@ -71,8 +79,11 @@ below are the procedure — run them in order, every task, not only the big ones
|
||||
the final judgment call, verification, merges, and anything that depends on context only you
|
||||
hold. Nothing else is yours by default.
|
||||
3. **Spawn every delegated unit first** — `fleet_spawn{profile, worktree:true, ticket}`, one per
|
||||
unit, *before* sending any. Pass `profile` explicitly: profiles differ in model and cost, not in
|
||||
tier, so the default is rarely what you want.
|
||||
unit, *before* sending any. Pass `profile` explicitly: profiles differ in model, cost and
|
||||
LIVENESS, not in tier, so the default is rarely what you want. The default is whatever the
|
||||
daemon reports, and on a host where it sits on an exhausted or withdrawn credential every
|
||||
unqualified spawn fails — sometimes loudly, sometimes as a member that spawns fine and then
|
||||
produces nothing. `fleet_profiles` reports the default; check it once per session.
|
||||
4. **Then send them all** — `fleet_send{sessionId, content, wait:false}`. Line 1 of every brief is
|
||||
`Load the <name> skill.` naming the worker's playbook; those skills are opt-in and that line is
|
||||
what makes them reliable. Where the project ships no such skill, spell the procedure out in the
|
||||
@@ -100,6 +111,12 @@ below are the procedure — run them in order, every task, not only the big ones
|
||||
without having read the diff yourself. A refusal is exactly when that shortcut is tempting,
|
||||
because no action is left that forces you to look, and taking it turns this step into
|
||||
forwarding a reviewer's verdict — which is delegating the merge by proxy, two lines above.
|
||||
**Test a refusal; do not read it off a permissions field.** A protected branch holds its merge
|
||||
rights separately from the repository permissions, so that field can say yes while the merge is
|
||||
refused, and still say no after a grant makes it work. Probe instead, with a request that cannot
|
||||
succeed on its merits, so a rejection can only mean the refusal. Treat a transport failure as a
|
||||
third answer that proves nothing: a timeout, a DNS error or a bad URL is not a refusal, and
|
||||
counting it as one makes you sure of something you never measured.
|
||||
|
||||
**Steps 3 and 4 are separate on purpose** — spawning and sending in one loop is how parallel work
|
||||
silently becomes serial, and it is the most common way this layer is wasted. For the same reason,
|
||||
@@ -120,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}` |
|
||||
|
||||
@@ -146,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
|
||||
|
||||
|
||||
@@ -192,7 +192,7 @@ Deliberately small; every one maps to a failure mode we have actually hit.
|
||||
| `fleet_send_duration_seconds` | histogram | delegated turn latency |
|
||||
| `fleet_replies_total{path}` | counter | path ∈ rendezvous\|inbox — how often a reply strands (CB-307's whole reason to exist) |
|
||||
| `fleet_inbox_depth{target}` | gauge | undrained replies; steady-state should be 0 |
|
||||
| `fleet_push_nudges_total{outcome}` | counter | outcome ∈ delivered\|exhausted — a rising `exhausted` means the primary is not draining |
|
||||
| `fleet_push_nudges_total{outcome}` | counter | outcome ∈ sent\|exhausted — a rising `exhausted` means the primary is not draining. `sent` was called `delivered` until fleetd #365; it counts the herdr paste-and-submit call returning, never a confirmation the pane read it |
|
||||
| `fleet_spawns_total{kind,outcome}` | counter | outcome ∈ ready\|timeout\|guard_rejected; per peer kind (CB-402) |
|
||||
| `fleet_sessions{state}` | gauge | SPAWNING/READY/BUSY/DONE census |
|
||||
| `fleet_herdr_calls_total{method,outcome}` | counter | socket health — the dependency everything rests on |
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -123,12 +123,21 @@ public final class Fleetd {
|
||||
// else can fail on a silently-empty one. A daemon started without a login shell (launchd)
|
||||
// boots fine either way — this is the only thing that says so out loud.
|
||||
reportRequiredSecrets(cfg);
|
||||
// fleetd #377: the git host value a member receives as GITEA_HOST is often a full URL
|
||||
// (scheme and trailing slash), not a host name. A member that assumes a bare host then
|
||||
// builds https://https://... and the request never leaves the machine. Report the shape
|
||||
// next to the secret report — shape only, never the value.
|
||||
reportGitHostShape(cfg);
|
||||
reportMemberTrustModel(cfg);
|
||||
// CB-596: an absent (or empty) memberCredentials: block blocks NOTHING — no credential
|
||||
// name is hardcoded any more to fall back on. Say so loudly, the same way a missing
|
||||
// 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
|
||||
@@ -139,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())
|
||||
@@ -579,8 +585,12 @@ public final class Fleetd {
|
||||
if (cfg.health() != null && cfg.health().isEnabled()) {
|
||||
// CB-580: a member found GONE/NEVER_READY must fail whatever ticket is waiting on it,
|
||||
// through the same idempotent target-wide operation CB-516 already uses on release.
|
||||
// fleetd #386: System::nanoTime freezes across a macOS sleep, so the stall check also
|
||||
// gets a wall-clock source to detect and correct for that freeze. Every other decision
|
||||
// in FleetHealthMonitor stays on the monotonic clock, unchanged.
|
||||
healthMonitor = new FleetHealthMonitor(agents, sessions::roster, messages, healthScheduler,
|
||||
System::nanoTime, cfg.health().intervalOrDefault(),
|
||||
System::nanoTime, () -> TimeUnit.MILLISECONDS.toNanos(System.currentTimeMillis()),
|
||||
cfg.health().intervalOrDefault(),
|
||||
cfg.health().workingSuspectAfterOrDefault(), messages::abandon);
|
||||
String coverage = FleetHealthMonitor.coverage(true,
|
||||
cfg.health().notifications() != null && cfg.health().notifications().configured());
|
||||
@@ -652,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();
|
||||
@@ -789,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).
|
||||
@@ -1195,6 +1216,76 @@ public final class Fleetd {
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #377: the env var names holding the git host value that members receive as
|
||||
* {@code GITEA_HOST}. {@code HerdrPeerLauncher.applyGitToken} injects {@code GITEA_HOST}
|
||||
* only for profiles that opted in via {@code gitTokenEnv} (CB-302), reading the value from
|
||||
* that profile's {@code gitHostEnv} (default {@code GITEA_HOST}), so the shape only matters
|
||||
* where a git token is opted in. A var used by more than one profile is one entry naming
|
||||
* every profile that reads it, the same shape as {@link #requiredSecretEnvVars}.
|
||||
*
|
||||
* <p>Package-private and pure (no I/O, no logging) so the derivation is unit-testable
|
||||
* without capturing log output; {@link #reportGitHostShape(FleetConfig)} is the logging caller.
|
||||
*/
|
||||
static Map<String, List<String>> gitHostEnvVars(FleetConfig cfg) {
|
||||
Map<String, List<String>> hostsBy = new LinkedHashMap<>();
|
||||
cfg.profiles().forEach((name, profile) -> {
|
||||
if (profile.hasGitToken()) {
|
||||
hostsBy.computeIfAbsent(profile.gitHostEnv(), _ -> new ArrayList<>())
|
||||
.add("profile '" + name + "' gitHostEnv");
|
||||
}
|
||||
});
|
||||
return hostsBy;
|
||||
}
|
||||
|
||||
/**
|
||||
* True when the value already starts with a URI scheme ({@code https://...}, {@code
|
||||
* http://...}). A bare host name and a host:port must both report {@code false} — the shape
|
||||
* this ticket exists for is a value that <em>looks like</em> a host but is a full URL, and
|
||||
* confusing those in the report would move the failure to the log instead of the network.
|
||||
*/
|
||||
static boolean startsWithScheme(String value) {
|
||||
return value.matches("[A-Za-z][A-Za-z0-9+.-]*://.*");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #377: log, on the same startup path as {@link #reportRequiredSecrets}, the SHAPE of
|
||||
* each git host value that members receive as {@code GITEA_HOST}: set or unset, its length,
|
||||
* whether it starts with a scheme, whether it ends with a slash. Never the value itself — the
|
||||
* same discipline as {@link #reportRequiredSecrets}, which logs by name only. Either shape is
|
||||
* legitimate: the value is passed to members unchanged, and a line that quietly rewrites it
|
||||
* would change what works on one host and breaks on another. The shape line only tells the
|
||||
* operator which URL form to expect from a member that builds on {@code GITEA_HOST}. An
|
||||
* unset variable is logged at INFO — useful information, not an error — and the daemon
|
||||
* starts on either way.
|
||||
*
|
||||
* <p>The env read lives in the overload below so a test can drive the line with a known
|
||||
* value and prove that value never reaches the log.
|
||||
*/
|
||||
static void reportGitHostShape(FleetConfig cfg) {
|
||||
reportGitHostShape(cfg, System.getenv());
|
||||
}
|
||||
|
||||
static void reportGitHostShape(FleetConfig cfg, Map<String, String> env) {
|
||||
Map<String, List<String>> hostsBy = gitHostEnvVars(cfg);
|
||||
if (hostsBy.isEmpty()) {
|
||||
log.info("startup git host: no profile sets a gitTokenEnv — nothing to check");
|
||||
return;
|
||||
}
|
||||
hostsBy.forEach((varName, sources) -> {
|
||||
String value = env.get(varName);
|
||||
if (value == null || value.isBlank()) {
|
||||
log.info("startup git host {}: unset ({}) — a member gets GITEA_TOKEN but no "
|
||||
+ "GITEA_HOST value", varName, String.join(", ", sources));
|
||||
} else {
|
||||
log.info("startup git host {}: set ({}) — length={}, startsWithScheme={}, "
|
||||
+ "trailingSlash={}",
|
||||
varName, String.join(", ", sources),
|
||||
value.length(), startsWithScheme(value), value.endsWith("/"));
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #184: state the member trust model at startup. Environment controls and worktrees do
|
||||
* not make a sandbox when fleetd and its members use the same OS user. A separate herdr may
|
||||
@@ -1244,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()) {
|
||||
|
||||
@@ -38,15 +38,46 @@ public final class FleetHealthMonitor {
|
||||
*/
|
||||
static final long ASK_LAPSE_RECHECK_DELAY_SECONDS = 120;
|
||||
|
||||
/**
|
||||
* fleetd #386: {@code System.nanoTime()} (or whatever {@link #clock} is) does not advance while
|
||||
* macOS sleeps, so a raw {@code nowNanos - lastActivityAtNanos} comparison freezes with the
|
||||
* host and can never cross {@link #workingSuspectAfterNanos}. This is a second, wall-clock
|
||||
* source used ONLY inside the stall check ({@link #stallElapsedNanos}) to detect and correct
|
||||
* for that freeze. Nothing else in this class reads it — every other decision (readiness grace,
|
||||
* the fault classification itself) stays exactly on {@link #clock}, as the ticket requires.
|
||||
*/
|
||||
private static final LongSupplier DEFAULT_REALTIME_CLOCK =
|
||||
() -> TimeUnit.MILLISECONDS.toNanos(System.currentTimeMillis());
|
||||
|
||||
private final AgentControl agents;
|
||||
private final Supplier<List<MemberSession>> roster;
|
||||
private final MessageService messages;
|
||||
private final ScheduledExecutorService scheduler;
|
||||
private final LongSupplier clock;
|
||||
private final LongSupplier realtimeClock;
|
||||
private final long intervalSeconds;
|
||||
private final long tickIntervalNanos;
|
||||
private final long workingSuspectAfterNanos;
|
||||
private final BiConsumer<String, String> failTarget;
|
||||
private final Map<String, HealthPrior> priors = new HashMap<>();
|
||||
/**
|
||||
* fleetd #386 clock-drift bookkeeping. {@code haveClockBaseline}/{@code lastTickMonoNanos}/
|
||||
* {@code lastTickRealNanos} track the previous tick's pair of readings so each new tick can
|
||||
* measure how far the two clocks moved apart since then. {@code accumulatedDriftNanos} is the
|
||||
* running total of every such divergence observed since this monitor started (never decreases —
|
||||
* the monotonic clock can only lag real time, never lead it). {@code busyDriftBaselineNanos}/
|
||||
* {@code busyBaselineActivityNanos} record, per target, the value of {@code accumulatedDriftNanos}
|
||||
* at the moment this monitor first saw that target's CURRENT {@code lastActivityAtNanos} while
|
||||
* BUSY — so {@link #stallElapsedNanos} adds back only the drift observed DURING this BUSY span,
|
||||
* never drift from a sleep that happened before the member went busy. All five fields are touched
|
||||
* only from {@code tick()}, like {@link #priors}.
|
||||
*/
|
||||
private boolean haveClockBaseline = false;
|
||||
private long lastTickMonoNanos;
|
||||
private long lastTickRealNanos;
|
||||
private long accumulatedDriftNanos = 0;
|
||||
private final Map<String, Long> busyDriftBaselineNanos = new HashMap<>();
|
||||
private final Map<String, Long> busyBaselineActivityNanos = new HashMap<>();
|
||||
/**
|
||||
* The live classification per member, and the only one of this class's three maps that more
|
||||
* than one scheduler task touches. {@code tick} writes it (and prunes it to the roster);
|
||||
@@ -88,12 +119,29 @@ public final class FleetHealthMonitor {
|
||||
public FleetHealthMonitor(AgentControl agents, Supplier<List<MemberSession>> roster, MessageService messages,
|
||||
ScheduledExecutorService scheduler, LongSupplier clock, long intervalSeconds,
|
||||
long workingSuspectAfterSeconds, BiConsumer<String, String> failTarget) {
|
||||
this(agents, roster, messages, scheduler, clock, DEFAULT_REALTIME_CLOCK, intervalSeconds,
|
||||
workingSuspectAfterSeconds, failTarget);
|
||||
}
|
||||
|
||||
/**
|
||||
* @param realtimeClock fleetd #386: a wall-clock nanosecond source (e.g.
|
||||
* {@code System.currentTimeMillis()} converted to nanos) that keeps
|
||||
* advancing while {@code clock} is frozen by a host sleep. Used only to
|
||||
* correct the stall check — see the class-level javadoc on the
|
||||
* clock-drift fields.
|
||||
*/
|
||||
public FleetHealthMonitor(AgentControl agents, Supplier<List<MemberSession>> roster, MessageService messages,
|
||||
ScheduledExecutorService scheduler, LongSupplier clock, LongSupplier realtimeClock,
|
||||
long intervalSeconds, long workingSuspectAfterSeconds,
|
||||
BiConsumer<String, String> failTarget) {
|
||||
this.agents = agents;
|
||||
this.roster = roster;
|
||||
this.messages = messages;
|
||||
this.scheduler = scheduler;
|
||||
this.clock = clock;
|
||||
this.realtimeClock = Objects.requireNonNull(realtimeClock, "realtimeClock");
|
||||
this.intervalSeconds = intervalSeconds;
|
||||
this.tickIntervalNanos = TimeUnit.SECONDS.toNanos(intervalSeconds);
|
||||
this.workingSuspectAfterNanos = TimeUnit.SECONDS.toNanos(workingSuspectAfterSeconds);
|
||||
this.failTarget = Objects.requireNonNull(failTarget, "failTarget");
|
||||
}
|
||||
@@ -123,6 +171,7 @@ public final class FleetHealthMonitor {
|
||||
for (Agent agent : agentsNow) live.put(agent.terminalId(), agent);
|
||||
HashSet<String> current = new HashSet<>();
|
||||
long nowNanos = clock.getAsLong();
|
||||
long driftBeforeThisTick = observeClockDrift(nowNanos);
|
||||
for (MemberSession session : rosterNow) {
|
||||
current.add(session.terminalId());
|
||||
Agent agent = live.get(session.terminalId());
|
||||
@@ -133,7 +182,7 @@ public final class FleetHealthMonitor {
|
||||
&& session.state() != MemberSession.State.SPAWNING;
|
||||
boolean readinessGraceElapsed = nowNanos - session.spawnedAtNanos() >= READINESS_GRACE_NANOS;
|
||||
boolean stalled = session.state() == MemberSession.State.BUSY
|
||||
&& nowNanos - session.lastActivityAtNanos() >= workingSuspectAfterNanos;
|
||||
&& stallElapsedNanos(session, nowNanos, driftBeforeThisTick) >= workingSuspectAfterNanos;
|
||||
// CB-643: the three message-layer facts CB-640 published. Read them here rather than
|
||||
// leaving them false — that constant is what made 8 of the 9 fault states dead.
|
||||
boolean queuedDelivery = messages.hasQueuedDelivery(session.terminalId());
|
||||
@@ -151,6 +200,8 @@ public final class FleetHealthMonitor {
|
||||
priors.keySet().retainAll(current);
|
||||
states.keySet().retainAll(current);
|
||||
orphanStreaks.keySet().retainAll(current);
|
||||
busyDriftBaselineNanos.keySet().retainAll(current);
|
||||
busyBaselineActivityNanos.keySet().retainAll(current);
|
||||
} catch (Throwable error) {
|
||||
// Any unclassified collection failure must never kill the monitor's only scheduler task.
|
||||
log.warn("fleet health collection failed; will retry next tick", error);
|
||||
@@ -175,6 +226,65 @@ public final class FleetHealthMonitor {
|
||||
return streak >= ORPHAN_CONFIRM_TICKS;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #386: compare this tick's monotonic and real-time readings against the previous
|
||||
* tick's, and fold any positive divergence into {@link #accumulatedDriftNanos} (a ratchet — it
|
||||
* never decreases, since the monotonic clock can only fall behind real time, never ahead of
|
||||
* it). Logs once, at WARN, when that single tick's divergence exceeds one full tick interval —
|
||||
* the signature of a host that slept between the two ticks (a tick literally cannot run while
|
||||
* the process itself is suspended, so the whole sleep duration lands inside one tick's gap).
|
||||
*
|
||||
* @return {@link #accumulatedDriftNanos} as it stood BEFORE this tick's divergence was folded
|
||||
* in — the baseline {@link #stallElapsedNanos} needs when a target is observed BUSY
|
||||
* for the first time this tick, so a sleep that happened before this member went busy
|
||||
* is not attributed to it.
|
||||
*/
|
||||
private long observeClockDrift(long nowNanos) {
|
||||
long nowRealNanos = realtimeClock.getAsLong();
|
||||
long driftBeforeThisTick = accumulatedDriftNanos;
|
||||
if (haveClockBaseline) {
|
||||
long monoDelta = nowNanos - lastTickMonoNanos;
|
||||
long realDelta = nowRealNanos - lastTickRealNanos;
|
||||
long tickDrift = realDelta - monoDelta;
|
||||
if (tickDrift > tickIntervalNanos) {
|
||||
log.warn("fleet health: the monotonic clock did not advance for about {}s that the "
|
||||
+ "real clock did since the last tick (host likely slept); the stall "
|
||||
+ "detector could not see that time", TimeUnit.NANOSECONDS.toSeconds(tickDrift));
|
||||
}
|
||||
if (tickDrift > 0) {
|
||||
accumulatedDriftNanos = driftBeforeThisTick + tickDrift;
|
||||
}
|
||||
}
|
||||
lastTickMonoNanos = nowNanos;
|
||||
lastTickRealNanos = nowRealNanos;
|
||||
haveClockBaseline = true;
|
||||
return driftBeforeThisTick;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #386: {@code nowNanos - lastActivityAtNanos} alone freezes across a host sleep, since
|
||||
* both come from the monotonic {@link #clock}. This adds back the real-time drift observed
|
||||
* since this BUSY span started — not the monitor's whole lifetime, so a sleep that happened
|
||||
* before this member went busy never leaks into its stall reading (see the class-level javadoc
|
||||
* on the drift fields). The baseline resets whenever {@code lastActivityAtNanos} changes (a new
|
||||
* turn) or the member is not currently BUSY.
|
||||
*/
|
||||
private long stallElapsedNanos(MemberSession session, long nowNanos, long driftBeforeThisTick) {
|
||||
String target = session.terminalId();
|
||||
if (session.state() != MemberSession.State.BUSY) {
|
||||
busyDriftBaselineNanos.remove(target);
|
||||
busyBaselineActivityNanos.remove(target);
|
||||
return nowNanos - session.lastActivityAtNanos();
|
||||
}
|
||||
Long baselineActivity = busyBaselineActivityNanos.get(target);
|
||||
if (baselineActivity == null || baselineActivity != session.lastActivityAtNanos()) {
|
||||
busyBaselineActivityNanos.put(target, session.lastActivityAtNanos());
|
||||
busyDriftBaselineNanos.put(target, driftBeforeThisTick);
|
||||
}
|
||||
long driftSinceBusyStart = accumulatedDriftNanos - busyDriftBaselineNanos.get(target);
|
||||
return (nowNanos - session.lastActivityAtNanos()) + driftSinceBusyStart;
|
||||
}
|
||||
|
||||
void reportTransition(String target, HealthState next) {
|
||||
HealthState previous = states.put(target, next);
|
||||
if (previous == next) return;
|
||||
|
||||
@@ -548,8 +548,19 @@ public final class CompletionResolver implements TurnListener {
|
||||
* fast completion. Runs the same backend-error classification the normal and raw-scrape paths
|
||||
* apply, against whatever is on screen right now: a match together with the too-fast crash
|
||||
* signature notifies {@link #backendErrorSink} (only on the resolution that wins the race). A
|
||||
* non-match stays the original generic too-fast failure, naming the member and both timings,
|
||||
* with whatever the pane shows appended so the caller sees the cause, not just "it failed".
|
||||
* non-match stays the generic too-fast failure, naming the member and both timings, with
|
||||
* whatever the pane shows appended so the caller sees the evidence, not just "it failed".
|
||||
*
|
||||
* <p>fleetd#376: <strong>this path must never resolve a completion.</strong> A fast backend can
|
||||
* genuinely answer inside the floor, so the failure is sometimes wrong — but it is wrong in the
|
||||
* loud direction, and the fix for that is honest wording, not a guess at the pane's meaning.
|
||||
* Reclassifying from the scrape was tried and rejected: there is no reliable positive marker for
|
||||
* "this is a real reply" across backends. {@link #lastAssistantBlock} falls back to the entire
|
||||
* pane when it finds no {@code ⏺} marker, so on a crash the candidate "reply" is the whole
|
||||
* screen; and {@code ⏺} itself is a Claude Code marker that an opencode pane never carries — the
|
||||
* very backend whose speed raised this ticket. Any weaker test (non-blank, or "contains sentence
|
||||
* punctuation") passes on almost every crash, because a pane holding a file path, a version
|
||||
* number or a hostname contains a full stop. That trades a loud wrong answer for a silent one.
|
||||
*/
|
||||
private void failTooFast(String target, InFlight turn, CompletableFuture<Rendezvous.Resolution> waiter,
|
||||
long elapsedNanos) {
|
||||
@@ -560,13 +571,13 @@ public final class CompletionResolver implements TurnListener {
|
||||
scrape = "";
|
||||
}
|
||||
String clippedScrape = clip(scrape);
|
||||
String baseReason = String.format(
|
||||
"member %s went BUSY -> DONE in %dms (floor %dms) — too fast to be real work, most "
|
||||
+ "likely a backend error before any work started",
|
||||
String timing = String.format(
|
||||
"member %s went BUSY -> DONE in %dms (floor %dms)",
|
||||
target, elapsedNanos / 1_000_000, MIN_TURN_NANOS / 1_000_000);
|
||||
String backendError = firstMatchingLine(scrape, backendErrorPatternOrFallback(target));
|
||||
if (backendError != null) {
|
||||
String reason = baseReason + ": " + clippedScrape;
|
||||
String reason = timing + " — too fast to be real work, and the pane carries a backend "
|
||||
+ "error: " + clippedScrape;
|
||||
if (rendezvous.resolveFailure(waiter, reason)) {
|
||||
inFlight.remove(target, turn);
|
||||
log.warn("failing send to {} via turn-stall fallback: {}", target, reason);
|
||||
@@ -575,7 +586,14 @@ public final class CompletionResolver implements TurnListener {
|
||||
}
|
||||
return;
|
||||
}
|
||||
fail(target, turn, clippedScrape.isBlank() ? baseReason : baseReason + ": " + clippedScrape);
|
||||
// fleetd#376: no pattern matched, so the cause is genuinely unknown. Say that, rather than
|
||||
// asserting a backend error the way this message used to — a fast backend really can finish
|
||||
// inside the floor, and a reader who trusts a wrong cause stops looking at the pane.
|
||||
String reason = timing + " — inside the floor. That is usually a backend error before any "
|
||||
+ "work started, but a fast backend can answer inside it too, and nothing here can "
|
||||
+ "tell those apart, so the turn is reported failed rather than guessed. Read the "
|
||||
+ "pane below before deciding which it was";
|
||||
fail(target, turn, clippedScrape.isBlank() ? reason : reason + ": " + clippedScrape);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -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,14 +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
|
||||
@@ -884,8 +915,8 @@ public final class FleetMcp {
|
||||
if (isBlank(content)) {
|
||||
return error("content is required");
|
||||
}
|
||||
messages.reply(callerTerminal, content);
|
||||
return text("delivered");
|
||||
MessageService.ReplyOutcome outcome = messages.reply(callerTerminal, content);
|
||||
return text(outcome.description());
|
||||
}
|
||||
|
||||
/** {@code fleet_ack}: acknowledge (remove) a specific reply from the inbox. */
|
||||
@@ -1104,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<>();
|
||||
@@ -1111,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 -> {
|
||||
@@ -1131,6 +1172,7 @@ public final class FleetMcp {
|
||||
});
|
||||
}
|
||||
}
|
||||
result.put("exhaustionDetectionArmed", exhaustionDetectionArmed);
|
||||
if (!quarantined.isEmpty()) {
|
||||
result.put("quarantined", quarantined);
|
||||
}
|
||||
@@ -1658,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()));
|
||||
}
|
||||
|
||||
|
||||
@@ -369,8 +369,7 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
// CB-547a: route the chosen profile but keep the caller's session identity — dropping it
|
||||
// here would silently sever the resume handle on every policy-routed spawn. CB-557: the
|
||||
// role rides along for the same reason, or a routed spawn would be labelled as a dev.
|
||||
SpawnRequest routedReq = new SpawnRequest(chosen.profile(), req.requestedCwd(), req.callerCwd(),
|
||||
req.sessionName(), req.resumeSessionId(), req.role());
|
||||
SpawnRequest routedReq = req.withProfile(chosen.profile());
|
||||
try {
|
||||
PeerHandle handle = d.spawn(routedReq);
|
||||
spawnedBy.put(handle.id(), d);
|
||||
|
||||
@@ -13,6 +13,7 @@ import java.nio.file.attribute.PosixFilePermissions;
|
||||
import java.time.Duration;
|
||||
import java.time.Instant;
|
||||
import java.util.ArrayList;
|
||||
import java.util.HashSet;
|
||||
import java.util.List;
|
||||
import java.util.Set;
|
||||
import java.util.stream.Stream;
|
||||
@@ -28,33 +29,85 @@ import java.util.stream.Stream;
|
||||
* control lacked: herdr applies that overlay BEFORE the shell starts, so any sourced file can undo
|
||||
* it — and did.
|
||||
*
|
||||
* <p><b>Which file is last depends on the platform, so the scrub runs from two of them.</b> zsh
|
||||
* <p><b>zsh reads its four startup files under three different conditions, so no single file is
|
||||
* guaranteed to run — the scrub has to cover the gap between them, not just the platforms.</b> zsh
|
||||
* reads {@code .zshenv} always, {@code .zprofile} and {@code .zlogin} only for a LOGIN shell, and
|
||||
* {@code .zshrc} only for an INTERACTIVE one. herdr does not open the same kind of shell
|
||||
* everywhere — measured on herdr 0.8.0: macOS panes run {@code -zsh} (login, so {@code .zlogin}
|
||||
* runs), Linux panes run a plain {@code /usr/bin/zsh} (interactive but NOT login, so
|
||||
* {@code .zlogin} never runs at all). A scrub in {@code .zlogin} alone is therefore a control that
|
||||
* silently does nothing on Linux — the exact failure this class exists to remove, one platform
|
||||
* over.
|
||||
* {@code .zshrc} only for an INTERACTIVE one. A pane shell that is at least one of login or
|
||||
* interactive is covered by sourcing the scrub from {@code .zshrc} and {@code .zlogin} (below), but
|
||||
* a pane shell that is NEITHER reads only {@code .zshenv} and stops — fleetd #388, measured: a herdr
|
||||
* pane can be neither login nor interactive, and such a pane read {@code .zshenv}, never reached
|
||||
* {@code scrub.zsh}, and left no report at all. A bare {@code argv[0]} of {@code /usr/bin/zsh}
|
||||
* proves the shell is NOT a login shell; it says nothing about whether it is interactive, so it
|
||||
* must never be read as "therefore interactive" — that wrong inference is what let #388 ship.
|
||||
*
|
||||
* <p>So both {@code .zshrc} and {@code .zlogin} source the same generated {@code scrub.zsh} after
|
||||
* sourcing their {@code $HOME} counterpart. On Linux only the first fires; on macOS both do, and
|
||||
* the second pass is deliberate rather than merely harmless — it re-scrubs anything the operator's
|
||||
* own {@code ~/.zlogin} exported after {@code .zshrc} had finished. Re-running is idempotent: a
|
||||
* name already blank is blanked again, and the report is rewritten with the same counts.
|
||||
* <p>So {@code .zshenv} carries a THIRD pass, guarded by the exact condition that defines the gap:
|
||||
* {@code [[ ! -o login && ! -o interactive ]]}. That guard is why this pass cannot double-scrub a
|
||||
* pane that {@code .zshrc} or {@code .zlogin} will also cover — one of {@code -o login}/
|
||||
* {@code -o interactive} is always true there, so the {@code .zshenv} pass never fires for them, and
|
||||
* their own unconditional sourcing is untouched. The guard also carries a sentinel
|
||||
* ({@value #SCRUB_SENTINEL}) so it fires once per PANE and not once per PROCESS: {@code .zshenv} is
|
||||
* read by every zsh a member's own tooling forks (a plain {@code zsh -c '...'} for a single
|
||||
* command is itself neither login nor interactive), and those children inherit variables their
|
||||
* parent deliberately set for them (git hooks get {@code GIT_DIR}, a venv gets
|
||||
* {@code VIRTUAL_ENV}, a build tool gets {@code NODE_OPTIONS} or {@code JAVA_TOOL_OPTIONS}).
|
||||
* Re-scrubbing every such child would blank all of that, and would also make the pane's own
|
||||
* {@code scrub-report.txt} — rewritten on every pass — describe whichever child exited last
|
||||
* instead of the pane. The sentinel is exported only AFTER {@code scrub.zsh} runs, so the pass
|
||||
* that sets it never sees it and cannot blank it; it must also be on the scrub's own allow-list
|
||||
* (see {@link #generate(Path, Set)}) so a later pass, in the same pane, cannot blank it back to
|
||||
* empty — an exported-but-empty sentinel reads as unset to the {@code -z} guard and would silently
|
||||
* re-enable scrubbing for every subsequent child of that pane.
|
||||
*
|
||||
* <p>So all three of {@code .zshenv} (gap only, guarded), {@code .zshrc}, and {@code .zlogin}
|
||||
* source the same generated {@code scrub.zsh} after sourcing their {@code $HOME} counterpart. A
|
||||
* login-and-interactive pane runs the {@code .zshrc} and {@code .zlogin} passes, and the second is
|
||||
* deliberate rather than merely harmless — it re-scrubs anything the operator's own
|
||||
* {@code ~/.zlogin} exported after {@code .zshrc} had finished. A pane that is neither runs only the
|
||||
* {@code .zshenv} pass. Re-running is idempotent: a name already blank is blanked again, and the
|
||||
* report is rewritten with the same counts.
|
||||
*
|
||||
* <p>Each generated file sources its {@code $HOME} counterpart FIRST, so {@code PATH} and every
|
||||
* toolchain binary still resolve exactly as the operator configured them; only afterwards does
|
||||
* {@code .zlogin} run the scrub: every EXPORTED variable not on the derived allow-list is re-exported
|
||||
* toolchain binary still resolve exactly as the operator configured them; only afterwards does the
|
||||
* scrub run: every EXPORTED variable not on the derived allow-list is re-exported
|
||||
* blank. Blank, not credential-shaped-pattern-filtered: a pattern list ({@code *TOKEN*}, …) is an
|
||||
* enumeration and misses what it did not think of — a username is the other half of a credential and
|
||||
* is shaped like none. Credential-SHAPED names among the blanked set go to the WARN log only,
|
||||
* 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 {
|
||||
|
||||
@@ -63,7 +116,11 @@ public final class EnvAllowListScrub {
|
||||
/** Name of the report file written into the generated directory by the scrub itself. */
|
||||
static final String REPORT_FILE = "scrub-report.txt";
|
||||
|
||||
/** The scrub body, generated once and sourced from both {@code .zshrc} and {@code .zlogin}. */
|
||||
/**
|
||||
* The scrub body, generated once and sourced from {@code .zshrc} and {@code .zlogin}
|
||||
* unconditionally, and from {@code .zshenv} when the pane shell is neither login nor
|
||||
* interactive (fleetd #388) — see the class javadoc.
|
||||
*/
|
||||
static final String SCRUB_FILE = "scrub.zsh";
|
||||
|
||||
/** Prefix of every generated directory — also what {@link #reapOrphans} matches on. */
|
||||
@@ -79,14 +136,42 @@ public final class EnvAllowListScrub {
|
||||
private static final String SOURCE_SCRUB =
|
||||
"source \"$ZDOTDIR/" + SCRUB_FILE + "\"\n";
|
||||
|
||||
/**
|
||||
* fleetd #388: marks a pane, not a process, as already scrubbed. Set only by the guarded
|
||||
* {@code .zshenv} pass (see {@link #NEITHER_LOGIN_NOR_INTERACTIVE_SCRUB}) after
|
||||
* {@code scrub.zsh} has run, so it must also be folded into that pass's own allow-list — see
|
||||
* the class javadoc's "must also be on the scrub's own allow-list" paragraph.
|
||||
*/
|
||||
static final String SCRUB_SENTINEL = "_CB633_SCRUBBED";
|
||||
|
||||
/**
|
||||
* Appended to {@code .zshenv}, after its {@code $HOME} source: the third pass, guarded on the
|
||||
* exact condition that defines the gap {@code .zshrc}/{@code .zlogin} do not cover — a shell
|
||||
* that is neither login nor interactive. The sentinel export happens only once the scrub has
|
||||
* already run, and only for as long as the current pane's environment has not been rebuilt from
|
||||
* scratch (a fresh {@code env -i} child would not inherit it — that is out of scope here, since
|
||||
* such a child is no longer running under the pane's own environment at all).
|
||||
*/
|
||||
private static final String NEITHER_LOGIN_NOR_INTERACTIVE_SCRUB =
|
||||
"if [[ ! -o login && ! -o interactive && -z \"${" + SCRUB_SENTINEL + ":-}\" ]]; then\n"
|
||||
+ " " + SOURCE_SCRUB
|
||||
+ " export " + SCRUB_SENTINEL + "=1\n"
|
||||
+ "fi\n";
|
||||
|
||||
private 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) {
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -105,11 +190,16 @@ public final class EnvAllowListScrub {
|
||||
reapOrphans(parentDir);
|
||||
Path dir = Files.createTempDirectory(parentDir, DIR_PREFIX);
|
||||
dir.toFile().deleteOnExit();
|
||||
// The report is written by zsh, after these hooks are registered, so register its path
|
||||
// too — otherwise the directory is non-empty at JVM exit and cannot be removed at all.
|
||||
dir.resolve(REPORT_FILE).toFile().deleteOnExit();
|
||||
write(dir, SCRUB_FILE, scrubScript(allowedNames));
|
||||
write(dir, ".zshenv", homeSourcingFile(".zshenv"));
|
||||
// zsh truncates this pre-created receipt after these hooks are registered. Register its
|
||||
// path too — otherwise the directory is non-empty at JVM exit and cannot be removed.
|
||||
Files.createFile(dir.resolve(REPORT_FILE)).toFile().deleteOnExit();
|
||||
// fleetd #388: scrub.zsh's OWN allow-list must also keep SCRUB_SENTINEL, or a later
|
||||
// pass in the same pane blanks it back to empty and the .zshenv guard below thinks it
|
||||
// was never scrubbed — see the class javadoc.
|
||||
Set<String> namesForScrubScript = new HashSet<>(allowedNames);
|
||||
namesForScrubScript.add(SCRUB_SENTINEL);
|
||||
write(dir, SCRUB_FILE, scrubScript(namesForScrubScript));
|
||||
write(dir, ".zshenv", homeSourcingFile(".zshenv") + NEITHER_LOGIN_NOR_INTERACTIVE_SCRUB);
|
||||
write(dir, ".zprofile", homeSourcingFile(".zprofile"));
|
||||
write(dir, ".zshrc", homeSourcingFile(".zshrc") + SOURCE_SCRUB);
|
||||
write(dir, ".zlogin", homeSourcingFile(".zlogin") + SOURCE_SCRUB);
|
||||
@@ -147,12 +237,11 @@ public final class EnvAllowListScrub {
|
||||
* into it: owner keeps full access, {@code group} gets traverse+read on the directory ({@code
|
||||
* rwxr-x---}, so a member process — a login shell reading it via {@code ZDOTDIR}, or another
|
||||
* process simply opening a file under it — running under that group can find and read the
|
||||
* files) and read-only on each file ({@code rw-r-----}) — deliberately no group WRITE anywhere,
|
||||
* since a member never needs to add or change what fleetd generated. (For the ZDOTDIR scrub
|
||||
* specifically, this also means the scrub script's own report write inside the pane fails
|
||||
* closed rather than open — see {@code scrub.zsh}'s trailing {@code 2>/dev/null} — which {@link
|
||||
* dev.ltms.fleet.member.HerdrPeerLauncher#releaseZdotdir} already treats as "cannot be
|
||||
* confirmed to have run" rather than success.)
|
||||
* files) and read-only on each file ({@code rw-r-----}), except the pre-created ZDOTDIR
|
||||
* {@code scrub-report.txt}. That receipt gets group write ({@code rw-rw----}), so
|
||||
* {@code scrub.zsh} can truncate and write it without granting group write on the directory.
|
||||
* If its optional permission change fails, the member cannot write a receipt and the launcher
|
||||
* keeps its existing WARN rather than failing the spawn.
|
||||
*
|
||||
* <p>Package-private and named generically on purpose: fleetd #213 built this for the ZDOTDIR
|
||||
* scrub directory, and fleetd #219 reuses it verbatim for {@link
|
||||
@@ -167,7 +256,15 @@ public final class EnvAllowListScrub {
|
||||
setGroupAndPermissions(dir, principal, "rwxr-x---");
|
||||
try (Stream<Path> entries = Files.list(dir)) {
|
||||
for (Path file : entries.toList()) {
|
||||
setGroupAndPermissions(file, principal, "rw-r-----");
|
||||
if (REPORT_FILE.equals(file.getFileName().toString())) {
|
||||
try {
|
||||
setGroupAndPermissions(file, principal, "rw-rw----");
|
||||
} catch (IOException | UnsupportedOperationException ignored) {
|
||||
// The receipt is optional. Its absence keeps the existing WARN path.
|
||||
}
|
||||
} else {
|
||||
setGroupAndPermissions(file, principal, "rw-r-----");
|
||||
}
|
||||
}
|
||||
}
|
||||
} catch (IOException e) {
|
||||
@@ -245,11 +342,13 @@ public final class EnvAllowListScrub {
|
||||
}
|
||||
return """
|
||||
# generated by fleetd (CB-633 memberCredentials policy=allow-list) — do not edit.
|
||||
# Sourced from .zshrc and again from .zlogin, each time AFTER that file has sourced
|
||||
# its $HOME counterpart — so this runs after everything the operator sourced, on a
|
||||
# login shell (macOS panes) and on a plain interactive one (Linux panes) alike.
|
||||
# Running twice is idempotent and deliberate: the second pass catches anything
|
||||
# ~/.zlogin exported after ~/.zshrc had finished.
|
||||
# Sourced unconditionally from .zshrc and again from .zlogin, each time AFTER that
|
||||
# file has sourced its $HOME counterpart — so this runs after everything the
|
||||
# operator sourced, on any pane that is login and/or interactive. Running twice is
|
||||
# idempotent and deliberate: the second pass catches anything ~/.zlogin exported
|
||||
# after ~/.zshrc had finished. Also sourced, once, from a guarded pass in .zshenv
|
||||
# (fleetd #388) when the pane shell is NEITHER login nor interactive — the one gap
|
||||
# those two files do not cover.
|
||||
|
||||
typeset -A _cb633_allowed
|
||||
for _cb633_n in %s; do _cb633_allowed[$_cb633_n]=1; done
|
||||
@@ -271,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);
|
||||
}
|
||||
|
||||
@@ -293,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);
|
||||
@@ -305,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();
|
||||
|
||||
@@ -56,10 +56,13 @@ public final class FleetMetrics {
|
||||
m.describe(REPLIES, "counter",
|
||||
"Worker replies by delivery path (rendezvous=resolved an open send, inbox=stranded and held).");
|
||||
m.describe(PUSH_NUDGES, "counter",
|
||||
"CB-307 push-loop nudges to the primary (delivered|exhausted).");
|
||||
"CB-307 push-loop nudges to the primary (sent|exhausted). fleetd #365: \"sent\" means "
|
||||
+ "the herdr paste-and-submit call succeeded, not that the pane read it — this "
|
||||
+ "layer has no read-receipt concept.");
|
||||
m.describe(HEARTBEAT_NUDGES, "counter",
|
||||
"CB-551 idle-lead heartbeat nudges (delivered|failed|exhausted). Quiet-cap exhaustion "
|
||||
+ "means the lead idled with nothing pending and was told to stand down.");
|
||||
"CB-551 idle-lead heartbeat nudges (sent|failed|exhausted). Quiet-cap exhaustion "
|
||||
+ "means the lead idled with nothing pending and was told to stand down. "
|
||||
+ "fleetd #365: \"sent\" means the herdr call succeeded, not that the lead read it.");
|
||||
m.describe(SPAWNS, "counter",
|
||||
"Worker spawn attempts by peer kind and outcome (ready|timeout|guard_rejected).");
|
||||
m.describe(HERDR_CALLS, "counter",
|
||||
|
||||
@@ -32,7 +32,11 @@ public interface LeadChannel {
|
||||
/** Non-destructive FIFO snapshot of the messages held for this daemon's own coord-id. */
|
||||
List<LeadMessage> peek();
|
||||
|
||||
/** Drop {@code msgId} from the held set and ack it on the broker. A no-op if it is not held. */
|
||||
/**
|
||||
* Drop {@code msgId} from the held set and ack it on the broker. A repeated ack that this
|
||||
* connection already completed may be a no-op. Any other unknown msgId must throw rather than
|
||||
* report an ack that did not reach the broker.
|
||||
*/
|
||||
void ack(String msgId);
|
||||
|
||||
/** This daemon's own lead coordination id — the mailbox it owns, and the {@code from} it sends as. */
|
||||
|
||||
@@ -5,6 +5,7 @@ import dev.ltms.fleet.herdr.AgentStatus;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
@@ -43,6 +44,9 @@ public final class LeadCoordLoop {
|
||||
|
||||
private static final Logger log = LoggerFactory.getLogger(LeadCoordLoop.class);
|
||||
|
||||
/** A bounded window is enough: redelivery can only follow a recent failed ack or recovery. */
|
||||
private static final int RECENT_DELIVERY_LIMIT = 1_024;
|
||||
|
||||
/** How an arriving peer message is rendered into the lead's pane — the sender's coord-id, then its text. */
|
||||
static final String DELIVERY_FORMAT = "[lead %s] %s";
|
||||
|
||||
@@ -51,6 +55,8 @@ public final class LeadCoordLoop {
|
||||
private final Supplier<Map<String, String>> leads;
|
||||
private final ScheduledExecutorService scheduler;
|
||||
private final long intervalMs;
|
||||
/** msgIds already written to the pane, so recovery redelivery is acked without another pane write. */
|
||||
private final LinkedHashMap<String, Boolean> delivered = new LinkedHashMap<>();
|
||||
|
||||
private volatile boolean running;
|
||||
|
||||
@@ -117,6 +123,13 @@ public final class LeadCoordLoop {
|
||||
if (held.isEmpty()) {
|
||||
return;
|
||||
}
|
||||
LeadMessage msg = held.getFirst();
|
||||
if (wasDelivered(msg.msgId())) {
|
||||
// This lives here, rather than in LeadMailbox, because only this loop knows a pane write
|
||||
// happened. The mailbox only knows broker delivery tags and must still redeliver after a crash.
|
||||
ackDelivered(msg);
|
||||
return;
|
||||
}
|
||||
String lead = resolveLocalLead();
|
||||
if (lead == null) {
|
||||
// Left unacked on purpose: the broker keeps holding it until a lead pane exists.
|
||||
@@ -136,7 +149,6 @@ public final class LeadCoordLoop {
|
||||
lead, status, held.size());
|
||||
return;
|
||||
}
|
||||
LeadMessage msg = held.getFirst();
|
||||
try {
|
||||
agents.send(lead, DELIVERY_FORMAT.formatted(msg.from(), msg.content()));
|
||||
} catch (RuntimeException e) {
|
||||
@@ -145,6 +157,13 @@ public final class LeadCoordLoop {
|
||||
msg.msgId(), msg.from(), lead, e.toString());
|
||||
return;
|
||||
}
|
||||
rememberDelivered(msg.msgId());
|
||||
if (ackDelivered(msg)) {
|
||||
log.debug("lead coordination: delivered message {} from {} to lead {}", msg.msgId(), msg.from(), lead);
|
||||
}
|
||||
}
|
||||
|
||||
private boolean ackDelivered(LeadMessage msg) {
|
||||
try {
|
||||
channel.ack(msg.msgId());
|
||||
} catch (RuntimeException e) {
|
||||
@@ -152,9 +171,24 @@ public final class LeadCoordLoop {
|
||||
// deliberate direction of this trade.
|
||||
log.warn("lead coordination: delivered message {} but could not ack it: {}",
|
||||
msg.msgId(), e.toString());
|
||||
return;
|
||||
return false;
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
private boolean wasDelivered(String msgId) {
|
||||
synchronized (delivered) {
|
||||
return delivered.containsKey(msgId);
|
||||
}
|
||||
}
|
||||
|
||||
private void rememberDelivered(String msgId) {
|
||||
synchronized (delivered) {
|
||||
delivered.put(msgId, Boolean.TRUE);
|
||||
if (delivered.size() > RECENT_DELIVERY_LIMIT) {
|
||||
delivered.remove(delivered.keySet().iterator().next());
|
||||
}
|
||||
}
|
||||
log.debug("lead coordination: delivered message {} from {} to lead {}", msg.msgId(), msg.from(), lead);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -96,7 +96,14 @@ public final class LeadHeartbeatLoop {
|
||||
this.metrics = metrics;
|
||||
}
|
||||
|
||||
/** Count one nudge outcome when a registry is wired; a no-op in unit tests. */
|
||||
/**
|
||||
* Count one nudge outcome when a registry is wired; a no-op in unit tests.
|
||||
*
|
||||
* <p>fleetd #365: the {@code "sent"} outcome (renamed from {@code "delivered"}) records only
|
||||
* that {@link #injectNudge} — a one-way herdr {@code agent.prompt} paste-and-submit — returned
|
||||
* without throwing, not that the lead's pane actually read or acted on the text. This layer has
|
||||
* no read-receipt concept, so "sent" is the honest word for what this call can ever establish.
|
||||
*/
|
||||
private void countNudge(String outcome) {
|
||||
if (metrics != null) {
|
||||
metrics.inc(FleetMetrics.HEARTBEAT_NUDGES, "outcome", outcome);
|
||||
@@ -245,7 +252,7 @@ public final class LeadHeartbeatLoop {
|
||||
agents.send(leadTerminal, fleet.nudgeText());
|
||||
log.debug("idle-heartbeat: nudge sent to lead {} (quiet nudges so far in this stretch: {})",
|
||||
leadTerminal, quietCount);
|
||||
countNudge("delivered");
|
||||
countNudge("sent");
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("idle-heartbeat: failed to nudge lead {}: {}", leadTerminal, e.toString());
|
||||
countNudge("failed");
|
||||
|
||||
@@ -59,7 +59,8 @@ import java.util.concurrent.TimeoutException;
|
||||
* <p><strong>Recovery.</strong> The connection is opened with automatic + topology recovery
|
||||
* enabled, mirroring {@code AmqpReplyInbox}: on reconnect the broker hands out fresh delivery tags,
|
||||
* so the held snapshot is cleared (dedup by {@code msgId} still prevents any double-queue on
|
||||
* redelivery) and any publish still awaiting its confirm is failed rather than left to idle out
|
||||
* redelivery). {@link LeadCoordLoop} separately deduplicates pane writes, since it alone knows
|
||||
* which messages reached a lead. Any publish still awaiting its confirm is failed rather than left to idle out
|
||||
* the confirm timeout against a sequence number that means nothing on the new channel.
|
||||
*/
|
||||
public final class LeadMailbox implements LeadChannel, AutoCloseable {
|
||||
@@ -85,6 +86,10 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable {
|
||||
private final Object channelLock = new Object();
|
||||
/** msgId → held delivery, for this mailbox's own queue only (there is exactly one). */
|
||||
private final LinkedHashMap<String, Held> held = new LinkedHashMap<>();
|
||||
/** Successful broker acks on this connection, retained only to make a repeated caller ack quiet. */
|
||||
private final LinkedHashMap<String, Boolean> recentlyAcked = new LinkedHashMap<>();
|
||||
/** Bounds {@link #recentlyAcked}: it is only an idempotency aid, never delivery state. */
|
||||
private static final int RECENT_ACK_LIMIT = 1_024;
|
||||
|
||||
/**
|
||||
* A dedicated channel for {@link #publish}, kept separate from {@link #channel} (consume + ack)
|
||||
@@ -162,16 +167,15 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable {
|
||||
}
|
||||
// On automatic recovery the broker redelivers unacked messages with FRESH delivery-tags; the
|
||||
// tags we were holding are now stale. Drop the held snapshot so the re-attached consumer
|
||||
// repopulates it with valid tags (dedup by msgId still prevents any double-queue). Any publish
|
||||
// repopulates it with valid tags (dedup by msgId still prevents any double-queue). LeadCoordLoop
|
||||
// remembers successful pane writes separately, so that redelivery cannot write a pane twice. Any publish
|
||||
// confirm still in flight when the connection dropped is equally stale — fail it now rather
|
||||
// than let it silently ride out CONFIRM_TIMEOUT_MS.
|
||||
if (connection instanceof Recoverable recoverable) {
|
||||
recoverable.addRecoveryListener(new RecoveryListener() {
|
||||
@Override
|
||||
public void handleRecovery(Recoverable recoverable) {
|
||||
synchronized (held) {
|
||||
held.clear();
|
||||
}
|
||||
clearHeldForRecovery();
|
||||
failPendingPublishesOnRecovery();
|
||||
log.info("AMQP lead mailbox connection recovered; cleared held messages for fresh redelivery");
|
||||
}
|
||||
@@ -350,15 +354,26 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable {
|
||||
return snapshot;
|
||||
}
|
||||
|
||||
/** Remove the held message {@code msgId} and ack it on the broker. No-op if not held. */
|
||||
/**
|
||||
* Remove the held message {@code msgId} and ack it on the broker.
|
||||
*
|
||||
* <p>A repeated ack that this connection already completed is a no-op, tracked in the bounded
|
||||
* {@link #recentlyAcked} set. Any other missing entry throws: recovery clears {@link #held} while
|
||||
* the broker still owns the unacked delivery, and quiet success there would hide a required retry.
|
||||
* The set is bounded because it only distinguishes a recent duplicate caller ack from an unknown
|
||||
* delivery; it is not a substitute for broker state across a reconnect.
|
||||
*/
|
||||
@Override
|
||||
public void ack(String msgId) {
|
||||
Held h;
|
||||
synchronized (held) {
|
||||
h = held.remove(msgId);
|
||||
if (h == null && recentlyAcked.containsKey(msgId)) {
|
||||
return;
|
||||
}
|
||||
}
|
||||
if (h == null) {
|
||||
return; // never held (or already acked) — no-op
|
||||
throw new IllegalStateException("cannot ack lead message " + msgId + ": it is not held");
|
||||
}
|
||||
try {
|
||||
synchronized (channelLock) {
|
||||
@@ -372,6 +387,19 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable {
|
||||
}
|
||||
throw new IllegalStateException("cannot ack lead message " + msgId, e);
|
||||
}
|
||||
synchronized (held) {
|
||||
recentlyAcked.put(msgId, Boolean.TRUE);
|
||||
if (recentlyAcked.size() > RECENT_ACK_LIMIT) {
|
||||
recentlyAcked.remove(recentlyAcked.keySet().iterator().next());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** Clear stale delivery tags after recovery; package-private so the recovery contract test drives this exact path. */
|
||||
void clearHeldForRecovery() {
|
||||
synchronized (held) {
|
||||
held.clear();
|
||||
}
|
||||
}
|
||||
|
||||
private DeliverCallback deliverCallback() {
|
||||
|
||||
@@ -138,6 +138,53 @@ public final class MessageService {
|
||||
public record AskResult(AskOutcome outcome, String answer) {
|
||||
}
|
||||
|
||||
/**
|
||||
* How a worker's {@code fleet_reply} ({@link #reply(String, String)}) actually landed
|
||||
* (fleetd #365) — the two doors that expose it, {@code fleet_reply} and {@code POST
|
||||
* /sessions/{id}/reply}, both used to report the single word "delivered" whichever of these
|
||||
* happened, so a caller could not tell an active handoff from a reply merely held for later
|
||||
* drain. Both are successes; they are not the same fact.
|
||||
*/
|
||||
public enum ReplyOutcome {
|
||||
/** Resolved a {@code fleet_send}/{@code fleet_ask} that was actively waiting on this reply. */
|
||||
RESOLVED_SEND("resolved_send", true,
|
||||
"delivered — resolved the fleet_send that was waiting for it"),
|
||||
/**
|
||||
* No live waiter was open, but the reply completed a parked async ticket directly
|
||||
* ({@link #askAnsweredAsyncTasks}) — a {@code fleet_poll} caller sees it immediately.
|
||||
*/
|
||||
RESOLVED_ASYNC_TICKET("resolved_async_ticket", true,
|
||||
"delivered — resolved a pending async ticket (visible to fleet_poll)"),
|
||||
/** Nothing was waiting; the reply was queued in the inbox for a later drain (CB-307). */
|
||||
QUEUED("queued", false,
|
||||
"queued — no send or ticket was waiting; held in the inbox for a later drain");
|
||||
|
||||
private final String wireName;
|
||||
private final boolean delivered;
|
||||
private final String description;
|
||||
|
||||
ReplyOutcome(String wireName, boolean delivered, String description) {
|
||||
this.wireName = wireName;
|
||||
this.delivered = delivered;
|
||||
this.description = description;
|
||||
}
|
||||
|
||||
/** Stable machine-readable name for a JSON/metrics label (REST's {@code outcome} field). */
|
||||
public String wireName() {
|
||||
return wireName;
|
||||
}
|
||||
|
||||
/** Whether something was actively waiting and received this reply right now. */
|
||||
public boolean delivered() {
|
||||
return delivered;
|
||||
}
|
||||
|
||||
/** Shared human-readable text — the one place both {@code fleet_reply} and REST word this. */
|
||||
public String description() {
|
||||
return description;
|
||||
}
|
||||
}
|
||||
|
||||
/** Lifecycle phase of an async delegation ticket. */
|
||||
public enum Phase {
|
||||
/** Delegated and in flight — queued for the worker or being worked. */
|
||||
@@ -439,16 +486,15 @@ public final class MessageService {
|
||||
* @throws IllegalArgumentException if {@code content} is {@code null} or blank — the caller must
|
||||
* report this as a client error (REST: 400 {@code bad_request}) rather than resolve
|
||||
* anything
|
||||
* @return always {@code true} — the reply resolved a live send, completed a parked ticket, or
|
||||
* was queued
|
||||
* @return which of the three ways (fleetd #365) the reply actually landed — never {@code null}
|
||||
*/
|
||||
public boolean reply(String session, String content) {
|
||||
public ReplyOutcome reply(String session, String content) {
|
||||
if (content == null || content.isBlank()) {
|
||||
throw new IllegalArgumentException("content is required");
|
||||
}
|
||||
if (rendezvous.resolve(session, content)) {
|
||||
count(FleetMetrics.REPLIES, "path", "rendezvous");
|
||||
return true; // a live send took it — unchanged fast path
|
||||
return ReplyOutcome.RESOLVED_SEND; // a live send took it — unchanged fast path
|
||||
}
|
||||
// #137/fleetd #307: no live rendezvous waiter, but this may be the worker's real fleet_reply resuming
|
||||
// a turn that either answer() (#137) or ask() (fleetd #307) already gave up waiting on:
|
||||
@@ -480,7 +526,7 @@ public final class MessageService {
|
||||
asyncTasksByTurn.remove(turnId, orphan);
|
||||
}
|
||||
count(FleetMetrics.REPLIES, "path", "async-recovered");
|
||||
return true; // the ticket itself took it — no inbox stranding at all
|
||||
return ReplyOutcome.RESOLVED_ASYNC_TICKET; // the ticket itself took it — no inbox stranding
|
||||
}
|
||||
} else if (candidates.size() > 1) {
|
||||
List<String> tickets = candidates.stream().map(t -> t.ticket).toList();
|
||||
@@ -498,7 +544,7 @@ public final class MessageService {
|
||||
if (pushLoop != null) {
|
||||
pushLoop.onReplyQueued(session);
|
||||
}
|
||||
return true; // held, not lost
|
||||
return ReplyOutcome.QUEUED; // held, not lost — but not delivered either
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1271,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 {
|
||||
|
||||
@@ -132,7 +132,15 @@ public final class ReplyPushLoop {
|
||||
this.metrics = metrics;
|
||||
}
|
||||
|
||||
/** Count one nudge outcome when a registry is wired; a no-op in unit tests. */
|
||||
/**
|
||||
* Count one nudge outcome when a registry is wired; a no-op in unit tests.
|
||||
*
|
||||
* <p>fleetd #365: the {@code "sent"} outcome (renamed from {@code "delivered"}) records only
|
||||
* that {@code agents.send} — a one-way herdr {@code agent.prompt} paste-and-submit — returned
|
||||
* without throwing. Nothing in this loop, or anywhere downstream of it, confirms the pane
|
||||
* actually read or acted on the text; there is no read-receipt concept at this layer. "Sent"
|
||||
* says exactly that; "delivered" claimed more than this call can ever establish.
|
||||
*/
|
||||
private void countNudge(String outcome) {
|
||||
if (metrics != null) {
|
||||
metrics.inc(FleetMetrics.PUSH_NUDGES, "outcome", outcome);
|
||||
@@ -222,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();
|
||||
}
|
||||
@@ -730,7 +758,7 @@ public final class ReplyPushLoop {
|
||||
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders,
|
||||
questionReminderCount + 1, maxReminders,
|
||||
replyTargets.size(), tickets.size(), questions.size());
|
||||
countNudge("delivered");
|
||||
countNudge("sent");
|
||||
for (PendingIncident incident : incidents) {
|
||||
if (pendingIncidents.remove(incident.key(), incident)) {
|
||||
deliveredIncidents.add(incident.key());
|
||||
|
||||
@@ -29,4 +29,9 @@ public record SpawnRequest(String profileName, String requestedCwd, String calle
|
||||
String sessionName, String resumeSessionId) {
|
||||
this(profileName, requestedCwd, callerCwd, sessionName, resumeSessionId, null);
|
||||
}
|
||||
|
||||
/** Return a copy of this request with {@code profileName} replaced by {@code profile}. */
|
||||
public SpawnRequest withProfile(String profile) {
|
||||
return new SpawnRequest(profile, requestedCwd, callerCwd, sessionName, resumeSessionId, role);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -696,6 +696,10 @@ public final class FleetApp {
|
||||
/**
|
||||
* The worker's structured reply ({@code fleet_reply}) — resolves the blocking send awaiting
|
||||
* on this session, or queues the reply in the inbox when no send is open (CB-307).
|
||||
*
|
||||
* <p>fleetd #365: the response body's {@code delivered} field used to be unconditionally
|
||||
* {@code true} for either case; it now reports whether a send/ticket was actually resolved,
|
||||
* with {@code outcome} naming which (see {@link MessageService.ReplyOutcome}).
|
||||
*/
|
||||
private void replyMessage(Context ctx) {
|
||||
String id = ctx.pathParam("id");
|
||||
@@ -719,13 +723,19 @@ public final class FleetApp {
|
||||
// a WRONG value instead of failing loudly. The check lives in MessageService.reply so both
|
||||
// this door and FleetMcp.reply inherit the same rule; this catch only translates it into the
|
||||
// {error, detail} envelope this file uses everywhere else.
|
||||
MessageService.ReplyOutcome outcome;
|
||||
try {
|
||||
messages.reply(id, content);
|
||||
outcome = messages.reply(id, content);
|
||||
} catch (IllegalArgumentException e) {
|
||||
ctx.status(400).json(Map.of("error", "bad_request", "detail", e.getMessage()));
|
||||
return;
|
||||
}
|
||||
ctx.status(200).json(Map.of("sessionId", id, "delivered", true));
|
||||
// fleetd #365: "delivered": true used to be unconditional here, whether the reply resolved
|
||||
// a waiting send or was merely queued in the inbox for a later drain — the same gap
|
||||
// FleetMcp.reply had over MCP. `delivered` now reflects which actually happened, and
|
||||
// `outcome` names the specific case (see MessageService.ReplyOutcome).
|
||||
ctx.status(200).json(Map.of("sessionId", id, "delivered", outcome.delivered(),
|
||||
"outcome", outcome.wireName()));
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -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");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,241 @@
|
||||
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 java.util.Map;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #377: the startup git host line must report the <em>shape</em> of the value a
|
||||
* member receives as GITEA_HOST — set or unset, length, scheme, trailing slash — and the
|
||||
* value itself must never reach the log.
|
||||
*/
|
||||
class GitHostShapeReportTest {
|
||||
|
||||
/** A profile that opts in to the git-forge token, so GITEA_HOST is the var that matters. */
|
||||
private static final String GIT_TOKEN_CONFIG = """
|
||||
profiles:
|
||||
local:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
gitTokenEnv: WORKER_GITEA_TOKEN
|
||||
""";
|
||||
|
||||
private static FleetConfig load(Path dir, String yaml) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, yaml);
|
||||
return FleetConfig.load(f);
|
||||
}
|
||||
|
||||
/**
|
||||
* The level this logger had before {@link #attach()} raised it, so {@link #detach} can put it
|
||||
* back. {@code null} is a real value here — it means "inherit from the parent" — and that is
|
||||
* exactly the state this logger starts in, so it must be restored as {@code null} rather than
|
||||
* as some concrete level.
|
||||
*/
|
||||
private static Level originalLevel;
|
||||
|
||||
/**
|
||||
* fleetd #377: the shape lines are logged at INFO, and {@code logback-test.xml} sets
|
||||
* {@code dev.ltms.fleet} to WARN — so INFO events are dropped by the level check BEFORE any
|
||||
* appender sees them. Attaching an appender is therefore not enough: without raising the level
|
||||
* the list stays empty and every assertion below fails against correct production code. The
|
||||
* sibling report tests do the same thing at each call site (see
|
||||
* {@code MemberTrustModelReportTest}); doing it here keeps it in one place.
|
||||
*/
|
||||
private static ListAppender<ILoggingEvent> attach() {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class);
|
||||
originalLevel = logger.getLevel();
|
||||
logger.setLevel(Level.INFO);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
return appender;
|
||||
}
|
||||
|
||||
private static void detach(ListAppender<ILoggingEvent> appender) {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class);
|
||||
logger.detachAppender(appender);
|
||||
logger.setLevel(originalLevel);
|
||||
}
|
||||
|
||||
private static List<String> messages(ListAppender<ILoggingEvent> appender) {
|
||||
return appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList();
|
||||
}
|
||||
|
||||
@Test
|
||||
void setLineReportsLengthSchemeAndTrailingSlash(@TempDir Path dir) throws Exception {
|
||||
FleetConfig cfg = load(dir, GIT_TOKEN_CONFIG);
|
||||
String value = "https://git.example.test/";
|
||||
|
||||
ListAppender<ILoggingEvent> appender = attach();
|
||||
try {
|
||||
Fleetd.reportGitHostShape(cfg, Map.of("GITEA_HOST", value));
|
||||
} finally {
|
||||
detach(appender);
|
||||
}
|
||||
|
||||
String expected = "startup git host GITEA_HOST: set (profile 'local' gitHostEnv) — "
|
||||
+ "length=" + value.length() + ", startsWithScheme=true, trailingSlash=true";
|
||||
assertTrue(messages(appender).contains(expected),
|
||||
"a value with a scheme and a trailing slash must be reported by shape only: "
|
||||
+ "its length, startsWithScheme=true, trailingSlash=true");
|
||||
}
|
||||
|
||||
@Test
|
||||
void bareHostLineReportsNoSchemeNoTrailingSlash(@TempDir Path dir) throws Exception {
|
||||
FleetConfig cfg = load(dir, GIT_TOKEN_CONFIG);
|
||||
String value = "git.example.test";
|
||||
|
||||
ListAppender<ILoggingEvent> appender = attach();
|
||||
try {
|
||||
Fleetd.reportGitHostShape(cfg, Map.of("GITEA_HOST", value));
|
||||
} finally {
|
||||
detach(appender);
|
||||
}
|
||||
|
||||
String expected = "startup git host GITEA_HOST: set (profile 'local' gitHostEnv) — "
|
||||
+ "length=" + value.length() + ", startsWithScheme=false, trailingSlash=false";
|
||||
assertTrue(messages(appender).contains(expected),
|
||||
"a bare host name must report startsWithScheme=false, trailingSlash=false");
|
||||
}
|
||||
|
||||
@Test
|
||||
void unsetVariableIsLoggedAtInfoAndDoesNotThrow(@TempDir Path dir) throws Exception {
|
||||
FleetConfig cfg = load(dir, GIT_TOKEN_CONFIG);
|
||||
|
||||
ListAppender<ILoggingEvent> appender = attach();
|
||||
try {
|
||||
// GITEA_HOST absent from the map entirely — the unset case must not throw.
|
||||
Fleetd.reportGitHostShape(cfg, Map.of());
|
||||
} finally {
|
||||
detach(appender);
|
||||
}
|
||||
|
||||
assertTrue(appender.list.stream().anyMatch(e ->
|
||||
e.getLevel() == Level.INFO
|
||||
&& e.getFormattedMessage()
|
||||
.startsWith("startup git host GITEA_HOST: unset (profile 'local' gitHostEnv)")),
|
||||
"an unset git host is useful information, not an error — say it at INFO level");
|
||||
}
|
||||
|
||||
/**
|
||||
* The important test: the value must NEVER appear in the log output. This test fails if the
|
||||
* line is ever changed to include the value, because it drives the real logging path with a
|
||||
* value that carries a marker no shape field could contain.
|
||||
*/
|
||||
@Test
|
||||
void theValueNeverAppearsInLogOutput(@TempDir Path dir) throws Exception {
|
||||
FleetConfig cfg = load(dir, GIT_TOKEN_CONFIG);
|
||||
String marker = "never-logged-host-shape-377";
|
||||
String value = "https://" + marker + "/";
|
||||
|
||||
ListAppender<ILoggingEvent> appender = attach();
|
||||
try {
|
||||
Fleetd.reportGitHostShape(cfg, Map.of("GITEA_HOST", value));
|
||||
} finally {
|
||||
detach(appender);
|
||||
}
|
||||
|
||||
List<String> msgs = messages(appender);
|
||||
assertFalse(msgs.stream().anyMatch(m -> m.contains(value)),
|
||||
"the full GITEA_HOST value must never reach the log");
|
||||
assertFalse(msgs.stream().anyMatch(m -> m.contains(marker)),
|
||||
"no fragment of the value may reach the log — a line that embeds the value"
|
||||
+ " would leak at least this marker");
|
||||
assertTrue(msgs.stream().anyMatch(m -> m.contains("startsWithScheme=true")
|
||||
&& m.contains("trailingSlash=true")),
|
||||
"the shape must still be reported although the value is not");
|
||||
}
|
||||
|
||||
@Test
|
||||
void unsetReportAlsoCarriesNoValue(@TempDir Path dir) throws Exception {
|
||||
FleetConfig cfg = load(dir, GIT_TOKEN_CONFIG);
|
||||
// A blank value is not injected either (putIfPresent skips it), so it must report unset
|
||||
// without echoing anything of it.
|
||||
ListAppender<ILoggingEvent> appender = attach();
|
||||
try {
|
||||
Fleetd.reportGitHostShape(cfg, Map.of("GITEA_HOST", " "));
|
||||
} finally {
|
||||
detach(appender);
|
||||
}
|
||||
|
||||
assertTrue(messages(appender).stream().anyMatch(m ->
|
||||
m.startsWith("startup git host GITEA_HOST: unset")),
|
||||
"a blank value is skipped by the launcher, so the line reports unset");
|
||||
}
|
||||
|
||||
@Test
|
||||
void defaultGitHostEnvIsGITEA_HOST(@TempDir Path dir) throws Exception {
|
||||
FleetConfig cfg = load(dir, GIT_TOKEN_CONFIG);
|
||||
|
||||
assertEquals(Map.of("GITEA_HOST", List.of("profile 'local' gitHostEnv")),
|
||||
Fleetd.gitHostEnvVars(cfg));
|
||||
}
|
||||
|
||||
@Test
|
||||
void explicitGitHostEnvReportsUnderItsOwnName(@TempDir Path dir) throws Exception {
|
||||
FleetConfig cfg = load(dir, """
|
||||
profiles:
|
||||
local:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
gitTokenEnv: WORKER_GITEA_TOKEN
|
||||
gitHostEnv: MY_FORGE_HOST
|
||||
""");
|
||||
String value = "https://forge.example.test/";
|
||||
|
||||
ListAppender<ILoggingEvent> appender = attach();
|
||||
try {
|
||||
Fleetd.reportGitHostShape(cfg, Map.of("MY_FORGE_HOST", value));
|
||||
} finally {
|
||||
detach(appender);
|
||||
}
|
||||
|
||||
assertTrue(messages(appender).contains(
|
||||
"startup git host MY_FORGE_HOST: set (profile 'local' gitHostEnv) — "
|
||||
+ "length=" + value.length()
|
||||
+ ", startsWithScheme=true, trailingSlash=true"),
|
||||
"an explicitly named gitHostEnv is reported under that name");
|
||||
}
|
||||
|
||||
@Test
|
||||
void noGitTokenMeansNoHostLineIsNeeded(@TempDir Path dir) throws Exception {
|
||||
FleetConfig cfg = load(dir, """
|
||||
profiles:
|
||||
local:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
""");
|
||||
|
||||
ListAppender<ILoggingEvent> appender = attach();
|
||||
try {
|
||||
Fleetd.reportGitHostShape(cfg, Map.of());
|
||||
} finally {
|
||||
detach(appender);
|
||||
}
|
||||
|
||||
assertEquals(List.of("startup git host: no profile sets a gitTokenEnv — nothing to check"),
|
||||
messages(appender));
|
||||
}
|
||||
|
||||
@Test
|
||||
void aHostWithAPortIsNotTreatedAsAScheme() {
|
||||
assertFalse(Fleetd.startsWithScheme("git.example.test"));
|
||||
assertFalse(Fleetd.startsWithScheme("git.example.test:3000"));
|
||||
assertFalse(Fleetd.startsWithScheme("git.example.test:3000/"));
|
||||
assertTrue(Fleetd.startsWithScheme("https://git.example.test"));
|
||||
assertTrue(Fleetd.startsWithScheme("https://git.example.test:3000/"));
|
||||
assertTrue(Fleetd.startsWithScheme("ssh://git.example.test"));
|
||||
}
|
||||
}
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
+172
@@ -0,0 +1,172 @@
|
||||
package dev.ltms.fleet.config;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.lang.reflect.Constructor;
|
||||
import java.lang.reflect.RecordComponent;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Arrays;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Locale;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
import java.util.Set;
|
||||
import java.util.TreeSet;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
|
||||
/**
|
||||
* Fleetd #358, the same "defect factory" #357 guarded on {@code FleetConfig.withDefaults()}
|
||||
* (see {@code FleetConfigWithDefaultsPreservesEveryComponentTest}), reproduced here on
|
||||
* {@link FleetConfig.Profile}. {@code Profile} carries a long back-compat constructor ladder — 8
|
||||
* constructors, re-counted directly against the source rather than trusted from the ticket, at
|
||||
* arities 25, 24, 22, 20, 18, 15, 14 and 12, against a canonical arity of 26 — and exactly ONE
|
||||
* rebuild site, {@link FleetConfig.Profile#withProfile(String)}, whose own
|
||||
* {@code return new Profile(...)} call is written at a literal 26-arg count. Add a 27th component
|
||||
* and its established back-compat constructor at the old (26-arg) arity, and {@code withProfile}'s
|
||||
* own call becomes a legal match for that new overload — silently dropping the new component every
|
||||
* time a profile's name is defaulted from its {@code workers:} key.
|
||||
*
|
||||
* <p>Builds one {@link FleetConfig.Profile} through the TRUE canonical constructor — resolved by
|
||||
* the record's own component types via {@code getDeclaredConstructor}, never by argument count —
|
||||
* with a real, distinctive, non-null value in every component, calls {@link
|
||||
* FleetConfig.Profile#withProfile(String)}, and asserts every component except {@code profile}
|
||||
* itself survives unchanged, while {@code profile} comes back as the new name it was given.
|
||||
*
|
||||
* <p>Every value here is chosen so {@code Profile}'s own compact constructor (which normalizes
|
||||
* several components — defaults {@code argv}/{@code kind}/{@code placement}/{@code workspace}/
|
||||
* {@code gitHostEnv}, nulls a handful of blank-checked strings, clamps {@code weight}, coerces
|
||||
* {@code subscription}) leaves it unchanged: every String is non-blank and already in the shape the
|
||||
* compact constructor would otherwise coerce it to (e.g. {@code placement} is already lowercase),
|
||||
* and every collection is non-empty. That is what makes "must survive unchanged" a valid assertion
|
||||
* for every component below, the same reasoning {@code FleetConfigWithDefaultsPreservesEveryComponentTest}
|
||||
* documents for {@code withDefaults()}.
|
||||
*
|
||||
* <p>{@link #EXCLUDED_FROM_SURVIVAL_CHECK} is kept deliberately empty and size-pinned by
|
||||
* {@link #exclusionListSizeIsPinned()} — a checker whose escape hatch can grow to silence a failure
|
||||
* is not a checker. Every one of {@code Profile}'s 26 current components has a real, non-null,
|
||||
* non-blank value here and none is excluded.
|
||||
*/
|
||||
class FleetConfigProfileWithProfilePreservesEveryComponentTest {
|
||||
|
||||
private static final RecordComponent[] COMPONENTS = FleetConfig.Profile.class.getRecordComponents();
|
||||
|
||||
/** Deliberately empty today; grow it only with a matching justification, and re-pin the size. */
|
||||
private static final Set<String> EXCLUDED_FROM_SURVIVAL_CHECK = Set.of();
|
||||
|
||||
/** One real, distinctive, non-null value per component, chosen to survive the compact ctor. */
|
||||
private static Map<String, Object> baseValues() {
|
||||
Map<String, Object> v = new LinkedHashMap<>();
|
||||
v.put("profile", "profile-guard");
|
||||
v.put("baseUrl", "https://guard.example/base");
|
||||
v.put("model", "model-guard");
|
||||
v.put("configDir", "/config/guard");
|
||||
v.put("tokenEnv", "GUARD_TOKEN");
|
||||
v.put("argv", List.of("guard-cmd"));
|
||||
v.put("placement", "guard-placement");
|
||||
v.put("workspace", "workspace-guard");
|
||||
v.put("tabLabel", "tab-guard");
|
||||
v.put("mcpUrl", "https://mcp.guard/");
|
||||
v.put("cwd", "/cwd/guard");
|
||||
v.put("parityOverlay", List.of(".guardrc"));
|
||||
v.put("gitTokenEnv", "GUARD_GIT_TOKEN");
|
||||
v.put("gitHostEnv", "GUARD_GIT_HOST");
|
||||
v.put("kind", "claude-code");
|
||||
v.put("env", Map.of("GUARD_ENV", "1"));
|
||||
v.put("weight", 2.5f);
|
||||
v.put("maxLoad", 4);
|
||||
v.put("subscription", Boolean.TRUE);
|
||||
v.put("exhaustedPattern", "pattern-guard");
|
||||
v.put("credentialId", "cred-guard");
|
||||
v.put("ideMcpUrl", "https://ide.guard/");
|
||||
v.put("ideProjectDir", "ide-project-guard");
|
||||
v.put("ideOpenCommand", "open-guard {dir}");
|
||||
v.put("autoCompactWindow", 150_000);
|
||||
v.put("errorPattern", "error-pattern-guard");
|
||||
assertNamesMatchComponents(v);
|
||||
return v;
|
||||
}
|
||||
|
||||
/**
|
||||
* Guards {@link #baseValues()} itself against drifting from the record's real shape — forgetting
|
||||
* to add a new component here fails this assertion by name, rather than silently checking one
|
||||
* component fewer than the record has.
|
||||
*/
|
||||
private static void assertNamesMatchComponents(Map<String, Object> values) {
|
||||
Set<String> names = new TreeSet<>();
|
||||
for (RecordComponent rc : COMPONENTS) {
|
||||
names.add(rc.getName());
|
||||
}
|
||||
assertEquals(names, new TreeSet<>(values.keySet()),
|
||||
"this test's value map has drifted from FleetConfig.Profile's actual components — "
|
||||
+ "update baseValues() alongside the record");
|
||||
}
|
||||
|
||||
/**
|
||||
* Builds a {@link FleetConfig.Profile} through the TRUE canonical constructor — resolved by the
|
||||
* record's own component types, not by argument count — so this never accidentally exercises a
|
||||
* back-compat overload the way a literal {@code new Profile(...)} call risks doing.
|
||||
*/
|
||||
private static FleetConfig.Profile profileOf(Map<String, Object> values) throws ReflectiveOperationException {
|
||||
Class<?>[] types = Arrays.stream(COMPONENTS).map(RecordComponent::getType).toArray(Class<?>[]::new);
|
||||
Object[] args = Arrays.stream(COMPONENTS).map(rc -> values.get(rc.getName())).toArray();
|
||||
Constructor<FleetConfig.Profile> ctor = FleetConfig.Profile.class.getDeclaredConstructor(types);
|
||||
return ctor.newInstance(args);
|
||||
}
|
||||
|
||||
@Test
|
||||
void exclusionListSizeIsPinned() {
|
||||
assertEquals(0, EXCLUDED_FROM_SURVIVAL_CHECK.size(),
|
||||
"EXCLUDED_FROM_SURVIVAL_CHECK grew from 0 — every entry needs a justification in "
|
||||
+ "this test class's javadoc AND this assertion re-pinned to the new size; a "
|
||||
+ "growing exclusion list that silences failures on its own is not a guard");
|
||||
}
|
||||
|
||||
/**
|
||||
* The mutation this is built to catch: make {@code withProfile(String)}'s final constructor call
|
||||
* literal at some arg count, add one more component to the record with a new back-compat
|
||||
* constructor at the old arity, and the stale call silently rebinds. Every component here is real
|
||||
* and non-null/non-blank, so none of it should be replaced by {@code withProfile}, except
|
||||
* {@code profile} itself, which the method is documented to replace.
|
||||
*/
|
||||
@Test
|
||||
void withProfilePreservesEveryOtherComponent() throws ReflectiveOperationException {
|
||||
Map<String, Object> base = baseValues();
|
||||
FleetConfig.Profile profile = profileOf(base);
|
||||
FleetConfig.Profile renamed = profile.withProfile("renamed-profile-guard");
|
||||
|
||||
List<String> dropped = new ArrayList<>();
|
||||
int checked = 0;
|
||||
for (RecordComponent rc : COMPONENTS) {
|
||||
String name = rc.getName();
|
||||
if (EXCLUDED_FROM_SURVIVAL_CHECK.contains(name)) {
|
||||
continue;
|
||||
}
|
||||
checked++;
|
||||
Object expected = "profile".equals(name) ? "renamed-profile-guard" : base.get(name);
|
||||
Object actual;
|
||||
try {
|
||||
actual = rc.getAccessor().invoke(renamed);
|
||||
} catch (ReflectiveOperationException e) {
|
||||
throw new RuntimeException("failed to read FleetConfig.Profile." + name + "()", e);
|
||||
}
|
||||
if (!Objects.equals(expected, actual)) {
|
||||
dropped.add(String.format(Locale.ROOT,
|
||||
"%s: withProfile() was expected to carry (%s) for '%s' but returned %s — a "
|
||||
+ "component silently dropped by withProfile(), the shape of the "
|
||||
+ "defect this test exists to catch (its final \"return new "
|
||||
+ "Profile(...)\" call binding to a back-compat constructor instead "
|
||||
+ "of the true canonical one)",
|
||||
name, expected, name, actual));
|
||||
}
|
||||
}
|
||||
|
||||
System.out.printf(Locale.ROOT,
|
||||
"FleetConfig.Profile.withProfile() component-survival coverage — %d components, %d "
|
||||
+ "checked, %d excluded, %d survived%n",
|
||||
COMPONENTS.length, checked, EXCLUDED_FROM_SURVIVAL_CHECK.size(), checked - dropped.size());
|
||||
assertEquals(List.of(), dropped,
|
||||
"withProfile() silently dropped these components: " + dropped);
|
||||
}
|
||||
}
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -30,6 +30,7 @@ import java.util.concurrent.ScheduledFuture;
|
||||
import java.util.concurrent.ScheduledThreadPoolExecutor;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
import java.util.function.BiConsumer;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
@@ -206,6 +207,139 @@ class FleetHealthMonitorTest {
|
||||
.workingSuspectAfterOrDefault());
|
||||
}
|
||||
|
||||
// --- fleetd #386: a stall detector whose only clock freezes with a sleeping host is worse
|
||||
// than a silent one — it reports "quiet" for a member that was genuinely busy for hours.
|
||||
|
||||
@Test void monotonicClockFrozenPastThresholdOnRealClockStillReportsStallSuspected() {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
AtomicLong mono = new AtomicLong(0);
|
||||
AtomicLong real = new AtomicLong(0);
|
||||
FakeHerdr herdr = new FakeHerdr().withAgent("busy", "term_busy", "pane_busy", "tab_busy");
|
||||
var scheduler = Executors.newSingleThreadScheduledExecutor();
|
||||
FleetHealthMonitor monitor = monitorWithClocks(herdr,
|
||||
List.of(member("term_busy", MemberSession.State.BUSY, 0, 0)), scheduler,
|
||||
mono::get, real::get, 60, 600, (_, _) -> { });
|
||||
|
||||
monitor.tick(); // establishes the clock baseline; nothing has diverged yet
|
||||
assertEquals(0, appender.list.stream().filter(event -> event.getFormattedMessage()
|
||||
.contains("state=STALL_SUSPECTED")).count());
|
||||
|
||||
// The host "sleeps": the monotonic clock stands completely still while the real clock
|
||||
// keeps moving, past the 600s stall threshold.
|
||||
real.set(TimeUnit.SECONDS.toNanos(700));
|
||||
monitor.tick();
|
||||
monitor.stop();
|
||||
|
||||
assertTrue(appender.list.stream().anyMatch(event -> event.getFormattedMessage()
|
||||
.contains("member=term_busy state=STALL_SUSPECTED")),
|
||||
"the real clock crossed the stall threshold even though the monotonic clock never moved");
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
|
||||
@Test void clockDivergenceIsLoggedOnceNotOncePerTick() {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
AtomicLong mono = new AtomicLong(0);
|
||||
AtomicLong real = new AtomicLong(0);
|
||||
var scheduler = Executors.newSingleThreadScheduledExecutor();
|
||||
FleetHealthMonitor monitor = monitorWithClocks(new FakeHerdr(), List.of(), scheduler,
|
||||
mono::get, real::get, 60, 600, (_, _) -> { });
|
||||
|
||||
monitor.tick(); // baseline: no divergence possible yet
|
||||
|
||||
// One sleep gap: the monotonic clock is frozen while the real clock jumps far past one
|
||||
// tick interval (60s).
|
||||
real.set(TimeUnit.SECONDS.toNanos(700));
|
||||
monitor.tick();
|
||||
|
||||
// The host is awake again: both clocks advance together from here, so no more divergence.
|
||||
mono.set(TimeUnit.SECONDS.toNanos(10));
|
||||
real.set(TimeUnit.SECONDS.toNanos(710));
|
||||
monitor.tick();
|
||||
mono.set(TimeUnit.SECONDS.toNanos(20));
|
||||
real.set(TimeUnit.SECONDS.toNanos(720));
|
||||
monitor.tick();
|
||||
monitor.stop();
|
||||
|
||||
assertEquals(1, appender.list.stream().filter(event -> event.getFormattedMessage()
|
||||
.contains("the monotonic clock did not advance"))
|
||||
.count(), "one sleep gap must produce exactly one divergence line, not one per tick");
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #386 follow-up, added on merge. The fix carries a PER-MEMBER drift baseline, so drift
|
||||
* from a sleep that happened BEFORE a member went busy is never charged to that member. The
|
||||
* two tests shipped with the fix both start with the member already BUSY, so a single global
|
||||
* baseline passes them — this one fails without the per-member map.
|
||||
*
|
||||
* <p>Order matters: the host sleeps while nothing is busy, and only then does a member take a
|
||||
* turn. Its stall clock must start at zero.
|
||||
*/
|
||||
@Test void driftFromASleepBeforeAMemberWentBusyIsNotChargedToThatMember() {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
AtomicLong mono = new AtomicLong(0);
|
||||
AtomicLong real = new AtomicLong(0);
|
||||
AtomicReference<List<MemberSession>> roster = new AtomicReference<>(List.of());
|
||||
FakeHerdr herdr = new FakeHerdr().withAgent("busy", "term_busy", "pane_busy", "tab_busy");
|
||||
var scheduler = Executors.newSingleThreadScheduledExecutor();
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
FleetHealthMonitor monitor = new FleetHealthMonitor(agents, roster::get,
|
||||
new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()),
|
||||
scheduler, mono::get, real::get, 60, 600, (_, _) -> { });
|
||||
|
||||
monitor.tick(); // baseline, no members yet
|
||||
|
||||
// The host sleeps for 700s with nobody busy: the monotonic clock stands still.
|
||||
real.set(TimeUnit.SECONDS.toNanos(700));
|
||||
monitor.tick();
|
||||
|
||||
// Awake again. Only NOW does a member start a turn, with a fresh activity stamp taken
|
||||
// from the monotonic clock. Both clocks advance together from here.
|
||||
mono.set(TimeUnit.SECONDS.toNanos(10));
|
||||
real.set(TimeUnit.SECONDS.toNanos(710));
|
||||
roster.set(List.of(member("term_busy", MemberSession.State.BUSY,
|
||||
0, TimeUnit.SECONDS.toNanos(10))));
|
||||
monitor.tick();
|
||||
|
||||
mono.set(TimeUnit.SECONDS.toNanos(20));
|
||||
real.set(TimeUnit.SECONDS.toNanos(720));
|
||||
monitor.tick();
|
||||
monitor.stop();
|
||||
|
||||
assertEquals(0, appender.list.stream().filter(event -> event.getFormattedMessage()
|
||||
.contains("member=term_busy state=STALL_SUSPECTED")).count(),
|
||||
"the member has been busy for 10s, not 710s — the earlier sleep is not its stall");
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
|
||||
private static FleetHealthMonitor monitorWithClocks(FakeHerdr herdr, List<MemberSession> roster,
|
||||
java.util.concurrent.ScheduledExecutorService scheduler, LongSupplier clock,
|
||||
LongSupplier realtimeClock, long intervalSeconds, long workingSuspectAfterSeconds,
|
||||
BiConsumer<String, String> failTarget) {
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
return new FleetHealthMonitor(agents, () -> roster,
|
||||
new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()),
|
||||
scheduler, clock, realtimeClock, intervalSeconds, workingSuspectAfterSeconds, failTarget);
|
||||
}
|
||||
|
||||
@Test void goneMemberRecoveryLogsOnceWithoutRefiringTargetFailure() {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class);
|
||||
Level previousLevel = logger.getLevel();
|
||||
|
||||
@@ -357,6 +357,54 @@ class CompletionResolverTest {
|
||||
"the failure carries whatever was on screen: " + waiter.getNow(null).text());
|
||||
}
|
||||
|
||||
@Test
|
||||
void aPlausibleLookingReplyInsideTheFloorStillFails() {
|
||||
// fleetd#376 guard. A fix was attempted that inspected the pane inside the floor and resolved
|
||||
// a COMPLETION when the text "looked like a real reply". Every cheap test for that is unsafe:
|
||||
// lastAssistantBlock falls back to the WHOLE pane when there is no ⏺ marker, and a crash pane
|
||||
// almost always contains sentence punctuation — in a file path, a version, or a hostname.
|
||||
// This pane is the trap: it reads like a finished answer and it is a backend failure.
|
||||
FakeHerdr herdr = new FakeHerdr().readText(
|
||||
"Error: connection reset while loading src/main/java/Foo.java v1.2.3\n❯ ");
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
long[] clock = {10_000_000_000L};
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
|
||||
ExhaustedPatternLookup.none(), ExhaustionSink.none(), () -> clock[0]);
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
var turn = new CompletionResolver.InFlight(waiter, null, clock[0]);
|
||||
clock[0] += CompletionResolver.MIN_TURN_NANOS - 1; // inside the floor
|
||||
|
||||
resolver.resolve("term_a", turn);
|
||||
|
||||
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(),
|
||||
"inside the floor the verdict is always FAILED — never guess a completion from pane text");
|
||||
}
|
||||
|
||||
@Test
|
||||
void theTooFastFailureDoesNotAssertACauseItCannotKnow() {
|
||||
// fleetd#376: the message used to say "most likely a backend error before any work started".
|
||||
// When no error pattern matches, that cause is a guess, and a reader who believes it stops
|
||||
// looking at the pane. The verdict stays FAILED; only the claim about WHY is withdrawn.
|
||||
FakeHerdr herdr = new FakeHerdr().readText("I am running on opencode/mimo-v2.5-free.\n❯ ");
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
long[] clock = {10_000_000_000L};
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
|
||||
ExhaustedPatternLookup.none(), ExhaustionSink.none(), () -> clock[0]);
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
var turn = new CompletionResolver.InFlight(waiter, null, clock[0]);
|
||||
clock[0] += CompletionResolver.MIN_TURN_NANOS - 1; // inside the floor
|
||||
|
||||
resolver.resolve("term_a", turn);
|
||||
|
||||
String text = waiter.getNow(null).text();
|
||||
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(), "still fails, still loud");
|
||||
assertFalse(text.contains("most likely a backend error"),
|
||||
"an unmatched fast turn must not assert a backend error: " + text);
|
||||
assertTrue(text.contains("mimo-v2.5-free"), "the pane is still carried: " + text);
|
||||
}
|
||||
|
||||
@Test
|
||||
void aBusyToDoneTransitionJustOutsideTheFloorResolvesNormally() {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("⏺ a real, if quick, answer\n❯ ");
|
||||
@@ -1117,8 +1165,16 @@ class CompletionResolverTest {
|
||||
assertTrue(waiter.isDone());
|
||||
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(),
|
||||
"still a failure — the floor itself, not the pattern, is why");
|
||||
assertTrue(waiter.getNow(null).text().contains("too fast to be real work"),
|
||||
"a non-match inside the floor stays the generic too-fast reason: " + waiter.getNow(null).text());
|
||||
// fleetd#376: this used to assert the phrase "too fast to be real work", which carried the
|
||||
// claim "most likely a backend error before any work started". With no pattern matched that
|
||||
// cause is a guess, so the wording was withdrawn. What this test really guards is unchanged:
|
||||
// the floor alone still fails the turn, it stays generic, and it never notifies the sink.
|
||||
String reason = waiter.getNow(null).text();
|
||||
assertTrue(reason.contains("inside the floor"),
|
||||
"a non-match inside the floor stays the generic floor reason: " + reason);
|
||||
assertFalse(reason.contains("most likely a backend error"),
|
||||
"a non-match must not assert a cause it did not establish: " + reason);
|
||||
assertTrue(reason.contains("still starting up"), "the pane is still carried: " + reason);
|
||||
assertTrue(notified.isEmpty(), "a non-match must never notify the typed sink");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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)));
|
||||
}
|
||||
|
||||
@@ -108,8 +110,10 @@ class FleetMcpTest {
|
||||
}
|
||||
assertTrue(rendezvous.isWaiting("term_a"), "send should have opened its waiter");
|
||||
|
||||
McpSchema.CallToolResult reply = FleetMcp.reply(messages, "term_a", "LGTM");
|
||||
assertEquals("delivered", textOf(reply));
|
||||
// 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", Role.WORKER, "LGTM");
|
||||
assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND.description(), textOf(reply));
|
||||
|
||||
McpSchema.CallToolResult res = send.get(6, TimeUnit.SECONDS);
|
||||
assertNotEquals(Boolean.TRUE, res.isError());
|
||||
@@ -134,8 +138,8 @@ class FleetMcpTest {
|
||||
}
|
||||
assertTrue(rendezvous.isWaiting("term_a"), "send should have opened its waiter");
|
||||
|
||||
McpSchema.CallToolResult reply = FleetMcp.reply(messages, "term_a", "async LGTM");
|
||||
assertEquals("delivered", textOf(reply));
|
||||
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.
|
||||
McpSchema.CallToolResult polled = FleetMcp.poll(messages, ticket, null);
|
||||
@@ -181,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);
|
||||
@@ -213,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");
|
||||
@@ -267,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(
|
||||
@@ -276,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)));
|
||||
}
|
||||
|
||||
@@ -328,9 +332,10 @@ class FleetMcpTest {
|
||||
@Test
|
||||
void replyWithNoPendingSendIsQueuedNotError() {
|
||||
// CB-307: a reply with no open send is now queued in the inbox, not an error.
|
||||
McpSchema.CallToolResult res = FleetMcp.reply(messages, "term_a", "orphan");
|
||||
// fleetd #365: it must also no longer claim "delivered" — nothing was waiting for it.
|
||||
McpSchema.CallToolResult res = FleetMcp.reply(messages, "term_a", Role.WORKER, "orphan");
|
||||
assertNotEquals(Boolean.TRUE, res.isError(), "a queued reply is not an error");
|
||||
assertEquals("delivered", textOf(res));
|
||||
assertEquals(MessageService.ReplyOutcome.QUEUED.description(), textOf(res));
|
||||
|
||||
// The reply is drainable by target.
|
||||
var drained = messages.drainReplies("term_a");
|
||||
@@ -338,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
|
||||
@@ -348,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"),
|
||||
@@ -361,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");
|
||||
@@ -412,8 +451,8 @@ 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");
|
||||
assertEquals("delivered", textOf(reply));
|
||||
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)));
|
||||
}
|
||||
|
||||
@@ -1251,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);
|
||||
}
|
||||
}
|
||||
@@ -8,6 +8,7 @@ import java.nio.charset.StandardCharsets;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.ArrayList;
|
||||
import java.util.HashMap;
|
||||
import java.util.HashSet;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
@@ -17,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;
|
||||
@@ -45,6 +47,15 @@ class EnvAllowListScrubTest {
|
||||
/** Env var names appearing in command output; anything else (prompts, wrapped lines) is noise. */
|
||||
private static final Pattern ENV_NAME = Pattern.compile("^([A-Za-z_][A-Za-z0-9_]*)$");
|
||||
|
||||
/**
|
||||
* fleetd #388: the sentinel name the generated {@code .zshenv} guard uses, kept here as a
|
||||
* literal rather than referencing {@link EnvAllowListScrub#SCRUB_SENTINEL} — the two tests that
|
||||
* use it must still compile and run against the pre-fix production class (which has no such
|
||||
* constant), so the revert-and-prove-it-fails step exercises a real assertion instead of a
|
||||
* compilation error.
|
||||
*/
|
||||
private static final String SCRUB_SENTINEL_NAME = "_CB633_SCRUBBED";
|
||||
|
||||
/**
|
||||
* The equality test. Expected survivors = baseline exports ∩ allowed — i.e. every survivor is
|
||||
* allowed AND every allowed name that existed survives. The operator's own secret-store exports
|
||||
@@ -108,6 +119,230 @@ 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 {
|
||||
assumeTrue(Files.isExecutable(ZSH), "/bin/zsh not present — nothing to prove here");
|
||||
Set<String> allowed = MemberEnvAllowList.derive(List.of());
|
||||
Path zdotdir = EnvAllowListScrub.generate(tmp, allowed, currentUserGroup());
|
||||
|
||||
Map<String, String> cleanParent = Map.of(
|
||||
"HOME", System.getProperty("user.home"),
|
||||
"PATH", "/usr/bin:/bin",
|
||||
"SHELL", "/bin/zsh");
|
||||
exportedNamesFromCleanParent(cleanParent, zdotdir);
|
||||
|
||||
EnvAllowListScrub.ScrubReport report = EnvAllowListScrub.readReport(zdotdir);
|
||||
assertNotNull(report, "a group-shared completed login shell must leave a report behind");
|
||||
assertTrue(report.allowed() >= 0 && report.total() >= report.allowed(),
|
||||
"allowed N of M with N <= M — the denominator is always reported");
|
||||
assertEquals("rw-rw----", java.nio.file.attribute.PosixFilePermissions.toString(
|
||||
Files.getPosixFilePermissions(zdotdir.resolve(EnvAllowListScrub.REPORT_FILE))),
|
||||
"the pre-created receipt must be group-writable");
|
||||
assertEquals("rwxr-x---", java.nio.file.attribute.PosixFilePermissions.toString(
|
||||
Files.getPosixFilePermissions(zdotdir)),
|
||||
"group sharing must not make the ZDOTDIR directory group-writable");
|
||||
}
|
||||
|
||||
/** Report parsing is lenient: absent file → null (no measurement), not an exception. */
|
||||
@Test
|
||||
void readReportReturnsNullForADirectoryWithoutOne(@TempDir Path dir) {
|
||||
@@ -224,6 +459,140 @@ class EnvAllowListScrubTest {
|
||||
+ "shell. A difference here means the scrub is dead on Linux.");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #388: the actual gap. zsh reads {@code .zshenv} always, {@code .zprofile}/
|
||||
* {@code .zlogin} only for a LOGIN shell, and {@code .zshrc} only for an INTERACTIVE one — so a
|
||||
* shell that is NEITHER (a bare {@code /bin/zsh} reading a script off a non-tty stdin, no
|
||||
* {@code -l}, no {@code -i}) reads only {@code .zshenv} and stops. Before this fix, that shell
|
||||
* never reached {@code scrub.zsh} at all: the decoy secret below would survive untouched. This
|
||||
* test injects that decoy directly into the process environment (not via a sourced dotfile,
|
||||
* since the whole point of the gap is that {@code .zshenv} is normally close to empty) so the
|
||||
* test does not depend on any real {@code ~/.zshrc} content existing on the host.
|
||||
*/
|
||||
@Test
|
||||
void scrubRunsInAShellThatIsNeitherLoginNorInteractive(@TempDir Path tmp) throws Exception {
|
||||
assumeTrue(Files.isExecutable(ZSH), "/bin/zsh not present — nothing to prove here");
|
||||
Set<String> allowed = MemberEnvAllowList.derive(List.of());
|
||||
Path zdotdir = EnvAllowListScrub.generate(tmp, allowed);
|
||||
|
||||
Map<String, String> cleanParent = new HashMap<>(Map.of(
|
||||
"HOME", System.getProperty("user.home"),
|
||||
"PATH", "/usr/bin:/bin",
|
||||
"SHELL", "/bin/zsh",
|
||||
"USER", System.getProperty("user.name", "nobody"),
|
||||
"TMPDIR", tmp.toString()));
|
||||
cleanParent.put("FLEETD_TEST_DECOY_SECRET", "x"); // not on any allow-list; must be blanked
|
||||
|
||||
List<String> neither = List.of(); // no -l, no -i; stdin is a pipe (never a tty) either way
|
||||
Set<String> baseline = exportedNamesFromCleanParent(cleanParent, null, neither);
|
||||
Set<String> scrubbed = exportedNamesFromCleanParent(cleanParent, zdotdir, neither);
|
||||
|
||||
Set<String> expected = new TreeSet<>();
|
||||
for (String name : baseline) {
|
||||
if (MemberEnvAllowList.keeps(allowed, name)) {
|
||||
expected.add(name);
|
||||
}
|
||||
}
|
||||
assertTrue(baseline.contains("FLEETD_TEST_DECOY_SECRET"),
|
||||
"sanity: the decoy must actually reach the un-scrubbed baseline, or this test proves "
|
||||
+ "nothing");
|
||||
expected.add("ZDOTDIR"); // the harness set it and it is infrastructure, so it must survive
|
||||
expected.add(SCRUB_SENTINEL_NAME); // set by the new .zshenv guard once scrubbed
|
||||
assertEquals(expected, scrubbed,
|
||||
"a pane shell that is NEITHER login nor interactive must still be scrubbed — its "
|
||||
+ "surviving exported names must EQUAL baseline ∩ allow-list, plus the "
|
||||
+ "sentinel the guard sets once it has run. FLEETD_TEST_DECOY_SECRET surviving "
|
||||
+ "here means the gap is still open.");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #388 invariants 3 and 4, which a name-set equality cannot show: a member's own tooling
|
||||
* forks plain, non-login, non-interactive zsh processes for a single command (the same shape as
|
||||
* the pane shell itself), and such a child must (a) keep whatever its parent deliberately set
|
||||
* for it, never (b) re-run the scrub and blank it, and never (c) overwrite the pane's own
|
||||
* {@code scrub-report.txt} with a description of itself instead of the pane. All three can only
|
||||
* be shown by actually running a child process from within the scrubbed pane shell.
|
||||
*
|
||||
* <p>The pane and the child both report presence via {@code ${NAME:+present}} — empty when a
|
||||
* name is unset OR blanked (exported empty), {@code present} when it is set and non-empty. No
|
||||
* value is ever printed, only these two shapes and the literal word {@code set}/{@code unset}
|
||||
* for the sentinel.
|
||||
*/
|
||||
@Test
|
||||
void neitherShellChildKeepsParentVariablesAndReceiptStillDescribesThePane(@TempDir Path tmp) throws Exception {
|
||||
assumeTrue(Files.isExecutable(ZSH), "/bin/zsh not present — nothing to prove here");
|
||||
Set<String> allowed = MemberEnvAllowList.derive(List.of());
|
||||
Path zdotdir = EnvAllowListScrub.generate(tmp, allowed);
|
||||
|
||||
Map<String, String> paneEnv = new HashMap<>(Map.of(
|
||||
"HOME", System.getProperty("user.home"),
|
||||
"PATH", "/usr/bin:/bin",
|
||||
"SHELL", "/bin/zsh",
|
||||
"USER", System.getProperty("user.name", "nobody"),
|
||||
"TMPDIR", tmp.toString()));
|
||||
paneEnv.put("FLEETD_TEST_DECOY_SECRET", "x"); // not allow-listed; the pane must blank it
|
||||
paneEnv.put("ZDOTDIR", zdotdir.toAbsolutePath().toString());
|
||||
|
||||
// The pane's own script reports what IT sees, then forks a plain non-login, non-interactive
|
||||
// child — the shape a member's own tooling uses — carrying a variable the "parent" (this
|
||||
// pane) deliberately set for it, the way git sets GIT_DIR for a hook.
|
||||
String outerScript = """
|
||||
print -r -- "PANE_SENTINEL=${%1$s:+set}"
|
||||
print -r -- "PANE_DECOY=${FLEETD_TEST_DECOY_SECRET:+present}"
|
||||
FLEETD_TEST_TOOL_VAR=keep /bin/zsh <<'CHILD'
|
||||
print -r -- "CHILD_LOGIN=$([[ -o login ]] && echo yes || echo no)"
|
||||
print -r -- "CHILD_INTERACTIVE=$([[ -o interactive ]] && echo yes || echo no)"
|
||||
print -r -- "CHILD_TOOL_VAR=${FLEETD_TEST_TOOL_VAR:+present}"
|
||||
print -r -- "CHILD_DECOY=${FLEETD_TEST_DECOY_SECRET:+present}"
|
||||
print -r -- "CHILD_SENTINEL=${%1$s:+set}"
|
||||
CHILD
|
||||
exit
|
||||
""".formatted(SCRUB_SENTINEL_NAME);
|
||||
|
||||
ProcessBuilder pb = new ProcessBuilder("/bin/zsh"); // no -l, no -i: the pane's own shape
|
||||
pb.environment().clear();
|
||||
pb.environment().putAll(paneEnv);
|
||||
pb.redirectError(ProcessBuilder.Redirect.DISCARD);
|
||||
Process zsh = pb.start();
|
||||
zsh.getOutputStream().write(outerScript.getBytes(StandardCharsets.UTF_8));
|
||||
zsh.getOutputStream().flush();
|
||||
zsh.getOutputStream().close();
|
||||
String stdout = new String(zsh.getInputStream().readAllBytes(), StandardCharsets.UTF_8);
|
||||
assertTrue(zsh.waitFor(60, java.util.concurrent.TimeUnit.SECONDS),
|
||||
"the pane+child probe did not exit within 60s");
|
||||
assertTrue(zsh.exitValue() == 0, "probe zsh exited non-zero: " + stdout);
|
||||
|
||||
Map<String, String> reported = new HashMap<>();
|
||||
for (String line : stdout.split("\n")) {
|
||||
int eq = line.indexOf('=');
|
||||
if (eq > 0) {
|
||||
reported.put(line.substring(0, eq).trim(), line.substring(eq + 1).trim());
|
||||
}
|
||||
}
|
||||
|
||||
assertEquals("set", reported.get("PANE_SENTINEL"),
|
||||
"the pane itself is neither login nor interactive, so the .zshenv guard must have "
|
||||
+ "run the scrub and exported the sentinel");
|
||||
assertEquals("", reported.get("PANE_DECOY"),
|
||||
"the pane must blank a non-allow-listed name — invariant 1");
|
||||
assertEquals("no", reported.get("CHILD_LOGIN"), "sanity: the child must also be non-login");
|
||||
assertEquals("no", reported.get("CHILD_INTERACTIVE"), "sanity: the child must also be non-interactive");
|
||||
assertEquals("present", reported.get("CHILD_TOOL_VAR"),
|
||||
"invariant 3: a variable the pane deliberately set for its child must survive — a "
|
||||
+ "child that re-ran the scrub would have blanked it");
|
||||
assertEquals("", reported.get("CHILD_DECOY"),
|
||||
"a name already blanked by the pane must stay blanked in the child, never resurrected");
|
||||
assertEquals("set", reported.get("CHILD_SENTINEL"),
|
||||
"the child must inherit the sentinel from the pane's environment, or it would re-scrub");
|
||||
|
||||
EnvAllowListScrub.ScrubReport report = EnvAllowListScrub.readReport(zdotdir);
|
||||
assertNotNull(report, "the pane's own scrub pass must leave a report behind");
|
||||
assertTrue(report.blanked().stream().noneMatch(n -> n.startsWith("FLEETD_TEST_TOOL_VAR")),
|
||||
"invariant 4: the receipt must still describe the PANE, not the child — a child that "
|
||||
+ "re-ran the scrub would have rewritten this file to list its own "
|
||||
+ "FLEETD_TEST_TOOL_VAR as blanked");
|
||||
}
|
||||
|
||||
/**
|
||||
* Run {@code /bin/zsh -l -i} from a clean parent and return the NAMES it has exported by prompt
|
||||
* time. With {@code zdotdir} non-null, {@code ZDOTDIR} points at a generated scrub directory, so
|
||||
|
||||
@@ -54,6 +54,40 @@ class LeadCoordLoopTest {
|
||||
assertTrue(channel.peek().isEmpty(), "and is no longer held");
|
||||
}
|
||||
|
||||
@Test
|
||||
void redeliveryOfAMessageAlreadyWrittenToThePaneIsAckedWithoutAnotherPaneWrite() {
|
||||
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "recover me"));
|
||||
var herdr = new FakeHerdr().agentStatus("idle");
|
||||
var loop = loop(channel, herdr, Map.of(LEAD_TERM, SELF));
|
||||
|
||||
loop.tick();
|
||||
channel.hold(new LeadMessage("m1", PEER, SELF, "recover me"));
|
||||
loop.tick();
|
||||
|
||||
assertEquals(1, prompts(herdr).size(), "a redelivery must not consume the lead pane twice");
|
||||
assertEquals(List.of("m1", "m1"), channel.acked(), "the redelivery still needs a fresh broker ack");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aRedeliveryIsAckedEvenWhileTheLeadIsMidTurn() {
|
||||
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "recover me"));
|
||||
var herdr = new FakeHerdr().agentStatus("idle");
|
||||
var loop = loop(channel, herdr, Map.of(LEAD_TERM, SELF));
|
||||
|
||||
loop.tick();
|
||||
channel.hold(new LeadMessage("m1", PEER, SELF, "recover me"));
|
||||
herdr.agentStatus("working");
|
||||
loop.tick();
|
||||
|
||||
assertEquals(1, prompts(herdr).size(), "the pane is still written exactly once");
|
||||
assertEquals(List.of("m1", "m1"), channel.acked(),
|
||||
"a message already written to the pane must be acked even mid-turn: the mid-turn "
|
||||
+ "gate exists to protect the pane, and this message needs no pane. Gating the ack "
|
||||
+ "on it leaves the message held on a lead that is busy most of the time, and every "
|
||||
+ "recovery redelivers it again — which is the loop this fix exists to stop");
|
||||
assertTrue(channel.peek().isEmpty(), "so it is no longer held");
|
||||
}
|
||||
|
||||
@Test
|
||||
void leavesTheMessageUnackedWhenTheLeadIsMidTurn() {
|
||||
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "hello"));
|
||||
|
||||
@@ -92,6 +92,33 @@ class LeadMailboxTest {
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void ackThrowsWhenRecoveryClearedTheHeldMessage() throws Exception {
|
||||
String to = coordId("lead-recovery-ack");
|
||||
try (LeadMailbox inbox = LeadMailbox.open(uri(), to)) {
|
||||
inbox.publish(to, new LeadMessage("recovery-ack", "lead-from", to, "in flight"));
|
||||
assertEquals(1, awaitPeek(inbox).size(), "the broker delivery must be held before recovery clears it");
|
||||
|
||||
inbox.clearHeldForRecovery();
|
||||
|
||||
assertThrows(IllegalStateException.class, () -> inbox.ack("recovery-ack"),
|
||||
"a cleared delivery has no valid tag, so ack must report that it did not reach the broker");
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void ackOfAMessageAlreadyAckedOnThisConnectionStaysQuiet() throws Exception {
|
||||
String to = coordId("lead-repeat-ack");
|
||||
try (LeadMailbox inbox = LeadMailbox.open(uri(), to)) {
|
||||
inbox.publish(to, new LeadMessage("repeat-ack", "lead-from", to, "once"));
|
||||
awaitPeek(inbox);
|
||||
|
||||
inbox.ack("repeat-ack");
|
||||
|
||||
inbox.ack("repeat-ack");
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void duplicateMsgIdIsNotDoubleQueued() throws Exception {
|
||||
String to = coordId("lead-dedup");
|
||||
|
||||
@@ -479,7 +479,8 @@ class MessageServiceTest {
|
||||
|
||||
// The worker resumes on its own (per the ask() contract) and eventually sends its real
|
||||
// fleet_reply; the async ticket must still resolve with it, not strand at PENDING.
|
||||
assertTrue(messages.reply(T, "real result"), "the worker's real reply must still be accepted");
|
||||
assertEquals(MessageService.ReplyOutcome.RESOLVED_ASYNC_TICKET, messages.reply(T, "real result"),
|
||||
"the worker's real reply must still be accepted, resolving the parked async ticket");
|
||||
} finally {
|
||||
messages.setAskTimeoutRaceHookForTest(null);
|
||||
}
|
||||
@@ -767,7 +768,9 @@ class MessageServiceTest {
|
||||
@Test
|
||||
void replyQueuesInInboxWhenNoSendIsOpen() {
|
||||
// No send is open for this session — reply should queue in the inbox.
|
||||
assertTrue(messages.reply(T, "queued-text"), "reply should succeed (queued)");
|
||||
// fleetd #365: this is the case that must read as QUEUED, not "delivered".
|
||||
assertEquals(MessageService.ReplyOutcome.QUEUED, messages.reply(T, "queued-text"),
|
||||
"reply should succeed but only as queued — nothing was waiting for it");
|
||||
|
||||
var drained = messages.drainReplies(T);
|
||||
assertEquals(1, drained.size());
|
||||
@@ -780,7 +783,9 @@ class MessageServiceTest {
|
||||
awaitUninterruptibly(T);
|
||||
|
||||
// An explicit reply resolves the open send.
|
||||
assertTrue(messages.reply(T, "send-resolved"), "reply should succeed (resolved live send)");
|
||||
// fleetd #365: this is the other case — RESOLVED_SEND, distinct from QUEUED above.
|
||||
assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND, messages.reply(T, "send-resolved"),
|
||||
"reply should succeed by resolving the live waiting send");
|
||||
|
||||
// The inbox should be empty — the reply went to the send, not the inbox.
|
||||
assertTrue(messages.drainReplies(T).isEmpty(), "no reply in the inbox");
|
||||
@@ -1094,8 +1099,11 @@ class MessageServiceTest {
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome(),
|
||||
"the primary's own bounded wait gives up before the worker finishes resuming");
|
||||
|
||||
// The worker keeps working past that window and only now calls fleet_reply.
|
||||
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
|
||||
// The worker keeps working past that window and only now calls fleet_reply. The forward
|
||||
// waiter answer() opened already timed out, so this resolves via the parked async ticket,
|
||||
// not a live send (fleetd #365).
|
||||
assertEquals(MessageService.ReplyOutcome.RESOLVED_ASYNC_TICKET,
|
||||
messages.reply(T, "PR opened: https://example/pulls/42"));
|
||||
|
||||
MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE);
|
||||
assertEquals("PR opened: https://example/pulls/42", done.reply(),
|
||||
@@ -1119,7 +1127,10 @@ class MessageServiceTest {
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome());
|
||||
|
||||
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
|
||||
// No live waiter (answer()'s own forward wait already timed out) — resolves the parked
|
||||
// async ticket instead (fleetd #365).
|
||||
assertEquals(MessageService.ReplyOutcome.RESOLVED_ASYNC_TICKET,
|
||||
messages.reply(T, "PR opened: https://example/pulls/42"));
|
||||
|
||||
// fleet_stop tears the worker's session down right after the reply landed — this must never
|
||||
// report the misleading "the worker session was released before it replied": a reply is
|
||||
@@ -1170,7 +1181,9 @@ class MessageServiceTest {
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
awaitWaiting(); // answer() opened its own forward waiter for the resumed worker turn
|
||||
|
||||
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
|
||||
// A live waiter is open (the forward wait above) — this resolves it directly (fleetd #365).
|
||||
assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND,
|
||||
messages.reply(T, "PR opened: https://example/pulls/42"));
|
||||
|
||||
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome(),
|
||||
"the lead's own answer() call must not throw because ask()'s timeout cleanup raced it");
|
||||
@@ -1221,8 +1234,9 @@ class MessageServiceTest {
|
||||
|
||||
// The worker keeps working past the timeout and only now calls fleet_reply — with no live
|
||||
// rendezvous waiter open (ask()'s timeout already closed it) and no new send() having
|
||||
// reopened one for this target.
|
||||
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
|
||||
// reopened one for this target. So it resolves the parked async ticket (fleetd #365).
|
||||
assertEquals(MessageService.ReplyOutcome.RESOLVED_ASYNC_TICKET,
|
||||
messages.reply(T, "PR opened: https://example/pulls/42"));
|
||||
|
||||
MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE);
|
||||
assertEquals("PR opened: https://example/pulls/42", done.reply(),
|
||||
@@ -1279,7 +1293,10 @@ class MessageServiceTest {
|
||||
// real reply — reproduce that interleaving directly instead of trying to win a real race.
|
||||
messages.forgetTurnForTest(turnId);
|
||||
|
||||
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
|
||||
// answer() is still waiting on its own forward waiter for the resumed turn — a live send —
|
||||
// so this resolves it directly, not the async ticket (fleetd #365).
|
||||
assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND,
|
||||
messages.reply(T, "PR opened: https://example/pulls/42"));
|
||||
|
||||
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome(),
|
||||
"the primary's own answer() call must still see the worker's real reply");
|
||||
@@ -1321,7 +1338,9 @@ class MessageServiceTest {
|
||||
|
||||
messages.setReplyOrphanTurnIdRaceHookForTest(() -> messages.forgetTurnForTest(turnId));
|
||||
try {
|
||||
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
|
||||
// No live waiter — resolves the parked async ticket (fleetd #365).
|
||||
assertEquals(MessageService.ReplyOutcome.RESOLVED_ASYNC_TICKET,
|
||||
messages.reply(T, "PR opened: https://example/pulls/42"));
|
||||
|
||||
MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE);
|
||||
assertEquals("PR opened: https://example/pulls/42", done.reply(),
|
||||
@@ -1353,7 +1372,8 @@ class MessageServiceTest {
|
||||
injectDelivery();
|
||||
assertEquals(MessageService.AskOutcome.TIMED_OUT, messages.ask(T, "Q2?", 200).outcome());
|
||||
|
||||
assertTrue(messages.reply(T, "which task does this answer?"));
|
||||
// Ambiguous — two candidates, so it must fall back to the inbox rather than guess (fleetd #365).
|
||||
assertEquals(MessageService.ReplyOutcome.QUEUED, messages.reply(T, "which task does this answer?"));
|
||||
|
||||
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket1).phase(),
|
||||
"an ambiguous reply must not guess ticket1");
|
||||
@@ -1782,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");
|
||||
|
||||
@@ -1897,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");
|
||||
@@ -1926,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));
|
||||
@@ -1936,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;
|
||||
@@ -1970,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
|
||||
@@ -2032,7 +2180,7 @@ class MessageServiceTest {
|
||||
CompletableFuture<MessageService.Reply> send = sendAsync();
|
||||
awaitWaiting();
|
||||
|
||||
assertTrue(messages.reply(T, "resolved-live"));
|
||||
assertEquals(MessageService.ReplyOutcome.RESOLVED_SEND, messages.reply(T, "resolved-live"));
|
||||
assertFalse(messages.hasStrandedReply(T), "a reply that resolved an open send is not stranded");
|
||||
|
||||
MessageService.Reply r = send.get(5, TimeUnit.SECONDS);
|
||||
@@ -2042,14 +2190,14 @@ class MessageServiceTest {
|
||||
@Test
|
||||
void hasStrandedReplyIsTrueWhenNoSendWasWaiting() {
|
||||
// No send is open for T — the reply queues into the inbox and is recorded as stranded.
|
||||
assertTrue(messages.reply(T, "nobody was waiting"));
|
||||
assertEquals(MessageService.ReplyOutcome.QUEUED, messages.reply(T, "nobody was waiting"));
|
||||
assertTrue(messages.hasStrandedReply(T),
|
||||
"a reply with no open send strands, even though it is safely queued in the inbox");
|
||||
}
|
||||
|
||||
@Test
|
||||
void hasStrandedReplyClearsOnceTheTargetsNextDeliveryIsAccepted() throws Exception {
|
||||
assertTrue(messages.reply(T, "stray"));
|
||||
assertEquals(MessageService.ReplyOutcome.QUEUED, messages.reply(T, "stray"));
|
||||
assertTrue(messages.hasStrandedReply(T));
|
||||
|
||||
// The next accepted delivery for T clears the stale stranding fact — the one case the
|
||||
@@ -2067,7 +2215,7 @@ class MessageServiceTest {
|
||||
|
||||
@Test
|
||||
void hasStrandedReplyClearsOnAbandon() {
|
||||
assertTrue(messages.reply(T, "stray"));
|
||||
assertEquals(MessageService.ReplyOutcome.QUEUED, messages.reply(T, "stray"));
|
||||
assertTrue(messages.hasStrandedReply(T));
|
||||
|
||||
messages.abandon(T, "session released");
|
||||
|
||||
@@ -984,7 +984,7 @@ class ReplyPushLoopTest {
|
||||
// --- metrics (CB-512) ----------------------------------------------------------------------
|
||||
|
||||
@Test
|
||||
void successfulNudgeIncrementsDelivered() throws Exception {
|
||||
void successfulNudgeIncrementsSent() throws Exception {
|
||||
var rec = recordingClient();
|
||||
agents = new AgentControl(rec);
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
@@ -994,11 +994,13 @@ class ReplyPushLoopTest {
|
||||
|
||||
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS),
|
||||
"one nudge (1 agent.prompt call) should have been sent");
|
||||
// The delivered count is bumped on the scheduler thread right after the send that releases
|
||||
// The sent count is bumped on the scheduler thread right after the send that releases
|
||||
// the latch — settle briefly so the counter is published before we read it.
|
||||
Thread.sleep(200);
|
||||
assertEquals(1, metrics.count(FleetMetrics.PUSH_NUDGES, "outcome", "delivered"),
|
||||
"a successfully sent nudge must count as delivered");
|
||||
// fleetd #365: "sent", not "delivered" — this only proves the herdr call succeeded, not
|
||||
// that the primary's pane read it.
|
||||
assertEquals(1, metrics.count(FleetMetrics.PUSH_NUDGES, "outcome", "sent"),
|
||||
"a successfully sent nudge must count as sent");
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -1013,11 +1015,11 @@ class ReplyPushLoopTest {
|
||||
|
||||
assertEquals(1, metrics.count(FleetMetrics.PUSH_NUDGES, "outcome", "exhausted"),
|
||||
"hitting the reminder cap must count as exhausted");
|
||||
assertEquals(0, metrics.count(FleetMetrics.PUSH_NUDGES, "outcome", "delivered"));
|
||||
assertEquals(0, metrics.count(FleetMetrics.PUSH_NUDGES, "outcome", "sent"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void successfulTicketNudgeIncrementsDelivered() throws Exception {
|
||||
void successfulTicketNudgeIncrementsSent() throws Exception {
|
||||
var rec = recordingClient();
|
||||
agents = new AgentControl(rec);
|
||||
Metrics metrics = new Metrics();
|
||||
@@ -1026,8 +1028,8 @@ class ReplyPushLoopTest {
|
||||
|
||||
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS), "one ticket nudge should have been sent");
|
||||
Thread.sleep(200);
|
||||
assertEquals(1, metrics.count(FleetMetrics.PUSH_NUDGES, "outcome", "delivered"),
|
||||
"a successfully sent ticket nudge must count as delivered, same metric as CB-307");
|
||||
assertEquals(1, metrics.count(FleetMetrics.PUSH_NUDGES, "outcome", "sent"),
|
||||
"a successfully sent ticket nudge must count as sent, same metric as CB-307");
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
+144
@@ -0,0 +1,144 @@
|
||||
package dev.ltms.fleet.peer;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.lang.reflect.Constructor;
|
||||
import java.lang.reflect.RecordComponent;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Arrays;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Locale;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
import java.util.Set;
|
||||
import java.util.TreeSet;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
|
||||
/**
|
||||
* Fleetd #382, the same "defect factory" #357 and #358 guarded on {@code FleetConfig.withDefaults()}
|
||||
* and {@link dev.ltms.fleet.session.MemberSession}'s rebuild sites — reproduced here on
|
||||
* {@link SpawnRequest#withProfile(String)}.
|
||||
*
|
||||
* <p>{@code SpawnRequest} has back-compat constructors at arity 3 and 5 alongside its canonical
|
||||
* arity-6 constructor. Before this ticket, {@code CompositePeerLauncher} routed a profile by
|
||||
* building a fresh {@code SpawnRequest} from a literal {@code new SpawnRequest(...)} call listing
|
||||
* six of the original request's own accessors. That call is only correct because it happens to
|
||||
* name exactly six arguments today — add a 7th component and the established back-compat pattern
|
||||
* (a new constructor at the old, now-shorter arity) and a call one argument short of the new
|
||||
* canonical arity would silently rebind to that back-compat constructor, dropping the new
|
||||
* component on every profile-routed spawn without any compile error. {@link #withProfile} replaces
|
||||
* that literal call, so this test guards the ONE rebuild site instead of a call site scattered
|
||||
* through a launcher.
|
||||
*
|
||||
* <p>The check below builds a {@link SpawnRequest} through the TRUE canonical constructor —
|
||||
* resolved by the record's own component types via {@code getDeclaredConstructor}, never by
|
||||
* argument count, so it can never itself land on a back-compat overload — with a real, distinctive,
|
||||
* non-null value in every component, calls {@link SpawnRequest#withProfile(String)}, and asserts
|
||||
* every component the method is not documented to change survives unchanged, while {@code
|
||||
* profileName} comes back as the new value it was given. A component that comes back anything else
|
||||
* was silently dropped — the shape of the defect this test exists to catch.
|
||||
*
|
||||
* <p>{@link #EXCLUDED_FROM_SURVIVAL_CHECK} is kept deliberately empty and size-pinned by
|
||||
* {@link #exclusionListSizeIsPinned()}, for the same reason the other two guards pin theirs at
|
||||
* zero: a checker whose escape hatch can grow to silence a failure is not a checker. Every one of
|
||||
* {@link SpawnRequest}'s 6 current components has a real, non-null value here and none is excluded.
|
||||
*/
|
||||
class SpawnRequestWithProfilePreservesEveryComponentTest {
|
||||
|
||||
private static final RecordComponent[] COMPONENTS = SpawnRequest.class.getRecordComponents();
|
||||
|
||||
/** Deliberately empty today; grow it only with a matching justification, and re-pin the size. */
|
||||
private static final Set<String> EXCLUDED_FROM_SURVIVAL_CHECK = Set.of();
|
||||
|
||||
/** One real, distinctive, non-null value per component — none of the 6 is excluded. */
|
||||
private static Map<String, Object> baseValues() {
|
||||
Map<String, Object> v = new LinkedHashMap<>();
|
||||
v.put("profileName", "profile-guard");
|
||||
v.put("requestedCwd", "/wt/requested-guard");
|
||||
v.put("callerCwd", "/wt/caller-guard");
|
||||
v.put("sessionName", "session-guard");
|
||||
v.put("resumeSessionId", "resume-guard");
|
||||
v.put("role", MemberRole.REVIEWER);
|
||||
assertNamesMatchComponents(v);
|
||||
return v;
|
||||
}
|
||||
|
||||
/**
|
||||
* Guards {@link #baseValues()} itself against drifting from the record's real shape — forgetting
|
||||
* to add a new component here fails this assertion by name, rather than silently checking one
|
||||
* component fewer than the record has.
|
||||
*/
|
||||
private static void assertNamesMatchComponents(Map<String, Object> values) {
|
||||
Set<String> names = new TreeSet<>();
|
||||
for (RecordComponent rc : COMPONENTS) {
|
||||
names.add(rc.getName());
|
||||
}
|
||||
assertEquals(names, new TreeSet<>(values.keySet()),
|
||||
"this test's value map has drifted from SpawnRequest's actual components — "
|
||||
+ "update baseValues() alongside the record");
|
||||
}
|
||||
|
||||
/**
|
||||
* Builds a {@link SpawnRequest} through the TRUE canonical constructor — resolved by the
|
||||
* record's own component types, not by argument count — so this never accidentally exercises a
|
||||
* back-compat overload the way a literal {@code new SpawnRequest(...)} call risks doing.
|
||||
*/
|
||||
private static SpawnRequest requestOf(Map<String, Object> values) throws ReflectiveOperationException {
|
||||
Class<?>[] types = Arrays.stream(COMPONENTS).map(RecordComponent::getType).toArray(Class<?>[]::new);
|
||||
Object[] args = Arrays.stream(COMPONENTS).map(rc -> values.get(rc.getName())).toArray();
|
||||
Constructor<SpawnRequest> ctor = SpawnRequest.class.getDeclaredConstructor(types);
|
||||
return ctor.newInstance(args);
|
||||
}
|
||||
|
||||
@Test
|
||||
void exclusionListSizeIsPinned() {
|
||||
assertEquals(0, EXCLUDED_FROM_SURVIVAL_CHECK.size(),
|
||||
"EXCLUDED_FROM_SURVIVAL_CHECK grew from 0 — every entry needs a justification in "
|
||||
+ "this test class's javadoc AND this assertion re-pinned to the new size; a "
|
||||
+ "growing exclusion list that silences failures on its own is not a guard");
|
||||
}
|
||||
|
||||
@Test
|
||||
void withProfilePreservesEveryOtherComponent() throws ReflectiveOperationException {
|
||||
Map<String, Object> base = baseValues();
|
||||
SpawnRequest request = requestOf(base);
|
||||
SpawnRequest result = request.withProfile("profile-updated");
|
||||
|
||||
Map<String, Object> expectedOverrides = Map.of("profileName", "profile-updated");
|
||||
|
||||
List<String> dropped = new ArrayList<>();
|
||||
int checked = 0;
|
||||
for (RecordComponent rc : COMPONENTS) {
|
||||
String name = rc.getName();
|
||||
if (EXCLUDED_FROM_SURVIVAL_CHECK.contains(name)) {
|
||||
continue;
|
||||
}
|
||||
checked++;
|
||||
Object expected = expectedOverrides.containsKey(name) ? expectedOverrides.get(name) : base.get(name);
|
||||
Object actual;
|
||||
try {
|
||||
actual = rc.getAccessor().invoke(result);
|
||||
} catch (ReflectiveOperationException e) {
|
||||
throw new RuntimeException("failed to read SpawnRequest." + name + "()", e);
|
||||
}
|
||||
if (!Objects.equals(expected, actual)) {
|
||||
dropped.add(String.format(Locale.ROOT,
|
||||
"%s: withProfile() was expected to carry (%s) for '%s' but returned %s — a "
|
||||
+ "component silently dropped by withProfile(), the shape of the defect "
|
||||
+ "this test exists to catch (its final \"return new SpawnRequest(...)\" "
|
||||
+ "call binding to a back-compat constructor instead of the true "
|
||||
+ "canonical one)",
|
||||
name, expected, name, actual));
|
||||
}
|
||||
}
|
||||
|
||||
System.out.printf(Locale.ROOT,
|
||||
"SpawnRequest.withProfile() component-survival coverage — %d components, %d checked, "
|
||||
+ "%d excluded, %d survived%n",
|
||||
COMPONENTS.length, checked, EXCLUDED_FROM_SURVIVAL_CHECK.size(), checked - dropped.size());
|
||||
assertEquals(List.of(), dropped,
|
||||
"withProfile() silently dropped these components: " + dropped);
|
||||
}
|
||||
}
|
||||
@@ -476,6 +476,11 @@ class FleetAppTest {
|
||||
Thread.sleep(200);
|
||||
HttpResponse<String> reply = postJson(port, "/sessions/term_a/reply", "{\"content\":\"LGTM ship it\"}");
|
||||
assertEquals(200, reply.statusCode());
|
||||
// fleetd #365: "delivered" used to be unconditionally true; a live send was actually waiting
|
||||
// here, so this is the case where it must genuinely read true, with outcome naming why.
|
||||
JsonNode replyBody = mapper.readTree(reply.body());
|
||||
assertEquals(true, replyBody.get("delivered").asBoolean());
|
||||
assertEquals("resolved_send", replyBody.get("outcome").asText());
|
||||
|
||||
HttpResponse<String> res = send.get(6, java.util.concurrent.TimeUnit.SECONDS);
|
||||
assertEquals(200, res.statusCode());
|
||||
@@ -527,6 +532,11 @@ class FleetAppTest {
|
||||
int port = startHealthy();
|
||||
HttpResponse<String> res = postJson(port, "/sessions/term_a/reply", "{\"content\":\"orphan\"}");
|
||||
assertEquals(200, res.statusCode());
|
||||
// fleetd #365: nothing was waiting, so "delivered" must now read false, not the old
|
||||
// unconditional true — outcome names this as queued.
|
||||
JsonNode resBody = mapper.readTree(res.body());
|
||||
assertEquals(false, resBody.get("delivered").asBoolean());
|
||||
assertEquals("queued", resBody.get("outcome").asText());
|
||||
|
||||
// The queued reply is drainable.
|
||||
HttpResponse<String> drain = req(port, "GET", "/sessions/term_a/replies");
|
||||
|
||||
@@ -1532,26 +1532,52 @@ class GitWorktreesTest {
|
||||
* for repo setup.
|
||||
*
|
||||
* <p>Scope, measured on the fleetd #369 merge and narrower than an earlier version of this
|
||||
* comment claimed: this protects the 5 {@link #seedingGitWorktrees} sites plus — through
|
||||
* comment claimed: this protects the {@link #seedingGitWorktrees} call sites plus — through
|
||||
* {@link #gitProcessBuilder} — every {@code git} subprocess the TEST itself starts. It does
|
||||
* NOT cover the other 53 {@code new GitWorktrees(...)} constructions in this file, which pass
|
||||
* NOT cover the {@code new GitWorktrees(...)} constructions elsewhere in this file that pass
|
||||
* no env override, so a production instance built that way still inherits the JVM's real
|
||||
* environment. Stripping this override from {@code seedingGitWorktrees} leaves the class green
|
||||
* both with and without the poison command above, so that half is currently unpinned.
|
||||
* environment. (Re-measured for fleetd #373, on this file as it stands here: 4 call sites go
|
||||
* through {@link #seedingGitWorktrees(Path, String, Path)} — not 5, an earlier count this
|
||||
* comment and fleetd #373's own ticket text both repeated without re-running it — out of 59
|
||||
* total {@code new GitWorktrees(...)} occurrences, one of which is the shared construction
|
||||
* inside {@link #seedingGitWorktrees(Path, String, Map)} itself. This class-wide count moves
|
||||
* every time a test is added, so treat any number here as a snapshot, not a fact to cite
|
||||
* without recounting.) Stripping the {@code gitEnv} override from a {@link
|
||||
* #seedingGitWorktrees} call site leaves the class green both with and without the poison
|
||||
* command above for that call site's OWN test, so that half was unpinned until fleetd #373
|
||||
* added {@link #seedingGitWorktreesResolvesTheExcludesFileFallbackInsideItsThrowawayDirectory}
|
||||
* below, which asserts the property directly instead of relying on a poisoned real machine.
|
||||
*/
|
||||
private static Map<String, String> hermeticGitEnv(Path tmp) {
|
||||
return hermeticGitEnvAt(tmp.resolve("hermetic-xdg-config-home-" + System.nanoTime()));
|
||||
}
|
||||
|
||||
/** Same isolation as {@link #hermeticGitEnv(Path)}, with an explicit {@code XDG_CONFIG_HOME}
|
||||
* instead of a fresh nanoTime-unique one under {@code tmp} — used by
|
||||
* {@link #seedingGitWorktreesResolvesTheExcludesFileFallbackInsideItsThrowawayDirectory}
|
||||
* (fleetd #373) so it can pre-populate that directory with a marker BEFORE the production
|
||||
* {@link GitWorktrees} instance reads it, something the random per-call name from
|
||||
* {@link #hermeticGitEnv(Path)} makes impossible to predict from outside. */
|
||||
private static Map<String, String> hermeticGitEnvAt(Path xdgConfigHome) {
|
||||
return Map.of(
|
||||
"GIT_CONFIG_GLOBAL", "/dev/null",
|
||||
"GIT_CONFIG_SYSTEM", "/dev/null",
|
||||
"GIT_TERMINAL_PROMPT", "0",
|
||||
"XDG_CONFIG_HOME", tmp.resolve("hermetic-xdg-config-home-" + System.nanoTime()).toString());
|
||||
"XDG_CONFIG_HOME", xdgConfigHome.toString());
|
||||
}
|
||||
|
||||
/** {@link GitWorktrees}'s full test seam, with a {@code memberSkillsSource} and no other
|
||||
* overrides — the shape every seeding test below needs, isolated via {@link #hermeticGitEnv}. */
|
||||
private static GitWorktrees seedingGitWorktrees(Path root, String memberSkillsSource, Path tmp) {
|
||||
return new GitWorktrees(root.toString(), null, _ -> {}, null, null, memberSkillsSource,
|
||||
hermeticGitEnv(tmp));
|
||||
return seedingGitWorktrees(root, memberSkillsSource, hermeticGitEnv(tmp));
|
||||
}
|
||||
|
||||
/** Same shape as {@link #seedingGitWorktrees(Path, String, Path)}, taking an already-built
|
||||
* {@code gitEnv} directly rather than computing one via {@link #hermeticGitEnv(Path)} — lets
|
||||
* fleetd #373's test drive the exact production construction a real member spawn uses, with a
|
||||
* {@code gitEnv} it has already pre-populated a marker into. */
|
||||
private static GitWorktrees seedingGitWorktrees(Path root, String memberSkillsSource, Map<String, String> gitEnv) {
|
||||
return new GitWorktrees(root.toString(), null, _ -> {}, null, null, memberSkillsSource, gitEnv);
|
||||
}
|
||||
|
||||
/** Acceptance criterion 2 (part 1): a worktree with no {@code .claude/} at all gets the skill
|
||||
@@ -1784,6 +1810,63 @@ class GitWorktreesTest {
|
||||
+ "after skill seeding ran — got:\n" + porcelain);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #373. Pins the production seam that fleetd #362 review finding 2 protects: {@link
|
||||
* GitWorktrees#previouslyEffectiveExcludesFileContent}'s XDG-fallback branch reads {@code
|
||||
* XDG_CONFIG_HOME}/{@code HOME} straight in Java, not through a {@code git} subprocess, so
|
||||
* {@code gitEnv} — the constructor seam every {@link #seedingGitWorktrees} instance in this
|
||||
* class is built with — is the ONLY thing that can isolate it. A mutation run during the
|
||||
* fleetd #372/#369 merge found this unpinned: replacing {@code hermeticGitEnv(tmp)} with
|
||||
* {@code null} in {@link #seedingGitWorktrees(Path, String, Path)} left every test in this
|
||||
* class green — the 56 tests that would fail against a real machine's poisoned {@code
|
||||
* XDG_CONFIG_HOME} were fixed by fleetd #369's subprocess-level isolation, but none of them
|
||||
* looks at what THIS Java-side read resolves, so deleting the override stays invisible.
|
||||
*
|
||||
* <p>This test asserts the PROPERTY, not the constructor argument: a {@link GitWorktrees}
|
||||
* built for seeding — through the very same {@link #seedingGitWorktrees(Path, String, Map)}
|
||||
* construction every other seeding test in this class goes through — must resolve the
|
||||
* excludes-file fallback inside its own throwaway {@code gitEnv}-supplied directory. It needs
|
||||
* NO externally-set poisoned environment variable: the marker pattern below is written ONLY
|
||||
* inside a throwaway {@code XDG_CONFIG_HOME} this test controls directly (bypassing {@link
|
||||
* #hermeticGitEnv(Path)}'s unpredictable nanoTime-named directory, via {@link
|
||||
* #hermeticGitEnvAt}, so the marker can be in place before the production instance ever reads
|
||||
* it), reachable ONLY through the {@code gitEnv} seam. If that seam is stripped, the
|
||||
* production code instead falls back to resolving the REAL {@code XDG_CONFIG_HOME}/{@code
|
||||
* HOME} of the machine running the test — which does not carry this marker — so the marker
|
||||
* file below shows up as untracked and the assertion fails on any machine, with no poison
|
||||
* command required. See the PR body for the pasted failure from actually running that
|
||||
* mutation (removing the {@code gitEnv} override from this test's own construction).
|
||||
*/
|
||||
@Test
|
||||
void seedingGitWorktreesResolvesTheExcludesFileFallbackInsideItsThrowawayDirectory(@TempDir Path tmp)
|
||||
throws Exception {
|
||||
Path xdgConfigHome = tmp.resolve("cb373-xdg-config-home");
|
||||
Files.createDirectories(xdgConfigHome.resolve("git"));
|
||||
Files.writeString(xdgConfigHome.resolve("git").resolve("ignore"), "cb373-xdg-fallback-marker\n");
|
||||
Map<String, String> gitEnv = hermeticGitEnvAt(xdgConfigHome);
|
||||
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
Path skillsSource = tmp.resolve("skills-src");
|
||||
writeSkill(skillsSource, "implementer", "IMPLEMENTER SKILL\n");
|
||||
GitWorktrees seeding = seedingGitWorktrees(tmp.resolve("wts"), skillsSource.toString(), gitEnv);
|
||||
|
||||
String wt = seeding.add(repo.toString(), "cb-373-xdg-seam", "HEAD");
|
||||
assertEquals("IMPLEMENTER SKILL\n",
|
||||
Files.readString(Path.of(wt, ".claude", "skills", "implementer", "SKILL.md")),
|
||||
"fixture check — the skill really was seeded, so previouslyEffectiveExcludesFileContent ran");
|
||||
|
||||
Files.writeString(Path.of(wt, "cb373-xdg-fallback-marker"),
|
||||
"would only be invisible to git status if the fallback resolved THIS throwaway "
|
||||
+ "XDG_CONFIG_HOME rather than the real machine's\n");
|
||||
|
||||
String porcelain = fullStatus(Path.of(wt));
|
||||
assertEquals("", porcelain,
|
||||
"the marker pattern lives only in this test's throwaway XDG_CONFIG_HOME; git "
|
||||
+ "status must still be empty, proving the production seam resolved the "
|
||||
+ "excludes-file fallback through the gitEnv seam rather than the JVM's "
|
||||
+ "real environment — got:\n" + porcelain);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #369, acceptance criterion 4 — make the fix hard to undo by accident. Every git
|
||||
* subprocess this class starts is required to go through {@link #gitProcessBuilder}, the one
|
||||
|
||||
+190
@@ -0,0 +1,190 @@
|
||||
package dev.ltms.fleet.session;
|
||||
|
||||
import dev.ltms.fleet.peer.CharterReceipt;
|
||||
import dev.ltms.fleet.peer.MemberRole;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.lang.reflect.Constructor;
|
||||
import java.lang.reflect.RecordComponent;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Arrays;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Locale;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
import java.util.Set;
|
||||
import java.util.TreeSet;
|
||||
import java.util.function.Function;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
|
||||
/**
|
||||
* Fleetd #358, the same "defect factory" #357 guarded on {@code FleetConfig.withDefaults()}
|
||||
* (see {@code FleetConfigWithDefaultsPreservesEveryComponentTest}), reproduced here on
|
||||
* {@link MemberSession} — the worse of the two sibling cases named in #358, because this record
|
||||
* has FIVE independent rebuild sites instead of one: {@link MemberSession#withState},
|
||||
* {@link MemberSession#withActivity}, {@link MemberSession#bumpTurn},
|
||||
* {@link MemberSession#withAgentSessionId} and {@link MemberSession#withFailureReason} each end in
|
||||
* their own literal {@code new MemberSession(...)} call. Add a 16th component and add the
|
||||
* established back-compat constructor at the old (15-arg) arity, and every one of those five
|
||||
* literal calls becomes a legal match for that new overload — silently dropping the new component,
|
||||
* independently, on whichever of the five paths a missed update leaves behind. That is harder to
|
||||
* spot than #357's single call site: the field would survive through some transitions and vanish
|
||||
* through others.
|
||||
*
|
||||
* <p>Each check below builds one {@link MemberSession} through the TRUE canonical constructor —
|
||||
* resolved by the record's own component types via {@code getDeclaredConstructor}, never by
|
||||
* argument count, so it can never itself land on a back-compat overload — with a real, distinctive,
|
||||
* non-null value in every component, calls the real rebuild method under test, and asserts every
|
||||
* component the method is not documented to change survives unchanged, while the component(s) it IS
|
||||
* documented to change come back as the new value it was given. A component that comes back
|
||||
* anything else was silently dropped or lost — the shape of the defect this test exists to catch.
|
||||
*
|
||||
* <p>{@link #EXCLUDED_FROM_SURVIVAL_CHECK} is kept deliberately empty and size-pinned by
|
||||
* {@link #exclusionListSizeIsPinned()}, for the same reason {@code FleetConfig}'s guard pins its own
|
||||
* exclusion list at zero: a checker whose escape hatch can grow to silence a failure is not a
|
||||
* checker. Every one of {@link MemberSession}'s 15 current components has a real, non-null,
|
||||
* non-blank value here and none is excluded.
|
||||
*/
|
||||
class MemberSessionRebuildPreservesEveryComponentTest {
|
||||
|
||||
private static final RecordComponent[] COMPONENTS = MemberSession.class.getRecordComponents();
|
||||
|
||||
/** Deliberately empty today; grow it only with a matching justification, and re-pin the size. */
|
||||
private static final Set<String> EXCLUDED_FROM_SURVIVAL_CHECK = Set.of();
|
||||
|
||||
/** One real, distinctive, non-null value per component — none of the 15 is excluded. */
|
||||
private static Map<String, Object> baseValues() {
|
||||
Map<String, Object> v = new LinkedHashMap<>();
|
||||
v.put("paneId", "pane-guard");
|
||||
v.put("terminalId", "term-guard");
|
||||
v.put("profile", "profile-guard");
|
||||
v.put("role", MemberRole.REVIEWER);
|
||||
v.put("cwd", "/wt/guard");
|
||||
v.put("ownerTerminal", "owner-guard");
|
||||
v.put("spawnedAtNanos", 111_111L);
|
||||
v.put("lastActivityAtNanos", 222_222L);
|
||||
v.put("turnCount", 7);
|
||||
v.put("state", MemberSession.State.BUSY);
|
||||
v.put("worktree", "/wt/guard-tree");
|
||||
v.put("branch", "worker/guard-branch");
|
||||
v.put("charterReceipt", new CharterReceipt(
|
||||
MemberRole.DEV, "profile-guard", "fleet.charters.dev", "deadbeefguard", 42));
|
||||
v.put("agentSessionId", "agent-guard");
|
||||
v.put("failureReason", "reason-guard");
|
||||
assertNamesMatchComponents(v);
|
||||
return v;
|
||||
}
|
||||
|
||||
/**
|
||||
* Guards {@link #baseValues()} itself against drifting from the record's real shape — forgetting
|
||||
* to add a new component here fails this assertion by name, rather than silently checking one
|
||||
* component fewer than the record has.
|
||||
*/
|
||||
private static void assertNamesMatchComponents(Map<String, Object> values) {
|
||||
Set<String> names = new TreeSet<>();
|
||||
for (RecordComponent rc : COMPONENTS) {
|
||||
names.add(rc.getName());
|
||||
}
|
||||
assertEquals(names, new TreeSet<>(values.keySet()),
|
||||
"this test's value map has drifted from MemberSession's actual components — "
|
||||
+ "update baseValues() alongside the record");
|
||||
}
|
||||
|
||||
/**
|
||||
* Builds a {@link MemberSession} through the TRUE canonical constructor — resolved by the
|
||||
* record's own component types, not by argument count — so this never accidentally exercises a
|
||||
* back-compat overload the way a literal {@code new MemberSession(...)} call risks doing.
|
||||
*/
|
||||
private static MemberSession sessionOf(Map<String, Object> values) throws ReflectiveOperationException {
|
||||
Class<?>[] types = Arrays.stream(COMPONENTS).map(RecordComponent::getType).toArray(Class<?>[]::new);
|
||||
Object[] args = Arrays.stream(COMPONENTS).map(rc -> values.get(rc.getName())).toArray();
|
||||
Constructor<MemberSession> ctor = MemberSession.class.getDeclaredConstructor(types);
|
||||
return ctor.newInstance(args);
|
||||
}
|
||||
|
||||
@Test
|
||||
void exclusionListSizeIsPinned() {
|
||||
assertEquals(0, EXCLUDED_FROM_SURVIVAL_CHECK.size(),
|
||||
"EXCLUDED_FROM_SURVIVAL_CHECK grew from 0 — every entry needs a justification in "
|
||||
+ "this test class's javadoc AND this assertion re-pinned to the new size; a "
|
||||
+ "growing exclusion list that silences failures on its own is not a guard");
|
||||
}
|
||||
|
||||
/**
|
||||
* Shared check for one rebuild site: build a base session with a real value in every component,
|
||||
* call {@code rebuild}, and assert every component comes back equal to {@code expectedOverrides}
|
||||
* when named there, or equal to the base value otherwise. Prints the same denominator style as
|
||||
* {@code FleetConfigWithDefaultsPreservesEveryComponentTest}.
|
||||
*/
|
||||
private void checkRebuildSite(String siteName, Function<MemberSession, MemberSession> rebuild,
|
||||
Map<String, Object> expectedOverrides) throws ReflectiveOperationException {
|
||||
Map<String, Object> base = baseValues();
|
||||
MemberSession session = sessionOf(base);
|
||||
MemberSession result = rebuild.apply(session);
|
||||
|
||||
List<String> dropped = new ArrayList<>();
|
||||
int checked = 0;
|
||||
for (RecordComponent rc : COMPONENTS) {
|
||||
String name = rc.getName();
|
||||
if (EXCLUDED_FROM_SURVIVAL_CHECK.contains(name)) {
|
||||
continue;
|
||||
}
|
||||
checked++;
|
||||
Object expected = expectedOverrides.containsKey(name) ? expectedOverrides.get(name) : base.get(name);
|
||||
Object actual;
|
||||
try {
|
||||
actual = rc.getAccessor().invoke(result);
|
||||
} catch (ReflectiveOperationException e) {
|
||||
throw new RuntimeException("failed to read MemberSession." + name + "()", e);
|
||||
}
|
||||
if (!Objects.equals(expected, actual)) {
|
||||
dropped.add(String.format(Locale.ROOT,
|
||||
"%s: %s() was expected to carry (%s) for '%s' but returned %s — a component "
|
||||
+ "silently dropped by %s(), the shape of the defect this test exists "
|
||||
+ "to catch (its final \"return new MemberSession(...)\" call binding "
|
||||
+ "to a back-compat constructor instead of the true canonical one)",
|
||||
name, siteName, expected, name, actual, siteName));
|
||||
}
|
||||
}
|
||||
|
||||
System.out.printf(Locale.ROOT,
|
||||
"MemberSession.%s() component-survival coverage — %d components, %d checked, %d "
|
||||
+ "excluded, %d survived%n",
|
||||
siteName, COMPONENTS.length, checked, EXCLUDED_FROM_SURVIVAL_CHECK.size(),
|
||||
checked - dropped.size());
|
||||
assertEquals(List.of(), dropped,
|
||||
siteName + "() silently dropped these components: " + dropped);
|
||||
}
|
||||
|
||||
@Test
|
||||
void withStatePreservesEveryOtherComponent() throws ReflectiveOperationException {
|
||||
checkRebuildSite("withState", s -> s.withState(MemberSession.State.FAILED),
|
||||
Map.of("state", MemberSession.State.FAILED));
|
||||
}
|
||||
|
||||
@Test
|
||||
void withActivityPreservesEveryOtherComponent() throws ReflectiveOperationException {
|
||||
checkRebuildSite("withActivity", s -> s.withActivity(999_999L),
|
||||
Map.of("lastActivityAtNanos", 999_999L));
|
||||
}
|
||||
|
||||
@Test
|
||||
void bumpTurnPreservesEveryOtherComponent() throws ReflectiveOperationException {
|
||||
checkRebuildSite("bumpTurn", s -> s.bumpTurn(999_999L),
|
||||
Map.of("lastActivityAtNanos", 999_999L, "turnCount", 8));
|
||||
}
|
||||
|
||||
@Test
|
||||
void withAgentSessionIdPreservesEveryOtherComponent() throws ReflectiveOperationException {
|
||||
checkRebuildSite("withAgentSessionId", s -> s.withAgentSessionId("agent-updated"),
|
||||
Map.of("agentSessionId", "agent-updated"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void withFailureReasonPreservesEveryOtherComponent() throws ReflectiveOperationException {
|
||||
checkRebuildSite("withFailureReason", s -> s.withFailureReason("reason-updated"),
|
||||
Map.of("failureReason", "reason-updated"));
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user