Compare commits
31 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 9e4e423ad6 | |||
| d4a2cd720c | |||
| 1fb6176783 | |||
| 49df79203c | |||
| 235644c0f0 | |||
| 7772b41993 | |||
| eccd0548ce | |||
| c4e23eebad | |||
| 7df503dfc2 | |||
| f5e02fedd6 | |||
| 29e7a06c49 | |||
| 92a96fcbd8 | |||
| 13482872bb | |||
| c1ca6273fc | |||
| 5f1b260c81 | |||
| 30d6872779 | |||
| 9d1306d442 | |||
| e54e3d87ea | |||
| bf027f10b9 | |||
| 3c5873dfe2 | |||
| e29227d5f4 | |||
| 45aca9eb3e | |||
| 2af13ab1ff | |||
| 0788d84be8 | |||
| 20c1094cbf | |||
| 9011c59b9f | |||
| c11ad71ed0 | |||
| bdcf285265 | |||
| d4f93a7b13 | |||
| 5289eb509f | |||
| 703a05db41 |
+10
-4
@@ -87,12 +87,18 @@ jobs:
|
||||
apt-get update && apt-get install -y --no-install-recommends maven
|
||||
mvn -version
|
||||
|
||||
# The `contract` profile clears the default-excludes group, so the @Tag("contract") AMQP test
|
||||
# runs against the RabbitMQ service container (AMQP_URI). Pinned to the one contract test to
|
||||
# avoid re-running the unit suite already covered by the `build` job.
|
||||
# The `contract` profile clears the default-excludes group, so `-Dgroups=contract` runs every
|
||||
# @Tag("contract") test and nothing from the unit suite the `build` job already covered — a
|
||||
# tag selects the whole group, so a test added to it later runs here automatically. A prior
|
||||
# version of this step pinned `-Dtest=AmqpReplyInboxContractTest` by class name instead: that
|
||||
# silently excluded every other contract test (including the herdr ones) from CI, and nobody
|
||||
# noticed until the herdr protocol drifted out from under a test that never ran here
|
||||
# (fleetd #449). If this runner has no herdr socket, the herdr-backed tests in the group
|
||||
# skip on their own `assumeTrue` and only the broker-backed ones actually run — check the
|
||||
# step output rather than assuming which.
|
||||
- name: Contract tests
|
||||
working-directory: fleetd
|
||||
run: mvn -B -Pcontract test -Dtest=AmqpReplyInboxContractTest
|
||||
run: mvn -B -Pcontract test -Dgroups=contract
|
||||
|
||||
- name: Failing test output
|
||||
if: failure()
|
||||
|
||||
@@ -58,8 +58,9 @@ and the sender silently receives nothing. Fail toward the recoverable error.
|
||||
re-send because a call looks slow — the bridge delivers when the peer is `idle`, `blocked` or
|
||||
`done`. A spawned member must **also** have mounted the bridge MCP: until it has, it is not
|
||||
deliverable, and a send waits on that gate for ~60s and then fails without ever reaching its pane.
|
||||
5. **Never drive the terminal multiplexer directly** (no `herdr` CLI, no socket). The bridge owns
|
||||
policy; the multiplexer owns PTYs. Going around the bridge bypasses every rule above.
|
||||
5. **Never move a fleet session, pane or peer except through the bridge.** The bridge owns policy;
|
||||
the multiplexer owns PTYs. Any route that changes fleet state without the bridge's checks
|
||||
bypasses every rule above — the `herdr` CLI and its socket are the usual example.
|
||||
|
||||
### Primary (lead) — run this on every task, in order
|
||||
|
||||
@@ -223,6 +224,19 @@ must obey belongs in the charter, not here.
|
||||
- **This repo is the bridge.** The daemon is `fleetd`, its MCP mount is `http://127.0.0.1:8765/mcp`,
|
||||
and the code behind the rules above is `mcp/FleetMcp` (tools), `auth/Authz` (the role table),
|
||||
`mcp/ConnectionIdentity` (connection→role), and `worker/*Launcher` (`REPLY_CHARTER`).
|
||||
- **Herdr socket tests (measured 2026-09-10).** In this repo, herdr is a subject under test. A
|
||||
worker assigned to herdr code, and the lead, may let a test open the herdr socket directly in a
|
||||
throwaway workspace that the test tears down. This only covers
|
||||
`fleetd/src/test/java/dev/ltms/fleet/herdr/AgentControlContractTest.java`,
|
||||
`fleetd/src/test/java/dev/ltms/fleet/herdr/HerdrContractTest.java`,
|
||||
`fleetd/src/test/java/dev/ltms/fleet/herdr/PaneLocatorContractTest.java`, and
|
||||
`fleetd/src/test/java/dev/ltms/fleet/herdr/WorkspacePlacementContractTest.java`. It is not a
|
||||
general licence. Using the herdr CLI or socket to move a real fleet session, pane, or peer stays
|
||||
banned. That is the control plane that invariant 5 protects. Re-measure with
|
||||
`grep -rl 'UnixSocketHerdrClient.connect()' fleetd/src/test/java --include='*.java'`. A non-empty
|
||||
result means tests still open the socket and this note still applies. An empty result means nobody
|
||||
does this any more; delete this section. Canonical invariant 5 restatement is tracked in #458 and
|
||||
is not part of this change.
|
||||
- **`fleet_profiles`/`fleet_list` report two separate outage states, and they are not the same
|
||||
thing.** *Quarantined* (CB-578) means the backend told us it is out of capacity — a long,
|
||||
1800s-default cooldown. *Cooling off* (fleetd #201/#227) means a profile's credential threw two
|
||||
@@ -379,6 +393,16 @@ print("in sync:", w[i:w.index("\n```\n", i) + 1] == block)
|
||||
PY
|
||||
```
|
||||
|
||||
**Only the lead can run that check (measured 2026-09-10).** A member's provisioned worktree has
|
||||
`wiki/` uninitialized, so the script dies with `FileNotFoundError: wiki/7-Use-Cases.md`. Measured
|
||||
in three worker worktrees: `git submodule status` printed a leading `-` and `wiki/` held 0
|
||||
entries; the primary's own clone printed a leading `+` and the file was there. So never make this
|
||||
check a member's acceptance criterion — it is unsatisfiable for them, and a brief that asks for it
|
||||
is asking a worker to invent a pass. A member told to check it must say it could not run it, and
|
||||
must never report it as passed. The lead runs it in the main clone before merging. Re-measure with
|
||||
`git submodule status` in a member's worktree: a leading `-` means this still applies; once it
|
||||
prints a commit with no `-`, delete this paragraph.
|
||||
|
||||
## IDE MCP tools & validation workflow (enforced)
|
||||
|
||||
> **Primary only.** Workers have no IDE MCP mount — if you are a worker, skip this section and
|
||||
|
||||
+32
-17
@@ -207,8 +207,12 @@ herdrSocket: ~/.config/herdr/herdr.sock
|
||||
# omit and this profile's completion fallback behaves exactly as before.
|
||||
# Every backend words its refusal differently, so this is config, never a
|
||||
# vendor string baked into fleetd itself.
|
||||
# DEFERRED: compiled once into a startup pattern map — editing it needs a
|
||||
# daemon restart, same as this profile's model/baseUrl/argv.
|
||||
# HOT (fleetd #446): read live, cached by profile name, at every
|
||||
# completion-fallback check AND by fleet_profiles' exhaustionDetectionArmed —
|
||||
# editing it and reloading arms or disarms usage-limit detection for this
|
||||
# profile with no daemon restart. (Before fleetd #446 this was DEFERRED,
|
||||
# compiled once into a startup pattern map like model/baseUrl/argv still are —
|
||||
# see errorPattern below, which is still deferred that way on purpose.)
|
||||
# credentialId → CB-578 stage B: the credential this profile quarantines WITH when a
|
||||
# BACKEND_EXHAUSTED classification fires. Two profiles that set the SAME
|
||||
# credentialId share one quarantine — the case this exists for is two models
|
||||
@@ -227,8 +231,9 @@ herdrSocket: ~/.config/herdr/herdr.sock
|
||||
# happens, just without a profile-specific match; every backend words its
|
||||
# failure differently, so a hardcoded sentence would only ever match one
|
||||
# of them.
|
||||
# DEFERRED: compiled once into a startup pattern map, same as exhaustedPattern
|
||||
# — editing it needs a daemon restart.
|
||||
# DEFERRED: compiled once into a startup pattern map — editing it needs a
|
||||
# daemon restart. Unlike exhaustedPattern above (made hot by fleetd #446),
|
||||
# errorPattern was scoped out of that ticket on purpose and stays deferred.
|
||||
# # errorPattern: "503 Service Unavailable" # opt-in: classify a backend outage
|
||||
#
|
||||
# What happens once a match fires (BackendOutagePolicy, credentialId-keyed,
|
||||
@@ -468,9 +473,10 @@ placement: weighted
|
||||
# when the reload happens — not about how important the key is:
|
||||
# HOT → takes effect on the next spawn: the whole `fleet:` block (every role pool,
|
||||
# `charters`, and `tabLabel`), `placement:`, and an existing profile's weight / maxLoad
|
||||
# / credentialId. Those are hot because the placement policy (and, for credentialId,
|
||||
# the CB-578 stage B quarantine check) reads them through a supplier — being config is
|
||||
# not by itself enough to make a key hot.
|
||||
# / credentialId / exhaustedPattern. Those are hot because the placement policy (and,
|
||||
# for credentialId, the CB-578 stage B quarantine check; for exhaustedPattern, fleetd
|
||||
# #446's LiveExhaustedPatterns) reads them through a supplier — being config is not by
|
||||
# itself enough to make a key hot.
|
||||
# EXCEPT `fleet.leaders`: Fleetd.main reads it once at startup to build the lead tab
|
||||
# scanner and launcher, and neither is rebuilt on reload. A changed/added/removed
|
||||
# `fleet.leaders` entry is silently accepted — the reload reports "config reloaded"
|
||||
@@ -482,10 +488,10 @@ placement: weighted
|
||||
# stage B — baked once into the quarantine tracker built at startup), ADDING or
|
||||
# REMOVING a profile (a new backend needs its own launcher, and launchers are built
|
||||
# once), AND an existing profile's launch settings — model, baseUrl, argv, env,
|
||||
# configDir, mcpUrl, tabLabel, exhaustedPattern, errorPattern (fleetd #201 / #227 —
|
||||
# compiled once into a startup pattern map the same way exhaustedPattern is). The
|
||||
# launcher takes a copy of `profiles:` at startup and resolves every spawn out of
|
||||
# that copy, so those never
|
||||
# configDir, mcpUrl, tabLabel, errorPattern (fleetd #201 / #227 — compiled once into
|
||||
# a startup pattern map; exhaustedPattern used to be compiled the same way until
|
||||
# fleetd #446 made it hot — see above). The launcher takes a copy of `profiles:` at
|
||||
# startup and resolves every spawn out of that copy, so those never
|
||||
# reach a launch until you restart. The reload logs them by name rather than
|
||||
# pretending they applied.
|
||||
# COLD → cannot change at all: `bind:`, `herdrSocket:`, `broker:` and `auth:`. The socket is
|
||||
@@ -787,12 +793,21 @@ guard:
|
||||
# .claude/skills/, so a member spawned against ANY repo — not only one that already ships its own
|
||||
# copy — can load a bridge skill (e.g. implementer). Unset (the default): no worktree is touched
|
||||
# beyond today's behaviour. A skill folder the target repo already carries under
|
||||
# .claude/skills/<name> is never overwritten — the repo's own copy always wins. Claude Code
|
||||
# members only; an opencode member reads a different path (.opencode/agent) this key does not
|
||||
# touch. Best-effort like worktreeGroup above: a missing/unreadable directory here is logged and
|
||||
# skipped, never a failed spawn. Every non-hidden subdirectory of this directory is copied
|
||||
# wholesale, with no per-file allowlist — don't park scratch files or drafts alongside the real
|
||||
# skill folders, they will be copied into every provisioned worktree too.
|
||||
# .claude/skills/<name> is never overwritten — the repo's own copy always wins. Best-effort like
|
||||
# worktreeGroup above: a missing/unreadable directory here is logged and skipped, never a failed
|
||||
# spawn. Every non-hidden subdirectory of this directory is copied wholesale, with no per-file
|
||||
# allowlist — don't park scratch files or drafts alongside the real skill folders, they will be
|
||||
# copied into every provisioned worktree too.
|
||||
#
|
||||
# fleetd #393: which member KINDS actually consume this once it is copied. kind: claude-code —
|
||||
# the Claude Code CLI discovers .claude/skills/ on its own; nothing else is needed. kind: opencode
|
||||
# — opencode has no such discovery, so OpenCodeLauncher reads whatever landed under
|
||||
# .claude/skills/ and appends each seeded skill's SKILL.md to the generated instructions[] file
|
||||
# (opencode's only channel for static guidance text; unlike Claude Code's Skill tool, the content
|
||||
# is always part of the system prompt, not loaded on demand). Both kinds are covered as of #393 —
|
||||
# earlier builds copied the files for every kind but only claude-code could read them, and the
|
||||
# seeding log said "N of M" regardless. Check the per-spawn launcher log (not just the seeding
|
||||
# log) to see what a given member actually got.
|
||||
# memberSkills: /path/to/fleetd/checkout/.claude/skills
|
||||
|
||||
# Session lifecycle limits (CB-303). All knobs are opt-in; omit or set to null to keep
|
||||
|
||||
@@ -18,6 +18,7 @@ import dev.ltms.fleet.inject.BackendErrorSink;
|
||||
import dev.ltms.fleet.inject.CompletionResolver;
|
||||
import dev.ltms.fleet.inject.ExhaustedPatternLookup;
|
||||
import dev.ltms.fleet.inject.ExhaustionSink;
|
||||
import dev.ltms.fleet.inject.LiveExhaustedPatterns;
|
||||
import dev.ltms.fleet.inject.Injector;
|
||||
import dev.ltms.fleet.inject.StatusPoller;
|
||||
import dev.ltms.fleet.inject.TurnListener;
|
||||
@@ -65,11 +66,13 @@ import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.ArrayList;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.LinkedHashSet;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Optional;
|
||||
import java.util.Set;
|
||||
import java.util.TreeSet;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
@@ -78,6 +81,7 @@ import java.util.function.Function;
|
||||
import java.util.function.Predicate;
|
||||
import java.util.function.Supplier;
|
||||
import java.util.regex.Pattern;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
/**
|
||||
* {@code fleetd} entry point. Wires the real herdr socket client to the REST app and
|
||||
@@ -378,28 +382,34 @@ public final class Fleetd {
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
// CB-578 stage A: classify a completion-fallback scrape that matches a profile's configured
|
||||
// usage-limit refusal as BACKEND_EXHAUSTED rather than handing it back as a real answer.
|
||||
// Compiled once at startup, keyed by profile name; a profile with no exhaustedPattern is
|
||||
// simply absent here, so its workers keep today's completion-fallback behaviour unchanged.
|
||||
Map<String, Pattern> exhaustedPatternsByProfile = new LinkedHashMap<>();
|
||||
cfg.profiles().forEach((name, profile) -> {
|
||||
if (profile.hasExhaustedPattern()) {
|
||||
exhaustedPatternsByProfile.put(name, Pattern.compile(profile.exhaustedPattern()));
|
||||
}
|
||||
});
|
||||
// fleetd #446: read live off `config` per lookup, cached by profile name — see
|
||||
// LiveExhaustedPatterns's class doc for why this replaces the old compiled-once-at-startup
|
||||
// map. A profile with no exhaustedPattern simply returns null here, so its workers keep
|
||||
// today's completion-fallback behaviour unchanged.
|
||||
LiveExhaustedPatterns liveExhaustedPatterns = new LiveExhaustedPatterns(() -> config.get().profiles());
|
||||
ExhaustedPatternLookup exhaustedPatterns = target -> sessions.roster().stream()
|
||||
.filter(session -> target.equals(session.terminalId()))
|
||||
.findFirst()
|
||||
.map(session -> exhaustedPatternsByProfile.get(session.profile()))
|
||||
.map(session -> liveExhaustedPatterns.patternFor(session.profile()))
|
||||
.orElse(null);
|
||||
// The startup coverage line still reports the boot-time snapshot only — it is printed once,
|
||||
// here, and a reload no longer needs to change what it said; exhaustionDetectionArmed (via
|
||||
// liveExhaustedPatterns.armed, wired into quarantineSource below) is what stays live.
|
||||
Set<String> exhaustedConfiguredAtStartup = cfg.profiles().entrySet().stream()
|
||||
.filter(e -> e.getValue().hasExhaustedPattern())
|
||||
.map(Map.Entry::getKey)
|
||||
.collect(Collectors.toCollection(LinkedHashSet::new));
|
||||
log.info("backend-exhausted classification (CB-578 stage A): {}",
|
||||
exhaustedPatternCoverageLine(cfg.profiles().keySet(), exhaustedPatternsByProfile.keySet()));
|
||||
exhaustedPatternCoverageLine(cfg.profiles().keySet(), exhaustedConfiguredAtStartup));
|
||||
// fleetd #201 Unit 5: classify a completion-fallback scrape that matches a profile's
|
||||
// configured backend-error refusal (a credential outage, a provider 5xx) as a backend error
|
||||
// rather than handing it back as a real answer. Compiled once at startup, keyed by profile
|
||||
// name, mirroring exhaustedPatternsByProfile above — a profile with no configured
|
||||
// errorPattern is simply absent here, so CompletionResolver falls back to its built-in
|
||||
// narrow {@code (?i)\bAPI Error\s*:} compatibility pattern for that profile's targets
|
||||
// (BackendErrorPatternLookup#legacy's contract — see backendErrorPatterns below).
|
||||
// rather than handing it back as a real answer. Still compiled once at startup, keyed by
|
||||
// profile name — unlike exhaustedPattern above (fleetd #446), errorPattern was left DEFERRED
|
||||
// on purpose: the ticket that made exhaustedPattern hot scoped errorPattern/cooling-off out
|
||||
// explicitly. A profile with no configured errorPattern is simply absent here, so
|
||||
// CompletionResolver falls back to its built-in narrow {@code (?i)\bAPI Error\s*:}
|
||||
// compatibility pattern for that profile's targets (BackendErrorPatternLookup#legacy's
|
||||
// contract — see backendErrorPatterns below).
|
||||
Map<String, Pattern> errorPatternsByProfile = new LinkedHashMap<>();
|
||||
cfg.profiles().forEach((name, profile) -> {
|
||||
if (profile.hasErrorPattern()) {
|
||||
@@ -417,46 +427,20 @@ public final class Fleetd {
|
||||
// CREDENTIAL — not the profile name — so a profile sharing that credential (e.g. two models
|
||||
// on one OpenAI account) is refused too, not just the one that happened to report it. Reads
|
||||
// the profile config live off `config`, so a credentialId edit is hot: no restart needed.
|
||||
//
|
||||
// fleetd #175: this sink is now also the quarantine target for OpenCodeLauncher's
|
||||
// model-mismatch check — it needs nothing profile-specific from the caller beyond `target`
|
||||
// (a herdr terminal id) and `reason`, so reusing it here is exactly "the existing
|
||||
// ExhaustionSink path", not a new mechanism.
|
||||
//
|
||||
// fleetd #234: that check fires from SessionAwareHandle.agentSessionId(), which runs during
|
||||
// SessionManager.acquire() BEFORE this session is registered in sessions.roster() — so the
|
||||
// roster-only lookup below used to find nothing, .ifPresent silently no-op'd, and the
|
||||
// ERROR the check had just logged ("quarantining this profile's credential") was a lie:
|
||||
// nothing was quarantined, and nothing said so. Two changes: (1) OpenCodeLauncher now
|
||||
// passes its OWN profile name via ExhaustionSink's 3-arg overload — it already has the
|
||||
// FleetConfig.Profile in hand and does not need the roster at all — used here as a
|
||||
// fallback whenever the roster lookup misses; (2) if a profile still cannot be resolved
|
||||
// (neither the roster nor the hint names a configured one), this logs loudly at ERROR
|
||||
// instead of silently doing nothing — a control that cannot act must say so.
|
||||
// fleetd #234, round 4: the 3-arg overload is now ExhaustionSink's single abstract method,
|
||||
// so this is safely a lambda — there is no separate 2-arg overload left for it to bind to
|
||||
// instead and silently drop profileHint (that was rounds 1-3's whole hazard).
|
||||
ExhaustionSink exhaustionSink = (target, reason, profileHint) -> {
|
||||
String profileName = sessions.roster().stream()
|
||||
.filter(session -> target.equals(session.terminalId()))
|
||||
.findFirst()
|
||||
.map(MemberSession::profile)
|
||||
.orElse(profileHint);
|
||||
FleetConfig.Profile profile = profileName == null ? null : config.get().profiles().get(profileName);
|
||||
if (profile == null) {
|
||||
log.error("quarantine requested for target '{}' ({}) but no profile could be "
|
||||
+ "resolved — the target is not (yet) in the roster, and {} — "
|
||||
+ "credential NOT quarantined (fleetd #234)",
|
||||
target, reason,
|
||||
profileHint == null ? "no profile hint was given"
|
||||
: "the hinted profile '" + profileHint + "' is not configured");
|
||||
return;
|
||||
}
|
||||
String credentialId = profile.effectiveCredentialId();
|
||||
quarantine.quarantine(credentialId);
|
||||
log.warn("credential '{}' quarantined for {}s (profile '{}'): {}", credentialId,
|
||||
cfg.quarantineCooldownSeconds(), profile.profile(), reason);
|
||||
};
|
||||
// fleetd #446 criterion 3: the backend text that triggered the most recent quarantine of
|
||||
// each credential, so fleet_profiles can report WHY a limit was hit, not only that it was.
|
||||
// Keyed by credential id — the same key BackendQuarantine's own remainingSeconds uses —
|
||||
// and written at the one call site that actually quarantines (inside exhaustionSink below),
|
||||
// so a reason can never be reported for a quarantine that never happened. Bounded the same
|
||||
// way BackendQuarantine's own internal map is documented to be: by the number of distinct
|
||||
// credentials ever exhausted, not by how often the config is edited.
|
||||
Map<String, String> quarantineReasonByCredential = new ConcurrentHashMap<>();
|
||||
// fleetd #446 follow-up (round 3): extracted into exhaustionSink(...) below — see that
|
||||
// method's javadoc for the full fleetd #175/#234/#446 history this used to carry inline —
|
||||
// so a dedicated test can drive the exact ExhaustionSink main() builds, not a hand-rebuilt
|
||||
// copy of its shape.
|
||||
ExhaustionSink exhaustionSink = exhaustionSink(sessions, config, quarantine,
|
||||
quarantineReasonByCredential, cfg);
|
||||
// fleetd #175: point the forwarding sink handed to OpenCodeLauncher above at the real one,
|
||||
// now that `sessions` exists to resolve target -> session -> profile.
|
||||
exhaustionSinkRef.set(exhaustionSink);
|
||||
@@ -673,7 +657,7 @@ public final class Fleetd {
|
||||
// / 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 = quarantineSource(config, quarantine,
|
||||
exhaustedPatternsByProfile);
|
||||
liveExhaustedPatterns, quarantineReasonByCredential);
|
||||
FleetMcp.OutageSource outageSource = new FleetMcp.OutageSource(profile -> {
|
||||
var configured = config.get().profiles().get(profile);
|
||||
return configured == null ? null : configured.effectiveCredentialId();
|
||||
@@ -805,16 +789,99 @@ public final class Fleetd {
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #404: production source for quarantine reporting. Credential IDs are hot, but
|
||||
* exhausted patterns are compiled once at startup for {@link CompletionResolver}, so the armed
|
||||
* field must use that same compiled map until restart.
|
||||
* fleetd #446 follow-up (round 3): the quarantine {@link ExhaustionSink} main() actually wires
|
||||
* — extracted out of {@code main} for the same reason {@link #quarantineSource} below and
|
||||
* {@link #usageLimitFixWarning}/{@link #usageLimitFixWarningNoModel} were. Round 2 pinned the
|
||||
* WARNING text by having {@code FleetdUsageLimitFixWarningTest} call those two methods
|
||||
* directly, but a mutation battery proved that left two gaps: nothing proved this sink's
|
||||
* {@code log.warn} call actually uses either method (replacing the whole ternary with a
|
||||
* literal string), and nothing proved it picks the right one for a profile WITH a configured
|
||||
* {@code model:} versus one WITHOUT (swapping the ternary's two branches) — both mutations
|
||||
* left round 2's full 1592-test suite green. {@code FleetdExhaustionSinkWarningTest} drives
|
||||
* THIS factory's return value directly and asserts on the real text a {@code ListAppender}
|
||||
* attached to this class's own logger captures — the only way to prove the caller, not just
|
||||
* the two callees in isolation.
|
||||
*
|
||||
* <p>Behaviourally unchanged from the inline lambda this replaces:
|
||||
* <ul>
|
||||
* <li>fleetd #175: also the quarantine target for {@code OpenCodeLauncher}'s model-mismatch
|
||||
* check — it needs nothing profile-specific beyond {@code target}/{@code reason}, so
|
||||
* reusing this sink is the existing {@code ExhaustionSink} path, not a new mechanism;</li>
|
||||
* <li>fleetd #234 (round 4): resolves {@code target} to a profile via the live session
|
||||
* roster first, falling back to {@code profileHint} — {@code
|
||||
* SessionAwareHandle.agentSessionId()} fires this sink during {@code
|
||||
* SessionManager.acquire()}, <em>before</em> that session is registered, so the
|
||||
* roster-only lookup alone used to silently miss it. A profile resolved neither way
|
||||
* logs loudly at ERROR instead of doing nothing;</li>
|
||||
* <li>fleetd #446 criterion 3: writes {@code quarantineReasonByCredential} at the one call
|
||||
* site that actually quarantines, keyed by credential id — see the field's declaration
|
||||
* in {@code main} for the bound on its size.</li>
|
||||
* </ul>
|
||||
*
|
||||
* <p>{@code cfg} is the startup {@link FleetConfig} snapshot, read only for {@code
|
||||
* quarantineCooldownSeconds()} in the log text — a Cold key (see {@code ConfigRef}'s class
|
||||
* doc), so reading it off the snapshot rather than {@code config.get()} makes no observable
|
||||
* difference and matches what the inline version already did.
|
||||
*/
|
||||
static ExhaustionSink exhaustionSink(SessionManager sessions, ConfigRef config, BackendQuarantine quarantine,
|
||||
Map<String, String> quarantineReasonByCredential, FleetConfig cfg) {
|
||||
return (target, reason, profileHint) -> {
|
||||
String profileName = sessions.roster().stream()
|
||||
.filter(session -> target.equals(session.terminalId()))
|
||||
.findFirst()
|
||||
.map(MemberSession::profile)
|
||||
.orElse(profileHint);
|
||||
FleetConfig.Profile profile = profileName == null ? null : config.get().profiles().get(profileName);
|
||||
if (profile == null) {
|
||||
log.error("quarantine requested for target '{}' ({}) but no profile could be "
|
||||
+ "resolved — the target is not (yet) in the roster, and {} — "
|
||||
+ "credential NOT quarantined (fleetd #234)",
|
||||
target, reason,
|
||||
profileHint == null ? "no profile hint was given"
|
||||
: "the hinted profile '" + profileHint + "' is not configured");
|
||||
return;
|
||||
}
|
||||
String credentialId = profile.effectiveCredentialId();
|
||||
quarantine.quarantine(credentialId);
|
||||
quarantineReasonByCredential.put(credentialId, reason);
|
||||
log.warn("credential '{}' quarantined for {}s (profile '{}'): {}", credentialId,
|
||||
cfg.quarantineCooldownSeconds(), profile.profile(), reason);
|
||||
// fleetd #446 criterion 2: name the fix, not just the fact — an operator reading this
|
||||
// should not have to work out which of several configured models to touch.
|
||||
String model = profile.model();
|
||||
log.warn(model != null && !model.isBlank()
|
||||
? usageLimitFixWarning(profile.profile(), model)
|
||||
: usageLimitFixWarningNoModel(profile.profile(), cfg.quarantineCooldownSeconds()));
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #404, superseded by fleetd #446: production source for quarantine reporting.
|
||||
* Credential IDs were already hot; {@code exhaustedPattern} used to be compiled once at
|
||||
* startup for {@link CompletionResolver}, so the armed field had to read that same frozen map
|
||||
* until restart — the exact asymmetry #446 exists to close. Both {@code exhaustedPatternArmed}
|
||||
* here and {@link CompletionResolver}'s own classification now read the ONE live {@link
|
||||
* LiveExhaustedPatterns} instance ({@code exhaustedPatterns.armed}/{@code .patternFor}), so a
|
||||
* reload that arms or disarms a profile's detection is visible to both at once — never two
|
||||
* independently-updated copies that could disagree, the fleetd #404 rule this method used to
|
||||
* violate on purpose and now upholds.
|
||||
*
|
||||
* <p>{@code modelFor} and {@code reasonFor} back fleetd #446 criterion 3 — {@code
|
||||
* fleet_profiles}'s quarantined row naming which model a quarantined profile runs
|
||||
* ({@code modelFor}, read live off {@code config} the same way {@code credentialIdFor} already
|
||||
* is) and the backend text that triggered the most recent quarantine of that credential
|
||||
* ({@code reasonFor}, backed by {@code reasonByCredential} — see its call site in {@code main}
|
||||
* for where that map is written, at the one place a credential is actually quarantined).
|
||||
*/
|
||||
static FleetMcp.QuarantineSource quarantineSource(ConfigRef config, BackendQuarantine quarantine,
|
||||
Map<String, Pattern> startupExhaustedPatterns) {
|
||||
LiveExhaustedPatterns exhaustedPatterns, Map<String, String> reasonByCredential) {
|
||||
return new FleetMcp.QuarantineSource(profile -> {
|
||||
var configured = config.get().profiles().get(profile);
|
||||
return configured == null ? null : configured.effectiveCredentialId();
|
||||
}, quarantine, profile -> startupExhaustedPatterns.containsKey(profile));
|
||||
}, quarantine, exhaustedPatterns::armed, profile -> {
|
||||
var configured = config.get().profiles().get(profile);
|
||||
return configured == null ? null : configured.model();
|
||||
}, reasonByCredential::get);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -833,6 +900,40 @@ public final class Fleetd {
|
||||
* with 0 errors and left all 1506 existing tests green before {@code
|
||||
* FleetdPatternCoverageLineTest} was added to catch exactly that swap.
|
||||
*/
|
||||
/**
|
||||
* fleetd #446 follow-up: the criterion-2 WARNING text — "name the fix, not just the fact" —
|
||||
* for a profile whose {@code model:} is configured. Extracted out of the {@code
|
||||
* exhaustionSink} lambda the same way {@link #exhaustedPatternCoverageLine} was extracted out
|
||||
* of {@code main}: the inline version compiled, ran, and read correctly, but nothing pinned
|
||||
* its text, so a later edit could silently stop naming the fix and every test would stay
|
||||
* green (measured: renaming the leading {@code "usage-limit fix:"} tag left all 1586
|
||||
* pre-follow-up tests passing). {@link FleetdUsageLimitFixWarningTest} asserts on the return
|
||||
* value of this method directly, which is exactly what the caller logs — {@code log.warn} is
|
||||
* called with this method's result as a single, already-formatted argument, so the string a
|
||||
* test sees here is byte-identical to what {@code fleetd.out} receives.
|
||||
*
|
||||
* <p>Deliberately asserts on three substrings, not the whole sentence (the profile name, the
|
||||
* model name, and the literal {@code enabled: false} / {@code models.allow} action) — see that
|
||||
* test's class doc for why a whole-sentence pin is the wrong granularity here.
|
||||
*/
|
||||
static String usageLimitFixWarning(String profile, String model) {
|
||||
return "usage-limit fix: profile '" + profile + "' runs model '" + model + "' — set "
|
||||
+ "`enabled: false` on that model's entry under models.allow in fleetd.yaml to "
|
||||
+ "stop new spawns landing on it (models: is hot, no restart needed); remove the "
|
||||
+ "line again once the subscription window resets";
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #446 follow-up: the {@link #usageLimitFixWarning} counterpart for a profile with no
|
||||
* {@code model:} configured — {@code models.allow} gates by model name, so there is nothing
|
||||
* for an operator to flip, and this says so instead of naming a fix that does not exist.
|
||||
*/
|
||||
static String usageLimitFixWarningNoModel(String profile, Integer quarantineCooldownSeconds) {
|
||||
return "usage-limit fix: profile '" + profile + "' has no model: configured, so "
|
||||
+ "models.allow cannot gate it by name — the " + quarantineCooldownSeconds
|
||||
+ "s quarantine above is the only thing keeping new spawns off it for now";
|
||||
}
|
||||
|
||||
static String exhaustedPatternCoverageLine(Set<String> allProfiles, Set<String> configuredProfiles) {
|
||||
return CompletionResolver.coverage("exhaustedPattern", CompletionResolver.UnsetMeaning.OFF,
|
||||
allProfiles, configuredProfiles);
|
||||
|
||||
@@ -51,7 +51,17 @@ import java.util.function.Supplier;
|
||||
* filter) and by {@code fleet_profiles}/{@code GET /profiles} (via
|
||||
* {@code PeerLauncher.disabledModels()}). Nothing about {@code models:} is baked into an
|
||||
* object built at startup, so — unlike the deferred keys below — there is no frozen half left
|
||||
* to report; it moved here from deferred rather than joining split.</li>
|
||||
* to report; it moved here from deferred rather than joining split. An existing profile's
|
||||
* {@code exhaustedPattern} (fleetd #446) joined this class the same way: it used to be
|
||||
* compiled once into {@code Fleetd.main}'s startup pattern map (see the Deferred bullet's old
|
||||
* wording, and {@code LiveExhaustedPatterns}'s class doc for the history), and is now read
|
||||
* live, cached by profile name, by both {@code CompletionResolver}'s classification (via
|
||||
* {@code LiveExhaustedPatterns.patternFor}) and {@code fleet_profiles}'s {@code
|
||||
* exhaustionDetectionArmed} (via {@code LiveExhaustedPatterns.armed}) — the one live object
|
||||
* both read, so a reload that arms or disarms a profile's usage-limit detection takes effect
|
||||
* on the next check with no restart. {@code errorPattern}, {@code exhaustedPattern}'s sibling
|
||||
* key for backend-error (not usage-limit) classification, was deliberately left OUT of this
|
||||
* fleetd #446 change and stays deferred below — the ticket scoped it out explicitly.</li>
|
||||
* <li><strong>Deferred</strong> — accepted into the new snapshot, but the wiring built at startup
|
||||
* keeps the old value until a restart: {@code lifecycle:}, {@code leadHeartbeat:},
|
||||
* {@code idleSleepGuard:} ({@code Fleetd.java} reads it once, at startup, to decide whether
|
||||
@@ -81,15 +91,17 @@ import java.util.function.Supplier;
|
||||
* (if any) simply keeps its old settings), adding or removing a profile (a new backend needs its own launcher,
|
||||
* which is constructed once), <em>and an existing profile's launch settings</em> —
|
||||
* {@code model}, {@code baseUrl}, {@code argv}, {@code env}, {@code mcpUrl},
|
||||
* {@code exhaustedPattern} (CB-578 stage A — compiled once into {@code Fleetd.main}'s
|
||||
* pattern map at startup), {@code errorPattern} (fleetd #201 Unit 5 — compiled once into
|
||||
* {@code Fleetd.main}'s backend-error pattern map at startup, the same way),
|
||||
* {@code errorPattern} (fleetd #201 Unit 5 — compiled once into
|
||||
* {@code Fleetd.main}'s backend-error pattern map at startup; deliberately NOT made hot
|
||||
* alongside {@code exhaustedPattern} by fleetd #446 — that ticket scoped {@code errorPattern}
|
||||
* and cooling-off out explicitly),
|
||||
* {@code ideProjectDir} / {@code ideOpenCommand} / {@code autoCompactWindow} (fleetd #323
|
||||
* instance 1 — all three are read at spawn off the same frozen profile map and were missing
|
||||
* from {@link #sameLaunchSettings}), and the rest of {@link #sameLaunchSettings}.
|
||||
* {@code credentialId} (CB-578 stage B) is NOT on
|
||||
* this list — it is read live off the config supplier at every quarantine check and
|
||||
* exhaustion event, exactly like {@code weight} / {@code maxLoad}, so it is hot instead.
|
||||
* {@code credentialId} (CB-578 stage B) and {@code exhaustedPattern} (fleetd #446) are NOT on
|
||||
* this list — both are read live off the config supplier at every quarantine check and
|
||||
* exhaustion event, exactly like {@code weight} / {@code maxLoad}, so both are hot instead
|
||||
* (see the Hot bullet above for {@code exhaustedPattern}'s history).
|
||||
* {@code HerdrPeerLauncher} takes {@code Map.copyOf(profiles)} at construction and resolves
|
||||
* each spawn out of that copy, so those never reach a launch until the daemon restarts. A
|
||||
* reload logs these rather than pretending they applied.</li>
|
||||
@@ -583,12 +595,16 @@ public final class ConfigRef implements Supplier<FleetConfig> {
|
||||
* {@link #sameLaunchSettings} because they are read <em>live</em>, not baked in at spawn — see
|
||||
* the class doc's <em>Hot</em> bullet. {@code weight} and {@code maxLoad} are read live by the
|
||||
* placement policy on every spawn; {@code credentialId} is read live by
|
||||
* {@code CompositePeerLauncher} and the CB-578 stage B exhaustion sink. Nothing else is
|
||||
* {@code CompositePeerLauncher} and the CB-578 stage B exhaustion sink; {@code exhaustedPattern}
|
||||
* (fleetd #446) is read live, cached by profile name, by {@code LiveExhaustedPatterns} — see
|
||||
* that class's doc and the class doc's <em>Hot</em> bullet for the history (it used to be
|
||||
* compared here, deferred, like its sibling {@code errorPattern} still is). Nothing else is
|
||||
* excluded — see {@code sameLaunchSettingsComparesEveryProfileComponentOrExcludesIt} in
|
||||
* {@code ConfigRefProfileCoverageTest}, which enumerates every {@code Profile} record component
|
||||
* by reflection and fails the build if one is neither compared below nor named here.
|
||||
*/
|
||||
static final Set<String> LAUNCH_SETTINGS_EXCLUDED = Set.of("weight", "maxLoad", "credentialId");
|
||||
static final Set<String> LAUNCH_SETTINGS_EXCLUDED =
|
||||
Set.of("weight", "maxLoad", "credentialId", "exhaustedPattern");
|
||||
|
||||
/**
|
||||
* Whether two versions of a profile would launch a peer identically.
|
||||
@@ -625,13 +641,15 @@ public final class ConfigRef implements Supplier<FleetConfig> {
|
||||
&& Objects.equals(a.kind(), b.kind())
|
||||
&& Objects.equals(a.env(), b.env())
|
||||
&& Objects.equals(a.subscription(), b.subscription())
|
||||
// CB-578 stage B: exhaustedPattern is compiled once into Fleetd.main's pattern map
|
||||
// at startup (see ExhaustedPatternLookup wiring) — a reload never re-reads it, so a
|
||||
// changed pattern must be reported as deferred, exactly like model/baseUrl/argv.
|
||||
&& Objects.equals(a.exhaustedPattern(), b.exhaustedPattern())
|
||||
// fleetd #446: exhaustedPattern moved to LAUNCH_SETTINGS_EXCLUDED — it is now read
|
||||
// live, cached by profile name, through LiveExhaustedPatterns (see that class's doc
|
||||
// and ConfigRef's class doc Hot bullet), so it must NOT be compared here any more: a
|
||||
// reload that changes only exhaustedPattern must report "config reloaded", not
|
||||
// "these changes need a restart".
|
||||
// fleetd #201 Unit 5: errorPattern is compiled once into Fleetd.main's backend-error
|
||||
// pattern map at startup (see BackendErrorPatternLookup wiring), the same way
|
||||
// exhaustedPattern is — a reload never re-reads it either.
|
||||
// pattern map at startup (see BackendErrorPatternLookup wiring) — unlike its sibling
|
||||
// exhaustedPattern above (fleetd #446), a reload still never re-reads it; fleetd #446
|
||||
// scoped errorPattern out on purpose (see the class doc's Hot bullet).
|
||||
&& Objects.equals(a.errorPattern(), b.errorPattern())
|
||||
// fleetd #323 instance 1: ideProjectDir and ideOpenCommand are read at spawn off the
|
||||
// same frozen profile map as ideMcpUrl above (ClaudeCodeLauncher.java:267/269,
|
||||
|
||||
@@ -424,6 +424,12 @@ public record FleetConfig(
|
||||
* {@code null}/blank ⇒ the classification never fires for this profile and
|
||||
* today's completion-fallback behaviour is unchanged. Every backend words
|
||||
* its refusal differently, so this is config, never a vendor string in code.
|
||||
* Read live off the current config through {@code LiveExhaustedPatterns},
|
||||
* cached by profile name (fleetd #446), so it is HOT: an operator can arm
|
||||
* or disarm this profile's usage-limit detection by editing this key and
|
||||
* reloading, with no restart. Before fleetd #446 this was compiled once at
|
||||
* daemon startup, the same deferred shape its sibling {@code errorPattern}
|
||||
* (below) still has.
|
||||
* @param credentialId the shared account this profile authenticates as (CB-578 stage B). Two
|
||||
* or more profiles setting the <em>same</em> non-blank value are quarantined
|
||||
* together by one {@code exhaustedPattern} classification on any one of
|
||||
|
||||
@@ -3,7 +3,7 @@ package dev.ltms.fleet.herdr;
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
|
||||
/**
|
||||
* Client face onto the herdr daemon (protocol 14, herdr 0.7.0).
|
||||
* Client face onto the herdr daemon (protocol 19, herdr 0.8.0).
|
||||
*
|
||||
* <p>This is the ONLY thing in {@code fleetd} that speaks to herdr. Every method
|
||||
* maps to a herdr JSON-RPC call over its Unix domain socket. Requests are
|
||||
|
||||
@@ -8,7 +8,7 @@ import com.fasterxml.jackson.databind.node.ObjectNode;
|
||||
import java.nio.charset.StandardCharsets;
|
||||
|
||||
/**
|
||||
* Wire codec for herdr's newline-delimited JSON-RPC (protocol 14).
|
||||
* Wire codec for herdr's newline-delimited JSON-RPC (protocol 19).
|
||||
*
|
||||
* <p>Split out from the socket so the framing rules — the ones that actually bit us
|
||||
* during the spike (id MUST be a string; response carries {@code result} or
|
||||
|
||||
@@ -0,0 +1,91 @@
|
||||
package dev.ltms.fleet.inject;
|
||||
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.function.Supplier;
|
||||
import java.util.regex.Pattern;
|
||||
|
||||
/**
|
||||
* fleetd #446: the single live source for "is {@code exhaustedPattern} configured for this
|
||||
* profile, and what does it compile to". Read fresh off the config supplier on every call — the
|
||||
* same reason {@code CompositePeerLauncher#models0} is a live supplier read rather than a value
|
||||
* captured at construction (see that class's doc) — so an operator can arm or disarm usage-limit
|
||||
* detection for a profile by editing {@code exhaustedPattern} and reloading, with no restart.
|
||||
*
|
||||
* <p>Before this class, {@code exhaustedPattern} was compiled once into a {@code Map<String,
|
||||
* Pattern>} built inside {@code Fleetd.main} at startup ({@code ConfigRef}'s class doc used to
|
||||
* list it under <em>Deferred</em>, CB-578 stage A) — the asymmetry fleetd #446 exists to close:
|
||||
* an operator could turn a model off at runtime ({@code models.allow}'s {@code enabled: false},
|
||||
* hot since fleetd #422) but could not arm the detector that would tell them to, without a
|
||||
* restart. That was backwards for a feature whose whole point is to react while the fleet runs.
|
||||
*
|
||||
* <p>{@link #patternFor} is what {@link CompletionResolver} enforces on (via the {@link
|
||||
* ExhaustedPatternLookup} production wiring in {@code Fleetd.main}); {@link #armed} is what
|
||||
* {@code fleet_profiles}/{@code GET /profiles} report as {@code exhaustionDetectionArmed}. Both
|
||||
* read this ONE object, so the report can never disagree with the behaviour — the same rule
|
||||
* {@code CompositePeerLauncher#modelGateState()}'s javadoc states for the model gate's own
|
||||
* armed/off pair (fleetd #404): "armed" and "which models are off" must come from one read of the
|
||||
* same accessor the gate enforces on.
|
||||
*
|
||||
* <h2>Cache eviction</h2>
|
||||
* Cached by profile NAME, not by pattern text. A {@code Map<String, Pattern>} keyed by the
|
||||
* pattern STRING would grow by one entry per distinct regex ever typed for any profile across the
|
||||
* daemon's uptime — unbounded in practice, because tuning a regex to match a backend's exact
|
||||
* wording is exactly the kind of edit an operator makes several times while getting it right, and
|
||||
* every edit-and-reload cycle would leave the previous attempt's compiled {@link Pattern} behind
|
||||
* forever. Keyed by profile name instead, this cache holds at most one entry per profile name that
|
||||
* has ever been looked up — and that key space is bounded by the (small, human-authored) set of
|
||||
* configured profiles, which changes only on a restart: adding or removing a profile is itself a
|
||||
* deferred key (a new backend needs its own launcher, built once — see {@code ConfigRef}'s class
|
||||
* doc), so profile names do not churn the way pattern text does. Re-editing an EXISTING profile's
|
||||
* {@code exhaustedPattern} — the case this class exists to make hot — simply overwrites that
|
||||
* profile's one cache entry; it never adds a new one.
|
||||
*/
|
||||
public final class LiveExhaustedPatterns {
|
||||
|
||||
/** A compiled pattern paired with the source string it was compiled from, for change detection. */
|
||||
private record Cached(String source, Pattern compiled) {
|
||||
}
|
||||
|
||||
private final Supplier<Map<String, FleetConfig.Profile>> profiles;
|
||||
private final ConcurrentHashMap<String, Cached> cache = new ConcurrentHashMap<>();
|
||||
|
||||
public LiveExhaustedPatterns(Supplier<Map<String, FleetConfig.Profile>> profiles) {
|
||||
this.profiles = Objects.requireNonNull(profiles, "profiles");
|
||||
}
|
||||
|
||||
/**
|
||||
* The compiled {@code exhaustedPattern} currently configured for {@code profileName}, or
|
||||
* {@code null} when that profile is unknown or has none configured. Recompiles only when the
|
||||
* live pattern text differs from what is cached for this profile name; {@code
|
||||
* FleetConfig#rejectMalformedProfilePatterns} already refuses a config (at load and at reload)
|
||||
* whose {@code exhaustedPattern} does not compile, so this is not expected to throw in
|
||||
* production — it is not defended against here for that reason, the same trust
|
||||
* {@code CompositePeerLauncher#models0} places in config validation having already run.
|
||||
*/
|
||||
public Pattern patternFor(String profileName) {
|
||||
if (profileName == null) {
|
||||
return null;
|
||||
}
|
||||
FleetConfig.Profile profile = profiles.get().get(profileName);
|
||||
if (profile == null || !profile.hasExhaustedPattern()) {
|
||||
return null;
|
||||
}
|
||||
String source = profile.exhaustedPattern();
|
||||
Cached cached = cache.get(profileName);
|
||||
if (cached != null && cached.source().equals(source)) {
|
||||
return cached.compiled();
|
||||
}
|
||||
Cached fresh = new Cached(source, Pattern.compile(source));
|
||||
cache.put(profileName, fresh);
|
||||
return fresh.compiled();
|
||||
}
|
||||
|
||||
/** Whether {@code profileName} currently has a usage-limit pattern configured (live). */
|
||||
public boolean armed(String profileName) {
|
||||
return patternFor(profileName) != null;
|
||||
}
|
||||
}
|
||||
@@ -121,17 +121,46 @@ public final class FleetMcp {
|
||||
* CB-578 stage B quarantine facts used by {@code fleet_profiles}: a profile → credential id
|
||||
* lookup, plus the shared {@link BackendQuarantine} to read remaining cooldowns off.
|
||||
*
|
||||
* @param exhaustedPatternArmed fleetd #395: profile → whether that profile's {@code
|
||||
* exhaustedPattern} is configured (see {@code
|
||||
* @param exhaustedPatternArmed fleetd #395, made live by fleetd #446: profile → whether that
|
||||
* profile's {@code exhaustedPattern} is configured right now (see {@code
|
||||
* FleetConfig.Profile#hasExhaustedPattern}), i.e. whether a backend refusal
|
||||
* on it can EVER be classified {@code BACKEND_EXHAUSTED} and quarantine its
|
||||
* credential. Bundled here, not a separate Source, because it answers the
|
||||
* exact question {@code fleet_profiles}'s quarantine facts already answer
|
||||
* for a QUARANTINED profile — "can this profile's usage limit ever be
|
||||
* caught?" — just for every profile, not only one currently caught.
|
||||
* caught?" — just for every profile, not only one currently caught. In
|
||||
* production this reads {@link dev.ltms.fleet.inject.LiveExhaustedPatterns#armed} —
|
||||
* the SAME live accessor {@link dev.ltms.fleet.inject.CompletionResolver}
|
||||
* enforces on — so a reload that arms or disarms detection is reported
|
||||
* correctly on the very next call, no restart (see that class's doc).
|
||||
* @param modelFor fleetd #446 criterion 3: profile → the {@code model:} it runs, or
|
||||
* {@code null} when the profile names none. Read live off the current
|
||||
* config, the same way {@code credentialIdFor} already is. Used to name,
|
||||
* in a QUARANTINED profile's row, which of the operator's configured models
|
||||
* the fix (\"set {@code enabled: false} on it under {@code models.allow}\")
|
||||
* actually applies to.
|
||||
* @param reasonFor fleetd #446 criterion 3: credential id → the backend text that triggered
|
||||
* its most recent quarantine, or {@code null} when none is known (e.g. an
|
||||
* inert source, or a quarantine recorded before this field existed). Never
|
||||
* consulted on its own to decide whether a credential is quarantined —
|
||||
* callers gate on {@link BackendQuarantine#remainingSeconds} first, exactly
|
||||
* like {@code credentialId} itself, so a stale reason left behind after a
|
||||
* quarantine expires is never surfaced.
|
||||
*/
|
||||
public record QuarantineSource(Function<String, String> credentialIdFor, BackendQuarantine quarantine,
|
||||
Function<String, Boolean> exhaustedPatternArmed) {
|
||||
Function<String, Boolean> exhaustedPatternArmed,
|
||||
Function<String, String> modelFor,
|
||||
Function<String, String> reasonFor) {
|
||||
/**
|
||||
* Backward-compatible 3-arg form, before fleetd #446 added {@code modelFor}/{@code
|
||||
* reasonFor} — neither is ever reported. Keeps every pre-existing call site (production and
|
||||
* test) compiling and behaving identically for the quarantine facts they actually asked for.
|
||||
*/
|
||||
public QuarantineSource(Function<String, String> credentialIdFor, BackendQuarantine quarantine,
|
||||
Function<String, Boolean> exhaustedPatternArmed) {
|
||||
this(credentialIdFor, quarantine, exhaustedPatternArmed, _ -> null, _ -> null);
|
||||
}
|
||||
|
||||
/**
|
||||
* Backward-compatible 2-arg form, before fleetd #395 added {@code exhaustedPatternArmed} —
|
||||
* reports every profile unarmed. Keeps every pre-existing call site (production and test)
|
||||
@@ -426,7 +455,8 @@ public final class FleetMcp {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
|
||||
leadSeats, callers == null ? Map.of() : callers.leads(),
|
||||
callerTerminal(exchange),
|
||||
new CoordinationSource(leadChannel, peers));
|
||||
new CoordinationSource(leadChannel, peers),
|
||||
coordinatorVisibleTo(principal(exchange)));
|
||||
};
|
||||
BiFunction<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> stopHandler =
|
||||
(exchange, req) -> {
|
||||
@@ -579,6 +609,20 @@ public final class FleetMcp {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #439: only the primary may read {@code fleet_list}'s {@code coordinator} row —
|
||||
* lead-to-lead coordination state (coord-ids, mailbox facts, held-message previews), never the
|
||||
* roster. Split out of the {@code fleet_list} handler, same reason as {@link #denyFor} and
|
||||
* {@link #recordPrimarySingleton}: the decision must be unit-testable without fabricating an
|
||||
* SDK {@code McpSyncServerExchange}, and the handler must call this named predicate rather than
|
||||
* inlining the check, so a future edit cannot silently pass a literal instead of asking who
|
||||
* called ({@code FleetMcpAuthzTest.theFleetListHandlerActuallyConsultsCoordinatorVisibleTo}
|
||||
* reads the source and asserts the handler calls this method by name, not a literal).
|
||||
*/
|
||||
static boolean coordinatorVisibleTo(Principal caller) {
|
||||
return caller.isPrimary();
|
||||
}
|
||||
|
||||
/** The worker identity resolved from this call's connection, or {@code null} if the primary. */
|
||||
private static String callerTerminal(McpSyncServerExchange exchange) {
|
||||
Object v = exchange.transportContext().get(CALLER_TERMINAL);
|
||||
@@ -994,7 +1038,11 @@ public final class FleetMcp {
|
||||
if (isBlank(target) || isBlank(msgId)) {
|
||||
return error("target and msgId are required");
|
||||
}
|
||||
messages.ackReply(target, msgId);
|
||||
if (!messages.ackReply(target, msgId)) {
|
||||
return error(msgId + " is not in " + target + "'s reply inbox (wrong id, wrong target, "
|
||||
+ "or already acked). Held lead-to-lead (peer) mail cannot be acked this way — "
|
||||
+ "read it with fleet_poll{coordId}.");
|
||||
}
|
||||
return text("acknowledged " + msgId);
|
||||
}
|
||||
|
||||
@@ -1206,13 +1254,21 @@ public final class FleetMcp {
|
||||
* #284 was, where one rule computed in two places was widened in only one and a single response
|
||||
* contradicted itself. Shared inputs do not make duplicated computation safe.
|
||||
*
|
||||
* <p>fleetd #395: also reports {@code exhaustionDetectionArmed}, one boolean per configured
|
||||
* profile — {@code true} when that profile's {@code exhaustedPattern} is set, {@code false}
|
||||
* when it is not, so an operator can tell "this profile is healthy" from "nothing can ever
|
||||
* quarantine this profile" without reading {@code fleetd.yaml}. Unlike {@code quarantined}/
|
||||
* {@code coolingOff}, this map always names every profile: an unarmed profile never enters a
|
||||
* transient state to be absent from, so silence here would read as "healthy" rather than "not
|
||||
* being watched at all".
|
||||
* <p>fleetd #395, made live by fleetd #446: also reports {@code exhaustionDetectionArmed}, one
|
||||
* boolean per configured profile — {@code true} when that profile's {@code exhaustedPattern} is
|
||||
* set RIGHT NOW, {@code false} when it is not, so an operator can tell "this profile is
|
||||
* healthy" from "nothing can ever quarantine this profile" without reading {@code fleetd.yaml}.
|
||||
* Unlike {@code quarantined}/{@code coolingOff}, this map always names every profile: an
|
||||
* unarmed profile never enters a transient state to be absent from, so silence here would read
|
||||
* as "healthy" rather than "not being watched at all". Hot since fleetd #446: a reload that
|
||||
* arms or disarms a profile's {@code exhaustedPattern} changes this map's answer on the very
|
||||
* next call, no restart — see {@link dev.ltms.fleet.inject.LiveExhaustedPatterns}'s class doc.
|
||||
*
|
||||
* <p>fleetd #446 criterion 3: each {@code quarantined} row also names {@code model} (the
|
||||
* profile's configured {@code model:}, omitted when the profile names none) and {@code reason}
|
||||
* (the backend text that triggered the most recent quarantine of that credential, omitted when
|
||||
* none is known) — so a lead can see WHICH model to turn off and WHY, without reading the
|
||||
* daemon log.
|
||||
*/
|
||||
public static Map<String, Object> profilesView(PeerLauncher workers, QuarantineSource quarantine, OutageSource outage) {
|
||||
Map<String, Object> result = new LinkedHashMap<>();
|
||||
@@ -1229,6 +1285,14 @@ public final class FleetMcp {
|
||||
Map<String, Object> row = new LinkedHashMap<>();
|
||||
row.put("credentialId", credentialId);
|
||||
row.put("quarantinedForSeconds", remaining);
|
||||
String model = quarantine.modelFor().apply(profile);
|
||||
if (model != null && !model.isBlank()) {
|
||||
row.put("model", model);
|
||||
}
|
||||
String reason = quarantine.reasonFor().apply(credentialId);
|
||||
if (reason != null && !reason.isBlank()) {
|
||||
row.put("reason", reason);
|
||||
}
|
||||
quarantined.put(profile, row);
|
||||
});
|
||||
}
|
||||
@@ -1340,12 +1404,48 @@ public final class FleetMcp {
|
||||
LeadSeatSource.none(), leads, selfTerm, coordination);
|
||||
}
|
||||
|
||||
/** As above, plus fleetd #176 lead-seat facts (see {@link LeadSeatSource}). */
|
||||
/**
|
||||
* As above, plus fleetd #176 lead-seat facts (see {@link LeadSeatSource}).
|
||||
*
|
||||
* <p>Assumes the caller is <strong>not</strong> the primary (fleetd #463) — every wrapper
|
||||
* overload above delegates here without carrying a caller identity, which is exactly right for
|
||||
* them: they exist for call sites (and unit tests) that have no {@link Principal} to hand over,
|
||||
* and a missing identity should fail closed rather than fail open onto lead-to-lead state. A
|
||||
* test that wants the {@code coordinator} row must call the canonical overload below with an
|
||||
* explicit {@code true}. The one call site that has a real caller ({@code fleet_list}'s MCP
|
||||
* handler) uses {@link #listFleet(PeerLauncher, SessionManager, MessageService, CapacitySource,
|
||||
* HealthCoverageSource, QuarantineSource, OutageSource, LeadSeatSource, Map, String,
|
||||
* CoordinationSource, boolean)} instead, so it can pass the true answer.
|
||||
*/
|
||||
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
|
||||
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
QuarantineSource quarantine, OutageSource outage,
|
||||
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
|
||||
CoordinationSource coordination) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
|
||||
leadSeats, leads, selfTerm, coordination, false);
|
||||
}
|
||||
|
||||
/**
|
||||
* As above, gated by the caller's role (fleetd #439). The {@code coordinator} row is
|
||||
* lead-to-lead coordination state — coordination between orchestrators, not roster
|
||||
* observation — so it is assembled and included only when {@code callerIsPrimary} is
|
||||
* {@code true}. A worker or an architect gets a result with the {@code coordinator} key
|
||||
* <strong>absent</strong>, never an empty or redacted one, and never pays the cost of
|
||||
* {@link #coordinatorView} probing peer mailboxes for a row it will not receive.
|
||||
*
|
||||
* @param callerIsPrimary whether the {@code fleet_list} caller is the primary; only the MCP
|
||||
* handler computes this from the real connection (see
|
||||
* {@code Principal#isPrimary()}) — every other overload passes
|
||||
* {@code false} (fleetd #463: a forgotten argument fails closed, not
|
||||
* open), so a test that wants the {@code coordinator} row must pass
|
||||
* an explicit {@code true}
|
||||
*/
|
||||
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
|
||||
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
QuarantineSource quarantine, OutageSource outage,
|
||||
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
|
||||
CoordinationSource coordination, boolean callerIsPrimary) {
|
||||
try {
|
||||
Map<String, Agent> live = workers.list().stream()
|
||||
.map(Agent.class::cast)
|
||||
@@ -1367,9 +1467,14 @@ public final class FleetMcp {
|
||||
Map<String, Object> result = new LinkedHashMap<>();
|
||||
result.put("leads", leadRows); result.put("members", out);
|
||||
result.put("healthCoverage", healthCoverage.value().get());
|
||||
Map<String, Object> coordinatorRow = coordinatorView(coordination);
|
||||
if (coordinatorRow != null) {
|
||||
result.put("coordinator", coordinatorRow);
|
||||
// fleetd #439: coordinator/coordinatorView is lead-to-lead coordination state and must
|
||||
// never reach a worker or an architect -- gate BEFORE assembling it, not after, so the
|
||||
// key is absent rather than present-and-empty.
|
||||
if (callerIsPrimary) {
|
||||
Map<String, Object> coordinatorRow = coordinatorView(coordination);
|
||||
if (coordinatorRow != null) {
|
||||
result.put("coordinator", coordinatorRow);
|
||||
}
|
||||
}
|
||||
if (capacity.available()) result.put("capacity", profiles.stream()
|
||||
.map(profile -> capacityView(profile, capacity.liveCount(), capacity.maxLoad(), roster, messages,
|
||||
@@ -1766,7 +1871,10 @@ public final class FleetMcp {
|
||||
+ "has processed a reply and wants to confirm it, leaving other pending replies "
|
||||
+ "in the inbox for later drain.",
|
||||
objectSchema(Map.of(
|
||||
"target", stringProp("Worker session id whose inbox to ack from"),
|
||||
"target", stringProp("Worker session id whose inbox to ack from. Must name a "
|
||||
+ "reply actually queued for it — an id in no inbox, or a coord-id "
|
||||
+ "(peer held mail, read with fleet_poll{coordId} instead), errors "
|
||||
+ "rather than reporting a false success"),
|
||||
"msgId", stringProp("The message id to acknowledge")),
|
||||
List.of("target", "msgId")));
|
||||
}
|
||||
@@ -1808,11 +1916,16 @@ public final class FleetMcp {
|
||||
"List the configured worker profiles (backends) and which one fleet_spawn uses by "
|
||||
+ "default. A 'quarantined' map is present when a backend-exhausted refusal put "
|
||||
+ "a profile's credential on cooldown — fleet_spawn onto it is refused until "
|
||||
+ "quarantinedForSeconds elapses; a profile sharing that credential is listed too. "
|
||||
+ "quarantinedForSeconds elapses; a profile sharing that credential is listed too, "
|
||||
+ "each row naming 'model' (the model that profile runs, when configured) and "
|
||||
+ "'reason' (the backend text that triggered the quarantine, when known) — the fix "
|
||||
+ "is usually `enabled: false` on that model under models.allow. "
|
||||
+ "'exhaustionDetectionArmed' reports, per profile, whether a usage-limit refusal "
|
||||
+ "on it can EVER be classified and quarantined (its exhaustedPattern is "
|
||||
+ "configured) — false means that profile's credential can never be quarantined "
|
||||
+ "by this mechanism, however many usage-limit refusals it sees.",
|
||||
+ "by this mechanism, however many usage-limit refusals it sees. Both "
|
||||
+ "exhaustedPattern and models.allow's on/off state are hot: editing fleetd.yaml "
|
||||
+ "and reloading arms/disarms detection or flips a model off with no restart.",
|
||||
objectSchema(Map.of(), List.of()));
|
||||
}
|
||||
|
||||
|
||||
@@ -331,11 +331,22 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
+ "form — opencode's per-model context limit could not be applied for this profile",
|
||||
cfg.profile(), cfg.model());
|
||||
}
|
||||
// fleetd #393: memberSkills seeding (GitWorktrees#seedSkills) copies skill folders into
|
||||
// EVERY provisioned worktree's .claude/skills/ regardless of which kind ultimately spawns
|
||||
// into it — that copy step cannot know the kind, only the caller of GitWorktrees#add does
|
||||
// (see that method's own javadoc). .claude/skills/ is a Claude Code CLI convention the CLI
|
||||
// discovers on its own; opencode has no such discovery, so without this, a seeded skill
|
||||
// never reaches an opencode member even though GitWorktrees logged it as seeded. Read
|
||||
// whatever landed under <cwd>/.claude/skills/ here — the one place in this launcher that
|
||||
// knows both the kind (opencode, by construction: this IS OpenCodeLauncher) and the cwd.
|
||||
List<Path> skillInstructionFiles = skillInstructionFiles(spec.cwd());
|
||||
// A config file is needed for the bridge MCP mount, a member charter, the IDE MCP (+ its
|
||||
// guidance overlay, CB-634), a pinned endpoint (CB-508), or a resolvable autoCompactWindow.
|
||||
// guidance overlay, CB-634), a pinned endpoint (CB-508), a resolvable autoCompactWindow, or
|
||||
// at least one seeded skill to deliver via instructions[] (fleetd #393).
|
||||
if (cfg.hasMcp() || cfg.hasIdeMcp() || spec.charter() != null || hasCustomProvider(cfg)
|
||||
|| wantsContextLimit) {
|
||||
workerEnv.put("OPENCODE_CONFIG", writeConfig(cfg, spec.charter(), spec.cwd()).toString());
|
||||
|| wantsContextLimit || !skillInstructionFiles.isEmpty()) {
|
||||
workerEnv.put("OPENCODE_CONFIG",
|
||||
writeConfig(cfg, spec.charter(), spec.cwd(), skillInstructionFiles).toString());
|
||||
}
|
||||
applyGitToken(workerEnv, cfg);
|
||||
List<String> argv = argvWithResume(argvWithModel(argvWithAuto(cfg), cfg), spec.resumeSessionId());
|
||||
@@ -358,6 +369,65 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
return withAgent;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #393: the {@code SKILL.md} paths under {@code <cwd>/.claude/skills/} this launcher can
|
||||
* turn into {@code instructions[]} entries, plus the honest log this ticket asks for — emitted
|
||||
* here, at the one point a skill's fate for THIS spawn is actually known, rather than trusting
|
||||
* {@code GitWorktrees#seedSkills}'s kind-blind "N of M" line to mean "and it will be read."
|
||||
*
|
||||
* <p>Every non-hidden subdirectory of {@code .claude/skills/} is a candidate, whether it got
|
||||
* there via {@code memberSkills:} seeding or because the target repo ships its own copy — this
|
||||
* launcher does not care which; it only cares what it can find at spawn time. A candidate with
|
||||
* a {@code SKILL.md} at its top level (the same shape {@link #writeConfig} already requires for
|
||||
* the charter and IDE-rules instructions entries) is delivered; anything else is a directory
|
||||
* this launcher cannot turn into a flat instructions entry, named explicitly in the log rather
|
||||
* than silently dropped, so a caller sees a real "cannot consume" reason and not just a smaller
|
||||
* number than {@code GitWorktrees}' own count.
|
||||
*
|
||||
* <p>No candidates at all (directory absent or empty) logs nothing — the same
|
||||
* no-log-when-nothing-to-say shape {@link #hasCustomProvider} and friends already follow, and
|
||||
* the shape {@code GitWorktrees#seedSkills} itself uses when {@code memberSkills:} is unset.
|
||||
* A failure to even list the directory is logged and treated as "nothing delivered" — best
|
||||
* effort, must never fail the spawn, matching {@code GitWorktrees#seedSkills}'s own contract.
|
||||
*/
|
||||
private List<Path> skillInstructionFiles(String cwd) {
|
||||
if (cwd == null || cwd.isBlank()) {
|
||||
return List.of();
|
||||
}
|
||||
Path skillsDir = Path.of(cwd, ".claude", "skills");
|
||||
if (!Files.isDirectory(skillsDir)) {
|
||||
return List.of();
|
||||
}
|
||||
List<Path> candidates;
|
||||
try (var listing = Files.list(skillsDir)) {
|
||||
candidates = listing.filter(Files::isDirectory)
|
||||
.filter(p -> !p.getFileName().toString().startsWith("."))
|
||||
.sorted()
|
||||
.toList();
|
||||
} catch (IOException e) {
|
||||
log.warn("could not scan {} for skill folders to deliver to this opencode member: {}",
|
||||
skillsDir, e.getMessage());
|
||||
return List.of();
|
||||
}
|
||||
if (candidates.isEmpty()) {
|
||||
return List.of();
|
||||
}
|
||||
List<Path> delivered = candidates.stream()
|
||||
.map(dir -> dir.resolve("SKILL.md"))
|
||||
.filter(Files::isRegularFile)
|
||||
.toList();
|
||||
List<String> undeliverable = candidates.stream()
|
||||
.filter(dir -> !Files.isRegularFile(dir.resolve("SKILL.md")))
|
||||
.map(dir -> dir.getFileName().toString())
|
||||
.toList();
|
||||
log.info("skill delivery: {} of {} skill folder(s) under {} reached this opencode member via "
|
||||
+ "instructions[] (opencode does not read .claude/skills/ natively, unlike "
|
||||
+ "Claude Code){}",
|
||||
delivered.size(), candidates.size(), skillsDir,
|
||||
undeliverable.isEmpty() ? "" : "; no SKILL.md, could not be delivered: " + undeliverable);
|
||||
return delivered;
|
||||
}
|
||||
|
||||
/**
|
||||
* True when this profile pins its own OpenAI-compatible endpoint (CB-508) rather than using
|
||||
* whatever provider opencode resolves by default.
|
||||
@@ -418,8 +488,15 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
* fresh per-spawn directory under {@link #configRoot}, and return the config file's path for
|
||||
* {@code OPENCODE_CONFIG}. The dir is unique per spawn so concurrent workers never race on it;
|
||||
* it is best-effort cleaned on JVM exit (worker config is disposable — regenerated every spawn).
|
||||
*
|
||||
* @param skillInstructionFiles fleetd #393: absolute {@code SKILL.md} paths from
|
||||
* {@link #skillInstructionFiles(String)}, appended to
|
||||
* {@code instructions[]} so a {@code memberSkills:}-seeded skill
|
||||
* reaches this opencode member the same way the charter and IDE
|
||||
* rules already do.
|
||||
*/
|
||||
private Path writeConfig(FleetConfig.Profile cfg, String charterText, String cwd) {
|
||||
private Path writeConfig(FleetConfig.Profile cfg, String charterText, String cwd,
|
||||
List<Path> skillInstructionFiles) {
|
||||
try {
|
||||
Path dir = Files.createTempDirectory(configParentDir(), "fleetd-opencode-");
|
||||
dir.toFile().deleteOnExit();
|
||||
@@ -447,7 +524,28 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
Files.writeString(charter, charterText);
|
||||
charter.toFile().deleteOnExit();
|
||||
|
||||
root.putArray("instructions").add(charter.toAbsolutePath().toString());
|
||||
// fleetd #393 follow-up: withArray, not putArray. putArray REPLACES whatever node
|
||||
// is already at "instructions" — harmless only as long as this block runs first
|
||||
// against a still-empty root, which is an ordering constraint nothing declared or
|
||||
// tested. The skills writer just below, and the IDE-rules writer further down,
|
||||
// both already use withArray (get-or-create) for exactly this reason; this was the
|
||||
// one straggler. Proven load-bearing on the fleetd #393 merge: flipping this one
|
||||
// call back to putArray left the whole suite green while silently deleting the
|
||||
// charter entry whenever skills or IDE rules ran after it — an opencode member
|
||||
// would launch with no role contract at all, worse than the bug #393 fixed, and
|
||||
// nothing caught it. See OpenCodeLauncherTest's
|
||||
// instructionsArrayHoldsCharterThenIdeRulesInOrder and
|
||||
// instructionsArrayHoldsCharterThenSkillsThenIdeRulesInOrder.
|
||||
root.withArray("instructions").add(charter.toAbsolutePath().toString());
|
||||
}
|
||||
|
||||
// fleetd #393: each seeded skill's SKILL.md, delivered as a plain instructions[] entry
|
||||
// — the only mechanism opencode has for static guidance text. Unlike Claude Code's
|
||||
// Skill tool, opencode cannot load one of these on demand by name; the content is just
|
||||
// always part of the system prompt from spawn. That is a real difference in HOW the
|
||||
// content reaches the member, not a reason to withhold it.
|
||||
for (Path skillFile : skillInstructionFiles) {
|
||||
root.withArray("instructions").add(skillFile.toAbsolutePath().toString());
|
||||
}
|
||||
|
||||
if (cfg.hasMcp() || cfg.hasIdeMcp()) {
|
||||
|
||||
@@ -394,17 +394,17 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void ack(String target, String msgId) {
|
||||
public boolean ack(String target, String msgId) {
|
||||
var perTarget = held.get(target);
|
||||
if (perTarget == null || perTarget == RELEASED) {
|
||||
return;
|
||||
return false;
|
||||
}
|
||||
Held h;
|
||||
synchronized (perTarget) {
|
||||
h = perTarget.remove(msgId);
|
||||
}
|
||||
if (h == null) {
|
||||
return; // never held (or already acked) — no-op
|
||||
return false; // never held (or already acked) — no-op
|
||||
}
|
||||
try {
|
||||
synchronized (channelLock) {
|
||||
@@ -418,6 +418,7 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
}
|
||||
throw new IllegalStateException("cannot ack reply " + msgId + " on " + queueName(target), e);
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
private DeliverCallback deliverCallback(String target) {
|
||||
|
||||
@@ -59,16 +59,17 @@ public final class InMemoryReplyInbox implements ReplyInbox {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void ack(String target, String msgId) {
|
||||
public boolean ack(String target, String msgId) {
|
||||
if (!owned.contains(target)) {
|
||||
return;
|
||||
return false;
|
||||
}
|
||||
var perTarget = store.get(target);
|
||||
if (perTarget != null) {
|
||||
//noinspection SynchronizationOnLocalVariableOrMethodParameter
|
||||
synchronized (perTarget) {
|
||||
perTarget.remove(msgId);
|
||||
}
|
||||
if (perTarget == null) {
|
||||
return false;
|
||||
}
|
||||
//noinspection SynchronizationOnLocalVariableOrMethodParameter
|
||||
synchronized (perTarget) {
|
||||
return perTarget.remove(msgId) != null;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -826,9 +826,13 @@ public final class MessageService {
|
||||
/**
|
||||
* Acknowledge a specific reply by {@code msgId} for {@code target}. Removes it from the inbox
|
||||
* so that a subsequent drain or peek no longer returns it.
|
||||
*
|
||||
* @return {@code true} if an entry was actually removed, {@code false} if {@code msgId} was not
|
||||
* in {@code target}'s inbox (wrong id, wrong target, or already acked). The caller —
|
||||
* {@link dev.ltms.fleet.mcp.FleetMcp#ack} — must not report success on {@code false}.
|
||||
*/
|
||||
public void ackReply(String target, String msgId) {
|
||||
inbox.ack(target, msgId);
|
||||
public boolean ackReply(String target, String msgId) {
|
||||
return inbox.ack(target, msgId);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -46,6 +46,13 @@ public interface ReplyInbox {
|
||||
/** Non-destructive snapshot of pending replies for {@code target} (FIFO), empty list if none. */
|
||||
List<InboxMessage> peek(String target);
|
||||
|
||||
/** Remove the reply {@code msgId} for {@code target} once the primary has taken it. No-op if absent. */
|
||||
void ack(String target, String msgId);
|
||||
/**
|
||||
* Remove the reply {@code msgId} for {@code target} once the primary has taken it.
|
||||
*
|
||||
* @return {@code true} if an entry was actually removed, {@code false} if there was nothing to
|
||||
* remove (unknown {@code target}, unowned {@code target}, or a {@code msgId} not held for
|
||||
* it). A {@code false} is not an error — acking a {@code target} this daemon does not own is
|
||||
* part of the normal contract, not a failure.
|
||||
*/
|
||||
boolean ack(String target, String msgId);
|
||||
}
|
||||
|
||||
@@ -168,6 +168,23 @@ public interface PeerLauncher {
|
||||
* answer for a launcher with no role-pool concept of its own (e.g. a single {@code
|
||||
* HerdrPeerLauncher} adapter, which is never reached this way in production: {@code
|
||||
* CompositePeerLauncher} always fronts it and resolves roles itself).
|
||||
*
|
||||
* <p>fleetd #453: this default is deliberately <em>not</em> abstract — unlike {@link
|
||||
* #spawn(SpawnRequest, PlacementDecision)} (fleetd #450), there is no live defect in inheriting
|
||||
* it today, and the only current single-adapter implementer ({@code HerdrPeerLauncher}) is
|
||||
* correct to do so. But it stays correct only as long as that holds: <strong>if a launcher ever
|
||||
* routes more than one profile per role, it MUST override this method</strong>, or every role
|
||||
* silently resolves to {@link #defaultProfile()} with no error and no log line. {@code
|
||||
* HerdrPeerLauncher.spawn(SpawnRequest, PlacementDecision)} — the override in {@code
|
||||
* dev.ltms.fleet.member}, not the declaration below — names this method and {@link #place}
|
||||
* explicitly as "unoverridden here" for exactly this reason. Read it before adding role-pool
|
||||
* routing to any {@code HerdrPeerLauncher} subclass.
|
||||
*
|
||||
* <p>Who is forced to read which paragraph, because it is not symmetric. A new class that
|
||||
* implements this interface directly must write a body for {@link #spawn(SpawnRequest,
|
||||
* PlacementDecision)}, which is abstract here, so it lands on this javadoc. A subclass of
|
||||
* {@code HerdrPeerLauncher} does not: that class already implements the method, and the
|
||||
* subclass inherits the body. So for a subclass this paragraph is advice, not a gate.
|
||||
*/
|
||||
default String defaultProfileFor(MemberRole role) {
|
||||
return defaultProfile();
|
||||
@@ -228,6 +245,13 @@ public interface PeerLauncher {
|
||||
* placement condition — the right answer for a launcher with no pool or placement-policy
|
||||
* concept of its own, matching {@link #defaultProfileFor}'s own default.
|
||||
*
|
||||
* <p>fleetd #453: same reasoning as {@link #defaultProfileFor}'s own #453 note — this default
|
||||
* is deliberately not abstract (no live defect today, correct for the sole single-adapter
|
||||
* implementer), but <strong>a launcher that ever routes more than one profile per role MUST
|
||||
* override this method too</strong>, or placement silently ignores {@code role} for it. See
|
||||
* {@code HerdrPeerLauncher.spawn(SpawnRequest, PlacementDecision)}'s javadoc, which names this
|
||||
* method as "unoverridden here" and why that is correct only for a single-profile adapter.
|
||||
*
|
||||
* @throws RuntimeException (implementation-specific, typically a placement exception) if no
|
||||
* candidate in {@code role}'s pool is currently placeable
|
||||
*/
|
||||
|
||||
@@ -690,14 +690,18 @@ public final class GitWorktrees implements Worktrees {
|
||||
* {@code fleet.seededSkillsNote}, readable with {@code git config --worktree --get-all
|
||||
* fleet.seededSkills}.
|
||||
*
|
||||
* <p><b>Claude Code specific by construction, not by a backend check here.</b> Only {@code
|
||||
* .claude/skills/<name>/SKILL.md} is a path any launcher reads today (opencode's equivalent is a
|
||||
* different shape under {@code .opencode/agent}, out of scope — see issue #362). This method
|
||||
* only copies files; like {@link #isolateToolSurface} — which neutralizes BOTH {@code .mcp.json}
|
||||
* and {@code opencode.json} unconditionally — it runs the same for every worktree regardless of
|
||||
* which backend ultimately spawns into it, because the backend is not yet chosen at {@link #add}
|
||||
* time. A seeded {@code .claude/skills/} directory in an opencode member's worktree is simply
|
||||
* never read by that launcher.
|
||||
* <p><b>Kind-blind by construction, not by a backend check here — this used to be a real gap
|
||||
* (fleetd #393).</b> This method only copies files; like {@link #isolateToolSurface} — which
|
||||
* neutralizes BOTH {@code .mcp.json} and {@code opencode.json} unconditionally — it runs the
|
||||
* same for every worktree regardless of which backend ultimately spawns into it, because no
|
||||
* caller of {@link #add} hands this class a kind to consult. Before fleetd #393, that made the
|
||||
* log line below a false claim of success for a {@code kind: opencode} member: opencode has no
|
||||
* built-in discovery of {@code .claude/skills/}, unlike the Claude Code CLI, so a seeded skill
|
||||
* never reached one. It now does — {@code OpenCodeLauncher#skillInstructionFiles} reads
|
||||
* whatever this method copied into {@code .claude/skills/} and appends each {@code SKILL.md} to
|
||||
* the generated {@code instructions[]} — but that delivery, and the log line that honestly
|
||||
* claims it (kind-aware, unlike this one), happens at the launcher, once the kind is actually
|
||||
* known, not here.
|
||||
*/
|
||||
private void seedSkills(String worktreePath) {
|
||||
if (memberSkillsSource == null) {
|
||||
@@ -739,7 +743,15 @@ public final class GitWorktrees implements Worktrees {
|
||||
if (detail.isEmpty()) {
|
||||
detail = "no skill folders found under " + source;
|
||||
}
|
||||
log.info("skill seeding: {} of {} candidate(s) from {} into {}/.claude/skills — {}",
|
||||
// fleetd #393: this only claims the copy step, deliberately — it cannot know the member
|
||||
// kind that will spawn into this worktree (see this method's own javadoc), so it must not
|
||||
// read as "and the member will act on it." Whether that is true depends on the kind: the
|
||||
// Claude Code CLI discovers .claude/skills/ on its own; OpenCodeLauncher logs its own
|
||||
// "skill delivery" line, once the kind is known, naming what it could and could not turn
|
||||
// into instructions[].
|
||||
log.info("skill seeding: {} of {} candidate(s) from {} into {}/.claude/skills — {} "
|
||||
+ "(whether the spawned member can act on this depends on its kind — see "
|
||||
+ "the launcher's own log for that)",
|
||||
seeded.size(), seeded.size() + kept.size(), source, worktreePath, detail);
|
||||
if (seeded.isEmpty()) {
|
||||
return;
|
||||
|
||||
@@ -2,6 +2,7 @@ package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.inject.LiveExhaustedPatterns;
|
||||
import dev.ltms.fleet.mcp.FleetMcp;
|
||||
import dev.ltms.fleet.placement.BackendQuarantine;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
@@ -11,16 +12,31 @@ import org.junit.jupiter.api.io.TempDir;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.Map;
|
||||
import java.util.regex.Pattern;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #404: {@code exhaustionDetectionArmed} must describe the startup pattern map, not the
|
||||
* reloaded config snapshot.
|
||||
* fleetd #446: {@code exhaustionDetectionArmed} must describe the LIVE config, not a startup
|
||||
* snapshot — the opposite of what fleetd #404 asked for. #404 pinned {@code exhaustedPattern} as
|
||||
* compiled once at startup (matching {@code CompletionResolver}'s then-frozen behaviour); #446's
|
||||
* whole point is that both sides — the classification {@code CompletionResolver} enforces and the
|
||||
* {@code exhaustionDetectionArmed} field this reports — now read the SAME live {@link
|
||||
* LiveExhaustedPatterns} instance, so a reload arms or disarms detection with no restart, and the
|
||||
* report can never disagree with the behaviour (the fleetd #404 rule, now upheld for real instead
|
||||
* of by freezing both sides).
|
||||
*
|
||||
* <p>This test needs a reload. At startup the two snapshots agree, so a test of only a newly
|
||||
* started daemon would not detect a live {@code config.get()} lookup in the report field.
|
||||
*
|
||||
* <p><b>The hotness proof criterion 1 demands:</b> run this against the fixed code (green) and
|
||||
* against the pre-#446 compile site — {@code Fleetd.main}'s {@code exhaustedPatternsByProfile}
|
||||
* compiled once into a {@code Map<String, Pattern>} at startup, handed to {@code quarantineSource}
|
||||
* — and show it fails there. See the class doc history above: that shape is exactly what
|
||||
* {@code reloadedPatternArmsDetectionWithNoRestart} below is written to catch, and it is the test
|
||||
* that used to assert the opposite (see git history for this file's pre-#446 version, which
|
||||
* asserted {@code assertFalse(...)} on the identical scenario this now asserts {@code assertTrue}
|
||||
* on).
|
||||
*/
|
||||
class FleetdExhaustionDetectionArmedWiringTest {
|
||||
|
||||
@@ -54,41 +70,68 @@ class FleetdExhaustionDetectionArmedWiringTest {
|
||||
""";
|
||||
|
||||
@Test
|
||||
@DisplayName("reloading an exhaustedPattern does not arm the startup detection source")
|
||||
void reloadedPatternDoesNotChangeTheArmedFieldUntilRestart(@TempDir Path dir) throws Exception {
|
||||
@DisplayName("reloading a profile's exhaustedPattern IN arms exhaustionDetectionArmed, no restart")
|
||||
void reloadedPatternArmsDetectionWithNoRestart(@TempDir Path dir) throws Exception {
|
||||
Path file = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(file, NO_PATTERN);
|
||||
ConfigRef config = new ConfigRef(file, FleetConfig.load(file));
|
||||
|
||||
FleetMcp.QuarantineSource before = Fleetd.quarantineSource(config, BackendQuarantine.none(),
|
||||
new LiveExhaustedPatterns(() -> config.get().profiles()), Map.of());
|
||||
assertFalse(before.exhaustedPatternArmed().apply("terra"),
|
||||
"no exhaustedPattern configured yet — must report unarmed");
|
||||
|
||||
Files.writeString(file, WITH_PATTERN);
|
||||
assertTrue(config.reload().applied());
|
||||
assertTrue(config.get().profiles().get("terra").hasExhaustedPattern());
|
||||
|
||||
FleetMcp.QuarantineSource source = Fleetd.quarantineSource(config, BackendQuarantine.none(),
|
||||
Map.of());
|
||||
assertFalse(source.exhaustedPatternArmed().apply("terra"),
|
||||
"exhaustionDetectionArmed must use the startup pattern map, not config.get()");
|
||||
// Re-read exhaustedPatternArmed off the SAME QuarantineSource built BEFORE the reload — no
|
||||
// new object, no new wiring — to prove the field itself is live, not merely that a freshly
|
||||
// built source would be. This is the exact axis the pre-#446 code fails on: the same
|
||||
// Function<String,Boolean> lambda, called again after a reload, sees the new answer only if
|
||||
// it re-reads config.get() on every call rather than a value captured earlier.
|
||||
assertTrue(before.exhaustedPatternArmed().apply("terra"),
|
||||
"exhaustionDetectionArmed must read the LIVE config on every call — fleetd #446 made "
|
||||
+ "this hot; it must reflect a reload with no restart and no new QuarantineSource");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a profile in the startup pattern map is reported as armed")
|
||||
void aProfileInTheStartupMapIsArmed(@TempDir Path dir) throws Exception {
|
||||
// fleetd #404, second direction. The test above only ever passes an EMPTY startup map, so
|
||||
// it cannot tell a correct lookup from one that is permanently off. Measured: replacing the
|
||||
// armed lambda with `profile -> false` left the whole suite green at 1475 tests. That
|
||||
// mutation would make #395's visibility feature dead — an operator fixing a detection gap
|
||||
// would be told the gap is still open after fixing it, forever. Both directions are needed:
|
||||
// this test is the only thing that fails when the field stops reporting armed at all.
|
||||
@DisplayName("reloading exhaustedPattern OUT disarms exhaustionDetectionArmed, no restart")
|
||||
void reloadedPatternRemovalDisarmsDetectionWithNoRestart(@TempDir Path dir) throws Exception {
|
||||
Path file = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(file, WITH_PATTERN);
|
||||
ConfigRef config = new ConfigRef(file, FleetConfig.load(file));
|
||||
|
||||
FleetMcp.QuarantineSource source = Fleetd.quarantineSource(config, BackendQuarantine.none(),
|
||||
Map.of("terra", Pattern.compile("usage limit")));
|
||||
new LiveExhaustedPatterns(() -> config.get().profiles()), Map.of());
|
||||
assertTrue(source.exhaustedPatternArmed().apply("terra"),
|
||||
"exhaustedPattern is configured from the start — must report armed");
|
||||
|
||||
Files.writeString(file, NO_PATTERN);
|
||||
assertTrue(config.reload().applied());
|
||||
|
||||
assertFalse(source.exhaustedPatternArmed().apply("terra"),
|
||||
"removing exhaustedPattern and reloading must disarm detection with no restart — an "
|
||||
+ "operator turning detection off (e.g. while debugging a false positive) "
|
||||
+ "must see that reflected immediately, exactly like arming it is");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a profile with no exhaustedPattern configured is reported unarmed")
|
||||
void aProfileWithNoPatternIsUnarmed(@TempDir Path dir) throws Exception {
|
||||
// fleetd #404's second direction, carried forward: this test alone fails if
|
||||
// exhaustedPatternArmed degenerates to a constant `profile -> true` — it needs at least one
|
||||
// profile that is genuinely unarmed to catch that.
|
||||
Path file = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(file, WITH_PATTERN);
|
||||
ConfigRef config = new ConfigRef(file, FleetConfig.load(file));
|
||||
|
||||
FleetMcp.QuarantineSource source = Fleetd.quarantineSource(config, BackendQuarantine.none(),
|
||||
new LiveExhaustedPatterns(() -> config.get().profiles()), Map.of());
|
||||
|
||||
assertTrue(source.exhaustedPatternArmed().apply("terra"),
|
||||
"a profile whose pattern was compiled at startup must report armed");
|
||||
"a profile whose exhaustedPattern is configured must report armed");
|
||||
assertFalse(source.exhaustedPatternArmed().apply("sonnet"),
|
||||
"a profile absent from the startup map must not report armed");
|
||||
"a profile absent from config entirely must not report armed");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,175 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.Logger;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
import dev.ltms.fleet.inject.ExhaustionSink;
|
||||
import dev.ltms.fleet.member.ClaudeCodeLauncher;
|
||||
import dev.ltms.fleet.placement.BackendQuarantine;
|
||||
import dev.ltms.fleet.session.SessionManager;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #446 follow-up (round 3): round 2's {@link FleetdUsageLimitFixWarningTest} calls {@link
|
||||
* Fleetd#usageLimitFixWarning}/{@link Fleetd#usageLimitFixWarningNoModel} directly, which pins
|
||||
* the two methods' TEXT but cannot see whether {@link Fleetd#exhaustionSink}'s {@code log.warn}
|
||||
* call actually invokes either one, or invokes the right one. A mutation battery run against
|
||||
* round 2's merge (1592 green tests) proved both gaps real:
|
||||
*
|
||||
* <ul>
|
||||
* <li><b>Cell A</b> — replacing the whole ternary result with a literal string
|
||||
* ({@code "usage-limit fix: MUTANT"}) left every test green;</li>
|
||||
* <li><b>Cell B</b> — swapping the ternary's two branches, so a profile WITH a {@code model:}
|
||||
* gets the no-model fallback message and vice versa, was measured unmeasured by the lead
|
||||
* but the existing test file shows no call that would notice it either.</li>
|
||||
* </ul>
|
||||
*
|
||||
* <p>This class drives {@link Fleetd#exhaustionSink} — the actual factory {@code main} wires,
|
||||
* since round 3's extraction pulled it out of the inline lambda for exactly this reason — and
|
||||
* asserts on the REAL text a {@link ListAppender} attached to {@code Fleetd}'s own logger
|
||||
* captures, mirroring {@code ExhaustedPatternGapReportTest}'s idiom. Chosen over a source-text
|
||||
* assertion (the {@code FleetMcpAuthzTest} {@code everyRegisteredToolHasItsHandlerActionPinned}
|
||||
* idiom) because that route can prove Cell A (the ternary is still there, still built from the
|
||||
* two named methods) but not Cell B (which branch a given profile actually reaches) — this route
|
||||
* answers both from one mechanism, since it reads what the sink actually logged for each shape of
|
||||
* profile.
|
||||
*
|
||||
* <p>Each test filters for the WARN line starting {@code "usage-limit fix:"} specifically — {@link
|
||||
* Fleetd#exhaustionSink} also logs a separate {@code "credential '...' quarantined for..."} WARN
|
||||
* on every call, and asserting against the wrong one would pass or fail for the wrong reason. That
|
||||
* prefix survives even under the M5 mutation above (the mutant keeps the tag, only the body
|
||||
* becomes {@code "MUTANT"}), so the filter itself is not what either mutation defeats.
|
||||
*/
|
||||
class FleetdExhaustionSinkWarningTest {
|
||||
|
||||
private static final String YAML = """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
herdrSocket: ~/.config/herdr/herdr.sock
|
||||
profiles:
|
||||
terra:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
model: claude-opus-5
|
||||
gx:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
guard:
|
||||
offSubscriptionHosts:
|
||||
- gx00.gw
|
||||
""";
|
||||
|
||||
private static SessionManager emptyRosterSessions() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
FleetConfig.Profile dummy = new FleetConfig.Profile(
|
||||
"dummy", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN", null,
|
||||
"tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
|
||||
ClaudeCodeLauncher launcher = new ClaudeCodeLauncher(new AgentControl(h), new WorkspaceControl(h),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(dummy.profile(), dummy), dummy.profile(), _ -> "tok");
|
||||
// Never acquires a session — exhaustionSink's roster lookup is expected to miss and fall
|
||||
// back to the profileHint argument, exactly like OpenCodeLauncher's real call site does
|
||||
// (fleetd #234) — so the sink under test never needs a populated roster.
|
||||
return new SessionManager(launcher);
|
||||
}
|
||||
|
||||
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);
|
||||
}
|
||||
|
||||
/** The one WARN line {@code exhaustionSink} builds from the ternary under test, or {@code null}. */
|
||||
private static String fixWarning(ListAppender<ILoggingEvent> appender) {
|
||||
List<String> warns = appender.list.stream()
|
||||
.filter(e -> e.getLevel() == Level.WARN)
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.filter(m -> m.startsWith("usage-limit fix:"))
|
||||
.toList();
|
||||
assertTrue(warns.size() == 1,
|
||||
"expected exactly one 'usage-limit fix:' WARN per onExhausted call, got " + warns.size()
|
||||
+ ": " + warns);
|
||||
return warns.getFirst();
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a profile WITH a configured model gets the actionable models.allow fix, not the fallback")
|
||||
void profileWithModelGetsTheActionableFix(@TempDir Path dir) throws Exception {
|
||||
Path file = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(file, YAML);
|
||||
FleetConfig cfg = FleetConfig.load(file);
|
||||
ConfigRef config = ConfigRef.fixed(cfg);
|
||||
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
|
||||
Map<String, String> reasonByCredential = new HashMap<>();
|
||||
ExhaustionSink sink = Fleetd.exhaustionSink(emptyRosterSessions(), config, quarantine,
|
||||
reasonByCredential, cfg);
|
||||
|
||||
ListAppender<ILoggingEvent> appender = attach();
|
||||
try {
|
||||
sink.onExhausted("term_x", "The usage limit has been reached", "terra");
|
||||
} finally {
|
||||
detach(appender);
|
||||
}
|
||||
|
||||
String warning = fixWarning(appender);
|
||||
assertTrue(warning.contains("terra"), "must name the profile: " + warning);
|
||||
assertTrue(warning.contains("claude-opus-5"), "must name the model: " + warning);
|
||||
assertTrue(warning.contains("enabled: false"), "must name the action: " + warning);
|
||||
assertTrue(warning.contains("models.allow"), "must name where the action goes: " + warning);
|
||||
assertFalse(warning.contains("has no model: configured"),
|
||||
"a profile WITH a model must not get the no-model fallback text: " + warning);
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a profile with NO configured model gets the fallback, never a fix that does not exist")
|
||||
void profileWithNoModelGetsTheFallbackNotAFix(@TempDir Path dir) throws Exception {
|
||||
Path file = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(file, YAML);
|
||||
FleetConfig cfg = FleetConfig.load(file);
|
||||
ConfigRef config = ConfigRef.fixed(cfg);
|
||||
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
|
||||
Map<String, String> reasonByCredential = new HashMap<>();
|
||||
ExhaustionSink sink = Fleetd.exhaustionSink(emptyRosterSessions(), config, quarantine,
|
||||
reasonByCredential, cfg);
|
||||
|
||||
ListAppender<ILoggingEvent> appender = attach();
|
||||
try {
|
||||
sink.onExhausted("term_y", "The usage limit has been reached", "gx");
|
||||
} finally {
|
||||
detach(appender);
|
||||
}
|
||||
|
||||
String warning = fixWarning(appender);
|
||||
assertTrue(warning.contains("gx"), "must name the profile: " + warning);
|
||||
assertTrue(warning.contains("has no model: configured"),
|
||||
"must say there is no model to gate by name: " + warning);
|
||||
assertFalse(warning.contains("enabled: false"),
|
||||
"a profile with no model: configured has no models.allow entry to flip — must "
|
||||
+ "not claim one exists: " + warning);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,75 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #446 follow-up (criterion 2): {@code exhaustionSink} used to build the "name the fix"
|
||||
* WARNING inline with two SLF4J {@code {}}-placeholder {@code log.warn} calls. Nothing pinned
|
||||
* that text — a mutation battery run against the merged PR #457 proved it: renaming the leading
|
||||
* {@code "usage-limit fix:"} tag to {@code "MUTANT no fix named:"} at both call sites left all
|
||||
* 1586 existing tests green. This class is what makes that mutation red.
|
||||
*
|
||||
* <p>{@link Fleetd#usageLimitFixWarning} and {@link Fleetd#usageLimitFixWarningNoModel} are the
|
||||
* extracted call sites {@code exhaustionSink} actually invokes — {@code log.warn} is called with
|
||||
* each method's return value as a single, already-formatted argument, so the string these tests
|
||||
* assert on is byte-identical to what {@code fleetd.out} receives. Same extracted-static-method +
|
||||
* dedicated-test idiom as {@link Fleetd#exhaustedPatternCoverageLine} /
|
||||
* {@link FleetdPatternCoverageLineTest}, chosen over the {@link ch.qos.logback.core.read.ListAppender}
|
||||
* idiom {@code ExhaustedPatternGapReportTest} uses — this warning is built once, at one call site,
|
||||
* from a single method with no branching inside the message itself, so there is no real log
|
||||
* plumbing left to prove; capturing an appender would only add setup/teardown around the same
|
||||
* assertion.
|
||||
*
|
||||
* <p>Each assertion below checks for a specific substring — the leading {@code "usage-limit fix:"}
|
||||
* tag an operator would grep the log for, the profile name, the model name, and the literal
|
||||
* action ({@code enabled: false} under {@code models.allow}) — rather than the whole sentence.
|
||||
* Criterion 2 promised those facts, not a particular wording of the prose that carries them; a
|
||||
* whole-sentence {@code assertEquals} breaks on every copy-edit of that surrounding prose and
|
||||
* teaches the next person to delete the test rather than read why it failed. The leading tag is
|
||||
* pinned on purpose, unlike the rest of the sentence: it is the identifying label the fleetd #446
|
||||
* follow-up's own mutation battery renamed ({@code "usage-limit fix:"} → {@code "MUTANT no fix
|
||||
* named:"}) to prove this class did not yet catch a tag rename — only the three-fact substrings
|
||||
* above would have missed it, since none of them mention the tag itself.
|
||||
*/
|
||||
class FleetdUsageLimitFixWarningTest {
|
||||
|
||||
@Test
|
||||
@DisplayName("usageLimitFixWarning names the profile, the model, and the models.allow fix")
|
||||
void usageLimitFixWarningNamesProfileModelAndFix() {
|
||||
String warning = Fleetd.usageLimitFixWarning("terra", "claude-opus-5");
|
||||
|
||||
assertTrue(warning.startsWith("usage-limit fix:"),
|
||||
"must carry the grep-able identifying tag: " + warning);
|
||||
assertTrue(warning.contains("terra"), "must name the profile: " + warning);
|
||||
assertTrue(warning.contains("claude-opus-5"), "must name the model: " + warning);
|
||||
assertTrue(warning.contains("enabled: false"), "must name the action: " + warning);
|
||||
assertTrue(warning.contains("models.allow"), "must name where the action goes: " + warning);
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("usageLimitFixWarning does not cross-name a different profile or model")
|
||||
void usageLimitFixWarningDoesNotCrossName() {
|
||||
String warning = Fleetd.usageLimitFixWarning("terra", "claude-opus-5");
|
||||
|
||||
// Guards against the mutation of swapping the two profile/model arguments at the call
|
||||
// site — a test that only checks "some name appears" cannot catch that.
|
||||
assertTrue(!warning.contains("sonnet"), "must not name an unrelated model: " + warning);
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("usageLimitFixWarningNoModel names the profile but no fix, since none exists")
|
||||
void usageLimitFixWarningNoModelNamesProfileNotAFix() {
|
||||
String warning = Fleetd.usageLimitFixWarningNoModel("gx", 1800);
|
||||
|
||||
assertTrue(warning.startsWith("usage-limit fix:"),
|
||||
"must carry the grep-able identifying tag: " + warning);
|
||||
assertTrue(warning.contains("gx"), "must name the profile: " + warning);
|
||||
assertTrue(warning.contains("1800"), "must name the quarantine duration it falls back to: " + warning);
|
||||
assertTrue(!warning.contains("enabled: false"),
|
||||
"a profile with no model: configured has no models.allow entry to flip — the "
|
||||
+ "fallback message must not claim one exists: " + warning);
|
||||
}
|
||||
}
|
||||
@@ -161,7 +161,15 @@ class ConfigRefProfileCoverageTest {
|
||||
// the #323 bug. Only ideProjectDir and worktreeGroup were saved by a behavioural test in
|
||||
// ConfigRefTest; the other two had none. So: growing this set now requires editing this
|
||||
// line as well, which is a visible, deliberate diff rather than a quiet one.
|
||||
assertEquals(Set.of("weight", "maxLoad", "credentialId"), excluded,
|
||||
//
|
||||
// fleetd #446: "exhaustedPattern" joined the set legitimately — it moved from compiled
|
||||
// once at Fleetd.main startup to read live, cached by profile name, through
|
||||
// LiveExhaustedPatterns (see that class's doc and ConfigRef's class doc). Unlike the #323
|
||||
// hazard this comment warns about, it is pinned by its OWN behavioural tests:
|
||||
// FleetdExhaustionDetectionArmedWiringTest proves the reload takes effect with no restart,
|
||||
// and the same test rerun against the pre-#446 compile site (see that ticket's report) is
|
||||
// what proves the assertion would have failed before the fix.
|
||||
assertEquals(Set.of("weight", "maxLoad", "credentialId", "exhaustedPattern"), excluded,
|
||||
"ConfigRef.LAUNCH_SETTINGS_EXCLUDED changed. A component belongs in it ONLY if it "
|
||||
+ "is read live off the config supplier, not baked into a launcher at "
|
||||
+ "startup. If you are adding one to silence this test, that is fleetd #323 "
|
||||
|
||||
@@ -339,12 +339,18 @@ class ConfigRefTest {
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-578 stage B: exhaustedPattern is compiled once into Fleetd.main's pattern map at startup
|
||||
* (see ExhaustedPatternLookup), so a reload never re-reads it — a changed pattern must be
|
||||
* reported deferred exactly like model/baseUrl, not silently claimed as applied.
|
||||
* fleetd #446: exhaustedPattern moved from deferred to hot — {@code LiveExhaustedPatterns}
|
||||
* reads {@code config.get()} fresh on every lookup, cached by profile name (not baked into a
|
||||
* frozen map inside {@code Fleetd.main} any more, the way it was before this ticket, and the
|
||||
* way its sibling {@code errorPattern} still is — see
|
||||
* {@code changingAProfilesErrorPatternIsReportedAsDeferred} right below for the contrast).
|
||||
* So a reload that changes ONLY exhaustedPattern must report a plain "config reloaded", with
|
||||
* nothing deferred — the opposite of what this test asserted before fleetd #446, when
|
||||
* ConfigRef.sameLaunchSettings still compared exhaustedPattern and this reload reported it as
|
||||
* one of the changes needing a restart.
|
||||
*/
|
||||
@Test
|
||||
void changingAProfilesExhaustedPatternIsReportedAsDeferred(@TempDir Path dir) throws Exception {
|
||||
void changingAProfilesExhaustedPatternTakesEffectWithNoRestart(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
@@ -379,17 +385,21 @@ class ConfigRefTest {
|
||||
ConfigRef.Outcome out = ref.reload();
|
||||
|
||||
assertTrue(out.applied());
|
||||
assertEquals(1, out.deferred().size(), out.deferred().toString());
|
||||
assertTrue(out.deferred().getFirst().contains("sonnet"), out.deferred().toString());
|
||||
assertTrue(out.deferred().getFirst().contains("launch settings"), out.deferred().toString());
|
||||
// The snapshot still carries the new value — a restart is what makes it take effect.
|
||||
assertEquals(0, out.deferred().size(),
|
||||
"exhaustedPattern is hot since fleetd #446 — changing only this key must not be "
|
||||
+ "reported as needing a restart: " + out.deferred());
|
||||
assertEquals(0, out.split().size(), out.split().toString());
|
||||
assertEquals("config reloaded", out.summary());
|
||||
// The snapshot carries the new value immediately — this IS the live value now, not
|
||||
// merely what a restart would eventually pick up.
|
||||
assertEquals("rate limit exceeded", ref.get().profiles().get("sonnet").exhaustedPattern());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #201 Unit 5: errorPattern is compiled once into Fleetd.main's backend-error pattern
|
||||
* map at startup (see BackendErrorPatternLookup), the same way exhaustedPattern is above — a
|
||||
* reload never re-reads it either, so a changed value must be reported deferred.
|
||||
* map at startup (see BackendErrorPatternLookup) — unlike its sibling exhaustedPattern above,
|
||||
* which fleetd #446 made hot, errorPattern was scoped OUT of that ticket on purpose and stays
|
||||
* deferred: a reload still never re-reads it, so a changed value must be reported deferred.
|
||||
*/
|
||||
@Test
|
||||
void changingAProfilesErrorPatternIsReportedAsDeferred(@TempDir Path dir) throws Exception {
|
||||
|
||||
@@ -17,15 +17,105 @@ import static org.junit.jupiter.api.Assumptions.assumeTrue;
|
||||
* SHELL directly (never {@code claude}, so no subscription/token involvement) and always tears
|
||||
* the throwaway space down.
|
||||
*
|
||||
* <p>The seed shell's own startup (restoring its session, printing its banner) is asynchronous
|
||||
* and its length is not a fleetd contract — measured here at ~2.5s on one host (fleetd #449).
|
||||
* The old version used two fixed sleeps: 1000ms before typing, then 800ms before reading. It
|
||||
* failed, and the pane showed the typed line followed by the startup banner with no command
|
||||
* output at all — which looks, at a glance, exactly like the env map never reaching the shell.
|
||||
* So this polls for a real signal instead of guessing a sleep length.
|
||||
*
|
||||
* <p><strong>What was measured, and what was not.</strong> Polling fixes it: 5 standalone runs
|
||||
* green. The load-bearing half is {@link #waitForText}, and one cell proves it. Keep the old
|
||||
* 1000ms write sleep and change only the read — the 800ms fixed sleep becomes a 5s poll — and
|
||||
* the test goes from 0 of 3 passing to 3 of 3. Removing the write wait instead
|
||||
* ({@link #SHELL_READY_TIMEOUT_MS} set to 0, so input is typed at once) also passes 3 of 3. So
|
||||
* the cause is the 800ms READ deadline, not the 1000ms write delay. The old version fails every
|
||||
* time, not sometimes, so "race" is the wrong word for it. The earlier explanation — that input
|
||||
* typed before the prompt is swallowed by the shell's startup — is not supported by anything
|
||||
* measured here. Please do not repeat it: if it were true, typing at 0ms would be worse than
|
||||
* typing at 1000ms, and it is not.
|
||||
*
|
||||
* <p><strong>Where this was measured.</strong> A 12-core macOS host, load average 2.6 to 5.9,
|
||||
* on commit 20c1094. The same four cells were also run under load, with 8 spinners on 12 cores.
|
||||
* The failure and the passes from that run are not worth the same. The cells ran in a fixed
|
||||
* order while the load climbed from 7 to 50. The old version ran last, at the top of that climb,
|
||||
* so it has a free explanation for failing and its 0 of 3 is discarded. A pass has no such free
|
||||
* explanation: a cell that survives a worse condition than a fair order would have given it is
|
||||
* evidence in the safe direction. So keep the three passes, each with the load it ran at: the
|
||||
* fixed version 3 of 3 at load 7.42 to 18.42, the 0ms-write cell 3 of 3 at 18.42 to 23.65, the
|
||||
* 1000ms-write cell 3 of 3 at 23.65 to 46.40. Above about load 20 everything here is slow for
|
||||
* reasons that have nothing to do with this seam, so read the positive claim — that the read
|
||||
* deadline was the whole cause — as "measured near idle on a 12-core host", and nothing
|
||||
* stronger. Do not carry that raw load average to another host either: load average counts
|
||||
* differently per core and per operating system, so only load per core compares. If this test
|
||||
* fails on a smaller or busier machine, raise {@link #OUTPUT_TIMEOUT_MS} before you suspect the
|
||||
* seam.
|
||||
*
|
||||
* <p>{@link #waitUntilSettled} is kept as cheap insurance against that swallow case, not because
|
||||
* anyone showed it was needed. The cell that tests the swallow case head-on is the one with
|
||||
* {@link #SHELL_READY_TIMEOUT_MS} at 0: input is typed at once, which is the worst case for
|
||||
* "typed before the prompt is ready". It passed 3 of 3 at load 18.42 to 23.65. The wider read
|
||||
* window cannot explain that pass away, because a swallowed keystroke is LOST, not late — the
|
||||
* command never runs, so no amount of polling makes its output appear. So the swallow mechanism
|
||||
* was tested and did not show up. If you want to delete this call, that is the cell to re-run.
|
||||
*
|
||||
* <p>Tagged {@code contract}; run with {@code mvn test -Pcontract}.
|
||||
*/
|
||||
@Tag("contract")
|
||||
class AgentControlContractTest {
|
||||
|
||||
private static final long POLL_INTERVAL_MS = 150;
|
||||
/** Bound for the seed shell to settle: observed ~2.5s three times running; this leaves headroom. */
|
||||
private static final long SHELL_READY_TIMEOUT_MS = 8_000;
|
||||
/** Bound for the typed command's output to appear once the shell is ready: observed ~0.2s. */
|
||||
private static final long OUTPUT_TIMEOUT_MS = 5_000;
|
||||
|
||||
private boolean noSocket() {
|
||||
return !Files.exists(UnixSocketHerdrClient.defaultSocketPath());
|
||||
}
|
||||
|
||||
private static String readPane(UnixSocketHerdrClient herdr, String paneId) {
|
||||
return herdr.call("pane.read", Map.of("pane_id", paneId, "source", "visible"))
|
||||
.path("read").path("text").asText("");
|
||||
}
|
||||
|
||||
/**
|
||||
* Poll {@code pane.read} until two consecutive reads come back identical — the shell's own
|
||||
* startup output (restore banner, prompt) has stopped changing — or {@code timeoutMs} elapses.
|
||||
* Never asserts by itself; the caller's own assertion is what actually verifies the outcome.
|
||||
* Setting {@code timeoutMs} to 0 skips the wait entirely and the test still passes here, so
|
||||
* treat this as insurance rather than as the fix — see the class javadoc.
|
||||
*/
|
||||
private static String waitUntilSettled(UnixSocketHerdrClient herdr, String paneId, long timeoutMs)
|
||||
throws InterruptedException {
|
||||
long deadline = System.currentTimeMillis() + timeoutMs;
|
||||
String previous = null;
|
||||
while (System.currentTimeMillis() < deadline) {
|
||||
Thread.sleep(POLL_INTERVAL_MS);
|
||||
String current = readPane(herdr, paneId);
|
||||
if (current.equals(previous) && !current.isBlank()) {
|
||||
return current;
|
||||
}
|
||||
previous = current;
|
||||
}
|
||||
return previous == null ? "" : previous;
|
||||
}
|
||||
|
||||
/** Poll {@code pane.read} until {@code needle} appears or {@code timeoutMs} elapses. */
|
||||
private static String waitForText(UnixSocketHerdrClient herdr, String paneId, String needle, long timeoutMs)
|
||||
throws InterruptedException {
|
||||
long deadline = System.currentTimeMillis() + timeoutMs;
|
||||
String last = "";
|
||||
while (System.currentTimeMillis() < deadline) {
|
||||
last = readPane(herdr, paneId);
|
||||
if (last.contains(needle)) {
|
||||
return last;
|
||||
}
|
||||
Thread.sleep(POLL_INTERVAL_MS);
|
||||
}
|
||||
return last;
|
||||
}
|
||||
|
||||
@Test
|
||||
void tabCreateInjectsEnvIntoTheSeedShell() throws Exception {
|
||||
assumeTrue(!noSocket(), "no herdr socket — skipping");
|
||||
@@ -36,15 +126,13 @@ class AgentControlContractTest {
|
||||
Map.of("ANTHROPIC_BASE_URL", "http://gx00.gw:8000"));
|
||||
try {
|
||||
assertNotNull(tab.rootPaneId(), "tab.create must return the seed pane");
|
||||
Thread.sleep(1000); // let the seed shell reach its prompt
|
||||
waitUntilSettled(herdr, tab.rootPaneId(), SHELL_READY_TIMEOUT_MS);
|
||||
herdr.call("pane.send_input", Map.of(
|
||||
"pane_id", tab.rootPaneId(),
|
||||
"text", "printf 'PROBE_BASE=[%s]\\n' \"$ANTHROPIC_BASE_URL\"",
|
||||
"keys", List.of("enter")));
|
||||
Thread.sleep(800);
|
||||
String visible = herdr.call("pane.read",
|
||||
Map.of("pane_id", tab.rootPaneId(), "source", "visible"))
|
||||
.path("read").path("text").asText("");
|
||||
String visible = waitForText(herdr, tab.rootPaneId(),
|
||||
"PROBE_BASE=[http://gx00.gw:8000]", OUTPUT_TIMEOUT_MS);
|
||||
assertTrue(visible.contains("PROBE_BASE=[http://gx00.gw:8000]"),
|
||||
"env map must reach the seed shell; saw: " + visible);
|
||||
} finally {
|
||||
|
||||
@@ -14,7 +14,7 @@ import static org.junit.jupiter.api.Assumptions.assumeTrue;
|
||||
* Contract test against a REAL running herdr. Tagged {@code contract} so it is
|
||||
* excluded from {@code mvn test}; run it with {@code mvn test -Pcontract}. It fails
|
||||
* loudly if herdr drifts from the protocol {@code fleetd} was built against
|
||||
* (0.7.0, protocol 14) — catching breakage that unit tests with canned frames cannot.
|
||||
* (0.8.0, protocol 19) — catching breakage that unit tests with canned frames cannot.
|
||||
*/
|
||||
@Tag("contract")
|
||||
class HerdrContractTest {
|
||||
@@ -24,13 +24,13 @@ class HerdrContractTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
void pingReturnsProtocol14() {
|
||||
void pingReturnsProtocol19() {
|
||||
assumeTrue(Files.exists(socket()), "no herdr socket at " + socket() + " — skipping");
|
||||
try (UnixSocketHerdrClient herdr = UnixSocketHerdrClient.connect()) {
|
||||
JsonNode pong = herdr.call("ping");
|
||||
assertEquals("pong", pong.get("type").asText());
|
||||
assertEquals(14, pong.get("protocol").asInt(),
|
||||
"fleetd is built against herdr protocol 14");
|
||||
assertEquals(19, pong.get("protocol").asInt(),
|
||||
"fleetd is built against herdr protocol 19");
|
||||
assertFalse(pong.get("version").asText().isBlank());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,138 @@
|
||||
package dev.ltms.fleet.inject;
|
||||
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.Map;
|
||||
import java.util.regex.Pattern;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotSame;
|
||||
import static org.junit.jupiter.api.Assertions.assertNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertSame;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #446: {@code LiveExhaustedPatterns} is the mechanism that makes {@code exhaustedPattern}
|
||||
* hot — read fresh off a config supplier per lookup, cached by profile name. This is the direct,
|
||||
* minimal unit test of that mechanism: read the pattern once, mutate the backing "config", read
|
||||
* again, and assert the second read saw the new value — the exact hotness proof shape.
|
||||
*
|
||||
* <p>This class alone does not prove {@code Fleetd.main} actually wires the live source into
|
||||
* production — {@link dev.ltms.fleet.FleetdExhaustionDetectionArmedWiringTest} and {@code
|
||||
* CompletionResolverTest}'s {@code ExhaustedPatternLookup} tests do that at the wiring level. What
|
||||
* this class proves is that the mechanism itself is genuinely live and genuinely cached, not that
|
||||
* it is used.
|
||||
*/
|
||||
class LiveExhaustedPatternsTest {
|
||||
|
||||
private static FleetConfig.Profile profileWithPattern(String pattern) {
|
||||
return new FleetConfig.Profile("terra", "http://gx00.gw:8000", "terra",
|
||||
null, null, null, "tab", "fleetd-workers", "w #{n}", null, null, null,
|
||||
null, null, null, null, null, null, null, pattern);
|
||||
}
|
||||
|
||||
/** Simple mutable holder standing in for {@code ConfigRef} — swapped, never mutated in place. */
|
||||
private static final class MutableProfiles {
|
||||
private volatile Map<String, FleetConfig.Profile> profiles;
|
||||
|
||||
MutableProfiles(Map<String, FleetConfig.Profile> initial) {
|
||||
this.profiles = initial;
|
||||
}
|
||||
|
||||
Map<String, FleetConfig.Profile> get() {
|
||||
return profiles;
|
||||
}
|
||||
|
||||
void set(Map<String, FleetConfig.Profile> fresh) {
|
||||
this.profiles = fresh;
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("patternFor reads the live config: read, change, read again, see the change — no restart")
|
||||
void patternForIsHotAcrossAConfigChange() {
|
||||
MutableProfiles live = new MutableProfiles(Map.of("terra", profileWithPattern("usage limit")));
|
||||
LiveExhaustedPatterns patterns = new LiveExhaustedPatterns(live::get);
|
||||
|
||||
Pattern before = patterns.patternFor("terra");
|
||||
assertTrue(before.matcher("your usage limit has been reached").find());
|
||||
|
||||
// Change the "config" — the exact thing a reload does to ConfigRef in production.
|
||||
live.set(Map.of("terra", profileWithPattern("rate limit exceeded")));
|
||||
|
||||
Pattern after = patterns.patternFor("terra");
|
||||
assertTrue(after.matcher("429: rate limit exceeded, try later").find(),
|
||||
"the second read must see the NEW pattern text with no restart");
|
||||
assertFalse(after.matcher("your usage limit has been reached").find(),
|
||||
"the second read must have actually recompiled — not just reused a stale match");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("armed follows the same live read: on, then off, after a config change — no restart")
|
||||
void armedIsHotAcrossAConfigChange() {
|
||||
MutableProfiles live = new MutableProfiles(Map.of("terra", profileWithPattern("usage limit")));
|
||||
LiveExhaustedPatterns patterns = new LiveExhaustedPatterns(live::get);
|
||||
|
||||
assertTrue(patterns.armed("terra"));
|
||||
|
||||
live.set(Map.of("terra", profileWithPattern(null)));
|
||||
|
||||
assertFalse(patterns.armed("terra"),
|
||||
"removing exhaustedPattern and 'reloading' must disarm detection with no restart");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("an unchanged pattern string is not recompiled — the cache is reused")
|
||||
void unchangedPatternTextReusesTheCompiledInstance() {
|
||||
MutableProfiles live = new MutableProfiles(Map.of("terra", profileWithPattern("usage limit")));
|
||||
LiveExhaustedPatterns patterns = new LiveExhaustedPatterns(live::get);
|
||||
|
||||
Pattern first = patterns.patternFor("terra");
|
||||
// A NEW Profile object, same pattern text — a reload of an unrelated key produces a fresh
|
||||
// FleetConfig even though this profile's own text did not change.
|
||||
live.set(Map.of("terra", profileWithPattern("usage limit")));
|
||||
Pattern second = patterns.patternFor("terra");
|
||||
|
||||
assertSame(first, second,
|
||||
"identical pattern text must reuse the cached compiled Pattern, not recompile it — "
|
||||
+ "recompiling a regex on every check is exactly the cost the cache exists to avoid");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a changed pattern string is recompiled, not silently reused from the cache")
|
||||
void changedPatternTextIsRecompiled() {
|
||||
MutableProfiles live = new MutableProfiles(Map.of("terra", profileWithPattern("usage limit")));
|
||||
LiveExhaustedPatterns patterns = new LiveExhaustedPatterns(live::get);
|
||||
|
||||
Pattern first = patterns.patternFor("terra");
|
||||
live.set(Map.of("terra", profileWithPattern("rate limit")));
|
||||
Pattern second = patterns.patternFor("terra");
|
||||
|
||||
assertNotSame(first, second, "changed pattern text must produce a freshly compiled Pattern");
|
||||
}
|
||||
|
||||
@Test
|
||||
void unknownProfileReturnsNullAndUnarmed() {
|
||||
LiveExhaustedPatterns patterns = new LiveExhaustedPatterns(Map::of);
|
||||
assertNull(patterns.patternFor("nope"));
|
||||
assertFalse(patterns.armed("nope"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void nullProfileNameReturnsNullAndUnarmed() {
|
||||
LiveExhaustedPatterns patterns = new LiveExhaustedPatterns(
|
||||
() -> Map.of("terra", profileWithPattern("usage limit")));
|
||||
assertNull(patterns.patternFor(null));
|
||||
assertFalse(patterns.armed(null));
|
||||
}
|
||||
|
||||
@Test
|
||||
void profileWithNoPatternConfiguredReturnsNull() {
|
||||
LiveExhaustedPatterns patterns = new LiveExhaustedPatterns(
|
||||
() -> Map.of("terra", profileWithPattern(null)));
|
||||
assertNull(patterns.patternFor("terra"));
|
||||
assertFalse(patterns.armed("terra"));
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,72 @@
|
||||
package dev.ltms.fleet.mcp;
|
||||
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.LinkedHashSet;
|
||||
import java.util.Set;
|
||||
import java.util.regex.Matcher;
|
||||
import java.util.regex.Pattern;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/** fleetd #464: launch charters must not name MCP tools the server does not register. */
|
||||
class CharterToolSurfaceTest {
|
||||
|
||||
private static final Path MCP_SOURCE = Path.of("src/main/java/dev/ltms/fleet/mcp/FleetMcp.java");
|
||||
|
||||
private static Set<String> matches(String text, String regex) {
|
||||
Matcher m = Pattern.compile(regex).matcher(text);
|
||||
Set<String> found = new LinkedHashSet<>();
|
||||
while (m.find()) {
|
||||
found.add(m.group(1));
|
||||
}
|
||||
return found;
|
||||
}
|
||||
|
||||
/** Every {@code fleet_*} or legacy {@code bridge_*} token in the configured launch charters. */
|
||||
private static Set<String> toolsNamedIn(FleetConfig config) {
|
||||
return matches(String.join("\n", config.fleet().charters().values()),
|
||||
"(fleet_[a-z_]+|bridge_[a-z_]+)");
|
||||
}
|
||||
|
||||
/** Every tool {@link FleetMcp} registers, read from its {@code tool("…")} calls. */
|
||||
private static Set<String> toolsTheServerRegisters() throws Exception {
|
||||
return matches(Files.readString(MCP_SOURCE), "tool\\(\\\"(fleet_[a-z_]+)\\\"");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SOURCE TEXT] every tool named in a configured charter is registered by the server")
|
||||
void configuredChartersNameOnlyRegisteredTools(@TempDir Path dir) throws Exception {
|
||||
Path configFile = dir.resolve("charters.yaml");
|
||||
Files.writeString(configFile, """
|
||||
fleet:
|
||||
charters:
|
||||
dev: |
|
||||
Send the final handoff through fleet_reply.
|
||||
reviewer: |
|
||||
Use fleet_ask only for the lead's decision.
|
||||
""");
|
||||
|
||||
FleetConfig config = FleetConfig.load(configFile);
|
||||
Set<String> named = toolsNamedIn(config);
|
||||
Set<String> registered = toolsTheServerRegisters();
|
||||
|
||||
assertTrue(!named.isEmpty(),
|
||||
"the charter fixture named no fleet_* or bridge_* tool. This test would check nothing; "
|
||||
+ "add charter text that names a tool before changing the extraction.");
|
||||
assertTrue(!registered.isEmpty(),
|
||||
"the FleetMcp registration scrape found no tools. This test would check nothing; "
|
||||
+ "repair the tool(\"…\") extraction before changing the assertion.");
|
||||
|
||||
Set<String> unknown = new LinkedHashSet<>(named);
|
||||
unknown.removeAll(registered);
|
||||
assertTrue(unknown.isEmpty(),
|
||||
"configured charter text names " + unknown + ", but FleetMcp does not register it. "
|
||||
+ "Checked " + named + " against " + registered + ". Fix the charter text or "
|
||||
+ "register the tool; do NOT weaken this test.");
|
||||
}
|
||||
}
|
||||
@@ -187,6 +187,71 @@ class FleetMcpAuthzTest {
|
||||
"no CallerResolver supplied ⇒ authorization not enforced (legacy behaviour)");
|
||||
}
|
||||
|
||||
// --- fleetd #439: who may see fleet_list's coordinator row ----------------------------------
|
||||
|
||||
/**
|
||||
* fleetd #439: {@link FleetMcp#coordinatorVisibleTo} is the whole policy decision for
|
||||
* {@code fleet_list}'s {@code coordinator} row — lead-to-lead coordination state, not roster
|
||||
* observation. Only the primary may see it; a worker, an architect, and (the case the previous
|
||||
* pass of this ticket did not cover) an anonymous caller must all be refused.
|
||||
*/
|
||||
@Test
|
||||
void onlyThePrimaryMaySeeTheCoordinatorRow() {
|
||||
assertTrue(FleetMcp.coordinatorVisibleTo(PRIMARY), "the primary must see its own coordination state");
|
||||
assertFalse(FleetMcp.coordinatorVisibleTo(WORKER_A), "a worker must not see lead-to-lead coordination state");
|
||||
assertFalse(FleetMcp.coordinatorVisibleTo(ARCH_DESIGN),
|
||||
"an architect holds READ today, but that must not extend to coordinator");
|
||||
assertFalse(FleetMcp.coordinatorVisibleTo(ANON), "authenticated as nothing must not see it either");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #439 / PR #462 review finding M2: the predicate above can be perfectly correct while
|
||||
* the one production call site (the {@code fleet_list} MCP handler) never actually asks it —
|
||||
* a literal {@code true} compiles, and the whole suite stayed green under that mutation because
|
||||
* every existing test drives {@link FleetMcp#listFleet} directly and supplies the boolean
|
||||
* itself. This test reads {@code FleetMcp.java}'s own source (same idiom as {@link
|
||||
* #toolsTheServerRegisters()} / {@link #everyRegisteredToolHasItsHandlerActionPinned()}) and
|
||||
* asserts the handler's call passes {@code coordinatorVisibleTo(principal(exchange))} — not a
|
||||
* literal {@code true} or {@code false} — as {@code listFleet}'s trailing argument.
|
||||
*
|
||||
* <p>Anchored on argument position, not a bare substring search: {@code true} appears many
|
||||
* times elsewhere in this file for unrelated reasons, so a plain {@code contains("true")}
|
||||
* check would prove nothing. The pattern requires the literal text immediately before the
|
||||
* closing {@code );} of the {@code listFleet(} call inside the handler block to be exactly
|
||||
* {@code coordinatorVisibleTo(principal(exchange))}.
|
||||
*/
|
||||
@Test
|
||||
void theFleetListHandlerActuallyConsultsCoordinatorVisibleTo() throws Exception {
|
||||
String source = Files.readString(MCP_SOURCE);
|
||||
|
||||
// Isolate the fleet_list handler block: from its declaration up to the next handler's
|
||||
// declaration. A change to variable naming would break this scrape loudly (see the control
|
||||
// assertion just below), rather than silently reporting "no violation found".
|
||||
int start = source.indexOf("listHandler =");
|
||||
assertTrue(start >= 0, "could not find the fleet_list handler (listHandler) in " + MCP_SOURCE
|
||||
+ " -- the scrape has stopped matching, fix the anchor before trusting this test");
|
||||
int end = source.indexOf("stopHandler =", start);
|
||||
assertTrue(end > start, "could not find the handler declared after listHandler to bound the scrape");
|
||||
String handlerBlock = source.substring(start, end);
|
||||
|
||||
// CONTROL: the block we scraped really does contain a call to listFleet(...) -- if this
|
||||
// fails, the anchors above moved and the assertion below would otherwise pass on nothing.
|
||||
assertTrue(handlerBlock.contains("listFleet("),
|
||||
"control failed: the scraped listHandler block contains no listFleet( call at all -- "
|
||||
+ "the anchors have drifted, this test is not testing what it claims to");
|
||||
|
||||
Pattern trailingArg = Pattern.compile(
|
||||
"listFleet\\([^;]*?,\\s*(coordinatorVisibleTo\\(principal\\(exchange\\)\\)|true|false)\\s*\\)\\s*;",
|
||||
Pattern.DOTALL);
|
||||
Matcher m = trailingArg.matcher(handlerBlock);
|
||||
assertTrue(m.find(), "could not locate listFleet(...)'s trailing boolean argument in the "
|
||||
+ "listHandler block -- the call shape changed, update this test's anchor: " + handlerBlock);
|
||||
String trailing = m.group(1);
|
||||
assertEquals("coordinatorVisibleTo(principal(exchange))", trailing,
|
||||
"the fleet_list handler must ask coordinatorVisibleTo(principal(exchange)) who is "
|
||||
+ "calling, not pass a literal boolean -- found: " + trailing);
|
||||
}
|
||||
|
||||
// --- which action each tool hands the gate (fleetd #272) ------------------------------------
|
||||
|
||||
/**
|
||||
|
||||
@@ -352,7 +352,7 @@ class FleetMcpTest {
|
||||
fail("a lead fleet_reply must not publish to the worker inbox");
|
||||
}
|
||||
@Override public List<InboxMessage> peek(String target) { return List.of(); }
|
||||
@Override public void ack(String target, String msgId) { }
|
||||
@Override public boolean ack(String target, String msgId) { return false; }
|
||||
};
|
||||
MessageService leadMessages = new MessageService(agents, new Injector(agents), new Rendezvous(),
|
||||
inboxThatRejectsPublishes);
|
||||
@@ -626,8 +626,8 @@ class FleetMcpTest {
|
||||
McpSchema.CallToolResult res = FleetMcp.listFleet(
|
||||
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
|
||||
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), Map.of(), "",
|
||||
new FleetMcp.CoordinationSource(channel, List.of()));
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
|
||||
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true);
|
||||
|
||||
String out = textOf(res);
|
||||
// fleetd #361: reports both which coord-id a peer must use to reach ME, and this daemon's
|
||||
@@ -655,8 +655,8 @@ class FleetMcpTest {
|
||||
McpSchema.CallToolResult res = FleetMcp.listFleet(
|
||||
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
|
||||
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), Map.of(), "",
|
||||
new FleetMcp.CoordinationSource(channel, List.of()));
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
|
||||
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true);
|
||||
|
||||
String out = textOf(res);
|
||||
assertTrue(out.contains("\"mailbox\":{\"status\":\"unknown\"}"), out);
|
||||
@@ -676,6 +676,118 @@ class FleetMcpTest {
|
||||
"an ordinary fleet's output must be unchanged by this feature");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #463: a compat overload called with no {@code callerIsPrimary} argument at all must
|
||||
* fail closed, not open. Before this fix the hidden default was {@code true}, so a caller that
|
||||
* forgot the argument silently got lead-to-lead coordination state. Lead coordination is fully
|
||||
* configured here (a real channel, a real mailbox) specifically so this is not conflated with
|
||||
* {@link #listOmitsTheCoordinatorRowWhenLeadCoordinationIsOff} -- the row is capable of being
|
||||
* assembled, and the missing argument is the only reason it is not.
|
||||
*/
|
||||
@Test
|
||||
void listCompatOverloadWithNoCallerIsPrimaryArgumentOmitsTheCoordinatorKey() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
|
||||
FakeLeadChannel channel = new FakeLeadChannel("mac-opus")
|
||||
.withMailbox("mac-opus", LeadChannel.MailboxState.exists("mac-opus", 0, 1));
|
||||
|
||||
McpSchema.CallToolResult res = FleetMcp.listFleet(
|
||||
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
|
||||
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), Map.of(), "",
|
||||
new FleetMcp.CoordinationSource(channel, List.of()));
|
||||
|
||||
String out = textOf(res);
|
||||
assertFalse(out.contains("\"coordinator\""),
|
||||
"no callerIsPrimary argument must fail closed (absent), not open (present): " + out);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #439: a worker calling {@code fleet_list} must get a result with the {@code
|
||||
* coordinator} key <strong>absent</strong> -- not an empty object, not a redacted one -- even
|
||||
* though lead coordination is fully configured and would otherwise report a row. This drives
|
||||
* the same {@code callerIsPrimary} value the MCP handler computes ({@code
|
||||
* Principal.worker(...).isPrimary()}), so it pins the real production boolean, not a literal.
|
||||
*/
|
||||
@Test
|
||||
void listOmitsTheCoordinatorKeyEntirelyForAWorkerEvenWhenLeadCoordinationIsOn() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
|
||||
FakeLeadChannel channel = new FakeLeadChannel("mac-opus")
|
||||
.withMailbox("mac-opus", LeadChannel.MailboxState.exists("mac-opus", 0, 1));
|
||||
boolean callerIsPrimary = Principal.worker("term_a", 1).isPrimary();
|
||||
|
||||
McpSchema.CallToolResult res = FleetMcp.listFleet(
|
||||
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
|
||||
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
|
||||
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), callerIsPrimary);
|
||||
|
||||
String out = textOf(res);
|
||||
assertFalse(out.contains("\"coordinator\""), "a worker must never see the coordinator key at all: " + out);
|
||||
assertFalse(out.contains("mac-opus"), "no fragment of the coordinator row may leak either: " + out);
|
||||
assertTrue(out.contains("\"leads\""), "the rest of the result must still be present: " + out);
|
||||
assertTrue(out.contains("\"members\""), out);
|
||||
assertTrue(out.contains("\"healthCoverage\""), out);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #439 acceptance criterion 2: an architect gets exactly the same treatment as a worker.
|
||||
* This is a real, executed test (not just reasoning by analogy) -- it drives the actual
|
||||
* {@code Principal.architect(...).isPrimary()} value the production handler would compute for
|
||||
* an architect caller, through the same gate a worker's call goes through.
|
||||
*/
|
||||
@Test
|
||||
void listOmitsTheCoordinatorKeyEntirelyForAnArchitectToo() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
|
||||
FakeLeadChannel channel = new FakeLeadChannel("mac-opus")
|
||||
.withMailbox("mac-opus", LeadChannel.MailboxState.exists("mac-opus", 0, 1));
|
||||
boolean callerIsPrimary = Principal.architect("lead-designer", "term_design", 400).isPrimary();
|
||||
|
||||
McpSchema.CallToolResult res = FleetMcp.listFleet(
|
||||
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
|
||||
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
|
||||
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), callerIsPrimary);
|
||||
|
||||
String out = textOf(res);
|
||||
assertFalse(out.contains("\"coordinator\""), "an architect must never see the coordinator key either: " + out);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #439 acceptance criterion 3: an explicitly-primary caller sees the coordinator row
|
||||
* fully assembled, with the same content #439 always produced for a primary.
|
||||
*
|
||||
* <p>fleetd #463 flipped the compat overloads' hidden default from {@code true} to
|
||||
* {@code false} (fail closed), so the old "pre-#439 overload" this test used to compare
|
||||
* against no longer stands in for a primary caller -- it is now exactly the implicit-default
|
||||
* path #463 closes. Verifying the primary path means calling the canonical overload with an
|
||||
* explicit {@code callerIsPrimary=true} directly, as the production {@code fleet_list} handler
|
||||
* does.
|
||||
*/
|
||||
@Test
|
||||
void listIsByteForByteUnchangedForThePrimaryCaller() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
|
||||
FakeLeadChannel channel = new FakeLeadChannel("mac-opus")
|
||||
.withMailbox("mac-opus", LeadChannel.MailboxState.exists("mac-opus", 0, 1));
|
||||
|
||||
String gatedAsPrimary = textOf(FleetMcp.listFleet(
|
||||
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
|
||||
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
|
||||
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true));
|
||||
|
||||
assertTrue(gatedAsPrimary.contains("\"coordinator\""), gatedAsPrimary);
|
||||
assertTrue(gatedAsPrimary.contains("\"selfId\":\"mac-opus\""), gatedAsPrimary);
|
||||
assertTrue(gatedAsPrimary.contains("\"mailbox\""), gatedAsPrimary);
|
||||
assertTrue(gatedAsPrimary.contains("\"heldCount\""), gatedAsPrimary);
|
||||
assertTrue(gatedAsPrimary.contains("\"heldDurable\""), gatedAsPrimary);
|
||||
assertTrue(gatedAsPrimary.contains("\"held\""), gatedAsPrimary);
|
||||
assertTrue(gatedAsPrimary.contains("\"peers\""), gatedAsPrimary);
|
||||
}
|
||||
|
||||
@Test
|
||||
void listReportsHeldMessagesWithATruncatedPreviewNeverTheFullBody() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
@@ -691,8 +803,8 @@ class FleetMcpTest {
|
||||
McpSchema.CallToolResult res = FleetMcp.listFleet(
|
||||
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
|
||||
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), Map.of(), "",
|
||||
new FleetMcp.CoordinationSource(channel, List.of()));
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
|
||||
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true);
|
||||
|
||||
String out = textOf(res);
|
||||
assertTrue(out.contains("\"msgId\":\"m1\""), out);
|
||||
@@ -723,8 +835,8 @@ class FleetMcpTest {
|
||||
McpSchema.CallToolResult res = FleetMcp.listFleet(
|
||||
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
|
||||
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), Map.of(), "",
|
||||
new FleetMcp.CoordinationSource(channel, List.of()));
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
|
||||
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true);
|
||||
|
||||
String out = textOf(res);
|
||||
assertTrue(out.contains("\"pending\":0"), out);
|
||||
@@ -752,8 +864,8 @@ class FleetMcpTest {
|
||||
McpSchema.CallToolResult res = FleetMcp.listFleet(
|
||||
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
|
||||
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), Map.of(), "",
|
||||
new FleetMcp.CoordinationSource(channel, List.of()));
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
|
||||
Map.of(), "", new FleetMcp.CoordinationSource(channel, List.of()), true);
|
||||
|
||||
String out = textOf(res);
|
||||
assertTrue(out.contains("\"heldDurable\":false"),
|
||||
@@ -820,8 +932,9 @@ class FleetMcpTest {
|
||||
McpSchema.CallToolResult res = FleetMcp.listFleet(
|
||||
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
|
||||
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), Map.of(), "",
|
||||
new FleetMcp.CoordinationSource(channel, List.of("fleet01-lead", "fleet02-lead", "fleet03-lead")));
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
|
||||
Map.of(), "",
|
||||
new FleetMcp.CoordinationSource(channel, List.of("fleet01-lead", "fleet02-lead", "fleet03-lead")), true);
|
||||
|
||||
String out = textOf(res);
|
||||
assertTrue(out.contains("\"coordId\":\"fleet01-lead\",\"status\":\"exists\",\"pending\":2,\"consumers\":1"), out);
|
||||
@@ -1386,6 +1499,10 @@ class FleetMcpTest {
|
||||
|
||||
@Test
|
||||
void bridgeAckReturnsConfirmationForValidArgs() {
|
||||
// fleet_ack only reports success for a msgId actually queued in the target's inbox
|
||||
// (fleetd #437) — publish one via the inbox directly rather than asserting on a
|
||||
// fabricated id nothing ever queued.
|
||||
inbox.publish("term_a", "msg-1", "queued reply");
|
||||
McpSchema.CallToolResult res = FleetMcp.ack(messages, "term_a", "msg-1");
|
||||
assertNotEquals(Boolean.TRUE, res.isError());
|
||||
assertTrue(textOf(res).contains("msg-1"), "response should mention the msgId");
|
||||
@@ -1398,6 +1515,26 @@ class FleetMcpTest {
|
||||
assertTrue(FleetMcp.ack(messages, " ", "msg-1").isError());
|
||||
}
|
||||
|
||||
@Test
|
||||
void bridgeAckOfAnIdInNoInboxIsAnError() {
|
||||
// fleetd #437: fleet_ack used to say "acknowledged <msgId>" for a message it never
|
||||
// touched, because nothing in the chain reported hit vs. miss. "never-queued" is in no
|
||||
// inbox at all, so this must error rather than claim success.
|
||||
McpSchema.CallToolResult res = FleetMcp.ack(messages, "term_a", "never-queued");
|
||||
assertTrue(res.isError());
|
||||
assertTrue(textOf(res).contains("never-queued"), textOf(res));
|
||||
}
|
||||
|
||||
@Test
|
||||
void bridgeAckOfACoordIdTargetIsAnErrorNamingFleetPoll() {
|
||||
// A coord-id names a peer lead's held mailbox (LeadChannel/LeadMailbox), never a
|
||||
// worker's ReplyInbox — fleet_ack has no route to it and must say so, pointing at
|
||||
// fleet_poll{coordId} instead of reporting a false "acknowledged".
|
||||
McpSchema.CallToolResult res = FleetMcp.ack(messages, "coord-some-peer", "msg-1");
|
||||
assertTrue(res.isError());
|
||||
assertTrue(textOf(res).contains("fleet_poll{coordId}"), textOf(res));
|
||||
}
|
||||
|
||||
@Test
|
||||
void bridgeAckRemovesSpecificReply() {
|
||||
// Queue a reply and capture its msgId.
|
||||
@@ -1411,8 +1548,18 @@ class FleetMcpTest {
|
||||
var peeked = messages.drainReplies("term_a");
|
||||
assertEquals(1, peeked.size(), "one fresh reply in the inbox");
|
||||
|
||||
// ackReply works (no-op since published with a different UUID, but callable).
|
||||
assertDoesNotThrow(() -> messages.ackReply("term_a", msgId));
|
||||
// fleetd #437: msgId was already drained above (a fresh UUID each publish), so it is no
|
||||
// longer in the inbox — ackReply must now report that miss instead of pretending to ack.
|
||||
assertFalse(messages.ackReply("term_a", msgId));
|
||||
}
|
||||
|
||||
@Test
|
||||
void bridgeAckRemovingARealQueuedReplyReportsSuccessAndRemovesIt() {
|
||||
// The worker path must not change behaviour: acking a reply that IS still in the inbox
|
||||
// still succeeds and still removes it (fleetd #437).
|
||||
inbox.publish("term_a", "real-1", "still queued");
|
||||
assertTrue(messages.ackReply("term_a", "real-1"), "ack of a real queued reply must report true");
|
||||
assertTrue(inbox.peek("term_a").isEmpty(), "the acked reply must be gone from the inbox");
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
+111
@@ -0,0 +1,111 @@
|
||||
package dev.ltms.fleet.mcp;
|
||||
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
import dev.ltms.fleet.member.ClaudeCodeLauncher;
|
||||
import dev.ltms.fleet.peer.PeerLauncher;
|
||||
import dev.ltms.fleet.placement.BackendQuarantine;
|
||||
import io.modelcontextprotocol.spec.McpSchema;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #446 follow-up (criterion 3): a mutation battery run against the merged PR #457 proved
|
||||
* that deleting either {@code row.put("model", model);} or {@code row.put("reason", reason);}
|
||||
* from {@link FleetMcp#profilesView} left all 1586 existing tests green — nothing exercised a
|
||||
* {@link FleetMcp.QuarantineSource} whose {@code modelFor}/{@code reasonFor} actually return a
|
||||
* value. This class is what makes both deletions red, and also pins the two "must be omitted, not
|
||||
* emitted as null/blank" branches those two {@code if} guards exist for — {@link
|
||||
* FleetMcpTest#profilesReportsAQuarantinedCredential} covers the credential/quarantine-duration
|
||||
* shape but supplies neither function, so it cannot distinguish "the field is missing" from "the
|
||||
* field was never asked for".
|
||||
*/
|
||||
class FleetProfilesQuarantineModelReasonFieldsTest {
|
||||
|
||||
private static final String PROFILE = "terra";
|
||||
private static final String CREDENTIAL = "cred-terra";
|
||||
|
||||
private static PeerLauncher launcher(FakeHerdr h) {
|
||||
FleetConfig.Profile profile = new FleetConfig.Profile(
|
||||
PROFILE, "http://gx00.gw:8000", "claude-opus-5", null, "FLEETD_WORKER_TOKEN", null,
|
||||
"tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
|
||||
return new ClaudeCodeLauncher(new AgentControl(h), new WorkspaceControl(h),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(PROFILE, profile), PROFILE, _ -> "tok");
|
||||
}
|
||||
|
||||
private static BackendQuarantine quarantined(String credentialId) {
|
||||
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
|
||||
quarantine.quarantine(credentialId);
|
||||
return quarantine;
|
||||
}
|
||||
|
||||
private static String textOf(McpSchema.CallToolResult r) {
|
||||
return ((McpSchema.TextContent) r.content().getFirst()).text();
|
||||
}
|
||||
|
||||
@Test
|
||||
void quarantinedRowNamesTheModelAndTheReason() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
FleetMcp.QuarantineSource source = new FleetMcp.QuarantineSource(
|
||||
_ -> CREDENTIAL, quarantined(CREDENTIAL), _ -> true,
|
||||
_ -> "claude-opus-5", _ -> "The usage limit has been reached");
|
||||
|
||||
McpSchema.CallToolResult res = FleetMcp.profiles(launcher(h), source);
|
||||
assertNotEquals(Boolean.TRUE, res.isError());
|
||||
String out = textOf(res);
|
||||
|
||||
assertTrue(out.contains("\"quarantined\""), out);
|
||||
assertTrue(out.contains("\"model\":\"claude-opus-5\""),
|
||||
"the quarantined row must name the model the fix applies to: " + out);
|
||||
assertTrue(out.contains("\"reason\":\"The usage limit has been reached\""),
|
||||
"the quarantined row must name why it was quarantined: " + out);
|
||||
}
|
||||
|
||||
@Test
|
||||
void quarantinedRowOmitsModelWhenNoneIsConfigured() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
// null and "" (blank) both exercise the same production guard (`!model.isBlank()`) —
|
||||
// covered together since both must produce the identical outcome: the key absent.
|
||||
for (String noModel : new String[] {null, "", " "}) {
|
||||
FleetMcp.QuarantineSource source = new FleetMcp.QuarantineSource(
|
||||
_ -> CREDENTIAL, quarantined(CREDENTIAL), _ -> true,
|
||||
_ -> noModel, _ -> "some reason");
|
||||
|
||||
McpSchema.CallToolResult res = FleetMcp.profiles(launcher(h), source);
|
||||
String out = textOf(res);
|
||||
|
||||
assertTrue(out.contains("\"quarantined\""), out);
|
||||
assertFalse(out.contains("\"model\""),
|
||||
"modelFor returned " + (noModel == null ? "null" : "'" + noModel + "'")
|
||||
+ " — the key must be ABSENT, not emitted as null or an empty string: " + out);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void quarantinedRowOmitsReasonWhenNoneIsKnown() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
for (String noReason : new String[] {null, "", " "}) {
|
||||
FleetMcp.QuarantineSource source = new FleetMcp.QuarantineSource(
|
||||
_ -> CREDENTIAL, quarantined(CREDENTIAL), _ -> true,
|
||||
_ -> "claude-opus-5", _ -> noReason);
|
||||
|
||||
McpSchema.CallToolResult res = FleetMcp.profiles(launcher(h), source);
|
||||
String out = textOf(res);
|
||||
|
||||
assertTrue(out.contains("\"quarantined\""), out);
|
||||
assertFalse(out.contains("\"reason\""),
|
||||
"reasonFor returned " + (noReason == null ? "null" : "'" + noReason + "'")
|
||||
+ " — the key must be ABSENT, not emitted as null or an empty string: " + out);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -762,6 +762,250 @@ class OpenCodeLauncherTest {
|
||||
"no IDE server when ideMcpUrl is unset");
|
||||
}
|
||||
|
||||
// --- fleetd #393: memberSkills seeding must actually reach an opencode member -------------------
|
||||
//
|
||||
// Before this fix, GitWorktrees#seedSkills copied skill folders into EVERY provisioned
|
||||
// worktree's .claude/skills/ and logged "skill seeding: N of M" regardless of which kind ended
|
||||
// up spawning into that worktree — a claim that held for kind: claude-code (the CLI discovers
|
||||
// that directory on its own) but was a guaranteed no-op for kind: opencode, which has no such
|
||||
// discovery. These tests drive the REAL GitWorktrees#add seeding path (not a hand-built
|
||||
// .claude/skills/ fixture), then spawn an opencode-kind member against the seeded worktree and
|
||||
// assert on what the member can actually consume — an instructions[] entry — not on the
|
||||
// seeding log alone. A minimal, non-hermetic git repo is enough here: unlike
|
||||
// GitWorktreesTest's own seeding tests, nothing in this file cares about core.excludesFile
|
||||
// composition, only about what lands in .claude/skills/ and whether OpenCodeLauncher reads it.
|
||||
|
||||
private static void git(Path cwd, String... args) throws Exception {
|
||||
Process p = new ProcessBuilder(prepend("git", args)).directory(cwd.toFile())
|
||||
.redirectErrorStream(true).start();
|
||||
String out = new String(p.getInputStream().readAllBytes());
|
||||
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git timed out: git " + String.join(" ", args));
|
||||
assertEquals(0, p.exitValue(), "git " + String.join(" ", args) + " failed:\n" + out);
|
||||
}
|
||||
|
||||
private static List<String> prepend(String head, String... rest) {
|
||||
List<String> cmd = new ArrayList<>();
|
||||
cmd.add(head);
|
||||
cmd.addAll(List.of(rest));
|
||||
return cmd;
|
||||
}
|
||||
|
||||
private static Path initRepo(Path dir) throws Exception {
|
||||
Files.createDirectories(dir);
|
||||
git(dir, "init", "-q", "-b", "main");
|
||||
git(dir, "config", "user.email", "test@example.invalid");
|
||||
git(dir, "config", "user.name", "Test");
|
||||
Files.writeString(dir.resolve("README.md"), "seed\n");
|
||||
git(dir, "add", "README.md");
|
||||
git(dir, "commit", "-q", "-m", "seed");
|
||||
return dir;
|
||||
}
|
||||
|
||||
@Test
|
||||
void aSeededSkillReachesTheOpencodeMembersInstructionsArray(@TempDir Path tmp) throws Exception {
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
Path skillsSource = tmp.resolve("skills-src");
|
||||
Path skillFile = skillsSource.resolve("implementer").resolve("SKILL.md");
|
||||
Files.createDirectories(skillFile.getParent());
|
||||
Files.writeString(skillFile, "IMPLEMENTER PROCEDURE\n");
|
||||
|
||||
// The real seeding path (fleetd #362), not a hand-built .claude/skills/ fixture — proves
|
||||
// OpenCodeLauncher reads what GitWorktrees#add actually produced.
|
||||
dev.ltms.fleet.session.GitWorktrees worktrees =
|
||||
new dev.ltms.fleet.session.GitWorktrees(tmp.resolve("wts").toString(), null, skillsSource.toString());
|
||||
String wt = worktrees.add(repo.toString(), "cb-393-opencode", "HEAD");
|
||||
Path seededSkillMd = Path.of(wt, ".claude", "skills", "implementer", "SKILL.md");
|
||||
assertTrue(Files.exists(seededSkillMd),
|
||||
"sanity: the real seeding step must have copied the skill into the worktree");
|
||||
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(OpenCodeLauncher.class);
|
||||
Level original = logger.getLevel();
|
||||
logger.setLevel(Level.INFO);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
|
||||
String cfgPath;
|
||||
try {
|
||||
// No mcpUrl, no ideUrl, no fleet (no charter): the seeded skill alone must be enough to
|
||||
// trigger OPENCODE_CONFIG — proves the gate itself was updated, not only writeConfig's body.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Path configRoot = Files.createDirectory(tmp.resolve("configs"));
|
||||
service(herdr, configRoot, opencodeIdeCfg(null, null, wt)).spawn();
|
||||
cfgPath = startEnv(herdr).get("OPENCODE_CONFIG");
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
logger.setLevel(original);
|
||||
}
|
||||
|
||||
assertNotNull(cfgPath, "a seeded skill with nothing else configured must still write a config");
|
||||
JsonNode json = new ObjectMapper().readTree(Path.of(cfgPath).toFile());
|
||||
List<String> instructions = new ArrayList<>();
|
||||
json.path("instructions").forEach(n -> instructions.add(n.asText()));
|
||||
assertTrue(instructions.contains(seededSkillMd.toAbsolutePath().toString()),
|
||||
"the seeded skill's SKILL.md must be an instructions[] entry — got: " + instructions);
|
||||
|
||||
List<String> infos = appender.list.stream()
|
||||
.filter(e -> e.getLevel() == Level.INFO)
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.toList();
|
||||
assertTrue(infos.stream().anyMatch(m -> m.contains("skill delivery") && m.contains("1 of 1")),
|
||||
"the launcher must log, kind-aware, that it delivered the skill — got:\n" + infos);
|
||||
}
|
||||
|
||||
@Test
|
||||
void aSkillFolderWithoutSkillMdIsNeverDeliveredAndTheLogNamesIt(@TempDir Path tmp) throws Exception {
|
||||
Path wt = Files.createDirectories(tmp.resolve("wt"));
|
||||
Path goodSkill = wt.resolve(".claude").resolve("skills").resolve("implementer");
|
||||
Files.createDirectories(goodSkill);
|
||||
Files.writeString(goodSkill.resolve("SKILL.md"), "GOOD\n");
|
||||
Path halfShipped = wt.resolve(".claude").resolve("skills").resolve("half-shipped");
|
||||
Files.createDirectories(halfShipped);
|
||||
Files.writeString(halfShipped.resolve("README.md"), "no SKILL.md here\n");
|
||||
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(OpenCodeLauncher.class);
|
||||
Level original = logger.getLevel();
|
||||
logger.setLevel(Level.INFO);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
|
||||
String cfgPath;
|
||||
try {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Path configRoot = Files.createDirectory(tmp.resolve("configs"));
|
||||
service(herdr, configRoot, opencodeIdeCfg(null, null, wt.toString())).spawn();
|
||||
cfgPath = startEnv(herdr).get("OPENCODE_CONFIG");
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
logger.setLevel(original);
|
||||
}
|
||||
|
||||
JsonNode json = new ObjectMapper().readTree(Path.of(cfgPath).toFile());
|
||||
List<String> instructions = new ArrayList<>();
|
||||
json.path("instructions").forEach(n -> instructions.add(n.asText()));
|
||||
assertTrue(instructions.contains(goodSkill.resolve("SKILL.md").toAbsolutePath().toString()),
|
||||
"the well-formed skill is still delivered alongside the malformed one");
|
||||
assertFalse(instructions.stream().anyMatch(i -> i.contains("half-shipped")),
|
||||
"a skill folder with no SKILL.md can never become an instructions[] entry");
|
||||
|
||||
List<String> infos = appender.list.stream()
|
||||
.filter(e -> e.getLevel() == Level.INFO)
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.toList();
|
||||
assertTrue(infos.stream().anyMatch(m -> m.contains("skill delivery") && m.contains("1 of 2")
|
||||
&& m.contains("half-shipped") && m.contains("could not be delivered")),
|
||||
"the log must say plainly which folder could not be consumed and why — got:\n" + infos);
|
||||
}
|
||||
|
||||
// --- fleetd #393 follow-up: instructions[] has three writers (charter, seeded skills, IDE
|
||||
// rules), and no test above ever exercises more than one or two of them together. A writer
|
||||
// that flips from withArray (get-or-create) to putArray (create-or-REPLACE) silently deletes
|
||||
// every entry written before it — proven live on this branch's merge: switching just the
|
||||
// skills writer to putArray left the entire suite (1603 tests) green while deleting the
|
||||
// charter entry an opencode member needs for its role contract. That hazard was found by the
|
||||
// fleet01 lead and independently verified against this branch; it is not a defect in the
|
||||
// skills-delivery or logging tests above, which both hold up under their own mutations — the
|
||||
// gap is that none of them combine all three writers in one config.
|
||||
//
|
||||
// These three tests assert instructions[] CONTENT as an exact, ordered list, not a size or a
|
||||
// "contains" check: a putArray mutation can replace N entries with a different N entries of
|
||||
// the same count, so only a content comparison can tell "all three paths present" apart from
|
||||
// "two paths present that replaced the earlier ones".
|
||||
|
||||
@Test
|
||||
void instructionsArrayHoldsExactlyTheCharterWhenNothingElseWritesToIt(@TempDir Path root) throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FleetConfig.Fleet fleet = new FleetConfig.Fleet(Map.of(), Map.of(), Map.of(), Map.of(),
|
||||
Map.of("dev", "role rule"), null);
|
||||
Path cwd = Files.createDirectory(root.resolve("checkout"));
|
||||
service(herdr, root, opencodeIdeCfg(null, null, cwd.toString()), () -> fleet).spawn();
|
||||
|
||||
String cfgPath = startEnv(herdr).get("OPENCODE_CONFIG");
|
||||
assertNotNull(cfgPath, "a role charter alone still writes a config");
|
||||
JsonNode json = new ObjectMapper().readTree(Path.of(cfgPath).toFile());
|
||||
Path charter = Path.of(cfgPath).resolveSibling("member-charter.md");
|
||||
assertTrue(Files.exists(charter), "the charter file was written");
|
||||
|
||||
List<String> instructions = new ArrayList<>();
|
||||
json.path("instructions").forEach(n -> instructions.add(n.asText()));
|
||||
assertEquals(List.of(charter.toAbsolutePath().toString()), instructions,
|
||||
"with only the charter writer active, instructions[] holds exactly one entry: the "
|
||||
+ "charter — got: " + instructions);
|
||||
}
|
||||
|
||||
@Test
|
||||
void instructionsArrayHoldsCharterThenIdeRulesInOrder(@TempDir Path root) throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FleetConfig.Fleet fleet = new FleetConfig.Fleet(Map.of(), Map.of(), Map.of(), Map.of(),
|
||||
Map.of("dev", "role rule"), null);
|
||||
Path cwd = Files.createDirectory(root.resolve("checkout"));
|
||||
service(herdr, root, opencodeIdeCfg(null,
|
||||
"http://127.0.0.1:29170/index-mcp/streamable-http", cwd.toString()), () -> fleet).spawn();
|
||||
|
||||
String cfgPath = startEnv(herdr).get("OPENCODE_CONFIG");
|
||||
assertNotNull(cfgPath, "charter + IDE rules still writes a config");
|
||||
JsonNode json = new ObjectMapper().readTree(Path.of(cfgPath).toFile());
|
||||
Path charter = Path.of(cfgPath).resolveSibling("member-charter.md");
|
||||
Path rules = Path.of(cfgPath).resolveSibling("ide-rules.md");
|
||||
assertTrue(Files.exists(charter), "the charter file was written");
|
||||
assertTrue(Files.exists(rules), "the ide-rules file was written");
|
||||
|
||||
List<String> instructions = new ArrayList<>();
|
||||
json.path("instructions").forEach(n -> instructions.add(n.asText()));
|
||||
assertEquals(List.of(charter.toAbsolutePath().toString(), rules.toAbsolutePath().toString()),
|
||||
instructions,
|
||||
"with charter + IDE-rules writers active, instructions[] holds both, charter first — "
|
||||
+ "got: " + instructions);
|
||||
}
|
||||
|
||||
@Test
|
||||
void instructionsArrayHoldsCharterThenSkillsThenIdeRulesInOrder(@TempDir Path tmp) throws Exception {
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
Path skillsSource = tmp.resolve("skills-src");
|
||||
Path skillFile = skillsSource.resolve("implementer").resolve("SKILL.md");
|
||||
Files.createDirectories(skillFile.getParent());
|
||||
Files.writeString(skillFile, "IMPLEMENTER PROCEDURE\n");
|
||||
|
||||
dev.ltms.fleet.session.GitWorktrees worktrees = new dev.ltms.fleet.session.GitWorktrees(
|
||||
tmp.resolve("wts").toString(), null, skillsSource.toString());
|
||||
String wt = worktrees.add(repo.toString(), "cb-393-follow-up", "HEAD");
|
||||
Path seededSkillMd = Path.of(wt, ".claude", "skills", "implementer", "SKILL.md");
|
||||
assertTrue(Files.exists(seededSkillMd), "sanity: the real seeding step copied the skill");
|
||||
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FleetConfig.Fleet fleet = new FleetConfig.Fleet(Map.of(), Map.of(), Map.of(), Map.of(),
|
||||
Map.of("dev", "role rule"), null);
|
||||
Path configRoot = Files.createDirectory(tmp.resolve("configs"));
|
||||
service(herdr, configRoot, opencodeIdeCfg(null,
|
||||
"http://127.0.0.1:29170/index-mcp/streamable-http", wt), () -> fleet).spawn();
|
||||
|
||||
String cfgPath = startEnv(herdr).get("OPENCODE_CONFIG");
|
||||
assertNotNull(cfgPath, "charter + skills + IDE rules still writes a config");
|
||||
JsonNode json = new ObjectMapper().readTree(Path.of(cfgPath).toFile());
|
||||
Path charter = Path.of(cfgPath).resolveSibling("member-charter.md");
|
||||
Path rules = Path.of(cfgPath).resolveSibling("ide-rules.md");
|
||||
assertTrue(Files.exists(charter), "the charter file was written");
|
||||
assertTrue(Files.exists(rules), "the ide-rules file was written");
|
||||
|
||||
List<String> instructions = new ArrayList<>();
|
||||
json.path("instructions").forEach(n -> instructions.add(n.asText()));
|
||||
// Exact ordered list, not size or "contains": a putArray mutation on any writer after the
|
||||
// charter replaces every entry written before it, and the replacement can still be a
|
||||
// plausible-looking array of a different shape. This is the one combination all three
|
||||
// writers are active for — and per fleet01 the realistic shape on a host where weighted
|
||||
// placement makes opencode the default for most members.
|
||||
assertEquals(List.of(charter.toAbsolutePath().toString(),
|
||||
seededSkillMd.toAbsolutePath().toString(),
|
||||
rules.toAbsolutePath().toString()),
|
||||
instructions,
|
||||
"with all three writers active, instructions[] must hold charter, then the seeded "
|
||||
+ "skill, then IDE rules — in that order and with nothing replaced. A "
|
||||
+ "putArray mutation on any writer after the charter would silently drop "
|
||||
+ "earlier entries here while still producing a same-shaped array — got: "
|
||||
+ instructions);
|
||||
}
|
||||
|
||||
// --- fleetd #219: config root + discovery root under memberHerdrSocket ------------------------
|
||||
|
||||
/** A config with {@code memberHerdrSocket:} set, and optionally {@code worktreeRoot:}/{@code worktreeGroup:}. */
|
||||
|
||||
@@ -16,6 +16,7 @@ import java.util.List;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertThrows;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
@@ -221,6 +222,31 @@ class AmqpReplyInboxContractTest {
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void ackReportsHitVsMissAgainstARealBroker() throws Exception {
|
||||
// fleetd #437: fleet_ack said "acknowledged <msgId>" for a message it never touched,
|
||||
// because ReplyInbox.ack() (void) could not tell a hit from a miss. Pin the fixed
|
||||
// boolean contract against a real broker — the adapter fleetd actually runs live.
|
||||
String target = "worker-ack-contract-" + System.nanoTime();
|
||||
try (AmqpReplyInbox inbox = AmqpReplyInbox.open(uri())) {
|
||||
inbox.own(target);
|
||||
|
||||
// Never held for this target at all: must report false, not throw.
|
||||
assertFalse(inbox.ack(target, "never-held"),
|
||||
"acking a msgId never held for an owned target must report false");
|
||||
|
||||
// A real message: first ack removes it and reports true...
|
||||
inbox.publish(target, "m1", "ack me");
|
||||
assertEquals(1, awaitPeek(inbox, target).size(), "the published reply should be held");
|
||||
assertTrue(inbox.ack(target, "m1"), "acking a held reply must report true");
|
||||
assertTrue(inbox.peek(target).isEmpty(), "an acked reply is dropped");
|
||||
|
||||
// ...and the second ack of the SAME msgId has nothing left to remove: false.
|
||||
assertFalse(inbox.ack(target, "m1"),
|
||||
"acking the same msgId twice must report false the second time");
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void confirmedPublishDeliversNormally() throws Exception {
|
||||
String target = "worker-confirm-" + System.nanoTime();
|
||||
|
||||
@@ -41,20 +41,20 @@ class InMemoryReplyInboxTest {
|
||||
@Test
|
||||
void ackRemovesTheMessage() {
|
||||
inbox.publish("term_a", "m1", "hello");
|
||||
inbox.ack("term_a", "m1");
|
||||
assertTrue(inbox.ack("term_a", "m1"), "fleetd #437: ack of a real entry must report true");
|
||||
assertTrue(inbox.peek("term_a").isEmpty(), "after ack, the message is gone");
|
||||
}
|
||||
|
||||
@Test
|
||||
void ackForUnknownMsgIdIsNoOp() {
|
||||
inbox.publish("term_a", "m1", "hello");
|
||||
inbox.ack("term_a", "no-such-id"); // no-op
|
||||
assertFalse(inbox.ack("term_a", "no-such-id"), "fleetd #437: a miss must report false"); // no-op
|
||||
assertEquals(1, inbox.peek("term_a").size(), "the published message is still there");
|
||||
}
|
||||
|
||||
@Test
|
||||
void ackForUnknownTargetIsNoOp() {
|
||||
inbox.ack("no-such-target", "m1"); // no-op, should not throw
|
||||
assertFalse(inbox.ack("no-such-target", "m1"), "fleetd #437: a miss must report false"); // no-op, should not throw
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -167,7 +167,7 @@ class InMemoryReplyInboxTest {
|
||||
@Test
|
||||
void peekAndAckAreNoOpsForUnownedTarget() {
|
||||
assertTrue(inbox.peek("term_not_owned").isEmpty());
|
||||
inbox.ack("term_not_owned", "m1"); // no-op, should not throw
|
||||
assertFalse(inbox.ack("term_not_owned", "m1"), "fleetd #437: a miss must report false"); // no-op, should not throw
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
@@ -0,0 +1,369 @@
|
||||
# Plan — move the fleet back onto fleet01
|
||||
|
||||
**Goal.** Stop the fleet depending on a laptop that sleeps.
|
||||
|
||||
**Written** 2026-09-05. **Rewritten the same day** after the operator pointed out that fleet01 is a VM
|
||||
and already has the repo. They were right, and my first draft was wrong in an important way: this is
|
||||
not a stand-up. **The whole fleet already ran on fleet01 in August.** It was abandoned, not attempted.
|
||||
|
||||
Every fact below was measured on 2026-09-05. Where I did not measure something, the text says so.
|
||||
|
||||
---
|
||||
|
||||
## Status, re-measured 2026-09-10 — phases 1 to 3 are DONE, and §9 is out of date
|
||||
|
||||
This plan is prior art now, not a to-do list. Every number in this section came from a read-only
|
||||
survey over `ssh fleet01` on 2026-09-10. **Delete this section and the plan once fleet01 is the
|
||||
fleet's only daemon** — at that point the plan has been executed and stops being useful.
|
||||
|
||||
| Plan item | State on 2026-09-10 | Command that re-measures it |
|
||||
|---|---|---|
|
||||
| Phase 1, refresh + build | **done, then drifted.** Checkout is on `main` at `4887731`, **77 commits behind** `origin/main`, 0 ahead. `fleetd/target/fleetd.jar` exists, built 2026-09-10 02:10 UTC. So it was rebuilt, and main has moved since. | `git -C ~/LTMS/fleetd rev-list --count HEAD..origin/main` |
|
||||
| Phase 2, port the config | **done.** `fleetd/fleetd.yaml` exists on the host, with `bind`, `profiles`, `configReload`, `fleet`, `health`, `lifecycle`, `guard`, `memberCredentials`, `broker` and `coordinator` all present. | `test -f ~/LTMS/fleetd/fleetd/fleetd.yaml` |
|
||||
| Phase 3, supervision | **done.** `~/.config/systemd/user/fleetd.service` and `herdr.service` both exist, both `active` and `enabled`. `loginctl show-user ltms -p Linger` prints `yes`. Exactly 1 java process, so the double-daemon problem in the old `restart.sh` is not present. | `systemctl --user is-active fleetd herdr` |
|
||||
| Phase 4-5, reachable and a member proven | **partly.** `/healthz` answers 200 on loopback and reports `{"protocol":19,"version":"0.8.0"}`, matching the herdr pin. 1 herdr socket present. I did not spawn a member from here, so "a member on fleet01 opens a PR" is still unproven by me. | `curl -s http://127.0.0.1:8765/healthz` on the host |
|
||||
| Phase 6-7, cutover and reboot proof | **not done.** The Mac still runs its own daemon and is still this fleet's lead. Uptime on fleet01 is 2 weeks 2 days, so no reboot proof has been taken since the units were installed. | `uptime -p` on the host |
|
||||
| §9 "not moving the lead yet" | **out of date.** `claude` is installed at `/home/ltms/.local/bin/claude` and a fleet01 lead is live — it reaches this session over the coordination channel. So the headless-login blocker named in §9 is solved. | `ssh fleet01 'command -v claude'` |
|
||||
|
||||
### The profiles fleet01 actually offers a member
|
||||
|
||||
Measured from the live `fleetd.yaml` on the host, with `placement: weighted`:
|
||||
|
||||
| profile | kind | model | weight | maxLoad |
|
||||
|---|---|---|---|---|
|
||||
| `gx` | opencode | `gx/deepseek-v4-flash` | 100 | 2 |
|
||||
| `xf` | opencode | `opencode/mimo-v2.5-free` | 80 | 5 |
|
||||
| `local` | claude-code | `deepseek-v4-flash` | 10 | 2 |
|
||||
| `opus` | claude-code | `claude-opus-5` | 0 | 1 |
|
||||
|
||||
`weighted` spreads by ratio across every profile with a free slot, so an **unqualified** spawn on
|
||||
fleet01 lands on `gx` or `xf` almost every time. Pass `profile` explicitly there, as the canonical
|
||||
block already says.
|
||||
|
||||
### The PATH split on fleet01, which is not the one I expected
|
||||
|
||||
I went looking for a defect and found the opposite, so this is written down to stop the next
|
||||
session repeating the search.
|
||||
|
||||
`opencode` is installed at `/home/ltms/.opencode/bin/opencode`. Whether a shell can see it depends
|
||||
on which kind of shell it is:
|
||||
|
||||
```
|
||||
zsh -ic 'command -v opencode' -> /home/ltms/.opencode/bin/opencode (1 PATH entry)
|
||||
zsh -lc 'command -v opencode' -> nothing (0 PATH entries)
|
||||
```
|
||||
|
||||
So it is on the **interactive** PATH (`.zshrc`), not the login one. Two consequences, and they
|
||||
point in opposite directions:
|
||||
|
||||
- A herdr pane on Linux is a plain non-login interactive zsh, so a pane **can** launch `opencode`.
|
||||
fleetd types the launch command into the pane rather than exec'ing it, so `gx` and `xf` are not
|
||||
broken by this. I have not spawned one to confirm, so that is inference from the shell
|
||||
measurement plus the typing behaviour, not an end-to-end result.
|
||||
- The daemon itself runs `ExecStart=/bin/zsh -lc "exec java -jar target/fleetd.jar fleetd.yaml"` —
|
||||
a **login** shell, on purpose, because credentials live in `.zprofile`. Its own PATH has 10
|
||||
entries and none contains `opencode`.
|
||||
|
||||
**The general shape: the login shell and the interactive shell see different PATHs, and which one
|
||||
matters depends on whether fleetd types a command or execs it.** Credentials live on the login
|
||||
side; `~/.opencode/bin` lives on the interactive side. Anything fleetd must exec itself is
|
||||
invisible to it if it lives only on the interactive PATH — and that is the mirror of the trap the
|
||||
login-shell `ExecStart` was added to fix.
|
||||
|
||||
---
|
||||
|
||||
---
|
||||
|
||||
## 1. My first draft was wrong — read this before the rest
|
||||
|
||||
I wrote "fleetd is not deployed on fleet01". That was wrong, and I got there by looking in one place
|
||||
and concluding about all of them. Three times:
|
||||
|
||||
| I checked | I concluded | What is actually true |
|
||||
|---|---|---|
|
||||
| `~/LTMS/claude-bridge` | "no repo checkout" | the repo is at `~/LTMS/fleetd` — the directory follows the **renamed** repo, and my own notes record that rename |
|
||||
| `~/.config/herdr/herdr.sock` | "herdr never ran" | the socket lives at `~/.config/herdr/sessions/fleet01/herdr.sock` — a **named session**, and its server log runs to Aug 28 |
|
||||
| both of the above | "fleetd is not deployed" | `~/LTMS/fleetd/fleetd-run/` holds `restart.sh`, `start-herdr.sh` and a `fleetd.out` from **Aug 24** |
|
||||
|
||||
The lesson is the one already written down here: enumerate one channel, conclude about all of them.
|
||||
A single-path check is not a survey.
|
||||
|
||||
---
|
||||
|
||||
## 2. What already worked on fleet01, proven from its own log
|
||||
|
||||
`~/LTMS/fleetd/fleetd-run/fleetd.out` covers 06:07 to 16:48 on 2026-08-24. It shows:
|
||||
|
||||
```
|
||||
4 distinct panes pane=c2504101-... w1:pC w1:pE
|
||||
4 worktrees created and removed /home/ltms/LTMS/.fleet-worktrees/{05f2a7-4,dd9f51-1,f498dd-2,fdb522-3}
|
||||
4 allow-list decisions "memberCredentials allow-list: pane w1:pC allowed 16 of 32 environment variables"
|
||||
AMQP on 127.0.0.1:5672 "AMQP connection recovered; cleared held replies for fresh redelivery"
|
||||
the full turn machinery SessionManager transitions, ReplyPushLoop, CompletionResolver fallback
|
||||
```
|
||||
|
||||
So on fleet01, already: fleetd listened, herdr made panes, members spawned into git worktrees, the
|
||||
credential allow-list fed them their environment, and the broker link was **loopback**.
|
||||
|
||||
That last point is the whole reason for this move. The AMQP resets in section 3 cannot happen to a
|
||||
loopback connection.
|
||||
|
||||
### The two traps I was going to design around are already solved there
|
||||
|
||||
My first draft named these as the biggest risks. Both were already handled in August:
|
||||
|
||||
**Trap A — systemd sources no login shell.** `restart.sh` already starts through one, and says why in
|
||||
its own comment:
|
||||
|
||||
```sh
|
||||
# 1. Start java from a LOGIN shell (zsh -lc). ~/.zprofile is where the credentials live, and a
|
||||
# non-login shell starts the daemon fine with an empty AI_GATEWAY_TOKEN -- a failure that
|
||||
# stays invisible until a member actually needs it.
|
||||
setsid zsh -lc "exec java -jar target/bridged.jar fleetd.yaml" < /dev/null > "$RUN/fleetd.out" 2>&1 &
|
||||
```
|
||||
|
||||
**Trap B — a Linux herdr pane is a plain zsh, so members get no credentials.** Not a problem, and not
|
||||
for the reason I assumed. Members do not inherit from the pane's shell profile — fleetd hands them an
|
||||
allow-list. The log proves it ran: `allowed 16 of 32 environment variables`, four times. The old
|
||||
`fleetd.yaml` has a `memberCredentials:` block that configures it.
|
||||
|
||||
**The headless pty trap is solved too.** `start-herdr.sh` carries the fix and the explanation:
|
||||
|
||||
```sh
|
||||
# Why the size matters: herdr creates each pane sized to the attached client's view. Started
|
||||
# under a pty with no winsize, the client reports 0x0, and every pane.split / workspace.create
|
||||
# then fails with "ghostty error -2" -- libghostty refusing a 0x0 surface.
|
||||
cat > /tmp/herdr-inner.sh <<'INNER'
|
||||
stty rows 50 cols 200 2>/dev/null || true
|
||||
exec herdr --session fleet01
|
||||
INNER
|
||||
setsid script -qfec /tmp/herdr-inner.sh /dev/null < /dev/null > /dev/null 2>&1 &
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## 3. Why we are moving
|
||||
|
||||
The daemon runs on a Mac laptop. On battery it idle-sleeps after **one minute**:
|
||||
|
||||
```
|
||||
pmset -g custom -> Battery Power: sleep 1
|
||||
AC Power: sleep 0
|
||||
```
|
||||
|
||||
Since the last restart: 16 `Connection reset` events and 6 ERROR lines, all on the AMQP link. I
|
||||
matched every one of the 16 to the nearest sleep or wake event in `pmset -g log`. Largest gap **55
|
||||
seconds**; most under 20. Not one reset lacked a nearby sleep or wake. The broker never sees a network
|
||||
fault — it sees the client stop sending heartbeats, then closes.
|
||||
|
||||
Every reconnect worked, so no message was lost. **The lost messages are not the problem.** The problem
|
||||
is that a member mid-turn freezes with the host, and a long worker turn with nobody typing is exactly
|
||||
the case that goes idle.
|
||||
|
||||
---
|
||||
|
||||
## 4. How the two hosts relate today
|
||||
|
||||
This answers "why is fleet01 related to this Mac at all?"
|
||||
|
||||
```
|
||||
Mac laptop fleet01 (KVM/QEMU guest)
|
||||
┌────────────────────────────┐ ┌──────────────────────┐
|
||||
│ Claude Code lead │ │ LavinMQ :5672 │
|
||||
│ fleetd 127.0.0.1:8765 │ AMQP over │ vhost /mac │
|
||||
│ herdr │──Tailscale────>│ vhost /fleet01 │
|
||||
│ members + worktrees │ utun4, 1280 │ │
|
||||
└────────────────────────────┘ │ (fleetd idle since │
|
||||
│ Aug 24) │
|
||||
└──────────────────────┘
|
||||
```
|
||||
|
||||
**Today fleet01 runs only the broker.** Everything else — daemon, herdr, members — is on the Mac. The
|
||||
single link between them is the Mac's fleetd opening AMQP to `10.10.20.13:5672` across Tailscale.
|
||||
|
||||
So the errors I reported were **the Mac's client dying when the Mac slept**, not fleet01 failing.
|
||||
fleet01 was healthy throughout: the container is up 11 days and its log shows a clean heartbeat
|
||||
timeout each time, which is what a broker sees when a client vanishes.
|
||||
|
||||
After the move that arrow becomes loopback and the whole class of problem is gone.
|
||||
|
||||
---
|
||||
|
||||
## 5. The shape we are restoring
|
||||
|
||||
```mermaid
|
||||
flowchart LR
|
||||
subgraph MAC["Mac laptop (free to sleep)"]
|
||||
LEAD["Claude Code lead"]
|
||||
TUN["ssh -N -L"]
|
||||
end
|
||||
subgraph F01["fleet01 (KVM guest, always on)"]
|
||||
FD["fleetd<br/>127.0.0.1:8765"]
|
||||
HD["herdr --session fleet01"]
|
||||
WT["members<br/>~/LTMS/.fleet-worktrees"]
|
||||
MQ["LavinMQ<br/>127.0.0.1:5672"]
|
||||
end
|
||||
LEAD --> TUN
|
||||
TUN -->|"ssh over Tailscale"| FD
|
||||
FD -->|"unix socket"| HD
|
||||
HD --> WT
|
||||
FD -->|"AMQP, loopback"| MQ
|
||||
WT -->|"AMQP, loopback"| MQ
|
||||
```
|
||||
|
||||
*The lead stays on the Mac. Everything that must survive a sleep is already able to run on fleet01.*
|
||||
|
||||
Two constraints fix this shape:
|
||||
|
||||
1. **fleetd, herdr and the worktrees must share one filesystem.** The herdr link is a Unix socket plus
|
||||
absolute path strings. `herdr --remote` is terminal attach, not a transport.
|
||||
2. **fleetd fails fast on a non-loopback bind without token auth**, by its own design. So we do not
|
||||
expose `:8765` to the `10.10.20.0/24` LAN. An SSH tunnel keeps the bind on loopback and needs no
|
||||
new secret.
|
||||
|
||||
**Why the lead stays on the Mac for now.** A session not in a herdr pane resolves as `primary`, so a
|
||||
Mac-side session over the tunnel works with nothing new. What we give up: async ticket nudges type
|
||||
into the lead's pane, and fleetd cannot type into a pane on another host — so `wait:false` tickets
|
||||
stop nudging and I poll instead. Moving the lead as well is section 9; its hard part is
|
||||
authenticating Claude Code on a headless box, which has nothing to do with sleep and must not block
|
||||
this.
|
||||
|
||||
---
|
||||
|
||||
## 6. The actual gap
|
||||
|
||||
Everything below is what stands between "it ran in August" and "it runs supervised today".
|
||||
|
||||
| # | Gap | Measured state |
|
||||
|---|---|---|
|
||||
| 1 | **Checkout is stale** | branch `cb-634-ide-mcp`, HEAD `7655f1b` (2026-08-24), **290 commits behind** `origin/main`, 0 ahead |
|
||||
| 2 | **Never rebuilt after the rename** | no `target/` anywhere; the scripts still say `bridged.jar` and `cd bridged`, but the tree is now `fleetd/` and `fleetd.jar` |
|
||||
| 3 | **No live `fleetd.yaml`** | gitignored, so not in git. Three backups exist under `bridged/` — the newest is `fleetd.yaml.bak-cb634-pin`, 5512 bytes, and it is **clean of inline secrets** (0 inline passwords, 6 uses of `uriEnv`/`tokenEnv`) |
|
||||
| 4 | **No supervision** | `Linger=no`; **zero** systemd user unit files. Only the hand-rolled `restart.sh` / `start-herdr.sh` |
|
||||
| 5 | **Nothing running now** | no herdr process, no fleetd, no answer on `:8765/healthz` |
|
||||
|
||||
Untracked files in the checkout: `.idea/`, `fleetd-run/`, `docs/CB-634-Worker-IDE-Worktree.md`, and the
|
||||
three yaml backups. All are **untracked, none modified**, and `7655f1b` is already an ancestor of
|
||||
`origin/main` — so nothing is lost by updating the branch. Keep `fleetd-run/` and the backups; they
|
||||
are the prior art this plan is built on.
|
||||
|
||||
### The old config's keys, which tell us what to port
|
||||
|
||||
```
|
||||
bind: herdrSocket: /home/ltms/.config/herdr/sessions/fleet01/herdr.sock
|
||||
profiles: placement: weighted
|
||||
configReload: fleet: health:
|
||||
lifecycle: guard: worktreeRoot: /home/ltms/LTMS/.fleet-worktrees
|
||||
memberCredentials: broker:
|
||||
```
|
||||
|
||||
`herdrSocket` already points at the **named-session** path, and `memberCredentials` is already
|
||||
configured. Those two are what made members work.
|
||||
|
||||
### What the repo already has for this
|
||||
|
||||
`deploy/fleetd.service` exists and is written for Linux. Three lines need fleet01's real paths:
|
||||
`ExecStart` names `/usr/lib/jvm/temurin-25-jdk/bin/java` (fleet01 has `/usr/bin/java`), the `PATH`
|
||||
names `/usr/share/maven/bin` (fleet01 has `/usr/bin/mvn`), and `WorkingDirectory` assumes
|
||||
`%h/src/claude-bridge`. It also declares `After=herdr.service` — **and no `herdr.service` exists in
|
||||
`deploy/`**. Writing that unit, from `start-herdr.sh`, is the one genuinely new piece of code here.
|
||||
|
||||
---
|
||||
|
||||
## 7. Phases
|
||||
|
||||
```mermaid
|
||||
flowchart TD
|
||||
P1["1. Refresh<br/>update + build"]
|
||||
P2["2. Config<br/>port fleetd.yaml"]
|
||||
P3["3. Supervise<br/>linger + 2 units"]
|
||||
P4["4. Reachability<br/>tunnel, primary"]
|
||||
P5["5. Prove a member"]
|
||||
P6["6. Cutover"]
|
||||
P7["7. Reboot proof"]
|
||||
P1 --> P2 --> P3 --> P4 --> P5 --> P6 --> P7
|
||||
```
|
||||
|
||||
**Phase 1 — refresh the checkout.** Update to `origin/main` (290 commits). Keep the untracked
|
||||
`fleetd-run/` and the yaml backups. Then `mvn clean install`, **unpiped** — a pipe hides a failure
|
||||
behind a zero exit. *Check:* `fleetd/target/fleetd.jar` exists and the suite is green.
|
||||
|
||||
**Phase 2 — port the config.** Write `fleetd/fleetd.yaml` from `bridged/fleetd.yaml.bak-cb634-pin`,
|
||||
updating it for the rename and 290 commits of config changes. Diff its keys against
|
||||
`fleetd/fleetd.example.yaml` on current main, key by key, and say what changed. *Check:* the daemon
|
||||
starts and `journalctl ... | grep 'startup secret'` reports **no** `MISSING`.
|
||||
|
||||
**Phase 3 — supervision.** This is the part that never existed. `loginctl enable-linger ltms`; write
|
||||
`deploy/herdr.service` from `start-herdr.sh`, keeping the `stty` sizing; fix the three path lines in
|
||||
`deploy/fleetd.service` and keep the login-shell `ExecStart`; install both under
|
||||
`~/.config/systemd/user/`. Secrets go in a `systemctl --user edit` drop-in or a 0600
|
||||
`EnvironmentFile`, never in the committed unit. *Check:* log out of every ssh session, log back in,
|
||||
and confirm the socket and the daemon are still there. That is what lingering is for, and it is the
|
||||
check people skip.
|
||||
|
||||
**Phase 4 — reachability.** `ssh -N -L <port>:127.0.0.1:8765 fleet01` from the Mac. Use a **different
|
||||
local port** for the first test so the Mac's own daemon on `:8765` is untouched and the whole test is
|
||||
reversible. Point the lead's `.mcp.json` at it — that file is `--skip-worktree` and must never be
|
||||
committed. Wrap the tunnel in `autossh` or a launchd `KeepAlive`, because the Mac still sleeps.
|
||||
*Check:* `fleet_whoami` answers `primary`. If it answers `worker`, the `fleet.leaders.*.tab` pin does
|
||||
not match — a known demotion, not a network fault.
|
||||
|
||||
**Phase 5 — prove a member.** `healthz` can be green while every spawn fails, so only a real spawn
|
||||
proves the herdr link. Spawn one member, then give it a real unit ending in a pushed PR — the push is
|
||||
what proves `WORKER_GITEA_TOKEN` resolved. Confirm the allow-list line still appears
|
||||
(`allowed N of M environment variables`) and that N is what you expect. *Check:* a member on fleet01
|
||||
opens a PR.
|
||||
|
||||
**Phase 6 — cutover.** Drain the Mac fleet properly first: `fleet_list`, `fleet_poll` anything still
|
||||
wanted, then `fleet_stop` each member — a restart drops in-flight tickets and a member's report is
|
||||
gone with its ticket. Then stop the Mac's launchd agent. This is the migration itself, not a change to
|
||||
the Mac's settings, and it is reversible in one command.
|
||||
|
||||
**Phase 7 — reboot proof.** Reboot fleet01. Without touching anything: socket present, healthz
|
||||
answering, `fleet_whoami` still `primary`, one spawn works. Until this passes, "supervised" is a claim.
|
||||
|
||||
---
|
||||
|
||||
## 8. Risks
|
||||
|
||||
| Risk | Why it bites | What this plan does |
|
||||
|---|---|---|
|
||||
| **290 commits of config drift** | the old yaml predates the rename and much else; a silently defaulted key turns a feature off with no error | phase 2 diffs key-by-key against current `fleetd.example.yaml` |
|
||||
| **A new config key gets silently dropped** | `FleetConfig`'s back-compat constructor ladder can absorb an arity change, so a new key compiles and is defaulted away | separate ticket already in flight; matters most here because fleet01 gets a hand-edited yaml |
|
||||
| **Nothing supervises herdr** | `deploy/fleetd.service` depends on a unit that does not exist | phase 3 writes it from the working script; phase 7 proves it |
|
||||
| **Linger left off** | everything dies at logout and looks fine until then | phase 3, checked by logging out |
|
||||
| **Wrong JDK/Maven path in the unit** | fleetd propagates its PATH to every member, so a bad PATH means no member can build | three lines fixed in phase 3, proven by phase 5 |
|
||||
| **Headless pty with no winsize** | `ghostty error -2`, reported three steps later as a spawn failure | the `stty` fix is carried into `herdr.service` |
|
||||
| **Port 8765 collides during the test** | both daemons want the same local port | phase 4 uses a different local port first |
|
||||
| **Tunnel dies when the Mac sleeps** | same sleep, far smaller blast radius — it interrupts my session, not members | `autossh`/launchd `KeepAlive` |
|
||||
| **Lead demoted to worker** | the tab pin no longer matches | phase 4's check is `fleet_whoami` |
|
||||
| **Upgrading herdr** | 0.8.0 is **protocol 19**, pinned on purpose — 0.8.2 is protocol 20 and fleetd has **no version handshake** | do not upgrade herdr during this work; both hosts measured at 0.8.0 today |
|
||||
| **`placement: tab` headless** | fails for the same 0x0 reason as `ghostty error -2` | fleet01 must keep `placement: pane`, as its August config did |
|
||||
| **Profile launch settings are deferred** | editing `placement:` and waiting for the 10s config watch does nothing — the launcher holds a startup snapshot | restart the daemon after those keys, do not wait for the reload |
|
||||
|
||||
---
|
||||
|
||||
## 9. Deliberately not doing
|
||||
|
||||
- **No changes to the Mac's power settings or host config.** The point is to stop depending on it.
|
||||
- **Not touching the leftover `bridged-lavinmq` container** on the Mac. It is unused and harmless.
|
||||
- **Not moving the broker.** It is already on fleet01 and already the durable one. After the move its
|
||||
connection becomes loopback, which is the fix.
|
||||
- **Not building a second fleet.** This is a move. Two daemons on one herdr session kill each other's
|
||||
members.
|
||||
- **Not moving the lead yet.** That needs Claude Code authenticated on a headless Ubuntu box and a
|
||||
`fleet.leaders.*.tab` pin on its pane. It buys back pane nudges. It has nothing to do with sleep, so
|
||||
it must not hold up phases 1–7. Note that fleet01 already carries an `opus` profile defined purely so the lead slot resolves; it cannot spawn until someone runs `claude` and completes `/login` on the host.
|
||||
|
||||
---
|
||||
|
||||
## 10. Open questions for the operator
|
||||
|
||||
1. **vhost** — keep `/mac`, or rename now the fleet is not on the Mac? Renaming loses the existing
|
||||
queues. (The old fleet01 config used its own; phase 2 must settle which this fleet owns.)
|
||||
2. **Fallback week** after cutover, or stop the Mac daemon for good?
|
||||
3. **Delete or keep the stale `cb-634-ide-mcp` branch** on fleet01 once the checkout is updated? Its
|
||||
tip is already in main, so nothing is lost either way.
|
||||
|
||||
The repo path question from the first draft is answered: **`/home/ltms/LTMS/fleetd`**, which already
|
||||
exists. `deploy/fleetd.service` should be pointed there rather than the reverse.
|
||||
Reference in New Issue
Block a user