Compare commits

..

1 Commits

Author SHA1 Message Date
Dai Ha 130962e8e3 fleetd must not let the host idle-sleep while members are live
CI / contract (pull_request) Successful in 54s
CI / build (pull_request) Successful in 1m54s
Adds a small IdleSleepGuard (dev.ltms.fleet.power) that holds a macOS
caffeinate -i child while at least one fleet member is live, and
releases it once none are. It hangs off SessionManager's existing
onAcquire/onRelease hooks and SessionManager#size() rather than
tracking members a second way. New idleSleepGuard: config block,
on by default, following the FleetConfig.Health/ConfigReload pattern.
2026-09-05 05:37:38 +07:00
43 changed files with 846 additions and 3391 deletions
+4 -4
View File
@@ -1,15 +1,15 @@
{
"name": "fleetd",
"name": "claude-bridge",
"description": "Tooling for orchestrating a fleet of delegated coding agents through the fleetd MCP gateway.",
"owner": {
"name": "LTMS"
},
"plugins": [
{
"name": "fleet",
"name": "claude-bridge",
"source": "./plugin",
"description": "Mount the fleetd MCP gateway and apply standard Claude Code settings so a session can orchestrate delegated workers. Ships no credentials.",
"version": "0.2.0",
"description": "Make a project bridge-ready: mount the fleetd MCP gateway and apply standard Claude Code settings so a session can orchestrate delegated workers. Ships no credentials.",
"version": "0.1.0",
"author": {
"name": "LTMS"
}
-11
View File
@@ -210,17 +210,6 @@ must obey belongs in the charter, not here.
- **Primary-side skills** (not delegation playbooks — a worker cannot use them):
`port-to-opencode` (make an OpenCode session a participant in this workspace) and
`fleets-status` (report every fleet that shares one LavinMQ instance).
- **This repo is also a Claude Code marketplace, and ships a plugin.** `.claude-plugin/marketplace.json`
points at `plugin/`, which carries the MCP mount and the `setup` skill
(`/claude-bridge:setup` — make any project bridge-ready). It was added in CB-527 and then went
unmentioned by every instruction file, so it drifted and a later session planned it from scratch
(#362). **Read `plugin/` before designing anything about onboarding a project.** Two limits are
structural, not bugs: a plugin cannot carry the role agent files, because
`ClaudeCodeLauncher.java:371` requires `<cwd>/.claude/agents/<role>.md` in the member's own
worktree; and a plugin cannot deliver anything to members at all, because
`ClaudeCodeLauncher.java:285` exports `CLAUDE_CONFIG_DIR` and every Claude profile here sets it,
so a member never reads the operator's plugin store. **The plugin is the lead-side surface;
member-facing assets travel in the worktree.**
- **Never commit** `.mcp.json` (the primary's local copy, flagged `--skip-worktree`) or `wiki/`
(a submodule with its own remote).
- **A provisioned worktree neutralizes `.mcp.json`, `opencode.json` and `.autoenv`** — the repo's
+8 -20
View File
@@ -110,6 +110,14 @@ bind:
# notifications:
# mode: disabled
# Idle-sleep guard: while at least one member is live, hold an OS-level assertion against idle
# sleep (macOS only — a `caffeinate -i` child; a no-op elsewhere or if caffeinate is missing), so
# an unattended host does not idle-sleep out from under a member's long turn. Unlike health/
# configReload above, this is ON BY DEFAULT — omitting the block entirely leaves it enabled, the
# same as `enabled: true`. Uncomment only to turn it off:
# idleSleepGuard:
# enabled: false
# herdr Unix socket. Omit to use the client default
# (${HERDR_SOCKET_PATH:-~/.config/herdr/herdr.sock}).
herdrSocket: ~/.config/herdr/herdr.sock
@@ -774,19 +782,6 @@ guard:
# so fleetd falls back to the weaker CB-596 sentinel overlay instead (a WARN names the gap).
# worktreeGroup: fleet-workers
# fleetd #362: a directory of skill folders (each a subdirectory holding a SKILL.md, the same
# shape as this repo's own .claude/skills/) copied into every PROVISIONED worktree's
# .claude/skills/, so a member spawned against ANY repo — not only one that already ships its own
# copy — can load a bridge skill (e.g. implementer). Unset (the default): no worktree is touched
# beyond today's behaviour. A skill folder the target repo already carries under
# .claude/skills/<name> is never overwritten — the repo's own copy always wins. Claude Code
# members only; an opencode member reads a different path (.opencode/agent) this key does not
# touch. Best-effort like worktreeGroup above: a missing/unreadable directory here is logged and
# skipped, never a failed spawn. Every non-hidden subdirectory of this directory is copied
# wholesale, with no per-file allowlist — don't park scratch files or drafts alongside the real
# skill folders, they will be copied into every provisioned worktree too.
# memberSkills: /path/to/fleetd/checkout/.claude/skills
# Session lifecycle limits (CB-303). All knobs are opt-in; omit or set to null to keep
# the feature disabled. By default the daemon never reaps, caps, or drains sessions.
# idleTtlSeconds → reap READY/DONE sessions idle longer than this (never BUSY/SPAWNING)
@@ -834,17 +829,10 @@ guard:
# across every daemon sharing this vhost.
# prefetch → consumer basicQos, capping how many unacked messages the mailbox holds in-heap.
# Default 32 when omitted.
# peers → fleetd #361: the coord-ids of the OTHER daemons on this vhost, declared by the
# operator (the daemon never guesses). fleet_list reports each one's live reachability
# (a passive queue check, never a presence protocol) alongside this daemon's own
# mailbox state. Omit, or leave empty, for a daemon with no known peers yet — an
# undeclared peer can still reach you and be reached by fleet_send, it just will not
# show up as a row in fleet_list.
# coordinator:
# uriEnv: LEAD_COORD_URI
# selfId: mac-opus
# prefetch: 32
# peers: [fleet01-lead]
# Active push-to-primary (CB-307 Stage 3). When a worker reply lands with no open fleet_send,
# the ReplyPushLoop injects a *drain nudge* (never the payload) into the primary's own herdr
@@ -55,6 +55,8 @@ import dev.ltms.fleet.member.MemberCredentialPolicyView;
import dev.ltms.fleet.member.OpenCodeLauncher;
import dev.ltms.fleet.placement.BackendOutagePolicy;
import dev.ltms.fleet.placement.BackendQuarantine;
import dev.ltms.fleet.power.CaffeinateSleepAssertionMechanism;
import dev.ltms.fleet.power.IdleSleepGuard;
import io.javalin.Javalin;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -248,11 +250,30 @@ public final class Fleetd {
contextCap = cfg.lifecycle().contextCap();
}
boolean clearAfterTurn = cfg.lifecycle() != null && cfg.lifecycle().clearAfterTurn();
SessionManager sessions = new SessionManager(workers,
new GitWorktrees(cfg.worktreeRoot(), cfg.worktreeGroup(), cfg.memberSkills()),
SessionManager sessions = new SessionManager(workers, new GitWorktrees(cfg.worktreeRoot(), cfg.worktreeGroup()),
System::nanoTime, contextCap, clearAfterTurn);
liveCountRef.set(profileName -> liveSessionCount(sessions.roster(), profileName));
// Idle-sleep guard: hold an OS-level assertion against idle sleep while at least one
// member is live, so an unattended host does not idle-sleep out from under a member's
// long turn (see FleetConfig.IdleSleepGuard / dev.ltms.fleet.power.IdleSleepGuard for the
// measurement that motivated this). Opt-out via idleSleepGuard.enabled: false; on by
// default. Hangs off SessionManager's own onAcquire/onRelease hooks (CB-520/CB-516,
// previously wired only to the reply inbox) and SessionManager#size() — the exact registry
// fleet_list's live/capacity numbers are themselves computed from — rather than tracking
// members a second way. No-op (never constructed) off macOS or when idleSleepGuard.enabled
// is explicitly false; the mechanism itself is additionally a no-op if 'caffeinate' cannot
// be started, so this can never fail a spawn, a release, or startup.
boolean idleSleepGuardEnabled = cfg.idleSleepGuard() == null || cfg.idleSleepGuard().isEnabled();
final IdleSleepGuard idleSleepGuard;
if (idleSleepGuardEnabled) {
idleSleepGuard = new IdleSleepGuard(new CaffeinateSleepAssertionMechanism(), sessions::size);
sessions.onAcquire(_ -> idleSleepGuard.recheck());
sessions.onRelease(_ -> idleSleepGuard.recheck());
} else {
idleSleepGuard = null;
}
// CB-303 part 1: idle-ttl reaper — only when configured, defaults to disabled.
final SessionReaper reaper;
if (cfg.lifecycle() != null
@@ -653,12 +674,7 @@ public final class Fleetd {
quarantineSource,
leadMailbox,
outageSource,
new FleetMcp.LeadSeatSource(leadSeatLookup(() -> config.get().profiles(), leaders, leads)),
// fleetd #361: the operator-declared peers this daemon's fleet_list should try to
// reach. Read from the SAME snapshot leadMailbox itself opened from (cfg.coordinator()),
// not the live config.get() — coordinator wiring is already boot-time-fixed (see
// leadMailbox above), so peers follows the same rule rather than half hot-reloading.
cfg.coordinator() == null ? List.of() : cfg.coordinator().peers());
new FleetMcp.LeadSeatSource(leadSeatLookup(() -> config.get().profiles(), leaders, leads)));
// CB-637: the receive half. Only constructed when a lead mailbox actually opened — with no
// coordinator (or an unreachable one) there is nothing to deliver, so no scheduler is
@@ -704,6 +720,11 @@ public final class Fleetd {
if (configWatcher != null) configWatcher.stop(); // CB-559: stop polling the config file
mcp.close();
if (reaper != null) reaper.stop();
// Idle-sleep guard: release unconditionally, even though sessions.close() above already
// drained every session (and each release already drove the live count to 0, which
// releases the guard's assertion on its own) — this is the backstop for a drain that was
// itself interrupted or threw, so no caffeinate child ever outlives the daemon.
if (idleSleepGuard != null) idleSleepGuard.close();
// Release the broker connection last among message resources (no-op for the in-memory inbox).
if (replyInbox instanceof AutoCloseable closeable) {
try {
@@ -37,13 +37,16 @@ import java.util.function.Supplier;
* makes {@code fleet:} split rather than hot — see below.</li>
* <li><strong>Deferred</strong> — accepted into the new snapshot, but the wiring built at startup
* keeps the old value until a restart: {@code lifecycle:}, {@code leadHeartbeat:},
* {@code idleSleepGuard:} ({@code Fleetd.java} reads it once, at startup, to decide whether
* to construct an {@code IdleSleepGuard} and wire {@code SessionManager}'s
* {@code onAcquire}/{@code onRelease} hooks to it — neither is rebuilt on reload, so a
* 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 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
* {@code worktreeGroup} missing from this list and from {@link #changedDeferredKeys};
* {@code memberSkills} (fleetd #362) followed the same shape), {@code primary:} (fleetd #326 — {@code Fleetd.java:506, 519,
* {@code guard:}, {@code worktreeRoot:} and {@code worktreeGroup:} (both baked once into the
* {@code GitWorktrees} built at {@code Fleetd.java:251} and never rebuilt — fleetd #323
* instance 2 found {@code worktreeGroup} missing from this list and from
* {@link #changedDeferredKeys}), {@code primary:} (fleetd #326 — {@code Fleetd.java:506, 519,
* 520} read {@code cfg.primary()} only off the startup snapshot to build {@code
* PrimaryRegistry} and size {@code ReplyPushLoop}'s reminder cap/backoff, and neither is
* rebuilt on reload. Say the consequence exactly: {@code primary.terminal} is DEPRECATED
@@ -130,10 +133,10 @@ import java.util.function.Supplier;
* five of COLD_KEYS" rather than re-listing them, so prose and set cannot drift again.</li>
* </ul>
*
* <p><strong>The denominator, measured on 2026-09-04 (fleetd #330; recounted for fleetd #333);
* recounted again for fleetd #362.</strong> {@code FleetConfig} has 23 top-level record components:
* 5 cold, 12 deferred, 3 split, 3 hot-excluded. Three of them are named nowhere in this file, and
* the reason is the same for all
* <p><strong>The denominator, measured on 2026-09-04 (fleetd #330; recounted for fleetd #333),
* recounted again after {@code idleSleepGuard:} was added.</strong>
* {@code FleetConfig} has 23 top-level record components: 5 cold, 12 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()}
* (placement through the {@code CompositePeerLauncher} supplier the Hot bullet names;
@@ -213,9 +216,9 @@ public final class ConfigRef implements Supplier<FleetConfig> {
* read {@link #COLD_KEYS} and {@link #SPLIT_KEYS}.
*/
static final Set<String> DEFERRED_KEYS = Set.of(
"guard", "worktreeRoot", "worktreeGroup", "memberSkills", "primary", "configReload",
"guard", "worktreeRoot", "worktreeGroup", "primary", "configReload",
"leadHeartbeat", "lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs",
"quarantineCooldownSeconds", "profiles");
"quarantineCooldownSeconds", "profiles", "idleSleepGuard");
private final Path path;
private final AtomicReference<FleetConfig> current;
@@ -404,13 +407,6 @@ public final class ConfigRef implements Supplier<FleetConfig> {
if (!Objects.equals(old.worktreeGroup(), fresh.worktreeGroup())) {
changed.add("worktreeGroup");
}
// fleetd #362: baked into the same GitWorktrees as worktreeRoot/worktreeGroup
// (Fleetd.java:251) and never rebuilt either — a reload that changes only memberSkills
// must be reported the same way, or a newly provisioned worktree keeps seeding from (or
// skipping) the old source directory with nothing telling the operator why.
if (!Objects.equals(old.memberSkills(), fresh.memberSkills())) {
changed.add("memberSkills");
}
// fleetd #326: Fleetd.java:506, 519, 520 read cfg.primary() only off the startup snapshot
// (PrimaryRegistry's pinned terminal, ReplyPushLoop's reminder cap and backoff) — neither is
// rebuilt on reload, so a changed value needs a restart. Note what it does NOT mean:
@@ -427,6 +423,14 @@ public final class ConfigRef implements Supplier<FleetConfig> {
if (!Objects.equals(old.configReload(), fresh.configReload())) {
changed.add("configReload");
}
// Fleetd.java reads cfg.idleSleepGuard() once, at startup, to decide whether to construct
// an IdleSleepGuard at all and wire SessionManager's onAcquire/onRelease hooks to it —
// neither is rebuilt on reload, so a running daemon keeps whatever this was at startup
// (armed or not) regardless of a later edit here. Not cold: nothing already-open goes
// inconsistent with the new value, an armed-or-not guard just keeps its original answer.
if (!Objects.equals(old.idleSleepGuard(), fresh.idleSleepGuard())) {
changed.add("idleSleepGuard");
}
if (!Objects.equals(old.spawnReadyTimeoutMs(), fresh.spawnReadyTimeoutMs())
|| !Objects.equals(old.spawnReadyPollMs(), fresh.spawnReadyPollMs())) {
changed.add("spawnReady*");
@@ -106,19 +106,11 @@ import java.util.regex.PatternSyntaxException;
* When {@code memberHerdrSocket} is NOT configured this field is never
* consulted at all; fleetd keeps reading its own {@code $SHELL}, exactly as
* before this field existed.
* @param memberSkills fleetd #362: nullable directory of skill folders (each a subdirectory
* holding a {@code SKILL.md}, the same shape as this repo's own {@code
* .claude/skills/}) copied into every provisioned worktree's {@code
* .claude/skills/}, so a member spawned against ANY repo — not only one that
* already ships its own copy — can load a bridge skill such as {@code
* implementer}. {@code null}/blank ⇒ off: no worktree is touched beyond
* today's behaviour. A skill folder the target repo already carries is never
* overwritten — see {@link dev.ltms.fleet.session.GitWorktrees}. Claude Code
* members only; an opencode member's equivalent lives under a different path
* ({@code .opencode/agent}) and is not covered by this key. Every non-hidden
* subdirectory of this directory is copied wholesale, with no per-file
* allowlist — do not park scratch files or drafts alongside the real skill
* folders, they will be copied into every provisioned worktree too.
* @param idleSleepGuard opt-in-by-default: hold an OS-level assertion against idle sleep while at
* least one member is live, so an unattended host does not sleep out from
* 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}.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record FleetConfig(
@@ -144,9 +136,9 @@ public record FleetConfig(
Coordinator coordinator,
String worktreeGroup,
String memberLoginShell,
String memberSkills) {
IdleSleepGuard idleSleepGuard) {
/** Back-compat form before the {@code memberSkills} key was added. */
/** Back-compat form before the {@code idleSleepGuard:} 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,
@@ -169,7 +161,7 @@ public record FleetConfig(
MemberCredentials memberCredentials, Coordinator coordinator, String worktreeGroup) {
this(bind, herdrSocket, memberHerdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth,
configReload, quarantineCooldownSeconds, memberCredentials, coordinator, worktreeGroup, null);
configReload, quarantineCooldownSeconds, memberCredentials, coordinator, worktreeGroup, null, null);
}
/** Back-compat form before the {@code worktreeGroup} key was added. */
@@ -916,18 +908,12 @@ public record FleetConfig(
* {@code null} ⇒ kept as {@code null} (no self id configured).
* @param prefetch the consumer's {@code basicQos} prefetch count. {@code null}/non-positive ⇒
* {@link LeadMailbox#DEFAULT_PREFETCH}.
* @param peers fleetd #361: the coord-ids the operator declares as this daemon's peers — the
* daemon never guesses who else exists. {@code fleet_list} reports each one's
* live reachability. Blank entries are dropped; {@code null} ⇒ an empty list, so
* a config written before this field existed still parses unchanged.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record Coordinator(String uri, String uriEnv, String selfId, Integer prefetch, List<String> peers) {
public record Coordinator(String uri, String uriEnv, String selfId, Integer prefetch) {
public Coordinator {
selfId = (selfId == null || selfId.isBlank()) ? null : selfId;
peers = peers == null ? List.of()
: peers.stream().filter(p -> p != null && !p.isBlank()).toList();
}
/** True when a {@code uriEnv} is configured by name, whether or not its variable resolves. */
@@ -1279,6 +1265,25 @@ public record FleetConfig(
}
}
/**
* Hold an OS-level assertion against idle sleep while at least one member is live (see
* {@link dev.ltms.fleet.power.IdleSleepGuard}).
*
* <p>Unlike most opt-in blocks in this file, this one defaults to <em>on</em>: an unattended
* host idle-sleeping mid-turn is a correctness problem (a dropped AMQP link, a frozen member),
* not a convenience, so the safer default is armed. An operator who wants the previous
* behaviour (no assertion held, ever) sets {@code enabled: false} explicitly.
*
* @param enabled {@code false} turns the guard off; {@code null} (the block omitted
* entirely) or {@code true} leaves it on
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record IdleSleepGuard(Boolean enabled) {
public boolean isEnabled() {
return !Boolean.FALSE.equals(enabled);
}
}
/**
* The terminal → lead-name map seeded from the legacy singular {@code primary:} pin (CB-530).
*
@@ -1529,7 +1534,7 @@ public record FleetConfig(
"bind", "herdrSocket", "memberHerdrSocket", "profiles", "guard", "worktreeRoot",
"lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs", "broker", "primary", "fleet",
"leadHeartbeat", "health", "placement", "auth", "configReload", "quarantineCooldownSeconds",
"memberCredentials", "coordinator", "worktreeGroup", "memberLoginShell", "memberSkills");
"memberCredentials", "coordinator", "worktreeGroup", "memberLoginShell", "idleSleepGuard");
/** Load and validate config from {@code path}. */
public static FleetConfig load(Path path) {
@@ -2206,12 +2211,14 @@ public record FleetConfig(
// memberLoginShell is left as-is (fleetd #213), like worktreeGroup: null/blank is "not
// configured", and there is no sane non-null default — a member's login shell is
// operator-specific and only meaningful when memberHerdrSocket is also set.
// memberSkills is left as-is (fleetd #362), like worktreeGroup/memberLoginShell: null/blank
// is "off", and there is no sane non-null default — the daemon may not even run from a
// checkout that ships its own .claude/skills/.
// idleSleepGuard is left as-is, like leadHeartbeat/configReload above, but for the opposite
// reason: it is on by default already (its own isEnabled() treats null the same as
// 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.
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);
quarantineCooldown, mc, coordinator, worktreeGroup, memberLoginShell, idleSleepGuard);
}
/**
@@ -5,7 +5,6 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.Collections;
import java.util.HashSet;
import java.util.LinkedHashMap;
import java.util.Locale;
import java.util.Map;
@@ -54,40 +53,17 @@ import java.util.function.Supplier;
* daemon and found again by this scan. The trust direction above is unaffected — fleetd writing a
* name for a lead it just started is not a pane promoting itself — but <em>staleness</em> becomes
* real: a label left behind by a session that has since died would read as a live lead forever.
* This scanner does not solve that (its job is naming, and a stale name costs nothing here); the
* launcher does, by requiring a running agent in the tab before it counts the lead as live. If you
* ever make a decision that <em>removes</em> something based on this map, add the same check.
* The remaining hazard is an <em>operator</em> one — a worker {@code tabLabel} template that
* happens to start with the same prefix would promote the whole fleet — and that is refused at
* startup by {@code FleetConfig.validateLeadTabPrefixes} rather than documented here.
*
* <p><strong>fleetd #359 — the staleness check this class used to skip.</strong> This used to say
* "a stale name costs nothing here" and leave liveness to {@code LeadLauncher}, on the theory that
* naming and removing are different decisions. That was wrong: {@code dev.ltms.fleet.msg.LeadCoordLoop}
* makes exactly the kind of removal decision the old javadoc warned about, by reading this map to
* pick which pane a peer message goes into — and a stale entry there is not free. On a host where
* {@code fleetd} had restarted more than once, a labelled-but-dead tab from a previous life was
* reported right alongside the live one; {@code LeadCoordLoop.resolveLocalLead()} saw more than one
* candidate and refused to guess (safe), but the fix the daemon's own WARN suggests — name a lead
* after {@code coordinator.selfId} — stops being safe once two tabs can share a label: step 1 of
* that resolution picks whichever matching entry it finds first, which can be the dead one, and
* typing a peer's message into a dead shell does not fail — it is silently gone instead of merely
* held. {@link #scan()} now cross-checks every labelled tab against {@code agent.list} (the same
* signal {@code LeadLauncher.countLeads} already trusts for the same purpose) and drops any tab
* with no agent running in it, so a dead tab is never in the map for a caller to pick at all.
*
* <p><strong>Caching.</strong> {@link #get()} is on the request path (every resolve), so the scan
* is TTL-cached and a stale-but-valid map is preferred to a herdr round-trip. A failed scan (herdr
* throws) keeps the previous answer instead of emptying it — a herdr hiccup must not silently
* demote a live lead mid-session.
*
* <p><strong>fleetd #359 review, finding 2 — a successful-but-wrong scan is the same hazard.</strong>
* The catch above only fires when a call throws. It does nothing for a call that returns 200 with an
* incomplete answer — exactly what the ticket's own evidence showed {@code agent.list} can do. Once
* this class started trusting that signal, an empty read would otherwise get cached as fact and
* silently drop a lead {@code CallerResolver} had, until then, correctly resolved — turning it into a
* {@code Role.WORKER}, which refuses every orchestration call. So a terminal this class already
* reported as live is not dropped the first time {@code agent.list} loses it: {@link #scan()} grants
* it one grace scan (see {@code gracedTerminals}) and only drops it if a <em>later</em> scan still
* finds no agent. A terminal never reported live before gets no grace — that would weaken the
* original #359 fix itself, which this class's own test suite already pins.
* is TTL-cached and a stale-but-valid map is preferred to a herdr round-trip. A failed scan keeps
* the previous answer instead of emptying it — a herdr hiccup must not silently demote a live lead
* mid-session.
*/
public final class LeadTabScanner implements Supplier<Map<String, String>> {
@@ -103,15 +79,6 @@ public final class LeadTabScanner implements Supplier<Map<String, String>> {
private long scannedAtNanos;
private boolean everScanned;
/**
* Terminals currently on their one grace scan: {@code cached} reported them live, the most
* recent {@link #scan()} found no agent for them, and they were re-included anyway. Cleared for
* a terminal the instant it is seen live again; a terminal still here on the <em>next</em> scan
* is finally dropped. Scoped separately from {@link #cached} so a graced terminal cannot renew
* its own grace forever just by staying in the exposed map (fleetd #359 review, finding 2).
*/
private Set<String> gracedTerminals = Set.of();
/**
* @param herdr the herdr client to query ({@code workspace.list},
* {@code tab.list}, {@code pane.list} — all read-only)
@@ -174,7 +141,7 @@ public final class LeadTabScanner implements Supplier<Map<String, String>> {
return cached;
}
/** One full pass: labelled tabs → live agents in them → those panes' terminals. */
/** One full pass: labelled tabs → their panes → those panes' terminals. */
private Map<String, String> scan() {
Map<String, String> nameByTab = new LinkedHashMap<>();
for (JsonNode w : herdr.call("workspace.list").path("workspaces")) {
@@ -191,47 +158,17 @@ public final class LeadTabScanner implements Supplier<Map<String, String>> {
}
}
if (nameByTab.isEmpty()) {
gracedTerminals = Set.of();
return Map.of();
}
// fleetd #359: a labelled tab is only a lead when herdr also reports a running agent in
// it — the same liveness signal LeadLauncher.countLeads trusts for the identical purpose.
// Without this, a tab left behind by a session that has since died reads as live forever.
Set<String> tabsWithAgent = new HashSet<>();
for (JsonNode a : herdr.call("agent.list").path("agents")) {
String tabId = a.path("tab_id").asText(null);
if (tabId != null) {
tabsWithAgent.add(tabId);
}
}
Map<String, String> byTerminal = new LinkedHashMap<>();
Set<String> stillGraced = new HashSet<>();
// One pane.list for every tab: panes carry tab_id, so the join is local.
for (JsonNode p : herdr.call("pane.list", Map.of()).path("panes")) {
String tabId = p.path("tab_id").asText(null);
String name = nameByTab.get(tabId);
String terminal = p.path("terminal_id").asText(null);
if (name == null || terminal == null || terminal.isBlank()) {
continue;
}
if (tabsWithAgent.contains(tabId)) {
byTerminal.put(terminal, name);
continue;
}
// No agent reported for this tab, but its tab/pane are still here — this is the
// ambiguous case review finding 2 named: a successful agent.list that came back short
// does not prove the lead is dead. Grant one grace scan to a terminal we had already
// reported as live; a terminal we never reported live gets none, so the original #359
// fix (a genuinely dead tab is never reported) is unaffected for the common case.
if (cached.containsKey(terminal) && !gracedTerminals.contains(terminal)) {
byTerminal.put(terminal, name);
stillGraced.add(terminal);
if (!nameByTab.isEmpty()) {
// One pane.list for every tab: panes carry tab_id, so the join is local.
for (JsonNode p : herdr.call("pane.list", Map.of()).path("panes")) {
String name = nameByTab.get(p.path("tab_id").asText(null));
String terminal = p.path("terminal_id").asText(null);
if (name != null && terminal != null && !terminal.isBlank()) {
byTerminal.put(terminal, name);
}
}
}
gracedTerminals = stillGraced;
return Collections.unmodifiableMap(byTerminal);
}
@@ -240,14 +177,12 @@ public final class LeadTabScanner implements Supplier<Map<String, String>> {
*
* <p>Exact match (case-insensitive, ends stripped) against {@link #tabToName} — no prefix
* stripping, so an operator's {@code "lead: something-else"} tab is never mistaken for a
* configured lead just because it shares a prefix. The match strips a trailing
* {@link PendingCloseMarker} first, so a tab {@code LeadLauncher} has flagged as maybe-dead but
* not yet closed keeps resolving normally while that reconcile is pending.
* configured lead just because it shares a prefix.
*/
private String leadNameOf(String label) {
if (label == null) {
return null;
}
return tabToName.get(PendingCloseMarker.strip(label).toLowerCase(Locale.ROOT));
return tabToName.get(label.strip().toLowerCase(Locale.ROOT));
}
}
@@ -1,42 +0,0 @@
package dev.ltms.fleet.herdr;
/**
* The suffix {@code dev.ltms.fleet.lead.LeadLauncher} appends to a lead tab's label the first time a
* reconcile finds no running agent in it, before it is sure enough to close the tab outright.
*
* <p><strong>fleetd #359 review, finding 1.</strong> The daemon's own evidence showed
* {@code agent.list} can read "no agent" for a tab that genuinely has one running — so a single such
* reading must never be treated as proof a tab is dead. {@code LeadLauncher} now writes this marker
* on the first miss, and only closes the tab if a <em>later</em>, independent reconcile still finds
* it dead while the marker is still there. Two consecutive misses, one restart apart, is a much
* stronger claim than one.
*
* <p>{@link LeadTabScanner} strips the same suffix before matching a label against a configured
* lead's {@code tab}, so a flagged-but-actually-still-live tab keeps resolving normally — the marker
* changes nothing about which pane {@code LeadCoordLoop} can reach while the flag is pending. Both
* classes must use exactly this suffix, which is why it lives here rather than as a private constant
* on either.
*/
public final class PendingCloseMarker {
public static final String SUFFIX = " [fleetd:pending-close]";
private PendingCloseMarker() {
}
/** The label with any trailing pending-close marker removed, for name matching. */
public static String strip(String label) {
if (label == null) {
return null;
}
String stripped = label.strip();
return stripped.endsWith(SUFFIX)
? stripped.substring(0, stripped.length() - SUFFIX.length()).strip()
: stripped;
}
/** Whether a label currently carries the marker. */
public static boolean isFlagged(String label) {
return label != null && label.strip().endsWith(SUFFIX);
}
}
@@ -4,7 +4,6 @@ import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.herdr.Agent;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.HerdrException;
import dev.ltms.fleet.herdr.PendingCloseMarker;
import dev.ltms.fleet.herdr.Tab;
import dev.ltms.fleet.herdr.Workspace;
import dev.ltms.fleet.herdr.WorkspaceControl;
@@ -13,11 +12,9 @@ import org.slf4j.LoggerFactory;
import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
import dev.ltms.fleet.peer.PeerLauncher;
/**
@@ -47,33 +44,8 @@ import dev.ltms.fleet.peer.PeerLauncher;
* privilege escalation (the tab label never granted anything a pane could take for itself; see that
* class's javadoc), but <em>staleness</em>: a label left behind by a crashed session would otherwise
* read as a live lead forever, and the lead would never be relaunched. So a lead counts as live only
* when herdr also reports a <em>running agent</em> in that tab — see {@link #countLeads}. A labelled
* when herdr also reports a <em>running agent</em> in that tab — see {@link #liveLeads}. A labelled
* tab with no agent in it is not a lead.
*
* <p><strong>fleetd #359 — a stale label used to pile up, not just mislead.</strong> Finding "not
* live" here used to mean only one thing: launch another. The old label was left exactly where it
* was, so a daemon that restarted enough times — or hit one herdr read that missed a genuinely
* running agent — accumulated one more identically-labelled dead tab per occurrence, and
* {@code LeadTabScanner} (before its own #359 fix) reported every one of them as a lead.
*
* <p><strong>Review finding 1 — closing on one reading is worse than the bug.</strong> The first
* version of this fix closed a name's dead tabs the moment a single {@link #countLeads} reading
* called them dead. The ticket's own live evidence rules that out: on a real host, {@code
* agent.list} was seen reporting "0 live" for a tab that a plain {@code ps} confirmed was running a
* real session. Closing on that reading would have destroyed the operator's actual lead — a worse
* failure than the extra tab it replaces. So {@link #ensureLeads()} now needs the same dead reading
* <em>twice</em>, one restart apart, before it closes anything: the first time a labelled tab reads
* dead, it is only flagged ({@link PendingCloseMarker}), left running, untouched; it is closed only
* if a <em>later</em>, independently-connected reconcile still finds it dead while the flag is
* still there. A transient miss self-heals — the next reconcile sees the agent again and clears the
* flag (see {@code toUnflag} below) — so the worst case for a single bad reading is one extra tab
* surviving one more restart, never a live session destroyed. Two alternatives were considered and
* rejected: corroborating {@code agent.list} against a second, truly independent signal was dropped
* because nothing else herdr exposes proves "is a process attached to this pane" any better — a
* second call to the same unreliable source is not independent evidence; capping the close to "all
* but the most recent dead tab" was dropped because "most recent" has no reliable ordering across
* tab ids and would leave the true failure mode (a name that is <em>never</em> reconfirmed) growing
* by one tab per bad reading forever, which is the exact defect this ticket exists to fix.
*/
public final class LeadLauncher {
@@ -107,9 +79,9 @@ public final class LeadLauncher {
return 0;
}
Map<String, LeadCount> live;
Map<String, Integer> live;
try {
live = countLeads(leaders);
live = liveLeads(leaders);
} catch (HerdrException e) {
// Counting is the whole safety mechanism against double-spawning. If we cannot count, we
// must not guess — spawning a second orchestrator is worse than starting none.
@@ -121,53 +93,9 @@ public final class LeadLauncher {
for (Map.Entry<String, FleetConfig.Leader> e : leaders.entrySet()) {
String name = e.getKey();
FleetConfig.Leader lead = e.getValue();
LeadCount state = live.getOrDefault(name, LeadCount.NONE);
int running = state.running();
int running = live.getOrDefault(name, 0);
int wanted = lead.instances();
// fleetd #359 review finding 1: a tab already flagged pending-close, still labelled for
// this lead, and STILL hosting no agent on this separate reconcile — two independent
// readings agree, so close it. A tab found dead for the first time is only flagged below,
// never closed on the spot.
for (String tabId : state.toClose()) {
log.info("lead '{}': closing tab {} — flagged pending-close on a previous reconcile "
+ "and still no agent running in it", name, tabId);
try {
spaces.closeTab(tabId);
} catch (RuntimeException cleanup) {
log.warn("could not close stale tab {} for lead '{}': {}",
tabId, name, cleanup.getMessage());
}
}
// A tab labelled for this lead, with no agent running in it, seen dead for the first
// time — flag it rather than closing it. One reading of `agent.list` is not enough
// evidence to destroy a tab that might genuinely be live (see the class javadoc).
for (String tabId : state.toFlag()) {
String flagged = lead.tabLabel() + PendingCloseMarker.SUFFIX;
log.info("lead '{}': tab {} has no agent running in it this reconcile — flagging it "
+ "'{}' rather than closing; it is only closed if a later reconcile still "
+ "finds it dead", name, tabId, flagged);
try {
spaces.renameTab(tabId, flagged);
} catch (RuntimeException cleanup) {
log.warn("could not flag stale tab {} for lead '{}': {}",
tabId, name, cleanup.getMessage());
}
}
// A previously-flagged tab that is running an agent again — the miss that flagged it was
// transient. Clear the flag so a future, unrelated miss starts its own two-reading count
// rather than closing on the strength of this one's already-spent flag.
for (String tabId : state.toUnflag()) {
log.info("lead '{}': tab {} is running an agent again — clearing its pending-close flag",
name, tabId);
try {
spaces.renameTab(tabId, lead.tabLabel());
} catch (RuntimeException cleanup) {
log.warn("could not clear the pending-close flag on tab {} for lead '{}': {}",
tabId, name, cleanup.getMessage());
}
}
if (running >= wanted) {
log.info("lead '{}': {} live, {} wanted — nothing to start", name, running, wanted);
continue;
@@ -198,29 +126,18 @@ public final class LeadLauncher {
}
/**
* How many live leads exist per configured name, and which of that name's labelled tabs are
* <em>not</em> live: a running agent in a tab labelled with that lead's exact {@code tab}
* (CB-579). Member workspaces are excluded, exactly as the scanner excludes them: a member must
* not be counted as a lead because it happens to sit in a matching tab.
* How many live leads exist per configured name: a running agent in a tab labelled with that
* lead's exact {@code tab} (CB-579). Member workspaces are excluded, exactly as the scanner
* excludes them: a member must not be counted as a lead because it happens to sit in a matching
* tab.
*
* <p>There used to be a second path here — a running agent on the terminal a
* {@code fleet.leaders.<name>.terminal} pin named, for a lead opened and pinned by hand. That
* pin is retired: {@code tab} is now the only field identity depends on, and {@link Agent}
* already carries {@link Agent#tabId()} directly, so a hand-opened lead is found the same way an
* auto-launched one is — by labelling its tab to match.
*
* <p>fleetd #359 review finding 1: a labelled tab with nothing running in it is split into
* {@code toClose} (already flagged pending-close by a previous reconcile, and still dead — two
* independent readings agree) and {@code toFlag} (dead for the first time — not enough evidence
* to close yet). {@code toUnflag} is the reverse: a tab flagged pending-close that is running an
* agent again, so the flag it carries no longer means anything and {@link #ensureLeads()} clears
* it.
*/
private record LeadCount(int running, List<String> toClose, List<String> toFlag, List<String> toUnflag) {
static final LeadCount NONE = new LeadCount(0, List.of(), List.of(), List.of());
}
private Map<String, LeadCount> countLeads(Map<String, FleetConfig.Leader> leaders) {
private Map<String, Integer> liveLeads(Map<String, FleetConfig.Leader> leaders) {
// A lead and the members share ONE workspace now (the operator asked for a single "session"
// with many tabs), so a workspace can no longer be excluded wholesale — the lead lives in the
// member workspace by design. The sole discriminator is the exact tab label: a lead carries
@@ -228,7 +145,6 @@ public final class LeadLauncher {
// profile's `worker: {profile} #{n}` template. These never collide, so an exact-label match
// separates them without needing to know which workspace anyone is in.
Map<String, String> nameByTab = new LinkedHashMap<>();
Set<String> flaggedTabIds = new LinkedHashSet<>();
for (Workspace ws : spaces.listWorkspaces()) {
if (ws.workspaceId() == null) {
continue;
@@ -237,65 +153,31 @@ public final class LeadLauncher {
String declared = leadNameOf(tab.label(), leaders);
if (declared != null && tab.tabId() != null) {
nameByTab.put(tab.tabId(), declared);
if (PendingCloseMarker.isFlagged(tab.label())) {
flaggedTabIds.add(tab.tabId());
}
}
}
}
Set<String> liveTabIds = new LinkedHashSet<>();
Map<String, Integer> counts = new LinkedHashMap<>();
for (Agent a : agents.list()) {
String name = nameByTab.get(a.tabId());
if (name != null) {
counts.merge(name, 1, Integer::sum);
liveTabIds.add(a.tabId());
}
}
Map<String, List<String>> toCloseByName = new LinkedHashMap<>();
Map<String, List<String>> toFlagByName = new LinkedHashMap<>();
Map<String, List<String>> toUnflagByName = new LinkedHashMap<>();
nameByTab.forEach((tabId, name) -> {
boolean live = liveTabIds.contains(tabId);
boolean flagged = flaggedTabIds.contains(tabId);
if (live) {
if (flagged) {
toUnflagByName.computeIfAbsent(name, k -> new ArrayList<>()).add(tabId);
}
return;
}
if (flagged) {
toCloseByName.computeIfAbsent(name, k -> new ArrayList<>()).add(tabId);
} else {
toFlagByName.computeIfAbsent(name, k -> new ArrayList<>()).add(tabId);
}
});
Map<String, LeadCount> out = new LinkedHashMap<>();
for (String name : leaders.keySet()) {
out.put(name, new LeadCount(counts.getOrDefault(name, 0),
toCloseByName.getOrDefault(name, List.of()),
toFlagByName.getOrDefault(name, List.of()),
toUnflagByName.getOrDefault(name, List.of())));
}
return out;
return counts;
}
/**
* The configured lead a tab label names, or {@code null} for a label that names none.
*
* <p>Matched exactly (case-insensitively) against each lead's configured {@code tab}, so an
* operator's {@code "lead: something-else"} tab is not mistaken for a configured lead. A
* trailing {@link PendingCloseMarker} is stripped first, so a tab this class flagged on a
* previous reconcile is still recognised as the same lead's tab on this one.
* operator's {@code "lead: something-else"} tab is not mistaken for a configured lead.
*/
private String leadNameOf(String label, Map<String, FleetConfig.Leader> leaders) {
if (label == null) {
return null;
}
String l = PendingCloseMarker.strip(label);
String l = label.strip();
for (Map.Entry<String, FleetConfig.Leader> e : leaders.entrySet()) {
String tab = e.getValue().tabLabel();
if (tab != null && l.equalsIgnoreCase(tab.strip())) {
@@ -44,7 +44,6 @@ import java.util.Objects;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
import java.util.function.BiFunction;
import java.util.function.Function;
import java.util.function.LongSupplier;
@@ -102,8 +101,6 @@ public final class FleetMcp {
private final LeadSeatSource leadSeats;
/** CB-637: this daemon's lead-to-lead channel; {@code null} when no coordinator is configured. */
private final LeadChannel leadChannel;
/** fleetd #361: {@code coordinator.peers} — see {@link CoordinationSource}. Empty when unset. */
private final List<String> peers;
/** Capacity facts used by {@code fleet_list}; production must supply the placement live count. */
public record CapacitySource(Function<String, Integer> liveCount, Function<String, Integer> maxLoad,
@@ -167,32 +164,6 @@ public final class FleetMcp {
public static LeadSeatSource none() { return new LeadSeatSource(_ -> 0); }
}
/**
* fleetd #361: peer-visibility facts for {@code fleet_list}'s {@code coordinator} row — this
* daemon's own {@link LeadChannel} (for its self mailbox state and held messages) plus the
* coord-ids the operator has declared as peers ({@code coordinator.peers}). Bundled as its own
* Source, the same idiom as {@link OutageSource}/{@link QuarantineSource}/{@link LeadSeatSource},
* so {@code listFleet}'s already-long overload chain gains exactly one new required parameter
* instead of a further bare positional argument.
*
* <p>Every mailbox look this triggers goes through {@link LeadChannel#inspect}, which is
* specified to run on its own disposable channel — never the channel {@link LeadChannel#publish}
* or the consume loop depends on — so a peer that happens to be down, or the coordination broker
* itself being unreachable, can never take {@code fleet_send}/{@code LeadCoordLoop}'s own path
* down with it. See {@code FleetMcp.probe} for the additional timeout bound on top of that.
*
* @param leadChannel this daemon's own channel, or {@code null} when no coordinator is configured
* @param peers the coord-ids declared under {@code coordinator.peers}, or empty
*/
public record CoordinationSource(LeadChannel leadChannel, List<String> peers) {
public CoordinationSource {
peers = peers == null ? List.of() : List.copyOf(peers);
}
/** Inert source — no coordinator row is ever reported. */
public static CoordinationSource none() { return new CoordinationSource(null, List.of()); }
}
/**
* @param callers resolves each call's {@link Principal}; {@code null} disables authorization.
* This surface needs its own enforcement: {@code /mcp} is a raw servlet on
@@ -240,34 +211,19 @@ public final class FleetMcp {
}
/**
* As above, with fleetd #176 lead-seat facts (see {@link LeadSeatSource}).
* As above, with fleetd #176 lead-seat facts (see {@link LeadSeatSource}). This is what
* {@code Fleetd.main} actually wires up.
*
* @param leadSeats required — pass {@link LeadSeatSource#none()} for a caller that does not want
* the feature, never a defaulting overload (the same rule {@code quarantine} and
* {@code outage} follow).
*/
public FleetMcp(MessageService messages, PeerLauncher workers, SessionManager sessions,
ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry,
CallerResolver callers, Metrics metrics, CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, LeadChannel leadChannel, OutageSource outage,
LeadSeatSource leadSeats) {
this(messages, workers, sessions, identity, presence, primaryRegistry, callers, metrics, capacity,
healthCoverage, quarantine, leadChannel, outage, leadSeats, List.of());
}
/**
* As above, with fleetd #361 {@code coordinator.peers} (see {@link CoordinationSource}). This is
* what {@code Fleetd.main} actually wires up.
*
* @param leadSeats required — pass {@link LeadSeatSource#none()} for a caller that does not want
* the feature, never a defaulting overload (the same rule {@code quarantine} and
* {@code outage} follow).
* @param peers the coord-ids declared under {@code coordinator.peers}; empty when unset or
* when {@code leadChannel} is {@code null}.
*/
public FleetMcp(MessageService messages, PeerLauncher workers, SessionManager sessions,
ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry,
CallerResolver callers, Metrics metrics, CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, LeadChannel leadChannel, OutageSource outage,
LeadSeatSource leadSeats, List<String> peers) {
this.leadChannel = leadChannel;
this.peers = peers == null ? List.of() : List.copyOf(peers);
this.capacity = capacity;
this.quarantine = Objects.requireNonNull(quarantine, "quarantine");
this.outage = Objects.requireNonNull(outage, "outage");
@@ -405,7 +361,7 @@ public final class FleetMcp {
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
leadSeats, callers == null ? Map.of() : callers.leads(),
callerTerminal(exchange),
new CoordinationSource(leadChannel, peers));
leadChannel == null ? null : leadChannel.selfCoordId());
};
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> stopHandler =
(exchange, req) -> {
@@ -741,15 +697,6 @@ public final class FleetMcp {
* turned into a tool error naming the coord-id. It is never allowed to escape as a crash: an
* unreachable peer is an ordinary outcome of addressing a fleet you do not control.
*
* <p><strong>fleetd #361: the success text is honest about what "durably confirmed" does and
* does not mean.</strong> The broker's publisher confirm proves the message is durably queued —
* it says nothing about whether the peer's pane has, or ever will, receive it. After a
* successful publish this looks at the target mailbox's consumer count (via
* {@link LeadChannel#inspect}, bounded and never allowed to fail the call — see {@link #probe})
* and appends a warning when it is zero: that is the observable form of "nobody is reading this
* right now". A zero-consumer publish is still reported as a SUCCESS, never an error — the
* message is safely queued and will be read once a daemon owning that coord-id connects.
*
* @param leadChannel this daemon's channel, or {@code null} when no coordinator is configured
*/
static McpSchema.CallToolResult sendToLead(LeadChannel leadChannel, String coordId, String content,
@@ -777,16 +724,7 @@ public final class FleetMcp {
+ ". Check that a daemon is running with coordinator.selfId=\"" + coordId
+ "\" and is connected to the same coordination broker.");
}
String result = "published to peer lead \"" + coordId + "\"'s mailbox and durably confirmed "
+ "by the broker (msgId " + msg.msgId() + ").";
LeadChannel.MailboxState state = probe(leadChannel, coordId);
if (state.exists() && state.consumers() == 0) {
result += " Warning: that mailbox currently has NO consumers attached — nobody is reading "
+ "it right now. The message is safely queued and will be delivered once a daemon "
+ "with coordinator.selfId=\"" + coordId + "\" is running and connected; until then "
+ "it will not reach that lead's pane.";
}
return text(result);
return text("delivered to peer lead " + coordId + " (msgId " + msg.msgId() + ")");
}
/**
@@ -1172,8 +1110,7 @@ public final class FleetMcp {
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, Map<String, String> leads, String selfTerm) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, leads, selfTerm,
CoordinationSource.none());
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, leads, selfTerm, null);
}
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
@@ -1182,33 +1119,33 @@ public final class FleetMcp {
QuarantineSource quarantine, OutageSource outage,
Map<String, String> leads, String selfTerm) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
LeadSeatSource.none(), leads, selfTerm, CoordinationSource.none());
LeadSeatSource.none(), leads, selfTerm, null);
}
/**
* As above, additionally reporting this daemon's own lead coordination state (CB-637, fleetd
* #361) when a coordinator is configured and its channel opened — see {@link #coordinatorView}
* for the shape. Omitted entirely when no coordinator is configured, so an ordinary fleet's
* output is unchanged.
* As above, additionally reporting this daemon's own lead coordination id (CB-637) when one is
* configured and its channel opened. There is no peer-discovery surface yet — a lead addresses a
* peer by a coord-id it was told — so this row exists to answer the one question the operator
* cannot answer any other way: what is MY coord-id, the one a peer must use to reach me. It is
* omitted entirely when no coordinator is configured, so an ordinary fleet's output is unchanged.
*
* @param coordination this daemon's lead channel plus its declared peers, or
* {@link CoordinationSource#none()} when lead coordination is off
* @param selfCoordId this daemon's coord-id, or {@code null} when lead coordination is off
*/
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, Map<String, String> leads, String selfTerm,
CoordinationSource coordination) {
String selfCoordId) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, OutageSource.none(),
LeadSeatSource.none(), leads, selfTerm, coordination);
LeadSeatSource.none(), leads, selfTerm, selfCoordId);
}
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, OutageSource outage,
Map<String, String> leads, String selfTerm, CoordinationSource coordination) {
Map<String, String> leads, String selfTerm, String selfCoordId) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
LeadSeatSource.none(), leads, selfTerm, coordination);
LeadSeatSource.none(), leads, selfTerm, selfCoordId);
}
/** As above, plus fleetd #176 lead-seat facts (see {@link LeadSeatSource}). */
@@ -1216,7 +1153,7 @@ public final class FleetMcp {
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, OutageSource outage,
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
CoordinationSource coordination) {
String selfCoordId) {
try {
Map<String, Agent> live = workers.list().stream()
.map(Agent.class::cast)
@@ -1238,9 +1175,8 @@ public final class FleetMcp {
Map<String, Object> result = new LinkedHashMap<>();
result.put("leads", leadRows); result.put("members", out);
result.put("healthCoverage", healthCoverage.value().get());
Map<String, Object> coordinatorRow = coordinatorView(coordination);
if (coordinatorRow != null) {
result.put("coordinator", coordinatorRow);
if (selfCoordId != null && !selfCoordId.isBlank()) {
result.put("coordinator", Map.of("selfId", selfCoordId, "configured", true));
}
if (capacity.available()) result.put("capacity", profiles.stream()
.map(profile -> capacityView(profile, capacity.liveCount(), capacity.maxLoad(), roster, messages,
@@ -1251,140 +1187,6 @@ public final class FleetMcp {
}
}
/**
* fleetd #361: the {@code coordinator} row — this daemon's own coord-id and mailbox state, the
* messages currently held for it, and the live reachability of every operator-declared peer.
* {@code null} (the row is then omitted entirely) when lead coordination is off, so an ordinary
* fleet's {@code fleet_list} output is byte-identical to before this feature existed.
*
* <p>Every peer/self mailbox look goes through {@link #probe}, which bounds each
* {@link LeadChannel#inspect} call to {@link #PEER_PROBE_TIMEOUT_MS} and never lets it throw —
* a coordination broker that is down or slow degrades this row toward "unreachable"/"unknown"
* counts, it can never make {@code fleet_list} itself slow or fail. {@code held} comes from
* {@link LeadChannel#peek}, a pure in-memory read with no broker round trip, so it is never
* subject to that bound.
*/
private static Map<String, Object> coordinatorView(CoordinationSource coordination) {
LeadChannel channel = coordination.leadChannel();
if (channel == null) {
return null;
}
String selfId = channel.selfCoordId();
Map<String, Object> row = new LinkedHashMap<>();
row.put("selfId", selfId);
row.put("configured", true);
row.put("mailbox", mailboxView(probe(channel, selfId)));
row.put("held", channel.peek().stream().map(FleetMcp::heldView).toList());
row.put("peers", coordination.peers().stream().map(p -> peerView(channel, p)).toList());
return row;
}
/**
* fleetd #361: render a {@link LeadChannel.MailboxState} without ever presenting an unmeasured
* fact as a measured one. {@code status} is the tri-state itself — {@code "exists"},
* {@code "absent"} (the broker positively confirmed no such queue), or {@code "unknown"} (the
* probe could not determine either way: down, unreachable, or timed out). {@code pending}/
* {@code consumers} are included ONLY when {@code status == "exists"} — a reader must never see
* them default to {@code 0} for a mailbox this call never actually measured. This is the fix for
* the review finding that a collapsed {@code absent()} rendered a self-probe timeout as
* "pending: 0, consumers: 0", indistinguishable from an actually-empty, actually-unread mailbox.
*/
private static Map<String, Object> mailboxView(LeadChannel.MailboxState state) {
Map<String, Object> row = new LinkedHashMap<>();
row.put("status", state.exists() ? "exists" : state.known() ? "absent" : "unknown");
if (state.exists()) {
row.put("pending", state.pending());
row.put("consumers", state.consumers());
}
return row;
}
/** One held-for-me message: enough to identify it and see roughly what it says, never the whole body. */
private static Map<String, Object> heldView(LeadMessage m) {
Map<String, Object> row = new LinkedHashMap<>();
row.put("msgId", m.msgId());
row.put("from", m.from());
row.put("preview", preview(m.content()));
return row;
}
/** Cap a held message's content to a short preview — {@code fleet_list} must never dump a full body. */
private static final int HELD_PREVIEW_MAX_CHARS = 80;
private static String preview(String content) {
if (content == null) {
return "";
}
return content.length() <= HELD_PREVIEW_MAX_CHARS
? content
: content.substring(0, HELD_PREVIEW_MAX_CHARS) + "…";
}
/**
* One declared peer's row: its coord-id, then the same tri-state {@link #mailboxView} shape.
* Deliberately no boolean "reachable" field — that collapsed "confirmed gone" and "could not
* check" into the same {@code false}, which is exactly the review finding this row now avoids:
* an operator reading {@code status} can tell "fleet01 is down" (a {@code coordinator.selfId}
* nobody has ever run) apart from "my own broker is slow or unreachable right now".
*/
private static Map<String, Object> peerView(LeadChannel channel, String coordId) {
Map<String, Object> row = new LinkedHashMap<>();
row.put("coordId", coordId);
row.putAll(mailboxView(probe(channel, coordId)));
return row;
}
/**
* fleetd #361: how long {@code fleet_list} waits on any single {@link LeadChannel#inspect} call
* before giving up on it — see {@link #probe}.
*/
private static final long PEER_PROBE_TIMEOUT_MS = 1_500L;
/**
* Dedicated pool for {@link LeadChannel#inspect} calls so a slow one blocks only its own virtual
* thread, never the MCP request thread calling {@code fleet_list}. Not a <em>bounded</em> pool —
* {@code newThreadPerTaskExecutor} starts a fresh virtual thread per call with no cap on how many
* run at once; virtual threads make that cheap, not bounded. What actually keeps a hung probe
* from accumulating forever is the {@link Future#cancel} in {@link #probe}, not a pool limit.
*/
private static final java.util.concurrent.ExecutorService PEER_PROBE_POOL =
java.util.concurrent.Executors.newThreadPerTaskExecutor(Thread.ofVirtual().name("fleet-peer-probe-", 0).factory());
/**
* fleetd #361: {@link LeadChannel#inspect}, bounded to {@link #PEER_PROBE_TIMEOUT_MS} and never
* allowed to throw or hang the caller — a coordination broker that is unreachable or slow
* degrades to {@link LeadChannel.MailboxState#unknown} (never {@code absent}: a timeout proves
* nothing about whether the mailbox exists) rather than making {@code fleet_list} slow or
* failing it. {@code inspect} itself is already specified to never throw, but this is the seam
* that also survives an implementation that does, or one that blocks indefinitely on a dead
* connection.
*
* <p><strong>A timeout cancels the orphaned task</strong> rather than abandoning it. Before this,
* {@code get(timeout)} on a hung {@code inspect} left the submitted task running forever on its
* own virtual thread, holding the AMQP channel it had already opened — against a broker that
* hangs rather than fails fast, every {@code fleet_list} call would orphan one more channel until
* the connection's channel-max (2047 by default) was exhausted, which would break {@link
* LeadChannel#publish} too. {@link Future#cancel(boolean) cancel(true)} interrupts the orphaned
* task's thread; {@link LeadMailbox#inspect} has no interruptible wait of its own to catch that,
* but the underlying AMQP RPC continuation does block on one, so the interrupt reaches it and the
* task's {@code finally} still closes the probe channel it opened rather than leaking it forever.
*/
private static LeadChannel.MailboxState probe(LeadChannel channel, String coordId) {
return probe(channel, coordId, PEER_PROBE_TIMEOUT_MS);
}
/** As {@link #probe(LeadChannel, String)}, with an explicit timeout — a seam for tests. */
static LeadChannel.MailboxState probe(LeadChannel channel, String coordId, long timeoutMs) {
java.util.concurrent.Future<LeadChannel.MailboxState> future =
PEER_PROBE_POOL.submit(() -> channel.inspect(coordId));
try {
return future.get(timeoutMs, TimeUnit.MILLISECONDS);
} catch (Exception e) {
future.cancel(true); // best-effort: don't leave a hung probe (and its channel) running forever
return LeadChannel.MailboxState.unknown(coordId);
}
}
/**
* Capacity is advisory only. {@code reclaimable} says there is no bridge work, not that fleetd
* may stop the member: the bridge has capacity facts but no work list, and choosing work needs
@@ -1575,12 +1377,8 @@ public final class FleetMcp {
+ "routes your answer back into the same turn (omit for a normal delegation)"),
"coordId", stringProp("A peer LEAD's coordination id — delivers content to that "
+ "lead's durable mailbox on the shared coordination broker, which works "
+ "across hosts. Mutually exclusive with sessionId and turnId. Success means "
+ "the message is durably queued and confirmed by the broker, with a warning "
+ "if that mailbox has no consumers attached right now (queued, but nobody is "
+ "reading it yet) — it does not mean the peer's pane has seen it. Your own "
+ "coordId, this daemon's mailbox state, and every coordinator.peers entry's "
+ "live reachability are reported by fleet_list's coordinator row.")),
+ "across hosts. Mutually exclusive with sessionId and turnId. Your own "
+ "coordId is reported by fleet_list.")),
List.of("content")));
}
@@ -1688,13 +1486,7 @@ public final class FleetMcp {
+ "for a quarantined profile's credential (see fleet_profiles), whatever its "
+ "maxLoad/live — with credentialId and quarantinedForSeconds naming the "
+ "quarantine, so 'free: 0, busy' can be told apart from 'free: 0, refusing "
+ "for N seconds'. When lead-to-lead coordination is configured, a 'coordinator' "
+ "object reports this daemon's own coord-id ('selfId') and mailbox state "
+ "('mailbox': pending/consumers), the messages currently held for it ('held': "
+ "msgId/from/preview, never the full body), and one row per coordinator.peers "
+ "coord-id ('peers': coordId/reachable, plus pending/consumers when reachable) — "
+ "this is peer DISCOVERY for cross-host leads, distinct from the local 'leads' "
+ "array above. It is omitted entirely when no coordinator is configured.",
+ "for N seconds'.",
objectSchema(Map.of(), List.of()));
}
@@ -8,7 +8,6 @@ import com.rabbitmq.client.DeliverCallback;
import com.rabbitmq.client.Recoverable;
import com.rabbitmq.client.RecoveryListener;
import com.rabbitmq.client.Return;
import com.rabbitmq.client.impl.DefaultExceptionHandler;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -154,22 +153,17 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
/** As {@link #open(String)}, with an explicit consumer prefetch (CB-527: caps the held backlog per target). */
public static AmqpReplyInbox open(String uri, int prefetch) {
try {
return new AmqpReplyInbox(connectionFactory(uri).newConnection(AmqpConnectionFailureLogger.REPLY_INBOX), prefetch);
ConnectionFactory factory = new ConnectionFactory();
factory.setUri(uri);
// Self-heal transient blips; topology recovery re-declares queues and re-attaches consumers.
factory.setAutomaticRecoveryEnabled(true);
factory.setTopologyRecoveryEnabled(true);
return new AmqpReplyInbox(factory.newConnection("fleetd-reply-inbox"), prefetch);
} catch (Exception e) {
throw new IllegalStateException("cannot connect to AMQP broker at " + uri, e);
}
}
static ConnectionFactory connectionFactory(String uri) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setUri(uri);
// Self-heal transient blips; topology recovery re-declares queues and re-attaches consumers.
factory.setAutomaticRecoveryEnabled(true);
factory.setTopologyRecoveryEnabled(true);
factory.setExceptionHandler(new AmqpConnectionFailureLogger(AmqpConnectionFailureLogger.REPLY_INBOX, log));
return factory;
}
/** Wrap an already-open connection with {@link #DEFAULT_PREFETCH} (injection seam for the contract test). */
AmqpReplyInbox(Connection connection) {
this(connection, DEFAULT_PREFETCH);
@@ -577,44 +571,3 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
}
}
}
/**
* Keeps RabbitMQ's forgiving exception behaviour while adding the connection identity that its
* default logger drops. Package-private so both AMQP connections use the same two names.
*/
final class AmqpConnectionFailureLogger extends DefaultExceptionHandler {
static final String REPLY_INBOX = "fleetd-reply-inbox";
static final String LEAD_MAILBOX = "fleetd-lead-mailbox";
private final String connectionName;
private final Logger logger;
AmqpConnectionFailureLogger(String connectionName, Logger logger) {
this.connectionName = connectionName;
this.logger = logger;
}
String connectionName() {
return connectionName;
}
@Override
protected void log(String message, Throwable cause) {
if (isSocketClosedOrConnectionReset(cause)) {
logger.warn("AMQP connection {}: {} (Exception message: {})", connectionName, message, cause.getMessage());
} else {
logger.error("AMQP connection {}: {}", connectionName, message, cause);
}
}
private static boolean isSocketClosedOrConnectionReset(Throwable cause) {
// Deliberate copy of ForgivingExceptionHandler's private static helper; check it on amqp-client upgrades.
if (!(cause instanceof IOException)) {
return false;
}
return "Connection reset".equals(cause.getMessage())
|| "Socket closed".equals(cause.getMessage())
|| "Connection reset by peer".equals(cause.getMessage());
}
}
@@ -4,8 +4,8 @@ import java.util.List;
/**
* The lead-to-lead message channel this daemon speaks, as its callers need it — one lead's own
* mailbox: publish to a peer's coord-id, look at what has arrived for me, ack what I have
* delivered, and (fleetd #361) inspect any coord-id's mailbox from the outside without owning it.
* mailbox: publish to a peer's coord-id, look at what has arrived for me, and ack what I have
* delivered.
*
* <p>Extracted from {@link LeadMailbox} purely as a seam. {@code LeadMailbox} is the one production
* implementation and owns a live AMQP connection, so a test that wanted to exercise the routing in
@@ -37,78 +37,4 @@ public interface LeadChannel {
/** This daemon's own lead coordination id — the mailbox it owns, and the {@code from} it sends as. */
String selfCoordId();
/**
* A non-destructive look at {@code coordId}'s mailbox — does it exist, how many messages are
* waiting on it, and how many consumers are attached — without owning, consuming, or otherwise
* changing it. {@code consumers == 0} on an existing mailbox is the observable form of "nobody
* is reading this right now": a publish to it will sit queued rather than reach a pane.
*
* <p><strong>Never throws</strong> — this is a best-effort fact-finding call, not an operation a
* caller must handle failing. But it must never turn "I could not check" into a false negative:
* {@link MailboxState#absent(String)} means the broker positively confirmed there is no such
* queue, and {@link MailboxState#unknown(String)} — a distinct value — means the look could not
* be completed at all (broker unreachable, timed out, connection closed). A caller that
* collapses those two into one, as fleetd #361 initially did, cannot tell "that peer is down"
* from "I could not check", and a reader of {@code pending}/{@code consumers} cannot tell a
* measured zero from a zero standing in for "not measured".
*
* <p><strong>Must never share fate with {@link #publish} or {@link #peek}/{@link #ack}.</strong>
* fleetd #361: in AMQP 0-9-1 a passive queue declare of a queue that does not exist closes the
* channel it was declared on with a 404. An implementation backed by a real broker connection
* must inspect on a channel it can afford to lose — never the channel {@link #publish} or the
* consume loop depends on — so that looking at a peer that happens to be down can never break
* this daemon's own send or receive path.
*/
MailboxState inspect(String coordId);
/**
* The result of {@link #inspect}. {@code presence} tells apart three states a caller must not
* conflate: a confirmed-existing mailbox ({@link Presence#EXISTS}, the only case where
* {@code pending}/{@code consumers} are measured facts), a confirmed-absent one
* ({@link Presence#ABSENT} — the broker positively said "no such queue"), and one this call
* simply could not determine ({@link Presence#UNKNOWN} — broker unreachable, timed out,
* connection closed). {@code pending}/{@code consumers} are always {@code 0} and meaningless
* outside {@link Presence#EXISTS}; a renderer must gate on {@link #exists()} (or {@code
* presence} directly), never present them as measured otherwise.
*
* @param coordId the coord-id inspected
* @param presence whether the mailbox is confirmed to exist, confirmed absent, or unknown
* @param pending messages ready for delivery but not yet in a consumer's hands (0 unless EXISTS)
* @param consumers how many consumers are attached (0 unless EXISTS)
*/
record MailboxState(String coordId, Presence presence, int pending, int consumers) {
/** Whether {@link #inspect} was able to reach a definite answer, of either kind. */
public enum Presence { EXISTS, ABSENT, UNKNOWN }
/** {@code true} only when the broker confirmed this exact queue is currently declared. */
public boolean exists() {
return presence == Presence.EXISTS;
}
/**
* {@code true} when {@link #inspect} reached a definite answer (exists or confirmed
* absent); {@code false} when it could not determine either way. A caller must never treat
* {@code !known()} the same as a confirmed absence — the mailbox may well exist.
*/
public boolean known() {
return presence != Presence.UNKNOWN;
}
/** The broker confirmed this queue exists, with these measured counts. */
public static MailboxState exists(String coordId, int pending, int consumers) {
return new MailboxState(coordId, Presence.EXISTS, pending, consumers);
}
/** The broker positively confirmed there is no such queue (e.g. a 404 on passive declare). */
public static MailboxState absent(String coordId) {
return new MailboxState(coordId, Presence.ABSENT, 0, 0);
}
/** The look could not be completed — broker unreachable, timed out, or connection closed. */
public static MailboxState unknown(String coordId) {
return new MailboxState(coordId, Presence.UNKNOWN, 0, 0);
}
}
}
@@ -10,7 +10,6 @@ import com.rabbitmq.client.DeliverCallback;
import com.rabbitmq.client.Recoverable;
import com.rabbitmq.client.RecoveryListener;
import com.rabbitmq.client.Return;
import com.rabbitmq.client.ShutdownSignalException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -123,22 +122,17 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable {
/** As {@link #open(String, String)}, with an explicit consumer prefetch. */
public static LeadMailbox open(String uri, String selfCoordId, int prefetch) {
try {
return new LeadMailbox(connectionFactory(uri).newConnection(AmqpConnectionFailureLogger.LEAD_MAILBOX), selfCoordId, prefetch);
ConnectionFactory factory = new ConnectionFactory();
factory.setUri(uri);
// Self-heal transient blips; topology recovery re-declares the queue and re-attaches the consumer.
factory.setAutomaticRecoveryEnabled(true);
factory.setTopologyRecoveryEnabled(true);
return new LeadMailbox(factory.newConnection("fleetd-lead-mailbox"), selfCoordId, prefetch);
} catch (Exception e) {
throw new IllegalStateException("cannot connect to AMQP coordination broker at " + uri, e);
}
}
static ConnectionFactory connectionFactory(String uri) throws Exception {
ConnectionFactory factory = new ConnectionFactory();
factory.setUri(uri);
// Self-heal transient blips; topology recovery re-declares queues and re-attaches consumers.
factory.setAutomaticRecoveryEnabled(true);
factory.setTopologyRecoveryEnabled(true);
factory.setExceptionHandler(new AmqpConnectionFailureLogger(AmqpConnectionFailureLogger.LEAD_MAILBOX, log));
return factory;
}
/** Wrap an already-open connection with {@link #DEFAULT_PREFETCH} (injection seam for tests). */
LeadMailbox(Connection connection, String selfCoordId) {
this(connection, selfCoordId, DEFAULT_PREFETCH);
@@ -260,89 +254,6 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable {
}
}
/**
* fleetd #361: look at {@code coordId}'s mailbox on a fresh, immediately-closed throwaway
* channel — never {@link #channel} (consume/ack) or {@link #publishChannel} (publish). A
* passive queue declare of a queue that does not exist closes the channel it was declared on
* with a 404; using a disposable probe channel means that closure can never touch either
* long-lived channel this instance depends on for {@link #publish} or the consume loop.
*
* <p>Classifies failures rather than collapsing them, both measured against a real broker in
* {@code LeadMailboxTest} rather than assumed from the AMQP 0-9-1 spec text:
* <ul>
* <li>a genuine 404 — an {@link IOException} wrapping a {@link ShutdownSignalException} whose
* {@link AMQP.Channel.Close#getReplyCode()} is {@code 404} — reports
* {@link MailboxState#absent}; every other declare failure reports
* {@link MailboxState#unknown} instead of quietly becoming the same "absent" value;
* <li>{@code catch (RuntimeException e)} on both attempts matters as much as the checked
* catches: a connection that is already closed makes {@link Connection#createChannel()}
* throw {@link com.rabbitmq.client.AlreadyClosedException} (a {@link RuntimeException},
* not an {@link IOException}) — an {@code inspect} that only caught {@code IOException}
* would let that escape, breaking the "never throws" contract this method promises.
* </ul>
*
* <p><strong>Honesty about which catch is measured and which is defensive:</strong> the
* {@code createChannel()} catch above is exercised end-to-end against a real broker by
* {@code LeadMailboxTest.inspectReportsUnknownRatherThanThrowingWhenTheConnectionIsAlreadyClosed}.
* The second {@code catch (RuntimeException e)}, around the passive declare itself — for the
* narrower race where the connection drops <em>between</em> {@code createChannel()} succeeding
* and the declare landing — has no such test; reaching it needs a connection that dies at that
* exact instant, which is not a scenario this suite drives on purpose. It stays purely
* defensive: correct by the same reasoning as the first catch, but unproven the way the first
* one is proven.
*/
@Override
public MailboxState inspect(String coordId) {
String queue = queueName(coordId);
Channel probe;
try {
probe = connection.createChannel();
} catch (IOException | RuntimeException e) {
log.debug("lead mailbox inspect: cannot open a probe channel for {}: {}", coordId, e.toString());
return MailboxState.unknown(coordId);
}
try {
AMQP.Queue.DeclareOk declared = probe.queueDeclarePassive(queue);
return MailboxState.exists(coordId, declared.getMessageCount(), declared.getConsumerCount());
} catch (IOException e) {
// The broker (or the client library) has already closed `probe` for us either way; only
// a confirmed 404 means "no such queue" — anything else (a different declare failure) is
// "could not determine", never silently reported as the same value as a genuine absence.
return isMissingQueue(e) ? MailboxState.absent(coordId) : MailboxState.unknown(coordId);
} catch (RuntimeException e) {
// E.g. the connection dropped between createChannel() and the declare landing.
log.debug("lead mailbox inspect: declare failed unexpectedly for {}: {}", coordId, e.toString());
return MailboxState.unknown(coordId);
} finally {
try {
if (probe.isOpen()) {
probe.close();
}
} catch (Exception e) {
log.debug("lead mailbox inspect: probe channel close for {}: {}", coordId, e.toString());
}
}
}
/**
* {@code true} only for the specific shape a missing-queue passive declare actually produces —
* measured against a real broker, not assumed from the spec text (see {@code
* LeadMailboxTest.passiveDeclareOfAMissingQueueThrowsAnIOExceptionWrappingA404ShutdownSignal}):
* an {@link IOException} whose cause is a {@link ShutdownSignalException} carrying an
* {@link AMQP.Channel.Close} reason with {@code replyCode == 404}. Any other shape (a different
* reply code, a {@code ShutdownSignalException} cause whose reason is not a
* {@code Channel.Close}, or no cause at all) is a declare failure of some other kind and must
* not be read as "confirmed absent" — pinned hermetically, with no broker needed, by
* {@code LeadMailboxIsMissingQueueTest} for exactly those three false shapes. Package-private
* (not {@code private}) so that test can call it directly.
*/
static boolean isMissingQueue(IOException e) {
if (!(e.getCause() instanceof ShutdownSignalException sse)) {
return false;
}
return sse.getReason() instanceof AMQP.Channel.Close close && close.getReplyCode() == AMQP.NOT_FOUND;
}
/** Convenience: {@link #peek} the current snapshot, then {@link #ack} every message in it. */
public List<LeadMessage> drain() {
List<LeadMessage> snapshot = peek();
@@ -0,0 +1,93 @@
package dev.ltms.fleet.power;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.util.Locale;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
/**
* Holds macOS idle sleep off by keeping a {@code caffeinate -i} child process alive for the life
* of the returned {@link SleepAssertion}.
*
* <p>{@code -i} asserts only against <em>idle</em> sleep — it does not stop the lid closing or an
* operator-requested sleep from taking effect. That is deliberate: this class exists to stop an
* unattended host from sleeping out from under a member's long turn, never to override the
* operator. {@code -s}/{@code -d} (which also block system/display sleep on demand) are
* intentionally not used here.
*
* <p>{@link #acquire()} never throws. It returns {@code null} — a no-op — off macOS, and again if
* starting the {@code caffeinate} child fails for any reason (binary missing, process table full,
* …); either case is logged once at INFO, not on every occurrence, so a daemon that runs for
* weeks with the tool unavailable does not fill its log.
*/
public final class CaffeinateSleepAssertionMechanism implements SleepAssertionMechanism {
private static final Logger log = LoggerFactory.getLogger(CaffeinateSleepAssertionMechanism.class);
private final AtomicBoolean loggedOnce = new AtomicBoolean(false);
/** {@code true} when running on macOS, the only platform {@code caffeinate} ships on. */
public static boolean isSupportedPlatform() {
return isSupportedPlatform(System.getProperty("os.name"));
}
/** Package-visible so a test can drive the platform check without touching a real property. */
static boolean isSupportedPlatform(String osName) {
return osName != null && osName.toLowerCase(Locale.ROOT).contains("mac");
}
@Override
public SleepAssertion acquire() {
if (!isSupportedPlatform()) {
logOnce("not running on macOS (os.name={}); the idle-sleep guard is a no-op on this platform",
System.getProperty("os.name"));
return null;
}
try {
Process process = new ProcessBuilder("caffeinate", "-i")
.redirectOutput(ProcessBuilder.Redirect.DISCARD)
.redirectError(ProcessBuilder.Redirect.DISCARD)
.start();
return new CaffeinateAssertion(process);
} catch (IOException | RuntimeException e) {
logOnce("could not start 'caffeinate -i' ({}); the host may idle-sleep while members are live",
e.toString());
return null;
}
}
private void logOnce(String format, Object arg) {
if (loggedOnce.compareAndSet(false, true)) {
log.info("idle-sleep guard: " + format, arg);
}
}
/** Wraps the live {@code caffeinate} child; {@link #close} force-destroys it, idempotently. */
private static final class CaffeinateAssertion implements SleepAssertion {
private final Process process;
CaffeinateAssertion(Process process) {
this.process = process;
}
@Override
public void close() {
if (!process.isAlive()) {
return;
}
process.destroy();
try {
if (!process.waitFor(2, TimeUnit.SECONDS)) {
process.destroyForcibly();
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
process.destroyForcibly();
}
}
}
}
@@ -0,0 +1,105 @@
package dev.ltms.fleet.power;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.function.IntSupplier;
/**
* Holds an OS-level assertion against idle sleep for exactly as long as at least one fleet
* member is live.
*
* <p><strong>Why this exists:</strong> a fleetd host was measured idle-sleeping after as little
* as one minute of inactivity (its {@code pmset -g custom} reports {@code sleep 1} on battery).
* Overnight the daemon's AMQP link to the broker dropped 13 times, and cross-checking every drop
* minute against {@code pmset -g log} found a sleep or wake event in the same minute or the one
* before, every time. The AMQP churn is only the visible symptom — the real problem is that a
* member mid-turn freezes with the host, and a long turn with nobody typing is exactly the case
* that goes idle.
*
* <p><strong>How it tracks "live":</strong> this is driven by {@code SessionManager}'s existing
* {@code onAcquire}/{@code onRelease} lifecycle hooks (added for CB-520/CB-516, previously wired
* to nothing but the reply inbox) rather than a second member count kept in parallel. Wire it as:
* <pre>{@code
* IdleSleepGuard guard = new IdleSleepGuard(mechanism, sessions::size);
* sessions.onAcquire(_ -> guard.recheck());
* sessions.onRelease(_ -> guard.recheck());
* }</pre>
* Every acquire/release event re-reads {@code SessionManager#size()} — the same registry {@code
* fleet_list}'s live/capacity numbers are themselves computed from — and only an actual 0→1 or
* 1→0 crossing touches the OS. A listener exception is already caught and logged by {@code
* SessionManager} itself (it must never let a listener failure block the acquire/release it is
* reacting to), so {@link #recheck()} does not need its own top-level try/catch to honor that.
*
* <p><strong>Failure posture:</strong> every method here is safe to call whether or not {@link
* SleepAssertionMechanism#acquire()} actually works. A mechanism that returns {@code null} (wrong
* platform, missing tool, spawn failure) simply means this guard never holds anything — it never
* throws and never blocks a spawn, a release, or shutdown.
*/
public final class IdleSleepGuard implements AutoCloseable {
private static final Logger log = LoggerFactory.getLogger(IdleSleepGuard.class);
private final SleepAssertionMechanism mechanism;
private final IntSupplier liveCount;
private final Object lock = new Object();
private SleepAssertion held;
public IdleSleepGuard(SleepAssertionMechanism mechanism, IntSupplier liveCount) {
this.mechanism = mechanism;
this.liveCount = liveCount;
}
/**
* Re-read the live count and acquire or release the held assertion to match: nothing held and
* at least one member live ⇒ acquire; something held and no member live ⇒ release. A steady
* count (still zero, still positive) is a no-op either way, so a single spawn or release only
* ever touches the OS on the crossing, not on every call.
*/
public void recheck() {
synchronized (lock) {
int live = liveCount.getAsInt();
if (live > 0 && held == null) {
held = mechanism.acquire();
if (held != null) {
log.debug("idle-sleep guard armed: {} live member(s)", live);
}
} else if (live == 0 && held != null) {
releaseHeldLocked();
}
}
}
/** {@code true} while an assertion is actually held. Exposed for tests. */
boolean isHeld() {
synchronized (lock) {
return held != null;
}
}
/**
* Release whatever is held, if anything. Idempotent and safe to call at any time, including
* repeatedly — a daemon shutdown hook calls this unconditionally so no assertion (and no
* {@code caffeinate} child) survives the process, even if the drain that would otherwise have
* driven the live count to zero was itself interrupted or threw.
*/
@Override
public void close() {
synchronized (lock) {
if (held != null) {
releaseHeldLocked();
}
}
}
/** Caller must hold {@link #lock}. */
private void releaseHeldLocked() {
try {
held.close();
} catch (RuntimeException e) {
log.warn("idle-sleep guard: failed to release its assertion cleanly: {}", e.toString());
} finally {
held = null;
}
}
}
@@ -0,0 +1,11 @@
package dev.ltms.fleet.power;
/**
* A held OS-level assertion against idle sleep. {@link #close} must be idempotent — safe to call
* more than once — and must never throw, matching {@link IdleSleepGuard}'s "never break the
* fleet" contract.
*/
public interface SleepAssertion extends AutoCloseable {
@Override
void close();
}
@@ -0,0 +1,21 @@
package dev.ltms.fleet.power;
/**
* The OS mechanism {@link IdleSleepGuard} uses to hold and release an idle-sleep assertion. This
* is the seam a test exercises instead of the real effect (a live {@code caffeinate} child) — see
* {@code IdleSleepGuardTest}.
*
* <p>Implementations must never throw. Every failure — wrong platform, missing tool, a spawn
* error — must show up as {@link #acquire()} returning {@code null}, so a caller can treat "no
* assertion held" and "the mechanism could not be used" identically and the fleet keeps running
* either way.
*/
public interface SleepAssertionMechanism {
/**
* Acquire a fresh assertion against idle sleep, or {@code null} when this mechanism is not
* usable right now (wrong platform, the tool is missing, the child process could not start).
* Never throws.
*/
SleepAssertion acquire();
}
@@ -92,15 +92,6 @@ public final class GitWorktrees implements Worktrees {
private final String configuredRoot;
/** OS group name for {@link #shareWithGroup} (fleetd #185 stage 3); {@code null} ⇒ feature off. */
private final String group;
/** Source directory of skill folders for {@link #seedSkills} (fleetd #362, {@code memberSkills:}
* in config); {@code null} ⇒ feature off, no worktree is touched beyond today's behaviour. */
private final String memberSkillsSource;
/** Extra environment merged into every {@code git} subprocess this instance runs. Always {@code
* Map.of()} from every production constructor. Test seam only (fleetd #362 review fix): lets
* {@code GitWorktreesTest} point {@code GIT_CONFIG_GLOBAL} at an isolated temp file so it can
* drive the real {@link #add} path against a controlled "operator's global git config" and
* prove the excludesFile composition below without ever touching the real machine's config. */
private final Map<String, String> gitEnv;
private final Consumer<String> afterWorktreeAdded;
/** How the initial {@code git worktree add} command runs. Package-private test seam for an
* interrupted command after Git has made worktree state. */
@@ -134,20 +125,7 @@ public final class GitWorktrees implements Worktrees {
* config); null/blank ⇒ {@link #shareWithGroup} is a no-op.
*/
public GitWorktrees(String configuredRoot, String group) {
this(configuredRoot, group, (String) null);
}
/**
* @param configuredRoot nullable absolute or relative path; null/blank derives a sibling of
* the repo root.
* @param group optional OS group name (fleetd #185 stage 3, {@code worktreeGroup:}
* in config); null/blank ⇒ {@link #shareWithGroup} is a no-op.
* @param memberSkillsSource fleetd #362: optional directory of skill folders ({@code
* memberSkills:} in config) copied into every provisioned worktree's
* {@code .claude/skills/}; null/blank ⇒ {@link #seedSkills} is a no-op.
*/
public GitWorktrees(String configuredRoot, String group, String memberSkillsSource) {
this(configuredRoot, group, _ -> {}, null, null, memberSkillsSource);
this(configuredRoot, group, _ -> {});
}
/** Test seam for changing a real worktree between its creation and its security check. */
@@ -157,20 +135,13 @@ public final class GitWorktrees implements Worktrees {
/** Test seam combining a configurable {@code group} with {@link #afterWorktreeAdded}. */
GitWorktrees(String configuredRoot, String group, Consumer<String> afterWorktreeAdded) {
this(configuredRoot, group, afterWorktreeAdded, null, null, null);
this(configuredRoot, group, afterWorktreeAdded, null, null);
}
/** Test seam for changing how {@link #shareWithGroup}'s processes run. */
GitWorktrees(String configuredRoot, String group, Consumer<String> afterWorktreeAdded,
Function<String[], String> shareGroupRunner) {
this(configuredRoot, group, afterWorktreeAdded, shareGroupRunner, null, null);
}
/** Test seam for changing how the initial {@code git worktree add} command runs, with no
* {@code memberSkillsSource} configured. */
GitWorktrees(String configuredRoot, String group, Consumer<String> afterWorktreeAdded,
Function<String[], String> shareGroupRunner, Function<String[], String> worktreeAddRunner) {
this(configuredRoot, group, afterWorktreeAdded, shareGroupRunner, worktreeAddRunner, null);
this(configuredRoot, group, afterWorktreeAdded, shareGroupRunner, null);
}
/**
@@ -180,30 +151,11 @@ public final class GitWorktrees implements Worktrees {
*
* @param shareGroupRunner {@code null} ⇒ the real {@link #exec(String...)}.
* @param worktreeAddRunner {@code null} ⇒ the real {@link #exec(String...)}.
* @param memberSkillsSource {@code null}/blank ⇒ {@link #seedSkills} is a no-op.
*/
GitWorktrees(String configuredRoot, String group, Consumer<String> afterWorktreeAdded,
Function<String[], String> shareGroupRunner, Function<String[], String> worktreeAddRunner,
String memberSkillsSource) {
this(configuredRoot, group, afterWorktreeAdded, shareGroupRunner, worktreeAddRunner,
memberSkillsSource, Map.of());
}
/**
* Full test seam, plus {@code gitEnv} (fleetd #362 review fix, verification only): extra
* environment merged into every {@code git} subprocess this instance runs, so a test can isolate
* something like {@code GIT_CONFIG_GLOBAL} from the real machine while still driving the real
* {@link #add} path end to end. Every production constructor above delegates here with {@code
* Map.of()}.
*/
GitWorktrees(String configuredRoot, String group, Consumer<String> afterWorktreeAdded,
Function<String[], String> shareGroupRunner, Function<String[], String> worktreeAddRunner,
String memberSkillsSource, Map<String, String> gitEnv) {
Function<String[], String> shareGroupRunner, Function<String[], String> worktreeAddRunner) {
this.configuredRoot = configuredRoot;
this.group = (group == null || group.isBlank()) ? null : group;
this.memberSkillsSource = (memberSkillsSource == null || memberSkillsSource.isBlank())
? null : memberSkillsSource;
this.gitEnv = gitEnv == null ? Map.of() : gitEnv;
this.afterWorktreeAdded = afterWorktreeAdded == null ? _ -> {} : afterWorktreeAdded;
this.shareGroupRunner = shareGroupRunner != null ? shareGroupRunner : this::exec;
this.worktreeAddRunner = worktreeAddRunner != null ? worktreeAddRunner : this::exec;
@@ -232,7 +184,6 @@ public final class GitWorktrees implements Worktrees {
configureEnvironmentCredentialHelper(repoRoot, wt);
configureHttpsUrlRewriteForSshOrigin(repoRoot, wt);
isolateToolSurface(wt);
seedSkills(wt);
} catch (RuntimeException e) {
cleanupAfterAddFailure(repoRoot, wt, branch, e);
throw e;
@@ -623,310 +574,6 @@ public final class GitWorktrees implements Worktrees {
return true;
}
/**
* fleetd #362: copy each skill folder from the configured {@link #memberSkillsSource} directory
* into {@code <worktreePath>/.claude/skills/}, so a member spawned against ANY repo — not only
* one that already ships its own {@code .claude/skills/} — can load a bridge skill such as
* {@code implementer}. Every brief this fleet sends starts with {@code "Load the <name>
* skill."}; outside a repo carrying its own copy that line was previously a no-op.
*
* <p><b>No-op — nothing read, nothing written, nothing logged</b> — when {@link
* #memberSkillsSource} is null/blank (today's default), the same off-switch shape as
* {@link #shareWithGroup}.
*
* <p><b>Invariant 1 — a repo's own skill wins.</b> A skill folder already present at
* {@code <worktreePath>/.claude/skills/<name>} — because the just-checked-out branch commits its
* own copy — is left completely untouched: never overwritten, and never even opened.
*
* <p><b>Invariant 2 — a seeded skill can never end up in a worker's commit.</b> Every path this
* writes is untracked in the target repo (that is the whole reason it is being seeded), so
* {@code git status} would otherwise show each one as a new, addable, committable path. The
* repo-wide {@code .git/info/exclude} is NOT used for this: measured against a real linked
* worktree, that file resolves to the repository's COMMON git dir even from a worktree (the
* same file {@link dev.ltms.fleet.member.ClaudeCodeLauncher#writeIdeOverlay writeIdeOverlay}
* appends {@code CLAUDE.local.md} to), so an entry written there would hide the seeded skill
* from {@code git status} in the PRIMARY's own checkout and every sibling worktree too — not
* only this one. Instead, {@link #excludeSeededSkillsFromGitStatus} points {@code
* core.excludesFile} at a file scoped {@code --worktree} (the same {@code
* extensions.worktreeConfig} mechanism {@link #configureEnvironmentCredentialHelper} already
* relies on) that itself lives under this worktree's own private git dir ({@code
* .git/worktrees/<nonce>/}, OUTSIDE the working tree) — invisible to this worktree's {@code git
* status} and structurally impossible for this worktree to commit, with no effect on any other
* worktree or the primary checkout. Proven with a real {@code git status --porcelain} in
* {@code GitWorktreesTest}, not by reasoning.
*
* <p><b>Compose, don't replace.</b> {@code core.excludesFile} is single-valued: the first cut of
* this method pointed it at fleetd's own file with {@code --replace-all}, which SHADOWS whatever
* the operator's own (global, or repo-local) {@code core.excludesFile} was already resolving to
* inside this worktree, rather than adding to it. Measured concretely: this repo's own {@code
* .gitignore} does not ignore {@code target/} — only an operator's global excludesFile does — so
* every worker's {@code mvn clean install} would otherwise make {@code target/} appear as
* untracked, and {@link #hasUncommitted}'s deliberately-untracked-inclusive {@code git status
* --porcelain} (CB-576) would then read every such worktree as dirty forever, so it is never
* cleaned up. {@link #excludeSeededSkillsFromGitStatus} now reads whatever {@code
* core.excludesFile} resolves to BEFORE writing anything (falling back to git's own documented
* default, {@code $XDG_CONFIG_HOME/git/ignore} or {@code $HOME/.config/git/ignore}, when the key
* is unset entirely — see {@code gitignore(5)}), and writes that content into its OWN exclude
* file ahead of the seeded skill patterns, so every operator-configured pattern keeps applying
* inside the seeded worktree exactly as it did before seeding ran.
*
* <p>Instead of using worktree-scoped-config as an add-then-append (a second key does not exist
* for {@code core.excludesFile} — it takes exactly one value), an actual second exclude source
* was ruled out because git resolves only ONE {@code core.excludesFile}; concatenating the prior
* content into fleetd's own file is what "compose" reduces to for a single-valued key.
*
* <p>Proven the same way as invariant 2's own leak check: {@code
* GitWorktreesTest#seedSkillsComposesWithAnAlreadyEffectiveGlobalExcludesFile} isolates a
* synthetic "operator's global config" via {@code GIT_CONFIG_GLOBAL} (never the real machine's),
* seeds a skill, and asserts {@code git status --porcelain} is still empty for a file matching
* that global config's own ignore pattern.
*
* <p><b>Invariant 3 — best-effort.</b> A missing/unreadable {@link #memberSkillsSource}, or a
* copy/exclude failure, is logged and skipped — it must never fail the spawn, the same contract
* {@link #overlayParity} and {@link #isolateToolSurface} already hold.
*
* <p>Recorded for the worker itself the same way fleetd #134 records {@code
* fleet.neutralizedConfig}: {@code fleet.seededSkills} (one value per seeded skill folder) and
* {@code fleet.seededSkillsNote}, readable with {@code git config --worktree --get-all
* fleet.seededSkills}.
*
* <p><b>Claude Code specific by construction, not by a backend check here.</b> Only {@code
* .claude/skills/<name>/SKILL.md} is a path any launcher reads today (opencode's equivalent is a
* different shape under {@code .opencode/agent}, out of scope — see issue #362). This method
* only copies files; like {@link #isolateToolSurface} — which neutralizes BOTH {@code .mcp.json}
* and {@code opencode.json} unconditionally — it runs the same for every worktree regardless of
* which backend ultimately spawns into it, because the backend is not yet chosen at {@link #add}
* time. A seeded {@code .claude/skills/} directory in an opencode member's worktree is simply
* never read by that launcher.
*/
private void seedSkills(String worktreePath) {
if (memberSkillsSource == null) {
return;
}
Path source = Path.of(memberSkillsSource).toAbsolutePath().normalize();
if (!Files.isDirectory(source)) {
log.warn("memberSkills source '{}' is not a directory — skipping skill seeding for worktree {}",
source, worktreePath);
return;
}
Path skillsRoot = Path.of(worktreePath).resolve(".claude").resolve("skills");
List<String> seeded = new ArrayList<>();
List<String> kept = new ArrayList<>();
try (var candidates = Files.list(source)) {
for (Path candidate : candidates
.filter(Files::isDirectory)
.filter(p -> !p.getFileName().toString().startsWith("."))
.sorted()
.toList()) {
String name = candidate.getFileName().toString();
Path dst = skillsRoot.resolve(name);
if (Files.exists(dst)) {
kept.add(name);
continue;
}
copySkillDirectory(candidate, dst);
seeded.add(name);
}
} catch (IOException | RuntimeException e) {
log.warn("failed to seed skills into worktree {} from memberSkills source '{}': {}",
worktreePath, source, e.getMessage());
return;
}
String detail = seeded.isEmpty() ? "" : "seeded: " + String.join(", ", seeded);
if (!kept.isEmpty()) {
detail += (detail.isEmpty() ? "" : "; ") + "kept the repo's own copy of: " + String.join(", ", kept);
}
if (detail.isEmpty()) {
detail = "no skill folders found under " + source;
}
log.info("skill seeding: {} of {} candidate(s) from {} into {}/.claude/skills — {}",
seeded.size(), seeded.size() + kept.size(), source, worktreePath, detail);
if (seeded.isEmpty()) {
return;
}
try {
excludeSeededSkillsFromGitStatus(worktreePath, seeded);
recordSeededSkillsForWorker(worktreePath, seeded);
} catch (RuntimeException e) {
log.warn("seeded skill(s) {} into {} but could not hide them from git status: {} — "
+ "they may show as untracked; never commit them", seeded, worktreePath, e.getMessage());
}
}
/** Recursively copy a skill folder ({@code src}) into a fresh destination ({@code dst}) that
* {@link #seedSkills} has already confirmed does not exist, preserving the directory structure
* (e.g. {@code implementer/SKILL.md}, {@code implementer/references/...}). */
private static void copySkillDirectory(Path src, Path dst) {
try (var walk = Files.walk(src)) {
for (Path path : walk.sorted().toList()) {
Path target = dst.resolve(src.relativize(path).toString());
if (Files.isDirectory(path)) {
Files.createDirectories(target);
} else {
Files.createDirectories(target.getParent());
Files.copy(path, target, StandardCopyOption.COPY_ATTRIBUTES);
}
}
} catch (IOException e) {
throw new WorktreeException("cannot copy skill directory " + src + " -> " + dst + ": "
+ e.getMessage(), e);
}
}
/**
* Make every path in {@code seededSkillNames} (each a name under {@code .claude/skills/})
* invisible to {@code git status} in THIS worktree only — see the invariant-2 discussion on
* {@link #seedSkills}. Sets {@code core.excludesFile} scoped {@code --worktree} to a file
* written under this worktree's own private git dir ({@code git rev-parse
* --absolute-git-dir}), which lives outside the working tree, so the exclude file itself can
* never be committed either.
*
* <p><b>Compose, don't replace.</b> {@code core.excludesFile} is single-valued, so pointing it at
* fleetd's own file would otherwise SHADOW whatever excludesFile this worktree was already
* resolving (an operator's global config, most commonly) rather than add to it — see the
* "Compose, don't replace" discussion on {@link #seedSkills}. {@link
* #previouslyEffectiveExcludesFileContent} is read BEFORE this method's own {@code --worktree}
* write below, so it still sees whatever was effective beforehand; that content is written into
* fleetd's own exclude file ahead of the seeded skill patterns, and the worktree-scoped override
* then points at that combined file — so every pattern the operator's own configuration already
* applied keeps applying, plus the seeded skill paths.
*
* <p><b>Assumes a fresh worktree — not idempotent.</b> {@link #seedSkills} only ever calls this
* from {@link #add}, which always creates a brand-new worktree, so {@code core.excludesFile} is
* never already worktree-scoped-set to fleetd's own file when this runs. A hypothetical second
* call on the SAME worktree would read fleetd's own already-composed file back as "previously
* effective" (worktree scope now wins) and append the seeded patterns a second time — harmless
* to {@code git status} (duplicate exclude lines are a no-op), but not something to rely on. No
* guard is added for this because the path does not exist today; if a future caller ever seeds
* the same worktree twice, it will need one.
*/
private void excludeSeededSkillsFromGitStatus(String worktreePath, List<String> seededSkillNames) {
exec("git", "-C", worktreePath, "config", "extensions.worktreeConfig", "true");
String previouslyEffective = previouslyEffectiveExcludesFileContent(worktreePath);
String gitDir = exec("git", "-C", worktreePath, "rev-parse", "--absolute-git-dir").trim();
Path excludeFile = Path.of(gitDir, "fleet-seeded-skills-exclude");
StringBuilder patterns = new StringBuilder();
if (!previouslyEffective.isEmpty()) {
patterns.append(previouslyEffective);
}
for (String name : seededSkillNames) {
patterns.append("/.claude/skills/").append(name).append('/').append(System.lineSeparator());
}
try {
Files.writeString(excludeFile, patterns.toString());
} catch (IOException e) {
throw new WorktreeException("cannot write skills exclude file " + excludeFile + ": "
+ e.getMessage(), e);
}
exec("git", "-C", worktreePath, "config", "--worktree", "--replace-all", "core.excludesFile",
excludeFile.toString());
}
/**
* The content of whatever {@code core.excludesFile} resolves to for {@code worktreePath} right
* now — BEFORE {@link #excludeSeededSkillsFromGitStatus} points that key at fleetd's own file —
* so it can be carried forward instead of shadowed. {@code --type=path} makes git itself perform
* {@code ~}/{@code ~user} expansion the same way it would when actually reading the key to build
* exclude rules, rather than handing back a raw, unexpanded config string.
*
* <p>When the key is unset entirely (exit code non-zero), falls back to git's own documented
* default excludes file — {@code $XDG_CONFIG_HOME/git/ignore}, or {@code
* $HOME/.config/git/ignore} when that variable is unset — per {@code gitignore(5)}: git applies
* that file even with no {@code core.excludesFile} configured at all, so skipping it here would
* silently drop patterns an operator never had to configure to get.
*
* <p>Never throws: a missing, unreadable, or unresolvable file is treated as "nothing to carry
* forward" (empty string) — this is a best-effort read in service of {@link #seedSkills}'s own
* invariant 3, not a new way for skill seeding to fail a spawn.
*
* <p><b>Review fix, finding 2.</b> The XDG-fallback branch below does not go through {@code git}
* at all, so a first cut of it read {@code XDG_CONFIG_HOME}/{@code HOME} straight from the JVM's
* own environment ({@link System#getenv} / {@code user.home}) — unlike every other value this
* class resolves, which goes through a {@code git} subprocess and therefore already honours
* {@link #gitEnv}. That meant no test could make this branch hermetic, and on any machine
* carrying a real {@code ~/.config/git/ignore} (this repo's own dev machine does), every
* skill-seeding test silently composed with that real file — correct in production, but
* machine-dependent in the test suite, and a future broader pattern in that real file could
* silently change what a seeded worktree's {@code git status} reports depending on whose home
* directory ran the test. {@link #resolveEnv} now checks {@link #gitEnv} first for both
* variables, falling back to the JVM's real environment only when the seam does not supply
* them — production behaviour (empty {@link #gitEnv}) is unchanged, and a test can now isolate
* this branch exactly as it already isolates every {@code git} subprocess call.
*
* <p><b>Snapshot, not a reference.</b> The content below is read once, at seeding time, and
* copied into fleetd's own exclude file. If the operator edits their global excludesFile
* afterward, an already-seeded worktree keeps the old copy — acceptable for a worktree's
* expected lifetime, but worth knowing before reading a stale pattern as a bug.
*
* @return the file's content, trailing-newline-normalized, or {@code ""} when there is nothing
* to compose with.
*/
private String previouslyEffectiveExcludesFileContent(String worktreePath) {
String resolvedPath;
if (exitCode("git", "-C", worktreePath, "config", "--get", "--type=path", "core.excludesFile") == 0) {
resolvedPath = exec("git", "-C", worktreePath, "config", "--get", "--type=path",
"core.excludesFile").trim();
} else {
String xdgConfigHome = resolveEnv("XDG_CONFIG_HOME");
Path fallback = (xdgConfigHome != null && !xdgConfigHome.isBlank())
? Path.of(xdgConfigHome, "git", "ignore")
: Path.of(resolveHome(), ".config", "git", "ignore");
resolvedPath = fallback.toString();
}
if (resolvedPath.isBlank()) {
return "";
}
Path file = Path.of(resolvedPath);
if (!Files.isRegularFile(file) || !Files.isReadable(file)) {
return "";
}
try {
String content = Files.readString(file);
return content.isBlank() ? "" : content.stripTrailing() + System.lineSeparator();
} catch (IOException e) {
log.warn("could not read previously-effective excludesFile {} while seeding skills into "
+ "{}: {} — its patterns will not carry forward into the seeded worktree",
file, worktreePath, e.getMessage());
return "";
}
}
/**
* Resolve environment variable {@code name} for {@link #previouslyEffectiveExcludesFileContent}'s
* XDG fallback, checking {@link #gitEnv} FIRST so a test can isolate this the same way it
* already isolates every {@code git} subprocess this class runs, and falling back to the JVM's
* real environment only when the seam does not supply it (always the case in production, where
* {@link #gitEnv} is {@code Map.of()}).
*/
private String resolveEnv(String name) {
String fromSeam = gitEnv.get(name);
return fromSeam != null ? fromSeam : System.getenv(name);
}
/** Same as {@link #resolveEnv(String)}, for {@code HOME} — falls back to {@code user.home}
* (rather than {@code System.getenv("HOME")}) when the seam does not supply it, matching this
* class's pre-existing behaviour for every other home-directory resolution. */
private String resolveHome() {
String fromSeam = gitEnv.get("HOME");
return fromSeam != null ? fromSeam : System.getProperty("user.home");
}
/**
* The worker-readable half of fleetd #362, mirroring {@link #recordNeutralizedConfigForWorker}:
* record which skill folders were seeded where the worker itself can read it, without a
* working-tree file that would show up in {@code git status}.
*/
private void recordSeededSkillsForWorker(String worktreePath, List<String> seeded) {
exec("git", "-C", worktreePath, "config", "extensions.worktreeConfig", "true");
for (String name : seeded) {
exec("git", "-C", worktreePath, "config", "--worktree", "--add", "fleet.seededSkills", name);
}
exec("git", "-C", worktreePath, "config", "--worktree", "fleet.seededSkillsNote",
"each fleet.seededSkills value names a skill folder fleetd copied into "
+ ".claude/skills/ because this repo did not already ship it; it is excluded "
+ "from git status (core.excludesFile, worktree-scoped) and must never be committed");
}
@Override
public void remove(String repoRoot, String worktreePath) {
Path p = Path.of(worktreePath);
@@ -1437,9 +1084,6 @@ public final class GitWorktrees implements Worktrees {
Process p;
try {
ProcessBuilder pb = new ProcessBuilder(command).redirectErrorStream(true);
if (!gitEnv.isEmpty()) {
pb.environment().putAll(gitEnv);
}
if (extraEnv != null && !extraEnv.isEmpty()) {
pb.environment().putAll(extraEnv);
}
@@ -1475,11 +1119,7 @@ public final class GitWorktrees implements Worktrees {
private int exitCode(String... command) {
Process p;
try {
ProcessBuilder pb = new ProcessBuilder(command).redirectErrorStream(true);
if (!gitEnv.isEmpty()) {
pb.environment().putAll(gitEnv);
}
p = pb.start();
p = new ProcessBuilder(command).redirectErrorStream(true).start();
} catch (IOException e) {
throw new WorktreeException("failed to start " + command[0] + ": " + e.getMessage(), e);
}
@@ -80,7 +80,7 @@ class FleetdLeadMailboxSelectionTest {
@Test
void opensTheMailboxWhenAUriAndSelfIdAreConfigured() {
var opener = new RecordingOpener();
var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, "mac-opus", null, null);
var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, "mac-opus", null);
Fleetd.openLeadMailbox(coordinator, Map.of(), opener);
@@ -94,7 +94,7 @@ class FleetdLeadMailboxSelectionTest {
void honoursUriEnvOverALiteralUri() {
var opener = new RecordingOpener();
var coordinator = new FleetConfig.Coordinator("amqp://stale:stale@old:5672/x", "COORD_URI",
"mac-opus", 8, null);
"mac-opus", 8);
Fleetd.openLeadMailbox(coordinator, Map.of("COORD_URI", RESOLVED_URI), opener);
@@ -105,7 +105,7 @@ class FleetdLeadMailboxSelectionTest {
@Test
void turnsOffWhenUriEnvDoesNotResolve() {
var opener = new RecordingOpener();
var coordinator = new FleetConfig.Coordinator(null, "COORD_URI", "mac-opus", null, null);
var coordinator = new FleetConfig.Coordinator(null, "COORD_URI", "mac-opus", null);
assertNull(Fleetd.openLeadMailbox(coordinator, Map.of(), opener));
@@ -116,7 +116,7 @@ class FleetdLeadMailboxSelectionTest {
void warnsAndStaysOffWhenSelfIdIsMissing() {
var appender = captureFleetdLogs();
var opener = new RecordingOpener();
var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, null, null, null);
var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, null, null);
assertNull(Fleetd.openLeadMailbox(coordinator, Map.of(), opener));
@@ -131,7 +131,7 @@ class FleetdLeadMailboxSelectionTest {
var appender = captureFleetdLogs();
var opener = new RecordingOpener();
opener.unreachable = true;
var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, "mac-opus", null, null);
var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, "mac-opus", null);
assertNull(Fleetd.openLeadMailbox(coordinator, Map.of(), opener),
"a down coordination broker turns the feature off; it must never take the daemon down");
@@ -104,10 +104,10 @@ class ConfigRefTopLevelReportingCoverageTest {
v.put("configReload", new FleetConfig.ConfigReload(true, 10));
v.put("quarantineCooldownSeconds", 1800);
v.put("memberCredentials", null);
v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-a", null, "self-a", 1, null));
v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-a", null, "self-a", 1));
v.put("worktreeGroup", "group-a");
v.put("memberLoginShell", null);
v.put("memberSkills", "/skills/a");
v.put("idleSleepGuard", new FleetConfig.IdleSleepGuard(true));
assertNamesMatchComponents(v);
return v;
}
@@ -145,10 +145,10 @@ class ConfigRefTopLevelReportingCoverageTest {
v.put("configReload", new FleetConfig.ConfigReload(false, 20));
v.put("quarantineCooldownSeconds", 3600);
v.put("memberCredentials", null);
v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-b", null, "self-b", 2, null));
v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-b", null, "self-b", 2));
v.put("worktreeGroup", "group-b");
v.put("memberLoginShell", null);
v.put("memberSkills", "/skills/b");
v.put("idleSleepGuard", new FleetConfig.IdleSleepGuard(false));
assertNamesMatchComponents(v);
return v;
}
@@ -1189,7 +1189,7 @@ class FleetConfigTest {
@Test
void coordinatorEffectiveUriHonorsUriEnv() {
FleetConfig.Coordinator withEnv = new FleetConfig.Coordinator(
"amqp://stale-clear-text@127.0.0.1:5672/coord", "LEAD_COORD_URI", "fleet01-lead", null, null);
"amqp://stale-clear-text@127.0.0.1:5672/coord", "LEAD_COORD_URI", "fleet01-lead", null);
assertEquals("amqp://from-env@127.0.0.1:5672/coord",
withEnv.effectiveUri(Map.of("LEAD_COORD_URI", "amqp://from-env@127.0.0.1:5672/coord")),
@@ -1200,61 +1200,12 @@ class FleetConfigTest {
"a blank uriEnv variable must not fall back to the literal uri");
FleetConfig.Coordinator noEnv = new FleetConfig.Coordinator(
"amqp://guest:guest@127.0.0.1:5672/coord", null, null, null, null);
"amqp://guest:guest@127.0.0.1:5672/coord", null, null, null);
assertEquals("amqp://guest:guest@127.0.0.1:5672/coord", noEnv.effectiveUri(Map.of()),
"the literal uri is used when no uriEnv is configured");
assertEquals(LeadMailbox.DEFAULT_PREFETCH, noEnv.prefetchOrDefault());
}
/**
* fleetd #361: the live block ships as just {@code uriEnv} + {@code selfId} (see
* {@code Fleetd.example.yaml} / the operator's real {@code fleetd.yaml}, gitignored). That exact
* shape, with no {@code peers:} key at all, must keep parsing unchanged after this field is added.
*/
@Test
void coordinatorBlockWithNoPeersKeyStillParses(@TempDir Path dir) throws Exception {
Path f = dir.resolve("coordinator-no-peers.yaml");
Files.writeString(f, """
bind:
port: 8080
coordinator:
uriEnv: COORD_AMQP_URI
selfId: mac
""");
FleetConfig cfg = FleetConfig.load(f);
assertNotNull(cfg.coordinator());
assertEquals("mac", cfg.coordinator().selfId());
assertEquals(List.of(), cfg.coordinator().peers(), "no peers: key means no configured peers, never null");
}
@Test
void coordinatorPeersParsesAndDropsBlankEntries(@TempDir Path dir) throws Exception {
Path f = dir.resolve("coordinator-peers.yaml");
Files.writeString(f, """
bind:
port: 8080
coordinator:
selfId: mac
uri: amqp://guest:guest@127.0.0.1:5672/coord
peers:
- fleet01
- ""
- fleet02
""");
FleetConfig cfg = FleetConfig.load(f);
assertEquals(List.of("fleet01", "fleet02"), cfg.coordinator().peers(),
"a blank peer entry must be dropped, never kept as an empty coord-id");
}
@Test
void coordinatorPeersDefaultsToEmptyWhenConstructedWithNull() {
FleetConfig.Coordinator c = new FleetConfig.Coordinator(
"amqp://guest:guest@127.0.0.1:5672/coord", null, "mac", null, null);
assertEquals(List.of(), c.peers(), "a null peers list must default to empty, never NPE downstream");
}
@Test
void absentWorktreeGroupLeavesItNull(@TempDir Path dir) throws Exception {
Path f = dir.resolve("no-worktree-group.yaml");
@@ -2629,4 +2580,58 @@ class FleetConfigTest {
"with no pool to choose from, every configured profile is a candidate and the "
+ "first one wins");
}
// ── idle-sleep guard: default-on config block ───────────────────────────────────────────────
@Test
void idleSleepGuardIsOnByDefaultWhenTheBlockIsEntirelyAbsent(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
bind:
host: 127.0.0.1
port: 8080
""");
FleetConfig cfg = FleetConfig.load(f);
assertNull(cfg.idleSleepGuard(), "an absent block parses to null, unlike most other blocks here");
// The block itself is absent, but the FEATURE stays on: Fleetd treats a null block the
// same as enabled: true (see FleetConfig.idleSleepGuard's javadoc) — this test only pins
// the parse result, the on-by-default behaviour is Fleetd's own null check.
}
@Test
void idleSleepGuardExplicitlyEnabledIsOn(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
idleSleepGuard:
enabled: true
""");
FleetConfig cfg = FleetConfig.load(f);
assertTrue(cfg.idleSleepGuard().isEnabled());
}
@Test
void idleSleepGuardExplicitlyDisabledIsOff(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
idleSleepGuard:
enabled: false
""");
FleetConfig cfg = FleetConfig.load(f);
assertFalse(cfg.idleSleepGuard().isEnabled());
}
@Test
void idleSleepGuardBlockPresentButEmptyDefaultsToEnabled(@TempDir Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
idleSleepGuard: {}
""");
FleetConfig cfg = FleetConfig.load(f);
assertTrue(cfg.idleSleepGuard().isEnabled(),
"unlike ConfigReload/Health, this block defaults to ON even when present but empty");
}
}
@@ -1,193 +0,0 @@
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;
/**
* Guards against a "defect factory" built into this file's own established pattern, found live
* while building the (parked) idle-sleep-guard PR: every time a component is added to
* {@link FleetConfig}, the record grows by one arg AND a new back-compat constructor is added at
* the OLD arity, so existing callers keep compiling. That is correct and required — see the
* constructor ladder just below the record header. But {@link #withDefaults()}'s own {@code return
* new FleetConfig(...)} call sits in this same file, written at a literal argument count. The very
* next time a component is added, the freshly-added back-compat constructor at the OLD arity
* silently captures that stale call, because it is now a legal overload at that arg count too. It
* compiles. Every other test passes, because nothing else exercises the new field. The new
* component is defaulted away — {@code null}, or whatever that back-compat overload defaults it to
* — on every {@link FleetConfig#load}. Measured, not theoretical: this exact sequence happened
* live when the {@code idleSleepGuard} component was added on a sibling branch; it was caught only
* because that branch's own new tests happened to assert on the new field's value.
*
* <p>This test proves the opposite property, and does it in a way that survives the next field
* being added without being rewritten: reflectively enumerate {@link FleetConfig}'s own record
* components (never a hardcoded count — the arity is exactly what changes over time), build one
* config through the true canonical constructor with a real, distinctive, non-null value in EVERY
* component (reusing the exact reflective-construction pattern
* {@link ConfigRefTopLevelReportingCoverageTest} already established for this file:
* {@code getDeclaredConstructor(exact record-component types)}, which resolves the canonical
* constructor by its true shape, not by binding to whichever overload happens to match arg count —
* the same way Jackson resolves it), call the real {@link FleetConfig#withDefaults()}, and assert
* every one of those values survives unchanged.
*
* <p>Why this is a valid check for every component, not just some: {@link #withDefaults()}'s own
* comments document that it only ever REPLACES a component when the incoming value is {@code null}
* (or blank, for {@code placement}) — {@code broker}/{@code primary}/{@code leadHeartbeat}/
* {@code configReload}/{@code coordinator}/{@code worktreeGroup}/{@code memberLoginShell}/
* {@code memberSkills} are left as-is unconditionally, and {@code bind}/{@code guard}/{@code lifecycle}/{@code auth}/
* {@code fleet}/{@code quarantineCooldownSeconds}/{@code memberCredentials}/{@code placement} are
* replaced only on null/blank input. A value that is never null or blank going in must therefore
* never change coming out, for every current component. No exclusion is needed today.
*
* <p>{@link #EXCLUDED_FROM_SURVIVAL_CHECK} exists anyway, kept deliberately empty and size-pinned
* by {@link #exclusionListSizeIsPinned()}: a future component that {@code withDefaults()} is
* <em>documented</em> to transform unconditionally (unlike every field today) would legitimately
* need one. Pinning the size at 0 means growing that set to make a failure go away is itself a
* visible diff to this test, not a silent one — a checker that can be silenced by adding to its
* own escape hatch is not a checker.
*/
class FleetConfigWithDefaultsPreservesEveryComponentTest {
private static final RecordComponent[] COMPONENTS = FleetConfig.class.getRecordComponents();
/** See the class javadoc — deliberately empty today; grow it only with a matching justification. */
private static final Set<String> EXCLUDED_FROM_SURVIVAL_CHECK = Set.of();
/** One real, distinctive, non-null (non-blank where blankness would mean "unset") value per component. */
private static Map<String, Object> baseValues() {
Map<String, Object> v = new LinkedHashMap<>();
v.put("bind", new FleetConfig.Bind("127.0.0.1", 8765));
v.put("herdrSocket", "~/.config/herdr/guard.sock");
v.put("memberHerdrSocket", "~/.config/herdr/member-guard.sock");
v.put("profiles", Map.of("sonnet", minimalProfile("sonnet")));
v.put("guard", new FleetConfig.Guard(List.of("host-guard")));
v.put("worktreeRoot", "/wt/guard");
v.put("lifecycle", new FleetConfig.Lifecycle(300, 5, 30, true));
v.put("spawnReadyTimeoutMs", 12_345);
v.put("spawnReadyPollMs", 234);
v.put("broker", new FleetConfig.Broker("amqp://guard", null, 7));
v.put("primary", new FleetConfig.Primary("term-guard", 4, 4000));
v.put("fleet", new FleetConfig.Fleet(
Map.of("opus", new FleetConfig.Leader("sonnet", "lead: opus-guard", 1, null, 10,
"claude", null, null, null)),
Map.of(), Map.of(), Map.of(), Map.of(), "{role}: {profile} #{n}"));
v.put("leadHeartbeat", new FleetConfig.LeadHeartbeat(301, 61_000L, 4));
v.put("health", new FleetConfig.Health(true, 31, 601, 61, null));
v.put("placement", "round-robin");
v.put("auth", new FleetConfig.Auth("loopback-trust", null));
v.put("configReload", new FleetConfig.ConfigReload(true, 11));
v.put("quarantineCooldownSeconds", 1801);
v.put("memberCredentials", new FleetConfig.MemberCredentials(
FleetConfig.MemberCredentials.POLICY_DENY_BY_DEFAULT,
List.of("git"), List.of("git", "ssh"), null));
v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-guard", null, "self-guard", 3, null));
v.put("worktreeGroup", "group-guard");
v.put("memberLoginShell", "/bin/zsh");
v.put("memberSkills", "/skills/guard");
assertNamesMatchComponents(v);
return v;
}
/** A minimal, otherwise-null {@link FleetConfig.Profile} — just enough to name one in a map. */
private static FleetConfig.Profile minimalProfile(String name) {
return new FleetConfig.Profile(name, null, null, null, null, null, null, null, null, null,
null, null, null, null, null, null, null, null, null, null, null, null, null, null,
null, null);
}
/**
* Guards {@link #baseValues()} itself against drifting from the record's real shape — the same
* assurance {@link ConfigRefTopLevelReportingCoverageTest} already relies on. This is what makes
* "no hardcoded arity" true in practice: 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's actual top-level components — "
+ "update baseValues() alongside the record");
}
/**
* Builds a {@link FleetConfig} 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 FleetConfig(...)} call risks doing.
*/
private static FleetConfig configOf(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> ctor = FleetConfig.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 withDefaults()}'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 for the one String — {@code placement} — where blank has
* meaning), so none of it should be replaced by {@code withDefaults()}; any component that
* comes back different was silently dropped.
*/
@Test
void everyComponentGivenARealValueSurvivesWithDefaults() throws ReflectiveOperationException {
Map<String, Object> base = baseValues();
FleetConfig config = configOf(base);
FleetConfig defaulted = config.withDefaults();
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 = base.get(name);
Object actual;
try {
actual = rc.getAccessor().invoke(defaulted);
} catch (ReflectiveOperationException e) {
throw new RuntimeException("failed to read FleetConfig." + name + "()", e);
}
if (!Objects.equals(expected, actual)) {
dropped.add(String.format(Locale.ROOT,
"%s: withDefaults() was given a real, non-null value (%s) for '%s' but "
+ "returned %s — a component silently dropped by withDefaults(), the "
+ "shape of the defect this test exists to catch (its final "
+ "\"return new FleetConfig(...)\" call binding to a back-compat "
+ "constructor instead of the true canonical one)",
name, expected, name, actual));
}
}
System.out.printf(Locale.ROOT,
"FleetConfig.withDefaults() 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,
"withDefaults() silently dropped these real, given components: " + dropped);
}
}
@@ -6,7 +6,6 @@ import org.junit.jupiter.api.Test;
import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
@@ -36,12 +35,6 @@ class LeadTabScannerTest {
final Map<String, String[]> tabs = new LinkedHashMap<>();
/** pane_id → [tab_id, terminal_id]. */
final Map<String, String[]> panes = new LinkedHashMap<>();
/**
* tab_id → whether herdr reports a running agent there. Defaults to {@code true} for every
* tab that has a pane, so every existing fixture keeps meaning "a live lead" unless a test
* says otherwise via {@link #deadAgent}.
*/
final Set<String> deadTabs = new LinkedHashSet<>();
int calls;
boolean failing;
@@ -60,18 +53,6 @@ class LeadTabScannerTest {
return this;
}
/** Mark {@code tabId} as labelled but agent-less — a dead lead's leftover tab (fleetd #359). */
TopologyHerdr deadAgent(String tabId) {
deadTabs.add(tabId);
return this;
}
/** Undo {@link #deadAgent} — models {@code agent.list} reporting the tab live again. */
TopologyHerdr reviveAgent(String tabId) {
deadTabs.remove(tabId);
return this;
}
@Override
public JsonNode call(String method, Object params) {
calls++;
@@ -102,19 +83,6 @@ class LeadTabScannerTest {
.formatted(id, p[0], p[1])));
return read("{\"panes\":[%s]}".formatted(String.join(",", items)));
}
case "agent.list" -> {
// One agent per distinct tab that has a pane and isn't marked dead — mirrors
// AgentControl.list()'s "agents" shape closely enough for the scanner's join,
// which only reads tab_id off each entry.
Set<String> seen = new LinkedHashSet<>();
panes.forEach((paneId, p) -> {
String tabId = p[0];
if (!deadTabs.contains(tabId) && seen.add(tabId)) {
items.add("{\"tab_id\":\"%s\"}".formatted(tabId));
}
});
return read("{\"agents\":[%s]}".formatted(String.join(",", items)));
}
default -> throw new AssertionError("unexpected herdr call: " + method);
}
}
@@ -240,114 +208,6 @@ class LeadTabScannerTest {
scanner(herdr, twoLeadsConfigured(), new AtomicLong()).get().get("term_opus_split"));
}
// ── fleetd #359: a labelled tab is only a lead when something is running in it ─────────────────
/**
* The core invariant this ticket restores: {@code get()} must never report a terminal for a
* lead whose pane no longer runs an agent. Before this fix the scanner joined labelled tabs to
* panes with no liveness check at all, so a tab left behind by a crashed/relaunched lead (see
* {@code LeadLauncher}'s own staleness handling) was reported as live forever — which is exactly
* what let duplicate lead tabs make {@code LeadCoordLoop.resolveLocalLead()} permanently unable
* to pick one. Mutate this away (drop the {@code agent.list} cross-check in {@link
* LeadTabScanner#scan()}) and this test must fail.
*/
@Test
void aLabelledTabWithNoRunningAgentIsNotReported() {
TopologyHerdr herdr = twoLeads().deadAgent("w1:t1"); // opus-5.0's tab is labelled but dead
Map<String, String> leads = scanner(herdr, twoLeadsConfigured(), new AtomicLong()).get();
assertFalse(leads.containsKey("term_opus"),
"a labelled tab with no running agent must never be reported as a live lead");
assertEquals("gpt-sol-5.6", leads.get("term_gpt"),
"the other, genuinely live lead must be unaffected");
}
/**
* The other direction, pinned separately so a fix cannot satisfy the test above by simply
* returning nothing: a labelled tab that DOES have a running agent must still be reported. A
* scanner that always comes back empty is worse than the bug it fixes.
*/
@Test
void aLabelledTabWithARunningAgentIsStillReported() {
Map<String, String> leads = scanner(twoLeads(), twoLeadsConfigured(), new AtomicLong()).get();
assertEquals("opus-5.0", leads.get("term_opus"));
assertEquals("gpt-sol-5.6", leads.get("term_gpt"));
}
/**
* fleetd #359 review, finding 2 — the exact scenario the ticket's own evidence showed:
* {@code agent.list} can come back successfully but short, without the herdr call ever throwing.
* A lead this class already reported as live must not be dropped on the strength of one such
* read: {@link LeadTabScanner#get()}'s "keep the cache on failure" contract only fires on an
* exception, so without a fix a single short {@code agent.list} silently empties the cached map —
* which would resolve that lead's pane as {@code Role.WORKER} downstream, refusing every
* orchestration call. Mutate this away (drop the one-scan grace in {@link
* LeadTabScanner#scan()}) and this test must fail.
*/
@Test
void aTransientAgentListMissDoesNotDemoteALeadAlreadyKnownLive() {
TopologyHerdr herdr = twoLeads();
AtomicLong clock = new AtomicLong();
LeadTabScanner s = scanner(herdr, twoLeadsConfigured(), clock);
assertTrue(s.get().containsKey("term_opus"), "opus must be known live before the miss");
// One scan where agent.list comes back without opus's tab, even though the tab and pane are
// completely unchanged — the tab/pane are still there, only the liveness read is short.
herdr.deadAgent("w1:t1");
clock.addAndGet(TTL);
assertTrue(s.get().containsKey("term_opus"),
"a single missed detection must not empty the cached map for a lead already known live");
assertEquals("gpt-sol-5.6", s.get().get("term_gpt"), "the unaffected lead is unchanged");
}
/**
* The other half of finding 2, so the grace above cannot be mistaken for permanent amnesty: the
* original #359 invariant (a genuinely dead tab is not reported forever) must still hold once a
* SECOND, independent scan agrees the agent is gone.
*/
@Test
void aLeadMissingFromAgentListOnTwoConsecutiveScansIsFinallyDropped() {
TopologyHerdr herdr = twoLeads();
AtomicLong clock = new AtomicLong();
LeadTabScanner s = scanner(herdr, twoLeadsConfigured(), clock);
assertTrue(s.get().containsKey("term_opus"));
herdr.deadAgent("w1:t1");
clock.addAndGet(TTL);
assertTrue(s.get().containsKey("term_opus"), "first miss is a grace period, not a verdict");
clock.addAndGet(TTL); // a second, independent scan — still no agent
assertFalse(s.get().containsKey("term_opus"),
"a second consecutive miss for the same terminal must finally drop it");
}
/** A lead that recovers between the two misses keeps its grace spent, not renewed for free. */
@Test
void aLeadThatRecoversBetweenMissesIsReportedNormallyAndResetsItsGrace() {
TopologyHerdr herdr = twoLeads();
AtomicLong clock = new AtomicLong();
LeadTabScanner s = scanner(herdr, twoLeadsConfigured(), clock);
assertTrue(s.get().containsKey("term_opus"));
herdr.deadAgent("w1:t1");
clock.addAndGet(TTL);
assertTrue(s.get().containsKey("term_opus"), "graced on the first miss");
herdr.reviveAgent("w1:t1"); // the miss really was transient
clock.addAndGet(TTL);
assertTrue(s.get().containsKey("term_opus"), "found live again — reported normally");
// A later, unrelated miss must get its own fresh grace scan rather than being dropped
// immediately because the earlier miss had already "used up" a slot for this terminal.
herdr.deadAgent("w1:t1");
clock.addAndGet(TTL);
assertTrue(s.get().containsKey("term_opus"),
"a miss after a genuine recovery is a new event and deserves its own grace scan");
}
/**
* CB-579 acceptance (6): this is the bug the ticket closes. A stale pin used to be merged back
* over every scan and never expire; now a scan is the whole answer, so a lead whose tab is gone
@@ -9,7 +9,6 @@ import org.junit.jupiter.api.Test;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
import static org.junit.jupiter.api.Assertions.*;
@@ -107,7 +106,6 @@ class LeadLauncherTest {
assertEquals(0, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads());
assertFalse(herdr.called("agent.start"), "the live lead must not be duplicated");
assertFalse(herdr.called("tab.close"), "a labelled tab WITH a live agent must never be closed");
}
/**
@@ -124,121 +122,6 @@ class LeadLauncherTest {
"a stale label is not a lead; the lead must be relaunched");
}
/**
* fleetd #359 review finding 1 — the exact scenario the ticket's own live evidence produced: on a
* real host, {@code agent.list} reported "0 live" for a tab a plain {@code ps} confirmed was
* running a real session. A single such reading must never close the tab outright — that would
* destroy the operator's actual lead, a worse failure than the stale-tab bug this ticket exists
* to fix. The first dead reading only flags the tab; mutate this away (make the first reading
* close instead of flag) and this test must fail.
*/
@Test
void aStaleLabelledTabIsFlaggedRatherThanClosedOnTheFirstReconcile() {
FakeHerdr herdr = new FakeHerdr()
.withWorkspace("wL", "fleet")
.withTab("wL", "wL:t1", "lead: opus"); // label only, first look — could be a live session agent.list missed
assertEquals(1, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads());
assertFalse(herdr.called("tab.close"),
"one missed reading must never close a tab that might be hosting a live session");
assertTrue(flaggedTabIds(herdr).contains("wL:t1"),
"the stale tab must be flagged pending-close so a later reconcile can confirm it");
}
/**
* The operator's own trace on fleet01 (#359): a daemon that restarted several times, each
* occasion finding "0 live" for whatever reason, had left several identically-labelled dead
* tabs sitting side by side. Every one of them is flagged on its first dead reading, not closed —
* none is more or less trustworthy than another.
*/
@Test
void allStaleLabelledTabsAreFlaggedRatherThanClosedOnTheFirstReconcile() {
FakeHerdr herdr = new FakeHerdr()
.withWorkspace("wL", "fleet")
.withTab("wL", "wL:t1", "lead: opus")
.withTab("wL", "wL:t2", "lead: opus")
.withTab("wL", "wL:t3", "lead: opus"); // three restarts' worth of debris, none confirmed twice yet
assertEquals(1, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads());
assertFalse(herdr.called("tab.close"), "no tab may be closed on its first dead reading");
assertEquals(Set.of("wL:t1", "wL:t2", "wL:t3"), flaggedTabIds(herdr),
"every dead labelled tab must be flagged, not just the first one found");
}
/**
* fleetd #359 review finding 1, the other half — a tab already flagged pending-close by an
* earlier reconcile, and STILL dead on this one, has now been read dead on two independent,
* separately-connected reconciles. That is strong enough evidence to actually close it. Mutate
* this away (never close a flagged tab) and this test must fail — the original #359 growth bug
* would come back for good.
*/
@Test
void aTabAlreadyFlaggedPendingCloseIsClosedWhenStillDeadOnALaterReconcile() {
FakeHerdr herdr = new FakeHerdr()
.withWorkspace("wL", "fleet")
.withTab("wL", "wL:t1", "lead: opus [fleetd:pending-close]"); // flagged last reconcile, still dead
assertEquals(1, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads());
assertTrue(herdr.called("tab.close"),
"a tab dead on two independent reconciles must finally be closed");
assertEquals("wL:t1", ((Map<?, ?>) herdr.lastCall("tab.close").params()).get("tab_id"));
}
/** All of several already-flagged, still-dead tabs are closed — not just the first found. */
@Test
void allTabsAlreadyFlaggedPendingCloseAreClosedWhenStillDead() {
FakeHerdr herdr = new FakeHerdr()
.withWorkspace("wL", "fleet")
.withTab("wL", "wL:t1", "lead: opus [fleetd:pending-close]")
.withTab("wL", "wL:t2", "lead: opus [fleetd:pending-close]")
.withTab("wL", "wL:t3", "lead: opus [fleetd:pending-close]");
assertEquals(1, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads());
List<Object> closedTabIds = herdr.calls.stream()
.filter(c -> c.method().equals("tab.close"))
.<Object>map(c -> ((Map<?, ?>) c.params()).get("tab_id"))
.toList();
assertEquals(3, closedTabIds.size(),
"every confirmed-dead labelled tab must be closed, not just the first one found");
assertEquals(Set.of("wL:t1", "wL:t2", "wL:t3"), Set.copyOf(closedTabIds));
}
/**
* A tab flagged pending-close on a previous reconcile that is running an agent again — the miss
* that flagged it was transient. It must never be closed, and its flag must be cleared so a
* future, unrelated miss starts its own two-reading count from zero.
*/
@Test
void aFlaggedTabRunningAnAgentAgainHasItsFlagClearedInsteadOfBeingClosed() {
FakeHerdr herdr = new FakeHerdr()
.withWorkspace("wL", "fleet")
.withTab("wL", "wL:t1", "lead: opus [fleetd:pending-close]")
.withAgent("lead-opus", "term_lead", "wL:p1", "wL:t1"); // it recovered — really alive now
assertEquals(0, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads(),
"the lead is live again — nothing to relaunch");
assertFalse(herdr.called("tab.close"), "a tab running an agent again must never be closed");
assertFalse(herdr.called("agent.start"), "the lead is live again — nothing to relaunch");
assertEquals("wL:t1", ((Map<?, ?>) herdr.lastCall("tab.rename").params()).get("tab_id"));
assertEquals("lead: opus", ((Map<?, ?>) herdr.lastCall("tab.rename").params()).get("label"),
"the pending-close flag must be cleared once the tab is confirmed live again");
}
/** Every {@code tab.rename} call whose label carries the pending-close marker, by tab id. */
private static Set<String> flaggedTabIds(FakeHerdr herdr) {
return herdr.calls.stream()
.filter(c -> c.method().equals("tab.rename"))
.filter(c -> String.valueOf(((Map<?, ?>) c.params()).get("label"))
.endsWith("[fleetd:pending-close]"))
.map(c -> String.valueOf(((Map<?, ?>) c.params()).get("tab_id")))
.collect(java.util.stream.Collectors.toUnmodifiableSet());
}
/**
* A lead the operator opened by hand is live once its tab carries the configured `tab:` label —
* CB-579 retired the `terminal:` pin, so a hand-opened lead is found the same way an
@@ -1,7 +1,6 @@
package dev.ltms.fleet.mcp;
import dev.ltms.fleet.msg.FakeLeadChannel;
import dev.ltms.fleet.msg.LeadChannel;
import dev.ltms.fleet.msg.LeadMessage;
import io.modelcontextprotocol.spec.McpSchema;
import org.junit.jupiter.api.Test;
@@ -94,56 +93,4 @@ class FleetMcpLeadCoordTest {
assertTrue(FleetMcp.sendToLead(channel, PEER, " ", null, null).isError());
assertEquals(0, channel.published().size());
}
/**
* fleetd #361: "delivered" overstated what publish actually proves — the broker's confirm means
* durably queued, not read. A zero-consumer target is the observable form of "this will not
* reach a pane right now", so the (still-successful) result must say so.
*/
@Test
void warnsWhenThePeerMailboxHasNoConsumersButStillReportsSuccess() {
var channel = new FakeLeadChannel(SELF)
.withMailbox(PEER, LeadChannel.MailboxState.exists(PEER, 0, 0));
McpSchema.CallToolResult res = FleetMcp.sendToLead(channel, PEER, "hi", null, null);
assertFalse(res.isError(), "a zero-consumer mailbox is still a successful, durably-queued publish");
String out = textOf(res);
assertTrue(out.contains("durably confirmed"), out);
assertTrue(out.toLowerCase().contains("no consumers"), () -> "must warn nobody is reading it: " + out);
}
@Test
void staysQuietAboutConsumersWhenThePeerMailboxHasOne() {
var channel = new FakeLeadChannel(SELF)
.withMailbox(PEER, LeadChannel.MailboxState.exists(PEER, 0, 1));
McpSchema.CallToolResult res = FleetMcp.sendToLead(channel, PEER, "hi", null, null);
assertFalse(res.isError());
String out = textOf(res);
assertTrue(out.contains("durably confirmed"), out);
assertFalse(out.toLowerCase().contains("no consumers"), () -> "a consumer IS attached: " + out);
}
/**
* fleetd #361 review finding 1: an unmeasured fact must never render as a definite one. When
* the post-publish probe could not determine the mailbox's consumer count at all (broker slow,
* unreachable, or the probe timed out — {@link LeadChannel.MailboxState#unknown}), the result
* must stay just as quiet as the has-a-consumer case — never assert "no consumers" for a mailbox
* this call never actually measured.
*/
@Test
void staysQuietAboutConsumersWhenThePeerMailboxStateIsUnknown() {
var channel = new FakeLeadChannel(SELF)
.withMailbox(PEER, LeadChannel.MailboxState.unknown(PEER));
McpSchema.CallToolResult res = FleetMcp.sendToLead(channel, PEER, "hi", null, null);
assertFalse(res.isError(), "an unresolved post-publish probe must never turn a durably-confirmed publish into an error");
String out = textOf(res);
assertTrue(out.contains("durably confirmed"), out);
assertFalse(out.toLowerCase().contains("no consumers"),
() -> "an unmeasured fact must never be reported as a definite zero-consumer mailbox: " + out);
}
}
@@ -10,9 +10,6 @@ import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.PaneLocator;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.inject.Injector;
import dev.ltms.fleet.msg.FakeLeadChannel;
import dev.ltms.fleet.msg.LeadChannel;
import dev.ltms.fleet.msg.LeadMessage;
import dev.ltms.fleet.msg.MessageService;
import dev.ltms.fleet.msg.Rendezvous;
import dev.ltms.fleet.session.FakeWorktrees;
@@ -37,9 +34,7 @@ import java.util.Map;
import java.util.EnumSet;
import java.util.Set;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.function.Function;
import static org.junit.jupiter.api.Assertions.*;
@@ -581,48 +576,17 @@ class FleetMcpTest {
void listReportsThisDaemonsOwnCoordIdWhenLeadCoordinationIsOn() {
FakeHerdr h = new FakeHerdr();
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
FakeLeadChannel channel = new FakeLeadChannel("mac-opus")
.withMailbox("mac-opus", LeadChannel.MailboxState.exists("mac-opus", 0, 1));
McpSchema.CallToolResult res = FleetMcp.listFleet(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), Map.of(), "",
new FleetMcp.CoordinationSource(channel, List.of()));
FleetMcp.QuarantineSource.none(), Map.of(), "", "mac-opus");
String out = textOf(res);
// fleetd #361: reports both which coord-id a peer must use to reach ME, and this daemon's
// own mailbox state (a self-diagnosis: is my own consumer actually attached?).
// There is no peer-discovery surface yet, so this row answers the one question an operator
// cannot answer any other way: which coord-id a peer must use to reach ME.
assertTrue(out.contains("\"coordinator\""), out);
assertTrue(out.contains("\"selfId\":\"mac-opus\""), out);
assertTrue(out.contains("\"mailbox\":{\"status\":\"exists\",\"pending\":0,\"consumers\":1}"), out);
assertTrue(out.contains("\"held\":[]"), out);
assertTrue(out.contains("\"peers\":[]"), out);
}
/**
* fleetd #361 review finding 1: a self-probe that could not complete (broker unreachable, timed
* out) must never render the same as a measured "0 pending, 0 consumers" — that was exactly the
* bug: a reader could not tell "my mailbox is empty and idle" from "I could not check", and the
* second one is the far more alarming state.
*/
@Test
void listReportsAnUnresolvedSelfProbeAsUnknownNeverAsAMeasuredZero() {
FakeHerdr h = new FakeHerdr();
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
FakeLeadChannel channel = new FakeLeadChannel("mac-opus")
.withMailbox("mac-opus", LeadChannel.MailboxState.unknown("mac-opus"));
McpSchema.CallToolResult res = FleetMcp.listFleet(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), Map.of(), "",
new FleetMcp.CoordinationSource(channel, List.of()));
String out = textOf(res);
assertTrue(out.contains("\"mailbox\":{\"status\":\"unknown\"}"), out);
assertFalse(out.contains("\"pending\""), "an unresolved probe must never carry a pending count at all: " + out);
assertFalse(out.contains("\"consumers\""), "an unresolved probe must never carry a consumers count at all: " + out);
}
@Test
@@ -637,110 +601,6 @@ class FleetMcpTest {
"an ordinary fleet's output must be unchanged by this feature");
}
@Test
void listReportsHeldMessagesWithATruncatedPreviewNeverTheFullBody() {
FakeHerdr h = new FakeHerdr();
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
String longContent = "x".repeat(200);
FakeLeadChannel channel = new FakeLeadChannel("mac-opus")
.hold(new LeadMessage("m1", "fleet01-lead", "mac-opus", longContent));
McpSchema.CallToolResult res = FleetMcp.listFleet(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), Map.of(), "",
new FleetMcp.CoordinationSource(channel, List.of()));
String out = textOf(res);
assertTrue(out.contains("\"msgId\":\"m1\""), out);
assertTrue(out.contains("\"from\":\"fleet01-lead\""), out);
assertFalse(out.contains(longContent), "fleet_list must never dump a held message's full body: " + out);
assertTrue(out.contains("x".repeat(80) + "…"), "expected an 80-char preview with an ellipsis: " + out);
}
@Test
void listReportsEachDeclaredPeersLiveReachability() {
FakeHerdr h = new FakeHerdr();
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
FakeLeadChannel channel = new FakeLeadChannel("mac-opus")
.withMailbox("fleet01-lead", LeadChannel.MailboxState.exists("fleet01-lead", 2, 1))
.withMailbox("fleet03-lead", LeadChannel.MailboxState.unknown("fleet03-lead"));
// "fleet02-lead" is declared as a peer but never configured on the fake — inspect() falls
// back to MailboxState.absent, exactly as a real down (never-run) peer would report.
// "fleet03-lead" IS configured, as unknown — a broker that could not be reached in time,
// which review finding 1 says must render distinctly from "fleet02-lead"'s confirmed absence.
McpSchema.CallToolResult res = FleetMcp.listFleet(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), Map.of(), "",
new FleetMcp.CoordinationSource(channel, List.of("fleet01-lead", "fleet02-lead", "fleet03-lead")));
String out = textOf(res);
assertTrue(out.contains("\"coordId\":\"fleet01-lead\",\"status\":\"exists\",\"pending\":2,\"consumers\":1"), out);
assertTrue(out.contains("\"coordId\":\"fleet02-lead\",\"status\":\"absent\""), out);
assertTrue(out.contains("\"coordId\":\"fleet03-lead\",\"status\":\"unknown\""), out);
assertFalse(out.contains("\"coordId\":\"fleet02-lead\",\"status\":\"absent\",\"pending\""),
"pending/consumers must be omitted, not faked as zero, for a confirmed-absent peer: " + out);
assertFalse(out.contains("\"coordId\":\"fleet03-lead\",\"status\":\"unknown\",\"pending\""),
"pending/consumers must be omitted, not faked as zero, for an unresolved peer probe: " + out);
}
/**
* fleetd #361 review finding 2: {@code get(timeout)} alone times out the CALLER but leaves the
* submitted {@link LeadChannel#inspect} task running forever on its own virtual thread — against
* a hung (not down) broker every probe would orphan one more thread holding an AMQP channel
* until the connection's channel-max is exhausted, which would break {@code publish} too. This
* proves {@link FleetMcp#probe(LeadChannel, String, long)} does not merely give up on a slow
* task: it interrupts it, so the task does not go on running unbounded after the caller has
* already moved on. No hung broker needed — a {@link LeadChannel} fake that blocks until
* interrupted is enough to observe the same mechanism.
*/
@Test
void aTimedOutProbeInterruptsTheOrphanedTaskRatherThanAbandoningIt() throws Exception {
CountDownLatch started = new CountDownLatch(1);
AtomicBoolean wasInterrupted = new AtomicBoolean(false);
LeadChannel hangs = new LeadChannel() {
@Override
public void publish(String toCoordId, LeadMessage m) { }
@Override
public List<LeadMessage> peek() { return List.of(); }
@Override
public void ack(String msgId) { }
@Override
public String selfCoordId() { return "mac-opus"; }
@Override
public MailboxState inspect(String coordId) {
started.countDown();
try {
Thread.sleep(60_000);
} catch (InterruptedException e) {
wasInterrupted.set(true);
Thread.currentThread().interrupt();
}
return MailboxState.unknown(coordId);
}
};
LeadChannel.MailboxState result = FleetMcp.probe(hangs, "fleet01-lead", 100L);
assertFalse(result.exists(), "a timed-out probe must never claim the mailbox exists");
assertFalse(result.known(), "a timed-out probe proves nothing either way — it must report unknown");
assertTrue(started.await(2, TimeUnit.SECONDS), "the probe task must actually have started");
// The interrupt is delivered asynchronously to the orphaned task's own thread — poll briefly
// rather than assume it has already landed the instant probe() returns.
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(2);
while (!wasInterrupted.get() && System.nanoTime() < deadline) {
Thread.sleep(20);
}
assertTrue(wasInterrupted.get(),
"probe() must cancel the orphaned task (interrupt it) instead of leaving it to run forever");
}
@Test
void capacityUsesThePlacementLiveCount() {
FakeHerdr h = new FakeHerdr();
@@ -1051,8 +911,7 @@ class FleetMcpTest {
String out = textOf(FleetMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")),
sessions, null, new FleetMcp.CapacitySource(profile -> 2, profile -> 3,
() -> Set.of("sonnet"), () -> 0), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), leadSeats, Map.of(), "",
FleetMcp.CoordinationSource.none()));
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), leadSeats, Map.of(), "", null));
assertTrue(out.contains("\"maxLoad\":3"), "maxLoad itself must be left untouched: " + out);
assertTrue(out.contains("\"live\":2"), out);
@@ -1075,8 +934,7 @@ class FleetMcpTest {
String out = textOf(FleetMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")),
sessions, null, new FleetMcp.CapacitySource(profile -> 0, profile -> 3,
() -> Set.of("sonnet"), () -> 0), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), leadSeats, Map.of(), "",
FleetMcp.CoordinationSource.none()));
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), leadSeats, Map.of(), "", null));
assertTrue(out.contains("\"live\":0"), out);
assertTrue(out.contains("\"free\":3"), "the real gate never subtracts the lead's seat: " + out);
@@ -1130,7 +988,7 @@ class FleetMcpTest {
String out = textOf(FleetMcp.listFleet(composite, sm, null,
new FleetMcp.CapacitySource(liveCount, p -> profiles.get(p).maxLoad(), profiles::keySet, () -> 0),
new FleetMcp.HealthCoverageSource(() -> "off"), FleetMcp.QuarantineSource.none(),
FleetMcp.OutageSource.none(), leadSeats, Map.of(), "", FleetMcp.CoordinationSource.none()));
FleetMcp.OutageSource.none(), leadSeats, Map.of(), "", null));
int reportedFree = extractInt(out, "free");
for (int i = 0; i < reportedFree; i++) {
@@ -1159,7 +1017,7 @@ class FleetMcpTest {
sessions, null, new FleetMcp.CapacitySource(profile -> 0, profile -> 2,
() -> Set.of("terra"), () -> 0), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
Map.of(), "", FleetMcp.CoordinationSource.none()));
Map.of(), "", null));
assertTrue(out.contains("\"free\":2"), out);
assertFalse(out.contains("leadSeats"), "no lead shares this profile's credential: " + out);
@@ -146,7 +146,7 @@ class MemberEnvAllowListTest {
return new FleetConfig(null, null, null, Map.of(), null, null, null, null, null,
new FleetConfig.Broker(null, brokerUriEnv, null), null, null, null, null, null,
null, null, null, null,
new FleetConfig.Coordinator(null, coordinatorUriEnv, null, null, null)).withDefaults();
new FleetConfig.Coordinator(null, coordinatorUriEnv, null, null)).withDefaults();
}
/** {@code LC_*} categories are infrastructure by prefix; everything else needs an exact match. */
@@ -1,134 +0,0 @@
package dev.ltms.fleet.msg;
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 com.rabbitmq.client.Channel;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.impl.DefaultExceptionHandler;
import org.junit.jupiter.api.Test;
import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.lang.reflect.Proxy;
import java.util.concurrent.atomic.AtomicInteger;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertInstanceOf;
import static org.junit.jupiter.api.Assertions.assertNotEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
class AmqpConnectionFailureLoggerTest {
@Test
void installedHandlersLogTheirOwnConnectionNamesAtErrorWithTheCause() throws Exception {
ConnectionFactory inboxFactory = AmqpReplyInbox.connectionFactory("amqp://127.0.0.1");
ConnectionFactory mailboxFactory = LeadMailbox.connectionFactory("amqp://127.0.0.1");
AmqpConnectionFailureLogger inboxHandler = installedStrictHandler(inboxFactory, "reply inbox");
AmqpConnectionFailureLogger mailboxHandler = installedStrictHandler(mailboxFactory, "lead mailbox");
assertEquals(AmqpConnectionFailureLogger.REPLY_INBOX, inboxHandler.connectionName());
assertEquals(AmqpConnectionFailureLogger.LEAD_MAILBOX, mailboxHandler.connectionName());
ListAppender<ILoggingEvent> inboxEvents = attach(AmqpReplyInbox.class);
ListAppender<ILoggingEvent> mailboxEvents = attach(LeadMailbox.class);
IllegalStateException inboxFailure = new IllegalStateException("inbox failure");
IllegalStateException mailboxFailure = new IllegalStateException("mailbox failure");
try {
inboxHandler.handleUnexpectedConnectionDriverException(null, inboxFailure);
mailboxHandler.handleConnectionRecoveryException(null, mailboxFailure);
assertError(inboxEvents, "AMQP connection fleetd-reply-inbox: An unexpected connection driver error occurred",
inboxFailure, "inbox failure line");
assertError(mailboxEvents, "AMQP connection fleetd-lead-mailbox: Caught an exception during connection recovery!",
mailboxFailure, "mailbox recovery line");
} finally {
detach(AmqpReplyInbox.class, inboxEvents);
detach(LeadMailbox.class, mailboxEvents);
}
}
@Test
void connectionResetKeepsForgivingHandlerWarningSemantics() {
AmqpConnectionFailureLogger handler = new AmqpConnectionFailureLogger(
AmqpConnectionFailureLogger.REPLY_INBOX, LoggerFactory.getLogger(AmqpReplyInbox.class));
ListAppender<ILoggingEvent> events = attach(AmqpReplyInbox.class);
try {
handler.handleUnexpectedConnectionDriverException(null, new IOException("Connection reset"));
assertEquals(1, events.list.size(), "the handler must still log a reset");
ILoggingEvent event = events.list.getFirst();
assertEquals(Level.WARN, event.getLevel(), "ForgivingExceptionHandler logs connection resets at WARN");
assertEquals("AMQP connection fleetd-reply-inbox: An unexpected connection driver error occurred "
+ "(Exception message: Connection reset)", event.getFormattedMessage());
assertTrue(event.getThrowableProxy() == null, "ForgivingExceptionHandler does not attach a reset stack trace");
} finally {
detach(AmqpReplyInbox.class, events);
}
}
@Test
void connectionNamesStayDistinct() {
assertNotEquals(AmqpConnectionFailureLogger.REPLY_INBOX, AmqpConnectionFailureLogger.LEAD_MAILBOX,
"reply-inbox and lead-mailbox failures must be distinguishable");
}
@Test
void strictConsumerExceptionStillClosesItsChannel() {
AtomicInteger closes = new AtomicInteger();
Channel channel = (Channel) Proxy.newProxyInstance(getClass().getClassLoader(), new Class<?>[] {Channel.class},
(_, method, _) -> switch (method.getName()) {
case "close" -> {
closes.incrementAndGet();
yield null;
}
case "toString" -> "test-channel";
default -> throw new UnsupportedOperationException(method.getName());
});
AmqpConnectionFailureLogger handler = new AmqpConnectionFailureLogger(
AmqpConnectionFailureLogger.REPLY_INBOX, LoggerFactory.getLogger(AmqpReplyInbox.class));
handler.handleConsumerException(channel, new IllegalStateException("consumer failed"), null, "tag", "handleDelivery");
assertEquals(1, closes.get(), "DefaultExceptionHandler must close a channel after a consumer exception");
}
@Test
void handlerOnlyChangesDefaultHandlerLogging() {
assertEquals(DefaultExceptionHandler.class,
AmqpConnectionFailureLogger.class.getSuperclass());
assertFalse(java.util.Arrays.stream(AmqpConnectionFailureLogger.class.getDeclaredMethods())
.anyMatch(method -> method.getName().startsWith("handle")),
"all exception-handling methods must remain inherited from DefaultExceptionHandler");
}
private static ListAppender<ILoggingEvent> attach(Class<?> owner) {
Logger logger = (Logger) LoggerFactory.getLogger(owner);
logger.setLevel(Level.DEBUG);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
return appender;
}
private static void detach(Class<?> owner, ListAppender<ILoggingEvent> appender) {
((Logger) LoggerFactory.getLogger(owner)).detachAppender(appender);
}
private static AmqpConnectionFailureLogger installedStrictHandler(ConnectionFactory factory, String connection) {
assertInstanceOf(DefaultExceptionHandler.class, factory.getExceptionHandler(),
connection + " must keep DefaultExceptionHandler: replacing the strict handler with a forgiving one "
+ "changes when a channel is closed");
return assertInstanceOf(AmqpConnectionFailureLogger.class, factory.getExceptionHandler());
}
private static void assertError(ListAppender<ILoggingEvent> events, String message, Throwable cause, String name) {
assertEquals(1, events.list.size(), name);
ILoggingEvent event = events.list.getFirst();
assertEquals(Level.ERROR, event.getLevel(), name);
assertEquals(message, event.getFormattedMessage(), name);
assertEquals(cause.toString(), event.getThrowableProxy().getClassName() + ": "
+ event.getThrowableProxy().getMessage(), name);
}
}
@@ -3,8 +3,6 @@ package dev.ltms.fleet.msg;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
/**
* Hermetic stand-in for {@link LeadChannel}: an in-memory mailbox that records what was published
@@ -26,19 +24,11 @@ public final class FakeLeadChannel implements LeadChannel {
private final List<String> acked = Collections.synchronizedList(new ArrayList<>());
/** When set, every {@link #publish} throws it — the unroutable/nacked/timed-out peer. */
private volatile IllegalStateException publishFailure;
/** Canned {@link #inspect} results by coord-id — absent for any coord-id not configured here. */
private final Map<String, MailboxState> mailboxes = new ConcurrentHashMap<>();
public FakeLeadChannel(String selfCoordId) {
this.selfCoordId = selfCoordId;
}
/** Make {@link #inspect(String)} return {@code state} for {@code coordId} instead of "absent". */
public FakeLeadChannel withMailbox(String coordId, MailboxState state) {
mailboxes.put(coordId, state);
return this;
}
/** Make every publish fail as an unreachable peer would. */
public FakeLeadChannel failPublishWith(String message) {
this.publishFailure = new IllegalStateException(message);
@@ -75,11 +65,6 @@ public final class FakeLeadChannel implements LeadChannel {
return selfCoordId;
}
@Override
public MailboxState inspect(String coordId) {
return mailboxes.getOrDefault(coordId, MailboxState.absent(coordId));
}
public List<LeadMessage> published() {
return List.copyOf(published);
}
@@ -1,77 +0,0 @@
package dev.ltms.fleet.msg;
import com.rabbitmq.client.ShutdownSignalException;
import com.rabbitmq.client.impl.AMQImpl;
import org.junit.jupiter.api.Test;
import java.io.IOException;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #361 review round 2: a mutation that made {@link LeadMailbox#isMissingQueue} return
* {@code true} unconditionally still left {@code mvn clean install} green — 1389 tests, 0
* failures — because nothing exercised its false branch. That branch is the whole discriminator
* between {@link LeadChannel.MailboxState#absent} and {@link LeadChannel.MailboxState#unknown};
* without a test pinning it, a future refactor that widens it back to "always true" (restoring the
* exact overstatement fleetd #361 exists to fix) would pass this suite.
*
* <p>Hermetic — no broker needed, per the review's own suggestion. {@code isMissingQueue} takes a
* plain {@link IOException}, so every input here is constructed directly rather than provoked from
* a live connection. The real 404 shape itself is still pinned against a real broker, in
* {@code LeadMailboxTest.passiveDeclareOfAMissingQueueThrowsAnIOExceptionWrappingA404ShutdownSignal}
* — this class covers the three false shapes {@link LeadMailbox#isMissingQueue}'s own javadoc
* lists, so both directions of the discriminator are proven somewhere.
*/
class LeadMailboxIsMissingQueueTest {
@Test
void aConfirmedMissingQueueIsRecognized() {
ShutdownSignalException sse = new ShutdownSignalException(true, false,
new AMQImpl.Channel.Close(404, "NOT_FOUND - no queue 'lead.x.inbox' in vhost '/'", 50, 10), null);
IOException e = new IOException("channel error", sse);
assertTrue(LeadMailbox.isMissingQueue(e), "a genuine 404 Channel.Close must be recognized as a missing queue");
}
@Test
void aDifferentReplyCodeIsNotAMissingQueue() {
// E.g. 403 ACCESS_REFUSED — the queue may well exist; this call was simply refused.
ShutdownSignalException sse = new ShutdownSignalException(true, false,
new AMQImpl.Channel.Close(403, "ACCESS_REFUSED", 50, 10), null);
IOException e = new IOException("channel error", sse);
assertFalse(LeadMailbox.isMissingQueue(e),
"a non-404 reply code must never be read as a confirmed absence — the mailbox's real state is unknown");
}
@Test
void aShutdownSignalWhoseReasonIsNotAChannelCloseIsNotAMissingQueue() {
// A Connection.Close (a whole different broker-level shutdown) is still a ShutdownSignalException,
// but its reason is not a Channel.Close at all — must not be misread as "no such queue".
ShutdownSignalException sse = new ShutdownSignalException(true, false,
new AMQImpl.Connection.Close(404, "coincidentally 404, but this is a CONNECTION close", 10, 50), null);
IOException e = new IOException("connection error", sse);
assertFalse(LeadMailbox.isMissingQueue(e),
"a ShutdownSignalException whose reason is not a Channel.Close must never be read as a missing queue,"
+ " even if its reply code happens to be 404");
}
@Test
void anIOExceptionWithNoCauseAtAllIsNotAMissingQueue() {
IOException e = new IOException("some other declare failure, no cause attached");
assertFalse(LeadMailbox.isMissingQueue(e),
"an IOException with no ShutdownSignalException cause must never be read as a confirmed absence");
}
@Test
void anIOExceptionWithAnUnrelatedCauseIsNotAMissingQueue() {
IOException e = new IOException("wrapped something else entirely", new RuntimeException("boom"));
assertFalse(LeadMailbox.isMissingQueue(e),
"a cause that isn't even a ShutdownSignalException must never be read as a confirmed absence");
}
}
@@ -1,9 +1,5 @@
package dev.ltms.fleet.msg;
import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ShutdownSignalException;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Tag;
import org.junit.jupiter.api.Test;
@@ -11,14 +7,11 @@ import org.testcontainers.containers.RabbitMQContainer;
import org.testcontainers.junit.jupiter.Testcontainers;
import org.testcontainers.utility.DockerImageName;
import java.io.IOException;
import java.util.List;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertInstanceOf;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
@@ -166,145 +159,6 @@ class LeadMailboxTest {
}
}
@Test
void inspectReportsAnOwnedMailboxAsExistingWithItsOwnConsumer() throws Exception {
// A LeadMailbox declares AND consumes its own queue the moment open() returns (see own()),
// so inspecting a coord-id this same process owns must always find exactly one consumer.
String self = coordId("lead-inspect-self");
try (LeadMailbox mailbox = LeadMailbox.open(uri(), self)) {
LeadChannel.MailboxState state = mailbox.inspect(self);
assertEquals(self, state.coordId());
assertTrue(state.exists(), "this daemon owns and has declared this exact queue");
assertEquals(0, state.pending(), "nothing has been published to it yet");
assertEquals(1, state.consumers(), "the mailbox's own constructor already attached a consumer");
}
}
@Test
void inspectReportsAMissingMailboxAsAbsentRatherThanThrowing() throws Exception {
String nobody = coordId("lead-inspect-nobody");
try (LeadMailbox mailbox = LeadMailbox.open(uri(), coordId("lead-inspect-caller"))) {
LeadChannel.MailboxState state = mailbox.inspect(nobody);
assertEquals(LeadChannel.MailboxState.absent(nobody), state,
"a queue nobody has ever declared must report absent, never throw");
assertTrue(state.known(), "a confirmed 404 IS a definite answer — this is not the unknown case");
}
}
/**
* fleetd #361 review finding 3: {@code inspect} is specified to never throw, but the original
* implementation caught only {@link IOException} — and {@link Connection#createChannel()} on an
* already-closed connection throws {@link com.rabbitmq.client.AlreadyClosedException}, an
* unchecked {@link RuntimeException} (pinned by {@code
* createChannelOnAnAlreadyClosedConnectionThrowsAnUncheckedException} above). This drives that
* exact scenario through the real {@link LeadMailbox#inspect} — not the raw client call — and
* checks both halves of finding 1 and finding 3 at once: no exception escapes, and the result is
* {@code UNKNOWN} rather than the wrong-but-plausible-looking {@code ABSENT}.
*/
@Test
void inspectReportsUnknownRatherThanThrowingWhenTheConnectionIsAlreadyClosed() throws Exception {
LeadMailbox mailbox = LeadMailbox.open(uri(), coordId("lead-inspect-dead-connection"));
mailbox.close(); // tears down the connection `inspect` will try to open a probe channel on
LeadChannel.MailboxState state = mailbox.inspect(coordId("lead-inspect-irrelevant-target"));
assertFalse(state.exists());
assertFalse(state.known(), "a dead connection proves nothing about the target mailbox — it must be unknown, not absent");
assertEquals(LeadChannel.MailboxState.Presence.UNKNOWN, state.presence());
}
@Test
void inspectReportsPendingMessagesAndZeroConsumersWhenNobodyIsReadingAnymore() throws Exception {
// Publish into a mailbox this test owns, then never consume from it, to prove `pending` and
// `consumers` really come off the broker rather than off this process's own in-memory state.
String to = coordId("lead-inspect-pending");
String observerId = coordId("lead-inspect-observer");
try (LeadMailbox owner = LeadMailbox.open(uri(), to);
LeadMailbox observer = LeadMailbox.open(uri(), observerId)) {
owner.publish(to, new LeadMessage("m1", "lead-from", to, "sitting in the queue"));
awaitPeek(owner); // make sure the broker has actually enqueued it before inspecting
} // `owner` closes here: its consumer disconnects, but the durable, unacked message stays queued.
try (LeadMailbox observer = LeadMailbox.open(uri(), coordId("lead-inspect-observer-2"))) {
// The broker requeues `owner`'s unacked delivery asynchronously once its connection drops,
// so poll rather than assume the very first passive declare already sees the settled state.
LeadChannel.MailboxState state = awaitInspect(observer, to, s -> s.consumers() == 0);
assertTrue(state.exists());
assertEquals(0, state.consumers(), "the only owner just closed — nobody is reading this anymore");
assertEquals(1, state.pending(), "the unacked message must be requeued, never dropped");
}
}
/**
* fleetd #361's central invariant, proved rather than assumed: a passive queue declare of a
* missing queue closes ITS channel with a 404 in AMQP 0-9-1. {@link LeadMailbox#inspect} is
* specified to run on its own disposable channel for exactly this reason — this test is the one
* that actually exercises the failure mode and shows {@link LeadMailbox#publish} on the SAME
* instance is unaffected by it.
*/
@Test
void inspectingAMissingMailboxNeverBreaksPublishOnTheSameInstance() throws Exception {
String self = coordId("lead-invariant-self");
try (LeadMailbox mailbox = LeadMailbox.open(uri(), self)) {
// Miss on a queue that has never existed — this is exactly the 404-closes-the-channel case.
LeadChannel.MailboxState missed = mailbox.inspect(coordId("lead-invariant-nobody-home"));
assertFalse(missed.exists());
assertTrue(missed.known(), "a genuine 404 on a queue that never existed is a confirmed fact, not an unknown");
// publish() must still work on THIS SAME instance: if inspect() had reused `publishChannel`
// (or `channel`), the broker's 404 would have closed it out from underneath publish().
LeadMessage sent = new LeadMessage("after-miss", "lead-from", self, "still alive");
mailbox.publish(self, sent);
List<LeadMessage> got = awaitPeek(mailbox);
assertEquals(1, got.size(), "publish must still reach this mailbox's own queue after a missed inspect");
assertEquals("after-miss", got.getFirst().msgId());
// And a second inspect() — of a mailbox that DOES exist this time — must also still work,
// proving the miss did not wedge inspect() itself either.
LeadChannel.MailboxState self2 = mailbox.inspect(self);
assertTrue(self2.exists());
}
}
/**
* Pins the exact exception shape {@link LeadMailbox#inspect} relies on to tell a genuine 404
* (mailbox confirmed absent) apart from everything else (mailbox state unknown) — measured
* against a real broker rather than assumed from the AMQP 0-9-1 spec text. If this ever fails,
* the classification in {@code inspect} is reading the wrong shape and must be revisited.
*/
@Test
void passiveDeclareOfAMissingQueueThrowsAnIOExceptionWrappingA404ShutdownSignal() throws Exception {
try (Connection conn = LeadMailbox.connectionFactory(uri()).newConnection()) {
Channel probe = conn.createChannel();
String missing = LeadMailbox.queueName(coordId("lead-404-shape"));
IOException thrown = assertThrows(IOException.class, () -> probe.queueDeclarePassive(missing));
assertInstanceOf(ShutdownSignalException.class, thrown.getCause(),
() -> "expected the IOException to wrap a ShutdownSignalException, got: " + thrown);
ShutdownSignalException sse = (ShutdownSignalException) thrown.getCause();
assertInstanceOf(AMQP.Channel.Close.class, sse.getReason(),
() -> "expected a Channel.Close reason: " + sse);
AMQP.Channel.Close close = (AMQP.Channel.Close) sse.getReason();
assertEquals(404, close.getReplyCode(), () -> "expected AMQP NOT_FOUND (404): " + close);
assertFalse(probe.isOpen(), "the 404 must have closed the channel the declare ran on");
}
}
/**
* The other half of the same measurement: calling {@code createChannel()} on an
* already-closed connection — the shape {@link LeadMailbox#inspect} hits when the broker
* connection itself is gone — throws {@link com.rabbitmq.client.AlreadyClosedException}, an
* unchecked {@link RuntimeException}, not an {@link IOException}. An {@code inspect} that only
* caught {@code IOException} here would let this escape instead of reporting "unknown".
*/
@Test
void createChannelOnAnAlreadyClosedConnectionThrowsAnUncheckedException() throws Exception {
Connection conn = LeadMailbox.connectionFactory(uri()).newConnection();
conn.close();
RuntimeException thrown = assertThrows(RuntimeException.class, conn::createChannel);
assertInstanceOf(com.rabbitmq.client.AlreadyClosedException.class, thrown,
() -> "expected AlreadyClosedException, got: " + thrown);
}
/** Poll peek until at least one message is held, or ~10s elapse (broker delivery is async). */
@SuppressWarnings("BusyWait")
private static List<LeadMessage> awaitPeek(LeadMailbox inbox) throws InterruptedException {
@@ -316,18 +170,4 @@ class LeadMailboxTest {
}
return msgs;
}
/** Poll inspect(coordId) until it satisfies {@code done}, or ~10s elapse (broker state settles async). */
@SuppressWarnings("BusyWait")
private static LeadChannel.MailboxState awaitInspect(
LeadMailbox observer, String coordId, java.util.function.Predicate<LeadChannel.MailboxState> done)
throws InterruptedException {
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(10);
LeadChannel.MailboxState state = observer.inspect(coordId);
while (!done.test(state) && System.nanoTime() < deadline) {
Thread.sleep(50);
state = observer.inspect(coordId);
}
return state;
}
}
@@ -0,0 +1,54 @@
package dev.ltms.fleet.power;
import org.junit.jupiter.api.Test;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* Platform-detection unit tests for {@link CaffeinateSleepAssertionMechanism}.
*
* <p>This deliberately never calls {@link CaffeinateSleepAssertionMechanism#acquire()} itself —
* doing so on a real macOS machine would actually start a live {@code caffeinate} child and hold
* a real idle-sleep assertion, which the ticket this class exists for explicitly forbids testing
* with. Instead this exercises the pure {@code isSupportedPlatform(String)} predicate that
* {@code acquire()} consults before ever touching {@link ProcessBuilder} — so it proves the
* platform check itself is correct on any CI OS, but it does <strong>not</strong> prove that a
* real {@code caffeinate -i} spawn succeeds or that its child is torn down correctly; that half is
* exercised indirectly by {@link IdleSleepGuardTest} against a {@link FakeSleepAssertionMechanism}
* instead, which is the seam invariant 2/3 in the ticket call for.
*/
class CaffeinateSleepAssertionMechanismTest {
@Test
void macOsNamesAreSupported() {
assertTrue(CaffeinateSleepAssertionMechanism.isSupportedPlatform("Mac OS X"));
assertTrue(CaffeinateSleepAssertionMechanism.isSupportedPlatform("macOS"));
assertTrue(CaffeinateSleepAssertionMechanism.isSupportedPlatform("MAC OS X"));
}
@Test
void nonMacNamesAreNotSupported() {
assertFalse(CaffeinateSleepAssertionMechanism.isSupportedPlatform("Linux"));
assertFalse(CaffeinateSleepAssertionMechanism.isSupportedPlatform("Windows 11"));
}
@Test
void nullOsNameIsNotSupported() {
assertFalse(CaffeinateSleepAssertionMechanism.isSupportedPlatform(null));
}
/**
* The overload {@code isSupportedPlatform()} (no args) reads the JVM's real {@code os.name} —
* proves the wiring is live, without asserting a specific answer (this suite itself must pass
* on both macOS and Linux CI).
*/
@Test
void noArgOverloadReadsRealSystemProperty() {
boolean expected = CaffeinateSleepAssertionMechanism
.isSupportedPlatform(System.getProperty("os.name"));
boolean actual = CaffeinateSleepAssertionMechanism.isSupportedPlatform();
assertEquals(expected, actual);
}
}
@@ -0,0 +1,56 @@
package dev.ltms.fleet.power;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.atomic.AtomicInteger;
/**
* Recording fake {@link SleepAssertionMechanism} — the seam behind the real OS effect (a live
* {@code caffeinate} child process). No test in this package ever spawns that real process; every
* assertion here is against this fake's own call log instead.
*
* <p>Each acquired {@link FakeAssertion} records its own {@code close()} calls, and every
* acquired instance is kept in {@link #acquired} so a test can inspect all of them, including
* ones {@link IdleSleepGuard} has already released.
*/
final class FakeSleepAssertionMechanism implements SleepAssertionMechanism {
/** Every {@link FakeAssertion} this mechanism has ever handed out, in order. */
final CopyOnWriteArrayList<FakeAssertion> acquired = new CopyOnWriteArrayList<>();
private final AtomicInteger acquireCalls = new AtomicInteger();
private volatile boolean unavailable = false;
/** Make the next (and every subsequent) {@link #acquire()} return {@code null}, like a missing tool. */
void makeUnavailable() {
unavailable = true;
}
int acquireCallCount() {
return acquireCalls.get();
}
@Override
public SleepAssertion acquire() {
acquireCalls.incrementAndGet();
if (unavailable) {
return null;
}
FakeAssertion a = new FakeAssertion();
acquired.add(a);
return a;
}
/** A held fake assertion; records how many times {@code close()} was actually called. */
static final class FakeAssertion implements SleepAssertion {
private final AtomicInteger closeCalls = new AtomicInteger();
int closeCallCount() {
return closeCalls.get();
}
@Override
public void close() {
closeCalls.incrementAndGet();
}
}
}
@@ -0,0 +1,118 @@
package dev.ltms.fleet.power;
import org.junit.jupiter.api.Test;
import java.util.concurrent.atomic.AtomicInteger;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* {@link IdleSleepGuard} against a {@link FakeSleepAssertionMechanism} — the seam that stands in
* for a real {@code caffeinate} child process. No test in this class ever spawns a real OS
* process or asserts against real idle sleep; every assertion is against the fake's call log
* (how many times {@code acquire()}/{@code close()} were actually called). That proves the
* <em>orchestration</em> — when the guard decides to hold or release an assertion, and that it
* never throws — but it does <strong>not</strong> prove that {@code caffeinate -i} itself
* actually stops macOS from idle-sleeping; that half is outside what a unit test can safely
* exercise (see {@link CaffeinateSleepAssertionMechanismTest}'s class doc).
*/
class IdleSleepGuardTest {
@Test
void acquiresOnZeroToOneAndReleasesOnOneToZero() {
FakeSleepAssertionMechanism mechanism = new FakeSleepAssertionMechanism();
AtomicInteger liveCount = new AtomicInteger(0);
IdleSleepGuard guard = new IdleSleepGuard(mechanism, liveCount::get);
assertFalse(guard.isHeld(), "nothing held before any member is live");
liveCount.set(1);
guard.recheck();
assertTrue(guard.isHeld(), "an assertion must be held once a member is live");
assertEquals(1, mechanism.acquired.size());
assertEquals(0, mechanism.acquired.get(0).closeCallCount());
liveCount.set(0);
guard.recheck();
assertFalse(guard.isHeld(), "the assertion must be released once the last member goes");
assertEquals(1, mechanism.acquired.get(0).closeCallCount(), "the SAME held assertion must be closed");
}
@Test
void steadyLiveCountDoesNotReacquireOrRerelease() {
FakeSleepAssertionMechanism mechanism = new FakeSleepAssertionMechanism();
AtomicInteger liveCount = new AtomicInteger(2);
IdleSleepGuard guard = new IdleSleepGuard(mechanism, liveCount::get);
guard.recheck(); // 0 -> 2 crossing: acquires
guard.recheck(); // still 2: must be a no-op
guard.recheck(); // still 2: must be a no-op
assertEquals(1, mechanism.acquireCallCount(), "only the crossing touches the mechanism");
liveCount.set(1); // 2 -> 1: still > 0, still a no-op
guard.recheck();
assertTrue(guard.isHeld());
assertEquals(0, mechanism.acquired.get(0).closeCallCount());
assertEquals(1, mechanism.acquireCallCount());
}
/**
* Invariant 2: a missing/unavailable mechanism must never throw, and the guard must simply
* hold nothing. {@link FakeSleepAssertionMechanism#makeUnavailable()} makes {@code acquire()}
* return {@code null}, exactly like {@link CaffeinateSleepAssertionMechanism} does off macOS
* or when the {@code caffeinate} binary is missing.
*/
@Test
void unavailableMechanismNeverThrowsAndHoldsNothing() {
FakeSleepAssertionMechanism mechanism = new FakeSleepAssertionMechanism();
mechanism.makeUnavailable();
AtomicInteger liveCount = new AtomicInteger(1);
IdleSleepGuard guard = new IdleSleepGuard(mechanism, liveCount::get);
guard.recheck(); // must not throw
assertFalse(guard.isHeld(), "acquire() returned null, so nothing is held");
assertEquals(1, mechanism.acquireCallCount());
// still must not throw or leak on release, even though nothing was ever actually held
liveCount.set(0);
guard.recheck();
assertFalse(guard.isHeld());
guard.close(); // teardown with nothing held must also be a safe no-op
}
/**
* Invariant 3 (teardown). This is the test the mutation testing step removes the production
* release call to fail: with {@code releaseHeldLocked()} not invoked from {@link
* IdleSleepGuard#close()}, the held fake assertion's {@code close()} would never be called and
* this assertion would fail.
*/
@Test
void closeReleasesAHeldAssertionEvenWithoutAZeroCrossing() {
FakeSleepAssertionMechanism mechanism = new FakeSleepAssertionMechanism();
AtomicInteger liveCount = new AtomicInteger(1);
IdleSleepGuard guard = new IdleSleepGuard(mechanism, liveCount::get);
guard.recheck();
assertTrue(guard.isHeld());
guard.close();
assertFalse(guard.isHeld(), "close() must release whatever is held, independent of live count");
assertEquals(1, mechanism.acquired.get(0).closeCallCount());
}
@Test
void closeIsIdempotent() {
FakeSleepAssertionMechanism mechanism = new FakeSleepAssertionMechanism();
AtomicInteger liveCount = new AtomicInteger(1);
IdleSleepGuard guard = new IdleSleepGuard(mechanism, liveCount::get);
guard.recheck();
guard.close();
guard.close(); // must not throw, must not double-release
assertEquals(1, mechanism.acquired.get(0).closeCallCount());
}
}
@@ -0,0 +1,72 @@
package dev.ltms.fleet.power;
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.session.MemberSession;
import dev.ltms.fleet.session.SessionManager;
import org.junit.jupiter.api.Test;
import java.util.List;
import java.util.Map;
import java.util.Set;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* Proves the wiring {@code Fleetd.main} actually performs — {@code
* sessions.onAcquire(_ -> guard.recheck())} / {@code sessions.onRelease(_ -> guard.recheck())} —
* not just {@link IdleSleepGuard}'s own orchestration logic in isolation
* ({@link IdleSleepGuardTest} already covers that in isolation, which on its own would not catch
* a wiring gap — e.g. an {@code onAcquire} call typo'd to a no-op lambda, or the listener wired to
* the wrong SessionManager instance — see fleetd's own "a test on the seam does not prove the
* caller" lesson). This test builds a real {@link SessionManager} exactly as
* {@code SessionManagerTest} does (a {@link FakeHerdr}-backed {@link ClaudeCodeLauncher}, no live
* herdr process), wires it to an {@link IdleSleepGuard} the same two lines {@code Fleetd.main}
* uses, and drives real {@link SessionManager#acquire} / {@link SessionManager#release} calls.
*/
class IdleSleepGuardWiringTest {
private SessionManager sessionManager(FakeHerdr herdr) {
FleetConfig.Profile cfg = new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
List.of("ccs", "ltms-local"), "tab", "fleetd-workers",
"worker: {profile} #{n}", null, null, null);
ClaudeCodeLauncher workers = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
return new SessionManager(workers);
}
@Test
void acquiringAndReleasingRealSessionsDrivesTheGuardThroughTheSameWiringFleetdUses() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
FakeSleepAssertionMechanism mechanism = new FakeSleepAssertionMechanism();
IdleSleepGuard guard = new IdleSleepGuard(mechanism, sessions::size);
// The exact two lines Fleetd.main wires up.
sessions.onAcquire(_ -> guard.recheck());
sessions.onRelease(_ -> guard.recheck());
assertFalse(guard.isHeld(), "no member yet: nothing held");
MemberSession a = sessions.acquire("ltms-local", "/a", "/caller", "ownerA");
assertTrue(guard.isHeld(), "0 -> 1: the first live member must arm the guard");
MemberSession b = sessions.acquire("ltms-local", "/b", "/caller", "ownerB");
assertEquals(1, mechanism.acquireCallCount(), "2nd member: still just 1 live-to-2 step, no new acquire");
sessions.release(a.paneId());
assertTrue(guard.isHeld(), "one member still live: the guard must stay armed");
assertEquals(0, mechanism.acquired.get(0).closeCallCount());
sessions.release(b.paneId());
assertFalse(guard.isHeld(), "1 -> 0: the last member releasing must disarm the guard");
assertEquals(1, mechanism.acquired.get(0).closeCallCount());
}
}
@@ -12,7 +12,6 @@ import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
@@ -54,56 +53,6 @@ class GitWorktreesTest {
/** A non-empty autoenv file — the form that would prompt for authorization in a worktree. */
private static final String AUTOENV_WITH_DIRECTIVE = "export HELLO=world\n";
/**
* fleetd #369. A throwaway directory that lives for the whole class (JUnit 5.4+ supports a
* static {@code @TempDir} field, created once and removed once every test in this class has
* run) — backing every raw {@code git} subprocess's {@code XDG_CONFIG_HOME} below. It only
* ever needs to exist and be guaranteed free of a {@code git/ignore} file; nothing writes
* inside it.
*/
@TempDir
private static Path CLASS_TMP;
/**
* fleetd #369 — the leak measured: {@code XDG_CONFIG_HOME=<dir with a `*` git/ignore> mvn test
* -Dtest=GitWorktreesTest} failed 56 of 59 tests on an unpatched checkout, because {@link
* #gitOutput} set {@code GIT_CONFIG_GLOBAL}/{@code GIT_CONFIG_SYSTEM}/{@code
* GIT_TERMINAL_PROMPT} but not {@code XDG_CONFIG_HOME}, and {@link #status}/{@link
* #fullStatus} (plus every other raw {@code git} subprocess this class started) set NOTHING at
* all — inheriting the JVM's whole real environment, including the operator's real {@code
* ~/.gitconfig} and real default excludes file ({@code $XDG_CONFIG_HOME/git/ignore} or {@code
* $HOME/.config/git/ignore}, applied by git with no {@code core.excludesFile} configured at
* all — see {@code gitignore(5)}). {@code GIT_CONFIG_GLOBAL=/dev/null} does not stop that
* default from applying; only setting {@code XDG_CONFIG_HOME} to a directory that provably
* carries no {@code git/ignore} does.
*
* <p>This is the same isolation {@link #hermeticGitEnv} already gives {@link
* #seedingGitWorktrees}'s production {@link GitWorktrees} instances (fleetd #362 review fix,
* finding 2), reused here for every subprocess the TEST ITSELF starts to drive and inspect
* those fixture repos.
*/
private static Map<String, String> hermeticEnv() {
return hermeticGitEnv(CLASS_TMP);
}
/**
* The one seam every git subprocess in this class is built through — see criterion 4's
* self-check, {@link #everyGitSubprocessGoesThroughTheHermeticFactory}, which fails the moment
* a future helper builds its own {@code git} subprocess directly instead of calling this, so
* the omission that caused fleetd #369 gets caught by name rather than rediscovered by a
* poisoned machine. The one deliberate exception is {@link
* #worktreeCredentialHelperCompletesWithoutUsingAnInheritedHelper}, which needs a
* non-hermetic, test-controlled global config to prove the credential helper ignores it — see
* the comment on that test.
*/
private static ProcessBuilder gitProcessBuilder(Path cwd, String... args) {
List<String> cmd = new java.util.ArrayList<>(List.of("git"));
cmd.addAll(List.of(args));
ProcessBuilder pb = new ProcessBuilder(cmd).directory(cwd.toFile()).redirectErrorStream(true);
pb.environment().putAll(hermeticEnv());
return pb;
}
private static Path initRepo(Path dir) throws Exception {
Files.createDirectories(dir);
git(dir, "init", "-q", "-b", "main");
@@ -121,16 +70,23 @@ class GitWorktreesTest {
}
private static String gitOutput(Path cwd, String... args) throws Exception {
Process p = gitProcessBuilder(cwd, args).start();
List<String> cmd = new java.util.ArrayList<>(List.of("git"));
cmd.addAll(List.of(args));
ProcessBuilder pb = new ProcessBuilder(cmd).directory(cwd.toFile()).redirectErrorStream(true);
pb.environment().put("GIT_CONFIG_GLOBAL", "/dev/null");
pb.environment().put("GIT_CONFIG_SYSTEM", "/dev/null");
pb.environment().put("GIT_TERMINAL_PROMPT", "0");
Process p = pb.start();
String out = new String(p.getInputStream().readAllBytes());
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git timed out: git " + String.join(" ", args));
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git timed out: " + String.join(" ", cmd));
assertEquals(0, p.exitValue(), "git " + String.join(" ", args) + " failed:\n" + out);
return out;
}
/** Pending changes to {@code file} in {@code cwd}, empty when git considers it unmodified. */
private static String status(Path cwd, String file) throws Exception {
Process p = gitProcessBuilder(cwd, "status", "--porcelain", "--", file).start();
Process p = new ProcessBuilder("git", "status", "--porcelain", "--", file)
.directory(cwd.toFile()).redirectErrorStream(true).start();
String out = new String(p.getInputStream().readAllBytes());
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git status timed out");
return out;
@@ -138,14 +94,16 @@ class GitWorktreesTest {
/** Every pending change in {@code cwd} — the whole-tree porcelain status, unlike {@link #status}. */
private static String fullStatus(Path cwd) throws Exception {
Process p = gitProcessBuilder(cwd, "status", "--porcelain").start();
Process p = new ProcessBuilder("git", "status", "--porcelain")
.directory(cwd.toFile()).redirectErrorStream(true).start();
String out = new String(p.getInputStream().readAllBytes());
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git status timed out");
return out;
}
private static String revParse(Path cwd, String ref) throws Exception {
Process p = gitProcessBuilder(cwd, "rev-parse", ref).start();
Process p = new ProcessBuilder("git", "-C", cwd.toString(), "rev-parse", ref)
.redirectErrorStream(true).start();
String out = new String(p.getInputStream().readAllBytes()).trim();
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git rev-parse timed out");
assertEquals(0, p.exitValue(), "git rev-parse " + ref + " failed:\n" + out);
@@ -154,7 +112,8 @@ class GitWorktreesTest {
/** The recursive file list of a commit's tree — used to check what a snapshot actually committed. */
private static String lsTree(Path cwd, String ref) throws Exception {
Process p = gitProcessBuilder(cwd, "ls-tree", "-r", "--name-only", ref).start();
Process p = new ProcessBuilder("git", "-C", cwd.toString(), "ls-tree", "-r", "--name-only", ref)
.redirectErrorStream(true).start();
String out = new String(p.getInputStream().readAllBytes());
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git ls-tree timed out");
assertEquals(0, p.exitValue(), "git ls-tree " + ref + " failed:\n" + out);
@@ -165,7 +124,8 @@ class GitWorktreesTest {
* snapshot's tree changed relative to its parent, the same shape {@code git status --porcelain}
* reports for the worktree it was taken from. */
private static Set<String> diffNameOnly(Path cwd, String from, String to) throws Exception {
Process p = gitProcessBuilder(cwd, "diff", "--name-only", from, to).start();
Process p = new ProcessBuilder("git", "-C", cwd.toString(), "diff", "--name-only", from, to)
.redirectErrorStream(true).start();
String out = new String(p.getInputStream().readAllBytes());
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git diff timed out");
assertEquals(0, p.exitValue(), "git diff " + from + ".." + to + " failed:\n" + out);
@@ -191,7 +151,8 @@ class GitWorktreesTest {
}
private static String forEachRef(Path cwd, String pattern) throws Exception {
Process p = gitProcessBuilder(cwd, "for-each-ref", pattern).start();
Process p = new ProcessBuilder("git", "-C", cwd.toString(), "for-each-ref", pattern)
.redirectErrorStream(true).start();
String out = new String(p.getInputStream().readAllBytes());
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git for-each-ref timed out");
assertEquals(0, p.exitValue(), "git for-each-ref " + pattern + " failed:\n" + out);
@@ -200,7 +161,8 @@ class GitWorktreesTest {
/** Write {@code content} as a blob into the object database; returns its sha. */
private static String blobOf(Path cwd, String content) throws Exception {
Process p = gitProcessBuilder(cwd, "hash-object", "-w", "--stdin").start();
Process p = new ProcessBuilder("git", "-C", cwd.toString(), "hash-object", "-w", "--stdin")
.redirectErrorStream(true).start();
p.getOutputStream().write(content.getBytes(StandardCharsets.UTF_8));
p.getOutputStream().close();
String out = new String(p.getInputStream().readAllBytes()).trim();
@@ -211,7 +173,8 @@ class GitWorktreesTest {
/** Build a single-file tree object from {@code blob}; returns the tree's sha. */
private static String treeOf(Path cwd, String path, String blob) throws Exception {
Process p = gitProcessBuilder(cwd, "mktree").start();
Process p = new ProcessBuilder("git", "-C", cwd.toString(), "mktree")
.redirectErrorStream(true).start();
p.getOutputStream().write(("100644 blob " + blob + "\t" + path + "\n").getBytes(StandardCharsets.UTF_8));
p.getOutputStream().close();
String out = new String(p.getInputStream().readAllBytes()).trim();
@@ -223,9 +186,10 @@ class GitWorktreesTest {
/** {@code git commit-tree} rooted at {@code tree} with a chosen committer date; returns the sha. */
private static String commitTree(Path cwd, String tree, String parent, String committerDate,
String message) throws Exception {
ProcessBuilder pb = gitProcessBuilder(cwd, "commit-tree", tree, "-p", parent, "-m", message);
ProcessBuilder pb = new ProcessBuilder("git", "-C", cwd.toString(), "commit-tree",
tree, "-p", parent, "-m", message);
pb.environment().put("GIT_COMMITTER_DATE", committerDate);
Process p = pb.start();
Process p = pb.redirectErrorStream(true).start();
String out = new String(p.getInputStream().readAllBytes()).trim();
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git commit-tree timed out");
assertEquals(0, p.exitValue(), "git commit-tree failed:\n" + out);
@@ -298,11 +262,6 @@ class GitWorktreesTest {
helper = !f() { printf 'username=%s\\npassword=%s\\n\\n' operator operator-secret; }; f
""");
// fleetd #369: the one deliberate exception to gitProcessBuilder. This test's whole point is
// that git must resolve `globalConfig` (a synthetic "operator's global config", never the
// real machine's) and then IGNORE it — so it cannot use the shared hermetic env, which would
// point GIT_CONFIG_GLOBAL at /dev/null and defeat the very thing under test. It never runs
// `git status`, so it does not need XDG_CONFIG_HOME isolation either.
ProcessBuilder pb = new ProcessBuilder("git", "credential", "fill")
.directory(Path.of(wt).toFile()).redirectErrorStream(true);
pb.environment().put("GIT_CONFIG_GLOBAL", globalConfig.toString());
@@ -369,16 +328,20 @@ class GitWorktreesTest {
Path worktree = Path.of(wt);
assertEquals("https://git.ltms.dev/akb/kb.git",
gitOutput(worktree, "remote", "get-url", "origin").trim());
assertEquals(1, gitExitCode(worktree, "config", "--worktree", "--get-regexp", "^url\\."),
assertEquals(1, exitCode("git", "-C", wt, "config", "--worktree", "--get-regexp", "^url\\."),
"no url.*.insteadOf rewrite should be added for an already-HTTPS origin");
}
/** Test-local exit-code probe, mirroring {@link GitWorktrees#exitCode} for an assertion the
* production class does not expose. */
private static int gitExitCode(Path cwd, String... args) throws Exception {
Process p = gitProcessBuilder(cwd, args).start();
private static int exitCode(String... command) throws Exception {
ProcessBuilder pb = new ProcessBuilder(command).redirectErrorStream(true);
pb.environment().put("GIT_CONFIG_GLOBAL", "/dev/null");
pb.environment().put("GIT_CONFIG_SYSTEM", "/dev/null");
pb.environment().put("GIT_TERMINAL_PROMPT", "0");
Process p = pb.start();
p.getInputStream().readAllBytes();
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "command timed out: git " + String.join(" ", args));
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "command timed out: " + String.join(" ", command));
return p.exitValue();
}
@@ -1506,346 +1469,4 @@ class GitWorktreesTest {
assertTrue(reportingAppender.list.isEmpty(),
"a null/empty overlay must log nothing, got:\n" + capturedMessages());
}
// ---- fleetd #362: seedSkills. Drives GitWorktrees#add end-to-end (not a bare worktree) so the
// real memberSkillsSource wiring is exercised, exactly like the credential-helper/origin tests
// above do for their own seams. ----
/** Write {@code content} as {@code <dir>/<skillName>/SKILL.md}, creating {@code dir} first. */
private static void writeSkill(Path dir, String skillName, String content) throws IOException {
Path skillFile = dir.resolve(skillName).resolve("SKILL.md");
Files.createDirectories(skillFile.getParent());
Files.writeString(skillFile, content);
}
/**
* fleetd #362 review fix, finding 2. {@code core.excludesFile}'s XDG-fallback branch
* ({@link GitWorktrees#previouslyEffectiveExcludesFileContent}) does not go through a {@code
* git} subprocess, so a first cut of it read {@code XDG_CONFIG_HOME}/{@code HOME} straight from
* the JVM's real environment — no test could isolate it, and on any machine carrying a real
* {@code ~/.config/git/ignore} (this repo's own dev machine does — measured, not assumed), every
* seeding test below silently composed with that real file instead of a controlled fixture.
* Every test that seeds at least one skill now constructs its {@link GitWorktrees} with this —
* an empty, machine-independent {@code XDG_CONFIG_HOME} (so the fallback resolves to a file that
* provably does not exist) plus the same {@code GIT_CONFIG_GLOBAL}/{@code GIT_CONFIG_SYSTEM}/
* {@code GIT_TERMINAL_PROMPT} isolation the {@link #git}/{@link #gitOutput} helpers already use
* for repo setup — so no test in this class can reach the real machine's home directory.
*/
private static Map<String, String> hermeticGitEnv(Path tmp) {
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());
}
/** {@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));
}
/** Acceptance criterion 2 (part 1): a worktree with no {@code .claude/} at all gets the skill
* copied in from the configured {@code memberSkillsSource}, structure and content intact. */
@Test
void seedSkillsCopiesIntoAWorktreeWithNoClaudeDirAtAll(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
Path skillsSource = tmp.resolve("skills-src");
writeSkill(skillsSource, "implementer", "IMPLEMENTER SKILL\n");
String wt = seedingGitWorktrees(tmp.resolve("wts"), skillsSource.toString(), tmp)
.add(repo.toString(), "cb-362-fresh", "HEAD");
assertEquals("IMPLEMENTER SKILL\n",
Files.readString(Path.of(wt, ".claude", "skills", "implementer", "SKILL.md")));
}
/** Acceptance criterion 2 (part 2) / invariant 1: a repo that already ships its own {@code
* implementer} skill keeps it byte-for-byte — fleetd's copy is never written over it, even
* though the configured source also carries a same-named skill with different content. */
@Test
void seedSkillsNeverOverwritesAReposOwnSkill(@TempDir Path tmp) throws Exception {
Path repo = tmp.resolve("repo");
Files.createDirectories(repo);
git(repo, "init", "-q", "-b", "main");
git(repo, "config", "user.email", "test@example.invalid");
git(repo, "config", "user.name", "Test");
writeSkill(repo.resolve(".claude/skills"), "implementer", "REPO OWN SKILL\n");
git(repo, "add", ".claude");
git(repo, "commit", "-q", "-m", "repo ships its own implementer skill");
Path skillsSource = tmp.resolve("skills-src");
writeSkill(skillsSource, "implementer", "FLEETD SKILL — must never land here\n");
String wt = new GitWorktrees(tmp.resolve("wts").toString(), null, skillsSource.toString())
.add(repo.toString(), "cb-362-repo-own", "HEAD");
assertEquals("REPO OWN SKILL\n",
Files.readString(Path.of(wt, ".claude", "skills", "implementer", "SKILL.md")),
"the repo's own committed skill must survive untouched");
}
/** Acceptance criterion 2 (part 3) / invariant 3: a misconfigured or missing {@code
* memberSkillsSource} must never fail the spawn — the worktree is still created. */
@Test
void seedSkillsIsBestEffortWhenSourceDoesNotExist(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
String missingSource = tmp.resolve("does-not-exist").toString();
String wt = new GitWorktrees(tmp.resolve("wts").toString(), null, missingSource)
.add(repo.toString(), "cb-362-missing-src", "HEAD");
assertTrue(Files.isDirectory(Path.of(wt)), "the spawn must still produce a worktree");
assertFalse(Files.exists(Path.of(wt, ".claude", "skills")),
"nothing should be seeded when the source directory does not exist");
assertTrue(capturedMessages().stream().anyMatch(m -> m.contains("is not a directory")),
"expected a warning naming the bad memberSkills source, got:\n" + capturedMessages());
}
/** Acceptance criterion 3: prove invariant 2 with a real git command — a freshly seeded skill
* must not appear in {@code git status --porcelain} for the worktree it was seeded into. */
@Test
void seedSkillsHidesSeededPathsFromGitStatus(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
Path skillsSource = tmp.resolve("skills-src");
writeSkill(skillsSource, "implementer", "IMPLEMENTER SKILL\n");
String wt = seedingGitWorktrees(tmp.resolve("wts"), skillsSource.toString(), tmp)
.add(repo.toString(), "cb-362-status", "HEAD");
assertEquals("", fullStatus(Path.of(wt)),
"a seeded skill must be invisible to git status, so it can never be staged or committed");
}
/** Invariant 2, the other direction: the exclude {@link #seedSkillsHidesSeededPathsFromGitStatus}
* proves is scoped to ONE worktree, not the whole repo. A second worktree of the same repo,
* provisioned with no {@code memberSkillsSource}, still reports an untracked {@code
* .claude/skills/} the ordinary way — proving the exclude did not leak in via the shared
* {@code .git/info/exclude} (which a linked worktree resolves to the repo's COMMON git dir). */
@Test
void seedSkillsExcludeDoesNotLeakIntoASiblingWorktree(@TempDir Path tmp) throws Exception {
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(), tmp);
// `plain` never seeds anything (memberSkillsSource is null, so seedSkills no-ops before it
// ever touches core.excludesFile), so it does not need the hermetic gitEnv seam.
GitWorktrees plain = new GitWorktrees(tmp.resolve("wts").toString());
String seededWt = seeding.add(repo.toString(), "cb-362-scope-a", "HEAD");
String plainWt = plain.add(repo.toString(), "cb-362-scope-b", "HEAD");
// Simulate the same untracked shape landing in the sibling worktree by hand, since `plain`
// was never configured with a memberSkillsSource to seed it itself.
writeSkill(Path.of(plainWt, ".claude", "skills"), "implementer", "unrelated untracked content\n");
assertEquals("", fullStatus(Path.of(seededWt)), "seeded worktree stays clean");
assertTrue(porcelainPaths(fullStatus(Path.of(plainWt))).contains(".claude/"),
"an unrelated worktree's own untracked .claude/ must still show up in its status — "
+ "the seeded worktree's exclude must not have leaked into it, got:\n"
+ fullStatus(Path.of(plainWt)));
}
/** Criterion 2's log shape, mirroring the {@code overlayParity} log assertions above: the
* denominator, what was seeded, and what was kept because the repo already had it. */
@Test
void seedSkillsLogsSeededAndKept(@TempDir Path tmp) throws Exception {
reportingLogger.setLevel(Level.INFO);
Path repo = tmp.resolve("repo");
Files.createDirectories(repo);
git(repo, "init", "-q", "-b", "main");
git(repo, "config", "user.email", "test@example.invalid");
git(repo, "config", "user.name", "Test");
writeSkill(repo.resolve(".claude/skills"), "hunter", "REPO OWN HUNTER\n");
git(repo, "add", ".");
git(repo, "commit", "-q", "-m", "repo ships hunter only");
Path skillsSource = tmp.resolve("skills-src");
writeSkill(skillsSource, "hunter", "FLEETD HUNTER\n");
writeSkill(skillsSource, "implementer", "FLEETD IMPLEMENTER\n");
seedingGitWorktrees(tmp.resolve("wts"), skillsSource.toString(), tmp)
.add(repo.toString(), "cb-362-log", "HEAD");
assertTrue(capturedMessages().stream().anyMatch(m ->
m.contains("seeded: implementer") && m.contains("kept the repo's own copy of: hunter")),
"expected a summary naming both the seeded and kept skills, got:\n" + capturedMessages());
}
/**
* fleetd #362 review fix — the "compose, don't replace" invariant, a third direction alongside
* {@link #seedSkillsHidesSeededPathsFromGitStatus} and
* {@link #seedSkillsExcludeDoesNotLeakIntoASiblingWorktree}. {@code core.excludesFile} is
* single-valued: the first cut of {@code excludeSeededSkillsFromGitStatus} pointed it at
* fleetd's own exclude file with {@code --replace-all}, which SHADOWS whatever excludesFile the
* worktree was already resolving (an operator's global config, most commonly) instead of adding
* to it. Concretely, this repo's own {@code .gitignore} does not ignore {@code target/} — only an
* operator's global excludesFile does — so every worker's {@code mvn clean install} would
* otherwise make {@code target/} appear as untracked, and CB-576's deliberately
* untracked-inclusive {@code hasUncommitted} would then read every such worktree as dirty
* forever, so {@code SessionManager} never cleans it up.
*
* <p>A synthetic "operator's global git config" is isolated via {@code GIT_CONFIG_GLOBAL}
* pointed at a throwaway temp file, passed to {@link GitWorktrees} through its {@code gitEnv}
* test seam — never the real machine's own git config. That global config ignores {@code
* target}. A skill is then seeded through the real {@link GitWorktrees#add} path, and a file
* named {@code target} is written into the worktree afterward: {@code git status --porcelain}
* must still be empty, proving the operator's own global pattern kept applying after seeding.
*/
@Test
void seedSkillsComposesWithAnAlreadyEffectiveGlobalExcludesFile(@TempDir Path tmp) throws Exception {
Path globalExcludes = tmp.resolve("operator-global-ignore");
Files.writeString(globalExcludes, "target\n");
Path globalConfig = tmp.resolve("operator-global.gitconfig");
Files.writeString(globalConfig, "[core]\n\texcludesFile = " + globalExcludes + "\n");
Map<String, String> gitEnv = Map.of(
"GIT_CONFIG_GLOBAL", globalConfig.toString(),
"GIT_CONFIG_SYSTEM", "/dev/null",
"GIT_TERMINAL_PROMPT", "0",
// core.excludesFile is explicitly set above, so the XDG fallback branch is never
// reached here — this is belt-and-braces so the test stays hermetic even if that
// ever changes, matching every other seeding test in this file.
"XDG_CONFIG_HOME", tmp.resolve("unused-xdg-config-home").toString());
Path repo = initRepo(tmp.resolve("repo"));
Path skillsSource = tmp.resolve("skills-src");
writeSkill(skillsSource, "implementer", "IMPLEMENTER SKILL\n");
GitWorktrees gitWorktrees = new GitWorktrees(tmp.resolve("wts").toString(), null, _ -> {},
null, null, skillsSource.toString(), gitEnv);
String wt = gitWorktrees.add(repo.toString(), "cb-362-global-compose", "HEAD");
assertEquals("IMPLEMENTER SKILL\n",
Files.readString(Path.of(wt, ".claude", "skills", "implementer", "SKILL.md")),
"fixture check — the skill really was seeded");
Files.writeString(Path.of(wt, "target"), "build output the operator's global config ignores\n");
String porcelain = fullStatus(Path.of(wt));
assertEquals("", porcelain,
"the operator's own global excludesFile pattern ('target') must still apply after "
+ "skill seeding ran — got:\n" + porcelain);
}
/**
* fleetd #362 review fix, finding 2: pins the XDG-fallback branch of {@link
* GitWorktrees#previouslyEffectiveExcludesFileContent}, exercised when {@code core.excludesFile}
* is unset entirely (no global, local, or worktree-scoped value at all) — the branch that used to
* read {@code XDG_CONFIG_HOME} straight from the JVM's own environment, unreachable by any test
* seam, and would silently compose with whatever real {@code ~/.config/git/ignore} happened to
* exist on the machine running the suite. {@code GIT_CONFIG_GLOBAL} points at an empty file (so
* {@code core.excludesFile} is genuinely unset, forcing the fallback branch to fire — not the
* "already configured" branch {@link #seedSkillsComposesWithAnAlreadyEffectiveGlobalExcludesFile}
* covers), and {@code XDG_CONFIG_HOME} is isolated through the {@code gitEnv} seam at a throwaway
* temp dir carrying a synthetic {@code git/ignore} that ignores {@code xdg-fallback-marker}. A
* skill is seeded through the real {@link GitWorktrees#add} path, and a file named {@code
* xdg-fallback-marker} is written into the worktree afterward: {@code git status --porcelain}
* must still be empty, proving the XDG-default pattern kept applying after seeding.
*
* <p>Deleting the fallback (so an unset key composes with {@code ""}) turns this test red with:
* {@code expected: <> but was: <?? xdg-fallback-marker\n>} — see the PR body for the pasted
* failure from actually running that mutation.
*/
@Test
void seedSkillsComposesWithTheXdgDefaultExcludesFileWhenNoneIsConfigured(@TempDir Path tmp) throws Exception {
Path xdgConfigHome = tmp.resolve("xdg-config-home");
Files.createDirectories(xdgConfigHome.resolve("git"));
Files.writeString(xdgConfigHome.resolve("git").resolve("ignore"), "xdg-fallback-marker\n");
Path emptyGlobalConfig = tmp.resolve("empty-global.gitconfig");
Files.writeString(emptyGlobalConfig, "");
Map<String, String> gitEnv = Map.of(
"GIT_CONFIG_GLOBAL", emptyGlobalConfig.toString(),
"GIT_CONFIG_SYSTEM", "/dev/null",
"GIT_TERMINAL_PROMPT", "0",
"XDG_CONFIG_HOME", xdgConfigHome.toString());
Path repo = initRepo(tmp.resolve("repo"));
Path skillsSource = tmp.resolve("skills-src");
writeSkill(skillsSource, "implementer", "IMPLEMENTER SKILL\n");
GitWorktrees gitWorktrees = new GitWorktrees(tmp.resolve("wts").toString(), null, _ -> {},
null, null, skillsSource.toString(), gitEnv);
String wt = gitWorktrees.add(repo.toString(), "cb-362-xdg-fallback", "HEAD");
assertEquals("IMPLEMENTER SKILL\n",
Files.readString(Path.of(wt, ".claude", "skills", "implementer", "SKILL.md")),
"fixture check — the skill really was seeded");
Files.writeString(Path.of(wt, "xdg-fallback-marker"),
"build output the XDG default ignore file (not core.excludesFile) covers\n");
String porcelain = fullStatus(Path.of(wt));
assertEquals("", porcelain,
"the XDG default excludesFile pattern ('xdg-fallback-marker') must still apply "
+ "after skill seeding ran — 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
* place {@link #hermeticEnv} is applied; a helper built directly, the way the original leak in
* {@link #status}/{@link #fullStatus} was, is now a source-level fact this test can catch by
* name instead of a machine-dependent failure someone has to rediscover.
*
* <p>This counts a literal marker in this very file's own source, split into three
* concatenated pieces below so the count is not thrown off by this method's own text — a
* plain, unsplit occurrence of the marker anywhere in this file (a helper's construction, or a
* comment that happens to spell it out contiguously) adds to the count the same way. Today
* there are exactly two: the factory itself, and the one documented exception in {@link
* #worktreeCredentialHelperCompletesWithoutUsingAnInheritedHelper}, which needs a
* non-hermetic, test-controlled global config to prove the credential helper ignores it. A
* third means a new helper was added the old, leak-prone way — route it through {@link
* #gitProcessBuilder} instead, or explain the new exception here and bump this number.
*/
@Test
void everyGitSubprocessGoesThroughTheHermeticFactory() throws Exception {
Path source = Path.of("src/test/java/dev/ltms/fleet/session/GitWorktreesTest.java");
String text = Files.readString(source);
String marker = "new " + "ProcessBuilder" + "(";
int count = 0;
for (int from = text.indexOf(marker); from >= 0; from = text.indexOf(marker, from + marker.length())) {
count++;
}
assertEquals(2, count,
"expected exactly 2 direct git-subprocess constructions in this file (the "
+ "gitProcessBuilder factory itself, plus the one documented exception in "
+ "worktreeCredentialHelperCompletesWithoutUsingAnInheritedHelper) — a "
+ "different count means a helper now bypasses the hermetic factory; route "
+ "it through gitProcessBuilder or document the new exception here");
}
/**
* fleetd #369 review round 2. {@link #everyGitSubprocessGoesThroughTheHermeticFactory} counts
* call sites, not behaviour — it catches a NEW helper built the old, leak-prone way, but it
* cannot catch {@link #gitProcessBuilder} itself being gutted: deleting {@code
* pb.environment().putAll(hermeticEnv())} from inside the factory leaves every call site
* unchanged, the count stays 2, and the whole unpoisoned suite stays green — the exact leak
* this ticket fixed would come back silently, with nothing but a human remembering to re-run
* the poison command to catch it. This test instead inspects what the factory actually hands
* to {@link ProcessBuilder#start()}, so it fails the moment the hermetic environment stops
* being applied, on any machine, with no poison needed.
*
* <p>The property under test: every git subprocess this class starts must run with an
* environment that cannot see the operator's real git configuration. A call-site count is a
* proxy for that; this is the thing itself.
*/
@Test
void gitProcessBuilderCarriesTheFullHermeticEnvironment(@TempDir Path tmp) {
Map<String, String> env = gitProcessBuilder(tmp, "status", "--porcelain").environment();
assertEquals("/dev/null", env.get("GIT_CONFIG_GLOBAL"),
"GIT_CONFIG_GLOBAL must be neutralized, or the operator's real ~/.gitconfig applies");
assertEquals("/dev/null", env.get("GIT_CONFIG_SYSTEM"),
"GIT_CONFIG_SYSTEM must be neutralized, or the machine's real /etc/gitconfig applies");
assertEquals("0", env.get("GIT_TERMINAL_PROMPT"),
"GIT_TERMINAL_PROMPT must be disabled, or a credential prompt can hang the subprocess");
String xdg = env.get("XDG_CONFIG_HOME");
assertNotNull(xdg,
"XDG_CONFIG_HOME must be set — left unset, git falls back to the operator's real "
+ "$HOME/.config/git/ignore (gitignore(5)), exactly fleetd #369's leak");
assertFalse(xdg.isBlank(), "XDG_CONFIG_HOME must not be blank — blank behaves like unset");
assertTrue(Path.of(xdg).startsWith(CLASS_TMP),
"XDG_CONFIG_HOME must point inside this test class's own throwaway directory, "
+ "never the operator's real one or the JVM's inherited value — got: " + xdg);
assertFalse(Files.exists(Path.of(xdg, "git", "ignore")),
"the resolved XDG default excludes file must provably not exist, or its contents "
+ "would silently apply to every git status this test class runs");
}
}
-337
View File
@@ -1,337 +0,0 @@
# Fleet as a Claude Code plugin — plan
Status: draft for architect review. Not implemented.
Author: primary (lead `opus`, Mac fleet). Date: 2026-09-05.
## 0. Correction — this already exists, and that changes the plan
I wrote sections below as if the plugin were new work. It is not. **This repo is already a Claude
Code marketplace and already ships a plugin**, added in `ef1e014` (CB-527) and last touched in
`2e138a1` (CB-634):
```
.claude-plugin/marketplace.json -> name "claude-bridge", plugins: [ ./plugin ]
plugin/.claude-plugin/plugin.json -> name "claude-bridge", version 0.1.0
plugin/.mcp.json -> mounts "fleetd" at http://127.0.0.1:8765/mcp
plugin/skills/setup/SKILL.md -> a full onboarding skill
plugin/README.md
```
The `setup` skill is good and covers most of what section 5 proposes: preflight, merge-not-clobber
into `.mcp.json`, read-only permissions only, credentials by env-var name, and a verify step that
insists on a **real spawn** because a green `/healthz` proves nothing.
So the operator's question — "can we pack things into plugins?" — is already answered *yes, and it
was built*. The real question is why it did nothing for the kb session. The answer is drift plus
invisibility.
### The drift, measured
| # | Finding | Evidence |
|---|---|---|
| 1 | **Mount name mismatch.** The plugin mounts the server as `fleetd`; the daemon's own constant is `fleet` | `plugin/.mcp.json` vs `PeerLauncher.java:34` `String MCP_MOUNT_NAME = "fleet"` |
| 2 | **URL hardcoded**, no env indirection, so one plugin cannot serve two hosts or ports | `plugin/.mcp.json` |
| 3 | **Ships no worker skills and no agents** | `plugin/` has 1 skill (`setup`); `.claude/skills/` has 5 and `.claude/agents/` has 3, none of them in `plugin/` |
| 4 | **Stale identity advice.** `setup` §5 tells the operator to pin `primary.terminal:` | CB-579 replaced that with `fleet.leaders.*.tab`. `record Primary` still exists (`FleetConfig.java:954`), so the advice is not dead — but it is no longer the mechanism |
| 5 | **Stale install path.** README says `/plugin marketplace add ltms/claude-bridge` | the repo is `fleet/fleetd` since CB-623 |
| 6 | **Stale names.** Plugin and marketplace are both `claude-bridge` | the project renamed to `fleetd` in CB-634 |
| 7 | **Nothing references it.** `grep -rn "plugin/" CLAUDE.md docs/*.md` returns nothing | so no session is ever told the plugin exists — which is exactly why I planned it from scratch |
Finding 7 is the root cause of the other six. A shipped capability that no instruction file
mentions gets no maintenance, and the next person rebuilds it. That is the same failure the
`CLAUDE.md` "Features" rule was written to stop.
Finding 3 is the one that matters most for the operator's actual problem. A worker spawned into a
**kb** worktree has no `implementer` skill, because only `claude-bridge` carries one in
`.claude/skills/`. Every brief that says "Load the implementer skill" is a no-op outside this repo.
The plugin is the right home for those skills and does not carry them yet.
## 1. The problem, restated
A project is "fleet-enabled" today by hand-edits nobody wrote down in one place:
- a `fleet.leaders.<name>` entry in a host's gitignored `fleetd.yaml`;
- the repo must carry `.claude/skills/*` for a worker to load `implementer` or `reviewer`;
- the repo must carry the canonical bridge block in its `CLAUDE.md`;
- the MCP mount arrives only because `LeadLauncher` and `ClaudeCodeLauncher` add `--mcp-config`
to the argv they build.
Shown live on 2026-09-05: an operator opened `claude` by hand in `/home/ltms/LTMS/kb` on fleet01
and the session had **no `fleet_*` tools at all**, because a hand-started agent never gets the
launcher's `--mcp-config`. The plugin would have fixed that — if it had been installed, and if it
had been mentioned anywhere.
## 2. What was verified, and how
| Claim | Evidence |
|---|---|
| A plugin can install globally | `~/.claude/plugins/installed_plugins.json` — scopes in use are `project` (5), `local` (4), **`user` (1)** |
| A plugin can mount an MCP server | `~/.claude/plugins/marketplaces/kb-alms/.mcp.json` mounts `memory` at `"url": "${KB_MEMORY_URL}"` |
| A plugin can carry skills, agents, commands, hooks | `kb-alms` ships `skills/` + `hooks/hooks.json`; `umputun-cc-thingz/plugins/planning` ships `agents/` + `skills/` |
| Env vars interpolate in a plugin's `.mcp.json` | same `kb-alms` file: `${KB_MEMORY_URL}`, `${MEMORY_MCP_TOKEN}` |
| A marketplace can be a plain git repo | `known_marketplaces.json` — `mgnl-code-review` has `"source": "git", "url": "https://..."` |
| **This repo is already such a marketplace** | `.claude-plugin/marketplace.json`, committed in `ef1e014` |
| A lead already binds to a project directory | `FleetConfig.java:1017` `record Leader(..., String workspace, String cwd)`; used at `LeadLauncher.java:193-199` |
| The lead's mount comes from argv, not config | `LeadLauncher.java:253` adds `--mcp-config` |
## 2b. Architect review + one measurement changed the design
The architect verified the plan against the code and returned **build it with these changes**. Two
of its findings are load-bearing. I checked both myself.
### A plugin cannot carry the agent definitions — confirmed
`ClaudeCodeLauncher.java:371` calls
`agentDefinitionFile(spec.cwd(), spec.role(), ".claude", "agents")`, and
`HerdrPeerLauncher.java:359-365` returns a path only when
`<cwd>/.claude/agents/<role>.md` `Files.isRegularFile`. `ClaudeCodeLauncher.java:391-393` adds
`--agent` **only** when that returns non-null. `OpenCodeLauncher` does the same for
`.opencode/agent`.
So the file must exist **in the member's worktree**. Moving `.claude/agents/*.md` into the plugin
would silently stop every member getting `--agent`. **The agents stay in the repo.** My plan had
this wrong.
### The plugin does not reach members at all — confirmed, and worse than the architect could see
`ClaudeCodeLauncher.java:285` does `putIfPresent(workerEnv, "CLAUDE_CONFIG_DIR", cfg.configDir())`.
A member with `configDir` set reads that directory, not the operator's `~/.claude`.
The architect could not check how far that goes, because `fleetd.yaml` is gitignored. I measured it:
- **All four Claude profiles on the live fleet set `configDir`** — `local`, `local-direct`, `opus`,
`sonnet` (`grep -c configDir fleetd/fleetd.yaml` = 4).
- Each `ccs` instance has its **own** `plugins/` directory: `gx10` (8 entries), `ltms` (11),
`ollama` (9), `work` (9).
- Those directories are **four separate real directories with four separate inodes**, and
`installed_plugins.json` in each is a **separate inode with an identical md5**
(`51c6e1c853e32e656b817e123fbbfcc5`). They are *copies made once*, not links.
So a plugin installed at user scope lands in exactly one instance's store. It would have to be
installed once per `CLAUDE_CONFIG_DIR`, and each copy would then drift. **The plugin is not a
delivery mechanism for member-facing assets on this host.**
A side effect worth recording: my own session's `CLAUDE_CONFIG_DIR` is set, so the
`~/.claude/plugins/*` evidence in section 2 is not even this session's store. The claims about what
a plugin *can* do still hold — they were read from real manifests — but the directory I read them
from is the wrong one for any conclusion about *this* session.
### The design that follows
Split by audience, not by mechanism:
| Audience | Delivered by | Carries |
|---|---|---|
| operator / lead (a human opening any project) | **the plugin**, per config dir | the MCP mount, `setup`, the bridge charter |
| member (a worker in a provisioned worktree) | **worktree provisioning** | `.claude/skills/*`, `.claude/agents/*` |
The second row is not a new idea — it is what the code already does for agents, and it is why
`agentDefinitionFile` looks in the worktree. Extending worktree provisioning to seed
`.claude/skills/` from a fleetd-owned source is the consistent move, and it is what actually fixes
"a worker in kb has no `implementer` skill". The plugin never could.
### Other review findings I accepted
- **Mount-name collision is real.** `PeerLauncher.java:34` is `fleet`; the plugin mounts `fleetd`.
A lead with both gets two mounts of one daemon and duplicate `fleet_*` tools. Rename the
plugin's server to `fleet`.
- **Keep `--mcp-config` in `LeadLauncher`.** It is config→argv from `profile.mcpUrl()`, not drift,
and it is the only path that works for a member with its own `configDir`.
- **I overstated #359.** `LeadCoordLoop.java:174-197` returns null and logs a warning that names the
fix; `tick()` leaves the message unacked, so the broker holds it and delivers once a lead is
named. It **stalls loudly and recovers** — it is not silent, and it is not data loss. My wording
in issue #361 needs the same correction.
- **Stage 1 ends with tools that mostly cannot be used until stage 3**, because authority still
comes from the tab. Section 3 already said this; the stage table did not.
## 3. What a plugin can and cannot do
This is the part that decides the design, so it is stated before the design.
**A plugin gives tools. It does not give authority.**
`fleet_whoami` resolves a caller's role from the connection, not from what is mounted. A
hand-started session in kb that mounts `fleet_*` through a plugin will be resolved as a **worker**
and refused on every orchestration call, because its pane is not in a tab matching
`fleet.leaders.*.tab`.
So the plugin alone does not make a project fleet-enabled. It makes it *tool*-enabled. Registering
the lead stays fleetd's job. Any plan that forgets this ships a plugin that looks installed and
does nothing.
```mermaid
flowchart TB
P["fleet plugin<br/>(user scope, every session)"] --> T["fleet_* tools mounted"]
D["fleetd.yaml<br/>leaders.kb {tab, cwd}"] --> A["role = primary"]
T --> W["can call fleet_*"]
A --> W2["calls are authorized"]
W --> OK["working lead"]
W2 --> OK
T --> NO["tools mounted, every call refused"]
classDef good fill:#2f855a,stroke:#22543d,color:#ffffff;
classDef bad fill:#9b2c2c,stroke:#63171b,color:#ffffff;
class OK good
class NO bad
```
*Both halves are needed. The plugin is the left half only.*
## 4. Proposed architecture — three layers
### Layer 1: the plugin — lead-side only, fix the one that exists
Keep it at `plugin/`, keep the marketplace at `.claude-plugin/marketplace.json`. Do not create a
second one, and do not put member-facing assets in it (see 2b).
```
.claude-plugin/marketplace.json -> rename to "fleetd"; keep source ./plugin
plugin/
.claude-plugin/plugin.json -> rename to "fleet"; bump version
.mcp.json -> mount name "fleet" (match PeerLauncher.MCP_MOUNT_NAME),
url "${FLEETD_MCP_URL}", 8765 default documented
skills/setup/SKILL.md -> EXISTS. fix the stale primary.terminal advice (§5)
skills/bridge-charter/SKILL.md -> NEW: the canonical CLAUDE.md block
README.md -> fix the install path (fleet/fleetd, not ltms/claude-bridge)
```
**Not in the plugin:** `agents/*.md` (the launcher requires them in the member's worktree —
`ClaudeCodeLauncher.java:371,391`) and the three worker skills (a member with `configDir` set never
reads the operator's plugin store — measured in 2b). Those belong to layer 1b.
**A rename is a breaking change for anyone who installed 0.1.0.** The mount name goes `fleetd` ->
`fleet`, so a project whose `.claude/settings.json` pre-allows `mcp__fleetd__fleet_whoami` stops
matching. Only this fleet has it installed today, so the cost is small now and grows. Decide once.
**This solves propagation of the charter.** `CLAUDE.md` says the bridge block "must stay
byte-identical with the template in the wiki" and that "other projects carrying the block need the
same edit" — a hand-copy the file itself admits is fragile, with a python snippet to check it. A
plugin skill turns that into a version bump.
### Layer 1b: worker skills reach members through the worktree, not the plugin
This is the change that actually fixes "a worker in kb cannot load `implementer`".
Worktree provisioning already writes into the member's tree — the parity overlay, the neutralised
`.mcp.json`, the IDE overlay. Add one more: seed `<worktree>/.claude/skills/` from a fleetd-owned
source directory, so every member gets `implementer`, `reviewer` and `hunter` whatever repo it is
working in. `.claude/agents/` is already required there by the launcher, so this follows the
grain of the design rather than cutting across it.
Open question for implementation: copy or symlink, and where the source lives (a config key such
as `memberSkills:`, or the plugin's own directory read by the daemon). A symlink is one source of
truth but breaks if the member's tree is archived; a copy drifts but is self-contained.
### Layer 2: host-global fleet settings
`~/.fleet/fleetd.yaml` — the things that are true for the **machine**, not the project:
- `broker:` and `coordinator:` (URIs come from env, no secrets in the file)
- `profiles:` — backends, models, credentials, weights
- `memberCredentials:` policy
- `worktreeRoot`, `worktreeGroup`
### Layer 3: per-project settings, committable
`<project>/.fleet/project.yaml` — the things that are true for the **repo**:
```yaml
lead:
tab: "lead: kb"
profile: opus
ide:
projectDir: "" # kb is a Python repo, everything at the root
worktree: true
```
**This split fixes a contradiction that exists today.** `ideProjectDir` is a property of a *repo*
(fleetd's Maven module is a subdirectory; kb's code is at the root) but the config key is
per-*profile*. One profile therefore cannot serve both repos — measured on fleet01 on 2026-09-05,
where the key had to be commented out to make kb work. Moving it to a project file removes the
contradiction rather than working around it.
It is also committable, because it holds no secrets. A project that has been fleet-enabled once
stays fleet-enabled for everyone who clones it.
## 5. The `fleet-setup` skill
What the operator actually asked for: one command that makes any project fleet-compatible.
```mermaid
sequenceDiagram
participant Op as Operator
participant Sk as fleet-setup skill
participant Fs as project files
participant Fd as fleetd
Op->>Sk: /fleet-setup (in any project)
Sk->>Fs: write .fleet/project.yaml
Sk->>Fs: add the bridge block to CLAUDE.md (if absent)
Sk->>Fd: register the lead (tab + cwd)
Fd-->>Sk: tab created, lead launched
Sk-->>Op: report what changed, and what is still manual
```
*The skill writes the project half and asks the daemon for the host half.*
The registration step needs something that does not exist yet: an MCP tool such as
`fleet_workspace_add{path, tab, profile}`, or a `fleetd` config include so a project file is picked
up without hand-editing the host file. **This is the one genuinely new piece of daemon work.**
## 6. Does this reduce fleetd's complexity?
Honestly: **partly**. Claiming more than this would be wrong.
**Yes, in three places.**
1. Skill and agent delivery leaves the daemon and the repos entirely.
2. Worktree config neutralisation (fleetd #134) gets safer. It blanks `.mcp.json` so the primary's
IDE and forge servers do not leak into a worker. Today the fleet mount survives only because
the launcher re-adds it by argv. With a user-scope plugin the fleet mount is outside the file
being neutralised, so the two concerns stop fighting.
3. The `ideProjectDir` per-profile/per-repo contradiction disappears.
**No, in the places that matter most.** fleetd still owns spawn, authorization, worktrees, herdr,
the broker, tickets, and identity. A plugin cannot do any of those. The plugin is a **distribution**
mechanism, not a replacement for the daemon.
**And it adds one new risk.** The mount URL becomes a second source of truth. `fleetd.yaml` has the
port; the plugin has the URL. Mitigation: the plugin reads `${FLEETD_MCP_URL}` only, and the host
env is the single place it is set.
## 7. Rollout stages
| Stage | Content | Ends with |
|---|---|---|
| 0 | **Make it visible.** One `CLAUDE.md` line and one Features entry saying the plugin exists and where | nobody re-plans it a third time |
| 1 | Fix the plugin's drift: mount name `fleet`, `${FLEETD_MCP_URL}`, names, README, stale `primary.terminal` advice | a lead in any project can install one plugin and get the mount |
| 1b | Seed `.claude/skills/` into provisioned worktrees | **a worker in *kb* can load `implementer`** |
| 2 | `.fleet/project.yaml` schema + `FleetConfig` reads it; `ideProjectDir` moves there | kb and fleetd both work off one profile |
| 3 | Fix #359, then config-include for lead registration, wired into the existing `setup` skill | `/fleet:setup` in a fresh project produces a working lead |
| 4 | Roll out to fleet01; retire the hand-copied CLAUDE.md block in favour of the skill | one `git pull` propagates the charter |
Stage 0 is minutes of work and is the one that stops this happening again, so it goes first.
**Stage 1b carries most of the value** and is independent of the plugin — it could ship first if the
plugin rename needs more thought. Stage 1 alone ends with tools a hand-started session mostly
cannot use, because authority still comes from the tab; that is fixed in stage 3, not stage 1.
#359 moves ahead of stage 3 on the architect's advice, because stage 3 is what creates the second
lead.
## 8. Questions for the architect
1. **Is the layer-2 / layer-3 split right?** Specifically: should `profiles:` stay host-global, or
should a project be able to pin which profiles it uses? Cost of getting this wrong is a config
that has to be re-split later.
2. **Config include, or a new MCP tool, for registering a project's lead?** An include is passive
and survives a restart; a tool is live but writes to a gitignored file the daemon owns.
3. **What happens when the plugin is absent?** Should `LeadLauncher` keep its `--mcp-config`
belt-and-braces, or is that the drift risk we should remove? Note opencode members cannot use
Claude plugins at all, so `OpenCodeLauncher` keeps its ephemeral config either way.
4. **Does a user-scope plugin mount leak into members in a way we do not want?** Members already
inherit user-scope MCP servers (`--mcp-config` adds, it does not replace). A worker getting
`fleet_*` is correct and already happens. Confirm nothing else in the plugin should be
worker-invisible.
5. **Two leads on one host both hold a subscription seat.** Is per-project leads the right unit, or
should one lead serve several projects by changing cwd?
6. **Blocking defect to fix first or alongside:** `LeadCoordLoop.resolveLocalLead()` (lines
174-190) routes a peer message to "the sole lead" when no lead is *named* after
`coordinator.selfId`. The moment a host has two leads — exactly what this plan encourages —
cross-host coordination silently stops. Tracked as #359.
+3 -3
View File
@@ -1,7 +1,7 @@
{
"name": "fleet",
"description": "Make a project fleet-ready: mount the fleetd MCP gateway and set up standard Claude Code settings so this session can orchestrate a fleet of delegated workers. Lead-side only — member skills and agents travel in the worktree. Ships no credentials.",
"version": "0.2.0",
"name": "claude-bridge",
"description": "Make a project bridge-ready: mount the fleetd MCP gateway and set up standard Claude Code settings so this session can orchestrate a fleet of delegated workers. Ships no credentials.",
"version": "0.1.0",
"author": {
"name": "LTMS"
},
+2 -2
View File
@@ -1,8 +1,8 @@
{
"mcpServers": {
"fleet": {
"fleetd": {
"type": "http",
"url": "${FLEETD_MCP_URL}"
"url": "http://127.0.0.1:8765/mcp"
}
}
}
+7 -33
View File
@@ -1,24 +1,11 @@
# fleet (Claude Code plugin)
# claude-bridge (Claude Code plugin)
Makes a project **fleet-ready**: mounts the `fleetd` MCP gateway and applies standard Claude Code
Makes a project **bridge-ready**: mounts the `fleetd` MCP gateway and applies standard Claude Code
settings, so the session can orchestrate a fleet of delegated workers.
**This plugin ships no credentials.** Every secret is referenced by environment-variable *name*;
the values stay with the user. Nothing the plugin writes is unsafe to commit.
## Scope — lead-side only
This plugin configures **the session you are sitting in**: a lead, or any human-started Claude Code
session that wants to talk to the daemon. It deliberately does **not** carry the worker playbook
skills or the role agent definitions, and it cannot:
- the launcher adds `--agent` only when `<worktree>/.claude/agents/<role>.md` exists in the
member's own tree (`ClaudeCodeLauncher.java:371,391`), so agent files must live in the repo;
- a member's `CLAUDE_CONFIG_DIR` points at its profile's config directory
(`ClaudeCodeLauncher.java:285`), so it never reads the operator's plugin store.
Member-facing assets travel in the worktree, not in this plugin. See fleetd #362.
## What it is not
The plugin is the **client-side setup**, not the bridge. `fleetd` is a separate daemon and `herdr`
@@ -29,35 +16,22 @@ not try to install system services on your behalf.
## Install
```shell
/plugin marketplace add https://git.ltms.dev/fleet/fleetd
/plugin install fleet@fleetd
```
Export the gateway URL — the plugin mounts `${FLEETD_MCP_URL}`, not a hardcoded address, so one
plugin serves hosts that run the daemon on different ports:
```shell
export FLEETD_MCP_URL=http://127.0.0.1:8765/mcp
/plugin marketplace add ltms/claude-bridge
/plugin install claude-bridge@claude-bridge
```
Then, in the project you want to onboard:
```shell
/fleet:setup
/claude-bridge:setup
```
## What you get
| Component | Effect |
|---|---|
| `.mcp.json` | mounts `fleet` at `${FLEETD_MCP_URL}` for any session with the plugin enabled |
| `skills/setup` | `/fleet:setup` — preflight, project settings, credential guidance, and verification |
The server is named **`fleet`** on purpose: that is `PeerLauncher.MCP_MOUNT_NAME` in the daemon and
the name a spawned member's own mount carries. Version 0.1.0 named it `fleetd`, which produced two
mounts of one daemon for anyone who also had a project-level `.mcp.json`. Upgrading from 0.1.0 is a
**breaking change** — a project that pre-allowed `mcp__fleetd__fleet_whoami` in
`.claude/settings.json` must be updated to `mcp__fleet__*`.
| `.mcp.json` | mounts `fleetd` at `http://127.0.0.1:8765/mcp` for any session with the plugin enabled |
| `skills/setup` | `/claude-bridge:setup` — preflight, project settings, credential guidance, and verification |
Because the plugin carries its own `.mcp.json`, an installed plugin needs no project-level MCP
file at all. The setup skill writes one only when you want the mount to work *without* the plugin —
+18 -43
View File
@@ -36,17 +36,9 @@ a time.
command -v herdr && herdr --version 2>&1 | head -1 || echo "MISSING: herdr"
command -v ccs && ccs version 2>&1 | head -1 || echo "MISSING: ccs (needed for worker profiles)"
command -v codex && codex --version 2>&1 | head -1 || echo "absent: codex (optional)"
curl -s -m 5 "${FLEETD_MCP_URL%/mcp}/healthz" 2>/dev/null \
|| curl -s -m 5 http://127.0.0.1:8765/healthz \
|| echo "MISSING: fleetd daemon is not reachable"
[ -n "$FLEETD_MCP_URL" ] && echo "FLEETD_MCP_URL is set" || echo "MISSING: FLEETD_MCP_URL"
curl -s -m 5 http://127.0.0.1:8765/healthz || echo "MISSING: fleetd daemon is not reachable"
```
**`FLEETD_MCP_URL` is required.** The plugin's own `.mcp.json` mounts `${FLEETD_MCP_URL}` rather
than a hardcoded address, so one plugin can serve hosts that run the daemon on different ports. If
it is unset the mount does not resolve. The usual value is `http://127.0.0.1:8765/mcp`; tell the
user to export it, do not write it into a file for them.
A healthy daemon answers with its status **and the herdr protocol it negotiated**:
```json
@@ -78,7 +70,7 @@ The entry to add, exactly:
```json
{
"mcpServers": {
"fleet": {
"fleetd": {
"type": "http",
"url": "http://127.0.0.1:8765/mcp"
}
@@ -86,13 +78,8 @@ The entry to add, exactly:
}
```
**The server must be named `fleet`.** That is `PeerLauncher.MCP_MOUNT_NAME` in the daemon, the name
a spawned member's mount carries, and the name the `mcp__fleet__*` role heuristic in `CLAUDE.md`
keys on. An earlier version of this plugin named it `fleetd`, which gave a lead with both a project
file and the plugin **two mounts of the same daemon** and a duplicated `fleet_*` tool set.
If `.mcp.json` already exists, add only the `fleet` key and leave every other server untouched.
If a `fleet` entry is already there with a different URL, **ask** rather than assuming yours is
If `.mcp.json` already exists, add only the `fleetd` key and leave every other server untouched.
If a `fleetd` entry is already there with a different URL, **ask** rather than assuming yours is
right — a non-default port usually means a deliberate second daemon.
> **If this plugin is installed, you can skip this step entirely.** The plugin ships its own
@@ -124,11 +111,11 @@ project already set.
"$schema": "https://json.schemastore.org/claude-code-settings.json",
"permissions": {
"allow": [
"mcp__fleet__fleet_whoami",
"mcp__fleet__fleet_list",
"mcp__fleet__fleet_status",
"mcp__fleet__fleet_profiles",
"mcp__fleet__fleet_poll"
"mcp__fleetd__fleet_whoami",
"mcp__fleetd__fleet_list",
"mcp__fleetd__fleet_status",
"mcp__fleetd__fleet_profiles",
"mcp__fleetd__fleet_poll"
]
}
}
@@ -174,31 +161,19 @@ fleet_whoami
```
- `{"role":"primary"}` — correct, you are done with this step.
- `{"role":"worker", …}` — **this is the trap.** If the lead runs inside a herdr pane whose tab the
daemon does not recognise, it is classified as a worker and refused on `spawn`/`send`/`stop`:
every verb an orchestrator exists to call. It is **self-locking**, because those are the same
calls that would tell the daemon who you are.
Identity is the **tab label**, matched exactly and case-insensitively:
- `{"role":"worker", …}` — **this is the trap.** If the primary runs inside a herdr pane, the
daemon resolves it to a terminal and classifies it as a worker, refusing `spawn`/`send`/`stop`:
every verb an orchestrator exists to call. It is **self-locking**, because the daemon can only
*learn* the primary's terminal from those same refused calls. The only way out is an
operator-set pin in the daemon's config:
```yaml
fleet:
leaders:
kb: # name it after coordinator.selfId if this host uses lead-to-lead
profile: opus
tab: "lead: kb" # the exact label of the tab this lead sits in
cwd: /path/to/the/project
primary:
terminal: term_xxxxxxxxxxxx # the terminalId fleet_whoami just reported
```
A tab label is stable across restarts of the agent inside it, which is why CB-579 replaced the
older `primary.terminal:` pin — a herdr `terminal_id` changed on every restart and cost a config
edit each time. `primary.terminal:` still parses, but it is no longer the mechanism; do not
reach for it.
The daemon reads `leaders:` **at boot**, so a new entry needs a restart. Two things to check
afterwards: that `fleet_whoami` now answers `primary`, and that no *stale* tab carries the same
label — duplicate lead tabs are their own failure (#359), and they stall lead-to-lead delivery
until one lead is named after `coordinator.selfId`.
The daemon reads this **at boot**, so it needs a restart. Re-pin whenever the primary moves
panes — a stale pin fails exactly as silently as no pin.
Then prove the fleet actually works, with a real spawn:
+6 -56
View File
@@ -116,50 +116,6 @@ check_log_path_matches_plist() {
ok "log path check: script and plist agree ($resolved_out)"
}
# Classify ERROR lines in one fresh log region. AMQP failure messages now include the connection
# name, so a recovery can clear only errors for its own connection. A candidate with neither name
# remains unexplained: it must never be quieted by a recovery on the other connection.
classify_amqp_connection_errors() {
local log_file="$1" line pending_inbox=0 pending_lead_mailbox=0
REDEPLOY_ERROR_COUNT=0
REDEPLOY_RECOVERED_AMQP_ERRORS=0
REDEPLOY_UNEXPLAINED_ERRORS=0
while IFS= read -r line || [ -n "$line" ]; do
case "$line" in
*' ERROR '*|*' SEVERE '*)
REDEPLOY_ERROR_COUNT=$((REDEPLOY_ERROR_COUNT + 1))
case "$line" in
*'AMQP connection fleetd-reply-inbox: An unexpected connection driver error occurred'*|*'AMQP connection fleetd-reply-inbox: Caught an exception during connection recovery!'*)
pending_inbox=$((pending_inbox + 1))
;;
*'AMQP connection fleetd-lead-mailbox: An unexpected connection driver error occurred'*|*'AMQP connection fleetd-lead-mailbox: Caught an exception during connection recovery!'*)
pending_lead_mailbox=$((pending_lead_mailbox + 1))
;;
*'AMQP connection'*'An unexpected connection driver error occurred'*|*'AMQP connection'*'Caught an exception during connection recovery!'*)
REDEPLOY_UNEXPLAINED_ERRORS=$((REDEPLOY_UNEXPLAINED_ERRORS + 1))
;;
*) REDEPLOY_UNEXPLAINED_ERRORS=$((REDEPLOY_UNEXPLAINED_ERRORS + 1)) ;;
esac
;;
*'AMQP connection recovered; cleared held replies for fresh redelivery'*)
if [ "$pending_inbox" -gt 0 ]; then
pending_inbox=$((pending_inbox - 1))
REDEPLOY_RECOVERED_AMQP_ERRORS=$((REDEPLOY_RECOVERED_AMQP_ERRORS + 1))
fi
;;
*'AMQP lead mailbox connection recovered; cleared held messages for fresh redelivery'*)
if [ "$pending_lead_mailbox" -gt 0 ]; then
pending_lead_mailbox=$((pending_lead_mailbox - 1))
REDEPLOY_RECOVERED_AMQP_ERRORS=$((REDEPLOY_RECOVERED_AMQP_ERRORS + 1))
fi
;;
esac
done < "$log_file"
REDEPLOY_UNEXPLAINED_ERRORS=$((REDEPLOY_UNEXPLAINED_ERRORS + pending_inbox + pending_lead_mailbox))
}
# CB-600: sourceable for testing. When this file is SOURCED (not executed) it stops here — nothing
# below runs — so a test harness can `source` it to call check_log_path_matches_plist (or the
# other pure helpers above) against a throwaway plist fixture without ever reaching the mutating
@@ -405,21 +361,15 @@ tail -n "+$((RESTART_MARK + 1))" "$OUT" 2>/dev/null \
| grep -iE 'deferred|classification:|fleet health:|coverage' | tail -8 | sed 's/^/ /' \
|| echo " (nothing reported)"
# Errors since the restart, anchored to the marker so old noise cannot leak in. Keep the fresh
# region in a file because the classifier must preserve the order of errors and recoveries.
FRESH_LOG="$(mktemp -t fleetd-fresh-log)"
trap 'rm -f "$FRESH_LOG"' EXIT
tail -n "+$((RESTART_MARK + 1))" "$OUT" > "$FRESH_LOG" 2>/dev/null || true
classify_amqp_connection_errors "$FRESH_LOG"
# Errors since the restart, anchored to the marker so old noise cannot leak in.
ERRS="$(tail -n "+$((RESTART_MARK + 1))" "$OUT" 2>/dev/null | grep -cE ' (ERROR|SEVERE) ' || true)"
say "result"
ok "pid $NEW_PID, jar $(jar_id)"
if [ "$REDEPLOY_ERROR_COUNT" -eq 0 ]; then
ok "no ERROR lines since restart"
elif [ "$REDEPLOY_UNEXPLAINED_ERRORS" -eq 0 ]; then
ok "$REDEPLOY_RECOVERED_AMQP_ERRORS AMQP connection reset ERROR lines recovered since restart"
if [ "${ERRS:-0}" -gt 0 ]; then
warn "$ERRS ERROR lines since restart:"
tail -n "+$((RESTART_MARK + 1))" "$OUT" | grep -E ' (ERROR|SEVERE) ' | tail -5 | sed 's/^/ /'
else
warn "$REDEPLOY_ERROR_COUNT ERROR lines since restart:"
grep -E ' (ERROR|SEVERE) ' "$FRESH_LOG" | tail -5 | sed 's/^/ /'
ok "no ERROR lines since restart"
fi
echo
echo " Next: call fleet_whoami and confirm it still answers 'primary'. A lead whose tab label"
-238
View File
@@ -1,238 +0,0 @@
#!/usr/bin/env bash
# Self-contained checks for the pure log classifier in redeploy-fleetd.sh.
set -euo pipefail
ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)"
TMP="$(mktemp -d "$ROOT/.redeploy-log-test.XXXXXX")"
trap 'rm -rf "$TMP"' EXIT
# Sourcing stops before redeploy-fleetd.sh can build, stop, or start the daemon.
source "$ROOT/scripts/redeploy-fleetd.sh"
fail() {
printf 'FAIL: %s\n' "$*" >&2
return 1
}
assert_equals() {
local expected="$1" actual="$2" description="$3"
[ "$expected" = "$actual" ] || fail "$description: expected $expected, got $actual"
}
classify_fixture() {
local name="$1"
classify_amqp_connection_errors "$TMP/$name"
}
test_no_errors() {
cat > "$TMP/no-errors.log" <<'LOG'
2026-09-05 12:00:00 INFO fleetd listening
LOG
classify_fixture no-errors.log
assert_equals 0 "$REDEPLOY_ERROR_COUNT" "no-errors total"
assert_equals 0 "$REDEPLOY_UNEXPLAINED_ERRORS" "no-errors unexplained"
}
test_recovery_patterns_match_source() {
grep -F 'AMQP connection {}: {}' "$ROOT/fleetd/src/main/java/dev/ltms/fleet/msg/AmqpReplyInbox.java" > /dev/null \
|| fail "AMQP failure pattern no longer matches source"
grep -F 'AMQP connection recovered; cleared held replies for fresh redelivery' \
"$ROOT/fleetd/src/main/java/dev/ltms/fleet/msg/AmqpReplyInbox.java" > /dev/null \
|| fail "reply-inbox recovery pattern no longer matches source"
grep -F 'AMQP lead mailbox connection recovered; cleared held messages for fresh redelivery' \
"$ROOT/fleetd/src/main/java/dev/ltms/fleet/msg/LeadMailbox.java" > /dev/null \
|| fail "lead-mailbox recovery pattern no longer matches source"
}
test_attributed_recovered_connection_error() {
cat > "$TMP/attributed-recovered.log" <<'LOG'
2026-09-05 12:00:00 ERROR [AMQP Connection broker:5672] d.l.fleet.msg.AmqpReplyInbox - AMQP connection fleetd-reply-inbox: An unexpected connection driver error occurred
2026-09-05 12:00:01 INFO [AMQP Connection broker:5672] d.l.fleet.msg.AmqpReplyInbox - AMQP connection recovered; cleared held replies for fresh redelivery
LOG
classify_fixture attributed-recovered.log
assert_equals 1 "$REDEPLOY_ERROR_COUNT" "attributed-recovered total"
assert_equals 1 "$REDEPLOY_RECOVERED_AMQP_ERRORS" "attributed-recovered errors"
assert_equals 0 "$REDEPLOY_UNEXPLAINED_ERRORS" "attributed-recovered unexplained"
}
test_source_derived_error_shapes_recover_by_connection() {
# These ERROR shapes come from AmqpConnectionFailureLogger on main. They need a live-log check
# after redeploy because the new code has not yet written a production line.
cat > "$TMP/source-derived.log" <<'LOG'
17:37:53.537 ERROR [AMQP Connection broker:5672] d.l.fleet.msg.AmqpReplyInbox - AMQP connection fleetd-reply-inbox: An unexpected connection driver error occurred
17:37:54.537 ERROR [AMQP Connection broker:5672] d.l.fleet.msg.AmqpReplyInbox - AMQP connection fleetd-reply-inbox: Caught an exception during connection recovery!
17:37:55.537 ERROR [AMQP Connection broker:5672] d.l.fleet.msg.AmqpReplyInbox - AMQP connection fleetd-reply-inbox: An unexpected connection driver error occurred
17:37:56.537 ERROR [AMQP Connection broker:5672] d.ltms.fleet.msg.LeadMailbox - AMQP connection fleetd-lead-mailbox: An unexpected connection driver error occurred
17:37:57.537 ERROR [AMQP Connection broker:5672] d.ltms.fleet.msg.LeadMailbox - AMQP connection fleetd-lead-mailbox: Caught an exception during connection recovery!
17:37:58.537 ERROR [AMQP Connection broker:5672] d.ltms.fleet.msg.LeadMailbox - AMQP connection fleetd-lead-mailbox: An unexpected connection driver error occurred
17:38:00.000 INFO [AMQP Connection broker:5672] d.l.fleet.msg.AmqpReplyInbox - AMQP connection recovered; cleared held replies for fresh redelivery
17:38:01.000 INFO [AMQP Connection broker:5672] d.l.fleet.msg.AmqpReplyInbox - AMQP connection recovered; cleared held replies for fresh redelivery
17:38:02.000 INFO [AMQP Connection broker:5672] d.l.fleet.msg.AmqpReplyInbox - AMQP connection recovered; cleared held replies for fresh redelivery
17:38:03.000 INFO [AMQP Connection broker:5672] d.ltms.fleet.msg.LeadMailbox - AMQP lead mailbox connection recovered; cleared held messages for fresh redelivery
17:38:04.000 INFO [AMQP Connection broker:5672] d.ltms.fleet.msg.LeadMailbox - AMQP lead mailbox connection recovered; cleared held messages for fresh redelivery
17:38:05.000 INFO [AMQP Connection broker:5672] d.ltms.fleet.msg.LeadMailbox - AMQP lead mailbox connection recovered; cleared held messages for fresh redelivery
LOG
classify_fixture source-derived.log
assert_equals 6 "$REDEPLOY_ERROR_COUNT" "source-derived total"
assert_equals 6 "$REDEPLOY_RECOVERED_AMQP_ERRORS" "source-derived recovered"
assert_equals 0 "$REDEPLOY_UNEXPLAINED_ERRORS" "source-derived unexplained"
}
test_cross_connection_unattributable_errors_stay_loud() {
# This candidate has neither stable connection name, so LeadMailbox recovery must not consume it.
cat > "$TMP/cross-unattributable.log" <<'LOG'
2026-09-05 12:00:00 ERROR [AMQP Connection broker:5672] unknown - AMQP connection: An unexpected connection driver error occurred
2026-09-05 12:00:01 ERROR [AMQP Connection broker:5672] unknown - AMQP connection: An unexpected connection driver error occurred
2026-09-05 12:00:02 INFO [AMQP Connection 10.10.20.13:5672] d.ltms.fleet.msg.LeadMailbox - AMQP lead mailbox connection recovered; cleared held messages for fresh redelivery
2026-09-05 12:00:03 INFO [AMQP Connection 10.10.20.13:5672] d.ltms.fleet.msg.LeadMailbox - AMQP lead mailbox connection recovered; cleared held messages for fresh redelivery
LOG
classify_fixture cross-unattributable.log
assert_equals 2 "$REDEPLOY_ERROR_COUNT" "cross-unattributable total"
assert_equals 0 "$REDEPLOY_RECOVERED_AMQP_ERRORS" "cross-unattributable recovered"
assert_equals 2 "$REDEPLOY_UNEXPLAINED_ERRORS" "cross-unattributable unexplained"
}
test_attributed_cross_connection_errors_stay_loud() {
# LeadMailbox recovery cannot heal AmqpReplyInbox errors.
cat > "$TMP/cross-attributed.log" <<'LOG'
2026-09-05 12:00:00 ERROR AmqpReplyInbox - AMQP connection fleetd-reply-inbox: An unexpected connection driver error occurred
2026-09-05 12:00:01 ERROR AmqpReplyInbox - AMQP connection fleetd-reply-inbox: An unexpected connection driver error occurred
2026-09-05 12:00:02 INFO LeadMailbox - AMQP lead mailbox connection recovered; cleared held messages for fresh redelivery
2026-09-05 12:00:03 INFO LeadMailbox - AMQP lead mailbox connection recovered; cleared held messages for fresh redelivery
LOG
classify_fixture cross-attributed.log
assert_equals 2 "$REDEPLOY_ERROR_COUNT" "cross-attributed total"
assert_equals 0 "$REDEPLOY_RECOVERED_AMQP_ERRORS" "cross-attributed recovered"
assert_equals 2 "$REDEPLOY_UNEXPLAINED_ERRORS" "cross-attributed unexplained"
}
test_attributed_unrecovered_connection_error() {
cat > "$TMP/unrecovered.log" <<'LOG'
2026-09-05 12:00:00 ERROR AmqpReplyInbox - AMQP connection fleetd-reply-inbox: An unexpected connection driver error occurred
LOG
classify_fixture unrecovered.log
assert_equals 1 "$REDEPLOY_ERROR_COUNT" "unrecovered total"
assert_equals 0 "$REDEPLOY_RECOVERED_AMQP_ERRORS" "unrecovered AMQP errors"
assert_equals 1 "$REDEPLOY_UNEXPLAINED_ERRORS" "unrecovered unexplained"
}
test_other_error_is_unexplained() {
cat > "$TMP/other-error.log" <<'LOG'
2026-09-05 12:00:00 ERROR dev.ltms.fleet.Fleetd - startup failed
2026-09-05 12:00:01 INFO dev.ltms.fleet.msg.AmqpReplyInbox - AMQP connection recovered; cleared held replies for fresh redelivery
LOG
classify_fixture other-error.log
assert_equals 1 "$REDEPLOY_ERROR_COUNT" "other-error total"
assert_equals 1 "$REDEPLOY_UNEXPLAINED_ERRORS" "other-error unexplained"
}
test_recovery_requirement_mutation_is_caught() {
classify_amqp_connection_errors() {
local log_file="$1" line
REDEPLOY_ERROR_COUNT=0
REDEPLOY_RECOVERED_AMQP_ERRORS=0
REDEPLOY_UNEXPLAINED_ERRORS=0
while IFS= read -r line || [ -n "$line" ]; do
case "$line" in
*' ERROR '*|*' SEVERE '*)
REDEPLOY_ERROR_COUNT=$((REDEPLOY_ERROR_COUNT + 1))
case "$line" in
*'AMQP connection fleetd-reply-inbox: An unexpected connection driver error occurred'*)
REDEPLOY_RECOVERED_AMQP_ERRORS=$((REDEPLOY_RECOVERED_AMQP_ERRORS + 1))
;;
*) REDEPLOY_UNEXPLAINED_ERRORS=$((REDEPLOY_UNEXPLAINED_ERRORS + 1)) ;;
esac
;;
esac
done < "$log_file"
}
if test_attributed_unrecovered_connection_error > "$TMP/mutation-output" 2>&1; then
fail "mutation accepted an unrecovered connection error"
fi
grep -F 'FAIL: unrecovered AMQP errors: expected 0, got 1' "$TMP/mutation-output" > /dev/null \
|| fail "mutation failed without the expected assertion"
printf 'Recovery mutation: FAIL: unrecovered AMQP errors: expected 0, got 1\n'
}
test_shared_counter_mutation_is_caught() {
classify_amqp_connection_errors() {
local log_file="$1" line pending=0
REDEPLOY_ERROR_COUNT=0
REDEPLOY_RECOVERED_AMQP_ERRORS=0
REDEPLOY_UNEXPLAINED_ERRORS=0
while IFS= read -r line || [ -n "$line" ]; do
case "$line" in
*' ERROR '*|*' SEVERE '*)
REDEPLOY_ERROR_COUNT=$((REDEPLOY_ERROR_COUNT + 1))
case "$line" in
*'AMQP connection'*'An unexpected connection driver error occurred'*|*'AMQP connection'*'Caught an exception during connection recovery!'*)
case "$line" in
*'fleetd-reply-inbox'*|*'fleetd-lead-mailbox'*) pending=$((pending + 1)) ;;
*) REDEPLOY_UNEXPLAINED_ERRORS=$((REDEPLOY_UNEXPLAINED_ERRORS + 1)) ;;
esac
;;
*) REDEPLOY_UNEXPLAINED_ERRORS=$((REDEPLOY_UNEXPLAINED_ERRORS + 1)) ;;
esac
;;
*'AMQP connection recovered; cleared held replies for fresh redelivery'*|*'AMQP lead mailbox connection recovered; cleared held messages for fresh redelivery'*)
if [ "$pending" -gt 0 ]; then
pending=$((pending - 1))
REDEPLOY_RECOVERED_AMQP_ERRORS=$((REDEPLOY_RECOVERED_AMQP_ERRORS + 1))
fi
;;
esac
done < "$log_file"
REDEPLOY_UNEXPLAINED_ERRORS=$((REDEPLOY_UNEXPLAINED_ERRORS + pending))
}
if test_attributed_cross_connection_errors_stay_loud > "$TMP/shared-mutation-output" 2>&1; then
fail "shared counter mutation accepted cross-connection recovery"
fi
grep -F 'FAIL: cross-attributed recovered: expected 0, got 2' "$TMP/shared-mutation-output" > /dev/null \
|| fail "shared counter mutation failed without the expected assertion"
printf 'Shared-counter mutation: FAIL: cross-attributed recovered: expected 0, got 2\n'
}
test_unattributable_quiet_mutation_is_caught() {
classify_amqp_connection_errors() {
local log_file="$1" line
REDEPLOY_ERROR_COUNT=0
REDEPLOY_RECOVERED_AMQP_ERRORS=0
REDEPLOY_UNEXPLAINED_ERRORS=0
while IFS= read -r line || [ -n "$line" ]; do
case "$line" in
*' ERROR '*|*' SEVERE '*)
REDEPLOY_ERROR_COUNT=$((REDEPLOY_ERROR_COUNT + 1))
case "$line" in
*'AMQP connection'*'An unexpected connection driver error occurred'*|*'AMQP connection'*'Caught an exception during connection recovery!'*)
REDEPLOY_RECOVERED_AMQP_ERRORS=$((REDEPLOY_RECOVERED_AMQP_ERRORS + 1))
;;
*) REDEPLOY_UNEXPLAINED_ERRORS=$((REDEPLOY_UNEXPLAINED_ERRORS + 1)) ;;
esac
;;
esac
done < "$log_file"
}
if test_cross_connection_unattributable_errors_stay_loud > "$TMP/unattributable-mutation-output" 2>&1; then
fail "unattributable mutation accepted an unknown connection"
fi
grep -F 'FAIL: cross-unattributable recovered: expected 0, got 2' "$TMP/unattributable-mutation-output" > /dev/null \
|| fail "unattributable mutation failed without the expected assertion"
printf 'Unattributable mutation: FAIL: cross-unattributable recovered: expected 0, got 2\n'
}
test_no_errors
test_recovery_patterns_match_source
test_attributed_recovered_connection_error
test_source_derived_error_shapes_recover_by_connection
test_cross_connection_unattributable_errors_stay_loud
test_attributed_cross_connection_errors_stay_loud
test_attributed_unrecovered_connection_error
test_other_error_is_unexplained
test_recovery_requirement_mutation_is_caught
test_shared_counter_mutation_is_caught
test_unattributable_quiet_mutation_is_caught
printf 'PASS: redeploy log classifier\n'