Compare commits
6 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| c50f5b2d61 | |||
| c796eac09c | |||
| 9d37f3aa29 | |||
| eaf89abaf6 | |||
| d895f02bc1 | |||
| 43206cac2f |
@@ -306,6 +306,35 @@ profiles:
|
||||
# GOTCHA 2 — `maxLoad` is the ONLY throttle you have here. There is no metering, no budget
|
||||
# and no refusal on cost; the cap on live members is the single thing standing between a
|
||||
# fan-out and your monthly limit. Set it deliberately and keep it small.
|
||||
#
|
||||
# GOTCHA 3 (fleetd #176) — `maxLoad` counts members, never the lead itself. The lead is a live
|
||||
# `claude` session on this SAME account (a lead is never moved off-subscription, whatever its
|
||||
# own profile says), so it already holds one seat before any member spawns. If a lead's
|
||||
# `fleet.leaders.<name>.profile` names THIS profile — or ANY OTHER `subscription: true`
|
||||
# profile that shares this one's account (see THE SENTINEL, just below, next to
|
||||
# `credentialId:`) — `fleet_list`'s `free` for this profile subtracts that lead's live
|
||||
# seat(s) automatically; see `profile:` under THE FLEET below. If no lead entry names a
|
||||
# profile sharing this account, fleetd has no way to know a lead holds a seat here, and `free`
|
||||
# will overstate what a fresh `fleet_spawn` actually gets by exactly the seats the lead is
|
||||
# quietly holding.
|
||||
#
|
||||
# THE SENTINEL (fleetd #176 stage 2, correcting an inert stage 1 fix): every `subscription:
|
||||
# true` profile that leaves `credentialId` unset shares ONE implicit account-wide credential
|
||||
# id with every other such profile on this host — because a subscription profile doesn't
|
||||
# authenticate with a credential of its own, it authenticates as the operator's own Claude
|
||||
# login, and there is exactly one of those. So on a typical host, `opus` (the lead's profile)
|
||||
# and `sonnet` (the members' profile) are linked automatically, with NOTHING to set here — that
|
||||
# is what makes GOTCHA 3 above work without also writing matching `credentialId:` values on
|
||||
# both. This linkage is not just cosmetic: it is the same key `BackendQuarantine`/cool-off use,
|
||||
# so a usage-limit hit on `opus` now quarantines `sonnet` too (and vice versa) — correct, since
|
||||
# they are one Claude account, but worth knowing before you wonder why an unrelated-looking
|
||||
# profile went quarantined.
|
||||
#
|
||||
# WHEN TO OVERRIDE — set explicit, DIFFERENT `credentialId:` values on two `subscription: true`
|
||||
# profiles only when they are genuinely two separate Claude logins on the same host (a real,
|
||||
# if unusual, setup). An explicit `credentialId` always wins over the sentinel, so this is the
|
||||
# one way to keep two subscription profiles from being treated as one account for lead-seat
|
||||
# counting AND for quarantine/cool-off grouping alike.
|
||||
# gitTokenEnv: GITEA_TOKEN # opt-in: let this profile's workers open their own PR (CB-302)
|
||||
# gitHostEnv: GITEA_HOST # defaults to GITEA_HOST; injected only with gitTokenEnv
|
||||
# exhaustedPattern: "usage limit has been reached" # opt-in: classify a usage-limit refusal (CB-578)
|
||||
@@ -503,6 +532,22 @@ fleet:
|
||||
# recognised: give it a `profile:` and the daemon launches the shortfall when fewer than
|
||||
# `instances` are live. Omit `profile:` and it is recognise-only, as before.
|
||||
#
|
||||
# `profile:` has a SECOND job as of fleetd #176, even for a recognise-only lead you never want
|
||||
# auto-launched: it is also how fleetd learns which account this lead's own session shares. A
|
||||
# `subscription: true` profile bills the operator's Claude account, and the lead itself is always
|
||||
# a live `claude` session on that same account — `maxLoad` never counted that seat. If a lead
|
||||
# entry here names a profile that shares a worker profile's account, `fleet_list`'s `free` for
|
||||
# that worker profile subtracts the lead's live seat(s) automatically. "Shares the account" is
|
||||
# decided by matching `effectiveCredentialId()`, which (fleetd #176 stage 2 — see THE SENTINEL,
|
||||
# next to `credentialId:`, in THE WORKERS above) means: an explicit, matching `credentialId:` on
|
||||
# both, OR — the common case, needing NO extra config — both being `subscription: true` with
|
||||
# `credentialId` left unset, since those all share one implicit account-wide id. A lead on `opus`
|
||||
# and workers on `sonnet` link automatically this way; they do NOT need the same profile name.
|
||||
# Setting `profile:` on an already-running, recognise-only lead is safe — the daemon only launches
|
||||
# the SHORTFALL below `instances`, so naming a profile here does not, by itself, start anything.
|
||||
# Omit it and fleetd has no way to derive the sharing — there is no other reliable signal on the
|
||||
# daemon's side — so that lead's seat goes uncounted, exactly as before this ticket.
|
||||
#
|
||||
# `tab:` (CB-579) is REQUIRED and is the only field identity depends on — the exact label of the
|
||||
# tab hosting the lead, matched case-insensitively. Label the tab yourself and put that same
|
||||
# string here, and the pane is recognised on the next rescan. Reopen the tab later, or the session
|
||||
|
||||
@@ -351,7 +351,8 @@ public final class Fleetd {
|
||||
.map(session -> exhaustedPatternsByProfile.get(session.profile()))
|
||||
.orElse(null);
|
||||
log.info("backend-exhausted classification (CB-578 stage A): {}",
|
||||
CompletionResolver.coverage(cfg.profiles().keySet(), exhaustedPatternsByProfile.keySet()));
|
||||
CompletionResolver.coverage("exhaustedPattern", cfg.profiles().keySet(),
|
||||
exhaustedPatternsByProfile.keySet()));
|
||||
// fleetd #201 Unit 5: classify a completion-fallback scrape that matches a profile's
|
||||
// configured backend-error refusal (a credential outage, a provider 5xx) as a backend error
|
||||
// rather than handing it back as a real answer. Compiled once at startup, keyed by profile
|
||||
@@ -365,13 +366,14 @@ public final class Fleetd {
|
||||
errorPatternsByProfile.put(name, Pattern.compile(profile.errorPattern()));
|
||||
}
|
||||
});
|
||||
BackendErrorPatternLookup backendErrorPatterns = target -> sessions.roster().stream()
|
||||
.filter(session -> target.equals(session.terminalId()))
|
||||
.findFirst()
|
||||
.map(session -> errorPatternsByProfile.get(session.profile()))
|
||||
.orElse(null);
|
||||
// fleetd #248: extracted to a static factory (see backendErrorPatternLookup below) so a
|
||||
// test can prove main() actually PASSES this into CompletionResolver, not only that the
|
||||
// lookup itself behaves correctly — the exact gap fleetd #248 exists to close.
|
||||
BackendErrorPatternLookup backendErrorPatterns = backendErrorPatternLookup(sessions::roster,
|
||||
errorPatternsByProfile);
|
||||
log.info("backend-error classification (fleetd #201 Unit 5): {}",
|
||||
CompletionResolver.coverage(cfg.profiles().keySet(), errorPatternsByProfile.keySet()));
|
||||
CompletionResolver.coverage("errorPattern", cfg.profiles().keySet(),
|
||||
errorPatternsByProfile.keySet()));
|
||||
// CB-578 stage B: on a classification that actually wins, quarantine the exhausted profile's
|
||||
// CREDENTIAL — not the profile name — so a profile sharing that credential (e.g. two models
|
||||
// on one OpenAI account) is refused too, not just the one that happened to report it. Reads
|
||||
@@ -424,59 +426,25 @@ public final class Fleetd {
|
||||
// construction-order cycle `exhaustionSinkRef` breaks above, broken the same way: a mutable
|
||||
// holder set once `pushLoop` exists, read lazily from inside the lambda built here.
|
||||
AtomicReference<ReplyPushLoop> pushLoopRef = new AtomicReference<>();
|
||||
// Order: (1) mark the member BACKEND_ERROR; (2) resolve profile/credential through the
|
||||
// roster — fail loud (never Optional.ifPresent, the fleetd #234 lesson applied to this new
|
||||
// sink) and notify the lead via onBackendTargetUnmapped when it cannot be resolved; (3)
|
||||
// record the error in BackendOutagePolicy; (4) on a NEW incident (the record() call that
|
||||
// actually crosses the threshold), tell the lead via onBackendIncident.
|
||||
BackendErrorSink backendErrorSink = (target, matchedLine, reason) -> {
|
||||
sessions.onBackendError(target, reason);
|
||||
|
||||
String profileName = sessions.roster().stream()
|
||||
.filter(session -> target.equals(session.terminalId()))
|
||||
.findFirst()
|
||||
.map(MemberSession::profile)
|
||||
.orElse(null);
|
||||
FleetConfig.Profile profile = profileName == null ? null : config.get().profiles().get(profileName);
|
||||
if (profile == null) {
|
||||
log.error("backend error on target '{}' ({}) but no profile could be resolved — the "
|
||||
+ "target is not (yet) in the roster — no cool-off applied (fleetd #201 Unit 5)",
|
||||
target, reason);
|
||||
ReplyPushLoop loop = pushLoopRef.get();
|
||||
if (loop != null) {
|
||||
loop.onBackendTargetUnmapped(target, reason);
|
||||
}
|
||||
return;
|
||||
}
|
||||
String credentialId = profile.effectiveCredentialId();
|
||||
Optional<BackendOutagePolicy.Incident> incident = outagePolicy.record(credentialId, target, reason);
|
||||
incident.ifPresent(inc -> {
|
||||
List<String> affectedProfiles = config.get().profiles().values().stream()
|
||||
.filter(p -> credentialId.equals(p.effectiveCredentialId()))
|
||||
.map(FleetConfig.Profile::profile)
|
||||
.sorted()
|
||||
.toList();
|
||||
log.warn("credential '{}' cooling off for {}s after backend errors on {} distinct "
|
||||
+ "target(s) (profile '{}'): {}", credentialId,
|
||||
inc.remainingCoolOffSeconds(), inc.evidenceCount(), profile.profile(), reason);
|
||||
ReplyPushLoop loop = pushLoopRef.get();
|
||||
if (loop != null) {
|
||||
loop.onBackendIncident(inc.id(), inc.targets(), credentialId, affectedProfiles,
|
||||
(int) inc.remainingCoolOffSeconds());
|
||||
}
|
||||
});
|
||||
};
|
||||
// fleetd #248: extracted to a static factory (see backendErrorSink below), public rather
|
||||
// than package-private like the other two factories here, so
|
||||
// dev.ltms.fleet.inject.BackendOutageFlowTest can exercise the REAL production sink
|
||||
// directly instead of a hand-mirrored copy of this lambda — that copy was precisely the
|
||||
// gap fleetd #248 exists to close (see that test's class doc for the history).
|
||||
BackendErrorSink backendErrorSink = backendErrorSink(sessions, () -> config.get().profiles(),
|
||||
outagePolicy, pushLoopRef::get);
|
||||
AgentControl agents = router.memberAgents();
|
||||
// Both fleetd#201 Unit 5 (backend-error patterns + sink) and fleetd#241 (the worktree/branch
|
||||
// lookup the fallback report names) land on this one call. The full constructor takes both,
|
||||
// so neither feature is dropped; nowNanos must be passed explicitly to reach it.
|
||||
//
|
||||
// fleetd #248: every argument built specifically for this call (backendErrorPatterns and
|
||||
// backendErrorSink above, and the worktree/branch lookup right here) now comes from a
|
||||
// static factory tested on its own; FleetdCompletionResolverWiringTest source-asserts that
|
||||
// THIS call actually passes them, which is the coverage that was missing before.
|
||||
CompletionResolver completion = new CompletionResolver(agents, rendezvous, exhaustedPatterns,
|
||||
exhaustionSink, backendErrorPatterns, backendErrorSink, System::nanoTime,
|
||||
target -> sessions.roster().stream()
|
||||
.filter(session -> target.equals(session.terminalId()))
|
||||
.findFirst()
|
||||
.map(session -> new CompletionResolver.WorktreeBranch(session.worktree(), session.branch()))
|
||||
.orElse(null));
|
||||
worktreeBranchLookup(sessions::roster));
|
||||
// CB-113: deliver only to an available worker (its MCP is connected), never its boot window.
|
||||
// CB-301: the manager's presence bridge records availability and drives SPAWNING → READY.
|
||||
MemberPresence presence = sessions.asPresence();
|
||||
@@ -672,7 +640,8 @@ public final class Fleetd {
|
||||
new FleetMcp.OutageSource(profile -> {
|
||||
var configured = config.get().profiles().get(profile);
|
||||
return configured == null ? null : configured.effectiveCredentialId();
|
||||
}, outagePolicy));
|
||||
}, outagePolicy),
|
||||
new FleetMcp.LeadSeatSource(leadSeatLookup(() -> config.get().profiles(), leaders, leads)));
|
||||
|
||||
// CB-637: the receive half. Only constructed when a lead mailbox actually opened — with no
|
||||
// coordinator (or an unreachable one) there is nothing to deliver, so no scheduler is
|
||||
@@ -770,6 +739,185 @@ public final class Fleetd {
|
||||
return target -> presence.isPresent(target) || leads.get().containsKey(target);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #248: package-private factory for the member worktree/branch lookup {@link
|
||||
* CompletionResolver} uses to name a fallback report's worktree and branch (fleetd#241).
|
||||
*
|
||||
* <p>Before this ticket the lookup was an anonymous lambda built inline inside {@code main}'s
|
||||
* {@code CompletionResolver} constructor call — provably untested wiring, the whole reason
|
||||
* fleetd #248 exists: dropping that one argument (passing {@code _ -> null} instead) compiled
|
||||
* clean and left every test green. Extracted here, {@code main} now calls this factory instead
|
||||
* of building the lambda inline, and a source assertion on that call site
|
||||
* ({@code FleetdCompletionResolverWiringTest}) proves the argument is still actually passed.
|
||||
*
|
||||
* <p>Takes the roster as a plain {@link Supplier} — not a {@link SessionManager} — so this is
|
||||
* directly testable with a hand-built session list; no real {@code SessionManager} (launcher,
|
||||
* worktrees, …) needs constructing. Follows the same {@code static} factory pattern as
|
||||
* {@link #deliverableTo} above.
|
||||
*
|
||||
* @param roster the live member roster, normally {@code sessions::roster}
|
||||
*/
|
||||
static Function<String, CompletionResolver.WorktreeBranch> worktreeBranchLookup(
|
||||
Supplier<List<MemberSession>> roster) {
|
||||
return target -> roster.get().stream()
|
||||
.filter(session -> target.equals(session.terminalId()))
|
||||
.findFirst()
|
||||
.map(session -> new CompletionResolver.WorktreeBranch(session.worktree(), session.branch()))
|
||||
.orElse(null);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #176: per-profile factory for {@link FleetMcp.LeadSeatSource} — how many seats a
|
||||
* profile's own live LEAD session(s) hold on the same Claude subscription.
|
||||
*
|
||||
* <p>{@code maxLoad} counts only members; the lead itself is a live {@code claude} session that
|
||||
* is never moved off-subscription ({@code LeadLauncher} strips {@code ANTHROPIC_BASE_URL}/
|
||||
* {@code AUTH_TOKEN} from a lead's env whatever its profile says), so a {@code subscription:
|
||||
* true} profile's real ceiling is lower than its configured {@code maxLoad} by exactly the
|
||||
* number of lead seats sharing that same account.
|
||||
*
|
||||
* <p><b>The derivation, and why this route was chosen over a new config key.</b> The link is
|
||||
* {@code fleet.leaders.<name>.profile} — the field the operator already sets to name which
|
||||
* {@code profiles:} entry a lead runs on (CB-557; see {@code fleetd.example.yaml}) — matched
|
||||
* against the profile passed in here via {@link FleetConfig.Profile#effectiveCredentialId()},
|
||||
* the same grouping key {@link dev.ltms.fleet.placement.BackendQuarantine} already uses to say
|
||||
* two profiles share one account. Nothing new is added to the config schema: this reuses a field
|
||||
* that already exists and already means "the profile this lead's own session runs on". A lead
|
||||
* entry that names no {@code profile:} (recognise-only, CB-558) says nothing about which account
|
||||
* it shares, and there is no other reliable signal on the daemon's side to derive that from — so
|
||||
* such a lead contributes no seats, exactly as before this ticket. Making that lead's seat count
|
||||
* requires the operator to add one line (`profile: sonnet` under its {@code fleet.leaders} entry)
|
||||
* — a config statement, not a code change, and the smallest one available given the field
|
||||
* already exists for a closely related purpose.
|
||||
*
|
||||
* <p>Only counts leads {@code liveLeadTerminals} currently reports — CB-531's live tab scan (or
|
||||
* the legacy {@code primary.terminal} pin) — never every configured lead: an entry whose
|
||||
* {@code instances} nobody has actually started is not really competing for a seat, and must not
|
||||
* shrink capacity for one that was never live.
|
||||
*
|
||||
* @param profiles the live profile map, normally {@code () -> config.get().profiles()}
|
||||
* in {@code main} — hot, like every other {@code maxLoad}/
|
||||
* {@code credentialId} read {@link FleetMcp.CapacitySource} already does
|
||||
* @param leaders {@code fleet.leaders}, read once at startup like the rest of that
|
||||
* block ({@code fleetd.example.yaml} notes it is not hot) — passed as a
|
||||
* plain map, never re-read from {@code config.get()}
|
||||
* @param liveLeadTerminals terminal_id → lead name for every CURRENTLY recognised lead, normally
|
||||
* the same supplier {@link dev.ltms.fleet.auth.CallerResolver#leads()}
|
||||
* and {@code LeadCoordLoop} already consult
|
||||
*/
|
||||
static Function<String, Integer> leadSeatLookup(Supplier<Map<String, FleetConfig.Profile>> profiles,
|
||||
Map<String, FleetConfig.Leader> leaders, Supplier<Map<String, String>> liveLeadTerminals) {
|
||||
return profileName -> {
|
||||
FleetConfig.Profile target = profiles.get().get(profileName);
|
||||
if (target == null || !target.isSubscription()) {
|
||||
return 0;
|
||||
}
|
||||
String targetCredential = target.effectiveCredentialId();
|
||||
int seats = 0;
|
||||
for (String leadName : liveLeadTerminals.get().values()) {
|
||||
FleetConfig.Leader lead = leaders.get(leadName);
|
||||
if (lead == null || lead.profile() == null || lead.profile().isBlank()) {
|
||||
continue;
|
||||
}
|
||||
FleetConfig.Profile leadProfile = profiles.get().get(lead.profile());
|
||||
if (leadProfile != null && targetCredential.equals(leadProfile.effectiveCredentialId())) {
|
||||
seats++;
|
||||
}
|
||||
}
|
||||
return seats;
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #248 / fleetd#201 Unit 5: package-private factory for the per-target backend-error
|
||||
* pattern lookup {@link CompletionResolver} classifies a pane scrape against. Closes over the
|
||||
* live roster (to resolve a target to a profile) and {@code errorPatternsByProfile} (each
|
||||
* profile's configured {@code errorPattern}, already compiled by the caller — the same map also
|
||||
* feeds the coverage log next to where this is called) — nothing else, so it is directly
|
||||
* testable. See {@link #worktreeBranchLookup} above for why this ticket exists and why the
|
||||
* factory takes a roster {@link Supplier} rather than a {@link SessionManager}.
|
||||
*
|
||||
* @param roster the live member roster, normally {@code sessions::roster}
|
||||
* @param errorPatternsByProfile every profile that has an {@code errorPattern} configured,
|
||||
* keyed by profile name
|
||||
*/
|
||||
static BackendErrorPatternLookup backendErrorPatternLookup(Supplier<List<MemberSession>> roster,
|
||||
Map<String, Pattern> errorPatternsByProfile) {
|
||||
return target -> roster.get().stream()
|
||||
.filter(session -> target.equals(session.terminalId()))
|
||||
.findFirst()
|
||||
.map(session -> errorPatternsByProfile.get(session.profile()))
|
||||
.orElse(null);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #248 / fleetd#201 Unit 5: factory for the production {@link BackendErrorSink} — the
|
||||
* collaborator {@link CompletionResolver} notifies when a pane-scrape classification actually
|
||||
* resolves a waiter as a backend error. Order: (1) mark the member BACKEND_ERROR; (2) resolve
|
||||
* profile/credential through the roster — fail loud (never {@code Optional.ifPresent}, the
|
||||
* fleetd #234 lesson applied to this sink) and notify the lead via {@code
|
||||
* onBackendTargetUnmapped} when it cannot be resolved; (3) record the error in {@code
|
||||
* outagePolicy}; (4) on a NEW incident (the {@code record()} call that actually crosses the
|
||||
* threshold), tell the lead via {@code onBackendIncident}.
|
||||
*
|
||||
* <p>{@code public}, unlike {@link #worktreeBranchLookup} and {@link #backendErrorPatternLookup}
|
||||
* above: {@code dev.ltms.fleet.inject.BackendOutageFlowTest} exercises this exact object as the
|
||||
* real, wired production path, replacing what its own class doc used to call out as a
|
||||
* hand-mirrored copy of this lambda ("mirrors {@code Fleetd.main}'s {@code backendErrorSink}
|
||||
* lambda line-for-line") — that copy proved only itself, never that {@code main} still wires the
|
||||
* real thing. That was precisely the gap fleetd #248 exists to close.
|
||||
*
|
||||
* @param sessions the session registry; both read (roster) and written (onBackendError)
|
||||
* @param profiles the live profile map, normally {@code () -> config.get().profiles()} in
|
||||
* {@code main}, or a fixed test map via {@code () -> profiles}
|
||||
* @param pushLoop the lead-nudge loop, read lazily: {@code main} builds this sink before the
|
||||
* real {@link ReplyPushLoop} exists (a genuine construction-order cycle, broken
|
||||
* the same way {@code exhaustionSinkRef} is a few lines above it), so a
|
||||
* {@link Supplier} reads whatever {@code main} has filled in by the time a real
|
||||
* backend error fires
|
||||
*/
|
||||
public static BackendErrorSink backendErrorSink(SessionManager sessions,
|
||||
Supplier<Map<String, FleetConfig.Profile>> profiles, BackendOutagePolicy outagePolicy,
|
||||
Supplier<ReplyPushLoop> pushLoop) {
|
||||
return (target, matchedLine, reason) -> {
|
||||
sessions.onBackendError(target, reason);
|
||||
|
||||
String profileName = sessions.roster().stream()
|
||||
.filter(session -> target.equals(session.terminalId()))
|
||||
.findFirst()
|
||||
.map(MemberSession::profile)
|
||||
.orElse(null);
|
||||
FleetConfig.Profile profile = profileName == null ? null : profiles.get().get(profileName);
|
||||
if (profile == null) {
|
||||
log.error("backend error on target '{}' ({}) but no profile could be resolved — the "
|
||||
+ "target is not (yet) in the roster — no cool-off applied (fleetd #201 Unit 5)",
|
||||
target, reason);
|
||||
ReplyPushLoop loop = pushLoop.get();
|
||||
if (loop != null) {
|
||||
loop.onBackendTargetUnmapped(target, reason);
|
||||
}
|
||||
return;
|
||||
}
|
||||
String credentialId = profile.effectiveCredentialId();
|
||||
Optional<BackendOutagePolicy.Incident> incident = outagePolicy.record(credentialId, target, reason);
|
||||
incident.ifPresent(inc -> {
|
||||
List<String> affectedProfiles = profiles.get().values().stream()
|
||||
.filter(p -> credentialId.equals(p.effectiveCredentialId()))
|
||||
.map(FleetConfig.Profile::profile)
|
||||
.sorted()
|
||||
.toList();
|
||||
log.warn("credential '{}' cooling off for {}s after backend errors on {} distinct "
|
||||
+ "target(s) (profile '{}'): {}", credentialId,
|
||||
inc.remainingCoolOffSeconds(), inc.evidenceCount(), profile.profile(), reason);
|
||||
ReplyPushLoop loop = pushLoop.get();
|
||||
if (loop != null) {
|
||||
loop.onBackendIncident(inc.id(), inc.targets(), credentialId, affectedProfiles,
|
||||
(int) inc.remainingCoolOffSeconds());
|
||||
}
|
||||
});
|
||||
};
|
||||
}
|
||||
|
||||
/** Injection seam for {@link #selectReplyInbox}: production binds {@link AmqpReplyInbox#open}. */
|
||||
@FunctionalInterface
|
||||
interface AmqpOpener {
|
||||
|
||||
@@ -644,14 +644,46 @@ public record FleetConfig(
|
||||
}
|
||||
|
||||
/**
|
||||
* The credential group this profile quarantines with (CB-578 stage B): the configured
|
||||
* {@link #credentialId} when set, else this profile's own name — so an unconfigured profile
|
||||
* quarantines alone, exactly as it did before this field existed. Two profiles that set the
|
||||
* same non-blank {@code credentialId} share one quarantine: a {@code BACKEND_EXHAUSTED}
|
||||
* classification on either one quarantines both.
|
||||
* The shared credential id every {@code subscription: true} profile falls back to when it
|
||||
* sets no explicit {@link #credentialId} (fleetd #176 stage 2, correcting an inert first cut
|
||||
* of that ticket). A subscription profile has no credential of its own to fall back to its
|
||||
* name for: it authenticates as the operator's own Claude login, and a host has exactly one
|
||||
* of those, whatever names the operator gives the profiles running on it. Falling back to the
|
||||
* profile's own name (the way an ordinary off-subscription profile does) would keep two
|
||||
* subscription profiles on one login apart from each other, which is the opposite of what
|
||||
* "one account" means.
|
||||
*
|
||||
* <p>Measured live and what it broke: a lead on profile {@code opus}, members on profile
|
||||
* {@code sonnet}, same Claude login, neither setting {@code credentialId}. Before this
|
||||
* sentinel, {@code opus.effectiveCredentialId()} was {@code "opus"} and {@code sonnet
|
||||
* .effectiveCredentialId()} was {@code "sonnet"} — so fleetd #176's lead-seat matcher (and,
|
||||
* this sentinel now also fixes, {@code CompositePeerLauncher}'s quarantine/cool-off spawn
|
||||
* refusal and {@code BackendOutagePolicy}'s incident grouping) silently never linked them: the
|
||||
* fix shipped, and stayed inert on the one host it was written for.
|
||||
*/
|
||||
public static final String SUBSCRIPTION_CREDENTIAL_ID = "<subscription>";
|
||||
|
||||
/**
|
||||
* The credential group this profile quarantines with (CB-578 stage B; extended fleetd #176
|
||||
* stage 2 — see {@link #SUBSCRIPTION_CREDENTIAL_ID}): the configured {@link #credentialId}
|
||||
* when set — that always wins, so an operator with two separate Claude logins on one host can
|
||||
* still keep them apart. Otherwise, a {@code subscription: true} profile falls back to
|
||||
* {@link #SUBSCRIPTION_CREDENTIAL_ID} rather than its own name; an ordinary off-subscription
|
||||
* profile falls back to its own name, exactly as it did before this field existed, so an
|
||||
* unconfigured off-subscription profile still quarantines alone.
|
||||
*
|
||||
* <p>A {@code BACKEND_EXHAUSTED} (or repeated backend-error) classification on any profile
|
||||
* sharing the result quarantines/cools off every profile that shares it — including, now,
|
||||
* every {@code subscription: true} profile with no explicit {@code credentialId}. That is
|
||||
* intended, not incidental: one Claude subscription hitting a usage limit really does take out
|
||||
* every profile running on it, the same way {@code credentialId: openai-shared} already lets
|
||||
* two OpenAI-backed profiles share one quarantine.
|
||||
*/
|
||||
public String effectiveCredentialId() {
|
||||
return (credentialId == null || credentialId.isBlank()) ? profile : credentialId;
|
||||
if (credentialId != null && !credentialId.isBlank()) {
|
||||
return credentialId;
|
||||
}
|
||||
return isSubscription() ? SUBSCRIPTION_CREDENTIAL_ID : profile;
|
||||
}
|
||||
|
||||
/** True when this profile's workers are granted a forge token to open their own PR (CB-302). */
|
||||
@@ -946,7 +978,15 @@ public record FleetConfig(
|
||||
* survives restarts of the agent inside it — so identity is now the tab label alone.
|
||||
*
|
||||
* @param profile the {@code profiles:} entry to launch this lead on when one must
|
||||
* be created; {@code null} ⇒ recognise-only, never create
|
||||
* be created; {@code null} ⇒ recognise-only, never create.
|
||||
* <p>fleetd #176: also the field {@code Fleetd.leadSeatLookup} reads
|
||||
* to learn which account this lead's own live session shares — set it
|
||||
* (safely, even on an already-running recognise-only lead: naming a
|
||||
* profile here never starts anything beyond {@code instances}) so a
|
||||
* {@code subscription: true} worker profile sharing its
|
||||
* {@code effectiveCredentialId()} has this lead's seat subtracted from
|
||||
* {@code fleet_list}'s {@code free}. {@code null} here also means this
|
||||
* lead's seat cannot be derived and is not counted.
|
||||
* @param tab the exact tab label hosting this lead, matched case-insensitively;
|
||||
* the only field identity depends on. Required — a lead with no
|
||||
* {@code tab} can never be discovered, launched or not
|
||||
|
||||
@@ -620,9 +620,9 @@ public final class CompletionResolver implements TurnListener {
|
||||
* @param allProfiles every configured profile name
|
||||
* @param configuredProfiles the subset of {@code allProfiles} that carry an exhausted pattern
|
||||
*/
|
||||
public static String coverage(Set<String> allProfiles, Set<String> configuredProfiles) {
|
||||
public static String coverage(String patternKey, Set<String> allProfiles, Set<String> configuredProfiles) {
|
||||
if (configuredProfiles.isEmpty()) {
|
||||
return "off (no profile has an exhaustedPattern configured; profiles: " + sorted(allProfiles) + ")";
|
||||
return "off (no profile has an " + patternKey + " configured; profiles: " + sorted(allProfiles) + ")";
|
||||
}
|
||||
Set<String> unconfigured = new TreeSet<>(allProfiles);
|
||||
unconfigured.removeAll(configuredProfiles);
|
||||
|
||||
@@ -96,6 +96,8 @@ public final class FleetMcp {
|
||||
private final QuarantineSource quarantine;
|
||||
/** fleetd #201 Unit 5: SEPARATE from {@link #quarantine} — see {@link OutageSource}'s doc. */
|
||||
private final OutageSource outage;
|
||||
/** fleetd #176: SEPARATE from both of the above — see {@link LeadSeatSource}'s doc. */
|
||||
private final LeadSeatSource leadSeats;
|
||||
/** CB-637: this daemon's lead-to-lead channel; {@code null} when no coordinator is configured. */
|
||||
private final LeadChannel leadChannel;
|
||||
|
||||
@@ -135,6 +137,24 @@ public final class FleetMcp {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #176: the seats a profile's own live LEAD session(s) hold on the same Claude
|
||||
* subscription — the third reason (alongside {@link QuarantineSource} and {@link OutageSource})
|
||||
* {@code free} can overstate what a fresh {@code fleet_spawn} would actually get.
|
||||
*
|
||||
* <p>{@code maxLoad} counts only <em>members</em>, never the lead itself. But a
|
||||
* {@code subscription: true} profile bills the operator's own Claude account, and the lead is
|
||||
* always a live {@code claude} session on that same account (it is never moved off-subscription
|
||||
* — see {@code LeadLauncher}). So a fan-out that fills every member slot still leaves the lead's
|
||||
* own seat unaccounted for, and the daemon reports a slot that was never really free. See
|
||||
* {@code Fleetd.leadSeatLookup} for how the count is derived — from {@code fleet.leaders.<name>
|
||||
* .profile} and each profile's {@code effectiveCredentialId()}, never a hardcoded constant.
|
||||
*/
|
||||
public record LeadSeatSource(Function<String, Integer> seatsFor) {
|
||||
/** Inert source — no profile is ever reported as sharing a seat with a lead. */
|
||||
public static LeadSeatSource none() { return new LeadSeatSource(_ -> 0); }
|
||||
}
|
||||
|
||||
/**
|
||||
* @param callers resolves each call's {@link Principal}; {@code null} disables authorization.
|
||||
* This surface needs its own enforcement: {@code /mcp} is a raw servlet on
|
||||
@@ -149,7 +169,7 @@ public final class FleetMcp {
|
||||
CallerResolver callers, Metrics metrics, CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
QuarantineSource quarantine) {
|
||||
this(messages, workers, sessions, identity, presence, primaryRegistry, callers, metrics, capacity,
|
||||
healthCoverage, quarantine, null, OutageSource.none());
|
||||
healthCoverage, quarantine, null, OutageSource.none(), LeadSeatSource.none());
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -163,12 +183,12 @@ public final class FleetMcp {
|
||||
CallerResolver callers, Metrics metrics, CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
QuarantineSource quarantine, LeadChannel leadChannel) {
|
||||
this(messages, workers, sessions, identity, presence, primaryRegistry, callers, metrics, capacity,
|
||||
healthCoverage, quarantine, leadChannel, OutageSource.none());
|
||||
healthCoverage, quarantine, leadChannel, OutageSource.none(), LeadSeatSource.none());
|
||||
}
|
||||
|
||||
/**
|
||||
* As above, with fleetd #201 Unit 5 cool-off facts for {@code fleet_list}/{@code fleet_profiles}
|
||||
* (see {@link OutageSource}). This is what {@code Fleetd.main} actually wires up.
|
||||
* (see {@link OutageSource}).
|
||||
*
|
||||
* @param outage required — pass {@link OutageSource#none()} for a caller that does not want the
|
||||
* feature, never a defaulting overload (the same rule {@code quarantine} follows).
|
||||
@@ -177,10 +197,28 @@ public final class FleetMcp {
|
||||
ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry,
|
||||
CallerResolver callers, Metrics metrics, CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
QuarantineSource quarantine, LeadChannel leadChannel, OutageSource outage) {
|
||||
this(messages, workers, sessions, identity, presence, primaryRegistry, callers, metrics, capacity,
|
||||
healthCoverage, quarantine, leadChannel, outage, LeadSeatSource.none());
|
||||
}
|
||||
|
||||
/**
|
||||
* As above, with fleetd #176 lead-seat facts (see {@link LeadSeatSource}). This is what
|
||||
* {@code Fleetd.main} actually wires up.
|
||||
*
|
||||
* @param leadSeats required — pass {@link LeadSeatSource#none()} for a caller that does not want
|
||||
* the feature, never a defaulting overload (the same rule {@code quarantine} and
|
||||
* {@code outage} follow).
|
||||
*/
|
||||
public FleetMcp(MessageService messages, PeerLauncher workers, SessionManager sessions,
|
||||
ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry,
|
||||
CallerResolver callers, Metrics metrics, CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
QuarantineSource quarantine, LeadChannel leadChannel, OutageSource outage,
|
||||
LeadSeatSource leadSeats) {
|
||||
this.leadChannel = leadChannel;
|
||||
this.capacity = capacity;
|
||||
this.quarantine = Objects.requireNonNull(quarantine, "quarantine");
|
||||
this.outage = Objects.requireNonNull(outage, "outage");
|
||||
this.leadSeats = Objects.requireNonNull(leadSeats, "leadSeats");
|
||||
this.healthCoverage = healthCoverage;
|
||||
McpJsonMapper json = new JacksonMcpJsonMapperSupplier().get();
|
||||
this.transport = HttpServletStreamableServerTransportProvider.builder()
|
||||
@@ -310,7 +348,7 @@ public final class FleetMcp {
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
|
||||
if (denied != null) return denied;
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
|
||||
callers == null ? Map.of() : callers.leads(),
|
||||
leadSeats, callers == null ? Map.of() : callers.leads(),
|
||||
callerTerminal(exchange),
|
||||
leadChannel == null ? null : leadChannel.selfCoordId());
|
||||
};
|
||||
@@ -994,7 +1032,8 @@ public final class FleetMcp {
|
||||
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
QuarantineSource quarantine, OutageSource outage,
|
||||
Map<String, String> leads, String selfTerm) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage, leads, selfTerm, null);
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
|
||||
LeadSeatSource.none(), leads, selfTerm, null);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1011,7 +1050,7 @@ public final class FleetMcp {
|
||||
QuarantineSource quarantine, Map<String, String> leads, String selfTerm,
|
||||
String selfCoordId) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, OutageSource.none(),
|
||||
leads, selfTerm, selfCoordId);
|
||||
LeadSeatSource.none(), leads, selfTerm, selfCoordId);
|
||||
}
|
||||
|
||||
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
|
||||
@@ -1019,6 +1058,16 @@ public final class FleetMcp {
|
||||
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
QuarantineSource quarantine, OutageSource outage,
|
||||
Map<String, String> leads, String selfTerm, String selfCoordId) {
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
|
||||
LeadSeatSource.none(), leads, selfTerm, selfCoordId);
|
||||
}
|
||||
|
||||
/** As above, plus fleetd #176 lead-seat facts (see {@link LeadSeatSource}). */
|
||||
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,
|
||||
String selfCoordId) {
|
||||
try {
|
||||
Map<String, Agent> live = workers.list().stream()
|
||||
.map(Agent.class::cast)
|
||||
@@ -1045,7 +1094,7 @@ public final class FleetMcp {
|
||||
}
|
||||
if (capacity.available()) result.put("capacity", profiles.stream()
|
||||
.map(profile -> capacityView(profile, capacity.liveCount(), capacity.maxLoad(), roster, messages,
|
||||
capacity.clock().getAsLong(), quarantine, outage)).toList());
|
||||
capacity.clock().getAsLong(), quarantine, outage, leadSeats)).toList());
|
||||
return text(json(result));
|
||||
} catch (HerdrException e) {
|
||||
return error("herdr error listing the fleet: " + e.getMessage());
|
||||
@@ -1084,20 +1133,33 @@ public final class FleetMcp {
|
||||
* {@code credentialId}/{@code coolingOffForSeconds}, but never {@code quarantinedForSeconds} —
|
||||
* that key is added only when exhaustion quarantine is ALSO active for this profile, since the
|
||||
* two checks are independent and either, both, or neither can be true.
|
||||
*
|
||||
* <p>fleetd #176: {@code maxLoad} counts panes, not subscription seats — it never counted the
|
||||
* lead's own seat on a {@code subscription: true} profile's account. {@link LeadSeatSource}
|
||||
* reports that count (0 for a non-subscription profile, or when no live lead shares its
|
||||
* credential), and it is subtracted from {@code free} the same way {@code live} already is —
|
||||
* {@code maxLoad} itself is left untouched, so the row still reports the configured cap. The
|
||||
* {@code leadSeats} key is added only when the count is positive, for the same
|
||||
* byte-identical-when-unused reason as the quarantine/cool-off keys above.
|
||||
*/
|
||||
private static Map<String, Object> capacityView(String profile, Function<String, Integer> liveCount,
|
||||
Function<String, Integer> maxLoad, List<MemberSession> roster,
|
||||
MessageService messages, long nowNanos, QuarantineSource quarantine,
|
||||
OutageSource outage) {
|
||||
OutageSource outage, LeadSeatSource leadSeats) {
|
||||
Integer cap = maxLoad.apply(profile);
|
||||
int live = liveCount.apply(profile);
|
||||
int leadSeatCount = leadSeats.seatsFor().apply(profile);
|
||||
int reclaimable = (int) roster.stream().filter(s -> profile.equals(s.profile()))
|
||||
.filter(s -> (s.state() == MemberSession.State.READY || s.state() == MemberSession.State.DONE))
|
||||
.filter(s -> messages == null || (!messages.hasAcceptedDelivery(s.terminalId()) && !messages.hasInboxMessage(s.terminalId())))
|
||||
.count();
|
||||
Map<String, Object> row = new LinkedHashMap<>();
|
||||
row.put("profile", profile); row.put("maxLoad", cap); row.put("live", live);
|
||||
row.put("free", cap == null ? null : Math.max(0, cap - live)); row.put("reclaimable", reclaimable);
|
||||
row.put("free", cap == null ? null : Math.max(0, cap - live - leadSeatCount));
|
||||
row.put("reclaimable", reclaimable);
|
||||
if (leadSeatCount > 0) {
|
||||
row.put("leadSeats", leadSeatCount);
|
||||
}
|
||||
String credentialId = quarantine.credentialIdFor().apply(profile);
|
||||
if (credentialId != null) {
|
||||
quarantine.quarantine().remainingSeconds(credentialId).ifPresent(remaining -> {
|
||||
|
||||
@@ -0,0 +1,59 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.inject.BackendErrorPatternLookup;
|
||||
import dev.ltms.fleet.peer.MemberRole;
|
||||
import dev.ltms.fleet.session.MemberSession;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.regex.Pattern;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertNull;
|
||||
|
||||
/**
|
||||
* fleetd #248 / fleetd#201 Unit 5: {@link Fleetd#backendErrorPatternLookup} is the factory that
|
||||
* replaced the local lambda {@code Fleetd.main} used to build {@code backendErrorPatterns} — one
|
||||
* of the two arguments {@code CompletionResolver} lost cleanly (0 compile errors, every test still
|
||||
* green) when this ticket's measurement dropped it alongside {@code backendErrorSink}. This class
|
||||
* proves the factory's own behaviour; {@code FleetdCompletionResolverWiringTest} proves {@code
|
||||
* main} still passes its result into {@code CompletionResolver}.
|
||||
*/
|
||||
class FleetdBackendErrorPatternLookupTest {
|
||||
|
||||
private static MemberSession session(String terminal, String profile) {
|
||||
return new MemberSession("pane-" + terminal, terminal, profile, MemberRole.DEV,
|
||||
"/cwd", null, 0L, 0L, 0, MemberSession.State.READY, null, null);
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a target on a profile with a configured pattern resolves to that pattern")
|
||||
void configuredProfileResolves() {
|
||||
Map<String, Pattern> byProfile = Map.of("terra", Pattern.compile("(?i)503"));
|
||||
BackendErrorPatternLookup lookup =
|
||||
Fleetd.backendErrorPatternLookup(() -> List.of(session("term1", "terra")), byProfile);
|
||||
|
||||
assertEquals("(?i)503", lookup.patternFor("term1").pattern());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a target on a profile with no configured pattern resolves to null")
|
||||
void unconfiguredProfileResolvesToNull() {
|
||||
Map<String, Pattern> byProfile = Map.of("terra", Pattern.compile("x"));
|
||||
BackendErrorPatternLookup lookup =
|
||||
Fleetd.backendErrorPatternLookup(() -> List.of(session("term1", "sol")), byProfile);
|
||||
|
||||
assertNull(lookup.patternFor("term1"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("an unknown target resolves to null")
|
||||
void unknownTargetResolvesToNull() {
|
||||
BackendErrorPatternLookup lookup =
|
||||
Fleetd.backendErrorPatternLookup(List::of, Map.of("terra", Pattern.compile("x")));
|
||||
|
||||
assertNull(lookup.patternFor("term_stranger"));
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,229 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
import dev.ltms.fleet.inject.BackendErrorSink;
|
||||
import dev.ltms.fleet.member.ClaudeCodeLauncher;
|
||||
import dev.ltms.fleet.member.CompositePeerLauncher;
|
||||
import dev.ltms.fleet.mcp.PrimaryRegistry;
|
||||
import dev.ltms.fleet.msg.InMemoryReplyInbox;
|
||||
import dev.ltms.fleet.msg.ReplyPushLoop;
|
||||
import dev.ltms.fleet.peer.Capability;
|
||||
import dev.ltms.fleet.peer.PeerHandle;
|
||||
import dev.ltms.fleet.peer.PeerLauncher;
|
||||
import dev.ltms.fleet.peer.SpawnRequest;
|
||||
import dev.ltms.fleet.placement.BackendOutagePolicy;
|
||||
import dev.ltms.fleet.placement.BackendQuarantine;
|
||||
import dev.ltms.fleet.placement.PlacementPolicies;
|
||||
import dev.ltms.fleet.session.MemberSession;
|
||||
import dev.ltms.fleet.session.SessionManager;
|
||||
import org.junit.jupiter.api.AfterEach;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.CopyOnWriteArrayList;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #248 / fleetd#201 Unit 5: {@link Fleetd#backendErrorSink} is the factory that replaced
|
||||
* the local lambda {@code Fleetd.main} used to build {@code backendErrorSink} — the other half of
|
||||
* the pair this ticket's measurement dropped cleanly (0 compile errors, every test still green).
|
||||
*
|
||||
* <p>Before this ticket, the closest thing to coverage was {@code
|
||||
* dev.ltms.fleet.inject.BackendOutageFlowTest}, whose own class doc said it "mirrors {@code
|
||||
* Fleetd.main}'s {@code backendErrorSink} lambda line-for-line" — a hand-copy that proves itself,
|
||||
* never that {@code main} still wires the real thing. This class exercises the actual production
|
||||
* factory instead. {@code FleetdCompletionResolverWiringTest} proves {@code main} still passes its
|
||||
* result into {@code CompletionResolver}.
|
||||
*/
|
||||
class FleetdBackendErrorSinkTest {
|
||||
|
||||
private final List<ScheduledExecutorService> schedulers = new ArrayList<>();
|
||||
|
||||
@AfterEach
|
||||
void tearDown() {
|
||||
schedulers.forEach(ScheduledExecutorService::shutdownNow);
|
||||
}
|
||||
|
||||
private static FleetConfig.Profile stubWorker(String profile, String credentialId) {
|
||||
return new FleetConfig.Profile(profile, "http://gx00.gw:8000", "coder",
|
||||
null, "FLEETD_WORKER_TOKEN", List.of("claude"), "tab", "fleetd-workers",
|
||||
"w #{n}", null, null, null, null, null, null, null, null, null,
|
||||
null, null, credentialId, null);
|
||||
}
|
||||
|
||||
private static Map<String, FleetConfig.Profile> orderedProfiles() {
|
||||
Map<String, FleetConfig.Profile> m = new LinkedHashMap<>();
|
||||
m.put("terra", stubWorker("terra", "shared-openai"));
|
||||
m.put("sol", stubWorker("sol", "shared-openai"));
|
||||
return m;
|
||||
}
|
||||
|
||||
/** Minimal recording {@code HerdrClient} for the LEAD pane — mirrors ReplyPushLoopTest's own. */
|
||||
private static final class RecordingLeadClient implements HerdrClient {
|
||||
private static final ObjectMapper MAPPER = new ObjectMapper();
|
||||
private final List<Object> prompts = new CopyOnWriteArrayList<>();
|
||||
volatile CountDownLatch sendLatch = new CountDownLatch(1);
|
||||
|
||||
@Override
|
||||
public JsonNode call(String method, Object params) {
|
||||
if ("agent.get".equals(method)) {
|
||||
return MAPPER.createObjectNode().set("agent", MAPPER.createObjectNode()
|
||||
.put("terminal_id", "term_primary").put("agent_status", "idle"));
|
||||
}
|
||||
if ("agent.prompt".equals(method)) {
|
||||
prompts.add(params);
|
||||
sendLatch.countDown();
|
||||
}
|
||||
return MAPPER.createObjectNode();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
}
|
||||
|
||||
int sendCount() {
|
||||
return prompts.size();
|
||||
}
|
||||
}
|
||||
|
||||
/** A {@link PeerLauncher} that never actually spawns — enough to construct a bare {@link SessionManager}. */
|
||||
private static final class NeverSpawnsLauncher implements PeerLauncher {
|
||||
@Override
|
||||
public Set<Capability> capabilities() {
|
||||
return Set.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<Capability> capabilitiesFor(String profileName) {
|
||||
return Set.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public PeerHandle spawn(SpawnRequest req) {
|
||||
throw new UnsupportedOperationException("not reachable — this test never acquires a session");
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<String> profiles() {
|
||||
return Set.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public String defaultProfile() {
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String effectiveCwd(SpawnRequest req) {
|
||||
throw new UnsupportedOperationException("not reachable — this test never acquires a session");
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<String> parityOverlay(String profileName) {
|
||||
return List.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<?> list() {
|
||||
return List.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public int reapOrphanWorkers() {
|
||||
return 0;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void stop(String id) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean clearContext(String id) {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a target with no resolvable profile logs and returns without recording an incident (never throws)")
|
||||
void unresolvableProfileDoesNotRecordOrThrow() {
|
||||
SessionManager sessions = new SessionManager(new NeverSpawnsLauncher());
|
||||
BackendOutagePolicy outagePolicy = new BackendOutagePolicy(() -> 0L);
|
||||
RecordingLeadClient leadClient = new RecordingLeadClient();
|
||||
InMemoryReplyInbox inbox = new InMemoryReplyInbox();
|
||||
PrimaryRegistry registry = new PrimaryRegistry(null);
|
||||
ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor();
|
||||
schedulers.add(scheduler);
|
||||
ReplyPushLoop pushLoop = new ReplyPushLoop(registry, new AgentControl(leadClient), inbox, scheduler, 3, 50);
|
||||
|
||||
BackendErrorSink sink = Fleetd.backendErrorSink(sessions, Map::of, outagePolicy, () -> pushLoop);
|
||||
sink.onBackendError("term_unmapped", "matched line", "503 Service Unavailable");
|
||||
|
||||
assertTrue(outagePolicy.remainingCoolOffSeconds("shared-openai").isEmpty(),
|
||||
"no credential is ever resolvable here, so nothing must be recorded");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("two distinct targets classified through the real factory start an incident and cool the credential")
|
||||
void twoDistinctTargetsStartAnIncident() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr()
|
||||
.readText("⏺ 503 Service Unavailable: upstream credential rejected\n❯ ");
|
||||
Map<String, FleetConfig.Profile> profiles = orderedProfiles();
|
||||
AtomicLong clockNanos = new AtomicLong(0L);
|
||||
BackendOutagePolicy outagePolicy = new BackendOutagePolicy(clockNanos::get);
|
||||
ClaudeCodeLauncher adapter = new ClaudeCodeLauncher(new AgentControl(herdr),
|
||||
new WorkspaceControl(herdr), new SubscriptionGuard(Set.of("gx00.gw")),
|
||||
profiles, "terra", _ -> "tok");
|
||||
CompositePeerLauncher workers = new CompositePeerLauncher(List.of(adapter), "terra", profiles,
|
||||
PlacementPolicies.weighted(), _ -> 0, null, BackendQuarantine.none(), outagePolicy);
|
||||
SessionManager sessions = new SessionManager(workers);
|
||||
MemberSession s1 = sessions.acquire("terra", null, null, null);
|
||||
MemberSession s2 = sessions.acquire("terra", null, null, null);
|
||||
|
||||
PrimaryRegistry registry = new PrimaryRegistry(null);
|
||||
registry.recordDelegation(s1.terminalId(), "term_primary");
|
||||
registry.recordDelegation(s2.terminalId(), "term_primary");
|
||||
InMemoryReplyInbox inbox = new InMemoryReplyInbox();
|
||||
inbox.own(s1.terminalId());
|
||||
inbox.own(s2.terminalId());
|
||||
RecordingLeadClient leadClient = new RecordingLeadClient();
|
||||
ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor();
|
||||
schedulers.add(scheduler);
|
||||
ReplyPushLoop pushLoop = new ReplyPushLoop(registry, new AgentControl(leadClient), inbox, scheduler, 3, 50);
|
||||
AtomicReference<ReplyPushLoop> pushLoopRef = new AtomicReference<>(pushLoop);
|
||||
|
||||
// The exact object under test: Fleetd's real production factory, not a hand copy.
|
||||
BackendErrorSink sink = Fleetd.backendErrorSink(sessions, () -> profiles, outagePolicy, pushLoopRef::get);
|
||||
|
||||
sink.onBackendError(s1.terminalId(), "matched line", "503 Service Unavailable");
|
||||
assertTrue(outagePolicy.remainingCoolOffSeconds("shared-openai").isEmpty(),
|
||||
"one distinct target must not start a cool-off");
|
||||
|
||||
sink.onBackendError(s2.terminalId(), "matched line", "503 Service Unavailable");
|
||||
|
||||
assertTrue(leadClient.sendLatch.await(3, TimeUnit.SECONDS),
|
||||
"the second distinct target must cross the threshold and nudge the lead");
|
||||
var remaining = outagePolicy.remainingCoolOffSeconds("shared-openai");
|
||||
assertTrue(remaining.isPresent(), "two distinct targets must start a cool-off");
|
||||
assertEquals(1, leadClient.sendCount());
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,89 @@
|
||||
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 #248: this is the test that was actually missing. {@code Fleetd.main} builds its {@code
|
||||
* CompletionResolver} from an 8-argument constructor, and the ticket's own measurement proved two
|
||||
* ways to silently unwire it — both compiled with 0 errors and left every existing test green:
|
||||
*
|
||||
* <ul>
|
||||
* <li>replacing the worktree/branch argument (the 8th) with {@code _ -> null} — drops
|
||||
* fleetd#241's fallback-report location entirely;</li>
|
||||
* <li>replacing {@code backendErrorPatterns, backendErrorSink} (5th/6th) with {@code
|
||||
* BackendErrorPatternLookup.legacy(), BackendErrorSink.none()} — drops fleetd#201 Unit 5's
|
||||
* backend-error classification and cool-off entirely.</li>
|
||||
* </ul>
|
||||
*
|
||||
* <p>Neither mutation could be caught by any test that constructs its own {@code
|
||||
* CompletionResolver} (every test before this one did exactly that) or by a test of {@link
|
||||
* Fleetd#worktreeBranchLookup}, {@link Fleetd#backendErrorPatternLookup}, or {@link
|
||||
* Fleetd#backendErrorSink} in isolation (see {@code FleetdWorktreeBranchLookupTest}, {@code
|
||||
* FleetdBackendErrorPatternLookupTest}, {@code FleetdBackendErrorSinkTest}) — those prove the
|
||||
* factories work, never that {@code main} still calls them. This class is a plain source-text
|
||||
* assertion on {@code Fleetd.java} — crude, but honest about what it checks, and it turns red the
|
||||
* instant the wiring is dropped, mirroring the same fallback shape {@link
|
||||
* FleetdFleetAppConstructionTest} already uses for a different constructor argument.
|
||||
*
|
||||
* <p><b>This test checks source text, not runtime behaviour.</b> It never constructs a {@code
|
||||
* CompletionResolver} and never runs {@code main}.
|
||||
*/
|
||||
class FleetdCompletionResolverWiringTest {
|
||||
|
||||
private static String fleetdSource() throws Exception {
|
||||
return Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SOURCE TEXT] CompletionResolver's construction call still names backendErrorPatterns and backendErrorSink")
|
||||
void backendErrorArgumentsAreStillNamedAtTheCallSite() throws Exception {
|
||||
String source = fleetdSource();
|
||||
assertTrue(source.contains(
|
||||
"exhaustionSink, backendErrorPatterns, backendErrorSink, System::nanoTime,"),
|
||||
"CompletionResolver's construction call must still pass backendErrorPatterns and "
|
||||
+ "backendErrorSink as its 5th/6th arguments. Replacing them with "
|
||||
+ "BackendErrorPatternLookup.legacy()/BackendErrorSink.none() (fleetd #248's measured "
|
||||
+ "mutation) compiles with 0 errors and leaves every behavioural test green — this "
|
||||
+ "source check is what must go red instead.");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SOURCE TEXT] CompletionResolver's construction call still passes worktreeBranchLookup(sessions::roster)")
|
||||
void worktreeBranchLookupIsStillPassedAtTheCallSite() throws Exception {
|
||||
String source = fleetdSource();
|
||||
assertTrue(source.contains("worktreeBranchLookup(sessions::roster)"),
|
||||
"CompletionResolver's construction call must still pass worktreeBranchLookup(sessions::roster) "
|
||||
+ "as its 8th (last) argument. Replacing it with the inert `_ -> null` (fleetd #248's "
|
||||
+ "other measured mutation) compiles with 0 errors and leaves every behavioural test "
|
||||
+ "green — this source check is what must go red instead.");
|
||||
assertFalse(source.contains("System::nanoTime,\n _ -> null"),
|
||||
"the worktree/branch argument must never regress to the inert `_ -> null` literal");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SOURCE TEXT] backendErrorPatterns is assigned from the extracted backendErrorPatternLookup(...) factory")
|
||||
void backendErrorPatternsComesFromTheFactory() throws Exception {
|
||||
String source = fleetdSource();
|
||||
assertTrue(source.contains(
|
||||
"BackendErrorPatternLookup backendErrorPatterns = backendErrorPatternLookup(sessions::roster,"),
|
||||
"backendErrorPatterns must be assigned from Fleetd.backendErrorPatternLookup(...), not an "
|
||||
+ "inline lambda that a source check on the CompletionResolver call alone cannot see "
|
||||
+ "through");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SOURCE TEXT] backendErrorSink is assigned from the extracted backendErrorSink(...) factory")
|
||||
void backendErrorSinkComesFromTheFactory() throws Exception {
|
||||
String source = fleetdSource();
|
||||
assertTrue(source.contains(
|
||||
"BackendErrorSink backendErrorSink = backendErrorSink(sessions, () -> config.get().profiles(),"),
|
||||
"backendErrorSink must be assigned from Fleetd.backendErrorSink(...), not an inline lambda "
|
||||
+ "that a source check on the CompletionResolver call alone cannot see through");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,156 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.Map;
|
||||
import java.util.function.Function;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
|
||||
/**
|
||||
* fleetd #176: {@link Fleetd#leadSeatLookup} is the factory {@code Fleetd.main} wires into {@code
|
||||
* FleetMcp.LeadSeatSource} so {@code fleet_list}'s {@code free} can subtract the seat(s) a
|
||||
* {@code subscription: true} profile's own live LEAD session holds on that same account —
|
||||
* {@code maxLoad} never counted the lead, only members. {@code FleetdLeadSeatWiringTest} proves
|
||||
* {@code main} still passes this factory's result in; this class proves the factory's own matching
|
||||
* logic: subscription-only, credential-matched, and counting only CURRENTLY LIVE leads.
|
||||
*/
|
||||
class FleetdLeadSeatLookupTest {
|
||||
|
||||
private static FleetConfig.Profile subscriptionProfile(String name, String credentialId) {
|
||||
return new FleetConfig.Profile(name, null, "claude-sonnet-5", null, null, null,
|
||||
"tab", "fleet", "w #{n}", null, null, null, null, null, null, null,
|
||||
null, 3, true, null, credentialId, null);
|
||||
}
|
||||
|
||||
private static FleetConfig.Profile offSubscriptionProfile(String name, String credentialId) {
|
||||
return new FleetConfig.Profile(name, "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
|
||||
null, "tab", "fleet", "w #{n}", null, null, null, null, null, null, null,
|
||||
null, 2, false, null, credentialId, null);
|
||||
}
|
||||
|
||||
private static FleetConfig.Leader leadOnProfile(String profile) {
|
||||
return new FleetConfig.Leader(profile, "lead: primary", 1, "lead:", 10, "claude", "claude-sonnet-5");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a live lead sharing the target profile's credential counts as one seat")
|
||||
void liveLeadSharingCredentialCountsAsOneSeat() {
|
||||
Map<String, FleetConfig.Profile> profiles = Map.of("sonnet", subscriptionProfile("sonnet", null));
|
||||
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("sonnet"));
|
||||
Function<String, Integer> lookup = Fleetd.leadSeatLookup(() -> profiles, leaders,
|
||||
() -> Map.of("term_primary", "primary"));
|
||||
|
||||
assertEquals(1, lookup.apply("sonnet"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("no live lead names this profile ⇒ zero seats, exactly as before this ticket")
|
||||
void noLiveLeadOnTheProfileCountsAsZero() {
|
||||
Map<String, FleetConfig.Profile> profiles = Map.of("sonnet", subscriptionProfile("sonnet", null));
|
||||
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("sonnet"));
|
||||
Function<String, Integer> lookup = Fleetd.leadSeatLookup(() -> profiles, leaders, Map::of);
|
||||
|
||||
assertEquals(0, lookup.apply("sonnet"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a lead entry with no `profile:` (recognise-only) contributes no seats — cannot be derived")
|
||||
void recogniseOnlyLeadWithNoProfileContributesNothing() {
|
||||
Map<String, FleetConfig.Profile> profiles = Map.of("sonnet", subscriptionProfile("sonnet", null));
|
||||
FleetConfig.Leader recogniseOnly = new FleetConfig.Leader(null, "lead: primary", 1, "lead:", 10,
|
||||
"claude", "claude-sonnet-5");
|
||||
Map<String, FleetConfig.Leader> leaders = Map.of("primary", recogniseOnly);
|
||||
Function<String, Integer> lookup = Fleetd.leadSeatLookup(() -> profiles, leaders,
|
||||
() -> Map.of("term_primary", "primary"));
|
||||
|
||||
assertEquals(0, lookup.apply("sonnet"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a non-subscription profile never has a lead seat subtracted, whatever the credential match")
|
||||
void nonSubscriptionProfileIsNeverAdjusted() {
|
||||
Map<String, FleetConfig.Profile> profiles = Map.of(
|
||||
"terra", offSubscriptionProfile("terra", "shared-openai"),
|
||||
"sonnet", subscriptionProfile("sonnet", "shared-openai"));
|
||||
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("sonnet"));
|
||||
Function<String, Integer> lookup = Fleetd.leadSeatLookup(() -> profiles, leaders,
|
||||
() -> Map.of("term_primary", "primary"));
|
||||
|
||||
assertEquals(0, lookup.apply("terra"), "terra is not subscription:true, so it must never be adjusted");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("explicit, different credentialIds still separate two subscription profiles (post fleetd #176 "
|
||||
+ "stage 2 sentinel)")
|
||||
void differentCredentialIsNotCounted() {
|
||||
Map<String, FleetConfig.Profile> profiles = Map.of(
|
||||
"sonnet", subscriptionProfile("sonnet", "claude-account-a"),
|
||||
"opus", subscriptionProfile("opus", "claude-account-b"));
|
||||
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("opus"));
|
||||
Function<String, Integer> lookup = Fleetd.leadSeatLookup(() -> profiles, leaders,
|
||||
() -> Map.of("term_primary", "primary"));
|
||||
|
||||
assertEquals(0, lookup.apply("sonnet"), "different accounts must never be conflated into one seat count "
|
||||
+ "— an explicit credentialId on both sides must still win over the subscription sentinel, so an "
|
||||
+ "operator with two separate Claude logins on one host can keep them apart");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #176 stage 2 — the exact live shape that shipped inert: a lead on subscription profile
|
||||
* {@code opus}, members on a DIFFERENTLY NAMED subscription profile {@code sonnet}, same Claude
|
||||
* login, and NEITHER profile sets {@code credentialId}. Every other test in this class puts the
|
||||
* lead on the SAME profile name as the target, which happened to keep working even with the old
|
||||
* fall-back-to-profile-name {@code effectiveCredentialId()} — this is the one that did not, and
|
||||
* its absence is what let the stage-1 fix ship without ever catching the bug it was filed for.
|
||||
*/
|
||||
@Test
|
||||
@DisplayName("[LIVE SHAPE] lead on a DIFFERENT subscription profile, same account, neither sets "
|
||||
+ "credentialId ⇒ still counts as a seat")
|
||||
void leadOnADifferentSubscriptionProfileSameAccountStillCountsAsASeat() {
|
||||
Map<String, FleetConfig.Profile> profiles = Map.of(
|
||||
"opus", subscriptionProfile("opus", null),
|
||||
"sonnet", subscriptionProfile("sonnet", null));
|
||||
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("opus"));
|
||||
Function<String, Integer> lookup = Fleetd.leadSeatLookup(() -> profiles, leaders,
|
||||
() -> Map.of("term_primary", "primary"));
|
||||
|
||||
assertEquals(1, lookup.apply("sonnet"), "opus and sonnet are both subscription:true with no explicit "
|
||||
+ "credentialId, so they share one Claude login and the lead's live seat on opus must be charged "
|
||||
+ "against sonnet too — this is the live host's actual shape (fleetd #176 stage 2)");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("two live instances of the same lead count as two seats")
|
||||
void twoLiveInstancesOfTheSameLeadCountAsTwoSeats() {
|
||||
Map<String, FleetConfig.Profile> profiles = Map.of("sonnet", subscriptionProfile("sonnet", null));
|
||||
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("sonnet"));
|
||||
Function<String, Integer> lookup = Fleetd.leadSeatLookup(() -> profiles, leaders,
|
||||
() -> Map.of("term_a", "primary", "term_b", "primary"));
|
||||
|
||||
assertEquals(2, lookup.apply("sonnet"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("an unconfigured target profile resolves to zero, not a thrown exception")
|
||||
void unconfiguredTargetProfileIsZero() {
|
||||
Function<String, Integer> lookup = Fleetd.leadSeatLookup(Map::of, Map.of(), Map::of);
|
||||
|
||||
assertEquals(0, lookup.apply("ghost"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("live leads are read through the supplier on every call, not snapshotted")
|
||||
void liveLeadsAreReadThroughOnEveryCall() {
|
||||
Map<String, FleetConfig.Profile> profiles = Map.of("sonnet", subscriptionProfile("sonnet", null));
|
||||
Map<String, FleetConfig.Leader> leaders = Map.of("primary", leadOnProfile("sonnet"));
|
||||
java.util.Map<String, String> live = new java.util.HashMap<>();
|
||||
Function<String, Integer> lookup = Fleetd.leadSeatLookup(() -> profiles, leaders, () -> live);
|
||||
|
||||
assertEquals(0, lookup.apply("sonnet"));
|
||||
live.put("term_primary", "primary");
|
||||
assertEquals(1, lookup.apply("sonnet"));
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,43 @@
|
||||
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.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #176: {@code Fleetd.main} builds its {@code FleetMcp} from a 14-argument constructor whose
|
||||
* last argument is a {@code FleetMcp.LeadSeatSource} wrapping {@link Fleetd#leadSeatLookup}. That
|
||||
* argument is exactly the kind of wiring fleetd #248 warned about: dropping it (or swapping it for
|
||||
* the inert {@code FleetMcp.LeadSeatSource.none()}) compiles with 0 errors and leaves every test
|
||||
* that builds its own {@code FleetMcp}/{@code CapacitySource} directly — every test that predates
|
||||
* this ticket — green, because none of them go through {@code main} at all.
|
||||
*
|
||||
* <p>{@link FleetdLeadSeatLookupTest} proves the factory's own matching logic; this class is the
|
||||
* plain source-text assertion that proves {@code main} still passes its result in, mirroring
|
||||
* {@code FleetdCompletionResolverWiringTest}'s approach for the same class of gap.
|
||||
*
|
||||
* <p><b>This test checks source text, not runtime behaviour.</b> It never constructs a
|
||||
* {@code FleetMcp} and never runs {@code main}.
|
||||
*/
|
||||
class FleetdLeadSeatWiringTest {
|
||||
|
||||
private static String fleetdSource() throws Exception {
|
||||
return Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SOURCE TEXT] FleetMcp's construction call still passes a LeadSeatSource built from leadSeatLookup(...)")
|
||||
void fleetMcpConstructionStillWiresLeadSeatLookup() throws Exception {
|
||||
String source = fleetdSource();
|
||||
assertTrue(source.contains("new FleetMcp.LeadSeatSource(leadSeatLookup(() -> config.get().profiles(), "
|
||||
+ "leaders, leads))"),
|
||||
"FleetMcp's construction call must still pass a LeadSeatSource built from "
|
||||
+ "Fleetd.leadSeatLookup(...). Dropping it or swapping in "
|
||||
+ "FleetMcp.LeadSeatSource.none() (fleetd #176's would-be silent regression, the same "
|
||||
+ "shape as fleetd #248's measured mutations) compiles with 0 errors and leaves every "
|
||||
+ "existing behavioural test green — this source check is what must go red instead.");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,64 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.inject.CompletionResolver;
|
||||
import dev.ltms.fleet.peer.MemberRole;
|
||||
import dev.ltms.fleet.session.MemberSession;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.function.Function;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertNull;
|
||||
|
||||
/**
|
||||
* fleetd #248: {@link Fleetd#worktreeBranchLookup} is the factory that replaced the anonymous
|
||||
* lambda {@code Fleetd.main} used to build inline, as the 8th (last) argument to {@code
|
||||
* CompletionResolver}'s constructor. Before this ticket that argument was untestable wiring:
|
||||
* replacing it with {@code _ -> null} compiled clean and every existing test stayed green, because
|
||||
* every existing test builds its own {@code CompletionResolver} directly rather than going through
|
||||
* {@code main}. This class proves the factory's own behaviour; {@code
|
||||
* FleetdCompletionResolverWiringTest} proves {@code main} still passes it in.
|
||||
*/
|
||||
class FleetdWorktreeBranchLookupTest {
|
||||
|
||||
private static MemberSession session(String terminal, String worktree, String branch) {
|
||||
return new MemberSession("pane-" + terminal, terminal, "terra", MemberRole.DEV,
|
||||
"/cwd", null, 0L, 0L, 0, MemberSession.State.READY, worktree, branch);
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("a known target resolves to its session's worktree and branch")
|
||||
void knownTargetResolves() {
|
||||
Function<String, CompletionResolver.WorktreeBranch> lookup =
|
||||
Fleetd.worktreeBranchLookup(() -> List.of(session("term1", "/wt/worker_x", "worker/x")));
|
||||
|
||||
CompletionResolver.WorktreeBranch resolved = lookup.apply("term1");
|
||||
|
||||
assertEquals("/wt/worker_x", resolved.worktree());
|
||||
assertEquals("worker/x", resolved.branch());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("an unknown target resolves to null, not a thrown exception")
|
||||
void unknownTargetResolvesToNull() {
|
||||
Function<String, CompletionResolver.WorktreeBranch> lookup =
|
||||
Fleetd.worktreeBranchLookup(() -> List.of(session("term1", "/wt/worker_x", "worker/x")));
|
||||
|
||||
assertNull(lookup.apply("term_stranger"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("the roster is read through the supplier on every call, not snapshotted")
|
||||
void rosterIsReadThroughOnEveryCall() {
|
||||
List<MemberSession> roster = new ArrayList<>();
|
||||
Function<String, CompletionResolver.WorktreeBranch> lookup = Fleetd.worktreeBranchLookup(() -> roster);
|
||||
|
||||
assertNull(lookup.apply("term_late"));
|
||||
roster.add(session("term_late", "/wt/late", "worker/late"));
|
||||
|
||||
assertEquals("/wt/late", lookup.apply("term_late").worktree());
|
||||
}
|
||||
}
|
||||
@@ -768,19 +768,37 @@ class CompletionResolverTest {
|
||||
@Test
|
||||
void coverageIsOffWhenNoProfileHasAPatternConfigured() {
|
||||
assertEquals("off (no profile has an exhaustedPattern configured; profiles: [terra])",
|
||||
CompletionResolver.coverage(Set.of("terra"), Set.of()));
|
||||
CompletionResolver.coverage("exhaustedPattern", Set.of("terra"), Set.of()));
|
||||
}
|
||||
|
||||
@Test
|
||||
void coverageIsFullWhenEveryProfileHasAPatternConfigured() {
|
||||
assertEquals("full (all profiles configured: [gx10, terra])",
|
||||
CompletionResolver.coverage(Set.of("terra", "gx10"), Set.of("terra", "gx10")));
|
||||
CompletionResolver.coverage("exhaustedPattern", Set.of("terra", "gx10"), Set.of("terra", "gx10")));
|
||||
}
|
||||
|
||||
@Test
|
||||
void coverageIsPartialAndNamesWhichProfilesAreConfigured() {
|
||||
assertEquals("partial (configured: [terra]; not configured: [gx10])",
|
||||
CompletionResolver.coverage(Set.of("terra", "gx10"), Set.of("terra")));
|
||||
CompletionResolver.coverage("exhaustedPattern", Set.of("terra", "gx10"), Set.of("terra")));
|
||||
}
|
||||
|
||||
/**
|
||||
* Found live on 2026-09-03, reading a real boot log rather than a test. {@code coverage} is
|
||||
* shared by two call sites — CB-578's {@code exhaustedPattern} line and fleetd#201 Unit 5's
|
||||
* {@code errorPattern} line — but its "off" branch hard-coded the word {@code exhaustedPattern}.
|
||||
* So a daemon with no {@code errorPattern} anywhere printed "no profile has an exhaustedPattern
|
||||
* configured" directly beneath a line reporting that two profiles DO have one. Both lines were
|
||||
* individually defensible and together they were nonsense, and the message sent an operator to
|
||||
* set the wrong key.
|
||||
*
|
||||
* <p>Every earlier test here passed the exhaustion case only, so none of them could see it. This
|
||||
* one pins that the message names the key the caller actually meant.
|
||||
*/
|
||||
@Test
|
||||
void coverageNamesTheConfigKeyItsCallerMeansRatherThanAlwaysSayingExhaustedPattern() {
|
||||
assertEquals("off (no profile has an errorPattern configured; profiles: [gx10, terra])",
|
||||
CompletionResolver.coverage("errorPattern", Set.of("terra", "gx10"), Set.of()));
|
||||
}
|
||||
|
||||
// --- fleetd#201 Unit 1: target-keyed backend-error pattern + typed sink ----------------------
|
||||
|
||||
@@ -651,6 +651,48 @@ class FleetMcpTest {
|
||||
assertEquals(2, out.split("\"free\":0", -1).length - 1, out);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #176 stage 2 (correcting the inert stage 1): two {@code subscription: true} profiles,
|
||||
* {@code opus} and {@code sonnet}, neither setting an explicit {@code credentialId} — the exact
|
||||
* shape measured on the live Mac fleet. This is INTENDED, not a regression: a real Claude usage
|
||||
* limit on the one login behind both profiles really does take out every profile running on it,
|
||||
* the same way {@code credentialId: openai-shared} already lets two OpenAI-backed profiles share
|
||||
* one quarantine (see {@code everyProfileSharingTheQuarantinedCredentialReportsZeroFree} above).
|
||||
* The {@code credentialIdFor} function here is built the same way {@code Fleetd.main} wires it —
|
||||
* {@code profile -> profiles.get(profile).effectiveCredentialId()} — so this proves the actual
|
||||
* config-driven behaviour, not just {@code capacityView}'s arithmetic with a hand-picked string.
|
||||
*/
|
||||
@Test
|
||||
void quarantiningOneSubscriptionProfileZeroesFreeOnTheOtherSharingTheSameAccount() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
|
||||
Map<String, FleetConfig.Profile> profiles = Map.of(
|
||||
"opus", new FleetConfig.Profile("opus", null, "claude-opus-4", null, null, null,
|
||||
"tab", "fleet", "w #{n}", null, null, null, null, null, null, null,
|
||||
null, 3, true, null, null, null),
|
||||
"sonnet", new FleetConfig.Profile("sonnet", null, "claude-sonnet-5", null, null, null,
|
||||
"tab", "fleet", "w #{n}", null, null, null, null, null, null, null,
|
||||
null, 3, true, null, null, null));
|
||||
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(20));
|
||||
quarantine.quarantine(FleetConfig.Profile.SUBSCRIPTION_CREDENTIAL_ID);
|
||||
FleetMcp.QuarantineSource source = new FleetMcp.QuarantineSource(profile -> {
|
||||
FleetConfig.Profile configured = profiles.get(profile);
|
||||
return configured == null ? null : configured.effectiveCredentialId();
|
||||
}, quarantine);
|
||||
|
||||
String out = textOf(FleetMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")),
|
||||
sessions, null, new FleetMcp.CapacitySource(profile -> 0, profile -> 3,
|
||||
() -> Set.of("opus", "sonnet"), () -> 0), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
source, Map.of(), ""));
|
||||
|
||||
assertEquals(2, out.split("\"free\":0", -1).length - 1,
|
||||
"opus and sonnet share one Claude login with neither setting credentialId, so quarantining "
|
||||
+ "opus's account must also zero sonnet's free — this is intended, not a side effect: "
|
||||
+ out);
|
||||
assertEquals(2, out.split("\"credentialId\":\"" + FleetConfig.Profile.SUBSCRIPTION_CREDENTIAL_ID + "\"", -1)
|
||||
.length - 1, out);
|
||||
}
|
||||
|
||||
/** fleetd #201 Unit 5: cool-off forces {@code free:0} but never adds {@code quarantinedForSeconds}. */
|
||||
@Test
|
||||
void coolingOffProfileReportsZeroFreeButNeverQuarantinedForSeconds() {
|
||||
@@ -767,6 +809,68 @@ class FleetMcpTest {
|
||||
assertFalse(out.contains("quarantinedForSeconds"), out);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #176: this is the exact shape measured on the Mac fleet — {@code maxLoad:3, live:2},
|
||||
* where one of the "free" three is really the lead's own seat on the same subscription. The old
|
||||
* formula ({@code max(0, cap - live)}) reported {@code free:1}; the real ceiling is {@code 0}
|
||||
* (two members plus the lead's own seat already fill all three).
|
||||
*/
|
||||
@Test
|
||||
void leadSeatSubtractsFromFreeTheSameWayLiveDoes() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
|
||||
FleetMcp.LeadSeatSource leadSeats = new FleetMcp.LeadSeatSource(
|
||||
profile -> "sonnet".equals(profile) ? 1 : 0);
|
||||
|
||||
String out = textOf(FleetMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")),
|
||||
sessions, null, new FleetMcp.CapacitySource(profile -> 2, profile -> 3,
|
||||
() -> Set.of("sonnet"), () -> 0), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), leadSeats, Map.of(), "", null));
|
||||
|
||||
assertTrue(out.contains("\"maxLoad\":3"), "maxLoad itself must be left untouched: " + out);
|
||||
assertTrue(out.contains("\"live\":2"), out);
|
||||
assertTrue(out.contains("\"free\":0"), "2 live + 1 lead seat fills all 3: " + out);
|
||||
assertTrue(out.contains("\"leadSeats\":1"), out);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #176: the OTHER measurement in the issue — a completely idle fleet still overstates
|
||||
* {@code free} by the lead's own seat. {@code maxLoad:3, live:0} must report {@code free:2}, the
|
||||
* real fan-out ceiling, not {@code 3}.
|
||||
*/
|
||||
@Test
|
||||
void leadSeatLowersFreeOnAnOtherwiseIdleSubscriptionProfile() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
|
||||
FleetMcp.LeadSeatSource leadSeats = new FleetMcp.LeadSeatSource(
|
||||
profile -> "sonnet".equals(profile) ? 1 : 0);
|
||||
|
||||
String out = textOf(FleetMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")),
|
||||
sessions, null, new FleetMcp.CapacitySource(profile -> 0, profile -> 3,
|
||||
() -> Set.of("sonnet"), () -> 0), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), leadSeats, Map.of(), "", null));
|
||||
|
||||
assertTrue(out.contains("\"live\":0"), out);
|
||||
assertTrue(out.contains("\"free\":2"), "an idle fleet's real ceiling is 3 minus the lead's own seat: " + out);
|
||||
assertTrue(out.contains("\"leadSeats\":1"), out);
|
||||
}
|
||||
|
||||
/** A profile with no lead seats reported must be byte-identical to before this ticket. */
|
||||
@Test
|
||||
void zeroLeadSeatsOmitsTheKeyAndLeavesFreeUnchanged() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
|
||||
|
||||
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"),
|
||||
FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(), FleetMcp.LeadSeatSource.none(),
|
||||
Map.of(), "", null));
|
||||
|
||||
assertTrue(out.contains("\"free\":2"), out);
|
||||
assertFalse(out.contains("leadSeats"), "no lead shares this profile's credential: " + out);
|
||||
}
|
||||
|
||||
@Test
|
||||
void listReportsLeadsAndFlagsTheCallersOwnRow() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
|
||||
Reference in New Issue
Block a user