Compare commits
18 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 8557289dc0 | |||
| 380eb63277 | |||
| 6c61355f8f | |||
| 96c406b968 | |||
| a5d81c3f70 | |||
| 92c0f164f1 | |||
| 9e813ec179 | |||
| 395b3b5c46 | |||
| 6e9e464d62 | |||
| 105c065615 | |||
| 01492059d4 | |||
| f84824ee29 | |||
| c4d40fbc2b | |||
| 7c684e40d3 | |||
| 6938f52155 | |||
| 0df34f3220 | |||
| 457458437f | |||
| 3759c41f99 |
@@ -1,15 +1,15 @@
|
||||
{
|
||||
"name": "claude-bridge",
|
||||
"name": "fleetd",
|
||||
"description": "Tooling for orchestrating a fleet of delegated coding agents through the fleetd MCP gateway.",
|
||||
"owner": {
|
||||
"name": "LTMS"
|
||||
},
|
||||
"plugins": [
|
||||
{
|
||||
"name": "claude-bridge",
|
||||
"name": "fleet",
|
||||
"source": "./plugin",
|
||||
"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",
|
||||
"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",
|
||||
"author": {
|
||||
"name": "LTMS"
|
||||
}
|
||||
|
||||
@@ -210,6 +210,17 @@ 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
|
||||
|
||||
+46
-46
@@ -1,72 +1,72 @@
|
||||
# CB-504 — systemd unit for fleetd (Linux).
|
||||
#
|
||||
# The macOS launchd agent (deploy/dev.ltms.fleetd.plist) is the supervision target for the
|
||||
# current single-host deployment. This unit exists for the per-host gateways CB-308 introduces,
|
||||
# which will run on Linux.
|
||||
# fleetd #360: the previous version of this file started clean and broke the daemon in three ways
|
||||
# that nothing logs (see the DO NOT block and the ExecStart/PrivateTmp comments below for what and
|
||||
# why). The unit below, plus its companion deploy/herdr.service, is the version that has actually
|
||||
# run on fleet01 without those failures. Do not "improve" it back toward the old shape without
|
||||
# re-reading why each line is the way it is.
|
||||
#
|
||||
# Install (user service — fleetd drives the user's herdr, not a system daemon):
|
||||
# mkdir -p ~/.config/systemd/user
|
||||
# cp deploy/fleetd.service ~/.config/systemd/user/
|
||||
# # edit ExecStart / WorkingDirectory / Environment below, then:
|
||||
# cp deploy/fleetd.service deploy/herdr.service ~/.config/systemd/user/
|
||||
# # edit WorkingDirectory / ExecStart below for your host's paths and java location
|
||||
# systemctl --user daemon-reload
|
||||
# systemctl --user enable --now fleetd
|
||||
# systemctl --user enable --now herdr fleetd
|
||||
# journalctl --user -u fleetd -f
|
||||
#
|
||||
# Secrets (AI_GATEWAY_TOKEN, WORKER_GITEA_TOKEN, LAVINMQ_URI, COORD_AMQP_URI, ...) are not set
|
||||
# here and need no systemd drop-in: ExecStart runs a login shell, so they come from wherever your
|
||||
# login shell already sources them (this host: ~/.fleet/secrets.sh via ~/.zprofile). If a token is
|
||||
# missing there, fleetd still starts — the daemon reports every secret a configured profile
|
||||
# references, by name, never by value:
|
||||
# journalctl --user -u fleetd | grep 'startup secret'
|
||||
# A resolved one logs "startup secret NAME: set (profile 'x' tokenEnv)"; a missing one logs
|
||||
# "startup secret NAME: MISSING" at WARN and the daemon starts anyway — the first visible symptom
|
||||
# is a member that cannot open a pull request, hours later and in a different component.
|
||||
|
||||
[Unit]
|
||||
Description=fleetd — claude-bridge message server
|
||||
Description=fleetd — fleet message server
|
||||
Documentation=https://git.ltms.dev/fleet/fleetd/wiki
|
||||
# Ordering only: herdr is a user process and its socket may appear after us. This is advisory —
|
||||
# fleetd retries the herdr socket rather than exiting, which is what actually makes a late
|
||||
# socket survivable. Do NOT add Requires=: a herdr restart must not take fleetd down with it.
|
||||
# Ordering only. fleetd retries the herdr socket rather than exiting, which is what actually makes
|
||||
# a late socket survivable. Do NOT add Requires=: a herdr restart must not take fleetd down too.
|
||||
After=herdr.service
|
||||
Wants=herdr.service
|
||||
|
||||
[Service]
|
||||
Type=simple
|
||||
WorkingDirectory=%h/src/claude-bridge/fleetd
|
||||
ExecStart=/usr/lib/jvm/temurin-25-jdk/bin/java -jar target/fleetd.jar fleetd.yaml
|
||||
WorkingDirectory=%h/LTMS/fleetd/fleetd
|
||||
|
||||
Environment=HERDR_SOCKET_PATH=%h/.config/herdr/herdr.sock
|
||||
# PATH matters more than it looks (CB-511): fleetd propagates its own PATH to every worker it
|
||||
# spawns, so this line decides whether the fleet can run a build at all. systemd does not source a
|
||||
# login shell, so without it the daemon — and every worker — gets a bare default with no JDK/Maven.
|
||||
Environment=PATH=/usr/lib/jvm/temurin-25-jdk/bin:/usr/share/maven/bin:/usr/local/bin:/usr/bin:/bin
|
||||
# Secrets are NOT set here — this file is committed. Put ALL three tokens in a private drop-in
|
||||
# that systemd reads with restrictive permissions. In `systemctl --user edit fleetd`, add:
|
||||
# [Service]
|
||||
# Environment=FLEETD_API_TOKEN=...
|
||||
# Environment=WORKER_GITEA_TOKEN=...
|
||||
# Environment=AI_GATEWAY_TOKEN=...
|
||||
# FLEETD_API_TOKEN protects fleetd's API. WORKER_GITEA_TOKEN lets members open pull requests; if
|
||||
# it is missing, fleetd still starts, but a member fails when it later tries to open a pull request.
|
||||
# AI_GATEWAY_TOKEN authenticates gateway profiles; if it is missing, fleetd still starts, but a
|
||||
# gateway profile later returns HTTP 401. Or, put the same three variables in a 0600 file and add:
|
||||
# EnvironmentFile=%h/.config/fleetd/env
|
||||
# After starting, check which of them actually resolved. The daemon reports every secret a
|
||||
# configured profile references, by name, never by value:
|
||||
# journalctl --user -u fleetd | grep 'startup secret'
|
||||
# A resolved one logs "startup secret NAME: set (profile 'x' tokenEnv)". A missing one logs
|
||||
# "startup secret NAME: MISSING" at WARN — and the daemon starts anyway, which is the whole
|
||||
# problem: without this grep the first sign is a member that cannot open a pull request, hours
|
||||
# later and in a different component.
|
||||
# Note what the report can and cannot tell you. It lists only names some profile actually
|
||||
# references (tokenEnv, gitTokenEnv, and the broker uriEnv). A secret nothing references is never
|
||||
# reported, because nothing needs it.
|
||||
# A LOGIN shell, not java directly. Every secret this daemon needs (AI_GATEWAY_TOKEN,
|
||||
# WORKER_GITEA_TOKEN, LAVINMQ_URI, COORD_AMQP_URI) lives in ~/.fleet/secrets.sh, which only
|
||||
# ~/.zprofile sources. systemd runs no login shell. Started any other way the daemon boots fine
|
||||
# and looks healthy, and the failure appears hours later as a member that cannot open a pull
|
||||
# request. exec keeps it one process, so systemd tracks the right PID.
|
||||
# This also avoids a SECOND copy of the secrets in a systemd drop-in: one source of truth.
|
||||
ExecStart=/bin/zsh -lc "exec java -jar target/fleetd.jar fleetd.yaml"
|
||||
|
||||
# PrivateTmp MUST stay false -- see herdr.service. fleetd creates the member ZDOTDIR scrub dir and
|
||||
# the opencode config dir under java.io.tmpdir, and the member pane (a herdr child, a different
|
||||
# unit) has to read them. A private /tmp turns the credential scrub into a silent no-op.
|
||||
PrivateTmp=false
|
||||
|
||||
Restart=on-failure
|
||||
RestartSec=10s
|
||||
# A bad config (e.g. a non-loopback bind without token auth) makes fleetd fail fast by design.
|
||||
# Give up rather than restart-loop on a permanent error.
|
||||
# A bad config makes fleetd fail fast by design. Give up rather than restart-loop forever.
|
||||
StartLimitBurst=5
|
||||
StartLimitIntervalSec=120
|
||||
|
||||
# The daemon reads the repo, writes worktrees, and talks to a Unix socket — it needs no more.
|
||||
# DO NOT add ProtectSystem=, ProtectHome=, ProtectKernelTunables= or ProtectControlGroups=.
|
||||
# Measured on fleet01 2026-09-05: each of those gives the unit its own mount namespace, and
|
||||
# fleetd resolves a caller role by running lsof to find the loopback peer PID
|
||||
# (mcp/LsofPeerPidLookup). Inside such a namespace lsof returns nothing, every caller falls back
|
||||
# to ANONYMOUS, and the primary is refused every orchestration call with
|
||||
# "unauthenticated: anonymous may not SPAWN".
|
||||
# The daemon still starts, healthz still returns ok and the secrets still resolve - the only
|
||||
# symptom is that the fleet cannot be driven at all. Verified by bisecting the directives:
|
||||
# no sandbox 3 lsof lines | ProtectSystem=strict 0 | ProtectHome=read-only 0
|
||||
# ProtectKernelTunables 0 | ProtectControlGroups 0 | RestrictSUIDSGID 3 | NoNewPrivileges 3
|
||||
# The two below add no mount namespace and are safe.
|
||||
NoNewPrivileges=true
|
||||
PrivateTmp=true
|
||||
ProtectSystem=strict
|
||||
ProtectHome=read-write
|
||||
ProtectKernelTunables=true
|
||||
ProtectControlGroups=true
|
||||
RestrictSUIDSGID=true
|
||||
|
||||
StandardOutput=journal
|
||||
|
||||
Executable
+19
@@ -0,0 +1,19 @@
|
||||
#!/bin/zsh
|
||||
# fleetd #360 — template for the script deploy/herdr.service's ExecStart wraps in a pty.
|
||||
#
|
||||
# `script -qfec <this> /dev/null` needs a real command to run, and that command has to be a LOGIN
|
||||
# shell script: herdr itself needs the same secrets fleetd.service's login shell picks up (this
|
||||
# host: ~/.fleet/secrets.sh via ~/.zprofile), because members it spawns inherit its environment.
|
||||
# systemd's own Environment= lines in herdr.service are not enough for that -- they set TERM and a
|
||||
# bare PATH so the pty starts at all, nothing more.
|
||||
#
|
||||
# Copy this file to the path deploy/herdr.service's ExecStart names
|
||||
# (%h/LTMS/fleetd/fleetd-run/herdr-inner.sh by default) and `chmod +x` it. Not committed under
|
||||
# that path itself because the session name below is host-specific.
|
||||
|
||||
# A 0x0 pty makes every pane spawn fail with "ghostty error -2" (see herdr-multi-instance-facts /
|
||||
# fleet01-headless-herdr-standup) -- give it a real size before herdr ever touches it.
|
||||
stty rows 50 cols 200
|
||||
|
||||
# -l: login shell, so herdr and everything it spawns gets the real secrets and PATH.
|
||||
exec zsh -lc 'exec herdr --session <name>'
|
||||
@@ -0,0 +1,40 @@
|
||||
# fleetd #360 — systemd unit for herdr (Linux), the terminal multiplexer fleetd drives.
|
||||
#
|
||||
# This is fleetd.service's companion: fleetd.service's After=/Wants=herdr.service assumes this
|
||||
# unit exists. Before this ticket it did not, so on a fresh host fleetd started against a herdr
|
||||
# that systemd never supervised at all.
|
||||
#
|
||||
# Install: see deploy/fleetd.service's header comment (both units install the same way).
|
||||
#
|
||||
# ExecStart below runs deploy/herdr-inner.sh (copy the template of that name from this directory
|
||||
# to the path in ExecStart, or point ExecStart at wherever you keep it, and make it executable).
|
||||
# It is a separate file rather than an inline command because it must itself be a login shell (see
|
||||
# its own header for why) and systemd's ExecStart does not run one.
|
||||
|
||||
[Unit]
|
||||
Description=herdr terminal multiplexer (fleet session)
|
||||
Documentation=https://git.ltms.dev/fleet/fleetd/wiki
|
||||
|
||||
[Service]
|
||||
Type=simple
|
||||
# script(1) gives herdr a real pty. Without it the client reports a 0x0 window and every pane
|
||||
# spawn fails with "ghostty error -2" -- which surfaces as a fleetd spawn failure, not a herdr one.
|
||||
ExecStart=/usr/bin/script -qfec %h/LTMS/fleetd/fleetd-run/herdr-inner.sh /dev/null
|
||||
StandardInput=null
|
||||
Environment=TERM=xterm-256color
|
||||
Environment=PATH=%h/.local/bin:/usr/local/bin:/usr/bin:/bin
|
||||
|
||||
# PrivateTmp MUST stay false. fleetd writes the member ZDOTDIR scrub dir and the opencode config
|
||||
# dir under its own java.io.tmpdir, and the member pane -- a child of THIS process -- has to read
|
||||
# them. A private /tmp here silently breaks the credential scrub instead of failing loudly.
|
||||
PrivateTmp=false
|
||||
|
||||
Restart=on-failure
|
||||
RestartSec=5s
|
||||
|
||||
StandardOutput=journal
|
||||
StandardError=journal
|
||||
SyslogIdentifier=herdr
|
||||
|
||||
[Install]
|
||||
WantedBy=default.target
|
||||
@@ -774,6 +774,19 @@ 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)
|
||||
@@ -821,10 +834,17 @@ 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
|
||||
|
||||
@@ -248,7 +248,8 @@ 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()),
|
||||
SessionManager sessions = new SessionManager(workers,
|
||||
new GitWorktrees(cfg.worktreeRoot(), cfg.worktreeGroup(), cfg.memberSkills()),
|
||||
System::nanoTime, contextCap, clearAfterTurn);
|
||||
liveCountRef.set(profileName -> liveSessionCount(sessions.roster(), profileName));
|
||||
|
||||
@@ -652,7 +653,12 @@ public final class Fleetd {
|
||||
quarantineSource,
|
||||
leadMailbox,
|
||||
outageSource,
|
||||
new FleetMcp.LeadSeatSource(leadSeatLookup(() -> config.get().profiles(), leaders, leads)));
|
||||
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());
|
||||
|
||||
// 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
|
||||
|
||||
@@ -39,10 +39,11 @@ import java.util.function.Supplier;
|
||||
* keeps the old value until a restart: {@code lifecycle:}, {@code leadHeartbeat:},
|
||||
* {@code spawnReadyTimeoutMs} / {@code spawnReadyPollMs}, {@code quarantineCooldownSeconds}
|
||||
* (CB-578 stage B — baked once into the {@code BackendQuarantine} built at startup),
|
||||
* {@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,
|
||||
* {@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,
|
||||
* 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
|
||||
@@ -129,9 +130,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).</strong>
|
||||
* {@code FleetConfig} has 22 top-level record components: 5 cold, 11 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 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
|
||||
* 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;
|
||||
@@ -211,7 +213,7 @@ 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", "primary", "configReload",
|
||||
"guard", "worktreeRoot", "worktreeGroup", "memberSkills", "primary", "configReload",
|
||||
"leadHeartbeat", "lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs",
|
||||
"quarantineCooldownSeconds", "profiles");
|
||||
|
||||
@@ -402,6 +404,13 @@ 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:
|
||||
|
||||
@@ -106,6 +106,19 @@ 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.
|
||||
*/
|
||||
@JsonIgnoreProperties(ignoreUnknown = true)
|
||||
public record FleetConfig(
|
||||
@@ -130,7 +143,22 @@ public record FleetConfig(
|
||||
MemberCredentials memberCredentials,
|
||||
Coordinator coordinator,
|
||||
String worktreeGroup,
|
||||
String memberLoginShell) {
|
||||
String memberLoginShell,
|
||||
String memberSkills) {
|
||||
|
||||
/** Back-compat form before the {@code memberSkills} key was added. */
|
||||
public FleetConfig(Bind bind, String herdrSocket, String memberHerdrSocket, Map<String, Profile> profiles,
|
||||
Guard guard, String worktreeRoot, Lifecycle lifecycle, Integer spawnReadyTimeoutMs,
|
||||
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
|
||||
LeadHeartbeat leadHeartbeat, Health health, String placement, Auth auth,
|
||||
ConfigReload configReload, Integer quarantineCooldownSeconds,
|
||||
MemberCredentials memberCredentials, Coordinator coordinator, String worktreeGroup,
|
||||
String memberLoginShell) {
|
||||
this(bind, herdrSocket, memberHerdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
|
||||
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth,
|
||||
configReload, quarantineCooldownSeconds, memberCredentials, coordinator, worktreeGroup,
|
||||
memberLoginShell, null);
|
||||
}
|
||||
|
||||
/** Back-compat form before the {@code memberLoginShell} key was added. */
|
||||
public FleetConfig(Bind bind, String herdrSocket, String memberHerdrSocket, Map<String, Profile> profiles,
|
||||
@@ -888,12 +916,18 @@ 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) {
|
||||
public record Coordinator(String uri, String uriEnv, String selfId, Integer prefetch, List<String> peers) {
|
||||
|
||||
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. */
|
||||
@@ -1495,7 +1529,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");
|
||||
"memberCredentials", "coordinator", "worktreeGroup", "memberLoginShell", "memberSkills");
|
||||
|
||||
/** Load and validate config from {@code path}. */
|
||||
public static FleetConfig load(Path path) {
|
||||
@@ -2172,9 +2206,12 @@ 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/.
|
||||
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);
|
||||
quarantineCooldown, mc, coordinator, worktreeGroup, memberLoginShell, memberSkills);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -5,6 +5,7 @@ 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;
|
||||
@@ -53,17 +54,40 @@ 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 keeps
|
||||
* the previous answer instead of emptying it — a herdr hiccup must not silently demote a live lead
|
||||
* mid-session.
|
||||
* 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.
|
||||
*/
|
||||
public final class LeadTabScanner implements Supplier<Map<String, String>> {
|
||||
|
||||
@@ -79,6 +103,15 @@ 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)
|
||||
@@ -141,7 +174,7 @@ public final class LeadTabScanner implements Supplier<Map<String, String>> {
|
||||
return cached;
|
||||
}
|
||||
|
||||
/** One full pass: labelled tabs → their panes → those panes' terminals. */
|
||||
/** One full pass: labelled tabs → live agents in them → those panes' terminals. */
|
||||
private Map<String, String> scan() {
|
||||
Map<String, String> nameByTab = new LinkedHashMap<>();
|
||||
for (JsonNode w : herdr.call("workspace.list").path("workspaces")) {
|
||||
@@ -158,17 +191,47 @@ public final class LeadTabScanner implements Supplier<Map<String, String>> {
|
||||
}
|
||||
}
|
||||
|
||||
Map<String, String> byTerminal = new LinkedHashMap<>();
|
||||
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);
|
||||
}
|
||||
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);
|
||||
}
|
||||
}
|
||||
gracedTerminals = stillGraced;
|
||||
return Collections.unmodifiableMap(byTerminal);
|
||||
}
|
||||
|
||||
@@ -177,12 +240,14 @@ 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.
|
||||
* 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.
|
||||
*/
|
||||
private String leadNameOf(String label) {
|
||||
if (label == null) {
|
||||
return null;
|
||||
}
|
||||
return tabToName.get(label.strip().toLowerCase(Locale.ROOT));
|
||||
return tabToName.get(PendingCloseMarker.strip(label).toLowerCase(Locale.ROOT));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,42 @@
|
||||
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,6 +4,7 @@ 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;
|
||||
@@ -12,9 +13,11 @@ 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;
|
||||
|
||||
/**
|
||||
@@ -44,8 +47,33 @@ 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 #liveLeads}. A labelled
|
||||
* when herdr also reports a <em>running agent</em> in that tab — see {@link #countLeads}. 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 {
|
||||
|
||||
@@ -79,9 +107,9 @@ public final class LeadLauncher {
|
||||
return 0;
|
||||
}
|
||||
|
||||
Map<String, Integer> live;
|
||||
Map<String, LeadCount> live;
|
||||
try {
|
||||
live = liveLeads(leaders);
|
||||
live = countLeads(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.
|
||||
@@ -93,9 +121,53 @@ public final class LeadLauncher {
|
||||
for (Map.Entry<String, FleetConfig.Leader> e : leaders.entrySet()) {
|
||||
String name = e.getKey();
|
||||
FleetConfig.Leader lead = e.getValue();
|
||||
int running = live.getOrDefault(name, 0);
|
||||
LeadCount state = live.getOrDefault(name, LeadCount.NONE);
|
||||
int running = state.running();
|
||||
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;
|
||||
@@ -126,18 +198,29 @@ public final class LeadLauncher {
|
||||
}
|
||||
|
||||
/**
|
||||
* 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.
|
||||
* 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.
|
||||
*
|
||||
* <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 Map<String, Integer> liveLeads(Map<String, FleetConfig.Leader> leaders) {
|
||||
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) {
|
||||
// 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
|
||||
@@ -145,6 +228,7 @@ 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;
|
||||
@@ -153,31 +237,65 @@ 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());
|
||||
}
|
||||
}
|
||||
return counts;
|
||||
|
||||
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;
|
||||
}
|
||||
|
||||
/**
|
||||
* 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.
|
||||
* 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.
|
||||
*/
|
||||
private String leadNameOf(String label, Map<String, FleetConfig.Leader> leaders) {
|
||||
if (label == null) {
|
||||
return null;
|
||||
}
|
||||
String l = label.strip();
|
||||
String l = PendingCloseMarker.strip(label);
|
||||
for (Map.Entry<String, FleetConfig.Leader> e : leaders.entrySet()) {
|
||||
String tab = e.getValue().tabLabel();
|
||||
if (tab != null && l.equalsIgnoreCase(tab.strip())) {
|
||||
|
||||
@@ -44,6 +44,7 @@ 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;
|
||||
@@ -101,6 +102,8 @@ 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,
|
||||
@@ -164,6 +167,32 @@ 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
|
||||
@@ -211,19 +240,34 @@ public final class FleetMcp {
|
||||
}
|
||||
|
||||
/**
|
||||
* 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).
|
||||
* As above, with fleetd #176 lead-seat facts (see {@link LeadSeatSource}).
|
||||
*/
|
||||
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");
|
||||
@@ -361,7 +405,7 @@ public final class FleetMcp {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
|
||||
leadSeats, callers == null ? Map.of() : callers.leads(),
|
||||
callerTerminal(exchange),
|
||||
leadChannel == null ? null : leadChannel.selfCoordId());
|
||||
new CoordinationSource(leadChannel, peers));
|
||||
};
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> stopHandler =
|
||||
(exchange, req) -> {
|
||||
@@ -697,6 +741,15 @@ 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,
|
||||
@@ -724,7 +777,16 @@ public final class FleetMcp {
|
||||
+ ". Check that a daemon is running with coordinator.selfId=\"" + coordId
|
||||
+ "\" and is connected to the same coordination broker.");
|
||||
}
|
||||
return text("delivered to peer lead " + coordId + " (msgId " + msg.msgId() + ")");
|
||||
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);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1110,7 +1172,8 @@ 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, null);
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, leads, selfTerm,
|
||||
CoordinationSource.none());
|
||||
}
|
||||
|
||||
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
|
||||
@@ -1119,33 +1182,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, null);
|
||||
LeadSeatSource.none(), leads, selfTerm, CoordinationSource.none());
|
||||
}
|
||||
|
||||
/**
|
||||
* 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.
|
||||
* 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.
|
||||
*
|
||||
* @param selfCoordId this daemon's coord-id, or {@code null} when lead coordination is off
|
||||
* @param coordination this daemon's lead channel plus its declared peers, or
|
||||
* {@link CoordinationSource#none()} 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,
|
||||
String selfCoordId) {
|
||||
CoordinationSource coordination) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, OutageSource.none(),
|
||||
LeadSeatSource.none(), leads, selfTerm, selfCoordId);
|
||||
LeadSeatSource.none(), leads, selfTerm, coordination);
|
||||
}
|
||||
|
||||
/** 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, String selfCoordId) {
|
||||
Map<String, String> leads, String selfTerm, CoordinationSource coordination) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
|
||||
LeadSeatSource.none(), leads, selfTerm, selfCoordId);
|
||||
LeadSeatSource.none(), leads, selfTerm, coordination);
|
||||
}
|
||||
|
||||
/** As above, plus fleetd #176 lead-seat facts (see {@link LeadSeatSource}). */
|
||||
@@ -1153,7 +1216,7 @@ public final class FleetMcp {
|
||||
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
QuarantineSource quarantine, OutageSource outage,
|
||||
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
|
||||
String selfCoordId) {
|
||||
CoordinationSource coordination) {
|
||||
try {
|
||||
Map<String, Agent> live = workers.list().stream()
|
||||
.map(Agent.class::cast)
|
||||
@@ -1175,8 +1238,9 @@ public final class FleetMcp {
|
||||
Map<String, Object> result = new LinkedHashMap<>();
|
||||
result.put("leads", leadRows); result.put("members", out);
|
||||
result.put("healthCoverage", healthCoverage.value().get());
|
||||
if (selfCoordId != null && !selfCoordId.isBlank()) {
|
||||
result.put("coordinator", Map.of("selfId", selfCoordId, "configured", true));
|
||||
Map<String, Object> coordinatorRow = coordinatorView(coordination);
|
||||
if (coordinatorRow != null) {
|
||||
result.put("coordinator", coordinatorRow);
|
||||
}
|
||||
if (capacity.available()) result.put("capacity", profiles.stream()
|
||||
.map(profile -> capacityView(profile, capacity.liveCount(), capacity.maxLoad(), roster, messages,
|
||||
@@ -1187,6 +1251,140 @@ 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
|
||||
@@ -1377,8 +1575,12 @@ 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. Your own "
|
||||
+ "coordId is reported by fleet_list.")),
|
||||
+ "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.")),
|
||||
List.of("content")));
|
||||
}
|
||||
|
||||
@@ -1486,7 +1688,13 @@ 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'.",
|
||||
+ "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.",
|
||||
objectSchema(Map.of(), List.of()));
|
||||
}
|
||||
|
||||
|
||||
@@ -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, and ack what I have
|
||||
* delivered.
|
||||
* 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.
|
||||
*
|
||||
* <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,4 +37,78 @@ 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,6 +10,7 @@ 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;
|
||||
|
||||
@@ -259,6 +260,89 @@ 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();
|
||||
|
||||
@@ -92,6 +92,15 @@ 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. */
|
||||
@@ -125,7 +134,20 @@ public final class GitWorktrees implements Worktrees {
|
||||
* config); null/blank ⇒ {@link #shareWithGroup} is a no-op.
|
||||
*/
|
||||
public GitWorktrees(String configuredRoot, String group) {
|
||||
this(configuredRoot, 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);
|
||||
}
|
||||
|
||||
/** Test seam for changing a real worktree between its creation and its security check. */
|
||||
@@ -135,13 +157,20 @@ 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);
|
||||
this(configuredRoot, group, afterWorktreeAdded, null, 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);
|
||||
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);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -151,11 +180,30 @@ 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) {
|
||||
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) {
|
||||
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;
|
||||
@@ -184,6 +232,7 @@ 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;
|
||||
@@ -574,6 +623,310 @@ 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);
|
||||
@@ -1084,6 +1437,9 @@ 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);
|
||||
}
|
||||
@@ -1119,7 +1475,11 @@ public final class GitWorktrees implements Worktrees {
|
||||
private int exitCode(String... command) {
|
||||
Process p;
|
||||
try {
|
||||
p = new ProcessBuilder(command).redirectErrorStream(true).start();
|
||||
ProcessBuilder pb = new ProcessBuilder(command).redirectErrorStream(true);
|
||||
if (!gitEnv.isEmpty()) {
|
||||
pb.environment().putAll(gitEnv);
|
||||
}
|
||||
p = pb.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);
|
||||
var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, "mac-opus", null, 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);
|
||||
"mac-opus", 8, null);
|
||||
|
||||
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);
|
||||
var coordinator = new FleetConfig.Coordinator(null, "COORD_URI", "mac-opus", null, 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);
|
||||
var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, 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);
|
||||
var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, "mac-opus", null, null);
|
||||
|
||||
assertNull(Fleetd.openLeadMailbox(coordinator, Map.of(), opener),
|
||||
"a down coordination broker turns the feature off; it must never take the daemon down");
|
||||
|
||||
+4
-2
@@ -104,9 +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));
|
||||
v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-a", null, "self-a", 1, null));
|
||||
v.put("worktreeGroup", "group-a");
|
||||
v.put("memberLoginShell", null);
|
||||
v.put("memberSkills", "/skills/a");
|
||||
assertNamesMatchComponents(v);
|
||||
return v;
|
||||
}
|
||||
@@ -144,9 +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));
|
||||
v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-b", null, "self-b", 2, null));
|
||||
v.put("worktreeGroup", "group-b");
|
||||
v.put("memberLoginShell", null);
|
||||
v.put("memberSkills", "/skills/b");
|
||||
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);
|
||||
"amqp://stale-clear-text@127.0.0.1:5672/coord", "LEAD_COORD_URI", "fleet01-lead", null, 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,12 +1200,61 @@ 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);
|
||||
"amqp://guest:guest@127.0.0.1:5672/coord", null, 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");
|
||||
|
||||
+4
-3
@@ -45,8 +45,8 @@ import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
* <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} are left
|
||||
* as-is unconditionally, and {@code bind}/{@code guard}/{@code lifecycle}/{@code auth}/
|
||||
* {@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.
|
||||
@@ -92,9 +92,10 @@ class FleetConfigWithDefaultsPreservesEveryComponentTest {
|
||||
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));
|
||||
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;
|
||||
}
|
||||
|
||||
@@ -0,0 +1,232 @@
|
||||
package dev.ltms.fleet.deploy;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.regex.Pattern;
|
||||
import java.util.stream.Stream;
|
||||
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.params.ParameterizedTest;
|
||||
import org.junit.jupiter.params.provider.MethodSource;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #360: {@code deploy/fleetd.service} started clean on fleet01 and broke the daemon in
|
||||
* three ways that nothing logs at startup.
|
||||
*
|
||||
* <ul>
|
||||
* <li>{@code ProtectSystem=}, {@code ProtectHome=}, {@code ProtectKernelTunables=} and
|
||||
* {@code ProtectControlGroups=} each give the unit its own mount namespace. fleetd resolves
|
||||
* a caller's role by running {@code lsof} to find the loopback peer PID (see
|
||||
* {@code mcp/LsofPeerPidLookup}). Inside such a namespace {@code lsof} returns nothing,
|
||||
* every caller falls back to ANONYMOUS, and the primary is refused every orchestration call
|
||||
* with "unauthenticated: anonymous may not SPAWN" -- while healthz still reports ok.
|
||||
* Measured by bisection on fleet01 2026-09-05 (lsof line count): no sandbox 3,
|
||||
* {@code ProtectSystem=strict} 0, {@code ProtectHome=read-only} 0,
|
||||
* {@code ProtectKernelTunables} 0, {@code ProtectControlGroups} 0,
|
||||
* {@code RestrictSUIDSGID} 3, {@code NoNewPrivileges} 3 -- so only the first four are
|
||||
* forbidden.
|
||||
* <li>{@code PrivateTmp=true} gives the unit its own {@code /tmp}. fleetd writes the member
|
||||
* ZDOTDIR credential-scrub directory and the opencode config directory under
|
||||
* {@code java.io.tmpdir}, and the member pane -- a child of a *different* unit
|
||||
* (herdr.service) -- has to read them back. A private {@code /tmp} on either unit turns the
|
||||
* credential scrub into a silent no-op.
|
||||
* <li>Running {@code java} directly from {@code ExecStart} skips the login shell. Every secret
|
||||
* fleetd needs (AI_GATEWAY_TOKEN, WORKER_GITEA_TOKEN, LAVINMQ_URI, COORD_AMQP_URI) lives in
|
||||
* a file only the login shell sources; systemd runs no login shell on its own. Started that
|
||||
* way the daemon boots fine with empty credentials, and the failure appears hours later as a
|
||||
* member that cannot open a pull request.
|
||||
* </ul>
|
||||
*
|
||||
* <p>Review round 2 on fleetd #360 found the same shape one file further down the chain:
|
||||
* {@code herdr.service}'s {@code ExecStart} only names {@code deploy/herdr-inner.sh}, so that
|
||||
* script is the ONLY place two more of these properties live, and nothing was reading it.
|
||||
*
|
||||
* <ul>
|
||||
* <li>If the script does not exec a login shell, herdr -- and every member pane it spawns as a
|
||||
* child -- starts with none of the secrets only a login shell sources, the same silent
|
||||
* empty-credentials failure as fleetd.service's {@code ExecStart}, one process further away.
|
||||
* <li>If the script does not set a real, non-zero pty size before starting herdr (with
|
||||
* {@code stty}), herdr reports a 0x0 window and every pane spawn fails with
|
||||
* "ghostty error -2" -- which surfaces as a fleetd spawn failure, not a herdr one.
|
||||
* </ul>
|
||||
*
|
||||
* <p>None of these five failures makes the unit (or the script) fail to start, or makes
|
||||
* {@code /healthz} report unhealthy, so nothing short of reading the files catches a regression.
|
||||
* This test is that read.
|
||||
*
|
||||
* <p><b>It checks source text, not behaviour</b> -- it cannot start systemd or fork an mount
|
||||
* namespace in a build sandbox. It parses the same two files an operator would install and fails
|
||||
* if a forbidden directive is active, exactly the way {@code McpContractDocTest} guards
|
||||
* {@code docs/MCP-Contract.md} against naming a tool that does not exist.
|
||||
*/
|
||||
class SystemdUnitSafetyTest {
|
||||
|
||||
/** Tests run with the module directory (fleetd/) as cwd; deploy/ is the repo-root sibling. */
|
||||
private static final Path FLEETD_SERVICE = Path.of("../deploy/fleetd.service");
|
||||
private static final Path HERDR_SERVICE = Path.of("../deploy/herdr.service");
|
||||
private static final Path HERDR_INNER_SCRIPT = Path.of("../deploy/herdr-inner.sh");
|
||||
|
||||
private static final List<String> NAMESPACING_DIRECTIVES = List.of(
|
||||
"ProtectSystem", "ProtectHome", "ProtectKernelTunables", "ProtectControlGroups");
|
||||
|
||||
/**
|
||||
* A shell invoked with a login flag, e.g. {@code zsh -lc '...'} or {@code /bin/sh -l}. The
|
||||
* property under test is the {@code -l}, not the specific shell or the rest of its flags, so
|
||||
* this matches a token ending in "sh" followed by a flag cluster containing "l" -- tighter
|
||||
* than a bare {@code contains("-lc")}, which a "-lc" anywhere in the line would also satisfy.
|
||||
*/
|
||||
private static final Pattern LOGIN_SHELL_INVOCATION =
|
||||
Pattern.compile("(?:^|\\s)\\S*sh\\s+-[A-Za-z]*l[A-Za-z]*\\b");
|
||||
|
||||
private static final Pattern POSITIVE_ROWS = Pattern.compile("\\brows\\s+([1-9]\\d*)\\b");
|
||||
private static final Pattern POSITIVE_COLS = Pattern.compile("\\bcols\\s+([1-9]\\d*)\\b");
|
||||
|
||||
static Stream<Path> bothUnits() {
|
||||
return Stream.of(FLEETD_SERVICE, HERDR_SERVICE);
|
||||
}
|
||||
|
||||
private static boolean invokesALoginShell(List<String> activeLines) {
|
||||
return activeLines.stream().anyMatch(l -> LOGIN_SHELL_INVOCATION.matcher(l).find());
|
||||
}
|
||||
|
||||
/** An active {@code stty} line naming both a positive row count and a positive column count. */
|
||||
private static boolean setsANonZeroPtySize(List<String> activeLines) {
|
||||
return activeLines.stream()
|
||||
.filter(l -> l.startsWith("stty"))
|
||||
.anyMatch(l -> POSITIVE_ROWS.matcher(l).find() && POSITIVE_COLS.matcher(l).find());
|
||||
}
|
||||
|
||||
/** Lines that are actually in force: comments and blank lines don't count. */
|
||||
private static List<String> activeLines(Path unit) throws Exception {
|
||||
List<String> active = new ArrayList<>();
|
||||
for (String line : Files.readAllLines(unit)) {
|
||||
String stripped = line.strip();
|
||||
if (!stripped.isEmpty() && !stripped.startsWith("#")) {
|
||||
active.add(stripped);
|
||||
}
|
||||
}
|
||||
return active;
|
||||
}
|
||||
|
||||
@ParameterizedTest
|
||||
@MethodSource("bothUnits")
|
||||
@DisplayName("[SOURCE TEXT] no active mount-namespacing directive -- it blinds lsof and turns every caller ANONYMOUS")
|
||||
void doesNotActivateAMountNamespace(Path unit) throws Exception {
|
||||
List<String> active = activeLines(unit);
|
||||
|
||||
for (String directive : NAMESPACING_DIRECTIVES) {
|
||||
Pattern activeDirective = Pattern.compile("^" + Pattern.quote(directive) + "\\s*=");
|
||||
List<String> hits = active.stream().filter(l -> activeDirective.matcher(l).find()).toList();
|
||||
assertTrue(hits.isEmpty(),
|
||||
unit + " sets " + directive + " (" + hits + "). That directive gives the unit its "
|
||||
+ "own mount namespace; inside it, fleetd's lsof-based caller lookup "
|
||||
+ "(mcp/LsofPeerPidLookup) returns nothing, so every MCP caller falls back to "
|
||||
+ "ANONYMOUS and the primary is refused every orchestration call with "
|
||||
+ "\"unauthenticated: anonymous may not SPAWN\" -- while the daemon still "
|
||||
+ "starts and /healthz still reports ok. Measured on fleet01 2026-09-05 "
|
||||
+ "(fleetd #360). A commented-out mention in the file's own DO-NOT-add block "
|
||||
+ "is fine; an active directive is not.");
|
||||
}
|
||||
}
|
||||
|
||||
@ParameterizedTest
|
||||
@MethodSource("bothUnits")
|
||||
@DisplayName("[SOURCE TEXT] PrivateTmp is not true -- it silently no-ops the credential scrub")
|
||||
void privateTmpIsNotTrue(Path unit) throws Exception {
|
||||
List<String> active = activeLines(unit);
|
||||
Pattern privateTmpTrue = Pattern.compile("^PrivateTmp\\s*=\\s*true\\b");
|
||||
|
||||
boolean hasPrivateTmpTrue = active.stream().anyMatch(l -> privateTmpTrue.matcher(l).find());
|
||||
assertFalse(hasPrivateTmpTrue,
|
||||
unit + " sets PrivateTmp=true. fleetd writes the member ZDOTDIR credential-scrub "
|
||||
+ "directory and the opencode config directory under java.io.tmpdir, and the "
|
||||
+ "member pane -- a child of a DIFFERENT unit -- has to read them back. A "
|
||||
+ "private /tmp on either unit turns the credential scrub into a silent no-op: "
|
||||
+ "no error, no log line, the scrub just never happens (fleetd #360).");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SOURCE TEXT] fleetd.service's ExecStart goes through a login shell -- otherwise every secret is empty")
|
||||
void execStartUsesALoginShell() throws Exception {
|
||||
List<String> active = activeLines(FLEETD_SERVICE);
|
||||
List<String> execStartLines = active.stream().filter(l -> l.startsWith("ExecStart=")).toList();
|
||||
|
||||
assertTrue(execStartLines.size() == 1,
|
||||
FLEETD_SERVICE + " must have exactly one active ExecStart= line; found "
|
||||
+ execStartLines.size() + ": " + execStartLines);
|
||||
|
||||
String execStart = execStartLines.get(0);
|
||||
assertTrue(LOGIN_SHELL_INVOCATION.matcher(execStart).find(),
|
||||
FLEETD_SERVICE + "'s ExecStart (" + execStart + ") does not run a login shell (a shell "
|
||||
+ "invoked with a \"-l\" flag, e.g. \"zsh -lc\"). Every secret fleetd needs "
|
||||
+ "(AI_GATEWAY_TOKEN, WORKER_GITEA_TOKEN, LAVINMQ_URI, COORD_AMQP_URI) lives in a "
|
||||
+ "file only the login shell sources; systemd runs no login shell on its own. "
|
||||
+ "Running java directly boots fine with every credential empty, and the failure "
|
||||
+ "surfaces hours later as a member that cannot open a pull request (fleetd #360).");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SOURCE TEXT] herdr-inner.sh execs a login shell -- otherwise herdr and every member it spawns start with empty credentials")
|
||||
void herdrInnerScriptUsesALoginShell() throws Exception {
|
||||
List<String> active = activeLines(HERDR_INNER_SCRIPT);
|
||||
|
||||
assertTrue(invokesALoginShell(active),
|
||||
HERDR_INNER_SCRIPT + " does not exec a login shell (a shell invoked with a \"-l\" flag, "
|
||||
+ "e.g. \"zsh -lc\"). herdr.service's ExecStart only names this script, so this "
|
||||
+ "is the ONLY place herdr's login-shell property lives. Without it, herdr -- and "
|
||||
+ "every member pane it spawns as a child of herdr -- starts with none of the "
|
||||
+ "secrets that only a login shell sources (this host: ~/.fleet/secrets.sh via "
|
||||
+ "~/.zprofile), and the failure surfaces hours later as a member with no "
|
||||
+ "credentials at all, not just fleetd (fleetd #360).");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SOURCE TEXT] herdr-inner.sh sets a non-zero pty size before starting herdr -- otherwise every pane spawn fails with \"ghostty error -2\"")
|
||||
void herdrInnerScriptSetsANonZeroPtySize() throws Exception {
|
||||
List<String> active = activeLines(HERDR_INNER_SCRIPT);
|
||||
|
||||
assertTrue(setsANonZeroPtySize(active),
|
||||
HERDR_INNER_SCRIPT + " does not run \"stty rows N cols M\" with both N and M positive "
|
||||
+ "before starting herdr. Without a real pty size, herdr reports a 0x0 window and "
|
||||
+ "every pane spawn fails with \"ghostty error -2\" -- which surfaces as a fleetd "
|
||||
+ "spawn failure, not a herdr one, and gives no hint that the actual cause is this "
|
||||
+ "script (fleetd #360).");
|
||||
}
|
||||
|
||||
/**
|
||||
* The denominator guard, same shape as {@code McpContractDocTest}'s: a check that scans for a
|
||||
* forbidden pattern passes trivially if it is handed nothing to scan. Pin that all three files
|
||||
* exist, are non-trivial, and that the DO-NOT block's commented mentions are still there -- so
|
||||
* the "comment survives, directive doesn't" distinction above is actually being exercised.
|
||||
*/
|
||||
@Test
|
||||
@DisplayName("[SOURCE TEXT] the namespacing check is not vacuous -- all three files exist and the DO-NOT block still mentions every forbidden directive in a comment")
|
||||
void theCheckActuallyHasSomethingToCheck() throws Exception {
|
||||
assertTrue(Files.size(FLEETD_SERVICE) > 200,
|
||||
FLEETD_SERVICE + " is missing or unexpectedly small -- the checks above would pass "
|
||||
+ "vacuously against an empty or absent file");
|
||||
assertTrue(Files.size(HERDR_SERVICE) > 200,
|
||||
HERDR_SERVICE + " is missing or unexpectedly small -- the checks above would pass "
|
||||
+ "vacuously against an empty or absent file");
|
||||
assertTrue(Files.size(HERDR_INNER_SCRIPT) > 50,
|
||||
HERDR_INNER_SCRIPT + " is missing or unexpectedly small -- the login-shell and "
|
||||
+ "pty-size checks above would either pass vacuously or fail with an opaque "
|
||||
+ "\"file not found\" instead of naming the actual consequence (empty "
|
||||
+ "credentials / \"ghostty error -2\") against an empty or absent file");
|
||||
|
||||
String fleetdService = Files.readString(FLEETD_SERVICE);
|
||||
for (String directive : NAMESPACING_DIRECTIVES) {
|
||||
assertTrue(fleetdService.contains(directive),
|
||||
FLEETD_SERVICE + " no longer mentions " + directive + " anywhere, not even in the "
|
||||
+ "DO-NOT-add comment block that explains why it must stay out. That comment "
|
||||
+ "is the whole point of fleetd #360 -- it is what stops the next edit from "
|
||||
+ "re-adding the directive without knowing why.");
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -6,6 +6,7 @@ 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;
|
||||
@@ -35,6 +36,12 @@ 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;
|
||||
|
||||
@@ -53,6 +60,18 @@ 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++;
|
||||
@@ -83,6 +102,19 @@ 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);
|
||||
}
|
||||
}
|
||||
@@ -208,6 +240,114 @@ 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,6 +9,7 @@ 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.*;
|
||||
|
||||
@@ -106,6 +107,7 @@ 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");
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -122,6 +124,121 @@ 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,6 +1,7 @@
|
||||
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;
|
||||
@@ -93,4 +94,56 @@ 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,6 +10,9 @@ 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;
|
||||
@@ -34,7 +37,9 @@ 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.*;
|
||||
@@ -576,17 +581,48 @@ 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(), "", "mac-opus");
|
||||
FleetMcp.QuarantineSource.none(), Map.of(), "",
|
||||
new FleetMcp.CoordinationSource(channel, List.of()));
|
||||
|
||||
String out = textOf(res);
|
||||
// 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.
|
||||
// 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?).
|
||||
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
|
||||
@@ -601,6 +637,110 @@ 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();
|
||||
@@ -911,7 +1051,8 @@ 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(), "", null));
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), leadSeats, Map.of(), "",
|
||||
FleetMcp.CoordinationSource.none()));
|
||||
|
||||
assertTrue(out.contains("\"maxLoad\":3"), "maxLoad itself must be left untouched: " + out);
|
||||
assertTrue(out.contains("\"live\":2"), out);
|
||||
@@ -934,7 +1075,8 @@ 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(), "", null));
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), leadSeats, Map.of(), "",
|
||||
FleetMcp.CoordinationSource.none()));
|
||||
|
||||
assertTrue(out.contains("\"live\":0"), out);
|
||||
assertTrue(out.contains("\"free\":3"), "the real gate never subtracts the lead's seat: " + out);
|
||||
@@ -988,7 +1130,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(), "", null));
|
||||
FleetMcp.OutageSource.none(), leadSeats, Map.of(), "", FleetMcp.CoordinationSource.none()));
|
||||
int reportedFree = extractInt(out, "free");
|
||||
|
||||
for (int i = 0; i < reportedFree; i++) {
|
||||
@@ -1017,7 +1159,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(), "", null));
|
||||
Map.of(), "", FleetMcp.CoordinationSource.none()));
|
||||
|
||||
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)).withDefaults();
|
||||
new FleetConfig.Coordinator(null, coordinatorUriEnv, null, null, null)).withDefaults();
|
||||
}
|
||||
|
||||
/** {@code LC_*} categories are infrastructure by prefix; everything else needs an exact match. */
|
||||
|
||||
@@ -3,6 +3,8 @@ 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
|
||||
@@ -24,11 +26,19 @@ 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);
|
||||
@@ -65,6 +75,11 @@ 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);
|
||||
}
|
||||
|
||||
@@ -0,0 +1,77 @@
|
||||
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,5 +1,9 @@
|
||||
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;
|
||||
@@ -7,11 +11,14 @@ 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;
|
||||
|
||||
@@ -159,6 +166,145 @@ 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 {
|
||||
@@ -170,4 +316,18 @@ 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;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -12,6 +12,7 @@ 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;
|
||||
@@ -1469,4 +1470,273 @@ 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);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,337 @@
|
||||
# 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.
|
||||
@@ -1,7 +1,7 @@
|
||||
{
|
||||
"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",
|
||||
"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",
|
||||
"author": {
|
||||
"name": "LTMS"
|
||||
},
|
||||
|
||||
+2
-2
@@ -1,8 +1,8 @@
|
||||
{
|
||||
"mcpServers": {
|
||||
"fleetd": {
|
||||
"fleet": {
|
||||
"type": "http",
|
||||
"url": "http://127.0.0.1:8765/mcp"
|
||||
"url": "${FLEETD_MCP_URL}"
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+33
-7
@@ -1,11 +1,24 @@
|
||||
# claude-bridge (Claude Code plugin)
|
||||
# fleet (Claude Code plugin)
|
||||
|
||||
Makes a project **bridge-ready**: mounts the `fleetd` MCP gateway and applies standard Claude Code
|
||||
Makes a project **fleet-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`
|
||||
@@ -16,22 +29,35 @@ not try to install system services on your behalf.
|
||||
## Install
|
||||
|
||||
```shell
|
||||
/plugin marketplace add ltms/claude-bridge
|
||||
/plugin install claude-bridge@claude-bridge
|
||||
/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
|
||||
```
|
||||
|
||||
Then, in the project you want to onboard:
|
||||
|
||||
```shell
|
||||
/claude-bridge:setup
|
||||
/fleet:setup
|
||||
```
|
||||
|
||||
## What you get
|
||||
|
||||
| Component | Effect |
|
||||
|---|---|
|
||||
| `.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 |
|
||||
| `.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__*`.
|
||||
|
||||
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 —
|
||||
|
||||
@@ -36,9 +36,17 @@ 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 http://127.0.0.1:8765/healthz || echo "MISSING: fleetd daemon is not reachable"
|
||||
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"
|
||||
```
|
||||
|
||||
**`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
|
||||
@@ -70,7 +78,7 @@ The entry to add, exactly:
|
||||
```json
|
||||
{
|
||||
"mcpServers": {
|
||||
"fleetd": {
|
||||
"fleet": {
|
||||
"type": "http",
|
||||
"url": "http://127.0.0.1:8765/mcp"
|
||||
}
|
||||
@@ -78,8 +86,13 @@ The entry to add, exactly:
|
||||
}
|
||||
```
|
||||
|
||||
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
|
||||
**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
|
||||
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
|
||||
@@ -111,11 +124,11 @@ project already set.
|
||||
"$schema": "https://json.schemastore.org/claude-code-settings.json",
|
||||
"permissions": {
|
||||
"allow": [
|
||||
"mcp__fleetd__fleet_whoami",
|
||||
"mcp__fleetd__fleet_list",
|
||||
"mcp__fleetd__fleet_status",
|
||||
"mcp__fleetd__fleet_profiles",
|
||||
"mcp__fleetd__fleet_poll"
|
||||
"mcp__fleet__fleet_whoami",
|
||||
"mcp__fleet__fleet_list",
|
||||
"mcp__fleet__fleet_status",
|
||||
"mcp__fleet__fleet_profiles",
|
||||
"mcp__fleet__fleet_poll"
|
||||
]
|
||||
}
|
||||
}
|
||||
@@ -161,19 +174,31 @@ fleet_whoami
|
||||
```
|
||||
|
||||
- `{"role":"primary"}` — correct, you are done with this step.
|
||||
- `{"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:
|
||||
- `{"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:
|
||||
|
||||
```yaml
|
||||
primary:
|
||||
terminal: term_xxxxxxxxxxxx # the terminalId fleet_whoami just reported
|
||||
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
|
||||
```
|
||||
|
||||
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.
|
||||
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`.
|
||||
|
||||
Then prove the fleet actually works, with a real spawn:
|
||||
|
||||
|
||||
Reference in New Issue
Block a user