Compare commits
23 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| e95ed99bf7 | |||
| 789b6a8716 | |||
| 01462c9695 | |||
| 5a467e1f8b | |||
| 1fb6176783 | |||
| 49df79203c | |||
| 235644c0f0 | |||
| eccd0548ce | |||
| c4e23eebad | |||
| 7df503dfc2 | |||
| f5e02fedd6 | |||
| 29e7a06c49 | |||
| 92a96fcbd8 | |||
| 13482872bb | |||
| c1ca6273fc | |||
| 5f1b260c81 | |||
| 9d1306d442 | |||
| e54e3d87ea | |||
| bf027f10b9 | |||
| 3c5873dfe2 | |||
| e29227d5f4 | |||
| 2af13ab1ff | |||
| 0788d84be8 |
@@ -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
|
||||
|
||||
@@ -452,6 +452,15 @@ placement: weighted
|
||||
# seconds, before a spawn may land on it again. Applies to every profile's effective credential
|
||||
# (its own name, or its credentialId if set above) — there is no per-profile override. Default
|
||||
# 1800 (30 minutes) when omitted or non-positive.
|
||||
#
|
||||
# fleetd #466: this is now only the BASE of an escalating backoff, not a flat retry rate. A
|
||||
# credential quarantined again within one base cooldown of the previous quarantine ending (still
|
||||
# reporting exhausted — e.g. a weekly subscription limit that hasn't reset) backs off further:
|
||||
# cooldown doubles each such time, capped at 12x this value (~6 hours at the 1800s default). A
|
||||
# quarantine that starts after a base-cooldown's worth of quiet resets back to this value. Not
|
||||
# configurable per se — the multiplier and ceiling are constants in BackendQuarantine, not new
|
||||
# YAML keys; see its class doc for the exact formula and why there is no automatic probe to clear
|
||||
# it early (the operator's own design constraint — a probe spends the quota it's measuring).
|
||||
# DEFERRED: baked once into the BackendQuarantine built at startup — a running quarantine keeps
|
||||
# its original cooldown regardless; a new value only applies to a quarantine that starts after a
|
||||
# restart. Editing this needs a daemon restart to take effect.
|
||||
|
||||
@@ -222,7 +222,13 @@ public final class Fleetd {
|
||||
// (checked at spawn) and the exhaustion sink wired in below (written on BACKEND_EXHAUSTED).
|
||||
// The cooldown is deferred (see FleetConfig#quarantineCooldownSeconds): it is read once
|
||||
// here, at startup, and a config reload only changes it for a daemon restart.
|
||||
BackendQuarantine quarantine = new BackendQuarantine(System::nanoTime,
|
||||
// fleetd #466: escalating, not flat — a credential that keeps reporting exhaustion (e.g. a
|
||||
// weekly subscription limit, which would otherwise be retried on every ~30-minute cooldown,
|
||||
// about 336 times across the week) backs off further each consecutive time, capped at
|
||||
// BackendQuarantine.DEFAULT_MAX_COOLDOWN_MULTIPLE x the base cooldown. See BackendQuarantine's
|
||||
// class doc for the mechanism, why this never fires on cooling-off (a separate, unescalated
|
||||
// mechanism — BackendOutagePolicy below), and the reset.
|
||||
BackendQuarantine quarantine = BackendQuarantine.withEscalation(System::nanoTime,
|
||||
TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds()));
|
||||
// fleetd #201 Unit 5: one outage-cool-off tracker for the whole daemon, shared between the
|
||||
// launcher (checked at spawn, like `quarantine` above) and the backend-error sink wired in
|
||||
|
||||
@@ -70,12 +70,14 @@ import java.util.regex.PatternSyntaxException;
|
||||
* {@code fixed} (default), {@code round-robin}, or {@code weighted}
|
||||
* @param auth API authentication mode ({@code null} → {@code loopback-trust}, the
|
||||
* historical behaviour), CB-501
|
||||
* @param quarantineCooldownSeconds how long a credential stays quarantined after a
|
||||
* @param quarantineCooldownSeconds the BASE cooldown a credential is quarantined for after a
|
||||
* {@code BACKEND_EXHAUSTED} classification (CB-578 stage B); {@code null}/{@code
|
||||
* <=0} → {@link #DEFAULT_QUARANTINE_COOLDOWN_SECONDS}. Baked once into the
|
||||
* {@code BackendQuarantine} built at startup, so it is DEFERRED: changing it
|
||||
* needs a restart, and a quarantine already running keeps whatever cooldown was
|
||||
* live when it started.
|
||||
* <=0} → {@link #DEFAULT_QUARANTINE_COOLDOWN_SECONDS}. Since fleetd #466 this is
|
||||
* only the first occurrence's length — a credential quarantined again shortly
|
||||
* after this cooldown ends backs off further, up to a ceiling; see {@code
|
||||
* BackendQuarantine}'s class doc. Baked once into the {@code BackendQuarantine}
|
||||
* built at startup, so it is DEFERRED: changing it needs a restart, and a
|
||||
* quarantine already running keeps whatever cooldown was live when it started.
|
||||
* @param memberCredentials deny-by-default policy (CB-596) for which of the operator's own host
|
||||
* credentials a spawned member's pane inherits. {@code null} (the block
|
||||
* omitted) blocks nothing — see {@link MemberCredentials}.
|
||||
|
||||
@@ -455,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) -> {
|
||||
@@ -608,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);
|
||||
@@ -1254,6 +1269,13 @@ public final class FleetMcp {
|
||||
* (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.
|
||||
*
|
||||
* <p>fleetd #466 scope item 2: each {@code quarantined} row also names {@code
|
||||
* quarantineAttempt} — 1 for a first-time exhaustion, 2 for the second in a row, and so on — so
|
||||
* an operator sees "this is the 5th time" instead of inferring it from {@code
|
||||
* quarantinedForSeconds} alone. Read off {@link BackendQuarantine#status(String)}, the same
|
||||
* one-call accessor {@code quarantinedForSeconds} itself comes from here (see its doc) — never a
|
||||
* separately derived count.
|
||||
*/
|
||||
public static Map<String, Object> profilesView(PeerLauncher workers, QuarantineSource quarantine, OutageSource outage) {
|
||||
Map<String, Object> result = new LinkedHashMap<>();
|
||||
@@ -1266,10 +1288,14 @@ public final class FleetMcp {
|
||||
exhaustionDetectionArmed.put(profile, quarantine.exhaustedPatternArmed().apply(profile));
|
||||
String credentialId = quarantine.credentialIdFor().apply(profile);
|
||||
if (credentialId != null) {
|
||||
quarantine.quarantine().remainingSeconds(credentialId).ifPresent(remaining -> {
|
||||
// fleetd #466 scope item 2: quarantinedForSeconds and quarantineAttempt come off the
|
||||
// ONE BackendQuarantine#status(credentialId) call — never a second, independent read
|
||||
// for the attempt count — so the two can never disagree about which streak this is.
|
||||
quarantine.quarantine().status(credentialId).ifPresent(status -> {
|
||||
Map<String, Object> row = new LinkedHashMap<>();
|
||||
row.put("credentialId", credentialId);
|
||||
row.put("quarantinedForSeconds", remaining);
|
||||
row.put("quarantinedForSeconds", status.remainingSeconds());
|
||||
row.put("quarantineAttempt", status.repeatCount());
|
||||
String model = quarantine.modelFor().apply(profile);
|
||||
if (model != null && !model.isBlank()) {
|
||||
row.put("model", model);
|
||||
@@ -1389,12 +1415,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)
|
||||
@@ -1416,9 +1478,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,
|
||||
@@ -1668,10 +1735,13 @@ public final class FleetMcp {
|
||||
}
|
||||
String credentialId = quarantine.credentialIdFor().apply(profile);
|
||||
if (credentialId != null) {
|
||||
quarantine.quarantine().remainingSeconds(credentialId).ifPresent(remaining -> {
|
||||
// fleetd #466 scope item 2: same one-call status() read as profilesView above — see its
|
||||
// comment for why this must not become two separate lookups.
|
||||
quarantine.quarantine().status(credentialId).ifPresent(status -> {
|
||||
row.put("free", 0);
|
||||
row.put("credentialId", credentialId);
|
||||
row.put("quarantinedForSeconds", remaining);
|
||||
row.put("quarantinedForSeconds", status.remainingSeconds());
|
||||
row.put("quarantineAttempt", status.repeatCount());
|
||||
});
|
||||
}
|
||||
String outageCredentialId = outage.credentialIdFor().apply(profile);
|
||||
|
||||
@@ -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
|
||||
*/
|
||||
|
||||
@@ -3,6 +3,7 @@ package dev.ltms.fleet.placement;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
import java.util.Optional;
|
||||
import java.util.OptionalLong;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.function.LongSupplier;
|
||||
@@ -21,46 +22,158 @@ import java.util.function.LongSupplier;
|
||||
* <p>The clock is injected ({@link LongSupplier}, conventionally {@code System::nanoTime} like
|
||||
* {@code FleetHealthMonitor}), never read inline, so a quarantine's expiry is testable without a
|
||||
* real sleep.
|
||||
*
|
||||
* <h2>Escalation (fleetd #466)</h2>
|
||||
* A flat cooldown does not fit every exhaustion. A backend that reports "out of capacity for the
|
||||
* rest of the hour" recovers in one cooldown; a weekly subscription limit does not — it keeps
|
||||
* reporting exhausted on every attempt made before the window resets, so a flat 30-minute cooldown
|
||||
* (the default {@code cooldownNanos}) means roughly 336 pointless spawn attempts across a week, one
|
||||
* every cooldown.
|
||||
*
|
||||
* <p><strong>This class only ever sees the exhaustion signal.</strong> Its only production caller is
|
||||
* {@code Fleetd.exhaustionSink}, wired to fire on a {@code BACKEND_EXHAUSTED} classification alone.
|
||||
* The daemon's other outage state — a credential "cooling off" after repeated non-exhaustion
|
||||
* backend errors (an HTTP 5xx storm, say) — is a separate mechanism, {@code BackendOutagePolicy},
|
||||
* with its own short fixed 60s cooldown and no repeat tracking. The two are never merged: escalating
|
||||
* on a cooling-off signal would turn a transient 5xx storm into a multi-hour backoff, which is
|
||||
* exactly the failure this ticket is not asking for. Confirmed by reading every call site of
|
||||
* {@link #quarantine} — {@code BackendOutagePolicy} has its own {@code coolOff} method and never
|
||||
* calls this one.
|
||||
*
|
||||
* <p><strong>Mechanism</strong> — the {@link #withEscalation} constructors track, per credential, how
|
||||
* many times in a row {@link #quarantine} has been called without an intervening "quiet" gap.
|
||||
* Each call computes {@code cooldownNanos * backoffMultiplier ^ (repeatCount - 1)}, capped at
|
||||
* {@code maxCooldownNanos}. A call counts as a continuation of the same streak — {@code repeatCount}
|
||||
* increments — when it arrives no more than one base {@code cooldownNanos} after the previous
|
||||
* quarantine's deadline (this covers both "still quarantined" and "quarantine just expired and it
|
||||
* was exhausted again immediately"); otherwise the streak resets and this call is treated as a fresh
|
||||
* first occurrence at the base cooldown.
|
||||
*
|
||||
* <p><strong>Reset, honestly stated.</strong> The ideal reset signal is "the cooldown expired and the
|
||||
* next attempt succeeded" — but nothing in this codebase reports a spawn success back to this class
|
||||
* (checked: {@code SessionManager} and {@code CompositePeerLauncher} never call any method here
|
||||
* except {@link #quarantine}/{@link #isQuarantined}/{@link #remainingSeconds}, none of which is a
|
||||
* success hook). Lacking that signal, the reset used here is a time-based proxy: a base-cooldown's
|
||||
* worth of quiet — no exhaustion report for that credential — since the last quarantine ended. It is
|
||||
* not proof the credential started working again, only the best available evidence without adding an
|
||||
* active probe, which is out of scope by the operator's own design constraint (no automatic probing
|
||||
* of a limited backend).
|
||||
*
|
||||
* <p><strong>Ceiling.</strong> {@code maxCooldownNanos} bounds the growth — an unbounded backoff is a
|
||||
* permanent, unrecoverable-without-a-restart outage, which would be worse than the flat-rate bug this
|
||||
* escalation fixes. {@link #withEscalation(LongSupplier, long)} defaults the ceiling to
|
||||
* {@value #DEFAULT_MAX_COOLDOWN_MULTIPLE}x the base cooldown (12x the 1800s default ≈ 6 hours), so a
|
||||
* chronically exhausted credential still gets re-tried roughly every 6 hours instead of every 30
|
||||
* minutes — about a dozen attempts a week instead of ~336.
|
||||
*
|
||||
* <p><strong>Backward compatibility.</strong> The original two-argument {@link #BackendQuarantine(
|
||||
* LongSupplier, long)} constructor is unchanged in behaviour: it is exactly {@code
|
||||
* withEscalation}'s mechanism with {@code backoffMultiplier = 1.0} and {@code maxCooldownNanos =
|
||||
* cooldownNanos}, which collapses the formula back to the original flat {@code now + cooldownNanos}
|
||||
* on every call regardless of history. Every existing call site (roughly 20 across the test suite,
|
||||
* plus {@link #none()}) keeps its current shape and behaviour unchanged.
|
||||
*/
|
||||
public final class BackendQuarantine {
|
||||
|
||||
private final ConcurrentHashMap<String, Long> quarantinedUntilNanos = new ConcurrentHashMap<>();
|
||||
/** Default growth per consecutive exhaustion streak — see the class doc's Mechanism section. */
|
||||
static final double DEFAULT_BACKOFF_MULTIPLIER = 2.0;
|
||||
/** Default ceiling, expressed as a multiple of the base cooldown — see the class doc's Ceiling section. */
|
||||
static final long DEFAULT_MAX_COOLDOWN_MULTIPLE = 12;
|
||||
|
||||
private final ConcurrentHashMap<String, QuarantineState> quarantines = new ConcurrentHashMap<>();
|
||||
private final LongSupplier nowNanos;
|
||||
private final long cooldownNanos;
|
||||
private final double backoffMultiplier;
|
||||
private final long maxCooldownNanos;
|
||||
/** True only for {@link #none()}. See {@link #quarantine} for why this exists. */
|
||||
private final boolean inert;
|
||||
|
||||
/** How many consecutive exhaustion reports a credential is on, and when the resulting cooldown ends. */
|
||||
private record QuarantineState(int repeatCount, long deadlineNanos) {
|
||||
}
|
||||
|
||||
/**
|
||||
* Flat cooldown, unchanged from before fleetd #466 — every {@link #quarantine} call blocks the
|
||||
* credential for exactly {@code cooldownNanos}, regardless of how many times it was called
|
||||
* before. Equivalent to {@link #withEscalation} with no growth ({@code backoffMultiplier = 1.0})
|
||||
* and a ceiling equal to the base cooldown, so it degrades to the identical {@code now +
|
||||
* cooldownNanos} formula every call. Kept for the existing call sites that want a fixed cooldown
|
||||
* (and for tests exercising the fixed-cooldown shape in isolation); production wiring uses
|
||||
* {@link #withEscalation} instead.
|
||||
*
|
||||
* @param nowNanos monotonic clock, injected for testability
|
||||
* @param cooldownNanos how long a fresh {@link #quarantine} call blocks the credential for;
|
||||
* must be positive
|
||||
*/
|
||||
public BackendQuarantine(LongSupplier nowNanos, long cooldownNanos) {
|
||||
this(nowNanos, cooldownNanos, false);
|
||||
this(nowNanos, cooldownNanos, 1.0, cooldownNanos, false);
|
||||
}
|
||||
|
||||
private BackendQuarantine(LongSupplier nowNanos, long cooldownNanos, boolean inert) {
|
||||
/**
|
||||
* Escalating cooldown (fleetd #466) — see the class doc's Mechanism/Reset/Ceiling sections.
|
||||
*
|
||||
* @param nowNanos monotonic clock, injected for testability
|
||||
* @param cooldownNanos base cooldown, applied to a fresh (non-streak) exhaustion; must be
|
||||
* positive
|
||||
* @param backoffMultiplier growth per consecutive exhaustion; must be {@code >= 1.0} ({@code 1.0}
|
||||
* disables growth and is exactly the flat two-argument constructor)
|
||||
* @param maxCooldownNanos ceiling on the escalated cooldown; must be {@code >= cooldownNanos}
|
||||
*/
|
||||
public BackendQuarantine(LongSupplier nowNanos, long cooldownNanos, double backoffMultiplier,
|
||||
long maxCooldownNanos) {
|
||||
this(nowNanos, cooldownNanos, backoffMultiplier, maxCooldownNanos, false);
|
||||
}
|
||||
|
||||
private BackendQuarantine(LongSupplier nowNanos, long cooldownNanos, double backoffMultiplier,
|
||||
long maxCooldownNanos, boolean inert) {
|
||||
this.nowNanos = Objects.requireNonNull(nowNanos, "nowNanos");
|
||||
if (cooldownNanos <= 0) {
|
||||
throw new IllegalArgumentException("cooldownNanos must be positive: " + cooldownNanos);
|
||||
}
|
||||
if (backoffMultiplier < 1.0) {
|
||||
throw new IllegalArgumentException("backoffMultiplier must be >= 1.0: " + backoffMultiplier);
|
||||
}
|
||||
if (maxCooldownNanos < cooldownNanos) {
|
||||
throw new IllegalArgumentException(
|
||||
"maxCooldownNanos must be >= cooldownNanos: " + maxCooldownNanos + " < " + cooldownNanos);
|
||||
}
|
||||
this.cooldownNanos = cooldownNanos;
|
||||
this.backoffMultiplier = backoffMultiplier;
|
||||
this.maxCooldownNanos = maxCooldownNanos;
|
||||
this.inert = inert;
|
||||
}
|
||||
|
||||
/**
|
||||
* Escalating cooldown with the fleetd #466 default shape: cooldown doubles
|
||||
* ({@value #DEFAULT_BACKOFF_MULTIPLIER}x) per consecutive exhaustion streak, capped at
|
||||
* {@value #DEFAULT_MAX_COOLDOWN_MULTIPLE}x the base cooldown. This is what production wiring
|
||||
* ({@code Fleetd.main}) uses.
|
||||
*
|
||||
* @param nowNanos monotonic clock, injected for testability
|
||||
* @param cooldownNanos base cooldown, applied to a fresh (non-streak) exhaustion; must be positive
|
||||
*/
|
||||
public static BackendQuarantine withEscalation(LongSupplier nowNanos, long cooldownNanos) {
|
||||
return new BackendQuarantine(nowNanos, cooldownNanos, DEFAULT_BACKOFF_MULTIPLIER,
|
||||
cooldownNanos * DEFAULT_MAX_COOLDOWN_MULTIPLE, false);
|
||||
}
|
||||
|
||||
/**
|
||||
* Inert quarantine — {@link #quarantine} does nothing on this instance, so nothing is ever
|
||||
* quarantined. The explicit stand-in a caller (or a test not exercising this feature) passes
|
||||
* instead of a defaulting overload, exactly like {@code ExhaustedPatternLookup.none()}.
|
||||
*/
|
||||
public static BackendQuarantine none() {
|
||||
return new BackendQuarantine(() -> 0L, 1, true);
|
||||
return new BackendQuarantine(() -> 0L, 1, 1.0, 1, true);
|
||||
}
|
||||
|
||||
/**
|
||||
* Quarantine {@code credentialId} for the configured cooldown, starting now. A repeat call while
|
||||
* already quarantined restarts the cooldown at full length — a fresh refusal is fresh evidence the
|
||||
* account is still exhausted, not a reason to let an earlier, shorter wait stand.
|
||||
* Quarantine {@code credentialId} starting now. On a flat instance (the two-argument
|
||||
* constructor) this always blocks for exactly {@code cooldownNanos}, restarting the cooldown at
|
||||
* full length on every call — a fresh refusal is fresh evidence the account is still exhausted,
|
||||
* not a reason to let an earlier, shorter wait stand. On an escalating instance ({@link
|
||||
* #withEscalation}) the cooldown grows with each call that arrives within one base cooldown of
|
||||
* the previous deadline, and resets to the base cooldown once a call arrives after a longer gap
|
||||
* — see the class doc.
|
||||
*
|
||||
* <p>On {@link #none()} this is a no-op. It has to be: that instance holds a clock frozen at 0,
|
||||
* so recording a deadline would produce a quarantine that never expires — a credential locked out
|
||||
@@ -73,7 +186,13 @@ public final class BackendQuarantine {
|
||||
if (inert) {
|
||||
return;
|
||||
}
|
||||
quarantinedUntilNanos.put(credentialId, nowNanos.getAsLong() + cooldownNanos);
|
||||
long now = nowNanos.getAsLong();
|
||||
quarantines.compute(credentialId, (id, prev) -> {
|
||||
int repeatCount = (prev == null || now - prev.deadlineNanos() > cooldownNanos)
|
||||
? 1
|
||||
: prev.repeatCount() + 1;
|
||||
return new QuarantineState(repeatCount, now + escalatedCooldownNanos(repeatCount));
|
||||
});
|
||||
}
|
||||
|
||||
/** Whether {@code credentialId} is quarantined right now. */
|
||||
@@ -87,6 +206,49 @@ public final class BackendQuarantine {
|
||||
return remaining > 0 ? OptionalLong.of(toSecondsRoundedUp(remaining)) : OptionalLong.empty();
|
||||
}
|
||||
|
||||
/**
|
||||
* Remaining seconds together with which consecutive exhaustion this is (fleetd #466 scope item
|
||||
* 2) — {@code repeatCount} 1 for a first occurrence, 2 for the second in a row, and so on; see
|
||||
* {@link #quarantine}'s class-doc Mechanism section for exactly when a call continues a streak
|
||||
* versus starts a fresh one.
|
||||
*
|
||||
* <p><strong>Read together, off the one {@link QuarantineState} entry {@link #quarantine} itself
|
||||
* wrote</strong> — a single {@code quarantines.get(credentialId)}, never a separate lookup or a
|
||||
* value re-derived from {@code remainingSeconds} (e.g. inverting {@link
|
||||
* #escalatedCooldownNanos}). That inversion is not just extra work to avoid: once a streak has
|
||||
* hit {@code maxCooldownNanos}, every further consecutive exhaustion reports the identical
|
||||
* cooldown, so a derivation that starts from the cooldown value cannot tell the 4th repeat from
|
||||
* the 9th — only the stored {@code repeatCount} can. This is the same rule {@code
|
||||
* CompositePeerLauncher.modelGateState()} documents for its own gate/report pair: the report
|
||||
* reads the exact accessor the behaviour reads, so it can never disagree with what actually
|
||||
* happened (the fleetd #404/#422 lesson). {@code fleet_profiles}/{@code fleet_list}/{@code GET
|
||||
* /profiles} all call this — never {@link #remainingSeconds} plus a second, independent count —
|
||||
* for exactly that reason.
|
||||
*
|
||||
* @return empty when {@code credentialId} is not currently quarantined (including on {@link
|
||||
* #none()}, which quarantines nothing)
|
||||
*/
|
||||
public Optional<Status> status(String credentialId) {
|
||||
QuarantineState state = quarantines.get(credentialId);
|
||||
if (state == null) {
|
||||
return Optional.empty();
|
||||
}
|
||||
long remaining = state.deadlineNanos() - nowNanos.getAsLong();
|
||||
return remaining > 0
|
||||
? Optional.of(new Status(toSecondsRoundedUp(remaining), state.repeatCount()))
|
||||
: Optional.empty();
|
||||
}
|
||||
|
||||
/**
|
||||
* @param remainingSeconds seconds left on the quarantine, identical to {@link
|
||||
* #remainingSeconds(String)}'s answer for the same credential at the
|
||||
* same instant
|
||||
* @param repeatCount 1 for a first occurrence, 2 for the second consecutive one, etc. —
|
||||
* see {@link #status(String)}
|
||||
*/
|
||||
public record Status(long remainingSeconds, int repeatCount) {
|
||||
}
|
||||
|
||||
/**
|
||||
* Every currently-quarantined credential id and its remaining seconds (CB-578 stage B fleet
|
||||
* reporting) — expired entries are never included. Not pruned from the backing map here: it stays
|
||||
@@ -95,8 +257,8 @@ public final class BackendQuarantine {
|
||||
*/
|
||||
public Map<String, Long> activeRemainingSeconds() {
|
||||
Map<String, Long> out = new LinkedHashMap<>();
|
||||
quarantinedUntilNanos.forEach((credentialId, deadline) -> {
|
||||
long remaining = deadline - nowNanos.getAsLong();
|
||||
quarantines.forEach((credentialId, state) -> {
|
||||
long remaining = state.deadlineNanos() - nowNanos.getAsLong();
|
||||
if (remaining > 0) {
|
||||
out.put(credentialId, toSecondsRoundedUp(remaining));
|
||||
}
|
||||
@@ -105,8 +267,14 @@ public final class BackendQuarantine {
|
||||
}
|
||||
|
||||
private long remainingNanos(String credentialId) {
|
||||
Long deadline = quarantinedUntilNanos.get(credentialId);
|
||||
return deadline == null ? 0L : deadline - nowNanos.getAsLong();
|
||||
QuarantineState state = quarantines.get(credentialId);
|
||||
return state == null ? 0L : state.deadlineNanos() - nowNanos.getAsLong();
|
||||
}
|
||||
|
||||
/** {@code cooldownNanos * backoffMultiplier ^ (repeatCount - 1)}, capped at {@code maxCooldownNanos}. */
|
||||
private long escalatedCooldownNanos(int repeatCount) {
|
||||
double raw = cooldownNanos * Math.pow(backoffMultiplier, repeatCount - 1);
|
||||
return raw >= (double) maxCooldownNanos ? maxCooldownNanos : (long) raw;
|
||||
}
|
||||
|
||||
private static long toSecondsRoundedUp(long nanos) {
|
||||
|
||||
@@ -0,0 +1,65 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #466 follow-up: {@code Fleetd.main} builds the daemon's one {@code BackendQuarantine}
|
||||
* from {@link dev.ltms.fleet.placement.BackendQuarantine#withEscalation(java.util.function.LongSupplier,
|
||||
* long)} — the escalating factory — rather than the plain two-argument constructor, which is still a
|
||||
* flat cooldown (kept for backward compatibility, see that class's doc). {@code
|
||||
* BackendQuarantineTest} proves {@code withEscalation} itself escalates, is ceilinged, and resets;
|
||||
* it says nothing about which one {@code main} actually calls.
|
||||
*
|
||||
* <p>Measured directly: reverting {@code main} to {@code new BackendQuarantine(System::nanoTime,
|
||||
* TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds()))} — the pre-#466 flat call — compiles
|
||||
* with 0 errors and leaves the entire 1608-test suite (including every {@code BackendQuarantineTest}
|
||||
* case) green, because no other test constructs its {@code BackendQuarantine} through {@code main};
|
||||
* every one of them builds its own instance directly. That silent regression is exactly the shape
|
||||
* {@link FleetdLeadSeatWiringTest} and {@link FleetdCompletionResolverWiringTest} already guard
|
||||
* against for their own constructor arguments — this is the same class of gap for fleetd #466's
|
||||
* factory choice, following their approach.
|
||||
*
|
||||
* <p><b>This test checks source text, not runtime behaviour.</b> It never constructs a {@code
|
||||
* BackendQuarantine} and never runs {@code main} — a green result here proves only that the exact
|
||||
* text {@code main} calls {@code BackendQuarantine.withEscalation(...)} rather than the flat
|
||||
* constructor. It does not prove that call actually executes at startup (no test here starts the
|
||||
* daemon), and it does not prove the escalation reaches a real backend or credential — only
|
||||
* {@code BackendQuarantineTest} proves the factory's own behaviour, and only a live daemon proves
|
||||
* the wiring runs.
|
||||
*/
|
||||
class FleetdBackendQuarantineWiringTest {
|
||||
|
||||
private static String fleetdSource() throws Exception {
|
||||
return Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SOURCE TEXT] main's BackendQuarantine local is still built from BackendQuarantine.withEscalation(...)")
|
||||
void mainStillWiresTheEscalatingQuarantineFactory() throws Exception {
|
||||
String source = fleetdSource();
|
||||
assertTrue(source.contains(
|
||||
"BackendQuarantine quarantine = BackendQuarantine.withEscalation(System::nanoTime,\n"
|
||||
+ " TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds()));"),
|
||||
"Fleetd.main's BackendQuarantine local must still be built from "
|
||||
+ "BackendQuarantine.withEscalation(System::nanoTime, "
|
||||
+ "TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds())). Reverting to the flat "
|
||||
+ "two-argument constructor (fleetd #466's measured regression) compiles with 0 errors "
|
||||
+ "and leaves the whole suite green, including every BackendQuarantineTest case that "
|
||||
+ "proves the escalation itself works — this source check is what must go red instead. "
|
||||
+ "A reverted daemon would go back to retrying a weekly subscription limit on every "
|
||||
+ "flat ~30-minute cooldown, about 336 times across the week.");
|
||||
|
||||
// Negative form of the same check: the pre-#466 flat call, if it ever reappears at this
|
||||
// declaration, must not be mistaken for the escalating one by a looser positive-only check.
|
||||
assertFalse(source.contains(
|
||||
"BackendQuarantine quarantine = new BackendQuarantine(System::nanoTime,\n"
|
||||
+ " TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds()));"),
|
||||
"main's BackendQuarantine local must never regress to the flat two-argument constructor");
|
||||
}
|
||||
}
|
||||
@@ -25,17 +25,39 @@ import static org.junit.jupiter.api.Assumptions.assumeTrue;
|
||||
* 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}. With {@link #SHELL_READY_TIMEOUT_MS}
|
||||
* set to 0 — so input is typed at once, with no settle wait at all — the test still passed 3 of
|
||||
* 3. So the proven cause is the 800ms READ deadline being too short, not the 1000ms write delay.
|
||||
* Note the direction, because it matters: typing at 0ms works where typing at 1000ms failed. The
|
||||
* earlier explanation for this test — that input typed before the prompt is swallowed by the
|
||||
* shell's startup — is therefore NOT supported by any measurement here. Please do not repeat it
|
||||
* as the reason; if it were true, 0ms would be worse than 1000ms, and it is better.
|
||||
* 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. If you want to delete it, the honest test is whether you can make
|
||||
* this test fail by typing early. Nobody has managed that yet.
|
||||
* 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}.
|
||||
*/
|
||||
|
||||
@@ -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) ------------------------------------
|
||||
|
||||
/**
|
||||
|
||||
@@ -574,6 +574,38 @@ class FleetMcpTest {
|
||||
assertTrue(out.contains("\"quarantinedForSeconds\":1800"), out);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #466 scope item 2: {@code fleet_profiles} must carry the repeat count beside the
|
||||
* remaining seconds, and the two must come off the one {@link BackendQuarantine#status} call so
|
||||
* they can never disagree about which streak this is (see {@code profilesView}'s javadoc).
|
||||
*/
|
||||
@Test
|
||||
void profilesReportsQuarantineAttemptBesideRemainingSeconds() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
java.util.concurrent.atomic.AtomicLong now = new java.util.concurrent.atomic.AtomicLong(0L);
|
||||
BackendQuarantine quarantine = BackendQuarantine.withEscalation(now::get, TimeUnit.MINUTES.toNanos(30));
|
||||
quarantine.quarantine("shared-openai"); // attempt 1: 1800s
|
||||
now.set(TimeUnit.MINUTES.toNanos(30));
|
||||
quarantine.quarantine("shared-openai"); // attempt 2: 3600s
|
||||
FleetMcp.QuarantineSource source = new FleetMcp.QuarantineSource(
|
||||
profile -> "ltms-local".equals(profile) ? "shared-openai" : null, quarantine);
|
||||
McpSchema.CallToolResult res = FleetMcp.profiles(
|
||||
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), source);
|
||||
String out = textOf(res);
|
||||
assertTrue(out.contains("\"quarantinedForSeconds\":3600"), out);
|
||||
assertTrue(out.contains("\"quarantineAttempt\":2"), out);
|
||||
}
|
||||
|
||||
/** A never-quarantined profile must not carry {@code quarantineAttempt} either. */
|
||||
@Test
|
||||
void profilesOmitsQuarantineAttemptWhenNotQuarantined() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
McpSchema.CallToolResult res = FleetMcp.profiles(
|
||||
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), FleetMcp.QuarantineSource.none());
|
||||
String out = textOf(res);
|
||||
assertFalse(out.contains("quarantineAttempt"), out);
|
||||
}
|
||||
|
||||
/** fleetd #201 Unit 5: {@code coolingOff} is a SEPARATE map from {@code quarantined}. */
|
||||
@Test
|
||||
void profilesReportsACoolingOffCredentialInASeparateMap() {
|
||||
@@ -626,8 +658,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 +687,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 +708,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 +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("\"msgId\":\"m1\""), out);
|
||||
@@ -723,8 +867,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 +896,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 +964,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);
|
||||
@@ -1008,6 +1153,35 @@ class FleetMcpTest {
|
||||
assertTrue(out.contains("\"quarantinedForSeconds\":1200"), out);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #466 scope item 2: {@code fleet_list}'s capacity rows (CB-583: they reuse quarantine)
|
||||
* must carry {@code quarantineAttempt} beside {@code quarantinedForSeconds}, off the same
|
||||
* {@link BackendQuarantine#status} call as {@code fleet_profiles} -- see {@code capacityView}'s
|
||||
* comment pointing back to {@code profilesView}.
|
||||
*/
|
||||
@Test
|
||||
void capacityRowReportsQuarantineAttemptBesideRemainingSeconds() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
|
||||
java.util.concurrent.atomic.AtomicLong now = new java.util.concurrent.atomic.AtomicLong(0L);
|
||||
BackendQuarantine quarantine = BackendQuarantine.withEscalation(now::get, TimeUnit.MINUTES.toNanos(20));
|
||||
quarantine.quarantine("shared-openai"); // attempt 1
|
||||
now.set(TimeUnit.MINUTES.toNanos(20));
|
||||
quarantine.quarantine("shared-openai"); // attempt 2
|
||||
now.set(TimeUnit.MINUTES.toNanos(60));
|
||||
quarantine.quarantine("shared-openai"); // attempt 3
|
||||
FleetMcp.QuarantineSource source = new FleetMcp.QuarantineSource(
|
||||
profile -> "terra".equals(profile) ? "shared-openai" : null, quarantine);
|
||||
|
||||
String out = textOf(FleetMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")),
|
||||
sessions, null, new FleetMcp.CapacitySource(profile -> 0, profile -> 2,
|
||||
() -> Set.of("terra"), () -> 0), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
source, Map.of(), ""));
|
||||
|
||||
assertTrue(out.contains("\"credentialId\":\"shared-openai\""), out);
|
||||
assertTrue(out.contains("\"quarantineAttempt\":3"), out);
|
||||
}
|
||||
|
||||
@Test
|
||||
void everyProfileSharingTheQuarantinedCredentialReportsZeroFree() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
|
||||
@@ -116,4 +116,262 @@ class BackendQuarantineTest {
|
||||
assertThrows(IllegalArgumentException.class, () -> new BackendQuarantine(() -> 0L, 0L));
|
||||
assertThrows(IllegalArgumentException.class, () -> new BackendQuarantine(() -> 0L, -1L));
|
||||
}
|
||||
|
||||
// --- fleetd #466: escalating cooldown -----------------------------------------------------
|
||||
//
|
||||
// Base cooldown 600s (10 min), multiplier 2.0, ceiling 2400s (4x base) — small round numbers
|
||||
// chosen so every deadline is an exact assertion, not just "greater than before". Each call
|
||||
// below lands at or before the previous deadline (a zero or negative gap), which is always
|
||||
// "no more than one base cooldown after the previous deadline" — i.e. every call continues the
|
||||
// same streak, matching a credential that keeps reporting exhausted with no lull.
|
||||
|
||||
@Test
|
||||
void anInvalidBackoffMultiplierIsRejected() {
|
||||
assertThrows(IllegalArgumentException.class,
|
||||
() -> new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30), 0.5, TimeUnit.HOURS.toNanos(6)));
|
||||
}
|
||||
|
||||
@Test
|
||||
void aCeilingBelowTheBaseCooldownIsRejected() {
|
||||
assertThrows(IllegalArgumentException.class,
|
||||
() -> new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30), 2.0, TimeUnit.MINUTES.toNanos(10)));
|
||||
}
|
||||
|
||||
@Test
|
||||
void repeatedExhaustionEscalatesTheCooldownByExactAmounts() {
|
||||
AtomicLong now = new AtomicLong(0L);
|
||||
BackendQuarantine q = new BackendQuarantine(now::get, TimeUnit.SECONDS.toNanos(600), 2.0,
|
||||
TimeUnit.SECONDS.toNanos(2400));
|
||||
|
||||
q.quarantine("shared-openai"); // 1st: base cooldown
|
||||
assertEquals(OptionalLong.of(600L), q.remainingSeconds("shared-openai"));
|
||||
|
||||
now.set(TimeUnit.SECONDS.toNanos(100)); // still inside the 1st quarantine (deadline 600s)
|
||||
q.quarantine("shared-openai"); // 2nd: 600 * 2^1 = 1200
|
||||
assertEquals(OptionalLong.of(1200L), q.remainingSeconds("shared-openai"),
|
||||
"a second consecutive exhaustion must double the cooldown, not just increase it");
|
||||
|
||||
now.set(TimeUnit.SECONDS.toNanos(1300)); // exactly the 2nd deadline (100 + 1200)
|
||||
q.quarantine("shared-openai"); // 3rd: 600 * 2^2 = 2400 (exactly at the ceiling)
|
||||
assertEquals(OptionalLong.of(2400L), q.remainingSeconds("shared-openai"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void escalationStopsAtTheCeiling() {
|
||||
AtomicLong now = new AtomicLong(0L);
|
||||
BackendQuarantine q = new BackendQuarantine(now::get, TimeUnit.SECONDS.toNanos(600), 2.0,
|
||||
TimeUnit.SECONDS.toNanos(2400));
|
||||
|
||||
q.quarantine("shared-openai"); // 1st: 600
|
||||
now.set(TimeUnit.SECONDS.toNanos(600));
|
||||
q.quarantine("shared-openai"); // 2nd: 1200, deadline 1800
|
||||
now.set(TimeUnit.SECONDS.toNanos(1800));
|
||||
q.quarantine("shared-openai"); // 3rd: 600 * 4 = 2400, at the ceiling, deadline 4200
|
||||
now.set(TimeUnit.SECONDS.toNanos(4200));
|
||||
q.quarantine("shared-openai"); // 4th: 600 * 8 = 4800 uncapped, must stay capped at 2400
|
||||
assertEquals(OptionalLong.of(2400L), q.remainingSeconds("shared-openai"),
|
||||
"the cooldown must never exceed the configured ceiling, however long the streak gets");
|
||||
|
||||
now.set(TimeUnit.SECONDS.toNanos(6600)); // 4th deadline
|
||||
q.quarantine("shared-openai"); // 5th: still capped
|
||||
assertEquals(OptionalLong.of(2400L), q.remainingSeconds("shared-openai"),
|
||||
"pushing well past the ceiling must not budge it");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aQuietGapLongerThanTheBaseCooldownResetsToTheBaseCooldown() {
|
||||
AtomicLong now = new AtomicLong(0L);
|
||||
BackendQuarantine q = new BackendQuarantine(now::get, TimeUnit.SECONDS.toNanos(600), 2.0,
|
||||
TimeUnit.SECONDS.toNanos(2400));
|
||||
|
||||
q.quarantine("shared-openai"); // 1st: 600, deadline 600
|
||||
now.set(TimeUnit.SECONDS.toNanos(600));
|
||||
q.quarantine("shared-openai"); // 2nd: 1200, deadline 1800
|
||||
now.set(TimeUnit.SECONDS.toNanos(1800));
|
||||
q.quarantine("shared-openai"); // 3rd: 2400, deadline 4200
|
||||
assertEquals(OptionalLong.of(2400L), q.remainingSeconds("shared-openai"));
|
||||
|
||||
// Quiet for well over one base cooldown (600s) past the 3rd deadline (4200s).
|
||||
now.set(TimeUnit.SECONDS.toNanos(20_000));
|
||||
q.quarantine("shared-openai"); // treated as a fresh occurrence
|
||||
assertEquals(OptionalLong.of(600L), q.remainingSeconds("shared-openai"),
|
||||
"a long quiet gap must reset the streak back to the base cooldown");
|
||||
}
|
||||
|
||||
@Test
|
||||
void escalatingOneCredentialDoesNotSlowAnother() {
|
||||
AtomicLong now = new AtomicLong(0L);
|
||||
BackendQuarantine q = new BackendQuarantine(now::get, TimeUnit.SECONDS.toNanos(600), 2.0,
|
||||
TimeUnit.SECONDS.toNanos(2400));
|
||||
|
||||
q.quarantine("shared-openai"); // 1st: 600
|
||||
now.set(TimeUnit.SECONDS.toNanos(600));
|
||||
q.quarantine("shared-openai"); // 2nd: 1200
|
||||
now.set(TimeUnit.SECONDS.toNanos(1800));
|
||||
q.quarantine("shared-openai"); // 3rd: 2400 — three-in-a-row streak on this credential only
|
||||
|
||||
q.quarantine("another-credential"); // its first and only exhaustion
|
||||
assertEquals(OptionalLong.of(600L), q.remainingSeconds("another-credential"),
|
||||
"an unrelated credential's cooldown must stay at the base rate, unaffected by a sibling's streak");
|
||||
}
|
||||
|
||||
@Test
|
||||
void withEscalationDefaultsToDoublingCappedAtTwelveTimesTheBase() {
|
||||
AtomicLong now = new AtomicLong(0L);
|
||||
BackendQuarantine q = BackendQuarantine.withEscalation(now::get, TimeUnit.MINUTES.toNanos(30));
|
||||
|
||||
q.quarantine("shared-openai");
|
||||
assertEquals(OptionalLong.of(1800L), q.remainingSeconds("shared-openai"),
|
||||
"the first occurrence must still use the base cooldown");
|
||||
|
||||
now.set(TimeUnit.MINUTES.toNanos(30));
|
||||
q.quarantine("shared-openai");
|
||||
assertEquals(OptionalLong.of(3600L), q.remainingSeconds("shared-openai"),
|
||||
"the default multiplier must be 2.0");
|
||||
}
|
||||
|
||||
// --- fleetd #466 scope item 2: status() reports repeatCount beside remainingSeconds ------------
|
||||
|
||||
@Test
|
||||
void statusReportsAttemptOneForAFirstOccurrenceNeverAbsentOrZero() {
|
||||
BackendQuarantine q = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30), 2.0,
|
||||
TimeUnit.HOURS.toNanos(6));
|
||||
q.quarantine("shared-openai");
|
||||
|
||||
BackendQuarantine.Status status = q.status("shared-openai").orElseThrow();
|
||||
assertEquals(1800L, status.remainingSeconds());
|
||||
assertEquals(1, status.repeatCount(),
|
||||
"a first-ever occurrence must report attempt 1, not 0 or absent -- 1 means unambiguously "
|
||||
+ "'the first time', where 0 would be indistinguishable from a bug that forgot to count");
|
||||
}
|
||||
|
||||
@Test
|
||||
void statusIsAbsentWhenTheCredentialIsNotQuarantined() {
|
||||
BackendQuarantine q = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30), 2.0,
|
||||
TimeUnit.HOURS.toNanos(6));
|
||||
assertTrue(q.status("shared-openai").isEmpty());
|
||||
}
|
||||
|
||||
/**
|
||||
* The acceptance criterion's strong form: the reported {@code repeatCount} must match the exact
|
||||
* step the cooldown's own growth implies, read off {@link BackendQuarantine#status}'s single
|
||||
* call -- not two independent reads that happen to agree in this easy case.
|
||||
*/
|
||||
@Test
|
||||
void statusReportsTheGrowingAttemptCountAlongsideTheEscalatingCooldown() {
|
||||
AtomicLong now = new AtomicLong(0L);
|
||||
BackendQuarantine q = new BackendQuarantine(now::get, TimeUnit.SECONDS.toNanos(600), 2.0,
|
||||
TimeUnit.SECONDS.toNanos(2400));
|
||||
|
||||
q.quarantine("shared-openai"); // 1st: 600s, attempt 1
|
||||
BackendQuarantine.Status first = q.status("shared-openai").orElseThrow();
|
||||
assertEquals(600L, first.remainingSeconds());
|
||||
assertEquals(1, first.repeatCount());
|
||||
|
||||
now.set(TimeUnit.SECONDS.toNanos(600));
|
||||
q.quarantine("shared-openai"); // 2nd: 1200s, attempt 2
|
||||
BackendQuarantine.Status second = q.status("shared-openai").orElseThrow();
|
||||
assertEquals(1200L, second.remainingSeconds());
|
||||
assertEquals(2, second.repeatCount());
|
||||
|
||||
now.set(TimeUnit.SECONDS.toNanos(1800));
|
||||
q.quarantine("shared-openai"); // 3rd: 2400s (at the ceiling), attempt 3
|
||||
BackendQuarantine.Status third = q.status("shared-openai").orElseThrow();
|
||||
assertEquals(2400L, third.remainingSeconds());
|
||||
assertEquals(3, third.repeatCount());
|
||||
}
|
||||
|
||||
/**
|
||||
* Once the cooldown hits its ceiling, every further consecutive exhaustion reports the SAME
|
||||
* {@code remainingSeconds} -- so a {@code repeatCount} re-derived from the cooldown value (e.g.
|
||||
* inverting {@code cooldownNanos * multiplier^(n-1)}) could not tell attempt 4 from attempt 9;
|
||||
* only the stored counter can. This is the scenario that makes "read the count off a second,
|
||||
* independent computation" provably wrong rather than just risky.
|
||||
*/
|
||||
@Test
|
||||
void repeatCountKeepsGrowingPastTheCeilingEvenThoughTheCooldownStaysFlat() {
|
||||
AtomicLong now = new AtomicLong(0L);
|
||||
BackendQuarantine q = new BackendQuarantine(now::get, TimeUnit.SECONDS.toNanos(600), 2.0,
|
||||
TimeUnit.SECONDS.toNanos(2400));
|
||||
|
||||
long t = 0L;
|
||||
for (int attempt = 1; attempt <= 5; attempt++) {
|
||||
now.set(t);
|
||||
q.quarantine("shared-openai");
|
||||
BackendQuarantine.Status status = q.status("shared-openai").orElseThrow();
|
||||
assertEquals(attempt, status.repeatCount(),
|
||||
"attempt " + attempt + " must be reported as exactly " + attempt
|
||||
+ ", not collapsed to whatever attempt first reached the ceiling");
|
||||
if (attempt >= 3) {
|
||||
assertEquals(2400L, status.remainingSeconds(), "attempt " + attempt + " must be capped");
|
||||
}
|
||||
t += status.remainingSeconds(); // land exactly on the next deadline: still the same streak
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void aQuietGapResetsTheReportedAttemptCountToOneToo() {
|
||||
AtomicLong now = new AtomicLong(0L);
|
||||
BackendQuarantine q = new BackendQuarantine(now::get, TimeUnit.SECONDS.toNanos(600), 2.0,
|
||||
TimeUnit.SECONDS.toNanos(2400));
|
||||
|
||||
q.quarantine("shared-openai");
|
||||
now.set(TimeUnit.SECONDS.toNanos(600));
|
||||
q.quarantine("shared-openai");
|
||||
now.set(TimeUnit.SECONDS.toNanos(1800));
|
||||
q.quarantine("shared-openai"); // attempt 3
|
||||
assertEquals(3, q.status("shared-openai").orElseThrow().repeatCount());
|
||||
|
||||
now.set(TimeUnit.SECONDS.toNanos(20_000)); // long quiet gap
|
||||
q.quarantine("shared-openai");
|
||||
assertEquals(1, q.status("shared-openai").orElseThrow().repeatCount(),
|
||||
"a reset streak must report attempt 1 again, matching the reset base cooldown");
|
||||
}
|
||||
|
||||
@Test
|
||||
void escalatingOneCredentialsAttemptCountDoesNotAffectAnother() {
|
||||
AtomicLong now = new AtomicLong(0L);
|
||||
BackendQuarantine q = new BackendQuarantine(now::get, TimeUnit.SECONDS.toNanos(600), 2.0,
|
||||
TimeUnit.SECONDS.toNanos(2400));
|
||||
|
||||
q.quarantine("shared-openai");
|
||||
now.set(TimeUnit.SECONDS.toNanos(600));
|
||||
q.quarantine("shared-openai");
|
||||
now.set(TimeUnit.SECONDS.toNanos(1800));
|
||||
q.quarantine("shared-openai"); // 3-in-a-row streak on this credential only
|
||||
|
||||
q.quarantine("another-credential");
|
||||
assertEquals(1, q.status("another-credential").orElseThrow().repeatCount(),
|
||||
"an unrelated credential's attempt count must stay at 1, unaffected by a sibling's streak");
|
||||
}
|
||||
|
||||
@Test
|
||||
void noneReportsNoStatusForAnything() {
|
||||
BackendQuarantine q = BackendQuarantine.none();
|
||||
q.quarantine("shared-openai"); // no-op on none(), same as every other mutator
|
||||
|
||||
assertTrue(q.status("shared-openai").isEmpty(),
|
||||
"none() quarantines nothing, so it must report no status at all -- never a fabricated "
|
||||
+ "attempt count for a credential that was never actually quarantined");
|
||||
}
|
||||
|
||||
/** The flat (non-escalating) two-argument constructor must still report a real, growing count. */
|
||||
@Test
|
||||
void aFlatTwoArgumentInstanceStillReportsAGrowingAttemptCountEvenThoughTheCooldownStaysFlat() {
|
||||
AtomicLong now = new AtomicLong(0L);
|
||||
BackendQuarantine q = new BackendQuarantine(now::get, TimeUnit.SECONDS.toNanos(600));
|
||||
|
||||
q.quarantine("shared-openai");
|
||||
BackendQuarantine.Status first = q.status("shared-openai").orElseThrow();
|
||||
assertEquals(600L, first.remainingSeconds());
|
||||
assertEquals(1, first.repeatCount());
|
||||
|
||||
now.set(TimeUnit.SECONDS.toNanos(100)); // still inside the 1st quarantine
|
||||
q.quarantine("shared-openai");
|
||||
BackendQuarantine.Status second = q.status("shared-openai").orElseThrow();
|
||||
assertEquals(600L, second.remainingSeconds(),
|
||||
"the flat constructor's cooldown must stay exactly the base length regardless of the streak");
|
||||
assertEquals(2, second.repeatCount(),
|
||||
"the flat constructor still counts the real streak -- it just does not scale the cooldown by it");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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