Compare commits
57 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 85c90d440a | |||
| ba51e0c6cc | |||
| 086c59848e | |||
| 0c10079755 | |||
| ece2091b53 | |||
| b3f917e6f5 | |||
| fa39a5f55e | |||
| d5128a1d35 | |||
| 61097e5cf0 | |||
| ef507bcd12 | |||
| 94f50e507a | |||
| 75b15086b0 | |||
| dab9645906 | |||
| e93b5f6512 | |||
| bfabe13e8f | |||
| 7e49c6eca2 | |||
| 11cbfa79b4 | |||
| c2c2746922 | |||
| 6ed70700a0 | |||
| 0087645da4 | |||
| 19cacf5b62 | |||
| 4dd12083ab | |||
| 30e21adec7 | |||
| 719b79f892 | |||
| 66e5247b6d | |||
| 1e41bd63b4 | |||
| 4887d03d88 | |||
| 7d5434455d | |||
| f71ee4926e | |||
| 282a2fc2b8 | |||
| f04e934b94 | |||
| 3fae35c357 | |||
| 18aecbfe67 | |||
| 5d75f72473 | |||
| e028a0ae54 | |||
| 2fa673d4c0 | |||
| 9020d01b40 | |||
| 1006805027 | |||
| 27aefbf9a0 | |||
| d42c2bc204 | |||
| 3916adc372 | |||
| fa97f598dd | |||
| ea9aa4fd77 | |||
| a2b8caf6b5 | |||
| b5ddbe5757 | |||
| d223a93039 | |||
| 96d8191149 | |||
| f0e7ac73d6 | |||
| e2fe861b4d | |||
| 279d6f5fbd | |||
| 34480cebef | |||
| 9d0bf14c46 | |||
| 5a3ab5764c | |||
| 3bad9f5785 | |||
| 8308c0b68f | |||
| 3b3063eb2b | |||
| df9086263d |
@@ -39,6 +39,19 @@ a worker made all 59 of its edits in the primary's tree and never noticed.
|
||||
test "$(git rev-parse --show-toplevel)" = "$PWD" || cd "$(git rev-parse --show-toplevel)"
|
||||
```
|
||||
|
||||
**Never run `git stash` (or `git stash pop`/`apply`/`drop`).** Your worktree is isolated, but the
|
||||
stash is **not**: `refs/stash` is one stack shared by the primary's checkout and every other
|
||||
worker's worktree of this repo. Measured on 2026-09-04 — `git stash list` from a worker's worktree
|
||||
and from the primary's tree returned byte-identical output. So a `git stash` you run can be popped
|
||||
into someone else's tree, and a `git stash pop` you run can drop **another worker's** uncommitted
|
||||
edits on top of yours. This has already happened here: two workers were running in parallel and one
|
||||
of them had its in-progress edit silently overwritten by the other's stash.
|
||||
|
||||
The branch is your isolation, so use it instead. To set work aside, commit it on your own branch
|
||||
(`git commit -m "wip: ..."`) and carry on; to try something and back out, use
|
||||
`git diff > /tmp/<your-branch>.patch` then `git checkout -- <file>`. Both stay inside your worktree.
|
||||
If you find a stash entry you did not create, leave it alone and say so in your report.
|
||||
|
||||
## 2. Implement
|
||||
|
||||
- Implement exactly the scope the lead named. Keep the diff focused; note anything out of scope
|
||||
|
||||
+20
-4
@@ -31,11 +31,27 @@ Environment=HERDR_SOCKET_PATH=%h/.config/herdr/herdr.sock
|
||||
# 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 the API/worker tokens in a private
|
||||
# drop-in that systemd reads with restrictive permissions:
|
||||
# systemctl --user edit fleetd → [Service] / Environment=FLEETD_API_TOKEN=...
|
||||
# or point EnvironmentFile at a 0600 file:
|
||||
# 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.
|
||||
|
||||
Restart=on-failure
|
||||
RestartSec=10s
|
||||
|
||||
+51
-29
@@ -307,16 +307,22 @@ profiles:
|
||||
# and no refusal on cost; the cap on live members is the single thing standing between a
|
||||
# fan-out and your monthly limit. Set it deliberately and keep it small.
|
||||
#
|
||||
# GOTCHA 3 (fleetd #176) — `maxLoad` counts members, never the lead itself. The lead is a live
|
||||
# `claude` session on this SAME account (a lead is never moved off-subscription, whatever its
|
||||
# own profile says), so it already holds one seat before any member spawns. If a lead's
|
||||
# `fleet.leaders.<name>.profile` names THIS profile — or ANY OTHER `subscription: true`
|
||||
# profile that shares this one's account (see THE SENTINEL, just below, next to
|
||||
# `credentialId:`) — `fleet_list`'s `free` for this profile subtracts that lead's live
|
||||
# seat(s) automatically; see `profile:` under THE FLEET below. If no lead entry names a
|
||||
# profile sharing this account, fleetd has no way to know a lead holds a seat here, and `free`
|
||||
# will overstate what a fresh `fleet_spawn` actually gets by exactly the seats the lead is
|
||||
# quietly holding.
|
||||
# GOTCHA 3 (fleetd #176, corrected by fleetd #257) — `maxLoad` counts members, never the lead
|
||||
# itself. The lead is a live `claude` session on this SAME account (a lead is never moved
|
||||
# off-subscription, whatever its own profile says), so it already holds one seat before any
|
||||
# member spawns. If a lead's `fleet.leaders.<name>.profile` names THIS profile — or ANY OTHER
|
||||
# `subscription: true` profile that shares this one's account (see THE SENTINEL, just below,
|
||||
# next to `credentialId:`) — `fleet_list` reports that seat count under `leadSeats`; see
|
||||
# `profile:` under THE FLEET below. `free` itself is NEVER reduced by `leadSeats`: `free` means
|
||||
# "what the real placement gate (`CompositePeerLauncher#enforceMaxLoad`) will actually grant a
|
||||
# fresh `fleet_spawn` right now", and that gate only ever compares live members against
|
||||
# `maxLoad` — it has no notion of the lead's own seat. An earlier cut of this feature
|
||||
# subtracted `leadSeats` from `free` on the theory it made `free` describe the true ceiling on
|
||||
# the account, but no backend seat ceiling shared with the lead has ever actually been
|
||||
# measured, and the subtraction just made `free` disagree with the one thing it is supposed to
|
||||
# describe — the fleetd #257 fix. `maxLoad: 3` means 3 member slots, full stop; a lead sharing
|
||||
# the account is a fact you can see in `leadSeats`, not a reason `free` undercounts spawns that
|
||||
# will, in practice, succeed.
|
||||
#
|
||||
# THE SENTINEL (fleetd #176 stage 2, correcting an inert stage 1 fix): every `subscription:
|
||||
# true` profile that leaves `credentialId` unset shares ONE implicit account-wide credential
|
||||
@@ -536,17 +542,18 @@ fleet:
|
||||
# auto-launched: it is also how fleetd learns which account this lead's own session shares. A
|
||||
# `subscription: true` profile bills the operator's Claude account, and the lead itself is always
|
||||
# a live `claude` session on that same account — `maxLoad` never counted that seat. If a lead
|
||||
# entry here names a profile that shares a worker profile's account, `fleet_list`'s `free` for
|
||||
# that worker profile subtracts the lead's live seat(s) automatically. "Shares the account" is
|
||||
# decided by matching `effectiveCredentialId()`, which (fleetd #176 stage 2 — see THE SENTINEL,
|
||||
# next to `credentialId:`, in THE WORKERS above) means: an explicit, matching `credentialId:` on
|
||||
# both, OR — the common case, needing NO extra config — both being `subscription: true` with
|
||||
# `credentialId` left unset, since those all share one implicit account-wide id. A lead on `opus`
|
||||
# and workers on `sonnet` link automatically this way; they do NOT need the same profile name.
|
||||
# Setting `profile:` on an already-running, recognise-only lead is safe — the daemon only launches
|
||||
# the SHORTFALL below `instances`, so naming a profile here does not, by itself, start anything.
|
||||
# Omit it and fleetd has no way to derive the sharing — there is no other reliable signal on the
|
||||
# daemon's side — so that lead's seat goes uncounted, exactly as before this ticket.
|
||||
# entry here names a profile that shares a worker profile's account, `fleet_list` reports the
|
||||
# lead's live seat(s) on that worker profile under `leadSeats` — informational only, as of fleetd
|
||||
# #257 it is NEVER subtracted from `free` (see GOTCHA 3, next to `maxLoad:`, in THE WORKERS above,
|
||||
# for why). "Shares the account" is decided by matching `effectiveCredentialId()`, which (fleetd
|
||||
# #176 stage 2 — see THE SENTINEL, next to `credentialId:`, in THE WORKERS above) means: an
|
||||
# explicit, matching `credentialId:` on both, OR — the common case, needing NO extra config — both
|
||||
# being `subscription: true` with `credentialId` left unset, since those all share one implicit
|
||||
# account-wide id. A lead on `opus` and workers on `sonnet` link automatically this way; they do
|
||||
# NOT need the same profile name. Setting `profile:` on an already-running, recognise-only lead is
|
||||
# safe — the daemon only launches the SHORTFALL below `instances`, so naming a profile here does
|
||||
# not, by itself, start anything. Omit it and fleetd has no way to derive the sharing — there is
|
||||
# no other reliable signal on the daemon's side — so that lead's seat never appears in `leadSeats`.
|
||||
#
|
||||
# `tab:` (CB-579) is REQUIRED and is the only field identity depends on — the exact label of the
|
||||
# tab hosting the lead, matched case-insensitively. Label the tab yourself and put that same
|
||||
@@ -669,20 +676,35 @@ guard:
|
||||
# every name here NOT also in `allow` is overlaid with a non-secret sentinel value before
|
||||
# the pane's login shell runs — real protection only for names that shell does not itself
|
||||
# re-export (see the ROUND-2 CORRECTION note above). Under allow-list: reporting only.
|
||||
# sshAuthSock → whether SSH_AUTH_SOCK may pass through under allow-list ("allow") or must be
|
||||
# blanked like any other non-derived name ("block", the default). This is a decision you
|
||||
# have to make explicitly: SSH_AUTH_SOCK is a handle to YOUR ssh-agent, and a member
|
||||
# holding it can sign with your keys — it sits in no secret file and looks like no
|
||||
# credential, which is why it slipped past three earlier tickets (gitea #110). Blocking
|
||||
# it breaks git over SSH inside members (push/fetch authenticate as you); use HTTPS
|
||||
# remotes or scoped deploy keys instead of allowing it lightly.
|
||||
# sshAgentEnv → whether SSH_AUTH_SOCK may pass through under allow-list ("inherit") or is omitted
|
||||
# from the member environment ("omit", the default). Omitting it only omits the
|
||||
# inherited ssh-agent path. It discourages automatic use of the operator's agent.
|
||||
# It does not deny same-user access to that socket. It also does not block SSH keys that
|
||||
# are readable on disk. Git over SSH may still work from inside a member. Keep the block:
|
||||
# it is correct and costs nothing, but it is not a control. A member runs as the same OS
|
||||
# user as the lead. Inside one uid, ordinary Unix permissions provide no meaningful
|
||||
# confidentiality boundary. A real boundary needs a different OS user or OS-level
|
||||
# confinement, such as a container or VM. That is the open question in fleetd #184.
|
||||
#
|
||||
# Still do not set this to "inherit" casually. SSH_AUTH_SOCK is a live handle to YOUR
|
||||
# ssh-agent, so a member holding it can sign with EVERY key the agent holds. It sits in
|
||||
# no secret file and looks like no credential, which is why it slipped past three
|
||||
# earlier tickets (gitea #110). Blocking it does not contain a member, but allowing it
|
||||
# hands one a signing capability for no gain — the block costs nothing, so keep it.
|
||||
#
|
||||
# Both halves of this are measured, not argued. 2026-08-28: a member with
|
||||
# SSH_AUTH_SOCK blanked pushed to the forge over SSH successfully, because `ssh -G`
|
||||
# resolves an IdentityFile outside ~/.ssh that is readable and has no passphrase. An
|
||||
# earlier version of this comment claimed blocking the socket BREAKS git over SSH. It
|
||||
# does not. That claim came from looking only in ~/.ssh, which holds nothing but four
|
||||
# `Include` lines — looking in one place and concluding about the whole host.
|
||||
#
|
||||
# HOT-RELOADABLE the same way `fleet:` is (CB-559): read fresh on every spawn, so editing this list
|
||||
# and reloading config (or restarting) changes what the NEXT spawn inherits; already-running members
|
||||
# are unaffected either way.
|
||||
# memberCredentials:
|
||||
# policy: deny-by-default # or "deny-list", or "allow-list" (CB-633) — see above
|
||||
# sshAuthSock: block # allow-list only; see the sshAuthSock note above
|
||||
# sshAgentEnv: omit # allow-list only; see the sshAgentEnv note above
|
||||
# allow:
|
||||
# - AI_GATEWAY_TOKEN # named in a profile's tokenEnv (local/gx) — a member reaching the
|
||||
# # gateway is by design, not a leak
|
||||
|
||||
@@ -29,6 +29,7 @@
|
||||
<commons-compress.version>1.27.1</commons-compress.version>
|
||||
<commons-lang3.version>3.18.0</commons-lang3.version>
|
||||
<sqlite-jdbc.version>3.53.4.0</sqlite-jdbc.version>
|
||||
<archunit.version>1.5.0</archunit.version>
|
||||
</properties>
|
||||
|
||||
<!--
|
||||
@@ -176,6 +177,14 @@
|
||||
<version>${testcontainers.version}</version>
|
||||
<scope>test</scope>
|
||||
</dependency>
|
||||
|
||||
<!-- fleetd #131: package-boundary and cycle enforcement (PackageCyclesTest). -->
|
||||
<dependency>
|
||||
<groupId>com.tngtech.archunit</groupId>
|
||||
<artifactId>archunit-junit5</artifactId>
|
||||
<version>${archunit.version}</version>
|
||||
<scope>test</scope>
|
||||
</dependency>
|
||||
</dependencies>
|
||||
|
||||
<build>
|
||||
|
||||
@@ -121,6 +121,7 @@ public final class Fleetd {
|
||||
// else can fail on a silently-empty one. A daemon started without a login shell (launchd)
|
||||
// boots fine either way — this is the only thing that says so out loud.
|
||||
reportRequiredSecrets(cfg);
|
||||
reportMemberTrustModel(cfg);
|
||||
// CB-596: an absent (or empty) memberCredentials: block blocks NOTHING — no credential
|
||||
// name is hardcoded any more to fall back on. Say so loudly, the same way a missing
|
||||
// secret is reported above, so upgrading past this commit never silently drops CB-592's
|
||||
@@ -249,9 +250,7 @@ public final class Fleetd {
|
||||
boolean clearAfterTurn = cfg.lifecycle() != null && cfg.lifecycle().clearAfterTurn();
|
||||
SessionManager sessions = new SessionManager(workers, new GitWorktrees(cfg.worktreeRoot(), cfg.worktreeGroup()),
|
||||
System::nanoTime, contextCap, clearAfterTurn);
|
||||
liveCountRef.set(profileName -> (int) sessions.roster().stream()
|
||||
.filter(s -> profileName.equals(s.profile()))
|
||||
.count());
|
||||
liveCountRef.set(profileName -> liveSessionCount(sessions.roster(), profileName));
|
||||
|
||||
// CB-303 part 1: idle-ttl reaper — only when configured, defaults to disabled.
|
||||
final SessionReaper reaper;
|
||||
@@ -592,7 +591,11 @@ public final class Fleetd {
|
||||
if (detail.agentSessionId() != null) {
|
||||
reason += " agentSessionId=" + detail.agentSessionId();
|
||||
}
|
||||
messages.abandon(detail.terminalId(), reason);
|
||||
// fleetd #275: this is an explicit teardown (fleet_stop, or the idle reaper) — the
|
||||
// worker's pane is being stopped right now, so an open fleet_ask has no turn left to
|
||||
// resume into. Sweep it too, unlike FleetHealthMonitor's health-classification call
|
||||
// (see MessageService.abandon's javadoc for why those two must differ).
|
||||
messages.abandon(detail.terminalId(), reason, true);
|
||||
replyInbox.release(detail.terminalId());
|
||||
primaryRegistry.forgetDelegation(detail.terminalId()); // CB-532: don't leak the lead binding
|
||||
});
|
||||
@@ -622,6 +625,19 @@ public final class Fleetd {
|
||||
log.info("auth: loopback-trust (any loopback non-worker caller is the primary)");
|
||||
}
|
||||
|
||||
// fleetd #297: named once and reused verbatim below for FleetApp's GET /profiles, rather than
|
||||
// built a second time — two independently-constructed sources reading the SAME BackendQuarantine
|
||||
// / BackendOutagePolicy would still be able to drift (e.g. a future edit to the credentialIdFor
|
||||
// closure in only one of the two places), exactly the shape #284 was.
|
||||
FleetMcp.QuarantineSource quarantineSource = new FleetMcp.QuarantineSource(profile -> {
|
||||
var configured = config.get().profiles().get(profile);
|
||||
return configured == null ? null : configured.effectiveCredentialId();
|
||||
}, quarantine);
|
||||
FleetMcp.OutageSource outageSource = new FleetMcp.OutageSource(profile -> {
|
||||
var configured = config.get().profiles().get(profile);
|
||||
return configured == null ? null : configured.effectiveCredentialId();
|
||||
}, outagePolicy);
|
||||
|
||||
FleetMcp mcp = new FleetMcp(messages, workers, sessions, identity, presence,
|
||||
primaryRegistry, callers, metrics, new FleetMcp.CapacitySource(profile -> liveCountRef.get().apply(profile),
|
||||
profile -> {
|
||||
@@ -633,15 +649,9 @@ public final class Fleetd {
|
||||
return FleetHealthMonitor.coverage(health != null && health.isEnabled(),
|
||||
health != null && health.notifications() != null && health.notifications().configured());
|
||||
}),
|
||||
new FleetMcp.QuarantineSource(profile -> {
|
||||
var configured = config.get().profiles().get(profile);
|
||||
return configured == null ? null : configured.effectiveCredentialId();
|
||||
}, quarantine),
|
||||
quarantineSource,
|
||||
leadMailbox,
|
||||
new FleetMcp.OutageSource(profile -> {
|
||||
var configured = config.get().profiles().get(profile);
|
||||
return configured == null ? null : configured.effectiveCredentialId();
|
||||
}, outagePolicy),
|
||||
outageSource,
|
||||
new FleetMcp.LeadSeatSource(leadSeatLookup(() -> config.get().profiles(), leaders, leads)));
|
||||
|
||||
// CB-637: the receive half. Only constructed when a lead mailbox actually opened — with no
|
||||
@@ -712,9 +722,12 @@ public final class Fleetd {
|
||||
// GET /sessions must merge across both, or a down/unpolled member daemon is invisible.
|
||||
// fleetd #111: live (re-read-per-request) memberCredentials view for GET /member-credentials —
|
||||
// same hot-reload shape as the memberCredentials supplier passed to ClaudeCodeLauncher above.
|
||||
// fleetd #297: quarantineSource/outageSource are the SAME instances passed to FleetMcp above —
|
||||
// GET /profiles must report the identical quarantine/cool-off facts as fleet_profiles.
|
||||
Javalin app = new FleetApp(herdr, memberHerdr, workers, sessions, messages, presence, mcp.servlet(),
|
||||
callers, metrics, deliverable,
|
||||
() -> MemberCredentialPolicyView.of(config.get().memberCredentials())).build();
|
||||
() -> MemberCredentialPolicyView.of(config.get().memberCredentials()),
|
||||
quarantineSource, outageSource).build();
|
||||
app.start(cfg.bind().host(), cfg.bind().port());
|
||||
log.info("fleetd listening on {}:{}, herdr socket {}",
|
||||
cfg.bind().host(), cfg.bind().port(), socket);
|
||||
@@ -854,6 +867,19 @@ public final class Fleetd {
|
||||
.orElse(null);
|
||||
}
|
||||
|
||||
/**
|
||||
* Count sessions that occupy a profile's spawn capacity. A {@code BACKEND_ERROR} or
|
||||
* {@code FAILED} session stays in the roster so {@code fleet_list} can show its failure, but a
|
||||
* member that cannot accept another delivery does not use a seat.
|
||||
*/
|
||||
static int liveSessionCount(List<MemberSession> roster, String profileName) {
|
||||
return (int) roster.stream()
|
||||
.filter(session -> profileName.equals(session.profile()))
|
||||
.filter(session -> session.state() != MemberSession.State.BACKEND_ERROR)
|
||||
.filter(session -> session.state() != MemberSession.State.FAILED)
|
||||
.count();
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #248 / fleetd#201 Unit 5: factory for the production {@link BackendErrorSink} — the
|
||||
* collaborator {@link CompletionResolver} notifies when a pane-scrape classification actually
|
||||
@@ -1136,6 +1162,24 @@ public final class Fleetd {
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #184: state the member trust model at startup. Environment controls and worktrees do
|
||||
* not make a sandbox when fleetd and its members use the same OS user. A separate herdr may
|
||||
* provide that boundary, but fleetd cannot inspect the uid at the other end of its socket.
|
||||
*/
|
||||
static void reportMemberTrustModel(FleetConfig cfg) {
|
||||
if (cfg.memberHerdrSocket() != null && !cfg.memberHerdrSocket().isBlank()) {
|
||||
log.info("member trust model: members are routed to a separate herdr through "
|
||||
+ "memberHerdrSocket. fleetd cannot see that herdr's uid, so confirm it runs "
|
||||
+ "as a different OS user before treating it as a boundary.");
|
||||
return;
|
||||
}
|
||||
log.info("member trust model: members run as the same OS user as fleetd, not in a sandbox. "
|
||||
+ "A member can read any file this user can read, including SSH keys and credential "
|
||||
+ "stores, whatever memberCredentials says. To add a real boundary, route members to "
|
||||
+ "a second herdr under a different OS user with memberHerdrSocket.");
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-596: {@code known:} empty (block absent entirely, or present but empty) means {@link
|
||||
* FleetConfig.MemberCredentials#blockedSet()} is empty too — every member pane inherits the
|
||||
|
||||
@@ -1,6 +1,8 @@
|
||||
package dev.ltms.fleet.config;
|
||||
|
||||
import com.fasterxml.jackson.annotation.JsonCreator;
|
||||
import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
|
||||
import com.fasterxml.jackson.annotation.JsonProperty;
|
||||
import com.fasterxml.jackson.core.JsonParser;
|
||||
import com.fasterxml.jackson.core.JsonToken;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
@@ -320,8 +322,9 @@ public record FleetConfig(
|
||||
* statement that does not stop being true just because the profile was
|
||||
* named directly. A negative value has no sane meaning (there is no
|
||||
* "excluded" to degrade to below zero) and is refused at config load
|
||||
* instead, naming the profile and the key. Live means any session the
|
||||
* registry still owns (acquired and not yet released), in any state.
|
||||
* instead, naming the profile and the key. Live means a session that can
|
||||
* receive another delivery. The roster keeps terminal {@code BACKEND_ERROR}
|
||||
* and {@code FAILED} sessions for diagnostics, but they do not use capacity.
|
||||
* @param kind which peer launcher spawns this profile: {@code "claude-code"} (default —
|
||||
* the {@link dev.ltms.fleet.member.ClaudeCodeLauncher}) or {@code "opencode"}.
|
||||
* The {@code CompositePeerLauncher} routes {@code spawn}/reap by this value, so
|
||||
@@ -1357,19 +1360,19 @@ public record FleetConfig(
|
||||
* deny-list/deny-by-default every name here that is NOT also in {@link #allow} is
|
||||
* overlaid with a non-secret sentinel value. Under allow-list this list is
|
||||
* reporting only.
|
||||
* @param sshAuthSock whether the member may inherit {@code SSH_AUTH_SOCK} under the allow-list
|
||||
* policy ({@code "allow"}) or must have it blanked ({@code "block"}, the default).
|
||||
* @param sshAgentEnv whether the member may inherit {@code SSH_AUTH_SOCK} under the allow-list
|
||||
* policy ({@code "inherit"}) or must have it omitted ({@code "omit"}, the default).
|
||||
* This is a DECISION, never a default: {@code SSH_AUTH_SOCK} is a handle to the
|
||||
* operator's ssh-agent, and a member holding it can sign with the operator's own
|
||||
* keys — but it appears in no secret file and is credential-shaped like nothing on
|
||||
* any list, which is why three earlier tickets missed it (gitea #110 / CB-607).
|
||||
* Blocking it breaks git over SSH inside the member; allow it only when members do
|
||||
* Omitting it does not prevent git over SSH inside the member; inherit it only when members do
|
||||
* not need to authenticate as the operator over SSH. Ignored under deny-list /
|
||||
* deny-by-default, which never touch the name.
|
||||
*/
|
||||
@JsonIgnoreProperties(ignoreUnknown = true)
|
||||
public record MemberCredentials(String policy, List<String> allow, List<String> known,
|
||||
String sshAuthSock) {
|
||||
String sshAgentEnv) {
|
||||
|
||||
/** Default policy: block every {@code known} name not in {@code allow}, via the env overlay. */
|
||||
public static final String POLICY_DENY_BY_DEFAULT = "deny-by-default";
|
||||
@@ -1386,11 +1389,25 @@ public record FleetConfig(
|
||||
*/
|
||||
public static final String POLICY_ALLOW_LIST = "allow-list";
|
||||
|
||||
/** The pre-CB-633 three-field form — {@code sshAuthSock} defaults to blocked. */
|
||||
/** The pre-CB-633 three-field form — {@code sshAgentEnv} defaults to omitted. */
|
||||
public MemberCredentials(String policy, List<String> allow, List<String> known) {
|
||||
this(policy, allow, known, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* Reads both the current {@code sshAgentEnv} key and the compatible {@code sshAuthSock} key.
|
||||
* When both keys are present, {@code sshAgentEnv} wins, even if its value is unrecognised.
|
||||
*/
|
||||
@JsonCreator
|
||||
public static MemberCredentials fromYaml(@JsonProperty("policy") String policy,
|
||||
@JsonProperty("allow") List<String> allow,
|
||||
@JsonProperty("known") List<String> known,
|
||||
@JsonProperty("sshAgentEnv") String sshAgentEnv,
|
||||
@JsonProperty("sshAuthSock") String sshAuthSock) {
|
||||
return new MemberCredentials(policy, allow, known,
|
||||
sshAgentEnv != null ? sshAgentEnv : sshAuthSock);
|
||||
}
|
||||
|
||||
public MemberCredentials {
|
||||
String normalizedPolicy = (policy == null || policy.isBlank())
|
||||
? POLICY_DENY_BY_DEFAULT : policy.toLowerCase(java.util.Locale.ROOT);
|
||||
@@ -1399,8 +1416,9 @@ public record FleetConfig(
|
||||
policy = POLICY_DENY_LIST.equals(normalizedPolicy) ? POLICY_DENY_BY_DEFAULT : normalizedPolicy;
|
||||
allow = allow == null ? List.of() : List.copyOf(allow);
|
||||
known = known == null ? List.of() : List.copyOf(known);
|
||||
sshAuthSock = (sshAuthSock != null && "allow".equalsIgnoreCase(sshAuthSock.trim()))
|
||||
? "allow" : "block";
|
||||
sshAgentEnv = (sshAgentEnv != null
|
||||
&& ("inherit".equalsIgnoreCase(sshAgentEnv.trim())
|
||||
|| "allow".equalsIgnoreCase(sshAgentEnv.trim()))) ? "inherit" : "omit";
|
||||
}
|
||||
|
||||
/** True when this block selects the CB-633 derived-allow-list policy. */
|
||||
@@ -1409,8 +1427,8 @@ public record FleetConfig(
|
||||
}
|
||||
|
||||
/** True when {@code SSH_AUTH_SOCK} may pass through under the allow-list policy. Default: no. */
|
||||
public boolean sshAuthSockAllowed() {
|
||||
return "allow".equals(sshAuthSock);
|
||||
public boolean sshAgentEnvInherited() {
|
||||
return "inherit".equals(sshAgentEnv);
|
||||
}
|
||||
|
||||
/** {@link #allow} as a set, for membership checks. */
|
||||
@@ -1489,7 +1507,7 @@ public record FleetConfig(
|
||||
rejectDuplicateMemberSlots(yaml);
|
||||
rejectNegativeMaxLoad(yaml);
|
||||
rejectAutoCompactWindowOutOfRange(yaml);
|
||||
rejectMalformedErrorPattern(yaml);
|
||||
rejectMalformedProfilePatterns(yaml);
|
||||
rejectUnknownKind(yaml);
|
||||
rejectUnknownAuthMode(yaml);
|
||||
rejectUnknownPlacement(yaml);
|
||||
@@ -1839,20 +1857,26 @@ public record FleetConfig(
|
||||
}
|
||||
|
||||
/**
|
||||
* Reject a profile whose {@code errorPattern} (fleetd #201 Unit 5) is not a valid Java regex,
|
||||
* naming the profile, the key, and the parser's own message.
|
||||
* Reject a profile whose {@code errorPattern} (fleetd #201 Unit 5) or {@code exhaustedPattern}
|
||||
* (CB-578 stage A) is not a valid Java regex, naming the profile, the key, and the parser's own
|
||||
* message.
|
||||
*
|
||||
* <p>Unset/{@code null} means "use {@code CompletionResolver}'s built-in {@code (?i)\bAPI
|
||||
* Error\s*:} compatibility pattern" and passes silently. A profile that DOES set the key gets it
|
||||
* compiled once at daemon startup ({@code Fleetd.main}, mirroring {@code exhaustedPattern}) — an
|
||||
* uncaught {@link java.util.regex.PatternSyntaxException} there crashes startup without naming
|
||||
* which profile or key is at fault. Validate eagerly here instead, at config load, the same
|
||||
* "fail loud at load, not lazily later" reasoning as {@link #rejectAutoCompactWindowOutOfRange}.
|
||||
* <p>Unset/{@code null} means, for {@code errorPattern}, "use {@code CompletionResolver}'s
|
||||
* built-in {@code (?i)\bAPI Error\s*:} compatibility pattern", and for {@code exhaustedPattern},
|
||||
* "opt out of that classification" — either way it passes silently. A profile that DOES set
|
||||
* either key gets it compiled once at daemon startup ({@code Fleetd.main}) — an uncaught
|
||||
* {@link java.util.regex.PatternSyntaxException} there crashes startup without naming which
|
||||
* profile or key is at fault (fleetd #273: this happened for {@code exhaustedPattern}, which had
|
||||
* no validator here even though its sibling {@code errorPattern} did). Validate eagerly here
|
||||
* instead, at config load, the same "fail loud at load, not lazily later" reasoning as
|
||||
* {@link #rejectAutoCompactWindowOutOfRange}. Both keys are checked from a single load, and any
|
||||
* failures from either are collected together into one message.
|
||||
*
|
||||
* @param yaml the raw config text
|
||||
* @throws IllegalStateException when any profile's {@code errorPattern} fails to compile
|
||||
* @throws IllegalStateException when any profile's {@code errorPattern} or
|
||||
* {@code exhaustedPattern} fails to compile
|
||||
*/
|
||||
static void rejectMalformedErrorPattern(String yaml) {
|
||||
static void rejectMalformedProfilePatterns(String yaml) {
|
||||
Map<?, ?> raw;
|
||||
try {
|
||||
raw = YAML.readValue(yaml, Map.class);
|
||||
@@ -1867,18 +1891,20 @@ public record FleetConfig(
|
||||
if (!(e.getValue() instanceof Map<?, ?> p)) {
|
||||
continue;
|
||||
}
|
||||
if (!(p.get("errorPattern") instanceof String pattern) || pattern.isBlank()) {
|
||||
continue;
|
||||
}
|
||||
try {
|
||||
Pattern.compile(pattern);
|
||||
} catch (PatternSyntaxException ex) {
|
||||
bad.add("profiles." + e.getKey() + ".errorPattern (\"" + pattern + "\"): " + ex.getMessage());
|
||||
for (String key : List.of("errorPattern", "exhaustedPattern")) {
|
||||
if (!(p.get(key) instanceof String pattern) || pattern.isBlank()) {
|
||||
continue;
|
||||
}
|
||||
try {
|
||||
Pattern.compile(pattern);
|
||||
} catch (PatternSyntaxException ex) {
|
||||
bad.add("profiles." + e.getKey() + "." + key + " (\"" + pattern + "\"): " + ex.getMessage());
|
||||
}
|
||||
}
|
||||
}
|
||||
bad.sort(String::compareTo);
|
||||
if (!bad.isEmpty()) {
|
||||
throw new IllegalStateException("refusing to start: malformed errorPattern — "
|
||||
throw new IllegalStateException("refusing to start: malformed pattern — "
|
||||
+ String.join("; ", bad));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -14,6 +14,7 @@ import java.util.HashSet;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.function.BiConsumer;
|
||||
@@ -28,6 +29,14 @@ public final class FleetHealthMonitor {
|
||||
static final int MAX_FAIL_TARGET_ATTEMPTS = 3;
|
||||
// CB-641: Match the injector's 60s readiness gate so health allows a full first boot.
|
||||
static final long READINESS_GRACE_NANOS = TimeUnit.SECONDS.toNanos(60);
|
||||
/**
|
||||
* fleetd #280: how long after a terminal transition to wait before the one bounded re-check
|
||||
* fires. Must exceed the worst-case reverse-rendezvous {@code fleet_ask} window (55-115s, see
|
||||
* {@code FleetMcp.ASK_DEFAULT_TIMEOUT_MS} / {@code FleetApp.MAX_ASK_TIMEOUT_MS}) so that, if the
|
||||
* target was genuinely {@code ASKING} when {@code state} was first observed, its own ask has had
|
||||
* time to lapse (clearing {@code Task#question} back to {@code null}) before this fires.
|
||||
*/
|
||||
static final long ASK_LAPSE_RECHECK_DELAY_SECONDS = 120;
|
||||
|
||||
private final AgentControl agents;
|
||||
private final Supplier<List<MemberSession>> roster;
|
||||
@@ -38,7 +47,16 @@ public final class FleetHealthMonitor {
|
||||
private final long workingSuspectAfterNanos;
|
||||
private final BiConsumer<String, String> failTarget;
|
||||
private final Map<String, HealthPrior> priors = new HashMap<>();
|
||||
private final Map<String, HealthState> states = new HashMap<>();
|
||||
/**
|
||||
* The live classification per member, and the only one of this class's three maps that more
|
||||
* than one scheduler task touches. {@code tick} writes it (and prunes it to the roster);
|
||||
* fleetd #280's delayed {@link #recheckTerminalTarget} reads it from its own separate scheduled
|
||||
* task. Both run on the single-threaded scheduler {@code Fleetd} passes in today, so they are
|
||||
* serialised — but nothing in this class enforces that, and an unsynchronised {@link HashMap}
|
||||
* read racing a resize can spin a CPU forever rather than fail visibly. {@code priors} and
|
||||
* {@code orphanStreaks} stay plain maps because {@code tick} is still their only toucher.
|
||||
*/
|
||||
private final Map<String, HealthState> states = new ConcurrentHashMap<>();
|
||||
/**
|
||||
* CB-643: consecutive ticks on which a target looked like an orphaned delegation. The fact
|
||||
* {@link MessageService#hasOrphanedDelegation} reports is a true snapshot, but it can read true
|
||||
@@ -171,6 +189,7 @@ public final class FleetHealthMonitor {
|
||||
// member stayed terminal.
|
||||
if (terminal(next)) {
|
||||
failTerminalTarget(target, next);
|
||||
scheduleTerminalRecheck(target, next);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -191,6 +210,48 @@ public final class FleetHealthMonitor {
|
||||
target, state, MAX_FAIL_TARGET_ATTEMPTS, last);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #280: schedule the one bounded, delayed follow-up for a terminal transition — never a
|
||||
* per-tick retry (CB-580 rejected that shape; {@link #reportTransition} still fires
|
||||
* {@link #failTerminalTarget} exactly once per transition, unconditionally on the tick loop).
|
||||
* This is a single one-shot task, scheduled once per transition into GONE/NEVER_READY, so a
|
||||
* member stuck terminal for the rest of its life gets exactly one extra attempt, not one per
|
||||
* tick. See {@link #recheckTerminalTarget} for why the extra attempt is safe.
|
||||
*/
|
||||
private void scheduleTerminalRecheck(String target, HealthState state) {
|
||||
if (scheduler.isShutdown()) return;
|
||||
try {
|
||||
scheduler.schedule(() -> recheckTerminalTarget(target, state),
|
||||
ASK_LAPSE_RECHECK_DELAY_SECONDS, TimeUnit.SECONDS);
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("fleet health: could not schedule terminal re-check for member={} state={}",
|
||||
target, state, e);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #280: the delayed re-check {@link #scheduleTerminalRecheck} scheduled for one terminal
|
||||
* transition. By now, a {@code fleet_ask} that was still open when {@code state} was first
|
||||
* observed has had time to lapse on its own (see {@link #ASK_LAPSE_RECHECK_DELAY_SECONDS}),
|
||||
* clearing {@code Task#question} back to {@code null} — which is exactly what
|
||||
* {@link MessageService#abandon(String, String, boolean)}'s {@code sweepAsking=false} filter
|
||||
* needs to finally match it. Calling {@link #failTerminalTarget} again is safe only because
|
||||
* {@code sweepAsking} stays {@code false}: a task genuinely still {@code ASKING} is skipped
|
||||
* exactly as it was on the very first attempt — this never fails a ticket whose ask has not yet
|
||||
* lapsed.
|
||||
*
|
||||
* <p><strong>Guarded on "target is still classified {@code state}."</strong> Without this guard,
|
||||
* a member that recovered (or was released and dropped from the roster) between the transition
|
||||
* and this re-check would still take a blind {@code failTarget} call — reaching into whatever
|
||||
* brand-new, unrelated turn it has since picked up and failing it too. {@link #states} already
|
||||
* carries the live classification (updated every tick, pruned to the current roster on release),
|
||||
* so a stale or recovered target simply reads as a mismatch here and this is a no-op.
|
||||
*/
|
||||
void recheckTerminalTarget(String target, HealthState state) {
|
||||
if (states.get(target) != state) return;
|
||||
failTerminalTarget(target, state);
|
||||
}
|
||||
|
||||
private static boolean terminal(HealthState state) {
|
||||
return state == HealthState.GONE || state == HealthState.NEVER_READY;
|
||||
}
|
||||
|
||||
@@ -139,16 +139,24 @@ public final class FleetMcp {
|
||||
|
||||
/**
|
||||
* fleetd #176: the seats a profile's own live LEAD session(s) hold on the same Claude
|
||||
* subscription — the third reason (alongside {@link QuarantineSource} and {@link OutageSource})
|
||||
* {@code free} can overstate what a fresh {@code fleet_spawn} would actually get.
|
||||
* subscription — a fact {@code fleet_list} reports alongside {@code free} via the
|
||||
* {@code leadSeats} key.
|
||||
*
|
||||
* <p>{@code maxLoad} counts only <em>members</em>, never the lead itself. But a
|
||||
* {@code subscription: true} profile bills the operator's own Claude account, and the lead is
|
||||
* always a live {@code claude} session on that same account (it is never moved off-subscription
|
||||
* — see {@code LeadLauncher}). So a fan-out that fills every member slot still leaves the lead's
|
||||
* own seat unaccounted for, and the daemon reports a slot that was never really free. See
|
||||
* {@code Fleetd.leadSeatLookup} for how the count is derived — from {@code fleet.leaders.<name>
|
||||
* .profile} and each profile's {@code effectiveCredentialId()}, never a hardcoded constant.
|
||||
* — see {@code LeadLauncher}). See {@code Fleetd.leadSeatLookup} for how the count is derived —
|
||||
* from {@code fleet.leaders.<name>.profile} and each profile's {@code effectiveCredentialId()},
|
||||
* never a hardcoded constant.
|
||||
*
|
||||
* <p>fleetd #257: this count is reported, never subtracted from {@code free}. An earlier cut of
|
||||
* this feature subtracted it, on the theory that it made {@code free} describe the real ceiling
|
||||
* on the account — but the real spawn gate ({@code CompositePeerLauncher#enforceMaxLoad}) never
|
||||
* read this count at all, so the subtraction made {@code free} disagree with the one thing it is
|
||||
* supposed to describe: what a fresh {@code fleet_spawn} will actually get. No backend seat
|
||||
* ceiling shared with the lead has been measured either — see {@code fleetd.example.yaml}'s
|
||||
* {@code maxLoad} docs. {@code free} now always equals {@code max(0, maxLoad - live)}, and
|
||||
* {@code leadSeats} is reported purely as a fact the caller may act on however it likes.
|
||||
*/
|
||||
public record LeadSeatSource(Function<String, Integer> seatsFor) {
|
||||
/** Inert source — no profile is ever reported as sharing a seat with a lead. */
|
||||
@@ -249,7 +257,7 @@ public final class FleetMcp {
|
||||
// Each handler is built once and wired to its fleet_* tool below.
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> sendHandler =
|
||||
(exchange, req) -> {
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.SEND,
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_send", req.arguments()),
|
||||
str(req.arguments(), "sessionId"));
|
||||
if (denied != null) return denied;
|
||||
String caller = callerTerminal(exchange);
|
||||
@@ -291,7 +299,7 @@ public final class FleetMcp {
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> replyHandler =
|
||||
(exchange, req) -> {
|
||||
String self = callerTerminal(exchange);
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.REPLY, self);
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_reply", req.arguments()), self);
|
||||
if (denied != null) return denied;
|
||||
return reply(messages, self, str(req.arguments(), "content"));
|
||||
};
|
||||
@@ -299,36 +307,38 @@ public final class FleetMcp {
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> askHandler =
|
||||
(exchange, req) -> {
|
||||
String self = callerTerminal(exchange);
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.ASK, self);
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_ask", req.arguments()), self);
|
||||
if (denied != null) return denied;
|
||||
return ask(messages, self, str(req.arguments(), "question"), timeoutMs(req.arguments()));
|
||||
};
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> statusHandler =
|
||||
(exchange, req) -> {
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_status", req.arguments()), null);
|
||||
if (denied != null) return denied;
|
||||
return status(messages, str(req.arguments(), "sessionId"));
|
||||
};
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> pollHandler =
|
||||
(exchange, req) -> {
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
|
||||
if (denied != null) return denied;
|
||||
Map<String, Object> a = req.arguments();
|
||||
return poll(messages, str(a, "ticket"), str(a, "target"));
|
||||
String target = str(a, "target");
|
||||
// The action depends on the ARGUMENTS, not on the tool name -- see pollAction.
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_poll", a), target);
|
||||
if (denied != null) return denied;
|
||||
return poll(messages, str(a, "ticket"), target);
|
||||
};
|
||||
// CB-307 Increment 3: per-msgId ack (not needed in v1 but supported by the inbox).
|
||||
// Acking removes a reply from the inbox, so it is a drain, not a read.
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> ackHandler =
|
||||
(exchange, req) -> {
|
||||
Map<String, Object> a = req.arguments();
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.DRAIN, str(a, "target"));
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_ack", a), str(a, "target"));
|
||||
if (denied != null) return denied;
|
||||
return ack(messages, str(a, "target"), str(a, "msgId"));
|
||||
};
|
||||
// Fleet management (CB-108): spawn/list/stop over ClaudeCodeLauncher.
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> spawnHandler =
|
||||
(exchange, req) -> {
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.SPAWN, null);
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_spawn", req.arguments()), null);
|
||||
if (denied != null) return denied;
|
||||
String caller = callerTerminal(exchange);
|
||||
// SPAWN is already auth-gated to PRIMARY (architects can never call it), but
|
||||
@@ -345,7 +355,7 @@ public final class FleetMcp {
|
||||
};
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> listHandler =
|
||||
(exchange, _) -> {
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_list", Map.of()), null);
|
||||
if (denied != null) return denied;
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
|
||||
leadSeats, callers == null ? Map.of() : callers.leads(),
|
||||
@@ -355,19 +365,19 @@ public final class FleetMcp {
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> stopHandler =
|
||||
(exchange, req) -> {
|
||||
String paneId = str(req.arguments(), "paneId");
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.STOP, paneId);
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_stop", req.arguments()), paneId);
|
||||
if (denied != null) return denied;
|
||||
return stop(sessions, paneId);
|
||||
};
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> profilesHandler =
|
||||
(exchange, _) -> {
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_profiles", Map.of()), null);
|
||||
if (denied != null) return denied;
|
||||
return profiles(workers, quarantine, outage);
|
||||
};
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> whoamiHandler =
|
||||
(exchange, _) -> {
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_whoami", Map.of()), null);
|
||||
if (denied != null) return denied;
|
||||
return whoami(principal(exchange), sessions);
|
||||
};
|
||||
@@ -716,6 +726,53 @@ public final class FleetMcp {
|
||||
return text("delivered to peer lead " + coordId + " (msgId " + msg.msgId() + ")");
|
||||
}
|
||||
|
||||
/**
|
||||
* Which authorization action a {@code fleet_poll} call needs, decided by its arguments
|
||||
* (fleetd #272).
|
||||
*
|
||||
* <p>{@code fleet_poll} is <strong>two operations behind one tool name</strong>. With {@code
|
||||
* ticket} it observes an async delegation and changes nothing, which is a {@link
|
||||
* Authz.Action#READ}. With {@code target} it calls {@link MessageService#drainReplies} on that
|
||||
* session -- the replies are removed from the inbox and a second call returns nothing -- so it
|
||||
* is a {@link Authz.Action#DRAIN}, the same gate {@code fleet_ack} already uses for removing a
|
||||
* single message, and the same one the REST path uses at {@code FleetApp.drainReplies}.
|
||||
*
|
||||
* <p>Until this method existed the handler passed a constant {@code READ} for both branches.
|
||||
* {@code READ} is open to every authenticated role, so any worker could read a peer's id out of
|
||||
* {@code fleet_list} and destroy the replies that peer had queued for the primary. The gate
|
||||
* failed open, and it did so because the required action is a function of the arguments while
|
||||
* the handler chose it before looking at them.
|
||||
*
|
||||
* <p>The choice lives in this method, and not inline in the handler, so that a test can assert
|
||||
* the mapping the handler actually uses. {@code FleetMcpAuthzTest} already checked every
|
||||
* {@link Authz.Action} against every {@link Role} and passed throughout -- it tested the policy
|
||||
* table, which was correct, while the defect was in which action the caller handed it.
|
||||
*
|
||||
* @param target the {@code target} argument of the call, or {@code null}/blank when absent
|
||||
*/
|
||||
static Authz.Action pollAction(String target) {
|
||||
return isBlank(target) ? Authz.Action.READ : Authz.Action.DRAIN;
|
||||
}
|
||||
|
||||
/**
|
||||
* The action a registered tool handler actually hands to the authorization gate.
|
||||
* Keeping this choice beside the registered-tool inventory makes a new tool fail the coverage
|
||||
* test until its action is pinned.
|
||||
*/
|
||||
static Authz.Action toolAction(String toolName, Map<String, Object> arguments) {
|
||||
return switch (toolName) {
|
||||
case "fleet_send" -> Authz.Action.SEND;
|
||||
case "fleet_reply" -> Authz.Action.REPLY;
|
||||
case "fleet_ask" -> Authz.Action.ASK;
|
||||
case "fleet_status", "fleet_list", "fleet_profiles", "fleet_whoami" -> Authz.Action.READ;
|
||||
case "fleet_poll" -> pollAction(str(arguments, "target"));
|
||||
case "fleet_ack" -> Authz.Action.DRAIN;
|
||||
case "fleet_spawn" -> Authz.Action.SPAWN;
|
||||
case "fleet_stop" -> Authz.Action.STOP;
|
||||
default -> throw new IllegalArgumentException("unregistered tool: " + toolName);
|
||||
};
|
||||
}
|
||||
|
||||
/** {@code fleet_poll}: check an async delegation by ticket, or drain a worker's inbox by target. */
|
||||
static McpSchema.CallToolResult poll(MessageService messages, String ticket, String target) {
|
||||
if (!isBlank(target)) {
|
||||
@@ -1110,15 +1167,35 @@ public final class FleetMcp {
|
||||
private static Map<String, Object> memberCapacityView(MemberSession session, Agent live,
|
||||
MessageService messages, long nowNanos) {
|
||||
Map<String, Object> row = SessionManager.rosterView(session, live);
|
||||
boolean open = messages != null && messages.hasAcceptedDelivery(session.terminalId());
|
||||
boolean inbox = messages != null && messages.hasInboxMessage(session.terminalId());
|
||||
boolean reclaimable = (session.state() == MemberSession.State.READY || session.state() == MemberSession.State.DONE)
|
||||
&& !open && !inbox;
|
||||
boolean reclaimable = reclaimable(session, messages);
|
||||
row.put("reclaimable", reclaimable);
|
||||
row.put("idleForSeconds", reclaimable ? Math.max(0, (nowNanos - session.lastActivityAtNanos()) / 1_000_000_000L) : null);
|
||||
return row;
|
||||
}
|
||||
|
||||
/**
|
||||
* The one definition of {@code reclaimable}: this member holds a spawn seat, and has no open
|
||||
* bridge work, so stopping it gives the seat back. Both views in a single {@code fleet_list}
|
||||
* response call it — the per-member flag in {@link #memberCapacityView} and the per-profile
|
||||
* count in {@link #capacityView} — because two copies of this rule in one response is how the
|
||||
* two numbers come to disagree.
|
||||
*
|
||||
* <p>fleetd #284: {@code BACKEND_ERROR} and {@code FAILED} are deliberately NOT reclaimable.
|
||||
* The ticket asked for them to be, and that half of the ticket was wrong. Once
|
||||
* {@code Fleetd.liveSessionCount} stopped counting a terminal session as live, that seat is
|
||||
* ALREADY in {@code free}; counting it here too reports the same seat twice, and
|
||||
* {@code free + reclaimable} then reads as more capacity than {@code maxLoad} allows. The dead
|
||||
* session stays visible either way: its roster row still carries {@code state:
|
||||
* "backend_error"} or {@code "failed"}, which is what tells the lead to stop it.
|
||||
*/
|
||||
static boolean reclaimable(MemberSession session, MessageService messages) {
|
||||
boolean holdsSeat = session.state() == MemberSession.State.READY
|
||||
|| session.state() == MemberSession.State.DONE;
|
||||
return holdsSeat && (messages == null
|
||||
|| (!messages.hasAcceptedDelivery(session.terminalId())
|
||||
&& !messages.hasInboxMessage(session.terminalId())));
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-583: {@code free} alone cannot tell a lead "busy, will free up" from "refusing, and
|
||||
* nothing changes for N seconds" — those need different decisions. So a quarantined profile
|
||||
@@ -1137,10 +1214,17 @@ public final class FleetMcp {
|
||||
* <p>fleetd #176: {@code maxLoad} counts panes, not subscription seats — it never counted the
|
||||
* lead's own seat on a {@code subscription: true} profile's account. {@link LeadSeatSource}
|
||||
* reports that count (0 for a non-subscription profile, or when no live lead shares its
|
||||
* credential), and it is subtracted from {@code free} the same way {@code live} already is —
|
||||
* {@code maxLoad} itself is left untouched, so the row still reports the configured cap. The
|
||||
* {@code leadSeats} key is added only when the count is positive, for the same
|
||||
* byte-identical-when-unused reason as the quarantine/cool-off keys above.
|
||||
* credential) via the {@code leadSeats} key, added only when the count is positive, for the
|
||||
* same byte-identical-when-unused reason as the quarantine/cool-off keys above.
|
||||
*
|
||||
* <p>fleetd #257: {@code leadSeatCount} is reported, never subtracted from {@code free}.
|
||||
* {@code free} means "what a fresh {@code fleet_spawn} on this profile will actually get", and
|
||||
* the real gate ({@code CompositePeerLauncher#enforceMaxLoad}) only ever compares {@code live}
|
||||
* against {@code maxLoad} — it has no notion of a lead's own seat. Subtracting
|
||||
* {@code leadSeatCount} here made {@code free} disagree with the gate it is supposed to
|
||||
* describe: it could report {@code free: 0} while a spawn on that exact profile still
|
||||
* succeeded. {@code leadSeats} stays in the row as a fact the caller can act on however it
|
||||
* likes, but it no longer changes what {@code free} means.
|
||||
*/
|
||||
private static Map<String, Object> capacityView(String profile, Function<String, Integer> liveCount,
|
||||
Function<String, Integer> maxLoad, List<MemberSession> roster,
|
||||
@@ -1150,12 +1234,16 @@ public final class FleetMcp {
|
||||
int live = liveCount.apply(profile);
|
||||
int leadSeatCount = leadSeats.seatsFor().apply(profile);
|
||||
int reclaimable = (int) roster.stream().filter(s -> profile.equals(s.profile()))
|
||||
.filter(s -> (s.state() == MemberSession.State.READY || s.state() == MemberSession.State.DONE))
|
||||
.filter(s -> messages == null || (!messages.hasAcceptedDelivery(s.terminalId()) && !messages.hasInboxMessage(s.terminalId())))
|
||||
.filter(s -> reclaimable(s, messages))
|
||||
.count();
|
||||
Map<String, Object> row = new LinkedHashMap<>();
|
||||
row.put("profile", profile); row.put("maxLoad", cap); row.put("live", live);
|
||||
row.put("free", cap == null ? null : Math.max(0, cap - live - leadSeatCount));
|
||||
// fleetd #257: free must report what the real spawn gate (CompositePeerLauncher#enforceMaxLoad)
|
||||
// will actually grant, and that gate never reads leadSeatCount — only maxLoad and live. Not
|
||||
// subtracting the lead's seat here used to make free UNDERSTATE what a fresh fleet_spawn would
|
||||
// get, so a lead believing free:0 gave up on a profile the gate would still spawn onto.
|
||||
// leadSeatCount is still reported via the leadSeats key below, just never subtracted from free.
|
||||
row.put("free", cap == null ? null : Math.max(0, cap - live));
|
||||
row.put("reclaimable", reclaimable);
|
||||
if (leadSeatCount > 0) {
|
||||
row.put("leadSeats", leadSeatCount);
|
||||
@@ -1360,8 +1448,13 @@ public final class FleetMcp {
|
||||
+ "member spawned without one shares its directory with others and never reports "
|
||||
+ "an id, however long it runs (fleetd #249). An empty 'members' "
|
||||
+ "means no members are spawned; it says nothing about peers. When capacity "
|
||||
+ "facts are configured, a 'capacity' row per profile also reports free: 0 for "
|
||||
+ "a quarantined profile's credential (see fleet_profiles), whatever its "
|
||||
+ "facts are configured, a 'capacity' row per profile reports 'free' — the "
|
||||
+ "slots a fresh fleet_spawn on that profile will actually be granted right "
|
||||
+ "now (max(0, maxLoad - live)), the same check the spawn gate itself runs. A "
|
||||
+ "'leadSeats' key, when present, reports how many of those live slots are a "
|
||||
+ "lead session sharing this profile's subscription — informational only, "
|
||||
+ "already NOT subtracted from 'free' (fleetd #257). It also reports free: 0 "
|
||||
+ "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'.",
|
||||
|
||||
@@ -19,6 +19,7 @@ import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.nio.file.StandardCopyOption;
|
||||
import java.nio.file.attribute.PosixFileAttributeView;
|
||||
import java.util.Arrays;
|
||||
import java.util.EnumSet;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
@@ -519,64 +520,212 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
|
||||
* the first. Both exist because of a real incident: see {@link HerdrPeerLauncher#isProvisionedWorktree}'s javadoc
|
||||
* and {@link #writeAtomically}'s javadoc.
|
||||
*
|
||||
* <p><b>Compare-and-swap against a writer the lock cannot reach (fleetd #247).</b>
|
||||
* {@code TRUST_JSON_LOCK} only serialises calls this launcher itself makes inside this one JVM.
|
||||
* It does nothing about a writer outside it — and on a host where the profile's
|
||||
* {@code configDir} is the operator's own {@code CLAUDE_CONFIG_DIR}, the file this method writes
|
||||
* IS the operator's own live Claude Code session's config file, being read and written by that
|
||||
* session while it runs. Measured 2026-09-03: its mtime moved minutes after a spawn while that
|
||||
* session was active. A plain read-modify-write there is a routine lost update, not a rare one:
|
||||
* fleetd reads v1, the operator's session reads v1 and writes v2 with their own change, fleetd's
|
||||
* {@code ATOMIC_MOVE} then lands v3 built from v1 — atomic, but v2's change is gone. So before
|
||||
* the move this method re-reads {@code target}'s exact bytes and compares them with the bytes it
|
||||
* built its update from; a mismatch means someone else wrote in between, and it discards its
|
||||
* work and rebuilds from the fresh bytes, up to {@link #MAX_TRUST_JSON_CAS_ATTEMPTS} times.
|
||||
* <b>Exhausting the retries writes nothing</b> — see the WARN at the end of the loop for why
|
||||
* that, not a last write-anyway, is the safe failure: the member shows the trust dialog and
|
||||
* fails to reach an injectable state, which is visible, logged and recoverable; overwriting the
|
||||
* operator's live config with a stale copy is neither. This narrows the lost-update window, it
|
||||
* does not close it — a write landing between the final re-read and the {@code ATOMIC_MOVE}
|
||||
* itself is still lost, because there is no OS-level compare-and-swap on a plain file, only this
|
||||
* cooperative narrowing of the gap.
|
||||
*
|
||||
* <p><b>fleetd #285: refuses under {@code memberHerdrSocket} rather than writing somewhere the
|
||||
* member cannot read.</b> Under {@code memberHerdrSocket:} the member pane runs as a
|
||||
* <em>different OS user with its own {@code $HOME}</em> — the same reason {@link
|
||||
* #writeCharterFile} routes the role/reply charter under {@code worktreeRoot} instead of
|
||||
* {@code java.io.tmpdir} and refuses the spawn when it cannot. This method has no equivalent
|
||||
* relocation available: unlike the charter (fleetd's own content, free to place anywhere and
|
||||
* hand to the peer via an argv flag), {@code .claude.json} is a file Claude Code looks up for
|
||||
* ITSELF at a fixed location — {@code CLAUDE_CONFIG_DIR/.claude.json}, or else the member OS
|
||||
* user's own {@code ~/.claude.json}, a path fleetd has no channel to learn. So when {@code
|
||||
* configDir} is unset, there is no member-readable target to seed at all — writing the
|
||||
* unqualified default would land in <em>fleetd's own</em> {@code ~/.claude.json} instead, the
|
||||
* exact defect this fix closes, not a workable fallback. And even with {@code configDir} set,
|
||||
* the file this method itself just wrote is {@code 0600} (owner-only — see {@link
|
||||
* #copyPosixPermissionsIfPresent}), unreadable by a different-uid member unless shared with
|
||||
* {@code worktreeGroup}, the same group {@link EnvAllowListScrub#shareWithGroup} already uses
|
||||
* for the ZDOTDIR scrub (fleetd #213) and the charter file (fleetd #219/#222). So under {@code
|
||||
* memberHerdrSocket} this method requires BOTH {@code configDir} and {@code worktreeGroup}
|
||||
* before it ever touches a file, and refuses the spawn — naming exactly which one is missing —
|
||||
* rather than silently corrupt fleetd's own home or hand the member an unreadable path. This
|
||||
* mirrors {@link #writeCharterFile}'s "refuse, don't degrade" decision: a member spawned without
|
||||
* a readable trust seed is not degraded, it sits on the interactive dialog forever and never
|
||||
* calls {@code fleet_reply} — exactly the failure fleetd #149 exists to prevent, so trading it
|
||||
* for "spawn something" is not worth it. When {@code configDir} and {@code worktreeGroup} are
|
||||
* both present, the write proceeds exactly as below and the resulting file is additionally
|
||||
* chgrp'd/chmod'd group-readable ({@code rw-r-----}) via {@link
|
||||
* EnvAllowListScrub#shareFileWithGroup(Path, String)} so the member's OS user can actually
|
||||
* open it — the
|
||||
* directory itself (unlike {@code worktreeRoot} or the charter's per-spawn directory) is not
|
||||
* fleetd-managed, so its own traversal permissions remain the operator's setup, same as they
|
||||
* already must be for the member to read anything else fleetd points {@code CLAUDE_CONFIG_DIR}
|
||||
* at. With {@code memberHerdrSocket} ABSENT (today's only live mode) every branch below is
|
||||
* byte-identical to before this fix.
|
||||
*
|
||||
* @param configDir the profile's {@code CLAUDE_CONFIG_DIR} ({@code cfg.configDir()}), or
|
||||
* {@code null}/blank to target the default {@code ~/.claude.json}
|
||||
* {@code null}/blank to target the default {@code ~/.claude.json} — refused
|
||||
* outright when {@code memberHerdrSocket} is configured, see above
|
||||
* @param cwd the spawn's resolved working directory — the exact key Claude Code will look
|
||||
* up for itself once it starts there
|
||||
* @throws IllegalStateException when {@code memberHerdrSocket} is configured but {@code
|
||||
* configDir} and/or {@code worktreeGroup} is not — the same
|
||||
* refusal shape as {@link #writeCharterFile}
|
||||
*/
|
||||
private static void seedTrustDialog(String configDir, String cwd) {
|
||||
private void seedTrustDialog(String configDir, String cwd) {
|
||||
if (!isProvisionedWorktree(cwd)) {
|
||||
return;
|
||||
}
|
||||
Path target = (configDir == null || configDir.isBlank())
|
||||
boolean unsetConfigDir = configDir == null || configDir.isBlank();
|
||||
boolean memberHerdrSocket = memberHerdrSocketConfigured();
|
||||
String group = memberHerdrSocket ? memberGroup() : null;
|
||||
if (memberHerdrSocket && (unsetConfigDir || group == null)) {
|
||||
throw new IllegalStateException("memberHerdrSocket is configured, so the workspace-trust "
|
||||
+ "seed (.claude.json, which gates Claude Code's interactive trust dialog) must be "
|
||||
+ "placed where the member's OS user can read it — configDir, shared via "
|
||||
+ "worktreeGroup — but " + (unsetConfigDir ? "configDir" : "worktreeGroup")
|
||||
+ " is not configured. Refusing to spawn rather than write fleetd's own default "
|
||||
+ "'~/.claude.json' or hand the member a config file it cannot read: that member "
|
||||
+ "would sit on the interactive trust dialog forever and never reach an "
|
||||
+ "injectable state. Configure configDir on this profile and worktreeGroup on the "
|
||||
+ "fleet to enable claude-code member spawns under memberHerdrSocket.");
|
||||
}
|
||||
Path target = unsetConfigDir
|
||||
? Path.of(System.getProperty("user.home"), ".claude.json")
|
||||
: Path.of(configDir, ".claude.json");
|
||||
if (unsetConfigDir) {
|
||||
// fleetd #247: configDir unset is the ONLY path that targets ~/.claude.json — the
|
||||
// operator's own home file, not a per-profile one — and it is the default, so a
|
||||
// profile that simply forgot to set configDir gets no signal at all short of the
|
||||
// operator noticing their own file changing. Say so loudly, every time it is about to
|
||||
// happen, rather than only once ever: each occurrence is a live write to a real
|
||||
// person's home config and deserves its own log line. (Reached only when
|
||||
// memberHerdrSocket is absent — the block above already refused otherwise.)
|
||||
log.warn("seedTrustDialog: profile has no configDir set, so the workspace-trust seed "
|
||||
+ "for cwd '{}' is about to write the operator's own default '{}' — set "
|
||||
+ "configDir on this profile to target a per-member config file instead",
|
||||
cwd, target);
|
||||
}
|
||||
synchronized (TRUST_JSON_LOCK) {
|
||||
boolean written = false;
|
||||
try {
|
||||
if (target.getParent() != null) {
|
||||
Files.createDirectories(target.getParent());
|
||||
}
|
||||
ObjectNode root = null;
|
||||
if (Files.isRegularFile(target)) {
|
||||
JsonNode existing = TRUST_JSON.readTree(target.toFile());
|
||||
if (existing instanceof ObjectNode existingObject) {
|
||||
root = existingObject;
|
||||
for (int attempt = 1; attempt <= MAX_TRUST_JSON_CAS_ATTEMPTS && !written; attempt++) {
|
||||
byte[] before = Files.isRegularFile(target) ? Files.readAllBytes(target) : null;
|
||||
ObjectNode root = parseTrustJsonOrEmpty(before);
|
||||
JsonNode projectsNode = root.get("projects");
|
||||
ObjectNode projects = projectsNode instanceof ObjectNode projectsObject
|
||||
? projectsObject : TRUST_JSON.createObjectNode();
|
||||
if (!(projectsNode instanceof ObjectNode)) {
|
||||
root.set("projects", projects);
|
||||
}
|
||||
JsonNode projectNode = projects.get(cwd);
|
||||
ObjectNode project = projectNode instanceof ObjectNode projectObject
|
||||
? projectObject : TRUST_JSON.createObjectNode();
|
||||
if (!(projectNode instanceof ObjectNode)) {
|
||||
projects.set(cwd, project);
|
||||
}
|
||||
// fleetd #247: ONLY hasTrustDialogAccepted. We used to write
|
||||
// hasCompletedProjectOnboarding beside it; do not put it back. Measured on
|
||||
// 2026-09-03, minutes after a live spawn seeded this file: 28 of 28 project
|
||||
// entries carried hasTrustDialogAccepted and 0 of 28 carried the onboarding key
|
||||
// — including the 27 entries Claude Code wrote for itself. Claude Code
|
||||
// normalises the whole file when it saves and drops that key every time, so
|
||||
// writing it achieved nothing except making the next reader think it mattered.
|
||||
// The member reached idle with the trust flag alone, which is the only outcome
|
||||
// this seed exists for. If a future Claude Code needs the second flag the
|
||||
// symptom returns as the trust dialog fleetd #149 describes — re-measure then,
|
||||
// do not restore it on a guess.
|
||||
project.put("hasTrustDialogAccepted", true);
|
||||
String newContent = TRUST_JSON.writerWithDefaultPrettyPrinter().writeValueAsString(root);
|
||||
|
||||
trustJsonCasTestHook.run();
|
||||
|
||||
// fleetd #247 CAS: re-read immediately before the move and compare with what
|
||||
// this attempt built its update from. A mismatch means another writer (most
|
||||
// plausibly the operator's own live Claude Code — see this method's javadoc)
|
||||
// landed a change in between; discard this attempt's work and rebuild from the
|
||||
// fresh bytes rather than blindly overwriting it.
|
||||
byte[] atMove = Files.isRegularFile(target) ? Files.readAllBytes(target) : null;
|
||||
if (!Arrays.equals(before, atMove)) {
|
||||
continue;
|
||||
}
|
||||
writeAtomically(target, newContent);
|
||||
written = true;
|
||||
}
|
||||
if (root == null) {
|
||||
root = TRUST_JSON.createObjectNode();
|
||||
if (!written) {
|
||||
// fleetd #247: deliberately do NOT write here. A member that starts without the
|
||||
// seed still starts — it may hit the trust dialog fleetd #149 describes and fail
|
||||
// to reach an injectable state, but that failure is visible (herdr reports it,
|
||||
// the spawn-readiness gate times out) and recoverable (retry the spawn). Writing
|
||||
// our stale copy over whatever the other writer left would be silent and, if that
|
||||
// other writer is the operator's own live session, could destroy real
|
||||
// configuration — fail toward the recoverable outcome, not the silent one.
|
||||
log.warn("seedTrustDialog: gave up seeding workspace-trust for cwd '{}' into '{}' "
|
||||
+ "after {} attempts — another writer (most plausibly the operator's own "
|
||||
+ "live Claude Code sharing this file) kept changing it faster than we "
|
||||
+ "could re-read it, so nothing was written; the member may show the "
|
||||
+ "trust dialog instead", cwd, target, MAX_TRUST_JSON_CAS_ATTEMPTS);
|
||||
}
|
||||
JsonNode projectsNode = root.get("projects");
|
||||
ObjectNode projects = projectsNode instanceof ObjectNode projectsObject
|
||||
? projectsObject : TRUST_JSON.createObjectNode();
|
||||
if (!(projectsNode instanceof ObjectNode)) {
|
||||
root.set("projects", projects);
|
||||
}
|
||||
JsonNode projectNode = projects.get(cwd);
|
||||
ObjectNode project = projectNode instanceof ObjectNode projectObject
|
||||
? projectObject : TRUST_JSON.createObjectNode();
|
||||
if (!(projectNode instanceof ObjectNode)) {
|
||||
projects.set(cwd, project);
|
||||
}
|
||||
// fleetd #247: ONLY hasTrustDialogAccepted. We used to write
|
||||
// hasCompletedProjectOnboarding beside it; do not put it back. Measured on
|
||||
// 2026-09-03, minutes after a live spawn seeded this file: 28 of 28 project
|
||||
// entries carried hasTrustDialogAccepted and 0 of 28 carried the onboarding key
|
||||
// — including the 27 entries Claude Code wrote for itself. Claude Code
|
||||
// normalises the whole file when it saves and drops that key every time, so
|
||||
// writing it achieved nothing except making the next reader think it mattered.
|
||||
// The member reached idle with the trust flag alone, which is the only outcome
|
||||
// this seed exists for. If a future Claude Code needs the second flag the
|
||||
// symptom returns as the trust dialog fleetd #149 describes — re-measure then,
|
||||
// do not restore it on a guess.
|
||||
project.put("hasTrustDialogAccepted", true);
|
||||
writeAtomically(target, TRUST_JSON.writerWithDefaultPrettyPrinter().writeValueAsString(root));
|
||||
} catch (Exception e) {
|
||||
log.debug("cannot seed workspace-trust entry for cwd '{}' into '{}'", cwd, target, e);
|
||||
return;
|
||||
}
|
||||
// fleetd #285: the write above lands as fleetd's own OS user; under memberHerdrSocket
|
||||
// that is NOT the member's OS user, so without this the member still cannot read the
|
||||
// file it exists to seed — a silent readiness timeout with the write looking "done".
|
||||
// Deliberately OUTSIDE the swallow-all catch above: a group that fails to resolve here
|
||||
// means the seed is unreadable despite a successful write, which must fail as loudly as
|
||||
// writeCharterFile's own EnvAllowListScrub.shareWithGroup call already does.
|
||||
if (written && memberHerdrSocket) {
|
||||
EnvAllowListScrub.shareFileWithGroup(target, group);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
/** Bound on {@link #seedTrustDialog}'s fleetd #247 compare-and-swap retry loop. */
|
||||
private static final int MAX_TRUST_JSON_CAS_ATTEMPTS = 5;
|
||||
|
||||
/**
|
||||
* Test-only seam for fleetd #247: invoked once per CAS attempt inside {@link #seedTrustDialog}'s
|
||||
* retry loop, after that attempt has read the target's bytes and built its replacement content,
|
||||
* but immediately before the final re-read/compare that decides whether to write. A no-op in
|
||||
* production. Package-visible (not {@code private}) so {@code ClaudeCodeLauncherTest} can install
|
||||
* a hook here that writes to the target file, deterministically simulating a writer racing
|
||||
* fleetd's own read-modify-write at the exact instant the CAS is meant to catch — the same race
|
||||
* an external process (most plausibly the operator's own live Claude Code) creates, without
|
||||
* depending on real thread scheduling to land the interleaving. A test that sets this MUST
|
||||
* restore it to the no-op default in a {@code finally} block — it is shared, static state.
|
||||
*/
|
||||
static Runnable trustJsonCasTestHook = () -> {};
|
||||
|
||||
/**
|
||||
* Parse {@code bytes} as a {@code .claude.json} tree, or hand back a fresh empty object when
|
||||
* {@code bytes} is {@code null} (no file yet) or does not parse to a JSON object — the same
|
||||
* missing-or-unreadable-is-empty fallback {@link #seedTrustDialog} always used, factored out so
|
||||
* the fleetd #247 CAS loop can call it once per attempt.
|
||||
*/
|
||||
private static ObjectNode parseTrustJsonOrEmpty(byte[] bytes) throws IOException {
|
||||
if (bytes == null) {
|
||||
return TRUST_JSON.createObjectNode();
|
||||
}
|
||||
JsonNode existing = TRUST_JSON.readTree(bytes);
|
||||
return existing instanceof ObjectNode existingObject ? existingObject : TRUST_JSON.createObjectNode();
|
||||
}
|
||||
|
||||
/**
|
||||
* Write {@code content} to {@code target} atomically: serialise to a sibling temp file in the
|
||||
* <strong>same directory</strong> as {@code target} (an atomic move is only guaranteed within
|
||||
|
||||
@@ -181,6 +181,33 @@ public final class EnvAllowListScrub {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #285: share ONE file with {@code group}, read-only ({@code rw-r-----}) — the same
|
||||
* per-file mode {@link #shareWithGroup} applies to a directory's entries, and the same error
|
||||
* shapes, but without touching a parent directory. Used for a file fleetd writes into a
|
||||
* directory it does NOT own — {@code configDir}'s own traversal permissions stay the
|
||||
* operator's setup — where the directory-wide {@link #shareWithGroup} would be wrong.
|
||||
*
|
||||
* @throws UncheckedIOException when {@code group} does not resolve on this host, the
|
||||
* filesystem has no POSIX group ownership, or a
|
||||
* group-ownership/permission call is refused
|
||||
*/
|
||||
static void shareFileWithGroup(Path file, String group) {
|
||||
try {
|
||||
GroupPrincipal principal = file.getFileSystem().getUserPrincipalLookupService()
|
||||
.lookupPrincipalByGroupName(group);
|
||||
setGroupAndPermissions(file, principal, "rw-r-----");
|
||||
} catch (IOException e) {
|
||||
throw new UncheckedIOException("cannot share generated file " + file + " with group '"
|
||||
+ group + "' — the group must exist, and the fleetd operator ("
|
||||
+ System.getProperty("user.name") + ") must be a member of it", e);
|
||||
} catch (UnsupportedOperationException e) {
|
||||
throw new UncheckedIOException("cannot share generated file " + file + " with group '"
|
||||
+ group + "' — this filesystem does not support POSIX group ownership",
|
||||
new IOException(e));
|
||||
}
|
||||
}
|
||||
|
||||
private static void setGroupAndPermissions(Path path, GroupPrincipal group, String perms) throws IOException {
|
||||
PosixFileAttributeView view = Files.getFileAttributeView(path, PosixFileAttributeView.class);
|
||||
if (view == null) {
|
||||
|
||||
@@ -117,9 +117,11 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
*
|
||||
* <p>fleetd #185 stage 2: that mirroring assumption holds only while the member pane runs under
|
||||
* the SAME OS user as the daemon. When {@code memberHerdrSocket:} is configured, member panes
|
||||
* run on a second herdr owned by a different user — different {@code $HOME}, different {@code
|
||||
* secrets.sh}, different environment entirely — so this field's data no longer describes what a
|
||||
* member pane inherits. See {@link #logCredentialGap} for how that mode is handled.
|
||||
* are routed to a second herdr, and fleetd has no channel to confirm what OS user that herdr
|
||||
* runs as — it may be a different user with a different {@code $HOME} and {@code secrets.sh},
|
||||
* or the same one the daemon runs as. Either way this field's data can no longer be trusted to
|
||||
* describe what a member pane inherits. See {@link #logCredentialGap} for how that mode is
|
||||
* handled.
|
||||
*/
|
||||
private final Supplier<Set<String>> hostEnvNames;
|
||||
|
||||
@@ -911,9 +913,14 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
* {@link #reapOrphanWorkers() orphan-reap} and spawn-gate-timeout paths, plus any caller that
|
||||
* passes a pane directly, keep working without an owning id.
|
||||
*
|
||||
* <p>Resolves the tab from the pane <em>before</em> closing it. An already-gone pane/tab
|
||||
* (repeated DELETE, crashed peer) is treated as success; any other failure propagates so a
|
||||
* genuinely failed teardown is not reported as done.
|
||||
* <p>Resolves the tab from the pane <em>before</em> closing it. {@code agents.close} (the pane)
|
||||
* is the one step whose failure means the teardown itself may not have happened: an already-gone
|
||||
* pane (repeated DELETE, crashed peer) is treated as success, but any other failure propagates so
|
||||
* a genuinely failed teardown is not reported as done. {@code spaces.closeTab} (fleetd #293) is
|
||||
* different — by the time it runs the pane is already closed, so it is cosmetic workspace tidying
|
||||
* rather than a real teardown failure, and a failure there is logged and never propagates, so it
|
||||
* cannot mask the two cleanups below it ({@link #releaseZdotdir}, and the caller's worktree
|
||||
* removal in {@code SessionManager.release}).
|
||||
*/
|
||||
@Override
|
||||
public void stop(String idOrPane) {
|
||||
@@ -932,7 +939,25 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
log.debug("pane.close({}) ignored — already gone: {}", paneId, e.getMessage());
|
||||
}
|
||||
if (loc != null && loc.tabPaneCount() == 1) {
|
||||
spaces.closeTab(loc.tabId());
|
||||
// fleetd #293: the pane above is already closed by this point, so a failing tab.close is
|
||||
// cosmetic workspace tidying, not a real teardown failure — it must not mask the two
|
||||
// cleanups below it (releaseZdotdir, and the caller's worktree removal). Unlike
|
||||
// agents.close above, this is not narrowed to "already gone": any failure here, whatever
|
||||
// its cause, is one we continue past, so we log it at WARN (not debug) with the tab id a
|
||||
// person can go close by hand.
|
||||
try {
|
||||
spaces.closeTab(loc.tabId());
|
||||
} catch (RuntimeException e) {
|
||||
// Caught as RuntimeException, not HerdrException, to match releaseZdotdir's own
|
||||
// guard five lines below. Today the two are the same set — HerdrCodec wraps every
|
||||
// encode/decode failure and UnixSocketHerdrClient wraps every IOException, so
|
||||
// HerdrException is all closeTab can actually throw. Narrowing to it anyway would
|
||||
// leave this step guarded against the expected failure and bare against any other,
|
||||
// which is the exact asymmetry fleetd #293 exists to remove. No behaviour change
|
||||
// today; it stops a later change inside WorkspaceControl.closeTab reopening it.
|
||||
log.warn("tab.close({}) failed — the pane is already torn down, so continuing; the "
|
||||
+ "tab may need manual cleanup: {}", loc.tabId(), e.getMessage());
|
||||
}
|
||||
} else if (loc != null) {
|
||||
log.debug("not closing tab {} — it holds {} panes (not a dedicated peer tab)",
|
||||
loc.tabId(), loc.tabPaneCount());
|
||||
@@ -1471,9 +1496,9 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
Set<String> brokerUriEnvNames = brokerUriEnvNames();
|
||||
Set<String> allowed = new java.util.TreeSet<>(
|
||||
MemberEnvAllowList.derive(profiles.values(), creds.allowSet(), brokerUriEnvNames));
|
||||
if (creds.sshAuthSockAllowed()) {
|
||||
if (creds.sshAgentEnvInherited()) {
|
||||
allowed.add(SSH_AUTH_SOCK);
|
||||
} // blocked by default: absent from the set ⇒ blanked by the scrub like any other name
|
||||
} // omitted by default: absent from the set ⇒ blanked by the scrub like any other name
|
||||
allowed.addAll(launch.env().keySet());
|
||||
allowed.removeAll(brokerUriEnvNames);
|
||||
return allowed;
|
||||
@@ -1495,10 +1520,27 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
* survive {@code allowed} (including the {@code LC_*} prefix rule). Neither number is a constant:
|
||||
* both come from the actual derived set and the actual environment this spawn sees. Never logs a
|
||||
* variable NAME or VALUE — only the counts.
|
||||
*
|
||||
* <p><strong>Under {@code memberHerdrSocket} the counts describe fleetd's own process, not the
|
||||
* member's</strong> (fleetd #269 follow-up), so the message says so rather than leaving the
|
||||
* reader to infer it from this javadoc, which the operator reading the log never sees.
|
||||
*/
|
||||
private void logAllowListCoverage(Set<String> allowed) {
|
||||
Set<String> hostNames = hostEnvNames.get();
|
||||
long kept = hostNames.stream().filter(name -> MemberEnvAllowList.keeps(allowed, name)).count();
|
||||
if (memberHerdrSocketConfigured()) {
|
||||
// fleetd #269 covered the sibling line below (logCredentialGap) and stopped there.
|
||||
// This line has the same problem: read plainly, "allowed 7 of 39" is a statement about
|
||||
// the member's pane, and under memberHerdrSocket it is not -- the pane is routed to a
|
||||
// second herdr whose environment fleetd cannot inspect. The counts stay useful, so
|
||||
// this is not a WARN and not a refusal; only the claim is narrowed to what is true.
|
||||
log.info("member credentials: allowed {} of {} names in fleetd's OWN environment — "
|
||||
+ "memberHerdrSocket is configured, so member panes are routed to a "
|
||||
+ "second herdr whose environment fleetd has no channel to inspect. "
|
||||
+ "These counts describe fleetd's process, NOT the member pane's.",
|
||||
kept, hostNames.size());
|
||||
return;
|
||||
}
|
||||
log.info("member credentials: allowed {} of {}", kept, hostNames.size());
|
||||
}
|
||||
|
||||
@@ -1533,21 +1575,24 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
/**
|
||||
* fleetd #213: {@code memberHerdrSocket} is configured and the member login shell IS zsh, but
|
||||
* {@code worktreeRoot} and/or {@code worktreeGroup} is missing, so the generated ZDOTDIR cannot
|
||||
* be placed anywhere the member's OS user can reach — {@code java.io.tmpdir} is fleetd's own
|
||||
* 0700 temp dir, unreadable by another uid, which is the exact gap this ticket exists to close.
|
||||
* Say so once per launcher instance, instead of either generating a directory nothing can read
|
||||
* (protection theatre) or refusing to spawn (turning a degraded credential control into an
|
||||
* outage for an opt-in feature).
|
||||
* be placed anywhere fleetd can be sure the member's OS user can reach — {@code java.io.tmpdir}
|
||||
* is fleetd's own 0700 temp dir, which is unreadable if the member pane runs as a different OS
|
||||
* user, and fleetd has no channel to confirm whether it does or not. Rather than gamble on that,
|
||||
* this treats memberHerdrSocket as reason enough to require an explicitly shared location, which
|
||||
* is the exact gap this ticket exists to close. Say so once per launcher instance, instead of
|
||||
* either generating a directory that might not be readable (protection theatre) or refusing to
|
||||
* spawn (turning a degraded credential control into an outage for an opt-in feature).
|
||||
*/
|
||||
private void warnCannotShareScrubDirectory() {
|
||||
if (cannotShareScrubDirWarned.compareAndSet(false, true)) {
|
||||
log.warn("memberCredentials policy=allow-list: memberHerdrSocket is configured and the "
|
||||
+ "member login shell is zsh, but worktreeRoot and/or worktreeGroup is not "
|
||||
+ "configured — the generated ZDOTDIR cannot be placed where the member's OS "
|
||||
+ "user can read it (java.io.tmpdir is fleetd's own, unreadable by another uid), "
|
||||
+ "so the scrub cannot be guaranteed to run. Falling back to the CB-596 sentinel "
|
||||
+ "overlay. Configure both worktreeRoot and worktreeGroup to enable the "
|
||||
+ "allow-list scrub under memberHerdrSocket.");
|
||||
+ "configured — the generated ZDOTDIR cannot be placed where fleetd can be sure "
|
||||
+ "the member's OS user can read it (java.io.tmpdir is fleetd's own 0700 dir, "
|
||||
+ "unreadable if the member runs as a different OS user — fleetd has no channel "
|
||||
+ "to confirm whether it does), so the scrub cannot be guaranteed to run. Falling "
|
||||
+ "back to the CB-596 sentinel overlay. Configure both worktreeRoot and "
|
||||
+ "worktreeGroup to enable the allow-list scrub under memberHerdrSocket.");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1626,10 +1671,12 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
private final AtomicBoolean unknownMemberEnvironmentWarned = new AtomicBoolean();
|
||||
|
||||
/**
|
||||
* fleetd #185 stage 2: whether {@code memberHerdrSocket:} is configured, i.e. member panes run
|
||||
* on a second herdr owned by a different OS user than the daemon's own process. Re-read from the
|
||||
* live config on every call (same hot-reload shape as {@link #memberCredentials}), never cached,
|
||||
* so a config reload takes effect on the next spawn without a restart.
|
||||
* fleetd #185 stage 2: whether {@code memberHerdrSocket:} is configured, i.e. member panes are
|
||||
* routed to a second herdr. This tests only that the config key is set — fleetd has no channel
|
||||
* to confirm what OS user that second herdr runs as, so a {@code true} result means "member
|
||||
* panes may run under a different OS user," not that they do. Re-read from the live config on
|
||||
* every call (same hot-reload shape as {@link #memberCredentials}), never cached, so a config
|
||||
* reload takes effect on the next spawn without a restart.
|
||||
*
|
||||
* <p>{@link #config} is {@code null} on any call site that never threaded the full config
|
||||
* through (every production {@code HerdrPeerLauncher} does; a handful of older tests do not) —
|
||||
@@ -1652,14 +1699,15 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
* fleetd #185 stage 2: the single replacement WARN for {@link #logCredentialGap}'s usual
|
||||
* conclusions when {@code memberHerdrSocket:} is configured. {@link #hostEnvNames} (and
|
||||
* everything derived from it — {@code known}/{@code allow} coverage, the allow-list scrub's
|
||||
* derived set) describes the DAEMON's own environment; under this config key member panes run as
|
||||
* a different OS user with a different environment entirely, so neither "every member pane
|
||||
* inherits them UNBLOCKED" nor "the scrub blanks them" is evidence-backed here — both would be
|
||||
* reporting on the wrong process. Logged once, names the config key, and states the honest
|
||||
* conclusion: the gap for member panes is UNKNOWN, not clean, so {@code memberCredentials} cannot
|
||||
* be verified from this daemon. The one count it does report is scoped explicitly to fleetd's own
|
||||
* environment, never presented as if it said anything about the member's — see {@link
|
||||
* #logCredentialGap}'s javadoc for why this branch exists.
|
||||
* derived set) describes the DAEMON's own environment; under this config key member panes are
|
||||
* routed to a second herdr, and fleetd has no channel to confirm what OS user that herdr runs
|
||||
* as or to read its environment, so neither "every member pane inherits them UNBLOCKED" nor
|
||||
* "the scrub blanks them" is evidence-backed here — both would be reporting on the wrong
|
||||
* process. Logged once, names the config key, and states the honest conclusion: the gap for
|
||||
* member panes is UNKNOWN, not clean, so {@code memberCredentials} cannot be verified from this
|
||||
* daemon. The one count it does report is scoped explicitly to fleetd's own environment, never
|
||||
* presented as if it said anything about the member's — see {@link #logCredentialGap}'s javadoc
|
||||
* for why this branch exists.
|
||||
*/
|
||||
private void warnUnknownMemberEnvironment(FleetConfig.MemberCredentials creds) {
|
||||
if (!unknownMemberEnvironmentWarned.compareAndSet(false, true)) {
|
||||
@@ -1672,13 +1720,13 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
.filter(name -> CREDENTIAL_SHAPED_NAME.matcher(name).matches())
|
||||
.filter(name -> !covered.contains(name))
|
||||
.count();
|
||||
log.warn("memberCredentials gap: memberHerdrSocket is configured, so member panes run under "
|
||||
+ "a different OS user than fleetd's own process, with a different environment "
|
||||
+ "entirely — fleetd has no channel to read that user's environment. {} of the "
|
||||
+ "{} names in fleetd's OWN environment are credential-shaped and not on "
|
||||
+ "known:/allow:, but that count describes fleetd's process, not the member "
|
||||
+ "herdr's. The credential gap for member panes is UNKNOWN, not clean, and "
|
||||
+ "memberCredentials cannot be verified from here.",
|
||||
log.warn("memberCredentials gap: memberHerdrSocket is configured, so member panes are routed "
|
||||
+ "to a second herdr — fleetd has no channel to confirm what OS user that herdr "
|
||||
+ "runs as, so it cannot tell whether those panes inherit its own environment or "
|
||||
+ "a different one entirely. {} of the {} names in fleetd's OWN environment are "
|
||||
+ "credential-shaped and not on known:/allow:, but that count describes fleetd's "
|
||||
+ "process, not the member herdr's. The credential gap for member panes is "
|
||||
+ "UNKNOWN, not clean, and memberCredentials cannot be verified from here.",
|
||||
gapInFleetdsOwnEnv, hostNames.size());
|
||||
}
|
||||
|
||||
|
||||
@@ -34,7 +34,7 @@ import java.util.TreeSet;
|
||||
*
|
||||
* <p>{@code SSH_AUTH_SOCK} is deliberately NOT here. It is a handle to the operator's ssh-agent — a
|
||||
* member holding it can sign with the operator's keys — so keeping it is a config decision
|
||||
* ({@code memberCredentials.sshAuthSock: allow}), not a derivation default.
|
||||
* ({@code memberCredentials.sshAgentEnv: inherit}), not a derivation default.
|
||||
*
|
||||
* <p><b>CB-633 follow-up:</b> the union also includes {@code memberCredentials.allow:} — the
|
||||
* operator's own explicit list. Before this, {@code policy: allow-list} silently ignored every name
|
||||
@@ -42,7 +42,7 @@ import java.util.TreeSet;
|
||||
* turning the policy on could blank credentials working members already depended on. {@code
|
||||
* SSH_AUTH_SOCK} and configured broker URI environment names are exceptions: even when the operator
|
||||
* lists them under {@code allow:}, they are excluded here. {@code SSH_AUTH_SOCK} is added back ONLY
|
||||
* by the caller when {@code sshAuthSock: allow} is explicitly set
|
||||
* by the caller when {@code sshAgentEnv: inherit} is explicitly set
|
||||
* (see {@link #SSH_AUTH_SOCK}'s javadoc) — it is a live handle to the operator's own ssh-agent, not
|
||||
* a value, so treating it like any other allow-listed name would hand a member every key the
|
||||
* operator's agent holds the moment they typed the name under {@code allow:} for an unrelated
|
||||
@@ -53,7 +53,7 @@ public final class MemberEnvAllowList {
|
||||
/**
|
||||
* The operator's ssh-agent socket path. Deliberately excluded from {@link #derive}'s union of
|
||||
* {@code memberCredentials.allow:} — see the class javadoc's CB-633 follow-up note. Governed
|
||||
* ONLY by {@code memberCredentials.sshAuthSock}, never by appearing in {@code allow:}.
|
||||
* ONLY by {@code memberCredentials.sshAgentEnv}, never by appearing in {@code allow:}.
|
||||
*/
|
||||
public static final String SSH_AUTH_SOCK = "SSH_AUTH_SOCK";
|
||||
|
||||
|
||||
@@ -21,6 +21,7 @@ import java.util.EnumSet;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
import java.util.function.BooleanSupplier;
|
||||
@@ -669,6 +670,16 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
private final AtomicBoolean discoveryUnavailableWarned =
|
||||
new AtomicBoolean();
|
||||
|
||||
/**
|
||||
* fleetd #267: one WARN per PROFILE (not per launcher instance — several profiles can each hit
|
||||
* this gap independently) for the model-mismatch check (fleetd #175) never getting to run
|
||||
* because the spawn was not given a fleetd-provisioned worktree (fleetd #249). Profile names
|
||||
* accumulate here for the life of this launcher instance and are never removed — the same
|
||||
* one-shot treatment {@link #discoveryUnavailableWarned} already gets, just keyed per profile
|
||||
* instead of globally.
|
||||
*/
|
||||
private final Set<String> modelCheckSkippedWarned = ConcurrentHashMap.newKeySet();
|
||||
|
||||
/** Add lazy on-disk session discovery to the base handle. */
|
||||
@Override
|
||||
public PeerHandle spawn(SpawnRequest req) {
|
||||
@@ -695,7 +706,8 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
// be running.
|
||||
FleetConfig.Profile cfg = requireProfile(req.profileName());
|
||||
return new SessionAwareHandle(inner, discovery, cwd, cfg,
|
||||
this::memberHerdrSocketConfigured, discoveryUnavailableWarned, exhaustionSink);
|
||||
this::memberHerdrSocketConfigured, discoveryUnavailableWarned,
|
||||
modelCheckSkippedWarned, exhaustionSink);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -720,6 +732,7 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
private final FleetConfig.Profile cfg;
|
||||
private final BooleanSupplier discoveryUnavailable;
|
||||
private final AtomicBoolean discoveryUnavailableWarned;
|
||||
private final Set<String> modelCheckSkippedWarned;
|
||||
private final ExhaustionSink exhaustionSink;
|
||||
/** CAS'd true the first (and only) time a model mismatch is reported for this handle. */
|
||||
private final AtomicBoolean modelMismatchReported = new AtomicBoolean();
|
||||
@@ -750,6 +763,7 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
FleetConfig.Profile cfg,
|
||||
BooleanSupplier discoveryUnavailable,
|
||||
AtomicBoolean discoveryUnavailableWarned,
|
||||
Set<String> modelCheckSkippedWarned,
|
||||
ExhaustionSink exhaustionSink) {
|
||||
this.delegate = delegate;
|
||||
this.discovery = discovery;
|
||||
@@ -757,6 +771,7 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
this.cfg = cfg;
|
||||
this.discoveryUnavailable = discoveryUnavailable;
|
||||
this.discoveryUnavailableWarned = discoveryUnavailableWarned;
|
||||
this.modelCheckSkippedWarned = modelCheckSkippedWarned;
|
||||
this.exhaustionSink = exhaustionSink;
|
||||
this.worktreeProvisioned = isProvisionedWorktree(cwd);
|
||||
}
|
||||
@@ -807,9 +822,29 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
// member's row apart from a sibling's in that case (measured: a three-day-old row from
|
||||
// a different profile). Refuse to guess — absent is the honest answer, and it is what
|
||||
// this codebase already returns elsewhere for absent evidence (fleetd #175's UNKNOWN).
|
||||
// No WARN here: unlike discoveryUnavailable above, this is the ordinary, expected shape
|
||||
// of the large majority of spawns (no worktree requested), not a configuration gap.
|
||||
// This IS the ordinary, expected shape of the large majority of spawns (no worktree
|
||||
// requested), not a configuration gap — but fleetd #267 found that same shape silently
|
||||
// switches off the fleetd #175 model-mismatch check for those spawns too, since
|
||||
// checkModelMatch's only call site is right below this gate. The check cannot be moved
|
||||
// off agentSessionId()'s resolved id: the id is the only safe way to key
|
||||
// actualModelForSessionId to THIS session's own row rather than "whatever is newest in
|
||||
// the shared directory" (fleetd #234) — re-deriving a second, independent answer via
|
||||
// `directory` here would reintroduce exactly the false-positive risk #234 fixed (a
|
||||
// sibling's differently-configured model looking like THIS profile's mismatch). So the
|
||||
// model genuinely is unknowable without a provisioned worktree, and unlike the silence
|
||||
// this branch used to keep, that gap now gets the same one-time, per-profile WARN
|
||||
// treatment discoveryUnavailable already gets above — but keyed by profile, since
|
||||
// several profiles can each hit this independently.
|
||||
if (!worktreeProvisioned) {
|
||||
if (cfg.model() != null && !cfg.model().isBlank()
|
||||
&& modelCheckSkippedWarned.add(cfg.profile())) {
|
||||
log.warn("opencode model-mismatch check (fleetd #175) cannot run for profile "
|
||||
+ "'{}': it was spawned without a fleetd-provisioned worktree (fleetd "
|
||||
+ "#249), so its cwd may be shared with other sessions and the actual "
|
||||
+ "model it is running cannot be safely told apart from a sibling's — "
|
||||
+ "spawn with worktree:true to enable the check for this profile.",
|
||||
cfg.profile());
|
||||
}
|
||||
return null;
|
||||
}
|
||||
// fleetd #234: once resolved, stay resolved. Re-deriving from `directory` on every call
|
||||
|
||||
@@ -360,6 +360,18 @@ public final class MessageService {
|
||||
* {@link #ask} clears the ticket's question and returns it to {@code PENDING}, but {@link #send}
|
||||
* already closed the forward waiter the instant the question surfaced, so the target has
|
||||
* neither an accepted nor a queued delivery left to show for it.
|
||||
*
|
||||
* <p><strong>Deliberately still {@code question == null} only (fleetd #275).</strong> This
|
||||
* method must not also report a still-{@link Phase#ASKING} task as orphaned: the worker may
|
||||
* genuinely be waiting on a live primary that is about to (or already mid-{@link #answer})
|
||||
* answer it, and {@link dev.ltms.fleet.health.FleetHealthMonitor} would classify that as
|
||||
* {@code DELEGATION_ORPHANED} on nothing more than an active, healthy conversation. {@link
|
||||
* #abandon(String, String, boolean)}'s {@code sweepAsking} path fixes the actual reachable gap
|
||||
* (a target torn down for good while genuinely {@code ASKING}) at the point of teardown itself,
|
||||
* by completing the task's future right there — so by the time this method would ever see it,
|
||||
* {@code task.future.isDone()} is already {@code true} and it is excluded regardless of this
|
||||
* guard. Widening this check instead of that one would trade a real fix for false positives on
|
||||
* every ordinary in-flight question.
|
||||
*/
|
||||
public boolean hasOrphanedDelegation(String target) {
|
||||
if (target == null || hasAcceptedDelivery(target) || hasQueuedDelivery(target)) {
|
||||
@@ -580,6 +592,39 @@ public final class MessageService {
|
||||
* reply — see the note above)
|
||||
*/
|
||||
public boolean abandon(String target, String reason) {
|
||||
return abandon(target, reason, false);
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #abandon(String, String)}, with control over whether a task still paused in
|
||||
* {@code fleet_ask} ({@link Phase#ASKING}) is swept too (fleetd #275).
|
||||
*
|
||||
* <p>{@code sweepAsking} must be {@code true} only when the caller has independent, certain
|
||||
* knowledge that {@code target} can never resume its turn — today that is only
|
||||
* {@code sessions.onRelease}'s teardown (an explicit {@code fleet_stop}, or the idle reaper):
|
||||
* the worker's pane is being stopped right now, so whatever it was mid-{@code fleet_ask} about
|
||||
* has no turn left to resume into. {@link dev.ltms.fleet.health.FleetHealthMonitor}'s
|
||||
* health-classification call keeps passing {@code false} (via {@link #abandon(String, String)}):
|
||||
* a GONE/NEVER_READY reading is the daemon's best guess from the live agent list, not a teardown
|
||||
* it performed itself, and {@code abandonDoesNotFailAnAsyncTicketWaitingForAnAnswer} documents
|
||||
* why an active ask must survive that guess — the primary may already be mid-{@link #answer} for
|
||||
* the very same turn, and completing it here first would preempt a real answer with a misleading
|
||||
* failure.
|
||||
*
|
||||
* <p><strong>Without {@code sweepAsking} on the release path, a target torn down while
|
||||
* genuinely {@code ASKING} was unrecoverable.</strong> {@link #resolveQuestion} had already
|
||||
* closed the forward waiter the instant the question surfaced (so the {@code waiter} branch
|
||||
* below finds nothing to fail), the {@code question == null} guard excluded the task from
|
||||
* {@code matching} (so the loop below skipped it too), and the worker's own {@code fleet_ask}
|
||||
* clears {@link Task#question} back to {@code null} only once it lapses (the reverse-rendezvous
|
||||
* window — up to {@code FleetMcp.ASK_DEFAULT_TIMEOUT_MS} / {@code FleetApp.MAX_ASK_TIMEOUT_MS},
|
||||
* 55–115s) — by which point the released session no longer appears in {@code sessions.roster()}
|
||||
* for {@link dev.ltms.fleet.health.FleetHealthMonitor} to ever re-observe, so nothing was ever
|
||||
* left to call {@link #abandon} on this target again. The ticket then sat in {@link #tasks}
|
||||
* forever: not terminal, so {@link #pruneTerminalTickets} never dropped it, and
|
||||
* {@code fleet_poll} reported it stuck at {@link Phase#PENDING} for good.
|
||||
*/
|
||||
public boolean abandon(String target, String reason, boolean sweepAsking) {
|
||||
boolean hadStrandedReply = hasStrandedReply(target);
|
||||
// CB-640: the session is gone — nothing will ever accept or deliver into it now.
|
||||
strandedReplies.remove(target);
|
||||
@@ -590,7 +635,8 @@ public final class MessageService {
|
||||
|
||||
List<Task> matching = new ArrayList<>();
|
||||
for (Task task : tasks.values()) {
|
||||
if (target.equals(task.target) && task.question == null && !task.future.isDone()) {
|
||||
if (target.equals(task.target) && (sweepAsking || task.question == null)
|
||||
&& !task.future.isDone()) {
|
||||
matching.add(task);
|
||||
}
|
||||
}
|
||||
@@ -607,11 +653,20 @@ public final class MessageService {
|
||||
for (Task task : matching) {
|
||||
boolean isRecovery = task == recoveryTask && recovered != null;
|
||||
Reply outcome = isRecovery ? recovered : new Reply(Outcome.WORKER_FAILED, reason);
|
||||
String turnId = task.turnId;
|
||||
if (task.future.complete(outcome)) {
|
||||
if (outcome.outcome() == Outcome.WORKER_FAILED) {
|
||||
asyncFailed = true;
|
||||
} else if (task.turnId != null) {
|
||||
asyncTasksByTurn.remove(task.turnId, task);
|
||||
}
|
||||
if (turnId != null) {
|
||||
// #275: whether this task was swept out of ASKING or was already answered and
|
||||
// only waiting on its resumed turn's real reply (#137), nothing will ever
|
||||
// complete this turnId now — drop it from this class's own bookkeeping AND the
|
||||
// reverse-rendezvous itself, so hasAsyncQuestion(target) stops reporting a turn
|
||||
// that is actually done, and a late answer() sees it as lapsed rather than
|
||||
// resolving a question nothing is listening for any more.
|
||||
asyncTasksByTurn.remove(turnId, task);
|
||||
rendezvous.closeAsk(turnId);
|
||||
}
|
||||
} else if (isRecovery) {
|
||||
// The recovered reply was already drained out of the inbox, but this task resolved
|
||||
@@ -872,7 +927,16 @@ public final class MessageService {
|
||||
}
|
||||
try {
|
||||
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(workerSession);
|
||||
// #282: mirror send()'s registration (:802) so a SECOND fleet_ask inside this same
|
||||
// resumed turn can re-associate the async ticket with its new turnId via
|
||||
// markAsyncQuestion — without this, that second ask has no Task to attach to, and
|
||||
// markAsyncQuestion silently returns null.
|
||||
Task task = asyncTasksByTurn.get(turnId);
|
||||
if (task != null) {
|
||||
asyncTasksByWaiter.put(reply, task);
|
||||
}
|
||||
if (!rendezvous.answerAsk(turnId, content)) {
|
||||
asyncTasksByWaiter.remove(reply);
|
||||
rendezvous.close(workerSession, reply);
|
||||
return new Reply(Outcome.STALE_TURN, null); // lapsed between the lookup and the unblock
|
||||
}
|
||||
@@ -880,7 +944,21 @@ public final class MessageService {
|
||||
try {
|
||||
Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS);
|
||||
Reply result = new Reply(outcomeOf(r.kind()), r.text(), r.turnId());
|
||||
finishAsyncTask(turnId, result);
|
||||
// #282: this waiter can resolve with a FRESH question rather than a terminal reply —
|
||||
// the worker chained a second fleet_ask before replying. Mirror sendAsync's own guard
|
||||
// (:1000) and leave the ticket open (markAsyncQuestion above already re-armed it under
|
||||
// the new turnId) instead of completing it here with a QUESTION "reply".
|
||||
// Measured when #282 was merged: this guard is DEFENCE IN DEPTH, not the thing
|
||||
// that makes the chained ask work. ask() calls markAsyncQuestion (:860) before
|
||||
// resolveQuestion (:861), so by the time this thread wakes, the task has already
|
||||
// moved to the new turnId and finishAsyncTask(oldTurnId, ...) finds nothing. Removing
|
||||
// this guard alone leaves the test green. Keep it anyway: it mirrors sendAsync's
|
||||
// sibling guard, and that sibling's own comment (:1017) warns the two orderings are
|
||||
// not something to rely on. Do NOT delete it as dead code without re-checking that
|
||||
// ordering, and do not treat it as the sole protection either.
|
||||
if (result.outcome() != Outcome.QUESTION) {
|
||||
finishAsyncTask(turnId, result);
|
||||
}
|
||||
return result;
|
||||
} catch (TimeoutException e) {
|
||||
// The worker resumed but hasn't replied yet — no completion fallback arms an answered
|
||||
@@ -893,6 +971,7 @@ public final class MessageService {
|
||||
Thread.currentThread().interrupt();
|
||||
throw new IllegalStateException("interrupted awaiting reply from " + workerSession, e);
|
||||
} finally {
|
||||
asyncTasksByWaiter.remove(reply);
|
||||
rendezvous.close(workerSession, reply);
|
||||
}
|
||||
} finally {
|
||||
@@ -1051,9 +1130,16 @@ public final class MessageService {
|
||||
private Task markAsyncQuestion(CompletableFuture<Rendezvous.Resolution> waiter, String text, String turnId) {
|
||||
Task task = waiter == null ? null : asyncTasksByWaiter.get(waiter);
|
||||
if (task != null) {
|
||||
String previousTurnId = task.turnId;
|
||||
task.question = new Reply(Outcome.QUESTION, text, turnId);
|
||||
task.turnId = turnId;
|
||||
asyncTasksByTurn.put(turnId, task);
|
||||
// #282: a second fleet_ask in the same resumed turn re-arms an already-answered task
|
||||
// (answer() re-registers it in asyncTasksByWaiter) under a FRESH turnId — drop the old
|
||||
// key so asyncTasksByTurn does not keep growing by one stale entry per chained ask.
|
||||
if (previousTurnId != null && !previousTurnId.equals(turnId)) {
|
||||
asyncTasksByTurn.remove(previousTurnId, task);
|
||||
}
|
||||
}
|
||||
return task;
|
||||
}
|
||||
|
||||
@@ -9,6 +9,7 @@ import dev.ltms.fleet.auth.Principal;
|
||||
import dev.ltms.fleet.guard.GuardException;
|
||||
import dev.ltms.fleet.metrics.FleetMetrics;
|
||||
import dev.ltms.fleet.metrics.Metrics;
|
||||
import dev.ltms.fleet.mcp.FleetMcp;
|
||||
import dev.ltms.fleet.herdr.Agent;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.herdr.HerdrException;
|
||||
@@ -47,6 +48,22 @@ import java.util.stream.Collectors;
|
||||
*/
|
||||
public final class FleetApp {
|
||||
|
||||
/** The authorization action the matching route handler hands to {@link #allow}. */
|
||||
static Authz.Action routeAction(String route) {
|
||||
return switch (route) {
|
||||
case "GET /metrics" -> Authz.Action.METRICS;
|
||||
case "POST /members" -> Authz.Action.SPAWN;
|
||||
case "DELETE /members/{paneId}" -> Authz.Action.STOP;
|
||||
case "POST /sessions/{id}/message" -> Authz.Action.SEND;
|
||||
case "POST /sessions/{id}/reply" -> Authz.Action.REPLY;
|
||||
case "GET /sessions/{id}/replies" -> Authz.Action.DRAIN;
|
||||
case "POST /sessions/{id}/ask" -> Authz.Action.ASK;
|
||||
case "GET /sessions", "GET /agents", "GET /members", "GET /profiles",
|
||||
"GET /member-credentials", "GET /sessions/{id}/status", "GET /tasks/{ticket}" -> Authz.Action.READ;
|
||||
default -> throw new IllegalArgumentException("route has no authorization gate: " + route);
|
||||
};
|
||||
}
|
||||
|
||||
/** Default blocking window for a message; kept under typical HTTP idle timeouts. */
|
||||
private static final long DEFAULT_MESSAGE_TIMEOUT_MS = 25_000;
|
||||
private static final long MAX_MESSAGE_TIMEOUT_MS = 120_000;
|
||||
@@ -70,6 +87,12 @@ public final class FleetApp {
|
||||
// absent() (the honest "no policy configured" view) for every constructor that does not wire
|
||||
// a real one, so existing legacy call sites keep building without knowing this field exists.
|
||||
private final Supplier<MemberCredentialPolicyView> memberCredentials;
|
||||
// fleetd #297: the SAME shared sources FleetMcp.profiles/fleet_profiles reads (BackendQuarantine
|
||||
// and BackendOutagePolicy are each one instance for the whole daemon — see Fleetd wiring) so
|
||||
// GET /profiles cannot drift from fleet_profiles about which profile is quarantined/cooling off.
|
||||
// .none() (the honest "feature not wired" view) for every constructor that does not pass one.
|
||||
private final FleetMcp.QuarantineSource quarantine;
|
||||
private final FleetMcp.OutageSource outage;
|
||||
private final ObjectMapper mapper = new ObjectMapper();
|
||||
|
||||
/**
|
||||
@@ -130,6 +153,23 @@ public final class FleetApp {
|
||||
MessageService messages, MemberPresence presence,
|
||||
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics,
|
||||
Predicate<String> deliverable, Supplier<MemberCredentialPolicyView> memberCredentials) {
|
||||
this(herdr, memberHerdr, workers, sessions, messages, presence, mcpServlet, auth, metrics,
|
||||
deliverable, memberCredentials, FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none());
|
||||
}
|
||||
|
||||
/**
|
||||
* @param quarantine the SAME {@link FleetMcp.QuarantineSource} instance passed to {@code
|
||||
* FleetMcp} (fleetd #297), so {@code GET /profiles} reports the identical
|
||||
* exhaustion-quarantine facts as {@code fleet_profiles} rather than a second,
|
||||
* independently-computed copy
|
||||
* @param outage the SAME {@link FleetMcp.OutageSource} instance passed to {@code FleetMcp} —
|
||||
* see {@code quarantine}; a SEPARATE check from it, never merged in
|
||||
*/
|
||||
public FleetApp(HerdrClient herdr, HerdrClient memberHerdr, PeerLauncher workers, SessionManager sessions,
|
||||
MessageService messages, MemberPresence presence,
|
||||
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics,
|
||||
Predicate<String> deliverable, Supplier<MemberCredentialPolicyView> memberCredentials,
|
||||
FleetMcp.QuarantineSource quarantine, FleetMcp.OutageSource outage) {
|
||||
this.herdr = herdr;
|
||||
this.memberHerdr = memberHerdr != null ? memberHerdr : herdr;
|
||||
this.workers = workers;
|
||||
@@ -140,6 +180,8 @@ public final class FleetApp {
|
||||
this.auth = auth;
|
||||
this.metrics = metrics;
|
||||
this.memberCredentials = memberCredentials != null ? memberCredentials : MemberCredentialPolicyView::absent;
|
||||
this.quarantine = quarantine != null ? quarantine : FleetMcp.QuarantineSource.none();
|
||||
this.outage = outage != null ? outage : FleetMcp.OutageSource.none();
|
||||
}
|
||||
|
||||
/** Wire routes onto a fresh, unstarted Javalin instance. Caller starts it. */
|
||||
@@ -222,7 +264,7 @@ public final class FleetApp {
|
||||
|
||||
/** Prometheus scrape endpoint (CB-502). */
|
||||
private void metrics(Context ctx) {
|
||||
if (!allow(ctx, Authz.Action.METRICS, null)) {
|
||||
if (!allow(ctx, routeAction("GET /metrics"), null)) {
|
||||
return;
|
||||
}
|
||||
ctx.status(200).contentType("text/plain; version=0.0.4; charset=utf-8").result(metrics.render());
|
||||
@@ -294,7 +336,7 @@ public final class FleetApp {
|
||||
* member workspace (they live on the member daemon only).
|
||||
*/
|
||||
private void sessions(Context ctx) {
|
||||
if (!allow(ctx, Authz.Action.READ, null)) {
|
||||
if (!allow(ctx, routeAction("GET /sessions"), null)) {
|
||||
return;
|
||||
}
|
||||
List<Map<String, Object>> out = new ArrayList<>();
|
||||
@@ -319,52 +361,104 @@ public final class FleetApp {
|
||||
|
||||
/** Discovery: every agent herdr tracks, keyed by its Claude session UUID. */
|
||||
private void agents(Context ctx) {
|
||||
if (!allow(ctx, Authz.Action.READ, null)) {
|
||||
if (!allow(ctx, routeAction("GET /agents"), null)) {
|
||||
return;
|
||||
}
|
||||
ctx.status(200).json(Map.of("agents",
|
||||
workers.list().stream().map(Agent.class::cast).map(FleetApp::view).toList()));
|
||||
try {
|
||||
ctx.status(200).json(Map.of("agents",
|
||||
workers.list().stream().map(Agent.class::cast).map(FleetApp::view).toList()));
|
||||
} catch (HerdrException e) {
|
||||
// fleetd #297: workers.list() reaches herdr — a transport failure must land in the same
|
||||
// {error, detail} envelope every other failure path here uses, not escape as a bare
|
||||
// exception and leave Javalin's default handling to respond outside the JSON contract.
|
||||
herdrError(ctx, e);
|
||||
}
|
||||
}
|
||||
|
||||
/** CB-304: bridge-owned roster merged with live herdr status by paneId. */
|
||||
private void listMembers(Context ctx) {
|
||||
if (!allow(ctx, Authz.Action.READ, null)) {
|
||||
if (!allow(ctx, routeAction("GET /members"), null)) {
|
||||
return;
|
||||
}
|
||||
// CB-519: the registry key is a host-unique id, not the pane coordinate — join on terminal.
|
||||
Map<String, Agent> live = workers.list().stream()
|
||||
.map(Agent.class::cast)
|
||||
.filter(a -> a.terminalId() != null)
|
||||
.collect(Collectors.toMap(Agent::terminalId, Function.identity(), (_, b) -> b));
|
||||
// fleetd #209: this REST roster reports agentSessionId via SessionManager.rosterView, so it
|
||||
// uses the resolving roster read (caller-driven, not a timer) rather than the plain one.
|
||||
List<Map<String, Object>> out = sessions.rosterResolved().stream()
|
||||
.map(s -> SessionManager.rosterView(s, live.get(s.terminalId())))
|
||||
.toList();
|
||||
Map<String, Object> body = new LinkedHashMap<>();
|
||||
// fleetd #199: the endpoint became /members in the CB-634 rename but the body key stayed
|
||||
// "workers", so a caller that read "members" saw an empty fleet and reported no members at
|
||||
// all. "members" is the canonical key; "workers" stays as a deprecated alias so an existing
|
||||
// REST consumer keeps working — the out-of-band path a lead falls back to when its MCP mount
|
||||
// drops reads this endpoint. Drop the alias once nothing reads it.
|
||||
body.put("members", out);
|
||||
body.put("workers", out);
|
||||
// CB-586: operator visibility for the refs/wip snapshot store without shelling into the
|
||||
// repo — how many snapshot refs exist and roughly what they cost. Present only once a
|
||||
// worktree session has established the repo, so a never-snapshotted fleet reports nothing.
|
||||
sessions.wipRefs().ifPresent(st -> body.put("wipRefs",
|
||||
Map.of("count", st.count(), "costBytes", st.costBytes())));
|
||||
ctx.status(200).json(body);
|
||||
try {
|
||||
// CB-519: the registry key is a host-unique id, not the pane coordinate — join on terminal.
|
||||
Map<String, Agent> live = workers.list().stream()
|
||||
.map(Agent.class::cast)
|
||||
.filter(a -> a.terminalId() != null)
|
||||
.collect(Collectors.toMap(Agent::terminalId, Function.identity(), (_, b) -> b));
|
||||
// fleetd #209: this REST roster reports agentSessionId via SessionManager.rosterView, so it
|
||||
// uses the resolving roster read (caller-driven, not a timer) rather than the plain one.
|
||||
List<Map<String, Object>> out = sessions.rosterResolved().stream()
|
||||
.map(s -> SessionManager.rosterView(s, live.get(s.terminalId())))
|
||||
.toList();
|
||||
Map<String, Object> body = new LinkedHashMap<>();
|
||||
// fleetd #199: the endpoint became /members in the CB-634 rename but the body key stayed
|
||||
// "workers", so a caller that read "members" saw an empty fleet and reported no members at
|
||||
// all. "members" is the canonical key; "workers" stays as a deprecated alias so an existing
|
||||
// REST consumer keeps working — the out-of-band path a lead falls back to when its MCP mount
|
||||
// drops reads this endpoint. Drop the alias once nothing reads it.
|
||||
body.put("members", out);
|
||||
body.put("workers", out);
|
||||
// CB-586: operator visibility for the refs/wip snapshot store without shelling into the
|
||||
// repo — how many snapshot refs exist and roughly what they cost. Present only once a
|
||||
// worktree session has established the repo, so a never-snapshotted fleet reports nothing.
|
||||
sessions.wipRefs().ifPresent(st -> body.put("wipRefs",
|
||||
Map.of("count", st.count(), "costBytes", st.costBytes())));
|
||||
ctx.status(200).json(body);
|
||||
} catch (HerdrException e) {
|
||||
// fleetd #297: same reasoning as agents() above — this is the out-of-band roster a lead
|
||||
// falls back to when its MCP mount drops, so it must stay inside the JSON error contract
|
||||
// exactly when herdr is briefly unreachable, not escape as a bare exception.
|
||||
herdrError(ctx, e);
|
||||
}
|
||||
}
|
||||
|
||||
/** The configured worker profiles and which one a no-argument spawn uses. */
|
||||
/**
|
||||
* The configured worker profiles, which one a no-argument spawn uses, and (fleetd #297) the two
|
||||
* outage states {@code fleet_profiles} already reports: {@code quarantined} (CB-578 stage B —
|
||||
* the backend reported it out of capacity) and {@code coolingOff} (fleetd #201 Unit 5 — the
|
||||
* credential threw repeated non-exhaustion backend errors). Both are read from the SAME shared
|
||||
* {@link FleetMcp.QuarantineSource}/{@link FleetMcp.OutageSource} instances {@code FleetMcp}
|
||||
* reads, never recomputed, so the two doors cannot disagree about which profile is down and why.
|
||||
* Independent checks, so a profile can appear in both maps at once; each map is present only
|
||||
* when at least one profile is in that state.
|
||||
*/
|
||||
private void profiles(Context ctx) {
|
||||
if (!allow(ctx, Authz.Action.READ, null)) {
|
||||
if (!allow(ctx, routeAction("GET /profiles"), null)) {
|
||||
return;
|
||||
}
|
||||
ctx.status(200).json(Map.of(
|
||||
"profiles", workers.profiles(),
|
||||
"default", workers.defaultProfile() == null ? "" : workers.defaultProfile()));
|
||||
Map<String, Object> body = new LinkedHashMap<>();
|
||||
body.put("profiles", workers.profiles());
|
||||
body.put("default", workers.defaultProfile() == null ? "" : workers.defaultProfile());
|
||||
Map<String, Object> quarantined = new LinkedHashMap<>();
|
||||
Map<String, Object> coolingOff = new LinkedHashMap<>();
|
||||
for (String profile : workers.profiles()) {
|
||||
String credentialId = quarantine.credentialIdFor().apply(profile);
|
||||
if (credentialId != null) {
|
||||
quarantine.quarantine().remainingSeconds(credentialId).ifPresent(remaining -> {
|
||||
Map<String, Object> row = new LinkedHashMap<>();
|
||||
row.put("credentialId", credentialId);
|
||||
row.put("quarantinedForSeconds", remaining);
|
||||
quarantined.put(profile, row);
|
||||
});
|
||||
}
|
||||
String outageCredentialId = outage.credentialIdFor().apply(profile);
|
||||
if (outageCredentialId != null) {
|
||||
outage.outagePolicy().remainingCoolOffSeconds(outageCredentialId).ifPresent(remaining -> {
|
||||
Map<String, Object> row = new LinkedHashMap<>();
|
||||
row.put("credentialId", outageCredentialId);
|
||||
row.put("coolingOffForSeconds", remaining);
|
||||
coolingOff.put(profile, row);
|
||||
});
|
||||
}
|
||||
}
|
||||
if (!quarantined.isEmpty()) {
|
||||
body.put("quarantined", quarantined);
|
||||
}
|
||||
if (!coolingOff.isEmpty()) {
|
||||
body.put("coolingOff", coolingOff);
|
||||
}
|
||||
ctx.status(200).json(body);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -376,7 +470,7 @@ public final class FleetApp {
|
||||
* name list, which is exactly what let the list drift silently behind the real policy.
|
||||
*/
|
||||
private void memberCredentials(Context ctx) {
|
||||
if (!allow(ctx, Authz.Action.READ, null)) {
|
||||
if (!allow(ctx, routeAction("GET /member-credentials"), null)) {
|
||||
return;
|
||||
}
|
||||
MemberCredentialPolicyView view = memberCredentials.get();
|
||||
@@ -396,7 +490,7 @@ public final class FleetApp {
|
||||
* the subscription boundary, 400 for an unknown profile.
|
||||
*/
|
||||
private void spawnMember(Context ctx) {
|
||||
if (!allow(ctx, Authz.Action.SPAWN, null)) {
|
||||
if (!allow(ctx, routeAction("POST /members"), null)) {
|
||||
return;
|
||||
}
|
||||
String role = ctx.queryParam("role");
|
||||
@@ -467,7 +561,7 @@ public final class FleetApp {
|
||||
/** Tear a worker down by pane id. */
|
||||
private void stopMember(Context ctx) {
|
||||
String paneId = ctx.pathParam("paneId");
|
||||
if (!allow(ctx, Authz.Action.STOP, paneId)) {
|
||||
if (!allow(ctx, routeAction("DELETE /members/{paneId}"), paneId)) {
|
||||
return;
|
||||
}
|
||||
sessions.release(paneId);
|
||||
@@ -482,7 +576,7 @@ public final class FleetApp {
|
||||
*/
|
||||
private void sendMessage(Context ctx) {
|
||||
String id = ctx.pathParam("id");
|
||||
if (!allow(ctx, Authz.Action.SEND, id)) {
|
||||
if (!allow(ctx, routeAction("POST /sessions/{id}/message"), id)) {
|
||||
return;
|
||||
}
|
||||
String content;
|
||||
@@ -569,7 +663,7 @@ public final class FleetApp {
|
||||
*/
|
||||
private void askMessage(Context ctx) {
|
||||
String id = ctx.pathParam("id");
|
||||
if (!allow(ctx, Authz.Action.ASK, id)) {
|
||||
if (!allow(ctx, routeAction("POST /sessions/{id}/ask"), id)) {
|
||||
return;
|
||||
}
|
||||
String question;
|
||||
@@ -608,7 +702,7 @@ public final class FleetApp {
|
||||
// The rule that matters: a worker may reply only as itself. Over MCP this was already true
|
||||
// structurally (identity comes from the connection, never an argument); over REST the path
|
||||
// id was simply trusted, so this is where the invariant actually gets enforced.
|
||||
if (!allow(ctx, Authz.Action.REPLY, id)) {
|
||||
if (!allow(ctx, routeAction("POST /sessions/{id}/reply"), id)) {
|
||||
return;
|
||||
}
|
||||
String content;
|
||||
@@ -629,7 +723,7 @@ public final class FleetApp {
|
||||
*/
|
||||
private void drainReplies(Context ctx) {
|
||||
String id = ctx.pathParam("id");
|
||||
if (!allow(ctx, Authz.Action.DRAIN, id)) {
|
||||
if (!allow(ctx, routeAction("GET /sessions/{id}/replies"), id)) {
|
||||
return;
|
||||
}
|
||||
var replies = messages.drainReplies(id);
|
||||
@@ -647,7 +741,7 @@ public final class FleetApp {
|
||||
*/
|
||||
private void sessionStatus(Context ctx) {
|
||||
String id = ctx.pathParam("id");
|
||||
if (!allow(ctx, Authz.Action.READ, id)) {
|
||||
if (!allow(ctx, routeAction("GET /sessions/{id}/status"), id)) {
|
||||
return;
|
||||
}
|
||||
try {
|
||||
@@ -672,7 +766,7 @@ public final class FleetApp {
|
||||
|
||||
/** Poll an async (wait:false) delegation by ticket. 404 for an unknown/expired ticket. */
|
||||
private void taskStatus(Context ctx) {
|
||||
if (!allow(ctx, Authz.Action.READ, null)) {
|
||||
if (!allow(ctx, routeAction("GET /tasks/{ticket}"), null)) {
|
||||
return;
|
||||
}
|
||||
MessageService.TaskView v = messages.poll(ctx.pathParam("ticket"));
|
||||
|
||||
@@ -167,14 +167,64 @@ public final class GitWorktrees implements Worktrees {
|
||||
log.info("adding worktree branch={} path={} base={}", branch, wt, base);
|
||||
removeUserInfoFromHttpsOrigin(repoRoot);
|
||||
exec("git", "-C", repoRoot, "worktree", "add", wt, "-b", branch, base);
|
||||
afterWorktreeAdded.accept(wt);
|
||||
requireCredentialFreeHttpsOrigin(wt);
|
||||
configureEnvironmentCredentialHelper(repoRoot, wt);
|
||||
configureHttpsUrlRewriteForSshOrigin(repoRoot, wt);
|
||||
isolateToolSurface(wt);
|
||||
try {
|
||||
afterWorktreeAdded.accept(wt);
|
||||
requireCredentialFreeHttpsOrigin(wt);
|
||||
configureEnvironmentCredentialHelper(repoRoot, wt);
|
||||
configureHttpsUrlRewriteForSshOrigin(repoRoot, wt);
|
||||
isolateToolSurface(wt);
|
||||
} catch (RuntimeException e) {
|
||||
cleanupAfterAddFailure(repoRoot, wt, branch, e);
|
||||
throw e;
|
||||
}
|
||||
return wt;
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code add()} has already created the worktree and its branch by the time any step from
|
||||
* {@link #afterWorktreeAdded} through {@link #isolateToolSurface} can throw — including
|
||||
* {@link #requireCredentialFreeHttpsOrigin}, an intended security refusal, not only an IO
|
||||
* accident. Without this, {@code add()} never returns, so its caller
|
||||
* ({@code SessionManager#acquireWithWorktree}) never receives a path to register or clean up:
|
||||
* its local {@code path} stays null, the {@code if (path != null)} guard in its own catch block
|
||||
* never runs, and the worktree directory and branch leak on disk forever with nothing tracking
|
||||
* them (fleetd #274).
|
||||
*
|
||||
* <p>Reuses {@link #remove} — the same {@code git worktree remove --force} path every other
|
||||
* cleanup exit in this class already goes through — rather than a bespoke removal. It
|
||||
* additionally deletes {@code branch}: {@link #remove} alone deliberately leaves a released
|
||||
* session's branch behind (a worker's branch is expected to outlive its worktree, for PRs and
|
||||
* recovery), but a branch that never finished provisioning has no session, no PR, and nothing
|
||||
* else pointing at it, so leaving it behind would just trade one leak for a smaller one. Forced
|
||||
* (`-D`) because the branch is new and unmerged by construction. The worktree is removed first:
|
||||
* a branch checked out by a worktree cannot be deleted until the worktree that holds it is gone.
|
||||
*
|
||||
* <p>Cleanup failure must never mask {@code original} — that is the exception that explains
|
||||
* what actually went wrong — so a failure here is only logged, matching the pattern already
|
||||
* used in {@code SessionManager#acquireWithWorktree}'s own catch block.
|
||||
*/
|
||||
private void cleanupAfterAddFailure(String repoRoot, String worktreePath, String branch, RuntimeException original) {
|
||||
log.warn("provisioning failed for branch={} path={}: {} — cleaning up before rethrowing",
|
||||
branch, worktreePath, original.getMessage());
|
||||
try {
|
||||
remove(repoRoot, worktreePath);
|
||||
} catch (RuntimeException cleanup) {
|
||||
log.warn("failed to remove leaked worktree {} after provisioning error: {}",
|
||||
worktreePath, cleanup.getMessage());
|
||||
}
|
||||
try {
|
||||
deleteBranch(repoRoot, branch);
|
||||
} catch (RuntimeException cleanup) {
|
||||
log.warn("failed to remove leaked branch {} after provisioning error: {}",
|
||||
branch, cleanup.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void deleteBranch(String repoRoot, String branch) {
|
||||
exec("git", "-C", repoRoot, "branch", "-D", branch);
|
||||
}
|
||||
|
||||
/**
|
||||
* A linked worktree shares its primary checkout's git config. Remove HTTPS user info before
|
||||
* adding one, so a credential accidentally embedded in that config cannot reach the member.
|
||||
|
||||
@@ -334,7 +334,20 @@ public final class SessionManager implements TurnListener {
|
||||
// slot that no longer appears in the roster and can never be reclaimed.
|
||||
launcher.stop(paneId);
|
||||
if (removed != null && !preserveWorktree && removed.worktree() != null) {
|
||||
worktrees.remove(worktrees.repoRoot(removed.cwd()), removed.worktree());
|
||||
// fleetd #283: this is the one cleanup step in this method that used to be bare. By the
|
||||
// time it runs, the registry entry, the retained handle, and the pane are all already
|
||||
// gone — so a throw here (a stale index lock, a slow filesystem, `remove`'s own 30s exec
|
||||
// timeout) must not escape release(): there is no retry path (a second stop on this
|
||||
// paneId is a no-op), and the caller would otherwise see a "failed stop" for a session
|
||||
// that is in fact fully torn down. Log and swallow, matching every sibling step above.
|
||||
try {
|
||||
worktrees.remove(worktrees.repoRoot(removed.cwd()), removed.worktree());
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("failed to remove worktree {} for pane={} terminal={} after release: the "
|
||||
+ "pane is already stopped and the session already deregistered, so this is "
|
||||
+ "not retryable — the directory must be reclaimed manually: {}",
|
||||
removed.worktree(), paneId, removed.terminalId(), e.toString());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -504,11 +517,26 @@ public final class SessionManager implements TurnListener {
|
||||
log.warn("spawn failed for profile={} role={} branch={} path={}: {}",
|
||||
preResolvedProfile, memberRole, branch, path, e.getMessage());
|
||||
if (path != null) {
|
||||
// fleetd #283: this catch covers every failure AFTER worktrees.add() returned —
|
||||
// overlayParity, shareWithGroup, launcher.spawn itself — so by this point `branch`
|
||||
// was actually created in git. #274 fixed the sibling failure INSIDE add() by having
|
||||
// GitWorktrees.cleanupAfterAddFailure delete both the worktree and the branch it
|
||||
// provisioned; this path removed only the worktree and left the branch orphaned. A
|
||||
// spawn failure here is routine (a quarantined credential, a backend refusal), so
|
||||
// every occurrence leaked a `worker/<slug>-<nonce>` branch nothing ever pointed at
|
||||
// again. Reuse the same Worktrees.deleteBranch GitWorktrees already has, rather than
|
||||
// a second copy of the git command. Best-effort and log-only, like the worktree
|
||||
// removal right above it — neither cleanup step may mask the original exception.
|
||||
try {
|
||||
worktrees.remove(repoRoot, path);
|
||||
} catch (RuntimeException cleanup) {
|
||||
log.warn("failed to clean up worktree {} after spawn error: {}", path, cleanup.getMessage());
|
||||
}
|
||||
try {
|
||||
worktrees.deleteBranch(repoRoot, branch);
|
||||
} catch (RuntimeException cleanup) {
|
||||
log.warn("failed to clean up branch {} after spawn error: {}", branch, cleanup.getMessage());
|
||||
}
|
||||
}
|
||||
throw e;
|
||||
}
|
||||
|
||||
@@ -11,6 +11,15 @@ public interface Worktrees {
|
||||
/** git -C <repoRoot> worktree remove --force <path>. Idempotent (already-gone tolerated). */
|
||||
void remove(String repoRoot, String worktreePath);
|
||||
|
||||
/**
|
||||
* git -C {@code repoRoot} branch -D {@code branch}. Force-deletes a branch that has no other
|
||||
* owner — used only on the failed-provisioning path (fleetd #274, #283), never on a normal
|
||||
* release: {@link SessionManager#release} deliberately leaves a released session's branch
|
||||
* behind so a lead can still recover the work, and this method must never be called from
|
||||
* that path.
|
||||
*/
|
||||
void deleteBranch(String repoRoot, String branch);
|
||||
|
||||
/**
|
||||
* True when the worktree holds uncommitted changes the bridge cannot see: tracked
|
||||
* modifications, staged files, or untracked files. {@code git status --porcelain} is the
|
||||
|
||||
@@ -226,4 +226,44 @@ class FleetdBackendErrorSinkTest {
|
||||
assertTrue(remaining.isPresent(), "two distinct targets must start a cool-off");
|
||||
assertEquals(1, leadClient.sendCount());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a backend-error session no longer blocks the real maxLoad spawn gate")
|
||||
void backendErrorSessionDoesNotBlockFreshSpawnAtMaxLoad() {
|
||||
SessionManager sessions = capacityLimitedSessions();
|
||||
MemberSession failed = sessions.acquire("terra", null, null, null);
|
||||
|
||||
assertTrue(sessions.onBackendError(failed.terminalId(), "backend exited"));
|
||||
|
||||
MemberSession fresh = sessions.acquire("terra", null, null, null);
|
||||
assertEquals("terra", fresh.profile(), "the real maxLoad gate grants a fresh spawn after a backend error");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a failed session no longer blocks the real maxLoad spawn gate")
|
||||
void failedSessionDoesNotBlockFreshSpawnAtMaxLoad() {
|
||||
SessionManager sessions = capacityLimitedSessions();
|
||||
MemberSession failed = sessions.acquire("terra", null, null, null);
|
||||
|
||||
sessions.onTurnFailed(failed.terminalId());
|
||||
|
||||
MemberSession fresh = sessions.acquire("terra", null, null, null);
|
||||
assertEquals("terra", fresh.profile(), "the real maxLoad gate grants a fresh spawn after a failed turn");
|
||||
}
|
||||
|
||||
private static SessionManager capacityLimitedSessions() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FleetConfig.Profile profile = new FleetConfig.Profile("terra", "http://gx00.gw:8000", "coder",
|
||||
null, "FLEETD_WORKER_TOKEN", List.of("claude"), "tab", "fleetd-workers",
|
||||
"w #{n}", null, null, null, null, null, null, null, 1.0f, 1);
|
||||
Map<String, FleetConfig.Profile> profiles = Map.of("terra", profile);
|
||||
ClaudeCodeLauncher adapter = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), profiles, "terra", _ -> "tok");
|
||||
AtomicReference<SessionManager> sessionsRef = new AtomicReference<>();
|
||||
CompositePeerLauncher workers = new CompositePeerLauncher(List.of(adapter), "terra", profiles,
|
||||
PlacementPolicies.fixed(), name -> Fleetd.liveSessionCount(sessionsRef.get().roster(), name));
|
||||
SessionManager sessions = new SessionManager(workers);
|
||||
sessionsRef.set(sessions);
|
||||
return sessions;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,73 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import ch.qos.logback.classic.Logger;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
|
||||
/**
|
||||
* fleetd #184: startup must state whether members share fleetd's OS user or use a separate herdr.
|
||||
*/
|
||||
class MemberTrustModelReportTest {
|
||||
|
||||
private static FleetConfig load(Path dir, String yaml) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, yaml);
|
||||
return FleetConfig.load(f);
|
||||
}
|
||||
|
||||
private static ListAppender<ILoggingEvent> attach() {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
return appender;
|
||||
}
|
||||
|
||||
private static void detach(ListAppender<ILoggingEvent> appender) {
|
||||
((Logger) LoggerFactory.getLogger(Fleetd.class)).detachAppender(appender);
|
||||
}
|
||||
|
||||
private static String report(FleetConfig cfg) {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class);
|
||||
ch.qos.logback.classic.Level original = logger.getLevel();
|
||||
logger.setLevel(ch.qos.logback.classic.Level.INFO);
|
||||
ListAppender<ILoggingEvent> appender = attach();
|
||||
try {
|
||||
Fleetd.reportMemberTrustModel(cfg);
|
||||
} finally {
|
||||
detach(appender);
|
||||
logger.setLevel(original);
|
||||
}
|
||||
return appender.list.getFirst().getFormattedMessage();
|
||||
}
|
||||
|
||||
@Test
|
||||
void unsetMemberHerdrSocketStatesThatMembersAreNotSandboxed(@TempDir Path dir) throws Exception {
|
||||
FleetConfig cfg = load(dir, "bind:\n host: 127.0.0.1\n port: 8765\n");
|
||||
|
||||
assertEquals("member trust model: members run as the same OS user as fleetd, not in a sandbox. "
|
||||
+ "A member can read any file this user can read, including SSH keys and credential "
|
||||
+ "stores, whatever memberCredentials says. To add a real boundary, route members to "
|
||||
+ "a second herdr under a different OS user with memberHerdrSocket.",
|
||||
report(cfg));
|
||||
}
|
||||
|
||||
@Test
|
||||
void configuredMemberHerdrSocketStatesThatFleetdCannotConfirmTheBoundary(@TempDir Path dir) throws Exception {
|
||||
FleetConfig cfg = load(dir, "memberHerdrSocket: /tmp/member-herdr.sock\n");
|
||||
|
||||
assertEquals("member trust model: members are routed to a separate herdr through "
|
||||
+ "memberHerdrSocket. fleetd cannot see that herdr's uid, so confirm it runs "
|
||||
+ "as a different OS user before treating it as a boundary.",
|
||||
report(cfg));
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,96 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import com.tngtech.archunit.base.DescribedPredicate;
|
||||
import com.tngtech.archunit.core.domain.JavaClass;
|
||||
import com.tngtech.archunit.core.domain.JavaClass.Predicates;
|
||||
import com.tngtech.archunit.core.importer.ClassFileImporter;
|
||||
import com.tngtech.archunit.core.importer.ImportOption;
|
||||
import com.tngtech.archunit.library.dependencies.SliceRule;
|
||||
import com.tngtech.archunit.library.dependencies.SlicesRuleDefinition;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
/**
|
||||
* fleetd #131 (CB-627): enforce package boundaries with an ArchUnit test instead of a
|
||||
* Maven module split.
|
||||
*
|
||||
* <p>This test fails the build the moment a NEW cycle appears between the top-level
|
||||
* {@code dev.ltms.fleet.*} packages. Today's cycles are recorded below as explicit,
|
||||
* narrow exceptions: each one ignores dependencies between exactly the two named
|
||||
* packages, in both directions, and nothing else. A cycle through any other pair of
|
||||
* packages -- or a brand new pair -- still fails this test.
|
||||
*
|
||||
* <p><b>Main code only.</b> The import excludes test classes
|
||||
* ({@link ImportOption.Predefined#DO_NOT_INCLUDE_TESTS}). Test code legitimately wires
|
||||
* across many packages for setup and mocking; that is not part of the shipped
|
||||
* architecture this rule protects. Verified: importing test classes too pulls in a much
|
||||
* larger, noisier cycle set -- {@code herdr}, {@code member}, {@code peer}, {@code
|
||||
* config}, {@code guard} and {@code placement} all show up in cycles that disappear the
|
||||
* moment test classes are excluded. Scanning off the classpath via {@code
|
||||
* importPackages(...)} (not a hardcoded {@code target/classes} path) also keeps this
|
||||
* test correct regardless of the working directory the build is invoked from.
|
||||
*
|
||||
* <p><b>No package moves here</b> -- ticket #131 is explicit that removing a cycle is
|
||||
* its own, later PR. See the comment on each exception below for which ticket step
|
||||
* removes it.
|
||||
*/
|
||||
class PackageCyclesTest {
|
||||
|
||||
@Test
|
||||
void packagesAreFreeOfCycles() {
|
||||
var classes = new ClassFileImporter()
|
||||
.withImportOption(ImportOption.Predefined.DO_NOT_INCLUDE_TESTS)
|
||||
.importPackages("dev.ltms.fleet");
|
||||
|
||||
SliceRule rule = SlicesRuleDefinition.slices()
|
||||
.matching("dev.ltms.fleet.(*)..")
|
||||
.should().beFreeOfCycles();
|
||||
|
||||
// fleetd #131 step 1: move ConnectionIdentity so authz stops depending on the
|
||||
// MCP layer. Evidence: auth/CallerResolver.java:3 imports mcp.ConnectionIdentity;
|
||||
// mcp/FleetMcp.java:3-7 imports auth.AuditLog, Authz, CallerResolver, Principal,
|
||||
// Role.
|
||||
rule = ignoreCycle(rule, "auth", "mcp");
|
||||
|
||||
// fleetd #131 step 2: PrimaryRegistry is used by loops in msg; move it, or put
|
||||
// an interface between msg and mcp. Evidence: msg/ReplyPushLoop.java:5 and
|
||||
// msg/LeadHeartbeatLoop.java:5 import mcp.PrimaryRegistry; mcp/FleetMcp.java:15-18
|
||||
// imports msg.LeadChannel, LeadMessage, MessageService, Rendezvous.
|
||||
rule = ignoreCycle(rule, "mcp", "msg");
|
||||
|
||||
// fleetd #131 -- found while implementing this test, NOT one of the ticket's
|
||||
// original three; it names its own follow-up step before removal. Evidence:
|
||||
// inject/CompletionResolver.java:4-5, inject/Injector.java:6 and
|
||||
// inject/TurnListener.java:3 import msg.Rendezvous / msg.TurnToken;
|
||||
// msg/MessageService.java:6 imports inject.Injector.
|
||||
rule = ignoreCycle(rule, "inject", "msg");
|
||||
|
||||
// fleetd #131 -- same as above, its own follow-up. Evidence:
|
||||
// metrics/FleetMetrics.java:3 imports msg.ReplyInbox; msg/MessageService.java:7-8,
|
||||
// msg/LeadHeartbeatLoop.java:6-7 and msg/ReplyPushLoop.java:6-7 import
|
||||
// metrics.FleetMetrics / metrics.Metrics.
|
||||
rule = ignoreCycle(rule, "metrics", "msg");
|
||||
|
||||
// fleetd #131 -- same as above, its own follow-up. Evidence:
|
||||
// session/SessionManager.java:7 imports msg.TurnToken;
|
||||
// msg/LeadHeartbeatLoop.java:8 imports session.MemberSession.
|
||||
rule = ignoreCycle(rule, "msg", "session");
|
||||
|
||||
rule.check(classes);
|
||||
}
|
||||
|
||||
/**
|
||||
* Accepts today's known cycle between two top-level packages, and nothing else.
|
||||
* Ignoring both directions removes exactly this pair from cycle detection; every
|
||||
* other dependency -- including any new one added later, between these same two
|
||||
* packages or any other pair -- is still checked.
|
||||
*/
|
||||
private static SliceRule ignoreCycle(SliceRule rule, String packageA, String packageB) {
|
||||
return rule
|
||||
.ignoreDependency(residesIn(packageA), residesIn(packageB))
|
||||
.ignoreDependency(residesIn(packageB), residesIn(packageA));
|
||||
}
|
||||
|
||||
private static DescribedPredicate<JavaClass> residesIn(String topLevelPackage) {
|
||||
return Predicates.resideInAPackage("dev.ltms.fleet." + topLevelPackage + "..");
|
||||
}
|
||||
}
|
||||
@@ -155,6 +155,90 @@ class FleetConfigTest {
|
||||
assertTrue(e.getMessage().contains("errorPattern"), "the offending key is named: " + e.getMessage());
|
||||
}
|
||||
|
||||
// ── fleetd #273: exhaustedPattern gets the same load-time validation as its sibling errorPattern ──
|
||||
|
||||
@Test
|
||||
void aProfileWithAMalformedExhaustedPatternIsRejectedAtLoadNamingTheProfileAndKey(@TempDir Path dir)
|
||||
throws Exception {
|
||||
Path f = dir.resolve("malformed-exhausted-pattern.yaml");
|
||||
Files.writeString(f, """
|
||||
profiles:
|
||||
ltms-local:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
exhaustedPattern: "["
|
||||
""");
|
||||
|
||||
IllegalStateException e = assertThrows(IllegalStateException.class, () -> FleetConfig.load(f));
|
||||
assertTrue(e.getMessage().contains("ltms-local"), "the offending profile is named: " + e.getMessage());
|
||||
assertTrue(e.getMessage().contains("exhaustedPattern"), "the offending key is named: " + e.getMessage());
|
||||
}
|
||||
|
||||
@Test
|
||||
void aMalformedErrorPatternAndAMalformedExhaustedPatternAreBothReportedFromOneLoad(@TempDir Path dir)
|
||||
throws Exception {
|
||||
Path f = dir.resolve("both-malformed.yaml");
|
||||
Files.writeString(f, """
|
||||
profiles:
|
||||
sonnet:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
errorPattern: "(unterminated["
|
||||
terra:
|
||||
baseUrl: http://gx01.gw:8000
|
||||
exhaustedPattern: "["
|
||||
""");
|
||||
|
||||
IllegalStateException e = assertThrows(IllegalStateException.class, () -> FleetConfig.load(f));
|
||||
assertTrue(e.getMessage().contains("sonnet"), "the errorPattern profile is named: " + e.getMessage());
|
||||
assertTrue(e.getMessage().contains("errorPattern"), e.getMessage());
|
||||
assertTrue(e.getMessage().contains("terra"), "the exhaustedPattern profile is named: " + e.getMessage());
|
||||
assertTrue(e.getMessage().contains("exhaustedPattern"), e.getMessage());
|
||||
}
|
||||
|
||||
@Test
|
||||
void validErrorPatternAndExhaustedPatternBothLoadFine(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("both-valid.yaml");
|
||||
Files.writeString(f, """
|
||||
profiles:
|
||||
ltms-local:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
errorPattern: "credential outage"
|
||||
exhaustedPattern: "usage limit has been reached"
|
||||
""");
|
||||
|
||||
FleetConfig.Profile w = FleetConfig.load(f).profiles().get("ltms-local");
|
||||
assertEquals("credential outage", w.errorPattern());
|
||||
assertEquals("usage limit has been reached", w.exhaustedPattern());
|
||||
}
|
||||
|
||||
@Test
|
||||
void aBlankExhaustedPatternNormalizesToNullJustLikeUnset(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("blank-exhausted-pattern.yaml");
|
||||
Files.writeString(f, """
|
||||
profiles:
|
||||
ltms-local:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
exhaustedPattern: " "
|
||||
""");
|
||||
|
||||
FleetConfig.Profile w = FleetConfig.load(f).profiles().get("ltms-local");
|
||||
assertNull(w.exhaustedPattern());
|
||||
assertFalse(w.hasExhaustedPattern());
|
||||
}
|
||||
|
||||
@Test
|
||||
void aProfileWithNoExhaustedPatternLoadsFine(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("no-exhausted-pattern.yaml");
|
||||
Files.writeString(f, """
|
||||
profiles:
|
||||
ltms-local:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
""");
|
||||
|
||||
FleetConfig.Profile w = FleetConfig.load(f).profiles().get("ltms-local");
|
||||
assertNull(w.exhaustedPattern());
|
||||
assertFalse(w.hasExhaustedPattern());
|
||||
}
|
||||
|
||||
@Test
|
||||
void withProfileCarriesErrorPatternThrough(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("with-profile-error-pattern.yaml");
|
||||
@@ -2010,22 +2094,93 @@ class FleetConfigTest {
|
||||
"deny-list normalizes onto the canonical deny-by-default value");
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-633: {@code SSH_AUTH_SOCK} is a decision, never a default — absent, blank, or misspelled,
|
||||
* it stays BLOCKED; only the literal "allow" (any case) passes it through. A typo like "alow"
|
||||
* failing safe here is the whole point of making it a knob.
|
||||
*/
|
||||
@Test
|
||||
void sshAuthSockDefaultsToBlockedAndOnlyExplicitAllowUnblocksIt() {
|
||||
assertTrue(new FleetConfig.MemberCredentials("allow-list", List.of(), List.of()).sshAuthSock().equals("block"),
|
||||
"absent knob blocks SSH_AUTH_SOCK");
|
||||
assertFalse(new FleetConfig.MemberCredentials("allow-list", List.of(), List.of()).sshAuthSockAllowed());
|
||||
assertFalse(new FleetConfig.MemberCredentials(null, null, null, "").sshAuthSockAllowed(),
|
||||
"blank knob blocks SSH_AUTH_SOCK");
|
||||
assertFalse(new FleetConfig.MemberCredentials(null, null, null, "alow").sshAuthSockAllowed(),
|
||||
"a misspelled value fails SAFE, not open");
|
||||
assertTrue(new FleetConfig.MemberCredentials(null, null, null, "ALLOW").sshAuthSockAllowed(),
|
||||
"the literal allow (case-insensitive) unblocks SSH_AUTH_SOCK");
|
||||
void sshAgentEnvAcceptsEveryCompatibleKeyAndValuePair(@TempDir Path dir) throws Exception {
|
||||
String[][] spellings = {
|
||||
{"sshAuthSock", "block", "omit"},
|
||||
{"sshAuthSock", "allow", "inherit"},
|
||||
{"sshAuthSock", "omit", "omit"},
|
||||
{"sshAuthSock", "inherit", "inherit"},
|
||||
{"sshAgentEnv", "block", "omit"},
|
||||
{"sshAgentEnv", "allow", "inherit"},
|
||||
{"sshAgentEnv", "omit", "omit"},
|
||||
{"sshAgentEnv", "inherit", "inherit"}
|
||||
};
|
||||
|
||||
for (int i = 0; i < spellings.length; i++) {
|
||||
Path file = dir.resolve("member-credentials-ssh-agent-" + i + ".yaml");
|
||||
Files.writeString(file, """
|
||||
bind:
|
||||
port: 8080
|
||||
memberCredentials:
|
||||
policy: allow-list
|
||||
""" + " " + spellings[i][0] + ": " + spellings[i][1] + "\n");
|
||||
|
||||
FleetConfig.MemberCredentials credentials = FleetConfig.load(file).memberCredentials();
|
||||
assertEquals(spellings[i][2], credentials.sshAgentEnv(),
|
||||
spellings[i][0] + ": " + spellings[i][1] + " must normalize correctly");
|
||||
assertEquals("inherit".equals(spellings[i][2]), credentials.sshAgentEnvInherited());
|
||||
}
|
||||
}
|
||||
|
||||
/** Protects the live {@code sshAuthSock: block} allow-list configuration during the rename. */
|
||||
@Test
|
||||
void legacySshAuthSockBlockKeepsLiveAllowListConfigOmitted(@TempDir Path dir) throws Exception {
|
||||
Path file = dir.resolve("live-member-credentials.yaml");
|
||||
Files.writeString(file, """
|
||||
bind:
|
||||
port: 8080
|
||||
memberCredentials:
|
||||
policy: allow-list
|
||||
sshAuthSock: block
|
||||
""");
|
||||
|
||||
FleetConfig.MemberCredentials credentials = FleetConfig.load(file).memberCredentials();
|
||||
|
||||
assertEquals("omit", credentials.sshAgentEnv());
|
||||
assertFalse(credentials.sshAgentEnvInherited(), "the live config must omit SSH_AUTH_SOCK");
|
||||
}
|
||||
|
||||
@Test
|
||||
void sshAgentEnvWinsWhenBothCompatibleKeysArePresent(@TempDir Path dir) throws Exception {
|
||||
Path file = dir.resolve("both-ssh-agent-keys.yaml");
|
||||
Files.writeString(file, """
|
||||
bind:
|
||||
port: 8080
|
||||
memberCredentials:
|
||||
policy: allow-list
|
||||
sshAuthSock: allow
|
||||
sshAgentEnv: omit
|
||||
""");
|
||||
|
||||
FleetConfig.MemberCredentials credentials = FleetConfig.load(file).memberCredentials();
|
||||
|
||||
assertEquals("omit", credentials.sshAgentEnv());
|
||||
assertFalse(credentials.sshAgentEnvInherited());
|
||||
}
|
||||
|
||||
@Test
|
||||
void sshAgentEnvDefaultsToOmitAndUnknownValuesFailClosed(@TempDir Path dir) throws Exception {
|
||||
Path absent = dir.resolve("member-credentials-ssh-agent-absent.yaml");
|
||||
Files.writeString(absent, """
|
||||
bind:
|
||||
port: 8080
|
||||
memberCredentials:
|
||||
policy: allow-list
|
||||
""");
|
||||
assertEquals("omit", FleetConfig.load(absent).memberCredentials().sshAgentEnv());
|
||||
|
||||
Path unknown = dir.resolve("member-credentials-ssh-agent-unknown.yaml");
|
||||
Files.writeString(unknown, """
|
||||
bind:
|
||||
port: 8080
|
||||
memberCredentials:
|
||||
policy: allow-list
|
||||
sshAgentEnv: inhert
|
||||
""");
|
||||
FleetConfig.MemberCredentials credentials = FleetConfig.load(unknown).memberCredentials();
|
||||
assertEquals("omit", credentials.sshAgentEnv());
|
||||
assertFalse(credentials.sshAgentEnvInherited(), "an unknown value must fail closed");
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -24,7 +24,9 @@ import org.slf4j.LoggerFactory;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledFuture;
|
||||
import java.util.concurrent.ScheduledThreadPoolExecutor;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
@@ -434,4 +436,120 @@ class FleetHealthMonitorTest {
|
||||
logger.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
|
||||
// --- fleetd #280: a GONE/NEVER_READY guess whose fleet_ask lapses AFTER the first sweep must
|
||||
// still be swept, without ever reaching into a target that has since recovered or left the
|
||||
// roster. See FleetHealthMonitor.recheckTerminalTarget's javadoc for the full reachability chain.
|
||||
|
||||
@Test void terminalTransitionSchedulesExactlyOneDelayedRecheck() {
|
||||
ScheduledThreadPoolExecutor scheduler = new ScheduledThreadPoolExecutor(1);
|
||||
FleetHealthMonitor monitor = monitor(new FakeHerdr(),
|
||||
List.of(member("term_a", MemberSession.State.BUSY, 0, 0)), scheduler, () -> 1, 600,
|
||||
(_, _) -> { });
|
||||
|
||||
monitor.reportTransition("term_a", HealthState.GONE);
|
||||
|
||||
// Driven via reportTransition directly (not tick()), so the queue holds only the recheck.
|
||||
assertEquals(1, scheduler.getQueue().size());
|
||||
ScheduledFuture<?> scheduled = (ScheduledFuture<?>) scheduler.getQueue().peek();
|
||||
assertTrue(scheduled.getDelay(TimeUnit.SECONDS) > 100,
|
||||
"the delay must clear the worst-case fleet_ask lapse window (up to 115s)");
|
||||
|
||||
// An unchanged tick must not queue a second one (CB-580's fire-once rule extends to this).
|
||||
monitor.reportTransition("term_a", HealthState.GONE);
|
||||
assertEquals(1, scheduler.getQueue().size());
|
||||
monitor.stop();
|
||||
}
|
||||
|
||||
@Test void recheckIsANoOpOnceTheTargetHasRecovered() {
|
||||
RecordingFailTarget failTarget = new RecordingFailTarget();
|
||||
FleetHealthMonitor monitor = monitorWith(failTarget);
|
||||
monitor.reportTransition("term_a", HealthState.GONE);
|
||||
assertEquals(1, failTarget.calls.size());
|
||||
|
||||
monitor.reportTransition("term_a", HealthState.IDLE); // recovered before the recheck fired
|
||||
monitor.recheckTerminalTarget("term_a", HealthState.GONE);
|
||||
|
||||
assertEquals(1, failTarget.calls.size(), "a recovered target must not be reached into again");
|
||||
monitor.stop();
|
||||
}
|
||||
|
||||
@Test void recheckIsANoOpForATargetItNeverObserved() {
|
||||
// Mirrors "left the roster": tick() prunes states.keySet() to the current roster on release
|
||||
// (see FleetHealthMonitor.tick), so a target this monitor never recorded is the same case.
|
||||
RecordingFailTarget failTarget = new RecordingFailTarget();
|
||||
FleetHealthMonitor monitor = monitorWith(failTarget);
|
||||
|
||||
monitor.recheckTerminalTarget("term_never_seen", HealthState.GONE);
|
||||
|
||||
assertEquals(0, failTarget.calls.size(), "an untracked/released target must not be reached into");
|
||||
monitor.stop();
|
||||
}
|
||||
|
||||
/**
|
||||
* The scenario from the ticket, end to end, driven through the real {@link MessageService}: a
|
||||
* target's ask is still genuinely open when health first observes GONE (sweep must skip it,
|
||||
* exactly as {@code abandonDoesNotFailAnAsyncTicketWaitingForAnAnswer} pins), the ask then lapses
|
||||
* on its own, an unchanged tick still must not refire, and only the delayed recheck sweeps the
|
||||
* now-lapsed ticket to FAILED.
|
||||
*/
|
||||
@Test void delayedRecheckSweepsATicketWhoseAskLapsedAfterGoneWasFirstObserved() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr().withAgent("worker", "term_a", "pane-term_a", "tab_a")
|
||||
.readText("$ prompt");
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
Injector injector = new Injector(agents);
|
||||
InMemoryReplyInbox inbox = new InMemoryReplyInbox();
|
||||
inbox.own("term_a");
|
||||
MessageService messages = new MessageService(agents, injector, rendezvous, inbox);
|
||||
FleetHealthMonitor monitor = new FleetHealthMonitor(agents,
|
||||
() -> List.of(member("term_a", MemberSession.State.READY, 0, 0)), messages,
|
||||
new ScheduledThreadPoolExecutor(1), () -> 1, 60, 600, messages::abandon);
|
||||
|
||||
String ticket = messages.sendAsync("term_a", "task that asks");
|
||||
long deadline = System.currentTimeMillis() + 2000;
|
||||
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertTrue(rendezvous.isWaiting("term_a"), "the async send should have opened its waiter");
|
||||
injector.onStatus("term_a", AgentStatus.IDLE); // deliver the task
|
||||
injector.onStatus("term_a", AgentStatus.WORKING); // the worker picks it up
|
||||
|
||||
// The worker asks, with a short timeout so its own fleet_ask lapses quickly in test time.
|
||||
CompletableFuture<MessageService.AskResult> ask = CompletableFuture.supplyAsync(
|
||||
() -> messages.ask("term_a", "which config?", 200));
|
||||
awaitPhase(messages, ticket, MessageService.Phase.ASKING);
|
||||
|
||||
monitor.reportTransition("term_a", HealthState.GONE);
|
||||
assertEquals(MessageService.Phase.ASKING, messages.poll(ticket).phase(),
|
||||
"the first sweep must not fail a ticket that is still genuinely being asked");
|
||||
|
||||
// The worker's own fleet_ask now lapses on its own — task.question clears to null.
|
||||
assertEquals(MessageService.AskOutcome.TIMED_OUT, ask.get(5, TimeUnit.SECONDS).outcome());
|
||||
|
||||
// An unchanged tick still must not refire (CB-580).
|
||||
monitor.reportTransition("term_a", HealthState.GONE);
|
||||
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket).phase());
|
||||
|
||||
// The delayed recheck scheduled for the original transition finally sweeps it.
|
||||
monitor.recheckTerminalTarget("term_a", HealthState.GONE);
|
||||
assertEquals(MessageService.Phase.FAILED, messages.poll(ticket).phase());
|
||||
|
||||
monitor.stop();
|
||||
}
|
||||
|
||||
private static MessageService.TaskView awaitPhase(MessageService messages, String ticket,
|
||||
MessageService.Phase phase) throws InterruptedException {
|
||||
long deadline = System.currentTimeMillis() + 2000;
|
||||
MessageService.TaskView view;
|
||||
do {
|
||||
view = messages.poll(ticket);
|
||||
if (view.phase() == phase) {
|
||||
return view;
|
||||
}
|
||||
Thread.sleep(5);
|
||||
} while (System.currentTimeMillis() < deadline);
|
||||
assertEquals(phase, view.phase());
|
||||
return view;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -7,6 +7,7 @@ import java.util.ArrayList;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.CopyOnWriteArrayList;
|
||||
|
||||
/**
|
||||
@@ -41,6 +42,9 @@ public final class FakeHerdr implements HerdrClient {
|
||||
private int agentPaneBusyFor = 0;
|
||||
private int workerTabPaneCount = 1;
|
||||
private String paneCloseErrorCode = null;
|
||||
private final Map<String, String> paneCloseErrorCodeFor = new ConcurrentHashMap<>();
|
||||
private String tabCloseErrorCode = null;
|
||||
private final Map<String, String> tabCloseErrorCodeFor = new ConcurrentHashMap<>();
|
||||
private String agentSendErrorCode = null;
|
||||
private boolean noPanes = false;
|
||||
private volatile String agentStatus = "idle"; // steady-state agent.get status
|
||||
@@ -86,12 +90,44 @@ public final class FakeHerdr implements HerdrClient {
|
||||
return this;
|
||||
}
|
||||
|
||||
/** Make {@code pane.close} fail with this herdr error code. */
|
||||
/** Make {@code pane.close} fail with this herdr error code, for every pane. */
|
||||
public FakeHerdr paneCloseFailsWith(String code) {
|
||||
this.paneCloseErrorCode = code;
|
||||
return this;
|
||||
}
|
||||
|
||||
/**
|
||||
* Make {@code pane.close} fail with this herdr error code, but only for the given {@code
|
||||
* pane_id} — every other pane's {@code pane.close} still succeeds. Unlike {@link
|
||||
* #paneCloseFailsWith}, which fails every call regardless of which pane it targets, this lets a
|
||||
* test reap/release several sessions at once and make exactly one of them fail to stop, so the
|
||||
* others' teardown can be asserted to proceed normally (fleetd #290).
|
||||
*/
|
||||
public FakeHerdr paneCloseFailsForPane(String paneId, String code) {
|
||||
this.paneCloseErrorCodeFor.put(paneId, code);
|
||||
return this;
|
||||
}
|
||||
|
||||
/** Make {@code tab.close} fail with this herdr error code, for every tab. */
|
||||
public FakeHerdr tabCloseFailsWith(String code) {
|
||||
this.tabCloseErrorCode = code;
|
||||
return this;
|
||||
}
|
||||
|
||||
/**
|
||||
* Make {@code tab.close} fail with this herdr error code, but only for the given {@code
|
||||
* tab_id} — every other tab's {@code tab.close} still succeeds. The {@code tab.close}
|
||||
* counterpart to {@link #paneCloseFailsForPane} (fleetd #290): lets a test make exactly one
|
||||
* session's tab teardown fail while proving the rest of {@code stop()} — {@code
|
||||
* releaseZdotdir}, and the caller's worktree removal — still runs (fleetd #293). Named "ForTab"
|
||||
* rather than "ForPane" (unlike its sibling) because {@code tab.close} keys on {@code tab_id},
|
||||
* not a pane id.
|
||||
*/
|
||||
public FakeHerdr tabCloseFailsForTab(String tabId, String code) {
|
||||
this.tabCloseErrorCodeFor.put(tabId, code);
|
||||
return this;
|
||||
}
|
||||
|
||||
/**
|
||||
* Make {@code pane.list} report no panes at all — models a second herdr daemon (CB-185) that
|
||||
* simply does not host the pane a {@link PaneLocator} is searching for.
|
||||
@@ -338,7 +374,17 @@ public final class FakeHerdr implements HerdrClient {
|
||||
.formatted(workerTabPaneCount,
|
||||
seeded.isEmpty() ? "" : "," + String.join(",", seeded)));
|
||||
}
|
||||
case "tab.close" -> mapper.readTree("{\"type\":\"ok\"}");
|
||||
case "tab.close" -> {
|
||||
Object tabIdParam = params instanceof Map<?, ?> m ? m.get("tab_id") : null;
|
||||
String perTabCode = tabIdParam == null ? null
|
||||
: tabCloseErrorCodeFor.get(String.valueOf(tabIdParam));
|
||||
String code = perTabCode != null ? perTabCode : tabCloseErrorCode;
|
||||
if (code != null) {
|
||||
throw new HerdrException("herdr error [" + code + "]: tab.close failed",
|
||||
code, null);
|
||||
}
|
||||
yield mapper.readTree("{\"type\":\"ok\"}");
|
||||
}
|
||||
case "pane.get" -> mapper.readTree("""
|
||||
{"type":"pane_info","pane":{"pane_id":"w9:pW","workspace_id":"w9",
|
||||
"tab_id":"w9:t2","agent_status":"idle"}}""");
|
||||
@@ -360,9 +406,13 @@ public final class FakeHerdr implements HerdrClient {
|
||||
"foreground_processes":[]}}""");
|
||||
}
|
||||
case "pane.close" -> {
|
||||
if (paneCloseErrorCode != null) {
|
||||
throw new HerdrException("herdr error [" + paneCloseErrorCode + "]: pane.close failed",
|
||||
paneCloseErrorCode, null);
|
||||
Object paneIdParam = params instanceof Map<?, ?> m ? m.get("pane_id") : null;
|
||||
String perPaneCode = paneIdParam == null ? null
|
||||
: paneCloseErrorCodeFor.get(String.valueOf(paneIdParam));
|
||||
String code = perPaneCode != null ? perPaneCode : paneCloseErrorCode;
|
||||
if (code != null) {
|
||||
throw new HerdrException("herdr error [" + code + "]: pane.close failed",
|
||||
code, null);
|
||||
}
|
||||
yield mapper.readTree("{\"type\":\"ok\"}");
|
||||
}
|
||||
|
||||
@@ -24,8 +24,13 @@ import io.modelcontextprotocol.spec.McpSchema;
|
||||
import org.junit.jupiter.api.AfterEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.LinkedHashSet;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.regex.Matcher;
|
||||
import java.util.regex.Pattern;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
@@ -44,6 +49,9 @@ import static org.junit.jupiter.api.Assertions.*;
|
||||
*/
|
||||
class FleetMcpAuthzTest {
|
||||
|
||||
private static final Path MCP_SOURCE = Path.of("src/main/java/dev/ltms/fleet/mcp/FleetMcp.java");
|
||||
private static final Pattern TOOL_REGISTRATION = Pattern.compile("tool\\(\\\"(fleet_[a-z_]+)\\\"");
|
||||
|
||||
private final FakeHerdr herdr = new FakeHerdr();
|
||||
private final AgentControl agents = new AgentControl(herdr);
|
||||
private Metrics metrics;
|
||||
@@ -179,6 +187,89 @@ class FleetMcpAuthzTest {
|
||||
"no CallerResolver supplied ⇒ authorization not enforced (legacy behaviour)");
|
||||
}
|
||||
|
||||
// --- which action each tool hands the gate (fleetd #272) ------------------------------------
|
||||
|
||||
/**
|
||||
* fleetd #272: {@code fleet_poll{target}} drains a session's reply inbox, so it needs
|
||||
* {@link Authz.Action#DRAIN} -- not the {@link Authz.Action#READ} the handler passed for both
|
||||
* of its branches until this ticket.
|
||||
*
|
||||
* <p>This asserts against {@link FleetMcp#pollAction}, the method the handler itself calls, so
|
||||
* the handler holds no separate copy of the rule that this test could miss. Every other test in
|
||||
* this class checks the policy table (is a worker allowed to DRAIN?) and all of them passed for
|
||||
* the whole time the defect was live -- the table was right, the action fed to it was wrong.
|
||||
*/
|
||||
@Test
|
||||
void pollingByTargetIsADrainAndPollingByTicketIsARead() {
|
||||
assertEquals(Authz.Action.DRAIN, FleetMcp.pollAction("term_b"),
|
||||
"poll by target removes the replies — that is a drain, not an observation");
|
||||
assertEquals(Authz.Action.READ, FleetMcp.pollAction(null),
|
||||
"poll by ticket changes nothing");
|
||||
assertEquals(Authz.Action.READ, FleetMcp.pollAction(" "),
|
||||
"a blank target is an absent target");
|
||||
}
|
||||
|
||||
@Test
|
||||
void everyRegisteredToolHasItsHandlerActionPinned() {
|
||||
Set<String> registered = toolsTheServerRegisters();
|
||||
assertTrue(registered.size() >= 10,
|
||||
"scraped only " + registered.size() + " tool registrations from FleetMcp (" + registered
|
||||
+ "); the server registers eleven, so the tool(\"…\") scrape has stopped matching");
|
||||
registered.forEach(tool -> assertDoesNotThrow(() -> FleetMcp.toolAction(tool, Map.of()),
|
||||
() -> tool + " is registered but has no pinned authorization action"));
|
||||
|
||||
assertEquals(Authz.Action.SEND, FleetMcp.toolAction("fleet_send", Map.of()));
|
||||
assertEquals(Authz.Action.REPLY, FleetMcp.toolAction("fleet_reply", Map.of()));
|
||||
assertEquals(Authz.Action.ASK, FleetMcp.toolAction("fleet_ask", Map.of()));
|
||||
assertEquals(Authz.Action.READ, FleetMcp.toolAction("fleet_status", Map.of()));
|
||||
assertEquals(Authz.Action.DRAIN, FleetMcp.toolAction("fleet_ack", Map.of()));
|
||||
assertEquals(Authz.Action.SPAWN, FleetMcp.toolAction("fleet_spawn", Map.of()));
|
||||
assertEquals(Authz.Action.READ, FleetMcp.toolAction("fleet_list", Map.of()));
|
||||
assertEquals(Authz.Action.STOP, FleetMcp.toolAction("fleet_stop", Map.of()));
|
||||
assertEquals(Authz.Action.READ, FleetMcp.toolAction("fleet_profiles", Map.of()));
|
||||
assertEquals(Authz.Action.READ, FleetMcp.toolAction("fleet_whoami", Map.of()));
|
||||
assertEquals(Authz.Action.READ, FleetMcp.toolAction("fleet_poll", Map.of("ticket", "task")));
|
||||
assertEquals(Authz.Action.DRAIN, FleetMcp.toolAction("fleet_poll", Map.of("target", "term_b")));
|
||||
}
|
||||
|
||||
private static Set<String> toolsTheServerRegisters() {
|
||||
try {
|
||||
Matcher matcher = TOOL_REGISTRATION.matcher(Files.readString(MCP_SOURCE));
|
||||
Set<String> tools = new LinkedHashSet<>();
|
||||
while (matcher.find()) {
|
||||
tools.add(matcher.group(1));
|
||||
}
|
||||
return tools;
|
||||
} catch (Exception e) {
|
||||
throw new AssertionError("could not scrape FleetMcp tool registrations", e);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void aWorkerMayNotDrainAnotherSessionsInboxByPolling() {
|
||||
FleetMcp m = mcp(true);
|
||||
|
||||
assertNotNull(m.denyFor(WORKER_A, FleetMcp.pollAction("term_b"), "term_b"),
|
||||
"a worker draining a peer's inbox would destroy replies queued for the primary");
|
||||
assertNotNull(m.denyFor(ARCH_DESIGN, FleetMcp.pollAction("term_b"), "term_b"),
|
||||
"an architect has no lifecycle rights either — same gate as fleet_ack");
|
||||
assertNull(m.denyFor(PRIMARY, FleetMcp.pollAction("term_b"), "term_b"),
|
||||
"collecting a held reply is the primary's job");
|
||||
}
|
||||
|
||||
/**
|
||||
* The tightening must not close the branch that legitimately serves non-primary callers: an
|
||||
* architect may {@code fleet_send}, so it owns tickets and must be able to poll them.
|
||||
*/
|
||||
@Test
|
||||
void pollingAnOwnTicketStaysOpenToWorkersAndArchitects() {
|
||||
FleetMcp m = mcp(true);
|
||||
|
||||
assertNull(m.denyFor(WORKER_A, FleetMcp.pollAction(null), null));
|
||||
assertNull(m.denyFor(ARCH_DESIGN, FleetMcp.pollAction(null), null),
|
||||
"an architect delegates with wait:false, so it must be able to poll its ticket");
|
||||
}
|
||||
|
||||
// --- identity reconstruction from the transport context ------------------------------------
|
||||
|
||||
@Test
|
||||
|
||||
@@ -31,9 +31,11 @@ import org.junit.jupiter.api.Test;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.EnumSet;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.function.Function;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
@@ -603,6 +605,60 @@ class FleetMcpTest {
|
||||
assertTrue(out.contains("\"reclaimable\":0"), out);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #284: {@code reclaimable} has exactly one definition, and this pins it over EVERY
|
||||
* {@link MemberSession.State} — so a state added later cannot slip through unconsidered. Both
|
||||
* views in a {@code fleet_list} response call {@link FleetMcp#reclaimable}, so they cannot
|
||||
* drift apart.
|
||||
*
|
||||
* <p>{@code BACKEND_ERROR} and {@code FAILED} are NOT reclaimable on purpose. The ticket asked
|
||||
* for them to be; that half of the ticket was wrong. Their seat is already out of
|
||||
* {@code live}, so it is already in {@code free} — counting it here too would report the same
|
||||
* seat twice.
|
||||
*/
|
||||
@Test
|
||||
void onlyReadyAndDoneSessionsAreReclaimable() {
|
||||
Set<MemberSession.State> expected = EnumSet.of(MemberSession.State.READY, MemberSession.State.DONE);
|
||||
for (MemberSession.State state : MemberSession.State.values()) {
|
||||
MemberSession session = new MemberSession("p1", "term1", "ltms-local", null,
|
||||
"/tmp", null, 0L, 0L, 0, state, null, null);
|
||||
assertEquals(expected.contains(state), FleetMcp.reclaimable(session, null),
|
||||
"state " + state + " must " + (expected.contains(state) ? "" : "not ")
|
||||
+ "count as reclaimable");
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #284, the operator-visible half. {@code liveCount} here is the value the real counter
|
||||
* ({@code Fleetd.liveSessionCount}, proven against the actual spawn gate in
|
||||
* {@code FleetdBackendErrorSinkTest}) produces for this roster: 0, because both sessions are
|
||||
* terminal. What this test pins is what {@code fleet_list} says around it — the two seats show
|
||||
* up once, in {@code free}, and are NOT counted a second time as {@code reclaimable}; neither
|
||||
* is any member row; and both dead sessions are still listed so the lead can see why.
|
||||
*/
|
||||
@Test
|
||||
void terminalFailureSessionsFreeTheirSeatWithoutBeingCountedReclaimable() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
|
||||
MemberSession backendError = sessions.acquire("ltms-local", null, null, null);
|
||||
MemberSession failed = sessions.acquire("ltms-local", null, null, null);
|
||||
assertTrue(sessions.onBackendError(backendError.terminalId(), "backend exited"));
|
||||
sessions.onTurnFailed(failed.terminalId());
|
||||
|
||||
String out = textOf(FleetMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")),
|
||||
sessions, null, new FleetMcp.CapacitySource(profile -> 0, profile -> 2,
|
||||
() -> Set.of("ltms-local"), () -> 0), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), Map.of(), ""));
|
||||
|
||||
assertTrue(out.contains("\"free\":2"), "both seats are back: " + out);
|
||||
assertTrue(out.contains("\"reclaimable\":0"),
|
||||
"the freed seats must not be counted a second time as reclaimable: " + out);
|
||||
assertFalse(out.contains("\"reclaimable\":true"),
|
||||
"no member row may claim a seat the profile count says is not held: " + out);
|
||||
assertTrue(out.contains("\"state\":\"backend_error\""), "the dead session stays visible: " + out);
|
||||
assertTrue(out.contains("\"state\":\"failed\""), "the failed session stays visible: " + out);
|
||||
}
|
||||
|
||||
@Test
|
||||
void inertCapacitySourceOmitsCapacityBlock() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
@@ -810,13 +866,16 @@ class FleetMcpTest {
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #176: this is the exact shape measured on the Mac fleet — {@code maxLoad:3, live:2},
|
||||
* where one of the "free" three is really the lead's own seat on the same subscription. The old
|
||||
* formula ({@code max(0, cap - live)}) reported {@code free:1}; the real ceiling is {@code 0}
|
||||
* (two members plus the lead's own seat already fill all three).
|
||||
* fleetd #257: this is the exact shape measured on the Mac fleet — {@code maxLoad:3, live:2},
|
||||
* one lead sharing the subscription. The OLD formula ({@code max(0, cap - live - leadSeats)})
|
||||
* reported {@code free:0} here, disagreeing with the real spawn gate (which never read
|
||||
* {@code leadSeats} and would still grant one more spawn — see
|
||||
* {@code freeMatchesWhatTheRealPlacementGateActuallyGrants} below for that proof against the
|
||||
* actual gate). {@code free} must report {@code 1}: {@code leadSeats} is carried as a fact, not
|
||||
* subtracted.
|
||||
*/
|
||||
@Test
|
||||
void leadSeatSubtractsFromFreeTheSameWayLiveDoes() {
|
||||
void leadSeatIsReportedButNeverSubtractedFromFree() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
|
||||
FleetMcp.LeadSeatSource leadSeats = new FleetMcp.LeadSeatSource(
|
||||
@@ -829,17 +888,17 @@ class FleetMcpTest {
|
||||
|
||||
assertTrue(out.contains("\"maxLoad\":3"), "maxLoad itself must be left untouched: " + out);
|
||||
assertTrue(out.contains("\"live\":2"), out);
|
||||
assertTrue(out.contains("\"free\":0"), "2 live + 1 lead seat fills all 3: " + out);
|
||||
assertTrue(out.contains("\"leadSeats\":1"), out);
|
||||
assertTrue(out.contains("\"free\":1"), "free is maxLoad - live only, never minus leadSeats: " + out);
|
||||
assertTrue(out.contains("\"leadSeats\":1"), "leadSeats is still reported, just not subtracted: " + out);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #176: the OTHER measurement in the issue — a completely idle fleet still overstates
|
||||
* {@code free} by the lead's own seat. {@code maxLoad:3, live:0} must report {@code free:2}, the
|
||||
* real fan-out ceiling, not {@code 3}.
|
||||
* fleetd #257: the OTHER measurement in the issue — a completely idle fleet with a lead sharing
|
||||
* the subscription. {@code maxLoad:3, live:0} must report {@code free:3}, matching what the
|
||||
* spawn gate (which has no notion of a lead's seat) would actually grant.
|
||||
*/
|
||||
@Test
|
||||
void leadSeatLowersFreeOnAnOtherwiseIdleSubscriptionProfile() {
|
||||
void leadSeatDoesNotLowerFreeOnAnOtherwiseIdleSubscriptionProfile() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
|
||||
FleetMcp.LeadSeatSource leadSeats = new FleetMcp.LeadSeatSource(
|
||||
@@ -851,10 +910,76 @@ class FleetMcpTest {
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), leadSeats, Map.of(), "", null));
|
||||
|
||||
assertTrue(out.contains("\"live\":0"), out);
|
||||
assertTrue(out.contains("\"free\":2"), "an idle fleet's real ceiling is 3 minus the lead's own seat: " + out);
|
||||
assertTrue(out.contains("\"free\":3"), "the real gate never subtracts the lead's seat: " + out);
|
||||
assertTrue(out.contains("\"leadSeats\":1"), out);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #257 — the defect, driven against the REAL placement gate, not a copy of its
|
||||
* arithmetic. {@code free} must equal exactly how many more spawns
|
||||
* {@link CompositePeerLauncher#spawn} (routed through the same {@link SessionManager} fleet_spawn
|
||||
* itself uses) will grant on this profile right now: this test reads whatever number
|
||||
* {@code fleet_list}'s real {@code listFleet} call reports, then drives that many spawns through
|
||||
* the REAL composite launcher and asserts every one succeeds, and the next one — one past what
|
||||
* fleet_list promised — is refused. A test that instead hand-computed {@code cap - live} and
|
||||
* compared it to {@code free} would pass even if both sides shared the same wrong formula (this
|
||||
* repo has been bitten by exactly that before); this one only passes when fleet_list's number and
|
||||
* the gate's real behaviour actually agree.
|
||||
*
|
||||
* <p>{@code liveCount} here is wired the same way {@code Fleetd.main} wires it in production: one
|
||||
* function, read by both the {@link CompositePeerLauncher}'s {@code maxLoad} gate and
|
||||
* {@code fleet_list}'s {@code CapacitySource}, off the SAME {@link SessionManager#roster()} — so
|
||||
* the two paths cannot silently drift on what "live" means.
|
||||
*/
|
||||
@Test
|
||||
void freeMatchesWhatTheRealPlacementGateActuallyGrants() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
FleetConfig.Profile wcfg = new FleetConfig.Profile(
|
||||
"sonnet", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN", null,
|
||||
"tab", "fleetd-workers", "worker: {profile} #{n}", null,
|
||||
null, null, null, null, null, null, null, 3);
|
||||
Map<String, FleetConfig.Profile> profiles = Map.of(wcfg.profile(), wcfg);
|
||||
ClaudeCodeLauncher delegate = new ClaudeCodeLauncher(
|
||||
new AgentControl(h), new WorkspaceControl(h), new SubscriptionGuard(Set.of("gx00.gw")),
|
||||
profiles, wcfg.profile(), k -> "FLEETD_WORKER_TOKEN".equals(k) ? "tok" : null);
|
||||
|
||||
java.util.concurrent.atomic.AtomicReference<SessionManager> smRef =
|
||||
new java.util.concurrent.atomic.AtomicReference<>();
|
||||
Function<String, Integer> liveCount = profile -> (int) smRef.get().roster().stream()
|
||||
.filter(s -> profile.equals(s.profile())).count();
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(
|
||||
List.of(delegate), wcfg.profile(), profiles, PlacementPolicies.fixed(), liveCount);
|
||||
SessionManager sm = new SessionManager(composite);
|
||||
smRef.set(sm);
|
||||
|
||||
// Two members already live — the exact shape measured in fleetd #257 (maxLoad:3, live:2).
|
||||
assertFalse(FleetMcp.spawn(sm, "sonnet").isError(), "setup: first live member must spawn cleanly");
|
||||
assertFalse(FleetMcp.spawn(sm, "sonnet").isError(), "setup: second live member must spawn cleanly");
|
||||
|
||||
// A lead session shares this profile's subscription — leadSeats:1, same as the ticket.
|
||||
FleetMcp.LeadSeatSource leadSeats = new FleetMcp.LeadSeatSource(p -> "sonnet".equals(p) ? 1 : 0);
|
||||
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));
|
||||
int reportedFree = extractInt(out, "free");
|
||||
|
||||
for (int i = 0; i < reportedFree; i++) {
|
||||
McpSchema.CallToolResult res = FleetMcp.spawn(sm, "sonnet");
|
||||
assertFalse(res.isError(), "fleet_list promised free:" + reportedFree + "; spawn #" + (i + 1)
|
||||
+ " of that many was refused by the real gate: " + textOf(res));
|
||||
}
|
||||
McpSchema.CallToolResult overflow = FleetMcp.spawn(sm, "sonnet");
|
||||
assertTrue(overflow.isError(), "fleet_list reported free:" + reportedFree
|
||||
+ " but the real placement gate granted at least one more spawn than that: " + textOf(overflow));
|
||||
}
|
||||
|
||||
private static int extractInt(String json, String key) {
|
||||
java.util.regex.Matcher m = java.util.regex.Pattern.compile("\"" + key + "\":(-?\\d+)").matcher(json);
|
||||
assertTrue(m.find(), "no \"" + key + "\" field in: " + json);
|
||||
return Integer.parseInt(m.group(1));
|
||||
}
|
||||
|
||||
/** A profile with no lead seats reported must be byte-identical to before this ticket. */
|
||||
@Test
|
||||
void zeroLeadSeatsOmitsTheKeyAndLeavesFreeUnchanged() {
|
||||
|
||||
@@ -26,6 +26,7 @@ import org.junit.jupiter.api.io.TempDir;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.io.UncheckedIOException;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.nio.file.attribute.PosixFileAttributeView;
|
||||
@@ -39,6 +40,7 @@ import java.util.UUID;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.Supplier;
|
||||
@@ -923,6 +925,71 @@ class ClaudeCodeLauncherTest {
|
||||
assertTrue(herdr.called("pane.close"), "stop via handle.id() must close the pane");
|
||||
}
|
||||
|
||||
/** A tab-placement launcher with {@code memberCredentials policy=allow-list} under a zsh shell — the
|
||||
* combination that makes {@link HerdrPeerLauncher#spawn} generate a real ZDOTDIR, so {@code
|
||||
* releaseZdotdir}'s effect (the directory's deletion) is observable from a test. */
|
||||
private ClaudeCodeLauncher serviceWithAllowList(FakeHerdr herdr) {
|
||||
FleetConfig.Profile cfg = new FleetConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
|
||||
List.of("ccs", "ltms-local"), "tab", "fleetd-workers",
|
||||
"worker: {profile} #{n}", null, null, null);
|
||||
Supplier<FleetConfig.MemberCredentials> creds = () -> new FleetConfig.MemberCredentials(
|
||||
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST, List.of(), List.of(), null);
|
||||
Function<String, String> env = name -> "SHELL".equals(name) ? "/bin/zsh" : null;
|
||||
return new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
|
||||
env, 0, System::currentTimeMillis, () -> { }, null, creds);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #293: {@code stop()} used to run {@code spaces.closeTab} bare — any non-{@code
|
||||
* *_not_found} herdr error propagated straight out of {@code stop()}, skipping {@code
|
||||
* releaseZdotdir} entirely (the pane was already closed by that point, so the tab-close failure
|
||||
* is cosmetic, not a real teardown failure). Proves both halves of the fix: {@code stop()} no
|
||||
* longer throws for this failure, and {@code releaseZdotdir} still runs — observed here by the
|
||||
* generated ZDOTDIR actually being deleted, since {@code releaseZdotdir}'s last line is {@code
|
||||
* EnvAllowListScrub.deleteRecursively(dir)}.
|
||||
*/
|
||||
@Test
|
||||
@SuppressWarnings("unchecked")
|
||||
void stopStillReleasesZdotdirWhenCloseTabFailsWithANonNotFoundCode() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
ClaudeCodeLauncher svc = serviceWithAllowList(herdr);
|
||||
PeerHandle handle = svc.spawn(new SpawnRequest(null, null, null));
|
||||
Map<String, Object> tabCreateParams = (Map<String, Object>) herdr.lastCall("tab.create").params();
|
||||
Map<String, String> tabEnv = (Map<String, String>) tabCreateParams.get("env");
|
||||
String zdotdir = tabEnv.get("ZDOTDIR");
|
||||
assertNotNull(zdotdir, "policy=allow-list under a zsh shell must have generated a ZDOTDIR: " + tabEnv);
|
||||
Path dir = Path.of(zdotdir);
|
||||
assertTrue(Files.isDirectory(dir), "the generated ZDOTDIR must exist before stop(): " + dir);
|
||||
herdr.tabCloseFailsForTab("w9:t2", "internal_error");
|
||||
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
assertDoesNotThrow(() -> svc.stop(handle.id()),
|
||||
"fleetd #293: a failing tab.close is cosmetic — it must not propagate out of stop()");
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
}
|
||||
|
||||
assertTrue(herdr.called("tab.close"), "tab.close was still attempted");
|
||||
assertFalse(Files.exists(dir),
|
||||
"releaseZdotdir must still run and delete the generated ZDOTDIR despite the tab.close "
|
||||
+ "failure: " + dir);
|
||||
String warn = appender.list.stream()
|
||||
.filter(e -> e.getLevel().equals(Level.WARN))
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.filter(m -> m.contains("tab.close") && m.contains("w9:t2"))
|
||||
.findFirst()
|
||||
.orElse(null);
|
||||
assertNotNull(warn, "the failing tab.close must be logged at WARN naming the tab id — a "
|
||||
+ "silently swallowed failure with no message is not an improvement. Log lines: "
|
||||
+ appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
|
||||
}
|
||||
|
||||
// --- CB-519: host-unique id, decoupled from the pane coordinate ------------------------------
|
||||
|
||||
@Test
|
||||
@@ -2594,4 +2661,272 @@ class ClaudeCodeLauncherTest {
|
||||
+ tornRead.get());
|
||||
assertEquals(newContent, Files.readString(target), "the final content must be the new content");
|
||||
}
|
||||
|
||||
// --- fleetd #247: CAS against a writer TRUST_JSON_LOCK cannot reach --------------------------
|
||||
//
|
||||
// fleetd #149's lock only serialises seedTrustDialog calls THIS launcher makes inside this one
|
||||
// JVM. It does nothing about the one writer that actually shares this file on a real host: the
|
||||
// operator's own live Claude Code, whose CLAUDE_CONFIG_DIR is routinely the very configDir this
|
||||
// profile is given. A plain read-modify-write there is a routine lost update — fleetd reads v1,
|
||||
// the operator's session writes v2, fleetd's ATOMIC_MOVE lands v3 built from v1, and v2 is gone,
|
||||
// atomically. The fix re-reads the file's exact bytes immediately before the move and compares
|
||||
// them with what the update was built from, retrying from fresh bytes on a mismatch.
|
||||
//
|
||||
// Both tests below drive the race through the real seedTrustDialog/spawn() path (not a
|
||||
// hand-rolled call to some extracted primitive) using trustJsonCasTestHook — a seam fired once
|
||||
// per CAS attempt, at the exact point between the read and the final compare, so the race is
|
||||
// deterministic instead of depending on real thread timing.
|
||||
|
||||
@Test
|
||||
void seedTrustDialogRetriesAndPreservesAConcurrentExternalWritersChange(
|
||||
@TempDir Path configDir, @TempDir Path worktree) throws Exception {
|
||||
markAsProvisionedWorktree(worktree);
|
||||
Path claudeJson = configDir.resolve(".claude.json");
|
||||
Files.writeString(claudeJson, "{\"projects\":{}}");
|
||||
|
||||
// Fires exactly once, on the first CAS attempt — simulating the operator's own live Claude
|
||||
// Code landing its own write to this SAME file in the gap between fleetd's read and write.
|
||||
AtomicBoolean fired = new AtomicBoolean(false);
|
||||
ClaudeCodeLauncher.trustJsonCasTestHook = () -> {
|
||||
if (fired.compareAndSet(false, true)) {
|
||||
try {
|
||||
Files.writeString(claudeJson,
|
||||
"{\"projects\":{\"/operator/own/project\":"
|
||||
+ "{\"hasTrustDialogAccepted\":true}}}");
|
||||
} catch (IOException e) {
|
||||
throw new UncheckedIOException(e);
|
||||
}
|
||||
}
|
||||
};
|
||||
try {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FleetConfig.Profile cfg = trustProfile(configDir.toString(), worktree.toString());
|
||||
new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null)
|
||||
.spawn();
|
||||
} finally {
|
||||
ClaudeCodeLauncher.trustJsonCasTestHook = () -> {};
|
||||
}
|
||||
|
||||
assertTrue(fired.get(), "the race hook must actually have fired during the spawn");
|
||||
JsonNode root = new ObjectMapper().readTree(claudeJson.toFile());
|
||||
assertTrue(root.path("projects").path("/operator/own/project")
|
||||
.path("hasTrustDialogAccepted").asBoolean(false),
|
||||
"the external writer's change, landed between fleetd's read and write, must SURVIVE "
|
||||
+ "— this is the whole point of the CAS: without it, fleetd's stale-built "
|
||||
+ "ATOMIC_MOVE would have silently discarded it. Final file: "
|
||||
+ Files.readString(claudeJson));
|
||||
assertTrue(root.path("projects").path(worktree.toString())
|
||||
.path("hasTrustDialogAccepted").asBoolean(false),
|
||||
"fleetd's own retry must still land its own trust entry, built from the fresh bytes");
|
||||
}
|
||||
|
||||
/**
|
||||
* Retry-exhaustion path: an external writer that changes the file on EVERY attempt (not just
|
||||
* once) exhausts all {@code MAX_TRUST_JSON_CAS_ATTEMPTS} retries. fleetd must then write nothing
|
||||
* at all — not even a partial/best-effort write — and log a WARN naming the file it gave up on.
|
||||
*/
|
||||
@Test
|
||||
void seedTrustDialogWritesNothingAndWarnsWhenCasRetriesAreExhausted(
|
||||
@TempDir Path configDir, @TempDir Path worktree) throws Exception {
|
||||
markAsProvisionedWorktree(worktree);
|
||||
Path claudeJson = configDir.resolve(".claude.json");
|
||||
Files.writeString(claudeJson, "{\"marker\":\"start\"}");
|
||||
|
||||
AtomicInteger hookCalls = new AtomicInteger();
|
||||
ClaudeCodeLauncher.trustJsonCasTestHook = () -> {
|
||||
try {
|
||||
Files.writeString(claudeJson,
|
||||
"{\"marker\":\"race-" + hookCalls.incrementAndGet() + "\"}");
|
||||
} catch (IOException e) {
|
||||
throw new UncheckedIOException(e);
|
||||
}
|
||||
};
|
||||
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(ClaudeCodeLauncher.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FleetConfig.Profile cfg = trustProfile(configDir.toString(), worktree.toString());
|
||||
new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null)
|
||||
.spawn();
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
ClaudeCodeLauncher.trustJsonCasTestHook = () -> {};
|
||||
}
|
||||
|
||||
assertEquals(5, hookCalls.get(), "the race hook must fire exactly once per CAS attempt");
|
||||
String finalContent = Files.readString(claudeJson);
|
||||
assertEquals("{\"marker\":\"race-" + hookCalls.get() + "\"}", finalContent,
|
||||
"the file must be left exactly as the external writer last left it — fleetd must not "
|
||||
+ "have written at all once retries are exhausted");
|
||||
assertFalse(finalContent.contains("hasTrustDialogAccepted"),
|
||||
"the trust entry must never appear — its presence would mean the CAS gave up and "
|
||||
+ "wrote anyway instead of skipping the seed");
|
||||
|
||||
assertTrue(appender.list.stream().anyMatch(e -> e.getLevel() == Level.WARN
|
||||
&& e.getFormattedMessage().contains("gave up seeding workspace-trust")
|
||||
&& e.getFormattedMessage().contains(worktree.toString())
|
||||
&& e.getFormattedMessage().contains(claudeJson.toString())),
|
||||
"a WARN naming both the cwd and the file it gave up on must be logged: "
|
||||
+ appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
|
||||
}
|
||||
|
||||
/**
|
||||
* The loud-default WARN (fleetd #247): a profile with no {@code configDir} targets the
|
||||
* operator's real {@code ~/.claude.json} (redirected here to a {@code @TempDir}), and every time
|
||||
* that happens must be logged, not just detected.
|
||||
*/
|
||||
@Test
|
||||
void seedTrustDialogWarnsEveryTimeItTargetsTheDefaultClaudeJson(
|
||||
@TempDir Path fakeHome, @TempDir Path worktree) throws Exception {
|
||||
markAsProvisionedWorktree(worktree);
|
||||
String originalHome = System.getProperty("user.home");
|
||||
System.setProperty("user.home", fakeHome.toString());
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(ClaudeCodeLauncher.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FleetConfig.Profile cfg = trustProfile(null, worktree.toString());
|
||||
new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null)
|
||||
.spawn();
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
System.setProperty("user.home", originalHome);
|
||||
}
|
||||
|
||||
assertTrue(appender.list.stream().anyMatch(e -> e.getLevel() == Level.WARN
|
||||
&& e.getFormattedMessage().contains("no configDir set")
|
||||
&& e.getFormattedMessage().contains(worktree.toString())),
|
||||
"a WARN naming the cwd must fire when the profile sets no configDir: "
|
||||
+ appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
|
||||
}
|
||||
|
||||
// --- fleetd #285: seedTrustDialog under memberHerdrSocket --------------------------------------
|
||||
//
|
||||
// seedTrustDialog gated only on isProvisionedWorktree(cwd) and, being static, could not see
|
||||
// memberHerdrSocketConfigured() at all — unlike its sibling writeCharterFile one method below,
|
||||
// which already refuses the spawn when it cannot place the charter where a different-uid member
|
||||
// can read it. Under memberHerdrSocket + configDir unset, seedTrustDialog wrote fleetd's OWN
|
||||
// ~/.claude.json (the operator's real file) while believing it was seeding the member's. The fix
|
||||
// makes the method an instance method so it can see memberHerdrSocketConfigured(), and applies
|
||||
// the same "refuse, don't silently write somewhere wrong" rule writeCharterFile already uses.
|
||||
|
||||
@Test
|
||||
void seedTrustDialogUnderMemberHerdrSocketSharesTheFileWithTheConfiguredGroup(
|
||||
@TempDir Path configDir, @TempDir Path worktree, @TempDir Path worktreeRoot) throws Exception {
|
||||
markAsProvisionedWorktree(worktree);
|
||||
String group = currentUserGroup();
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
// trustProfile always sets a bridge mcpUrl, so a reply charter is generated too, which
|
||||
// means writeCharterFile ALSO runs under memberHerdrSocket and needs its own worktreeRoot —
|
||||
// pass one so this test isolates the trust-seed behaviour instead of tripping that refusal.
|
||||
FleetConfig.Profile cfg = trustProfile(configDir.toString(), worktree.toString());
|
||||
serviceWithConfig(herdr, cfg, () -> configWithMemberHerdrSocket(worktreeRoot.toString(), group)).spawn();
|
||||
|
||||
Path claudeJson = configDir.resolve(".claude.json");
|
||||
assertTrue(Files.exists(claudeJson), "still seeded into <configDir>/.claude.json under memberHerdrSocket");
|
||||
JsonNode project = new ObjectMapper().readTree(claudeJson.toFile())
|
||||
.path("projects").path(worktree.toString());
|
||||
assertTrue(project.path("hasTrustDialogAccepted").asBoolean(false));
|
||||
|
||||
assertEquals("rw-r-----", PosixFilePermissions.toString(Files.getPosixFilePermissions(claudeJson)),
|
||||
"under memberHerdrSocket the file must be shared group-readable, mirroring the "
|
||||
+ "charter file's own per-file mode (fleetd #222) — a 0600 file (Claude "
|
||||
+ "Code's own default) is unreadable by the member's different OS user");
|
||||
String actualGroup = Files.getFileAttributeView(claudeJson, PosixFileAttributeView.class)
|
||||
.readAttributes().group().getName();
|
||||
assertEquals(group, actualGroup, "the file must be chgrp'd to the configured worktreeGroup");
|
||||
}
|
||||
|
||||
/**
|
||||
* The bug itself: {@code memberHerdrSocket} configured, {@code configDir} unset. Before the fix
|
||||
* this wrote fleetd's own default {@code ~/.claude.json} (here redirected to {@code fakeHome} so
|
||||
* a reintroduced bug still cannot touch the real operator file); after the fix it must refuse the
|
||||
* spawn instead, naming {@code configDir} as the missing key, before the member is ever started.
|
||||
*/
|
||||
@Test
|
||||
void seedTrustDialogUnderMemberHerdrSocketRefusesWhenConfigDirUnset(
|
||||
@TempDir Path fakeHome, @TempDir Path worktree) throws Exception {
|
||||
markAsProvisionedWorktree(worktree);
|
||||
String originalHome = System.getProperty("user.home");
|
||||
System.setProperty("user.home", fakeHome.toString());
|
||||
try {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FleetConfig.Profile cfg = trustProfile(null, worktree.toString());
|
||||
ClaudeCodeLauncher launcher = serviceWithConfig(herdr, cfg,
|
||||
() -> configWithMemberHerdrSocket(null, "some-group"));
|
||||
|
||||
IllegalStateException ex = assertThrows(IllegalStateException.class, launcher::spawn,
|
||||
"memberHerdrSocket + no configDir must refuse the spawn, not write fleetd's own "
|
||||
+ "default ~/.claude.json");
|
||||
assertTrue(ex.getMessage().contains("configDir"),
|
||||
"the refusal must name the missing config key — got: " + ex.getMessage());
|
||||
assertFalse(Files.exists(fakeHome.resolve(".claude.json")),
|
||||
"nothing may be written to fleetd's own default home — this is the exact fleetd "
|
||||
+ "#285 defect: writing the operator's own home instead of the member's");
|
||||
assertFalse(herdr.called("agent.start"),
|
||||
"the spawn must be refused BEFORE the member is ever started — got calls: " + herdr.calls);
|
||||
} finally {
|
||||
System.setProperty("user.home", originalHome);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code configDir} alone is not enough — without {@code worktreeGroup} the file fleetd writes
|
||||
* stays {@code 0600} and the member's different OS user still cannot read it, so this must also
|
||||
* refuse, naming {@code worktreeGroup} this time.
|
||||
*/
|
||||
@Test
|
||||
void seedTrustDialogUnderMemberHerdrSocketRefusesWhenWorktreeGroupUnset(
|
||||
@TempDir Path configDir, @TempDir Path worktree) throws Exception {
|
||||
markAsProvisionedWorktree(worktree);
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FleetConfig.Profile cfg = trustProfile(configDir.toString(), worktree.toString());
|
||||
ClaudeCodeLauncher launcher = serviceWithConfig(herdr, cfg,
|
||||
() -> configWithMemberHerdrSocket(null, null));
|
||||
|
||||
IllegalStateException ex = assertThrows(IllegalStateException.class, launcher::spawn,
|
||||
"memberHerdrSocket + no worktreeGroup must refuse the spawn, not write an unreadable file");
|
||||
assertTrue(ex.getMessage().contains("worktreeGroup"),
|
||||
"configDir alone is not enough — got: " + ex.getMessage());
|
||||
assertFalse(Files.exists(configDir.resolve(".claude.json")),
|
||||
"nothing may be written when the file cannot be shared with the member's group");
|
||||
assertFalse(herdr.called("agent.start"),
|
||||
"the spawn must be refused BEFORE the member is ever started");
|
||||
}
|
||||
|
||||
/**
|
||||
* Regression proof: with {@code memberHerdrSocket} ABSENT — even given a LIVE, non-null {@code
|
||||
* config} supplier (not merely {@code config == null}, which every other seedTrustDialog test in
|
||||
* this file already exercises) — the write must stay byte-identical to before this fix: existing
|
||||
* {@code 0600} permissions preserved, no chgrp/chmod attempted. This is the proof the ticket asks
|
||||
* for: the path the whole live fleet uses today is unchanged by this fix.
|
||||
*/
|
||||
@Test
|
||||
void seedTrustDialogPreservesExisting0600PermissionsWhenMemberHerdrSocketAbsentEvenWithALiveConfigSupplier(
|
||||
@TempDir Path configDir, @TempDir Path worktree) throws Exception {
|
||||
markAsProvisionedWorktree(worktree);
|
||||
Path claudeJson = configDir.resolve(".claude.json");
|
||||
Files.writeString(claudeJson, "{}");
|
||||
assumeTrue(Files.getFileAttributeView(claudeJson, PosixFileAttributeView.class) != null,
|
||||
"no POSIX permissions on this filesystem — skipping rather than failing");
|
||||
Files.setPosixFilePermissions(claudeJson, PosixFilePermissions.fromString("rw-------"));
|
||||
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FleetConfig.Profile cfg = trustProfile(configDir.toString(), worktree.toString());
|
||||
FleetConfig config = new FleetConfig(null, null, null, Map.of(), null, null, null, null, null,
|
||||
null, null, null, null, null, null, null, null, null, null, null, null, null).withDefaults();
|
||||
serviceWithConfig(herdr, cfg, () -> config).spawn();
|
||||
|
||||
assertEquals("rw-------", PosixFilePermissions.toString(Files.getPosixFilePermissions(claudeJson)),
|
||||
"with memberHerdrSocket absent — even given a live config supplier — the seed must "
|
||||
+ "stay byte-identical to before this fix: no chgrp/chmod attempted");
|
||||
}
|
||||
}
|
||||
|
||||
+81
-3
@@ -170,10 +170,10 @@ class HerdrPeerLauncherAllowListWiringTest {
|
||||
/**
|
||||
* {@code SSH_AUTH_SOCK} is a live ssh-agent handle, not a value — it must stay blocked under
|
||||
* {@code allow-list} even when the operator lists it under {@code allow:}, because {@code
|
||||
* sshAuthSock} defaults to blocked. Governed ONLY by {@code memberCredentials.sshAuthSock}.
|
||||
* sshAgentEnv} defaults to omit. Governed ONLY by {@code memberCredentials.sshAgentEnv}.
|
||||
*/
|
||||
@Test
|
||||
void sshAuthSockStaysBlockedEvenWhenListedInMemberCredentialsAllow() {
|
||||
void sshAgentEnvStaysOmittedEvenWhenListedInMemberCredentialsAllow() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
WiringLauncher launcher = new WiringLauncher(herdr,
|
||||
allowListWithAllow(List.of("SSH_AUTH_SOCK")));
|
||||
@@ -184,7 +184,7 @@ class HerdrPeerLauncherAllowListWiringTest {
|
||||
String scrub = readAll(dir.resolve(EnvAllowListScrub.SCRUB_FILE));
|
||||
assertFalse(scrub.contains("'SSH_AUTH_SOCK'"),
|
||||
"SSH_AUTH_SOCK must not be on the derived allow-list just because the operator put "
|
||||
+ "it under allow: — sshAuthSock is unset here, so it defaults to block");
|
||||
+ "it under allow: — sshAgentEnv is unset here, so it defaults to omit");
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -250,6 +250,57 @@ class HerdrPeerLauncherAllowListWiringTest {
|
||||
+ appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #269 follow-up: the same overclaim the WARN in {@code logCredentialGap} was fixed for,
|
||||
* in the INFO line beside it. With {@code memberHerdrSocket} configured, member panes are routed
|
||||
* to a second herdr whose environment fleetd has no channel to inspect, so the counts come from
|
||||
* fleetd's OWN environment. The bare line "member credentials: allowed 1 of 3" reads as a fact
|
||||
* about the member's pane, and there it is not one.
|
||||
*
|
||||
* <p>#269 reworded four sites and stopped at the sibling below; this pins the pair together so
|
||||
* a future edit cannot fix one and leave the other. Real path: asserted after a real {@link
|
||||
* HerdrPeerLauncher#spawn}, reading the log production actually emits.
|
||||
*/
|
||||
@Test
|
||||
void theAllowedCountLineSaysWhoseEnvironmentItCountedWhenMemberHerdrSocketIsSet(@TempDir Path worktreeRoot)
|
||||
throws IOException {
|
||||
String group = currentUserGroup();
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Set<String> hostEnvNames = Set.of(INJECTED, "SOME_UNRELATED_NAME", "ANOTHER_UNRELATED_NAME");
|
||||
WiringLauncher launcher = new WiringLauncher(herdr, allowList(), "/bin/bash-should-be-ignored",
|
||||
() -> hostEnvNames,
|
||||
() -> configWithMemberHerdrSocketRootAndGroup("/tmp/other-user-herdr.sock", "/bin/zsh",
|
||||
worktreeRoot.toString(), group));
|
||||
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
|
||||
Level original = logger.getLevel();
|
||||
logger.setLevel(Level.INFO);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
logger.setLevel(original);
|
||||
}
|
||||
|
||||
List<String> lines = appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList();
|
||||
String coverage = lines.stream()
|
||||
.filter(l -> l.startsWith("member credentials: allowed "))
|
||||
.findFirst()
|
||||
.orElse(null);
|
||||
assertNotNull(coverage, "the coverage line must still be logged — narrowing the claim must "
|
||||
+ "not silently delete the line: " + lines);
|
||||
assertTrue(coverage.contains("fleetd's OWN environment"),
|
||||
"the line must say whose environment it counted: " + coverage);
|
||||
assertTrue(coverage.contains("NOT the member pane's"),
|
||||
"and must say plainly that it is not the member's: " + coverage);
|
||||
// The counts themselves stay real — narrowing the claim must not turn them into constants.
|
||||
assertTrue(coverage.startsWith("member credentials: allowed 1 of 3"),
|
||||
"the real counts must survive the rewording: " + coverage);
|
||||
}
|
||||
|
||||
/**
|
||||
* Lead-review fix: on a NON-zsh shell no scrub ever runs (bash ignores {@code ZDOTDIR}), so the
|
||||
* "allowed N of M" line — which describes what the scrub does — must not be printed there either.
|
||||
@@ -387,6 +438,33 @@ class HerdrPeerLauncherAllowListWiringTest {
|
||||
+ appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #184 item 5: the unknown-environment WARN must state the honest reason for the
|
||||
* UNKNOWN conclusion — fleetd has no channel to confirm what OS user the second herdr runs
|
||||
* as — and must NOT assert as fact that member panes run under a different OS user just
|
||||
* because {@code memberHerdrSocket} is configured. An operator may point it at a second herdr
|
||||
* running as the SAME user, for pane isolation alone; in that case members DO inherit fleetd's
|
||||
* environment, and asserting otherwise would tell the operator to disregard a real, known gap.
|
||||
*/
|
||||
@Test
|
||||
void unknownEnvironmentWarnStatesUncertaintyNotAnAssertedDifferentUser() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Set<String> hostEnvNames = Set.of("FLEETD_WORKER_TOKEN", "SOME_UNKNOWN_SECRET_TOKEN");
|
||||
WiringLauncher launcher = new WiringLauncher(herdr, allowList(), "/bin/zsh", () -> hostEnvNames,
|
||||
() -> configWithMemberHerdrSocket("/tmp/other-user-herdr.sock", "/bin/zsh"));
|
||||
|
||||
List<String> messages = spawnAndCaptureLogs(launcher);
|
||||
|
||||
assertTrue(messages.stream().anyMatch(m -> m.contains("memberHerdrSocket")
|
||||
&& m.contains("no channel to confirm what OS user that herdr runs as")),
|
||||
"expected the WARN to name the actual uncertainty (no channel to confirm the "
|
||||
+ "herdr's uid), got: " + messages);
|
||||
assertFalse(messages.stream().anyMatch(m -> m.contains("member panes run under a different OS "
|
||||
+ "user than fleetd's own process")),
|
||||
"the WARN must not assert as fact that members run under a different OS user just "
|
||||
+ "because memberHerdrSocket is configured — got: " + messages);
|
||||
}
|
||||
|
||||
/**
|
||||
* Hard constraint: the gap detector must never log an env var VALUE, only its NAME. {@code
|
||||
* SOME_UNKNOWN_SECRET_TOKEN} resolves to a distinctive canary value through the same {@code env}
|
||||
|
||||
@@ -109,11 +109,11 @@ class MemberEnvAllowListTest {
|
||||
/**
|
||||
* {@code SSH_AUTH_SOCK} is a live handle to the operator's ssh-agent, never a value — so it must
|
||||
* stay excluded from the derived set even when the operator lists it under {@code allow:} for an
|
||||
* unrelated reason. It is governed ONLY by {@code memberCredentials.sshAuthSock}, applied
|
||||
* unrelated reason. It is governed ONLY by {@code memberCredentials.sshAgentEnv}, applied
|
||||
* separately by the caller ({@code HerdrPeerLauncher}).
|
||||
*/
|
||||
@Test
|
||||
void sshAuthSockInMemberCredentialsAllowIsStillExcluded() {
|
||||
void sshAgentEnvInMemberCredentialsAllowIsStillExcluded() {
|
||||
Set<String> derived = MemberEnvAllowList.derive(List.of(), Set.of("SSH_AUTH_SOCK", "OTHER_NAME"));
|
||||
|
||||
assertFalse(derived.contains("SSH_AUTH_SOCK"),
|
||||
|
||||
@@ -1440,4 +1440,131 @@ class OpenCodeLauncherTest {
|
||||
"the profile hint must survive the Fleetd-style forwarding hop and reach the real "
|
||||
+ "sink — a lambda forwarder drops it and this must go red");
|
||||
}
|
||||
|
||||
// --- fleetd #267: the #175 check never ran for the ordinary (no-worktree) spawn shape --------
|
||||
|
||||
/**
|
||||
* fleetd #267 acceptance criterion 2, half 1 — a regression guard for the NEW code path only:
|
||||
* a spawn WITH a fleetd-provisioned worktree must keep running the fleetd #175 model check
|
||||
* exactly as before (already proven thoroughly above), and must now ALSO never emit the new
|
||||
* fleetd #267 "cannot run" WARN, since the check is not skipped in this shape. Driven through
|
||||
* the real {@code SessionManager.acquire()}/{@code get()} late-resolve path (fleetd #209),
|
||||
* the same path the existing #175 tests already exercise.
|
||||
*/
|
||||
@Test
|
||||
void aProvisionedWorktreeSpawnRunsTheModelCheckThroughSessionManagerAndNeverLogsTheSkipWarn(
|
||||
@TempDir Path configRoot, @TempDir Path discRoot) throws Exception {
|
||||
String workDir = provisionedWorkDir(configRoot);
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FleetConfig.Profile cfg = opencodeCfg("opencode/nemotron-3-ultra-free", null, null);
|
||||
List<String> exhausted = new ArrayList<>();
|
||||
ExhaustionSink sink = (target, reason, profile) -> exhausted.add(reason);
|
||||
OpenCodeLauncher launcher = serviceWithSink(herdr, configRoot, discRoot, cfg, sink);
|
||||
SessionManager sessions = new SessionManager(launcher);
|
||||
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(OpenCodeLauncher.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
MemberSession acquired = sessions.acquire(cfg.profile(), workDir, null, null);
|
||||
assertNull(acquired.agentSessionId(), "no opencode row yet");
|
||||
|
||||
OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "ses_x", workDir, 1000L,
|
||||
"{\"id\":\"gpt-5.6-sol\",\"providerID\":\"openai\"}");
|
||||
|
||||
Optional<MemberSession> after = sessions.get(acquired.paneId());
|
||||
assertEquals("ses_x", after.get().agentSessionId());
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
}
|
||||
|
||||
assertEquals(1, exhausted.size(),
|
||||
"the mismatch check still runs on the real path with a provisioned worktree: " + exhausted);
|
||||
boolean cannotRunWarn = appender.list.stream()
|
||||
.filter(e -> e.getLevel() == Level.WARN)
|
||||
.anyMatch(e -> e.getFormattedMessage().contains("cannot run"));
|
||||
assertFalse(cannotRunWarn, "a provisioned-worktree spawn must never log the fleetd #267 "
|
||||
+ "'cannot run' WARN — the check ran, it was not skipped: " + appender.list);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #267's central defect, reproduced and fixed: {@code OpenCodeLauncher.SessionAwareHandle
|
||||
* .agentSessionId()} is the ONLY caller of {@code checkModelMatch}, and it sits behind the
|
||||
* fleetd #249 worktree gate — so a plain {@code fleet_spawn} with no {@code worktree:true}
|
||||
* (the ticket's "ordinary, expected shape of the large majority of spawns") never reached
|
||||
* {@code checkModelMatch} at all. A test that called {@code checkModelMatch} directly, or built
|
||||
* a {@link OpenCodeLauncher.SessionAwareHandle}/{@link PeerHandle} in isolation, would have
|
||||
* passed on every single day this gap existed — it never drives {@code agentSessionId()}
|
||||
* through the worktree gate the way production does. This test instead drives the REAL
|
||||
* late-resolve path: {@code SessionManager.acquire()} (which calls {@code handle
|
||||
* .agentSessionId()} to build the very first {@code MemberSession}) and a re-poll via {@code
|
||||
* SessionManager.get()} (fleetd #209's retained-handle mechanism) — the exact sequence a live
|
||||
* pane goes through.
|
||||
*
|
||||
* <p>The fix chosen (see {@code OpenCodeLauncher}'s javadoc on the {@code !worktreeProvisioned}
|
||||
* branch) is the WARN path, not a decoupled check: {@code actualModelForSessionId} can only be
|
||||
* keyed safely by a RESOLVED session id (fleetd #234's fix for exactly this false-positive
|
||||
* risk), and without a provisioned worktree no id can ever be safely resolved (fleetd #249) —
|
||||
* re-deriving "whatever is newest in this shared directory" here would silently reintroduce the
|
||||
* false-positive risk #234 fixed. This test proves both halves: a plausible-looking mismatch
|
||||
* row for the shared, non-provisioned cwd never quarantines anything, AND the new one-time,
|
||||
* per-profile WARN replaces the old total silence.
|
||||
*/
|
||||
@Test
|
||||
void aSpawnWithoutAProvisionedWorktreeNeverRunsTheModelCheckButWarnsOncePerProfile(
|
||||
@TempDir Path configRoot, @TempDir Path discRoot) throws Exception {
|
||||
// Deliberately NOT provisionedWorkDir(...) / markAsProvisionedWorktree(...): a plain
|
||||
// directory with no .git marker — the exact "fleet_spawn with no worktree:" shape fleetd
|
||||
// #267 is about, and the ordinary shape the ticket says most spawns actually take.
|
||||
Path workDir = Files.createDirectories(configRoot.resolve("shared-cwd"));
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FleetConfig.Profile cfg = opencodeCfg("opencode/nemotron-3-ultra-free", null, null);
|
||||
List<String> exhausted = new ArrayList<>();
|
||||
ExhaustionSink sink = (target, reason, profile) -> exhausted.add(reason);
|
||||
OpenCodeLauncher launcher = serviceWithSink(herdr, configRoot, discRoot, cfg, sink);
|
||||
SessionManager sessions = new SessionManager(launcher);
|
||||
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(OpenCodeLauncher.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
// The real production entrypoint: acquire() calls handle.agentSessionId() itself to
|
||||
// build the very first MemberSession, BEFORE any row exists.
|
||||
MemberSession acquired = sessions.acquire(cfg.profile(), workDir.toString(), null, null);
|
||||
assertNull(acquired.agentSessionId(),
|
||||
"still refuses to guess an identity for a shared, non-provisioned cwd (fleetd #249)");
|
||||
|
||||
// A row for this exact (shared) directory appears, running a model that WOULD look like
|
||||
// a mismatch against cfg.model() if fleetd trusted the shared-directory heuristic —
|
||||
// exactly the false-positive shape fleetd #234 fixed for the id-resolved case.
|
||||
OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "ses_sibling", workDir.toString(), 1000L,
|
||||
"{\"id\":\"gpt-5.6-sol\",\"providerID\":\"openai\"}");
|
||||
|
||||
// Re-drive the SAME real late-resolve path (fleetd #209) — repeatedly, to also prove
|
||||
// the new WARN fires at most once per profile, not once per poll.
|
||||
Optional<MemberSession> resolved = sessions.get(acquired.paneId());
|
||||
assertNull(resolved.get().agentSessionId(), "still no identity — the gate never opens");
|
||||
sessions.get(acquired.paneId());
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
}
|
||||
|
||||
assertTrue(exhausted.isEmpty(),
|
||||
"must never quarantine off a shared-directory row it cannot trust as this session's "
|
||||
+ "own — fleetd #234's exact concern, now also for the model check: " + exhausted);
|
||||
|
||||
List<String> skipWarnings = appender.list.stream()
|
||||
.filter(e -> e.getLevel() == Level.WARN)
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.filter(m -> m.contains("cannot run"))
|
||||
.toList();
|
||||
assertEquals(1, skipWarnings.size(),
|
||||
"exactly one 'cannot run' WARN across acquire() + two get() re-polls — the old code "
|
||||
+ "logged NOTHING here, which is the bug this ticket fixes; got: " + skipWarnings);
|
||||
assertTrue(skipWarnings.get(0).contains(cfg.profile()),
|
||||
"the WARN must name the profile, same treatment discoveryUnavailable already gets: "
|
||||
+ skipWarnings.get(0));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -828,6 +828,56 @@ class MessageServiceTest {
|
||||
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
|
||||
}
|
||||
|
||||
// --- fleetd #275: a target torn down FOR GOOD while genuinely ASKING must not orphan --------
|
||||
//
|
||||
// sessions.onRelease (fleet_stop, or the idle reaper) is the one abandon() caller that knows
|
||||
// for certain the target can never resume: its pane is being stopped right now. Unlike the
|
||||
// health-classification caller above (a GONE/NEVER_READY guess, not a teardown it performed),
|
||||
// it must sweep an ASKING ticket right here — see MessageService.abandon(String, String,
|
||||
// boolean)'s javadoc for the full reachability chain this closes: without this, the forward
|
||||
// waiter is already closed by the time the question surfaces, the ASKING guard skips the task,
|
||||
// and by the time the worker's own fleet_ask lapses (~55-115s later) the released session no
|
||||
// longer appears in FleetHealthMonitor's roster for anything to ever sweep it again — leaving
|
||||
// fleet_poll{ticket} stuck PENDING forever.
|
||||
|
||||
@Test
|
||||
void abandonWithSweepAskingFailsATornDownTargetsAskingTicket() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "task that asks");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
CompletableFuture<MessageService.AskResult> ask =
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 300));
|
||||
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
|
||||
|
||||
assertTrue(messages.abandon(T, "the worker session was released before it replied", true),
|
||||
"a released target's open ask can never resume, so it must fail right here");
|
||||
|
||||
MessageService.TaskView failed = awaitTicketPhase(ticket, MessageService.Phase.FAILED);
|
||||
assertEquals("the worker session was released before it replied", failed.detail());
|
||||
|
||||
// The reverse-rendezvous ask is torn down too: the worker's still-blocked fleet_ask rides
|
||||
// out its own timeout (nothing completed its answer future), and a late answer() for the
|
||||
// same turnId must see it as lapsed rather than resolving a question nobody is waiting on.
|
||||
assertEquals(MessageService.AskOutcome.TIMED_OUT, ask.get(5, TimeUnit.SECONDS).outcome());
|
||||
assertEquals(MessageService.Outcome.STALE_TURN,
|
||||
messages.answer(asking.turnId(), "config.yaml", 200).outcome());
|
||||
}
|
||||
|
||||
@Test
|
||||
void abandonWithoutSweepAskingBehavesLikeTheTwoArgOverload() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "task that asks");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
|
||||
awaitTicketPhase(ticket, MessageService.Phase.ASKING);
|
||||
|
||||
assertFalse(messages.abandon(T, "agent target term_a not found", false),
|
||||
"sweepAsking=false must match the plain abandon(target, reason) overload");
|
||||
assertEquals(MessageService.Phase.ASKING, messages.poll(ticket).phase());
|
||||
}
|
||||
|
||||
// --- #137: a fleet_ask round-trip must not orphan the ticket's own reply -------------------
|
||||
//
|
||||
// The primary's fleet_send{turnId} answer call is itself bounded (a real MCP call, capped well
|
||||
@@ -925,6 +975,64 @@ class MessageServiceTest {
|
||||
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #282: a worker that chains a SECOND {@code fleet_ask} inside the same resumed turn —
|
||||
* before it ever calls {@code fleet_reply} — used to kill its own async ticket. {@code answer()}
|
||||
* opens a fresh forward waiter but (unlike {@code send()}) never registered it in
|
||||
* {@code asyncTasksByWaiter}, so the second ask's {@code markAsyncQuestion} found no {@code Task}
|
||||
* to re-associate. That waiter still resolved with the second {@code QUESTION} once the worker
|
||||
* asked again, and {@code answer()} completed the ticket's future with that QUESTION "reply"
|
||||
* unconditionally — so {@code fleet_poll} reported FAILED while the worker was still alive and
|
||||
* the primary was mid-conversation with it.
|
||||
*
|
||||
* <p>Driven entirely through {@code MessageService}'s public API (sendAsync/ask/answer/poll) —
|
||||
* never by reaching into {@link Rendezvous} or the task maps directly, so this test cannot pass
|
||||
* for a reason unrelated to the real bug.
|
||||
*/
|
||||
@Test
|
||||
void secondFleetAskInTheSameResumedTurnDoesNotKillTheAsyncTicket() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "task that asks twice");
|
||||
awaitWaiting();
|
||||
|
||||
// The worker's first fleet_ask.
|
||||
CompletableFuture<MessageService.AskResult> ask1 =
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "Q1", 5000));
|
||||
MessageService.TaskView asking1 = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
|
||||
assertEquals("Q1", asking1.reply());
|
||||
|
||||
// The primary answers it — answer() resumes the turn and blocks for what comes next.
|
||||
CompletableFuture<MessageService.Reply> answer1 = CompletableFuture.supplyAsync(
|
||||
() -> messages.answer(asking1.turnId(), "a1", 5000));
|
||||
assertEquals("a1", ask1.get(5, TimeUnit.SECONDS).answer());
|
||||
|
||||
// Still in the SAME resumed turn — before replying — the worker asks again.
|
||||
CompletableFuture<MessageService.AskResult> ask2 =
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "Q2", 5000));
|
||||
|
||||
// answer1's own call unblocks with the second QUESTION (documented QUESTION-chaining
|
||||
// behaviour — see FleetMcp.answer's javadoc: "Answer it by calling fleet_send again with
|
||||
// turnId=..."). The bug: this used to also kill the async ticket in the process.
|
||||
MessageService.Reply firstAnswerResult = answer1.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(MessageService.Outcome.QUESTION, firstAnswerResult.outcome());
|
||||
String turnId2 = firstAnswerResult.turnId();
|
||||
|
||||
MessageService.TaskView asking2 = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
|
||||
assertEquals("Q2", asking2.reply(),
|
||||
"the ticket must surface the SECOND question, not be dead/FAILED");
|
||||
assertEquals(turnId2, asking2.turnId());
|
||||
|
||||
// The primary answers the second question; the worker finally sends its real fleet_reply.
|
||||
CompletableFuture<MessageService.Reply> answer2 = CompletableFuture.supplyAsync(
|
||||
() -> messages.answer(turnId2, "a2", 5000));
|
||||
assertEquals("a2", ask2.get(5, TimeUnit.SECONDS).answer());
|
||||
awaitWaiting();
|
||||
assertTrue(rendezvous.resolve(T, "done"));
|
||||
assertEquals(MessageService.Outcome.REPLIED, answer2.get(5, TimeUnit.SECONDS).outcome());
|
||||
|
||||
MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE);
|
||||
assertEquals("done", done.reply());
|
||||
}
|
||||
|
||||
// --- CB-582: fleet_status pendingAsk() ------------------------------------------------------
|
||||
|
||||
@Test
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package dev.ltms.fleet.rest;
|
||||
|
||||
import dev.ltms.fleet.auth.CallerResolver;
|
||||
import dev.ltms.fleet.auth.Authz;
|
||||
import dev.ltms.fleet.auth.MemberRegistry;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
@@ -25,8 +26,15 @@ import java.net.URI;
|
||||
import java.net.http.HttpClient;
|
||||
import java.net.http.HttpRequest;
|
||||
import java.net.http.HttpResponse;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.LinkedHashSet;
|
||||
import java.util.Locale;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.regex.Matcher;
|
||||
import java.util.regex.Pattern;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
@@ -36,6 +44,10 @@ import static org.junit.jupiter.api.Assertions.*;
|
||||
*/
|
||||
class FleetAppAuthTest {
|
||||
|
||||
private static final Path REST_SOURCE = Path.of("src/main/java/dev/ltms/fleet/rest/FleetApp.java");
|
||||
private static final Pattern ROUTE_REGISTRATION =
|
||||
Pattern.compile("app\\.(get|post|delete|put|patch)\\(\\s*\"([^\"]+)\"");
|
||||
|
||||
private final HttpClient http = HttpClient.newHttpClient();
|
||||
private Javalin app;
|
||||
private Metrics metrics;
|
||||
@@ -92,6 +104,49 @@ class FleetAppAuthTest {
|
||||
return http.send(b.build(), HttpResponse.BodyHandlers.ofString());
|
||||
}
|
||||
|
||||
@Test
|
||||
void everyRegisteredRouteHasItsHandlerActionPinned() {
|
||||
Set<String> registered = routesTheServerRegisters();
|
||||
assertTrue(registered.size() >= 15,
|
||||
"scraped only " + registered.size() + " route registrations from FleetApp (" + registered
|
||||
+ "); the app.<verb>(\"…\") scrape has stopped matching");
|
||||
// Liveness must work before credentials can be checked, so this route is deliberately open.
|
||||
assertTrue(registered.remove("GET /healthz"), "GET /healthz must stay an explicit ungated exception");
|
||||
registered.forEach(route -> assertDoesNotThrow(() -> FleetApp.routeAction(route),
|
||||
() -> route + " is registered but has no pinned authorization action"));
|
||||
assertEquals(Authz.Action.METRICS, FleetApp.routeAction("GET /metrics"));
|
||||
assertEquals(Authz.Action.SPAWN, FleetApp.routeAction("POST /members"));
|
||||
assertEquals(Authz.Action.STOP, FleetApp.routeAction("DELETE /members/{paneId}"));
|
||||
assertEquals(Authz.Action.SEND, FleetApp.routeAction("POST /sessions/{id}/message"));
|
||||
assertEquals(Authz.Action.REPLY, FleetApp.routeAction("POST /sessions/{id}/reply"));
|
||||
assertEquals(Authz.Action.DRAIN, FleetApp.routeAction("GET /sessions/{id}/replies"));
|
||||
assertEquals(Authz.Action.ASK, FleetApp.routeAction("POST /sessions/{id}/ask"));
|
||||
for (String route : Set.of("GET /sessions", "GET /agents", "GET /members", "GET /profiles",
|
||||
"GET /member-credentials", "GET /sessions/{id}/status", "GET /tasks/{ticket}")) {
|
||||
assertEquals(Authz.Action.READ, FleetApp.routeAction(route), route);
|
||||
}
|
||||
assertThrows(IllegalArgumentException.class, () -> FleetApp.routeAction("GET /healthz"));
|
||||
}
|
||||
|
||||
private static Set<String> routesTheServerRegisters() {
|
||||
try {
|
||||
String source = Files.readString(REST_SOURCE).lines()
|
||||
.filter(line -> {
|
||||
String stripped = line.stripLeading();
|
||||
return !(stripped.startsWith("//") || stripped.startsWith("*") || stripped.startsWith("/*"));
|
||||
})
|
||||
.collect(Collectors.joining("\n"));
|
||||
Matcher matcher = ROUTE_REGISTRATION.matcher(source);
|
||||
Set<String> routes = new LinkedHashSet<>();
|
||||
while (matcher.find()) {
|
||||
routes.add(matcher.group(1).toUpperCase(Locale.ROOT) + " " + matcher.group(2));
|
||||
}
|
||||
return routes;
|
||||
} catch (Exception e) {
|
||||
throw new AssertionError("could not scrape FleetApp route registrations", e);
|
||||
}
|
||||
}
|
||||
|
||||
// --- loopback-trust: the caller is the primary -------------------------------------------
|
||||
|
||||
@Test
|
||||
|
||||
@@ -19,6 +19,10 @@ import dev.ltms.fleet.session.SessionManager;
|
||||
import dev.ltms.fleet.session.Worktrees;
|
||||
import dev.ltms.fleet.member.ClaudeCodeLauncher;
|
||||
import dev.ltms.fleet.member.CompositePeerLauncher;
|
||||
import dev.ltms.fleet.member.MemberCredentialPolicyView;
|
||||
import dev.ltms.fleet.mcp.FleetMcp;
|
||||
import dev.ltms.fleet.placement.BackendOutagePolicy;
|
||||
import dev.ltms.fleet.placement.BackendQuarantine;
|
||||
import dev.ltms.fleet.placement.PlacementPolicies;
|
||||
import io.javalin.Javalin;
|
||||
import org.junit.jupiter.api.AfterEach;
|
||||
@@ -32,6 +36,7 @@ import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.UUID;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.function.Predicate;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
@@ -70,6 +75,19 @@ class FleetAppTest {
|
||||
|
||||
private int start(FakeHerdr herdr, String workerBaseUrl, Set<String> allow, String placement,
|
||||
Worktrees worktrees, Predicate<String> deliverable) {
|
||||
return start(herdr, workerBaseUrl, allow, placement, worktrees, deliverable,
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #297: same wiring as above, plus the two SAME shared sources {@code GET /profiles}
|
||||
* must read — lets a test prove the quarantined/coolingOff facts it reports come from a real
|
||||
* {@link dev.ltms.fleet.placement.BackendQuarantine}/{@link
|
||||
* dev.ltms.fleet.placement.BackendOutagePolicy}, exactly like {@code fleet_profiles}'s own tests.
|
||||
*/
|
||||
private int start(FakeHerdr herdr, String workerBaseUrl, Set<String> allow, String placement,
|
||||
Worktrees worktrees, Predicate<String> deliverable,
|
||||
FleetMcp.QuarantineSource quarantine, FleetMcp.OutageSource outage) {
|
||||
FleetConfig.Profile wcfg = new FleetConfig.Profile(
|
||||
"ltms-local", workerBaseUrl, "coder", null, "FLEETD_WORKER_TOKEN", null,
|
||||
placement, "fleet", "worker: {profile} #{n}", null, null, null);
|
||||
@@ -90,8 +108,9 @@ class FleetAppTest {
|
||||
// it directly so the inbox contract holds for those endpoints.
|
||||
inbox.own("term_a");
|
||||
MessageService messages = new MessageService(agents, injector, rendezvous, inbox);
|
||||
app = new FleetApp(herdr, workers, sessions, messages, this.presence, null,
|
||||
null, null, id -> this.presence.isPresent(id) || deliverable.test(id))
|
||||
app = new FleetApp(herdr, herdr, workers, sessions, messages, this.presence, null,
|
||||
null, null, id -> this.presence.isPresent(id) || deliverable.test(id),
|
||||
MemberCredentialPolicyView::absent, quarantine, outage)
|
||||
.build().start("127.0.0.1", 0);
|
||||
return app.port();
|
||||
}
|
||||
@@ -164,6 +183,22 @@ class FleetAppTest {
|
||||
assertEquals("idle", agents.get(0).get("status").asText());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #297 gap 1: {@code workers.list()} reaches herdr, and a transport failure there must
|
||||
* land in the same {@code {error, detail}} envelope every other failure path in this file uses
|
||||
* (see {@code herdrError}), not escape as a bare exception outside the JSON contract.
|
||||
*/
|
||||
@Test
|
||||
void agentsMapsAHerdrFailureToTheJsonErrorEnvelope() throws Exception {
|
||||
FakeHerdr down = new FakeHerdr().healthy(false);
|
||||
int port = start(down, "http://gx00.gw:8000", Set.of("gx00.gw"));
|
||||
HttpResponse<String> res = req(port, "GET", "/agents");
|
||||
assertEquals(502, res.statusCode(), res.body());
|
||||
JsonNode body = mapper.readTree(res.body());
|
||||
assertEquals("herdr_error", body.get("error").asText());
|
||||
assertTrue(body.has("detail"), res.body());
|
||||
}
|
||||
|
||||
@Test
|
||||
void spawnWorkerLandsInOwnTabInWorkerSpaceAndInjectsBaseUrl() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
@@ -203,6 +238,42 @@ class FleetAppTest {
|
||||
JsonNode body = mapper.readTree(req(port, "GET", "/profiles").body());
|
||||
assertEquals("ltms-local", body.get("default").asText());
|
||||
assertEquals("ltms-local", body.get("profiles").get(0).asText());
|
||||
assertFalse(body.has("quarantined"), "nothing is quarantined, so the key is omitted: " + body);
|
||||
assertFalse(body.has("coolingOff"), "nothing is cooling off, so the key is omitted: " + body);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #297 gap 2: {@code GET /profiles} must report the same two outage states {@code
|
||||
* fleet_profiles} does — CB-578 stage B exhaustion quarantine and fleetd #201 Unit 5 cool-off —
|
||||
* reading the SAME shared {@link BackendQuarantine}/{@link BackendOutagePolicy} instances rather
|
||||
* than recomputing them. The two checks are independent, and this profile is deliberately put in
|
||||
* both states at once, matching {@code FleetMcpTest}'s own coverage of that overlap.
|
||||
*/
|
||||
@Test
|
||||
void profilesReportsQuarantineAndCoolingOffFromTheSameSharedSources() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
|
||||
quarantine.quarantine("shared-openai");
|
||||
FleetMcp.QuarantineSource quarantineSource = new FleetMcp.QuarantineSource(
|
||||
profile -> "ltms-local".equals(profile) ? "shared-openai" : null, quarantine);
|
||||
BackendOutagePolicy outagePolicy = new BackendOutagePolicy(() -> 0L);
|
||||
outagePolicy.record("shared-openai", "t1", "API Error: rate limited");
|
||||
outagePolicy.record("shared-openai", "t2", "API Error: rate limited"); // 2nd distinct target starts the incident
|
||||
FleetMcp.OutageSource outageSource = new FleetMcp.OutageSource(
|
||||
profile -> "ltms-local".equals(profile) ? "shared-openai" : null, outagePolicy);
|
||||
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"), "tab", new GitWorktrees(),
|
||||
ignored -> false, quarantineSource, outageSource);
|
||||
|
||||
JsonNode body = mapper.readTree(req(port, "GET", "/profiles").body());
|
||||
assertTrue(body.has("quarantined"), body.toString());
|
||||
assertEquals("shared-openai",
|
||||
body.get("quarantined").get("ltms-local").get("credentialId").asText());
|
||||
assertEquals(1800,
|
||||
body.get("quarantined").get("ltms-local").get("quarantinedForSeconds").asLong());
|
||||
assertTrue(body.has("coolingOff"), body.toString());
|
||||
assertEquals("shared-openai",
|
||||
body.get("coolingOff").get("ltms-local").get("credentialId").asText());
|
||||
assertEquals(60, body.get("coolingOff").get("ltms-local").get("coolingOffForSeconds").asLong());
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -239,6 +310,23 @@ class FleetAppTest {
|
||||
"liveStatus is unknown when herdr has no matching pane");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #297 gap 1: same reasoning as {@code agentsMapsAHerdrFailureToTheJsonErrorEnvelope} —
|
||||
* {@code GET /members} is the endpoint's own comment names as "the out-of-band path a lead falls
|
||||
* back to when its MCP mount drops", so it must stay inside the {@code {error, detail}} envelope
|
||||
* exactly when herdr is briefly unreachable.
|
||||
*/
|
||||
@Test
|
||||
void membersMapsAHerdrFailureToTheJsonErrorEnvelope() throws Exception {
|
||||
FakeHerdr down = new FakeHerdr().healthy(false);
|
||||
int port = start(down, "http://gx00.gw:8000", Set.of("gx00.gw"));
|
||||
HttpResponse<String> res = req(port, "GET", "/members");
|
||||
assertEquals(502, res.statusCode(), res.body());
|
||||
JsonNode body = mapper.readTree(res.body());
|
||||
assertEquals("herdr_error", body.get("error").asText());
|
||||
assertTrue(body.has("detail"), res.body());
|
||||
}
|
||||
|
||||
@Test
|
||||
void spawnWithACwdParamRootsTheWorkerThere() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
|
||||
@@ -7,6 +7,7 @@ import java.util.Locale;
|
||||
import java.util.Set;
|
||||
import java.util.regex.Matcher;
|
||||
import java.util.regex.Pattern;
|
||||
import java.util.stream.Collectors;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
@@ -68,6 +69,32 @@ class RestRouteInventoryTest {
|
||||
"GET /tasks/{ticket}"
|
||||
);
|
||||
|
||||
private static final Pattern ROUTE_CALL =
|
||||
Pattern.compile("app\\.(get|post|delete|put|patch)\\(\\s*\"([^\"]+)\"");
|
||||
|
||||
/**
|
||||
* Drop whole-line comments before scraping. Without this the scrape reads commented-out code as
|
||||
* live: a registration disabled with {@code //} still matched, so the route stayed in the
|
||||
* inventory while the server no longer served it — a silent false PASS, measured on 2026-09-03
|
||||
* by commenting out {@code app.get("/tasks/{ticket}", ...)} and watching this test stay green.
|
||||
* Deleting the same line was caught correctly, so only the commented-out shape was blind.
|
||||
*
|
||||
* <p>Only lines whose first non-blank characters are {@code //}, {@code *} or {@code /*} are
|
||||
* dropped — deliberately NOT every {@code //} anywhere on a line, because that would also cut a
|
||||
* string literal containing {@code //} (a URL) and could silently delete a real registration
|
||||
* sharing that line. The remaining gap is a trailing comment on the same line as real code; no
|
||||
* registration in this file has that shape, and the vacuity test below would catch a scrape that
|
||||
* lost registrations wholesale.
|
||||
*/
|
||||
private static String withoutCommentLines(String source) {
|
||||
return source.lines()
|
||||
.filter(line -> {
|
||||
String t = line.stripLeading();
|
||||
return !(t.startsWith("//") || t.startsWith("*") || t.startsWith("/*"));
|
||||
})
|
||||
.collect(Collectors.joining("\n"));
|
||||
}
|
||||
|
||||
/**
|
||||
* Every {@code app.<verb>("<path>")} call in {@link FleetApp}'s source, as {@code "VERB path"}.
|
||||
* This matches inside an {@code if (...) { ... }} block just as well as a top-level statement —
|
||||
@@ -75,8 +102,7 @@ class RestRouteInventoryTest {
|
||||
* catches {@code GET /metrics} (registered conditionally on {@code metrics != null}).
|
||||
*/
|
||||
private static Set<String> routesTheServerRegisters() throws Exception {
|
||||
String source = Files.readString(REST_SOURCE);
|
||||
Matcher m = Pattern.compile("app\\.(get|post|delete|put|patch)\\(\\s*\"([^\"]+)\"").matcher(source);
|
||||
Matcher m = ROUTE_CALL.matcher(withoutCommentLines(Files.readString(REST_SOURCE)));
|
||||
Set<String> found = new LinkedHashSet<>();
|
||||
while (m.find()) {
|
||||
found.add(m.group(1).toUpperCase(Locale.ROOT) + " " + m.group(2));
|
||||
|
||||
@@ -17,6 +17,9 @@ public final class FakeWorktrees implements Worktrees {
|
||||
public record RemoveCall(String repoRoot, String worktreePath) {
|
||||
}
|
||||
|
||||
public record DeleteBranchCall(String repoRoot, String branch) {
|
||||
}
|
||||
|
||||
public record OverlayCall(String repoRoot, String worktreePath,
|
||||
List<String> requested, List<String> copied, List<String> skipWorktree) {
|
||||
}
|
||||
@@ -35,6 +38,7 @@ public final class FakeWorktrees implements Worktrees {
|
||||
|
||||
private final List<AddCall> addCalls = new CopyOnWriteArrayList<>();
|
||||
private final List<RemoveCall> removeCalls = new CopyOnWriteArrayList<>();
|
||||
private final List<DeleteBranchCall> deleteBranchCalls = new CopyOnWriteArrayList<>();
|
||||
private final List<OverlayCall> overlayCalls = new CopyOnWriteArrayList<>();
|
||||
private final List<RepoRootCall> repoRootCalls = new CopyOnWriteArrayList<>();
|
||||
private final List<SnapshotCall> snapshotCalls = new CopyOnWriteArrayList<>();
|
||||
@@ -45,9 +49,14 @@ public final class FakeWorktrees implements Worktrees {
|
||||
private final List<String> overlayShareOrder = new CopyOnWriteArrayList<>();
|
||||
private final Set<String> existingPaths = ConcurrentHashMap.newKeySet();
|
||||
private final Set<String> trackedPaths = ConcurrentHashMap.newKeySet();
|
||||
/** Worktree paths that currently exist, mirroring GitWorktrees' {@code Files.exists} check for
|
||||
* the already-gone case (CB-576 review, fleetd #116). */
|
||||
private final Set<String> worktreePaths = ConcurrentHashMap.newKeySet();
|
||||
private final AtomicLong snapshotSeq = new AtomicLong();
|
||||
private volatile RuntimeException addFailure;
|
||||
private volatile RuntimeException snapshotFailure;
|
||||
private volatile RuntimeException removeFailure;
|
||||
private volatile RuntimeException overlayFailure;
|
||||
private volatile boolean dirty = false;
|
||||
private volatile String repoRoot = "/repo";
|
||||
private volatile String prefix = "/worktrees";
|
||||
@@ -95,6 +104,21 @@ public final class FakeWorktrees implements Worktrees {
|
||||
return this;
|
||||
}
|
||||
|
||||
/** Make subsequent {@link #remove} calls throw (fleetd #283: a stale index lock, a slow
|
||||
* filesystem, or {@code remove}'s own 30s exec timeout escaping the last, previously bare,
|
||||
* step of {@link SessionManager#release}). */
|
||||
public FakeWorktrees failRemove(String message) {
|
||||
this.removeFailure = new WorktreeException(message);
|
||||
return this;
|
||||
}
|
||||
|
||||
/** Make subsequent {@link #overlayParity} calls throw (simulates a post-{@code add()} spawn
|
||||
* failure — fleetd #283 defect 2 — so the {@code acquireWithWorktree} catch runs). */
|
||||
public FakeWorktrees failOverlay(String message) {
|
||||
this.overlayFailure = new WorktreeException(message);
|
||||
return this;
|
||||
}
|
||||
|
||||
/** Configure the value returned by {@link #wipRefs}. */
|
||||
public FakeWorktrees withWipRefs(WipRefStats stats) {
|
||||
this.wipRefs = stats;
|
||||
@@ -115,21 +139,45 @@ public final class FakeWorktrees implements Worktrees {
|
||||
}
|
||||
// The branch already carries a unique nonce, so the derived path is distinct per acquire
|
||||
// without an extra counter — keep it a pure function of the branch the test can predict.
|
||||
return prefix + "/" + branch.replace('/', '_');
|
||||
String path = prefix + "/" + branch.replace('/', '_');
|
||||
worktreePaths.add(path);
|
||||
return path;
|
||||
}
|
||||
|
||||
/** Model an operator / {@code git worktree prune} removing the worktree before release. */
|
||||
public FakeWorktrees markGone(String worktreePath) {
|
||||
worktreePaths.remove(worktreePath);
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void remove(String repoRoot, String worktreePath) {
|
||||
removeCalls.add(new RemoveCall(repoRoot, worktreePath));
|
||||
if (removeFailure != null) {
|
||||
throw removeFailure;
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void deleteBranch(String repoRoot, String branch) {
|
||||
deleteBranchCalls.add(new DeleteBranchCall(repoRoot, branch));
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean hasUncommitted(String worktreePath) {
|
||||
// A path that was never added, or was marked gone, is reported clean, mirroring
|
||||
// GitWorktrees' already-gone guard — never an error, so teardown still completes.
|
||||
if (!worktreePaths.contains(worktreePath)) {
|
||||
return false;
|
||||
}
|
||||
return dirty;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void overlayParity(String repoRoot, String worktreePath, List<String> overlay) {
|
||||
if (overlayFailure != null) {
|
||||
throw overlayFailure;
|
||||
}
|
||||
List<String> copied = new java.util.ArrayList<>();
|
||||
List<String> skipped = new java.util.ArrayList<>();
|
||||
for (String rel : overlay) {
|
||||
@@ -202,6 +250,14 @@ public final class FakeWorktrees implements Worktrees {
|
||||
return removeCalls.isEmpty() ? null : removeCalls.getLast();
|
||||
}
|
||||
|
||||
public List<DeleteBranchCall> deleteBranchCalls() {
|
||||
return List.copyOf(deleteBranchCalls);
|
||||
}
|
||||
|
||||
public DeleteBranchCall lastDeleteBranch() {
|
||||
return deleteBranchCalls.isEmpty() ? null : deleteBranchCalls.getLast();
|
||||
}
|
||||
|
||||
public OverlayCall lastOverlay() {
|
||||
return overlayCalls.isEmpty() ? null : overlayCalls.getLast();
|
||||
}
|
||||
|
||||
@@ -21,6 +21,7 @@ import java.util.Map;
|
||||
import java.util.Optional;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
@@ -362,6 +363,46 @@ class GitWorktreesTest {
|
||||
assertEquals("worktree origin contains HTTPS user info; refusing provision", error.getMessage());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #274. {@code add()} creates the worktree and its branch, then runs several more steps
|
||||
* that can throw — {@code requireCredentialFreeHttpsOrigin} among them, an intended security
|
||||
* refusal, not an IO accident. Before the fix, any exception from those later steps left
|
||||
* {@code add()} never returning, so its caller never learned the path and the worktree
|
||||
* directory plus its branch leaked on disk forever with nothing tracking them.
|
||||
*
|
||||
* <p>This drives the exact same {@code afterWorktreeAdded} test seam as
|
||||
* {@link #provisioningRefusesAWorktreeWhoseOriginStillHasHttpsUserInfo} — a mutation applied
|
||||
* right after {@code git worktree add}, so the step that throws
|
||||
* ({@code requireCredentialFreeHttpsOrigin}, reached moments later inside {@code add()} itself)
|
||||
* runs strictly after the worktree and branch already exist, not downstream of {@code add()}
|
||||
* in some other caller. {@code afterWorktreeAdded} also hands back the created path, so the
|
||||
* assertions below don't have to guess the generated nonce.
|
||||
*/
|
||||
@Test
|
||||
void addCleansUpTheWorktreeAndBranchWhenAPostCreationStepThrows(@TempDir Path tmp) throws Exception {
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
git(repo, "remote", "add", "origin", "https://git.ltms.dev/akb/kb.git");
|
||||
String branch = "cb-274-leak";
|
||||
AtomicReference<String> createdPath = new AtomicReference<>();
|
||||
GitWorktrees worktrees = new GitWorktrees(tmp.resolve("wts").toString(), worktreePath -> {
|
||||
createdPath.set(worktreePath);
|
||||
try {
|
||||
git(Path.of(worktreePath), "remote", "set-url", "origin",
|
||||
"https://synthetic-test-token@git.ltms.dev/akb/kb.git");
|
||||
} catch (Exception e) {
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
});
|
||||
|
||||
assertThrows(WorktreeException.class, () -> worktrees.add(repo.toString(), branch, "HEAD"));
|
||||
|
||||
assertNotNull(createdPath.get(), "afterWorktreeAdded must have run with the created path");
|
||||
assertFalse(Files.exists(Path.of(createdPath.get())),
|
||||
"the worktree directory leaked after a post-creation step threw");
|
||||
String heads = forEachRef(repo, "refs/heads/" + branch);
|
||||
assertTrue(heads.isBlank(), "the branch leaked after a post-creation step threw:\n" + heads);
|
||||
}
|
||||
|
||||
// ---- CB-189: broader remote-URL coverage — every remote, both fetch and push URLs, any
|
||||
// non-SSH scheme. Reporting only, additive to the origin/https strip-and-refuse tests above. ----
|
||||
|
||||
|
||||
@@ -74,6 +74,7 @@ class SessionManagerTest {
|
||||
*/
|
||||
private static final class RecordingWorktrees implements Worktrees {
|
||||
private final List<String> removeCalls = new java.util.ArrayList<>();
|
||||
private final List<String> deleteBranchCalls = new java.util.ArrayList<>();
|
||||
private final List<String> snapshotCalls = new java.util.ArrayList<>();
|
||||
private final java.util.Set<String> failRemoveFor = new java.util.HashSet<>();
|
||||
private volatile boolean dirty = false;
|
||||
@@ -114,6 +115,11 @@ class SessionManagerTest {
|
||||
removeCalls.add(worktreePath);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void deleteBranch(String repoRoot, String branch) {
|
||||
deleteBranchCalls.add(branch);
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean hasUncommitted(String worktreePath) {
|
||||
if (hasUncommittedFailure != null) {
|
||||
@@ -158,6 +164,10 @@ class SessionManagerTest {
|
||||
return List.copyOf(removeCalls);
|
||||
}
|
||||
|
||||
List<String> deleteBranchCalls() {
|
||||
return List.copyOf(deleteBranchCalls);
|
||||
}
|
||||
|
||||
List<String> snapshotCalls() {
|
||||
return List.copyOf(snapshotCalls);
|
||||
}
|
||||
@@ -989,8 +999,21 @@ class SessionManagerTest {
|
||||
+ "dirty check threw");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #283 defect 1 changed this test's own premise, so its assertions are updated along
|
||||
* with the production fix. Before #283, the middle session's worktree-removal failure escaped
|
||||
* {@code release()} uncaught, and this test proved {@code reapIdle}'s own per-session try/catch
|
||||
* (CB-581) kept the rest of the pass going regardless. Now that {@code release()} itself catches
|
||||
* a worktree-removal failure (matching every sibling cleanup step in that method) and only logs
|
||||
* a WARN, {@code release()} no longer throws for this reason — so all three idle sessions are
|
||||
* released and counted, and the middle one's removal failure is now visible only as the WARN
|
||||
* {@code release()} itself logs, not as a reap-loop catch. {@code reapIdle}'s own guard (for a
|
||||
* failure {@code release()} still cannot swallow, e.g. from {@code launcher.stop}) is untouched
|
||||
* by this ticket. This test no longer exercises that guard — reaching it now needs a failure
|
||||
* that #283 does not catch inside {@code release()} itself.
|
||||
*/
|
||||
@Test
|
||||
void reapIdleSurvivesOneSessionThatFailsToRelease() {
|
||||
void reapIdleCountsAllThreeSessionsWhenOnlyItsWorktreeRemovalFails() {
|
||||
long[] clock = {0};
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
RecordingWorktrees worktrees = new RecordingWorktrees();
|
||||
@@ -1004,8 +1027,8 @@ class SessionManagerTest {
|
||||
sessions.asPresence().markPresent(a.terminalId());
|
||||
sessions.asPresence().markPresent(b.terminalId());
|
||||
sessions.asPresence().markPresent(c.terminalId());
|
||||
// The middle session's worktree removal fails — release() propagates that, so this is the
|
||||
// one call reapIdle's per-session guard must survive without skipping the rest of the pass.
|
||||
// The middle session's worktree removal fails — fleetd #283 makes release() catch and log
|
||||
// this itself, so it no longer propagates out of release() at all.
|
||||
worktrees.failRemoveFor(b.worktree());
|
||||
|
||||
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
|
||||
@@ -1026,14 +1049,16 @@ class SessionManagerTest {
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.filter(m -> m.contains(b.paneId()))
|
||||
.findFirst()
|
||||
.orElse("no reap-failure WARN logged");
|
||||
.orElse("no worktree-removal-failure WARN logged");
|
||||
assertTrue(warn.contains(b.terminalId()), "the WARN names the failed session's terminal: " + warn);
|
||||
assertTrue(warn.contains(b.worktree()), "the WARN names the failed session's worktree: " + warn);
|
||||
} finally {
|
||||
sessionLog.detachAppender(appender);
|
||||
}
|
||||
|
||||
assertEquals(2, reaped, "the middle session's failure is logged, not counted as reaped");
|
||||
assertEquals(3, reaped,
|
||||
"fleetd #283: release() no longer throws for a worktree-removal failure, so reapIdle "
|
||||
+ "counts all three idle sessions as reaped");
|
||||
assertTrue(sessions.get(a.paneId()).isEmpty(), "the first session is still released");
|
||||
assertTrue(sessions.get(c.paneId()).isEmpty(), "the third session is still released");
|
||||
assertTrue(sessions.get(b.paneId()).isEmpty(),
|
||||
@@ -1044,6 +1069,79 @@ class SessionManagerTest {
|
||||
assertEquals(1, paneCloseCallsFor(herdr, "w9:pRoot_3"), "the third pane is stopped");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #290: the #283 fix above closed the one trigger this suite used for {@code
|
||||
* reapIdle}'s own per-session try/catch (CB-581) — a worktree-removal failure is now caught
|
||||
* and logged inside {@code release()} itself, so it never reaches {@code reapIdle}'s guard at
|
||||
* all. This test restores coverage of that guard using the trigger the ticket names: {@code
|
||||
* release()} calls {@code launcher.stop(paneId)} with no try/catch around it, so a failing
|
||||
* {@code pane.close} propagates straight out of {@code release()} uncaught. {@link
|
||||
* FakeHerdr#paneCloseFailsForPane} (added for this ticket) makes exactly the middle session's
|
||||
* stop fail, while the other two still succeed, so this proves {@code reapIdle} keeps reaping
|
||||
* the rest of the roster rather than aborting the whole pass.
|
||||
*/
|
||||
@Test
|
||||
void reapIdleSurvivesOneSessionWhoseLauncherStopFails() {
|
||||
long[] clock = {0};
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
RecordingWorktrees worktrees = new RecordingWorktrees();
|
||||
SessionManager sessions = sessionManager(herdr, worktrees, () -> clock[0]);
|
||||
MemberSession a = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("cb-290a", null));
|
||||
MemberSession b = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("cb-290b", null));
|
||||
MemberSession c = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("cb-290c", null));
|
||||
sessions.asPresence().markPresent(a.terminalId());
|
||||
sessions.asPresence().markPresent(b.terminalId());
|
||||
sessions.asPresence().markPresent(c.terminalId());
|
||||
// Only the middle session's herdr pane fails to close — a and c stop normally. This is the
|
||||
// trigger reapIdle's own guard is for, now that #283 closed the worktree-removal trigger.
|
||||
herdr.paneCloseFailsForPane("w9:pRoot_2", "internal_error");
|
||||
|
||||
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
|
||||
ch.qos.logback.classic.Logger sessionLog =
|
||||
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(SessionManager.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.setContext(ctx);
|
||||
appender.start();
|
||||
sessionLog.addAppender(appender);
|
||||
sessionLog.setLevel(Level.WARN);
|
||||
int reaped;
|
||||
try {
|
||||
clock[0] = 100;
|
||||
reaped = sessions.reapIdle(10);
|
||||
|
||||
String warn = appender.list.stream()
|
||||
.filter(e -> e.getLevel().equals(Level.WARN))
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.filter(m -> m.contains("reap failed") && m.contains(b.paneId()))
|
||||
.findFirst()
|
||||
.orElse("no reap-failed WARN logged for the failing session");
|
||||
assertTrue(warn.contains(b.terminalId()), "the WARN names the failed session's terminal: " + warn);
|
||||
} finally {
|
||||
sessionLog.detachAppender(appender);
|
||||
}
|
||||
|
||||
assertEquals(2, reaped,
|
||||
"the middle session's launcher.stop failure is not counted as reaped, but must not "
|
||||
+ "abort reaping the other two");
|
||||
assertTrue(sessions.get(a.paneId()).isEmpty(), "the first session is still released");
|
||||
assertTrue(sessions.get(c.paneId()).isEmpty(),
|
||||
"the third session is still reached and released — proves the pass did not abort "
|
||||
+ "when the middle session's release() threw");
|
||||
assertTrue(sessions.get(b.paneId()).isEmpty(),
|
||||
"the middle session is still deregistered — release() removes it from the registry "
|
||||
+ "before launcher.stop() runs, regardless of whether stop() then throws");
|
||||
assertEquals(1, paneCloseCallsFor(herdr, "w9:pRoot_1"), "the first pane is stopped");
|
||||
assertEquals(1, paneCloseCallsFor(herdr, "w9:pRoot_2"),
|
||||
"the middle pane's stop was attempted, even though it failed");
|
||||
assertEquals(1, paneCloseCallsFor(herdr, "w9:pRoot_3"), "the third pane is stopped");
|
||||
assertEquals(List.of(a.worktree(), c.worktree()), worktrees.removeCalls().stream().sorted().toList(),
|
||||
"the middle session's worktree removal never runs — release() throws before reaching "
|
||||
+ "it — while the other two, unaffected, still have theirs removed");
|
||||
}
|
||||
|
||||
@Test
|
||||
void unchangedRegressionCleanCompletedReleaseStillRemovesTheWorktree() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
@@ -1058,6 +1156,41 @@ class SessionManagerTest {
|
||||
"COMPLETED release of a clean worktree still removes it");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #293: {@code HerdrPeerLauncher.stop()} used to run {@code spaces.closeTab} bare — a
|
||||
* failing {@code tab.close} (any code other than {@code *_not_found}) propagated straight out
|
||||
* of {@code stop()}. {@code SessionManager.release} calls {@code launcher.stop(paneId)} with
|
||||
* no try/catch (fleetd #283 wrapped the WORKTREE-removal step further down, not this one), so
|
||||
* the throw happened <em>before</em> that worktree-removal step ever ran — and by then {@code
|
||||
* registry.remove(paneId)} had already run, so a second {@code stop} is a no-op: the worktree
|
||||
* leaked with no retry path. The pane itself is already closed by the time {@code tab.close}
|
||||
* runs, so its failure is cosmetic workspace tidying, not a real teardown failure. The fix
|
||||
* wraps {@code closeTab} inside {@code stop()} so it no longer throws for this reason; this
|
||||
* test proves both halves at once: {@code release()} does not throw, and it still removes the
|
||||
* worktree.
|
||||
*/
|
||||
@Test
|
||||
void releaseStillRemovesTheWorktreeWhenCloseTabFails() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
RecordingWorktrees worktrees = new RecordingWorktrees();
|
||||
SessionManager sessions = sessionManager(herdr, worktrees);
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("cb-293a", null));
|
||||
// FakeHerdr's pane.get always answers with tab_id "w9:t2" for a tab-placement spawn.
|
||||
herdr.tabCloseFailsForTab("w9:t2", "internal_error");
|
||||
|
||||
assertDoesNotThrow(() -> sessions.release(s.paneId()),
|
||||
"a failing tab.close is cosmetic (the pane is already closed by then) — it must not "
|
||||
+ "propagate out of release()");
|
||||
|
||||
assertTrue(herdr.called("tab.close"), "tab.close was still attempted");
|
||||
assertEquals(List.of(s.worktree()), worktrees.removeCalls(),
|
||||
"release() must still remove the worktree even though tab.close failed — this is "
|
||||
+ "the leak fleetd #293 reports: before the fix, release() never reached this "
|
||||
+ "step at all");
|
||||
assertTrue(sessions.get(s.paneId()).isEmpty(), "the session is still deregistered");
|
||||
}
|
||||
|
||||
@Test
|
||||
void unchangedRegressionDirtyCompletedReleaseStillPreservesTheWorktree() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
|
||||
@@ -212,9 +212,41 @@ class WorktreeSessionManagerTest {
|
||||
assertEquals("/repo", remove.repoRoot());
|
||||
assertEquals(s.worktree(), remove.worktreePath());
|
||||
// The fake records no branch-delete calls because Worktrees.remove only removes the checkout.
|
||||
assertTrue(worktrees.deleteBranchCalls().isEmpty(),
|
||||
"a normal release must NEVER delete the branch — it is the worker's only recoverable "
|
||||
+ "copy of committed work, and only the failed-provisioning path may remove it");
|
||||
assertTrue(sessions.get(paneId).isEmpty(), "released session is no longer retrievable");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #283 defect 1. Every other cleanup step in {@code release()} is wrapped in try/catch,
|
||||
* because {@code exec()} can throw on a non-zero exit or its own 30s timeout — this was the one
|
||||
* step left bare. By the time it runs, the registry entry, the retained handle, and the pane are
|
||||
* all already gone, so a throw here used to escape {@code release()} after the session was
|
||||
* already fully torn down: a second stop on the same paneId is a no-op (nothing left to find),
|
||||
* so there was no retry path, and the caller saw a failed stop for a session that was in fact
|
||||
* gone. This test makes the worktree removal throw and asserts release() still completes with
|
||||
* the pane stopped and the registry clean.
|
||||
*/
|
||||
@Test
|
||||
void releaseCompletesAndStopsPaneEvenWhenWorktreeRemovalThrows() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt")
|
||||
.failRemove("stale index lock");
|
||||
SessionManager sessions = new SessionManager(workerService(herdr), worktrees);
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("cb-283-1", null));
|
||||
String paneId = s.paneId();
|
||||
|
||||
sessions.release(paneId); // must not throw
|
||||
|
||||
assertTrue(herdr.called("pane.close"), "the pane is still stopped despite the removal failure");
|
||||
assertEquals(1, worktrees.removeCalls().size(), "worktree removal was still attempted");
|
||||
assertTrue(sessions.get(paneId).isEmpty(),
|
||||
"the session is deregistered regardless of the removal failure");
|
||||
assertEquals(0, sessions.size(), "the registry is left clean");
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-576. A normal {@code COMPLETED} release whose worktree holds uncommitted work must NOT
|
||||
* remove it — {@code --force} would destroy the worker's only copy. The bridge cannot see
|
||||
@@ -258,6 +290,38 @@ class WorktreeSessionManagerTest {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-576 review (fleetd #116). A worktree that is already gone (operator cleanup,
|
||||
* {@code git worktree prune}, an earlier half-completed release) must not break teardown.
|
||||
* {@code hasUncommitted} reports the missing path clean (mirroring {@code GitWorktrees}), so
|
||||
* {@code release} still runs {@code notifyReleased} (the CB-516 fast-fail for a blocked
|
||||
* {@code fleet_send} caller) and {@code launcher.stop} (so the pane is not orphaned), and falls
|
||||
* through to the already-gone-tolerant {@code remove}. Since CB-581 this all happens because the
|
||||
* notify-and-stop work sits in {@code release}'s {@code finally}/post-try block rather than a
|
||||
* checked branch — this test pins that shape by construction.
|
||||
*/
|
||||
@Test
|
||||
void releaseStillStopsPaneAndNotifiesWhenWorktreeIsGone() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt");
|
||||
SessionManager sessions = new SessionManager(workerService(herdr), worktrees);
|
||||
List<SessionManager.ReleaseDetail> released = new java.util.concurrent.CopyOnWriteArrayList<>();
|
||||
sessions.onRelease(released::add);
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("cb-576g", null));
|
||||
|
||||
worktrees.markGone(s.worktree());
|
||||
sessions.release(s.paneId());
|
||||
|
||||
assertEquals(1, released.size(),
|
||||
"notifyReleased must still fire when the worktree is already gone (CB-516)");
|
||||
assertEquals(s.terminalId(), released.getFirst().terminalId());
|
||||
assertTrue(herdr.called("pane.close"),
|
||||
"the pane must still be stopped when the worktree is already gone");
|
||||
assertEquals(1, worktrees.removeCalls().size(),
|
||||
"release still calls the already-gone-tolerant remove");
|
||||
}
|
||||
|
||||
@Test
|
||||
void drainAllPreservesWorktreeOfIdleSession() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
@@ -320,6 +384,39 @@ class WorktreeSessionManagerTest {
|
||||
assertEquals(0, sessions.size(), "failed acquire leaves no registry entry");
|
||||
assertFalse(herdr.called("agent.start"), "spawn is never reached when add fails");
|
||||
assertTrue(worktrees.removeCalls().isEmpty(), "no worktree was added, so none is removed");
|
||||
assertTrue(worktrees.deleteBranchCalls().isEmpty(),
|
||||
"add() itself never created the branch in git, so there is nothing to delete");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #283 defect 2. {@code acquireWithWorktree}'s catch covers every failure AFTER
|
||||
* {@code worktrees.add()} returns — {@code overlayParity}, {@code shareWithGroup},
|
||||
* {@code launcher.spawn} itself — so by the time it runs, {@code branch} was actually created in
|
||||
* git. It removed only the worktree and forgot the branch, leaking a {@code worker/<slug>-<nonce>}
|
||||
* branch on every routine spawn failure (a quarantined credential, a backend refusal). This test
|
||||
* makes {@code overlayParity} (a post-add() step) throw and asserts the branch is deleted, the
|
||||
* same way #274 already does for the sibling failure inside {@code add()} itself.
|
||||
*/
|
||||
@Test
|
||||
void spawnFailureAfterAddDeletesTheOrphanedBranch() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt")
|
||||
.failOverlay("overlayParity failed");
|
||||
SessionManager sessions = new SessionManager(workerService(herdr), worktrees);
|
||||
|
||||
assertThrows(WorktreeException.class, () ->
|
||||
sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("cb-283-2", null)));
|
||||
|
||||
assertEquals(0, sessions.size(), "failed acquire leaves no registry entry");
|
||||
assertFalse(herdr.called("agent.start"), "spawn is never reached when overlayParity fails");
|
||||
assertEquals(1, worktrees.removeCalls().size(), "the worktree checkout is still removed");
|
||||
assertEquals(1, worktrees.deleteBranchCalls().size(),
|
||||
"the orphaned branch that add() actually created must also be deleted");
|
||||
FakeWorktrees.DeleteBranchCall del = worktrees.lastDeleteBranch();
|
||||
assertEquals("/repo", del.repoRoot());
|
||||
FakeWorktrees.AddCall add = worktrees.lastAdd();
|
||||
assertEquals(add.branch(), del.branch(), "the branch deleted is the exact one add() created");
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
Reference in New Issue
Block a user